import { createHash, randomUUID } from "node:crypto"; import { SUBAGENT_OBSERVE_EVENT, SUBAGENT_RPC_MAX_BYTES, type SubagentObserveRequest, type SubagentRpcNode, type SubagentRpcProjection, type SubagentRpcLifecycle, } from "@selesai/code"; import type { SubagentState } from "../shared/types.ts"; import { sanitizeDisplayText, truncateDisplayText } from "../shared/display-text.ts"; interface Events { on(name: string, handler: (payload: unknown) => void): (() => void) | void } interface WaitingChild { runId: string; childIndex: number; expectsReply: boolean } const statuses = new Set(["queued", "pending", "running", "complete", "completed", "failed", "partial", "paused", "stopped", "rejected", "detached"]); const record = (v: unknown): v is Record => Boolean(v && typeof v === "object" && !Array.isArray(v)); const count = (v: unknown): v is number => typeof v === "number" && Number.isSafeInteger(v) && v >= 0; function text(v: unknown, max: number): string | undefined { if (typeof v !== "string") return undefined; return truncateDisplayText(sanitizeDisplayText(v.slice(0, 4096)), max) || undefined; } function identity(v: unknown): string | undefined { return typeof v === "string" && v.length > 0 && v.length <= 128 && !/[\\/\r\n]/.test(v) ? v : undefined; } function tokens(v: unknown) { if (!record(v) || !count(v.input) || !count(v.output) || !count(v.total)) return undefined; return { input: v.input, output: v.output, total: v.total, ...(count(v.window) ? { window: v.window } : {}), ...(count(v.windowPeak) ? { windowPeak: v.windowPeak } : {}) }; } /** Simple adapter over runtime-owned memory. No artifact reads, executors, or new run state. */ export function registerSubagentObservability(options: { events: Events; state: SubagentState; getHostSessionId(): string | null; getWaitingChildren(): Iterable; getFleet(): SubagentRpcProjection["fleet"]; }): { dispose(): void } { const salt = randomUUID(); const key = (internal: string) => `node-${createHash("sha256").update(salt).update(internal).digest("hex").slice(0, 24)}`; const subscribers = new Set<(projection: SubagentRpcProjection | null) => void>(); let timer: ReturnType | undefined; let disposed = false; let boundSession: string | null = null; const snapshot = (): SubagentRpcProjection => { if (disposed) throw new Error("Subagent observation source disposed"); const { state } = options; const sessionId = state.currentSessionId; const waiting = new Set([...options.getWaitingChildren()].filter(r => r.expectsReply).map(r => `${r.runId}:${r.childIndex}`)); let nodes = 0; const representedNativeIds = new Set(); const omitted = { runs: 0, children: 0 }; const project = (raw: unknown, internal: string, depth: number, runId?: string, childIndex?: number, retained = false): SubagentRpcNode | undefined => { if (!record(raw) || !statuses.has((raw.status ?? raw.state) as string) || nodes >= 64 || depth > 4) return undefined; nodes++; const node: SubagentRpcNode = { key: key(internal), status: (raw.status ?? raw.state) as SubagentRpcLifecycle, children: [], omittedChildren: 0 }; const nativeRunId = identity(raw.runId ?? raw.asyncId ?? raw.id); if (nativeRunId) { node.runId = nativeRunId; representedNativeIds.add(nativeRunId); } const workflowId = identity(raw.parentWorkflowRunId ?? (raw.mode === "workflow" ? nativeRunId : undefined)); if (workflowId) node.workflowId = workflowId; const agent = text(raw.agent ?? raw.currentAgent, 96); if (agent) node.agent = agent; if (waiting.has(`${runId ?? nativeRunId}:${childIndex ?? raw.index ?? raw.currentIndex ?? 0}`)) node.activity = "waiting_supervisor"; else if ((raw.activityState ?? raw.currentActivityState) === "needs_attention") node.activity = "needs_attention"; for (const [field, value] of [["startedAt", raw.startedAt], ["updatedAt", raw.updatedAt ?? raw.lastUpdate], ["endedAt", raw.endedAt]] as const) if (count(value)) node[field] = value; const progress = text(Array.isArray(raw.recentOutput) ? raw.recentOutput.slice(-1)[0] : raw.description, 256); if (progress) node.progress = progress; const currentTool = text(raw.currentTool, 96); if (currentTool) node.currentTool = currentTool; const usage = tokens(raw.totalTokens ?? raw.tokens) ?? (count(raw.inputTokens) && count(raw.outputTokens) && count(raw.tokens) ? tokens({ input: raw.inputTokens, output: raw.outputTokens, total: raw.tokens, window: raw.window, windowPeak: raw.windowPeak }) : undefined); if (usage) node.tokens = usage; const error = text(raw.error, 256); if (error) node.error = error; const reason = text(raw.reason ?? raw.detachedReason, 128); if (reason) node.reason = reason; const execution = raw.execution; if (record(execution) && statuses.has(execution.status as string)) node.outcome = { status: execution.status as SubagentRpcLifecycle, ...(typeof execution.success === "boolean" ? { success: execution.success } : {}), ...(typeof execution.exitCode === "number" && Number.isSafeInteger(execution.exitCode) ? { exitCode: execution.exitCode } : {}), }; const proof = raw.processTerminal; if (record(proof) && ["pending", "observed", "unknown", "not-started"].includes(proof.state as string)) node.processTerminal = { state: proof.state as NonNullable["state"], ...(count(proof.observedAt) ? { observedAt: proof.observedAt } : {}), ...(text(proof.reason, 128) ? { reason: text(proof.reason, 128) } : {}), }; const lists = [raw.activeChildren instanceof Map ? [...raw.activeChildren.values()].map(c => ({ ...c, status: "running" })) : [], Array.isArray(raw.steps) ? raw.steps : [], Array.isArray(raw.children) ? raw.children : [], Array.isArray(raw.nestedChildren) ? raw.nestedChildren : []]; let offset = 0; for (const [listIndex, list] of lists.entries()) for (const child of list) { const index = record(child) && count(child.index) ? child.index : offset; const childId = record(child) ? identity(child.runId ?? child.id ?? child.workflowKey) : undefined; if (retained && childId && representedNativeIds.has(childId)) { offset++; continue; } const childNode = node.children.length < 16 ? project(child, `${internal}:${listIndex}:${childId ?? index}`, depth + 1, nativeRunId ?? runId, index, retained) : undefined; offset++; if (childNode) node.children.push(childNode); else { node.omittedChildren++; omitted.children++; } } return node; }; const runs: SubagentRpcNode[] = []; const add = (raw: unknown, internal: string, retained = false) => { const node = runs.length < 32 ? project(raw, internal, 0, undefined, undefined, retained) : undefined; if (node) runs.push(node); else omitted.runs++; }; if (sessionId) { for (const control of state.foregroundControls.values()) if (control.sessionId === sessionId) add({ ...control, status: "running" }, `foreground:${control.runId}`); for (const run of state.foregroundRuns?.values() ?? []) if (run.sessionId === sessionId && !state.foregroundControls.has(run.runId)) { // Retained per-child lifecycle is authoritative; the run has no aggregate verdict. for (const child of run.children) add({ ...child, runId: run.runId }, `foreground:${run.runId}:${child.index}`); } const seen = new Set(); for (const jobs of [state.asyncJobs, state.fleetJobs]) for (const job of jobs?.values() ?? []) { if (job.sessionId !== sessionId || seen.has(job.asyncId)) continue; seen.add(job.asyncId); add(job, `async:${job.asyncId}`); } } for (const [rootRunId, cached] of state.retainedForegroundNestedChildren ?? []) { if (cached.sessionId !== sessionId || !state.retainedForegroundNestedRoutes?.has(rootRunId)) continue; for (const descendant of cached.children) { if (representedNativeIds.has(descendant.id)) continue; // These are actual nested runs; the settled foreground root has no invented aggregate state. add({ ...descendant, parentWorkflowRunId: cached.workflowId }, `retained:${descendant.id}`, true); } } const projection = { fleet: options.getFleet(), runs, omitted }; while (Buffer.byteLength(JSON.stringify(projection), "utf8") > SUBAGENT_RPC_MAX_BYTES - 1024 && runs.length) { const removed = runs.pop()!; const removeOmissions = (n: SubagentRpcNode) => { omitted.children -= n.omittedChildren; n.children.forEach(removeOmissions); }; removeOmissions(removed); omitted.runs++; } return projection; }; let previous = ""; const unsubscribe = options.events.on(SUBAGENT_OBSERVE_EVENT, raw => { const request = raw as SubagentObserveRequest; if (disposed || request?.version !== 1 || request.sessionId !== options.getHostSessionId() || typeof request.accept !== "function") return; boundSession = request.sessionId; request.accept({ snapshot, subscribe(listener) { subscribers.add(listener); if (!timer) { previous = JSON.stringify(snapshot()); timer = setInterval(() => { if (options.getHostSessionId() !== boundSession) return; const projection = snapshot(); const current = JSON.stringify(projection); if (current === previous) return; previous = current; for (const subscriber of subscribers) subscriber(projection); }, 250); timer.unref?.(); } return () => { subscribers.delete(listener); if (!subscribers.size && timer) { clearInterval(timer); timer = undefined; } }; } }); }); return { dispose() { if (disposed) return; disposed = true; if (typeof unsubscribe === "function") unsubscribe(); if (timer) clearInterval(timer); for (const subscriber of subscribers) subscriber(null); subscribers.clear(); } }; }