diff --git a/packages/sdk/src/client.ts b/packages/sdk/src/client.ts index 60dae0694..43e3d9196 100644 --- a/packages/sdk/src/client.ts +++ b/packages/sdk/src/client.ts @@ -810,7 +810,7 @@ async function readDurableRunEvents( request: (method: string, params: Record) => Promise, reference: { sessionId: string; runId: string; projectKey?: string }, lastSeq: number, -): Promise }>> { +): Promise }> | undefined> { const events: Array<{ event?: Record }> = []; let afterSeq = 0; while (afterSeq < lastSeq) { @@ -820,13 +820,20 @@ async function readDurableRunEvents( ...(reference.projectKey !== undefined ? { projectKey: reference.projectKey } : {}), afterSeq, limit: 500, - }) as { events?: Array<{ event?: Record; seq?: number }>; nextSeq?: number }; + }) as { events?: Array<{ event?: Record; 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; } @@ -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 { diff --git a/packages/sdk/test/transport.test.ts b/packages/sdk/test/transport.test.ts index 3c2d07be5..6191dfcb5 100644 --- a/packages/sdk/test/transport.test.ts +++ b/packages/sdk/test/transport.test.ts @@ -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 }> = 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 = { + "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 = [];