diff --git a/docs/design/COMPLETION_INBOX.md b/docs/design/COMPLETION_INBOX.md new file mode 100644 index 00000000..f8392d19 --- /dev/null +++ b/docs/design/COMPLETION_INBOX.md @@ -0,0 +1,77 @@ +# Background Completion Inbox + +- Status: validated +- Created: 2026-09-04 +- Verified: 2026-09-04 +- Source boundary: OpenPI implementation and tests in the pull request that closes issue #160 +- Related issue: https://github.com/openpi-dev/openpi/issues/160 +- Related pull request: https://github.com/openpi-dev/openpi/pull/382 +- Supersedes: the three producer-local consumption maps, not their execution state machines + +## Boundary + +Direct Subagent, Background Terminal, and Workflow keep independent execution +lifecycles and canonical terminal records. The shared completion inbox is a +small in-process delivery mechanism. It does not own execution status, result +bytes, artifacts, cancellation, Goals, Tasks, model judgment, or UI state. + +The four relevant projections stay distinct: + +1. `SubagentSnapshot`, `TerminalSnapshot`, and `WorkflowDetails` are canonical + execution facts. +2. `CompletionEnvelope` carries only delivery identity, owner, producer, + terminal reference, wake policy, and an in-process payload. +3. Pi's existing `followUp` and `nextTurn` messages create model-visible + context. +4. Existing TUI and Web surfaces render producer state; they do not infer + completion from the inbox. + +This uses Pi's SessionManager identity and message-delivery APIs rather than +adding a second Session mailbox, scheduler, or agent runtime. + +## Envelope and ownership + +Every envelope has a stable `deliveryId`, `{sessionId, epoch}` owner, +`producer`, `producerId`, `terminalRef`, and producer-selected `wake` policy. +All producers observing the same Pi SessionManager object share one +process-local epoch. Replacing that object, or changing its Session id, creates +a new epoch. A claim with a missing or mismatched owner becomes an inspectable +dead letter and is never redirected to another transcript. + +Workflow alone has durable producer state. Its delivery owner is persisted in +`WorkflowDetails`. On restart, restoring a pending terminal artifact is the +explicit revival boundary: the same Session id is rebound to the current +process-local epoch and persisted; a different Session remains pending in +canonical Workflow state and is dead-lettered by the inbox. Direct Subagents +and Background Terminals do not survive Pi's `session_shutdown`, so their +process-local inbox entries are cleared with their existing runtime lifecycle. + +## Consumption and receipts + +Pending and in-flight maps form one atomic consumption gate. Explicit +`wait`/`status` consumption and automatic delivery race on that gate; the first +claim wins. A transport failure restores the exact in-flight envelopes ahead +of newer work. Independent batches can be acknowledged or retried without +overwriting each other. + +The contract deliberately does not claim distributed exactly-once delivery: + +- Direct and Background transports consume after synchronous acceptance and + restore on synchronous rejection. +- Workflow keeps its stable per-run receipt and durable at-least-once recovery. + If transport succeeds but receipt persistence fails, the same delivery id + may replay so the parent can identify it. +- Partial Workflow receipts retry only unacknowledged siblings. + +Wake behavior remains producer-owned: Direct always follows up at its parent +boundary, Background preserves its idle-follow-up versus busy-next-turn +policy, and Workflow preserves its client-aware delivery adapter. + +## Validation + +Focused tests cover atomic explicit/automatic consumption, independent +in-flight batches, retry ordering, stable producer identities, partial +receipts, persistence failure, same-Session revival, Session id and epoch +switches, owner loss, and dead letters. Existing integration tests retain the +producer-specific wake, shutdown, process-restart, and canonical-state +behavior. Repository gates are `bun run check` and `bun run test` on Node 22. diff --git a/docs/design/README.md b/docs/design/README.md index 0283a682..ea37c7e1 100644 --- a/docs/design/README.md +++ b/docs/design/README.md @@ -13,5 +13,6 @@ These records predate [`Decision 0001`](../decisions/0001-documentation-and-evid - [`WORKFLOW_INVOCATION_GRAPH.md`](WORKFLOW_INVOCATION_GRAPH.md) — durable invocation facts, same-run handoff refs, reusable operators, and derived graph semantics - [`CHILD_TOOL_ACTIVITY.md`](CHILD_TOOL_ACTIVITY.md) — shared compact and Pi-native expanded evidence projection for Direct Subagent and Workflow child transcripts - [`OPENPI_WEB_ARCHITECTURE.md`](OPENPI_WEB_ARCHITECTURE.md) — draft architecture, protocol boundaries, delivery phases, and visual direction for the local Web workbench +- [`COMPLETION_INBOX.md`](COMPLETION_INBOX.md) — shared owner, epoch, consumption, retry, and receipt contract for background completions 开发与热更新流程见 [`docs/development/OPENPI_WEB_DEVELOPMENT.md`](../development/OPENPI_WEB_DEVELOPMENT.md)。 diff --git a/extensions/background-terminals/index.ts b/extensions/background-terminals/index.ts index 5d38a846..938186ea 100644 --- a/extensions/background-terminals/index.ts +++ b/extensions/background-terminals/index.ts @@ -32,6 +32,7 @@ import { OPENPI_TOOL_SURFACE, patchOwnedTools, } from "../shared/tool-surface.ts"; +import { completionOwnerFor } from "../shared/completion-inbox.ts"; import { projectBackgroundTerminalCapability, registerWebCapability, @@ -105,7 +106,12 @@ export default function (pi: ExtensionAPI) { let ui: ExtensionUIContext | undefined; let unsubStatus: (() => void) | undefined; let startReservations = 0; - const resultDelivery = createDeferredResultDelivery(); + const resultDelivery = createDeferredResultDelivery({ + owner: () => + sessionContext + ? completionOwnerFor(sessionContext.sessionManager) + : undefined, + }); const hideLifecycleTools = () => patchOwnedTools(pi, "background", { disable: OPENPI_TOOL_SURFACE.background.deferred, @@ -245,6 +251,7 @@ export default function (pi: ExtensionAPI) { const flushResults = (wake: boolean) => { const snaps = resultDelivery.drain(MAX_RUNNING); if (!deliverResults(snaps, wake)) resultDelivery.restore(snaps); + else resultDelivery.acknowledge(snaps); }; const idleResultBatcher = createIdleResultBatcher({ diff --git a/extensions/background-terminals/src/result-delivery.ts b/extensions/background-terminals/src/result-delivery.ts index 72a12b5c..b5d9599a 100644 --- a/extensions/background-terminals/src/result-delivery.ts +++ b/extensions/background-terminals/src/result-delivery.ts @@ -1,44 +1,64 @@ import type { ConsumableResultDeliveryQueue } from "../../shared/result-delivery.ts"; +import { + type CompletionOwner, + createCompletionInbox, +} from "../../shared/completion-inbox.ts"; /** - * Deferred one-shot delivery map (same semantics as subagents'): a settled - * terminal's result is held here until it is either drained into a follow-up - * message or consumed by a tool call (bg_kill / bg_status) that already - * returned the settlement itself. Keyed by id, so double delivery is - * structurally impossible — whoever drains first wins. + * Deferred one-shot delivery adapter (same semantics as subagents'): a + * settled terminal's result is held in the shared inbox until it is either + * drained into a follow-up message or consumed by a tool call (bg_kill / + * bg_status) that already returned the settlement itself. Stable ids make + * double delivery structurally impossible — whoever claims first wins. */ -export function createDeferredResultDelivery() { - const pending = new Map(); +export function createDeferredResultDelivery( + options: { readonly owner?: () => CompletionOwner | undefined } = {}, +) { + const inbox = createCompletionInbox(); + const owner = options.owner ?? (() => ({ sessionId: "test", epoch: 0 })); const queue = { defer(result: T) { - pending.set(result.id, result); - return pending.size; + const currentOwner = owner(); + inbox.defer( + { + deliveryId: `background:${result.id}`, + owner: currentOwner ?? { sessionId: "unowned", epoch: 0 }, + producer: "background", + producerId: result.id, + terminalRef: { kind: "terminal-snapshot", id: result.id }, + wake: "producer-policy", + payload: result, + }, + currentOwner, + ); + return inbox.size(); }, consume(ids: Iterable) { - for (const id of ids) pending.delete(id); + inbox.consume("background", ids); }, drain(maxResults = Number.POSITIVE_INFINITY) { - const results: T[] = []; - for (const [id, result] of pending) { - if (results.length >= maxResults) break; - results.push(result); - pending.delete(id); - } - return results; + return inbox + .claim(owner(), maxResults) + .map((envelope) => envelope.payload); }, restore(results: readonly T[]) { - const current = [...pending.values()]; - pending.clear(); - for (const result of results) pending.set(result.id, result); - for (const result of current) pending.set(result.id, result); + inbox.retryClaimed( + "background", + results.map((result) => result.id), + owner(), + ); + }, + acknowledge(results: readonly T[]) { + inbox.acknowledge(results.map((result) => `background:${result.id}`)); }, size() { - return pending.size; + return inbox.size(); }, clear() { - pending.clear(); + inbox.clear(); }, + inspectDeadLetters: inbox.inspectDeadLetters, }; return queue satisfies ConsumableResultDeliveryQueue; } diff --git a/extensions/shared/completion-inbox.ts b/extensions/shared/completion-inbox.ts new file mode 100644 index 00000000..28f9984f --- /dev/null +++ b/extensions/shared/completion-inbox.ts @@ -0,0 +1,188 @@ +export type CompletionProducer = "subagent" | "workflow" | "background"; +export type CompletionWakePolicy = + | "follow-up" + | "next-turn" + | "producer-policy"; + +export interface CompletionOwner { + readonly sessionId: string; + readonly epoch: number; +} + +export interface CompletionSessionIdentity { + getSessionId(): string; +} + +/** Transport metadata only; producer state remains the terminal authority. */ +export interface CompletionEnvelope { + readonly deliveryId: string; + readonly owner: CompletionOwner; + readonly producer: CompletionProducer; + readonly producerId: string; + readonly terminalRef: unknown; + readonly wake: CompletionWakePolicy; + readonly payload: T; +} + +export interface CompletionDeadLetter { + readonly deliveryId: string; + readonly producer: CompletionProducer; + readonly producerId: string; + readonly failure: "owner-unavailable" | "stale-owner"; +} + +function sameOwner(left: CompletionOwner, right: CompletionOwner) { + return left.sessionId === right.sessionId && left.epoch === right.epoch; +} + +let nextOwnerEpoch = 1; +const ownerBySessionIdentity = new WeakMap(); + +/** + * Bind one process-local generation to a Pi SessionManager identity. + * + * All OpenPI producers observing the same manager share an owner. A replaced + * manager, or a manager whose Session id changes, receives a new epoch so late + * callbacks cannot target the replacement transcript. + */ +export function completionOwnerFor( + identity: CompletionSessionIdentity, +): CompletionOwner { + const sessionId = identity.getSessionId(); + const existing = ownerBySessionIdentity.get(identity); + if (existing?.sessionId === sessionId) return existing; + const owner = { sessionId, epoch: nextOwnerEpoch++ }; + ownerBySessionIdentity.set(identity, owner); + return owner; +} + +/** + * One atomic consumption gate shared by all background producers. + * + * Claim removes before transport; a failed transport retries the exact + * envelopes. A successful transport leaves them consumed. No execution facts + * or result bytes are stored here. + */ +export function createCompletionInbox() { + const pending = new Map>(); + const inFlight = new Map>(); + const deadLetters: CompletionDeadLetter[] = []; + + const reject = ( + envelope: CompletionEnvelope, + failure: CompletionDeadLetter["failure"], + ) => { + deadLetters.push({ + deliveryId: envelope.deliveryId, + producer: envelope.producer, + producerId: envelope.producerId, + failure, + }); + return false; + }; + + const admit = ( + envelope: CompletionEnvelope, + owner: CompletionOwner | undefined, + ) => { + if (!owner) return reject(envelope, "owner-unavailable"); + if (!sameOwner(envelope.owner, owner)) { + return reject(envelope, "stale-owner"); + } + pending.set(envelope.deliveryId, envelope); + return true; + }; + + /** Restore a failed attempt ahead of completions that arrived meanwhile. */ + const retry = ( + envelopes: readonly CompletionEnvelope[], + owner: CompletionOwner | undefined, + ) => { + const current = [...pending.values()]; + pending.clear(); + for (const envelope of envelopes) { + inFlight.delete(envelope.deliveryId); + admit(envelope, owner); + } + for (const envelope of current) admit(envelope, owner); + }; + + return { + defer(envelope: CompletionEnvelope, owner: CompletionOwner | undefined) { + return admit(envelope, owner); + }, + + /** Explicit status/wait and automatic delivery atomically race here. */ + consume(producer: CompletionProducer, producerIds: Iterable) { + const ids = new Set(producerIds); + for (const [deliveryId, envelope] of pending) { + if (envelope.producer === producer && ids.has(envelope.producerId)) { + pending.delete(deliveryId); + inFlight.delete(deliveryId); + } + } + }, + + consumeDeliveryIds(deliveryIds: Iterable) { + for (const deliveryId of deliveryIds) { + pending.delete(deliveryId); + inFlight.delete(deliveryId); + } + }, + + claim( + owner: CompletionOwner | undefined, + maximum = Number.POSITIVE_INFINITY, + ) { + const claimed: CompletionEnvelope[] = []; + for (const [deliveryId, envelope] of pending) { + if (claimed.length >= maximum) break; + pending.delete(deliveryId); + if (!owner) { + reject(envelope, "owner-unavailable"); + continue; + } + if (!sameOwner(envelope.owner, owner)) { + reject(envelope, "stale-owner"); + continue; + } + inFlight.set(deliveryId, envelope); + claimed.push(envelope); + } + return claimed; + }, + + retry, + + retryClaimed( + producer: CompletionProducer, + producerIds: Iterable, + owner: CompletionOwner | undefined, + ) { + const ids = new Set(producerIds); + const envelopes = [...inFlight.values()].filter( + (envelope) => + envelope.producer === producer && ids.has(envelope.producerId), + ); + retry(envelopes, owner); + }, + + acknowledge(deliveryIds: Iterable) { + for (const deliveryId of deliveryIds) inFlight.delete(deliveryId); + }, + + size() { + return pending.size; + }, + + inspectDeadLetters() { + return [...deadLetters]; + }, + + clear() { + pending.clear(); + inFlight.clear(); + deadLetters.length = 0; + }, + }; +} diff --git a/extensions/subagents/index.ts b/extensions/subagents/index.ts index 7bcb0340..7bfbc7df 100644 --- a/extensions/subagents/index.ts +++ b/extensions/subagents/index.ts @@ -60,6 +60,7 @@ import { resolveStandaloneChildProjectTrust, } from "../shared/child-session.ts"; import { formatContextUtilization } from "../shared/context-utilization.ts"; +import { completionOwnerFor } from "../shared/completion-inbox.ts"; import { registerEditorLayer, removeEditorLayer, @@ -514,6 +515,10 @@ export default function ( ); const resultDelivery = createSubagentResultDelivery({ isIdle: () => sessionContext?.isIdle() === true, + owner: () => + sessionContext + ? completionOwnerFor(sessionContext.sessionManager) + : undefined, // Every unconsumed fire-and-forget result must reach the parent. The // delivery coordinator batches results that settled while it was busy. deliver: dispatchResults, diff --git a/extensions/subagents/src/result-delivery.ts b/extensions/subagents/src/result-delivery.ts index eb31f6e4..777fe404 100644 --- a/extensions/subagents/src/result-delivery.ts +++ b/extensions/subagents/src/result-delivery.ts @@ -1,10 +1,16 @@ import type { ConsumableResultDeliveryQueue } from "../../shared/result-delivery.ts"; +import { + type CompletionOwner, + createCompletionInbox, +} from "../../shared/completion-inbox.ts"; export interface SubagentResultDeliveryOptions { /** True only when the parent has no run or queued continuation in flight. */ readonly isIdle: () => boolean; /** Deliver one drained batch and wake the parent. */ readonly deliver: (results: readonly T[]) => void; + /** Current Pi Session transcript owner. */ + readonly owner?: () => CompletionOwner | undefined; } /** @@ -22,50 +28,63 @@ export interface SubagentResultDeliveryOptions { * The parent boundary wakes even if an earlier extension handler has already * started another turn: Pi queues the follow-up into that active run. * - * The Map is the one-shot gate: `subagent_wait` may consume a result before it - * is delivered, and whichever path drains first prevents duplicate delivery. + * The shared inbox is the one-shot gate: `subagent_wait` may consume a result + * before it is delivered, and whichever path claims first prevents duplicate + * delivery. */ export function createSubagentResultDelivery( options: SubagentResultDeliveryOptions, ) { - const pending = new Map(); + const inbox = createCompletionInbox(); + const owner = options.owner ?? (() => ({ sessionId: "test", epoch: 0 })); const flush = () => { - if (pending.size === 0) return; - const results = [...pending.values()]; - pending.clear(); + const envelopes = inbox.claim(owner()); + if (envelopes.length === 0) return; + const results = envelopes.map((envelope) => envelope.payload); try { options.deliver(results); + inbox.acknowledge(envelopes.map((envelope) => envelope.deliveryId)); } catch (error) { // A synchronous session teardown may reject append/send. Preserve the // original batch ahead of anything deferred re-entrantly while delivery // ran, so a later boundary can retry without loss or reordering. - const current = [...pending.values()]; - pending.clear(); - for (const result of results) pending.set(result.id, result); - for (const result of current) pending.set(result.id, result); + inbox.retry(envelopes, owner()); throw error; } }; const queue = { defer(result: T) { - pending.set(result.id, result); + const currentOwner = owner(); + inbox.defer( + { + deliveryId: `subagent:${result.id}`, + owner: currentOwner ?? { sessionId: "unowned", epoch: 0 }, + producer: "subagent", + producerId: result.id, + terminalRef: { kind: "subagent-snapshot", id: result.id }, + wake: "follow-up", + payload: result, + }, + currentOwner, + ); if (options.isIdle()) flush(); }, consume(ids: Iterable) { - for (const id of ids) pending.delete(id); + inbox.consume("subagent", ids); }, /** Flush at the authoritative parent boundary. */ parentSettled() { flush(); }, clear() { - pending.clear(); + inbox.clear(); }, size() { - return pending.size; + return inbox.size(); }, + inspectDeadLetters: inbox.inspectDeadLetters, }; return queue satisfies ConsumableResultDeliveryQueue; } diff --git a/extensions/workflows/dashboard.ts b/extensions/workflows/dashboard.ts index 2a9b34ef..dc214b19 100644 --- a/extensions/workflows/dashboard.ts +++ b/extensions/workflows/dashboard.ts @@ -220,6 +220,14 @@ function normalizeDelivery(value: unknown): WorkflowDetails["delivery"] { : 0; return { id: sanitizeLine(record.id, 256), + ...(typeof record.ownerSessionId === "string" && record.ownerSessionId + ? { ownerSessionId: sanitizeLine(record.ownerSessionId, 256) } + : {}), + ...(typeof record.ownerEpoch === "number" && + Number.isSafeInteger(record.ownerEpoch) && + record.ownerEpoch >= 0 + ? { ownerEpoch: record.ownerEpoch } + : {}), state, attempts, updatedAt, diff --git a/extensions/workflows/index.ts b/extensions/workflows/index.ts index 6e41341c..6abb1b78 100644 --- a/extensions/workflows/index.ts +++ b/extensions/workflows/index.ts @@ -55,6 +55,7 @@ import { import { fitNavigationSides } from "../shared/below-editor-navigation.ts"; import { waitBounded } from "../shared/child-session.ts"; import { contextPercent } from "../shared/context-utilization.ts"; +import { completionOwnerFor } from "../shared/completion-inbox.ts"; import { registerEditorLayer, removeEditorLayer, @@ -857,6 +858,8 @@ export default function workflows( }; const resultDelivery = createWorkflowResultDelivery({ isIdle: () => lastContext?.isIdle() ?? false, + owner: () => + lastContext ? completionOwnerFor(lastContext.sessionManager) : undefined, persist: (details) => { if (!details.delivery) throw new Error("Workflow delivery identity is missing"); @@ -1276,6 +1279,8 @@ export default function workflows( agents: [], delivery: { id: `workflow:${runId}:terminal`, + ownerSessionId: completionOwnerFor(ctx.sessionManager).sessionId, + ownerEpoch: completionOwnerFor(ctx.sessionManager).epoch, state: launchMode === "inline" ? "held-for-inline" : "none", attempts: 0, updatedAt: now, diff --git a/extensions/workflows/model.ts b/extensions/workflows/model.ts index dcb4f102..445296a9 100644 --- a/extensions/workflows/model.ts +++ b/extensions/workflows/model.ts @@ -67,6 +67,10 @@ export type WorkflowDeliveryState = export interface WorkflowDelivery { /** Stable per-run idempotency identity, never a transport-batch id. */ id: string; + /** Destination transcript identity; legacy records fall back to run.sessionId. */ + ownerSessionId?: string; + /** Process-local generation of the Pi SessionManager owner. */ + ownerEpoch?: number; state: WorkflowDeliveryState; attempts: number; updatedAt: number; diff --git a/extensions/workflows/result-delivery.ts b/extensions/workflows/result-delivery.ts index 4cac0de0..61fcf3ad 100644 --- a/extensions/workflows/result-delivery.ts +++ b/extensions/workflows/result-delivery.ts @@ -3,6 +3,11 @@ import type { DurableResultDeliveryQueue, DurableResultDeliveryReceipt, } from "../shared/result-delivery.ts"; +import { + type CompletionEnvelope, + type CompletionOwner, + createCompletionInbox, +} from "../shared/completion-inbox.ts"; export interface WorkflowCompletionEnvelope { deliveryId: string; @@ -12,6 +17,7 @@ export interface WorkflowCompletionEnvelope { export interface WorkflowResultDeliveryOptions { isIdle: () => boolean; + owner?: () => CompletionOwner | undefined; persist: (details: WorkflowDetails) => void; deliver: ( envelopes: readonly WorkflowCompletionEnvelope[], @@ -34,7 +40,9 @@ function errorText(error: unknown) { export function createWorkflowResultDelivery( options: WorkflowResultDeliveryOptions, ) { - const pending = new Map(); + const inbox = createCompletionInbox(); + const currentOwner = + options.owner ?? (() => ({ sessionId: "test", epoch: 0 })); let flushing: Promise | undefined; let flushRequested = false; let wakeRequested = false; @@ -57,10 +65,35 @@ export function createWorkflowResultDelivery( options.persist(details); }; - const enqueue = (envelope: WorkflowCompletionEnvelope) => { - pending.set(envelope.deliveryId, envelope); + const inboxEnvelope = ( + envelope: WorkflowCompletionEnvelope, + ): CompletionEnvelope => { + const owner = currentOwner(); + return { + deliveryId: envelope.deliveryId, + owner: { + sessionId: + envelope.details.delivery?.ownerSessionId ?? + envelope.details.sessionId ?? + owner?.sessionId ?? + "unowned", + epoch: envelope.details.delivery?.ownerEpoch ?? 0, + }, + producer: "workflow", + producerId: envelope.runId, + terminalRef: { + kind: "workflow-terminal", + runId: envelope.runId, + status: envelope.details.status, + }, + wake: "producer-policy", + payload: envelope, + }; }; + const enqueue = (envelope: WorkflowCompletionEnvelope) => + inbox.defer(inboxEnvelope(envelope), currentOwner()); + const retainPending = ( envelope: WorkflowCompletionEnvelope, patch: Partial>, @@ -69,7 +102,7 @@ export function createWorkflowResultDelivery( // Memory owns the retry before persistence is attempted. A broken disk // must not make this envelope, or any sibling after it, disappear from the // current process. - enqueue(envelope); + const admitted = enqueue(envelope); try { persistState(envelope.details, "pending", patch); } catch (error) { @@ -83,6 +116,7 @@ export function createWorkflowResultDelivery( }; } } + return admitted; }; const flush = async (wake: boolean) => { @@ -91,17 +125,16 @@ export function createWorkflowResultDelivery( wakeRequested ||= wake; return flushing; } - if (pending.size === 0) return; + if (inbox.size() === 0) return; flushing = (async () => { let passWake = wake; - while (pending.size > 0) { + while (inbox.size() > 0) { flushRequested = false; wakeRequested = false; - const envelopes = [...pending.values()]; - for (const envelope of envelopes) { - pending.delete(envelope.deliveryId); - } + const claimed = inbox.claim(currentOwner()); + const envelopes = claimed.map((envelope) => envelope.payload); + if (envelopes.length === 0) break; let receipts: readonly DurableResultDeliveryReceipt[] | undefined; try { @@ -109,6 +142,7 @@ export function createWorkflowResultDelivery( } catch (error) { const message = errorText(error); for (const envelope of envelopes) { + inbox.acknowledge([envelope.deliveryId]); retainPending( envelope, { @@ -127,6 +161,7 @@ export function createWorkflowResultDelivery( for (const envelope of envelopes) { const receipt = byId.get(envelope.deliveryId); if (receipt?.delivered) { + inbox.acknowledge([envelope.deliveryId]); try { persistState(envelope.details, "delivered", { attempts: (envelope.details.delivery?.attempts ?? 0) + 1, @@ -150,6 +185,7 @@ export function createWorkflowResultDelivery( } continue; } + inbox.acknowledge([envelope.deliveryId]); retainPending( envelope, { @@ -183,7 +219,7 @@ export function createWorkflowResultDelivery( /** Terminal won the wait/abort arbitration and will be returned inline. */ consumeInline(details: WorkflowDetails) { - if (details.delivery) pending.delete(details.delivery.id); + if (details.delivery) inbox.consumeDeliveryIds([details.delivery.id]); persistState(details, "consumed-inline", { deliveredAt: Date.now(), lastError: undefined, @@ -214,18 +250,37 @@ export function createWorkflowResultDelivery( restore(envelope: WorkflowCompletionEnvelope) { const state = envelope.details.delivery?.state; if (state !== "pending" && state !== "held-for-inline") return false; - // A process restart cannot still own the inline waiter. Deterministically - // reconstruct pending delivery from the terminal artifact. - if (state === "held-for-inline") { - retainPending( + const owner = currentOwner(); + const storedSessionId = + envelope.details.delivery?.ownerSessionId ?? + envelope.details.sessionId ?? + owner?.sessionId; + + // Restoring canonical producer state is the explicit owner-revival + // boundary. Rebind only the same transcript to this process-local + // SessionManager generation; a different Session still dead-letters. + if (owner && storedSessionId === owner.sessionId) { + envelope.details.delivery = { + ...envelope.details.delivery!, + ownerSessionId: owner.sessionId, + ownerEpoch: owner.epoch, + }; + return retainPending( envelope, - { lastError: "Inline waiter was not active after session restart" }, + { + ownerSessionId: owner.sessionId, + ownerEpoch: owner.epoch, + lastError: + state === "held-for-inline" + ? "Inline waiter was not active after session restart" + : undefined, + }, "Restored delivery state persistence failed", ); - } else { - enqueue(envelope); } - return true; + // Keep the canonical terminal artifact pending, but record that this + // process is not its transcript owner instead of redirecting it. + return enqueue(envelope); }, retryPending() { @@ -242,12 +297,13 @@ export function createWorkflowResultDelivery( }, size() { - return pending.size; + return inbox.size(); }, clear() { - pending.clear(); + inbox.clear(); }, + inspectDeadLetters: inbox.inspectDeadLetters, }; return queue satisfies DurableResultDeliveryQueue; } diff --git a/extensions/workflows/retention.ts b/extensions/workflows/retention.ts index 9d5e7a8d..66fd383a 100644 --- a/extensions/workflows/retention.ts +++ b/extensions/workflows/retention.ts @@ -237,6 +237,12 @@ function makeProjection( ? { delivery: { id: details.delivery.id, + ...(details.delivery.ownerSessionId + ? { ownerSessionId: details.delivery.ownerSessionId } + : {}), + ...(details.delivery.ownerEpoch !== undefined + ? { ownerEpoch: details.delivery.ownerEpoch } + : {}), state: details.delivery.state, attempts: details.delivery.attempts, updatedAt: details.delivery.updatedAt, diff --git a/tests/extensions/background-terminals/result-delivery.test.ts b/tests/extensions/background-terminals/result-delivery.test.ts index 1c84d17c..c6f422b3 100644 --- a/tests/extensions/background-terminals/result-delivery.test.ts +++ b/tests/extensions/background-terminals/result-delivery.test.ts @@ -90,6 +90,18 @@ test("a drained result can be retained for retry after delivery fails", () => { assert.deepEqual(delivery.drain(), [result]); }); +test("a deferred terminal completion cannot cross a Session switch", () => { + let sessionId = "session-1"; + const delivery = createDeferredResultDelivery<{ id: string }>({ + owner: () => ({ sessionId, epoch: 0 }), + }); + delivery.defer({ id: "bt-1" }); + + sessionId = "session-2"; + assert.deepEqual(delivery.drain(), []); + assert.equal(delivery.inspectDeadLetters()[0]?.failure, "stale-owner"); +}); + test("idle result batching uses one fixed bounded window", () => { let callback: (() => void) | undefined; let scheduled = 0; diff --git a/tests/extensions/shared/completion-inbox.test.ts b/tests/extensions/shared/completion-inbox.test.ts new file mode 100644 index 00000000..b2760798 --- /dev/null +++ b/tests/extensions/shared/completion-inbox.test.ts @@ -0,0 +1,148 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { + completionOwnerFor, + type CompletionEnvelope, + type CompletionOwner, + createCompletionInbox, +} from "../../../extensions/shared/completion-inbox.ts"; + +const owner: CompletionOwner = { sessionId: "session-1", epoch: 1 }; + +function envelope( + producerId: string, + patch: Partial> = {}, +): CompletionEnvelope<{ value: string }> { + return { + deliveryId: `subagent:${producerId}:terminal`, + owner, + producer: "subagent", + producerId, + terminalRef: { id: producerId, status: "done" }, + wake: "follow-up", + payload: { value: producerId }, + ...patch, + }; +} + +test("all producers share one owner generation for a Pi Session identity", () => { + let sessionId = "session-1"; + const identity = { getSessionId: () => sessionId }; + const first = completionOwnerFor(identity); + assert.equal(completionOwnerFor(identity), first); + + sessionId = "session-2"; + const switched = completionOwnerFor(identity); + assert.equal(switched.sessionId, "session-2"); + assert.notEqual(switched.epoch, first.epoch); + + const replacement = completionOwnerFor({ getSessionId: () => sessionId }); + assert.notEqual(replacement.epoch, switched.epoch); +}); + +test("explicit consumption and automatic claim share one atomic gate", () => { + const inbox = createCompletionInbox<{ value: string }>(); + inbox.defer(envelope("sa-1"), owner); + inbox.consume("subagent", ["sa-1"]); + assert.deepEqual(inbox.claim(owner), []); + + inbox.defer(envelope("sa-2"), owner); + assert.deepEqual( + inbox.claim(owner).map((item) => item.producerId), + ["sa-2"], + ); + inbox.consume("subagent", ["sa-2"]); + assert.deepEqual(inbox.claim(owner), []); +}); + +test("a failed transport restores exact envelopes ahead of newer work", () => { + const inbox = createCompletionInbox<{ value: string }>(); + const first = envelope("sa-1"); + const second = envelope("sa-2"); + inbox.defer(first, owner); + const claimed = inbox.claim(owner); + inbox.defer(second, owner); + inbox.retry(claimed, owner); + + assert.deepEqual( + inbox.claim(owner).map((item) => item.deliveryId), + [first.deliveryId, second.deliveryId], + ); +}); + +test("separate in-flight batches can be acknowledged or retried independently", () => { + const inbox = createCompletionInbox<{ value: string }>(); + const first = envelope("sa-1"); + const second = envelope("sa-2"); + const third = envelope("sa-3"); + for (const item of [first, second, third]) inbox.defer(item, owner); + + assert.deepEqual( + inbox.claim(owner, 2).map((item) => item.producerId), + ["sa-1", "sa-2"], + ); + assert.deepEqual( + inbox.claim(owner, 1).map((item) => item.producerId), + ["sa-3"], + ); + inbox.acknowledge([third.deliveryId]); + inbox.retryClaimed("subagent", [first.producerId, second.producerId], owner); + + assert.deepEqual( + inbox.claim(owner).map((item) => item.producerId), + ["sa-1", "sa-2"], + ); +}); + +test("a Session epoch switch dead-letters stale completions", () => { + const inbox = createCompletionInbox<{ value: string }>(); + inbox.defer(envelope("sa-1"), owner); + const nextOwner = { ...owner, epoch: owner.epoch + 1 }; + + assert.deepEqual(inbox.claim(nextOwner), []); + assert.deepEqual(inbox.inspectDeadLetters(), [ + { + deliveryId: "subagent:sa-1:terminal", + producer: "subagent", + producerId: "sa-1", + failure: "stale-owner", + }, + ]); +}); + +test("a Session identity switch cannot redirect a completion", () => { + const inbox = createCompletionInbox<{ value: string }>(); + inbox.defer(envelope("sa-1"), owner); + + assert.deepEqual( + inbox.claim({ sessionId: "session-2", epoch: owner.epoch }), + [], + ); + assert.equal(inbox.inspectDeadLetters()[0]?.failure, "stale-owner"); +}); + +test("producer identity prevents cross-capability consumption", () => { + const inbox = createCompletionInbox<{ value: string }>(); + inbox.defer(envelope("same"), owner); + inbox.defer( + envelope("same", { + deliveryId: "background:same:terminal", + producer: "background", + wake: "next-turn", + }), + owner, + ); + + inbox.consume("subagent", ["same"]); + assert.deepEqual( + inbox.claim(owner).map((item) => item.producer), + ["background"], + ); +}); + +test("an unavailable owner becomes an inspectable dead letter", () => { + const inbox = createCompletionInbox<{ value: string }>(); + assert.equal(inbox.defer(envelope("sa-1"), undefined), false); + assert.equal(inbox.size(), 0); + assert.equal(inbox.inspectDeadLetters()[0]?.failure, "owner-unavailable"); +}); diff --git a/tests/extensions/subagents/result-delivery.test.ts b/tests/extensions/subagents/result-delivery.test.ts index 7256defd..6ae29264 100644 --- a/tests/extensions/subagents/result-delivery.test.ts +++ b/tests/extensions/subagents/result-delivery.test.ts @@ -124,3 +124,20 @@ test("a synchronous delivery failure restores the batch in order", () => { assert.deepEqual(delivered, [["sa-1", "sa-2", "sa-3"]]); }); + +test("a completion cannot cross a parent Session switch", () => { + let sessionId = "session-1"; + const delivered: string[] = []; + const delivery = createSubagentResultDelivery<{ id: string }>({ + isIdle: () => false, + owner: () => ({ sessionId, epoch: 0 }), + deliver: (results) => delivered.push(...results.map((result) => result.id)), + }); + + delivery.defer({ id: "sa-1" }); + sessionId = "session-2"; + delivery.parentSettled(); + + assert.deepEqual(delivered, []); + assert.equal(delivery.inspectDeadLetters()[0]?.failure, "stale-owner"); +}); diff --git a/tests/extensions/workflows/execute.e2e.test.ts b/tests/extensions/workflows/execute.e2e.test.ts index 7db782a3..867b7430 100644 --- a/tests/extensions/workflows/execute.e2e.test.ts +++ b/tests/extensions/workflows/execute.e2e.test.ts @@ -782,17 +782,14 @@ test("shutdown preserves a failed completion for reload recovery", async () => { "delivered", ); - // The original instance represents the process that was replaced. Drain - // its intentionally retained in-memory retry so later tests do not batch it - // with an unrelated completion; the assertion above was reached solely via - // the fresh instance and persisted session state. + // The original instance represents the process that was replaced. Its + // owner was disposed at shutdown, so the shared inbox must dead-letter the + // stale in-memory retry rather than append a duplicate to the revived + // Session. The fresh instance recovered solely from canonical durable state. for (const handler of handlers.get("agent_settled") ?? []) { await handler({}, ctx); } - await waitFor( - () => sentMessages.length === 2, - "discarded predecessor instance retry", - ); + assert.equal(sentMessages.length, 1); sentMessages.length = 0; }); diff --git a/tests/extensions/workflows/result-delivery.test.ts b/tests/extensions/workflows/result-delivery.test.ts index 27fe0adb..63eac1a7 100644 --- a/tests/extensions/workflows/result-delivery.test.ts +++ b/tests/extensions/workflows/result-delivery.test.ts @@ -197,6 +197,65 @@ test("stale held inline completion restores as pending", () => { assert.equal(delivery.size(), 1); }); +test("a restored workflow completion cannot enter another Session", () => { + const run = details("wf_stale_owner"); + run.sessionId = "session-1"; + run.delivery = { + ...run.delivery!, + ownerSessionId: "session-1", + ownerEpoch: 0, + state: "pending", + }; + const delivery = createWorkflowResultDelivery({ + isIdle: () => false, + owner: () => ({ sessionId: "session-2", epoch: 0 }), + persist: () => {}, + deliver: async () => [], + }); + + assert.equal( + delivery.restore({ + deliveryId: run.delivery.id, + runId: run.runId, + details: run, + }), + false, + ); + assert.equal(delivery.size(), 0); + assert.equal(delivery.inspectDeadLetters()[0]?.failure, "stale-owner"); + assert.equal(run.delivery.state, "pending"); +}); + +test("same-Session restore explicitly revives a pending completion", () => { + const run = details("wf_revived_owner"); + run.sessionId = "session-1"; + run.delivery = { + ...run.delivery!, + ownerSessionId: "session-1", + ownerEpoch: 2, + state: "pending", + }; + const delivery = createWorkflowResultDelivery({ + isIdle: () => false, + owner: () => ({ sessionId: "session-1", epoch: 8 }), + persist: () => {}, + deliver: async () => [], + }); + + assert.equal( + delivery.restore({ + deliveryId: run.delivery.id, + runId: run.runId, + details: run, + }), + true, + ); + assert.equal(run.delivery.ownerEpoch, 8); + assert.equal(run.delivery.state, "pending"); + assert.equal(delivery.size(), 1); + assert.deepEqual(delivery.inspectDeadLetters(), []); +}); + test("a receipt persistence failure retains the same delivery for at-least-once recovery", async () => { const run = details("wf_receipt"); const delivery = createWorkflowResultDelivery({