// Direct-parent coordination for Nested Subagents. This module is shared by // the parent-side launcher/watch and the child-side Completion publisher when // both extensions run in the same Pi process. import { isNestedSubagentProcess } from "./operation-env.ts"; import type { OperationArtifacts } from "./operation-artifacts.ts"; import { getProcessParentRuntimeId, type ParentRuntimeId, } from "./operation-identity.ts"; import { retainAcrossReload } from "./reload-boundary.ts"; export interface NestedDescendant { operationId: string; /** Private operation namespace used for descendant control artifacts. */ artifacts: OperationArtifacts; } /** Parent-side control used to cancel one adopted descendant. */ export type DescendantCancellationAction = () => void; /** A result accepted by a Direct parent before its next continuation starts. */ export interface DescendantResultHandoff { descendant: NestedDescendant; /** The independent child outcome; never becomes the parent's final result. */ outcome: T; } export type CompletionPublicationRequest = | "published" | "deferred" | "already-requested" | "cancelled"; export type DescendantHandoffRequest = "accepted" | "already-accepted" | "suppressed"; export type CancellationRequest = "cancelled" | "already-cancelled" | "already-completed"; export type NestedLifecycleState = "active" | "draining" | "cancelled" | "completed"; /** Stable lifecycle data consumed by reload-safe status presentation. */ export interface NestedLifecycleSnapshot { state: NestedLifecycleState; selfSettled: boolean; cancellationRequested: boolean; activeLaunchReservations: number; activeDescendants: number; pendingHandoffs: number; } export interface NestedLifecycleCoordinator { /** Process-local owner identity; stable across the Supported reload boundary. */ readonly parentRuntimeId: ParentRuntimeId; /** Count an accepted Native startup before Child adoption. */ reserveDescendant(descendant: NestedDescendant, cancel?: DescendantCancellationAction): void; /** Promote a launch reservation into an adopted direct descendant. */ adoptDescendant(descendant: NestedDescendant, cancel?: DescendantCancellationAction): boolean; /** Release a reservation after startup fails before adoption. */ releaseDescendant(descendant: NestedDescendant): void; /** Atomically accept a result behind the next parent continuation. */ acceptDescendantHandoff( descendant: NestedDescendant, result: DescendantResultHandoff, ): DescendantHandoffRequest; /** Admit one terminal Descendant result; flat parents always admit. */ admitDescendantResult( descendant: NestedDescendant, result: DescendantResultHandoff, ): DescendantHandoffRequest; /** Reject an unaccepted delivery while keeping lifecycle bookkeeping coherent. */ rejectDescendantHandoff(descendant: NestedDescendant): void; /** Mark the descendant's delivery/reclamation arc as complete. */ completeDescendant(descendant: NestedDescendant): void; /** * Settle one handoff after its single delivery attempt: a consumed * delivery clears the pending handoff against the self-settled boundary; * a lost delivery rejects it so deferred Completion proceeds. A duplicate * admission leaves the original pending handoff standing. */ settleDescendantDelivery( descendant: NestedDescendant, request: DescendantHandoffRequest, consumed: boolean, ): void; /** Mark a continuation turn as started and consume its coalesced batch. */ continuationStarted(): readonly DescendantResultHandoff[]; /** Request this process's terminal Completion publication. */ requestCompletion(publish: () => void): CompletionPublicationRequest; /** * Establish the Cancellation cutover and notify each active descendant or * launch reservation exactly once. Cancellation is control flow and never * becomes a child Completion outcome. */ requestCancellation(cancelDescendant?: (descendant: NestedDescendant) => void): CancellationRequest; /** Whether the Cancellation cutover has already won this lifecycle. */ isCancellationRequested(): boolean; /** Current lifecycle control state, including a parent waiting for drain. */ state(): NestedLifecycleState; /** Whether the current agent loop reached its provisional settled boundary. */ isSelfSettled(): boolean; /** Whether a reservation, active child, or pending handoff still owns a descendant. */ ownsDescendant(descendant: NestedDescendant): boolean; /** Read the reload-safe ownership and drain state for status projection. */ snapshot(): NestedLifecycleSnapshot; /** Resolve after all active descendants and launch reservations terminate. */ waitForDescendantDrain(): Promise; /** Number of accepted launch reservations awaiting adoption. */ activeLaunchReservationCount(): number; /** Number of adopted direct descendants that have not completed delivery. */ activeDescendantCount(): number; /** Number of accepted result handoffs awaiting a continuation. */ pendingHandoffCount(): number; /** Whether the parent lifecycle is still waiting on descendants or handoffs. */ isDraining(): boolean; } function descendantKey(descendant: NestedDescendant): string { // Completion operation IDs are unique per Fresh attempt. The artifact // namespace locates protocol data but is not part of parent-local ownership. return descendant.operationId; } class Coordinator implements NestedLifecycleCoordinator { readonly parentRuntimeId: ParentRuntimeId; private readonly reservations = new Set(); constructor(parentRuntimeId: ParentRuntimeId) { this.parentRuntimeId = parentRuntimeId; } private readonly active = new Set(); private readonly handoffs = new Map(); private readonly settled = new Set(); private readonly cancellationNotified = new Set(); private readonly cancellationActions = new Map(); private pendingPublication?: { publish: () => void }; private publicationCompleted = false; private selfSettled = false; private replacePendingPublication = false; private cancellationRequested = false; private publicationInProgress = false; private cancellationHandler?: (descendant: NestedDescendant) => void; private drainWaiters: Array<() => void> = []; reserveDescendant(descendant: NestedDescendant, cancel?: DescendantCancellationAction): void { if (this.publicationCompleted || this.publicationInProgress) { throw new Error("Cannot launch a descendant after parent Completion publication."); } if (this.cancellationRequested) { throw new Error("Cannot launch a descendant after Cancellation cutover."); } const key = descendantKey(descendant); if (this.settled.has(key) || this.active.has(key) || this.reservations.has(key)) { // A second launch with the same operation identity is invalid. Reject it // instead of letting a later release clear the first launch's reservation. throw new Error("Completion operation identity is already owned by this parent."); } this.descendants.set(key, descendant); if (cancel) this.cancellationActions.set(key, cancel); this.reservations.add(key); } adoptDescendant(descendant: NestedDescendant, cancel?: DescendantCancellationAction): boolean { if (this.publicationCompleted || this.publicationInProgress) return false; const key = descendantKey(descendant); if (this.settled.has(key) || !this.reservations.has(key)) return false; this.descendants.set(key, descendant); if (cancel) this.cancellationActions.set(key, cancel); this.reservations.delete(key); if (this.cancellationRequested) { this.settled.add(key); this.notifyCancellation(key, descendant); this.resolveDrainIfReady(); return false; } this.active.add(key); return true; } releaseDescendant(descendant: NestedDescendant): void { const key = descendantKey(descendant); // Completion operation identity is the only lifecycle lookup key. The // optional path/namespace may differ across reconstructed references. const preserveAcceptedHandoff = this.settled.has(key) && this.handoffs.has(key); this.reservations.delete(key); this.active.delete(key); if (!preserveAcceptedHandoff) this.handoffs.delete(key); this.settled.add(key); this.cancellationActions.delete(key); this.descendants.delete(key); this.tryPublish(); this.resolveDrainIfReady(); } acceptDescendantHandoff( descendant: NestedDescendant, result: DescendantResultHandoff, ): DescendantHandoffRequest { const key = descendantKey(descendant); if (descendantKey(result.descendant) !== key) return "suppressed"; if (this.publicationCompleted || this.publicationInProgress || this.cancellationRequested) return "suppressed"; if (!this.active.has(key)) return "suppressed"; if (this.handoffs.has(key)) return "already-accepted"; this.handoffs.set(key, result); return "accepted"; } admitDescendantResult( descendant: NestedDescendant, result: DescendantResultHandoff, ): DescendantHandoffRequest { // A flat parent owns no descendant handoff state; its results are // always admitted because delivery is unconditional. if (!hasNestedLifecycleContext()) return "accepted"; return this.acceptDescendantHandoff(descendant, result); } settleDescendantDelivery( descendant: NestedDescendant, request: DescendantHandoffRequest, consumed: boolean, ): void { if (request === "already-accepted") return; if (request !== "accepted" || !consumed) { this.rejectDescendantHandoff(descendant); return; } // A consumed delivery no longer holds the drain open. Pi owns the // continuation after acceptance: a foreground result arrives through the // parent's own running loop, while a background steer is queued into the // parent conversation, and Pi drains queued messages before // `agent_settled`. Waiting for a separate `agent_start` would strand an // in-run delivery with no remaining transition to clear it. this.deliveredHandoffProcessed(descendant); } rejectDescendantHandoff(descendant: NestedDescendant): void { const key = descendantKey(descendant); // Once the descendant has completed, its accepted handoff is the only // remaining ownership that keeps the parent alive. A late delivery // from a pre-reload Completion delivery must not erase that handoff and // lose the continuation. Rejection is only authoritative while the descendant // is still active (or while a launch reservation is being handed off). if (this.settled.has(key) && this.handoffs.has(key)) return; this.handoffs.delete(key); this.tryPublish(); } completeDescendant(descendant: NestedDescendant): void { const key = descendantKey(descendant); this.reservations.delete(key); this.active.delete(key); this.settled.add(key); this.cancellationActions.delete(key); this.descendants.delete(key); this.tryPublish(); this.resolveDrainIfReady(); } private deliveredHandoffProcessed(descendant: NestedDescendant): void { const removed = this.handoffs.delete(descendantKey(descendant)); if (removed && this.pendingPublication) this.supersedeSelfSettledBoundary(); } continuationStarted(): readonly DescendantResultHandoff[] { // Multiple child results accepted before Pi starts a turn form one // Descendant continuation batch. A new continuation supersedes the // previous self-settled boundary, so its own settled event is required // before the parent's deferred Completion can be published. const hadPendingHandoff = this.handoffs.size > 0; const batch = Array.from(this.handoffs.values()); this.handoffs.clear(); if (hadPendingHandoff && this.pendingPublication) this.supersedeSelfSettledBoundary(); return batch; } private supersedeSelfSettledBoundary(): void { this.selfSettled = false; this.replacePendingPublication = true; } requestCompletion(publish: () => void): CompletionPublicationRequest { if (this.cancellationRequested) return "cancelled"; if (this.publicationCompleted || this.publicationInProgress) return "already-requested"; if (this.pendingPublication) { if (this.replacePendingPublication) { // A continuation has superseded the provisional result from the // previous self-settled loop. this.pendingPublication = { publish }; this.replacePendingPublication = false; } this.selfSettled = true; return this.tryPublish() ? "published" : "deferred"; } this.selfSettled = true; this.pendingPublication = { publish }; if (this.tryPublish()) return "published"; return "deferred"; } requestCancellation(cancelDescendant?: (descendant: NestedDescendant) => void): CancellationRequest { if (this.publicationCompleted || this.publicationInProgress) return "already-completed"; if (this.cancellationRequested) { // A child-side module may install the parent-side cancellation hook // after the cutover (for example during /reload). Replay only the // descendants that have not already been notified. if (!this.cancellationHandler && cancelDescendant) { this.cancellationHandler = cancelDescendant; for (const key of [...this.reservations, ...this.active]) { this.notifyCancellation(key, this.descendantForOperation(key)); } } return "already-cancelled"; } this.cancellationRequested = true; this.cancellationHandler = cancelDescendant; this.pendingPublication = undefined; this.selfSettled = false; this.replacePendingPublication = false; this.handoffs.clear(); for (const key of [...this.reservations, ...this.active]) { this.notifyCancellation(key, this.descendantForOperation(key)); } this.resolveDrainIfReady(); return "cancelled"; } isCancellationRequested(): boolean { return this.cancellationRequested; } isSelfSettled(): boolean { return this.selfSettled; } ownsDescendant(descendant: NestedDescendant): boolean { const key = descendantKey(descendant); return this.reservations.has(key) || this.active.has(key) || this.handoffs.has(key); } snapshot(): NestedLifecycleSnapshot { return { state: this.state(), selfSettled: this.selfSettled, cancellationRequested: this.cancellationRequested, activeLaunchReservations: this.reservations.size, activeDescendants: this.active.size, pendingHandoffs: this.handoffs.size, }; } state(): NestedLifecycleState { if (this.publicationCompleted) return "completed"; if (this.cancellationRequested) return this.isDraining() ? "draining" : "cancelled"; if (this.isDraining()) return "draining"; return "active"; } waitForDescendantDrain(): Promise { if (!this.isDraining()) return Promise.resolve(); return new Promise((resolve) => this.drainWaiters.push(resolve)); } activeLaunchReservationCount(): number { return this.reservations.size; } activeDescendantCount(): number { return this.active.size; } pendingHandoffCount(): number { return this.handoffs.size; } isDraining(): boolean { return this.reservations.size > 0 || this.active.size > 0 || this.handoffs.size > 0; } private descendantForOperation(key: string): NestedDescendant { // The key is only used for cancellation notification. Keep the original // operation reference beside the active sets so artifact paths are // preserved for the callback without reconstructing them from the key. return this.descendants.get(key)!; } private readonly descendants = new Map(); private notifyCancellation(key: string, descendant: NestedDescendant): void { if (this.cancellationNotified.has(key)) return; this.cancellationNotified.add(key); this.descendants.set(key, descendant); try { const action = this.cancellationActions.get(key); if (action) action(); else this.cancellationHandler?.(descendant); } catch { // Cancellation is best effort; termination/release remains authoritative. } } private resolveDrainIfReady(): void { if (this.isDraining()) return; const waiters = this.drainWaiters; this.drainWaiters = []; for (const resolve of waiters) resolve(); } private tryPublish(): boolean { if (this.cancellationRequested || this.publicationInProgress) return false; if (!this.pendingPublication || !this.selfSettled || this.isDraining()) return false; const pending = this.pendingPublication; // Linearize valid final evidence before invoking the publisher. A // re-entrant cancellation request must not turn one operation into both // a normal Completion and a cancelled lifecycle. this.publicationInProgress = true; try { pending.publish(); } catch (error) { this.publicationInProgress = false; // The pending request remains protected for the next lifecycle signal. throw error; } this.pendingPublication = undefined; this.publicationCompleted = true; this.publicationInProgress = false; this.resolveDrainIfReady(); return true; } } export function createNestedLifecycleCoordinator( parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): NestedLifecycleCoordinator { return new Coordinator(parentRuntimeId); } const GLOBAL_KEY = Symbol.for("pi-subagents/nested-lifecycle-coordinators"); type CoordinatorStore = Map; function globalStore(): CoordinatorStore { return retainAcrossReload(GLOBAL_KEY, () => new Map()); } /** * Return the process-local coordinator for a Parent Runtime ID. The identity * is independent of Pi persistence, so an Ephemeral Child parent still gets * an isolated lifecycle coordinator and the same identity survives /reload. */ export function getNestedLifecycleCoordinator( parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): NestedLifecycleCoordinator { const store = globalStore(); const existing = store.get(parentRuntimeId); if (existing) return existing; const coordinator = createNestedLifecycleCoordinator(parentRuntimeId); store.set(parentRuntimeId, coordinator); return coordinator; } /** * Drop a completed parent coordinator before a new parent runtime starts. * Active coordinators must remain installed across /reload and are never * reset by this helper. */ export function resetNestedLifecycleCoordinator( parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): void { globalStore().delete(parentRuntimeId); } /** * Whether this process runs as a Direct parent with a nested-lifecycle * context. The Completion operation identity is the only child-runtime marker. */ export function hasNestedLifecycleContext(): boolean { return isNestedSubagentProcess(); }