fix(runner): await cancellation proof after clean event EOF

Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip committed 2026-09-29 22:25:41 -05:00
1 parent 8b6136bf66
commit e822b614fc
3 files changed
+58 -1

No files matched your search

@@ -1395,6 +1395,34 @@ describe("ACPX runtime host", () => {
);
});
it("keeps host shutdown independent of a concurrent managed turn cancellation", async () => {
const fixture = await hostFixture();
const raw = runtimeTurn(); // Native cancel acknowledges, but never resolves result.
const runtime = runtimePort({ startTurn: () => raw });
const host = await AcpxRuntimeHost.open(
{
...fixture.options,
agent: "codex",
model: "gpt-5.6-sol",
permissionMode: "approve-reads",
environment: { PAPERCLIP_ACPX_CODEX_AUTH_JSON_SECRET: "{}" },
},
fixture.dependencies({ openRuntime: async () => runtime }),
);
const managed = host.startTurn({ text: "Work", requestId: "turn-1" });
// The host retains raw.cancel, so its shutdown cannot wait on the managed
// cancellation which falls back to that same host's close promise.
const closed = host.close({ reason: "controller shutdown" });
const cancelled = managed.cancel({ reason: "operator stop" });
await expect(Promise.all([closed, cancelled])).resolves.toEqual([undefined, undefined]);
await expect(managed.result).resolves.toEqual({
status: "cancelled", stopReason: "cancelled_after_runtime_close",
});
expect(host.isClosed()).toBe(true);
expect(runtime.close).toHaveBeenCalledExactlyOnceWith({ reason: "controller shutdown" });
expect(fixture.commandClose).toHaveBeenCalledOnce();
});
it("clones ephemeral capabilities and fences steering controls to an acknowledged active turn", async () => {
const fixture = await hostFixture();
const turn = runtimeTurn();
@@ -50,6 +50,29 @@ it("never publishes cancellation when owned provider cleanup fails", async () =>
await Promise.all([result, next, stopped]);
});
it.each(["completed", "failed"] as const)("waits after clean EOF for %s owned cleanup", async outcome => {
vi.useFakeTimers(); const f = fixture();
const next = f.turn.events[Symbol.asyncIterator]().next();
const stopped = f.turn.cancel();
const nextOutcome = next.then(() => "ended", () => "failed");
let streamSettled = false;
void nextOutcome.then(() => { streamSettled = true; });
const result = outcome === "failed"
? expect(f.turn.result).rejects.toThrow("cleanup failed")
: expect(f.turn.result).resolves.toMatchObject({ status: "cancelled" });
const cancellation = outcome === "failed"
? expect(stopped).rejects.toThrow("cleanup failed")
: expect(stopped).resolves.toBeUndefined();
f.stream.resolve({ done: true, value: undefined });
await vi.advanceTimersByTimeAsync(20);
expect(f.close).toHaveBeenCalledOnce();
expect(streamSettled).toBe(false);
if (outcome === "failed") f.cleanup.reject(new Error("cleanup failed"));
else f.cleanup.resolve();
await Promise.all([result, cancellation]);
await expect(nextOutcome).resolves.toBe(outcome === "failed" ? "failed" : "ended");
});
it.each(["completed", "cancelled", "failed"] as const)("settles a racing provider %s terminal without waiting for the run timeout", async status => {
const f = fixture();
const stopped = f.turn.cancel();
@@ -76,7 +76,13 @@ export function withAcpxTurnCancellation(
await cancellation;
return { kind: "retired" as const };
});
if (next.kind === "retired" || next.value.done) return;
if (next.kind === "retired") return;
if (next.value.done) {
// Clean EOF is no stronger than a transport error: when Stop owns
// settlement, keep the stream open until terminal/cleanup proof.
if (cancellation) await cancellation;
return;
}
yield next.value.value;
}
} finally {