fix(stats): tolerate pending user column

This commit is contained in:
Adam
2026-06-20 15:56:36 -05:00
parent 9d59639f9d
commit 51afb5aef7
4 changed files with 221 additions and 135 deletions
+49 -37
View File
@@ -8,6 +8,8 @@ import {
chunks, chunks,
collapseRows, collapseRows,
inserted, inserted,
isMissingUniqueUsersColumn,
omitUniqueUsers,
rankRowsWithMarketShare, rankRowsWithMarketShare,
statPeriodKey, statPeriodKey,
statRowScope, statRowScope,
@@ -136,49 +138,59 @@ export class GeoStatRepo extends Context.Service<GeoStatRepo, GeoStatRepo.Servic
chunks(rows, UPSERT_CHUNK_SIZE), chunks(rows, UPSERT_CHUNK_SIZE),
(chunk) => (chunk) =>
Effect.tryPromise({ Effect.tryPromise({
try: () => try: async () => {
db try {
.insert(geoStat) return await upsertGeoChunk(chunk, true)
.values(chunk) } catch (cause) {
.onDuplicateKeyUpdate({ if (!isMissingUniqueUsersColumn(cause)) throw cause
set: { return upsertGeoChunk(chunk, false)
continent: inserted("continent"), }
sessions: inserted("sessions"), },
requests: inserted("requests"),
unique_users: inserted("unique_users"),
input_tokens: inserted("input_tokens"),
output_tokens: inserted("output_tokens"),
reasoning_tokens: inserted("reasoning_tokens"),
cache_read_tokens: inserted("cache_read_tokens"),
total_tokens: inserted("total_tokens"),
input_cost_microcents: inserted("input_cost_microcents"),
output_cost_microcents: inserted("output_cost_microcents"),
total_cost_microcents: inserted("total_cost_microcents"),
avg_duration_ms: inserted("avg_duration_ms"),
p50_duration_ms: inserted("p50_duration_ms"),
p95_duration_ms: inserted("p95_duration_ms"),
avg_ttfb_ms: inserted("avg_ttfb_ms"),
p50_ttfb_ms: inserted("p50_ttfb_ms"),
p95_ttfb_ms: inserted("p95_ttfb_ms"),
avg_output_tps: inserted("avg_output_tps"),
success_count: inserted("success_count"),
error_count: inserted("error_count"),
sample_count: inserted("sample_count"),
market_share_tokens: inserted("market_share_tokens"),
market_share_requests: inserted("market_share_requests"),
market_share_sessions: inserted("market_share_sessions"),
rank_by_tokens: inserted("rank_by_tokens"),
rank_by_requests: inserted("rank_by_requests"),
rank_by_sessions: inserted("rank_by_sessions"),
rank_by_cost: inserted("rank_by_cost"),
},
}),
catch: (cause) => DatabaseError.make({ cause }), catch: (cause) => DatabaseError.make({ cause }),
}), }),
{ discard: true }, { discard: true },
) )
}) })
function upsertGeoChunk(chunk: GeoStatRow[], includeUniqueUsers: boolean) {
return db
.insert(geoStat)
.values(includeUniqueUsers ? chunk : omitUniqueUsers(chunk))
.onDuplicateKeyUpdate({
set: {
continent: inserted("continent"),
sessions: inserted("sessions"),
requests: inserted("requests"),
...(includeUniqueUsers ? { unique_users: inserted("unique_users") } : {}),
input_tokens: inserted("input_tokens"),
output_tokens: inserted("output_tokens"),
reasoning_tokens: inserted("reasoning_tokens"),
cache_read_tokens: inserted("cache_read_tokens"),
total_tokens: inserted("total_tokens"),
input_cost_microcents: inserted("input_cost_microcents"),
output_cost_microcents: inserted("output_cost_microcents"),
total_cost_microcents: inserted("total_cost_microcents"),
avg_duration_ms: inserted("avg_duration_ms"),
p50_duration_ms: inserted("p50_duration_ms"),
p95_duration_ms: inserted("p95_duration_ms"),
avg_ttfb_ms: inserted("avg_ttfb_ms"),
p50_ttfb_ms: inserted("p50_ttfb_ms"),
p95_ttfb_ms: inserted("p95_ttfb_ms"),
avg_output_tps: inserted("avg_output_tps"),
success_count: inserted("success_count"),
error_count: inserted("error_count"),
sample_count: inserted("sample_count"),
market_share_tokens: inserted("market_share_tokens"),
market_share_requests: inserted("market_share_requests"),
market_share_sessions: inserted("market_share_sessions"),
rank_by_tokens: inserted("rank_by_tokens"),
rank_by_requests: inserted("rank_by_requests"),
rank_by_sessions: inserted("rank_by_sessions"),
rank_by_cost: inserted("rank_by_cost"),
},
})
}
const deleteRetiredDimensions = Effect.fn("GeoStatRepo.deleteRetiredDimensions")(function* (rows: GeoStatRow[]) { const deleteRetiredDimensions = Effect.fn("GeoStatRepo.deleteRetiredDimensions")(function* (rows: GeoStatRow[]) {
const scope = statRowScope(rows) const scope = statRowScope(rows)
if (!scope) return if (!scope) return
+103 -62
View File
@@ -8,6 +8,8 @@ import {
chunks, chunks,
collapseRows, collapseRows,
inserted, inserted,
isMissingUniqueUsersColumn,
omitUniqueUsers,
rankBy, rankBy,
statPeriodKey, statPeriodKey,
statRowScope, statRowScope,
@@ -56,35 +58,55 @@ export class ModelStatRepo extends Context.Service<ModelStatRepo, ModelStatRepo.
const listDaily = Effect.fn("ModelStatRepo.listDaily")(function* () { const listDaily = Effect.fn("ModelStatRepo.listDaily")(function* () {
return yield* Effect.tryPromise({ return yield* Effect.tryPromise({
try: () => try: async () => {
db try {
.select({ return await db
periodKey: modelStat.period_key, .select({
updatedAt: modelStat.updated_at, periodKey: modelStat.period_key,
tier: modelStat.tier, updatedAt: modelStat.updated_at,
provider: modelStat.provider, tier: modelStat.tier,
model: modelStat.model, provider: modelStat.provider,
sessions: modelStat.sessions, model: modelStat.model,
uniqueUsers: modelStat.unique_users, sessions: modelStat.sessions,
inputTokens: modelStat.input_tokens, uniqueUsers: modelStat.unique_users,
outputTokens: modelStat.output_tokens, inputTokens: modelStat.input_tokens,
reasoningTokens: modelStat.reasoning_tokens, outputTokens: modelStat.output_tokens,
cacheReadTokens: modelStat.cache_read_tokens, reasoningTokens: modelStat.reasoning_tokens,
totalTokens: modelStat.total_tokens, cacheReadTokens: modelStat.cache_read_tokens,
inputCostMicrocents: modelStat.input_cost_microcents, totalTokens: modelStat.total_tokens,
outputCostMicrocents: modelStat.output_cost_microcents, inputCostMicrocents: modelStat.input_cost_microcents,
totalCostMicrocents: modelStat.total_cost_microcents, outputCostMicrocents: modelStat.output_cost_microcents,
}) totalCostMicrocents: modelStat.total_cost_microcents,
.from(modelStat) })
.where( .from(modelStat)
and( .where(modelDailyScope())
eq(modelStat.grain, "day"), .orderBy(asc(modelStat.period_key))
eq(modelStat.client, "all"), } catch (cause) {
eq(modelStat.source, "all"), if (!isMissingUniqueUsersColumn(cause)) throw cause
inArray(modelStat.tier, ["Go", "go"]), return (
), await db
) .select({
.orderBy(asc(modelStat.period_key)), periodKey: modelStat.period_key,
updatedAt: modelStat.updated_at,
tier: modelStat.tier,
provider: modelStat.provider,
model: modelStat.model,
sessions: modelStat.sessions,
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))
).map((row) => ({ ...row, uniqueUsers: 0 }))
}
},
catch: (cause) => DatabaseError.make({ cause }), catch: (cause) => DatabaseError.make({ cause }),
}) })
}) })
@@ -94,45 +116,55 @@ export class ModelStatRepo extends Context.Service<ModelStatRepo, ModelStatRepo.
chunks(rows, UPSERT_CHUNK_SIZE), chunks(rows, UPSERT_CHUNK_SIZE),
(chunk) => (chunk) =>
Effect.tryPromise({ Effect.tryPromise({
try: () => try: async () => {
db try {
.insert(modelStat) return await upsertModelChunk(chunk, true)
.values(chunk) } catch (cause) {
.onDuplicateKeyUpdate({ if (!isMissingUniqueUsersColumn(cause)) throw cause
set: { return upsertModelChunk(chunk, false)
provider_model: inserted("provider_model"), }
sessions: inserted("sessions"), },
requests: inserted("requests"),
unique_users: inserted("unique_users"),
input_tokens: inserted("input_tokens"),
output_tokens: inserted("output_tokens"),
reasoning_tokens: inserted("reasoning_tokens"),
cache_read_tokens: inserted("cache_read_tokens"),
total_tokens: inserted("total_tokens"),
input_cost_microcents: inserted("input_cost_microcents"),
output_cost_microcents: inserted("output_cost_microcents"),
total_cost_microcents: inserted("total_cost_microcents"),
avg_duration_ms: inserted("avg_duration_ms"),
p50_duration_ms: inserted("p50_duration_ms"),
p95_duration_ms: inserted("p95_duration_ms"),
avg_ttfb_ms: inserted("avg_ttfb_ms"),
p50_ttfb_ms: inserted("p50_ttfb_ms"),
p95_ttfb_ms: inserted("p95_ttfb_ms"),
avg_output_tps: inserted("avg_output_tps"),
success_count: inserted("success_count"),
error_count: inserted("error_count"),
sample_count: inserted("sample_count"),
rank_by_tokens: inserted("rank_by_tokens"),
rank_by_requests: inserted("rank_by_requests"),
rank_by_cost: inserted("rank_by_cost"),
},
}),
catch: (cause) => DatabaseError.make({ cause }), catch: (cause) => DatabaseError.make({ cause }),
}), }),
{ discard: true }, { discard: true },
) )
}) })
function upsertModelChunk(chunk: ModelStatRow[], includeUniqueUsers: boolean) {
return db
.insert(modelStat)
.values(includeUniqueUsers ? chunk : omitUniqueUsers(chunk))
.onDuplicateKeyUpdate({
set: {
provider_model: inserted("provider_model"),
sessions: inserted("sessions"),
requests: inserted("requests"),
...(includeUniqueUsers ? { unique_users: inserted("unique_users") } : {}),
input_tokens: inserted("input_tokens"),
output_tokens: inserted("output_tokens"),
reasoning_tokens: inserted("reasoning_tokens"),
cache_read_tokens: inserted("cache_read_tokens"),
total_tokens: inserted("total_tokens"),
input_cost_microcents: inserted("input_cost_microcents"),
output_cost_microcents: inserted("output_cost_microcents"),
total_cost_microcents: inserted("total_cost_microcents"),
avg_duration_ms: inserted("avg_duration_ms"),
p50_duration_ms: inserted("p50_duration_ms"),
p95_duration_ms: inserted("p95_duration_ms"),
avg_ttfb_ms: inserted("avg_ttfb_ms"),
p50_ttfb_ms: inserted("p50_ttfb_ms"),
p95_ttfb_ms: inserted("p95_ttfb_ms"),
avg_output_tps: inserted("avg_output_tps"),
success_count: inserted("success_count"),
error_count: inserted("error_count"),
sample_count: inserted("sample_count"),
rank_by_tokens: inserted("rank_by_tokens"),
rank_by_requests: inserted("rank_by_requests"),
rank_by_cost: inserted("rank_by_cost"),
},
})
}
const deleteRetiredDimensions = Effect.fn("ModelStatRepo.deleteRetiredDimensions")(function* ( const deleteRetiredDimensions = Effect.fn("ModelStatRepo.deleteRetiredDimensions")(function* (
rows: ModelStatRow[], rows: ModelStatRow[],
) { ) {
@@ -165,6 +197,15 @@ export class ModelStatRepo extends Context.Service<ModelStatRepo, ModelStatRepo.
) )
} }
function modelDailyScope() {
return and(
eq(modelStat.grain, "day"),
eq(modelStat.client, "all"),
eq(modelStat.source, "all"),
inArray(modelStat.tier, ["Go", "go"]),
)
}
export function rowsFromAggregates(aggregates: ModelStatAggregate[]) { export function rowsFromAggregates(aggregates: ModelStatAggregate[]) {
return rankRows([ return rankRows([
...synthesizeAllTierRows( ...synthesizeAllTierRows(
+48 -36
View File
@@ -8,6 +8,8 @@ import {
chunks, chunks,
collapseRows, collapseRows,
inserted, inserted,
isMissingUniqueUsersColumn,
omitUniqueUsers,
rankRowsWithMarketShare, rankRowsWithMarketShare,
statRowScope, statRowScope,
synthesizeAllTierRows, synthesizeAllTierRows,
@@ -107,48 +109,58 @@ export class ProviderStatRepo extends Context.Service<ProviderStatRepo, Provider
chunks(rows, UPSERT_CHUNK_SIZE), chunks(rows, UPSERT_CHUNK_SIZE),
(chunk) => (chunk) =>
Effect.tryPromise({ Effect.tryPromise({
try: () => try: async () => {
db try {
.insert(providerStat) return await upsertProviderChunk(chunk, true)
.values(chunk) } catch (cause) {
.onDuplicateKeyUpdate({ if (!isMissingUniqueUsersColumn(cause)) throw cause
set: { return upsertProviderChunk(chunk, false)
sessions: inserted("sessions"), }
requests: inserted("requests"), },
unique_users: inserted("unique_users"),
input_tokens: inserted("input_tokens"),
output_tokens: inserted("output_tokens"),
reasoning_tokens: inserted("reasoning_tokens"),
cache_read_tokens: inserted("cache_read_tokens"),
total_tokens: inserted("total_tokens"),
input_cost_microcents: inserted("input_cost_microcents"),
output_cost_microcents: inserted("output_cost_microcents"),
total_cost_microcents: inserted("total_cost_microcents"),
avg_duration_ms: inserted("avg_duration_ms"),
p50_duration_ms: inserted("p50_duration_ms"),
p95_duration_ms: inserted("p95_duration_ms"),
avg_ttfb_ms: inserted("avg_ttfb_ms"),
p50_ttfb_ms: inserted("p50_ttfb_ms"),
p95_ttfb_ms: inserted("p95_ttfb_ms"),
avg_output_tps: inserted("avg_output_tps"),
success_count: inserted("success_count"),
error_count: inserted("error_count"),
sample_count: inserted("sample_count"),
market_share_tokens: inserted("market_share_tokens"),
market_share_requests: inserted("market_share_requests"),
market_share_sessions: inserted("market_share_sessions"),
rank_by_tokens: inserted("rank_by_tokens"),
rank_by_requests: inserted("rank_by_requests"),
rank_by_sessions: inserted("rank_by_sessions"),
rank_by_cost: inserted("rank_by_cost"),
},
}),
catch: (cause) => DatabaseError.make({ cause }), catch: (cause) => DatabaseError.make({ cause }),
}), }),
{ discard: true }, { discard: true },
) )
}) })
function upsertProviderChunk(chunk: ProviderStatRow[], includeUniqueUsers: boolean) {
return db
.insert(providerStat)
.values(includeUniqueUsers ? chunk : omitUniqueUsers(chunk))
.onDuplicateKeyUpdate({
set: {
sessions: inserted("sessions"),
requests: inserted("requests"),
...(includeUniqueUsers ? { unique_users: inserted("unique_users") } : {}),
input_tokens: inserted("input_tokens"),
output_tokens: inserted("output_tokens"),
reasoning_tokens: inserted("reasoning_tokens"),
cache_read_tokens: inserted("cache_read_tokens"),
total_tokens: inserted("total_tokens"),
input_cost_microcents: inserted("input_cost_microcents"),
output_cost_microcents: inserted("output_cost_microcents"),
total_cost_microcents: inserted("total_cost_microcents"),
avg_duration_ms: inserted("avg_duration_ms"),
p50_duration_ms: inserted("p50_duration_ms"),
p95_duration_ms: inserted("p95_duration_ms"),
avg_ttfb_ms: inserted("avg_ttfb_ms"),
p50_ttfb_ms: inserted("p50_ttfb_ms"),
p95_ttfb_ms: inserted("p95_ttfb_ms"),
avg_output_tps: inserted("avg_output_tps"),
success_count: inserted("success_count"),
error_count: inserted("error_count"),
sample_count: inserted("sample_count"),
market_share_tokens: inserted("market_share_tokens"),
market_share_requests: inserted("market_share_requests"),
market_share_sessions: inserted("market_share_sessions"),
rank_by_tokens: inserted("rank_by_tokens"),
rank_by_requests: inserted("rank_by_requests"),
rank_by_sessions: inserted("rank_by_sessions"),
rank_by_cost: inserted("rank_by_cost"),
},
})
}
const deleteRetiredDimensions = Effect.fn("ProviderStatRepo.deleteRetiredDimensions")(function* ( const deleteRetiredDimensions = Effect.fn("ProviderStatRepo.deleteRetiredDimensions")(function* (
rows: ProviderStatRow[], rows: ProviderStatRow[],
) { ) {
+21
View File
@@ -147,6 +147,18 @@ export function combineRows<T extends StatBaseRow>(left: T, right: T): T {
} }
} }
export function isMissingUniqueUsersColumn(cause: unknown): boolean {
return errorText(cause).includes("Unknown column 'unique_users'")
}
export function omitUniqueUsers<T extends { unique_users?: number }>(rows: T[]) {
return rows.map((row) => {
const result = { ...row }
delete result.unique_users
return result
})
}
export function statPeriodKey(row: StatBaseRow) { export function statPeriodKey(row: StatBaseRow) {
return [row.grain, row.period_key, row.dataset, row.tier, row.client, row.source].join("\u0000") return [row.grain, row.period_key, row.dataset, row.tier, row.client, row.source].join("\u0000")
} }
@@ -242,6 +254,15 @@ export function inserted(column: string) {
return sql.raw(`values(\`${column}\`)`) return sql.raw(`values(\`${column}\`)`)
} }
function errorText(cause: unknown): string {
if (cause instanceof Error) return `${cause.message} ${errorText((cause as { cause?: unknown }).cause)}`
if (typeof cause === "object" && cause)
return Object.values(cause as Record<string, unknown>)
.map(errorText)
.join(" ")
return String(cause)
}
export function weightedAverage( export function weightedAverage(
left: number | null | undefined, left: number | null | undefined,
leftWeight = 0, leftWeight = 0,