// Running subagent registry: the parent-side collection of adopted // RunningSubagent records. It owns the RunningSubagent record type, its // observation and cancellation verbs, and cancellation target resolution; it // yields structured resolutions, not presentation. See CONTEXT.md // "Running subagent registry". import { createLifecycle, markCompleted, markCompletionDetected, markFailed, markInterruptRequested, markProcessRunning, observeActivity, observeNativeAgentObservation, type SubagentLifecycle, } from "./lifecycle.ts"; import { readSubagentActivityFile, type ActivityReadResult, type SubagentActivityState, } from "./activity.ts"; import { publishCancellationRequest } from "./cancellation-sidecar.ts"; import { closeHerdrSurface, sendHerdrAgentEscape } from "./herdr.ts"; import type { OperationArtifacts } from "./operation-artifacts.ts"; import type { CompletionOperation } from "./completion-sidecar.ts"; import { createOperationReference } from "./operation-identity.ts"; import { completionResultForOutcome, type CompletionControlOutcome, } from "./completion.ts"; import type { CompletionHandoffResult } from "./completion-handoff.ts"; import type { NativeAgentObservation } from "./native-supervision.ts"; import type { ResolvedRuntimePlan } from "./runtime-routing.ts"; /** * State for a launched (but not yet completed) subagent. * * Ownership: identity fields (RunningSubagentIdentity) are set at adoption * and never mutated; runtime fields are written only by this module's verbs * — observation and interrupt verbs plus the Completion watch arc's outcome * verbs — and the coordination flags carry the parent shutdown policy's * writes through noteExternalCancellation. */ export interface RunningSubagent extends RunningSubagentIdentity { // ── Watch-managed runtime: only registry verbs write these ── activity?: SubagentActivityState; activityRead?: { ok: boolean; reason?: "missing" | "invalid" | "wrong-id" | "transient"; error?: string; }; abortController?: AbortController; lifecycle: SubagentLifecycle; // ── Coordination flags: shutdown policy and cancellation verbs write, the watch reads ── /** Set when lifecycle cancellation already performed native cleanup. */ cancellationRequested?: boolean; /** Set when the user's cancel tool requested this cancellation, not a derived cascade. */ userCancelRequested?: boolean; } export function ensureLifecycle(running: RunningSubagent): SubagentLifecycle { if (running.lifecycle) return running.lifecycle; // Runtime entries surviving a /reload from a pre-lifecycle build carry no // lifecycle; present them as a running process until Completion evidence // arrives. Current launches always set lifecycle at adoption time. running.lifecycle = markProcessRunning(createLifecycle(running.startTime), running.startTime); return running.lifecycle; } export function observeRunningSubagent(running: RunningSubagent, observedAt = Date.now()) { ensureLifecycle(running); const activityFile = running.activityFile; const read: ActivityReadResult = activityFile ? readSubagentActivityFile(activityFile, running.id) : { ok: false, reason: "missing" }; running.activityRead = read.ok ? { ok: true } : { ok: false, reason: read.reason, error: read.error }; if (read.ok) running.activity = read.activity; running.lifecycle = observeActivity(ensureLifecycle(running), read, observedAt); } export function interruptAndCloseSubagent( running: Pick, interrupt: (surface: string) => void = sendHerdrAgentEscape, close: (surface: string) => void = closeHerdrSurface, ): void { try { interrupt(running.surface); } catch { // The native surface may already be gone; cleanup still gets a close attempt. } try { close(running.surface); } catch { // Cleanup is best effort after cancellation. } } /** * Establish lifecycle cancellation for one adopted descendant. The native * surface is cleaned up before aborting its watcher, so cancellation remains * control flow and the watcher can release the record only after its arc * observes the cancellation. */ export function requestSubagentCancellation( running: RunningSubagent, interrupt: (surface: string) => void = sendHerdrAgentEscape, ): void { if (running.cancellationRequested) return; running.cancellationRequested = true; try { publishCancellationRequest(operationForRunningSubagent(running)); } catch { // The native signal still establishes local control cancellation; the // watcher remains tracked if the in-runtime request cannot be written. } try { interrupt(running.surface); } catch { // A missing surface is handled by Completion supervision; the lifecycle // remains cancelled and its registry entry stays tracked until drain. } } /** Outcome of a user-facing cancellation request. */ export type UserCancelResult = | { kind: "accepted" } | { kind: "already-cancelling" } | { kind: "finalizing" }; /** * User-facing cancellation: the cancel tool's request to terminate one * Running Subagent. Marks the record as user-requested so the Completion * watch can report the terminal cancellation to the parent conversation, * then establishes lifecycle cancellation. The child's surface is reclaimed * by the watch arc after the child's cancellation acknowledgement, and a * Subagent already past its Completion evidence is never cancelled. */ export function requestUserCancellation( running: RunningSubagent, interrupt: (surface: string) => void = sendHerdrAgentEscape, ): UserCancelResult { const processKind = ensureLifecycle(running).process.kind; if (processKind === "finalizing" || processKind === "completed" || processKind === "failed") { return { kind: "finalizing" }; } if (running.cancellationRequested) { return { kind: "already-cancelling" }; } running.userCancelRequested = true; requestSubagentCancellation(running, interrupt); return { kind: "accepted" }; } // ── Watch-arc verbs: the Completion watch drives runtime state through these ── /** Project one Native supervision tick: native observation, activity read, and lifecycle fold. */ export function observeSupervision( running: RunningSubagent, observation: NativeAgentObservation, observedAt: number, ): void { ensureLifecycle(running); running.lifecycle = observation.kind === "present" ? observeNativeAgentObservation( running.lifecycle, { ...observation, observedAt }, observedAt, ) : observeNativeAgentObservation(running.lifecycle, observation, observedAt); // Pi activity is optional enrichment; Native supervision stays authoritative. observeRunningSubagent(running, observedAt); } /** Note the outcome Completion supervision reported, before handoff enrichment. */ export function noteOutcomeObserved( running: RunningSubagent, outcome: CompletionControlOutcome, detectedAt: number, ): void { if (outcome.kind === "cancelled") { running.lifecycle = markInterruptRequested(ensureLifecycle(running), detectedAt); return; } const completion = completionResultForOutcome(outcome); running.lifecycle = markCompletionDetected(ensureLifecycle(running), completion, detectedAt); } /** Settle a handoff result: mark the terminal lifecycle and fold the observed runtime plan. */ export function settleOutcome(running: RunningSubagent, result: CompletionHandoffResult): void { ensureLifecycle(running); if (result.runtimePlan) running.runtimePlan = result.runtimePlan; running.lifecycle = result.exitCode === 0 ? markCompleted(running.lifecycle, Date.now()) : markFailed( running.lifecycle, result.errorMessage ?? result.summary, Date.now(), result.exitCode, result.failureCategory, ); } /** Attach the watcher AbortController; the registry owns the watcher slot. */ export function attachWatcher(running: RunningSubagent, watcher: AbortController): void { running.abortController = watcher; } // ── Cancellation-cleanup policy ── /** Who reclaims the native surface when this record's Completion watch ends aborted. */ export type CancelledWatchCleanup = "structured-cascade" | "close-surface"; /** * Single owner of the cancellation-cleanup conjunction: exactly one party * closes a given native surface, exactly once. A not-yet-cancelled foreground * record still owns its own surface; everything else — already-requested * cancellation, background records — belongs to the structured Descendant * cancellation cascade. */ export function resolveCancelledWatchCleanup( running: Pick, ): CancelledWatchCleanup { if (!running.cancellationRequested && (running.cleanupOnCancellation ?? true)) { return "close-surface"; } return "structured-cascade"; } // ── External-writer verbs: shutdown policy and launch admission ── /** * Record that cancellation already took effect elsewhere (cutover-won * adoption or the parent shutdown policy); no native action is taken here. */ export function noteExternalCancellation( running: Pick, ): void { running.cancellationRequested = true; } // ── Record construction ── /** Identity fields resolved before adoption; set once and never mutated. */ export interface RunningSubagentIdentity { /** Completion operation ID; the sole registry identity for this attempt. */ id: string; name: string; task: string; agent?: string; surface: string; startTime: number; /** Private operation namespace for Completion and cancellation artifacts. */ artifacts: OperationArtifacts; activityFile?: string; /** Foreground cancellation owns cleanup for the active child. */ cleanupOnCancellation?: boolean; /** Parent-resolved model/thinking selection and provenance. */ runtimePlan: ResolvedRuntimePlan | undefined; } /** * Construct an adopted record with its lifecycle already running from the * operation start; the registry owns which lifecycle step adoption means. */ export function createRunningSubagent(identity: RunningSubagentIdentity): RunningSubagent { return { ...identity, lifecycle: markProcessRunning(createLifecycle(identity.startTime), identity.startTime), }; } /** * Build the current Completion operation reference for a RunningSubagent. * The operation ID and private artifact namespace are the complete control-plane * identity; callers never reconstruct an external path pair. */ export function operationForRunningSubagent( running: Pick, ): CompletionOperation { return createOperationReference(running.id, running.artifacts); } /** Structured interrupt-target resolution; message wording stays with callers. */ export type RunningRegistryTargetResolution = | { kind: "resolved"; running: RunningSubagent } | { kind: "ambiguous"; matches: readonly RunningSubagent[] } | { kind: "not-found" } | { kind: "missing-target" }; export interface RunningRegistry { /** * Adopt one RunningSubagent. Registry identity fields are set at adoption * and never mutated; re-registering the same record replaces it. */ register(running: RunningSubagent): void; /** Release a record by its Completion operation identity. */ remove(running: RunningSubagent): void; /** Resolve an interrupt target by exact operation id or display name. */ resolveTarget(params: { id?: string; name?: string }): RunningRegistryTargetResolution; /** Current records as read-only views; runtime state changes cross the verbs. */ values(): Iterable>; readonly size: number; /** Clear every record; used by the parent shutdown policy. */ clear(): void; } /** * Create the parent runtime's registry. */ export function createRunningRegistry(): RunningRegistry { const entries = new Map(); function register(running: RunningSubagent): void { // Completion operation IDs are allocated independently and are the sole // logical identity of a RunningSubagent. Artifact paths are only // protocol locations and must never create a second registry entry. entries.set(running.id, running); } function remove(running: RunningSubagent): void { // The operation identity is stable across the watch arc; callers may // release it with a different record object after a reload. entries.delete(running.id); } function resolveTarget(params: { id?: string; name?: string }): RunningRegistryTargetResolution { const requestedId = params.id?.trim(); if (requestedId) { const matches = Array.from(entries.values()).filter((running) => running.id === requestedId); if (matches.length === 1) return { kind: "resolved", running: matches[0] }; if (matches.length > 1) return { kind: "ambiguous", matches }; return { kind: "not-found" }; } const requestedName = params.name?.trim(); if (!requestedName) return { kind: "missing-target" }; const matches = Array.from(entries.values()).filter((running) => running.name === requestedName); if (matches.length === 1) return { kind: "resolved", running: matches[0] }; if (matches.length > 1) return { kind: "ambiguous", matches }; return { kind: "not-found" }; } return { register, remove, resolveTarget, values: () => entries.values(), get size() { return entries.size; }, clear: () => entries.clear(), }; }