From 431bfa87c266863079799d1b3503ede7b00db87a Mon Sep 17 00:00:00 2001 From: Dylan Steele Date: Mon, 8 Jun 2026 21:02:17 -0400 Subject: [PATCH] Harden starred import resilience --- apps/web/src/lib/api.test.ts | 29 ++++++++++ apps/web/src/lib/api.ts | 22 ++++++- apps/web/src/lib/db.test.ts | 1 + apps/web/src/lib/import-messages.test.ts | 3 + apps/web/src/lib/import-messages.ts | 11 +++- apps/web/src/lib/import-pipeline.ts | 22 +++++-- apps/web/src/lib/import-retry.test.ts | 54 +++++++++++++++++ apps/web/src/lib/import-retry.ts | 74 ++++++++++++++++++++++++ docs/12-import-pipeline.md | 6 ++ packages/core/src/import-state.ts | 5 ++ packages/core/test/import-state.test.mjs | 3 + packages/shared/src/import.ts | 1 + 12 files changed, 223 insertions(+), 8 deletions(-) create mode 100644 apps/web/src/lib/api.test.ts create mode 100644 apps/web/src/lib/import-retry.test.ts create mode 100644 apps/web/src/lib/import-retry.ts diff --git a/apps/web/src/lib/api.test.ts b/apps/web/src/lib/api.test.ts new file mode 100644 index 0000000..af2b948 --- /dev/null +++ b/apps/web/src/lib/api.test.ts @@ -0,0 +1,29 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { WorkerApi, type WorkerApiError } from "./api"; + +describe("WorkerApi", () => { + afterEach(() => { + vi.unstubAllGlobals(); + }); + + it("preserves retry-after metadata from worker errors", async () => { + vi.stubGlobal( + "fetch", + vi.fn().mockResolvedValue( + new Response( + JSON.stringify({ + error: "Too many requests. Try again shortly.", + retry_after_seconds: 42, + }), + { status: 429 }, + ), + ), + ); + + await expect(new WorkerApi("https://worker.example").getSession()).rejects.toMatchObject({ + name: "WorkerApiError", + status: 429, + retryAfterSeconds: 42, + } satisfies Partial); + }); +}); diff --git a/apps/web/src/lib/api.ts b/apps/web/src/lib/api.ts index 61c9567..314c4b6 100644 --- a/apps/web/src/lib/api.ts +++ b/apps/web/src/lib/api.ts @@ -53,6 +53,7 @@ export class WorkerApiError extends Error { message: string, readonly status: number, readonly rateLimit: GitHubRateLimitSnapshot | null, + readonly retryAfterSeconds: number | null = null, ) { super(message); this.name = "WorkerApiError"; @@ -154,12 +155,30 @@ export class WorkerApi { typeof payload === "object" && payload && "rate_limit" in payload ? (payload.rate_limit as GitHubRateLimitSnapshot) : null; - throw new WorkerApiError(message, response.status, rateLimit); + throw new WorkerApiError( + message, + response.status, + rateLimit, + retryAfterSecondsFromResponse(response, payload), + ); } return payload as T; } } +function retryAfterSecondsFromResponse(response: Response, payload: unknown) { + const payloadRetryAfter = + typeof payload === "object" && payload && "retry_after_seconds" in payload + ? Number(payload.retry_after_seconds) + : NaN; + if (Number.isFinite(payloadRetryAfter) && payloadRetryAfter >= 0) return payloadRetryAfter; + + const headerRetryAfter = Number(response.headers.get("retry-after")); + if (Number.isFinite(headerRetryAfter) && headerRetryAfter >= 0) return headerRetryAfter; + + return null; +} + export function createImportEvent(): ImportEvent { return { id: crypto.randomUUID(), @@ -169,6 +188,7 @@ export function createImportEvent(): ImportEvent { pages: 0, repositories: 0, rate_limits: [], + retry_after_seconds: null, errors: [], }; } diff --git a/apps/web/src/lib/db.test.ts b/apps/web/src/lib/db.test.ts index ef04da2..14fdf43 100644 --- a/apps/web/src/lib/db.test.ts +++ b/apps/web/src/lib/db.test.ts @@ -191,6 +191,7 @@ function createImportEvent(id: string, startedAt: string): ImportEvent { pages: 1, repositories: 1, rate_limits: [], + retry_after_seconds: null, errors: [], }; } diff --git a/apps/web/src/lib/import-messages.test.ts b/apps/web/src/lib/import-messages.test.ts index 16a72d6..02635d9 100644 --- a/apps/web/src/lib/import-messages.test.ts +++ b/apps/web/src/lib/import-messages.test.ts @@ -44,6 +44,9 @@ describe("import worker messages", () => { expect(getImportTerminalText({ ...importedPage, status: "cancelled" })).toBe( "Import cancelled after 1 page(s) and 100 repositories.", ); + expect( + getImportTerminalText({ ...importedPage, status: "rate_limited", retry_after_seconds: 120 }), + ).toBe("Import paused by GitHub rate limits after 1 page(s). Try again in about 2 minute(s)."); expect(getImportTerminalText({ ...importedPage, status: "rate_limited" })).toBe( "Import paused by GitHub rate limits after 1 page(s). Try again later.", ); diff --git a/apps/web/src/lib/import-messages.ts b/apps/web/src/lib/import-messages.ts index bcbdddc..657e27d 100644 --- a/apps/web/src/lib/import-messages.ts +++ b/apps/web/src/lib/import-messages.ts @@ -20,7 +20,8 @@ export function getImportTerminalText(importRun: ImportRunState) { return `Import cancelled after ${importRun.pages} page(s) and ${importRun.repositories} repositories.`; } if (importRun.status === "rate_limited") { - return `Import paused by GitHub rate limits after ${importRun.pages} page(s). Try again later.`; + const retryText = formatRetryAfter(importRun.retry_after_seconds); + return `Import paused by GitHub rate limits after ${importRun.pages} page(s).${retryText}`; } if (importRun.status === "failed") { return importRun.errors[0] || "Import failed."; @@ -31,3 +32,11 @@ export function getImportTerminalText(importRun: ImportRunState) { export function sortObservedFieldNames(observedFieldNames: Iterable) { return Array.from(observedFieldNames).sort(); } + +function formatRetryAfter(retryAfterSeconds: number | null) { + if (retryAfterSeconds === null) return " Try again later."; + if (retryAfterSeconds < 60) return ` Try again in about ${retryAfterSeconds} seconds.`; + + const minutes = Math.ceil(retryAfterSeconds / 60); + return ` Try again in about ${minutes} minute(s).`; +} diff --git a/apps/web/src/lib/import-pipeline.ts b/apps/web/src/lib/import-pipeline.ts index ae102a2..cda3e70 100644 --- a/apps/web/src/lib/import-pipeline.ts +++ b/apps/web/src/lib/import-pipeline.ts @@ -8,13 +8,14 @@ import { rateLimitImport, recordImportPage, } from "@forage/core"; -import { WorkerApi, WorkerApiError } from "./api"; +import { type StarredPageResponse, WorkerApi, WorkerApiError } from "./api"; import { reconcileImportedRepositories, saveAnalysisResults, saveLocalLibraryProfile, saveRepositories, } from "./db"; +import { getRateLimitRetryAfterSeconds, runImportRequestWithRetry } from "./import-retry"; import type { ImportWorkerPhase, StartRepositoryImportInput } from "./import-worker"; export interface ImportPipelineInput extends StartRepositoryImportInput { @@ -47,18 +48,22 @@ export async function runRepositoryImportPipeline( let page: number | null = 1; while (page) { - onProgress({ importRun, phase: "importing", page, observedFieldNames }); - const result = await api.getStarredPage(page, 100, input.signal); + const currentPage: number = page; + onProgress({ importRun, phase: "importing", page: currentPage, observedFieldNames }); + const result = await runImportRequestWithRetry( + () => api.getStarredPage(currentPage, 100, input.signal), + { signal: input.signal }, + ); await saveRepositories(result.repositories); for (const repository of result.repositories) { importedRepositoryIds.add(repository.github_id); } - onProgress({ importRun, phase: "analyzing", page, observedFieldNames }); + onProgress({ importRun, phase: "analyzing", page: currentPage, observedFieldNames }); await saveAnalysisResults(analyzeRepositories(result.repositories)); importRun = recordImportPage(importRun, { - page, + page: currentPage, repositories: result.repositories.length, rate_limit: result.rate_limit, }); @@ -82,7 +87,12 @@ export async function runRepositoryImportPipeline( if (input.isCancelRequested() || isAbortError(error)) { importRun = cancelImport(importRun); } else if (isRateLimitError(error)) { - importRun = rateLimitImport(importRun, errorMessage, error.rateLimit); + importRun = rateLimitImport( + importRun, + errorMessage, + error.rateLimit, + getRateLimitRetryAfterSeconds(error), + ); } else { importRun = failImport(importRun, errorMessage); } diff --git a/apps/web/src/lib/import-retry.test.ts b/apps/web/src/lib/import-retry.test.ts new file mode 100644 index 0000000..c459014 --- /dev/null +++ b/apps/web/src/lib/import-retry.test.ts @@ -0,0 +1,54 @@ +import { describe, expect, it, vi } from "vitest"; +import { WorkerApiError } from "./api"; +import { + getImportRetryDelayMs, + getRateLimitRetryAfterSeconds, + runImportRequestWithRetry, + shouldRetryImportRequest, +} from "./import-retry"; + +const rateLimit = { + limit: "5000", + remaining: "0", + reset: "1780000000", + resource: "core", + used: "5000", +}; + +describe("import retry helpers", () => { + it("retries bounded transient failures", async () => { + const request = vi + .fn() + .mockRejectedValueOnce(new WorkerApiError("unavailable", 503, null)) + .mockResolvedValueOnce("ok"); + const sleep = vi.fn().mockResolvedValue(undefined); + + await expect( + runImportRequestWithRetry(request, { signal: new AbortController().signal, sleep }), + ).resolves.toBe("ok"); + + expect(request).toHaveBeenCalledTimes(2); + expect(sleep).toHaveBeenCalledWith(750, expect.any(AbortSignal)); + }); + + it("does not retry auth, validation, or rate-limit failures", () => { + expect(shouldRetryImportRequest(new WorkerApiError("auth", 401, null), 1)).toBe(false); + expect(shouldRetryImportRequest(new WorkerApiError("rate", 429, null, 60), 1)).toBe(false); + expect(shouldRetryImportRequest(new WorkerApiError("temporary", 502, null), 3)).toBe(false); + }); + + it("uses retry-after metadata before exponential fallback", () => { + expect(getImportRetryDelayMs(new WorkerApiError("temporary", 503, null, 2), 1)).toBe(2_000); + expect(getImportRetryDelayMs(new WorkerApiError("temporary", 503, null), 2)).toBe(1_500); + }); + + it("estimates rate-limit retry timing from GitHub reset metadata", () => { + expect( + getRateLimitRetryAfterSeconds( + new WorkerApiError("limited", 403, rateLimit), + 1_779_999_940_000, + ), + ).toBe(60); + expect(getRateLimitRetryAfterSeconds(new WorkerApiError("limited", 429, null, 10))).toBe(10); + }); +}); diff --git a/apps/web/src/lib/import-retry.ts b/apps/web/src/lib/import-retry.ts new file mode 100644 index 0000000..3706d99 --- /dev/null +++ b/apps/web/src/lib/import-retry.ts @@ -0,0 +1,74 @@ +import { WorkerApiError } from "./api"; + +export const importRetryPolicy = { + maxAttempts: 3, + baseDelayMs: 750, + maxDelayMs: 5_000, + retryableStatuses: new Set([408, 500, 502, 503, 504]), +}; + +interface RetryOptions { + signal: AbortSignal; + sleep?: (delayMs: number, signal: AbortSignal) => Promise; +} + +export async function runImportRequestWithRetry( + request: () => Promise, + { signal, sleep = sleepWithAbort }: RetryOptions, +) { + for (let attempt = 1; ; attempt += 1) { + throwIfAborted(signal); + + try { + return await request(); + } catch (error) { + if (!shouldRetryImportRequest(error, attempt)) throw error; + await sleep(getImportRetryDelayMs(error, attempt), signal); + } + } +} + +export function shouldRetryImportRequest(error: unknown, attempt: number) { + if (attempt >= importRetryPolicy.maxAttempts) return false; + if (error instanceof WorkerApiError) { + return importRetryPolicy.retryableStatuses.has(error.status); + } + return error instanceof TypeError; +} + +export function getImportRetryDelayMs(error: unknown, attempt: number) { + if (error instanceof WorkerApiError && error.retryAfterSeconds !== null) { + return Math.min(error.retryAfterSeconds * 1000, importRetryPolicy.maxDelayMs); + } + + const exponentialDelay = importRetryPolicy.baseDelayMs * 2 ** (attempt - 1); + return Math.min(exponentialDelay, importRetryPolicy.maxDelayMs); +} + +export function getRateLimitRetryAfterSeconds(error: WorkerApiError, nowMs = Date.now()) { + if (error.retryAfterSeconds !== null) return error.retryAfterSeconds; + const resetSeconds = Number(error.rateLimit?.reset); + if (!Number.isFinite(resetSeconds)) return null; + + return Math.max(0, Math.ceil((resetSeconds * 1000 - nowMs) / 1000)); +} + +function sleepWithAbort(delayMs: number, signal: AbortSignal) { + return new Promise((resolve, reject) => { + const timeout = globalThis.setTimeout(resolve, delayMs); + signal.addEventListener( + "abort", + () => { + globalThis.clearTimeout(timeout); + reject(new DOMException("Import cancelled.", "AbortError")); + }, + { once: true }, + ); + }); +} + +function throwIfAborted(signal: AbortSignal) { + if (signal.aborted) { + throw new DOMException("Import cancelled.", "AbortError"); + } +} diff --git a/docs/12-import-pipeline.md b/docs/12-import-pipeline.md index 08465ba..f204328 100644 --- a/docs/12-import-pipeline.md +++ b/docs/12-import-pipeline.md @@ -33,6 +33,12 @@ Import state contract: - `packages/core` owns the pure import run state helpers. - Current terminal states are `completed`, `failed`, `cancelled`, and `rate_limited`. - The web app uses this contract for visible import progress, cancellation, failure, and rate-limit terminal states. +- Rate-limited runs preserve retry timing when the Worker or GitHub response provides `Retry-After`, `retry_after_seconds`, or GitHub rate-limit reset metadata. + +Retry model: +- The web import worker may retry narrow transient failures: request timeout, network fetch failure, and 5xx worker/GitHub proxy responses. +- Authentication failures, validation failures, user cancellation, and rate-limit responses are not retried automatically. +- Retry delays are bounded by the web import retry policy so background import work does not quietly wait for long windows. Browser Web Worker responsibilities: - Pagination orchestration diff --git a/packages/core/src/import-state.ts b/packages/core/src/import-state.ts index 7926542..2c71f5f 100644 --- a/packages/core/src/import-state.ts +++ b/packages/core/src/import-state.ts @@ -7,6 +7,7 @@ export interface ImportRunState { pages: number; repositories: number; rate_limits: GitHubRateLimitSnapshot[]; + retry_after_seconds: number | null; errors: string[]; current_page: number | null; can_cancel: boolean; @@ -26,6 +27,7 @@ export function createImportRunState(startedAt = new Date().toISOString()): Impo pages: 0, repositories: 0, rate_limits: [], + retry_after_seconds: null, errors: [], current_page: 1, can_cancel: true, @@ -72,12 +74,14 @@ export function rateLimitImport( state: ImportRunState, error: string, rateLimit: GitHubRateLimitSnapshot | null, + retryAfterSeconds: number | null = null, completedAt = new Date().toISOString(), ): ImportRunState { assertRunning(state); return { ...terminalState(state, "rate_limited", completedAt, error), rate_limits: rateLimit ? [...state.rate_limits, rateLimit] : state.rate_limits, + retry_after_seconds: retryAfterSeconds, }; } @@ -90,6 +94,7 @@ export function importRunStateToEvent(id: string, state: ImportRunState): Import pages: state.pages, repositories: state.repositories, rate_limits: state.rate_limits, + retry_after_seconds: state.retry_after_seconds, errors: state.errors, }; } diff --git a/packages/core/test/import-state.test.mjs b/packages/core/test/import-state.test.mjs index 278d966..9a690f5 100644 --- a/packages/core/test/import-state.test.mjs +++ b/packages/core/test/import-state.test.mjs @@ -42,6 +42,7 @@ test("records page progress and converts to import event", () => { pages: 1, repositories: 100, rate_limits: [rateLimit], + retry_after_seconds: null, errors: [], }); }); @@ -57,6 +58,7 @@ test("captures failed, cancelled, and rate-limited terminal states", () => { createImportRunState(), "Rate limit exceeded", rateLimit, + 60, "2026-06-05T00:00:00.000Z", ); @@ -66,6 +68,7 @@ test("captures failed, cancelled, and rate-limited terminal states", () => { assert.equal(rateLimited.status, "rate_limited"); assert.deepEqual(rateLimited.errors, ["Rate limit exceeded"]); assert.deepEqual(rateLimited.rate_limits, [rateLimit]); + assert.equal(rateLimited.retry_after_seconds, 60); }); test("rejects updates after terminal states", () => { diff --git a/packages/shared/src/import.ts b/packages/shared/src/import.ts index 87b4452..c7d5dcf 100644 --- a/packages/shared/src/import.ts +++ b/packages/shared/src/import.ts @@ -16,5 +16,6 @@ export interface ImportEvent { pages: number; repositories: number; rate_limits: GitHubRateLimitSnapshot[]; + retry_after_seconds: number | null; errors: string[]; }