Skip to content

Commit 99bf1bc

Browse files
committed
feat(sdk): add iterator resume on transient errors
1 parent ad1f86f commit 99bf1bc

3 files changed

Lines changed: 118 additions & 12 deletions

File tree

docs/GETTING_STARTED.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -96,7 +96,9 @@ try {
9696
for await (const event of client.iterateJobEvents({
9797
jobId: "22222222-2222-2222-2222-222222222222",
9898
limit: 10,
99-
maxEvents: 50
99+
maxEvents: 50,
100+
resumeOnTransientError: true,
101+
maxResumeAttempts: 3
100102
})) {
101103
console.log(event.kind, event.at);
102104
}

typescript/packages/sdk/src/index.ts

Lines changed: 39 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,10 @@ export type GetJobEventsAllPagesRequest = Omit<GetJobEventsRequest, "before"> &
9999
onPage?: (page: JobEventsResponse, pageIndex: number) => void;
100100
};
101101

102-
export type IterateJobEventsRequest = GetJobEventsAllPagesRequest;
102+
export type IterateJobEventsRequest = GetJobEventsAllPagesRequest & {
103+
resumeOnTransientError?: boolean;
104+
maxResumeAttempts?: number;
105+
};
103106

104107
export type JobEvent = {
105108
at: string;
@@ -285,20 +288,45 @@ export class SpatiadClient {
285288
return;
286289
}
287290

291+
const resumeOnTransientError = request.resumeOnTransientError ?? false;
292+
const maxResumeAttempts = Math.max(0, request.maxResumeAttempts ?? 3);
288293
let cursor: string | undefined;
289294
let yielded = 0;
295+
let resumeAttempts = 0;
290296

291297
for (let page = 0; page < maxPages; page += 1) {
292-
const current = await this.getJobEvents({
293-
jobId: request.jobId,
294-
limit: request.limit,
295-
cursor,
296-
kinds: request.kinds,
297-
dispatcherToken: request.dispatcherToken,
298-
dispatcherAuthMode: request.dispatcherAuthMode,
299-
signal: request.signal,
300-
retry: request.retry
301-
});
298+
let current: JobEventsResponse;
299+
try {
300+
current = await this.getJobEvents({
301+
jobId: request.jobId,
302+
limit: request.limit,
303+
cursor,
304+
kinds: request.kinds,
305+
dispatcherToken: request.dispatcherToken,
306+
dispatcherAuthMode: request.dispatcherAuthMode,
307+
signal: request.signal,
308+
retry: request.retry
309+
});
310+
} catch (error) {
311+
if (
312+
resumeOnTransientError
313+
&& error instanceof SpatiadApiError
314+
&& error.retryable
315+
&& resumeAttempts < maxResumeAttempts
316+
) {
317+
resumeAttempts += 1;
318+
const backoffBase = Math.max(0, request.retry?.backoffMs ?? 150);
319+
const backoffMax = Math.max(backoffBase, request.retry?.maxBackoffMs ?? 2000);
320+
const waitMs = Math.min(backoffMax, backoffBase * (2 ** (resumeAttempts - 1)));
321+
await waitWithSignal(waitMs, request.signal);
322+
page -= 1;
323+
continue;
324+
}
325+
326+
throw error;
327+
}
328+
329+
resumeAttempts = 0;
302330

303331
request.onPage?.(current, page);
304332

typescript/packages/sdk/test/client.test.js

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -531,3 +531,79 @@ test("iterateJobEvents respects maxEvents", async () => {
531531
globalThis.fetch = originalFetch;
532532
}
533533
});
534+
535+
test("iterateJobEvents can resume on transient api errors", async () => {
536+
const originalFetch = globalThis.fetch;
537+
let calls = 0;
538+
539+
globalThis.fetch = async () => {
540+
calls += 1;
541+
if (calls === 1) {
542+
return makeJsonResponse(503, {
543+
error: "service_unavailable",
544+
message: "temporary outage"
545+
});
546+
}
547+
548+
return makeJsonResponse(200, {
549+
job_id: "job-iter-recover",
550+
events: [
551+
{ at: "2026-03-20T10:00:00Z", kind: "offer_created", offer_id: "o1", driver_id: "d1", status: "pending" }
552+
],
553+
next_cursor: null,
554+
next_before_cursor: null
555+
});
556+
};
557+
558+
try {
559+
const client = new SpatiadClient("http://localhost:3000");
560+
const streamed = [];
561+
562+
for await (const event of client.iterateJobEvents({
563+
jobId: "job-iter-recover",
564+
resumeOnTransientError: true,
565+
maxResumeAttempts: 1,
566+
retry: { maxAttempts: 1, backoffMs: 1, retryOnStatuses: [503] }
567+
})) {
568+
streamed.push(event.kind);
569+
}
570+
571+
assert.deepEqual(streamed, ["offer_created"]);
572+
assert.equal(calls, 2);
573+
} finally {
574+
globalThis.fetch = originalFetch;
575+
}
576+
});
577+
578+
test("iterateJobEvents fails without resumeOnTransientError", async () => {
579+
const originalFetch = globalThis.fetch;
580+
581+
globalThis.fetch = async () =>
582+
makeJsonResponse(503, {
583+
error: "service_unavailable",
584+
message: "temporary outage"
585+
});
586+
587+
try {
588+
const client = new SpatiadClient("http://localhost:3000");
589+
590+
await assert.rejects(
591+
async () => {
592+
for await (const _event of client.iterateJobEvents({
593+
jobId: "job-iter-fail",
594+
retry: { maxAttempts: 1, retryOnStatuses: [503] }
595+
})) {
596+
// no-op
597+
}
598+
},
599+
(error) => {
600+
assert.ok(error instanceof SpatiadApiError);
601+
assert.equal(error.status, 503);
602+
assert.equal(error.retryable, true);
603+
return true;
604+
}
605+
);
606+
} finally {
607+
globalThis.fetch = originalFetch;
608+
}
609+
});

0 commit comments

Comments
 (0)