Merge pull request #7551 from kmccleary3301/refactor/agent-hub-fullscreen

refactor(coding-agent): polish fullscreen Agent Hub
This commit is contained in:
Can Bölük
2026-08-05 09:20:39 +02:00
committed by GitHub
34 changed files with 3296 additions and 455 deletions
+3
View File
@@ -14,6 +14,7 @@
### Changed
- Reworked the Ctrl+S Agent Hub into a responsive fullscreen roster and selected-agent inspector with aggregate status/usage, per-agent task/model/activity/usage/lineage details, roster and spawn-tree views, stable ordering, bounded large-roster rendering, asynchronous persisted-session discovery, restored task/timestamp metadata for historical agents, and consistent keyboard and mouse navigation.
- Restored the legacy project-scoped session directory naming scheme and removed its automatic migration ([#7646](https://github.com/can1357/oh-my-pi/issues/7646)).
- Routed Bun install-cache pruning in `update-cli` through the shared `compareVersions` utility (`@oh-my-pi/pi-utils`), removing a duplicate local comparator that rounded large numeric version identifiers via `Number`.
@@ -41,6 +42,7 @@
- Fixed Herdr rejecting the macOS development launcher because its foreground process was reported as `bun` instead of `omp`.
- Completed usage-aware model fallback across startup, queued turns, same-turn tool continuations, ACP/TUI confirmation cancellation, eligible account reselection, cooldown restoration, and isolated subagent settings so low-usage handoffs remain lossless and cannot consume cancelled queued work.
- Fixed Agent Hub opening and selection becoming O(all rows) on large rosters: row rendering is now lazy around the selected viewport, and observer lookup is O(1) by id instead of copy-sorting every session per row.
- Fixed persisted Agent Hub rows dropping an explicit caller model role when a subagent used a model override, preserving role provenance after restart.
- Fixed the bash interceptor blocking `grep`/`cat`/`find` used as a downstream pipeline stage (e.g. `printf 'x\n' | grep x`); a stage consuming piped stdin cannot be replaced by a path-based dedicated tool, so it is no longer matched, while standalone and first-stage searches stay intercepted ([#7496](https://github.com/can1357/oh-my-pi/issues/7496)).
- Fixed floating rejections from cmux browser guest JavaScript terminating the main process and every active session; attributable rejections now fail the browser run as tool errors while unrelated process rejections retain the fatal path ([#7365](https://github.com/can1357/oh-my-pi/issues/7365)).
- Fixed the Windows bash tool silently taking down the whole omp process when a command blocked until its timeout: cancelling a timed-out run walked the spawned child's descendant tree from raw `th32ParentProcessID` links, and a recycled pid matching the harness's stale recorded parent pid could enumerate omp as a false descendant and `TerminateProcess` it, killing the session with no `session_exit` record. Run-cancellation sweeps now refuse to signal the harness or any process collected beneath it, while still reaping the timed-out target when it owns a recycled ancestor pid ([#7452](https://github.com/can1357/oh-my-pi/issues/7452)).
@@ -63,6 +65,7 @@
### Changed
- Replaced arktype with @oh-my-pi/omptype for tool parameter and config schemas, significantly improving startup performance with ~100x faster schema construction. Config schema errors are now reported via OmpErrors using the same path/problem structure.
- Replaced arktype with `@oh-my-pi/omptype` across all tool parameter and config schemas: ~100x faster schema construction removes the arktype startup tax (the `scope({}, { jitless: true })` workarounds are gone). Config schema errors now report via `OmpErrors` entries with the same `path`/`problem` shape.
### Fixed
@@ -932,6 +932,28 @@ function normalizeModelPatternList(value: string | string[] | undefined): string
return patterns.map(pattern => pattern.trim()).filter(Boolean);
}
/**
* Extract the first explicit model-role alias from a raw model selection.
*
* This intentionally runs before role expansion so callers can retain the
* source identity (`@smol`, `pi/slow`, or `*`) even when it resolves to a
* concrete provider/model or inherited fallback. Bare role names and explicit
* provider/model selectors are not role aliases.
*/
export function resolveExplicitModelRole(
value: string | string[] | undefined,
settings?: ModelRoleLookup,
): string | undefined {
for (const pattern of normalizeModelPatternList(value)) {
const prefixLength = modelRoleAliasPrefixLength(pattern);
if (prefixLength === undefined) continue;
const { base } = splitThinkingSuffix(pattern, prefixLength, MAX_THINKING_SUFFIX_OPTIONS);
const role = getModelRoleAlias(base, settings);
if (role) return role;
}
return undefined;
}
function isSessionInheritedAgentPattern(value: string): boolean {
return (
value === DEFAULT_MODEL_ROLE ||
@@ -1068,6 +1090,8 @@ export function resolveConfiguredModelPatterns(
});
}
export interface AgentModelPatternResolutionOptions {
/** Highest-priority request selector, when supplied by a caller. */
requestModel?: string | string[];
settingsOverride?: string | string[];
agentModel?: string | string[];
settings?: Settings;
@@ -1075,11 +1099,25 @@ export interface AgentModelPatternResolutionOptions {
fallbackModelPattern?: string;
}
export function resolveAgentModelPatterns(options: AgentModelPatternResolutionOptions): string[] {
const { settingsOverride, agentModel, settings, activeModelPattern, fallbackModelPattern } = options;
interface EffectiveAgentModelSelection {
source?: string | string[];
patterns: string[];
}
function resolveEffectiveAgentModelSelection(
options: AgentModelPatternResolutionOptions,
): EffectiveAgentModelSelection {
const { requestModel, settingsOverride, agentModel, settings, activeModelPattern, fallbackModelPattern } = options;
const requestPatterns = resolveConfiguredModelPatterns(requestModel, settings);
if (requestPatterns.length > 0) {
return { source: requestModel, patterns: requestPatterns };
}
const overridePatterns = resolveConfiguredModelPatterns(settingsOverride, settings);
if (overridePatterns.length > 0) return overridePatterns;
if (overridePatterns.length > 0) {
return { source: settingsOverride, patterns: overridePatterns };
}
const normalizedAgentPatterns = normalizeModelPatternList(agentModel);
const configuredAgentPatterns = resolveConfiguredModelPatterns(agentModel, settings);
@@ -1090,14 +1128,23 @@ export function resolveAgentModelPatterns(options: AgentModelPatternResolutionOp
singleAgentPattern === formatModelRoleAlias("task") ||
singleAgentPattern === `${LEGACY_MODEL_ROLE_ALIAS_PREFIX}task`
) {
return configuredAgentPatterns;
return { source: agentModel, patterns: configuredAgentPatterns };
}
if (!agentInheritsSessionModel) return configuredAgentPatterns;
if (!agentInheritsSessionModel) return { source: agentModel, patterns: configuredAgentPatterns };
}
const fallback =
activeModelPattern?.trim() || fallbackModelPattern?.trim() || settings?.getModelRole("default")?.trim() || "";
return resolveConfiguredModelPatterns(fallback, settings);
return { patterns: resolveConfiguredModelPatterns(fallback, settings) };
}
/** Return the raw selector source that supplies the effective agent patterns. */
export function resolveAgentModelSource(options: AgentModelPatternResolutionOptions): string | string[] | undefined {
return resolveEffectiveAgentModelSelection(options).source;
}
export function resolveAgentModelPatterns(options: AgentModelPatternResolutionOptions): string[] {
return resolveEffectiveAgentModelSelection(options).patterns;
}
/** Default prewalk hand-off target when no explicit target is configured. */
export const DEFAULT_PREWALK_TARGET = "@smol";
@@ -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}`);
}
File diff suppressed because it is too large Load Diff
@@ -110,6 +110,23 @@ const MANUAL_LOGIN_PROMPT = "Paste the authorization code (or full redirect URL)
export class SelectorController {
constructor(private ctx: InteractiveModeContext) {}
/**
* Mount a primary fullscreen menu through the one polished modal path shared
* by Settings, Model Hub, and Agent Hub.
*/
#showFullscreenMenu(component: Component): OverlayHandle {
const handle = this.ctx.ui.showOverlay(component, {
anchor: "bottom-center",
width: "100%",
maxHeight: "100%",
margin: 0,
fullscreen: true,
});
this.ctx.ui.setFocus(component);
this.ctx.ui.requestRender();
return handle;
}
#defaultRoleMutationTail = Promise.resolve();
async #acquireDefaultRoleMutation(): Promise<() => void> {
@@ -242,15 +259,7 @@ export class SelectorController {
},
},
);
overlayHandle = this.ctx.ui.showOverlay(selector, {
anchor: "bottom-center",
width: "100%",
maxHeight: "100%",
margin: 0,
fullscreen: true,
});
this.ctx.ui.setFocus(selector);
this.ctx.ui.requestRender();
overlayHandle = this.#showFullscreenMenu(selector);
});
}
@@ -1023,15 +1032,7 @@ export class SelectorController {
initialProviderId: hubOptions.initialProviderId,
},
);
overlayHandle = this.ctx.ui.showOverlay(hub, {
anchor: "bottom-center",
width: "100%",
maxHeight: "100%",
margin: 0,
fullscreen: true,
});
this.ctx.ui.setFocus(hub);
this.ctx.ui.requestRender();
overlayHandle = this.#showFullscreenMenu(hub);
}
/** /login round-trip for a locked provider; reopen the hub on that provider only after a successful login. */
@@ -2000,27 +2001,23 @@ export class SelectorController {
...this.ctx.keybindings.getKeys("app.agents.hub"),
...this.ctx.keybindings.getKeys("app.session.observe"),
];
let hub: AgentHubOverlayComponent | undefined;
let overlayHandle: OverlayHandle | undefined;
let closed = false;
// Render the hub inline in the editor slot — the same anchored region
// every other selector (model, session, tree, the `ask` tool) uses —
// rather than a floating overlay. A non-fullscreen overlay composited over
// a live transcript strands a stale copy in native scrollback every time a
// running subagent's progress grows the frame and scrolls the window; the
// hub is opened mid-run, so those copies stacked into a wall of duplicate
// "Agent Hub" frames bleeding the task tree behind them. As an editor-slot
// component it rides the normal append-only commit path: the transcript
// commits above it exactly once and the hub repaints in place.
const done = () => {
hub?.dispose();
this.ctx.editorContainer.clear();
this.ctx.editorContainer.addChild(this.ctx.editor);
this.ctx.ui.setFocus(this.ctx.editor);
if (closed) return;
closed = true;
hub.dispose();
overlayHandle?.hide();
// 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();
};
hub = new AgentHubOverlayComponent({
const hub = new AgentHubOverlayComponent({
observers,
settings: this.ctx.settings,
hubKeys,
expandKeys: this.ctx.keybindings.getKeys("app.tools.expand"),
onDone: done,
@@ -2038,30 +2035,24 @@ export class SelectorController {
});
const showReadyHub = () => {
// The double-← gesture passes requireContent so it stays inert when
// neither live nor persisted subagents are available. Persisted rows now
// load asynchronously, so defer the gate until that scan has refreshed the
// hub instead of treating the initial empty table as authoritative.
if (closed) return;
// The double-← gesture stays inert when neither live nor persisted
// subagents are available, so wait for discovery before making the gate.
if (options?.requireContent && hub.isEmpty) {
hub.dispose();
done();
return;
}
this.ctx.editorContainer.clear();
this.ctx.editorContainer.addChild(hub);
this.ctx.ui.setFocus(hub);
// When the hub was raised by the editor's double-← gesture, prime its own
// close detector so the *next* single ← dismisses it — the two taps that
// opened it were consumed by the editor's detector (issue #4780).
// Prime the detector before the first frame when the editor's double-←
// gesture opened the hub, so the next single ← dismisses it.
if (options?.armCloseTap) hub.armCloseTap();
this.ctx.ui.requestRender();
overlayHandle = this.#showFullscreenMenu(hub);
};
if (options?.requireContent && hub.isEmpty) {
void hub.persistedSubagentsReady.then(showReadyHub);
return;
} else {
showReadyHub();
}
showReadyHub();
}
}
@@ -20,6 +20,7 @@
* a superseded revive) can never clobber a newer same-id ref.
*/
import * as fs from "node:fs/promises";
import { logger, untilAborted } from "@oh-my-pi/pi-utils";
import type { AgentSession } from "../session/agent-session";
import { trackLateCleanup } from "../utils/late-cleanup";
@@ -27,6 +28,7 @@ import {
type AgentRef,
type AgentRefExpectation,
AgentRegistry,
getAgentTombstonePath,
MAIN_AGENT_ID,
type RegistryEvent,
} from "./agent-registry";
@@ -35,6 +37,14 @@ export type AgentReviver = (expected: AgentRef) => Promise<AgentSession>;
const AGENT_RELEASE_GRACE_MS = 5000;
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
@@ -382,13 +392,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
@@ -30,6 +39,35 @@ export type AgentStatus = "running" | "idle" | "parked" | "aborted";
*/
export type AgentKind = "main" | "sub" | "advisor";
/** Persisted per-agent totals reconstructed from the child session transcript. */
export interface AgentMetricsSummary {
tokens: number;
requests: number;
tools: number;
cost: number;
durationMs: number;
durationKind?: AgentDurationKind;
contextTokens?: number;
contextWindow?: number;
}
/** Historical identity and telemetry that remain available after the live session is disposed. */
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;
/** Durable task output artifact, when the executor wrote one. */
outputPath?: string;
/** Captured isolated-worktree patch, when patch capture succeeded. */
patchPath?: string;
/** Isolated branch identity, when branch-mode capture succeeded. */
branchName?: string;
}
export interface AgentRef {
id: string;
displayName: string;
@@ -43,6 +81,8 @@ export interface AgentRef {
lastActivity: number;
/** Short gist of what the agent is currently doing (latest intent or tool), for the work-aware roster. Display-only. */
activity?: string;
/** Persisted identity and telemetry restored after the live observer is gone. */
history?: AgentHistorySummary;
}
export type AgentRefExpectation = AgentRef | AgentSession;
@@ -50,6 +90,7 @@ export type AgentRefExpectation = AgentRef | AgentSession;
export type RegistryEvent =
| { type: "registered"; ref: AgentRef }
| { type: "status_changed"; ref: AgentRef }
| { type: "metadata_changed"; ref: AgentRef }
| { type: "removed"; ref: AgentRef };
type RegistryListener = (event: RegistryEvent) => void;
@@ -62,6 +103,14 @@ export interface RegisterInput {
session: AgentSession | null;
sessionFile?: string | null;
status?: AgentStatus;
/** Last persisted task summary, when restoring a historical agent. */
activity?: string;
/** Original registration timestamp, when known from persisted history. */
createdAt?: number;
/** Last transcript activity timestamp, when known from persisted history. */
lastActivity?: number;
/** Persisted identity and telemetry restored after the live observer is gone. */
history?: AgentHistorySummary;
}
export class AgentRegistry {
@@ -96,8 +145,10 @@ export class AgentRegistry {
status: input.status ?? "running",
session: input.session,
sessionFile: input.sessionFile ?? null,
createdAt: now,
lastActivity: now,
createdAt: input.createdAt ?? now,
lastActivity: input.lastActivity ?? now,
activity: input.activity,
history: input.history,
};
this.#refs.set(ref.id, ref);
this.#emit({ type: "registered", ref });
@@ -116,6 +167,18 @@ export class AgentRegistry {
return current === expected && current.status === "parked" && !current.session ? current : undefined;
}
/** Attach transcript-derived identity and telemetry without changing lifecycle state. */
setHistory(id: string, history: AgentHistorySummary, expectedSessionFile?: string): boolean {
const ref = this.#refs.get(id);
if (!ref || (expectedSessionFile !== undefined && ref.sessionFile !== expectedSessionFile)) return false;
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;
}
setStatus(id: string, status: AgentStatus, expected?: AgentRefExpectation): boolean {
const ref = this.#refs.get(id);
if (!ref || !this.#matchesExpected(ref, expected)) return false;
@@ -1,38 +1,312 @@
import * as fs from "node:fs";
import * as path from "node:path";
import { ADVISOR_TRANSCRIPT_FILENAME, isAdvisorTranscriptName } from "../advisor/transcript-recorder";
import { SessionManager } from "../session/session-manager";
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 AgentRegistry, MAIN_AGENT_ID } from "./agent-registry";
import {
type AgentHistorySummary,
type AgentMetricsSummary,
type AgentRegistry,
getAgentTombstonePath,
MAIN_AGENT_ID,
} from "./agent-registry";
/** Maximum prefix entries inspected for task metadata. */
const MAX_METADATA_LINES = 64;
interface PersistedAgentMetadata {
activity?: string;
createdAt?: number;
lastActivity?: number;
history?: AgentHistorySummary;
}
interface PersistedTranscript {
id: string;
sessionFile: string;
createdAt?: number;
lastActivity?: number;
}
function recordOf(value: unknown): Record<string, unknown> | undefined {
return typeof value === "object" && value !== null && !Array.isArray(value)
? (value as Record<string, unknown>)
: undefined;
}
function timestampOf(value: unknown): number | undefined {
if (typeof value !== "string") return undefined;
const timestamp = Date.parse(value);
return Number.isFinite(timestamp) ? timestamp : undefined;
}
function summarizePersistedTask(task: string): string | undefined {
const withoutPreamble = task.replace(/^Complete the assignment below,\s*thoroughly:\s*/i, "");
const lines = withoutPreamble.split(/\r?\n/);
const targetIndex = lines.findIndex(line => line.trim().toLowerCase() === "# target");
const targetLines: string[] = [];
if (targetIndex >= 0) {
for (const line of lines.slice(targetIndex + 1)) {
if (line.trimStart().startsWith("# ")) break;
targetLines.push(line);
}
}
const summary = (targetLines.length > 0 ? targetLines : lines).join(" ").replace(/\s+/g, " ").trim();
return summary ? summary.slice(0, 1_000) : undefined;
}
function finiteNumber(value: unknown): number {
return typeof value === "number" && Number.isFinite(value) ? value : 0;
}
function inferBundledAgent(systemPrompt: string): { agent?: string; modelRole?: string; readOnly?: boolean } {
const matches = loadBundledAgents().filter(agent => {
const rolePrompt = agent.systemPrompt.trim();
return rolePrompt.length > 0 && systemPrompt.includes(rolePrompt);
});
// `task` and `sonic` intentionally share a prompt body. Ambiguous historical
// prompts stay unlabelled rather than inventing provenance.
if (matches.length !== 1) return {};
const [agent] = matches;
return {
agent: agent.name,
modelRole: resolveExplicitModelRole(agent.model),
readOnly: isReadOnlyAgent(agent),
};
}
function usageTokens(usage: Record<string, unknown>): number {
const computed = finiteNumber(usage.input) + finiteNumber(usage.output) + finiteNumber(usage.cacheWrite);
return computed > 0 ? computed : finiteNumber(usage.totalTokens);
}
interface AssistantMetrics {
tokens: number;
tools: number;
cost: number;
contextTokens?: number;
resolvedModel?: string;
}
function assistantMetrics(message: Record<string, unknown>): AssistantMetrics {
const usage = recordOf(message.usage) ?? {};
const cost = recordOf(usage.cost);
const content = Array.isArray(message.content) ? message.content : [];
const provider = typeof message.provider === "string" ? message.provider : undefined;
const model = typeof message.model === "string" ? message.model : undefined;
return {
tokens: usageTokens(usage),
tools: content.filter(part => recordOf(part)?.type === "toolCall").length,
cost: finiteNumber(cost?.total),
contextTokens: finiteNumber(usage.totalTokens) || undefined,
resolvedModel: provider && model ? `${provider}/${model}` : undefined,
};
}
async function readPersistedAgentHistory(
transcript: PersistedTranscript,
shouldContinue: () => boolean,
): 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 {
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 parsedTimestamp = timestampOf(record.timestamp);
if (parsedTimestamp !== undefined) leafTimestamp = parsedTimestamp;
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));
},
{ shouldContinue },
);
} catch {
return {};
}
const metrics: AgentMetricsSummary = {
tokens: 0,
requests: 0,
tools: 0,
cost: 0,
durationMs: Math.max(
0,
(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;
}
if (contextTokens !== undefined) metrics.contextTokens = contextTokens;
return {
...(metrics.requests > 0 ? { metrics } : {}),
...(resolvedModel ? { resolvedModel, resolvedModelIsFallback } : {}),
...(modelRole ? { modelRole } : {}),
};
}
/**
* 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.
* Read only the small session prefix needed by the Hub. A subagent's first
* `session_init` is written before its conversation, so this never walks a
* multi-megabyte historical transcript just to populate one roster row.
*/
async function readPersistedVibeChildIds(sessionFile: string): Promise<Set<string>> {
let sessionManager: SessionManager;
async function readPersistedAgentMetadata(sessionFile: string): Promise<PersistedAgentMetadata> {
const stat = fs.promises.stat(sessionFile).catch(() => undefined);
const artifactBase = sessionFile.slice(0, -".jsonl".length);
const outputPath = `${artifactBase}.md`;
const patchPath = `${artifactBase}.patch`;
const artifactFiles = Promise.all([Bun.file(outputPath).exists(), Bun.file(patchPath).exists()]);
let createdAt: number | undefined;
let activity: string | undefined;
let history: AgentHistorySummary = {};
try {
sessionManager = await SessionManager.open(sessionFile, undefined, undefined, { suppressBreadcrumb: true });
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.
}
const [file, [hasOutput, hasPatch]] = await Promise.all([stat, artifactFiles]);
return {
activity,
createdAt: createdAt ?? file?.birthtimeMs,
lastActivity: file?.mtimeMs,
history: {
...history,
...(hasOutput ? { outputPath } : {}),
...(hasPatch ? { patchPath } : {}),
},
};
}
async function readPersistedVibeChildIds(sessionFile: string, shouldContinue: () => boolean): Promise<Set<string>> {
const ids = new Set<string>();
try {
await visitEntriesFromFileStream(
sessionFile,
entry => {
for (const id of persistedVibeChildIds([entry])) ids.add(id);
},
{ shouldContinue },
);
return ids;
} catch {
return new Set();
}
try {
return persistedVibeChildIds(sessionManager.getEntries());
} finally {
await sessionManager.close();
}
}
/** Register persisted subagent and advisor transcripts as parked registry refs. */
export async function registerPersistedSubagents(
registry: AgentRegistry,
sessionFile: string | null | undefined,
options: { shouldContinue?: () => boolean } = {},
): Promise<void> {
if (!sessionFile?.endsWith(".jsonl")) return;
const vibeOwnedIds = await readPersistedVibeChildIds(sessionFile);
const shouldContinue = options.shouldContinue ?? (() => true);
if (!shouldContinue()) return;
const vibeOwnedIds = await readPersistedVibeChildIds(sessionFile, shouldContinue);
if (!shouldContinue()) return;
const root = sessionFile.slice(0, -6);
await registerPersistedSubagentsFromDir(registry, root, undefined, vibeOwnedIds);
const transcripts: PersistedTranscript[] = [];
await registerPersistedSubagentsFromDir(registry, root, undefined, vibeOwnedIds, transcripts, shouldContinue);
if (!shouldContinue()) return;
let nextTranscript = 0;
const workers = Array.from({ length: Math.min(4, transcripts.length) }, async () => {
for (;;) {
if (!shouldContinue()) return;
const index = nextTranscript++;
const transcript = transcripts[index];
if (!transcript) return;
const history = await readPersistedAgentHistory(transcript, shouldContinue);
if (!shouldContinue()) return;
registry.setHistory(transcript.id, history, transcript.sessionFile);
}
});
await Promise.all(workers);
}
async function registerPersistedSubagentsFromDir(
@@ -40,14 +314,25 @@ async function registerPersistedSubagentsFromDir(
dir: string,
parentId: string | undefined,
vibeOwnedIds: ReadonlySet<string>,
transcripts: PersistedTranscript[],
shouldContinue: () => boolean,
): Promise<void> {
if (!shouldContinue()) return;
let entries: fs.Dirent[];
try {
entries = await fs.promises.readdir(dir, { withFileTypes: true });
} catch {
return;
}
if (!shouldContinue()) return;
let entriesSinceYield = 0;
for (const entry of entries) {
if (!shouldContinue()) return;
if (++entriesSinceYield >= 16) {
entriesSinceYield = 0;
await Bun.sleep(0);
}
if (!shouldContinue()) return;
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
@@ -66,6 +351,8 @@ async function registerPersistedSubagentsFromDir(
// user task literally named `<owner>/advisor`): leave it, skip the advisor.
if (existing && existing.kind !== "advisor") continue;
if (existing?.sessionFile !== sessionFile) {
const metadata = await readPersistedAgentMetadata(sessionFile);
if (!shouldContinue()) return;
// The id is reused across `/new`; refresh it to the current session's file.
if (existing) registry.unregister(advisorId);
registry.register({
@@ -75,14 +362,34 @@ async function registerPersistedSubagentsFromDir(
parentId: owner,
session: null,
sessionFile,
activity: metadata.activity,
createdAt: metadata.createdAt,
lastActivity: metadata.lastActivity,
history: { ...metadata.history, readOnly: true },
status: "parked",
});
transcripts.push({
id: advisorId,
sessionFile,
createdAt: metadata.createdAt,
lastActivity: metadata.lastActivity,
});
}
continue;
}
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 (!shouldContinue()) return;
if (!registry.get(id)) {
const metadata = await readPersistedAgentMetadata(sessionFile);
if (!shouldContinue()) return;
registry.register({
id,
displayName: id,
@@ -90,9 +397,27 @@ async function registerPersistedSubagentsFromDir(
parentId: parentId ?? MAIN_AGENT_ID,
session: null,
sessionFile,
status: "parked",
activity: metadata.activity,
createdAt: metadata.createdAt,
lastActivity: metadata.lastActivity,
history: metadata.history,
status: tombstoned ? "aborted" : "parked",
});
const ref = registry.get(id);
transcripts.push({
id,
sessionFile,
createdAt: ref?.createdAt,
lastActivity: ref?.lastActivity,
});
}
await registerPersistedSubagentsFromDir(registry, path.join(dir, id), id, vibeOwnedIds);
await registerPersistedSubagentsFromDir(
registry,
path.join(dir, id),
id,
vibeOwnedIds,
transcripts,
shouldContinue,
);
}
}
@@ -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 {
@@ -207,6 +209,14 @@ export interface SessionInitEntry extends SessionEntryBase {
task: string;
/** Tools available to the agent */
tools: string[];
/** Agent definition name (for example `scout` or `reviewer`). */
agent?: string;
/** Semantic model role declared by the agent, retained even after concrete model resolution. */
modelRole?: string;
/** Initially resolved provider/model selector for historical display. */
resolvedModel?: string;
/** Whether the agent definition is read-only, allowing an exact zero-LoC attribution. */
readOnly?: boolean;
/** Output schema if structured output was requested. */
outputSchema?: unknown;
/** Enforcement policy recorded with the output schema for faithful revival. */
@@ -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;
}
@@ -2058,6 +2065,10 @@ export class SessionManager {
systemPrompt: string;
task: string;
tools: string[];
agent?: string;
modelRole?: string;
resolvedModel?: string;
readOnly?: boolean;
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
restrictToolNames?: boolean;
@@ -2547,6 +2558,9 @@ export class SessionManager {
systemPrompt: string;
task: string;
tools: string[];
agent?: string;
modelRole?: string;
resolvedModel?: string;
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
restrictToolNames?: boolean;
@@ -2567,6 +2581,9 @@ export class SessionManager {
systemPrompt: string;
task: string;
tools: string[];
agent?: string;
modelRole?: string;
resolvedModel?: string;
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
restrictToolNames?: boolean;
@@ -2580,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,
@@ -1302,7 +1302,7 @@ export class TurnRecovery {
return false;
}
if (this.#host.model() !== candidate) return false;
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) {
@@ -1422,7 +1422,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",
+44 -9
View File
@@ -17,6 +17,7 @@ import {
formatModelStringWithRouting,
resolveAgentPrewalkPattern,
resolveConfiguredModelPatterns,
resolveExplicitModelRole,
resolveModelOverride,
resolveModelOverrideWithAuthFallback,
} from "../config/model-resolver";
@@ -59,6 +60,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,
@@ -344,6 +346,8 @@ export interface ExecutorOptions {
*/
detached?: boolean;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
/**
* Active model selector of the parent session, used as an auth-aware fallback
* if the resolved subagent model has no working credentials. See #985.
@@ -890,6 +894,8 @@ interface RunMonitorArgs {
/** Parent settings for tiny-model label generation. */
settings?: Settings;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
signal?: AbortSignal;
onProgress?: (progress: AgentProgress) => void;
eventBus?: EventBus;
@@ -1005,6 +1011,7 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
cost: 0,
durationMs: 0,
modelOverride: args.modelOverride,
modelRole: args.modelRole,
};
const outputChunks: string[] = [];
@@ -2055,6 +2062,8 @@ interface FinalizeRunArgs {
task: string;
assignment?: string;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
outputSchemaSource?: StructuredSubagentSchemaSource;
@@ -2074,7 +2083,7 @@ interface FinalizeRunArgs {
* event.
*/
async function finalizeRunResult(args: FinalizeRunArgs): Promise<SingleResult> {
const { monitor, done, index, id, agent, task, assignment, signal, modelOverride } = args;
const { monitor, done, index, id, agent, task, assignment, signal, modelOverride, modelRole } = args;
const progress = monitor.progress;
let exitCode = done.exitCode;
let stderr = done.error ?? "";
@@ -2201,6 +2210,7 @@ async function finalizeRunResult(args: FinalizeRunArgs): Promise<SingleResult> {
contextTokens: progress.contextTokens,
contextWindow: progress.contextWindow,
modelOverride,
modelRole,
resolvedModel: progress.resolvedModel,
resolvedModelIsFallback: progress.resolvedModelIsFallback,
error: exitCode !== 0 && stderr ? stderr : undefined,
@@ -2222,6 +2232,8 @@ export interface IrcWakeTurnMonitorOptions {
agent: AgentDefinition;
description?: string;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
eventBus?: EventBus;
parentToolCallId?: string;
/** Fallback session file when the registry ref carries none. */
@@ -2266,6 +2278,7 @@ export function attachIrcWakeTurnMonitor(session: AgentSession, options: IrcWake
task: ircTask,
description: options.description,
modelOverride: options.modelOverride,
modelRole: options.modelRole,
eventBus: options.eventBus,
parentToolCallId: options.parentToolCallId,
detached: true,
@@ -2323,6 +2336,7 @@ export function attachIrcWakeTurnMonitor(session: AgentSession, options: IrcWake
agent,
task: ircTask,
modelOverride: options.modelOverride,
modelRole: options.modelRole,
outputSchema: options.outputSchema,
outputSchemaMode: options.outputSchemaMode,
outputSchemaSource: options.outputSchemaSource,
@@ -2391,14 +2405,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;
}
@@ -2447,6 +2467,8 @@ export interface FollowUpTurnOptions {
message: string;
index?: number;
description?: string;
/** Explicit pre-expansion model role alias retained from the original run. */
modelRole?: string;
/** Structured-output state retained from the original invocation. */
outputSchema?: unknown;
outputSchemaMode?: StructuredSubagentSchemaMode;
@@ -2486,6 +2508,7 @@ export async function runSubagentFollowUpTurn(options: FollowUpTurnOptions): Pro
agent,
task: message,
description: options.description,
modelRole: options.modelRole,
signal,
onProgress: options.onProgress,
eventBus: options.eventBus,
@@ -2535,6 +2558,7 @@ export async function runSubagentFollowUpTurn(options: FollowUpTurnOptions): Pro
id,
agent,
task: message,
modelRole: options.modelRole,
outputSchema: options.outputSchema,
outputSchemaMode: options.outputSchemaMode,
outputSchemaSource: options.outputSchemaSource,
@@ -2561,6 +2585,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
id,
worktree,
modelOverride,
modelRole,
thinkingLevel,
outputSchema,
enableLsp,
@@ -2590,6 +2615,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
tokens: 0,
requests: 0,
modelOverride,
modelRole,
error: "Cancelled before start",
aborted: true,
abortReason: "Cancelled before start",
@@ -2680,6 +2706,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
modelRegistry: options.modelRegistry,
settings,
modelOverride,
modelRole,
signal,
onProgress,
eventBus: options.eventBus,
@@ -2713,6 +2740,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
agent,
description: options.description,
modelOverride,
modelRole,
eventBus: options.eventBus,
parentToolCallId: options.parentToolCallId,
sessionFile: subtaskSessionFile,
@@ -3098,6 +3126,10 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
systemPrompt: session.agent.state.systemPrompt.join("\n\n"),
task,
tools: session.getActiveToolNames(),
agent: agent.name,
modelRole: modelRole ?? resolveExplicitModelRole(modelOverride ?? agent.model, subagentSettings),
resolvedModel: progress.resolvedModel,
readOnly: isReadOnlyAgent(agent),
spawns: spawnsEnv,
readSummarize: agent.readSummarize,
outputSchema,
@@ -3358,7 +3390,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
const done = await runSubagent();
monitor.finish();
return finalizeRunResult({
const result = await finalizeRunResult({
monitor,
done,
index,
@@ -3367,6 +3399,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
task,
assignment,
modelOverride,
modelRole,
outputSchema,
outputSchemaMode: options.outputSchemaMode,
outputSchemaSource: options.outputSchemaSource,
@@ -3378,4 +3411,6 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
sessionFile: subtaskSessionFile,
startTime,
});
AgentRegistry.global().setHistory(id, { outputPath: result.outputPath });
return result;
}
+9 -28
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
@@ -860,6 +836,7 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
id: agentId,
agent: agentType,
agentSource,
modelRole: policy.modelRole,
status: "pending",
task: renderSubagentUserPrompt(assignment),
assignment,
@@ -1143,8 +1120,10 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
// polling row reflects the resolved model, reasoning level,
// and running counters without reverting the "running"
// 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 =
nextProgress.resolvedModelIsFallback ?? progress.resolvedModelIsFallback;
progress.tokens = nextProgress.tokens;
progress.requests = nextProgress.requests;
progress.contextTokens = nextProgress.contextTokens;
@@ -1187,9 +1166,11 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
progress.extractedToolData = singleResult?.extractedToolData;
progress.retryFailure = singleResult?.retryFailure;
progress.retryState = undefined;
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;
@@ -20,6 +20,7 @@
*/
import * as path from "node:path";
import type * as natives from "@oh-my-pi/pi-natives";
import { AgentRegistry } from "../registry/agent-registry";
import type { ToolSession } from "../tools";
import { generateCommitMessage } from "../utils/commit-message-generator";
import * as git from "../utils/git";
@@ -44,6 +45,15 @@ import {
type IsoBackendKind = natives.IsoBackendKind;
function rememberAgentArtifacts(result: SingleResult): SingleResult {
AgentRegistry.global().setHistory(result.id, {
outputPath: result.outputPath,
patchPath: result.patchPath,
branchName: result.branchName,
});
return result;
}
/** Resolved repo + baseline used by every isolated spawn in a single call. */
export interface IsolationContext {
repoRoot: string;
@@ -172,14 +182,14 @@ export async function runIsolatedSubprocess(opts: IsolatedRunOptions): Promise<S
opts.description,
opts.buildCommitMessage?.(),
);
return {
return rememberAgentArtifacts({
...result,
branchName: commitResult?.branchName,
branchBaseSha: commitResult?.baseSha,
nestedPatches: commitResult?.nestedPatches,
};
});
} catch (mergeErr) {
// Agent succeeded but branch commit failed — clean up stale branch
// Agent succeeded but branch commit failed — clean up stale branch.
const branchName = `omp/task/${opts.agentId}`;
await git.branch.tryDelete(opts.context.repoRoot, branchName);
const msg = mergeErr instanceof Error ? mergeErr.message : String(mergeErr);
@@ -190,34 +200,37 @@ export async function runIsolatedSubprocess(opts: IsolatedRunOptions): Promise<S
opts.artifactsDir,
opts.agentId,
);
return {
return rememberAgentArtifacts({
...result,
patchPath: patchResult.patchPath,
nestedPatches: patchResult.nestedPatches,
error: `Merge failed: ${msg}`,
};
});
} catch (patchErr) {
const patchMsg = patchErr instanceof Error ? patchErr.message : String(patchErr);
return { ...result, error: `Merge failed: ${msg}; patch capture failed: ${patchMsg}` };
return rememberAgentArtifacts({
...result,
error: `Merge failed: ${msg}; patch capture failed: ${patchMsg}`,
});
}
}
}
if (result.exitCode === 0) {
try {
const patchResult = await writeIsolationPatch(isolationDir, taskBaseline, opts.artifactsDir, opts.agentId);
return {
return rememberAgentArtifacts({
...result,
patchPath: patchResult.patchPath,
nestedPatches: patchResult.nestedPatches,
};
});
} catch (patchErr) {
const msg = patchErr instanceof Error ? patchErr.message : String(patchErr);
return { ...result, error: `Patch capture failed: ${msg}` };
return rememberAgentArtifacts({ ...result, error: `Patch capture failed: ${msg}` });
}
}
return result;
return rememberAgentArtifacts(result);
} catch (err) {
return opts.buildFailureResult(err);
return rememberAgentArtifacts(opts.buildFailureResult(err));
} finally {
if (handle) {
const isolationHandle = handle;
@@ -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));
}
@@ -8,7 +8,7 @@ import * as fs from "node:fs/promises";
import * as os from "node:os";
import path from "node:path";
import { $env, prompt, Snowflake } from "@oh-my-pi/pi-utils";
import { resolveAgentModelPatterns } from "../config/model-resolver";
import { resolveAgentModelPatterns, resolveAgentModelSource, resolveExplicitModelRole } from "../config/model-resolver";
import type { LocalProtocolOptions } from "../internal-urls";
import { registerArtifactsDir } from "../internal-urls/registry-helpers";
import { MCPManager } from "../mcp/manager";
@@ -123,6 +123,8 @@ export interface EffectiveSubagentPolicy {
agent: AgentDefinition;
effectiveAgent: AgentDefinition;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
parentActiveModelPattern?: string;
schema: StructuredSubagentSchemaResolution;
planMode: boolean;
@@ -278,13 +280,18 @@ export async function resolveEffectiveSubagentPolicy(
}
const agentModelOverrides = request.session.settings.get("task.agentModelOverrides");
const parentActiveModelPattern = request.session.getActiveModelString?.();
const modelOverride = resolveAgentModelPatterns({
settingsOverride: request.model ?? agentModelOverrides[agentName],
const modelResolution = {
requestModel: request.model,
settingsOverride: agentModelOverrides[agentName],
agentModel: effectiveAgent.model,
settings: request.session.settings,
activeModelPattern: parentActiveModelPattern,
fallbackModelPattern: request.session.getModelString?.(),
});
};
// Keep role identity from the same effective non-empty source that supplies
// model selection: caller request, settings override, then agent definition.
const modelRole = resolveExplicitModelRole(resolveAgentModelSource(modelResolution), request.session.settings);
const modelOverride = resolveAgentModelPatterns(modelResolution);
const isolationMode = request.session.settings.get("task.isolation.mode");
const isIsolated = request.isolation?.requested === true;
if (isIsolated && isolationMode === "none") {
@@ -299,6 +306,7 @@ export async function resolveEffectiveSubagentPolicy(
agent,
effectiveAgent,
modelOverride,
modelRole,
parentActiveModelPattern,
schema,
planMode,
@@ -394,6 +402,7 @@ function buildExecutorOptions(
invokedAt: request.invokedAt,
acquiredAt: request.acquiredAt,
modelOverride: policy.modelOverride,
modelRole: policy.modelRole,
parentActiveModelPattern: policy.parentActiveModelPattern,
thinkingLevel: policy.effectiveAgent.thinkingLevel,
effort: request.effort,
@@ -476,6 +485,7 @@ function buildFailureResult(
tokens: 0,
requests: 0,
modelOverride: policy.modelOverride,
modelRole: policy.modelRole,
error: message,
};
};
+4
View File
@@ -426,6 +426,8 @@ export interface AgentProgress {
cost: number;
durationMs: number;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
/** Resolved model display string in the form `<provider>/<id>`, optionally suffixed with `:<thinkingLevel>` when the level was set explicitly. Undefined when the model could not be resolved. */
resolvedModel?: string;
/** True when {@link resolvedModel} is the target of an active retry fallback (not the originally configured model). Lets observer-only UIs (collab guests, Agent Hub rows with no live session) flag the fallback and keep the provider. */
@@ -495,6 +497,8 @@ export interface SingleResult {
/** Model's context window in tokens, when known. */
contextWindow?: number;
modelOverride?: string | string[];
/** Explicit pre-expansion model role alias selected for this run. */
modelRole?: string;
/** Resolved model display string in the form `<provider>/<id>`, optionally suffixed with `:<thinkingLevel>` when the level was set explicitly. Omitted from tool-result JSON when undefined to keep wire payloads small. */
resolvedModel?: string;
/** True when {@link resolvedModel} is the target of an active retry fallback. Mirrors {@link AgentProgress.resolvedModelIsFallback} onto the settled result. */
@@ -889,7 +889,7 @@ export async function applyStealthPatches(
export function systemChromiumCandidatesForTest(
platform: NodeJS.Platform = process.platform,
home?: string,
which?: (name: string) => string | undefined,
which?: (name: string) => string | null | undefined,
): string[] {
return systemChromiumCandidates(platform, home, which);
}
+5 -5
View File
@@ -25,7 +25,6 @@ import { MCPManager } from "../mcp/manager";
import vibeTurnResultTemplate from "../prompts/tools/vibe-turn-result.md" with { type: "text" };
import { AgentLifecycleManager } from "../registry/agent-lifecycle";
import { type AgentRef, AgentRegistry, MAIN_AGENT_ID } from "../registry/agent-registry";
import type { SessionEntry } from "../session/session-entries";
import { SessionManager, SessionPersistenceIndeterminateError } from "../session/session-manager";
import { getBundledAgent } from "../task/agents";
import { type ExecutorOptions, runSubagentFollowUpTurn, runSubprocess } from "../task/executor";
@@ -363,11 +362,12 @@ function parseLifecycleEvent(value: unknown): VibeLifecycleEvent | undefined {
return undefined;
}
/** Child ids claimed by any valid Vibe spawn event, independent of current parent scope. */
export function persistedVibeChildIds(entries: Iterable<SessionEntry>): Set<string> {
/** Child ids claimed by valid Vibe spawn records from untrusted persisted JSON. */
export function persistedVibeChildIds(entries: Iterable<unknown>): Set<string> {
const ids = new Set<string>();
for (const entry of entries) {
if (entry.type !== "custom" || entry.customType !== VIBE_LIFECYCLE_CUSTOM_TYPE) continue;
for (const value of entries) {
const entry = objectRecord(value);
if (entry?.type !== "custom" || entry.customType !== VIBE_LIFECYCLE_CUSTOM_TYPE) continue;
const event = parseLifecycleEvent(entry.data);
if (
event?.action === "spawn" &&
@@ -15,7 +15,9 @@ 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";
const AGENT_ID = "Worker";
@@ -36,6 +38,7 @@ function makeHub(focusAgent: (id: string) => Promise<void>) {
const done = Promise.withResolvers<void>();
const renderRequested = Promise.withResolvers<void>();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {
@@ -50,6 +53,30 @@ function makeHub(focusAgent: (id: string) => Promise<void>) {
return { hub, doneCalls: () => doneCalls, done: done.promise, renderRequested: renderRequested.promise };
}
const ROSTER_ENTRY_PATTERN = /^(❯| ) (\S+) (?:(?:(?:│ {3}| {4})*)(?:├── |└── ))?(\S+)/u;
function renderedRosterEntry(hub: AgentHubOverlayComponent, id: string, width: number): string {
const cells = hub.render(width).map(raw => {
const line = Bun.stripANSI(raw);
if (!line.startsWith("│ ")) return undefined;
const divider = line.indexOf("│", Math.max(2, Math.floor(line.length / 3)));
return divider < 0 ? undefined : line.slice(2, Math.max(2, divider - 1));
});
const start = cells.findIndex(cell => {
const match = cell ? ROSTER_ENTRY_PATTERN.exec(cell) : null;
return match?.[3] === id;
});
expect(start).toBeGreaterThanOrEqual(0);
const entry: string[] = [];
for (let i = start; i < cells.length; i++) {
const cell = cells[i];
if (cell === undefined || cell.trim().length === 0) break;
if (i > start && ROSTER_ENTRY_PATTERN.test(cell)) break;
entry.push(cell.trimEnd());
}
return entry.join("\n");
}
describe("Agent hub Enter activation", () => {
beforeAll(() => {
initTheme();
@@ -99,6 +126,7 @@ describe("Agent hub Enter activation", () => {
await Bun.write(workerSessionFile, "");
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
@@ -110,13 +138,226 @@ describe("Agent hub Enter activation", () => {
});
await hub.persistedSubagentsReady;
const rendered = Bun.stripANSI(hub.render(120).join("\n"));
expect(rendered).toContain("Worker");
expect(rendered).toContain("parked");
const workerEntry = renderedRosterEntry(hub, "Worker", 120);
expect(workerEntry).toContain("○ Worker");
expect(agents.get("Worker")?.sessionFile).toBe(workerSessionFile);
hub.dispose();
});
it("stops persisted discovery when the Hub is disposed", async () => {
using tempDir = TempDir.createSync("@omp-agent-hub-disposed-scan-");
const sessionFile = path.join(tempDir.path(), "main.jsonl");
await Bun.write(sessionFile, "");
await Bun.write(path.join(tempDir.path(), "main", "Worker.jsonl"), "");
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
sessionFile,
});
hub.dispose();
await hub.persistedSubagentsReady;
expect(agents.get("Worker")).toBeUndefined();
});
it("restores nested parent lineage after restart", async () => {
using tempDir = TempDir.createSync("@omp-agent-hub-persisted-tree-");
const sessionFile = path.join(tempDir.path(), "main.jsonl");
const parentSessionFile = path.join(tempDir.path(), "main", "Parent.jsonl");
const childSessionFile = path.join(tempDir.path(), "main", "Parent", "Child.jsonl");
await Bun.write(sessionFile, "");
await Bun.write(parentSessionFile, "");
await Bun.write(childSessionFile, "");
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
sessionFile,
});
await hub.persistedSubagentsReady;
expect(agents.get("Parent")?.parentId).toBe("Main");
expect(agents.get("Child")?.parentId).toBe("Parent");
hub.handleInput("t");
expect(Bun.stripANSI(renderedRosterEntry(hub, "Child", 120))).toContain("└── Child");
hub.dispose();
});
it("restores saved task metadata and timestamps for completed agents", async () => {
using tempDir = TempDir.createSync("@omp-agent-hub-persisted-metadata-");
const sessionFile = path.join(tempDir.path(), "main.jsonl");
const workerSessionFile = path.join(tempDir.path(), "main", "Worker.jsonl");
const createdAt = "2026-07-30T01:13:37.835Z";
const lastActivity = new Date("2026-07-30T01:15:00.000Z");
await Bun.write(sessionFile, "");
await Bun.write(
workerSessionFile,
[
JSON.stringify({ type: "session", version: 3, id: "worker-session", timestamp: createdAt, cwd: TEST_CWD }),
JSON.stringify({
type: "session_init",
id: "init",
parentId: null,
timestamp: createdAt,
systemPrompt: "system",
task: "Complete the assignment below, thoroughly:\n\n# Target\nInspect dependency boundaries and report unsafe coupling.\n\n# Change\nRead the implementation.",
tools: ["read"],
}),
].join("\n"),
);
await fs.utimes(workerSessionFile, lastActivity, lastActivity);
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
sessionFile,
});
await hub.persistedSubagentsReady;
expect(agents.get("Worker")).toMatchObject({
activity: "Inspect dependency boundaries and report unsafe coupling.",
createdAt: Date.parse(createdAt),
lastActivity: lastActivity.getTime(),
status: "parked",
});
const workerEntry = renderedRosterEntry(hub, "Worker", 120);
expect(workerEntry).toContain("Inspect dependency boundaries and report unsafe coupling.");
expect(workerEntry.replace(/\s+/g, " ")).toContain("usage —");
expect(workerEntry).not.toContain("$0.000");
hub.dispose();
});
it("restores persisted model role, usage, spend, and tool totals", async () => {
using tempDir = TempDir.createSync("@omp-agent-hub-persisted-usage-");
const sessionFile = path.join(tempDir.path(), "main.jsonl");
const workerSessionFile = path.join(tempDir.path(), "main", "Worker.jsonl");
const createdAt = "2026-07-30T01:13:30.000Z";
const lastActivity = new Date("2026-07-30T01:15:00.000Z");
await Bun.write(sessionFile, "");
await Bun.write(
workerSessionFile,
[
JSON.stringify({ type: "session", version: 3, id: "worker-session", timestamp: createdAt, cwd: TEST_CWD }),
JSON.stringify({
type: "model_change",
id: "model",
parentId: null,
timestamp: createdAt,
model: "openai-codex/gpt-5.6-luna",
// Historical concrete overrides did not persist a model-role field.
}),
JSON.stringify({
type: "session_init",
id: "init",
parentId: "model",
timestamp: createdAt,
systemPrompt: `base prompt\n\nROLE\n====\n${getBundledAgent("scout")?.systemPrompt}`,
task: "Inspect persisted telemetry.",
tools: ["read", "grep"],
}),
JSON.stringify({
type: "message",
id: "assistant",
parentId: "init",
timestamp: lastActivity.toISOString(),
message: {
role: "assistant",
timestamp: lastActivity.getTime(),
content: [
{ type: "toolCall", id: "read-call", name: "read", arguments: { path: "src/a.ts" } },
{ type: "toolCall", id: "grep-call", name: "grep", arguments: { pattern: "needle" } },
],
usage: {
input: 100,
output: 25,
cacheRead: 200,
cacheWrite: 10,
totalTokens: 335,
cost: { input: 0.01, output: 0.1, cacheRead: 0.01, cacheWrite: 0.003, total: 0.123 },
},
},
}),
].join("\n"),
);
await fs.utimes(workerSessionFile, lastActivity, lastActivity);
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
sessionFile,
});
await hub.persistedSubagentsReady;
const workerEntry = renderedRosterEntry(hub, "Worker", 120).replace(/\s+/g, " ");
expect(workerEntry).toContain("SMOL");
expect(workerEntry).toContain("$0.123");
expect(workerEntry).toContain("1m30s");
expect(workerEntry).toContain("1 req");
expect(workerEntry).toContain("2 tools");
expect(workerEntry).toContain("135 tok");
expect(Bun.stripANSI(hub.render(120).join("\n"))).toContain("Read-only · 0 LoC");
hub.dispose();
});
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(), "session.jsonl");
const entry = JSON.stringify({
type: "message",
id: "entry",
parentId: null,
timestamp: "2026-07-30T01:13:30.000Z",
message: { role: "user", content: [{ type: "text", text: "small" }] },
});
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;
});
try {
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 {
vi.useRealTimers();
}
});
it("does not generically revive active or tombstoned Vibe children copied by a post-exit fork", async () => {
using tempDir = TempDir.createSync("@omp-agent-hub-vibe-fork-");
const manager = SessionManager.create(tempDir.path(), tempDir.path());
@@ -160,6 +401,7 @@ describe("Agent hub Enter activation", () => {
const agents = new AgentRegistry();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
@@ -191,21 +433,22 @@ describe("Agent hub Enter activation", () => {
const editor = {};
let capturedHub: AgentHubOverlayComponent | undefined;
let editorRestoredCount = 0;
const focusedIds: string[] = [];
const focusResolved = Promise.withResolvers<void>();
const editorFocused = Promise.withResolvers<void>();
const focusTargets: unknown[] = [];
const editorContainer = {
children: [editor],
clear: () => {},
addChild: (child: unknown) => {
if (child === editor) editorRestoredCount++;
else capturedHub = child as AgentHubOverlayComponent;
},
addChild: () => {},
};
const ctx = {
keybindings: { getKeys: () => [] },
ui: {
showOverlay: (component: AgentHubOverlayComponent) => {
capturedHub = component;
return { hide: () => {} };
},
setFocus: (target: unknown) => {
focusTargets.push(target);
if (target === editor) editorFocused.resolve();
@@ -235,7 +478,6 @@ describe("Agent hub Enter activation", () => {
await editorFocused.promise;
expect(focusedIds).toEqual([AGENT_ID]);
expect(editorRestoredCount).toBe(1);
expect(focusTargets.at(-1)).toBe(editor);
capturedHub!.dispose();
});
@@ -252,12 +494,19 @@ describe("Agent hub double-← gating", () => {
function setup(agents: AgentRegistry, sessionFile: string | null = null) {
let shown: AgentHubOverlayComponent | undefined;
let overlayOptions: Record<string, unknown> | undefined;
const shownReady = Promise.withResolvers<AgentHubOverlayComponent>();
const editor = {};
const focusTargets: unknown[] = [];
const ctx = {
keybindings: { getKeys: () => [] },
ui: {
showOverlay: (component: AgentHubOverlayComponent, options: Record<string, unknown>) => {
shown = component;
overlayOptions = options;
shownReady.resolve(component);
return { hide: () => {} };
},
setFocus: (target: unknown) => {
focusTargets.push(target);
},
@@ -265,13 +514,9 @@ describe("Agent hub double-← gating", () => {
},
editor,
editorContainer: {
children: [editor],
clear: () => {},
addChild: (child: unknown) => {
if (child !== editor) {
shown = child as AgentHubOverlayComponent;
shownReady.resolve(shown);
}
},
addChild: () => {},
},
collabGuest: { agentRegistry: agents, hubRemote: undefined },
focusAgentSession: async () => {},
@@ -285,6 +530,7 @@ describe("Agent hub double-← gating", () => {
editor,
shown: () => shown,
shownReady: shownReady.promise,
overlayOptions: () => overlayOptions,
focusTargets,
};
}
@@ -347,14 +593,24 @@ describe("Agent hub double-← gating", () => {
shownHub!.dispose();
});
it("the explicit hub key opens the empty roster even with no subagents", () => {
it("the explicit hub opens fullscreen before persisted subagents load", async () => {
using tempDir = TempDir.createSync("@omp-agent-hub-explicit-");
const sessionFile = path.join(tempDir.path(), "main.jsonl");
await Bun.write(sessionFile, "");
await Bun.write(path.join(tempDir.path(), "main", "Worker.jsonl"), "");
const agents = new AgentRegistry();
const { controller, shown } = setup(agents);
const { controller, shown, overlayOptions } = setup(agents, sessionFile);
controller.showAgentHub(new SessionObserverRegistry());
expect(shown()).toBeDefined();
shown()!.dispose();
const hub = shown();
expect(hub).toBeDefined();
expect(overlayOptions()).toMatchObject({ width: "100%", maxHeight: "100%", margin: 0, fullscreen: true });
expect(agents.get("Worker")).toBeUndefined();
expect(Bun.stripANSI(hub!.render(120).join("\n"))).toContain("Loading saved agents");
await hub!.persistedSubagentsReady;
expect(agents.get("Worker")?.status).toBe("parked");
hub!.dispose();
});
it("armCloseTap lets a single ← dismiss the hub the opening ←← raised", () => {
@@ -404,6 +660,7 @@ describe("Agent hub data refresh coalescing", () => {
const observers = new SessionObserverRegistry();
const requestRender = vi.fn();
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers,
hubKeys: [],
onDone: () => {},
@@ -447,4 +704,120 @@ describe("Agent hub data refresh coalescing", () => {
vi.useRealTimers();
}
});
it("refreshes direct-session fallback stats on the age cadence, not paints or heartbeats", async () => {
vi.useFakeTimers();
const agents = new AgentRegistry();
const observers = new SessionObserverRegistry();
const requestRender = vi.fn();
let inputTokens = 100;
let assistantMessages = 1;
const getSessionStats = vi.fn(() => ({
sessionFile: undefined,
sessionId: "sdk-agent",
userMessages: 1,
assistantMessages,
toolCalls: 2,
toolResults: 2,
totalMessages: 6,
tokens: {
input: inputTokens,
output: 50,
reasoning: 0,
cacheRead: 20,
cacheWrite: 0,
total: inputTokens + 70,
},
premiumRequests: 0,
cost: 0.1,
}));
agents.register({
id: "SdkAgent",
displayName: "SDK agent",
kind: "sub",
parentId: "Main",
session: { getSessionStats, subscribe: () => () => {} } as unknown as AgentSession,
status: "running",
});
const hub = new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers,
hubKeys: [],
onDone: () => {},
requestRender,
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
});
try {
await hub.persistedSubagentsReady;
expect(getSessionStats).toHaveBeenCalledTimes(1);
for (let i = 0; i < 4; i++) hub.render(120);
expect(getSessionStats).toHaveBeenCalledTimes(1);
expect(Bun.stripANSI(hub.render(120).join("\n"))).toContain("150 tok");
inputTokens = 400;
assistantMessages = 2;
agents.setActivity("SdkAgent", "heartbeat");
vi.advanceTimersByTime(100);
expect(getSessionStats).toHaveBeenCalledTimes(1);
expect(Bun.stripANSI(hub.render(120).join("\n"))).toContain("150 tok");
vi.advanceTimersByTime(4_899);
expect(getSessionStats).toHaveBeenCalledTimes(1);
vi.advanceTimersByTime(1);
expect(getSessionStats).toHaveBeenCalledTimes(2);
const refreshed = Bun.stripANSI(hub.render(120).join("\n"));
expect(refreshed).toContain("450 tok");
expect(refreshed).toContain("2 req");
expect(refreshed).toContain("1/1");
expect(refreshed).toContain("measured");
hub.render(120);
expect(getSessionStats).toHaveBeenCalledTimes(2);
} finally {
hub.dispose();
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();
}
});
});
@@ -6,10 +6,12 @@
* agents that appear while the hub is open are appended at the end.
*/
import { afterEach, beforeAll, describe, expect, it, setSystemTime, vi } from "bun:test";
import { ThinkingLevel } from "@oh-my-pi/pi-agent-core";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { IrcBus } from "@oh-my-pi/pi-coding-agent/irc/bus";
import { AgentHubOverlayComponent } from "@oh-my-pi/pi-coding-agent/modes/components/agent-hub";
import { type AgentHubDeps, AgentHubOverlayComponent } from "@oh-my-pi/pi-coding-agent/modes/components/agent-hub";
import { SessionObserverRegistry } from "@oh-my-pi/pi-coding-agent/modes/session-observer-registry";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import { initTheme, theme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
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 { visibleWidth } from "@oh-my-pi/pi-tui/utils";
@@ -40,8 +42,9 @@ function stubStdoutGeometry(cols: number): GeometryStub {
};
}
function makeHub(agents: AgentRegistry) {
function makeHub(agents: AgentRegistry, overrides: Partial<AgentHubDeps> = {}) {
return new AgentHubOverlayComponent({
settings: Settings.isolated(),
observers: new SessionObserverRegistry(),
hubKeys: [],
onDone: () => {},
@@ -49,18 +52,78 @@ function makeHub(agents: AgentRegistry) {
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
...overrides,
});
}
function renderedAgentIds(hub: AgentHubOverlayComponent): string[] {
// Entry first lines are ` <cursor> <status-glyph> <id> …`; task lines are
// indented deeper and chrome lines never carry the cursor slot.
const ids: string[] = [];
for (const raw of hub.render(120)) {
const match = /^ (?:❯| ) (\S+) (\S+)/u.exec(Bun.stripANSI(raw));
if (match) ids.push(match[2]!);
interface RenderedAgentRow {
id: string;
selected: boolean;
}
const ROSTER_ENTRY_PATTERN = /^(❯| ) (\S+) (?:(?:(?:│ {3}| {4})*)(?:├── |└── ))?(\S+)/u;
function rosterCell(raw: string): string | undefined {
const line = Bun.stripANSI(raw);
if (!line.startsWith("│ ")) return undefined;
const divider = line.indexOf("│", Math.max(2, Math.floor(line.length / 3)));
if (divider < 0) return undefined;
return line.slice(2, Math.max(2, divider - 1));
}
function renderedAgentRows(hub: AgentHubOverlayComponent, width = 120): RenderedAgentRow[] {
// Roster entry first cells are
// `<cursor> <status-glyph> [tree-prefix] <id> …`; task cells are
// indented deeper and never match the cursor/status slots.
const rows: RenderedAgentRow[] = [];
for (const raw of hub.render(width)) {
const cell = rosterCell(raw);
const match = cell ? ROSTER_ENTRY_PATTERN.exec(cell) : null;
if (match) rows.push({ id: match[3]!, selected: match[1] === "❯" });
}
return ids;
return rows;
}
function renderedAgentIds(hub: AgentHubOverlayComponent): string[] {
return renderedAgentRows(hub).map(row => row.id);
}
function selectedAgentId(hub: AgentHubOverlayComponent): string | undefined {
return renderedAgentRows(hub).find(row => row.selected)?.id;
}
function renderedRosterEntry(hub: AgentHubOverlayComponent, id: string, width: number): string {
const cells = hub.render(width).map(rosterCell);
const start = cells.findIndex(cell => {
const match = cell ? ROSTER_ENTRY_PATTERN.exec(cell) : null;
return match?.[3] === id;
});
expect(start).toBeGreaterThanOrEqual(0);
const entry: string[] = [];
for (let i = start; i < cells.length; i++) {
const cell = cells[i];
if (cell === undefined || cell.trim().length === 0) break;
if (i > start && ROSTER_ENTRY_PATTERN.test(cell)) break;
entry.push(cell.trimEnd());
}
return entry.join("\n");
}
function renderedRosterHeaderLineRaw(hub: AgentHubOverlayComponent, id: string, width: number): string {
const line = hub.render(width).find(raw => {
const cell = rosterCell(raw);
const match = cell ? ROSTER_ENTRY_PATTERN.exec(cell) : null;
return match?.[3] === id;
});
if (!line) throw new Error(`No rendered roster header for ${id}`);
return line;
}
function leftClick(row1Based: number): string {
return `\x1b[<0;4;${row1Based}M`;
}
function wheel(direction: "up" | "down"): string {
return `\x1b[<${direction === "down" ? 65 : 64};4;4M`;
}
describe("Agent hub row ordering", () => {
@@ -79,6 +142,20 @@ describe("Agent hub row ordering", () => {
AgentRegistry.resetGlobalForTests();
});
it("renders a useful empty state before any task agents exist", () => {
geometry = stubStdoutGeometry(120);
const hub = makeHub(new AgentRegistry());
try {
const rendered = Bun.stripANSI(hub.render(120).join("\n"));
expect(rendered).toContain("No agents in this session");
expect(rendered).toContain("Finished, parked, and killed subagents remain with the session");
expect(rendered).toContain("Resume that session with omp-dev --continue, or spawn a task here.");
} finally {
hub.dispose();
}
});
it("freezes the initial lastActivity order while the hub is open", () => {
vi.useFakeTimers();
let hub: AgentHubOverlayComponent | undefined;
@@ -99,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);
@@ -147,9 +225,9 @@ describe("Agent hub row ordering", () => {
getSessions.mockClear();
getSession.mockClear();
const visibleIds = renderedAgentIds(hub);
// rows=12 → line budget 5 for single-line parked rows; plus one failed
// boundary probe past the window. Must not scale with the 10_000 roster.
expect(visibleIds).toHaveLength(5);
// rows=12 → line budget 5; unknown usage is an explicit second line,
// so two complete entries fit while rendering remains viewport-bounded.
expect(visibleIds).toHaveLength(2);
expect(getSessions).not.toHaveBeenCalled();
expect(getSession.mock.calls.length).toBeLessThanOrEqual(8);
expect(getSession.mock.calls.length).toBeGreaterThan(0);
@@ -163,7 +241,8 @@ describe("Agent hub row ordering", () => {
getSession.mockClear();
hub.handleInput("j");
const afterMove = renderedAgentIds(hub);
expect(afterMove).toHaveLength(5);
expect(afterMove.length).toBeGreaterThan(0);
expect(afterMove.length).toBeLessThanOrEqual(2);
expect(afterMove).toContain(visibleIds[1]!);
expect(getSessions).not.toHaveBeenCalled();
expect(getSession.mock.calls.length).toBeLessThanOrEqual(8);
@@ -229,42 +308,110 @@ 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,
});
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSession").mockReturnValue({
id: "RevAgentStream",
kind: "subagent",
label: "Subagent",
status: "active",
description: "Complete the assignment below, thoroughly:\n- check performance\n- check leaks",
lastUpdate: Date.now(),
});
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "RevAgentStream",
kind: "subagent",
label: "Subagent",
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 = new AgentHubOverlayComponent({
observers,
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
});
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);
expect(cleanLine.includes("\r")).toBe(false);
const width = visibleWidth(line);
expect(width).toBeLessThanOrEqual(78);
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();
});
it("fits the fullscreen table to short terminals and windows large registries", () => {
geometry = stubStdoutGeometry(80);
geometry.setRows(10);
const agents = new AgentRegistry();
for (let i = 0; i < 50; i++) {
agents.register({
id: `Agent${i}`,
displayName: `Agent ${i}`,
kind: "sub",
session: {} as AgentSession,
});
}
const observers = new SessionObserverRegistry();
const sessions = vi.spyOn(observers, "getSessions").mockReturnValue([]);
const hub = makeHub(agents, { observers });
try {
const lines = hub.render(80);
expect(lines.length).toBe(10);
expect(sessions.mock.calls.length).toBeLessThan(agents.list().length);
expect(Bun.stripANSI(lines.join("\n"))).toContain("…");
} finally {
hub.dispose();
}
});
it("matches fullscreen menu mouse selection, wheel, and activation", async () => {
geometry = stubStdoutGeometry(120);
const agents = new AgentRegistry();
setSystemTime(1_000);
agents.register({ id: "Alpha", displayName: "Alpha", kind: "sub", session: {} as AgentSession });
setSystemTime(2_000);
agents.register({ id: "Beta", displayName: "Beta", kind: "sub", session: {} as AgentSession });
setSystemTime(3_000);
agents.register({ id: "Gamma", displayName: "Gamma", kind: "sub", session: {} as AgentSession });
const focused: string[] = [];
const done = vi.fn();
const hub = makeHub(agents, {
onDone: done,
focusAgent: async id => {
focused.push(id);
},
});
try {
expect(selectedAgentId(hub)).toBe("Gamma");
hub.handleInput(wheel("down"));
expect(selectedAgentId(hub)).toBe("Beta");
const frame = hub.render(120);
const alphaRow = frame.findIndex(line => /^│ {3}\S+ Alpha/u.test(Bun.stripANSI(line)));
expect(alphaRow).toBeGreaterThanOrEqual(0);
hub.handleInput(`\x1b[<0;110;${alphaRow + 1}M`);
expect(selectedAgentId(hub)).toBe("Beta");
expect(focused).toEqual([]);
hub.handleInput(leftClick(alphaRow + 1));
await Promise.resolve();
expect(selectedAgentId(hub)).toBe("Alpha");
expect(focused).toEqual(["Alpha"]);
expect(done).toHaveBeenCalledTimes(1);
} finally {
hub.dispose();
}
});
it("flags a fallback badge for observer-only rows with no live session", () => {
geometry = stubStdoutGeometry(120);
@@ -286,15 +433,7 @@ describe("Agent hub row ordering", () => {
} as never,
});
const hub = new AgentHubOverlayComponent({
observers,
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
});
const hub = makeHub(agents, { observers });
try {
expect(Bun.stripANSI(hub.render(120).join("\n"))).toContain("fallback → openai/gpt-4o");
@@ -326,15 +465,7 @@ describe("Agent hub row ordering", () => {
} as never,
});
const hub = new AgentHubOverlayComponent({
observers,
hubKeys: [],
onDone: () => {},
requestRender: () => {},
registry: agents,
irc: new IrcBus(agents),
focusAgent: async () => {},
});
const hub = makeHub(agents, { observers });
try {
expect(Bun.stripANSI(hub.render(120).join("\n"))).toContain("fallback → fireworks/kimi-k2");
@@ -342,4 +473,487 @@ describe("Agent hub row ordering", () => {
hub.dispose();
}
});
it("retains the live thinking level unless progress has an explicit suffix", () => {
geometry = stubStdoutGeometry(140);
const agents = new AgentRegistry();
const inheritedSession = { thinkingLevel: ThinkingLevel.High } as unknown as AgentSession;
const explicitSession = { thinkingLevel: ThinkingLevel.High } as unknown as AgentSession;
agents.register({
id: "InheritedLevel",
displayName: "Inherited level",
kind: "sub",
session: inheritedSession,
});
agents.register({
id: "ExplicitLevel",
displayName: "Explicit level",
kind: "sub",
session: explicitSession,
});
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "InheritedLevel",
kind: "subagent",
label: "Inherited level",
status: "active",
lastUpdate: Date.now(),
progress: { resolvedModel: "openai/gpt-5.4" } as never,
},
{
id: "ExplicitLevel",
kind: "subagent",
label: "Explicit level",
status: "active",
lastUpdate: Date.now(),
progress: { resolvedModel: "openai/gpt-5.4:low" } as never,
},
]);
const hub = makeHub(agents, { observers });
try {
const inherited = renderedRosterEntry(hub, "InheritedLevel", 140);
expect(inherited).toContain("gpt-5.4");
expect(inherited).toContain(theme.thinking.high);
const explicit = renderedRosterEntry(hub, "ExplicitLevel", 140);
expect(explicit).toContain("gpt-5.4");
expect(explicit).toContain(theme.thinking.low);
expect(explicit).not.toContain(theme.thinking.high);
} finally {
hub.dispose();
}
});
it("renders aggregate usage and a selected-agent inspector without inventing change attribution", () => {
geometry = stubStdoutGeometry(140);
geometry.setRows(28);
const agents = new AgentRegistry();
agents.register({
id: "Reviewer",
displayName: "Security Reviewer",
kind: "sub",
parentId: "Main",
session: null,
history: {
outputPath: "/tmp/Reviewer.md",
patchPath: "/tmp/Reviewer.patch",
branchName: "omp/task/Reviewer",
},
});
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "Reviewer",
kind: "subagent",
label: "Reviewer",
description: "Review the session lifecycle and produce actionable findings",
status: "active",
lastUpdate: Date.now(),
progress: {
id: "Reviewer",
index: 0,
agent: "reviewer",
agentSource: "bundled",
status: "running",
task: "Review the session lifecycle",
currentTool: "read",
currentToolArgs: "src/session/agent-session.ts",
recentTools: [],
recentOutput: [],
toolCount: 27,
requests: 12,
tokens: 18_400,
contextTokens: 31_000,
contextWindow: 128_000,
cost: 0.2134,
durationMs: 134_000,
resolvedModel: "openai/gpt-5.4:high",
} as never,
},
]);
const hub = makeHub(agents, { observers });
try {
const rendered = Bun.stripANSI(hub.render(140).join("\n"));
expect(rendered).toContain("1 running");
expect(rendered).toContain("Flat");
expect(rendered).toContain("By parent");
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%");
expect(rendered).toContain("Registered ");
expect(rendered).toContain("Shared workspace · per-agent LoC not attributable");
expect(rendered).toContain("Output /tmp/Reviewer.md");
expect(rendered).toContain("Patch /tmp/Reviewer.patch");
hub.handleInput("\x1b[6~");
expect(Bun.stripANSI(hub.render(140).join("\n"))).toContain("Worktree branch omp/task/Reviewer");
} finally {
hub.dispose();
}
});
it("shows dense measured usage for running and completed progress with aggregate coverage", () => {
geometry = stubStdoutGeometry(160);
geometry.setRows(32);
const agents = new AgentRegistry();
agents.register({ id: "Running", displayName: "Running", kind: "sub", session: null, status: "running" });
agents.register({ id: "Completed", displayName: "Completed", kind: "sub", session: null, status: "idle" });
agents.register({
id: "Historical",
displayName: "Historical",
kind: "sub",
session: null,
status: "parked",
activity: "Restored task",
});
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "Running",
kind: "subagent",
label: "Running",
status: "active",
lastUpdate: Date.now(),
progress: {
id: "Running",
index: 0,
agent: "worker",
agentSource: "bundled",
status: "running",
task: "Run checks",
recentTools: [],
recentOutput: [],
toolCount: 4,
requests: 3,
tokens: 1_200,
cost: 0.1234,
durationMs: 6_500,
} as never,
},
{
id: "Completed",
kind: "subagent",
label: "Completed",
status: "completed",
lastUpdate: Date.now(),
progress: {
id: "Completed",
index: 1,
agent: "worker",
agentSource: "bundled",
status: "completed",
task: "Finish checks",
recentTools: [],
recentOutput: [],
toolCount: 8,
requests: 5,
tokens: 2_500,
cost: 0.4567,
durationMs: 125_000,
} as never,
},
]);
const hub = makeHub(agents, { observers });
try {
const rendered = Bun.stripANSI(hub.render(160).join("\n"));
expect(rendered).toContain("2/3 measured");
expect(rendered).toContain("$0.580");
expect(rendered).toContain("3.7K tok");
expect(rendered).toContain("8 req");
expect(rendered).toContain("12 tools");
expect(rendered).toContain("2m11s active agent time");
const running = renderedRosterEntry(hub, "Running", 160);
expect(running).toContain("$0.123");
expect(running).toContain("6.5s");
expect(running).toContain("3 req");
expect(running).toContain("4 tools");
expect(running).toContain("1.2K tok");
const completed = renderedRosterEntry(hub, "Completed", 160);
expect(completed).toContain("$0.457");
expect(completed).toContain("2m5s");
expect(completed).toContain("5 req");
expect(completed).toContain("8 tools");
expect(completed).toContain("2.5K tok");
const historical = renderedRosterEntry(hub, "Historical", 160);
expect(historical).toContain("Restored task");
expect(historical).toContain("usage —");
expect(historical).not.toContain("$0.000");
} finally {
hub.dispose();
}
});
it("treats incomplete and non-finite progress usage as unknown", () => {
geometry = stubStdoutGeometry(160);
const agents = new AgentRegistry();
const getSessionStats = vi.fn(() => ({
sessionFile: undefined,
sessionId: "incomplete",
userMessages: 1,
assistantMessages: 9,
toolCalls: 4,
toolResults: 4,
totalMessages: 18,
tokens: { input: 100, output: 50, reasoning: 0, cacheRead: 0, cacheWrite: 0, total: 150 },
premiumRequests: 0,
cost: 0.25,
}));
agents.register({
id: "Incomplete",
displayName: "Incomplete",
kind: "sub",
session: { getSessionStats } as unknown as AgentSession,
});
agents.register({ id: "NonFinite", displayName: "Non-finite", kind: "sub", session: null });
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "Incomplete",
kind: "subagent",
label: "Incomplete",
status: "active",
lastUpdate: Date.now(),
progress: {
tokens: 100,
toolCount: 2,
cost: 0.1,
durationMs: 1_000,
} as never,
},
{
id: "NonFinite",
kind: "subagent",
label: "Non-finite",
status: "active",
lastUpdate: Date.now(),
progress: {
tokens: Number.NaN,
requests: 2,
toolCount: 2,
cost: 0.1,
durationMs: 1_000,
} as never,
},
]);
const hub = makeHub(agents, { observers });
try {
const rendered = Bun.stripANSI(hub.render(160).join("\n"));
expect(rendered).toContain("0/2 measured");
expect(renderedRosterEntry(hub, "Incomplete", 160)).toContain("usage —");
expect(renderedRosterEntry(hub, "NonFinite", 160)).toContain("usage —");
expect(getSessionStats).not.toHaveBeenCalled();
} finally {
hub.dispose();
}
});
it("shows configured role text beside a resolved model but not for an explicit selector", () => {
geometry = stubStdoutGeometry(160);
const agents = new AgentRegistry();
agents.register({ id: "RoleAgent", displayName: "Role Agent", kind: "sub", session: null });
agents.register({ id: "ExplicitAgent", displayName: "Explicit Agent", kind: "sub", session: null });
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "RoleAgent",
kind: "subagent",
label: "Role Agent",
status: "active",
lastUpdate: Date.now(),
progress: {
id: "RoleAgent",
index: 0,
agent: "worker",
agentSource: "bundled",
status: "running",
task: "Run with the configured role",
recentTools: [],
recentOutput: [],
toolCount: 0,
requests: 1,
tokens: 10,
cost: 0,
durationMs: 100,
modelRole: "rapid",
resolvedModel: "openai/gpt-4o",
} as never,
},
{
id: "ExplicitAgent",
kind: "subagent",
label: "Explicit Agent",
status: "active",
lastUpdate: Date.now(),
progress: {
id: "ExplicitAgent",
index: 1,
agent: "worker",
agentSource: "bundled",
status: "running",
task: "Run with an explicit selector",
recentTools: [],
recentOutput: [],
toolCount: 0,
requests: 1,
tokens: 10,
cost: 0,
durationMs: 100,
resolvedModel: "openai/gpt-4o",
} as never,
},
]);
const hub = makeHub(agents, {
observers,
settings: Settings.isolated({
modelRoles: { rapid: "openai/gpt-4o" },
modelTags: { rapid: { name: "Quick", color: "warning" } },
}),
});
try {
const roleBlock = renderedRosterEntry(hub, "RoleAgent", 160);
expect(roleBlock).toContain("Quick");
expect(roleBlock).toContain("gpt-4o");
expect(roleBlock.indexOf("Quick")).toBeLessThan(roleBlock.indexOf("gpt-4o"));
const explicitBlock = renderedRosterEntry(hub, "ExplicitAgent", 160);
expect(explicitBlock).toContain("gpt-4o");
expect(explicitBlock).not.toContain("Quick");
} finally {
hub.dispose();
}
});
it("switches between inline Flat and By parent projections with selection preserved", () => {
vi.useFakeTimers();
geometry = stubStdoutGeometry(120);
const agents = new AgentRegistry();
setSystemTime(1_000);
agents.register({ id: "Parent", displayName: "Parent", kind: "sub", parentId: "Main", session: null });
setSystemTime(2_000);
agents.register({ id: "Peer", displayName: "Peer", kind: "sub", parentId: "Main", session: null });
setSystemTime(3_000);
agents.register({ id: "Child", displayName: "Child", kind: "sub", parentId: "Parent", session: null });
const hub = makeHub(agents);
try {
expect(renderedAgentIds(hub)).toEqual(["Child", "Peer", "Parent"]);
expect(selectedAgentId(hub)).toBe("Child");
const flat = Bun.stripANSI(hub.render(120).join("\n"));
expect(flat).toContain("Flat");
expect(flat).toContain("By parent");
hub.setHoverIndex(0);
expect(renderedRosterHeaderLineRaw(hub, "Child", 120)).toContain(theme.getBgAnsi("selectedBg"));
hub.handleInput("t");
const byParentIds = renderedAgentIds(hub);
expect(byParentIds).toEqual(["Parent", "Child", "Peer"]);
expect(selectedAgentId(hub)).toBe("Child");
const byParent = Bun.stripANSI(hub.render(120).join("\n"));
expect(byParent).toContain("Flat");
expect(byParent).toContain("By parent");
expect(byParentIds.indexOf("Parent")).toBeLessThan(byParentIds.indexOf("Child"));
expect(renderedRosterHeaderLineRaw(hub, "Parent", 120)).not.toContain(theme.getBgAnsi("selectedBg"));
expect(renderedRosterHeaderLineRaw(hub, "Child", 120)).not.toContain(theme.getBgAnsi("selectedBg"));
hub.handleInput("t");
expect(renderedAgentIds(hub)).toEqual(["Child", "Peer", "Parent"]);
expect(selectedAgentId(hub)).toBe("Child");
} finally {
hub.dispose();
vi.useRealTimers();
setSystemTime();
}
});
it("renders parent lineage with bash-style tree connectors", () => {
geometry = stubStdoutGeometry(120);
geometry.setRows(32);
const agents = new AgentRegistry();
agents.register({ id: "Parent", displayName: "Parent", kind: "sub", parentId: "Main", session: null });
agents.register({ id: "First", displayName: "First", kind: "sub", parentId: "Parent", session: null });
agents.register({ id: "Grandchild", displayName: "Grandchild", kind: "sub", parentId: "First", session: null });
agents.register({ id: "Last", displayName: "Last", kind: "sub", parentId: "Parent", session: null });
const hub = makeHub(agents);
try {
hub.handleInput("t");
expect(Bun.stripANSI(renderedRosterHeaderLineRaw(hub, "First", 120))).toContain("├── First");
expect(Bun.stripANSI(renderedRosterHeaderLineRaw(hub, "Grandchild", 120))).toContain("│ └── Grandchild");
expect(Bun.stripANSI(renderedRosterHeaderLineRaw(hub, "Last", 120))).toContain("└── Last");
} finally {
hub.dispose();
}
});
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(12);
const agents = new AgentRegistry();
agents.register({ id: "NarrowAgent", displayName: "Narrow Agent", kind: "sub", session: null });
const observers = new SessionObserverRegistry();
vi.spyOn(observers, "getSessions").mockReturnValue([
{
id: "NarrowAgent",
kind: "subagent",
label: "Narrow Agent",
status: "active",
lastUpdate: Date.now(),
progress: {
id: "NarrowAgent",
status: "running",
task: "Inspect responsive behavior",
recentTools: [],
recentOutput: [],
toolCount: 3,
requests: 2,
tokens: 900,
cost: 0,
durationMs: 2_000,
} as never,
},
]);
const hub = makeHub(agents, { observers });
try {
const roster = Bun.stripANSI(hub.render(80).join("\n"));
expect(roster).toContain("Tab:details");
expect(roster).not.toContain("Registered ");
hub.handleInput("\t");
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 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");
expect(Bun.stripANSI(hub.render(80).join("\n"))).toContain("Roster");
} finally {
hub.dispose();
}
});
});
@@ -10,9 +10,11 @@ import {
parseModelString,
pickDefaultAvailableModel,
resolveAgentModelPatterns,
resolveAgentModelSource,
resolveAgentPrewalkPattern,
resolveAllowedModels,
resolveCliModel,
resolveExplicitModelRole,
resolveModelFromString,
resolveModelOverride,
resolveModelRoleValue,
@@ -844,6 +846,43 @@ describe("resolveAgentPrewalkPattern", () => {
});
});
describe("resolveAgentModelPatterns", () => {
test("selects the first non-empty source and skips aliases with no patterns", () => {
const settings = Settings.isolated({
modelRoles: {
empty: "",
override: "openai/gpt-4o",
definition: "anthropic/claude-sonnet-4-5",
},
});
const emptyRequest = {
requestModel: "",
settingsOverride: "@override",
agentModel: ["@definition"],
settings,
};
expect(resolveAgentModelPatterns(emptyRequest)).toEqual(["openai/gpt-4o"]);
expect(resolveAgentModelSource(emptyRequest)).toBe("@override");
const emptyAlias = {
requestModel: "@empty",
settingsOverride: ",,",
agentModel: ["@definition"],
settings,
};
expect(resolveAgentModelPatterns(emptyAlias)).toEqual(["anthropic/claude-sonnet-4-5"]);
expect(resolveAgentModelSource(emptyAlias)).toEqual(["@definition"]);
const concreteRequest = {
requestModel: "openai/gpt-4o",
settingsOverride: "@override",
agentModel: ["@definition"],
settings,
};
expect(resolveAgentModelSource(concreteRequest)).toBe("openai/gpt-4o");
expect(resolveExplicitModelRole(resolveAgentModelSource(concreteRequest), settings)).toBeUndefined();
});
test("falls back to the active session model when @task is unset", () => {
const settings = Settings.isolated({
modelRoles: { default: "anthropic/claude-sonnet-4-5" },
@@ -1699,6 +1738,29 @@ describe("resolveModelFromString", () => {
});
});
describe("resolveExplicitModelRole", () => {
test("extracts built-in, custom, legacy, default, and thinking-suffixed aliases before expansion", () => {
const settings = Settings.isolated({
modelRoles: {
reviewer: "openai/gpt-4o",
},
});
expect(resolveExplicitModelRole("@task", settings)).toBe("task");
expect(resolveExplicitModelRole("pi/reviewer:high", settings)).toBe("reviewer");
expect(resolveExplicitModelRole("@reviewer:xhigh", settings)).toBe("reviewer");
expect(resolveExplicitModelRole("*:low", settings)).toBe("default");
});
test("does not infer a role from an explicit model selector", () => {
const settings = Settings.isolated({ modelRoles: { reviewer: "openai/gpt-4o" } });
expect(resolveExplicitModelRole("openai/gpt-4o", settings)).toBeUndefined();
expect(resolveExplicitModelRole("openai/gpt-4o:high", settings)).toBeUndefined();
expect(resolveExplicitModelRole("openai/gpt-4o:max", settings)).toBeUndefined();
expect(resolveExplicitModelRole(["openai/gpt-4o", "@reviewer:high"], settings)).toBe("reviewer");
});
});
describe("expandRoleAlias", () => {
test("expands @vision to configured vision role", () => {
const settings = Settings.isolated();
@@ -614,10 +614,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" });
@@ -365,4 +365,25 @@ describe("runSubprocess parent-discovery pass-through (issue #2190)", () => {
const forwarded = spy.mock.calls[0]?.[0];
expect(forwarded?.thinkingLevel).toBe(ThinkingLevel.Low);
});
it("persists an explicit role from a caller model override", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
if (!model) throw new Error("Expected claude-sonnet-4-5 model to exist");
const settings = Settings.isolated({
modelRoles: { reviewer: `${model.provider}/${model.id}` },
});
const session = yieldEmittingSession();
const initSpy = vi.spyOn(session.sessionManager, "appendSessionInit");
vi.spyOn(sdkModule, "createAgentSession").mockResolvedValue(createSessionResult(session));
const result = await runSubprocess({
...baseOptions,
id: "subagent-model-override-role",
modelOverride: "@reviewer",
settings,
modelRegistry: createModelRegistry(model),
});
expect(result.exitCode).toBe(0);
expect(initSpy).toHaveBeenCalledWith(expect.objectContaining({ modelRole: "reviewer" }));
});
});
@@ -2,6 +2,7 @@ import { afterEach, describe, expect, it, vi } from "bun:test";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
import * as executorModule from "@oh-my-pi/pi-coding-agent/task/executor";
import {
applyEligibleNestedPatches,
@@ -73,6 +74,7 @@ async function seedFooRepo(finalContent: string): Promise<{ repoRoot: string; pa
describe("runIsolatedSubprocess", () => {
afterEach(async () => {
vi.restoreAllMocks();
AgentRegistry.resetGlobalForTests();
await Promise.all(tempRoots.splice(0).map(tempRoot => fs.rm(tempRoot, { force: true, recursive: true })));
});
@@ -107,6 +109,13 @@ describe("runIsolatedSubprocess", () => {
nestedPatches: [],
});
const cleanupSpy = vi.spyOn(worktreeModule, "cleanupIsolation").mockResolvedValue();
AgentRegistry.global().register({
id: "PreserveBranchFailure",
displayName: "PreserveBranchFailure",
kind: "sub",
session: null,
status: "parked",
});
const deleteSpy = vi.spyOn(gitModule.branch, "tryDelete").mockResolvedValue(true);
const outcome = await runIsolatedSubprocess({
@@ -138,6 +147,7 @@ describe("runIsolatedSubprocess", () => {
expect(captureSpy).toHaveBeenCalledWith(isolationDir, baseline);
expect(deleteSpy).toHaveBeenCalledWith(repoRoot, "omp/task/PreserveBranchFailure");
expect(cleanupSpy).toHaveBeenCalledTimes(1);
expect(AgentRegistry.global().get("PreserveBranchFailure")?.history?.patchPath).toBe(patchPath);
});
it("keeps an isolated worktree until deferred child cleanup settles", async () => {
@@ -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();
@@ -37,6 +37,7 @@ function session(
maxDepth?: number;
isolationMode?: "none" | "worktree";
isolationApply?: boolean;
modelRoles?: Record<string, string>;
} = {},
): ToolSession {
return {
@@ -47,6 +48,7 @@ function session(
"task.maxRecursionDepth": options.maxDepth ?? 2,
"task.isolation.mode": options.isolationMode ?? "none",
"task.enableLsp": true,
...(options.modelRoles ? { modelRoles: options.modelRoles } : {}),
...(options.isolationApply !== undefined ? { "task.isolation.apply": options.isolationApply } : {}),
}),
getSessionFile: () => null,
@@ -165,6 +167,112 @@ describe("structured subagent primitive", () => {
).rejects.toThrow("isolation, apply, and merge controls are unavailable in plan mode");
expect(discover).not.toHaveBeenCalled();
});
it("propagates a custom thinking-suffixed role alias through policy, dispatch, and settlement", async () => {
const customAgent = { ...AGENT, model: ["@reviewer:high"] };
mockDiscovery(customAgent);
const childSession = session({ modelRoles: { reviewer: "openai/gpt-4o" } });
const dispatched: executorModule.ExecutorOptions[] = [];
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
dispatched.push(options);
return { ...result(), modelRole: options.modelRole };
});
const settled = await runStructuredSubagent(
request({ session: childSession, agent: "worker", retainArtifacts: true }),
);
expect(settled.policy.modelRole).toBe("reviewer");
expect(dispatched[0]?.modelRole).toBe("reviewer");
expect(settled.result.modelRole).toBe("reviewer");
await fs.rm(settled.artifactsDir, { recursive: true, force: true });
});
it("derives modelRole from the raw selector source in request, override, definition order", async () => {
const customAgent = { ...AGENT, model: ["@definition"] };
mockDiscovery(customAgent);
const roleSession = session({
modelRoles: {
request: "openai/gpt-4o",
override: "openai/gpt-4o",
definition: "openai/gpt-4o",
},
});
roleSession.settings.override("task.agentModelOverrides", { worker: "@override" });
const requestPolicy = await resolveEffectiveSubagentPolicy(request({ session: roleSession, model: "@request" }));
expect(requestPolicy.modelRole).toBe("request");
const overridePolicy = await resolveEffectiveSubagentPolicy(request({ session: roleSession }));
expect(overridePolicy.modelRole).toBe("override");
const concreteOverrideSession = session({
modelRoles: {
override: "openai/gpt-4o",
definition: "openai/gpt-4o",
},
});
concreteOverrideSession.settings.override("task.agentModelOverrides", { worker: "openai/gpt-4o" });
const concreteOverridePolicy = await resolveEffectiveSubagentPolicy(
request({ session: concreteOverrideSession }),
);
expect(concreteOverridePolicy.modelRole).toBeUndefined();
const definitionPolicy = await resolveEffectiveSubagentPolicy(
request({ session: session({ modelRoles: { definition: "openai/gpt-4o" } }) }),
);
expect(definitionPolicy.modelRole).toBe("definition");
});
it("falls through an empty request selector to the agent definition role", async () => {
const customAgent = { ...AGENT, model: ["@definition"] };
mockDiscovery(customAgent);
const childSession = session({ modelRoles: { definition: "openai/gpt-4o" } });
const policy = await resolveEffectiveSubagentPolicy(request({ session: childSession, model: "" }));
expect(policy.modelRole).toBe("definition");
expect(policy.modelOverride).toEqual(["openai/gpt-4o"]);
});
it("falls through an empty configured override to the agent definition role", async () => {
const customAgent = { ...AGENT, model: ["@definition"] };
mockDiscovery(customAgent);
const childSession = session({ modelRoles: { definition: "openai/gpt-4o" } });
childSession.settings.override("task.agentModelOverrides", { worker: "" });
const policy = await resolveEffectiveSubagentPolicy(request({ session: childSession }));
expect(policy.modelRole).toBe("definition");
expect(policy.modelOverride).toEqual(["openai/gpt-4o"]);
});
it("falls through a configured alias that expands to no patterns", async () => {
const customAgent = { ...AGENT, model: ["@definition"] };
mockDiscovery(customAgent);
const childSession = session({ modelRoles: { empty: "", definition: "openai/gpt-4o" } });
childSession.settings.override("task.agentModelOverrides", { worker: "@empty" });
const policy = await resolveEffectiveSubagentPolicy(request({ session: childSession }));
expect(policy.modelRole).toBe("definition");
expect(policy.modelOverride).toEqual(["openai/gpt-4o"]);
});
it("does not assign a role when a child uses an explicit model selector", async () => {
mockDiscovery();
const childSession = session({ modelRoles: { reviewer: "openai/gpt-4o" } });
const dispatched: executorModule.ExecutorOptions[] = [];
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
dispatched.push(options);
return result();
});
const settled = await runStructuredSubagent(
request({ session: childSession, model: "openai/gpt-4o", retainArtifacts: true }),
);
expect(settled.policy.modelRole).toBeUndefined();
expect(dispatched[0]?.modelRole).toBeUndefined();
expect(settled.result.modelRole).toBeUndefined();
await fs.rm(settled.artifactsDir, { recursive: true, force: true });
});
it("leases temporary artifacts for a retained invocation and registers them for agent URLs", async () => {
mockDiscovery();
@@ -0,0 +1,13 @@
import { describe, expect, it } from "bun:test";
import { nativeVersionFromExports } from "./native-version";
describe("native addon release sentinel", () => {
it("normalizes the unique version sentinel", () => {
expect(nativeVersionFromExports(["load", "__piNativesV17_2_6", "other"])).toBe("17.2.6");
});
it("rejects missing or ambiguous sentinels", () => {
expect(nativeVersionFromExports(["load"])).toBeUndefined();
expect(nativeVersionFromExports(["__piNativesV17_2_6", "__piNativesV17_2_7"])).toBeUndefined();
});
});
+22
View File
@@ -0,0 +1,22 @@
import { createRequire } from "node:module";
const VERSION_SENTINEL_RE = /^__piNativesV(\d+)_(\d+)_(\d+)$/;
/** Return the sole release version advertised by a native addon's exports. */
export function nativeVersionFromExports(exports: readonly string[]): string | undefined {
const versions = exports
.map(name => VERSION_SENTINEL_RE.exec(name))
.filter((match): match is RegExpExecArray => match !== null)
.map(match => `${match[1]}.${match[2]}.${match[3]}`);
return versions.length === 1 ? versions[0] : undefined;
}
if (import.meta.main) {
const addonPath = process.argv[2];
if (!addonPath) throw new Error("Usage: bun scripts/install-tests/native-version.ts <addon-path>");
const require = createRequire(import.meta.url);
const bindings = require(addonPath) as Record<string, unknown>;
const version = nativeVersionFromExports(Object.keys(bindings));
if (!version) throw new Error(`Native addon has no unique release version sentinel: ${addonPath}`);
process.stdout.write(version);
}
+41 -1
View File
@@ -7,7 +7,15 @@ WORK_DIR="$(mktemp -d)"
TMP_WORK_DIR="$WORK_DIR/tmp"
mkdir -p "$TMP_WORK_DIR"
export TMPDIR="$TMP_WORK_DIR"
trap 'rm -rf "$WORK_DIR"' EXIT
NATIVES_PACKAGE="$ROOT_DIR/packages/natives/package.json"
NATIVES_PACKAGE_INITIAL="$WORK_DIR/natives-package.initial.json"
cp "$NATIVES_PACKAGE" "$NATIVES_PACKAGE_INITIAL"
restore_workspace() {
cp "$NATIVES_PACKAGE_INITIAL" "$NATIVES_PACKAGE"
rm -rf "$WORK_DIR"
}
trap restore_workspace EXIT
section() {
echo ""
@@ -42,10 +50,42 @@ find_tarball() {
echo "${matches[0]}"
}
align_native_manifest() {
local addon_version=""
local addon
local candidate_version
local candidates=()
shopt -s nullglob
candidates=("$ROOT_DIR"/packages/natives/native/pi_natives.*.node)
shopt -u nullglob
if [ "${#candidates[@]}" -eq 0 ]; then
echo "No native addon found for install smoke" >&2
exit 1
fi
for addon in "${candidates[@]}"; do
candidate_version="$(bun "$ROOT_DIR/scripts/install-tests/native-version.ts" "$addon")" || exit 1
if [ -z "$addon_version" ]; then
addon_version="$candidate_version"
elif [ "$addon_version" != "$candidate_version" ]; then
echo "Native addon version mismatch: $addon_version vs $candidate_version ($addon)" >&2
exit 1
fi
done
local declared_version
declared_version="$(jq -r '.version' "$NATIVES_PACKAGE")"
if [ "$declared_version" = "$addon_version" ]; then return; fi
echo "Aligning install smoke native manifest $declared_version → $addon_version"
jq --arg version "$addon_version" '.version = $version' "$NATIVES_PACKAGE" > "$WORK_DIR/natives-package.aligned.json"
mv "$WORK_DIR/natives-package.aligned.json" "$NATIVES_PACKAGE"
}
section "Binary install smoke"
if [ "${OMP_INSTALL_TEST_SKIP_NATIVE_BUILD:-0}" != "1" ]; then
bun --cwd=packages/natives run build
fi
align_native_manifest
bun --cwd=packages/coding-agent run build
BINARY_DIR="$WORK_DIR/binary-bin"