File size: 4,411 Bytes
4bbfe8b
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
import { Context, Effect, Layer, Schema } from "effect"
import { Resource } from "sst/resource"

const R2_SQL_MAX_ROWS = 10_000
const R2_SQL_TIMEOUT_MS = 15 * 60_000
const R2SqlValue = Schema.Union([Schema.String, Schema.Number, Schema.Boolean, Schema.Null])
const R2SqlResponse = Schema.Struct({
  success: Schema.Boolean,
  result: Schema.optional(
    Schema.NullOr(
      Schema.Struct({
        request_id: Schema.String,
        rows: Schema.Array(Schema.Record(Schema.String, R2SqlValue)),
      }),
    ),
  ),
  errors: Schema.Array(Schema.Unknown),
})
const decodeResponse = Schema.decodeUnknownEffect(Schema.fromJsonString(R2SqlResponse))

export type R2SqlData = Record<string, string>

export class R2SqlQueryError extends Error {
  readonly _tag = "R2SqlQueryError"
  readonly requestId?: string
  readonly status?: number

  constructor(input: { message: string; requestId?: string; status?: 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
  }
}

export declare namespace R2Sql {
  export interface Service {
    readonly query: (query: string) => Effect.Effect<R2SqlData[], R2SqlQueryError>
  }
}

export class R2Sql extends Context.Service<R2Sql, R2Sql.Service>()("@opencode/stats/R2Sql") {
  static readonly layer: Layer.Layer<R2Sql> = Layer.succeed(
    R2Sql,
    R2Sql.of({
      query: Effect.fn("R2Sql.query")(function* (query: string) {
        const startedAt = Date.now()
        const response = yield* Effect.tryPromise({
          try: (signal) => {
            const options = {
              method: "POST",
              // Analytical queries can exceed Bun's default five-minute idle timer.
              // Bound the whole request and cancel it when the sync is interrupted.
              timeout: false,
              signal: AbortSignal.any([signal, AbortSignal.timeout(R2_SQL_TIMEOUT_MS)]),
              headers: {
                Authorization: `Bearer ${Resource.R2SqlAuthToken.value}`,
                "Content-Type": "application/json",
              },
              body: JSON.stringify({ query }),
            }
            return Bun.fetch(
              `https://api.sql.cloudflarestorage.com/api/v1/accounts/${Resource.R2Sql.accountId}/r2-sql/query/${Resource.R2Sql.bucket}`,
              options,
            )
          },
          catch: (cause) =>
            new R2SqlQueryError({
              message: `Failed to run R2 SQL stats query after ${Date.now() - startedAt}ms`,
              cause,
            }),
        })
        const body = yield* Effect.tryPromise({
          try: () => response.text(),
          catch: (cause) =>
            new R2SqlQueryError({ message: "Failed to read R2 SQL stats response", status: response.status, cause }),
        })
        const decoded = yield* decodeResponse(body).pipe(
          Effect.mapError(
            (cause) =>
              new R2SqlQueryError({
                message: "R2 SQL returned an invalid stats response",
                status: response.status,
                cause,
              }),
          ),
        )
        if (!response.ok || !decoded.success || !decoded.result)
          return yield* Effect.fail(
            new R2SqlQueryError({
              message: `R2 SQL stats query failed: ${JSON.stringify(decoded.errors)}`,
              requestId: decoded.result?.request_id,
              status: response.status,
            }),
          )

        // R2 SQL has no OFFSET support and caps LIMIT at 10,000. Each stats
        // query is scoped to one day or week, and reaching the cap is treated as
        // an error so a newly high-cardinality period can never be truncated.
        if (decoded.result.rows.length >= R2_SQL_MAX_ROWS)
          return yield* Effect.fail(
            new R2SqlQueryError({
              message: `R2 SQL stats query reached the ${R2_SQL_MAX_ROWS} row limit`,
              requestId: decoded.result.request_id,
              status: response.status,
            }),
          )

        return decoded.result.rows.map((row) =>
          Object.fromEntries(
            Object.entries(row).flatMap(([key, value]) => (value === null ? [] : [[key, String(value)]])),
          ),
        )
      }),
    }),
  )
}