Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 15 additions & 7 deletions packages/sdk/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -810,7 +810,7 @@ async function readDurableRunEvents(
request: (method: string, params: Record<string, unknown>) => Promise<unknown>,
reference: { sessionId: string; runId: string; projectKey?: string },
lastSeq: number,
): Promise<Array<{ event?: Record<string, unknown> }>> {
): Promise<Array<{ event?: Record<string, unknown> }> | undefined> {
const events: Array<{ event?: Record<string, unknown> }> = [];
let afterSeq = 0;
while (afterSeq < lastSeq) {
Expand All @@ -820,13 +820,20 @@ async function readDurableRunEvents(
...(reference.projectKey !== undefined ? { projectKey: reference.projectKey } : {}),
afterSeq,
limit: 500,
}) as { events?: Array<{ event?: Record<string, unknown>; seq?: number }>; nextSeq?: number };
}) as { events?: Array<{ event?: Record<string, unknown>; seq?: number }>; gap?: boolean };
if (result.gap) return undefined;
const page = result.events ?? [];
if (page.length === 0) break;
events.push(...page);
const nextSeq = result.nextSeq ?? page.at(-1)?.seq;
if (typeof nextSeq !== "number" || nextSeq <= afterSeq) break;
afterSeq = nextSeq;
const previousSeq = afterSeq;
for (const item of page) {
const seq = item.seq;
if (typeof seq !== "number" || !Number.isSafeInteger(seq) || seq <= 0) return undefined;
if (seq <= afterSeq) continue;
if (seq !== afterSeq + 1) return undefined;
events.push(item);
afterSeq = seq;
if (afterSeq >= lastSeq) break;
}
if (afterSeq === previousSeq) return undefined;
}
return events;
}
Expand Down Expand Up @@ -3489,6 +3496,7 @@ export function createPilotDeckClientWithTransportFactory(
let output: unknown;
if (Number.isSafeInteger(record.lastSeq) && Number(record.lastSeq) > 0) {
const events = await readDurableRunEvents(request, reference, Number(record.lastSeq));
if (!events) return { status: "result_unknown", recovery: { sessionId: reference.sessionId, runId: reference.runId } };
output = durableRunOutput(events);
}
return {
Expand Down
55 changes: 55 additions & 0 deletions packages/sdk/test/transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3830,6 +3830,61 @@ test("runs.result reconstructs durable output across paged event history", async
await client.close();
});

for (const scenario of [
"page duplicates", "cross-page duplicates", "nextSeq ahead",
"explicit gap", "missing sequence", "empty first page", "truncated tail",
"no progress", "missing seq field", "zero seq", "negative seq",
"fractional seq", "unsafe seq",
]) {
test(`runs.result handles durable history ${scenario}`, async () => {
class ReplayWebSocket extends DurableResultWebSocket {
override send(raw: string): void {
const frame = JSON.parse(raw);
if (frame.method !== "run_events") return super.send(raw);
DurableResultWebSocket.requests.push({ method: frame.method, params: frame.params });
const afterSeq = frame.params.afterSeq ?? 0;
let events: Array<{ seq?: number; event: Record<string, unknown> }> = DurableResultWebSocket.events
.filter((item) => item.seq > afterSeq).slice(0, 500);
let nextSeq = events.at(-1)?.seq;
if (scenario === "page duplicates") events = events.flatMap((item) => [item, item]);
if (scenario === "cross-page duplicates" && afterSeq > 0) events.unshift(DurableResultWebSocket.events[0]!);
if (scenario === "nextSeq ahead") nextSeq = 501;
if (scenario === "missing sequence" && afterSeq === 0) events.splice(1, 1);
if (scenario === "empty first page" || (scenario === "truncated tail" && afterSeq > 0)) events = [];
if (scenario === "no progress" && afterSeq > 0) {
events = [DurableResultWebSocket.events[0]!];
nextSeq = afterSeq;
}
if (afterSeq === 0) {
const invalidSeq: Record<string, number | undefined> = {
"missing seq field": undefined, "zero seq": 0, "negative seq": -1,
"fractional seq": 1.5, "unsafe seq": Number.MAX_SAFE_INTEGER + 1,
};
if (scenario in invalidSeq) events[0] = { ...events[0]!, seq: invalidSeq[scenario] };
}
this.emitForTest({ type: "response", id: frame.id, ok: true,
result: { events, nextSeq, ...(scenario === "explicit gap" ? { gap: true } : {}) } });
}
}
DurableResultWebSocket.requests = [];
(globalThis as any).WebSocket = ReplayWebSocket;
const client = createPilotDeckClient({ gatewayUrl: "ws://fake", authToken: "token" });
const reference = { sessionId: "session-existing", runId: "durable-completed" };
try {
const result = await client.runs.result(reference);
if (["page duplicates", "cross-page duplicates", "nextSeq ahead"].includes(scenario)) {
assert.deepEqual(result, { status: "completed", output: "hello world", usage: { inputTokens: 2 }, finishReason: "stop" });
assert.deepEqual(DurableResultWebSocket.requests.filter((frame) => frame.method === "run_events").map((frame) => frame.params.afterSeq), [0, 500]);
} else {
assert.deepEqual(result, { status: "result_unknown", recovery: reference });
}
assert.equal(DurableResultWebSocket.requests.some((frame) => frame.method === "submit_turn"), false);
} finally {
await client.close();
}
});
}

test("client exposes Gateway-authoritative session, run and resource facades", async () => {
FakeSdkClientWebSocket.connections = 0;
FakeSdkClientWebSocket.requests = [];
Expand Down
Loading