fix(coding-agent): harden Agent Hub lifecycle and persistence

This commit is contained in:
Kyle McCleary
2026-08-04 16:29:15 -07:00
parent 91467c2f27
commit 8e5f619502
21 changed files with 1115 additions and 646 deletions
@@ -0,0 +1,248 @@
import type { AgentMetricsSummary, AgentRef, AgentStatus } from "../../registry/agent-registry";
import { MAIN_AGENT_ID } from "../../registry/agent-registry";
import type { ObservableSession } from "../session-observer-registry";
export type AgentMetrics = AgentMetricsSummary;
export interface AggregateMetrics extends AgentMetrics {
reportedAgents: number;
/** Rows whose duration is an observer-measured active runtime. */
activeDurationAgents: number;
}
export interface AgentTreeProjection {
rows: AgentRef[];
depthById: Map<string, number>;
parentById: Map<string, string>;
lastSiblingById: Map<string, boolean>;
}
export const STATUS_ORDER: Record<AgentStatus, number> = { running: 0, idle: 1, parked: 2, aborted: 3 };
function finiteMetric(value: number | undefined): number {
return typeof value === "number" && Number.isFinite(value) ? value : 0;
}
/** Exact observer usage for one roster entry. */
export function progressMetrics(observed: ObservableSession | undefined): AgentMetrics | undefined {
const progress = observed?.progress;
if (!progress) return undefined;
const { tokens, requests, toolCount: tools, cost, durationMs } = progress;
if (
typeof tokens !== "number" ||
!Number.isFinite(tokens) ||
typeof requests !== "number" ||
!Number.isFinite(requests) ||
typeof tools !== "number" ||
!Number.isFinite(tools) ||
typeof cost !== "number" ||
!Number.isFinite(cost) ||
typeof durationMs !== "number" ||
!Number.isFinite(durationMs)
) {
return undefined;
}
return {
tokens,
requests,
tools,
cost,
durationMs,
durationKind: "active",
contextTokens:
typeof progress.contextTokens === "number" && Number.isFinite(progress.contextTokens)
? progress.contextTokens
: undefined,
contextWindow:
typeof progress.contextWindow === "number" && Number.isFinite(progress.contextWindow)
? progress.contextWindow
: undefined,
};
}
/**
* Read direct assistant usage from a live session. SessionStats also includes
* usage embedded in completed `task` tool results, so using it for a parent
* row would double-count child rows in the aggregate.
*/
export function readSessionMetrics(session: NonNullable<AgentRef["session"]>): AgentMetrics | undefined {
try {
const stats = session.getSessionStats();
const messages = session.agent?.state?.messages;
if (!Array.isArray(messages)) {
return {
tokens: stats.tokens.input + stats.tokens.output + stats.tokens.cacheWrite,
requests: stats.assistantMessages,
tools: stats.toolCalls,
cost: stats.cost,
durationMs: 0,
durationKind: "unknown",
contextTokens: stats.contextUsage?.tokens,
contextWindow: stats.contextUsage?.contextWindow,
};
}
let tokens = 0;
let requests = 0;
let tools = 0;
let cost = 0;
for (const message of messages) {
if (message.role !== "assistant") continue;
requests++;
tokens += message.usage.input + message.usage.output + message.usage.cacheWrite;
tools += message.content.filter(content => content.type === "toolCall").length;
cost += message.usage.cost.total;
}
return {
tokens,
requests,
tools,
cost,
durationMs: 0,
durationKind: "unknown",
contextTokens: stats.contextUsage?.tokens,
contextWindow: stats.contextUsage?.contextWindow,
};
} catch {
// Render-only doubles and sessions being torn down may not expose a
// complete statistics host. Missing metrics are preferable to a broken hub.
return undefined;
}
}
export function aggregateMetrics(args: {
rows: readonly AgentRef[];
observedById: ReadonlyMap<string, ObservableSession>;
metricsFor: (ref: AgentRef, observed: ObservableSession | undefined) => AgentMetrics | undefined;
fallbackStatsSession: (
ref: AgentRef,
observed: ObservableSession | undefined,
) => NonNullable<AgentRef["session"]> | undefined;
sessionMetrics: WeakMap<object, { metrics: AgentMetrics | undefined }>;
refreshFallback: boolean;
}): { metrics: AggregateMetrics; hasFallbackLiveSessions: boolean } {
const total: AggregateMetrics = {
tokens: 0,
requests: 0,
tools: 0,
cost: 0,
durationMs: 0,
durationKind: "active",
reportedAgents: 0,
activeDurationAgents: 0,
};
let hasFallbackLiveSessions = false;
const countedFallbackSessions = new Set<NonNullable<AgentRef["session"]>>();
for (const ref of args.rows) {
const observed = args.observedById.get(ref.id);
const fallbackSession = args.fallbackStatsSession(ref, observed);
if (fallbackSession) {
hasFallbackLiveSessions = true;
if (args.refreshFallback || !args.sessionMetrics.has(fallbackSession)) {
args.sessionMetrics.set(fallbackSession, { metrics: readSessionMetrics(fallbackSession) });
}
}
const metrics = args.metricsFor(ref, observed);
if (!metrics || (fallbackSession && countedFallbackSessions.has(fallbackSession))) continue;
if (fallbackSession) countedFallbackSessions.add(fallbackSession);
total.reportedAgents++;
total.tokens += finiteMetric(metrics.tokens);
total.requests += finiteMetric(metrics.requests);
total.tools += finiteMetric(metrics.tools);
total.cost += finiteMetric(metrics.cost);
if (metrics.durationKind === "active") {
total.durationMs += finiteMetric(metrics.durationMs);
total.activeDurationAgents++;
}
}
return { metrics: total, hasFallbackLiveSessions };
}
/** Parent-before-child projection preserving the roster's stable sibling order. */
export function projectAgentTree(refs: readonly AgentRef[]): AgentTreeProjection {
const ids = new Set<string>();
const operationalIndex = new Map<string, number>();
for (let i = 0; i < refs.length; i++) {
ids.add(refs[i].id);
operationalIndex.set(refs[i].id, i);
}
const parentById = new Map<string, string>();
const children = new Map<string, AgentRef[]>();
for (const ref of refs) {
const parent =
ref.parentId && ref.parentId !== MAIN_AGENT_ID && ids.has(ref.parentId) ? ref.parentId : MAIN_AGENT_ID;
parentById.set(ref.id, parent);
const siblings = children.get(parent);
if (siblings) siblings.push(ref);
else children.set(parent, [ref]);
}
// A tree group occupies the position of its earliest operational row.
// Compute subtree minima iteratively so pathological lineage depth remains stack-safe.
const subtreeOrder = new Map<string, number>();
const visiting = new Set<string>();
const ranked = new Set<string>();
for (const start of refs) {
if (ranked.has(start.id)) continue;
const stack: Array<{ ref: AgentRef; expanded: boolean }> = [{ ref: start, expanded: false }];
while (stack.length > 0) {
const current = stack.pop();
if (!current) continue;
if (current.expanded) {
let order = operationalIndex.get(current.ref.id) ?? Number.MAX_SAFE_INTEGER;
for (const child of children.get(current.ref.id) ?? []) {
order = Math.min(order, subtreeOrder.get(child.id) ?? Number.MAX_SAFE_INTEGER);
}
subtreeOrder.set(current.ref.id, order);
visiting.delete(current.ref.id);
ranked.add(current.ref.id);
continue;
}
if (ranked.has(current.ref.id) || visiting.has(current.ref.id)) continue;
visiting.add(current.ref.id);
stack.push({ ref: current.ref, expanded: true });
const descendants = children.get(current.ref.id);
if (!descendants) continue;
for (let i = descendants.length - 1; i >= 0; i--) {
const child = descendants[i];
if (!ranked.has(child.id) && !visiting.has(child.id)) stack.push({ ref: child, expanded: false });
}
}
}
for (const siblings of children.values()) {
siblings.sort(
(a, b) =>
(subtreeOrder.get(a.id) ?? Number.MAX_SAFE_INTEGER) - (subtreeOrder.get(b.id) ?? Number.MAX_SAFE_INTEGER) ||
(operationalIndex.get(a.id) ?? Number.MAX_SAFE_INTEGER) -
(operationalIndex.get(b.id) ?? Number.MAX_SAFE_INTEGER),
);
}
const lastSiblingById = new Map<string, boolean>();
for (const siblings of children.values()) {
for (let i = 0; i < siblings.length; i++) lastSiblingById.set(siblings[i].id, i === siblings.length - 1);
}
const rows: AgentRef[] = [];
const visited = new Set<string>();
const depthById = new Map<string, number>();
const visit = (root: AgentRef, rootDepth: number): void => {
const stack: Array<{ ref: AgentRef; depth: number }> = [{ ref: root, depth: rootDepth }];
while (stack.length > 0) {
const current = stack.pop();
if (!current || visited.has(current.ref.id)) continue;
visited.add(current.ref.id);
depthById.set(current.ref.id, current.depth);
rows.push(current.ref);
const descendants = children.get(current.ref.id);
if (!descendants) continue;
for (let i = descendants.length - 1; i >= 0; i--)
stack.push({ ref: descendants[i], depth: current.depth + 1 });
}
};
for (const root of children.get(MAIN_AGENT_ID) ?? []) visit(root, 0);
// Corrupt persisted parent cycles remain visible as roots instead of disappearing.
for (const ref of refs) visit(ref, 0);
return { rows, depthById, parentById, lastSiblingById };
}
@@ -0,0 +1,194 @@
import { ThinkingLevel } from "@oh-my-pi/pi-agent-core";
import { Ellipsis, visibleWidth } from "@oh-my-pi/pi-tui";
import { formatDuration, formatNumber, sanitizeText } from "@oh-my-pi/pi-utils";
import { getRoleInfo } from "../../config/model-roles";
import type { Settings } from "../../config/settings";
import { type AgentRef, MAIN_AGENT_ID } from "../../registry/agent-registry";
import { parseThinkingLevel } from "../../thinking";
import { replaceTabs, TRUNCATE_LENGTHS, truncateToWidth } from "../../tools/render-utils";
import type { ObservableSession } from "../session-observer-registry";
import { theme } from "../theme/theme";
import type { AgentMetrics } from "./agent-hub-projection";
export interface RosterRender {
lines: string[];
hitRows: Array<number | undefined>;
}
/** Legacy progress snapshots may omit counters; snapshot absence remains distinct. */
export function metricNumber(value: number | undefined): number {
return typeof value === "number" && Number.isFinite(value) ? value : 0;
}
/** Compute the max content width for the current terminal, accounting for chrome. */
export function contentWidth(): number {
return Math.max(TRUNCATE_LENGTHS.SHORT, (process.stdout.columns || 80) - 6);
}
/** Remove terminal controls and normalize a value before it reaches the TUI. */
export function sanitizeDisplayText(text: string): string {
return replaceTabs(sanitizeText(text)).replace(/[\r\n]+/g, " ");
}
/** Sanitize a line for TUI display and truncate it to the viewport width. */
export function sanitizeLine(text: string, maxWidth?: number): string {
return truncateToWidth(sanitizeDisplayText(text), maxWidth ?? contentWidth());
}
export function clampHubLine(line: string, width: number): string {
return truncateToWidth(line.replace(/[\r\n]+/g, " "), Math.max(1, width), Ellipsis.Omit);
}
/** Status glyph, colored per theme status conventions. The title-line counts spell out the words. */
export function statusGlyph(status: AgentRef["status"]): string {
switch (status) {
case "running":
return theme.fg("accent", theme.status.running);
case "idle":
return theme.fg("success", theme.status.enabled);
case "parked":
return theme.fg("muted", theme.status.shadowed);
case "aborted":
return theme.fg("error", theme.status.aborted);
}
}
export function statusText(status: AgentRef["status"], text: string): string {
switch (status) {
case "running":
return theme.fg("accent", text);
case "idle":
return theme.fg("success", text);
case "parked":
return theme.fg("muted", text);
case "aborted":
return theme.fg("error", text);
}
}
/** Model id + thinking level (`sonnet-4-6 ◒ high`), level colored per theme. */
export function formatModelBadge(modelId: string, level: ThinkingLevel | undefined): string {
const model = theme.fg("muted", sanitizeDisplayText(modelId));
if (!level || level === ThinkingLevel.Off || level === ThinkingLevel.Inherit) return model;
const display = theme.thinking[level as keyof typeof theme.thinking] ?? level;
return `${model} ${theme.getThinkingBorderColor(level)(display)}`;
}
/** Textual model-role tag; color reinforces (but never replaces) the label. */
export function formatRoleBadge(role: string, settings: Settings): string {
const info = getRoleInfo(role, settings);
return theme.fg(info.color ?? "muted", sanitizeDisplayText(info.tag ?? info.name ?? role));
}
/** Format a resolved selector, preserving provider identity when requested. */
export function formatResolvedModelBadge(
resolved: string,
preserveProvider = false,
fallbackLevel?: ThinkingLevel,
): string {
const cleanResolved = sanitizeDisplayText(resolved);
// Model ids may themselves contain colons (`qwen3:14b`), so only treat the
// suffix as a thinking level when it parses as one.
const colon = cleanResolved.lastIndexOf(":");
const explicitLevel = colon >= 0 ? parseThinkingLevel(cleanResolved.slice(colon + 1)) : undefined;
const selector = explicitLevel !== undefined ? cleanResolved.slice(0, colon) : cleanResolved;
const label = preserveProvider ? selector : selector.slice(selector.indexOf("/") + 1);
return formatModelBadge(label, explicitLevel ?? fallbackLevel);
}
/**
* Resolved model + reasoning level for a hub row. Exact executor progress is
* authoritative (and survives completion); direct live sessions are the
* fallback for agents without an observer snapshot.
*/
export function modelBadge(ref: AgentRef, observed: ObservableSession | undefined): string | undefined {
const progress = observed?.progress;
const liveThinkingLevel = ref.session?.thinkingLevel;
const fallbackSelector =
ref.session?.retryFallbackModel ??
(progress?.resolvedModelIsFallback ? progress.resolvedModel : undefined) ??
(ref.history?.resolvedModelIsFallback ? ref.history.resolvedModel : undefined);
if (fallbackSelector) {
return `${theme.fg("warning", "fallback →")} ${formatResolvedModelBadge(fallbackSelector, true, liveThinkingLevel)}`;
}
const resolvedModel = progress?.resolvedModel ?? ref.history?.resolvedModel;
if (resolvedModel) return formatResolvedModelBadge(resolvedModel, false, liveThinkingLevel);
const model = ref.session?.model;
if (!model) return undefined;
const level = model.thinking ? liveThinkingLevel : undefined;
return formatModelBadge(model.id, level);
}
export function formatMetricDuration(metrics: AgentMetrics): string | undefined {
const durationMs = metricNumber(metrics.durationMs);
if (durationMs <= 0) return undefined;
const label = metrics.durationKind === "active" ? "active" : metrics.durationKind === "span" ? "span" : "duration";
return `${formatDuration(durationMs)} ${label}`;
}
export function formatCost(cost: number): string {
const amount = metricNumber(cost);
if (amount < 0.01) return `$${amount.toFixed(4)}`;
if (amount < 1) return `$${amount.toFixed(3)}`;
return `$${amount.toFixed(2)}`;
}
export function formatMetrics(metrics: AgentMetrics): string {
return [
formatCost(metrics.cost),
formatMetricDuration(metrics) ?? "time —",
`${formatNumber(metrics.requests)} req`,
`${formatNumber(metrics.tools)} tools`,
`${formatNumber(metrics.tokens)} tok`,
].join(theme.sep.dot);
}
export function contextGauge(tokens: number, window: number): string {
const ratio = Math.max(0, Math.min(1, tokens / window));
const filled = Math.round(ratio * 10);
return `${theme.fg("accent", "━".repeat(filled))}${theme.fg("dim", "─".repeat(10 - filled))} ${formatNumber(tokens)}/${formatNumber(window)} ${Math.round(ratio * 100)}%`;
}
/** Fit a child-id preview without joining an arbitrarily large child set. */
export function formatChildIds(children: readonly AgentRef[], width: number): string {
const max = Math.max(1, width);
let shown = 0;
let text = "";
while (shown < children.length) {
const id = sanitizeLine(children[shown].id, max);
const candidate = text ? `${text}, ${id}` : id;
const remaining = children.length - shown - 1;
const suffix = remaining > 0 ? `, … +${remaining}` : "";
if (visibleWidth(candidate + suffix) > max) {
const includesCurrent = text.length === 0;
const omitted = children.length - shown - Number(includesCurrent);
return truncateToWidth(`${includesCurrent ? id : text}${omitted > 0 ? `, … +${omitted}` : ""}`, max);
}
text = candidate;
shown++;
}
return text;
}
/** Bash `tree`-style ancestry prefix, clipped from the left on pathological depth. */
export function treeBranch(
ref: AgentRef,
maxWidth: number,
depthById: ReadonlyMap<string, number>,
parentById: ReadonlyMap<string, string>,
lastSiblingById: ReadonlyMap<string, boolean>,
): string {
if ((depthById.get(ref.id) ?? 0) === 0) return "";
const segments: string[] = [lastSiblingById.get(ref.id) ? "└── " : "├── "];
const ancestry = new Set<string>();
let parent = parentById.get(ref.id);
while (parent && parentById.get(parent) !== MAIN_AGENT_ID && !ancestry.has(parent)) {
ancestry.add(parent);
segments.push(lastSiblingById.get(parent) ? " " : "│ ");
parent = parentById.get(parent);
}
const maxSegments = Math.max(1, Math.floor(Math.max(4, maxWidth - 2) / 4));
const omitted = Math.max(0, segments.length - maxSegments);
const prefix = segments.slice(0, maxSegments).reverse().join("");
return theme.fg("dim", `${omitted > 0 ? "… " : ""}${prefix}`);
}
@@ -13,10 +13,9 @@
*
* Replaces the old SessionObserverOverlayComponent (ctrl+s observer).
*/
import { type AgentTool, ThinkingLevel } from "@oh-my-pi/pi-agent-core";
import type { AgentTool } from "@oh-my-pi/pi-agent-core";
import {
Container,
Ellipsis,
matchesKey,
type OverlayHandle,
padding,
@@ -27,27 +26,43 @@ import {
visibleWidth,
wrapTextWithAnsi,
} from "@oh-my-pi/pi-tui";
import { formatAge, formatDuration, formatNumber, getProjectDir, logger } from "@oh-my-pi/pi-utils";
import { formatAge, formatNumber, getProjectDir, logger } from "@oh-my-pi/pi-utils";
import type { KeyId } from "../../config/keybindings";
import { getRoleInfo } from "../../config/model-roles";
import type { Settings } from "../../config/settings";
import type { MessageRenderer } from "../../extensibility/extensions/types";
import { IrcBus } from "../../irc/bus";
import { AgentLifecycleManager } from "../../registry/agent-lifecycle";
import {
type AgentMetricsSummary,
type AgentRef,
AgentRegistry,
type AgentStatus,
MAIN_AGENT_ID,
} from "../../registry/agent-registry";
import { type AgentRef, AgentRegistry, type AgentStatus, MAIN_AGENT_ID } from "../../registry/agent-registry";
import { registerPersistedSubagents } from "../../registry/persisted-agents";
import { USER_INTERRUPT_LABEL } from "../../session/messages";
import { parseThinkingLevel } from "../../thinking";
import { replaceTabs, TRUNCATE_LENGTHS, truncateToWidth } from "../../tools/render-utils";
import { truncateToWidth } from "../../tools/render-utils";
import type { ObservableSession, SessionObserverRegistry } from "../session-observer-registry";
import { theme } from "../theme/theme";
import { matchesSelectDown, matchesSelectUp } from "../utils/keybinding-matchers";
import {
type AgentMetrics,
type AggregateMetrics,
aggregateMetrics,
progressMetrics,
projectAgentTree,
STATUS_ORDER,
} from "./agent-hub-projection";
import {
clampHubLine,
contextGauge,
formatChildIds,
formatCost,
formatMetricDuration,
formatMetrics,
formatRoleBadge,
modelBadge,
type RosterRender,
sanitizeDisplayText,
sanitizeLine,
statusGlyph,
statusText,
treeBranch,
} from "./agent-hub-renderer";
import { AgentTranscriptViewer } from "./agent-transcript-viewer";
import {
bottomBorder,
@@ -60,207 +75,18 @@ import {
topBorderSplit,
} from "./overlay-box";
/** Two-pane mode needs a useful roster and a readable inspector. */
const SPLIT_MIN_WIDTH = 96;
const DETAIL_MIN_WIDTH = 34;
const ROSTER_MIN_WIDTH = 48;
type HubViewMode = "roster" | "tree";
type AgentMetrics = AgentMetricsSummary;
interface AggregateMetrics extends AgentMetrics {
reportedAgents: number;
}
interface RosterRender {
lines: string[];
hitRows: Array<number | undefined>;
}
/** Legacy progress snapshots may omit counters; snapshot absence remains distinct. */
function metricNumber(value: number | undefined): number {
return typeof value === "number" && Number.isFinite(value) ? value : 0;
}
/** Refresh cadence for the relative-time column */
/** Refresh cadence for the relative-time column. */
const AGE_TICK_MS = 5_000;
const DATA_CHANGE_RENDER_COALESCE_MS = 100;
/** Double-tap window for the table's left-left "close hub" gesture. */
const LEFT_TAP_WINDOW_MS = 500;
/** Compute the max content width for the current terminal, accounting for chrome. */
function contentWidth(): number {
return Math.max(TRUNCATE_LENGTHS.SHORT, (process.stdout.columns || 80) - 6);
}
/** Sanitize a line for TUI display: replace tabs, then truncate to viewport width. */
function sanitizeLine(text: string, maxWidth?: number): string {
const singleLine = replaceTabs(text).replace(/[\r\n]+/g, " ");
return truncateToWidth(singleLine, maxWidth ?? contentWidth());
}
function clampHubLine(line: string, width: number): string {
return truncateToWidth(line.replace(/[\r\n]+/g, " "), Math.max(1, width), Ellipsis.Omit);
}
const STATUS_ORDER: Record<AgentStatus, number> = { running: 0, idle: 1, parked: 2, aborted: 3 };
/** Status glyph, colored per theme status conventions. The title-line counts spell out the words. */
function statusGlyph(status: AgentStatus): string {
switch (status) {
case "running":
return theme.fg("accent", theme.status.running);
case "idle":
return theme.fg("success", theme.status.enabled);
case "parked":
return theme.fg("muted", theme.status.shadowed);
case "aborted":
return theme.fg("error", theme.status.aborted);
}
}
function statusText(status: AgentStatus, text: string): string {
switch (status) {
case "running":
return theme.fg("accent", text);
case "idle":
return theme.fg("success", text);
case "parked":
return theme.fg("muted", text);
case "aborted":
return theme.fg("error", text);
}
}
/** Model id + thinking level (`sonnet-4-6 ◒ high`), level colored per theme. */
function formatModelBadge(modelId: string, level: ThinkingLevel | undefined): string {
const model = theme.fg("muted", replaceTabs(modelId));
if (!level || level === ThinkingLevel.Off || level === ThinkingLevel.Inherit) return model;
const display = theme.thinking[level as keyof typeof theme.thinking] ?? level;
return `${model} ${theme.getThinkingBorderColor(level)(display)}`;
}
/** Textual model-role tag; color reinforces (but never replaces) the label. */
function formatRoleBadge(role: string, settings: Settings): string {
const info = getRoleInfo(role, settings);
return theme.fg(info.color ?? "muted", replaceTabs(info.tag ?? info.name ?? role));
}
/** Format a resolved selector, preserving provider identity when requested. */
function formatResolvedModelBadge(resolved: string, preserveProvider = false, fallbackLevel?: ThinkingLevel): string {
// Model ids may themselves contain colons (`qwen3:14b`), so only treat the
// suffix as a thinking level when it parses as one.
const colon = resolved.lastIndexOf(":");
const explicitLevel = colon >= 0 ? parseThinkingLevel(resolved.slice(colon + 1)) : undefined;
const selector = explicitLevel !== undefined ? resolved.slice(0, colon) : resolved;
const label = preserveProvider ? selector : selector.slice(selector.indexOf("/") + 1);
return formatModelBadge(label, explicitLevel ?? fallbackLevel);
}
/**
* Resolved model + reasoning level for a hub row. Exact executor progress is
* authoritative (and survives completion); direct live sessions are the
* fallback for agents without an observer snapshot. Active retry fallbacks
* retain provider identity and carry an explicit marker.
*/
function modelBadge(ref: AgentRef, observed: ObservableSession | undefined): string | undefined {
const progress = observed?.progress;
const liveThinkingLevel = ref.session?.thinkingLevel;
// The executor fallback flag also covers retries that do not populate the
// live session's retryFallbackModel (for example Fireworks Fast → base).
const fallbackSelector =
ref.session?.retryFallbackModel ?? (progress?.resolvedModelIsFallback ? progress.resolvedModel : undefined);
if (fallbackSelector) {
return `${theme.fg("warning", "fallback →")} ${formatResolvedModelBadge(fallbackSelector, true, liveThinkingLevel)}`;
}
const resolvedModel = progress?.resolvedModel ?? ref.history?.resolvedModel;
if (resolvedModel) {
return formatResolvedModelBadge(resolvedModel, false, liveThinkingLevel);
}
const model = ref.session?.model;
if (!model) return undefined;
const level = model.thinking ? liveThinkingLevel : undefined;
return formatModelBadge(model.id, level);
}
/** Exact observer usage for one roster entry. */
function progressMetrics(observed: ObservableSession | undefined): AgentMetrics | undefined {
const progress = observed?.progress;
if (!progress) return undefined;
const { tokens, requests, toolCount: tools, cost, durationMs } = progress;
if (
typeof tokens !== "number" ||
!Number.isFinite(tokens) ||
typeof requests !== "number" ||
!Number.isFinite(requests) ||
typeof tools !== "number" ||
!Number.isFinite(tools) ||
typeof cost !== "number" ||
!Number.isFinite(cost) ||
typeof durationMs !== "number" ||
!Number.isFinite(durationMs)
) {
return undefined;
}
return {
tokens,
requests,
tools,
cost,
durationMs,
contextTokens:
typeof progress.contextTokens === "number" && Number.isFinite(progress.contextTokens)
? progress.contextTokens
: undefined,
contextWindow:
typeof progress.contextWindow === "number" && Number.isFinite(progress.contextWindow)
? progress.contextWindow
: undefined,
};
}
function formatCost(cost: number): string {
const amount = metricNumber(cost);
if (amount < 0.01) return `$${amount.toFixed(4)}`;
if (amount < 1) return `$${amount.toFixed(3)}`;
return `$${amount.toFixed(2)}`;
}
function formatMetrics(metrics: AgentMetrics): string {
return [
formatCost(metrics.cost),
formatDuration(metrics.durationMs),
`${formatNumber(metrics.requests)} req`,
`${formatNumber(metrics.tools)} tools`,
`${formatNumber(metrics.tokens)} tok`,
].join(theme.sep.dot);
}
function contextGauge(tokens: number, window: number): string {
const ratio = Math.max(0, Math.min(1, tokens / window));
const filled = Math.round(ratio * 10);
return `${theme.fg("accent", "━".repeat(filled))}${theme.fg("dim", "─".repeat(10 - filled))} ${formatNumber(tokens)}/${formatNumber(window)} ${Math.round(ratio * 100)}%`;
}
/** Fit a child-id preview without joining an arbitrarily large child set. */
function formatChildIds(children: readonly AgentRef[], width: number): string {
const max = Math.max(1, width);
let shown = 0;
let text = "";
while (shown < children.length) {
const id = sanitizeLine(children[shown].id, max);
const candidate = text ? `${text}, ${id}` : id;
const remaining = children.length - shown - 1;
const suffix = remaining > 0 ? `, … +${remaining}` : "";
if (visibleWidth(candidate + suffix) > max) {
const includesCurrent = text.length === 0;
const omitted = children.length - shown - Number(includesCurrent);
return truncateToWidth(`${includesCurrent ? id : text}${omitted > 0 ? `, … +${omitted}` : ""}`, max);
}
text = candidate;
shown++;
}
return text;
}
/** Two-pane mode needs a useful roster and a readable inspector. */
const SPLIT_MIN_WIDTH = 96;
const DETAIL_MIN_WIDTH = 34;
const ROSTER_MIN_WIDTH = 48;
/** Result of one host-backed transcript read for the Agent Hub viewer. */
export interface AgentHubRemoteTranscript {
text: string;
@@ -343,6 +169,7 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
#notice: string | undefined;
/** Captured row order from the first refresh; keeps the hub stable while open. */
#rowOrder: Map<string, number> | undefined;
#nextRowOrder = 0;
/** Double-tap window state for the table's left-left "close hub" gesture. */
#lastLeftTap = 0;
/** Operational ordering by default; tree mode groups descendants under their spawner. */
@@ -358,7 +185,9 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
tools: 0,
cost: 0,
durationMs: 0,
durationKind: "active",
reportedAgents: 0,
activeDurationAgents: 0,
};
#statusCounts: Record<AgentStatus, number> = { running: 0, idle: 0, parked: 0, aborted: 0 };
#childrenByParent = new Map<string, AgentRef[]>();
@@ -369,6 +198,10 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
/** On narrow terminals Tab replaces the roster with the selected-agent inspector. */
#narrowDetailsOpen = false;
#lastRenderWasSplit = false;
#lastSplitRosterWidth: number | undefined;
/** Scroll offset for the selected-agent inspector when its content overflows. */
#detailScrollOffset = 0;
#detailAgentId: string | undefined;
// Transcript-viewer launch deps (passed through to AgentTranscriptViewer).
#ui: TUI;
@@ -416,7 +249,7 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
this.#unsubscribers.push(this.#observers.onChange(() => this.#scheduleDataChange()));
this.#ageTimer = setInterval(() => {
if (this.#hasFallbackLiveSessions) {
this.#aggregate = this.#aggregateMetrics(this.#observedById, true);
this.#refreshAggregate(true);
}
this.#requestRender();
}, AGE_TICK_MS);
@@ -476,7 +309,15 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
}
handleInput(keyData: string): void {
if (routeSgrMouseInput(keyData, event => routeSelectListMouse(this, event, event.row))) return;
if (
routeSgrMouseInput(keyData, event => {
const split = this.#lastSplitRosterWidth;
if (split !== undefined && event.wheel === null && event.col > split + 2) return false;
return routeSelectListMouse(this, event, event.row);
})
) {
return;
}
// The hub/observe keys always close the overlay (toggle semantics)
for (const key of this.#hubKeys) {
@@ -508,11 +349,12 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
* restored when the viewer closes. No-op without a real TUI (render-only test stub).
*/
openChat(id: string): void {
if (!this.#registry.get(id)) return;
if (this.#disposed || !this.#registry.get(id)) return;
if (typeof this.#ui.showOverlay !== "function") return;
this.#closeTranscriptOverlay();
this.#notice = undefined;
const viewer = new AgentTranscriptViewer({
let viewer: AgentTranscriptViewer;
viewer = new AgentTranscriptViewer({
agentId: id,
registry: this.#registry,
remote: this.#remote,
@@ -527,10 +369,11 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
expandKeys: this.#expandKeys,
hubKeys: this.#hubKeys,
requestRender: this.#requestRender,
onClose: () => this.#closeTranscriptOverlay(),
onClose: () => this.#closeTranscriptOverlay(viewer),
onHubClose: () => {
this.#closeTranscriptOverlay();
this.#onDone();
if (this.#disposed) return;
this.#closeTranscriptOverlay(viewer);
if (!this.#disposed) this.#onDone();
},
});
this.#transcriptViewer = viewer;
@@ -540,13 +383,19 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
}
/** Close and dispose the transcript overlay, restoring focus to the hub table. */
#closeTranscriptOverlay(): void {
this.#transcriptOverlay?.hide();
#closeTranscriptOverlay(expectedViewer?: AgentTranscriptViewer): void {
if (expectedViewer && this.#transcriptViewer !== expectedViewer) return;
const overlay = this.#transcriptOverlay;
const viewer = this.#transcriptViewer;
if (!overlay && !viewer) return;
overlay?.hide();
this.#transcriptOverlay = undefined;
this.#transcriptViewer?.dispose();
viewer?.dispose();
this.#transcriptViewer = undefined;
if (typeof this.#ui.setFocus === "function") this.#ui.setFocus(this);
this.#requestRender();
if (!this.#disposed) {
if (typeof this.#ui.setFocus === "function") this.#ui.setFocus(this);
this.#requestRender();
}
}
// ========================================================================
@@ -572,37 +421,45 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
const refs = this.#registry.list().filter(ref => ref.id !== MAIN_AGENT_ID);
this.#observedById = new Map();
for (const session of this.#observers.getSessions()) this.#observedById.set(session.id, session);
const rowOrder = this.#rowOrder;
let rosterRows: AgentRef[];
if (!rowOrder) {
// First refresh (usually the constructor): order by status, then recency.
rosterRows = refs.sort(
(a, b) => STATUS_ORDER[a.status] - STATUS_ORDER[b.status] || b.lastActivity - a.lastActivity,
(a, b) =>
STATUS_ORDER[a.status] - STATUS_ORDER[b.status] ||
b.lastActivity - a.lastActivity ||
a.id.localeCompare(b.id),
);
this.#rowOrder = new Map();
for (let i = 0; i < rosterRows.length; i++) this.#rowOrder.set(rosterRows[i].id, i);
for (const ref of rosterRows) this.#rowOrder.set(ref.id, this.#nextRowOrder++);
} else {
// After the hub is open, freeze relative order within each lifecycle
// group. New agents append instead of moving the user's selection.
rosterRows = refs.sort((a, b) => {
const statusDiff = STATUS_ORDER[a.status] - STATUS_ORDER[b.status];
if (statusDiff !== 0) return statusDiff;
return (rowOrder.get(a.id) ?? Number.MAX_SAFE_INTEGER) - (rowOrder.get(b.id) ?? Number.MAX_SAFE_INTEGER);
});
rosterRows = refs.sort(
(a, b) => (rowOrder.get(a.id) ?? Number.MAX_SAFE_INTEGER) - (rowOrder.get(b.id) ?? Number.MAX_SAFE_INTEGER),
);
for (const ref of rosterRows) {
if (!rowOrder.has(ref.id)) rowOrder.set(ref.id, rowOrder.size);
if (!rowOrder.has(ref.id)) rowOrder.set(ref.id, this.#nextRowOrder++);
}
}
this.#rows = this.#viewMode === "tree" ? this.#orderAsTree(rosterRows) : rosterRows;
if (this.#viewMode === "roster") {
if (this.#viewMode === "tree") {
const tree = projectAgentTree(rosterRows);
this.#rows = tree.rows;
this.#treeDepthById = tree.depthById;
this.#treeParentById = tree.parentById;
this.#treeLastSiblingById = tree.lastSiblingById;
} else {
this.#rows = rosterRows;
this.#treeDepthById.clear();
this.#treeParentById.clear();
this.#treeLastSiblingById.clear();
}
const keptIndex = selectedId ? this.#rows.findIndex(ref => ref.id === selectedId) : -1;
this.#selectedRow = keptIndex >= 0 ? keptIndex : Math.min(this.#selectedRow, Math.max(0, this.#rows.length - 1));
const detailAgentId = this.#rows[this.#selectedRow]?.id;
if (detailAgentId !== this.#detailAgentId) {
this.#detailAgentId = detailAgentId;
this.#detailScrollOffset = 0;
}
this.#childrenByParent.clear();
for (const ref of rosterRows) {
@@ -613,12 +470,10 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
}
this.#statusCounts = { running: 0, idle: 0, parked: 0, aborted: 0 };
for (const ref of rosterRows) this.#statusCounts[ref.status]++;
this.#aggregate = this.#aggregateMetrics(this.#observedById);
this.#refreshAggregate();
}
#metricsFor(ref: AgentRef, observed: ObservableSession | undefined): AgentMetrics | undefined {
// An observer snapshot is authoritative. Legacy/incomplete snapshots remain
// unknown rather than being silently replaced by unrelated transcript stats.
if (observed?.progress) return progressMetrics(observed);
if (ref.history?.metrics) return ref.history.metrics;
const session = this.#fallbackStatsSession(ref, observed);
@@ -634,140 +489,6 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
return session && typeof session.getSessionStats === "function" ? session : undefined;
}
#readSessionMetrics(session: NonNullable<AgentRef["session"]>): AgentMetrics | undefined {
try {
const stats = session.getSessionStats();
return {
// Match AgentProgress lifetime billing volume: cache reads are
// deliberately excluded because every turn rereads cached context.
tokens: stats.tokens.input + stats.tokens.output + stats.tokens.cacheWrite,
requests: stats.assistantMessages,
tools: stats.toolCalls,
cost: stats.cost,
durationMs: 0,
contextTokens: stats.contextUsage?.tokens,
contextWindow: stats.contextUsage?.contextWindow,
};
} catch {
// Render-only doubles and sessions being torn down may not expose a
// complete statistics host. Missing metrics are preferable to a broken hub.
return undefined;
}
}
/** Parent-before-child projection preserving the roster's stable sibling order. */
#orderAsTree(refs: AgentRef[]): AgentRef[] {
this.#treeParentById.clear();
this.#treeLastSiblingById.clear();
const ids = new Set<string>();
const operationalIndex = new Map<string, number>();
for (let i = 0; i < refs.length; i++) {
ids.add(refs[i].id);
operationalIndex.set(refs[i].id, i);
}
const children = new Map<string, AgentRef[]>();
for (const ref of refs) {
const parent =
ref.parentId && ref.parentId !== MAIN_AGENT_ID && ids.has(ref.parentId) ? ref.parentId : MAIN_AGENT_ID;
this.#treeParentById.set(ref.id, parent);
const siblings = children.get(parent);
if (siblings) siblings.push(ref);
else children.set(parent, [ref]);
}
// A tree group occupies the position of its earliest operational row.
// Otherwise a recent child could leave its older parent group below an
// unrelated peer (`child, peer, parent` → `peer, parent, child`).
// Compute subtree minima iteratively so even pathological lineage depth
// remains stack-safe.
const subtreeOrder = new Map<string, number>();
const visiting = new Set<string>();
const ranked = new Set<string>();
for (const start of refs) {
if (ranked.has(start.id)) continue;
const stack: Array<{ ref: AgentRef; expanded: boolean }> = [{ ref: start, expanded: false }];
while (stack.length > 0) {
const current = stack.pop();
if (!current) continue;
if (current.expanded) {
let order = operationalIndex.get(current.ref.id) ?? Number.MAX_SAFE_INTEGER;
for (const child of children.get(current.ref.id) ?? []) {
order = Math.min(order, subtreeOrder.get(child.id) ?? Number.MAX_SAFE_INTEGER);
}
subtreeOrder.set(current.ref.id, order);
visiting.delete(current.ref.id);
ranked.add(current.ref.id);
continue;
}
if (ranked.has(current.ref.id) || visiting.has(current.ref.id)) continue;
visiting.add(current.ref.id);
stack.push({ ref: current.ref, expanded: true });
const descendants = children.get(current.ref.id);
if (!descendants) continue;
for (let i = descendants.length - 1; i >= 0; i--) {
const child = descendants[i];
if (!ranked.has(child.id) && !visiting.has(child.id)) {
stack.push({ ref: child, expanded: false });
}
}
}
}
for (const siblings of children.values()) {
siblings.sort(
(a, b) =>
(subtreeOrder.get(a.id) ?? Number.MAX_SAFE_INTEGER) -
(subtreeOrder.get(b.id) ?? Number.MAX_SAFE_INTEGER) ||
(operationalIndex.get(a.id) ?? Number.MAX_SAFE_INTEGER) -
(operationalIndex.get(b.id) ?? Number.MAX_SAFE_INTEGER),
);
}
for (const siblings of children.values()) {
for (let i = 0; i < siblings.length; i++) {
this.#treeLastSiblingById.set(siblings[i].id, i === siblings.length - 1);
}
}
const ordered: AgentRef[] = [];
const visited = new Set<string>();
this.#treeDepthById.clear();
const visit = (root: AgentRef, rootDepth: number): void => {
const stack: Array<{ ref: AgentRef; depth: number }> = [{ ref: root, depth: rootDepth }];
while (stack.length > 0) {
const current = stack.pop();
if (!current || visited.has(current.ref.id)) continue;
visited.add(current.ref.id);
this.#treeDepthById.set(current.ref.id, current.depth);
ordered.push(current.ref);
const descendants = children.get(current.ref.id);
if (!descendants) continue;
for (let i = descendants.length - 1; i >= 0; i--) {
stack.push({ ref: descendants[i], depth: current.depth + 1 });
}
}
};
for (const root of children.get(MAIN_AGENT_ID) ?? []) visit(root, 0);
// Corrupt persisted parent cycles remain visible as roots instead of
// disappearing from the operational surface.
for (const ref of refs) visit(ref, 0);
return ordered;
}
/** Bash `tree`-style ancestry prefix, clipped from the left on pathological depth. */
#treeBranch(ref: AgentRef, maxWidth: number): string {
if ((this.#treeDepthById.get(ref.id) ?? 0) === 0) return "";
const segments: string[] = [this.#treeLastSiblingById.get(ref.id) ? "└── " : "├── "];
let parent = this.#treeParentById.get(ref.id);
while (parent && this.#treeParentById.get(parent) !== MAIN_AGENT_ID) {
segments.push(this.#treeLastSiblingById.get(parent) ? " " : "│ ");
parent = this.#treeParentById.get(parent);
}
const maxSegments = Math.max(1, Math.floor(Math.max(4, maxWidth - 2) / 4));
const omitted = Math.max(0, segments.length - maxSegments);
const prefix = segments.slice(0, maxSegments).reverse().join("");
return theme.fg("dim", `${omitted > 0 ? "… " : ""}${prefix}`);
}
// ========================================================================
// Table view
// ========================================================================
@@ -778,6 +499,7 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
const observedById = this.#observedById;
const split = this.#splitRosterWidth(width);
this.#lastRenderWasSplit = split !== undefined;
this.#lastSplitRosterWidth = split;
const selected = this.#rows[this.#selectedRow];
const lines: string[] = [];
@@ -826,12 +548,15 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
#footer(showingNarrowDetails: boolean, availableWidth: number): string {
const nextView = this.#viewMode === "roster" ? "by parent" : "flat";
if (showingNarrowDetails) {
return theme.fg("dim", `Tab:roster Enter:open t:${nextView} Esc:roster`);
return theme.fg("dim", `Tab:roster PgUp/PgDn:scroll Enter:open t:${nextView} Esc:roster`);
}
if (availableWidth < 96) {
return theme.fg("dim", `j/k:select Enter:open t:${nextView} Tab:details r/x:manage Esc:close`);
}
return theme.fg("dim", `j/k/wheel:select Enter/click:open t:${nextView} r:revive x:kill Esc:close`);
return theme.fg(
"dim",
`j/k/wheel:select PgUp/PgDn:details Enter/click:open t:${nextView} r:revive x:kill Esc:close`,
);
}
#renderRosterPanel(width: number, rows: number, observedById: ReadonlyMap<string, ObservableSession>): RosterRender {
@@ -986,12 +711,14 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
);
return lines;
}
const activeTime = formatMetricDuration(metrics);
const usage = [
theme.fg("statusLineCost", formatCost(metrics.cost)),
theme.fg("dim", `${formatDuration(metrics.durationMs)} agent time`),
theme.fg("dim", activeTime ? `${activeTime} agent time` : "agent time —"),
theme.fg("dim", `${formatNumber(metrics.requests)} req`),
theme.fg("dim", `${formatNumber(metrics.tools)} tools`),
theme.fg("dim", `${formatNumber(metrics.tokens)} tok`),
theme.fg("dim", `${metrics.activeDurationAgents}/${metrics.reportedAgents} timed`),
theme.fg("dim", `${metrics.reportedAgents}/${this.#rows.length} measured`),
].join(theme.fg("dim", theme.sep.dot));
lines.push(...wrapTextWithAnsi(usage, Math.max(1, width)));
@@ -1007,36 +734,17 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
return parts.join(theme.sep.dot);
}
#aggregateMetrics(observedById: ReadonlyMap<string, ObservableSession>, refreshFallback = false): AggregateMetrics {
const total: AggregateMetrics = {
tokens: 0,
requests: 0,
tools: 0,
cost: 0,
durationMs: 0,
reportedAgents: 0,
};
let hasFallbackLiveSessions = false;
for (const ref of this.#rows) {
const observed = observedById.get(ref.id);
const fallbackSession = this.#fallbackStatsSession(ref, observed);
if (fallbackSession) {
hasFallbackLiveSessions = true;
if (refreshFallback || !this.#sessionMetrics.has(fallbackSession)) {
this.#sessionMetrics.set(fallbackSession, { metrics: this.#readSessionMetrics(fallbackSession) });
}
}
const metrics = this.#metricsFor(ref, observed);
if (!metrics) continue;
total.reportedAgents++;
total.tokens += metrics.tokens;
total.requests += metrics.requests;
total.tools += metrics.tools;
total.cost += metrics.cost;
total.durationMs += metrics.durationMs;
}
this.#hasFallbackLiveSessions = hasFallbackLiveSessions;
return total;
#refreshAggregate(refreshFallback = false): void {
const result = aggregateMetrics({
rows: this.#rows,
observedById: this.#observedById,
metricsFor: (ref, observed) => this.#metricsFor(ref, observed),
fallbackStatsSession: (ref, observed) => this.#fallbackStatsSession(ref, observed),
sessionMetrics: this.#sessionMetrics,
refreshFallback,
});
this.#aggregate = result.metrics;
this.#hasFallbackLiveSessions = result.hasFallbackLiveSessions;
}
#renderDetailPanel(
@@ -1052,20 +760,20 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
const children = this.#childrenByParent.get(ref.id) ?? [];
const lines: string[] = [];
const add = (line = ""): void => {
if (lines.length < rows) lines.push(truncateToWidth(line, width));
lines.push(truncateToWidth(line, width));
};
const addWrapped = (text: string, maxRows = 2): void => {
for (const wrapped of wrapTextWithAnsi(sanitizeLine(text), Math.max(1, width)).slice(0, maxRows)) add(wrapped);
};
const section = (label: string): void => {
if (lines.length > 0) add();
const section = (label: string, contentRows = 0): void => {
if (lines.length > 0 && lines.length + 1 + contentRows < rows) add();
add(theme.bold(theme.fg("accent", label)));
};
add(`${statusGlyph(ref.status)} ${theme.bold(replaceTabs(ref.displayName || ref.id))}`);
if (ref.displayName && ref.displayName !== ref.id) add(theme.fg("dim", ref.id));
add(`${statusGlyph(ref.status)} ${theme.bold(sanitizeDisplayText(ref.displayName || ref.id))}`);
if (ref.displayName && ref.displayName !== ref.id) add(theme.fg("dim", sanitizeDisplayText(ref.id)));
const lifecycleDetails = [
metrics?.durationMs ? formatDuration(metrics.durationMs) : undefined,
metrics ? formatMetricDuration(metrics) : undefined,
`active ${formatAge(Math.max(1, Math.round((Date.now() - ref.lastActivity) / 1000)))}`,
].filter(Boolean);
add(
@@ -1095,7 +803,7 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
}
}
section("Usage");
section("Usage", 1);
if (metrics) {
addWrapped(formatMetrics(metrics), 3);
if (metrics.contextTokens !== undefined && metrics.contextWindow) {
@@ -1107,7 +815,7 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
section("Lineage");
add(
`Spawned by ${replaceTabs(ref.parentId ?? MAIN_AGENT_ID)}${children.length > 0 ? ` · ${children.length} children` : ""}`,
`Spawned by ${sanitizeDisplayText(ref.parentId ?? MAIN_AGENT_ID)}${children.length > 0 ? ` · ${children.length} children` : ""}`,
);
if (children.length > 0) add(theme.fg("dim", formatChildIds(children, width)));
add(theme.fg("dim", `Registered ${new Date(ref.createdAt).toISOString().slice(0, 16).replace("T", " ")}Z`));
@@ -1122,8 +830,11 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
),
);
while (lines.length < rows) lines.push("");
return lines.slice(0, rows);
const maxScroll = Math.max(0, lines.length - rows);
this.#detailScrollOffset = Math.min(this.#detailScrollOffset, maxScroll);
const visible = lines.slice(this.#detailScrollOffset, this.#detailScrollOffset + rows);
while (visible.length < rows) visible.push("");
return visible;
}
/**
@@ -1141,15 +852,18 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
const max = Math.max(1, width);
const cursor = selected ? theme.fg("accent", theme.nav.cursor) : " ";
const depth = this.#viewMode === "tree" ? (this.#treeDepthById.get(ref.id) ?? 0) : 0;
const branch = this.#viewMode === "tree" ? this.#treeBranch(ref, max) : "";
const id = replaceTabs(ref.id);
const branch =
this.#viewMode === "tree"
? treeBranch(ref, max, this.#treeDepthById, this.#treeParentById, this.#treeLastSiblingById)
: "";
const id = sanitizeDisplayText(ref.id);
const styledId = selected ? theme.bold(theme.fg("accent", id)) : theme.bold(id);
const fields: string[] = [`${cursor} ${statusGlyph(ref.status)} ${branch}${styledId}`];
if (ref.displayName && ref.displayName !== ref.id) {
fields.push(theme.fg("dim", replaceTabs(ref.displayName)));
fields.push(theme.fg("dim", sanitizeDisplayText(ref.displayName)));
}
if (this.#viewMode === "roster" && ref.parentId && ref.parentId !== MAIN_AGENT_ID) {
fields.push(theme.fg("dim", `↳ ${replaceTabs(ref.parentId)}`));
fields.push(theme.fg("dim", `↳ ${sanitizeDisplayText(ref.parentId)}`));
}
if (ref.kind === "advisor") {
fields.push(theme.fg("warning", "read-only"));
@@ -1198,10 +912,23 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
});
}
#scrollDetails(direction: -1 | 1): void {
this.#detailScrollOffset = Math.max(0, this.#detailScrollOffset + direction * 5);
this.#requestRender();
}
#selectRow(index: number): void {
if (index !== this.#selectedRow) {
this.#detailScrollOffset = 0;
this.#detailAgentId = this.#rows[index]?.id;
}
this.#selectedRow = index;
}
handleWheel(delta: -1 | 1): void {
this.#hoveredRow = null;
if (this.#rows.length > 0) {
this.#selectedRow = Math.max(0, Math.min(this.#selectedRow + delta, this.#rows.length - 1));
this.#selectRow(Math.max(0, Math.min(this.#selectedRow + delta, this.#rows.length - 1)));
}
this.#requestRender();
}
@@ -1217,14 +944,12 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
}
clickItem(index: number): void {
const selected = this.#rows[index];
if (!selected) return;
this.#hoveredRow = index;
if (index === this.#selectedRow) {
const selected = this.#rows[index];
if (selected) this.#activateAgent(selected);
return;
}
this.#selectedRow = index;
this.#selectRow(index);
this.#requestRender();
this.#activateAgent(selected);
}
#handleTableInput(keyData: string): void {
@@ -1242,6 +967,16 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
this.#requestRender();
return;
}
if (this.#lastRenderWasSplit || this.#narrowDetailsOpen) {
if (matchesKey(keyData, "pageUp")) {
this.#scrollDetails(-1);
return;
}
if (matchesKey(keyData, "pageDown")) {
this.#scrollDetails(1);
return;
}
}
if (keyData === "t") {
this.#hoveredRow = null;
this.#viewMode = this.#viewMode === "roster" ? "tree" : "roster";
@@ -1267,14 +1002,14 @@ export class AgentHubOverlayComponent extends Container implements SelectListMou
this.#hoveredRow = null;
if (matchesKey(keyData, "j") || matchesSelectDown(keyData)) {
if (this.#rows.length > 0) {
this.#selectedRow = Math.min(this.#selectedRow + 1, this.#rows.length - 1);
this.#selectRow(Math.min(this.#selectedRow + 1, this.#rows.length - 1));
}
this.#requestRender();
return;
}
if (matchesKey(keyData, "k") || matchesSelectUp(keyData)) {
if (this.#rows.length > 0) {
this.#selectedRow = Math.max(this.#selectedRow - 1, 0);
this.#selectRow(Math.max(this.#selectedRow - 1, 0));
}
this.#requestRender();
return;
@@ -2009,7 +2009,9 @@ export class SelectorController {
closed = true;
hub.dispose();
overlayHandle?.hide();
this.focusActiveEditorArea();
// A gated empty Hub may never have been mounted. Restoring editor
// focus in that case would steal focus from a menu opened meanwhile.
if (overlayHandle) this.focusActiveEditorArea();
this.ctx.ui.requestRender();
};
@@ -20,18 +20,28 @@
* a superseded revive) can never clobber a newer same-id ref.
*/
import * as fs from "node:fs/promises";
import { logger } from "@oh-my-pi/pi-utils";
import type { AgentSession } from "../session/agent-session";
import {
type AgentRef,
type AgentRefExpectation,
AgentRegistry,
getAgentTombstonePath,
MAIN_AGENT_ID,
type RegistryEvent,
} from "./agent-registry";
export type AgentReviver = (expected: AgentRef) => Promise<AgentSession>;
async function persistAgentTombstone(sessionFile: string): Promise<void> {
try {
await fs.writeFile(getAgentTombstonePath(sessionFile), "", { encoding: "utf8", flag: "wx", mode: 0o600 });
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error;
}
}
/**
* Builds a reviver for a `parked` ref restored from disk (Agent Hub scan,
* collab mirror, resumed process) that carries a sessionFile but no in-memory
@@ -379,13 +389,10 @@ export class AgentLifecycleManager {
}
if (options?.tombstone) {
// Explicit kill: mark the ref terminal `aborted` and detach the session
// BEFORE disposing. aborted refs must satisfy the AgentRef invariant
// (session === null) so ensureLive / hub focus can never route into a
// disposed session; setting the terminal status first also makes
// createAgentSession's dispose wrapper (unregisterUnlessParked) preserve
// the ref instead of removing it, so a later persisted-subagent rescan
// skips it. The transcript file is left intact (history://<id>).
// Persist the terminal decision before detaching the session. The
// sidecar prevents a later discovery pass from reviving this transcript
// as a fresh parked ref.
if (ref.sessionFile) await persistAgentTombstone(ref.sessionFile);
this.#registry.setStatus(id, "aborted", ref);
}
const live = this.#registry.get(id) === ref ? ref.session : null;
@@ -14,6 +14,13 @@ import { oneLineLabel } from "../task/types";
export const MAIN_AGENT_ID = "Main";
/** Sidecar marker retained beside a child transcript after an explicit kill. */
export const AGENT_TOMBSTONE_SUFFIX = ".tombstone";
export function getAgentTombstonePath(sessionFile: string): string {
return `${sessionFile}${AGENT_TOMBSTONE_SUFFIX}`;
}
/**
* - `running`: a turn is in flight.
* - `idle`: live AgentSession in memory, awaiting work. Finished agents are
@@ -22,6 +29,8 @@ export const MAIN_AGENT_ID = "Main";
* - `aborted`: hard-killed, terminal.
*/
export type AgentStatus = "running" | "idle" | "parked" | "aborted";
/** Provenance of a displayed duration: active runtime, transcript span, or unavailable. */
export type AgentDurationKind = "active" | "span" | "unknown";
/**
* - `main`/`sub`: the user-facing agent tree (driving agent + task subagents).
* - `advisor`: a passive review transcript persisted like a subagent for usage
@@ -37,6 +46,7 @@ export interface AgentMetricsSummary {
tools: number;
cost: number;
durationMs: number;
durationKind?: AgentDurationKind;
contextTokens?: number;
contextWindow?: number;
}
@@ -46,6 +56,8 @@ export interface AgentHistorySummary {
agent?: string;
modelRole?: string;
resolvedModel?: string;
/** Whether the last resolved model was selected by retry fallback routing. */
resolvedModelIsFallback?: boolean;
metrics?: AgentMetricsSummary;
readOnly?: boolean;
}
@@ -153,7 +165,10 @@ export class AgentRegistry {
setHistory(id: string, history: AgentHistorySummary, expectedSessionFile?: string): boolean {
const ref = this.#refs.get(id);
if (!ref || (expectedSessionFile !== undefined && ref.sessionFile !== expectedSessionFile)) return false;
ref.history = { ...ref.history, ...history };
const definedHistory = Object.fromEntries(
Object.entries(history).filter(([, value]) => value !== undefined),
) as AgentHistorySummary;
ref.history = { ...ref.history, ...definedHistory };
this.#emit({ type: "metadata_changed", ref });
return true;
}
@@ -1,23 +1,21 @@
import * as fs from "node:fs";
import * as path from "node:path";
import { readLines } from "@oh-my-pi/pi-utils";
import { ADVISOR_TRANSCRIPT_FILENAME, isAdvisorTranscriptName } from "../advisor/transcript-recorder";
import { resolveExplicitModelRole } from "../config/model-resolver";
import { EPHEMERAL_MODEL_CHANGE_ROLE } from "../session/session-entries";
import { visitEntriesFromFileStream } from "../session/session-loader";
import { loadBundledAgents } from "../task/agents";
import { isReadOnlyAgent } from "../task/read-only-policy";
import { persistedVibeChildIds } from "../vibe/runtime";
import {
type AgentHistorySummary,
type AgentMetricsSummary,
type AgentRegistry,
getAgentTombstonePath,
MAIN_AGENT_ID,
} from "./agent-registry";
/**
* Child ids owned by the Vibe roster persisted in this session file. Vibe
* workers are revived through the Vibe registry's own journal, so the generic
* persisted-subagent scan must not register them as plain `sub` refs.
*/
const VIBE_LIFECYCLE_MARKER = Buffer.from('"vibe-session-lifecycle"');
/** Maximum prefix entries inspected for task metadata. */
const MAX_METADATA_LINES = 64;
interface PersistedAgentMetadata {
@@ -61,15 +59,6 @@ function summarizePersistedTask(task: string): string | undefined {
return summary ? summary.slice(0, 1_000) : undefined;
}
const READ_ONLY_AGENT_TOOLS: Record<string, true> = {
read: true,
grep: true,
glob: true,
web_search: true,
ast_grep: true,
yield: true,
};
function finiteNumber(value: unknown): number {
return typeof value === "number" && Number.isFinite(value) ? value : 0;
}
@@ -86,7 +75,7 @@ function inferBundledAgent(systemPrompt: string): { agent?: string; modelRole?:
return {
agent: agent.name,
modelRole: resolveExplicitModelRole(agent.model),
readOnly: !!agent.tools?.length && agent.tools.every(tool => READ_ONLY_AGENT_TOOLS[tool] === true),
readOnly: isReadOnlyAgent(agent),
};
}
@@ -121,29 +110,32 @@ function assistantMetrics(message: Record<string, unknown>): AssistantMetrics {
async function readPersistedAgentHistory(transcript: PersistedTranscript): Promise<AgentHistorySummary> {
const parents = new Map<string, string | undefined>();
const assistantById = new Map<string, AssistantMetrics>();
const modelChangeById = new Map<string, { model: string; role?: string; resolvedModelIsFallback: boolean }>();
let leafId: string | undefined;
let leafTimestamp: number | undefined;
try {
for await (const bytes of readLines(Bun.file(transcript.sessionFile).stream())) {
const line = Buffer.isBuffer(bytes) ? bytes : Buffer.from(bytes.buffer, bytes.byteOffset, bytes.byteLength);
const prefix = line.subarray(0, Math.min(line.byteLength, 512)).toString("utf8");
const id = /"id":"([^"]+)"/.exec(prefix)?.[1];
if (!id) continue;
const parentMatch = /"parentId":(?:"([^"]+)"|null)/.exec(prefix);
parents.set(id, parentMatch?.[1]);
await visitEntriesFromFileStream(transcript.sessionFile, entry => {
const record = recordOf(entry);
if (!record) return;
const id = typeof record.id === "string" ? record.id : undefined;
if (!id) return;
const parentId = typeof record.parentId === "string" ? record.parentId : undefined;
parents.set(id, parentId);
leafId = id;
const entryTimestamp = /"timestamp":"([^"]+)"/.exec(prefix)?.[1];
const parsedTimestamp = timestampOf(entryTimestamp);
const parsedTimestamp = timestampOf(record.timestamp);
if (parsedTimestamp !== undefined) leafTimestamp = parsedTimestamp;
if (!prefix.includes('"type":"message"') || !line.includes(Buffer.from('"role":"assistant"'))) continue;
try {
const entry = recordOf(JSON.parse(line.toString("utf8")));
const message = recordOf(entry?.message);
if (message?.role === "assistant") assistantById.set(id, assistantMetrics(message));
} catch {
// One malformed historical entry must not erase valid totals.
if (record.type === "model_change" && typeof record.model === "string") {
modelChangeById.set(id, {
model: record.model,
role: typeof record.role === "string" ? record.role : undefined,
resolvedModelIsFallback: record.resolvedModelIsFallback === true,
});
return;
}
}
if (record.type !== "message") return;
const message = recordOf(record.message);
if (message?.role === "assistant") assistantById.set(id, assistantMetrics(message));
});
} catch {
return {};
}
@@ -158,23 +150,42 @@ async function readPersistedAgentHistory(transcript: PersistedTranscript): Promi
(leafTimestamp ?? transcript.lastActivity ?? transcript.createdAt ?? 0) -
(transcript.createdAt ?? leafTimestamp ?? 0),
),
durationKind: "span",
};
let resolvedModel: string | undefined;
let resolvedModelIsFallback: boolean | undefined;
let modelRole: string | undefined;
let contextTokens: number | undefined;
let modelChangeFound = false;
const visited = new Set<string>();
for (let id = leafId; id && !visited.has(id); id = parents.get(id)) {
visited.add(id);
const modelChange = modelChangeById.get(id);
if (modelChange && !modelChangeFound) {
modelChangeFound = true;
resolvedModel = modelChange.model;
resolvedModelIsFallback = modelChange.resolvedModelIsFallback;
if (modelChange.role && modelChange.role !== EPHEMERAL_MODEL_CHANGE_ROLE) {
modelRole = modelChange.role;
}
}
const assistant = assistantById.get(id);
if (!assistant) continue;
if (!modelChangeFound && resolvedModel === undefined && assistant.resolvedModel) {
resolvedModel = assistant.resolvedModel;
}
metrics.requests++;
metrics.tokens += assistant.tokens;
metrics.tools += assistant.tools;
metrics.cost += assistant.cost;
contextTokens ??= assistant.contextTokens;
resolvedModel ??= assistant.resolvedModel;
}
metrics.contextTokens = contextTokens;
return { metrics: metrics.requests > 0 ? metrics : undefined, resolvedModel };
if (contextTokens !== undefined) metrics.contextTokens = contextTokens;
return {
...(metrics.requests > 0 ? { metrics } : {}),
...(resolvedModel ? { resolvedModel, resolvedModelIsFallback } : {}),
...(modelRole ? { modelRole } : {}),
};
}
/**
@@ -188,43 +199,42 @@ async function readPersistedAgentMetadata(sessionFile: string): Promise<Persiste
let activity: string | undefined;
let history: AgentHistorySummary = {};
try {
let linesRead = 0;
for await (const bytes of readLines(Bun.file(sessionFile).stream())) {
if (linesRead++ >= MAX_METADATA_LINES) break;
const line = Buffer.isBuffer(bytes) ? bytes : Buffer.from(bytes.buffer, bytes.byteOffset, bytes.byteLength);
let entry: Record<string, unknown> | undefined;
try {
entry = recordOf(JSON.parse(line.toString("utf8")));
} catch {
continue;
}
if (entry?.type === "session") {
createdAt ??= timestampOf(entry.timestamp);
continue;
}
if (entry?.type === "model_change") {
if (typeof entry.model === "string") history.resolvedModel = entry.model;
if (typeof entry.role === "string") history.modelRole = entry.role;
continue;
}
if (entry?.type !== "session_init") continue;
createdAt ??= timestampOf(entry.timestamp);
if (typeof entry.task === "string") activity = summarizePersistedTask(entry.task);
const inferred =
typeof entry.systemPrompt === "string"
? inferBundledAgent(entry.systemPrompt)
: ({} satisfies AgentHistorySummary);
history = {
...history,
...inferred,
agent: typeof entry.agent === "string" ? entry.agent : inferred.agent,
modelRole:
typeof entry.modelRole === "string" ? entry.modelRole : (history.modelRole ?? inferred.modelRole),
resolvedModel: typeof entry.resolvedModel === "string" ? entry.resolvedModel : history.resolvedModel,
readOnly: typeof entry.readOnly === "boolean" ? entry.readOnly : inferred.readOnly,
};
break;
}
await visitEntriesFromFileStream(
sessionFile,
entry => {
const record = recordOf(entry);
if (!record) return;
if (record.type === "session") {
createdAt ??= timestampOf(record.timestamp);
return;
}
if (record.type === "model_change") {
if (typeof record.model === "string") history.resolvedModel = record.model;
if (typeof record.role === "string" && record.role !== EPHEMERAL_MODEL_CHANGE_ROLE) {
history.modelRole = record.role;
}
if (typeof record.resolvedModelIsFallback === "boolean") {
history.resolvedModelIsFallback = record.resolvedModelIsFallback;
}
return;
}
if (record.type !== "session_init") return;
createdAt ??= timestampOf(record.timestamp);
if (typeof record.task === "string") activity = summarizePersistedTask(record.task);
const inferred = typeof record.systemPrompt === "string" ? inferBundledAgent(record.systemPrompt) : {};
history = {
...history,
...inferred,
agent: typeof record.agent === "string" ? record.agent : inferred.agent,
modelRole:
typeof record.modelRole === "string" ? record.modelRole : (history.modelRole ?? inferred.modelRole),
resolvedModel: typeof record.resolvedModel === "string" ? record.resolvedModel : history.resolvedModel,
readOnly: typeof record.readOnly === "boolean" ? record.readOnly : inferred.readOnly,
};
return false;
},
{ maxRecords: MAX_METADATA_LINES },
);
} catch {
// A readable transcript is still useful even when its optional metadata
// prefix is malformed.
@@ -241,17 +251,9 @@ async function readPersistedAgentMetadata(sessionFile: string): Promise<Persiste
async function readPersistedVibeChildIds(sessionFile: string): Promise<Set<string>> {
const ids = new Set<string>();
try {
for await (const bytes of readLines(Bun.file(sessionFile).stream())) {
const line = Buffer.isBuffer(bytes) ? bytes : Buffer.from(bytes.buffer, bytes.byteOffset, bytes.byteLength);
if (line.indexOf(VIBE_LIFECYCLE_MARKER) === -1) continue;
try {
const entry: unknown = JSON.parse(line.toString("utf8"));
for (const id of persistedVibeChildIds([entry])) ids.add(id);
} catch {
// Match lenient session loading: one malformed line must not hide
// valid lifecycle entries later in the transcript.
}
}
await visitEntriesFromFileStream(sessionFile, entry => {
for (const id of persistedVibeChildIds([entry])) ids.add(id);
});
return ids;
} catch {
return new Set();
@@ -294,7 +296,12 @@ async function registerPersistedSubagentsFromDir(
} catch {
return;
}
let entriesSinceYield = 0;
for (const entry of entries) {
if (++entriesSinceYield >= 16) {
entriesSinceYield = 0;
await Bun.sleep(0);
}
if (!entry.isFile() || !entry.name.endsWith(".jsonl") || entry.name.includes(".bak")) continue;
const sessionFile = path.join(dir, entry.name);
// The advisor transcript is observability-only: register it as a non-peer
@@ -340,6 +347,13 @@ async function registerPersistedSubagentsFromDir(
}
const id = entry.name.slice(0, -6);
if (vibeOwnedIds.has(id) && registry.get(id)?.sessionFile !== sessionFile) continue;
let tombstoned = false;
try {
await fs.promises.access(getAgentTombstonePath(sessionFile));
tombstoned = true;
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "ENOENT") continue;
}
if (!registry.get(id)) {
const metadata = await readPersistedAgentMetadata(sessionFile);
registry.register({
@@ -353,7 +367,7 @@ async function registerPersistedSubagentsFromDir(
createdAt: metadata.createdAt,
lastActivity: metadata.lastActivity,
history: metadata.history,
status: "parked",
status: tombstoned ? "aborted" : "parked",
});
const ref = registry.get(id);
transcripts.push({
@@ -84,6 +84,8 @@ export interface ModelChangeEntry extends SessionEntryBase {
model: string;
/** Role: "default", "smol", "slow", etc. Undefined treated as "default" */
role?: string;
/** True when this transition selected a retry-fallback model rather than the configured model. */
resolvedModelIsFallback?: boolean;
}
export interface ServiceTierChangeEntry extends SessionEntryBase {
@@ -21,6 +21,19 @@ import {
} from "./session-title-slot";
const STREAM_LOAD_THRESHOLD_BYTES = 8 * 1024 * 1024;
const STREAM_YIELD_BYTES = 1 * 1024 * 1024;
const STREAM_YIELD_ENTRIES = 8_192;
export interface VisitEntriesFromFileStreamOptions {
/** Stop after the visitor returns `false`. */
shouldContinue?: () => boolean;
/** Stop after this many valid or malformed JSONL records have been consumed. */
maxRecords?: number;
/** Yield to the macrotask queue after this many bytes have been consumed. */
yieldEveryBytes?: number;
/** Yield to the macrotask queue after this many entries have been visited. */
yieldEveryEntries?: number;
}
function splitTitleSlot(content: string): { body: string; slot: SessionTitleUpdate | undefined } {
const slot = titleUpdateFromSlot(parseTitleSlotFromContent(content));
@@ -59,11 +72,19 @@ export function parseSessionContent(content: string): {
/** Parse session JSONL and visit each entry without retaining prior entries. */
export async function visitEntriesFromFileStream(
filePath: string,
visit: (entry: FileEntry) => void,
visit: (entry: FileEntry) => void | boolean,
options: VisitEntriesFromFileStreamOptions = {},
): Promise<SessionTitleUpdate | undefined> {
let titleSlot: SessionTitleUpdate | undefined;
let sawFirstLine = false;
let bytesSinceYield = 0;
let entriesSinceYield = 0;
let recordsSeen = 0;
const maxRecords = Math.max(0, options.maxRecords ?? Number.POSITIVE_INFINITY);
let stopped = false;
let visitorThrew = false;
const yieldEveryBytes = Math.max(0, options.yieldEveryBytes ?? STREAM_YIELD_BYTES);
const yieldEveryEntries = Math.max(0, options.yieldEveryEntries ?? STREAM_YIELD_ENTRIES);
// Byte buffer (NOT a decoded string): multibyte UTF-8 sequences that straddle
// a stream-chunk boundary stay intact, and Bun.JSONL.parseChunk accepts typed
// arrays directly. Only the unconsumed remainder is held (≤ one record + a
@@ -72,22 +93,62 @@ export async function visitEntriesFromFileStream(
let buffer: Uint8Array = new Uint8Array();
const decoder = new TextDecoder();
const drain = () => {
while (buffer.length > 0) {
const yieldToMacrotask = async (): Promise<void> => {
if (yieldEveryBytes === 0 && yieldEveryEntries === 0) return;
const bytesReady = yieldEveryBytes === 0 || bytesSinceYield < yieldEveryBytes;
const entriesReady = yieldEveryEntries === 0 || entriesSinceYield < yieldEveryEntries;
if (bytesReady && entriesReady) {
return;
}
bytesSinceYield = 0;
entriesSinceYield = 0;
await Bun.sleep(0);
};
const drain = async (): Promise<void> => {
while (buffer.length > 0 && !stopped) {
if (recordsSeen >= maxRecords) {
stopped = true;
break;
}
const { values, error, read, done } = Bun.JSONL.parseChunk(buffer);
for (const value of values) {
if (recordsSeen >= maxRecords) {
stopped = true;
break;
}
if (options.shouldContinue && !options.shouldContinue()) {
stopped = true;
break;
}
try {
visit(value as FileEntry);
if (visit(value as FileEntry) === false) {
stopped = true;
break;
}
recordsSeen++;
entriesSinceYield++;
if (recordsSeen >= maxRecords) {
stopped = true;
break;
}
} catch (err) {
visitorThrew = true;
throw err;
}
await yieldToMacrotask();
}
if (stopped) break;
if (error) {
// Malformed record: skip past the next newline and continue.
const nextNewline = buffer.indexOf(0x0a, read);
if (nextNewline === -1) break; // rest of the bad line not yet received
recordsSeen++;
buffer = buffer.subarray(nextNewline + 1);
if (recordsSeen >= maxRecords) {
stopped = true;
break;
}
continue;
}
if (read === 0) break; // incomplete record awaiting more data
@@ -101,6 +162,8 @@ export async function visitEntriesFromFileStream(
try {
for await (const chunk of Bun.file(filePath).stream()) {
if (stopped) break;
bytesSinceYield += chunk.byteLength;
buffer = buffer.length === 0 ? chunk : Buffer.concat([buffer, chunk]);
// The optional fixed-width title slot is a physical first line that is
// NOT JSON; peel it before the parser would (correctly) reject it. The
@@ -121,14 +184,15 @@ export async function visitEntriesFromFileStream(
}
}
}
drain();
await drain();
await yieldToMacrotask();
}
// A trailing record without a final newline: terminate it so the parser
// can complete it (readline yielded it; parseChunk needs the delimiter).
if (buffer.length > 0 && buffer[buffer.length - 1] !== 0x0a) {
if (!stopped && buffer.length > 0 && buffer[buffer.length - 1] !== 0x0a) {
buffer = Buffer.concat([buffer, new Uint8Array([0x0a])]);
await drain();
}
drain();
} catch (err) {
if (visitorThrew) throw err;
if (isEnoent(err)) return undefined;
@@ -144,7 +208,9 @@ export async function loadEntriesFromFileStream(filePath: string): Promise<{
titleSlot: SessionTitleUpdate | undefined;
}> {
const entries: FileEntry[] = [];
const titleSlot = await visitEntriesFromFileStream(filePath, entry => entries.push(entry));
const titleSlot = await visitEntriesFromFileStream(filePath, entry => {
entries.push(entry);
});
return { entries: foldTitleSlot(entries, titleSlot), titleSlot };
}
@@ -2047,9 +2047,16 @@ export class SessionManager {
* Append a model change as a child of the current leaf, then advance the leaf.
* @param model Model in "provider/modelId" format
* @param role Optional role (default: "default")
* @param resolvedModelIsFallback Whether this transition selected a retry-fallback model
*/
appendModelChange(model: string, role?: string): string {
const entry: ModelChangeEntry = { type: "model_change", ...this.#freshEntryFields(), model, role };
appendModelChange(model: string, role?: string, resolvedModelIsFallback = false): string {
const entry: ModelChangeEntry = {
type: "model_change",
...this.#freshEntryFields(),
model,
role,
resolvedModelIsFallback,
};
this.#recordEntry(entry);
return entry.id;
}
@@ -2551,6 +2558,9 @@ export class SessionManager {
systemPrompt: string;
task: string;
tools: string[];
agent?: string;
modelRole?: string;
resolvedModel?: string;
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
restrictToolNames?: boolean;
@@ -2571,6 +2581,9 @@ export class SessionManager {
systemPrompt: string;
task: string;
tools: string[];
agent?: string;
modelRole?: string;
resolvedModel?: string;
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
restrictToolNames?: boolean;
@@ -2584,6 +2597,9 @@ export class SessionManager {
systemPrompt: entry.systemPrompt,
task: entry.task,
tools: entry.tools,
agent: entry.agent,
modelRole: entry.modelRole,
resolvedModel: entry.resolvedModel,
outputSchema: entry.outputSchema,
outputSchemaMode: entry.outputSchemaMode,
restrictToolNames: entry.restrictToolNames,
@@ -1102,7 +1102,7 @@ export class TurnRecovery {
: clampThinkingLevelToCeiling(candidate, requestedThinkingLevel, this.#host.thinkingLevelCeiling());
const candidateSelector = formatModelStringWithRouting(candidate);
await this.#host.setModelWithProviderSessionReset(candidate);
this.#host.sessionManager.appendModelChange(candidateSelector, EPHEMERAL_MODEL_CHANGE_ROLE);
this.#host.sessionManager.appendModelChange(candidateSelector, EPHEMERAL_MODEL_CHANGE_ROLE, true);
this.#host.settings.getStorage()?.recordModelUsage(candidateSelector);
this.#host.setThinkingLevel(nextThinkingLevel);
if (!this.#activeRetryFallback) {
@@ -1220,7 +1220,7 @@ export class TurnRecovery {
if (!apiKey) return false;
const baseSelector = formatModelStringWithRouting(baseModel);
await this.#host.setModelWithProviderSessionReset(baseModel);
this.#host.sessionManager.appendModelChange(baseSelector, EPHEMERAL_MODEL_CHANGE_ROLE);
this.#host.sessionManager.appendModelChange(baseSelector, EPHEMERAL_MODEL_CHANGE_ROLE, true);
this.#host.settings.getStorage()?.recordModelUsage(baseSelector);
await this.#host.emitSessionEvent({
type: "retry_fallback_applied",
+16 -8
View File
@@ -59,6 +59,7 @@ import { buildNamedToolChoice } from "../utils/tool-choice";
import type { WorkspaceTree } from "../workspace-tree";
import { generateTaskLabel } from "./label";
import { resolveAgentPrewalkDefault } from "./prewalk";
import { isReadOnlyAgent } from "./read-only-policy";
import { subprocessToolRegistry } from "./subprocess-tool-registry";
import {
type AgentDefinition,
@@ -2378,14 +2379,20 @@ export async function finalizeSubagentLifecycle(args: {
args.abortKind === "budget" && args.keepAlive && !args.isolated && args.reviveSession !== null;
if (args.aborted && !resumableAbort) {
if (ref && ownsRef) {
// Terminal hard kill: mark `aborted` and detach the session before
// disposing so the ref satisfies the AgentRef invariant (session null
// when aborted) — ensureLive/hub focus must treat it as terminal, never
// route into the disposed session.
registry.setStatus(args.id, "aborted", ref);
registry.detachSession(args.id, ref);
// Route hard kills through the lifecycle owner so the terminal
// decision is durable and a restart cannot rediscover the transcript
// as a revivable parked agent.
try {
await AgentLifecycleManager.global().release(args.id, ref, { tombstone: true });
} catch (error) {
logger.warn("runSubagent: failed to persist kill tombstone", { id: args.id, error: String(error) });
registry.setStatus(args.id, "aborted", ref);
registry.detachSession(args.id, ref);
await disposeSession();
}
} else {
await disposeSession();
}
await disposeSession();
return;
}
@@ -3094,8 +3101,9 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
task,
tools: session.getActiveToolNames(),
agent: agent.name,
modelRole: modelRole ?? resolveExplicitModelRole(agent.model),
modelRole: modelRole ?? resolveExplicitModelRole(agent.model, subagentSettings),
resolvedModel: progress.resolvedModel,
readOnly: isReadOnlyAgent(agent),
spawns: spawnsEnv,
readSummarize: agent.readSummarize,
outputSchema,
+6 -27
View File
@@ -27,6 +27,7 @@ import { TASK_EFFORTS, type TaskEffort } from "../thinking";
import { truncateForPrompt } from "../tools/approval";
import { isIrcEnabled } from "../tools/hub";
import { formatBytes, formatDuration } from "../tools/render-utils";
import { isReadOnlyAgent } from "./read-only-policy";
import { isScoutSpawnable, resolveSpawnPolicy } from "./spawn-policy";
import {
type AgentDefinition,
@@ -102,6 +103,7 @@ export { loadBundledAgents as BUNDLED_AGENTS } from "./agents";
export { discoverCommands, expandCommand, getCommand } from "./commands";
export { discoverAgents, getAgent } from "./discovery";
export { AgentOutputManager } from "./output-manager";
export * from "./read-only-policy";
export type {
AgentDefinition,
AgentProgress,
@@ -119,32 +121,6 @@ export {
taskSchema,
} from "./types";
// Built-in tools whose approval tier is "read" (see tool classes' `approval`).
// An agent is read-only iff its declared tools are a non-empty subset of this set.
// Fail-safe: any unknown tool makes the agent not read-only.
export const READ_ONLY_TOOL_NAMES: ReadonlySet<string> = new Set([
"read",
"grep",
"glob",
"web_search",
"ast_grep",
"yield",
"hub",
"ask",
"todo",
"recall",
"reflect",
"retain",
"memory_edit",
"inspect_image",
"checkpoint",
"rewind",
]);
export function isReadOnlyAgent(agent: AgentDefinition): boolean {
return !!agent.tools?.length && agent.tools.every(tool => READ_ONLY_TOOL_NAMES.has(tool));
}
/**
* Preview text for a child result. Falls back to "(no output)" — annotated
* with the request count when the child actually did work, so the parent can
@@ -1146,6 +1122,8 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
// status back to the subagent's initial "pending" snapshot.
progress.modelRole = nextProgress.modelRole ?? progress.modelRole;
progress.resolvedModel = nextProgress.resolvedModel;
progress.resolvedModelIsFallback =
nextProgress.resolvedModelIsFallback ?? progress.resolvedModelIsFallback;
progress.tokens = nextProgress.tokens;
progress.requests = nextProgress.requests;
progress.contextTokens = nextProgress.contextTokens;
@@ -1191,7 +1169,8 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
progress.modelRole = singleResult?.modelRole ?? progress.modelRole;
if (singleResult?.resolvedModel) {
progress.resolvedModel = singleResult.resolvedModel;
progress.resolvedModelIsFallback = singleResult.resolvedModelIsFallback;
progress.resolvedModelIsFallback =
singleResult.resolvedModelIsFallback ?? progress.resolvedModelIsFallback;
} else {
delete progress.resolvedModel;
delete progress.resolvedModelIsFallback;
@@ -1,6 +1,6 @@
import * as fs from "node:fs/promises";
import type { ModelRegistry } from "../config/model-registry";
import { formatModelRoleAlias } from "../config/model-roles";
import type { Settings } from "../config/settings";
import { MCPManager } from "../mcp/manager";
import type { PersistedSubagentReviverFactory } from "../registry/agent-lifecycle";
@@ -79,6 +79,14 @@ export function createPersistedSubagentReviverFactory(
taskDepth++;
parentId = registry.get(parentId)?.parentId;
}
const subagentSettings = createSubagentSettings(
ctx.settings,
init.readSummarize === false ? { "read.summarize.enabled": false } : undefined,
);
const persistedModelPattern =
init.modelRole && init.modelRole !== "default"
? [formatModelRoleAlias(init.modelRole), ...(init.resolvedModel ? [init.resolvedModel] : [])]
: init.resolvedModel;
return async expectedRef => {
// Re-open fresh on every revive: park closes the writer, so this takes
// the single-writer lock cleanly and restores the full message history.
@@ -96,10 +104,9 @@ export function createPersistedSubagentReviverFactory(
cwd: ctx.session.sessionManager.getCwd(),
authStorage: ctx.authStorage,
modelRegistry: ctx.modelRegistry,
settings: createSubagentSettings(
ctx.settings,
init.readSummarize === false ? { "read.summarize.enabled": false } : undefined,
),
...(persistedModelPattern ? { modelPattern: persistedModelPattern } : {}),
modelPatternAuthFallback: init.resolvedModel,
settings: subagentSettings,
sessionManager: reopened,
agentId: ref.id,
agentDisplayName: ref.displayName,
@@ -0,0 +1,27 @@
import type { AgentDefinition } from "./types";
// Built-in tools whose approval tier is "read" (see tool classes' `approval`).
// An agent is read-only iff its declared tools are a non-empty subset of this set.
// Fail-safe: any unknown tool makes the agent not read-only.
export const READ_ONLY_TOOL_NAMES: ReadonlySet<string> = new Set([
"read",
"grep",
"glob",
"web_search",
"ast_grep",
"yield",
"hub",
"ask",
"todo",
"recall",
"reflect",
"retain",
"memory_edit",
"inspect_image",
"checkpoint",
"rewind",
]);
export function isReadOnlyAgent(agent: AgentDefinition): boolean {
return !!agent.tools?.length && agent.tools.every(tool => READ_ONLY_TOOL_NAMES.has(tool));
}
@@ -15,6 +15,7 @@ import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
import type { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { visitEntriesFromFileStream } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getBundledAgent } from "@oh-my-pi/pi-coding-agent/task/agents";
import { TempDir } from "@oh-my-pi/pi-utils";
@@ -296,62 +297,43 @@ describe("Agent hub Enter activation", () => {
expect(Bun.stripANSI(hub.render(120).join("\n"))).toContain("Read-only · 0 LoC");
hub.dispose();
});
it("avoids multi-second event-loop stalls while discovering agents from a large session", async () => {
it("yields to a macrotask while streaming a large session", async () => {
vi.useFakeTimers();
using tempDir = TempDir.createSync("@omp-agent-hub-responsive-");
const sessionFile = path.join(tempDir.path(), "main.jsonl");
const header = JSON.stringify({
type: "session",
version: 3,
id: "responsive-test",
timestamp: new Date().toISOString(),
cwd: tempDir.path(),
});
const payload = "x".repeat(250_000);
const message = JSON.stringify({
const sessionFile = path.join(tempDir.path(), "session.jsonl");
const entry = JSON.stringify({
type: "message",
id: "large-message",
id: "entry",
parentId: null,
timestamp: new Date().toISOString(),
message: {
role: "user",
content: [{ type: "text", text: payload }],
timestamp: Date.now(),
},
timestamp: "2026-07-30T01:13:30.000Z",
message: { role: "user", content: [{ type: "text", text: "small" }] },
});
await Bun.write(sessionFile, `${header}\n${Array.from({ length: 600 }, () => message).join("\n")}\n`);
Bun.gc(true);
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
await Bun.write(sessionFile, `${entry}\n`.repeat(8_193));
let complete = false;
let yieldedBeforeComplete = false;
let visited = 0;
const visit = visitEntriesFromFileStream(
sessionFile,
() => {
visited++;
if (visited !== 8_192) return;
setTimeout(() => {
if (!complete) yieldedBeforeComplete = true;
}, 0);
},
{ yieldEveryBytes: 0, yieldEveryEntries: 8_192 },
).finally(() => {
complete = true;
});
let maxLagMs = 0;
let previous = performance.now();
const timer = setInterval(() => {
const now = performance.now();
maxLagMs = Math.max(maxLagMs, now - previous - 2);
previous = now;
}, 2);
try {
await hub.persistedSubagentsReady;
await Bun.sleep(10);
for (let i = 0; i < 20_000 && visited < 8_192 && !complete; i++) await Promise.resolve();
expect(visited).toBeGreaterThanOrEqual(8_192);
vi.runOnlyPendingTimers();
await visit;
expect(yieldedBeforeComplete).toBe(true);
} finally {
clearInterval(timer);
hub.dispose();
vi.useRealTimers();
}
// Guard the original seconds-long freeze without treating shared-runner
// scheduling jitter as an Agent Hub regression.
expect(maxLagMs).toBeLessThan(500);
});
it("does not generically revive active or tombstoned Vibe children copied by a post-exit fork", async () => {
@@ -767,7 +749,8 @@ describe("Agent hub data refresh coalescing", () => {
const refreshed = Bun.stripANSI(hub.render(120).join("\n"));
expect(refreshed).toContain("450 tok");
expect(refreshed).toContain("2 req");
expect(refreshed).toContain("1/1 measured");
expect(refreshed).toContain("1/1");
expect(refreshed).toContain("measured");
hub.render(120);
expect(getSessionStats).toHaveBeenCalledTimes(2);
} finally {
@@ -775,4 +758,44 @@ describe("Agent hub data refresh coalescing", () => {
vi.useRealTimers();
}
});
it("counts shared fallback session usage once across parent and descendant rows", () => {
const agents = new AgentRegistry();
const getSessionStats = vi.fn(() => ({
tokens: { input: 100, output: 40, cacheRead: 10, cacheWrite: 10, total: 160 },
assistantMessages: 1,
toolCalls: 2,
cost: 0.1,
contextUsage: undefined,
}));
const session = { getSessionStats } as unknown as AgentSession;
agents.register({ id: "Parent", displayName: "Parent", kind: "sub", session, status: "idle" });
agents.register({
id: "Child",
displayName: "Child",
kind: "sub",
parentId: "Parent",
session,
status: "idle",
});
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
});
try {
const rendered = Bun.stripANSI(hub.render(120).join("\n"));
expect(rendered).toContain("150 tok");
expect(rendered).toContain("1/2");
expect(rendered).toContain("measured");
expect(getSessionStats).toHaveBeenCalledTimes(1);
} finally {
hub.dispose();
}
});
});
@@ -176,17 +176,18 @@ describe("Agent hub row ordering", () => {
hub = makeHub(agents);
expect(renderedAgentIds(hub)).toEqual(["C", "B", "A"]);
// Bump A's lastActivity far ahead of the others. The hub is already open,
// so the captured order must not change.
// Bump A's lastActivity far ahead of the others; captured order wins.
setSystemTime(4000);
agents.setActivity("A", "still running");
// Registering a new agent schedules a coalesced row refresh; the
// existing rows must stay put once the scheduled refresh runs.
// Status changes must not reorder the captured roster either.
agents.setStatus("B", "idle");
// Registering a new agent schedules a coalesced row refresh; even a
// different status is appended after all rows captured on open.
setSystemTime(5000);
const sessionD = {} as AgentSession;
agents.register({ id: "D", displayName: "Delta", kind: "sub", session: sessionD });
agents.register({ id: "D", displayName: "Delta", kind: "sub", session: sessionD, status: "parked" });
expect(renderedAgentIds(hub)).toEqual(["C", "B", "A"]);
vi.advanceTimersByTime(100);
@@ -204,7 +205,7 @@ describe("Agent hub row ordering", () => {
const sessionA = {} as AgentSession;
agents.register({
id: "RevAgentStream",
displayName: "Agent runtime + compaction reviewer",
displayName: "Agent runtime + compaction reviewer\u0007",
kind: "sub",
session: sessionA,
});
@@ -218,12 +219,17 @@ describe("Agent hub row ordering", () => {
status: "active",
description: "Complete the assignment below, thoroughly:\n- check performance\n- check leaks",
lastUpdate: Date.now(),
progress: {
currentTool: "bash",
currentToolArgs: "\x1b[2Jdangerous args",
} as never,
},
]);
const hub = makeHub(agents, { observers });
const lines = hub.render(80);
expect(lines.join("\n")).not.toContain("\u0007");
for (const line of lines) {
const cleanLine = Bun.stripANSI(line);
expect(cleanLine.includes("\n")).toBe(false);
@@ -231,6 +237,10 @@ describe("Agent hub row ordering", () => {
const width = visibleWidth(line);
expect(width).toBeLessThanOrEqual(80);
}
hub.handleInput("\t");
const details = hub.render(80).join("\n");
expect(details).not.toContain("\x1b[2J");
expect(details).toContain("dangerous args");
hub.dispose();
});
@@ -287,14 +297,12 @@ describe("Agent hub row ordering", () => {
const frame = hub.render(120);
const alphaRow = frame.findIndex(line => /^│ {3}\S+ Alpha/u.test(Bun.stripANSI(line)));
expect(alphaRow).toBeGreaterThanOrEqual(0);
hub.handleInput(leftClick(alphaRow + 1));
expect(selectedAgentId(hub)).toBe("Alpha");
hub.handleInput(`\x1b[<0;110;${alphaRow + 1}M`);
expect(selectedAgentId(hub)).toBe("Beta");
expect(focused).toEqual([]);
const selectedFrame = hub.render(120);
const selectedAlphaRow = selectedFrame.findIndex(line => /^│ ❯ \S+ Alpha/u.test(Bun.stripANSI(line)));
hub.handleInput(leftClick(selectedAlphaRow + 1));
hub.handleInput(leftClick(alphaRow + 1));
await Promise.resolve();
expect(selectedAgentId(hub)).toBe("Alpha");
expect(focused).toEqual(["Alpha"]);
expect(done).toHaveBeenCalledTimes(1);
} finally {
@@ -468,7 +476,7 @@ describe("Agent hub row ordering", () => {
expect(rendered).toContain("1 running");
expect(rendered).toContain("Flat");
expect(rendered).toContain("By parent");
expect(rendered).toContain("$0.213 · 2m14s · 12 req · 27 tools · 18K tok");
expect(rendered).toContain("$0.213 · 2m14s active · 12 req · 27 tools · 18K tok");
expect(rendered).toContain("Security Reviewer");
expect(rendered).toContain("read · src/session/agent-session.ts");
expect(rendered).toContain("31K/128K 24%");
@@ -548,7 +556,7 @@ describe("Agent hub row ordering", () => {
expect(rendered).toContain("3.7K tok");
expect(rendered).toContain("8 req");
expect(rendered).toContain("12 tools");
expect(rendered).toContain("2m11s agent time");
expect(rendered).toContain("2m11s active agent time");
const running = renderedRosterEntry(hub, "Running", 160);
expect(running).toContain("$0.123");
@@ -774,9 +782,25 @@ describe("Agent hub row ordering", () => {
}
});
it("keeps cyclic parent links renderable in tree mode", () => {
const agents = new AgentRegistry();
agents.register({ id: "CycleA", displayName: "Cycle A", kind: "sub", parentId: "CycleB", session: null });
agents.register({ id: "CycleB", displayName: "Cycle B", kind: "sub", parentId: "CycleA", session: null });
const hub = makeHub(agents);
try {
hub.handleInput("t");
const rendered = Bun.stripANSI(hub.render(120).join("\n"));
expect(rendered).toContain("CycleA");
expect(rendered).toContain("CycleB");
} finally {
hub.dispose();
}
});
it("opens the selected-agent inspector as a narrow-terminal fallback", () => {
geometry = stubStdoutGeometry(80);
geometry.setRows(28);
geometry.setRows(12);
const agents = new AgentRegistry();
agents.register({ id: "NarrowAgent", displayName: "Narrow Agent", kind: "sub", session: null });
const observers = new SessionObserverRegistry();
@@ -812,8 +836,10 @@ describe("Agent hub row ordering", () => {
const details = Bun.stripANSI(hub.render(80).join("\n"));
expect(details).toContain("Agent Hub · NarrowAgent");
expect(details).toContain("Usage");
expect(details).toContain("$0.0000 · 2.0s · 2 req · 3 tools · 900 tok");
expect(details).toContain("$0.0000 · 2.0s active · 2 req · 3 tools · 900 tok");
expect(details).toContain("Tab:roster");
hub.handleInput("\x1b[6~");
expect(Bun.stripANSI(hub.render(80).join("\n"))).toContain("Changes");
for (const line of hub.render(80)) expect(visibleWidth(line)).toBeLessThanOrEqual(80);
hub.handleInput("\x1b");
@@ -597,10 +597,11 @@ describe("AgentLifecycleManager", () => {
// session (the ref carries session === null), it treats it as unrevivable.
await expect(lifecycle.ensureLive(workerId)).rejects.toThrow(/aborted/);
// Reopening the Agent Hub rescans on-disk transcripts. The surviving
// `.jsonl` must not be re-adopted as a fresh `parked` row, because the
// id is still present in the registry.
await registerPersistedSubagents(registry, rootSessionFile);
expect(registry.get(workerId)?.status).toBe("aborted");
// Reopening after the original registry is gone must preserve the terminal
// decision from the sidecar, not infer a fresh parked agent from the JSONL.
expect(await Bun.file(`${workerSessionFile}.tombstone`).exists()).toBe(true);
const restoredRegistry = new AgentRegistry();
await registerPersistedSubagents(restoredRegistry, rootSessionFile);
expect(restoredRegistry.get(workerId)?.status).toBe("aborted");
});
});
@@ -86,7 +86,9 @@ describe("loadEntriesFromFileStream (Bun.JSONL parity)", () => {
].join("\n");
const file = await writeTemp(content);
const visited: FileEntry[] = [];
const titleSlot = await sessionLoader.visitEntriesFromFileStream(file, entry => visited.push(entry));
const titleSlot = await sessionLoader.visitEntriesFromFileStream(file, entry => {
visited.push(entry);
});
expect(titleSlot?.title).toBe("Visitor");
expect(entryIds(visited)).toEqual(["s1", "m1", "m2"]);
@@ -101,11 +103,35 @@ describe("loadEntriesFromFileStream (Bun.JSONL parity)", () => {
const file = await writeTemp(content);
const visited: FileEntry[] = [];
await sessionLoader.visitEntriesFromFileStream(file, entry => visited.push(entry));
await sessionLoader.visitEntriesFromFileStream(file, entry => {
visited.push(entry);
});
expect(entryIds(visited)).toEqual(["s1", "m1", "m2"]);
});
it("bounds visitor scans by physical records, including malformed lines", async () => {
const content = [
JSON.stringify(HEADER),
"{ malformed one",
"{ malformed two",
"{ malformed three",
JSON.stringify(msg("after-bad", "s1", "must not be visited")),
].join("\n");
const file = await writeTemp(content);
const visited: FileEntry[] = [];
await sessionLoader.visitEntriesFromFileStream(
file,
entry => {
visited.push(entry);
},
{ maxRecords: 2 },
);
expect(entryIds(visited)).toEqual(["s1"]);
});
it("propagates ENOENT errors thrown by the visitor", async () => {
const file = await writeTemp(`${JSON.stringify(HEADER)}\n`);
const failure = Object.assign(new Error("visitor failed"), { code: "ENOENT" });
@@ -62,7 +62,7 @@ function createRevivedSession(activeToolNames: string[][]): RevivedSessionHandle
return { session, observer: () => observer };
}
async function createPersistedSession(cwd: string, restrictToolNames?: boolean): Promise<string> {
async function createPersistedSession(cwd: string, restrictToolNames?: boolean, modelRole?: string): Promise<string> {
const manager = SessionManager.create(cwd, path.join(cwd, "sessions"));
const sessionFile = manager.getSessionFile();
if (!sessionFile) throw new Error("Expected a persisted session file");
@@ -71,6 +71,8 @@ async function createPersistedSession(cwd: string, restrictToolNames?: boolean):
task: "persisted task",
tools: ["read", "yield"],
restrictToolNames,
modelRole,
resolvedModel: modelRole ? "anthropic/claude-sonnet-4-5" : undefined,
});
manager.appendMessage({
role: "assistant",
@@ -179,6 +181,42 @@ describe("persisted subagent revival", () => {
expect(capturedOptions?.customTools?.map(tool => tool.name)).toEqual(["mcp__server_read"]);
});
it("restores the persisted custom model role before reopening the session", async () => {
const cwd = makeTempDir("@pi-custom-role-revive-");
const sessionFile = await createPersistedSession(cwd, false, "review-fast");
let capturedOptions: CreateAgentSessionOptions | undefined;
vi.spyOn(sdkModule, "createAgentSession").mockImplementation(async options => {
capturedOptions = options;
return { session: createRevivedSession([]).session } as CreateAgentSessionResult;
});
const ref = createRef(sessionFile);
const reviver = await createFactory(cwd)(ref);
if (!reviver) throw new Error("Expected a persisted reviver");
await reviver(ref);
expect(capturedOptions?.modelPattern).toEqual(["@review-fast", "anthropic/claude-sonnet-4-5"]);
expect(capturedOptions?.modelPatternAuthFallback).toBe("anthropic/claude-sonnet-4-5");
});
it("pins the persisted concrete model when the default role is revived", async () => {
const cwd = makeTempDir("@pi-default-role-revive-");
const sessionFile = await createPersistedSession(cwd, false, "default");
let capturedOptions: CreateAgentSessionOptions | undefined;
vi.spyOn(sdkModule, "createAgentSession").mockImplementation(async options => {
capturedOptions = options;
return { session: createRevivedSession([]).session } as CreateAgentSessionResult;
});
const ref = createRef(sessionFile);
const reviver = await createFactory(cwd)(ref);
if (!reviver) throw new Error("Expected a persisted reviver");
await reviver(ref);
expect(capturedOptions?.modelPattern).toBe("anthropic/claude-sonnet-4-5");
expect(capturedOptions?.modelPatternAuthFallback).toBe("anthropic/claude-sonnet-4-5");
});
it("installs an IRC wake monitor that emits cold-revive lifecycle frames on the shared bus", async () => {
AgentRegistry.resetGlobalForTests();
AgentLifecycleManager.resetGlobalForTests();
@@ -1952,6 +1952,41 @@ describe("vibe session registry", () => {
}
});
it("finalizer persists a hard-abort tombstone before disposing the live ref", async () => {
const parentManager = await createPersistedParent();
await parentManager.ensureOnDisk();
const parentSessionFile = parentManager.getSessionFile();
if (!parentSessionFile) throw new Error("Expected a persisted parent session file");
const workerId = "finalizer-hard-abort";
const workerSessionFile = path.join(parentSessionFile.slice(0, -6), `${workerId}.jsonl`);
await Bun.write(workerSessionFile, "");
const worker = createFakeWorkerSession();
const ref = AgentRegistry.global().register({
id: workerId,
displayName: workerId,
kind: "sub",
parentId: "Main",
session: worker.session,
sessionFile: workerSessionFile,
status: "running",
});
await executorModule.finalizeSubagentLifecycle({
id: workerId,
session: worker.session,
aborted: true,
keepAlive: true,
isolated: false,
agentIdleTtlMs: 0,
reviveSession: null,
});
expect(worker.isDisposed()).toBe(true);
expect(AgentRegistry.global().get(workerId)).toMatchObject({ status: "aborted", session: null });
expect(AgentRegistry.global().get(workerId)).toBe(ref);
expect(await fileExists(`${workerSessionFile}.tombstone`)).toBe(true);
});
it("keeps a persisted in-flight kill terminal when the old executor finalizes late", async () => {
let worker: ReturnType<typeof createFakeWorkerSession> | undefined;
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {