honeycomb-backfill.ts 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008
  1. import { Client } from "@planetscale/database"
  2. import { readdir } from "node:fs/promises"
  3. import path from "node:path"
  4. import { drizzle } from "drizzle-orm/planetscale-serverless"
  5. import { geoStat, modelStat, providerStat } from "./database/schema"
  6. import { statModel, statProvider } from "./domain/model-normalization"
  7. import {
  8. chunks,
  9. collapseRows,
  10. inserted,
  11. isoWeekId,
  12. normalizeCountry,
  13. normalizeTier,
  14. periodKeyFor,
  15. rankBy,
  16. rankRowsWithMarketShare,
  17. startOfIsoWeek,
  18. startOfUtcDay,
  19. statPeriodKey,
  20. synthesizeAllTierRows,
  21. toStatBaseRow,
  22. type StatBaseAggregate,
  23. } from "./domain/stat"
  24. const DAY_MS = 86_400_000
  25. const DEFAULT_UPSERT_CHUNK_SIZE = 100
  26. const DEFAULT_TIERS = ["Go", "Free", "Paid"]
  27. const FREE_MODELS = new Set(["gpt-5-nano", "grok-code", "big-pickle"])
  28. type Grain = "day" | "week"
  29. type MetricDimension = "model" | "provider" | "geo" | "geo-model"
  30. type LookupDimension = "model-provider-model" | "geo-continent"
  31. type ImportKey = `${MetricDimension | LookupDimension}-${Grain}`
  32. type QuerySpec = {
  33. name: string
  34. importKey: ImportKey
  35. importFlag: `--${ImportKey}`
  36. query: ReturnType<typeof metricQuery>
  37. }
  38. type RawRow = Record<string, string>
  39. type ImportOptions = {
  40. dataset: string
  41. databaseUrl: string | undefined
  42. directories: string[]
  43. dryRun: boolean
  44. periodStart: Date | undefined
  45. upsertChunkSize: number
  46. files: Partial<Record<ImportKey, string[]>>
  47. }
  48. type ModelAggregate = StatBaseAggregate & { provider: string; model: string; provider_model: string }
  49. type ProviderAggregate = StatBaseAggregate & { provider: string }
  50. type GeoAggregate = StatBaseAggregate & { provider: string; model: string; country: string; continent: string }
  51. type ModelStatRow = typeof modelStat.$inferInsert
  52. type ProviderStatRow = typeof providerStat.$inferInsert
  53. type GeoStatRow = typeof geoStat.$inferInsert
  54. const inputKeys = [
  55. "model-day",
  56. "model-week",
  57. "model-provider-model-day",
  58. "model-provider-model-week",
  59. "provider-day",
  60. "provider-week",
  61. "geo-day",
  62. "geo-week",
  63. "geo-model-day",
  64. "geo-model-week",
  65. "geo-continent-day",
  66. "geo-continent-week",
  67. ] as const satisfies ImportKey[]
  68. if (import.meta.main) await main()
  69. async function main() {
  70. const command = process.argv[2]
  71. if (command === "queries") return printQueries(process.argv.slice(3))
  72. if (command === "import") return importFiles(process.argv.slice(3))
  73. usage()
  74. }
  75. function printQueries(args: string[]) {
  76. const flags = parseFlags(args)
  77. const limit = parseIntegerFlag(flags, "limit") ?? 1000
  78. const tiers = parseListFlag(flags, "tiers") ?? DEFAULT_TIERS
  79. const queries = buildQueries(limit, tiers)
  80. const only = flags.get("only")?.[0]
  81. if (only) {
  82. const item = queries.find((query) => query.name === only)
  83. if (!item) fail(`Unknown --only ${only}. Expected one of: ${queries.map((query) => query.name).join(", ")}`)
  84. console.log(JSON.stringify(item.query, null, 2))
  85. return
  86. }
  87. console.log(
  88. JSON.stringify(
  89. {
  90. tiers,
  91. import_hint: "bun src/honeycomb-backfill.ts import --dir downloads",
  92. queries,
  93. },
  94. null,
  95. 2,
  96. ),
  97. )
  98. }
  99. async function importFiles(args: string[]) {
  100. const parsed = parseImportOptions(args)
  101. const opts = { ...parsed, files: mergeFiles(parsed.files, await discoverFiles(parsed.directories)) }
  102. if (!inputKeys.some((key) => opts.files[key]?.length)) fail("No CSV or JSON import files were provided or discovered")
  103. const providerModelLookup = new Map([
  104. ...(await lookupRows(opts.files["model-provider-model-day"], "day", opts, modelProviderModelLookup)),
  105. ...(await lookupRows(opts.files["model-provider-model-week"], "week", opts, modelProviderModelLookup)),
  106. ])
  107. const continentLookup = new Map([
  108. ...(await lookupRows(opts.files["geo-continent-day"], "day", opts, geoContinentLookup)),
  109. ...(await lookupRows(opts.files["geo-continent-week"], "week", opts, geoContinentLookup)),
  110. ])
  111. const modelAggregates = [
  112. ...(await metricRows(opts.files["model-day"], "day", opts, (row, base) =>
  113. modelAggregate(row, base, providerModelLookup),
  114. )),
  115. ...(await metricRows(opts.files["model-week"], "week", opts, (row, base) =>
  116. modelAggregate(row, base, providerModelLookup),
  117. )),
  118. ]
  119. const modelRows = modelRowsFromAggregates(modelAggregates)
  120. const providerRows = providerRowsFromAggregates([
  121. ...(await metricRows(opts.files["provider-day"], "day", opts, (row, base) => ({
  122. ...base,
  123. provider: provider(row) ?? "unknown",
  124. }))),
  125. ...(await metricRows(opts.files["provider-week"], "week", opts, (row, base) => ({
  126. ...base,
  127. provider: provider(row) ?? "unknown",
  128. }))),
  129. ])
  130. const geoRows = geoRowsFromAggregates([
  131. ...(await metricRows(opts.files["geo-day"], "day", opts, (row, base) => ({
  132. ...base,
  133. provider: "all",
  134. model: "all",
  135. country: country(row),
  136. continent: continentLookup.get(lookupKey(base, country(row))) ?? continent(row),
  137. }))),
  138. ...(await metricRows(opts.files["geo-week"], "week", opts, (row, base) => ({
  139. ...base,
  140. provider: "all",
  141. model: "all",
  142. country: country(row),
  143. continent: continentLookup.get(lookupKey(base, country(row))) ?? continent(row),
  144. }))),
  145. ...(await metricRows(opts.files["geo-model-day"], "day", opts, (row, base) =>
  146. geoModelAggregate(row, base, continentLookup),
  147. )),
  148. ...(await metricRows(opts.files["geo-model-week"], "week", opts, (row, base) =>
  149. geoModelAggregate(row, base, continentLookup),
  150. )),
  151. ])
  152. console.log(
  153. JSON.stringify(
  154. {
  155. inputs: Object.fromEntries(
  156. inputKeys.flatMap((key) => (opts.files[key]?.length ? [[key, opts.files[key].length]] : [])),
  157. ),
  158. modelRows: modelRows.length,
  159. providerRows: providerRows.length,
  160. geoRows: geoRows.length,
  161. dryRun: opts.dryRun,
  162. upsertChunkSize: opts.upsertChunkSize,
  163. },
  164. null,
  165. 2,
  166. ),
  167. )
  168. if (opts.dryRun) return
  169. if (!opts.databaseUrl) fail("DATABASE_URL is required unless --dry-run is set")
  170. const db = drizzle({ client: new Client({ url: opts.databaseUrl }) })
  171. await upsertModelRows(db, modelRows, opts.upsertChunkSize)
  172. await upsertProviderRows(db, providerRows, opts.upsertChunkSize)
  173. await upsertGeoRows(db, geoRows, opts.upsertChunkSize)
  174. }
  175. function buildQueries(limit: number, tiers: string[]): QuerySpec[] {
  176. const daily = tiers.flatMap((tier) => [
  177. querySpec(
  178. "model-day",
  179. tier,
  180. metricQuery(["date", "tier", "stat_provider_2", "stat_model_2"], limit, tierFilters(tier)),
  181. ),
  182. querySpec("provider-day", tier, metricQuery(["date", "tier", "stat_provider_2"], limit, tierFilters(tier))),
  183. querySpec("geo-day", tier, metricQuery(["date", "tier", "country", "continent"], limit, tierFilters(tier))),
  184. querySpec(
  185. "geo-model-day",
  186. tier,
  187. metricQuery(
  188. ["date", "tier", "stat_provider_2", "stat_model_2", "country", "continent"],
  189. limit,
  190. tierFilters(tier),
  191. ),
  192. ),
  193. ])
  194. const weekly = tiers.flatMap((tier) => [
  195. querySpec(
  196. "model-week",
  197. tier,
  198. metricQuery(["week", "tier", "stat_provider_2", "stat_model_2"], limit, tierFilters(tier)),
  199. ),
  200. querySpec("provider-week", tier, metricQuery(["week", "tier", "stat_provider_2"], limit, tierFilters(tier))),
  201. querySpec("geo-week", tier, metricQuery(["week", "tier", "country", "continent"], limit, tierFilters(tier))),
  202. querySpec(
  203. "geo-model-week",
  204. tier,
  205. metricQuery(
  206. ["week", "tier", "stat_provider_2", "stat_model_2", "country", "continent"],
  207. limit,
  208. tierFilters(tier),
  209. ),
  210. ),
  211. ])
  212. return [...daily, ...weekly]
  213. }
  214. function querySpec(importKey: ImportKey, tier: string, query: ReturnType<typeof metricQuery>) {
  215. return {
  216. name: `${importKey}-${queryNameSegment(tier)}`,
  217. importKey,
  218. importFlag: `--${importKey}` as const,
  219. query,
  220. }
  221. }
  222. function metricQuery(breakdowns: string[], limit: number, filters: ReturnType<typeof commonFilters> = []) {
  223. return {
  224. granularity: 0,
  225. breakdowns,
  226. calculations: [
  227. { op: "COUNT_DISTINCT", column: "session" },
  228. { op: "COUNT" },
  229. { op: "COUNT_DISTINCT", column: "workspace" },
  230. { op: "SUM", column: "tokens.input" },
  231. { op: "SUM", column: "tokens.output" },
  232. { op: "SUM", column: "tokens.reasoning" },
  233. { op: "SUM", column: "tokens.cache_read" },
  234. { op: "SUM", column: "tokens" },
  235. { op: "SUM", column: "cost.input.microcents" },
  236. { op: "SUM", column: "cost.output.microcents" },
  237. { op: "SUM", column: "cost.total.microcents" },
  238. { op: "AVG", column: "duration" },
  239. { op: "P50", column: "duration" },
  240. { op: "P95", column: "duration" },
  241. { op: "AVG", column: "time_to_first_byte" },
  242. { op: "P50", column: "time_to_first_byte" },
  243. { op: "P95", column: "time_to_first_byte" },
  244. { op: "AVG", column: "tps.output" },
  245. ],
  246. filters: [...commonFilters(), ...filters],
  247. filter_combination: "AND",
  248. orders: [{ column: "tokens", op: "SUM", order: "descending" }],
  249. havings: [],
  250. limit,
  251. formulas: [],
  252. }
  253. }
  254. function tierFilters(tier: string) {
  255. if (tier === "all") return []
  256. return [{ column: "tier", op: "=", value: tier }]
  257. }
  258. function queryNameSegment(value: string) {
  259. return (
  260. value
  261. .toLowerCase()
  262. .replace(/[^a-z0-9]+/g, "-")
  263. .replace(/^-|-$/g, "") || "all"
  264. )
  265. }
  266. function commonFilters() {
  267. return [
  268. { column: "event_type", op: "=", value: "completions" },
  269. { column: "model", op: "exists" },
  270. { column: "model", op: "!=", value: "" },
  271. { column: "model", op: "!=", value: "alpha-gpt-next" },
  272. ]
  273. }
  274. function metricRows<T extends StatBaseAggregate>(
  275. files: string[] | undefined,
  276. grain: Grain,
  277. opts: ImportOptions,
  278. map: (row: RawRow, base: StatBaseAggregate) => T | T[],
  279. ) {
  280. if (!files) return Promise.resolve([])
  281. return readFiles(files).then((rows) => rows.flatMap((row) => map(row, baseAggregate(row, grain, opts))))
  282. }
  283. function lookupRows(
  284. files: string[] | undefined,
  285. grain: Grain,
  286. opts: ImportOptions,
  287. map: (row: RawRow, grain: Grain, opts: ImportOptions) => readonly (readonly [string, string])[],
  288. ) {
  289. if (!files) return Promise.resolve([])
  290. return readFiles(files).then((rows) =>
  291. Array.from(
  292. rows
  293. .flatMap((row) => map(row, grain, opts))
  294. .reduce((result, [key, value]) => {
  295. if (value && value > (result.get(key) ?? "")) result.set(key, value)
  296. return result
  297. }, new Map<string, string>()),
  298. ),
  299. )
  300. }
  301. async function readFiles(files: string[]) {
  302. return (await Promise.all(files.map(readRows))).flat()
  303. }
  304. async function discoverFiles(directories: string[]) {
  305. const classified = await Promise.all(
  306. (await Promise.all(directories.map(filesInDirectory))).flat().map(async (file) => ({
  307. file,
  308. key: classifyRows(file, await readRows(file)),
  309. })),
  310. )
  311. return classified.reduce<Partial<Record<ImportKey, string[]>>>((result, item) => {
  312. return { ...result, [item.key]: [...(result[item.key] ?? []), item.file] }
  313. }, {})
  314. }
  315. async function filesInDirectory(directory: string): Promise<string[]> {
  316. return (
  317. await Promise.all(
  318. (await readdir(directory, { withFileTypes: true })).map((entry) => {
  319. const file = path.join(directory, entry.name)
  320. if (entry.isDirectory()) return filesInDirectory(file)
  321. if (entry.isFile() && /\.(csv|json)$/i.test(entry.name)) return Promise.resolve([file])
  322. return Promise.resolve([])
  323. }),
  324. )
  325. ).flat()
  326. }
  327. function classifyRows(file: string, rows: RawRow[]): ImportKey {
  328. if (rows.length === 0) fail(`Cannot classify empty export: ${file}`)
  329. const headers = new Set(rows.flatMap((row) => Object.keys(row).map(normalizeHeader)))
  330. const grain: Grain = headers.has("date") ? "day" : "week"
  331. if (hasHeader(headers, ["country", "cf.country"])) {
  332. if (hasHeader(headers, ["model", "stat_model", "stat_model_2"]) && hasMetricHeaders(headers))
  333. return `geo-model-${grain}`
  334. return hasMetricHeaders(headers) ? `geo-${grain}` : `geo-continent-${grain}`
  335. }
  336. if (hasHeader(headers, ["model", "stat_model", "stat_model_2"]))
  337. return hasMetricHeaders(headers) ? `model-${grain}` : `model-provider-model-${grain}`
  338. if (
  339. hasHeader(headers, [
  340. "provider",
  341. "provider.normalized",
  342. "stat_provider",
  343. "stat_provider_2",
  344. "provider.model",
  345. "provider_model",
  346. ])
  347. )
  348. return `provider-${grain}`
  349. fail(`Cannot classify export from columns in ${file}`)
  350. }
  351. function hasMetricHeaders(headers: Set<string>) {
  352. return [
  353. "sumtokens",
  354. "sumtokensinput",
  355. "inputtokens",
  356. "totaltokens",
  357. "avgduration",
  358. "countdistinctsession",
  359. "countdistinctworkspace",
  360. ].some((header) => headers.has(header))
  361. }
  362. function hasHeader(headers: Set<string>, names: string[]) {
  363. return names.some((name) => headers.has(normalizeHeader(name)))
  364. }
  365. function mergeFiles(left: Partial<Record<ImportKey, string[]>>, right: Partial<Record<ImportKey, string[]>>) {
  366. return inputKeys.reduce<Partial<Record<ImportKey, string[]>>>((result, key) => {
  367. const files = [...(left[key] ?? []), ...(right[key] ?? [])]
  368. if (files.length === 0) return result
  369. return { ...result, [key]: files }
  370. }, {})
  371. }
  372. function modelProviderModelLookup(row: RawRow, grain: Grain, opts: ImportOptions): [string, string][] {
  373. const base = basePeriod(row, grain, opts)
  374. const value = providerModel(row)
  375. const author = provider(row)
  376. if (!value || !author) return []
  377. return [[lookupKey({ ...base, dataset: opts.dataset, tier: tier(row), grain }, author, model(row)), value]]
  378. }
  379. function modelAggregate(
  380. row: RawRow,
  381. base: StatBaseAggregate,
  382. providerModelLookup: Map<string, string>,
  383. ): ModelAggregate[] {
  384. const author = provider(row)
  385. if (!author) return []
  386. return [
  387. {
  388. ...base,
  389. provider: author,
  390. model: model(row),
  391. provider_model: providerModelLookup.get(lookupKey(base, author, model(row))) ?? providerModel(row),
  392. },
  393. ]
  394. }
  395. function geoContinentLookup(row: RawRow, grain: Grain, opts: ImportOptions): [string, string][] {
  396. const base = basePeriod(row, grain, opts)
  397. const value = continent(row)
  398. if (!value) return []
  399. return [[lookupKey({ ...base, dataset: opts.dataset, tier: tier(row), grain }, country(row)), value]]
  400. }
  401. function geoModelAggregate(row: RawRow, base: StatBaseAggregate, continentLookup: Map<string, string>): GeoAggregate[] {
  402. const author = provider(row)
  403. if (!author) return []
  404. return [
  405. {
  406. ...base,
  407. provider: author,
  408. model: model(row),
  409. country: country(row),
  410. continent: continentLookup.get(lookupKey(base, country(row))) ?? continent(row),
  411. },
  412. ]
  413. }
  414. function baseAggregate(row: RawRow, grain: Grain, opts: ImportOptions): StatBaseAggregate {
  415. return {
  416. ...basePeriod(row, grain, opts),
  417. grain,
  418. dataset: opts.dataset,
  419. tier: tier(row),
  420. sessions: integer(row, "sessions", ["COUNT_DISTINCT(session)"]),
  421. requests: integer(row, "requests", ["COUNT", "COUNT()"]),
  422. unique_users: integer(row, "unique_users", [
  423. "COUNT_DISTINCT(user_id)",
  424. "COUNT_DISTINCT(workspace)",
  425. "COUNT_DISTINCT(api_key)",
  426. ]),
  427. input_tokens: integer(row, "input_tokens", ["SUM(tokens.input)", "SUM(tokens_input)"]),
  428. output_tokens: integer(row, "output_tokens", ["SUM(tokens.output)", "SUM(tokens_output)"]),
  429. reasoning_tokens: integer(row, "reasoning_tokens", ["SUM(tokens.reasoning)", "SUM(tokens_reasoning)"]),
  430. cache_read_tokens: integer(row, "cache_read_tokens", ["SUM(tokens.cache_read)", "SUM(tokens_cache_read)"]),
  431. total_tokens: integer(row, "total_tokens", ["SUM(stat_tokens_total)", "SUM(tokens)", "SUM(tokens_total)"]),
  432. input_cost_microcents: integer(row, "input_cost_microcents", [
  433. "SUM(cost.input.microcents)",
  434. "SUM(stat_cost_input_microcents)",
  435. ]),
  436. output_cost_microcents: integer(row, "output_cost_microcents", [
  437. "SUM(cost.output.microcents)",
  438. "SUM(stat_cost_output_microcents)",
  439. ]),
  440. total_cost_microcents: integer(row, "total_cost_microcents", [
  441. "SUM(cost.total.microcents)",
  442. "SUM(stat_cost_total_microcents)",
  443. ]),
  444. avg_duration_ms: nullableNumber(row, "avg_duration_ms", ["AVG(duration)", "AVG(duration_ms)"]),
  445. p50_duration_ms: nullableInteger(row, "p50_duration_ms", ["P50(duration)", "P50(duration_ms)"]),
  446. p95_duration_ms: nullableInteger(row, "p95_duration_ms", ["P95(duration)", "P95(duration_ms)"]),
  447. avg_ttfb_ms: nullableNumber(row, "avg_ttfb_ms", ["AVG(time_to_first_byte)", "AVG(ttfb_ms)"]),
  448. p50_ttfb_ms: nullableInteger(row, "p50_ttfb_ms", ["P50(time_to_first_byte)", "P50(ttfb_ms)"]),
  449. p95_ttfb_ms: nullableInteger(row, "p95_ttfb_ms", ["P95(time_to_first_byte)", "P95(ttfb_ms)"]),
  450. avg_output_tps: nullableNumber(row, "avg_output_tps", ["AVG(tps.output)", "AVG(stat_output_tps)"]),
  451. success_count: integer(row, "success_count", ["SUM(success)", "SUM(is_success)", "SUM(stat_success)"]),
  452. error_count: integer(row, "error_count", ["SUM(error)", "SUM(is_error)", "SUM(stat_error)"]),
  453. sample_count: integer(row, "sample_count", ["COUNT", "COUNT()"]),
  454. }
  455. }
  456. function basePeriod(row: RawRow, grain: Grain, opts: ImportOptions) {
  457. return { period_key: periodKey(row, grain, opts) }
  458. }
  459. function periodKey(row: RawRow, grain: Grain, opts: ImportOptions) {
  460. if (grain === "week") {
  461. const week = parseWeek(row)
  462. if (week) return week
  463. fail("weekly imports require a week or period_key column")
  464. }
  465. const time = parseTime(row)
  466. const start = time ? startOfUtcDay(time) : opts.periodStart
  467. if (!start) fail("daily imports require a time column or --period-start")
  468. return periodKeyFor("day", start)
  469. }
  470. function modelRowsFromAggregates(aggregates: ModelAggregate[]) {
  471. return rankModelRows([
  472. ...synthesizeAllTierRows(
  473. collapseRows(aggregates.filter((item) => item.grain === "week").map(toModelRow), modelDimensionKey),
  474. modelDimensionKey,
  475. ),
  476. ...synthesizeAllTierRows(
  477. collapseRows(aggregates.filter((item) => item.grain === "day").map(toModelRow), modelDimensionKey),
  478. modelDimensionKey,
  479. ),
  480. ])
  481. }
  482. function providerRowsFromAggregates(aggregates: ProviderAggregate[]) {
  483. return rankRowsWithMarketShare([
  484. ...synthesizeAllTierRows(
  485. collapseRows(aggregates.filter((item) => item.grain === "week").map(toProviderRow), providerDimensionKey),
  486. providerDimensionKey,
  487. ),
  488. ...synthesizeAllTierRows(
  489. collapseRows(aggregates.filter((item) => item.grain === "day").map(toProviderRow), providerDimensionKey),
  490. providerDimensionKey,
  491. ),
  492. ])
  493. }
  494. function geoRowsFromAggregates(aggregates: GeoAggregate[]) {
  495. return rankRowsWithMarketShare(
  496. [
  497. ...synthesizeAllTierRows(
  498. collapseRows(aggregates.filter((item) => item.grain === "week").map(toGeoRow), geoDimensionKey),
  499. geoDimensionKey,
  500. ),
  501. ...synthesizeAllTierRows(
  502. collapseRows(aggregates.filter((item) => item.grain === "day").map(toGeoRow), geoDimensionKey),
  503. geoDimensionKey,
  504. ),
  505. ],
  506. geoMarketShareKey,
  507. )
  508. }
  509. function toModelRow(data: ModelAggregate): ModelStatRow {
  510. return { ...toStatBaseRow(data), provider: data.provider, model: data.model, provider_model: data.provider_model }
  511. }
  512. function toProviderRow(data: ProviderAggregate): ProviderStatRow {
  513. return { ...toStatBaseRow(data), provider: data.provider }
  514. }
  515. function toGeoRow(data: GeoAggregate): GeoStatRow {
  516. return {
  517. ...toStatBaseRow(data),
  518. provider: data.provider,
  519. model: data.model,
  520. country: data.country,
  521. continent: data.continent,
  522. }
  523. }
  524. function rankModelRows(rows: ModelStatRow[]) {
  525. return Object.values(
  526. rows.reduce<Record<string, ModelStatRow[]>>((result, row) => {
  527. const key = statPeriodKey(row)
  528. result[key] = [...(result[key] ?? []), row]
  529. return result
  530. }, {}),
  531. ).flatMap((group) => {
  532. const tokenRanks = rankBy(group, (row) => row.total_tokens ?? 0)
  533. const requestRanks = rankBy(group, (row) => row.requests ?? 0)
  534. const costRanks = rankBy(group, (row) => row.total_cost_microcents ?? 0)
  535. return group.map((row) => ({
  536. ...row,
  537. rank_by_tokens: tokenRanks.get(row) ?? null,
  538. rank_by_requests: requestRanks.get(row) ?? null,
  539. rank_by_cost: costRanks.get(row) ?? null,
  540. }))
  541. })
  542. }
  543. function modelDimensionKey(row: ModelStatRow) {
  544. return [row.provider, row.model].join("\u0000")
  545. }
  546. function providerDimensionKey(row: ProviderStatRow) {
  547. return row.provider
  548. }
  549. function geoDimensionKey(row: GeoStatRow) {
  550. return [row.provider, row.model, row.country].join("\u0000")
  551. }
  552. function geoMarketShareKey(row: GeoStatRow) {
  553. return [statPeriodKey(row), row.provider, row.model].join("\u0000")
  554. }
  555. function lookupKey(base: { grain: string; period_key: string; dataset: string; tier: string }, ...dimension: string[]) {
  556. return [base.grain, base.period_key, base.dataset, base.tier, ...dimension].join("\u0000")
  557. }
  558. function tier(row: RawRow) {
  559. return normalizeTier(cell(row, ["stat_tier", "tier"]) || deriveTier(row))
  560. }
  561. function deriveTier(row: RawRow) {
  562. const source = cell(row, ["source"])
  563. const value = model(row)
  564. if (source === "lite") return "Go"
  565. if (FREE_MODELS.has(value) || /-free(:global)?$/.test(rawModel(row))) return "Free"
  566. return "Zen"
  567. }
  568. function provider(row: RawRow) {
  569. return statProvider(model(row), providerModel(row), cell(row, ["stat_provider_2", "stat_provider"]))
  570. }
  571. function model(row: RawRow) {
  572. return statModel(cell(row, ["stat_model_2", "stat_model"]) || rawModel(row), providerModel(row))
  573. }
  574. function rawModel(row: RawRow) {
  575. return cell(row, ["model"]) || "unknown"
  576. }
  577. function providerModel(row: RawRow) {
  578. return cell(row, ["provider.model", "provider_model"]) || ""
  579. }
  580. function country(row: RawRow) {
  581. return normalizeCountry(cell(row, ["stat_country", "cf.country", "cf_country", "country"]))
  582. }
  583. function continent(row: RawRow) {
  584. return cell(row, ["cf.continent", "cf_continent", "continent"]) || ""
  585. }
  586. function integer(row: RawRow, name: string, aliases: string[] = []) {
  587. return Math.round(number(row, name, aliases))
  588. }
  589. function nullableInteger(row: RawRow, name: string, aliases: string[] = []) {
  590. if (!hasCell(row, [name, ...aliases])) return null
  591. return Math.round(number(row, name, aliases))
  592. }
  593. function nullableNumber(row: RawRow, name: string, aliases: string[] = []) {
  594. if (!hasCell(row, [name, ...aliases])) return null
  595. return Number(number(row, name, aliases).toFixed(2))
  596. }
  597. function number(row: RawRow, name: string, aliases: string[] = []) {
  598. const value = Number(cell(row, [name, ...aliases]).replace(/,/g, ""))
  599. return Number.isFinite(value) ? value : 0
  600. }
  601. function hasCell(row: RawRow, names: string[]) {
  602. return names.some((name) => row[name] !== undefined && row[name] !== "")
  603. }
  604. function cell(row: RawRow, names: string[]) {
  605. const normalized = normalizedCells(row)
  606. return (
  607. names.flatMap((name) => [row[name], normalized.get(normalizeHeader(name))]).find((value) => value !== undefined) ??
  608. ""
  609. )
  610. }
  611. function normalizedCells(row: RawRow) {
  612. return new Map(Object.entries(row).map(([key, value]) => [normalizeHeader(key), value]))
  613. }
  614. function normalizeHeader(value: string) {
  615. return value.toLowerCase().replace(/[^a-z0-9]+/g, "")
  616. }
  617. function parseTime(row: RawRow) {
  618. const value = cell(row, ["date", "time", "timestamp", "datetime", "bucket"])
  619. if (!value) return undefined
  620. const numeric = Number(value)
  621. const date = Number.isFinite(numeric)
  622. ? new Date(numeric > 10_000_000_000 ? numeric : numeric * 1000)
  623. : new Date(value)
  624. if (Number.isNaN(date.getTime())) fail(`Invalid time value: ${value}`)
  625. return date
  626. }
  627. function parseWeek(row: RawRow) {
  628. const value = cell(row, ["period_key", "week", "stat_week"])
  629. if (!value) return undefined
  630. const match = /^(\d{4})-W(\d{1,2})$/.exec(value)
  631. if (!match) fail(`Invalid week value: ${value}`)
  632. const year = Number(match[1])
  633. const week = Number(match[2])
  634. if (week < 1 || week > 53) fail(`Invalid week value: ${value}`)
  635. const start = new Date(startOfIsoWeek(new Date(Date.UTC(year, 0, 4))).getTime() + (week - 1) * 7 * DAY_MS)
  636. const id = `${year}-W${String(week).padStart(2, "0")}`
  637. if (isoWeekId(start) !== id) fail(`Invalid week value: ${value}`)
  638. return id
  639. }
  640. async function readRows(file: string) {
  641. const text = await Bun.file(file).text()
  642. if (file.toLowerCase().endsWith(".json")) {
  643. const parsed: unknown = JSON.parse(text)
  644. return rowsFromJson(parsed)
  645. }
  646. return rowsFromCsv(text)
  647. }
  648. function rowsFromJson(value: unknown): RawRow[] {
  649. if (Array.isArray(value)) return value.flatMap(rowFromUnknown)
  650. if (!isRecord(value)) fail("JSON imports must be an array of rows or an object with results/data/rows")
  651. const rows = [value.results, value.data, value.rows].flatMap((candidate) =>
  652. Array.isArray(candidate) ? candidate.flatMap(rowFromUnknown) : [],
  653. )
  654. if (rows.length === 0) fail("JSON import did not contain rows")
  655. return rows
  656. }
  657. function rowFromUnknown(value: unknown): RawRow[] {
  658. if (!isRecord(value)) return []
  659. const nested = isRecord(value.data) ? value.data : {}
  660. return [
  661. Object.fromEntries(
  662. Object.entries({ ...value, ...nested }).flatMap(([key, item]) => {
  663. if (key === "data") return []
  664. return [[key, cellValue(item)]]
  665. }),
  666. ),
  667. ]
  668. }
  669. function rowsFromCsv(text: string): RawRow[] {
  670. const [headers, ...rows] = csvRecords(text).filter((row) => row.some((value) => value.trim() !== ""))
  671. if (!headers) return []
  672. return rows.map((row) =>
  673. Object.fromEntries(headers.map((header, index) => [header.trim(), row[index]?.trim() ?? ""])),
  674. )
  675. }
  676. function csvRecords(text: string) {
  677. const rows: string[][] = []
  678. let row: string[] = []
  679. let field = ""
  680. let quoted = false
  681. for (let index = 0; index < text.length; index++) {
  682. const char = text[index]
  683. const next = text[index + 1]
  684. if (quoted) {
  685. if (char === '"' && next === '"') {
  686. field += '"'
  687. index++
  688. continue
  689. }
  690. if (char === '"') {
  691. quoted = false
  692. continue
  693. }
  694. field += char
  695. continue
  696. }
  697. if (char === '"') {
  698. quoted = true
  699. continue
  700. }
  701. if (char === ",") {
  702. row.push(field)
  703. field = ""
  704. continue
  705. }
  706. if (char === "\n") {
  707. row.push(field)
  708. rows.push(row)
  709. row = []
  710. field = ""
  711. continue
  712. }
  713. if (char === "\r") continue
  714. field += char
  715. }
  716. row.push(field)
  717. rows.push(row)
  718. return rows
  719. }
  720. function cellValue(value: unknown) {
  721. if (value === null || value === undefined) return ""
  722. if (typeof value === "string") return value
  723. if (typeof value === "number" || typeof value === "boolean" || typeof value === "bigint") return String(value)
  724. return JSON.stringify(value) ?? ""
  725. }
  726. function isRecord(value: unknown): value is Record<string, unknown> {
  727. return typeof value === "object" && value !== null && !Array.isArray(value)
  728. }
  729. async function upsertModelRows(db: ReturnType<typeof drizzle>, rows: ModelStatRow[], chunkSize: number) {
  730. const batches = chunks(rows, chunkSize)
  731. console.log(JSON.stringify({ table: "model_stat", batches: batches.length, chunkSize }))
  732. for (const chunk of batches) {
  733. await db
  734. .insert(modelStat)
  735. .values(chunk)
  736. .onDuplicateKeyUpdate({
  737. set: {
  738. provider_model: inserted("provider_model"),
  739. sessions: inserted("sessions"),
  740. requests: inserted("requests"),
  741. unique_users: inserted("unique_users"),
  742. input_tokens: inserted("input_tokens"),
  743. output_tokens: inserted("output_tokens"),
  744. reasoning_tokens: inserted("reasoning_tokens"),
  745. cache_read_tokens: inserted("cache_read_tokens"),
  746. total_tokens: inserted("total_tokens"),
  747. input_cost_microcents: inserted("input_cost_microcents"),
  748. output_cost_microcents: inserted("output_cost_microcents"),
  749. total_cost_microcents: inserted("total_cost_microcents"),
  750. avg_duration_ms: inserted("avg_duration_ms"),
  751. p50_duration_ms: inserted("p50_duration_ms"),
  752. p95_duration_ms: inserted("p95_duration_ms"),
  753. avg_ttfb_ms: inserted("avg_ttfb_ms"),
  754. p50_ttfb_ms: inserted("p50_ttfb_ms"),
  755. p95_ttfb_ms: inserted("p95_ttfb_ms"),
  756. avg_output_tps: inserted("avg_output_tps"),
  757. success_count: inserted("success_count"),
  758. error_count: inserted("error_count"),
  759. sample_count: inserted("sample_count"),
  760. rank_by_tokens: inserted("rank_by_tokens"),
  761. rank_by_requests: inserted("rank_by_requests"),
  762. rank_by_cost: inserted("rank_by_cost"),
  763. },
  764. })
  765. }
  766. }
  767. async function upsertProviderRows(db: ReturnType<typeof drizzle>, rows: ProviderStatRow[], chunkSize: number) {
  768. const batches = chunks(rows, chunkSize)
  769. console.log(JSON.stringify({ table: "provider_stat", batches: batches.length, chunkSize }))
  770. for (const chunk of batches) {
  771. await db
  772. .insert(providerStat)
  773. .values(chunk)
  774. .onDuplicateKeyUpdate({
  775. set: {
  776. sessions: inserted("sessions"),
  777. requests: inserted("requests"),
  778. unique_users: inserted("unique_users"),
  779. input_tokens: inserted("input_tokens"),
  780. output_tokens: inserted("output_tokens"),
  781. reasoning_tokens: inserted("reasoning_tokens"),
  782. cache_read_tokens: inserted("cache_read_tokens"),
  783. total_tokens: inserted("total_tokens"),
  784. input_cost_microcents: inserted("input_cost_microcents"),
  785. output_cost_microcents: inserted("output_cost_microcents"),
  786. total_cost_microcents: inserted("total_cost_microcents"),
  787. avg_duration_ms: inserted("avg_duration_ms"),
  788. p50_duration_ms: inserted("p50_duration_ms"),
  789. p95_duration_ms: inserted("p95_duration_ms"),
  790. avg_ttfb_ms: inserted("avg_ttfb_ms"),
  791. p50_ttfb_ms: inserted("p50_ttfb_ms"),
  792. p95_ttfb_ms: inserted("p95_ttfb_ms"),
  793. avg_output_tps: inserted("avg_output_tps"),
  794. success_count: inserted("success_count"),
  795. error_count: inserted("error_count"),
  796. sample_count: inserted("sample_count"),
  797. market_share_tokens: inserted("market_share_tokens"),
  798. market_share_requests: inserted("market_share_requests"),
  799. market_share_sessions: inserted("market_share_sessions"),
  800. rank_by_tokens: inserted("rank_by_tokens"),
  801. rank_by_requests: inserted("rank_by_requests"),
  802. rank_by_sessions: inserted("rank_by_sessions"),
  803. rank_by_cost: inserted("rank_by_cost"),
  804. },
  805. })
  806. }
  807. }
  808. async function upsertGeoRows(db: ReturnType<typeof drizzle>, rows: GeoStatRow[], chunkSize: number) {
  809. const batches = chunks(rows, chunkSize)
  810. console.log(JSON.stringify({ table: "geo_stat", batches: batches.length, chunkSize }))
  811. for (const chunk of batches) {
  812. await db
  813. .insert(geoStat)
  814. .values(chunk)
  815. .onDuplicateKeyUpdate({
  816. set: {
  817. continent: inserted("continent"),
  818. sessions: inserted("sessions"),
  819. requests: inserted("requests"),
  820. unique_users: inserted("unique_users"),
  821. input_tokens: inserted("input_tokens"),
  822. output_tokens: inserted("output_tokens"),
  823. reasoning_tokens: inserted("reasoning_tokens"),
  824. cache_read_tokens: inserted("cache_read_tokens"),
  825. total_tokens: inserted("total_tokens"),
  826. input_cost_microcents: inserted("input_cost_microcents"),
  827. output_cost_microcents: inserted("output_cost_microcents"),
  828. total_cost_microcents: inserted("total_cost_microcents"),
  829. avg_duration_ms: inserted("avg_duration_ms"),
  830. p50_duration_ms: inserted("p50_duration_ms"),
  831. p95_duration_ms: inserted("p95_duration_ms"),
  832. avg_ttfb_ms: inserted("avg_ttfb_ms"),
  833. p50_ttfb_ms: inserted("p50_ttfb_ms"),
  834. p95_ttfb_ms: inserted("p95_ttfb_ms"),
  835. avg_output_tps: inserted("avg_output_tps"),
  836. success_count: inserted("success_count"),
  837. error_count: inserted("error_count"),
  838. sample_count: inserted("sample_count"),
  839. market_share_tokens: inserted("market_share_tokens"),
  840. market_share_requests: inserted("market_share_requests"),
  841. market_share_sessions: inserted("market_share_sessions"),
  842. rank_by_tokens: inserted("rank_by_tokens"),
  843. rank_by_requests: inserted("rank_by_requests"),
  844. rank_by_sessions: inserted("rank_by_sessions"),
  845. rank_by_cost: inserted("rank_by_cost"),
  846. },
  847. })
  848. }
  849. }
  850. function parseImportOptions(args: string[]): ImportOptions {
  851. const flags = parseFlags(args)
  852. const files = inputKeys.reduce<Partial<Record<ImportKey, string[]>>>((result, key) => {
  853. const values = flags.get(key)
  854. if (!values) return result
  855. return { ...result, [key]: values }
  856. }, {})
  857. return {
  858. dataset: flags.get("dataset")?.[0] ?? "zen",
  859. databaseUrl: flags.get("database-url")?.[0] ?? process.env.DATABASE_URL,
  860. directories: flags.get("dir") ?? flags.get("directory") ?? [],
  861. dryRun: flags.has("dry-run"),
  862. periodStart: parseDateFlag(flags, "period-start"),
  863. upsertChunkSize: parseIntegerFlag(flags, "upsert-chunk-size") ?? DEFAULT_UPSERT_CHUNK_SIZE,
  864. files,
  865. }
  866. }
  867. function parseFlags(args: string[]) {
  868. const result = new Map<string, string[]>()
  869. for (let index = 0; index < args.length; index++) {
  870. const arg = args[index]
  871. if (!arg.startsWith("--")) fail(`Unexpected argument: ${arg}`)
  872. const name = arg.slice(2)
  873. if (name === "dry-run" || name === "include-weekly") {
  874. result.set(name, ["true"])
  875. continue
  876. }
  877. const nextFlag = args.findIndex((value, valueIndex) => valueIndex > index && value.startsWith("--"))
  878. const values = args.slice(index + 1, nextFlag === -1 ? args.length : nextFlag)
  879. if (values.length === 0) fail(`Missing value for --${name}`)
  880. result.set(name, [...(result.get(name) ?? []), ...values])
  881. index += values.length
  882. }
  883. return result
  884. }
  885. function parseDateFlag(flags: Map<string, string[]>, name: string) {
  886. const value = flags.get(name)?.[0]
  887. if (!value) return undefined
  888. const date = new Date(value)
  889. if (Number.isNaN(date.getTime())) fail(`Invalid --${name}: ${value}`)
  890. return date
  891. }
  892. function parseIntegerFlag(flags: Map<string, string[]>, name: string) {
  893. const value = flags.get(name)?.[0]
  894. if (!value) return undefined
  895. const parsed = Number(value)
  896. if (!Number.isInteger(parsed) || parsed <= 0) fail(`Invalid --${name}: ${value}`)
  897. return parsed
  898. }
  899. function parseListFlag(flags: Map<string, string[]>, name: string) {
  900. const value = flags.get(name)?.[0]
  901. if (!value) return undefined
  902. if (value === "all") return ["all"]
  903. return value
  904. .split(",")
  905. .map((item) => item.trim())
  906. .filter(Boolean)
  907. }
  908. function usage(): never {
  909. fail(`Usage:
  910. bun src/honeycomb-backfill.ts queries [--tiers Go,Free,Paid] [--limit 1000]
  911. bun src/honeycomb-backfill.ts import [--dry-run] [--upsert-chunk-size 100] [--database-url URL] --dir downloads
  912. bun src/honeycomb-backfill.ts import [--dry-run] [--upsert-chunk-size 100] [--database-url URL] --model-day file.csv [--model-day more.csv] ...`)
  913. }
  914. function fail(message: string): never {
  915. console.error(message)
  916. process.exit(1)
  917. }