1- import { Context , Effect , Layer , Schema } from "effect"
1+ import { Context , Effect , Layer , Schedule , Schema } from "effect"
22import { Resource } from "sst/resource"
33
44const 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 } )
1920const decodeResponse = Schema . decodeUnknownEffect ( Schema . fromJsonString ( R2SqlResponse ) )
2021
2122export 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