import * as fs from "node:fs"; import * as path from "node:path"; import type { ReleaseBlocker, ReleaseReadinessRequest, StopAllBackgroundItem, StopAllBackgroundRequest } from "@selesai/code"; import { hasLiveSubagentWork } from "../integrations/pi-web-session-liveness.ts"; import { reconcileAsyncRun } from "../runs/background/stale-run-reconciler.ts"; import { resultCandidateFilesForSession, resultPayloadPathForSessionRun } from "../runs/background/result-files.ts"; import { hasLiveNestedDescendants } from "../runs/shared/nested-events.ts"; import { DIRS, type AsyncJobState, type NestedRunSummary, type SubagentState } from "../shared/types.ts"; interface Events { on(name: string, handler: (payload: unknown) => void): (() => void) | void } /** * Mirrors `RELEASE_READINESS_EVENT` / `answerReleaseReadiness` from the host package. Kept local (types only are imported) so this * module stays loadable under the test shim, like the other event-name constants in this extension; a host-side test pins them together. */ export const RELEASE_READINESS_EVENT = "release:readiness:v1"; function answerReleaseReadiness(raw: unknown, name: string, boundSessionId: string | null | undefined, compute: () => ReleaseBlocker[]): void { const request = raw as Partial | null | undefined; if (!request || request.version !== 1 || typeof request.contribute !== "function" || !boundSessionId || request.sessionId !== boundSessionId) return; let blockers: ReleaseBlocker[]; try { blockers = compute(); } catch (error) { blockers = [{ source: name, kind: "contributor_error", detail: (error instanceof Error ? error.message : String(error)).slice(0, 512) || "contributor failed" }]; } request.contribute(name, blockers); } const TERMINAL = new Set(["complete", "failed", "partial", "paused", "stopped", "rejected"]); const MAX_RESULT_BLOCKERS = 32; const message = (error: unknown) => error instanceof Error ? error.message : String(error); const stamp = (value: unknown): number | undefined => typeof value === "number" && Number.isFinite(value) && value >= 0 ? Math.trunc(value) : undefined; export interface SubagentReleaseDeps { state: SubagentState; /** `completionNotifier.hasPendingDelivery()`: completions queued for the parent but not yet delivered. */ hasPendingDelivery(): boolean; /** Same predicate the result watcher uses to decide whether delivery still matters (`hasResultDeliveryDemand`). */ hasDeliveryDemand?(): boolean; /** Run ids with a persisted result this process still owes the parent. */ listUndeliveredResults(): string[]; resultsDir?: string; /** Injectable for tests; defaults to the real dead-runner repair. */ reconcile?: (asyncDir: string, options: Parameters[1]) => ReturnType; } export interface SchedulerReleaseDeps { armedSchedules(): Array<{ id: string; nextRunAt?: number; sessionOnly?: boolean }>; observedCompletionRunIds(): Iterable; } /** * Persisted results for the given sessions that this process still owes the parent. Results owned by another * live (or unprovable) process are skipped; a result whose previous owner is provably dead is reported because this * process will take it over and deliver it (`canTakeOver` is read-only: nothing is claimed here). Unreadable results * fail closed. */ export function findUndeliveredResults(options: { resultsDir: string; sessionIds: Array; owns(sessionId: string, completionOwnerId: unknown): boolean; canTakeOver?(sessionId: string, completionOwnerId: unknown, runKey: string): boolean; }): string[] { const runIds = new Set(); for (const sessionId of options.sessionIds) { if (!sessionId) continue; for (const file of resultCandidateFilesForSession(options.resultsDir, sessionId)) { const runId = file.replace(/\.json$/i, ""); let data: Record | undefined; try { const payload = resultPayloadPathForSessionRun(options.resultsDir, sessionId, runId); if (payload) data = JSON.parse(fs.readFileSync(payload, "utf-8")) as Record; } catch { runIds.add(runId); continue; } if (!data || typeof data !== "object") continue; if (typeof data.notificationDeliveredAt === "number") continue; if (typeof data.sessionId === "string" && !options.owns(data.sessionId, data.completionOwnerId) && !options.canTakeOver?.(data.sessionId, data.completionOwnerId, runId)) continue; runIds.add(runId); } } return [...runIds]; } /** * Blockers for background subagent work. Reconciles dead runners first (decision D5: a dead runner is not a * blocker; its repaired failure result is an `undelivered_result` until delivered). Fail closed everywhere. */ export function subagentReleaseBlockers(deps: SubagentReleaseDeps): ReleaseBlocker[] { const { state } = deps; const blockers: ReleaseBlocker[] = []; const add = (kind: string, detail: string, id?: string, since?: number) => { const at = stamp(since); blockers.push({ source: "subagents", kind, detail: detail.slice(0, 512), ...(id ? { id: id.slice(0, 256) } : {}), ...(at === undefined ? {} : { since: at }) }); }; const guard = (label: string, run: () => void) => { try { run(); } catch (error) { add("contributor_error", `${label}: ${message(error)}`); } }; const reconcile = deps.reconcile ?? reconcileAsyncRun; const resultsDir = deps.resultsDir ?? DIRS.results; const live = new Map(); const blockedRunIds = new Set(); const reconciledTerminal = new Set(); for (const job of state.asyncJobs.values()) { if (job.status !== "queued" && job.status !== "running") continue; let settled = false; try { const reconciled = reconcile(job.asyncDir, { resultsDir, startedRun: { runId: job.asyncId, pid: job.pid, sessionId: job.sessionId, completionOwnerId: job.completionOwnerId, mode: job.mode, agents: job.agents, chainStepCount: job.chainStepCount, parallelGroups: job.parallelGroups, startedAt: job.startedAt, sessionFile: job.sessionFile, }, }); if (reconciled.status && TERMINAL.has(reconciled.status.state)) { settled = true; reconciledTerminal.add(job.asyncId); add("undelivered_result", reconciled.repaired ? `Async run ${job.asyncId} ended without a live runner (${reconciled.status.state}); its result has not been delivered` : `Async run ${job.asyncId} finished (${reconciled.status.state}); the parent session has not observed it yet`, job.asyncId, reconciled.status.endedAt ?? job.updatedAt); blockedRunIds.add(job.asyncId); } } catch (error) { add("contributor_error", `Failed to reconcile async run ${job.asyncId}: ${message(error)}`, job.asyncId); } if (!settled) live.set(job.asyncId, job); } for (const job of live.values()) { add("async_run", `Async ${job.mode ?? "run"} ${job.status}${job.agents?.length ? `: ${job.agents.join(", ")}` : ""}`, job.asyncId, job.startedAt); blockedRunIds.add(job.asyncId); } guard("nested descendants", () => { for (const job of state.asyncJobs.values()) { if (live.has(job.asyncId) || !hasLiveNestedDescendants(job.nestedChildren)) continue; add("async_run", `Async run ${job.asyncId} has live nested descendants`, job.asyncId, job.startedAt); blockedRunIds.add(job.asyncId); } }); for (const control of state.foregroundControls.values()) { add("foreground_run", `Foreground ${control.mode} run in progress${control.currentAgent ? `: ${control.currentAgent}` : ""}`, control.runId, control.startedAt); } for (const rootRunId of state.retainedForegroundNestedRoutes?.keys() ?? []) { add("async_run", `Retained nested route for foreground run ${rootRunId} may still have live descendants`, rootRunId); } guard("completion notifier", () => { if (deps.hasPendingDelivery()) add("undelivered_result", "Completion notifications are queued for delivery to the parent session"); }); guard("undelivered results", () => { const runIds = deps.listUndeliveredResults().filter((runId) => !blockedRunIds.has(runId)); for (const runId of runIds.slice(0, MAX_RESULT_BLOCKERS)) add("undelivered_result", `Result for run ${runId} has not been delivered to the parent session`, runId); if (runIds.length > MAX_RESULT_BLOCKERS) add("undelivered_result", `${runIds.length - MAX_RESULT_BLOCKERS} further undelivered results`); }); // Catch-all so the total is never under-reported against the liveness predicate used by pi-web // (`hasLiveSubagentWork(state) || hasPendingDelivery()`), evaluated after reconciliation. guard("liveness", () => { const afterReconcile = { asyncJobs: new Map([...state.asyncJobs].filter(([id]) => !reconciledTerminal.has(id))), foregroundControls: state.foregroundControls, retainedForegroundNestedRoutes: state.retainedForegroundNestedRoutes }; const unattributed = hasLiveSubagentWork(afterReconcile) && !blockers.some((b) => b.kind === "async_run" || b.kind === "foreground_run"); if (unattributed) add("async_run", "Live subagent work reported by the liveness predicate could not be attributed to a specific run"); }); guard("delivery demand", () => { if (!blockers.length && deps.hasDeliveryDemand?.()) add("undelivered_result", "Result delivery is still pending for this session"); }); return blockers; } /** * Armed SESSION-ONLY schedules (decisions D7 + D11: project schedules persist on disk and do not pin a process) and * scheduled runs still awaiting completion (live work, whatever the schedule kind). An armed timer whose record could not * be read (`sessionOnly` undefined) fails closed and blocks. */ export function schedulerReleaseBlockers(deps: SchedulerReleaseDeps): ReleaseBlocker[] { const blockers: ReleaseBlocker[] = []; for (const schedule of deps.armedSchedules()) { if (schedule.sessionOnly === false) continue; const next = stamp(schedule.nextRunAt); blockers.push({ source: "scheduler", kind: "armed_schedule", id: schedule.id.slice(0, 256), detail: `Armed ${schedule.sessionOnly === true ? "session-only" : "unreadable"} schedule${next === undefined ? "" : `, next run ${new Date(next).toISOString()}`}`, ...(next === undefined ? {} : { since: next }), }); } for (const runId of deps.observedCompletionRunIds()) { blockers.push({ source: "scheduler", kind: "scheduled_run_pending", id: runId.slice(0, 256), detail: `Scheduled run ${runId} is awaiting completion` }); } return blockers; } /** Mirrors `STOP_ALL_BACKGROUND_EVENT` / `answerStopAllBackground` from the host package (same reason as the readiness mirror above). */ export const STOP_ALL_BACKGROUND_EVENT = "release:stop-all:v1"; function answerStopAllBackground( raw: unknown, name: string, boundSessionId: string | null | undefined, compute: (request: { deadline: number }) => StopAllBackgroundItem[] | PromiseLike, ): void { const request = raw as Partial | null | undefined; if (!request || request.version !== 1 || typeof request.contribute !== "function" || !boundSessionId || request.sessionId !== boundSessionId) return; const deadline = typeof request.deadline === "number" && Number.isFinite(request.deadline) ? request.deadline : Date.now() + 5_000; let items: StopAllBackgroundItem[] | PromiseLike; try { items = compute({ deadline }); } catch (error) { const detail = message(error).slice(0, 400) || "contributor failed"; items = [{ source: name, kind: "contributor", detail: `Stop failed: ${detail}`, ok: false, error: detail.slice(0, 256) }]; } request.contribute(name, items); } export interface SubagentStopDeps { state: SubagentState; /** Run ids a scheduled run is still awaiting (`ScheduledRunManager.observedCompletionRunIds`): the `scheduled_run_pending` readiness blockers. */ observedCompletionRunIds(): Iterable; /** * Stops ONE async run through the machinery behind the subagents RPC `stop` method (`stopAsyncRun` in rpc.ts): checks that the * run belongs to the active session and drops a stop request into the run's control inbox, which the detached runner consumes by * itself (nothing here depends on this process staying alive). Throws when the request is refused. */ stopRun(target: { dir: string }): unknown; asyncDirRoot?: string; resultsDir?: string; reconcile?: SubagentReleaseDeps["reconcile"]; now?(): number; sleep?(ms: number): Promise; pollIntervalMs?: number; } export interface SchedulerStopDeps { /** `ScheduledRunManager.disarmSessionOnlySchedules`: pauses (persisted) and un-arms this session's session-only schedules; project schedules untouched. */ disarmSessionSchedules(): Array<{ id: string; ok: boolean; detail: string; error?: string }>; } const liveNested = (children: NestedRunSummary[] | undefined, found: Array<{ id: string; asyncDir: string }> = []): Array<{ id: string; asyncDir: string }> => { for (const child of children ?? []) { if (!TERMINAL.has(child.state) && child.asyncDir) found.push({ id: child.id, asyncDir: child.asyncDir }); liveNested(child.children, found); liveNested(child.steps?.flatMap((step) => step.children ?? []), found); } return found; }; /** * Kill switch for background subagent work of this session. Stops exactly what the readiness contributor reports as * `async_run` / `scheduled_run_pending`: every live async run of the session (including detached runners), runs a scheduled * run is awaiting, and live nested descendants; then waits (bounded by `request.deadline`) until each run's status is terminal. * Foreground runs belong to the parent turn and are killed by the core abort (their tool signal); they are only awaited. * The synchronous prefix (enumeration + delivering every stop request) runs before the first `await`, so the requests are on disk * even if the process exits right after. Per-run failures are `ok:false` items and never throw. Runs that already ended * (including a dead runner repaired by reconciliation) are not reported: there is nothing left to stop, which also makes the call idempotent. */ export async function stopSubagentWork(deps: SubagentStopDeps, request: { deadline: number }): Promise { const { state } = deps; const items: StopAllBackgroundItem[] = []; const add = (kind: string, detail: string, ok: boolean, id?: string, error?: string) => { items.push({ source: "subagents", kind, detail: detail.slice(0, 512), ok, ...(id ? { id: id.slice(0, 256) } : {}), ...(error ? { error: error.slice(0, 256) } : {}) }); }; const now = deps.now ?? Date.now; const sleep = deps.sleep ?? ((ms: number) => new Promise((resolve) => setTimeout(resolve, ms))); const reconcile = deps.reconcile ?? reconcileAsyncRun; const resultsDir = deps.resultsDir ?? DIRS.results; const asyncRoot = deps.asyncDirRoot ?? DIRS.async; /** "live" and "unknown" both get a stop attempt (fail closed); only a provably finished run is skipped. */ const check = (asyncDir: string, job?: AsyncJobState): { settled: boolean; state?: string } => { try { const reconciled = reconcile(asyncDir, job ? { resultsDir, startedRun: { runId: job.asyncId, pid: job.pid, sessionId: job.sessionId, completionOwnerId: job.completionOwnerId, mode: job.mode, agents: job.agents, chainStepCount: job.chainStepCount, parallelGroups: job.parallelGroups, startedAt: job.startedAt, sessionFile: job.sessionFile, } } : { resultsDir }); const status = reconciled.status; if (!status) return { settled: job === undefined }; return { settled: TERMINAL.has(status.state), state: status.state }; } catch { return { settled: false }; } }; const candidates = new Map(); for (const job of state.asyncJobs.values()) { if (job.status !== "queued" && job.status !== "running") continue; if (job.sessionId && state.currentSessionId && job.sessionId !== state.currentSessionId) { add("async_run", `Async run ${job.asyncId} belongs to another session; left untouched`, false, job.asyncId, "foreign_session"); continue; } candidates.set(job.asyncId, { asyncDir: job.asyncDir, kind: "async_run", label: `Async ${job.mode ?? "run"}${job.agents?.length ? ` (${job.agents.join(", ")})` : ""} ${job.asyncId}` }); } for (const runId of deps.observedCompletionRunIds()) { if (candidates.has(runId)) continue; candidates.set(runId, { asyncDir: path.join(asyncRoot, runId), kind: "scheduled_run", label: `Scheduled run ${runId}` }); } const nested = [ ...[...state.asyncJobs.values()].flatMap((job) => liveNested(job.nestedChildren)), ...[...(state.retainedForegroundNestedChildren?.values() ?? [])].flatMap((retained) => liveNested(retained.children)), ]; for (const child of nested) { if (!candidates.has(child.id)) candidates.set(child.id, { asyncDir: child.asyncDir, kind: "nested_run", label: `Nested run ${child.id}` }); } const waiting: Array<{ id: string; asyncDir: string; label: string; kind: string }> = []; for (const [id, candidate] of candidates) { const job = state.asyncJobs.get(id); if (check(candidate.asyncDir, job).settled) continue; try { deps.stopRun({ dir: candidate.asyncDir }); waiting.push({ id, ...candidate }); } catch (error) { add(candidate.kind, `Failed to stop ${candidate.label}: ${message(error)}`, false, id, message(error)); } } const foreground = [...state.foregroundControls.values()].map((control) => ({ runId: control.runId, label: `Foreground ${control.mode} run${control.currentAgent ? ` (${control.currentAgent})` : ""} ${control.runId}` })); const poll = deps.pollIntervalMs ?? 50; let pendingForeground = foreground; let pendingRuns = waiting; while ((pendingRuns.length > 0 || pendingForeground.length > 0) && now() < request.deadline) { await sleep(poll); pendingRuns = pendingRuns.filter((run) => { const result = check(run.asyncDir); if (!result.settled) return true; add(run.kind, result.state === "stopped" || result.state === undefined ? `Stopped ${run.label}` : `${run.label} ended (${result.state}) while stopping`, true, run.id); return false; }); pendingForeground = pendingForeground.filter((run) => { if (state.foregroundControls.has(run.runId)) return true; add("foreground_run", `${run.label} ended with the aborted turn`, true, run.runId); return false; }); } for (const run of pendingRuns) add(run.kind, `Stop request delivered for ${run.label}, but the runner has not acknowledged it yet`, false, run.id, "timeout"); for (const run of pendingForeground) add("foreground_run", `${run.label} is still active after the turn abort`, false, run.runId, "timeout"); return items; } /** Disarms this session's session-only schedules (see `ScheduledRunManager.disarmSessionOnlySchedules`). */ export function stopSchedules(deps: SchedulerStopDeps): StopAllBackgroundItem[] { return deps.disarmSessionSchedules().map((entry) => ({ source: "scheduler", kind: "schedule", id: entry.id.slice(0, 256), detail: entry.detail.slice(0, 512), ok: entry.ok, ...(entry.error ? { error: entry.error.slice(0, 256) } : {}), })); } /** * Registers the `subagents` and `scheduler` release-readiness contributors and, when `stopAll` is given, their kill-switch * handlers on the same bus. All answer only for the session this extension hosts (`supervisorOwnerSessionId`, not * `currentSessionId`); other sessions' requests are ignored. */ export function registerReleaseReadiness(options: { events: Events; getHostSessionId(): string | null; subagents: SubagentReleaseDeps; scheduler: SchedulerReleaseDeps; stopAll?: { subagents: SubagentStopDeps; scheduler: SchedulerStopDeps }; }): { dispose(): void } { let disposed = false; const unsubscribes = [ options.events.on(RELEASE_READINESS_EVENT, (raw) => { if (!disposed) answerReleaseReadiness(raw, "subagents", options.getHostSessionId(), () => subagentReleaseBlockers(options.subagents)); }), options.events.on(RELEASE_READINESS_EVENT, (raw) => { if (!disposed) answerReleaseReadiness(raw, "scheduler", options.getHostSessionId(), () => schedulerReleaseBlockers(options.scheduler)); }), ]; const stopAll = options.stopAll; if (stopAll) { // Schedules first: their timers are cleared synchronously, before the run stops below can complete a scheduled run and re-arm anything. unsubscribes.push( options.events.on(STOP_ALL_BACKGROUND_EVENT, (raw) => { if (!disposed) answerStopAllBackground(raw, "scheduler", options.getHostSessionId(), () => stopSchedules(stopAll.scheduler)); }), options.events.on(STOP_ALL_BACKGROUND_EVENT, (raw) => { if (!disposed) answerStopAllBackground(raw, "subagents", options.getHostSessionId(), (request) => stopSubagentWork(stopAll.subagents, request)); }), ); } return { dispose() { disposed = true; for (const unsubscribe of unsubscribes) if (typeof unsubscribe === "function") unsubscribe(); } }; }