{"version":3,"file":"subagent-wait.d.ts","sourceRoot":"","sources":["../../../../src/runs/background/subagent-wait.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;GAsCG;AAGH,OAAO,KAAK,EAAE,eAAe,EAAE,MAAM,yBAAyB,CAAC;AAC/D,OAAO,EACN,KAAK,sBAAsB,EAI3B,MAAM,8BAA8B,CAAC;AAEtC,OAAO,EACN,KAAK,OAAO,EASZ,KAAK,aAAa,EAClB,MAAM,uBAAuB,CAAC;AAG/B,OAAO,EAAE,KAAK,sBAAsB,EAAE,qBAAqB,EAAE,qBAAqB,EAAE,MAAM,kBAAkB,CAAC;AAS7G,MAAM,WAAW,kBAAkB;IAClC,uGAAuG;IACvG,EAAE,CAAC,EAAE,MAAM,CAAC;IACZ,qFAAqF;IACrF,WAAW,CAAC,EAAE,OAAO,CAAC;IACtB;;;;;OAKG;IACH,GAAG,CAAC,EAAE,OAAO,CAAC;IACd,oEAAoE;IACpE,SAAS,CAAC,EAAE,MAAM,CAAC;CACnB;AAED,wEAAwE;AACxE,MAAM,WAAW,YAAY;IAC5B,EAAE,CAAC,OAAO,EAAE,MAAM,EAAE,OAAO,EAAE,CAAC,IAAI,EAAE,OAAO,KAAK,IAAI,GAAG,MAAM,IAAI,CAAC;CAClE;AAED,MAAM,WAAW,gBAAgB;IAChC,KAAK,EAAE,aAAa,CAAC;IACrB,0DAA0D;IAC1D,QAAQ,CAAC,EAAE,CAAC,MAAM,EAAE,eAAe,CAAC,OAAO,CAAC,KAAK,IAAI,CAAC;IACtD,YAAY,CAAC,EAAE,MAAM,CAAC;IACtB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,IAAI,CAAC,EAAE,CAAC,GAAG,EAAE,MAAM,EAAE,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,GAAG,CAAC,KAAK,OAAO,CAAC;IAC7D,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,cAAc,CAAC,EAAE,MAAM,CAAC;IACxB,kFAAkF;IAClF,OAAO,CAAC,EAAE,OAAO,CAAC;IAClB,kCAAkC;IAClC,KAAK,CAAC,EAAE,CAAC,EAAE,EAAE,MAAM,EAAE,MAAM,CAAC,EAAE,WAAW,KAAK,OAAO,CAAC,IAAI,CAAC,CAAC;IAC5D,qEAAqE;IACrE,eAAe,CAAC,EAAE,OAAO,CAAC;IAC1B,iFAAiF;IACjF,gBAAgB,CAAC,EAAE,OAAO,CAAC;IAC3B,wFAAwF;IACxF,SAAS,CAAC,EAAE,CAAC,KAAK,EAAE;QACnB,UAAU,EAAE,OAAO,GAAG,YAAY,CAAC;QACnC,KAAK,EAAE,MAAM,CAAC;QACd,WAAW,EAAE,MAAM,CAAC;QACpB,SAAS,EAAE,MAAM,CAAC;KAClB,KAAK;QAAE,KAAK,EAAE,MAAM,CAAC;QAAC,SAAS,EAAE,MAAM,CAAA;KAAE,CAAC;IAC3C,qEAAqE;IACrE,cAAc,CAAC,EAAE;QAChB,QAAQ,CAAC,SAAS,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,GAAG,sBAAsB,CAAC;QACnE,YAAY,IAAI,SAAS,MAAM,EAAE,CAAC;KAClC,CAAC;IACF;;;;;OAKG;IACH,MAAM,CAAC,EAAE,YAAY,CAAC;CACtB;AA8YD;;;;GAIG;AACH,wBAAsB,gBAAgB,CACrC,MAAM,EAAE,kBAAkB,EAC1B,MAAM,EAAE,WAAW,GAAG,SAAS,EAC/B,IAAI,EAAE,gBAAgB,GACpB,OAAO,CAAC,eAAe,CAAC,OAAO,CAAC,CAAC,CA+MnC","sourcesContent":["/**\n * `subagent_wait` tool: block the current turn until outstanding async runs\n * or a named remembered detached foreground run finishes.\n *\n * Background subagent runs are detached. In an interactive session the parent\n * can end its turn and Pi will wake it with a completion notification. That\n * does not work when the parent is a skill that must run to completion, and it\n * cannot work at all non-interactively (`pi -p ...`), where the run is a single\n * turn: once the turn ends there is nothing left to receive the notification.\n *\n * `subagent_wait` closes that gap. It keeps the turn alive until a tracked async\n * run for this session reaches a terminal state (complete / failed / paused),\n * the caller-supplied timeout elapses, or the turn is aborted. Because it awaits\n * inside the turn, the completion the model was told to wait for is actually\n * observed before the tool returns.\n *\n * By default `subagent_wait` returns as soon as ONE run finishes, so a fleet\n * manager can use it in a rolling-replacement loop: launch N workers, wait for\n * the next one to finish, spawn its replacement, then call `subagent_wait`\n * again — keeping N in flight instead of draining to zero between batches.\n * Pass `all: true` to block until every tracked async run is terminal, or `id`\n * to block on one specific async or remembered detached foreground run.\n *\n * `subagent_wait` also returns when a run needs attention — not just on\n * completion. A child that goes idle or blocks for a decision surfaces\n * `needs_attention` (the same signal Pi shows as a control notice and,\n * interactively, wakes the parent with). Since `subagent_wait` is used exactly\n * where there is no next turn to receive that notice, it must break on it too,\n * or a stuck child would stall the loop until the timeout. Attention runs are\n * reported so the caller can inspect / nudge / resume / interrupt them.\n *\n * Wake mechanism: when given Pi's event bus (`deps.events`), `subagent_wait`\n * subscribes to the subagent completion/control channels and wakes the instant\n * any fires, rather than waiting out a fixed poll interval. A poll still runs\n * on the interval as a reconciliation fallback (crashed runners, missed\n * events), and the poll is the source of truth for what actually changed — the\n * event only ends the sleep early. With no bus, `subagent_wait` degrades to pure\n * polling.\n */\n\nimport * as fs from \"node:fs\";\nimport type { AgentToolResult } from \"@lpb-work/pi-agent-core\";\nimport {\n\ttype BackgroundWorkSnapshot,\n\tlistBackgroundWorkWakeChannels,\n\ttype RegisteredBackgroundWorkItem,\n\tsnapshotBackgroundWork,\n} from \"../../api/background-work.ts\";\nimport { formatDuration, shortenPath } from \"../../shared/formatters.ts\";\nimport {\n\ttype Details,\n\tDIRS,\n\ttype ForegroundResumeRun,\n\tINTERCOM_DETACH_REQUEST_EVENT,\n\tSUBAGENT_ASYNC_COMPLETE_EVENT,\n\tSUBAGENT_CONTROL_EVENT,\n\tSUBAGENT_CONTROL_INTERCOM_EVENT,\n\tSUBAGENT_FOREGROUND_COMPLETE_EVENT,\n\tSUBAGENT_RESULT_INTERCOM_EVENT,\n\ttype SubagentState,\n} from \"../../shared/types.ts\";\nimport { type AsyncRunSummary, formatAsyncRunList, listAsyncRuns } from \"./async-status.ts\";\n\nexport { type ResolvedWaitToolConfig, resolveWaitToolConfig, WAIT_TOOL_ENABLED_ENV } from \"./wait-config.ts\";\n\n/** States that mean a run is still in flight (not yet resolved). */\nconst ACTIVE_STATES: ReadonlyArray<AsyncRunSummary[\"state\"]> = [\"queued\", \"running\"];\n\nconst DEFAULT_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes\nconst MIN_POLL_INTERVAL_MS = 250;\nconst DEFAULT_POLL_INTERVAL_MS = 1000;\n\nexport interface SubagentWaitParams {\n\t/** Optional run id/prefix to wait for. When omitted, waits across every active run in this session. */\n\tid?: string;\n\t/** Arm a durable exact-run wake subscription and return immediately. Requires id. */\n\tnonBlocking?: boolean;\n\t/**\n\t * When true, block until EVERY active run in this session (or matching `id`)\n\t * is terminal. Default false: return as soon as the first run finishes, so a\n\t * fleet manager can spawn a replacement and wait again. Ignored when `id`\n\t * targets a single run.\n\t */\n\tall?: boolean;\n\t/** Give up after this many milliseconds. Defaults to 30 minutes. */\n\ttimeoutMs?: number;\n}\n\n/** Minimal event-bus surface wait subscribes to (matches pi.events). */\nexport interface WaitEventBus {\n\ton(channel: string, handler: (data: unknown) => void): () => void;\n}\n\nexport interface SubagentWaitDeps {\n\tstate: SubagentState;\n\t/** Stream live wait status into Pi's pending tool row. */\n\tonUpdate?: (result: AgentToolResult<Details>) => void;\n\tasyncDirRoot?: string;\n\tresultsDir?: string;\n\tkill?: (pid: number, signal?: NodeJS.Signals | 0) => boolean;\n\tnow?: () => number;\n\tpollIntervalMs?: number;\n\t/** False makes the tool return immediately without blocking active async runs. */\n\tenabled?: boolean;\n\t/** Injectable sleep for tests. */\n\tsleep?: (ms: number, signal?: AbortSignal) => Promise<void>;\n\t/** Internal auto-drain mode waits through needs-attention states. */\n\tstopOnAttention?: boolean;\n\t/** Internal auto-drain mode surfaces failed terminal subagent runs as errors. */\n\tfailOnFailedRuns?: boolean;\n\t/** Arm a durable exact-target wait subscription in a long-lived interactive runtime. */\n\tsubscribe?: (input: {\n\t\ttargetKind: \"async\" | \"foreground\";\n\t\trunId: string;\n\t\trequestedId: string;\n\t\ttimeoutMs: number;\n\t}) => { token: string; expiresAt: number };\n\t/** Injectable provider protocol surfaces for deterministic tests. */\n\tbackgroundWork?: {\n\t\tsnapshot(sessionId: string, nowMs: number): BackgroundWorkSnapshot;\n\t\twakeChannels(): readonly string[];\n\t};\n\t/**\n\t * Optional event bus (pi.events). When provided, wait wakes immediately on a\n\t * subagent completion/control event instead of waiting out the poll interval;\n\t * the poll then remains as a reconciliation fallback (crashed runners, missed\n\t * events). Omit in tests that want pure poll behavior.\n\t */\n\tevents?: WaitEventBus;\n}\n\n/** Bus channels that indicate a run changed state or needs attention. */\nconst WAKE_CHANNELS = [\n\tINTERCOM_DETACH_REQUEST_EVENT,\n\tSUBAGENT_ASYNC_COMPLETE_EVENT,\n\tSUBAGENT_FOREGROUND_COMPLETE_EVENT,\n\tSUBAGENT_CONTROL_EVENT,\n\tSUBAGENT_CONTROL_INTERCOM_EVENT,\n\tSUBAGENT_RESULT_INTERCOM_EVENT,\n];\n\nfunction defaultSleep(ms: number, signal?: AbortSignal): Promise<void> {\n\treturn new Promise((resolve) => {\n\t\tif (signal?.aborted) {\n\t\t\tresolve();\n\t\t\treturn;\n\t\t}\n\t\tconst timer = setTimeout(() => {\n\t\t\tsignal?.removeEventListener(\"abort\", onAbort);\n\t\t\tresolve();\n\t\t}, ms);\n\t\tconst onAbort = () => {\n\t\t\tclearTimeout(timer);\n\t\t\tresolve();\n\t\t};\n\t\tsignal?.addEventListener(\"abort\", onAbort, { once: true });\n\t});\n}\n\n/**\n * Sleep up to `ms`, but wake early if a subagent event fires on the bus (or the\n * turn aborts). Returns when the first of those happens. With no bus this is a\n * plain sleep, so the poll interval alone drives progress.\n */\nfunction waitForWake(ms: number, signal: AbortSignal | undefined, deps: SubagentWaitDeps): Promise<void> {\n\tconst sleep = deps.sleep ?? defaultSleep;\n\tconst events = deps.events;\n\tif (!events) return sleep(ms, signal);\n\tconst providerChannels = deps.backgroundWork?.wakeChannels() ?? listBackgroundWorkWakeChannels();\n\treturn new Promise((resolve, reject) => {\n\t\tlet settled = false;\n\t\tconst unsubs: Array<() => void> = [];\n\t\tconst wakeController = new AbortController();\n\t\tconst done = () => {\n\t\t\tif (settled) return;\n\t\t\tsettled = true;\n\t\t\twakeController.abort();\n\t\t\tsignal?.removeEventListener(\"abort\", done);\n\t\t\tfor (const u of unsubs) {\n\t\t\t\ttry {\n\t\t\t\t\tu();\n\t\t\t\t} catch {\n\t\t\t\t\t/* best effort */\n\t\t\t\t}\n\t\t\t}\n\t\t\tresolve();\n\t\t};\n\t\tif (signal?.aborted) {\n\t\t\tdone();\n\t\t\treturn;\n\t\t}\n\t\tsignal?.addEventListener(\"abort\", done, { once: true });\n\t\ttry {\n\t\t\tfor (const channel of [...new Set([...WAKE_CHANNELS, ...providerChannels])]) {\n\t\t\t\tunsubs.push(events.on(channel, done));\n\t\t\t}\n\t\t} catch (error) {\n\t\t\tsignal?.removeEventListener(\"abort\", done);\n\t\t\tfor (const unsubscribe of unsubs) {\n\t\t\t\ttry {\n\t\t\t\t\tunsubscribe();\n\t\t\t\t} catch {\n\t\t\t\t\t/* best effort cleanup */\n\t\t\t\t}\n\t\t\t}\n\t\t\treject(error);\n\t\t\treturn;\n\t\t}\n\t\t// Poll-interval fallback so we still reconcile even if no event arrives.\n\t\t// The local signal cancels that fallback timer when an event wakes us first.\n\t\tvoid sleep(ms, wakeController.signal).then(done);\n\t});\n}\n\nfunction matchesId(run: AsyncRunSummary, id: string): boolean {\n\treturn run.id === id || run.id.startsWith(id);\n}\n\nfunction activeDetachedForegroundRuns(params: SubagentWaitParams, deps: SubagentWaitDeps): ForegroundResumeRun[] {\n\tif (!params.id || !deps.state.foregroundRuns) return [];\n\tconst sessionId = deps.state.currentSessionId;\n\tif (!sessionId) return [];\n\treturn [...deps.state.foregroundRuns.values()].filter(\n\t\t(run) =>\n\t\t\t(run.runId === params.id || run.runId.startsWith(params.id!)) &&\n\t\t\trun.sessionId === sessionId &&\n\t\t\trun.children.some((child) => child.status === \"detached\"),\n\t);\n}\n\nfunction summarizeForegroundChildren(run: ForegroundResumeRun, indices: Set<number>): string {\n\tconst counts = new Map<string, number>();\n\tfor (const child of run.children) {\n\t\tif (!indices.has(child.index) || child.status === \"detached\") continue;\n\t\tcounts.set(child.status, (counts.get(child.status) ?? 0) + 1);\n\t}\n\treturn [...counts.entries()].map(([status, count]) => `${count} ${status}`).join(\", \");\n}\n\nfunction foregroundChildrenNeedingAttention(run: ForegroundResumeRun, indices: Set<number>) {\n\treturn run.children.filter(\n\t\t(child) => indices.has(child.index) && child.status === \"detached\" && child.activityState === \"needs_attention\",\n\t);\n}\n\nfunction formatForegroundAttention(\n\trun: ForegroundResumeRun,\n\tchildren: ReturnType<typeof foregroundChildrenNeedingAttention>,\n\telapsedMs: number,\n): AgentToolResult<Details> {\n\tconst childList = children\n\t\t.map((child) => `${child.agent}${child.index !== undefined ? `#${child.index}` : \"\"}`)\n\t\t.join(\", \");\n\treturn result(\n\t\t`Waited ${formatDuration(elapsedMs)} for remembered detached foreground run \"${run.runId}\"; attention required. ${children.length} child run(s) need attention: ${childList}. Reply to any pending supervisor request, then call subagent_wait({ id: \"${run.runId}\" }) again or inspect status; do not resume or launch a replacement while it remains detached.`,\n\t);\n}\n\n/** A running run that has flagged it needs the parent's attention. */\nfunction needsAttention(run: AsyncRunSummary): boolean {\n\treturn run.activityState === \"needs_attention\" || run.steps.some((step) => step.activityState === \"needs_attention\");\n}\n\nfunction hasSupervisorTool(run: AsyncRunSummary): boolean {\n\treturn (\n\t\trun.currentTool === \"contact_supervisor\" ||\n\t\trun.currentTool === \"intercom\" ||\n\t\trun.steps.some((step) => step.currentTool === \"contact_supervisor\" || step.currentTool === \"intercom\")\n\t);\n}\n\nfunction backgroundWorkIdentity(item: RegisteredBackgroundWorkItem): string {\n\treturn `${item.provider}\\0${item.sessionId}\\0${item.id}`;\n}\n\nfunction backgroundWorkForSession(deps: SubagentWaitDeps, nowMs: number): BackgroundWorkSnapshot {\n\tconst sessionId = deps.state.currentSessionId;\n\tif (!sessionId)\n\t\tthrow new Error(\"subagent_wait requires an active session identity to scope background work safely.\");\n\treturn deps.backgroundWork?.snapshot(sessionId, nowMs) ?? snapshotBackgroundWork(sessionId, nowMs);\n}\n\n/** Queued/running runs from this session, including runs that need attention. */\nfunction activeRunsForSession(params: SubagentWaitParams, deps: SubagentWaitDeps): AsyncRunSummary[] {\n\tconst asyncDirRoot = deps.asyncDirRoot ?? DIRS.async;\n\tconst resultsDir = deps.resultsDir ?? DIRS.results;\n\tconst runs = listAsyncRuns(asyncDirRoot, {\n\t\tstates: [...ACTIVE_STATES],\n\t\tsessionId: deps.state.currentSessionId ?? undefined,\n\t\tresultsDir,\n\t\tkill: deps.kill,\n\t\tnow: deps.now,\n\t\t...(params.id ? { runId: params.id } : {}),\n\t});\n\treturn params.id ? runs.filter((run) => matchesId(run, params.id!)) : runs;\n}\n\n/** Runs (from the initial set) currently flagged needs_attention, for reporting. */\nfunction attentionRunsForSession(\n\tparams: SubagentWaitParams,\n\tdeps: SubagentWaitDeps,\n\tinitialIds: Set<string>,\n): AsyncRunSummary[] {\n\treturn activeRunsForSession(params, deps).filter((run) => needsAttention(run) && initialIds.has(run.id));\n}\n\n/** All runs (any state) for this session, for the final summary. */\nfunction allRunsForSession(params: SubagentWaitParams, deps: SubagentWaitDeps): AsyncRunSummary[] {\n\tconst asyncDirRoot = deps.asyncDirRoot ?? DIRS.async;\n\tconst resultsDir = deps.resultsDir ?? DIRS.results;\n\tconst runs = listAsyncRuns(asyncDirRoot, {\n\t\tsessionId: deps.state.currentSessionId ?? undefined,\n\t\tresultsDir,\n\t\tkill: deps.kill,\n\t\tnow: deps.now,\n\t\t...(params.id ? { runId: params.id } : {}),\n\t});\n\treturn params.id ? runs.filter((run) => matchesId(run, params.id!)) : runs;\n}\n\nfunction summarizeTerminalRuns(runs: AsyncRunSummary[], providerFinishedCount = 0): string {\n\tif (runs.length === 0 && providerFinishedCount === 0) return \"\";\n\tconst counts = { complete: 0, failed: 0, paused: 0 } as Record<string, number>;\n\tfor (const run of runs) {\n\t\tconst count = counts[run.state];\n\t\tif (count !== undefined) counts[run.state] = count + 1;\n\t}\n\tconst parts: string[] = [];\n\tif (counts.complete) parts.push(`${counts.complete} complete`);\n\tif (counts.failed) parts.push(`${counts.failed} failed`);\n\tif (counts.paused) parts.push(`${counts.paused} paused`);\n\tif (providerFinishedCount > 0) parts.push(`${providerFinishedCount} provider item(s) finished`);\n\treturn parts.join(\", \");\n}\n\nfunction result(text: string, isError = false): AgentToolResult<Details> {\n\treturn {\n\t\tcontent: [{ type: \"text\", text }],\n\t\t...(isError ? { isError: true } : {}),\n\t\tdetails: { mode: \"management\", results: [] },\n\t};\n}\n\n/** Build the live status shown while async work keeps subagent_wait blocked. */\nfunction asyncWaitUpdate(runs: AsyncRunSummary[], providerCount: number, elapsedMs: number): AgentToolResult<Details> {\n\tconst activity = runs.flatMap((run) => {\n\t\tconst activeSteps = run.steps.filter((step) => step.status === \"pending\" || step.status === \"running\");\n\t\tif (activeSteps.length === 0) {\n\t\t\treturn [`${run.id}: ${run.state}`];\n\t\t}\n\t\treturn activeSteps.map((step) => {\n\t\t\tconst current = step.currentTool ?? (step.status === \"pending\" ? \"queued\" : \"thinking…\");\n\t\t\treturn `${step.agent}: ${current}${step.currentPath ? ` ${shortenPath(step.currentPath)}` : \"\"}`;\n\t\t});\n\t});\n\tconst headline = [\n\t\t`Waiting ${formatDuration(elapsedMs)} for ${runs.length} async run(s) and ${providerCount} provider item(s).`,\n\t\t...activity,\n\t].join(\" · \");\n\treturn result([headline, runs.length > 0 ? formatAsyncRunList(runs) : \"\"].filter(Boolean).join(\"\\n\"));\n}\n\nconst TRANSCRIPT_TAIL_BYTES = 128 * 1024;\nconst TRANSCRIPT_PREVIEW_LINES = 3;\nconst TRANSCRIPT_PREVIEW_WIDTH = 220;\n\ninterface TranscriptActivity {\n\tcurrentTool?: string;\n\tcurrentToolArgs?: string;\n\tlatestAt?: number;\n\trecent: string[];\n}\n\nfunction compactTranscriptText(value: unknown, maxWidth = TRANSCRIPT_PREVIEW_WIDTH): string | undefined {\n\tif (typeof value !== \"string\") return undefined;\n\tconst compact = value.replace(/\\s+/g, \" \").trim();\n\tif (!compact) return undefined;\n\treturn compact.length <= maxWidth ? compact : `${compact.slice(0, Math.max(1, maxWidth - 1))}…`;\n}\n\nfunction readTranscriptActivity(transcriptPath: string | undefined): TranscriptActivity | undefined {\n\tif (!transcriptPath) return undefined;\n\tlet fd: number | undefined;\n\ttry {\n\t\tconst stat = fs.statSync(transcriptPath);\n\t\tif (!stat.isFile() || stat.size <= 0) return undefined;\n\t\tconst start = Math.max(0, stat.size - TRANSCRIPT_TAIL_BYTES);\n\t\tconst length = stat.size - start;\n\t\tconst buffer = Buffer.allocUnsafe(length);\n\t\tfd = fs.openSync(transcriptPath, \"r\");\n\t\tconst bytesRead = fs.readSync(fd, buffer, 0, length, start);\n\t\tlet lines = buffer.subarray(0, bytesRead).toString(\"utf-8\").split(/\\r?\\n/);\n\t\tif (start > 0) lines = lines.slice(1); // The first tail line may be partial JSON.\n\n\t\tlet currentTool: string | undefined;\n\t\tlet currentToolArgs: string | undefined;\n\t\tlet latestAt: number | undefined;\n\t\tconst recent: string[] = [];\n\t\tfor (const line of lines) {\n\t\t\tif (!line.trim()) continue;\n\t\t\tlet record: Record<string, unknown>;\n\t\t\ttry {\n\t\t\t\trecord = JSON.parse(line) as Record<string, unknown>;\n\t\t\t} catch {\n\t\t\t\tcontinue; // Concurrent append can leave the final line temporarily incomplete.\n\t\t\t}\n\t\t\tif (typeof record.ts === \"number\" && Number.isFinite(record.ts)) latestAt = record.ts;\n\t\t\tif (record.recordType === \"tool_start\" && typeof record.toolName === \"string\") {\n\t\t\t\tcurrentTool = record.toolName;\n\t\t\t\tcurrentToolArgs = compactTranscriptText(record.argsPreview, 140);\n\t\t\t\tconst preview = currentToolArgs ? `${currentTool}: ${currentToolArgs}` : currentTool;\n\t\t\t\trecent.push(`tool: ${preview}`);\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tif (record.recordType === \"tool_end\") {\n\t\t\t\tcurrentTool = undefined;\n\t\t\t\tcurrentToolArgs = undefined;\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tif (record.recordType === \"message\") {\n\t\t\t\tif (record.sourceEventType === \"initial_prompt\" || record.role === \"user\") continue;\n\t\t\t\tconst text = compactTranscriptText(record.text);\n\t\t\t\tif (text) recent.push(`${typeof record.role === \"string\" ? record.role : \"message\"}: ${text}`);\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tif (record.recordType === \"stdout\" || record.recordType === \"stderr\") {\n\t\t\t\tconst text = compactTranscriptText(record.text);\n\t\t\t\tif (text) recent.push(`${record.recordType}: ${text}`);\n\t\t\t}\n\t\t}\n\t\treturn { currentTool, currentToolArgs, latestAt, recent: recent.slice(-TRANSCRIPT_PREVIEW_LINES) };\n\t} catch {\n\t\treturn undefined;\n\t} finally {\n\t\tif (fd !== undefined) {\n\t\t\ttry {\n\t\t\t\tfs.closeSync(fd);\n\t\t\t} catch {\n\t\t\t\t/* best-effort close */\n\t\t\t}\n\t\t}\n\t}\n}\n\nfunction detachedForegroundWaitUpdate(\n\trun: ForegroundResumeRun,\n\tpendingIndices: Set<number>,\n\tnowMs: number,\n\telapsedMs: number,\n): AgentToolResult<Details> {\n\tconst lines = [`Waiting for detached foreground run \"${run.runId}\" · ${formatDuration(elapsedMs)}`];\n\tfor (const child of run.children) {\n\t\tif (!pendingIndices.has(child.index) || child.status !== \"detached\") continue;\n\t\tconst activity = readTranscriptActivity(child.transcriptPath);\n\t\tconst age =\n\t\t\tactivity?.latestAt !== undefined\n\t\t\t\t? ` · activity ${formatDuration(Math.max(0, nowMs - activity.latestAt))} ago`\n\t\t\t\t: \"\";\n\t\tlines.push(`${child.agent} · working after supervisor handoff${age}`);\n\t\tif (activity?.currentTool) {\n\t\t\tlines.push(\n\t\t\t\t`  current: ${activity.currentTool}${activity.currentToolArgs ? `: ${activity.currentToolArgs}` : \"\"}`,\n\t\t\t);\n\t\t}\n\t\tfor (const preview of activity?.recent ?? []) lines.push(`  ${preview}`);\n\t\tif (!activity && child.transcriptPath) lines.push(\"  live transcript has no new activity yet\");\n\t\tif (!child.transcriptPath) lines.push(\"  live transcript unavailable; waiting for completion event\");\n\t}\n\treturn result(lines.join(\"\\n\"));\n}\n\nasync function waitForDetachedForegroundRun(\n\trun: ForegroundResumeRun,\n\tsignal: AbortSignal | undefined,\n\tdeps: SubagentWaitDeps,\n\tstartedAt: number,\n\tnow: () => number,\n\tpollIntervalMs: number,\n\ttimeoutMs: number,\n): Promise<AgentToolResult<Details>> {\n\tconst initialDetachedIndices = new Set(\n\t\trun.children.filter((child) => child.status === \"detached\").map((child) => child.index),\n\t);\n\twhile (true) {\n\t\tif (deps.state.currentSessionId !== run.sessionId) {\n\t\t\treturn result(\n\t\t\t\t`Wait stopped because the active session changed while remembered foreground run \"${run.runId}\" was still detached. Return to the originating session to inspect or wait for it. Reply to any pending supervisor request before resuming or launching a replacement.`,\n\t\t\t\ttrue,\n\t\t\t);\n\t\t}\n\t\tconst current = deps.state.foregroundRuns?.get(run.runId);\n\t\tif (!current || current.sessionId !== run.sessionId) {\n\t\t\treturn result(\n\t\t\t\t`Remembered foreground run \"${run.runId}\" disappeared before a terminal child result was recorded. Completion cannot be confirmed; do not launch a replacement without checking the originating child session.`,\n\t\t\t\ttrue,\n\t\t\t);\n\t\t}\n\t\tconst pending = current.children.filter(\n\t\t\t(child) => initialDetachedIndices.has(child.index) && child.status === \"detached\",\n\t\t);\n\t\tconst attention =\n\t\t\tdeps.stopOnAttention === false ? [] : foregroundChildrenNeedingAttention(current, initialDetachedIndices);\n\t\tif (attention.length > 0) return formatForegroundAttention(current, attention, now() - startedAt);\n\t\tif (pending.length === 0) {\n\t\t\tconst outcome = summarizeForegroundChildren(current, initialDetachedIndices);\n\t\t\treturn result(\n\t\t\t\t`Waited ${formatDuration(now() - startedAt)} for remembered detached foreground run \"${run.runId}\"; done. Outcome: ${outcome || \"no recovered child status\"}. Completion event observed; inspect with subagent({ action: \"status\", id: \"${run.runId}\" }) for recovered output.`,\n\t\t\t);\n\t\t}\n\t\tconst updateNow = now();\n\t\tdeps.onUpdate?.(detachedForegroundWaitUpdate(current, initialDetachedIndices, updateNow, updateNow - startedAt));\n\t\tif (signal?.aborted) {\n\t\t\treturn result(\n\t\t\t\t`Wait aborted after ${formatDuration(now() - startedAt)}. Remembered foreground run \"${run.runId}\" remains detached. Reply to any pending supervisor request before resuming or launching a replacement.`,\n\t\t\t\ttrue,\n\t\t\t);\n\t\t}\n\t\tif (now() - startedAt >= timeoutMs) {\n\t\t\treturn result(\n\t\t\t\t`Wait timed out after ${formatDuration(timeoutMs)} with remembered foreground run \"${run.runId}\" still detached. Reply to any pending supervisor request, then call subagent_wait({ id: \"${run.runId}\" }) again or inspect status; do not resume or launch a replacement while it remains detached.`,\n\t\t\t\ttrue,\n\t\t\t);\n\t\t}\n\t\tawait waitForWake(pollIntervalMs, signal, deps);\n\t}\n}\n\n/**\n * Block until the targeted async or remembered detached foreground run finishes,\n * the timeout elapses, or the turn is aborted. Resolves with a short\n * human-readable summary either way.\n */\nexport async function waitForSubagents(\n\tparams: SubagentWaitParams,\n\tsignal: AbortSignal | undefined,\n\tdeps: SubagentWaitDeps,\n): Promise<AgentToolResult<Details>> {\n\tif (deps.enabled === false) {\n\t\treturn result(\n\t\t\t'subagent_wait is disabled by config.waitTool or PI_SUBAGENT_WAIT_TOOL_ENABLED; returning immediately without blocking background work. Active work keeps going, and you can inspect subagents with subagent({ action: \"status\" }) or rely on completion notifications.',\n\t\t);\n\t}\n\tif (!deps.state.currentSessionId) {\n\t\treturn result(\"subagent_wait requires an active session identity to scope background work safely.\", true);\n\t}\n\n\tconst now = deps.now ?? Date.now;\n\tconst pollIntervalMs = Math.max(MIN_POLL_INTERVAL_MS, deps.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS);\n\tconst timeoutMs = params.timeoutMs !== undefined && params.timeoutMs > 0 ? params.timeoutMs : DEFAULT_TIMEOUT_MS;\n\tconst startedAt = now();\n\tconst waitForAll = params.id ? true : params.all === true;\n\tif (params.nonBlocking && !params.id) {\n\t\treturn result(\n\t\t\t\"Non-blocking wait subscriptions require id so the registration can bind one exact run identity.\",\n\t\t\ttrue,\n\t\t);\n\t}\n\tif (params.nonBlocking && params.all) {\n\t\treturn result(\"nonBlocking cannot be combined with all; subscribe to one exact run id.\", true);\n\t}\n\n\tlet active: AsyncRunSummary[];\n\tlet foreground: ForegroundResumeRun[];\n\tlet providerSnapshot: BackgroundWorkSnapshot;\n\ttry {\n\t\tactive = activeRunsForSession(params, deps);\n\t\tforeground = activeDetachedForegroundRuns(params, deps);\n\t\tproviderSnapshot = params.id ? { providers: [], items: [] } : backgroundWorkForSession(deps, startedAt);\n\t} catch (error) {\n\t\treturn result(error instanceof Error ? error.message : String(error), true);\n\t}\n\n\tif (params.id) {\n\t\tconst candidates = [\n\t\t\t...active.map((run) => ({ kind: \"async\" as const, id: run.id, run })),\n\t\t\t...foreground.map((run) => ({ kind: \"foreground\" as const, id: run.runId, run })),\n\t\t];\n\t\tconst exact = candidates.filter((candidate) => candidate.id === params.id);\n\t\tconst matches = exact.length > 0 ? exact : candidates;\n\t\tif (matches.length > 1) {\n\t\t\treturn result(\n\t\t\t\t`Ambiguous subagent run id prefix \"${params.id}\" matched ${matches.length} active runs: ${matches.map((candidate) => candidate.id).join(\", \")}. Pass a longer id.`,\n\t\t\t\ttrue,\n\t\t\t);\n\t\t}\n\t\tconst selected = matches[0];\n\t\tif (selected && params.nonBlocking) {\n\t\t\tif (!deps.subscribe) {\n\t\t\t\treturn result(\n\t\t\t\t\t\"Non-blocking wait subscriptions require a long-lived interactive subagent runtime; this runtime can only use blocking subagent_wait calls.\",\n\t\t\t\t\ttrue,\n\t\t\t\t);\n\t\t\t}\n\t\t\ttry {\n\t\t\t\tconst registration = deps.subscribe({\n\t\t\t\t\ttargetKind: selected.kind,\n\t\t\t\t\trunId: selected.id,\n\t\t\t\t\trequestedId: params.id,\n\t\t\t\t\ttimeoutMs,\n\t\t\t\t});\n\t\t\t\treturn result(\n\t\t\t\t\t`Armed wait subscription ${registration.token} for exact ${selected.kind} run ${selected.id}. Returning immediately; this session will be woken on completion, failure, attention, reconciliation failure, or timeout. Inspect armed subscriptions with subagent({ action: \"status\" }).`,\n\t\t\t\t);\n\t\t\t} catch (error) {\n\t\t\t\treturn result(error instanceof Error ? error.message : String(error), true);\n\t\t\t}\n\t\t}\n\t\tif (selected?.kind === \"foreground\") {\n\t\t\treturn waitForDetachedForegroundRun(selected.run, signal, deps, startedAt, now, pollIntervalMs, timeoutMs);\n\t\t}\n\t\tactive = selected?.kind === \"async\" ? [selected.run] : [];\n\t}\n\n\tlet providerActive = providerSnapshot.items;\n\tif (active.length === 0 && providerActive.length === 0) {\n\t\treturn result(\n\t\t\tparams.id\n\t\t\t\t? `No active run matched \"${params.id}\". Nothing to wait for.`\n\t\t\t\t: \"No active async runs or registered provider work in this session. Nothing to wait for.\",\n\t\t);\n\t}\n\tconst waitParams = params.id ? { ...params, id: active[0]!.id } : params;\n\tconst initialAsyncIds = new Set(active.map((run) => run.id));\n\tconst initialProviderIds = new Set(providerActive.map(backgroundWorkIdentity));\n\tconst initialProviderNames = new Set(providerActive.map((item) => item.provider));\n\tconst initialCount = initialAsyncIds.size + initialProviderIds.size;\n\tconst stopOnAttention = deps.stopOnAttention !== false;\n\tlet attention = active.filter((run) => needsAttention(run));\n\n\tconst isDone = (): boolean => {\n\t\tif (stopOnAttention && attention.some((run) => initialAsyncIds.has(run.id))) return true;\n\t\tconst activeAsyncIds = new Set(active.map((run) => run.id));\n\t\tconst activeProviderIds = new Set(providerActive.map(backgroundWorkIdentity));\n\t\tif (waitForAll) {\n\t\t\treturn (\n\t\t\t\t[...initialAsyncIds].every((id) => !activeAsyncIds.has(id)) &&\n\t\t\t\t[...initialProviderIds].every((id) => !activeProviderIds.has(id))\n\t\t\t);\n\t\t}\n\t\treturn (\n\t\t\t[...initialAsyncIds].some((id) => !activeAsyncIds.has(id)) ||\n\t\t\t[...initialProviderIds].some((id) => !activeProviderIds.has(id))\n\t\t);\n\t};\n\n\twhile (!isDone()) {\n\t\tconst activeInitialRuns = active.filter((run) => initialAsyncIds.has(run.id));\n\t\tconst activeInitialProviderItems = providerActive.filter((item) =>\n\t\t\tinitialProviderIds.has(backgroundWorkIdentity(item)),\n\t\t);\n\t\tconst stillActive = [\n\t\t\t...activeInitialRuns.map((run) => `${run.id} (${run.state})`),\n\t\t\t...activeInitialProviderItems.map((item) => `${item.provider}/${item.id}`),\n\t\t].join(\", \");\n\t\tdeps.onUpdate?.(asyncWaitUpdate(activeInitialRuns, activeInitialProviderItems.length, now() - startedAt));\n\t\tif (signal?.aborted) {\n\t\t\treturn result(`Wait aborted after ${formatDuration(now() - startedAt)}. Still active: ${stillActive}.`, true);\n\t\t}\n\t\tif (now() - startedAt >= timeoutMs) {\n\t\t\treturn result(\n\t\t\t\t`Wait timed out after ${formatDuration(timeoutMs)} with ${activeInitialRuns.length} async run(s) and ${activeInitialProviderItems.length} provider item(s) still active: ${stillActive}. The work keeps going; call subagent_wait again or inspect subagent status.`,\n\t\t\t\ttrue,\n\t\t\t);\n\t\t}\n\t\ttry {\n\t\t\tawait waitForWake(pollIntervalMs, signal, deps);\n\t\t\tactive = activeRunsForSession(waitParams, deps);\n\t\t\tattention = attentionRunsForSession(waitParams, deps, initialAsyncIds);\n\t\t\tproviderSnapshot = params.id ? providerSnapshot : backgroundWorkForSession(deps, now());\n\t\t\tfor (const provider of initialProviderNames) {\n\t\t\t\tif (!providerSnapshot.providers.includes(provider)) {\n\t\t\t\t\treturn result(\n\t\t\t\t\t\t`Background-work provider '${provider}' disappeared while subagent_wait was tracking its active work; completion cannot be confirmed.`,\n\t\t\t\t\t\ttrue,\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t}\n\t\t\tproviderActive = providerSnapshot.items;\n\t\t} catch (error) {\n\t\t\treturn result(error instanceof Error ? error.message : String(error), true);\n\t\t}\n\t}\n\n\tlet terminalSummary: string;\n\tlet finishedAsyncCount: number;\n\tlet failedAsyncCount: number;\n\tconst activeProviderIds = new Set(providerActive.map(backgroundWorkIdentity));\n\tconst providerFinishedCount = [...initialProviderIds].filter((id) => !activeProviderIds.has(id)).length;\n\ttry {\n\t\tconst allNow = allRunsForSession(waitParams, deps);\n\t\tconst terminal = allNow.filter((run) => !ACTIVE_STATES.includes(run.state) && initialAsyncIds.has(run.id));\n\t\tfinishedAsyncCount = terminal.length;\n\t\tfailedAsyncCount = terminal.filter((run) => run.state === \"failed\").length;\n\t\tterminalSummary = summarizeTerminalRuns(terminal, providerFinishedCount);\n\t} catch (error) {\n\t\treturn result(error instanceof Error ? error.message : String(error), true);\n\t}\n\n\tconst relevantAttention = attention.filter((run) => initialAsyncIds.has(run.id));\n\tconst supervisorAttentionHint = relevantAttention.some(hasSupervisorTool)\n\t\t? ' Reply to any pending supervisor request. If subagent_supervisor({ action: \"pending\" }) is empty, check intercom({ action: \"pending\" }) because an external intercom tool may own the request.'\n\t\t: \"\";\n\tconst attentionNote =\n\t\trelevantAttention.length > 0\n\t\t\t? ` ${relevantAttention.length} run(s) need attention: ${relevantAttention.map((run) => run.id).join(\", \")} —${supervisorAttentionHint} inspect with subagent({ action: \"status\" }) then steer a top-level live async child, resume a paused/completed/failed child, or interrupt explicitly.`\n\t\t\t: \"\";\n\tconst stillRunning =\n\t\tactive.filter((run) => initialAsyncIds.has(run.id)).length +\n\t\tproviderActive.filter((item) => initialProviderIds.has(backgroundWorkIdentity(item))).length;\n\tconst elapsed = formatDuration(now() - startedAt);\n\tconst outcome = terminalSummary ? ` Outcome: ${terminalSummary}.` : \"\";\n\n\tif (waitForAll) {\n\t\tconst scope = params.id\n\t\t\t? `run \"${params.id}\"`\n\t\t\t: initialProviderIds.size === 0\n\t\t\t\t? `${initialAsyncIds.size} async run(s)`\n\t\t\t\t: `${initialAsyncIds.size} async run(s) and ${initialProviderIds.size} provider item(s)`;\n\t\tconst status = relevantAttention.length > 0 ? \"attention required\" : \"done\";\n\t\treturn result(\n\t\t\t`Waited ${elapsed} for ${scope}; ${status}.${outcome}${attentionNote} Completion/control events have been observed; inspect status if a notification is not visible yet.`,\n\t\t\tdeps.failOnFailedRuns === true && failedAsyncCount > 0,\n\t\t);\n\t}\n\n\tconst finishedCount = finishedAsyncCount + providerFinishedCount;\n\tconst subject = initialProviderIds.size === 0 ? \"run(s)\" : \"item(s)\";\n\tconst remainder =\n\t\tstillRunning > 0\n\t\t\t? ` ${stillRunning} ${subject} still in flight — call subagent_wait again to catch the next one.`\n\t\t\t: relevantAttention.length > 0\n\t\t\t\t? \" No other work is waitable until attention is handled.\"\n\t\t\t\t: initialProviderIds.size === 0\n\t\t\t\t\t? \" No runs remain in flight.\"\n\t\t\t\t\t: \" No work remains in flight.\";\n\tconst progress =\n\t\trelevantAttention.length > 0 && finishedCount === 0\n\t\t\t? `${relevantAttention.length} of ${initialCount} ${subject} need attention`\n\t\t\t: `${finishedCount} of ${initialCount} ${subject} finished`;\n\treturn result(\n\t\t`Waited ${elapsed}; ${progress}.${outcome}${attentionNote}${remainder} Relevant completion/control events have been observed; inspect status if a notification is not visible yet.`,\n\t\tdeps.failOnFailedRuns === true && failedAsyncCount > 0,\n\t);\n}\n"]}