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
29 changes: 29 additions & 0 deletions apps/web/src/lib/api.test.ts
Original file line number Diff line number Diff line change
@@ -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<WorkerApiError>);
});
});
22 changes: 21 additions & 1 deletion apps/web/src/lib/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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(),
Expand All @@ -169,6 +188,7 @@ export function createImportEvent(): ImportEvent {
pages: 0,
repositories: 0,
rate_limits: [],
retry_after_seconds: null,
errors: [],
};
}
1 change: 1 addition & 0 deletions apps/web/src/lib/db.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,7 @@ function createImportEvent(id: string, startedAt: string): ImportEvent {
pages: 1,
repositories: 1,
rate_limits: [],
retry_after_seconds: null,
errors: [],
};
}
3 changes: 3 additions & 0 deletions apps/web/src/lib/import-messages.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.",
);
Expand Down
11 changes: 10 additions & 1 deletion apps/web/src/lib/import-messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.";
Expand All @@ -31,3 +32,11 @@ export function getImportTerminalText(importRun: ImportRunState) {
export function sortObservedFieldNames(observedFieldNames: Iterable<string>) {
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).`;
}
22 changes: 16 additions & 6 deletions apps/web/src/lib/import-pipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<StarredPageResponse>(
() => 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,
});
Expand All @@ -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);
}
Expand Down
54 changes: 54 additions & 0 deletions apps/web/src/lib/import-retry.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
74 changes: 74 additions & 0 deletions apps/web/src/lib/import-retry.ts
Original file line number Diff line number Diff line change
@@ -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<void>;
}

export async function runImportRequestWithRetry<T>(
request: () => Promise<T>,
{ 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<void>((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");
}
}
6 changes: 6 additions & 0 deletions docs/12-import-pipeline.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions packages/core/src/import-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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,
Expand Down Expand Up @@ -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,
};
}

Expand All @@ -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,
};
}
Expand Down
3 changes: 3 additions & 0 deletions packages/core/test/import-state.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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: [],
});
});
Expand All @@ -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",
);

Expand All @@ -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", () => {
Expand Down
1 change: 1 addition & 0 deletions packages/shared/src/import.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,5 +16,6 @@ export interface ImportEvent {
pages: number;
repositories: number;
rate_limits: GitHubRateLimitSnapshot[];
retry_after_seconds: number | null;
errors: string[];
}
Loading