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 };