import { spawn } from "node:child_process"; import * as fs from "node:fs"; import * as path from "node:path"; import type { AgentConfig } from "../../agents/agent-config.ts"; import { loadConfig } from "../../config/config.ts"; import { getCrewEnv } from "../../config/env-vars.ts"; import type { PiTeamsConfig } from "../../config/types.ts"; import { registerChildProcess, unregisterChildProcess } from "../../extension/crew-cleanup.ts"; import type { WorkerExitStatus } from "../../state/types.ts"; import { runEventBus } from "../../ui/run-event-bus.ts"; import { logInternalError } from "../../utils/internal-error.ts"; import { redactSecretString } from "../../utils/redaction.ts"; import { getActiveBrokerIssuer } from "../broker/broker-issuer.ts"; import { BoundedTail } from "../compaction/compact-stages/bounded-tail.ts"; import { EventLogTailSource } from "../event-log-tail-source.ts"; import { bridgeEventFromJsonEvent } from "../event-stream-bridge.ts"; import { classifyOnExit, getSurfaceRuntimeController, makeTerminalEventProbe } from "../surface/degrade.ts"; import type { SurfaceGateEnvSnapshot } from "../surface/resolve-surface.ts"; import type { SurfaceProvider } from "../surface/surface-provider.ts"; import { prepareSurfaceSpawn, type SurfaceExitInfo, type SurfaceSpawnOutcome, waitForSurfaceExit } from "../surface/surface-spawn.ts"; import { FINAL_DRAIN_MS, HARD_KILL_MS, POST_EXIT_STDIO_GUARD_MS, RESPONSE_TIMEOUT_MS } from "./child-pi-constants.ts"; import { clearHardKillTimer, killProcessTree, registerActiveChild, unregisterActiveChild } from "./child-pi-kill.ts"; import { buildFinalChildPiSpawnOptions, prepareSpawnContext } from "./child-pi-spawn.ts"; import { ChildPiSteeringController } from "./child-pi-steering.ts"; // Internal helpers for active-child bookkeeping (extracted to child-pi-kill.ts). import { ChildPiLineObserver } from "./child-pi-streams.ts"; // Phase 2.3: the six timer constructs moved to child-pi-timers.ts (pure motion). import { createChildPiTimers } from "./child-pi-timers.ts"; import { runMockChildPi } from "./mock-fixtures.ts"; // W2 (session-file recovery): deterministic worker session flags + crash-path // tail-replay of the session JSONL (see session-recovery.ts header). import { appendWorkerSessionArgs, deriveSessionPaths, recoverLastAssistantFromSession, resolveSessionRecoveryEnabled, type SessionRecoveryInfo, shouldAttemptSessionRecovery, } from "./session-recovery.ts"; // ── Re-exports from child-pi-kill.ts (H-7 decomposition step 2) ── // killProcessTree is internal (not previously exported) — keep that invariant. export { killProcessPid, terminateActiveChildPiProcesses, } from "./child-pi-kill.ts"; // ── Re-export from child-pi-spawn.ts (H-7 decomposition step 6) ── // buildChildPiSpawnOptions was previously exported from child-pi.ts. Keep the // public API surface stable by re-exporting from the new module. // buildFinalChildPiSpawnOptions (BLOCKER 2 / S5) — composed spawn helper that // owns the canary + filter + spread sequence; lives in child-pi-spawn.ts. export { buildChildPiSpawnOptions, buildFinalChildPiSpawnOptions } from "./child-pi-spawn.ts"; // ── Re-export from child-pi-streams.ts (H-7 decomposition step 4) ── export { ChildPiLineObserver } from "./child-pi-streams.ts"; import { checkCrewDepth, cleanupTempDir, resolveHermeticWorkers } from "../model/pi-args.ts"; import { attachPostExitStdioGuard, trySignalChild } from "../process/post-exit-stdio-guard.ts"; import { classifyProcessCrash } from "../recovery/crash-classification.ts"; import { probePiVersion } from "./pi-version.ts"; /** Maximum size (bytes) for the ChildPiLineObserver's line accumulation buffer. * When exceeded, the buffer is force-flushed to prevent unbounded memory growth * from chatty child processes that produce output without newlines. * (Constant moved to child-pi-constants.ts.) */ // Periodic cleanup of dead child process entries to prevent memory leaks. /** * SEC-1: Extract a redacted stderr/stdout excerpt for embedding in lifecycle * events and error messages. The in-memory stdout/stderr accumulators receive * RAW worker output (only structurally compacted via compactChildPiEvent — * NOT secret-redacted), so any slice embedded into a persisted event must be * redacted here. Otherwise worker-emitted secrets (API keys, tokens returned * from a tool call) leak through diagnostic logs that bypass artifact-store * redaction. * * Extracted as a single helper (8 call sites were duplicating this) so the * redaction boundary is unit-testable directly. The real spawn error/timeout * paths are integration-level and NOT reachable via PI_TEAMS_MOCK_CHILD_PI * (the mock returns before the lifecycle-event handlers run), so a behavior * test must target this helper rather than the full runChildPi path. */ export function redactStderrExcerpt(stderr: string, maxChars: number): string { return redactSecretString(stderr.slice(-maxChars)); } /** * B6: spawn taskkill and attach an 'error' listener. spawn() emits ENOENT/EACCES * asynchronously via the 'error' event (not as a throw), so an unlistened spawn * can crash the parent as an uncaught exception. taskkill is a standard Windows * binary so this is defensive, but the listener keeps failures bounded. */ /** Structured lifecycle event emitted by child-pi for critical transitions. */ export interface ChildPiLifecycleEvent { /** * Event discriminator. `surface_spawned` / `surface_closed` are emitted by * the MuxSurface spawn branch (spec §13.1) — no stdio process exists, so * `pid` is absent there and surface fields carry the pane identity. */ type: | "spawned" | "spawn_error" | "response_timeout" | "final_drain" | "hard_kill" | "exit" | "close" | "surface_spawned" | "surface_closed" /** T11 (spec §7): a REAL surface spawn attempt failed → headless fallback * + the run-scoped spawn-fail lockout counter (3 consecutive → OFF). */ | "surface_spawn_failed" /** FINDING-3 (report 10tier 2026-08-27): a gate rejected surface BEFORE * touching the mux (mode-off/async-run/depth/pane-cap/role-not-visible/ * no-mux). Emitted ONLY when the user opted into surface * (visibleAgents non-empty) — default runs stay silent. */ | "surface_gate_blocked"; /** Process ID when available. */ pid?: number; /** Pane id (surface events). */ paneId?: string; /** Surface backend kind (surface events). */ surfaceKind?: "tmux" | "herdr"; /** Why the pane ended / degraded (surface_closed): "pane-closed" | "mux-dead" | "detached". */ paneExitReason?: string; /** Which gate rejected surface (surface_gate_blocked): one of the SurfaceGateName values. */ gate?: string; /** Env signals the gate saw (surface_gate_blocked) — bounded snapshot, never the whole env. */ env?: SurfaceGateEnvSnapshot; /** Exit code for exit/close events. */ exitCode?: number | null; /** Error message for error events. */ error?: string; /** Stderr captured at timeout moment (for response_timeout events). */ stderr?: string; /** Last N chars of stderr for error context (exit/error events). */ stderrExcerpt?: string; /** Timestamp (ISO). */ ts: string; /** F12: optional cause for `final_drain` events. `"stdout-quiet"` indicates * the drain was triggered by the quiet-window early-exit rather than the * default 5 s ceiling. Other drain reasons (default) leave this undefined. */ reason?: "stdout-quiet"; /** Phase-0 diagnostic (HB-003a): the signal that killed the child (when * available). Was previously discarded after building the error string. */ signal?: string; /** Phase-0 diagnostic (HB-003a): final-drain race timing, present only on * exit events where a drain timer was armed. Surfaces the exit-null race. */ diagnostic?: { finalDrainArmed: boolean; forcedFinalDrain: boolean; finalDrainFiredMonotonicMs?: number; finalAssistantEventMonotonicMs?: number; exitMonotonicMs: number; }; } export interface ChildPiRunInput { cwd: string; task: string; agent: AgentConfig; model?: string; skillPaths?: string[]; signal?: AbortSignal; transcriptPath?: string; onStdoutLine?: (line: string) => void; onJsonEvent?: (event: unknown) => void; onSpawn?: (pid: number) => void; /** Structured lifecycle events for durable logging (spawn, crash, timeout, kill, exit). */ onLifecycleEvent?: (event: ChildPiLifecycleEvent) => void; /** F1 fix (battery 2026-09-10): task-state activity bridge for SURFACE workers. * Surface workers have no stdout pipe, so onJsonEvent/onStdoutLine never fire * and task.heartbeat/agentProgress starve - the heartbeat-watcher then flags a * healthy pane worker dead (~340s) and the stale-reconciler cancels the run * (observed live: run team_20260910161234, worker productive 7min in pane * w2:p8Z, 0 task.progress events). This callback receives every event tailed * from the per-agent recorder log so the CALLER can update task liveness state. * Contract: the callback MUST NOT append to the per-agent events file being * tailed (watch->append loop - see event-log-tail-source.ts header); task-state * persist + run-level task.progress are safe (different files). */ onSurfaceActivity?: (event: Record) => void; maxDepth?: number; finalDrainMs?: number; /** F12: early-exit the drain when stdout has been silent for this many ms * after the final assistant event. Set to ≥ finalDrainMs to disable. */ finalDrainQuietMs?: number; hardKillMs?: number; responseTimeoutMs?: number; /** Soft limit on assistant turns — inject steer at this count. */ maxTurns?: number; /** Extra turns after soft limit before hard abort. Default: 5. */ graceTurns?: number; /** Parent conversation context to inherit when inheritContext is true. */ parentContext?: string; /** When true, prepend parentContext to the task prompt. */ inheritContext?: boolean; /** Pass to pi to mark certain commands as context-excluded. Default: false */ excludeContextBash?: boolean; /** pi session ID for session naming (aligns with pi-crew run ID). W2: now * WIRED — flows to `--session-id` (deterministic worker session files for * crash-path tail-recovery; see session-recovery.ts). */ sessionId?: string; /** W2 (session-file recovery): explicit session directory override. When * absent, derived per-worker as `/sessions//` so * session JSONL is cleaned up together with run artifacts. */ sessionDir?: string; /** W2 (session-file recovery): explicit on/off (runtime.sessionRecovery * config). Env PI_CREW_SESSION_RECOVERY overrides in either direction; * default ON. Only consulted on crash-ish settles (exitCode null / killed). */ sessionRecovery?: boolean; /** Path to steering JSONL file for real-time steer injection. */ steeringFile?: string; /** Run ID for cleanup tracking */ runId?: string; /** Agent ID for cleanup tracking */ agentId?: string; /** Role for tool restrictions (from role-tools.ts) */ role?: string; /** Team-role thinking override (takes precedence over agent.thinking). */ thinkingOverride?: string; /** R3-1: run-static worker header for the --append-system-prompt channel * (renderTaskPrompt().stablePrefix). Threaded by child-executor so the * Protocol/mailbox/workspace/runtime-context block survives compaction * instead of living in the summarizable user-message span. */ systemPromptAppend?: string; /** R3-19/D5 (hermetic worker spawns): thread runtime.hermeticWorkers to * buildPiWorkerArgs (--no-extensions). Env PI_CREW_HERMETIC_WORKERS * overrides either way; default TRUE. */ hermeticWorkers?: boolean; /** U9 (pi 1.0.4 version floor): host pi binary version for the argv * builder's version-gated flags (--no-mcp at the 1.0.4 floor). runChildPi * OVERWRITES this with the single-flight probePiVersion() result whenever * this spawn can consume a version-gated flag (parity non-hermetic + not * pre-aborted); on hermetic spawns (the default — no version-gated flag * consulted) the field passes through untouched and the probe is SKIPPED, * keeping default-path spawns (and unit-test surface paths) free of a real * `pi --version` child. See child-pi/pi-version.ts for the floor registry. */ piVersion?: string | null; /** Root directory for artifacts (used to validate transcriptPath). */ artifactsRoot?: string; /** I5: run events JSONL path — threaded to the worker so its scratchpad * execute handler can append fire-and-forget metric events. Optional; * absent in non-team contexts (emission silently skipped). */ eventsPath?: string; /** Phase 1 scratchpad: model-fallback attempt index (0-based) for per-attempt * snapshot relativePath `scratchpad/.attempt-.snapshot.json` * (C3). Optional — worker defaults to attempt 0 when unset (e.g. custom * agent spawn outside child-executor). */ attempt?: number; /** * Optional broker spawn context (Phase 0 inter-pi broker). When present, * `prepareSpawnContext` injects `PI_CREW_BROKER_SOCKET` and * `PI_CREW_BROKER_TOKEN` into the child env (control-namespace keys, * accepted by `assertOnlyControlEnvKeys`). The token is a 128-bit random * value issued by the parent's CrewBroker at spawn time; it lives only * in the parent's in-memory `Map` and in the child's env. * NEVER persisted to manifest, events, mailbox, or any run-dir file. */ brokerSpawn?: { socketPath: string; token: string }; /** * Optional Phase 0 broker credentials issuer. Called only when * `brokerSpawn` is not already set on the input. The default no-op * issuer returns undefined, leaving the child to run without broker * credentials (today's behavior — file paths remain authoritative). * The production wiring passes a closure that delegates to the * session's `CrewBrokerLifecycleController.issueForChild`. * * ADR-5 §4: the third parameter carries the child's depth for * delegate-spawned grandchildren — the issuer gates token minting on * `childDepth < resolved maxDepth` (see BrokerIssuer). */ brokerIssuer?: (runId: string, taskId?: string, childDepth?: number) => Promise<{ socketPath: string; token: string } | undefined>; /** * ADR-5 §3 (governed nesting): explicit child depth for delegate-spawned * grandchildren. The spawn policy computes this from the PARENT TASK * RECORD (task.depth) — never from this process's env (the root has depth * 0, so env-derived depth would wrongly be 1). Expressed as the parent's * depth in the base env so the existing `parentDepth + 1` spawn math and * the `checkCrewDepth` gate both see the parent's true depth. Undefined = * normal worker spawn (unchanged env-derived behavior). */ depthOverride?: number; /** Base env for depth computation (defaults to process.env). Advanced — * used with depthOverride by the root-side delegate handler. */ env?: NodeJS.ProcessEnv; /** * MuxSurface A1 (spec §13.1): opt-in/hint for booting this worker in a * multiplexer pane instead of a stdio pipe. Undefined → the branch still * consults the discovered project config, which defaults surface OFF * (`visibleAgents: []`), so behaviour stays headless unless explicitly * configured. Any failure degrades to the normal headless spawn path (§3 * fail-closed) — NEVER throws. * * - `providers`: pre-resolved provider instances — dispatch-owned (team-runner * T11 passes its resolved provider + livePaneCount) and test injection. * Hard gates (async run / host depth) still apply. * - `config`: overrides discovered config (tests). * - `livePaneCount`: panes already live for this run — `MAX_SURFACE_WORKERS` * cap check. Defaults 0 until T11 threads manifest state. * - `baseDir`: launch-script dir override (default getPiTempBase()). */ surface?: { providers?: { tmux?: SurfaceProvider; herdr?: SurfaceProvider }; config?: PiTeamsConfig; livePaneCount?: number; baseDir?: string; }; } export interface ChildPiRunResult { exitCode: number | null; stdout: string; stderr: string; error?: string; /** RAW (uncapped) final assistant text, captured at stream-parse time BEFORE * the 16K transcript compaction. This is the AUTHORITATIVE worker output — * it becomes results/.txt so downstream dependencies are not bounded by * the transcript's telemetry cap. Undefined when no assistant text was seen * (mock paths, error paths) — callers MUST fall back to transcript-derived * finalText. See research-findings/output-handling-deep-dive.md §A. */ rawFinalText?: string; exitStatus?: WorkerExitStatus; /** True if the agent was hard-aborted (max_turns + grace exceeded). */ aborted?: boolean; /** * DR5 (RPC double-execution guard, deep-review 2026-10-05 §2): set by the * RPC transport when a session event (message_start/message_end/…) was * observed after the prompt was sent but the turn never settled — i.e. * the agent turn HAD begun when the transport failed, even though no * assistant text was captured. `isEarlyRpcTransportFailure` treats this * as non-retryable: the turn may already have performed side effects, and * double-executing a task that already produced work is worse than * surfacing the transport error. An explicit result field (not an * error-string marker) so the contract is testable and greppable; the * stdio transport never sets it. */ rpcAgentStarted?: boolean; /** True if the agent was steered to wrap up (hit soft turn limit) but finished in time. */ steered?: boolean; /** #7 hardening: bounded digest of intermediate findings (last N tool results or * assistant text lines) from the run. Populated by ChildPiLineObserver so that * workers that exhaust their budget on tool calls (never emit final assistant * text) still produce a non-empty result. Consumers should prefer rawFinalText * first — this is a last-resort fallback. */ intermediateFindings?: string; /** W2 (session-file recovery): present when the worker died WITHOUT a * final assistant event (exitCode null / killed) and a COMPLETE assistant * record was tail-recovered from the worker's session JSONL (deterministic * `--session-id`/`--session-dir` under the run artifacts root). Consumers * treat it as a crash-path supplement to rawFinalText — never a replacement * for clean-path output, and manifest polling is untouched. */ recoveredFromSession?: SessionRecoveryInfo; /** * MuxSurface A1: present ONLY when the worker booted inside a multiplexer * pane instead of a stdio pipe (spec §13.1). Downstream (T9 EventLogTailSource, * T11 degrade) consumes the pane identity; the empty stdout/stderr is EXPECTED * in surface mode — events arrive through the per-agent event log, not a pipe. * * `degraded` is set when classifyOnExit found NO worker.completed within the * 2s window after the pane died (spec §7 D3) — task-runner hands that to * finalize as `surfaceLost` instead of reading it as an empty/false result. */ surface?: { kind: "tmux" | "herdr"; paneId: string; scriptPath: string; degraded?: { cause: "pane-closed" | "mux-dead"; exitReason: string; classifiedAt: string; }; }; } // Base allowlist of non-provider env vars always passed to child workers. // ── Transcript batching + compaction (H-7 decomposition step 1) ──────── // Extracted to ./child-pi-transcript.ts. Re-exported here to preserve the // existing public API surface. export { appendTranscript, compactString, compactValue, flushPendingTranscriptWrites, resetTranscriptBatchState, } from "./child-pi-transcript.ts"; /** Mock-only path — real code path reuses a single observer. * OPT-06 follow-up: returns a Promise so callers can await the transcript * drain before resolving runChildPi. Without this, mock-mode callers that * read the transcript file post-run see ENOENT (the async file handle had * not yet been opened). */ async function observeStdoutChunk(input: ChildPiRunInput, text: string): Promise { const observer = new ChildPiLineObserver(input); observer.observe(text); await observer.flush(); } function asRecord(value: unknown): Record | undefined { return value && typeof value === "object" && !Array.isArray(value) ? (value as Record) : undefined; } function isFinalAssistantEvent(event: unknown): boolean { const obj = asRecord(event); if (obj?.type !== "message_end") return false; const message = asRecord(obj.message); const role = message?.role; if (role !== undefined && role !== "assistant") return false; const stopReason = typeof message?.stopReason === "string" ? message.stopReason : typeof obj.stopReason === "string" ? obj.stopReason : undefined; if (stopReason !== undefined && stopReason !== "stop") return false; const content = Array.isArray(message?.content) ? message.content : []; return !content.some((part) => asRecord(part)?.type === "toolCall"); } /** * Discover the project crew config for the surface gate. The discovered config * DEFAULTS surface off (`visibleAgents: []`) — headless unless explicitly opted * in. Any discovery failure → default config (fail-closed §3), never throws. */ function safeLoadSurfaceConfig(cwd: string): PiTeamsConfig { try { return loadConfig(cwd).config; } catch (error) { logInternalError( "child-pi.surface-config", error instanceof Error ? error : new Error(String(error)), "surface gate uses default (off)", ); return {}; } } /** * D1/D5: flags that apply to HEADLESS worker spawns only. The surface TUI * branch strips them before booting a pane — a pane hosts a full interactive * pi session where the interactive trust prompt must stay available and the * ambient extension stack is part of the human-visible session (task mandate: * worker spawn CHỈ, KHÔNG đụng surface TUI spawn). */ function stripHeadlessOnlyFlags(args: string[]): string[] { return args.filter((arg) => arg !== "--no-approve" && arg !== "--no-extensions"); } /** * MuxSurface A1 spawn branch (spec §13.1). Returns a ChildPiRunResult when the * worker was booted INSIDE a multiplexer pane (no stdio process is spawned and * completion is awaited through the pane's own lifetime), or null to fall * through to the normal stdio spawn. * * NEVER throws: any preparation failure degrades to null (headless §3). */ async function trySurfaceBranch( input: ChildPiRunInput, depthEnv: NodeJS.ProcessEnv | undefined, builtArgs: string[], mergedEnv: NodeJS.ProcessEnv, builtEnv: Record, tempDir: string | undefined, ): Promise { // A pane needs a task id for its title/event-log path — an anonymous worker // has nothing meaningful to boot there, so stay on the headless road. if (!input.agentId) return null; const surfaceOpts = input.surface; const preResolved = surfaceOpts?.providers?.tmux ?? surfaceOpts?.providers?.herdr; // T11 (spec §7): per-run degrade controller — team-runner owns the policy // (lockout counters, pane cap accounting); this layer only consults it. // Absent controller (delegate spawn outside a team run, unit tests) keeps // today's behavior untouched. const taskId = input.agentId; const degradeController = getSurfaceRuntimeController(input.runId); if (degradeController && !degradeController.shouldAttemptSurface()) return null; const surfaceConfig = surfaceOpts?.config ?? safeLoadSurfaceConfig(input.cwd); let outcome: SurfaceSpawnOutcome; try { outcome = await prepareSurfaceSpawn({ // Detection env must be the HOST's base env (depth 0 tier-1), NOT the // child's built env whose PI_CREW_DEPTH is parentDepth+1 — layer-1 // guard would misfire on our own tier-1 workers. env: depthEnv ?? process.env, workerEnv: buildFinalChildPiSpawnOptions(input.cwd, mergedEnv, builtEnv, input.model).env as Record, config: surfaceConfig, role: input.role ?? input.agent.name, livePaneCount: degradeController ? degradeController.livePaneCount() : (surfaceOpts?.livePaneCount ?? 0), taskId, cwd: input.cwd, piArgs: stripHeadlessOnlyFlags(builtArgs), stateRoot: input.eventsPath ? path.dirname(input.eventsPath) : "", baseDir: surfaceOpts?.baseDir, deps: preResolved ? { provider: preResolved } : {}, }); } catch (error) { // Defense in depth — prepareSurfaceSpawn is no-throw by contract (§3). logInternalError("child-pi.surface-unexpected", error instanceof Error ? error : new Error(String(error)), `taskId=${taskId}`); return null; } if (outcome.mode !== "surface") { // Spawn-fail ≠ flap (spec §7) but it feeds ITS OWN lockout counter. Only // `attempted` outcomes count — gate/resolve-null rejections are the mux // simply being absent, not the multiplexer misbehaving. if (outcome.mode === "headless" && outcome.attempted) { degradeController?.notifySpawnFailed({ taskId, reason: outcome.reason ?? "surface spawn failed" }); input.onLifecycleEvent?.({ type: "surface_spawn_failed", surfaceKind: preResolved?.kind, error: outcome.reason, ts: new Date().toISOString(), }); } // FINDING-3 (report 10tier 2026-08-27 adjudication): the gate-null path // used to be completely silent — a misconfigured surface run (wrong // visibleAgents, mode typo, …) was indistinguishable from a by-design // headless run in events.jsonl. Emit ONLY when the user opted into // surface (visibleAgents non-empty) so default runs gain zero noise. // The lockout path above stays silent here — its cause already emitted // surface.degraded events on the way to locking. if (outcome.mode === "headless" && outcome.gateRejected) { const visibleAgents = surfaceConfig.runtime?.surface?.visibleAgents ?? []; if (visibleAgents.length > 0) { input.onLifecycleEvent?.({ type: "surface_gate_blocked", gate: outcome.gateRejected.gate, error: outcome.gateRejected.reason, env: outcome.gateRejected.env, ts: new Date().toISOString(), }); } } return null; } const surfaceMeta: NonNullable = { kind: outcome.kind, paneId: outcome.paneId, scriptPath: outcome.scriptPath, }; degradeController?.notifySpawned({ taskId, paneId: outcome.paneId, provider: outcome.kind, // Task 5 (tab-layout): tab của run lên controller → manifest surface.tabs // — run end (closeRunTabs) đóng đúng tab, worker xong KHÔNG đóng. ...(outcome.tabKey && outcome.tabId ? { tabKey: outcome.tabKey, tabId: outcome.tabId } : {}), }); input.onLifecycleEvent?.({ type: "surface_spawned", surfaceKind: outcome.kind, paneId: outcome.paneId, ts: new Date().toISOString(), }); // T9 (spec §5.3): surface worker không có stdout JSON stream — worker-side // recorder ghi per-agent events.jsonl; host tail file đó và bridge thẳng ra // run event bus cho dashboard/sidebar ("như cũ"). KHÔNG đưa vào // input.onJsonEvent: callback task-runner gọi appendCrewAgentEventBuffered // vào CHÍNH file này (worker đã ghi sẵn) → ghi đôi + vòng tự kích hoạt // watch→append→watch không dừng (xem header event-log-tail-source.ts). const eventSource = outcome.eventsPath ? new EventLogTailSource({ eventsPath: outcome.eventsPath }) : undefined; if (eventSource && input.runId && input.agentId) { const runId = input.runId; const taskId = input.agentId; const controllerForBridge = degradeController; eventSource.onEvent((event) => { // T11: the worker's own `worker.started` self-report carries the pid + // session path that doctor/degrade-resume need (§12.2). Bridge it into // the run-scoped controller BEFORE the dashboard shim — it is the only // host-side consumer of that payload. if ((event as { type?: unknown } | undefined)?.type === "worker.started") { const started = event as { pid?: unknown; sessionPath?: unknown }; controllerForBridge?.notifyWorkerStarted({ taskId, pid: typeof started.pid === "number" ? started.pid : undefined, sessionPath: typeof started.sessionPath === "string" ? started.sessionPath : undefined, }); } const bridgeEvent = bridgeEventFromJsonEvent(runId, taskId, event); if (bridgeEvent) runEventBus.emit({ type: "worker_status", runId, taskId, data: bridgeEvent }); // F1 fix (battery 2026-09-10): feed the same tailed event to the task-state // bridge so child-executor can keep task.heartbeat/agentProgress fresh for // surface workers (see ChildPiRunInput.onSurfaceActivity doc). input.onSurfaceActivity?.(event as Record); }); } let exitInfo: SurfaceExitInfo; try { // Pane lifetime IS the worker lifetime in A1: T8's auto-exit closes the // pane at turn end (spec §13.2); cancel closes it here via AbortSignal and // the response deadline force-closes a WEDGED worker (fix round 1/F2 — // parity with headless minimum timeout so a stuck pane cannot hold its // worker slot forever; A1 has no stream-activity signal so this is ONE // hard deadline). exitInfo = await waitForSurfaceExit(outcome, { signal: input.signal, deadlineMs: input.responseTimeoutMs ?? RESPONSE_TIMEOUT_MS, }); } finally { eventSource?.close(); } // T11 classify (spec §7): D7 report-before-dying means a normally-finishing // worker already wrote `worker.completed` to the RUN event log before its // pane closed, so the probe finds it on the first scan. Nothing within the // 2s window → the worker never finished → the pane death is a DEGRADE and // team-runner re-dispatches this unit headless. Cancel/deadline force-closes // are HOST-initiated — classify is skipped there (never counts as mux flap), // but the pane release itself is reported in EVERY case so the run's live // pane cap stays accurate. let completedPayload: Record | undefined; if (!exitInfo.cancelledByAbort && !exitInfo.timedOut) { const probe = makeTerminalEventProbe({ eventsPath: input.eventsPath ?? "", taskId, runId: input.runId ?? undefined }); const verdict = await classifyOnExit(outcome.handle, probe); const completed = verdict === "completed"; if (completed) { const payload = probe.foundPayload(); completedPayload = payload && typeof payload.result === "string" && payload.result.trim().length > 0 ? payload : undefined; } else { surfaceMeta.degraded = { cause: exitInfo.reason === "mux-dead" ? "mux-dead" : "pane-closed", exitReason: exitInfo.reason, classifiedAt: new Date().toISOString(), }; } degradeController?.notifyPaneExited({ taskId, paneId: outcome.paneId, completed, exitReason: exitInfo.reason, }); } else { degradeController?.notifyPaneExited({ taskId, paneId: outcome.paneId, completed: false, exitReason: exitInfo.reason, cancelledByAbort: exitInfo.cancelledByAbort, timedOut: exitInfo.timedOut, }); } // Result text captured from `worker.completed` — WITHOUT it every surface // task would look empty downstream (rawFinalText/stdout are the only result // channels post-execution trusts) and flip to failed via the empty-output // guard even though the worker SUCCEEDED. const surfaceResultText = typeof completedPayload?.result === "string" ? completedPayload.result : ""; input.onLifecycleEvent?.({ type: "surface_closed", surfaceKind: outcome.kind, paneId: outcome.paneId, paneExitReason: exitInfo.reason, ts: new Date().toISOString(), }); // F1: the pane is gone — nothing can still need the prompt/task temp dir. // The headless branch cleans this up in settle(); surface returns early. cleanupTempDir(tempDir); // Empty stdout/stderr is EXPECTED — events flow through the per-agent event // log (worker-side recorder, spec §5.3), not through a pipe we own. The only // text the worker produced surfaces via `worker.completed` (surfaceResultText). return { exitCode: exitInfo.cancelledByAbort || exitInfo.timedOut ? null : 0, stdout: "", stderr: "", ...(exitInfo.cancelledByAbort ? { error: `Cancelled while running in ${outcome.kind} pane ${outcome.paneId} (${exitInfo.reason})` } : {}), ...(exitInfo.timedOut ? { error: `Surface worker in ${outcome.kind} pane ${outcome.paneId} produced no completion within ` + `${input.responseTimeoutMs ?? RESPONSE_TIMEOUT_MS}ms response timeout; pane was force-closed.`, } : {}), rawFinalText: surfaceResultText, intermediateFindings: "", ...(exitInfo.cancelledByAbort ? { aborted: true } : {}), surface: surfaceMeta, exitStatus: { exitCode: exitInfo.cancelledByAbort || exitInfo.timedOut ? null : 0, cancelled: exitInfo.cancelledByAbort, timedOut: exitInfo.timedOut, killed: false, cleanupErrors: [], finalDrainMs: 0, }, }; } /** * Test seam — trySurfaceBranch là internal wiring giữa prepareSurfaceSpawn và * lifecycle events (surface_gate_blocked). Mock mode intercepts runChildPi * BEFORE this branch, so unit tests exercise the wiring through this export. */ export const __test__trySurfaceBranch = trySurfaceBranch; export async function runChildPi(input: ChildPiRunInput): Promise { // Phase 1 (live-session parity): prepend parent context when inheritContext is true. // This mirrors the effectivePrompt logic in live-session-runtime.ts so that // child-process workers receive the same inherited-context treatment. const effectiveTask = input.inheritContext === true && input.parentContext ? `${input.parentContext}\n\n---\n# Child Worker Task\n${input.task}` : input.task; // ADR-5 §3 depthOverride: delegate-spawned grandchildren carry their depth // from the parent task RECORD. Expressed as the parent's depth in the base // env so BOTH the depth gate and the spawn env builder (parentDepth + 1) // see the parent's true depth — the root process's env (depth 0) is never // consulted for grandchildren. const depthEnv = input.depthOverride !== undefined ? { ...(input.env ?? process.env), PI_CREW_DEPTH: String(input.depthOverride - 1), PI_TEAMS_DEPTH: String(input.depthOverride - 1), } : undefined; const depth = checkCrewDepth(input.maxDepth, depthEnv); if (depth.blocked) return { exitCode: 1, stdout: "", stderr: `pi-crew depth guard blocked child worker: depth ${depth.depth} >= max ${depth.maxDepth}`, }; // H3 phase 3 (2026-08-10): mock-mode fixtures extracted to mock-fixtures.ts. // Returns undefined when mock mode is NOT active → fall through to spawn. const mockResult = await runMockChildPi(input, effectiveTask, observeStdoutChunk); if (mockResult) return mockResult; // H-7 step 6: spawn/env/args preparation extracted to child-pi-spawn.ts. // prepareSpawnContext builds the worker args, attaches the steering file env, // and handles the pre-spawn abort check (returns an immediate-abort result // if the parent signal has already fired). // // Phase 0 broker: if the caller did not pre-fill `brokerSpawn`, ask the // optional issuer for one. The issuer is gated by the lifecycle controller // (root-session + flag); a no-op issuer yields undefined (no credentials). let brokerSpawn = input.brokerSpawn; const brokerIssuer = input.brokerIssuer ?? getActiveBrokerIssuer(); if (!brokerSpawn && brokerIssuer && input.runId) { try { // ADR-5 §4: thread the child depth so the issuer can contain broker // credentials at depths that may not delegate (a child at the default // maxDepth=4 cap gets NO socket/token — env containment AC). brokerSpawn = await brokerIssuer(input.runId, input.agentId, input.depthOverride); } catch (error) { // H8 (2026-08-10): surface the silent degradation. Previously this // swallowed ALL issuer failures (token-rotation race, broker socket // down, key fetch network error) with zero observability — the child // spawned without broker credentials and the run just looked "slow". // logInternalError writes to the internal-error channel (sampled + // bounded); it does NOT propagate, so the child still runs without // broker acceleration (the durable-first invariant is preserved). logInternalError( "child-pi.broker-issuer-failed", error instanceof Error ? error : new Error(String(error)), `runId=${input.runId} agentId=${input.agentId ?? "?"} — child will spawn without broker credentials`, ); brokerSpawn = undefined; } } // U9 (pi 1.0.4 version floor): resolve the host pi binary's version ONCE per // process (probePiVersion — single-flight memoized, injectable exec seam, // never rejects) and thread it into the argv builder via input. Probed ONLY // when this spawn can consume a version-gated flag: the parity // (non-hermetic) argv applies the --no-mcp floor at 1.0.4, while hermetic // spawns (the default) are already MCP-clean via --no-extensions and consult // no version-gated flag — skipping the probe there keeps every default-path // spawn — and the unit-test surface paths — free of a real `pi --version` // child. An already-aborted signal skips too: prepareSpawnContext's B5 // guard below returns before the argv is ever consumed. const versionedInput: ChildPiRunInput = !input.signal?.aborted && !resolveHermeticWorkers(input.hermeticWorkers) ? { ...input, piVersion: await probePiVersion() } : input; const spawnPrep = prepareSpawnContext(brokerSpawn ? { ...versionedInput, brokerSpawn } : versionedInput, effectiveTask, depthEnv); if (spawnPrep.kind === "aborted") return spawnPrep.result; const { spawnSpec, mergedEnv, tempDir, builtEnv, builtArgs } = spawnPrep.ctx; // W2 (session-file recovery): resolve the worker's deterministic session // identity (per-worker dir under the run artifacts root) and arm it on BOTH // argv views — the headless spawnSpec.args ([script, ...builtArgs]) and the // raw builtArgs (also consumed by the surface branch). prepareSpawnContext // does not forward session identity to the arg builder (ownership boundary // with child-pi-spawn.ts — see the lane report), so this is the in-zone // activation point; appendWorkerSessionArgs is idempotent, making the move // to builder-side forwarding later a zero-diff no-op. const workerSession = deriveSessionPaths(input); const sessionRecoveryEnabled = resolveSessionRecoveryEnabled(input.sessionRecovery); if (workerSession) { try { fs.mkdirSync(workerSession.sessionDir, { recursive: true }); } catch { // Unwritable target → spawn proceeds; recovery simply finds no files. } appendWorkerSessionArgs(spawnSpec.args, builtArgs, workerSession); } // MuxSurface A1 (spec §13.1): attempt booting the worker in a mux pane BEFORE // spawning a stdio pipe. Returns a finished result in surface mode (pane exit // awaited) or null → fall through to the classic headless spawn below. const surfaceResult = await trySurfaceBranch(input, depthEnv, builtArgs, mergedEnv, builtEnv, tempDir); if (surfaceResult) return surfaceResult; try { return await new Promise((resolve) => { // Compose the final SpawnOptions: canary + filter + spread are now // owned by buildFinalChildPiSpawnOptions (see child-pi-spawn.ts, BLOCKER 2 / S5). const spawnOptions = buildFinalChildPiSpawnOptions(input.cwd, mergedEnv, builtEnv, input.model); const child = spawn(spawnSpec.command, spawnSpec.args, spawnOptions); if (child.pid) { registerActiveChild(child.pid, child); input.onSpawn?.(child.pid); input.onLifecycleEvent?.({ type: "spawned", pid: child.pid, ts: new Date().toISOString(), }); // Register with cleanup handler for graceful shutdown (RT-11): // ALWAYS register — every spawned child must be visible to host-SIGTERM // cleanup. Use synthetic IDs when the caller omitted runId/agentId so the // PID is still tracked for killProcessPid on session shutdown. registerChildProcess( child.pid, input.runId ?? `untracked-run-${child.pid}`, input.agentId ?? `untracked-agent-${child.pid}`, ); } else { input.onLifecycleEvent?.({ type: "spawn_error", error: "spawn returned no pid", ts: new Date().toISOString(), }); } // P0-1: O(1)-amortized bounded accumulators (segment ring) — replaces the O(n²) // appendBoundedTail rebuild-per-line pattern that re-scanned 512 KiB every line. const stdoutTail = new BoundedTail(); const stderrTail = new BoundedTail(); let settled = false; let childExited = false; let postExitGuardCleanup: (() => void) | undefined; const finalDrainMs = input.finalDrainMs ?? FINAL_DRAIN_MS; const hardKillMs = input.hardKillMs ?? HARD_KILL_MS; // Phase-0 diagnostic (HB-003a): track the final-drain race that produces // `exit null` for ctx.agent({disableTools:true}). These vars are READ-ONLY // instrumentation — no behavior change. finalDrainArmed lets the close // handler know a drain timer existed even after clearFinalDrainTimers() ran; // spawnMonotonicMs gives us relative timing to distinguish a race from a crash. let finalDrainArmed = false; // F12: monotonic timestamp of the last stdout JSON event (any event — // we want to know when stdout *stopped*, not when the final assistant // event arrived). Updated on every onJsonEvent dispatch. let lastStdoutActivityMonotonicMs = performance.now(); let finalDrainFiredMonotonicMs: number | undefined; const spawnMonotonicMs = performance.now(); let finalAssistantEventMonotonicMs: number | undefined; // FIX (Round 14): Bound the env-controlled response timeout to // [1_000ms, 3_600_000ms] (1s–1h) so a hostile or accidental value // (e.g. 1, or 999_999_999) cannot disable the timeout or cause // instant kills. Out-of-range values fall back to the input or // built-in default. const RESPONSE_TIMEOUT_MIN_MS = 1_000; const RESPONSE_TIMEOUT_MAX_MS = 3_600_000; const responseTimeoutEnv = Number.parseInt(getCrewEnv("PI_TEAMS_CHILD_RESPONSE_TIMEOUT_MS") ?? "", 10); const envInRange = Number.isFinite(responseTimeoutEnv) && responseTimeoutEnv >= RESPONSE_TIMEOUT_MIN_MS && responseTimeoutEnv <= RESPONSE_TIMEOUT_MAX_MS; const responseTimeoutMs = envInRange ? responseTimeoutEnv : (input.responseTimeoutMs ?? RESPONSE_TIMEOUT_MS); let responseTimeoutHit = false; let forcedFinalDrain = false; let abortRequested = input.signal?.aborted === true; let hardKilled = false; const cleanupErrors: string[] = []; const steeringController = new ChildPiSteeringController(input.maxTurns, input.graceTurns); let abortDueToParentSignal = false; // CP-1: track whether the turn-limit hard-abort has been initiated. Once // true, we must NOT restart the no-response timer — the child is already // being killed via killProcessTree (SIGTERM → SIGKILL after 3s), and // restarting the timer would delay detection of a SIGTERM-ignoring child. // Round 27 (BUG 4): extract to a named handler so settle() can remove it. // The previous anonymous listener was never removed → on runs with >10 // tasks sharing one AbortSignal (background-runner), Node emitted // MaxListenersExceededWarning and each leaked listener pinned the task's // stack frame (abortDueToParentSignal closure) in memory. { once: true } // only auto-removes AFTER the signal fires; on normal completion it leaks. const onParentAbort = (): void => { abortDueToParentSignal = true; }; input.signal?.addEventListener("abort", onParentAbort, { once: true, }); // Phase 2.3: the six timer constructs (noResponseTimer, finalDrainTimer, // hardKillTimer, safetyTimer, cancelHardKill, pollHandle) live in // child-pi-timers.ts. Mutable flags are shared via getters/setters so // the timer callbacks observe the same values as the handlers below. // The methods are destructured so the call sites keep their original // source shape (HB-003a source-contract test asserts the steering guard // directly precedes `restartNoResponseTimer()`). const { restartNoResponseTimer, clearNoResponseTimer, clearFinalDrainTimers, armFinalDrain, hasFinalDrainTimer, armCancelHardKill, clearAll, } = createChildPiTimers({ child, input, responseTimeoutMs, finalDrainMs, hardKillMs, stdoutTail, stderrTail, cleanupErrors, getSettle: () => settle, redactStderrExcerpt, state: { getSettled: () => settled, getChildExited: () => childExited, setResponseTimeoutHit: (value) => { responseTimeoutHit = value; }, getHardKilled: () => hardKilled, setHardKilled: (value) => { hardKilled = value; }, setForcedFinalDrain: (value) => { forcedFinalDrain = value; }, getLastStdoutActivityMonotonicMs: () => lastStdoutActivityMonotonicMs, setFinalDrainFiredMonotonicMs: (value) => { finalDrainFiredMonotonicMs = value; }, getAbortRequested: () => abortRequested, }, }); restartNoResponseTimer(); const lineObserver = new ChildPiLineObserver({ ...input, onStdoutLine: (line) => { if (!steeringController.isHardAbortInitiated()) restartNoResponseTimer(); stdoutTail.push(`${line}\n`); input.onStdoutLine?.(line); }, onJsonEvent: (event) => { if (!steeringController.isHardAbortInitiated()) restartNoResponseTimer(); // Turn-count-based steering: soft limit steer + hard abort after graceTurns if (event && typeof event === "object" && !Array.isArray(event)) { const obj = event as Record; if (obj.type === "turn_end") { // H-7 step 5: steering state machine extracted to ChildPiSteeringController. const action = steeringController.onTurnEnd(child.pid, child, input.steeringFile); if (action.kind === "hardAbort") killProcessTree(action.pid, action.child); } } // F12: capture monotonic timestamp BEFORE dispatching — any stdout // JSON event counts as activity. This lets the quiet-window // detection measure "time since last byte of stdout" accurately // regardless of what onJsonEvent does. lastStdoutActivityMonotonicMs = performance.now(); input.onJsonEvent?.(event); if (!isFinalAssistantEvent(event) || childExited || settled || hasFinalDrainTimer()) return; finalAssistantEventMonotonicMs = performance.now(); finalDrainArmed = true; // Phase-0 diagnostic: track that a drain timer was created. armFinalDrain(); }, }); const clearPostExitGuard = (): void => { if (postExitGuardCleanup) { postExitGuardCleanup(); postExitGuardCleanup = undefined; } }; const clearChildPiTimeouts = (): void => { // R6-F1: clearAll() covers all six timer constructs // (incl. cancelHardKill — previously a local const inside abort() // that leaked past settle on short-lived children). clearAll(); clearPostExitGuard(); }; // W2 (session-file recovery): crash-path augment — when the worker died // WITHOUT a final assistant event (exitCode null / killed), tail-replay // its session JSONL and surface the last COMPLETE assistant record on the // resolved result. Best-effort and deadline-bounded (session-recovery.ts) // so settle can never hang on pathological IO. This AUGMENTS the result — // manifest polling (heartbeat-watcher / detached-run-results) is untouched. const recoverForSettle = async (settleResult: ChildPiRunResult): Promise => { if (!workerSession || !sessionRecoveryEnabled) return undefined; if (!shouldAttemptSessionRecovery(settleResult, hardKilled)) return undefined; try { return (await recoverLastAssistantFromSession(workerSession.sessionDir, workerSession.sessionId)) ?? undefined; } catch { return undefined; } }; const settle = (result: ChildPiRunResult): Promise => { if (settled) return Promise.resolve(); settled = true; clearChildPiTimeouts(); // OPT-06 follow-up: lineObserver.flush() is now async (returns // Promise) and drains the module-scoped transcript batch buffer // before resolving. We must await it before calling `resolve()` // below so callers that read the transcript file post-`runChildPi` // see all written lines. Caller invocations of `settle` from // sync event handlers (`child.on('close'|'exit'|'error')`, // safety timer) use `void settle(...)` — errors are caught and // logged inside, and `resolve()` only fires after the drain // completes, so runChildPi's outer Promise resolves with a // durable transcript on disk. return lineObserver .flush() .then(async () => { input.signal?.removeEventListener("abort", abort); input.signal?.removeEventListener("abort", onParentAbort); try { cleanupTempDir(tempDir); } catch (error) { cleanupErrors.push(error instanceof Error ? error.message : String(error)); } // W2: crash-path session tail-recovery (await BEFORE resolve so the // recovered record rides the same settled result object). const recoveredFromSession = await recoverForSettle(result); // Catch all errors from settle to prevent unhandled rejection from propagating try { resolve({ ...result, rawFinalText: lineObserver.getRawFinalText(), intermediateFindings: lineObserver.getIntermediateFindings(), ...(recoveredFromSession ? { recoveredFromSession } : {}), exitStatus: result.exitStatus ?? { exitCode: result.exitCode, cancelled: abortRequested, timedOut: responseTimeoutHit, killed: hardKilled, // Phase-0 diagnostic (HB-003a): surface the final-drain race state. // finalDrainArmed lets Phase 1 decide whether a signal-death (exitCode=null) // should be treated as a forced final drain. READ-ONLY for now. ...(finalDrainArmed || forcedFinalDrain ? { finalDrainArmed, forcedFinalDrain, finalDrainFiredMonotonicMs, } : {}), cleanupErrors, finalDrainMs, }, }); } catch (resolveError) { logInternalError( "child-pi.settle-resolve", resolveError, `result=${JSON.stringify({ exitCode: result.exitCode })}`, ); } }) .catch(async (flushError) => { // Drain failed — log and still resolve so runChildPi doesn't hang. logInternalError( "child-pi.settle-flush-failed", flushError, `result=${JSON.stringify({ exitCode: result.exitCode })}`, ); input.signal?.removeEventListener("abort", abort); input.signal?.removeEventListener("abort", onParentAbort); try { cleanupTempDir(tempDir); } catch (error) { cleanupErrors.push(error instanceof Error ? error.message : String(error)); } // W2: crash-path session tail-recovery — same augment on the // drain-failed path (a torn flush is itself crash-adjacent). const recoveredFromSession = await recoverForSettle(result); try { resolve({ ...result, rawFinalText: lineObserver.getRawFinalText(), intermediateFindings: lineObserver.getIntermediateFindings(), ...(recoveredFromSession ? { recoveredFromSession } : {}), exitStatus: result.exitStatus ?? { exitCode: result.exitCode, cancelled: abortRequested, timedOut: responseTimeoutHit, killed: hardKilled, ...(finalDrainArmed || forcedFinalDrain ? { finalDrainArmed, forcedFinalDrain, finalDrainFiredMonotonicMs, } : {}), cleanupErrors, finalDrainMs, }, }); } catch (resolveError) { logInternalError( "child-pi.settle-resolve", resolveError, `result=${JSON.stringify({ exitCode: result.exitCode })}`, ); } }); }; const abort = (): void => { abortRequested = true; clearNoResponseTimer(); killProcessTree(child.pid, child); if (process.platform !== "win32") { trySignalChild(child, "SIGTERM"); } try { child.kill(process.platform === "win32" ? undefined : "SIGTERM"); } catch { // Ignore kill races. } // 3.5 — fast-escalate to SIGKILL within 200ms on explicit cancel // so /team-cancel completes round-trip well under the operator // expectation. The standard finalDrainMs / HARD_KILL_MS paths // are for graceful drain, not user-initiated cancel. R6-F1: the // timer handle is owned by child-pi-timers.ts and cleared by // clearAll() on settle (was previously an unreachable local const). armCancelHardKill(); }; input.signal?.addEventListener("abort", abort, { once: true }); // 3.1 — soft watermark backpressure. When inbound stdout exceeds // 256KB before the next macrotask, pause for 50ms so the line // observer + ancillary handlers get to drain. Prevents the runaway // case where a chatty child saturates the parent event loop. const BACKPRESSURE_HIGH = 256 * 1024; let backpressureBytes = 0; const releaseBackpressure = (): void => { backpressureBytes = 0; try { child.stdout?.resume(); } catch { /* ignore */ } }; child.stdout?.on("data", (chunk: Buffer) => { if (!steeringController.isHardAbortInitiated()) restartNoResponseTimer(); const text = chunk.toString("utf-8"); backpressureBytes += text.length; try { lineObserver.observe(text); } catch (err) { logInternalError("child-pi.line-observer-observe", err, `text=${text.slice(0, 100)}`); } if (backpressureBytes > BACKPRESSURE_HIGH && child.stdout && !child.stdout.isPaused()) { try { child.stdout.pause(); } catch { /* ignore */ } const timer = setTimeout(releaseBackpressure, 50); timer.unref(); } }); child.stderr?.on("data", (chunk: Buffer) => { if (!steeringController.isHardAbortInitiated()) restartNoResponseTimer(); stderrTail.push(chunk.toString("utf-8")); }); child.on("error", (error) => { // P0-1: snapshot the bounded accumulators once for this handler. const stdout = stdoutTail.value(); const stderr = stderrTail.value(); // SEC-1: redact stderr secrets embedded in the error message + excerpt. const processError = new Error( `Child Pi process error: ${error.message}. Stderr: ${redactStderrExcerpt(stderr, 500) || "(none)"}`, ); try { input.onLifecycleEvent?.({ type: "spawn_error", pid: child.pid, error: processError.message, ts: new Date().toISOString(), stderrExcerpt: redactStderrExcerpt(stderr, 500) || undefined, }); } catch (err) { logInternalError("child-pi.on-lifecycle-event", err, `event=error, pid=${child.pid}`); } void settle({ exitCode: null, stdout, stderr, error: processError.message, exitStatus: { exitCode: null, cancelled: abortRequested, timedOut: responseTimeoutHit, killed: false, cleanupErrors, finalDrainMs, crashClass: classifyProcessCrash({ exitCode: null, cancelled: abortRequested, timedOut: responseTimeoutHit, spawnError: error, stderrSnippet: stderr ? redactStderrExcerpt(stderr, 1000) : undefined, }).crashClass, }, }); }); child.on("exit", (code, signal) => { // P0-1: snapshot the bounded stderr accumulator once for this handler. const stderr = stderrTail.value(); if (child.pid) { unregisterActiveChild(child.pid); clearHardKillTimer(child.pid); // Unregister from cleanup handler unregisterChildProcess(child.pid); } // Build comprehensive exit error for unexpected exits // Round-10 test fix: also require non-zero exit code OR a known abnormal condition. // Previously fired "exited unexpectedly" on every clean exit (code=0) because the // OS-level 'exit' event fires BEFORE pi's 'agent_end' JSON event reaches the line // observer (race). Worker actually succeeded but onLifecycleEvent reported an error. const abnormalExit = code !== 0 && code !== null; const isUnexpectedExit = !childExited && !settled && !responseTimeoutHit && !abortRequested && abnormalExit; const exitError = isUnexpectedExit ? new Error( `Child Pi process exited unexpectedly (code=${code ?? "null"} signal=${signal ?? "null"}). ` + `Stderr: ${redactStderrExcerpt(stderr, 1000) || "(none)"}`, ) : null; try { // Phase-0 diagnostic (HB-003a): capture signal + drain timing in the // exit lifecycle event so the exit-null race is diagnosable instead of // opaque. `signal` was previously discarded after building the error msg. input.onLifecycleEvent?.({ type: "exit", pid: child.pid, exitCode: code, ts: new Date().toISOString(), error: exitError?.message, stderrExcerpt: isUnexpectedExit ? redactStderrExcerpt(stderr, 1000) || undefined : undefined, // Phase-0 diagnostic fields (kept optional — no type change required). ...(signal ? { signal } : {}), ...(finalDrainArmed || forcedFinalDrain ? { diagnostic: { finalDrainArmed, forcedFinalDrain, finalDrainFiredMonotonicMs, finalAssistantEventMonotonicMs, exitMonotonicMs: performance.now() - spawnMonotonicMs, }, } : {}), }); } catch (err) { logInternalError("child-pi.on-lifecycle-event", err, `event=exit, pid=${child.pid}`); } childExited = true; clearNoResponseTimer(); clearFinalDrainTimers(); if (!postExitGuardCleanup) { postExitGuardCleanup = attachPostExitStdioGuard(child, { idleMs: POST_EXIT_STDIO_GUARD_MS, hardMs: HARD_KILL_MS, }); } }); child.on("close", (exitCode) => { // P0-1: snapshot the bounded accumulators once for this handler. const stdout = stdoutTail.value(); const stderr = stderrTail.value(); if (child.pid) { unregisterActiveChild(child.pid); clearHardKillTimer(child.pid); // Unregister from cleanup handler unregisterChildProcess(child.pid); } try { input.onLifecycleEvent?.({ type: "close", pid: child.pid, exitCode, ts: new Date().toISOString(), }); } catch (err) { logInternalError("child-pi.on-lifecycle-event", err, `event=close, pid=${child.pid}`); } const timeoutError = responseTimeoutHit && !stderr.trim() ? { error: `Child Pi produced no new output for ${responseTimeoutMs}ms; process was terminated as unresponsive.`, } : responseTimeoutHit && stderr.trim() ? { error: `Child Pi timed out after ${responseTimeoutMs}ms with stderr: ${redactStderrExcerpt(stderr, 500)}`, } : undefined; // M6 fix: log when forced final drain converts non-zero exit to 0. // This is expected in normal operation (child finished cleanly but linger was killed), // but the telemetry helps detect regressions where crashes are hidden. if (forcedFinalDrain && !timeoutError && exitCode !== 0) { logInternalError( "child-pi.final-drain-zero-exit", new Error(`Child exit code overridden to 0 after forced final drain (original=${exitCode})`), `pid=${child.pid}, finalDrainMs=${finalDrainMs}`, ); } const finalExitCode = forcedFinalDrain && !timeoutError ? 0 : exitCode; const wasGraceAborted = steeringController.isSoftLimitReached() && steeringController.getTurnCount() >= (steeringController.getMaxTurns() ?? 0) + (steeringController.getGraceTurns() ?? 5); const wasParentAborted = abortDueToParentSignal && !wasGraceAborted; // P0 crash taxonomy: classify the exit so callers/dashboards can bucket // failure modes (timeout vs cancel vs native panic vs signal …). // The classifier is a pure function; this is the single integration point. const crashClassification = classifyProcessCrash({ exitCode: finalExitCode, signal: child.signalCode ?? undefined, cancelled: abortRequested, timedOut: responseTimeoutHit, killed: hardKilled, spawnError: undefined, stderrSnippet: stderr ? redactStderrExcerpt(stderr, 1000) : undefined, }); void settle({ exitCode: finalExitCode, stdout, stderr, ...(timeoutError ? { error: timeoutError.error } : {}), aborted: wasGraceAborted || wasParentAborted, steered: steeringController.isSoftLimitReached() && !wasGraceAborted, exitStatus: { exitCode: finalExitCode, cancelled: abortRequested, timedOut: responseTimeoutHit, killed: hardKilled, cleanupErrors, finalDrainMs, crashClass: crashClassification.crashClass, }, }); }); }); } finally { // cleanupTempDir is already called inside settle(), but guard against // the case where settle() was never reached (spawn throws synchronously). if (tempDir && fs.existsSync(tempDir)) { cleanupTempDir(tempDir); } } }