diff --git a/packages/stats/core/src/r2-sql.test.ts b/packages/stats/core/src/r2-sql.test.ts index bb19ca393cd..1077a143f5c 100644 --- a/packages/stats/core/src/r2-sql.test.ts +++ b/packages/stats/core/src/r2-sql.test.ts @@ -77,3 +77,35 @@ test("propagates a later page failure instead of returning partial aggregates", expect(result).toBe(error) expect(calls).toHaveLength(2) }) + +test("retries a transient timeout on the same page without repeating earlier rows", async () => { + const calls: string[] = [] + const rows = await Effect.runPromise( + queryR2SqlPages("SELECT * FROM aggregates", ["model"], (query) => { + calls.push(query) + if (calls.length === 1) + return Effect.succeed(Array.from({ length: 10000 }, (_, index) => ({ model: String(index).padStart(5, "0") }))) + if (calls.length === 2) + return Effect.fail(new R2SqlQueryError({ message: "query timeout", status: 400, code: 40005 })) + return Effect.succeed([{ model: "10000" }]) + }), + ) + expect(rows).toHaveLength(10001) + expect(new Set(rows.map((row) => row.model)).size).toBe(10001) + expect(calls).toHaveLength(3) + expect(calls[1]).toBe(calls[2]) + expect(calls[0]).not.toBe(calls[1]) +}, 10000) + +test("bounds transient retries and preserves the final error", async () => { + const calls: string[] = [] + const error = new R2SqlQueryError({ message: "unavailable", status: 503 }) + const result = await Effect.runPromise( + queryR2SqlPages("SELECT * FROM aggregates", undefined, (query) => { + calls.push(query) + return Effect.fail(error) + }).pipe(Effect.flip), + ) + expect(calls).toHaveLength(3) + expect(result).toBe(error) +}, 20000) diff --git a/packages/stats/core/src/r2-sql.ts b/packages/stats/core/src/r2-sql.ts index b9f3af44976..a3509c2f2c7 100644 --- a/packages/stats/core/src/r2-sql.ts +++ b/packages/stats/core/src/r2-sql.ts @@ -1,4 +1,4 @@ -import { Context, Effect, Layer, Schema } from "effect" +import { Context, Effect, Layer, Schedule, Schema } from "effect" import { Resource } from "sst/resource" const R2_SQL_MAX_ROWS = 10_000 @@ -16,6 +16,7 @@ const R2SqlResponse = Schema.Struct({ ), errors: Schema.Array(Schema.Unknown), }) +const R2SqlApiError = Schema.Struct({ code: Schema.Number, message: Schema.String }) const decodeResponse = Schema.decodeUnknownEffect(Schema.fromJsonString(R2SqlResponse)) export type R2SqlData = Record @@ -24,14 +25,16 @@ export class R2SqlQueryError extends Error { readonly _tag = "R2SqlQueryError" readonly requestId?: string readonly status?: number + readonly code?: number - constructor(input: { message: string; requestId?: string; status?: number; cause?: unknown }) { + constructor(input: { message: string; requestId?: string; status?: number; code?: number; cause?: unknown }) { super(input.cause instanceof Error ? `${input.message}: ${input.cause.toString()}` : input.message, { cause: input.cause, }) this.name = "R2SqlQueryError" this.requestId = input.requestId this.status = input.status + this.code = input.code } } @@ -60,7 +63,18 @@ export const queryR2SqlPages = Effect.fn("R2Sql.query")(function* ( const rows: R2SqlData[] = [] let cursor: R2SqlData | undefined while (true) { - const page = yield* fetchRows(columns?.length ? pageQuery(query, columns, cursor) : query) + const page = yield* Effect.suspend(() => + fetchRows(columns?.length ? pageQuery(query, columns, cursor) : query), + ).pipe( + // Retry only this page; a transient R2 timeout must not discard the whole + // display-window backfill. Syntax, auth, and row-limit errors still fail. + Effect.retry({ + times: 2, + schedule: Schedule.exponential("5 seconds"), + while: (error) => + error.code === 40005 || error.status === 429 || (error.status !== undefined && error.status >= 500), + }), + ) if (page.length >= R2_SQL_MAX_ROWS && !columns?.length) return yield* Effect.fail( new R2SqlQueryError({ message: `R2 SQL stats query reached the ${R2_SQL_MAX_ROWS} row limit` }), @@ -134,6 +148,7 @@ const fetchRows = Effect.fn("R2Sql.fetchRows")(function* (query: string) { message: `R2 SQL stats query failed: ${JSON.stringify(decoded.errors)}`, requestId: decoded.result?.request_id, status: response.status, + code: decoded.errors.find(Schema.is(R2SqlApiError))?.code, }), )