paperclip/packages/paperclip-runner/src/native-session-runtime.test.ts

8305 lines
257 KiB
TypeScript

import { describe, expect, it, vi } from "vitest";
import type { ControlPlanePort } from "./contracts/control-plane-port.js";
import type { NativeExecutionInputV1 } from "./contracts/native-execution.js";
import type { NativeRunIdentity } from "./contracts/types.js";
import type {
NativeSession,
NativeSessionBackend,
PersistedNativeSession,
} from "./contracts/native-session-backend.js";
import {
NativeSessionCloseUnrecoverableError,
NativeSessionCleanupQuarantinedError,
NativeSessionProtocolIntegrityError,
} from "./contracts/native-session-backend.js";
import type {
PrpEvent,
PrpStructuredRunResult,
PrpTerminalState,
} from "./protocol/replay-contract.js";
import {
NATIVE_RUNTIME_ASSET_SCHEMA,
PAPERCLIP_EXECUTION_PROMPT,
PAPERCLIP_EXECUTION_PROMPT_REVISION,
canonicalNativeRuntimeContextDigest,
nativeRuntimePromptDigest,
} from "./contracts/runtime-context.js";
import {
executeNativeSession,
completeTerminatedRemoteNativeSessionCleanup,
completeTerminatedLocalNativeSessionCleanup,
type ExecuteNativeSessionOptions,
} from "./native-session-runtime.js";
const identity = {
runId: "run-recovery",
sessionId: "session-recovery",
companyId: "company-recovery",
issueId: "issue-recovery",
agentId: "agent-recovery",
};
const result: PrpStructuredRunResult = {
schema: "paperclip.run_result.v1",
reportedWorkDisposition: "done",
summary: "Recovered native work completed.",
completionClaim: {
contractRevision: "1",
objectiveSatisfied: true,
criteria: [
{ criterionId: "objective", status: "satisfied", evidenceRefs: [] },
],
remainingWork: [],
},
evidence: [],
verification: [{ commandOrCheck: "recovery", status: "passed" }],
attentionRequests: [],
artifacts: [],
};
const terminal: PrpTerminalState = {
schema: "paperclip.prp.terminal.v1",
turnTerminalState: "completed",
runTerminalState: "succeeded",
reportedWorkDisposition: "done",
};
const yieldedResult: PrpStructuredRunResult = {
schema: "paperclip.run_result.v1",
reportedWorkDisposition: "yielded",
summary: "Waiting for the requested response.",
completionClaim: {
contractRevision: "1",
objectiveSatisfied: false,
criteria: [
{
criterionId: "objective",
status: "unknown",
evidenceRefs: ["interaction:pending"],
},
],
remainingWork: [
{ description: "Resume after the response.", blocksCompletion: true },
],
},
evidence: [{ ref: "interaction:pending" }],
verification: [],
attentionRequests: [],
artifacts: [{ kind: "issue_thread_interaction", ref: "interaction:pending" }],
continuation: {
kind: "response_wake",
summary: "Resume from the answer.",
idempotencyKey: "interaction-response:pending",
},
};
const input: NativeExecutionInputV1 = {
schema: "paperclip.native-execution-input.v1",
binding: {
companyId: identity.companyId,
runId: identity.runId,
issueId: identity.issueId,
agentId: identity.agentId,
executionWorkspaceId: "workspace-recovery",
},
task: {
identifier: "PAP-RECOVERY",
title: "Recover native work",
description: null,
prompt: "# PAP-RECOVERY: Recover native work",
workMode: "standard",
},
workspace: {
cwd: "/workspace",
repoUrl: null,
repoRef: null,
branchName: null,
},
session: {
normalizedSessionId: identity.sessionId,
driverKind: "codex_app_server",
protocolVersion: 1,
},
provider: { kind: "codex", model: null },
completionContract: {
id: "contract-recovery",
sha256: "contract-recovery-sha",
schemaVersion: "paperclip.completion-contract.v1",
contract: {
revision: "1",
objective: "Recover native work",
criteria: [{ id: "objective", requirement: "Complete after recovery" }],
},
},
interactionResponses: [],
credentialBindings: [],
};
function controlEvent(
sourceSeq: number,
eventType: PrpEvent["eventType"],
payload: Record<string, unknown>,
): PrpEvent {
return {
schema: "paperclip.prp.event.v1",
sourceEventId: `control-recovery:${identity.runId}:${sourceSeq}`,
sourceSeq,
sourceInstanceId: "control-recovery",
sourceKind: "control_plane",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-recovery",
eventType,
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:00.000Z",
payload,
};
}
function canonicalTestJson(value: unknown): string {
if (Array.isArray(value)) {
return `[${value.map(canonicalTestJson).join(",")}]`;
}
if (typeof value === "object" && value !== null) {
const record = value as Record<string, unknown>;
return `{${Object.keys(record)
.sort()
.map((key) => `${JSON.stringify(key)}:${canonicalTestJson(record[key])}`)
.join(",")}}`;
}
return JSON.stringify(value) ?? "undefined";
}
function runnerEvent(
sourceSeq: number,
eventType: PrpEvent["eventType"],
payload: Record<string, unknown> = {},
): PrpEvent {
return {
schema: "paperclip.prp.event.v1",
sourceEventId: `runner-recovery:${identity.runId}:${sourceSeq}`,
sourceSeq,
sourceInstanceId: "runner-recovery",
sourceKind: "runner",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-recovery",
eventType,
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:00.000Z",
payload,
};
}
function highestContiguous(events: PrpEvent[]): number {
const sequences = new Set(events.map((event) => event.sourceSeq));
let cursor = 0;
while (sequences.has(cursor + 1)) cursor += 1;
return cursor;
}
describe("executeNativeSession recovery", () => {
it.each([undefined, 0, 7 * 24 * 60 * 60 * 1000, 30 * 24 * 60 * 60 * 1000])(
"honors long-lived turn duration independently of operation bounds (%s)",
async (turnTimeoutMs) => {
vi.useFakeTimers();
let finish = () => {};
const waiting = new Promise<void>((resolve) => { finish = resolve; });
const capabilities = {
resume: true, typedEvents: true, steering: false,
interruption: true, structuredResult: true,
};
const cancel = vi.fn(async () => { finish(); });
const close = vi.fn(async () => { finish(); });
const session: NativeSession = {
identity: () => identity,
async capabilities() { return capabilities; },
async *events() {
yield runnerEvent(1, "tool.execution.started", { name: "wait for CI", status: "running" });
await waiting;
yield runnerEvent(2, "turn.completed");
},
async startTurn() { return { turnId: "turn-recovery" }; },
async result() { return { result, terminal, turnId: "turn-recovery" }; },
async snapshot() {
return { backendKind: "mock", sessionId: "driver-recovery", identity,
providerSessionId: "provider-recovery", cursor: null, activeTurnId: null,
pendingRuntimeRequests: [], lineage: [] };
},
cancel, close,
};
const appendEvent = vi.fn<ControlPlanePort["appendEvent"]>(async event => ({
cursor: event.sourceSeq, highestContiguousSourceSeq: event.sourceSeq,
disposition: "committed",
}));
const completeRun = vi.fn(async () => {});
try {
const execution = executeNativeSession({
input, turnTimeoutMs,
backend: {
async descriptor() { return { kind: "mock", name: "long-lived", version: "1", capabilities }; },
async openSession() { return session; },
},
controlPlane: {
async openRun() {}, async checkpointSession() {}, appendEvent,
async replayEvents() { return { events: [], highestContiguousSourceSeq: 0 }; },
completeRun,
},
runnerInstanceId: "runner-recovery", controlPlaneInstanceId: "control-recovery",
});
// Observe rejection before advancing the clock, including on the old implementation.
let failure: unknown;
const observed = execution.catch(error => { failure = error; return null; });
await vi.advanceTimersByTimeAsync(0);
expect(appendEvent).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(6 * 24 * 60 * 60 * 1000);
expect(failure).toBeUndefined();
expect(cancel).not.toHaveBeenCalled();
expect(close).not.toHaveBeenCalled();
expect(completeRun).not.toHaveBeenCalled();
if (turnTimeoutMs) {
// Includes a 30-day bound beyond Node's single-timer maximum.
await vi.advanceTimersByTimeAsync(turnTimeoutMs - 6 * 24 * 60 * 60 * 1000 - 1);
expect(failure).toBeUndefined();
await vi.advanceTimersByTimeAsync(1);
await observed;
expect(failure).toBeInstanceOf(Error);
expect((failure as Error).message).toContain(`native session timed out after ${turnTimeoutMs}ms`);
expect(cancel).toHaveBeenCalledOnce();
expect(completeRun).not.toHaveBeenCalled();
} else {
await vi.advanceTimersByTimeAsync(2 * 24 * 60 * 60 * 1000);
finish();
expect(await observed).toMatchObject({ result });
expect(completeRun).toHaveBeenCalledOnce();
expect(cancel).not.toHaveBeenCalled();
}
} finally {
finish();
vi.useRealTimers();
}
},
);
it.each([false, true])("preserves a durable session failure when its stream closes (throws=%s)", async throws => {
const capabilities = { resume: true, typedEvents: true, steering: false, interruption: true, structuredResult: true };
const session: NativeSession = {
identity: () => identity,
async capabilities() { return capabilities; },
async *events() {
yield runnerEvent(1, "session.failed", { error: { code: "notification_transport_failed" }, recoverable: false });
if (throws) throw new Error("provider stdout closed");
},
async startTurn() { return { turnId: "turn-recovery" }; },
async result() { return null; },
async snapshot() { return { backendKind: "mock", sessionId: "driver-recovery", identity, providerSessionId: "provider-recovery", cursor: null, activeTurnId: null, pendingRuntimeRequests: [], lineage: [] }; },
async close() {},
};
const completeRun = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() { return { kind: "mock", name: "failure-fixture", version: "1", capabilities }; },
async openSession() { return session; },
};
const port: ControlPlanePort = {
async openRun() {}, async checkpointSession() {},
async appendEvent() { return { cursor: 1, highestContiguousSourceSeq: 1, disposition: "committed" }; },
async replayEvents() { return { events: [], highestContiguousSourceSeq: 0 }; }, completeRun,
};
await expect(executeNativeSession({ input, backend, controlPlane: port, runnerInstanceId: "runner-recovery", controlPlaneInstanceId: "control-recovery" }))
.rejects.toMatchObject({ code: "native_provider_terminal_failed", providerCode: "notification_transport_failed", recoverable: false });
expect(completeRun).not.toHaveBeenCalled();
});
it.each((["complete", "paused", "blocked", "limited", "usageLimited", "budgetLimited"] as const)
.flatMap((status) => [false, true].map((snapshotBeforeUpdate) => ({ status, snapshotBeforeUpdate }))))(
"handles a new chat turn instead of completing it from an existing $status goal (snapshot: $snapshotBeforeUpdate)", async ({ status, snapshotBeforeUpdate }) => {
const oldGoal = {
threadId: "provider-recovery",
objective: "Say hello",
status,
tokenBudget: null,
tokensUsed: 500,
timeUsedSeconds: 2,
createdAt: Date.parse("2026-08-09T00:00:00.000Z"),
updatedAt: Date.parse("2026-08-09T00:00:02.000Z"),
};
const reply = { ...result, summary: "Said bye in response to the new message." };
const capabilities = {
resume: true, typedEvents: true, steering: false,
interruption: true, structuredResult: true,
};
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: oldGoal.threadId,
activeTurnId: null,
semanticResult: null,
terminal: null,
terminalTurns: [],
pendingRuntimeRequests: [],
goal: { ...oldGoal, createdAt: oldGoal.createdAt / 1000, updatedAt: oldGoal.updatedAt / 1000 },
};
const startTurn = vi.fn<NativeSession["startTurn"]>(async () => ({ turnId: "turn-recovery" }));
const goal = vi.fn(async () => oldGoal);
const session: NativeSession = {
identity: () => identity,
async capabilities() { return capabilities; },
async *events() {
let seq = 0;
// The resume snapshot is durable UI state, not work for this prompt.
if (snapshotBeforeUpdate) yield runnerEvent(++seq, "session.goal.snapshot", {
goal: {
...oldGoal,
createdAt: new Date(oldGoal.createdAt).toISOString(),
elapsedSeconds: oldGoal.timeUsedSeconds,
},
workingNow: false,
});
// Codex replays the unchanged goal as an update during resume too.
// Usage-only changes do not make an inactive goal own a new prompt.
yield runnerEvent(++seq, "session.goal.updated", {
goal: {
...oldGoal,
createdAt: new Date(oldGoal.createdAt).toISOString(),
tokensUsed: 600,
updatedAt: new Date(oldGoal.updatedAt + 1000).toISOString(),
},
workingNow: false,
});
yield runnerEvent(++seq, "turn.started");
yield runnerEvent(++seq, "run.result.proposed", reply);
yield runnerEvent(++seq, "turn.completed");
},
startTurn,
goal,
async result() { return { result: reply, terminal, turnId: "turn-recovery" }; },
async snapshot() { return checkpoint; },
async close() {},
};
const appended: PrpEvent[] = [];
const completed = await executeNativeSession({
input: { ...input, task: { ...input.task, prompt: "Say bye" } },
backend: {
async descriptor() {
return { kind: "mock", name: "chat-after-goal", version: "1", capabilities };
},
async openSession() { throw new Error("must resume the same provider session"); },
async recoverSession() { return { recovered: true, session }; },
},
persistedSession: checkpoint,
controlPlane: {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
appended.push(event);
return { cursor: event.sourceSeq, highestContiguousSourceSeq: event.sourceSeq, disposition: "committed" };
},
async replayEvents() { return { events: [], highestContiguousSourceSeq: 0 }; },
async completeRun() {},
},
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 1000,
});
expect(startTurn).toHaveBeenCalledOnce();
expect(JSON.parse(startTurn.mock.calls[0]![0]!.message.text).task.prompt).toBe("Say bye");
expect(goal).not.toHaveBeenCalled();
expect(completed.providerSessionId).toBe(oldGoal.threadId);
expect(completed.result).toEqual(reply);
expect(appended.map((event) => event.eventType)).toEqual([
...(snapshotBeforeUpdate ? ["session.goal.snapshot"] : []),
"session.goal.updated", "turn.started", "run.result.proposed", "turn.completed",
"run.result.accepted", "run.terminal",
]);
});
it("applies a session goal control without starting an ordinary turn", async () => {
const activeGoal = {
threadId: "provider-recovery",
objective: "Verify goal mode",
status: "active" as const,
tokenBudget: 12_000,
tokensUsed: 100,
timeUsedSeconds: 1,
createdAt: Date.parse("2026-08-09T00:00:00.000Z"),
updatedAt: Date.parse("2026-08-09T00:00:01.000Z"),
};
const completeGoal = {
...activeGoal,
status: "complete" as const,
tokensUsed: 500,
timeUsedSeconds: 2,
updatedAt: Date.parse("2026-08-09T00:00:02.000Z"),
};
const startTurn = vi.fn(async () => ({ turnId: "turn-recovery" }));
const goal = vi.fn(async (operation: Parameters<NonNullable<NativeSession["goal"]>>[0]) =>
operation.action === "get" ? completeGoal : activeGoal,
);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "session.goal.snapshot", {
goal: null,
workingNow: false,
});
yield runnerEvent(2, "session.goal.updated", {
requestId: "goal-create",
goal: {
objective: activeGoal.objective,
status: activeGoal.status,
tokenBudget: activeGoal.tokenBudget,
tokensUsed: activeGoal.tokensUsed,
elapsedSeconds: activeGoal.timeUsedSeconds,
},
workingNow: false,
});
yield runnerEvent(3, "turn.started");
yield runnerEvent(4, "run.result.proposed", result);
yield runnerEvent(5, "turn.completed");
yield runnerEvent(6, "session.goal.snapshot", {
goal: {
objective: completeGoal.objective,
status: completeGoal.status,
tokenBudget: completeGoal.tokenBudget,
tokensUsed: completeGoal.tokensUsed,
elapsedSeconds: completeGoal.timeUsedSeconds,
},
workingNow: false,
});
},
startTurn,
goal,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
goal: completeGoal,
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "goal-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const appended: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
appended.push(event);
return {
cursor: event.sourceSeq,
highestContiguousSourceSeq: event.sourceSeq,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const completed = await executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
sessionGoalControl: {
requestId: "goal-create",
action: "create",
objective: activeGoal.objective,
tokenBudget: activeGoal.tokenBudget,
},
});
expect(startTurn).not.toHaveBeenCalled();
expect(goal).toHaveBeenNthCalledWith(1, {
action: "set",
objective: activeGoal.objective,
status: "active",
requestId: "goal-create",
tokenBudget: activeGoal.tokenBudget,
});
expect(goal).toHaveBeenNthCalledWith(2, { action: "get" });
expect(appended.map((event) => event.eventType)).toContain("session.goal.snapshot");
expect(completed.result).toMatchObject({
reportedWorkDisposition: "done",
summary: result.summary,
completionClaim: { objectiveSatisfied: true },
});
});
it.each((["complete", "paused", "blocked", "limited", "usageLimited", "budgetLimited"] as const)
.flatMap((status) => [false, true].map((lateGoal) => ({ status, lateGoal }))))(
"reconciles an out-of-band $status goal (after semantic result: $lateGoal)", async ({ status, lateGoal }) => {
const activeGoal = {
threadId: "provider-agent-goal",
objective: "Finish autonomous work",
status: "active" as const,
tokenBudget: null,
tokensUsed: 100,
timeUsedSeconds: 1,
createdAt: Date.parse("2026-08-09T00:00:00.000Z"),
updatedAt: Date.parse("2026-08-09T00:00:01.000Z"),
};
const completeGoal = {
...activeGoal,
status,
tokensUsed: 500,
timeUsedSeconds: 2,
updatedAt: Date.parse("2026-08-09T00:00:02.000Z"),
};
const startTurn = vi.fn(async () => ({ turnId: "turn-agent-goal" }));
const goal = vi.fn(async () => completeGoal);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
if (lateGoal) yield runnerEvent(1, "run.result.proposed", result);
yield runnerEvent(lateGoal ? 2 : 1, "session.goal.updated", {
goal: {
objective: activeGoal.objective,
status: activeGoal.status,
tokenBudget: activeGoal.tokenBudget,
tokensUsed: activeGoal.tokensUsed,
elapsedSeconds: activeGoal.timeUsedSeconds,
},
workingNow: true,
});
if (!lateGoal) yield runnerEvent(2, "run.result.proposed", result);
// A newly observed goal owns the lifetime even when the completion
// proposal arrived first.
await new Promise((resolve) => setTimeout(resolve, 25));
yield runnerEvent(3, "session.goal.updated", {
goal: {
objective: completeGoal.objective,
status: completeGoal.status,
tokenBudget: completeGoal.tokenBudget,
tokensUsed: completeGoal.tokensUsed,
elapsedSeconds: completeGoal.timeUsedSeconds,
},
workingNow: true,
});
yield runnerEvent(4, "turn.completed");
},
startTurn,
goal,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-agent-goal",
identity,
providerSessionId: activeGoal.threadId,
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
goal: activeGoal,
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "agent-goal-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const appended: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
appended.push(event);
return {
cursor: event.sourceSeq,
highestContiguousSourceSeq: event.sourceSeq,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const completed = await executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-agent-goal",
controlPlaneInstanceId: "control-agent-goal",
});
expect(startTurn).toHaveBeenCalledOnce();
expect(goal).not.toHaveBeenCalled();
expect(appended.map((event) => event.eventType)).toEqual([
...(lateGoal ? ["run.result.proposed", "session.goal.updated"] : ["session.goal.updated", "run.result.proposed"]),
"session.goal.updated",
"turn.completed",
"run.result.accepted",
"run.terminal",
]);
expect(completed.result).toMatchObject({
reportedWorkDisposition: status === "complete" ? "done" : status === "blocked" ? "blocked" : "yielded",
...(status === "complete" ? { summary: result.summary } : {}),
completionClaim: { objectiveSatisfied: status === "complete" },
});
});
it.each([
{ keepSessionOpen: false, recoveryOnly: false },
{ keepSessionOpen: true, recoveryOnly: false },
{ keepSessionOpen: false, recoveryOnly: true },
{ keepSessionOpen: true, recoveryOnly: true },
].flatMap((options) => [false, true].map((cleared) => ({ ...options, cleared }))))("reconciles goal recovery without replaying completed controls (warm=$keepSessionOpen, recovery=$recoveryOnly, cleared=$cleared)", async ({ keepSessionOpen, recoveryOnly, cleared }) => {
const pausedGoal = {
threadId: "provider-recovery",
objective: "Verify recovered goal control",
status: recoveryOnly ? "complete" as const : "paused" as const,
tokenBudget: null,
tokensUsed: 250,
timeUsedSeconds: 2,
createdAt: Date.parse("2026-08-09T00:00:00.000Z"),
updatedAt: Date.parse("2026-08-09T00:00:02.000Z"),
};
const startTurn = vi.fn(async () => ({ turnId: "turn-recovery" }));
const goal = vi.fn(async () => cleared ? null : pausedGoal);
const requestId = recoveryOnly ? `recovery_${input.binding.runId}` : cleared ? "goal-clear" : "goal-pause";
const close = vi.fn(async () => {});
let snapshotCount = 0;
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, cleared ? "session.goal.cleared" : "session.goal.updated", {
requestId,
goal: cleared ? null : {
objective: pausedGoal.objective,
status: pausedGoal.status,
tokenBudget: pausedGoal.tokenBudget,
tokensUsed: pausedGoal.tokensUsed,
elapsedSeconds: pausedGoal.timeUsedSeconds,
},
workingNow: !recoveryOnly,
});
yield runnerEvent(2, "turn.completed");
yield runnerEvent(3, "session.goal.snapshot", {
goal: cleared ? null : {
objective: pausedGoal.objective,
status: pausedGoal.status,
tokenBudget: pausedGoal.tokenBudget,
tokensUsed: pausedGoal.tokensUsed,
elapsedSeconds: pausedGoal.timeUsedSeconds,
},
workingNow: false,
});
},
startTurn,
goal,
async result() {
return null;
},
async snapshot() {
snapshotCount += 1;
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: !recoveryOnly && snapshotCount === 1 ? "turn-recovery" : null,
pendingRuntimeRequests: [],
goal: cleared ? null : pausedGoal,
lineage: [],
};
},
close,
};
const persistedSession: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: "turn-recovery",
pendingRuntimeRequests: [],
goal: { ...pausedGoal, status: "active" },
lineage: [],
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "goal-recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("fresh session must not be opened");
},
async recoverSession() {
return { recovered: true, session };
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
return {
cursor: event.sourceSeq,
highestContiguousSourceSeq: event.sourceSeq,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const completed = await executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
persistedSession,
keepSessionOpen,
requireSessionCloseBeforeReturn: true,
resumeSessionGoalHeartbeat: recoveryOnly,
sessionGoalControl: recoveryOnly ? null : {
requestId,
action: cleared ? "clear" : "pause",
},
});
expect(startTurn).not.toHaveBeenCalled();
expect(goal).toHaveBeenNthCalledWith(1, {
action: recoveryOnly ? "get" : cleared ? "clear" : "pause",
requestId,
});
expect(goal).toHaveBeenCalledTimes(cleared && !recoveryOnly ? 2 : 1);
if (cleared && !recoveryOnly) {
expect(goal).toHaveBeenNthCalledWith(2, { action: "get" });
}
expect(close).toHaveBeenCalledTimes(1);
expect(completed.result).toMatchObject({
reportedWorkDisposition: recoveryOnly && !cleared ? "done" : "yielded",
completionClaim: { objectiveSatisfied: recoveryOnly && !cleared },
});
});
it.each([
{
recoverable: false,
message:
"There's an issue with the selected model (custom-model). It may not exist or you may not have access to it.",
modelRejected: true,
},
{
recoverable: true,
message:
"There's an issue with the selected model (custom-model). It may not exist or you may not have access to it.",
modelRejected: false,
},
{
recoverable: false,
message: "The model service failed while processing output.",
modelRejected: false,
},
])(
"preserves structured provider failure and model retry classification ($recoverable, $modelRejected)",
async ({ recoverable, message, modelRejected }) => {
const capabilities = {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
const close = vi.fn(async () => {});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return capabilities;
},
async *events() {
yield runnerEvent(1, "turn.failed", {
error: { code: "RUNTIME", recoverable, message },
});
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "model-rejection",
version: "1",
capabilities,
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const result = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await expect(result).rejects.toThrow(message);
await expect(result).rejects.toMatchObject({
code: "native_provider_terminal_failed",
providerCode: "RUNTIME",
recoverable,
});
if (modelRejected) {
await expect(result).rejects.toThrow("native_provider_model_rejected:");
} else {
await expect(result).rejects.not.toThrow(
"native_provider_model_rejected",
);
}
expect(close).toHaveBeenCalled();
},
);
it("keeps governed-wait discovery synchronous", () => {
type GovernedWaitResolver = NonNullable<
ExecuteNativeSessionOptions["resolveGovernedWait"]
>;
const resolver: GovernedWaitResolver = () => null;
// An async resolver could retain control-plane mutation authority after
// execution settles, so the public boundary rejects it at compile time.
// @ts-expect-error governed-wait discovery must not return a promise
const asynchronousResolver: GovernedWaitResolver = async () => null;
expect(resolver).toBeTypeOf("function");
void asynchronousResolver;
});
it("preserves durable success while quarantined cleanup stays bounded", async () => {
vi.useFakeTimers();
try {
let executionCloseCount = 0;
let quarantineAttempt = 0;
const close = vi.fn(({ reason }: { reason: string }) => {
if (reason === "native session quarantined cleanup recovery") {
quarantineAttempt += 1;
const attempt = quarantineAttempt;
return new Promise<void>((resolve, reject) => {
setTimeout(() => {
if (attempt < 3)
reject(new Error("transient quarantine failure"));
else resolve();
}, 6_500);
});
}
if (reason === "native session execution complete") {
executionCloseCount += 1;
if (executionCloseCount > 1) return Promise.resolve();
}
return Promise.reject(new Error("persistent close failure"));
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await expect(execute()).resolves.toMatchObject({ result });
expect(close).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(3_000);
expect(close).toHaveBeenCalledTimes(5);
// Admission inherits the already-running three-attempt recovery and
// waits through its bounded attempts instead of timing out after one.
const recoveredExecution = execute();
let admissionSettled = false;
void recoveredExecution.then(
() => {
admissionSettled = true;
},
() => {
admissionSettled = true;
},
);
await vi.advanceTimersByTimeAsync(20_000);
expect(admissionSettled).toBe(false);
expect(openSession).toHaveBeenCalledTimes(1);
expect(close).toHaveBeenCalledTimes(7);
await vi.advanceTimersByTimeAsync(2_000);
await expect(recoveredExecution).resolves.toMatchObject({ result });
expect(quarantineAttempt).toBe(3);
expect(close).toHaveBeenCalledTimes(8);
expect(openSession).toHaveBeenCalledTimes(2);
} finally {
vi.useRealTimers();
}
});
it("awaits the quarantine owner that replaces an exhausted close recovery", async () => {
vi.useFakeTimers();
try {
let executionCloseCount = 0;
const close = vi.fn(({ reason }: { reason: string }) => {
if (reason === "native session execution complete") {
executionCloseCount += 1;
return executionCloseCount === 1
? Promise.reject(new Error("initial close failed"))
: Promise.resolve();
}
if (
reason.startsWith(
"native session cleanup recovery after close failure",
)
) {
return Promise.reject(new Error("bounded close recovery failed"));
}
if (reason === "native session quarantined cleanup recovery") {
return new Promise<void>((resolve) => setTimeout(resolve, 50));
}
return Promise.resolve();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await expect(execute()).resolves.toMatchObject({ result });
const admitted = execute();
await vi.advanceTimersByTimeAsync(3_100);
await expect(admitted).resolves.toMatchObject({ result });
expect(openSession).toHaveBeenCalledTimes(2);
expect(close).toHaveBeenCalledWith({
reason: "native session quarantined cleanup recovery",
});
} finally {
vi.useRealTimers();
}
});
it("covers maximum-duration retained and replacement recovery phases", async () => {
vi.useFakeTimers();
try {
let executionCloseCount = 0;
let retainedRecoveryAttempt = 0;
let quarantineRecoveryAttempt = 0;
const close = vi.fn(({ reason }: { reason: string }) => {
if (reason === "native session execution complete") {
executionCloseCount += 1;
if (executionCloseCount > 1) return Promise.resolve();
return new Promise<void>((_resolve, reject) => {
setTimeout(() => reject(new Error("initial close failed")), 6_900);
});
}
if (
reason.startsWith(
"native session cleanup recovery after close failure",
)
) {
retainedRecoveryAttempt += 1;
return new Promise<void>((_resolve, reject) => {
setTimeout(
() => reject(new Error("bounded retained recovery failed")),
6_900,
);
});
}
if (reason === "native session quarantined cleanup recovery") {
quarantineRecoveryAttempt += 1;
const attempt = quarantineRecoveryAttempt;
return new Promise<void>((resolve, reject) => {
setTimeout(() => {
if (attempt < 3)
reject(new Error("transient quarantine recovery failure"));
else resolve();
}, 6_500);
});
}
return Promise.resolve();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
const firstExecution = execute();
await vi.advanceTimersByTimeAsync(100);
await expect(firstExecution).resolves.toMatchObject({ result });
const admitted = execute();
let admissionSettled = false;
void admitted.then(
() => {
admissionSettled = true;
},
() => {
admissionSettled = true;
},
);
// The initial close plus three near-bound retries exceed the old 23s
// admission grace but remain within the configured 31s owner phase.
await vi.advanceTimersByTimeAsync(23_100);
expect(retainedRecoveryAttempt).toBe(2);
expect(quarantineRecoveryAttempt).toBe(0);
expect(admissionSettled).toBe(false);
expect(openSession).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(7_600);
expect(retainedRecoveryAttempt).toBe(3);
expect(quarantineRecoveryAttempt).toBe(1);
expect(admissionSettled).toBe(false);
expect(openSession).toHaveBeenCalledOnce();
// The replacement then receives its complete three-attempt bound rather
// than inheriting only the remainder of the retained owner's deadline.
await vi.advanceTimersByTimeAsync(21_600);
await expect(admitted).resolves.toMatchObject({ result });
expect(quarantineRecoveryAttempt).toBe(3);
expect(openSession).toHaveBeenCalledTimes(2);
} finally {
vi.useRealTimers();
}
});
it("runs a full admission batch after an inherited scheduled attempt fails", async () => {
vi.useFakeTimers();
try {
let executionCloseCount = 0;
let scheduledRecoveryAttempt = 0;
let admissionRecoveryAttempt = 0;
const close = vi.fn(({ reason }: { reason: string }) => {
if (reason === "native session execution complete") {
executionCloseCount += 1;
return executionCloseCount === 1
? Promise.reject(new Error("initial close failed"))
: Promise.resolve();
}
if (
reason.startsWith(
"native session cleanup recovery after close failure",
)
) {
return Promise.reject(new Error("bounded retained recovery failed"));
}
if (reason === "native session quarantined cleanup recovery") {
return Promise.reject(
new Error("bounded quarantine recovery failed"),
);
}
if (
reason === "native session scheduled quarantined cleanup recovery"
) {
scheduledRecoveryAttempt += 1;
return new Promise<void>((_resolve, reject) => {
setTimeout(
() => reject(new Error("scheduled cleanup failed")),
6_500,
);
});
}
if (reason === "native session quarantined admission recovery") {
admissionRecoveryAttempt += 1;
const attempt = admissionRecoveryAttempt;
return new Promise<void>((resolve, reject) => {
setTimeout(() => {
if (attempt === 1)
reject(new Error("transient admission cleanup failure"));
else resolve();
}, 6_500);
});
}
return Promise.resolve();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await expect(execute()).resolves.toMatchObject({ result });
// Exhaust retained and initial quarantine batches, then enter the slow
// autonomous one-attempt recovery scheduled sixty seconds later.
await vi.advanceTimersByTimeAsync(65_001);
expect(scheduledRecoveryAttempt).toBe(1);
const admitted = execute();
let admissionSettled = false;
void admitted.then(
() => {
admissionSettled = true;
},
() => {
admissionSettled = true;
},
);
await vi.advanceTimersByTimeAsync(6_600);
expect(admissionRecoveryAttempt).toBe(1);
expect(admissionSettled).toBe(false);
expect(openSession).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(14_100);
await expect(admitted).resolves.toMatchObject({ result });
expect(admissionRecoveryAttempt).toBe(2);
expect(openSession).toHaveBeenCalledTimes(2);
} finally {
vi.useRealTimers();
}
});
it("stops scheduled cleanup after the quarantine lifetime budget", async () => {
vi.useFakeTimers();
try {
let executionCloseCount = 0;
let admissionAttempt = 0;
let scheduledAttempt = 0;
const lifetimeIdentity = {
...identity,
companyId: "company-cleanup-lifetime",
};
const close = vi.fn(({ reason }: { reason: string }) => {
if (reason === "native session execution complete") {
executionCloseCount += 1;
return executionCloseCount === 1
? Promise.reject(new Error("initial close failed"))
: Promise.resolve();
}
if (
reason.startsWith(
"native session cleanup recovery after close failure",
)
) {
return Promise.reject(new Error("bounded close recovery failed"));
}
if (reason === "native session quarantined cleanup recovery") {
return Promise.reject(new Error("quarantined close recovery failed"));
}
if (reason === "native session quarantined admission recovery") {
admissionAttempt += 1;
return Promise.reject(new Error("admission cleanup failed"));
}
if (
reason === "native session scheduled quarantined cleanup recovery"
) {
scheduledAttempt += 1;
return Promise.reject(new Error("scheduled cleanup failed"));
}
return Promise.resolve();
});
const session: NativeSession = {
identity: () => lifetimeIdentity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity: lifetimeIdentity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input: {
...input,
binding: {
...input.binding,
companyId: lifetimeIdentity.companyId,
},
},
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await expect(execute()).resolves.toMatchObject({ result });
await vi.advanceTimersByTimeAsync(6_000);
const recoveredExecution = expect(execute()).rejects.toThrow(
"prior session cleanup remains incomplete",
);
await vi.advanceTimersByTimeAsync(2_100);
await recoveredExecution;
expect(admissionAttempt).toBe(3);
expect(scheduledAttempt).toBe(0);
await vi.advanceTimersByTimeAsync(180_000);
expect(scheduledAttempt).toBe(3);
await vi.advanceTimersByTimeAsync(300_000);
expect(scheduledAttempt).toBe(3);
expect(openSession).toHaveBeenCalledOnce();
} finally {
vi.useRealTimers();
}
});
it("fails admission within a bound when quarantined close never settles", async () => {
vi.useFakeTimers();
try {
let releaseBlockedClose = () => {};
let blockClose = false;
const close = vi.fn(() => {
if (blockClose) {
return new Promise<void>((resolve) => {
releaseBlockedClose = resolve;
});
}
return Promise.reject(new Error("persistent close failure"));
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await expect(execute()).resolves.toMatchObject({ result });
await vi.advanceTimersByTimeAsync(6_000);
blockClose = true;
const blockedAdmission = execute();
const blockedResult = expect(blockedAdmission).rejects.toThrow(
"prior session cleanup exceeded the admission grace",
);
await vi.advanceTimersByTimeAsync(31_100);
await blockedResult;
expect(openSession).toHaveBeenCalledOnce();
const isolatedSession: NativeSession = {
...session,
identity: () => ({ ...identity, companyId: "company-isolated" }),
close: vi.fn(async () => undefined),
};
const isolatedOpenSession = vi.fn(async () => isolatedSession);
await expect(
executeNativeSession({
input: {
...input,
binding: {
...input.binding,
companyId: "company-isolated",
},
},
backend: { ...backend, openSession: isolatedOpenSession },
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({ result });
expect(isolatedOpenSession).toHaveBeenCalledOnce();
blockClose = false;
releaseBlockedClose();
close.mockResolvedValue(undefined);
await vi.runAllTimersAsync();
await expect(execute()).resolves.toMatchObject({ result });
} finally {
vi.useRealTimers();
}
});
it("quarantines a pending first close before another provider session can open", async () => {
vi.useFakeTimers();
try {
let releaseClose = () => {};
const pendingClose = new Promise<void>((resolve) => {
releaseClose = resolve;
});
const close = vi.fn(() => pendingClose);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const openSession = vi.fn(async () => session);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
const firstExecution = execute();
await vi.advanceTimersByTimeAsync(100);
await expect(firstExecution).resolves.toMatchObject({ result });
expect(close).toHaveBeenCalledOnce();
const blockedAdmission = execute();
const blockedResult = expect(blockedAdmission).rejects.toThrow(
"prior session cleanup exceeded the admission grace",
);
await vi.advanceTimersByTimeAsync(31_100);
await blockedResult;
expect(openSession).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
releaseClose();
await vi.runAllTimersAsync();
} finally {
vi.useRealTimers();
}
});
it("fails closed before launch when a v3 driver does not declare complete native context realization", async () => {
const digest = "0".repeat(64);
const context = {
prompt: {
revision: PAPERCLIP_EXECUTION_PROMPT_REVISION,
text: PAPERCLIP_EXECUTION_PROMPT,
digest: nativeRuntimePromptDigest(),
},
instructions: {
entryPath: "AGENTS.md",
bundle: {
schema: NATIVE_RUNTIME_ASSET_SCHEMA,
digest,
manifestDigest: digest,
rootPath: "/paperclip/context/instructions",
fileCount: 1,
totalBytes: 1,
},
},
skills: [],
mcp: { assignmentSetId: "none", digest, bindingId: null },
} as const;
const openSession = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "future-provider",
name: "future-provider",
version: "1",
capabilities: {
resume: false,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input: {
...input,
schema: "paperclip.native-execution-input.v3",
executionMode: "default",
planningContext: null,
runtimeContext: {
...context,
aggregateDigest: canonicalNativeRuntimeContextDigest(context),
},
},
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).rejects.toThrow("does not natively realize instructions, skills, mcp");
expect(openSession).not.toHaveBeenCalled();
});
it("does not admit a fresh run when provider session initialization fails", async () => {
const providerFailure = new Error("provider initialization failed");
const openSession = vi.fn(async () => {
throw providerFailure;
});
const openRun = vi.fn(async () => undefined);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "fresh-backend",
version: "1",
capabilities: {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).rejects.toBe(providerFailure);
expect(openSession).toHaveBeenCalledWith(
expect.objectContaining({
identity,
workingDirectory: input.workspace.cwd,
signal: expect.any(AbortSignal),
}),
);
expect(openRun).not.toHaveBeenCalled();
});
it("retains late bootstrap cleanup through the next admission", async () => {
vi.useFakeTimers();
let releaseClose = () => {};
try {
let resolveBootstrap = (_value: NativeSession) => {};
const stalledBootstrap = new Promise<NativeSession>((resolve) => {
resolveBootstrap = resolve;
});
let markBootstrapStarted = () => {};
const bootstrapStarted = new Promise<void>((resolve) => {
markBootstrapStarted = resolve;
});
let markCloseStarted = () => {};
const closeStarted = new Promise<void>((resolve) => {
markCloseStarted = resolve;
});
const closeReleased = new Promise<void>((resolve) => {
releaseClose = resolve;
});
const close = vi.fn(async () => {
markCloseStarted();
await closeReleased;
});
const lateSession: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
};
},
async *events() {},
async startTurn() {
throw new Error("late fresh session must not start");
},
async result() {
return null;
},
async snapshot() {
throw new Error("late fresh session must not snapshot");
},
close,
};
let bootstrapSignal: AbortSignal | undefined;
let bootstrapCount = 0;
const openSession = vi.fn(
(bootstrapInput: {
identity: NativeRunIdentity;
workingDirectory?: string;
signal?: AbortSignal;
}) => {
bootstrapCount += 1;
if (bootstrapCount > 1) {
throw new Error("replacement bootstrap launched");
}
bootstrapSignal = bootstrapInput.signal;
bootstrapInput.signal?.addEventListener(
"abort",
() => resolveBootstrap(lateSession),
{ once: true },
);
markBootstrapStarted();
return stalledBootstrap;
},
);
const openRun = vi.fn(async () => undefined);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "fresh-stalled-backend",
version: "1",
capabilities: await lateSession.capabilities(),
};
},
openSession,
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 5,
});
const execution = execute();
const rejection = expect(execution).rejects.toThrow(
"native session bootstrap timed out after 5ms",
);
await bootstrapStarted;
expect(bootstrapSignal?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(5);
await rejection;
expect(bootstrapSignal?.aborted).toBe(true);
expect(openRun).not.toHaveBeenCalled();
await closeStarted;
expect(close).toHaveBeenCalledWith({
reason: "native session bootstrap timed out",
});
// Even after the detached disposer exhausts its short settlement grace,
// the exact close remains admission-visible until it releases provider
// resources. A replacement bootstrap cannot start concurrently.
await vi.advanceTimersByTimeAsync(101);
const blockedAdmission = execute();
const blockedAdmissionRejection = expect(
blockedAdmission,
).rejects.toThrow("replacement bootstrap launched");
await vi.advanceTimersByTimeAsync(1_000);
expect(openSession).toHaveBeenCalledOnce();
releaseClose();
await vi.advanceTimersByTimeAsync(0);
await blockedAdmissionRejection;
expect(openSession).toHaveBeenCalledTimes(2);
expect(openRun).not.toHaveBeenCalled();
} finally {
releaseClose();
vi.useRealTimers();
}
});
it("waits for controller ownership publication before dispatching a turn", async () => {
let release!: () => void;
const published = new Promise<void>((resolve) => { release = resolve; });
const snapshotFailure = new Error("stop after ownership publication");
const snapshot = vi.fn(async () => { throw snapshotFailure; });
const session: NativeSession = {
identity: () => identity,
async capabilities() { return { resume: false, typedEvents: true, steering: false, interruption: true }; },
async *events() {},
async startTurn() { throw new Error("unexpected turn"); },
async result() { return null; },
snapshot,
close: vi.fn(async () => undefined),
};
const onSession = vi.fn(async (current: NativeSession | null) => { if (current) await published; });
const running = executeNativeSession({
input,
backend: { async descriptor() { return { kind: "mock", name: "owner-barrier", version: "1", capabilities: await session.capabilities() }; }, async openSession() { return session; } },
controlPlane: { async openRun() {}, async checkpointSession() {}, async appendEvent() { throw new Error("unexpected event"); }, async replayEvents() { return { events: [], highestContiguousSourceSeq: 0 }; }, async completeRun() {} },
runnerInstanceId: "runner-recovery", controlPlaneInstanceId: "control-recovery", onSession,
});
const rejected = expect(running).rejects.toBe(snapshotFailure);
await vi.waitFor(() => expect(onSession).toHaveBeenCalledWith(session));
expect(snapshot).not.toHaveBeenCalled();
release();
await rejected;
});
it("closes the provider when owner quarantine notification throws", async () => {
const snapshotFailure = new Error("snapshot failed");
const close = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
};
},
async *events() {},
async startTurn() {
throw new Error("unexpected turn");
},
async result() {
return null;
},
async snapshot() {
throw snapshotFailure;
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "owner-notification-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const retainedSessions: Array<NativeSession | null> = [];
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
keepSessionOpen: true,
onSession(current) {
retainedSessions.push(current);
if (current === null) throw new Error("owner notification failed");
},
}),
).rejects.toBe(snapshotFailure);
expect(retainedSessions).toEqual([session, null]);
expect(close).toHaveBeenCalledOnce();
});
it.each(["control-plane checkpoint", "owner checkpoint"] as const)(
"aborts consumption and closes the provider when the startup %s never settles",
async (stalledBoundary) => {
vi.useFakeTimers();
let releaseStream = () => {};
try {
const streamReleased = new Promise<void>((resolve) => {
releaseStream = resolve;
});
const never = new Promise<void>(() => undefined);
let markCheckpointStalled = () => {};
const checkpointStalled = new Promise<void>((resolve) => {
markCheckpointStalled = resolve;
});
let checkpointSignal: AbortSignal | undefined;
let controlPlaneCheckpointCount = 0;
const checkpointSession: NonNullable<
ControlPlanePort["checkpointSession"]
> = async (_snapshot, checkpointOptions) => {
controlPlaneCheckpointCount += 1;
if (
stalledBoundary === "control-plane checkpoint" &&
controlPlaneCheckpointCount === 2
) {
checkpointSignal = checkpointOptions?.signal;
markCheckpointStalled();
await never;
}
};
let ownerCheckpointCount = 0;
const onCheckpoint: NonNullable<
ExecuteNativeSessionOptions["onCheckpoint"]
> = async (_snapshot, checkpointOptions) => {
ownerCheckpointCount += 1;
if (
stalledBoundary === "owner checkpoint" &&
ownerCheckpointCount === 2
) {
checkpointSignal = checkpointOptions?.signal;
markCheckpointStalled();
await never;
}
};
const startTurn = vi.fn(async () => ({
turnId: "turn-checkpoint-stalled",
}));
const close = vi.fn(async () => {
releaseStream();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
},
async *events() {
await streamReleased;
},
startTurn,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-checkpoint-stalled",
cursor: "0",
activeTurnId: "turn-checkpoint-stalled",
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "checkpoint-stalled-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const openRun = vi.fn(async () => undefined);
const port: ControlPlanePort = {
openRun,
checkpointSession,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const retainedSessions: Array<NativeSession | null> = [];
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 1_000,
checkpointTimeoutMs: 1,
keepSessionOpen: true,
onCheckpoint,
onSession: (current) => retainedSessions.push(current),
});
const rejection = expect(execution).rejects.toThrow(
"native session checkpoint timed out after 1ms",
);
await checkpointStalled;
expect(checkpointSignal?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(1);
await rejection;
expect(checkpointSignal?.aborted).toBe(true);
expect(openRun).toHaveBeenCalledOnce();
expect(startTurn).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
expect(retainedSessions).toEqual([session, null]);
} finally {
releaseStream();
vi.useRealTimers();
}
},
);
it.each([
"provider result",
"completion checkpoint",
"control-plane replay",
"final event append",
"run completion",
] as const)(
"bounds post-terminal finalization when %s never settles",
async (stalledBoundary) => {
vi.useFakeTimers();
try {
const never = new Promise<never>(() => undefined);
let markFinalizationStalled = () => {};
const finalizationStalled = new Promise<void>((resolve) => {
markFinalizationStalled = resolve;
});
let stalledSignal: AbortSignal | undefined;
let resultCalls = 0;
let resultResolved = false;
let completionCheckpointCalls = 0;
const close = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
resultCalls += 1;
if (stalledBoundary === "provider result") {
markFinalizationStalled();
return await never;
}
resultResolved = true;
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "finalization-timeout-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession(_snapshot, operationOptions) {
if (stalledBoundary === "completion checkpoint" && resultResolved) {
completionCheckpointCalls += 1;
stalledSignal = operationOptions?.signal;
markFinalizationStalled();
await never;
}
},
async appendEvent(event, operationOptions) {
if (
stalledBoundary === "final event append" &&
(event as PrpEvent).sourceKind === "control_plane"
) {
stalledSignal = operationOptions?.signal;
markFinalizationStalled();
return await never;
}
return {
cursor: 1,
highestContiguousSourceSeq: (event as PrpEvent).sourceSeq,
disposition: "committed",
};
},
async replayEvents(_replay, operationOptions) {
if (stalledBoundary === "control-plane replay") {
stalledSignal = operationOptions?.signal;
markFinalizationStalled();
return await never;
}
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun(_completion, operationOptions) {
if (stalledBoundary === "run completion") {
stalledSignal = operationOptions?.signal;
markFinalizationStalled();
await never;
}
},
};
const retainedSessions: Array<NativeSession | null> = [];
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 10,
keepSessionOpen: true,
onSession: (current) => retainedSessions.push(current),
});
const rejection = expect(execution).rejects.toThrow(
"native session finalization timed out after 10ms",
);
await finalizationStalled;
if (stalledBoundary !== "provider result") {
expect(stalledSignal?.aborted).toBe(false);
}
await vi.advanceTimersByTimeAsync(20);
await rejection;
if (stalledBoundary !== "provider result") {
expect(stalledSignal?.aborted).toBe(true);
}
expect(resultCalls).toBe(1);
if (stalledBoundary === "completion checkpoint") {
expect(completionCheckpointCalls).toBe(1);
}
expect(close).toHaveBeenCalledOnce();
expect(retainedSessions).toEqual([session, null]);
} finally {
vi.useRealTimers();
}
},
);
it.each(["final event append", "run completion"] as const)(
"confirms durable completion when %s commits before its acknowledgement stalls",
async (stalledBoundary) => {
vi.useFakeTimers();
try {
const never = new Promise<never>(() => undefined);
let markFinalizationStalled = () => {};
const finalizationStalled = new Promise<void>((resolve) => {
markFinalizationStalled = resolve;
});
let stalledOnce = false;
let stalledSignal: AbortSignal | undefined;
const events: PrpEvent[] = [];
let durableCompletion: unknown = null;
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "durable-finalization-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event, operationOptions) {
const appended = structuredClone(event as PrpEvent);
const existing = events.find(
(candidate) =>
candidate.sourceInstanceId === appended.sourceInstanceId &&
candidate.sourceSeq === appended.sourceSeq,
);
if (existing === undefined) events.push(appended);
if (
stalledBoundary === "final event append" &&
appended.sourceKind === "control_plane" &&
!stalledOnce
) {
stalledOnce = true;
stalledSignal = operationOptions?.signal;
markFinalizationStalled();
return await never;
}
const sourceEvents = events.filter(
(candidate) =>
candidate.sourceInstanceId === appended.sourceInstanceId,
);
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(sourceEvents),
disposition: existing === undefined ? "committed" : "duplicate",
};
},
async replayEvents(replay) {
const sourceEvents = events.filter(
(event) => event.sourceInstanceId === replay.sourceInstanceId,
);
return {
events: structuredClone(
sourceEvents.filter(
(event) => event.sourceSeq > replay.afterSourceSeq,
),
),
highestContiguousSourceSeq: highestContiguous(sourceEvents),
};
},
async completeRun(completion, operationOptions) {
if (durableCompletion === null) {
durableCompletion = structuredClone(completion);
} else {
expect(completion).toEqual(durableCompletion);
}
if (stalledBoundary === "run completion" && !stalledOnce) {
stalledOnce = true;
stalledSignal = operationOptions?.signal;
markFinalizationStalled();
await never;
}
},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 10,
});
await finalizationStalled;
expect(stalledSignal?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(10);
await expect(execution).resolves.toMatchObject({
result,
terminal,
nativeEventCount: 3,
});
expect(stalledSignal?.aborted).toBe(true);
expect(durableCompletion).toMatchObject({ result, terminal });
expect(
events.filter((event) => event.sourceKind === "control_plane"),
).toHaveLength(2);
} finally {
vi.useRealTimers();
}
},
);
it.each([
{ boundary: "control-plane replay", fault: "typed", admitted: false },
{ boundary: "final event append", fault: "typed", admitted: false },
{
boundary: "lost final event acknowledgement",
fault: "typed",
admitted: false,
},
{ boundary: "control-plane replay", fault: "lookalike", admitted: true },
{ boundary: "control-plane replay", fault: "generic", admitted: true },
{
boundary: "completion invoked before commit",
fault: "typed",
admitted: true,
},
{
boundary: "lost completion acknowledgement",
fault: "typed",
admitted: true,
},
] as const)(
"observes integrity before completion admission without revoking admitted completion (%j)",
async ({ boundary, fault, admitted }) => {
vi.useFakeTimers();
let releaseBoundary = () => {};
const boundaryReleased = new Promise<void>((resolve) => {
releaseBoundary = resolve;
});
let markBoundaryReached = () => {};
const boundaryReached = new Promise<void>((resolve) => {
markBoundaryReached = resolve;
});
const failure =
fault === "typed"
? new NativeSessionProtocolIntegrityError(
"semantic_input_digest_mismatch",
)
: fault === "lookalike"
? Object.assign(new Error("untrusted transport error"), {
code: "native_event_replay_conflict",
reason: "semantic_input_digest_mismatch",
})
: new Error("snapshot temporarily unavailable");
let latchedFailure: Error | null = null;
let boundaryBlocked = false;
const waitAtBoundary = async () => {
if (boundaryBlocked) return;
boundaryBlocked = true;
markBoundaryReached();
await boundaryReleased;
};
const events: PrpEvent[] = [];
let durableCompletion: unknown = null;
const close = vi.fn(async () => undefined);
const resolveResult = vi.fn(async () => ({
result,
terminal,
turnId: "turn-recovery",
}));
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
result: resolveResult,
async snapshot() {
if (latchedFailure !== null) throw latchedFailure;
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: `completion-integrity-${boundary}-${fault}`,
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const completeRun = vi.fn<ControlPlanePort["completeRun"]>(
async (completion) => {
if (boundary === "completion invoked before commit")
await waitAtBoundary();
if (durableCompletion === null)
durableCompletion = structuredClone(completion);
else expect(completion).toEqual(durableCompletion);
if (
boundary === "lost completion acknowledgement" &&
!boundaryBlocked
) {
await waitAtBoundary();
await new Promise<never>(() => undefined);
}
},
);
const checkpointSession = vi.fn<ControlPlanePort["checkpointSession"]>(
async () => undefined,
);
const port: ControlPlanePort = {
async openRun() {},
checkpointSession,
async appendEvent(rawEvent) {
const event = structuredClone(rawEvent as PrpEvent);
const existing = events.some(
(candidate) =>
candidate.sourceInstanceId === event.sourceInstanceId &&
candidate.sourceSeq === event.sourceSeq,
);
if (!existing) events.push(event);
if (
boundary === "final event append" &&
event.eventType === "run.terminal"
) {
await waitAtBoundary();
}
if (
boundary === "lost final event acknowledgement" &&
event.eventType === "run.terminal" &&
!boundaryBlocked
) {
await waitAtBoundary();
await new Promise<never>(() => undefined);
}
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(
events.filter(
(candidate) =>
candidate.sourceInstanceId === event.sourceInstanceId,
),
),
disposition: existing ? "duplicate" : "committed",
};
},
async replayEvents(replay) {
if (
boundary === "control-plane replay" &&
replay.sourceInstanceId === "control-recovery"
) {
await waitAtBoundary();
}
const sourceEvents = events.filter(
(event) => event.sourceInstanceId === replay.sourceInstanceId,
);
return {
events: structuredClone(
sourceEvents.filter(
(event) => event.sourceSeq > replay.afterSourceSeq,
),
),
highestContiguousSourceSeq: highestContiguous(sourceEvents),
};
},
completeRun,
};
const outcome = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 10,
requireSessionCloseBeforeReturn: true,
}).then(
(value) => ({ value, error: null }),
(error: unknown) => ({ value: null, error }),
);
try {
await boundaryReached;
const checkpointsBeforeFault = checkpointSession.mock.calls.length;
if (boundary === "completion invoked before commit") {
expect(completeRun).toHaveBeenCalledOnce();
expect(durableCompletion).toBeNull();
}
latchedFailure = failure;
releaseBoundary();
if (
boundary === "lost completion acknowledgement" ||
boundary === "lost final event acknowledgement"
) {
await vi.advanceTimersByTimeAsync(10);
}
const settled = await outcome;
if (!admitted) {
expect(settled.error).toBe(failure);
expect(settled.value).toBeNull();
expect(completeRun).not.toHaveBeenCalled();
expect(durableCompletion).toBeNull();
} else {
expect(settled.error).toBeNull();
expect(settled.value).toMatchObject({
result,
terminal,
nativeEventCount: 3,
});
expect(durableCompletion).toMatchObject({ result, terminal });
expect(completeRun).toHaveBeenCalledTimes(
boundary === "lost completion acknowledgement" ? 2 : 1,
);
}
expect(
events.filter((event) => event.sourceKind === "control_plane"),
).toHaveLength(2);
expect(checkpointSession).toHaveBeenCalledTimes(checkpointsBeforeFault);
expect(resolveResult).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
} finally {
releaseBoundary();
await vi.advanceTimersByTimeAsync(30);
await outcome;
vi.useRealTimers();
}
},
);
it.each([
"provider snapshot",
"post-completion checkpoint",
"provider usage",
] as const)(
"preserves durable completion when %s never settles",
async (stalledBoundary) => {
vi.useFakeTimers();
try {
const never = new Promise<never>(() => undefined);
let markEnrichmentStalled = () => {};
const enrichmentStalled = new Promise<void>((resolve) => {
markEnrichmentStalled = resolve;
});
let stalledSignal: AbortSignal | undefined;
let runCompleted = false;
const close = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot(snapshotOptions) {
if (stalledBoundary === "provider snapshot" && runCompleted) {
stalledSignal = snapshotOptions?.signal;
markEnrichmentStalled();
return await never;
}
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
async usage() {
if (stalledBoundary === "provider usage" && runCompleted) {
markEnrichmentStalled();
return await never;
}
return { driverVersion: "2" };
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "post-completion-enrichment-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const completeRun = vi.fn(async () => {
runCompleted = true;
});
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession(_snapshot, operationOptions) {
if (
stalledBoundary === "post-completion checkpoint" &&
runCompleted
) {
stalledSignal = operationOptions?.signal;
markEnrichmentStalled();
await never;
}
},
async appendEvent(event) {
return {
cursor: 1,
highestContiguousSourceSeq: (event as PrpEvent).sourceSeq,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
completeRun,
};
const retainedSessions: Array<NativeSession | null> = [];
const enrichmentFailures: Array<"checkpoint" | "usage"> = [];
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 10,
keepSessionOpen: true,
onSession: (current) => retainedSessions.push(current),
onPostCompletionEnrichmentFailure: ({ stage }) =>
enrichmentFailures.push(stage),
});
await enrichmentStalled;
if (stalledSignal !== undefined)
expect(stalledSignal.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(10);
await expect(execution).resolves.toMatchObject({
result,
terminal,
providerSessionId: "provider-recovery",
driverVersion: "1",
usage: null,
});
if (stalledSignal !== undefined)
expect(stalledSignal.aborted).toBe(true);
expect(completeRun).toHaveBeenCalledOnce();
expect(enrichmentFailures).toEqual([
stalledBoundary === "provider usage" ? "usage" : "checkpoint",
]);
if (stalledBoundary === "provider usage") {
expect(close).not.toHaveBeenCalled();
expect(retainedSessions).toEqual([session]);
} else {
expect(close).toHaveBeenCalledOnce();
expect(retainedSessions).toEqual([session, null]);
}
} finally {
vi.useRealTimers();
}
},
);
it("contains a consumer rejection when starting the turn fails first", async () => {
let markAppendStarted = () => {};
const appendStarted = new Promise<void>((resolve) => {
markAppendStarted = resolve;
});
let releaseAppend = () => {};
const appendReleased = new Promise<void>((resolve) => {
releaseAppend = resolve;
});
let appendCommitted = false;
const close = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.started");
},
async startTurn() {
await appendStarted;
throw new Error("start turn failed");
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(_event, options) {
markAppendStarted();
await Promise.race([
appendReleased,
new Promise<never>((_resolve, reject) => {
const rejectAbort = () =>
reject(options?.signal.reason ?? new Error("append aborted"));
if (options?.signal.aborted) rejectAbort();
else
options?.signal.addEventListener("abort", rejectAbort, {
once: true,
});
}),
]);
appendCommitted = true;
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed" as const,
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await appendStarted;
await expect(execution).rejects.toThrow("start turn failed");
expect(close).toHaveBeenCalled();
expect(appendCommitted).toBe(false);
releaseAppend();
await new Promise<void>((resolve) => setImmediate(resolve));
expect(appendCommitted).toBe(false);
});
it("stops and closes a timed-out consumer even when the caller requested a warm session", async () => {
let markAppendStarted = () => {};
const appendStarted = new Promise<void>((resolve) => {
markAppendStarted = resolve;
});
let releaseAppend = () => {};
const appendReleased = new Promise<void>((resolve) => {
releaseAppend = resolve;
});
let releaseTeardown = () => {};
const teardownReleased = new Promise<void>((resolve) => {
releaseTeardown = resolve;
});
const iteratorTeardown = vi.fn();
let appendCommitted = false;
const appendEvent = vi.fn(
async (_event: PrpEvent, options?: { signal: AbortSignal }) => {
markAppendStarted();
await Promise.race([
appendReleased,
new Promise<never>((_resolve, reject) => {
const rejectAbort = () =>
reject(options?.signal.reason ?? new Error("append aborted"));
if (options?.signal.aborted) rejectAbort();
else
options?.signal.addEventListener("abort", rejectAbort, {
once: true,
});
}),
]);
appendCommitted = true;
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed" as const,
};
},
);
const cancel = vi.fn(() => {
releaseTeardown();
return { cleanup: Promise.resolve() };
});
const close = vi.fn(async () => {
releaseTeardown();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
try {
yield runnerEvent(1, "turn.completed");
} finally {
iteratorTeardown();
await teardownReleased;
}
},
async startTurn() {
return { turnId: "turn-recovery" };
},
cancel,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
appendEvent,
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 1,
keepSessionOpen: true,
});
const rejection = expect(execution).rejects.toThrow(
"native session timed out",
);
await appendStarted;
await vi.waitFor(() => expect(iteratorTeardown).toHaveBeenCalledOnce());
await rejection;
expect(cancel).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalled();
expect(appendEvent).toHaveBeenCalledOnce();
expect(appendCommitted).toBe(false);
releaseAppend();
await new Promise<void>((resolve) => setImmediate(resolve));
expect(appendCommitted).toBe(false);
});
it("closes a failed session while retaining an uncancellable event read", async () => {
let releaseStream = () => {};
const streamReleased = new Promise<void>((resolve) => {
releaseStream = resolve;
});
const close = vi.fn(async () => {
releaseStream();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
},
async *events() {
await streamReleased;
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 1,
keepSessionOpen: true,
}),
).rejects.toThrow("native session timed out");
expect(close).toHaveBeenCalledOnce();
});
it("commits cancellation before bounding failed provider cleanup", async () => {
let releaseStream = () => {};
const streamReleased = new Promise<void>((resolve) => {
releaseStream = resolve;
});
let releaseCancellation = () => {};
const cancellationReleased = new Promise<void>((resolve) => {
releaseCancellation = resolve;
});
const interrupt = vi.fn(() => cancellationReleased);
let cancellationCommitted = false;
const cancel = vi.fn(() => {
cancellationCommitted = true;
return { cleanup: cancellationReleased };
});
const close = vi.fn(async () => {
releaseStream();
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
await streamReleased;
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
interrupt,
cancel,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 1,
keepSessionOpen: true,
});
await vi.waitFor(() => expect(close).toHaveBeenCalledOnce());
expect(interrupt).not.toHaveBeenCalled();
expect(cancel).toHaveBeenCalledOnce();
await expect(execution).rejects.toThrow("native session timed out");
expect(cancellationCommitted).toBe(true);
releaseCancellation();
});
it("bounds failure when iterator teardown and provider close never settle", async () => {
const never = new Promise<void>(() => undefined);
let releaseClose = () => {};
const pendingClose = new Promise<void>((resolve) => {
releaseClose = resolve;
});
const close = vi.fn(() => pendingClose);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
},
async *events() {
await never;
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 1,
keepSessionOpen: true,
}),
).rejects.toThrow("native session timed out");
expect(close).toHaveBeenCalledOnce();
releaseClose();
await pendingClose;
});
it("preserves durable success when provider close never settles", async () => {
let releaseClose = () => {};
const pendingClose = new Promise<void>((resolve) => {
releaseClose = resolve;
});
const close = vi.fn(() => pendingClose);
const completeRun = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
completeRun,
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({ result, terminal });
expect(completeRun).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
releaseClose();
await pendingClose;
});
it.each([false, true])(
"waits for a required backend checkpoint close when enrichment failure=%s",
async (enrichmentFails) => {
let releaseClose = () => {};
const pendingClose = new Promise<void>((resolve) => {
releaseClose = resolve;
});
const close = vi.fn(() => pendingClose);
let runCompleted = false;
let enrichmentFailureObserved = false;
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
// Exercise only the best-effort enrichment snapshot after the
// control plane has durably committed the run result.
if (enrichmentFails && runCompleted) {
enrichmentFailureObserved = true;
throw new Error("checkpoint enrichment failed");
}
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {
runCompleted = true;
},
};
let resolved = false;
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
requireSessionCloseBeforeReturn: true,
}).then((value) => {
resolved = true;
return value;
});
await vi.waitFor(() => expect(close).toHaveBeenCalledOnce());
await new Promise((resolve) => setTimeout(resolve, 150));
expect(resolved).toBe(false);
releaseClose();
await expect(execution).resolves.toMatchObject({ result, terminal });
expect(resolved).toBe(true);
expect(enrichmentFailureObserved).toBe(enrichmentFails);
},
);
it.each([
{ typed: true, closeFails: false },
{ typed: true, closeFails: true },
{ typed: false, closeFails: true },
{ typed: true, closeFails: true, startupRace: true },
])(
"preserves a permanent integrity failure through required cleanup (%j)",
async ({ typed, closeFails, startupRace = false }) => {
const failure = typed
? new NativeSessionProtocolIntegrityError(
"semantic_input_digest_mismatch",
)
: Object.assign(new Error("ordinary provider connection failed"), {
code: "native_event_replay_conflict",
recovery: "operator_required",
});
const closeFailure = new NativeSessionCloseUnrecoverableError();
let observeClose = () => {};
const closeStarted = new Promise<void>((resolve) => {
observeClose = resolve;
});
const close = vi.fn(async () => {
observeClose();
if (closeFails) throw closeFailure;
});
const capabilities = {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
const session: NativeSession = {
identity: () => identity,
capabilities: async () => capabilities,
async *events() {
throw failure;
},
startTurn: async () => {
if (startupRace) {
await closeStarted;
throw new Error(
"provider_transport_failed: startup raced with close",
);
}
return { turnId: "turn-recovery" };
},
result: vi.fn(async () => null),
snapshot: async () => ({
backendKind: "mock",
sessionId: "driver-integrity",
identity,
providerSessionId: "provider-integrity",
cursor: "0",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
}),
close,
};
const backend: NativeSessionBackend = {
descriptor: async () => ({
kind: "mock",
name: `integrity-${typed}-${closeFails}-${startupRace}`,
version: "1",
capabilities,
}),
openSession: vi.fn(async () => session),
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
openRun: async () => {},
checkpointSession: async () => {},
appendEvent: async (event) => {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
replayEvents: async () => ({
events: [],
highestContiguousSourceSeq: 0,
}),
completeRun: vi.fn(async () => {}),
};
const onSession = vi.fn();
const onSessionAdmission = vi.fn(async () => {});
const options: ExecuteNativeSessionOptions = {
onSessionAdmission,
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
onSession,
requireSessionCloseBeforeReturn: true,
timeoutMs: 900_000,
};
await expect(executeNativeSession(options)).rejects.toBe(
typed ? failure : closeFailure,
);
expect(close).toHaveBeenCalledOnce();
expect(onSession).toHaveBeenLastCalledWith(null);
expect(port.completeRun).not.toHaveBeenCalled();
expect(session.result).not.toHaveBeenCalled();
expect(events.some((event) => event.sourceKind === "runner")).toBe(false);
if (closeFails) {
await expect(executeNativeSession(options)).rejects.toBeInstanceOf(
NativeSessionCleanupQuarantinedError,
);
expect(backend.openSession).toHaveBeenCalledOnce();
expect(onSessionAdmission).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
}
},
);
it("retires only the terminated remote resource, including two sandboxes for one run", async () => {
const scopedIdentity = { ...identity, companyId: "remote-stop-company", runId: "remote-stop-run" };
const binding = { ...scopedIdentity, remoteCleanupScope: "first-sandbox" };
const scopedInput = { ...input, binding: { ...input.binding, companyId: scopedIdentity.companyId, runId: scopedIdentity.runId } };
const failure = new NativeSessionCloseUnrecoverableError();
const capabilities = { resume: true, typedEvents: true, steering: false, interruption: false, structuredResult: true };
const session: NativeSession = {
identity: () => scopedIdentity,
capabilities: async () => capabilities,
async *events() { throw new Error("cancelled remote transport"); },
startTurn: async () => ({ turnId: "remote-turn" }),
result: async () => null,
close: vi.fn(async () => { throw failure; }),
};
const backend: NativeSessionBackend = {
descriptor: async () => ({ kind: "remote", name: "remote-stop-test", version: "1", capabilities }),
openSession: vi.fn(async () => session),
};
const controlPlane: ControlPlanePort = {
openRun: async () => {}, checkpointSession: async () => {},
appendEvent: async () => ({ cursor: 0, highestContiguousSourceSeq: 0, disposition: "committed" }),
replayEvents: async () => ({ events: [], highestContiguousSourceSeq: 0 }),
completeRun: vi.fn(async () => {}),
};
const options = { input: scopedInput, backend, controlPlane, runnerInstanceId: "remote-runner",
controlPlaneInstanceId: "control", requireSessionCloseBeforeReturn: true,
remoteCleanupScope: binding.remoteCleanupScope };
await expect(executeNativeSession(options)).rejects.toBe(failure);
expect(completeTerminatedRemoteNativeSessionCleanup({ ...binding, runId: "other-run" })).toBe(true);
expect(completeTerminatedRemoteNativeSessionCleanup({ ...binding, companyId: "other-company" })).toBe(true);
await expect(executeNativeSession(options)).rejects.toBeInstanceOf(NativeSessionCleanupQuarantinedError);
expect(backend.openSession).toHaveBeenCalledOnce();
// A separate sandbox can start without inheriting this process quarantine.
const independent = { ...scopedIdentity, sessionId: "other-sandbox-session" };
const independentSession = { ...session, identity: () => independent };
const independentBackend = { ...backend, openSession: vi.fn(async () => independentSession) };
await expect(executeNativeSession({ ...options, remoteCleanupScope: "other-sandbox",
input: { ...scopedInput, binding: { ...scopedInput.binding, runId: independent.runId } },
backend: independentBackend })).rejects.toBe(failure);
expect(independentBackend.openSession).toHaveBeenCalledOnce();
expect(completeTerminatedRemoteNativeSessionCleanup(binding)).toBe(true);
// Same company/run, different sandbox: its quarantine must remain intact.
await expect(executeNativeSession({ ...options, remoteCleanupScope: "other-sandbox",
backend: independentBackend })).rejects.toBeInstanceOf(NativeSessionCleanupQuarantinedError);
expect(independentBackend.openSession).toHaveBeenCalledOnce();
// Reopening is now possible; the old failure/result was never rewritten.
await expect(executeNativeSession(options)).rejects.toBe(failure);
expect(backend.openSession).toHaveBeenCalledTimes(2);
expect(controlPlane.completeRun).not.toHaveBeenCalled();
completeTerminatedRemoteNativeSessionCleanup(binding);
completeTerminatedRemoteNativeSessionCleanup({ ...binding, remoteCleanupScope: "other-sandbox" });
});
it("retires local quarantine only for the stopped run and runner instance", async () => {
const scopedIdentity = { ...identity, companyId: "local-stop-company", runId: "local-stop-run" };
const binding = { ...scopedIdentity, runnerInstanceId: "local-runner" };
const scopedInput = { ...input, binding: { ...input.binding, companyId: scopedIdentity.companyId, runId: scopedIdentity.runId } };
const failure = new NativeSessionCloseUnrecoverableError();
const capabilities = { resume: true, typedEvents: true, steering: false, interruption: false, structuredResult: true };
const session: NativeSession = {
identity: () => scopedIdentity, capabilities: async () => capabilities,
async *events() { throw new Error("local transport stopped"); },
startTurn: async () => ({ turnId: "local-turn" }), result: async () => null,
close: vi.fn(async () => { throw failure; }),
};
const backend: NativeSessionBackend = {
descriptor: async () => ({ kind: "local", name: "local-stop-test", version: "1", capabilities }),
openSession: vi.fn(async () => session),
};
const controlPlane: ControlPlanePort = {
openRun: async () => {}, checkpointSession: async () => {},
appendEvent: async () => ({ cursor: 0, highestContiguousSourceSeq: 0, disposition: "committed" }),
replayEvents: async () => ({ events: [], highestContiguousSourceSeq: 0 }), completeRun: vi.fn(async () => {}),
};
const options = { input: scopedInput, backend, controlPlane, runnerInstanceId: binding.runnerInstanceId,
controlPlaneInstanceId: "control", requireSessionCloseBeforeReturn: true };
await expect(executeNativeSession(options)).rejects.toBe(failure);
completeTerminatedLocalNativeSessionCleanup({ ...binding, companyId: "other-company" });
completeTerminatedLocalNativeSessionCleanup({ ...binding, runId: "other-run" });
expect(completeTerminatedLocalNativeSessionCleanup({ ...binding, runnerInstanceId: "other-runner" })).toBe(false);
await expect(executeNativeSession(options)).rejects.toBeInstanceOf(NativeSessionCleanupQuarantinedError);
expect(backend.openSession).toHaveBeenCalledOnce();
expect(completeTerminatedLocalNativeSessionCleanup(binding)).toBe(true);
await expect(executeNativeSession(options)).rejects.toBe(failure);
expect(backend.openSession).toHaveBeenCalledTimes(2);
expect(controlPlane.completeRun).not.toHaveBeenCalled();
completeTerminatedLocalNativeSessionCleanup(binding);
});
it("propagates an exhausted required backend checkpoint close", async () => {
vi.useFakeTimers();
try {
const closeFailure = new Error("required remote checkpoint close failed");
const close = vi.fn(({ reason }: { reason: string }) =>
reason === "native session quarantined cleanup recovery"
? Promise.resolve()
: Promise.reject(closeFailure),
);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: false,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
requireSessionCloseBeforeReturn: true,
}),
).rejects.toThrow(closeFailure);
await vi.advanceTimersByTimeAsync(3_000);
await execution;
expect(close).toHaveBeenCalledTimes(5);
} finally {
vi.useRealTimers();
}
});
it("closes after a synchronous governed-wait probe returns no result", async () => {
const resolveGovernedWait = vi.fn(() => null);
const lifecycle: string[] = [];
const close = vi.fn(async () => {
lifecycle.push("closed");
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "item.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: "turn-recovery",
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 5,
resolveGovernedWait,
});
await expect(execution).rejects.toThrow("before a turn terminal fact");
expect(resolveGovernedWait).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalled();
expect(lifecycle).toEqual(["closed"]);
});
it("commits a governed wait without waiting for abort-insensitive provider cleanup", async () => {
const lifecycle: string[] = [];
let cancellationSignal: AbortSignal | undefined;
const cancel = vi.fn(({ signal }: { signal: AbortSignal }) => {
cancellationSignal = signal;
lifecycle.push("cancelled");
return { cleanup: new Promise<void>(() => undefined) };
});
const close = vi.fn(async () => {
lifecycle.push("closed");
});
const events: PrpEvent[] = [];
const retainedSessions: Array<NativeSession | null> = [];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "item.completed");
},
async startTurn() {
return { turnId: "turn-recovery" };
},
cancel,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: null,
activeTurnId: "turn-recovery",
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: event.sourceSeq,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 5,
resolveGovernedWait: () => yieldedResult,
keepSessionOpen: true,
onSession: (current) => retainedSessions.push(current),
});
await expect(execution).resolves.toMatchObject({ result: yieldedResult });
expect(cancel).toHaveBeenCalledOnce();
expect(cancellationSignal?.aborted).toBe(true);
expect(close).toHaveBeenCalledOnce();
expect(lifecycle).toEqual(["cancelled", "closed"]);
expect(retainedSessions.at(-1)).toBeNull();
expect(events.map((event) => event.eventType)).toEqual([
"item.completed",
"run.result.accepted",
"run.terminal",
]);
});
it("uses the execution timeout when a completion report has no provider terminal", async () => {
const lifecycle: string[] = [];
let releaseProvider = () => {};
const providerReleased = new Promise<void>((resolve) => {
releaseProvider = resolve;
});
const cancel = vi.fn(() => {
lifecycle.push("cancelled");
return {
cleanup: new Promise<void>((resolve) =>
setTimeout(() => {
releaseProvider();
resolve();
}, 150),
),
};
});
const providerResult = vi.fn(async () => null);
const close = vi.fn(async () => {
lifecycle.push("closed");
});
const events: PrpEvent[] = [];
const retainedSessions: Array<NativeSession | null> = [];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "run.result.proposed", result);
await providerReleased;
},
async startTurn() {
return { turnId: "turn-recovery" };
},
cancel,
result: providerResult,
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: "turn-recovery",
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "semantic-result-terminal-stall-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
const sourceEvents = events.filter(
(candidate) => candidate.sourceInstanceId === event.sourceInstanceId,
);
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(sourceEvents),
disposition: "committed",
};
},
async replayEvents(replay) {
const sourceEvents = events.filter(
(event) => event.sourceInstanceId === replay.sourceInstanceId,
);
return {
events: structuredClone(
sourceEvents.filter(
(event) => event.sourceSeq > replay.afterSourceSeq,
),
),
highestContiguousSourceSeq: highestContiguous(sourceEvents),
};
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
timeoutMs: 20,
keepSessionOpen: true,
onSession: (current) => retainedSessions.push(current),
}),
).rejects.toThrow("native session timed out after 20ms");
expect(cancel).toHaveBeenCalledWith({
reason: "Native session event consumption failed.",
signal: expect.any(AbortSignal),
});
expect(providerResult).not.toHaveBeenCalled();
expect(events.map((event) => event.eventType)).toEqual([
"run.result.proposed",
]);
expect(lifecycle).toContain("cancelled");
expect(close).toHaveBeenCalledOnce();
expect(retainedSessions.at(-1)).toBeNull();
});
it("retains a reusable session while a semantic terminal releases its remote subscription", async () => {
const close = vi.fn(async () => undefined);
const retainedSessions: Array<NativeSession | null> = [];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
try {
yield runnerEvent(1, "run.result.proposed", result);
yield {
...runnerEvent(2, "turn.completed"),
turnId: "turn-recovery",
};
} finally {
await new Promise<void>((resolve) => setTimeout(resolve, 150));
}
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "2",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "semantic-terminal-subscription-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
keepSessionOpen: true,
onSession: (current) => retainedSessions.push(current),
}),
).resolves.toMatchObject({ result, terminal });
expect(close).not.toHaveBeenCalled();
expect(retainedSessions).toEqual([session]);
});
it.each([
["turn.completed", "succeeded"],
["turn.failed", "failed"],
["turn.cancelled", "cancelled"],
["turn.interrupted", "cancelled"],
["stream.closed", null],
] as const)("persists a delayed final answer before settling %s", async (eventType, runTerminalState) => {
vi.useFakeTimers();
try {
let releasePersistence!: () => void;
const persistence = new Promise<void>((resolve) => { releasePersistence = resolve; });
let persistingAnswer = false;
const cancel = vi.fn(() => ({ cleanup: Promise.resolve() }));
const close = vi.fn(async () => undefined);
const events: PrpEvent[] = [];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "run.result.proposed", result);
await new Promise((resolve) => setTimeout(resolve, 6_000));
yield runnerEvent(2, "item.completed", {
kind: "agentMessage", channel: "final", text: "Final response.",
});
if (eventType !== "stream.closed") yield runnerEvent(3, eventType);
},
async startTurn() {
return { turnId: "turn-recovery" };
},
cancel,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-recovery",
cursor: "3",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "semantic-result-final-response-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
if (event.eventType === "item.completed") {
persistingAnswer = true;
await persistence;
}
events.push(structuredClone(event as PrpEvent));
const sourceEvents = events.filter(
(candidate) => candidate.sourceInstanceId === event.sourceInstanceId,
);
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(sourceEvents),
disposition: "committed",
};
},
async replayEvents(replay) {
const sourceEvents = events.filter(
(event) => event.sourceInstanceId === replay.sourceInstanceId,
);
return {
events: structuredClone(
sourceEvents.filter(
(event) => event.sourceSeq > replay.afterSourceSeq,
),
),
highestContiguousSourceSeq: highestContiguous(sourceEvents),
};
},
async completeRun() {},
};
let completed = false;
const execution = executeNativeSession({
input, backend, controlPlane: port,
runnerInstanceId: "runner-recovery", controlPlaneInstanceId: "control-recovery",
}).then((value) => { completed = true; return value; });
await vi.waitFor(() => expect(events).toHaveLength(1));
await vi.advanceTimersByTimeAsync(6_001);
expect(persistingAnswer).toBe(true);
expect(completed).toBe(false);
expect(events.map((event) => event.eventType)).toEqual(["run.result.proposed"]);
expect(cancel).not.toHaveBeenCalled();
releasePersistence();
if (eventType === "stream.closed") {
await expect(execution).rejects.toThrow("native event stream closed before a turn terminal fact");
expect(events.map((event) => event.eventType)).toEqual(["run.result.proposed", "item.completed"]);
return;
}
await expect(execution).resolves.toMatchObject({ result, terminal: { runTerminalState } });
expect(cancel).not.toHaveBeenCalled();
expect(close).toHaveBeenCalledOnce();
expect(events.map((event) => event.eventType)).toEqual([
"run.result.proposed",
"item.completed",
eventType,
"run.result.accepted",
"run.terminal",
]);
} finally {
vi.useRealTimers();
}
});
it("rejects a mismatched checkpoint before it mutates control-plane state", async () => {
const openRun = vi.fn(async () => undefined);
const checkpointSession = vi.fn(async () => undefined);
const openSession = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
openRun,
async loadSessionCheckpoint() {
return {
backendKind: "mock",
sessionId: "driver-recovery",
identity: { ...identity, companyId: "other-company" },
};
},
checkpointSession,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).rejects.toThrow("native_session_checkpoint_binding_mismatch");
expect(openRun).not.toHaveBeenCalled();
expect(checkpointSession).not.toHaveBeenCalled();
expect(openSession).not.toHaveBeenCalled();
});
it("rejects a mismatched existing session before opening control-plane state", async () => {
const openRun = vi.fn(async () => undefined);
const attachRun = vi.fn(async () => undefined);
const existingSession: NativeSession = {
identity: () => ({ ...identity, companyId: "other-company" }),
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
attachRun,
async *events() {},
async startTurn() {
return { turnId: "unexpected" };
},
async result() {
return null;
},
async snapshot() {
throw new Error("unexpected snapshot");
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "existing-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("unexpected open");
},
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
existingSession,
}),
).rejects.toThrow("native_session_attach_binding_mismatch");
expect(openRun).not.toHaveBeenCalled();
expect(attachRun).not.toHaveBeenCalled();
});
it("quarantines a retained session when attachment partially mutates then fails", async () => {
const attachmentFailure = new Error("provider attachment failed");
const openRun = vi.fn(async () => undefined);
let retainedIdentity = { ...identity, runId: "run-previous" };
const attachRun = vi.fn(async (input: { identity: NativeRunIdentity }) => {
retainedIdentity = structuredClone(input.identity);
throw attachmentFailure;
});
const close = vi.fn(() => new Promise<void>(() => {}));
const startTurn = vi.fn(async () => ({ turnId: "unexpected" }));
const onSession = vi.fn();
const existingSession: NativeSession = {
identity: () => structuredClone(retainedIdentity),
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
attachRun,
async *events() {},
startTurn,
async result() {
return null;
},
async snapshot() {
throw new Error("unexpected snapshot");
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "attachment-failure-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("unexpected open");
},
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
existingSession,
onSession,
}),
).rejects.toBe(attachmentFailure);
expect(attachRun).toHaveBeenCalledWith({ identity });
expect(retainedIdentity).toEqual(identity);
expect(onSession).toHaveBeenCalledOnce();
expect(onSession).toHaveBeenCalledWith(null);
expect(close).toHaveBeenCalledWith({
reason: "native session attachment failed",
});
expect(openRun).not.toHaveBeenCalled();
expect(startTurn).not.toHaveBeenCalled();
});
it("quarantines an attached session when control-plane run admission fails", async () => {
const admissionFailure = new Error("control-plane admission failed");
const openRun = vi.fn(async () => {
throw admissionFailure;
});
const attachRun = vi.fn(async () => undefined);
const close = vi.fn(async () => undefined);
const startTurn = vi.fn(async () => ({ turnId: "unexpected" }));
const onSession = vi.fn();
const existingSession: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
attachRun,
async *events() {},
startTurn,
async result() {
return null;
},
async snapshot() {
throw new Error("unexpected snapshot");
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "control-plane-admission-failure-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("unexpected open");
},
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
existingSession,
onSession,
}),
).rejects.toBe(admissionFailure);
expect(attachRun).toHaveBeenCalledWith({ identity });
expect(openRun).toHaveBeenCalledOnce();
expect(onSession).toHaveBeenCalledOnce();
expect(onSession).toHaveBeenCalledWith(null);
expect(close).toHaveBeenCalledWith({
reason: "native control-plane run admission failed",
});
expect(startTurn).not.toHaveBeenCalled();
});
it("rejects checkpoint adoption when the requested session id is absent", async () => {
const openRun = vi.fn(async () => undefined);
const openSession = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
};
const port: ControlPlanePort = {
openRun,
async loadSessionCheckpoint() {
return {
backendKind: "mock",
sessionId: "driver-other-session",
identity: { ...identity, sessionId: "other-session" },
};
},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input: {
...input,
session: { ...input.session, normalizedSessionId: null },
},
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).rejects.toThrow("native_session_checkpoint_binding_mismatch");
expect(openRun).not.toHaveBeenCalled();
expect(openSession).not.toHaveBeenCalled();
});
it("proves required provider recovery before re-opening the durable run", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-unrecoverable",
identity,
providerSessionId: "provider-unrecoverable",
providerRecoveryPolicy: "same_session_only",
cursor: "0",
activeTurnId: "turn-unrecoverable",
pendingRuntimeRequests: [],
lineage: [],
};
const openRun = vi.fn(async () => undefined);
const completeRun = vi.fn(async () => undefined);
const recoverSession = vi.fn(async () => ({
recovered: false as const,
reason: "provider session no longer exists",
}));
const openSession = vi.fn(async () => {
throw new Error("replacement is forbidden");
});
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
recoverSession,
};
const port: ControlPlanePort = {
openRun,
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
completeRun,
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).rejects.toThrow(
"native_session_recovery_failed: provider session no longer exists",
);
expect(recoverSession).toHaveBeenCalledOnce();
expect(openRun).not.toHaveBeenCalled();
expect(completeRun).not.toHaveBeenCalled();
expect(openSession).not.toHaveBeenCalled();
});
it("bounds paginated recovery replay with one signal and observes a late rejection", async () => {
vi.useFakeTimers();
let rejectStalledReplay = (_error: Error) => {};
try {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-replay-stalled",
identity,
providerSessionId: "provider-replay-stalled",
providerRecoveryPolicy: "same_session_only",
cursor: "0",
activeTurnId: "turn-replay-stalled",
pendingRuntimeRequests: [],
lineage: [],
};
const stalledReplay = new Promise<never>((_resolve, reject) => {
rejectStalledReplay = reject;
});
let markSecondPageStarted = () => {};
const secondPageStarted = new Promise<void>((resolve) => {
markSecondPageStarted = resolve;
});
const replaySignals: AbortSignal[] = [];
const replayEvents = vi.fn<ControlPlanePort["replayEvents"]>(
async (replay, operationOptions) => {
expect(operationOptions?.signal).toBeInstanceOf(AbortSignal);
replaySignals.push(operationOptions!.signal);
if (replay.afterSourceSeq === 0) {
return {
events: [runnerEvent(1, "item.completed", { kind: "progress" })],
highestContiguousSourceSeq: 1,
};
}
markSecondPageStarted();
return await stalledReplay;
},
);
const recoverSession = vi.fn();
const openRun = vi.fn(async () => undefined);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("unexpected replacement");
},
recoverSession,
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
replayEvents,
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
persistedSession: structuredClone(checkpoint),
timeoutMs: 5,
});
const rejection = expect(execution).rejects.toThrow(
"native session recovery replay timed out after 5ms",
);
await secondPageStarted;
expect(replayEvents).toHaveBeenCalledTimes(2);
expect(replaySignals).toHaveLength(2);
expect(replaySignals[1]).toBe(replaySignals[0]);
expect(replaySignals[0]?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(5);
await rejection;
expect(replaySignals[0]?.aborted).toBe(true);
expect(recoverSession).not.toHaveBeenCalled();
expect(openRun).not.toHaveBeenCalled();
// A broken adapter may ignore abort and reject later. The bounded helper
// keeps that losing operation observed after execution already rejected.
rejectStalledReplay(new Error("late recovery replay failure"));
await Promise.resolve();
} finally {
vi.useRealTimers();
}
});
it("bounds provider recovery and closes a session returned after timeout", async () => {
vi.useFakeTimers();
try {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery-stalled",
identity,
providerSessionId: "provider-recovery-stalled",
providerRecoveryPolicy: "same_session_only",
cursor: "0",
activeTurnId: "turn-recovery-stalled",
pendingRuntimeRequests: [],
lineage: [],
};
let resolveRecovery = (_value: {
recovered: true;
session: NativeSession;
}) => {};
const stalledRecovery = new Promise<{
recovered: true;
session: NativeSession;
}>((resolve) => {
resolveRecovery = resolve;
});
let markRecoveryStarted = () => {};
const recoveryStarted = new Promise<void>((resolve) => {
markRecoveryStarted = resolve;
});
let markCloseStarted = () => {};
const closeStarted = new Promise<void>((resolve) => {
markCloseStarted = resolve;
});
const close = vi.fn(async () => {
markCloseStarted();
});
const lateSession: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
};
},
async *events() {},
async startTurn() {
throw new Error("late recovery session must not start");
},
async result() {
return null;
},
async snapshot() {
return structuredClone(checkpoint);
},
close,
};
let recoverySignal: AbortSignal | undefined;
const recoverSession = vi.fn(
(
_checkpoint: PersistedNativeSession,
recoveryOptions: { signal: AbortSignal },
) => {
recoverySignal = recoveryOptions.signal;
recoveryOptions.signal.addEventListener(
"abort",
() => resolveRecovery({ recovered: true, session: lateSession }),
{ once: true },
);
markRecoveryStarted();
return stalledRecovery;
},
);
const openRun = vi.fn(async () => undefined);
const onSession = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
},
};
},
async openSession() {
throw new Error("unexpected replacement");
},
recoverSession,
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
persistedSession: structuredClone(checkpoint),
timeoutMs: 5,
onSession,
});
const rejection = expect(execution).rejects.toThrow(
"native session provider recovery timed out after 5ms",
);
await recoveryStarted;
expect(recoverySignal?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(5);
await rejection;
expect(recoverySignal?.aborted).toBe(true);
expect(openRun).not.toHaveBeenCalled();
expect(onSession).not.toHaveBeenCalled();
await closeStarted;
expect(close).toHaveBeenCalledWith({
reason: "native session provider recovery timed out",
});
expect(onSession).not.toHaveBeenCalled();
} finally {
vi.useRealTimers();
}
});
it("bounds replacement bootstrap and closes a session returned after timeout", async () => {
vi.useFakeTimers();
try {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-replacement-stalled",
identity,
providerSessionId: "provider-replacement-stalled",
providerRecoveryPolicy: "allow_replacement_after_resume_failure",
cursor: "0",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
let resolveReplacement = (_value: NativeSession) => {};
const stalledReplacement = new Promise<NativeSession>((resolve) => {
resolveReplacement = resolve;
});
let markReplacementStarted = () => {};
const replacementStarted = new Promise<void>((resolve) => {
markReplacementStarted = resolve;
});
let markCloseStarted = () => {};
const closeStarted = new Promise<void>((resolve) => {
markCloseStarted = resolve;
});
const close = vi.fn(async () => {
markCloseStarted();
});
const lateSession: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
};
},
async *events() {},
async startTurn() {
throw new Error("late replacement session must not start");
},
async result() {
return null;
},
async snapshot() {
return structuredClone(checkpoint);
},
close,
};
let replacementSignal: AbortSignal | undefined;
const openReplacementSession = vi.fn(
(replacementInput: {
identity: NativeRunIdentity;
workingDirectory?: string;
signal?: AbortSignal;
}) => {
replacementSignal = replacementInput.signal;
replacementInput.signal?.addEventListener(
"abort",
() => resolveReplacement(lateSession),
{ once: true },
);
markReplacementStarted();
return stalledReplacement;
},
);
const openRun = vi.fn(async () => undefined);
const onSession = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "replacement-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
},
};
},
async openSession() {
throw new Error("replacement seam must be used");
},
async recoverSession() {
return { recovered: false, reason: "provider session is missing" };
},
openReplacementSession,
};
const port: ControlPlanePort = {
openRun,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
persistedSession: structuredClone(checkpoint),
timeoutMs: 5,
onSession,
});
const rejection = expect(execution).rejects.toThrow(
"native session replacement bootstrap timed out after 5ms",
);
await replacementStarted;
expect(replacementSignal?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(5);
await rejection;
expect(replacementSignal?.aborted).toBe(true);
expect(openRun).not.toHaveBeenCalled();
expect(onSession).not.toHaveBeenCalled();
await closeStarted;
expect(close).toHaveBeenCalledWith({
reason: "native session replacement bootstrap timed out",
});
expect(onSession).not.toHaveBeenCalled();
} finally {
vi.useRealTimers();
}
});
it("does not persist a reconciled recovery cursor when provider recovery rejects", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery-rejects",
identity,
providerSessionId: "provider-recovery-rejects",
providerRecoveryPolicy: "same_session_only",
cursor: "0",
activeTurnId: "turn-recovery-rejects",
pendingRuntimeRequests: [],
lineage: [],
};
const recoveryFailure = new Error("provider recovery rejected");
const openRun = vi.fn(async () => undefined);
const checkpointSession = vi.fn(async () => undefined);
const onCheckpoint = vi.fn(async () => undefined);
const recoverSession = vi.fn(
async (recoveryCheckpoint: PersistedNativeSession) => {
expect(recoveryCheckpoint.cursor).toBe("1");
throw recoveryFailure;
},
);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("unexpected replacement");
},
recoverSession,
};
const port: ControlPlanePort = {
openRun,
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
checkpointSession,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents(replay) {
return {
events: replay.afterSourceSeq === 0 ? [runnerEvent(1)] : [],
highestContiguousSourceSeq: 1,
};
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
onCheckpoint,
}),
).rejects.toBe(recoveryFailure);
expect(recoverSession).toHaveBeenCalledOnce();
expect(openRun).not.toHaveBeenCalled();
expect(checkpointSession).not.toHaveBeenCalled();
expect(onCheckpoint).not.toHaveBeenCalled();
});
it("does not persist a reconciled recovery cursor when run admission rejects", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-admission-rejects",
identity,
providerSessionId: "provider-admission-rejects",
providerRecoveryPolicy: "same_session_only",
cursor: "0",
activeTurnId: "turn-admission-rejects",
pendingRuntimeRequests: [],
lineage: [],
};
const admissionFailure = new Error("control-plane admission rejected");
const openRun = vi.fn(async () => {
throw admissionFailure;
});
const checkpointSession = vi.fn(async () => undefined);
const onCheckpoint = vi.fn(async () => undefined);
const onSession = vi.fn();
const close = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {},
async startTurn() {
throw new Error("unexpected turn");
},
async result() {
return null;
},
async snapshot() {
throw new Error("unexpected snapshot");
},
close,
};
const recoverSession = vi.fn(
async (recoveryCheckpoint: PersistedNativeSession) => {
expect(recoveryCheckpoint.cursor).toBe("1");
return { recovered: true as const, session };
},
);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("unexpected replacement");
},
recoverSession,
};
const port: ControlPlanePort = {
openRun,
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
checkpointSession,
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents(replay) {
return {
events: replay.afterSourceSeq === 0 ? [runnerEvent(1)] : [],
highestContiguousSourceSeq: 1,
};
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
onCheckpoint,
onSession,
}),
).rejects.toBe(admissionFailure);
expect(recoverSession).toHaveBeenCalledOnce();
expect(openRun).toHaveBeenCalledOnce();
expect(checkpointSession).not.toHaveBeenCalled();
expect(onCheckpoint).not.toHaveBeenCalled();
expect(onSession).toHaveBeenCalledOnce();
expect(onSession).toHaveBeenCalledWith(null);
expect(close).toHaveBeenCalledWith({
reason: "native control-plane run admission failed",
});
});
it("continues a provider-reported active turn without starting a duplicate turn", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "0",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
const providerSnapshot: PersistedNativeSession = {
...checkpoint,
cursor: "1",
activeTurnId: "turn-recovery",
};
const terminalEvent: PrpEvent = {
schema: "paperclip.prp.event.v1",
sourceEventId: "provider-recovery:1",
sourceSeq: 1,
sourceInstanceId: "provider-recovery",
sourceKind: "provider",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-recovery",
eventType: "turn.completed",
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:00.000Z",
payload: {},
};
const bySource = new Map<string, PrpEvent[]>();
const startTurn = vi.fn(async () => ({ turnId: "duplicate-turn" }));
const openSession = vi.fn(async () => {
throw new Error("must recover the provider session");
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield terminalEvent;
},
startTurn,
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return structuredClone(providerSnapshot);
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
async recoverSession() {
return { recovered: true, session };
},
};
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent(event) {
const list = bySource.get(event.sourceInstanceId) ?? [];
list.push(structuredClone(event));
bySource.set(event.sourceInstanceId, list);
return {
cursor: list.length,
highestContiguousSourceSeq: highestContiguous(list),
disposition: "committed",
};
},
async replayEvents(replay) {
const list = bySource.get(replay.sourceInstanceId) ?? [];
return {
events: structuredClone(
list.filter((event) => event.sourceSeq > replay.afterSourceSeq),
),
highestContiguousSourceSeq: highestContiguous(list),
};
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({
turnId: "turn-recovery",
providerSessionId: "provider-recovery",
});
expect(openSession).not.toHaveBeenCalled();
expect(startTurn).not.toHaveBeenCalled();
});
it.each([
{ checkpointCursor: "12", expectedCursor: "41", terminalSequence: 42 },
{ checkpointCursor: "50", expectedCursor: "50", terminalSequence: 51 },
])(
"seeds recovery from the larger of checkpoint $checkpointCursor and the persisted source high-water mark",
async ({ checkpointCursor, expectedCursor, terminalSequence }) => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: checkpointCursor,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
const runnerEvents = [
runnerEvent(13, "item.completed", { kind: "progress" }),
runnerEvent(41, "item.completed", { kind: "progress" }),
];
const terminalEvent = runnerEvent(terminalSequence, "turn.completed");
const controlEvents: PrpEvent[] = [];
const checkpoints: PersistedNativeSession[] = [];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield terminalEvent;
},
async startTurn() {
return { turnId: "turn-recovery" };
},
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return {
...checkpoint,
cursor: String(terminalSequence),
activeTurnId: null,
};
},
async close() {},
};
const recoverSession = vi.fn(
async (recoveryCheckpoint: PersistedNativeSession) => {
expect(recoveryCheckpoint.cursor).toBe(expectedCursor);
return { recovered: true, session };
},
);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("must recover the provider session");
},
recoverSession,
};
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession(snapshot) {
checkpoints.push(structuredClone(snapshot));
},
async appendEvent(event) {
const target =
event.sourceInstanceId === "runner-recovery"
? runnerEvents
: controlEvents;
if (
target.some((existing) => existing.sourceSeq === event.sourceSeq)
) {
throw new Error(`native_event_replay_conflict:${event.sourceSeq}`);
}
target.push(structuredClone(event));
return {
cursor: target.length,
highestContiguousSourceSeq: highestContiguous(target),
disposition: "committed",
};
},
async replayEvents(replay) {
const source =
replay.sourceInstanceId === "runner-recovery"
? runnerEvents
: controlEvents;
const events = source
.filter((event) => event.sourceSeq > replay.afterSourceSeq)
.sort((left, right) => left.sourceSeq - right.sourceSeq)
.slice(0, replay.limit);
return {
events: structuredClone(events),
highestContiguousSourceSeq: highestContiguous(source),
};
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({ turnId: "turn-recovery" });
expect(recoverSession).toHaveBeenCalledOnce();
expect(
runnerEvents.some((event) => event.sourceSeq === terminalSequence),
).toBe(true);
if (checkpointCursor === "12") {
expect(checkpoints[0]).toMatchObject({ cursor: "41" });
}
},
);
it("attempts exact recovery before opening an observable replacement session", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-old",
identity,
providerSessionId: "provider-old",
providerRecoveryPolicy: "allow_replacement_after_resume_failure",
cursor: null,
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
const replacementSnapshot: PersistedNativeSession = {
...checkpoint,
sessionId: "driver-new",
providerSessionId: "provider-new",
providerRecoveryPolicy: "same_session_only",
};
const replacementSession: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
async startTurn() {
return { turnId: "turn-replacement" };
},
async result() {
return { result, terminal, turnId: "turn-replacement" };
},
async snapshot() {
return structuredClone(replacementSnapshot);
},
async close() {},
};
const recoverSession = vi.fn(async () => ({
recovered: false as const,
reason: "provider reported the prior session missing",
}));
const openReplacementSession = vi.fn(async () => replacementSession);
const onContinuityBreak = vi.fn(async () => undefined);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "replacement-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("replacement seam must be used");
},
recoverSession,
openReplacementSession,
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-replacement",
controlPlaneInstanceId: "control-replacement",
onContinuityBreak,
}),
).resolves.toMatchObject({ providerSessionId: "provider-new" });
expect(recoverSession).toHaveBeenCalledOnce();
expect(openReplacementSession).toHaveBeenCalledOnce();
expect(onContinuityBreak).toHaveBeenCalledWith({
reason: "provider reported the prior session missing",
previousDriverSessionId: "driver-old",
previousProviderSessionId: "provider-old",
replacementDriverSessionId: "driver-new",
replacementProviderSessionId: "provider-new",
});
});
it("does not replace a failed provider session when recovery policy forbids it", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-failed-constrained",
identity,
providerSessionId: "provider-failed-constrained",
providerRecoveryPolicy: "same_session_only",
cursor: null,
activeTurnId: null,
semanticResult: null,
terminal: {
schema: "paperclip.prp.terminal.v1",
turnTerminalState: "failed",
runTerminalState: "failed",
reportedWorkDisposition: "yielded",
},
terminalTurns: [{ turnId: "turn-failed", fingerprint: "failed" }],
pendingRuntimeRequests: [],
lineage: [],
};
const openRun = vi.fn(async () => undefined);
const openSession = vi.fn(async () => {
throw new Error("replacement is forbidden");
});
const recoverSession = vi.fn();
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "constrained-recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
recoverSession,
};
const port: ControlPlanePort = {
openRun,
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async appendEvent() {
throw new Error("unexpected event");
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-constrained-recovery",
controlPlaneInstanceId: "control-constrained-recovery",
}),
).rejects.toMatchObject({
code: "native_provider_terminal_failed",
providerCode: "provider_checkpoint_failed_terminal",
recoverable: false,
});
expect(recoverSession).not.toHaveBeenCalled();
expect(openSession).not.toHaveBeenCalled();
expect(openRun).not.toHaveBeenCalled();
});
it("replaces a provider session that already ended with a failed terminal", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-failed",
identity,
providerSessionId: "provider-failed",
providerRecoveryPolicy: "allow_replacement_after_resume_failure",
cursor: null,
activeTurnId: null,
semanticResult: null,
terminal: {
schema: "paperclip.prp.terminal.v1",
turnTerminalState: "failed",
runTerminalState: "failed",
reportedWorkDisposition: "yielded",
},
terminalTurns: [{ turnId: "turn-failed", fingerprint: "failed" }],
pendingRuntimeRequests: [],
lineage: [],
};
const replacementSnapshot: PersistedNativeSession = {
...checkpoint,
sessionId: "driver-replacement",
providerSessionId: "provider-replacement",
terminal: null,
terminalTurns: [],
};
const startTurn = vi.fn(async () => ({ turnId: "turn-replacement" }));
const replacementSession: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield runnerEvent(1, "turn.completed");
},
startTurn,
async result() {
return { result, terminal, turnId: "turn-replacement" };
},
async snapshot() {
return structuredClone(replacementSnapshot);
},
async close() {},
};
const recoverSession = vi.fn(async () => ({
recovered: true as const,
session: replacementSession,
}));
const openReplacementSession = vi.fn(async () => replacementSession);
const onContinuityBreak = vi.fn(async () => undefined);
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "replacement-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("replacement seam must be used");
},
recoverSession,
openReplacementSession,
};
const events: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-replacement",
controlPlaneInstanceId: "control-replacement",
onContinuityBreak,
}),
).resolves.toMatchObject({ providerSessionId: "provider-replacement" });
expect(recoverSession).not.toHaveBeenCalled();
expect(openReplacementSession).toHaveBeenCalledOnce();
const replacementEnvelope = JSON.parse(
startTurn.mock.calls[0]![0].message.text,
) as { task: { prompt: string } };
expect(replacementEnvelope.task.prompt).toBe(input.task.prompt);
expect(onContinuityBreak).toHaveBeenCalledWith({
reason: "provider session ended with a failed terminal",
previousDriverSessionId: "driver-failed",
previousProviderSessionId: "provider-failed",
replacementDriverSessionId: "driver-replacement",
replacementProviderSessionId: "provider-replacement",
});
});
it("only replays the original ACPX envelope for a proven effect-free initial turn", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
driverKind: "acpx_runtime",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
terminalTurns: [
{ turnId: "turn-work", fingerprint: "terminal-fingerprint" },
],
dispositionOnlyRecoveryConsumed: true,
dispositionOnlyRecoveryTurnId: "turn-missing-disposition",
pendingRuntimeRequests: [],
lineage: [],
};
const recoveredSnapshot: PersistedNativeSession = {
...checkpoint,
dispositionOnlyRecoveryConsumed: false,
dispositionOnlyRecoveryTurnId: null,
};
const terminalEvent: PrpEvent = {
schema: "paperclip.prp.event.v1",
sourceEventId: "provider-recovery:2",
sourceSeq: 2,
sourceInstanceId: "provider-recovery",
sourceKind: "provider",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-continuation",
eventType: "turn.completed",
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:01.000Z",
payload: {},
};
const startTurn = vi.fn(async () => ({ turnId: "turn-continuation" }));
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield terminalEvent;
},
startTurn,
async result() {
return { result, terminal, turnId: "turn-continuation" };
},
async snapshot() {
return structuredClone(recoveredSnapshot);
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("must recover the provider session");
},
async recoverSession() {
return { recovered: true, session };
},
};
const bySource = new Map<string, PrpEvent[]>();
const replayedPages: PrpEvent[][] = [];
const replayEvents = vi.fn(
async (replay: Parameters<ControlPlanePort["replayEvents"]>[0]) => {
const list = bySource.get(replay.sourceInstanceId) ?? [];
const events = structuredClone(
list.filter((event) => event.sourceSeq > replay.afterSourceSeq),
);
replayedPages.push(events);
return {
events,
highestContiguousSourceSeq: highestContiguous(list),
};
},
);
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent(event) {
const list = bySource.get(event.sourceInstanceId) ?? [];
list.push(structuredClone(event));
bySource.set(event.sourceInstanceId, list);
return {
cursor: list.length,
highestContiguousSourceSeq: highestContiguous(list),
disposition: "committed",
};
},
replayEvents,
async completeRun() {},
};
const submittedTurn = runnerEvent(1, "turn.submitted");
delete submittedTurn.turnId;
const effectFreeTurn = [
submittedTurn,
{
...runnerEvent(2, "turn.started", { status: "inProgress" }),
turnId: "turn-work",
},
{ ...runnerEvent(3, "turn.accepted"), turnId: "turn-work" },
{
...runnerEvent(4, "item.completed", {
kind: "usage",
usage: {
total: {
requests: 1,
inputTokens: 0,
outputTokens: 0,
activeSeconds: 0,
providerCostUsd: 0,
cacheReadTokens: 0,
cacheWriteTokens: 0,
},
runDelta: {
requests: 1,
inputTokens: 0,
outputTokens: 0,
activeSeconds: 0,
providerCostUsd: 0,
cacheReadTokens: 0,
cacheWriteTokens: 0,
},
},
}),
turnId: "turn-work",
},
{
...runnerEvent(5, "turn.completed", {
status: "completed",
error: null,
}),
turnId: "turn-work",
},
];
bySource.set("runner-recovery", effectFreeTurn);
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({
turnId: "turn-continuation",
providerSessionId: "provider-recovery",
});
expect(startTurn).toHaveBeenCalledOnce();
expect(replayEvents).toHaveBeenCalledWith({
runId: identity.runId,
sourceInstanceId: "runner-recovery",
afterSourceSeq: 0,
limit: 1_000,
});
expect(
replayedPages.some(
(events) =>
events.length === effectFreeTurn.length &&
events.every(
(event, index) =>
event.sourceSeq === effectFreeTurn[index]!.sourceSeq,
),
),
).toBe(true);
const recoveryEnvelope = JSON.parse(
startTurn.mock.calls[0]![0].message.text,
) as { task: { prompt: string } };
expect(recoveryEnvelope.task.prompt).toBe(input.task.prompt);
startTurn.mockClear();
bySource.set("runner-recovery", [
...effectFreeTurn.slice(0, 3),
{
...runnerEvent(4, "item.completed", {
kind: "agentMessage",
text: "Work may already have been performed.",
}),
turnId: "turn-work",
},
{
...runnerEvent(5, "turn.completed", {
status: "completed",
error: null,
}),
turnId: "turn-work",
},
]);
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({
turnId: "turn-continuation",
providerSessionId: "provider-recovery",
});
const dispositionEnvelope = JSON.parse(
startTurn.mock.calls[0]![0].message.text,
) as { task: { prompt: string } };
expect(dispositionEnvelope.task.prompt).toContain(
"semantic-result recovery for a prior completed provider turn",
);
expect(dispositionEnvelope.task.prompt).toContain(
"Do not repeat implementation, tests, research, or the final answer",
);
expect(dispositionEnvelope.task.prompt).not.toContain(input.task.prompt);
checkpoint.dispositionOnlyRecoveryTurnId = undefined;
recoveredSnapshot.dispositionOnlyRecoveryTurnId = undefined;
startTurn.mockClear();
bySource.clear();
bySource.set("runner-recovery", [
{
...terminalEvent,
sourceEventId: "runner-recovery:stale-terminal",
turnId: "turn-stale-unbound",
},
]);
await expect(
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
}),
).resolves.toMatchObject({
turnId: "turn-continuation",
providerSessionId: "provider-recovery",
});
expect(startTurn).toHaveBeenCalledOnce();
});
it.each(
[true, false].flatMap((dispositionRecovery) =>
(["turn.completed", "turn.interrupted"] as const).flatMap(
(terminalType) =>
[false, true].map((failInitialAppend) => ({
dispositionRecovery,
terminalType,
failInitialAppend,
})),
),
),
)(
"consumes adopted $terminalType without resending (disposition: $dispositionRecovery, failed first append: $failInitialAppend)",
async ({ dispositionRecovery, terminalType, failInitialAppend }) => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: dispositionRecovery ? null : "turn-disposition",
terminalTurns: dispositionRecovery
? [{ turnId: "turn-work", fingerprint: "work-terminal" }]
: [],
dispositionOnlyRecoveryConsumed: false,
pendingRuntimeRequests: [],
lineage: [],
};
const recoveredSnapshot: PersistedNativeSession = {
...checkpoint,
cursor: "2",
activeTurnId: null,
terminalTurns: [
...checkpoint.terminalTurns!,
{ turnId: "turn-disposition", fingerprint: "disposition-terminal" },
],
dispositionOnlyRecoveryConsumed: dispositionRecovery,
};
const terminalEvent: PrpEvent = {
schema: "paperclip.prp.event.v1",
sourceEventId: "provider-recovery:2",
sourceSeq: 2,
sourceInstanceId: "provider-recovery",
sourceKind: "provider",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-disposition",
eventType: terminalType,
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:01.000Z",
payload: {},
};
const startTurn = vi.fn(async () => ({ turnId: "unexpected-turn" }));
const close = vi.fn(async () => undefined);
const completeRun = vi.fn(async () => undefined);
const appendFailure = new Error("adopted terminal append failed");
let failNextAppend = failInitialAppend;
let durableCheckpoint = structuredClone(checkpoint);
const recoveryCheckpoints: PersistedNativeSession[] = [];
let dispositionTerminalCommitted = false;
let prematureDispositionCheckpoint = false;
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield terminalEvent;
},
startTurn,
async result() {
return null;
},
async snapshot() {
return structuredClone(recoveredSnapshot);
},
close,
};
const events: PrpEvent[] = [];
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("must recover the provider session");
},
async recoverSession(snapshot) {
recoveryCheckpoints.push(structuredClone(snapshot));
return { recovered: true, session: { ...session } };
},
};
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(durableCheckpoint);
},
async checkpointSession(snapshot) {
if (
snapshot.terminalTurns?.some(
(turn) => turn.turnId === "turn-disposition",
) &&
!dispositionTerminalCommitted
)
prematureDispositionCheckpoint = true;
durableCheckpoint = structuredClone(snapshot);
},
async appendEvent(event) {
if (event.eventType === terminalType && failNextAppend) {
failNextAppend = false;
throw appendFailure;
}
events.push(structuredClone(event));
if (
event.eventType === terminalType &&
event.turnId === "turn-disposition"
) {
dispositionTerminalCommitted = true;
}
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents(replay) {
return {
events: structuredClone(
events.filter(
(event) =>
event.sourceInstanceId === replay.sourceInstanceId &&
event.sourceSeq > replay.afterSourceSeq,
),
),
highestContiguousSourceSeq: highestContiguous(events),
};
},
completeRun,
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
resolveMissingResult: async () => result,
});
if (failInitialAppend) {
await expect(execute()).rejects.toBe(appendFailure);
expect(events).toEqual([]);
expect(startTurn).not.toHaveBeenCalled();
expect(completeRun).not.toHaveBeenCalled();
expect(close).toHaveBeenCalledOnce();
expect(prematureDispositionCheckpoint).toBe(false);
expect(durableCheckpoint.activeTurnId).toBe(checkpoint.activeTurnId);
expect(durableCheckpoint.terminalTurns).toEqual(
checkpoint.terminalTurns,
);
expect(durableCheckpoint.identity).toEqual(checkpoint.identity);
}
await expect(execute()).resolves.toMatchObject({
result,
turnId: "turn-disposition",
terminal: {
turnTerminalState:
terminalType === "turn.completed" ? "completed" : "interrupted",
runTerminalState:
terminalType === "turn.completed" ? "succeeded" : "cancelled",
},
});
expect(recoveryCheckpoints).toHaveLength(failInitialAppend ? 2 : 1);
for (const recoveredCheckpoint of recoveryCheckpoints) {
expect(recoveredCheckpoint.activeTurnId).toBe(checkpoint.activeTurnId);
expect(recoveredCheckpoint.terminalTurns).toEqual(
checkpoint.terminalTurns,
);
expect(recoveredCheckpoint.identity).toEqual(checkpoint.identity);
}
expect(startTurn).not.toHaveBeenCalled();
expect(completeRun).toHaveBeenCalledOnce();
expect(prematureDispositionCheckpoint).toBe(false);
expect(events.map((event) => event.eventType)).toEqual([
terminalType,
"run.result.accepted",
"run.terminal",
]);
},
);
it("resolves a proposal-less durable disposition terminal through control-plane policy", async () => {
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "1",
activeTurnId: null,
terminalTurns: [{ turnId: "turn-work", fingerprint: "work-terminal" }],
dispositionOnlyRecoveryConsumed: true,
dispositionOnlyRecoveryTurnId: "turn-disposition",
pendingRuntimeRequests: [],
lineage: [],
};
const terminalEvent: PrpEvent = {
schema: "paperclip.prp.event.v1",
sourceEventId: "runner-recovery:run-native:4",
sourceSeq: 4,
sourceInstanceId: "runner-recovery",
sourceKind: "provider",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-disposition",
eventType: "turn.completed",
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:01.000Z",
payload: {},
};
const resultProposalEvent: PrpEvent = {
...terminalEvent,
sourceEventId: "runner-recovery:run-native:3",
sourceSeq: 3,
eventType: "run.result.proposed",
payload: result,
};
const originalTaskTerminal: PrpEvent = {
...terminalEvent,
sourceEventId: "runner-recovery:run-native:2",
sourceSeq: 2,
turnId: "turn-work",
};
const originalTaskProposal: PrpEvent = {
...resultProposalEvent,
sourceEventId: "runner-recovery:run-native:1",
sourceSeq: 1,
turnId: "turn-work",
};
const startTurn = vi.fn(async () => ({ turnId: "unexpected-turn" }));
let recoveredSubmissionOwned = true;
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield structuredClone(terminalEvent);
},
startTurn,
async result() {
return null;
},
async snapshot() {
return {
...structuredClone(checkpoint),
dispositionOnlyRecoveryConsumed: recoveredSubmissionOwned,
terminalTurns: recoveredSubmissionOwned
? [
...structuredClone(checkpoint.terminalTurns ?? []),
{
turnId: "turn-disposition",
fingerprint: "disposition-terminal",
},
]
: structuredClone(checkpoint.terminalTurns),
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("must recover the provider session");
},
async recoverSession() {
return { recovered: true, session };
},
};
const bySource = new Map<string, PrpEvent[]>([
[
"runner-recovery",
[
structuredClone(originalTaskProposal),
structuredClone(originalTaskTerminal),
],
],
]);
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent(event) {
const list = bySource.get(event.sourceInstanceId) ?? [];
list.push(structuredClone(event));
bySource.set(event.sourceInstanceId, list);
return {
cursor: list.length,
highestContiguousSourceSeq: highestContiguous(list),
disposition: "committed",
};
},
async replayEvents(replay) {
const list = bySource.get(replay.sourceInstanceId) ?? [];
return {
events: structuredClone(
list.filter((event) => event.sourceSeq > replay.afterSourceSeq),
),
highestContiguousSourceSeq: highestContiguous(list),
};
},
async completeRun() {},
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
resolveMissingResult: async ({ terminalEvent: replayed }) => {
expect(replayed).toEqual(terminalEvent);
return result;
},
});
await expect(execute()).resolves.toMatchObject({
result,
turnId: "turn-disposition",
});
bySource.set("runner-recovery", [
structuredClone(originalTaskProposal),
structuredClone(originalTaskTerminal),
structuredClone(resultProposalEvent),
structuredClone(terminalEvent),
]);
// Provider recovery may clear a legacy pre-acceptance marker when thread
// history has no matching turn. Durable replay remains authoritative and
// must still prevent a duplicate disposition submission.
recoveredSubmissionOwned = false;
await expect(execute()).resolves.toMatchObject({
result,
turnId: "turn-disposition",
});
expect(startTurn).not.toHaveBeenCalled();
expect(bySource.get("runner-recovery")).toEqual([
originalTaskProposal,
originalTaskTerminal,
resultProposalEvent,
terminalEvent,
]);
expect(
bySource.get("control-recovery")?.map((event) => event.eventType),
).toEqual(["run.result.accepted", "run.terminal"]);
});
it("resolves a checkpointed result-less disposition without resubmitting when its terminal event is missing", async () => {
const workProposal: PrpEvent = {
...runnerEvent(1, "run.result.proposed", result),
turnId: "turn-work",
};
const workTerminal: PrpEvent = {
...runnerEvent(2, "turn.completed"),
turnId: "turn-work",
};
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "3",
activeTurnId: null,
terminalTurns: [
{ turnId: "turn-work", fingerprint: "work-terminal" },
{ turnId: "turn-disposition", fingerprint: "disposition-terminal" },
],
dispositionOnlyRecoveryConsumed: true,
dispositionOnlyRecoveryTurnId: "turn-disposition",
pendingRuntimeRequests: [],
lineage: [],
};
const recoveredCheckpoint: PersistedNativeSession =
structuredClone(checkpoint);
const startTurn = vi.fn(async () => ({ turnId: "unexpected-turn" }));
const events = vi.fn(() =>
(async function* () {
throw new Error("checkpoint fallback must not consume provider events");
})(),
);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
events,
startTurn,
async result() {
return null;
},
async snapshot() {
return structuredClone(recoveredCheckpoint);
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("must recover the provider session");
},
async recoverSession() {
return { recovered: true, session };
},
};
const bySource = new Map<string, PrpEvent[]>([
["runner-recovery", [workProposal, workTerminal]],
]);
const completeRun = vi.fn(async () => undefined);
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession() {},
async appendEvent(event) {
const list = bySource.get(event.sourceInstanceId) ?? [];
list.push(structuredClone(event));
bySource.set(event.sourceInstanceId, list);
return {
cursor: list.length,
highestContiguousSourceSeq: highestContiguous(list),
disposition: "committed",
};
},
async replayEvents(replay) {
const list = bySource.get(replay.sourceInstanceId) ?? [];
return {
events: structuredClone(
list.filter((event) => event.sourceSeq > replay.afterSourceSeq),
),
highestContiguousSourceSeq: highestContiguous(list),
};
},
completeRun,
};
const execute = () =>
executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
resolveMissingResult: async ({ turnId, terminalEvent }) => {
expect(turnId).toBe("turn-disposition");
expect(terminalEvent).toMatchObject({
sourceInstanceId: "control-recovery",
sourceKind: "control_plane",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-disposition",
eventType: "turn.completed",
payload: {
recovery: "checkpointed_resultless_disposition",
terminalFingerprint: "disposition-terminal",
},
});
return result;
},
});
await expect(execute()).resolves.toMatchObject({
result,
turnId: "turn-disposition",
});
expect(startTurn).not.toHaveBeenCalled();
expect(events).not.toHaveBeenCalled();
expect(completeRun).toHaveBeenCalledOnce();
expect(bySource.get("runner-recovery")).toEqual([
workProposal,
workTerminal,
]);
expect(
bySource.get("control-recovery")?.map((event) => event.eventType),
).toEqual(["run.result.accepted", "run.terminal"]);
recoveredCheckpoint.terminalTurns![1]!.fingerprint = "conflicting-terminal";
await expect(execute()).rejects.toThrow(
"native_disposition_recovery_checkpoint_conflict",
);
expect(startTurn).not.toHaveBeenCalled();
expect(events).not.toHaveBeenCalled();
expect(completeRun).toHaveBeenCalledOnce();
});
it("keeps a reconstructed semantic result on its matched terminal turn", async () => {
const semanticFingerprint = canonicalTestJson(result);
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-recovery",
identity,
providerSessionId: "provider-recovery",
cursor: "4",
semanticResult: result,
terminal,
activeTurnId: null,
terminalTurns: [
{
turnId: "turn-with-result",
fingerprint: JSON.stringify({
status: "completed",
semanticResult: semanticFingerprint,
}),
},
{
turnId: "turn-later-failed",
fingerprint: JSON.stringify({ status: "failed" }),
},
],
pendingRuntimeRequests: [],
lineage: [],
};
const events = [
{
...controlEvent(1, "run.result.accepted", { result }),
turnId: "turn-with-result",
},
];
const checkpoints: PersistedNativeSession[] = [];
const completeRun = vi.fn(async () => undefined);
const startTurn = vi.fn(async () => ({ turnId: "unexpected-turn" }));
const openSession = vi.fn(async () => {
throw new Error(
"a recovered run must not open a second provider session",
);
});
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {},
startTurn,
async result() {
return { result, terminal, turnId: "turn-recovery" };
},
async snapshot() {
return structuredClone(checkpoint);
},
async close() {},
};
const recoverSession = vi.fn(async () => ({ recovered: true, session }));
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
recoverSession,
};
const port: ControlPlanePort = {
async openRun() {},
async loadSessionCheckpoint() {
return structuredClone(checkpoint);
},
async checkpointSession(snapshot) {
checkpoints.push(structuredClone(snapshot));
},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents(replay) {
const replayed = events.filter(
(event) => event.sourceSeq > replay.afterSourceSeq,
);
return {
events: structuredClone(replayed),
highestContiguousSourceSeq: highestContiguous(events),
};
},
completeRun,
};
const completed = await executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
expect(openSession).not.toHaveBeenCalled();
expect(recoverSession).toHaveBeenCalledOnce();
expect(startTurn).not.toHaveBeenCalled();
expect(events.map((event) => event.eventType)).toEqual([
"run.result.accepted",
"run.terminal",
]);
expect(events.map((event) => event.sourceSeq)).toEqual([1, 2]);
expect(events.map((event) => event.turnId)).toEqual([
"turn-with-result",
"turn-with-result",
]);
expect(completeRun).toHaveBeenCalledOnce();
expect(completeRun).toHaveBeenCalledWith(
expect.objectContaining({ turnId: "turn-with-result" }),
expect.anything(),
);
expect(completed).toMatchObject({
nativeEventCount: 1,
highestContiguousSourceSeq: 2,
});
expect(checkpoints.at(-1)).toMatchObject({
semanticResult: result,
terminal,
});
});
it("accepts a control-plane governed wait when a completed turn omitted its semantic result", async () => {
const terminalEvent: PrpEvent = {
schema: "paperclip.prp.event.v1",
sourceEventId: "provider-recovery:1",
sourceSeq: 1,
sourceInstanceId: "provider-recovery",
sourceKind: "provider",
runId: identity.runId,
normalizedSessionId: identity.sessionId,
turnId: "turn-waiting",
eventType: "turn.completed",
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T00:00:00.000Z",
payload: {},
};
const yielded: PrpStructuredRunResult = {
schema: "paperclip.run_result.v1",
reportedWorkDisposition: "yielded",
summary: "Waiting for the requested response.",
completionClaim: {
contractRevision: "1",
objectiveSatisfied: false,
criteria: [
{
criterionId: "objective",
status: "unknown",
evidenceRefs: ["interaction:pending"],
},
],
remainingWork: [
{ description: "Resume after the response.", blocksCompletion: true },
],
},
evidence: [{ ref: "interaction:pending" }],
verification: [],
attentionRequests: [],
artifacts: [],
continuation: {
kind: "response_wake",
summary: "Resume from the answer.",
idempotencyKey: "interaction-response:pending",
},
};
const events: PrpEvent[] = [];
const completeRun = vi.fn(async () => undefined);
const resolveMissingResult = vi.fn(async () => yielded);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield terminalEvent;
},
async startTurn() {
return { turnId: "turn-waiting" };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-waiting",
cursor: "1",
activeTurnId: "turn-waiting",
pendingRuntimeRequests: [],
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "governed-wait-backend",
version: "1",
capabilities: {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents(replay) {
const replayed = events.filter(
(event) =>
event.sourceInstanceId === replay.sourceInstanceId &&
event.sourceSeq > replay.afterSourceSeq,
);
return {
events: structuredClone(replayed),
highestContiguousSourceSeq: highestContiguous(replayed),
};
},
completeRun,
};
const completed = await executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
resolveMissingResult,
});
expect(resolveMissingResult).toHaveBeenCalledWith({
turnId: "turn-waiting",
terminalEvent,
});
expect(completed).toMatchObject({
result: yielded,
terminal: {
runTerminalState: "succeeded",
reportedWorkDisposition: "yielded",
},
turnId: "turn-waiting",
});
expect(completeRun).toHaveBeenCalledWith(
expect.objectContaining({ result: yielded }),
{ signal: expect.any(AbortSignal) },
);
expect(events.map((event) => event.eventType)).toEqual([
"turn.completed",
"run.result.accepted",
"run.terminal",
]);
});
it("parks a provider turn immediately after a durable governed wait appears", async () => {
const yielded: PrpStructuredRunResult = {
schema: "paperclip.run_result.v1",
reportedWorkDisposition: "yielded",
summary: "Waiting for the requested response.",
completionClaim: {
contractRevision: "1",
objectiveSatisfied: false,
criteria: [
{
criterionId: "objective",
status: "unknown",
evidenceRefs: ["interaction:pending"],
},
],
remainingWork: [
{ description: "Resume after the response.", blocksCompletion: true },
],
},
evidence: [{ ref: "interaction:pending" }],
verification: [],
attentionRequests: [],
artifacts: [],
continuation: {
kind: "response_wake",
summary: "Resume from the answer.",
idempotencyKey: "interaction-response:pending",
},
};
const itemCompleted: PrpEvent = {
...controlEvent(1, "item.completed", {
kind: "dynamicToolCall",
item: { id: "ask-1", name: "ask_user_questions" },
}),
sourceEventId: "provider-recovery:1",
sourceInstanceId: "provider-recovery",
sourceKind: "provider",
turnId: "turn-waiting",
};
const turnInterrupted: PrpEvent = {
...controlEvent(2, "turn.interrupted", { reason: "governed_wait" }),
sourceEventId: "provider-recovery:2",
sourceInstanceId: "provider-recovery",
sourceKind: "provider",
turnId: "turn-waiting",
};
let releaseCancelled!: () => void;
const cancelled = new Promise<void>((resolve) => {
releaseCancelled = resolve;
});
const cancel = vi.fn(() => {
releaseCancelled();
return { cleanup: Promise.resolve() };
});
const events: PrpEvent[] = [];
const completeRun = vi.fn(async () => undefined);
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield itemCompleted;
await cancelled;
yield turnInterrupted;
},
async startTurn() {
return { turnId: "turn-waiting" };
},
cancel,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-waiting",
cursor: "2",
activeTurnId: "turn-waiting",
pendingRuntimeRequests: [],
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "governed-wait-backend",
version: "1",
capabilities: {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(events),
disposition: "committed",
};
},
async replayEvents(replay) {
const replayed = events.filter(
(event) =>
event.sourceInstanceId === replay.sourceInstanceId &&
event.sourceSeq > replay.afterSourceSeq,
);
return {
events: structuredClone(replayed),
highestContiguousSourceSeq: highestContiguous(replayed),
};
},
completeRun,
};
const completed = await executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
resolveGovernedWait: ({ event }) =>
event.eventType === "item.completed" ? yielded : null,
});
expect(cancel).toHaveBeenCalledOnce();
expect(completed).toMatchObject({
result: yielded,
terminal: {
turnTerminalState: "completed",
runTerminalState: "succeeded",
reportedWorkDisposition: "yielded",
},
});
expect(events.map((event) => event.eventType)).toEqual([
"item.completed",
"run.result.accepted",
"run.terminal",
]);
});
it("hands a committed structured input to the durable wait after its live window", async () => {
vi.useFakeTimers();
try {
const questionSet = {
schema: "paperclip.question_set.v1" as const,
questions: [
{
id: "region",
prompt: "Which region?",
required: true,
answerMode: "single_select" as const,
options: [
{ id: "us", label: "US" },
{ id: "eu", label: "Europe" },
],
},
],
};
const request = {
schema: "paperclip.runtime_request.v2",
requestKind: "runtime",
requestId: "input-1",
type: "input",
status: "pending",
prompt: "Which region?",
input: questionSet,
origin: { adapter: "mock" },
turnId: "turn-waiting",
itemId: "input-1",
};
const created = {
...runnerEvent(1, "runtime_request.created", { request }),
turnId: "turn-waiting",
};
const expired = {
...runnerEvent(2, "runtime_request.expired", {
requestId: "input-1",
requestKind: "runtime",
turnId: "turn-waiting",
itemId: "input-1",
reason: "durable_handoff",
replayAllowed: false,
requestType: "input",
request,
}),
turnId: "turn-waiting",
};
const interrupted = {
...runnerEvent(3, "turn.interrupted", { reason: "governed_wait" }),
turnId: "turn-waiting",
};
let releaseHandoff!: () => void;
const handedOff = new Promise<void>((resolve) => {
releaseHandoff = resolve;
});
let releaseCancelled!: () => void;
const cancelled = new Promise<void>((resolve) => {
releaseCancelled = resolve;
});
let releaseCreated!: () => void;
const createdCommitted = new Promise<void>((resolve) => {
releaseCreated = resolve;
});
const handoffRuntimeRequest = vi.fn(() => {
releaseHandoff();
return { result: "handed_off" as const, cleanup: Promise.resolve() };
});
const cancel = vi.fn(() => {
releaseCancelled();
return { cleanup: Promise.resolve() };
});
const events: PrpEvent[] = [];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
runtimeRequestHandoff: true,
};
},
async *events() {
yield created;
await handedOff;
yield expired;
await cancelled;
yield interrupted;
},
async startTurn() {
return { turnId: "turn-waiting" };
},
handoffRuntimeRequest,
cancel,
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-waiting",
cursor: "3",
activeTurnId: "turn-waiting",
pendingRuntimeRequests: [],
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "runtime-input-wait-backend",
version: "1",
capabilities: {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
runtimeRequestHandoff: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
events.push(structuredClone(event as PrpEvent));
if (event.eventType === "runtime_request.created") releaseCreated();
const sourceEvents = events.filter(
(candidate) =>
candidate.sourceInstanceId === event.sourceInstanceId,
);
return {
cursor: events.length,
highestContiguousSourceSeq: highestContiguous(sourceEvents),
disposition: "committed",
};
},
async replayEvents(replay) {
const replayed = events.filter(
(event) =>
event.sourceInstanceId === replay.sourceInstanceId &&
event.sourceSeq > replay.afterSourceSeq,
);
return {
events: structuredClone(replayed),
highestContiguousSourceSeq: highestContiguous(replayed),
};
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
runtimeInputLiveWindowMs: 120,
resolveGovernedWait: ({ event }) =>
event.eventType === "runtime_request.expired" ? yieldedResult : null,
});
await createdCommitted;
expect(handoffRuntimeRequest).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(119);
expect(handoffRuntimeRequest).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1);
await expect(execution).resolves.toMatchObject({ result: yieldedResult });
expect(handoffRuntimeRequest).toHaveBeenCalledWith({
requestId: "input-1",
turnId: "turn-waiting",
reason: "durable_handoff",
signal: expect.any(AbortSignal),
});
expect(cancel).toHaveBeenCalledOnce();
expect(events.map((event) => event.eventType)).toContain(
"runtime_request.expired",
);
} finally {
vi.useRealTimers();
}
});
it("aborts and bounds a durable handoff that never settles", async () => {
const request = {
schema: "paperclip.runtime_request.v2",
requestKind: "runtime",
requestId: "input-stalled",
type: "input",
status: "pending",
prompt: "Which region?",
input: {
schema: "paperclip.question_set.v1",
questions: [
{
id: "region",
prompt: "Which region?",
required: true,
answerMode: "text",
},
],
},
origin: { adapter: "mock" },
turnId: "turn-stalled",
itemId: "input-stalled",
};
const created = {
...runnerEvent(1, "runtime_request.created", { request }),
turnId: "turn-stalled",
};
let releaseEvents = () => {};
const eventsReleased = new Promise<void>((resolve) => {
releaseEvents = resolve;
});
let markHandoffStarted = () => {};
const handoffStarted = new Promise<void>((resolve) => {
markHandoffStarted = resolve;
});
let handoffSignal: AbortSignal | undefined;
let releaseHandoff = () => {};
const close = vi.fn(async () => releaseEvents());
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
runtimeRequestHandoff: true,
};
},
async *events() {
yield created;
await eventsReleased;
},
async startTurn() {
return { turnId: "turn-stalled" };
},
handoffRuntimeRequest(input) {
handoffSignal = input.signal;
markHandoffStarted();
return {
result: "handed_off",
cleanup: new Promise<void>((resolve) => {
releaseHandoff = resolve;
}),
};
},
cancel() {
releaseEvents();
return { cleanup: Promise.resolve() };
},
async result() {
return null;
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-stalled",
activeTurnId: "turn-stalled",
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "stalled-handoff-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
async completeRun() {},
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
runtimeInputLiveWindowMs: 1,
timeoutMs: 25,
keepSessionOpen: true,
});
await handoffStarted;
await vi.waitFor(() => expect(close).toHaveBeenCalledOnce());
expect(handoffSignal?.aborted).toBe(true);
await expect(execution).rejects.toThrow("native session timed out");
releaseHandoff();
expect(close).toHaveBeenCalledOnce();
});
it("preserves terminal success while iterator teardown remains pending", async () => {
vi.useFakeTimers();
try {
let releaseTeardown = () => {};
const teardownStarted = vi.fn();
const close = vi.fn(async () => undefined);
const readResult = vi.fn(async () => ({
result,
terminal,
turnId: "turn-terminal",
}));
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
try {
yield {
...runnerEvent(1, "turn.completed"),
turnId: "turn-terminal",
};
} finally {
teardownStarted();
await new Promise<void>((resolve) => {
releaseTeardown = resolve;
});
}
},
async startTurn() {
return { turnId: "turn-terminal" };
},
result: readResult,
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-terminal",
cursor: "1",
activeTurnId: null,
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "slow-teardown-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const completeRun = vi.fn(async () => undefined);
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent() {
return {
cursor: 1,
highestContiguousSourceSeq: 1,
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
completeRun,
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
});
await vi.waitFor(() => expect(teardownStarted).toHaveBeenCalledOnce());
await vi.advanceTimersByTimeAsync(100);
await expect(execution).resolves.toMatchObject({ result });
expect(readResult).toHaveBeenCalledOnce();
expect(completeRun).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
releaseTeardown();
await Promise.resolve();
} finally {
vi.useRealTimers();
}
});
it("preserves terminal success while quarantining stalled handoff cleanup", async () => {
const request = {
schema: "paperclip.runtime_request.v2",
requestKind: "runtime",
requestId: "input-terminal",
type: "input",
status: "pending",
prompt: "Which region?",
input: {
schema: "paperclip.question_set.v1",
questions: [
{
id: "region",
prompt: "Which region?",
required: true,
answerMode: "text",
},
],
},
origin: { adapter: "mock" },
turnId: "turn-terminal",
itemId: "input-terminal",
};
let markHandoffStarted = () => {};
const handoffStarted = new Promise<void>((resolve) => {
markHandoffStarted = resolve;
});
let releaseHandoff = () => {};
let handoffSignal: AbortSignal | undefined;
const close = vi.fn(async () => undefined);
const onSession = vi.fn();
const completeRun = vi.fn(async () => undefined);
const readResult = vi.fn(async () => ({
result,
terminal,
turnId: "turn-terminal",
}));
const providerEvents = [
{
...runnerEvent(1, "runtime_request.created", { request }),
turnId: "turn-terminal",
},
{ ...runnerEvent(2, "turn.completed"), turnId: "turn-terminal" },
];
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
runtimeRequestHandoff: true,
};
},
async *events() {
yield providerEvents[0]!;
await handoffStarted;
yield providerEvents[1]!;
},
async startTurn() {
return { turnId: "turn-terminal" };
},
handoffRuntimeRequest(input) {
handoffSignal = input.signal;
markHandoffStarted();
return {
result: "handed_off",
cleanup: new Promise<void>((resolve) => {
releaseHandoff = resolve;
}),
};
},
result: readResult,
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-terminal",
cursor: "2",
activeTurnId: null,
};
},
close,
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "terminal-handoff-backend",
version: "1",
capabilities: await session.capabilities(),
};
},
async openSession() {
return session;
},
};
const appended: PrpEvent[] = [];
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
appended.push(structuredClone(event as PrpEvent));
return {
cursor: appended.length,
highestContiguousSourceSeq: highestContiguous(appended),
disposition: "committed",
};
},
async replayEvents() {
return { events: [], highestContiguousSourceSeq: 0 };
},
completeRun,
};
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
runtimeInputLiveWindowMs: 1,
keepSessionOpen: true,
onSession,
});
await handoffStarted;
await vi.waitFor(() => expect(handoffSignal?.aborted).toBe(true));
await expect(execution).resolves.toMatchObject({ result });
expect(readResult).toHaveBeenCalledOnce();
expect(completeRun).toHaveBeenCalledOnce();
expect(close).toHaveBeenCalledOnce();
expect(onSession).toHaveBeenLastCalledWith(null);
releaseHandoff();
await Promise.resolve();
expect(close).toHaveBeenCalledOnce();
expect(completeRun).toHaveBeenCalledOnce();
});
it("keeps a settling structured input in the original turn while its append crosses expiry", async () => {
const request = {
schema: "paperclip.runtime_request.v2",
requestKind: "runtime",
requestId: "input-live",
type: "input",
status: "pending",
prompt: "Which region?",
input: {
schema: "paperclip.question_set.v1",
questions: [
{
id: "region",
prompt: "Which region?",
required: true,
answerMode: "text",
},
],
},
origin: { adapter: "mock" },
turnId: "turn-live",
itemId: "input-live",
};
const providerEvents = [
{
...runnerEvent(1, "runtime_request.created", { request }),
turnId: "turn-live",
},
{
...runnerEvent(2, "runtime_request.resolved", {
requestId: "input-live",
requestKind: "user_input",
turnId: "turn-live",
itemId: "input-live",
action: "submit",
requestType: "input",
}),
turnId: "turn-live",
},
{ ...runnerEvent(3, "turn.completed"), turnId: "turn-live" },
];
const appended: PrpEvent[] = [];
let markSettlementAppendStarted!: () => void;
const settlementAppendStarted = new Promise<void>((resolve) => {
markSettlementAppendStarted = resolve;
});
let releaseSettlementAppend!: () => void;
const settlementAppendReleased = new Promise<void>((resolve) => {
releaseSettlementAppend = resolve;
});
const handoffRuntimeRequest = vi.fn(() => ({
result: "handed_off" as const,
cleanup: Promise.resolve(),
}));
const session: NativeSession = {
identity: () => identity,
async capabilities() {
return {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
runtimeRequestHandoff: true,
};
},
async *events() {
yield* providerEvents;
},
async startTurn() {
return { turnId: "turn-live" };
},
handoffRuntimeRequest,
async result() {
return { result, terminal, turnId: "turn-live" };
},
async snapshot() {
return {
backendKind: "mock",
sessionId: identity.sessionId,
identity,
providerSessionId: "provider-live",
cursor: "3",
activeTurnId: null,
pendingRuntimeRequests: [],
lineage: [],
};
},
async close() {},
};
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "runtime-input-live-backend",
version: "1",
capabilities: {
resume: false,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
runtimeRequestHandoff: true,
},
};
},
async openSession() {
return session;
},
};
const port: ControlPlanePort = {
async openRun() {},
async checkpointSession() {},
async appendEvent(event) {
if (event.eventType === "runtime_request.resolved") {
markSettlementAppendStarted();
await settlementAppendReleased;
}
appended.push(structuredClone(event as PrpEvent));
const sourceEvents = appended.filter(
(candidate) => candidate.sourceInstanceId === event.sourceInstanceId,
);
return {
cursor: appended.length,
highestContiguousSourceSeq: highestContiguous(sourceEvents),
disposition: "committed",
};
},
async replayEvents(replay) {
const replayed = appended.filter(
(event) =>
event.sourceInstanceId === replay.sourceInstanceId &&
event.sourceSeq > replay.afterSourceSeq,
);
return {
events: structuredClone(replayed),
highestContiguousSourceSeq: highestContiguous(replayed),
};
},
async completeRun() {},
};
const originalSetTimeout = globalThis.setTimeout;
let queuedHandoffCallback: (() => void) | null = null;
const timeoutSpy = vi.spyOn(globalThis, "setTimeout").mockImplementation(((
callback,
delay,
...args
) => {
if (delay === 123_456) {
queuedHandoffCallback = () => callback(...args);
const handle = originalSetTimeout(() => undefined, 60_000);
handle.unref?.();
return handle;
}
return originalSetTimeout(callback, delay, ...args);
}) as typeof setTimeout);
try {
const execution = executeNativeSession({
input,
backend,
controlPlane: port,
runnerInstanceId: "runner-recovery",
controlPlaneInstanceId: "control-recovery",
runtimeInputLiveWindowMs: 123_456,
});
await settlementAppendStarted;
expect(queuedHandoffCallback).not.toBeNull();
queuedHandoffCallback?.();
await Promise.resolve();
expect(handoffRuntimeRequest).not.toHaveBeenCalled();
releaseSettlementAppend();
await expect(execution).resolves.toMatchObject({ result });
queuedHandoffCallback?.();
await Promise.resolve();
expect(handoffRuntimeRequest).not.toHaveBeenCalled();
expect(
appended.some((event) => event.eventType === "runtime_request.expired"),
).toBe(false);
} finally {
releaseSettlementAppend();
timeoutSpy.mockRestore();
}
});
});