From f9cbd8e44ed53a55d39c4edee79a5c398fd81f0e Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 14 Sep 2026 15:33:17 -0700 Subject: [PATCH] Fail fast when an approval delivery is never accepted A correlated approval decision waited for reactor acceptance with no deadline, so one delivery with no observed stream event wedged the serial sessionOps tail and every later approval timed out behind it. The wait now races a 30s deadline that abandons the waiter and throws diagnostics naming the correlation, stage, outcome, and that the action may still have been applied; retries never re-deliver a handed-over decision. --- src/tui/approval-delivery.test.ts | 187 +++++++++++++++++++++++++ src/tui/approval-delivery.ts | 167 ++++++++++++++++++++++ src/tui/correlation-acceptance.test.ts | 16 +++ src/tui/correlation-acceptance.ts | 11 ++ src/tui/runner/session.ts | 22 ++- 5 files changed, 390 insertions(+), 13 deletions(-) create mode 100644 src/tui/approval-delivery.test.ts create mode 100644 src/tui/approval-delivery.ts diff --git a/src/tui/approval-delivery.test.ts b/src/tui/approval-delivery.test.ts new file mode 100644 index 000000000..b7694b291 --- /dev/null +++ b/src/tui/approval-delivery.test.ts @@ -0,0 +1,187 @@ +import { describe, expect, test } from "bun:test"; +import type { InboundMessage } from "@intx/types/runtime"; + +import { createSessionOperationQueue } from "./delivery-queue.js"; +import { + ApprovalDeliveryTimeoutError, + createApprovalDeliverer, +} from "./approval-delivery.js"; + +function decisionMessage( + correlationId: string | undefined, + outcome: "approved" | "rejected", +): InboundMessage { + return { + ref: { uid: 0, mailbox: "approval" }, + headers: { + from: "approval@local", + to: ["agent@local"], + date: new Date().toISOString(), + messageId: `approval-${correlationId ?? "uncorrelated"}`, + ...(correlationId !== undefined + ? { interchangeCorrelationId: correlationId } + : {}), + }, + flags: [], + content: JSON.stringify({ outcome }), + signatureStatus: "missing", + } satisfies InboundMessage; +} + +function controllableAcceptance() { + const waiters = new Map void>(); + return { + wait(correlationId: string): Promise { + return new Promise((resolve) => { + waiters.set(correlationId, () => { + waiters.delete(correlationId); + resolve(); + }); + }); + }, + settle(correlationId: string): void { + waiters.get(correlationId)?.(); + }, + abandon(correlationId: string): void { + waiters.delete(correlationId); + }, + pendingCount(): number { + return waiters.size; + }, + }; +} + +describe("approval delivery acceptance bound", () => { + test("a stuck acceptance wait fails fast and the sessionOps tail advances", async () => { + const acceptance = controllableAcceptance(); + const delivered: InboundMessage[] = []; + const deliverer = createApprovalDeliverer({ + deliverToAgent: (message) => { + delivered.push(message); + }, + acceptance, + timeoutMs: 15, + }); + const sessionOps = createSessionOperationQueue(); + const started = Date.now(); + + const stuck = sessionOps.enqueue(() => + deliverer.deliver(decisionMessage("corr-stuck", "approved")), + ); + let tailAdvanced = false; + const next = sessionOps.enqueue(async () => { + tailAdvanced = true; + }); + + const failure = await stuck.then( + () => null, + (err: unknown) => err, + ); + await next; + + expect(failure).toBeInstanceOf(ApprovalDeliveryTimeoutError); + const err = failure as ApprovalDeliveryTimeoutError; + expect(err.correlationId).toBe("corr-stuck"); + expect(err.stage).toBe("reactor-acceptance"); + expect(err.timeoutMs).toBe(15); + expect(err.mayStillApply).toBe(true); + expect(err.message).toContain("corr-stuck"); + expect(err.message).toContain("reactor-acceptance"); + expect(err.message).toContain("may still"); + expect(tailAdvanced).toBe(true); + expect(Date.now() - started).toBeLessThan(5000); + expect(delivered).toHaveLength(1); + expect(acceptance.pendingCount()).toBe(0); + }); + + test("retry after an acceptance timeout does not re-deliver a handed-over decision", async () => { + const acceptance = controllableAcceptance(); + const delivered: InboundMessage[] = []; + const deliverer = createApprovalDeliverer({ + deliverToAgent: (message) => { + delivered.push(message); + }, + acceptance, + timeoutMs: 10, + }); + const message = decisionMessage("corr-retry", "approved"); + + await expect(deliverer.deliver(message)).rejects.toBeInstanceOf( + ApprovalDeliveryTimeoutError, + ); + expect(delivered).toHaveLength(1); + + const retry = deliverer.deliver(message); + acceptance.settle("corr-retry"); + await retry; + + expect(delivered).toHaveLength(1); + }); + + test("retry after a deliver throw re-delivers: recovery from a failed send", async () => { + const acceptance = controllableAcceptance(); + const delivered: InboundMessage[] = []; + let failDeliver = true; + const deliverer = createApprovalDeliverer({ + deliverToAgent: (message) => { + if (failDeliver) throw new Error("agent is done"); + delivered.push(message); + }, + acceptance, + timeoutMs: 50, + }); + const message = decisionMessage("corr-failed-send", "rejected"); + + await expect(deliverer.deliver(message)).rejects.toThrow("agent is done"); + expect(delivered).toHaveLength(0); + + failDeliver = false; + const retry = deliverer.deliver(message); + acceptance.settle("corr-failed-send"); + await retry; + + expect(delivered).toHaveLength(1); + }); + + test("an uncorrelated decision delivers without waiting for acceptance", async () => { + const acceptance = controllableAcceptance(); + const delivered: InboundMessage[] = []; + const deliverer = createApprovalDeliverer({ + deliverToAgent: (message) => { + delivered.push(message); + }, + acceptance, + timeoutMs: 10, + }); + + await deliverer.deliver(decisionMessage(undefined, "rejected")); + + expect(delivered).toHaveLength(1); + expect(acceptance.pendingCount()).toBe(0); + }); + + test("timeout diagnostics name the rejected outcome and its uncertainty", async () => { + const acceptance = controllableAcceptance(); + const delivered: InboundMessage[] = []; + const deliverer = createApprovalDeliverer({ + deliverToAgent: (message) => { + delivered.push(message); + }, + acceptance, + timeoutMs: 10, + }); + + const failure = await deliverer + .deliver(decisionMessage("corr-reject", "rejected")) + .then( + () => null, + (err: unknown) => err, + ); + + expect(failure).toBeInstanceOf(ApprovalDeliveryTimeoutError); + const err = failure as ApprovalDeliveryTimeoutError; + expect(err.outcome).toBe("rejected"); + expect(err.mayStillApply).toBe(true); + expect(err.message).toContain("rejected"); + }); +}); diff --git a/src/tui/approval-delivery.ts b/src/tui/approval-delivery.ts new file mode 100644 index 000000000..6b7e9721c --- /dev/null +++ b/src/tui/approval-delivery.ts @@ -0,0 +1,167 @@ +/** + * Bounded deliver-and-await-acceptance for approval decisions. + * + * The reactor accepts a correlated decision asynchronously after deliver() + * returns, so the sessionOps tail waits for the correlation-acceptance signal. + * That wait was unbounded: any delivery that produces no observed stream event + * (reactor deliver() silently drops when done, an approved hold with no + * tool.start, a missed correlation event) wedged the serial tail forever, and + * every later approval — send_input answers, interrupt_agent releases, + * ask_operator escalations — queued behind it until the reactor approval + * timeout. The vendored reactor surface ({ start, deliver, abort }) exposes no + * liveness query, so absent acceptance is treated as delivery failure: bound + * the wait with a deadline race that settles the waiter, logs, and lets the + * tail advance. + * + * Retry safety: once deliver() has returned, the reactor may hold the decision + * even though no acceptance was observed. A retry for the same correlationId + * must not hand the decision over twice (a duplicate approved decision could + * re-dispatch the parked call), so it re-awaits acceptance only. A retry after + * a deliver() throw re-delivers, because nothing was handed over. + */ + +import type { InboundMessage } from "@intx/types/runtime"; + +import { getLogger } from "@intx/log"; + +import { LOG_NAMESPACE_ROOT } from "../branding.js"; + +const logger = getLogger([LOG_NAMESPACE_ROOT, "approval-delivery"]); + +/** + * Deadline for the reactor to observably accept a delivered approval decision. + * Far above normal tick latency (acceptance usually lands within a reactor + * tick, including the approved hold until tool.start) and far below the + * vendored reactor approval timeout, so one stuck delivery fails fast instead + * of wedging the tail. + */ +export const APPROVAL_ACCEPTANCE_TIMEOUT_MS = 30_000; + +export type ApprovalDecisionOutcome = "approved" | "rejected" | "unknown"; + +export class ApprovalDeliveryTimeoutError extends Error { + readonly correlationId: string; + readonly stage = "reactor-acceptance"; + readonly timeoutMs: number; + readonly outcome: ApprovalDecisionOutcome; + /** deliver() returned, so the reactor may still act on the decision. */ + readonly mayStillApply = true; + + constructor(args: { + correlationId: string; + timeoutMs: number; + outcome: ApprovalDecisionOutcome; + }) { + const action = + args.outcome === "approved" + ? "the approved action may still run" + : args.outcome === "rejected" + ? "the parked call may still be answered by this decision" + : "the decision may still be applied"; + super( + `approval delivery timed out waiting for reactor acceptance ` + + `(correlationId=${args.correlationId}, stage=reactor-acceptance, ` + + `timeoutMs=${args.timeoutMs}, outcome=${args.outcome}). ` + + `The decision was handed to the reactor but no acceptance event was ` + + `observed, so ${action}. Do not re-deliver: the delivery tail has ` + + `advanced. Check run state, then use interrupt_agent to release the ` + + `parked worker if it is still parked.`, + ); + this.name = "ApprovalDeliveryTimeoutError"; + this.correlationId = args.correlationId; + this.timeoutMs = args.timeoutMs; + this.outcome = args.outcome; + } +} + +function decisionOutcome(message: InboundMessage): ApprovalDecisionOutcome { + if (message.content === undefined) return "unknown"; + let raw: unknown; + try { + raw = JSON.parse(message.content); + } catch { + return "unknown"; + } + if (raw === null || typeof raw !== "object") return "unknown"; + const outcome = (raw as { outcome?: unknown }).outcome; + return outcome === "approved" || outcome === "rejected" ? outcome : "unknown"; +} + +export interface ApprovalAcceptance { + wait(correlationId: string): Promise; + settle(correlationId: string): void; + abandon(correlationId: string): void; +} + +export interface ApprovalDelivererDeps { + deliverToAgent: (message: InboundMessage) => void; + acceptance: ApprovalAcceptance; + timeoutMs?: number; + onTimeout?: (err: ApprovalDeliveryTimeoutError) => void; +} + +export interface ApprovalDeliverer { + deliver: (message: InboundMessage) => Promise; +} + +export function createApprovalDeliverer( + deps: ApprovalDelivererDeps, +): ApprovalDeliverer { + const timeoutMs = deps.timeoutMs ?? APPROVAL_ACCEPTANCE_TIMEOUT_MS; + const delivered = new Set(); + + const awaitAcceptance = async ( + correlationId: string, + accepted: Promise, + outcome: ApprovalDecisionOutcome, + ): Promise => { + let timer: ReturnType | undefined; + try { + await Promise.race([ + accepted, + new Promise((_, reject) => { + timer = setTimeout(() => { + deps.acceptance.abandon(correlationId); + const err = new ApprovalDeliveryTimeoutError({ + correlationId, + timeoutMs, + outcome, + }); + logger.warn`approval acceptance timed out correlation=${correlationId} timeoutMs=${timeoutMs} outcome=${outcome}`; + deps.onTimeout?.(err); + reject(err); + }, timeoutMs); + }), + ]); + } finally { + clearTimeout(timer); + } + }; + + return { + deliver: async (message) => { + const correlationId = message.headers.interchangeCorrelationId; + if (correlationId === undefined) { + deps.deliverToAgent(message); + return; + } + if (delivered.has(correlationId)) { + await awaitAcceptance( + correlationId, + deps.acceptance.wait(correlationId), + decisionOutcome(message), + ); + return; + } + const accepted = deps.acceptance.wait(correlationId); + try { + deps.deliverToAgent(message); + } catch (err) { + deps.acceptance.settle(correlationId); + throw err; + } + delivered.add(correlationId); + await awaitAcceptance(correlationId, accepted, decisionOutcome(message)); + }, + }; +} diff --git a/src/tui/correlation-acceptance.test.ts b/src/tui/correlation-acceptance.test.ts index 2fe0c7360..3dcc7b62b 100644 --- a/src/tui/correlation-acceptance.test.ts +++ b/src/tui/correlation-acceptance.test.ts @@ -37,6 +37,22 @@ describe("createCorrelationAcceptance", () => { acceptance.settle("missing"); }); + test("abandon drops the waiter without resolving it", async () => { + const acceptance = createCorrelationAcceptance(); + const pending = acceptance.wait("corr-1"); + let settled = false; + void pending.then(() => { + settled = true; + }); + acceptance.abandon("corr-1"); + await Promise.resolve(); + await new Promise((resolve) => setTimeout(resolve, 5)); + expect(settled).toBe(false); + const retry = acceptance.wait("corr-1"); + acceptance.settle("corr-1"); + await retry; + }); + test("an approved correlation does not settle until tool.start", async () => { const acceptance = createCorrelationAcceptance(); const pending = acceptance.wait("corr-1"); diff --git a/src/tui/correlation-acceptance.ts b/src/tui/correlation-acceptance.ts index 8a65af72d..9d50eeda4 100644 --- a/src/tui/correlation-acceptance.ts +++ b/src/tui/correlation-acceptance.ts @@ -78,6 +78,16 @@ export function createCorrelationAcceptance() { for (const correlationId of [...holdUntilToolStart]) settle(correlationId); }; + /** + * Drop a waiter without resolving it. The acceptance deadline path uses + * this instead of settle: settling would fulfill the very promise the + * deadline race is trying to reject, so the timeout could never fire. + */ + const abandon = (correlationId: string): void => { + holdUntilToolStart.delete(correlationId); + waiters.delete(correlationId); + }; + return { wait(correlationId: string): Promise { const pending = waiters.get(correlationId); @@ -89,6 +99,7 @@ export function createCorrelationAcceptance() { }); }, settle, + abandon, settleAll(): void { holdUntilToolStart.clear(); for (const correlationId of [...waiters.keys()]) settle(correlationId); diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index 2940f5c69..3c28cafc9 100644 --- a/src/tui/runner/session.ts +++ b/src/tui/runner/session.ts @@ -82,6 +82,7 @@ import { createSessionOperationQueue, } from "../delivery-queue.js"; import { createCorrelationAcceptance } from "../correlation-acceptance.js"; +import { createApprovalDeliverer } from "../approval-delivery.js"; import { createAgentToolset, type MCPServerState, @@ -506,6 +507,13 @@ export async function assembleTUISession( const sessionOps = createSessionOperationQueue(); // No resolveParkedCallId: the vendored reactor exposes no // correlationId-to-call lookup, so the history heuristic is the path. + // The deliverer bounds the acceptance wait so one stuck delivery fails fast + // with diagnostics instead of wedging the sessionOps tail for every later + // approval (send_input answers, interrupt_agent releases, ask_operator). + const approvalDeliverer = createApprovalDeliverer({ + deliverToAgent: (message) => liveAgent(state).deliver(message), + acceptance: correlationAcceptance, + }); const approvalResume = createApprovalResume({ getAgent: () => state.currentAgent, captureGeneration: deliveryGeneration.capture, @@ -517,19 +525,7 @@ export async function assembleTUISession( return sessionOps.enqueue(async () => { if (!stillCurrent()) return; if (state.fatalBuildError !== null) throw state.fatalBuildError; - const correlationId = message.headers.interchangeCorrelationId; - const accepted = - correlationId === undefined - ? undefined - : correlationAcceptance.wait(correlationId); - try { - liveAgent(state).deliver(message); - await accepted; - } catch (err) { - if (correlationId !== undefined) - correlationAcceptance.settle(correlationId); - throw err; - } + await approvalDeliverer.deliver(message); }); }, gate: permissionGate,