From 5fdf869f11e329cf67f177c08cd63f997bb42ccf Mon Sep 17 00:00:00 2001 From: nickyleach <331803+nickyleach@users.noreply.github.com> Date: Fri, 11 Sep 2026 19:55:17 +0000 Subject: [PATCH] fix(acpx-engine): drain the permission observer's pending log writes at finalization The observer starts each durable log write and does not await it on the permission critical path. finalizeRun could return before its own acpx.permission_unsettled or acpx.permission_observer_truncated write, or an earlier still-pending write, reached the log sink. finalizeRun now tracks every write it starts and waits for the full set to settle before it returns, so a caller that awaits finalizeRun sees every queued record land first. Co-authored-by: Paperclip --- .../adapter-utils/src/acpx-engine/execute.ts | 16 ++- .../acpx-engine/permission-observer.test.ts | 125 +++++++++++++++--- .../src/acpx-engine/permission-observer.ts | 50 +++++-- 3 files changed, 156 insertions(+), 35 deletions(-) diff --git a/packages/adapter-utils/src/acpx-engine/execute.ts b/packages/adapter-utils/src/acpx-engine/execute.ts index 795f1adcc8..35779d257b 100644 --- a/packages/adapter-utils/src/acpx-engine/execute.ts +++ b/packages/adapter-utils/src/acpx-engine/execute.ts @@ -4124,11 +4124,12 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { transport: !observedExecutionTarget || observedExecutionTarget.kind === "local" ? "local" : observedExecutionTarget.transport, - emitLog: (payload) => { - // Fire-and-forget: the observer's caller must not await a log - // write on the permission critical path. - void emitAcpxLog(ctx, payload).catch(() => {}); - }, + // Return the write's promise; do not swallow it here. The observer + // never awaits this on the permission critical path, but it keeps + // the promise so `finalizeRun` can drain every pending write before + // it returns, so the run never finalizes with a write still in + // flight. + emitLog: (payload) => emitAcpxLog(ctx, payload), }); // Capture the run's staging lease release now that the runtime built. The // run root `finally` releases it as the final settlement act. @@ -5342,8 +5343,9 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { recordDispositionReport(report); if (childStderrState) flushChildStderr(childStderrState); // Report every permission request the run never saw settle. Not on - // the permission critical path. `emitLog` returns `void`, so this - // call awaits no log write. + // the permission critical path. This call waits for every queued + // log write, including one still in flight from an earlier + // permission event, to reach durable storage before it returns. await permissionObserver?.finalizeRun(); }, reproduceResult: async (): Promise => { diff --git a/packages/adapter-utils/src/acpx-engine/permission-observer.test.ts b/packages/adapter-utils/src/acpx-engine/permission-observer.test.ts index 3fedca7c70..bc41f57989 100644 --- a/packages/adapter-utils/src/acpx-engine/permission-observer.test.ts +++ b/packages/adapter-utils/src/acpx-engine/permission-observer.test.ts @@ -47,7 +47,9 @@ describe("createAcpPermissionObserver — handlePermissionRequest", () => { it("resolves to undefined for a normal request", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -58,7 +60,9 @@ describe("createAcpPermissionObserver — handlePermissionRequest", () => { it("resolves to undefined for a request that carries unknown fields", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -115,7 +119,9 @@ describe("createAcpPermissionObserver — handlePermissionRequest", () => { it("emits only allow-listed scalar fields, never the raw request payload", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -149,7 +155,9 @@ describe("createAcpPermissionObserver — handlePermissionRequest", () => { it("maps an unmapped tool kind to exactly 'unknown', dropping the source string", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -165,7 +173,9 @@ describe("createAcpPermissionObserver — ledger lifecycle", () => { let clock = 1_000; const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", now: () => clock, @@ -193,7 +203,9 @@ describe("createAcpPermissionObserver — ledger lifecycle", () => { let clock = 0; const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", now: () => clock, @@ -226,7 +238,9 @@ describe("createAcpPermissionObserver — ledger lifecycle", () => { it("does not close an entry on a non-terminal tool_call status", () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -237,7 +251,9 @@ describe("createAcpPermissionObserver — ledger lifecycle", () => { it("ignores a tool_call event for a toolCallId it never opened", () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -248,7 +264,9 @@ describe("createAcpPermissionObserver — ledger lifecycle", () => { it("does not open a ledger entry for a request with no tool-call identifier", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -269,7 +287,9 @@ describe("createAcpPermissionObserver — ledger lifecycle", () => { it("does not let two requests with a missing session identifier share one ledger entry", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -308,7 +328,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("caps the ledger and the observed-event count when more than 256 requests never settle", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -331,7 +353,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("bounds cumulative settled emissions at 256 across more than 1000 open-and-settle cycles", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -354,7 +378,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("keeps the unsettled-event budget reachable after a normal run spends the settled budget", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -379,7 +405,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("bounds cumulative unsettled emissions at 256 when more than 256 requests never settle", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -397,7 +425,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("emits exactly one truncated summary event across two finalizeRun calls", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -415,7 +445,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("keeps the settled-event budget reachable after the observed budget is spent", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -444,7 +476,9 @@ describe("createAcpPermissionObserver — per-run budgets", () => { it("emits a summary event with only the four counters and the type, using an exact key match", async () => { const events: PermissionObserverLogEvent[] = []; const observer = createAcpPermissionObserver({ - emitLog: (event) => events.push(event), + emitLog: (event) => { + events.push(event); + }, permissionMode: "approve-all", transport: "sandbox", }); @@ -487,3 +521,60 @@ describe("createAcpPermissionObserver — per-run budgets", () => { await expect(observer.finalizeRun()).resolves.toBeUndefined(); }); }); + +describe("createAcpPermissionObserver — asynchronous log persistence", () => { + it("waits for a write still pending from an earlier event before finalizeRun resolves", async () => { + const durable: PermissionObserverLogEvent[] = []; + const releaseWrite: Array<() => void> = []; + const observer = createAcpPermissionObserver({ + // A durable log sink that starts a write and only completes it when the + // test calls the matching release function. This stands in for a real + // write to storage that takes more than one microtask turn. + emitLog: (event) => + new Promise((resolve) => { + releaseWrite.push(() => { + durable.push(event); + resolve(); + }); + }), + permissionMode: "approve-all", + transport: "sandbox", + }); + const signal = new AbortController().signal; + // This request's "observed" write starts but does not finish. The tool + // call it opened never settles, so the entry is still open at + // finalization. + await observer.handlePermissionRequest(buildRequest(), { signal }); + expect(durable).toHaveLength(0); + + const finalizePromise = observer.finalizeRun(); + // finalizeRun has started its own "unsettled" write for the still-open + // entry. That write is also pending. Flush a few microtask turns: if + // finalizeRun dropped either write's promise, it would resolve here even + // though neither write has completed. + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + expect(durable).toHaveLength(0); + + // Release both writes, the one queued before finalizeRun ran and the one + // finalizeRun queued itself. + for (const release of releaseWrite) release(); + await finalizePromise; + + expect(durable.map((event) => event.type).sort()).toEqual( + ["acpx.permission_observed", "acpx.permission_unsettled"].sort(), + ); + }); + + it("resolves finalizeRun without throwing when a queued write rejects", async () => { + const observer = createAcpPermissionObserver({ + emitLog: () => Promise.reject(new Error("durable sink unavailable")), + permissionMode: "approve-all", + transport: "sandbox", + }); + const signal = new AbortController().signal; + await observer.handlePermissionRequest(buildRequest(), { signal }); + await expect(observer.finalizeRun()).resolves.toBeUndefined(); + }); +}); diff --git a/packages/adapter-utils/src/acpx-engine/permission-observer.ts b/packages/adapter-utils/src/acpx-engine/permission-observer.ts index 822f2b9cf2..33652b4fb5 100644 --- a/packages/adapter-utils/src/acpx-engine/permission-observer.ts +++ b/packages/adapter-utils/src/acpx-engine/permission-observer.ts @@ -136,11 +136,12 @@ export interface PermissionObserverLogEvent { export interface AcpPermissionObserverOptions { /** - * Starts the durable log write. The observer never awaits this on the - * permission critical path — call it and return, do not `await` it inside - * `handlePermissionRequest`. + * Starts the durable log write and returns a promise for it. The observer + * never awaits this promise on the permission critical path — call it and + * return, do not `await` it inside `handlePermissionRequest`. The observer + * still tracks the promise so `finalizeRun` can wait for it later. */ - emitLog: (event: PermissionObserverLogEvent) => void; + emitLog: (event: PermissionObserverLogEvent) => Promise | void; /** The engine's effective permission mode for this run. */ permissionMode: unknown; /** The run's execution transport. */ @@ -166,9 +167,12 @@ export interface AcpPermissionObserver { */ noteToolCallEvent: (sessionId: string | undefined, event: PermissionObserverToolCallEvent) => void; /** - * Emit one `acpx.permission_unsettled` event per entry still open. Call - * this once, at run finalization. Not on the permission critical path. - * `emitLog` returns `void`, so this call awaits no log write. + * Emit one `acpx.permission_unsettled` event per entry still open, then + * wait for every log write this observer started — including a write + * still in flight from an earlier `handlePermissionRequest` or + * `noteToolCallEvent` call — to finish. Call this once, at run + * finalization. A caller that awaits this method sees every one of the + * observer's log writes land before it returns. */ finalizeRun: () => Promise; } @@ -191,6 +195,26 @@ export function createAcpPermissionObserver(options: AcpPermissionObserverOption let suppressedUnsettledEvents = 0; let hasFinalized = false; + // Every write `emitLog` starts stays in this set until it settles. A write + // started from the permission critical path (`handlePermissionRequest`, + // `noteToolCallEvent`) is never awaited there, so it can still be pending + // when `finalizeRun` runs. `finalizeRun` drains this whole set before it + // returns, so a caller that awaits `finalizeRun` never sees it resolve + // before every queued write, including its own unsettled and truncation + // records, has reached the log sink. + const pendingWrites = new Set>(); + + const queueLog = (event: PermissionObserverLogEvent): void => { + const write = (async () => { + await options.emitLog(event); + })().catch(() => { + // The log sink is diagnostic only; a write failure must never surface + // into the permission critical path or into `finalizeRun`. + }); + pendingWrites.add(write); + void write.finally(() => pendingWrites.delete(write)); + }; + const ledgerKey = (sessionId: string, toolCallId: string) => sessionId + "\u0000" + toolCallId; const handlePermissionRequest: AcpPermissionObserver["handlePermissionRequest"] = async (request) => { @@ -226,7 +250,7 @@ export function createAcpPermissionObserver(options: AcpPermissionObserverOption // an agent cannot refill the budget by settling old requests. if (observedEventCount < MAX_OBSERVED_EVENTS) { observedEventCount += 1; - options.emitLog({ + queueLog({ type: "acpx.permission_observed", sessionId, toolCallId, @@ -268,7 +292,7 @@ export function createAcpPermissionObserver(options: AcpPermissionObserverOption ledger.delete(key); if (settledEventCount < MAX_SETTLED_EVENTS) { settledEventCount += 1; - options.emitLog({ + queueLog({ type: "acpx.permission_settled", sessionId: entry.sessionId, toolCallId: entry.toolCallId, @@ -296,7 +320,7 @@ export function createAcpPermissionObserver(options: AcpPermissionObserverOption for (const entry of openEntries) { if (unsettledEventCount < MAX_UNSETTLED_EVENTS) { unsettledEventCount += 1; - options.emitLog({ + queueLog({ type: "acpx.permission_unsettled", sessionId: entry.sessionId, toolCallId: entry.toolCallId, @@ -317,7 +341,7 @@ export function createAcpPermissionObserver(options: AcpPermissionObserverOption suppressedSettledEvents > 0 || suppressedUnsettledEvents > 0 ) { - options.emitLog({ + queueLog({ type: "acpx.permission_observer_truncated", suppressedLedgerEntries, suppressedObservedEvents, @@ -325,6 +349,10 @@ export function createAcpPermissionObserver(options: AcpPermissionObserverOption suppressedUnsettledEvents, }); } + // Wait for every write this observer started, including a write still in + // flight from an earlier call, so the caller sees every record reach the + // log sink before this method returns. + await Promise.allSettled([...pendingWrites]); }; return { handlePermissionRequest, noteToolCallEvent, finalizeRun };