diff --git a/packages/stats/core/src/domain/inference.test.ts b/packages/stats/core/src/domain/inference.test.ts index 4f51599d71..8267b3c860 100644 --- a/packages/stats/core/src/domain/inference.test.ts +++ b/packages/stats/core/src/domain/inference.test.ts @@ -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({ diff --git a/packages/stats/core/src/domain/inference.ts b/packages/stats/core/src/domain/inference.ts index bdba90b11f..2b695e4ad9 100644 --- a/packages/stats/core/src/domain/inference.ts +++ b/packages/stats/core/src/domain/inference.ts @@ -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( diff --git a/packages/stats/core/src/r2-sql.test.ts b/packages/stats/core/src/r2-sql.test.ts new file mode 100644 index 0000000000..b02f429b16 --- /dev/null +++ b/packages/stats/core/src/r2-sql.test.ts @@ -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") +}) diff --git a/packages/stats/core/src/r2-sql.ts b/packages/stats/core/src/r2-sql.ts index 91093643c0..1941d439ec 100644 --- a/packages/stats/core/src/r2-sql.ts +++ b/packages/stats/core/src/r2-sql.ts @@ -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()("@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(), diff --git a/packages/stats/core/src/stat-sync.ts b/packages/stats/core/src/stat-sync.ts index aca7fbc6a4..ed71f2fee6 100644 --- a/packages/stats/core/src/stat-sync.ts +++ b/packages/stats/core/src/stat-sync.ts @@ -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),