8305 lines
257 KiB
TypeScript
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();
|
|
}
|
|
});
|
|
});
|