Await provider reply delivery in the direct ACPX driver

Fail closed without reporting resolution when the pipe write is ambiguous. Preserve provider-supported Copilot session permissions.

Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip committed 2026-09-28 12:03:50 -05:00
1 parent a3a71ebdb1
commit d4e953e443
5 files changed
+66 -18

No files matched your search

@@ -1,7 +1,7 @@
#!/usr/bin/env node
import { createHash } from "node:crypto";
import { createInterface } from "node:readline";
import { deliverAcpxResponse, requireAcpxResponseDelivery } from "./acpx-response-delivery.js";
import { deliverAcpxResponse, requireAcpxResponseDelivery } from "../drivers/acpx/response-delivery.js";
import type {
AcpElicitationContext,
@@ -1513,7 +1513,7 @@ describe("Codex ACPX harness driver", () => {
sessionId: "agent-session-1", toolCall: { toolCallId: "tool-permission", title: "Run validation" },
options: [{ optionId: "once", kind: "allow_once", name: "Allow once" }],
},
} as Parameters<typeof callback>[0], { signal: new AbortController().signal });
} as Parameters<typeof callback>[0], { signal: new AbortController().signal, responseDelivery: Promise.resolve() });
const events = await created;
const request = session.pendingRuntimeRequests!()[0]!;
expect(events.at(-1)?.payload).toMatchObject({ request: {
@@ -1534,6 +1534,39 @@ describe("Codex ACPX harness driver", () => {
await session.close({ reason: "permission verified" });
});
it.each(["written", "failed"] as const)("waits for the exact provider permission reply receipt: %s", async outcome => {
const fixture = driverFixture({ agent: "copilot", model: "explicit-test-model", providerPolicy: { readOnly: false } });
const session = await fixture.driver.openSession({ runId: "run-receipt", normalizedSessionId: "session-1", workingDirectory: "/workspace" });
const created = collectUntil(session.events(), "runtime_request.created");
const { turnId } = await session.startTurn({ message: { role: "user", text: "Approve with receipt." } });
const callback = fixture.host.startTurn.mock.calls[0]![0].onPermissionRequest!;
const receipt = deferred<void>();
const providerResponse = callback({ inferredKind: "execute", raw: {
sessionId: "agent-session-1", toolCall: { toolCallId: "receipt-tool", title: "Run validation" },
options: [{ optionId: "session", kind: "allow_always", name: "Allow for session" }],
} } as Parameters<typeof callback>[0], { signal: new AbortController().signal, responseDelivery: receipt.promise });
await created;
const request = session.pendingRuntimeRequests!()[0]!;
expect(request.details).toMatchObject({ choices: [{ key: "accept_for_session" }, { key: "cancel" }] });
let acknowledged = false;
const emitted = collectUntil(session.events(), outcome === "written" ? "runtime_request.resolved" : "runtime_request.expired");
const resolution = session.resolveRuntimeRequest!({ requestId: request.requestId, turnId, resolution: { action: "accept_for_session" } })
.then(() => { acknowledged = true; });
await expect(providerResponse).resolves.toEqual({ outcome: "allow_always" });
expect(acknowledged).toBe(false);
if (outcome === "written") { receipt.resolve(); await resolution; }
else {
const rejected = expect(resolution).rejects.toThrow("pipe failed");
receipt.reject(new Error("pipe failed")); await rejected;
}
const events = await emitted;
expect(events.filter(event => event.eventType === "runtime_request.resolved")).toHaveLength(outcome === "written" ? 1 : 0);
if (outcome === "failed") expect(events.at(-1)?.payload).toMatchObject({ replayAllowed: false, reason: "response_delivery_failed" });
expect(session.pendingRuntimeRequests!()).toHaveLength(0);
if (outcome === "written") fixture.finishTurn({ status: "completed", stopReason: "end_turn" });
await session.close({ reason: "receipt checked" });
});
it("round-trips a provider-neutral ACP form through the runtime request boundary", async () => {
const fixture = driverFixture();
const session = await fixture.driver.openSession({
@@ -1568,7 +1601,7 @@ describe("Codex ACPX harness driver", () => {
},
},
},
{ requestId: "rpc-question-1", signal: controller.signal },
{ requestId: "rpc-question-1", signal: controller.signal, responseDelivery: Promise.resolve() },
);
const events = await createdEvent;
const request = session.pendingRuntimeRequests!()[0]!;
@@ -1650,7 +1683,7 @@ describe("Codex ACPX harness driver", () => {
},
{
requestId: "rpc-question-created-pressure",
signal: new AbortController().signal,
signal: new AbortController().signal, responseDelivery: Promise.resolve(),
},
),
).resolves.toEqual({ action: "cancel" });
@@ -1684,7 +1717,7 @@ describe("Codex ACPX harness driver", () => {
properties: { value: { type: "string" } },
},
},
{ requestId: "rpc-question-abort", signal: controller.signal },
{ requestId: "rpc-question-abort", signal: controller.signal, responseDelivery: Promise.resolve() },
);
await createdEvent;
const cancelledEvent = collectUntil(
@@ -1733,7 +1766,7 @@ describe("Codex ACPX harness driver", () => {
},
{
requestId: "rpc-question-stream-failure",
signal: new AbortController().signal,
signal: new AbortController().signal, responseDelivery: Promise.resolve(),
},
);
await vi.waitFor(() => {
@@ -1779,7 +1812,7 @@ describe("Codex ACPX harness driver", () => {
},
{
requestId: "rpc-question-handoff",
signal: new AbortController().signal,
signal: new AbortController().signal, responseDelivery: Promise.resolve(),
},
);
await createdEvent;
@@ -1847,7 +1880,7 @@ describe("Codex ACPX harness driver", () => {
},
{
requestId: "rpc-question-handoff-pressure",
signal: new AbortController().signal,
signal: new AbortController().signal, responseDelivery: Promise.resolve(),
},
);
await vi.waitFor(() => {
@@ -1,3 +1,4 @@
import { requireAcpxResponseDelivery } from "./response-delivery.js";
import { acpxProfileClientCapabilities, bindAcpxExtensionTurn, validateAcpxRichEvent, createAcpxProfileExtensionAdapter, type AcpxExtensionInput } from "./profile-extensions.js";
import { createHash, randomBytes } from "node:crypto";
@@ -115,6 +116,7 @@ const QUARANTINED_HOST_ADMISSION_GRACE_MS =
interface PendingAcpxRuntimeRequest {
request: HarnessRuntimeRequest;
responseDelivery: Promise<void>;
prepareResolution(resolution: HarnessRuntimeRequestResolution): () => void;
cancel(): void;
cleanup(): void;
@@ -894,7 +896,7 @@ class CodexAcpxSession implements HarnessSession {
onExtensionRequest: extensions.onExtensionRequest,
onExtensionNotification: extensions.onExtensionNotification,
onPermissionRequest: (request, context) =>
this.#handlePermission(turnId, request, context.signal),
this.#handlePermission(turnId, request, context),
onElicitation: (request, context) =>
this.#handleElicitation(turnId, request, context),
});
@@ -988,6 +990,7 @@ class CodexAcpxSession implements HarnessSession {
);
}
pending.settling = true;
let dispatched = false;
try {
const resolution = parseHarnessRuntimeRequestResolution(
pending.request.requestKind,
@@ -997,6 +1000,9 @@ class CodexAcpxSession implements HarnessSession {
const deliver = pending.prepareResolution(resolution);
if (!this.#pendingRuntimeRequests.delete(input.requestId)) return;
pending.cleanup();
dispatched = true;
deliver();
await pending.responseDelivery;
this.#emit(
"runtime_request.resolved",
harnessRuntimeRequestOutcome(pending.request, {
@@ -1007,9 +1013,14 @@ class CodexAcpxSession implements HarnessSession {
}),
{ turnId: input.turnId, itemId: pending.request.itemId },
);
deliver();
} catch (error) {
pending.settling = false;
if (!dispatched) { pending.settling = false; throw error; }
this.#emit("runtime_request.expired", {
...harnessRuntimeRequestOutcome(pending.request, { reason: "response_delivery_failed" }),
replayAllowed: false,
}, { turnId: input.turnId, itemId: pending.request.itemId });
try { await this.close({ reason: "ACP response delivery failed" }); }
catch (cleanupError) { throw new AggregateError([error, cleanupError], "ACP response delivery and provider cleanup failed"); }
throw error;
}
}
@@ -1640,10 +1651,11 @@ class CodexAcpxSession implements HarnessSession {
async #handleExtensionInput(
turnId: string,
input: AcpxExtensionInput,
context: { requestId: string | number; signal: AbortSignal },
context: { requestId: string | number; signal: AbortSignal; responseDelivery?: Promise<void> },
): Promise<Record<string, unknown>> {
if (this.#closed || this.#activeTurnId !== turnId || context.signal.aborted
|| this.#pendingRuntimeRequests.size >= MAX_PENDING_RUNTIME_REQUESTS) return input.cancel();
const responseDelivery = requireAcpxResponseDelivery(context);
const requestId = stableId("acpx-request", `${turnId}:${++this.#runtimeRequestSequence}:${typeof context.requestId}:${context.requestId}`);
const request: HarnessRuntimeRequest = {
requestId, requestKind: "elicitation", method: input.method, turnId, itemId: requestId,
@@ -1662,7 +1674,7 @@ class CodexAcpxSession implements HarnessSession {
settle(input.cancel());
};
this.#pendingRuntimeRequests.set(requestId, {
request,
request, responseDelivery,
prepareResolution: resolution => {
const response = input.resolve(resolution);
return () => settle(response);
@@ -1680,13 +1692,15 @@ class CodexAcpxSession implements HarnessSession {
async #handlePermission(
turnId: string,
request: AcpPermissionRequest,
signal: AbortSignal,
context: { signal: AbortSignal; responseDelivery?: Promise<void> },
): Promise<AcpPermissionDecision> {
const { signal } = context;
if (this.#closed || this.#activeTurnId !== turnId || signal.aborted
|| this.#pendingRuntimeRequests.size >= MAX_PENDING_RUNTIME_REQUESTS) {
return { outcome: "cancel" };
}
const normalized = normalizeAcpxPermission(request, this.#agent === "pi" ? { allowAlwaysScope: "session" } : {});
const responseDelivery = requireAcpxResponseDelivery(context);
const normalized = normalizeAcpxPermission(request, ["pi", "copilot"].includes(this.#agent) ? { allowAlwaysScope: "session" } : {});
const requestId = stableId("acpx-permission", `${turnId}:${++this.#runtimeRequestSequence}:${normalized.toolCallId}`);
const runtimeRequest: HarnessRuntimeRequest = {
requestId, requestKind: "permission_approval", method: "session/request_permission",
@@ -1705,7 +1719,7 @@ class CodexAcpxSession implements HarnessSession {
settle({ outcome: "cancel" });
};
this.#pendingRuntimeRequests.set(requestId, {
request: runtimeRequest,
request: runtimeRequest, responseDelivery,
prepareResolution: (resolution) => {
const response = normalized.resolve(resolution);
return () => settle(response);
@@ -1792,6 +1806,7 @@ class CodexAcpxSession implements HarnessSession {
);
return { action: "cancel" };
}
const responseDelivery = requireAcpxResponseDelivery(context);
const requestId = stableId(
"acpx-request",
`${turnId}:${++this.#runtimeRequestSequence}:${typeof context.requestId}:${String(context.requestId)}`,
@@ -1844,7 +1859,7 @@ class CodexAcpxSession implements HarnessSession {
};
context.signal.addEventListener("abort", cancel, { once: true });
this.#pendingRuntimeRequests.set(requestId, {
request: runtimeRequest,
request: runtimeRequest, responseDelivery,
prepareResolution: (resolution) => {
const response = acpElicitationResponse(normalized, resolution);
return () => settle(response);
@@ -1,5 +1,5 @@
import { describe, expect, it, vi } from "vitest";
import { deliverAcpxResponse, requireAcpxResponseDelivery } from "./acpx-response-delivery.js";
import { deliverAcpxResponse, requireAcpxResponseDelivery } from "./response-delivery.js";
describe("ACP interaction response delivery", () => {
it("does not acknowledge callback settlement before the provider pipe write", async () => {