Skip to content

Commit d870e22

Browse files
authored
fix(stats): retry transient query failures (anomalyco#50205)
1 parent 45ad8dc commit d870e22

2 files changed

Lines changed: 50 additions & 3 deletions

File tree

‎packages/stats/core/src/r2-sql.test.ts‎

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,3 +77,35 @@ test("propagates a later page failure instead of returning partial aggregates",
7777
expect(result).toBe(error)
7878
expect(calls).toHaveLength(2)
7979
})
80+
81+
test("retries a transient timeout on the same page without repeating earlier rows", async () => {
82+
const calls: string[] = []
83+
const rows = await Effect.runPromise(
84+
queryR2SqlPages("SELECT * FROM aggregates", ["model"], (query) => {
85+
calls.push(query)
86+
if (calls.length === 1)
87+
return Effect.succeed(Array.from({ length: 10000 }, (_, index) => ({ model: String(index).padStart(5, "0") })))
88+
if (calls.length === 2)
89+
return Effect.fail(new R2SqlQueryError({ message: "query timeout", status: 400, code: 40005 }))
90+
return Effect.succeed([{ model: "10000" }])
91+
}),
92+
)
93+
expect(rows).toHaveLength(10001)
94+
expect(new Set(rows.map((row) => row.model)).size).toBe(10001)
95+
expect(calls).toHaveLength(3)
96+
expect(calls[1]).toBe(calls[2])
97+
expect(calls[0]).not.toBe(calls[1])
98+
}, 10000)
99+
100+
test("bounds transient retries and preserves the final error", async () => {
101+
const calls: string[] = []
102+
const error = new R2SqlQueryError({ message: "unavailable", status: 503 })
103+
const result = await Effect.runPromise(
104+
queryR2SqlPages("SELECT * FROM aggregates", undefined, (query) => {
105+
calls.push(query)
106+
return Effect.fail(error)
107+
}).pipe(Effect.flip),
108+
)
109+
expect(calls).toHaveLength(3)
110+
expect(result).toBe(error)
111+
}, 20000)

‎packages/stats/core/src/r2-sql.ts‎

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { Context, Effect, Layer, Schema } from "effect"
1+
import { Context, Effect, Layer, Schedule, Schema } from "effect"
22
import { Resource } from "sst/resource"
33

44
const R2_SQL_MAX_ROWS = 10_000
@@ -16,6 +16,7 @@ const R2SqlResponse = Schema.Struct({
1616
),
1717
errors: Schema.Array(Schema.Unknown),
1818
})
19+
const R2SqlApiError = Schema.Struct({ code: Schema.Number, message: Schema.String })
1920
const decodeResponse = Schema.decodeUnknownEffect(Schema.fromJsonString(R2SqlResponse))
2021

2122
export type R2SqlData = Record<string, string>
@@ -24,14 +25,16 @@ export class R2SqlQueryError extends Error {
2425
readonly _tag = "R2SqlQueryError"
2526
readonly requestId?: string
2627
readonly status?: number
28+
readonly code?: number
2729

28-
constructor(input: { message: string; requestId?: string; status?: number; cause?: unknown }) {
30+
constructor(input: { message: string; requestId?: string; status?: number; code?: number; cause?: unknown }) {
2931
super(input.cause instanceof Error ? `${input.message}: ${input.cause.toString()}` : input.message, {
3032
cause: input.cause,
3133
})
3234
this.name = "R2SqlQueryError"
3335
this.requestId = input.requestId
3436
this.status = input.status
37+
this.code = input.code
3538
}
3639
}
3740

@@ -60,7 +63,18 @@ export const queryR2SqlPages = Effect.fn("R2Sql.query")(function* (
6063
const rows: R2SqlData[] = []
6164
let cursor: R2SqlData | undefined
6265
while (true) {
63-
const page = yield* fetchRows(columns?.length ? pageQuery(query, columns, cursor) : query)
66+
const page = yield* Effect.suspend(() =>
67+
fetchRows(columns?.length ? pageQuery(query, columns, cursor) : query),
68+
).pipe(
69+
// Retry only this page; a transient R2 timeout must not discard the whole
70+
// display-window backfill. Syntax, auth, and row-limit errors still fail.
71+
Effect.retry({
72+
times: 2,
73+
schedule: Schedule.exponential("5 seconds"),
74+
while: (error) =>
75+
error.code === 40005 || error.status === 429 || (error.status !== undefined && error.status >= 500),
76+
}),
77+
)
6478
if (page.length >= R2_SQL_MAX_ROWS && !columns?.length)
6579
return yield* Effect.fail(
6680
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) {
134148
message: `R2 SQL stats query failed: ${JSON.stringify(decoded.errors)}`,
135149
requestId: decoded.result?.request_id,
136150
status: response.status,
151+
code: decoded.errors.find(Schema.is(R2SqlApiError))?.code,
137152
}),
138153
)
139154

0 commit comments

Comments
 (0)