diff --git a/packages/stats/core/src/domain/geo.ts b/packages/stats/core/src/domain/geo.ts index b75e08a34b..4b18446999 100644 --- a/packages/stats/core/src/domain/geo.ts +++ b/packages/stats/core/src/domain/geo.ts @@ -8,7 +8,6 @@ import { chunks, collapseRows, inserted, - isMissingUniqueUsersColumn, omitUniqueUsers, rankRowsWithMarketShare, statPeriodKey, @@ -16,6 +15,7 @@ import { synthesizeAllTierRows, toStatBaseRow, UPSERT_CHUNK_SIZE, + withUniqueUsersFallback, type StatBaseAggregate, } from "./stat" @@ -138,14 +138,7 @@ export class GeoStatRepo extends Context.Service Effect.tryPromise({ - try: async () => { - try { - return await upsertGeoChunk(chunk, true) - } catch (cause) { - if (!isMissingUniqueUsersColumn(cause)) throw cause - return upsertGeoChunk(chunk, false) - } - }, + try: () => withUniqueUsersFallback((includeUniqueUsers) => upsertGeoChunk(chunk, includeUniqueUsers)), catch: (cause) => DatabaseError.make({ cause }), }), { discard: true }, diff --git a/packages/stats/core/src/domain/model.ts b/packages/stats/core/src/domain/model.ts index ebaf93e1c8..2ad266eab8 100644 --- a/packages/stats/core/src/domain/model.ts +++ b/packages/stats/core/src/domain/model.ts @@ -8,7 +8,6 @@ import { chunks, collapseRows, inserted, - isMissingUniqueUsersColumn, omitUniqueUsers, rankBy, statPeriodKey, @@ -16,6 +15,7 @@ import { synthesizeAllTierRows, toStatBaseRow, UPSERT_CHUNK_SIZE, + withUniqueUsersFallback, type StatBaseAggregate, } from "./stat" @@ -59,31 +59,31 @@ export class ModelStatRepo extends Context.Service { - try { - return await db - .select({ - periodKey: modelStat.period_key, - updatedAt: modelStat.updated_at, - tier: modelStat.tier, - provider: modelStat.provider, - model: modelStat.model, - sessions: modelStat.sessions, - uniqueUsers: modelStat.unique_users, - inputTokens: modelStat.input_tokens, - outputTokens: modelStat.output_tokens, - reasoningTokens: modelStat.reasoning_tokens, - cacheReadTokens: modelStat.cache_read_tokens, - totalTokens: modelStat.total_tokens, - inputCostMicrocents: modelStat.input_cost_microcents, - outputCostMicrocents: modelStat.output_cost_microcents, - totalCostMicrocents: modelStat.total_cost_microcents, - }) - .from(modelStat) - .where(modelDailyScope()) - .orderBy(asc(modelStat.period_key)) - } catch (cause) { - if (!isMissingUniqueUsersColumn(cause)) throw cause + try: () => + withUniqueUsersFallback(async (includeUniqueUsers) => { + if (includeUniqueUsers) + return await db + .select({ + periodKey: modelStat.period_key, + updatedAt: modelStat.updated_at, + tier: modelStat.tier, + provider: modelStat.provider, + model: modelStat.model, + sessions: modelStat.sessions, + uniqueUsers: modelStat.unique_users, + inputTokens: modelStat.input_tokens, + outputTokens: modelStat.output_tokens, + reasoningTokens: modelStat.reasoning_tokens, + cacheReadTokens: modelStat.cache_read_tokens, + totalTokens: modelStat.total_tokens, + inputCostMicrocents: modelStat.input_cost_microcents, + outputCostMicrocents: modelStat.output_cost_microcents, + totalCostMicrocents: modelStat.total_cost_microcents, + }) + .from(modelStat) + .where(modelDailyScope()) + .orderBy(asc(modelStat.period_key)) + return ( await db .select({ @@ -106,8 +106,7 @@ export class ModelStatRepo extends Context.Service ({ ...row, uniqueUsers: 0 })) - } - }, + }), catch: (cause) => DatabaseError.make({ cause }), }) }) @@ -125,14 +124,7 @@ export class ModelStatRepo extends Context.Service Effect.tryPromise({ - try: async () => { - try { - return await upsertModelChunk(chunk, true) - } catch (cause) { - if (!isMissingUniqueUsersColumn(cause)) throw cause - return upsertModelChunk(chunk, false) - } - }, + try: () => withUniqueUsersFallback((includeUniqueUsers) => upsertModelChunk(chunk, includeUniqueUsers)), catch: (cause) => DatabaseError.make({ cause }), }), { discard: true }, diff --git a/packages/stats/core/src/domain/provider.ts b/packages/stats/core/src/domain/provider.ts index 1b177e8e83..c1292f0dd6 100644 --- a/packages/stats/core/src/domain/provider.ts +++ b/packages/stats/core/src/domain/provider.ts @@ -8,13 +8,13 @@ import { chunks, collapseRows, inserted, - isMissingUniqueUsersColumn, omitUniqueUsers, rankRowsWithMarketShare, statRowScope, synthesizeAllTierRows, toStatBaseRow, UPSERT_CHUNK_SIZE, + withUniqueUsersFallback, type StatBaseAggregate, } from "./stat" @@ -109,14 +109,8 @@ export class ProviderStatRepo extends Context.Service Effect.tryPromise({ - try: async () => { - try { - return await upsertProviderChunk(chunk, true) - } catch (cause) { - if (!isMissingUniqueUsersColumn(cause)) throw cause - return upsertProviderChunk(chunk, false) - } - }, + try: () => + withUniqueUsersFallback((includeUniqueUsers) => upsertProviderChunk(chunk, includeUniqueUsers)), catch: (cause) => DatabaseError.make({ cause }), }), { discard: true }, diff --git a/packages/stats/core/src/domain/stat.test.ts b/packages/stats/core/src/domain/stat.test.ts new file mode 100644 index 0000000000..38ce1c59b2 --- /dev/null +++ b/packages/stats/core/src/domain/stat.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, test } from "bun:test" +import { withUniqueUsersFallback } from "./stat" + +describe("withUniqueUsersFallback", () => { + test("writes with unique users first", async () => { + const calls: boolean[] = [] + + const result = await withUniqueUsersFallback(async (includeUniqueUsers) => { + calls.push(includeUniqueUsers) + return "written" + }) + + expect(result).toBe("written") + expect(calls).toEqual([true]) + }) + + test("retries without unique users only for the missing column error", async () => { + const calls: boolean[] = [] + + const result = await withUniqueUsersFallback(async (includeUniqueUsers) => { + calls.push(includeUniqueUsers) + if (includeUniqueUsers) throw new Error("Unknown column 'unique_users' in 'field list'") + return "written" + }) + + expect(result).toBe("written") + expect(calls).toEqual([true, false]) + }) + + test("rethrows unrelated errors without retrying", async () => { + const calls: boolean[] = [] + const error = new Error("connection failed") + + const result = withUniqueUsersFallback(async (includeUniqueUsers) => { + calls.push(includeUniqueUsers) + throw error + }) + + await expect(result).rejects.toBe(error) + expect(calls).toEqual([true]) + }) + + test("rethrows an error from the fallback write", async () => { + const calls: boolean[] = [] + const error = new Error("fallback failed") + + const result = withUniqueUsersFallback(async (includeUniqueUsers) => { + calls.push(includeUniqueUsers) + if (includeUniqueUsers) throw new Error("Unknown column 'unique_users' in 'field list'") + throw error + }) + + await expect(result).rejects.toBe(error) + expect(calls).toEqual([true, false]) + }) +}) diff --git a/packages/stats/core/src/domain/stat.ts b/packages/stats/core/src/domain/stat.ts index 553c91b9a5..4473df7e49 100644 --- a/packages/stats/core/src/domain/stat.ts +++ b/packages/stats/core/src/domain/stat.ts @@ -151,6 +151,15 @@ export function isMissingUniqueUsersColumn(cause: unknown): boolean { return errorText(cause).includes("Unknown column 'unique_users'") } +export async function withUniqueUsersFallback(write: (includeUniqueUsers: boolean) => Promise) { + try { + return await write(true) + } catch (cause) { + if (!isMissingUniqueUsersColumn(cause)) throw cause + return write(false) + } +} + export function omitUniqueUsers(rows: T[]) { return rows.map((row) => { const result = { ...row }