fix(stats): prevent full retention sync timeouts
This commit is contained in:
parent
b7ca4f91d9
commit
43e89ea167
5 changed files with 125 additions and 30 deletions
|
|
@ -205,8 +205,8 @@ describe("inference stat normalization", () => {
|
|||
dataset: "zen",
|
||||
})
|
||||
|
||||
expect(queries).toHaveLength(1)
|
||||
expect(queries[0]?.cohortDates).toEqual(["2026-08-10", "2026-08-17"])
|
||||
expect(queries).toHaveLength(2)
|
||||
expect(queries.map((query) => query.cohortDates)).toEqual([["2026-08-10"], ["2026-08-17"]])
|
||||
expect(queries[0]?.query).toContain("AND product = 'go'")
|
||||
expect(queries[0]?.query).toContain("AND lower(model) NOT IN ('alpha-gpt-next')")
|
||||
expect(queries[0]?.query).toContain("CASE\n WHEN lower(model) IN ('omen-alpha') THEN 'unknown'\n")
|
||||
|
|
@ -222,9 +222,12 @@ describe("inference stat normalization", () => {
|
|||
)
|
||||
expect(queries[0]?.query).not.toContain(" OVER (")
|
||||
expect(queries[0]?.query).toContain("WHEN '2026-08-17' THEN '2026-08-10'")
|
||||
expect(queries[0]?.query).toContain("WHEN '2026-08-24' THEN '2026-08-17'")
|
||||
expect(queries[0]?.query).not.toContain("WHEN '2026-08-24' THEN '2026-08-17'")
|
||||
expect(queries[1]?.query).toContain("WHEN '2026-08-24' THEN '2026-08-17'")
|
||||
expect(queries[0]?.query).toContain("started_at >= '2026-08-10T00:00:00.000Z'")
|
||||
expect(queries[0]?.query).toContain("started_at < '2026-08-31T00:00:00.000Z'")
|
||||
expect(queries[0]?.query).toContain("started_at < '2026-08-24T00:00:00.000Z'")
|
||||
expect(queries[1]?.query).toContain("started_at >= '2026-08-17T00:00:00.000Z'")
|
||||
expect(queries[1]?.query).toContain("started_at < '2026-08-31T00:00:00.000Z'")
|
||||
expect(queries[0]?.query).toContain("LEFT JOIN returned ON primary_models.user_key = returned.user_key")
|
||||
expect(queries[0]?.query).toContain("primary_models.cohort_date = returned.cohort_date")
|
||||
expect(queries[0]?.query).toContain("'Go' AS tier")
|
||||
|
|
@ -232,6 +235,29 @@ describe("inference stat normalization", () => {
|
|||
expect(queries[0]?.query).toContain("LIMIT 10000")
|
||||
})
|
||||
|
||||
test("splits a full retention window without dropping or duplicating cohorts", () => {
|
||||
const source = { namespace: "inference", table: "generation", dataset: "zen" }
|
||||
const queries = buildRetentionQueries(new Date("2026-07-16T19:00:00Z"), new Date("2026-09-10T00:00:00Z"), source)
|
||||
|
||||
expect(queries.flatMap((query) => query.cohortDates)).toEqual([
|
||||
"2026-07-13",
|
||||
"2026-07-20",
|
||||
"2026-07-27",
|
||||
"2026-08-03",
|
||||
"2026-08-10",
|
||||
"2026-08-17",
|
||||
"2026-08-24",
|
||||
])
|
||||
queries.forEach((query) => {
|
||||
const start = new Date(`${query.cohortDates[0]}T00:00:00Z`)
|
||||
const end = new Date(start.getTime() + 14 * 86_400_000)
|
||||
expect(query).toEqual(buildRetentionQueries(start, end, source)[0])
|
||||
})
|
||||
expect(buildRetentionQueries(new Date("2026-08-31T00:00:00Z"), new Date("2026-09-10T00:00:00Z"), source)).toEqual(
|
||||
[],
|
||||
)
|
||||
})
|
||||
|
||||
test("maps retention query results", () => {
|
||||
expect(
|
||||
toRetentionAggregate({
|
||||
|
|
|
|||
|
|
@ -56,13 +56,12 @@ export function buildRetentionQueries(periodStart: Date, periodEnd: Date, input?
|
|||
dataset: Resource.StatsSyncConfig.dataset,
|
||||
}
|
||||
const periods = retentionPeriods(periodStart, periodEnd)
|
||||
if (periods.length === 0) return []
|
||||
return [
|
||||
{
|
||||
cohortDates: periods.map((period) => period.start.toISOString().slice(0, 10)),
|
||||
query: buildRetentionQuery(periods, source),
|
||||
},
|
||||
]
|
||||
// Bound the user-level joins to one activity week and its return week.
|
||||
// Combining the entire display window makes full syncs much more expensive.
|
||||
return periods.map((period) => ({
|
||||
cohortDates: [period.start.toISOString().slice(0, 10)],
|
||||
query: buildRetentionQuery([period], source),
|
||||
}))
|
||||
}
|
||||
|
||||
function buildRetentionQuery(
|
||||
|
|
|
|||
18
packages/stats/core/src/r2-sql.test.ts
Normal file
18
packages/stats/core/src/r2-sql.test.ts
Normal file
|
|
@ -0,0 +1,18 @@
|
|||
import { expect, test } from "bun:test"
|
||||
import { R2SqlQueryError } from "./r2-sql"
|
||||
|
||||
test("includes the transport failure in the message used by sync logging", () => {
|
||||
const cause = new DOMException("The operation timed out.", "TimeoutError")
|
||||
const error = new R2SqlQueryError({ message: "Failed to run R2 SQL stats query", cause })
|
||||
|
||||
expect(error.message).toBe("Failed to run R2 SQL stats query: TimeoutError: The operation timed out.")
|
||||
expect(error.cause).toBe(cause)
|
||||
})
|
||||
|
||||
test("preserves an R2 response error and request details", () => {
|
||||
const error = new R2SqlQueryError({ message: "R2 SQL stats query failed", status: 400, requestId: "test-request" })
|
||||
|
||||
expect(error.message).toBe("R2 SQL stats query failed")
|
||||
expect(error.status).toBe(400)
|
||||
expect(error.requestId).toBe("test-request")
|
||||
})
|
||||
|
|
@ -2,6 +2,7 @@ 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,
|
||||
|
|
@ -25,7 +26,9 @@ export class R2SqlQueryError extends Error {
|
|||
readonly status?: number
|
||||
|
||||
constructor(input: { message: string; requestId?: string; status?: number; cause?: unknown }) {
|
||||
super(input.message, { cause: input.cause })
|
||||
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
|
||||
|
|
@ -43,20 +46,31 @@ export class R2Sql extends Context.Service<R2Sql, R2Sql.Service>()("@opencode/st
|
|||
R2Sql,
|
||||
R2Sql.of({
|
||||
query: Effect.fn("R2Sql.query")(function* (query: string) {
|
||||
const startedAt = Date.now()
|
||||
const response = yield* Effect.tryPromise({
|
||||
try: () =>
|
||||
Bun.fetch(
|
||||
`https://api.sql.cloudflarestorage.com/api/v1/accounts/${Resource.R2Sql.accountId}/r2-sql/query/${Resource.R2Sql.bucket}`,
|
||||
{
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `Bearer ${Resource.R2SqlAuthToken.value}`,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
body: JSON.stringify({ query }),
|
||||
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",
|
||||
},
|
||||
),
|
||||
catch: (cause) => new R2SqlQueryError({ message: "Failed to run R2 SQL stats query", cause }),
|
||||
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(),
|
||||
|
|
|
|||
|
|
@ -47,9 +47,29 @@ export const syncStats: (options?: {
|
|||
|
||||
yield* logRuntimeCheck()
|
||||
|
||||
const rows = yield* Effect.forEach(buildStatsQueries(periodStart, periodEnd), r2Sql.query, {
|
||||
concurrency: 4,
|
||||
}).pipe(Effect.map((batches) => batches.flat()))
|
||||
const queries = buildStatsQueries(periodStart, periodEnd)
|
||||
yield* Effect.logInfo(
|
||||
`stats sync started ${JSON.stringify({ full: options?.full ?? false, periodStart, periodEnd, queries: queries.length })}`,
|
||||
)
|
||||
const rows = yield* Effect.forEach(
|
||||
queries,
|
||||
(query, index) =>
|
||||
r2Sql.query(query).pipe(
|
||||
Effect.tap((rows) =>
|
||||
Effect.logInfo(
|
||||
`stats query complete ${JSON.stringify({ index, total: queries.length, rows: rows.length })}`,
|
||||
),
|
||||
),
|
||||
Effect.tapError((error) =>
|
||||
Effect.logError(
|
||||
`stats query failed ${JSON.stringify({ index, total: queries.length, error: error.message })}`,
|
||||
),
|
||||
),
|
||||
),
|
||||
{
|
||||
concurrency: 4,
|
||||
},
|
||||
).pipe(Effect.map((batches) => batches.flat()))
|
||||
const modelRows = modelRowsFromAggregates(rows.filter((row) => row.dimension === "model").flatMap(toModelAggregate))
|
||||
const providerRows = providerRowsFromAggregates(
|
||||
rows.filter((row) => row.dimension === "provider").flatMap(toProviderAggregate),
|
||||
|
|
@ -68,12 +88,30 @@ export const syncStats: (options?: {
|
|||
startOfUtcDay(periodEnd),
|
||||
)
|
||||
: []
|
||||
yield* Effect.logInfo(`stats sync querying retention ${JSON.stringify({ queries: retentionQueries.length })}`)
|
||||
const retentionRows = retentionRowsFromAggregates(
|
||||
yield* Effect.forEach(retentionQueries, (item) => r2Sql.query(item.query), { concurrency: 4 }).pipe(
|
||||
Effect.map((batches) => batches.flatMap((batch) => batch.flatMap(toRetentionAggregate))),
|
||||
),
|
||||
yield* Effect.forEach(
|
||||
retentionQueries,
|
||||
(item) =>
|
||||
r2Sql.query(item.query).pipe(
|
||||
Effect.tap((rows) =>
|
||||
Effect.logInfo(
|
||||
`retention query complete ${JSON.stringify({ cohortDates: item.cohortDates, rows: rows.length })}`,
|
||||
),
|
||||
),
|
||||
Effect.tapError((error) =>
|
||||
Effect.logError(
|
||||
`retention query failed ${JSON.stringify({ cohortDates: item.cohortDates, error: error.message })}`,
|
||||
),
|
||||
),
|
||||
),
|
||||
{ concurrency: 4 },
|
||||
).pipe(Effect.map((batches) => batches.flatMap((batch) => batch.flatMap(toRetentionAggregate)))),
|
||||
)
|
||||
|
||||
yield* Effect.logInfo(
|
||||
`stats sync writing aggregates ${JSON.stringify({ modelRows: modelRows.length, providerRows: providerRows.length, geoRows: geoRows.length, retentionRows: retentionRows.length })}`,
|
||||
)
|
||||
yield* Effect.all(
|
||||
[
|
||||
modelStats.upsert(modelRows),
|
||||
|
|
|
|||
Loading…
Reference in a new issue