stat-sync.ts 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101
  1. import { DateTime, Effect } from "effect"
  2. import { Resource } from "sst/resource"
  3. import { Athena, AthenaQueryError, AthenaQueryTimeoutError } from "./athena"
  4. import { DatabaseError } from "./database"
  5. import { GeoStatRepo, rowsFromAggregates as geoRowsFromAggregates } from "./domain/geo"
  6. import { buildStatsQuery, toGeoAggregate, toModelAggregate, toProviderAggregate } from "./domain/inference"
  7. import { ModelStatRepo, rowsFromAggregates as modelRowsFromAggregates } from "./domain/model"
  8. import { ProviderStatRepo, rowsFromAggregates as providerRowsFromAggregates } from "./domain/provider"
  9. import { startOfIsoWeek } from "./domain/stat"
  10. const DATALAKE_INGESTION_LAG_MS = 5 * 60_000
  11. const STATS_DATA_START_MS = new Date("2026-05-28T00:00:00.000Z").getTime()
  12. const WEEK_MS = 7 * 86_400_000
  13. export type SyncStatsResult = { ok: true; rows: number; startedAt: string; periodStart: string; periodEnd: string }
  14. export type SyncStatsError = AthenaQueryError | AthenaQueryTimeoutError | DatabaseError
  15. export const syncStats: () => Effect.Effect<
  16. SyncStatsResult,
  17. SyncStatsError,
  18. Athena | ModelStatRepo | ProviderStatRepo | GeoStatRepo
  19. > = Effect.fn("StatSync.sync")(function* () {
  20. const startedAt = yield* DateTime.nowAsDate
  21. const periodEnd = new Date(Math.floor((startedAt.getTime() - DATALAKE_INGESTION_LAG_MS) / 60_000) * 60_000)
  22. // May 27 was partial, so keep Athena stats anchored at the first complete day.
  23. const periodStart = new Date(Math.max(startOfIsoWeek(periodEnd).getTime() - WEEK_MS, STATS_DATA_START_MS))
  24. const athena = yield* Athena
  25. const modelStats = yield* ModelStatRepo
  26. const providerStats = yield* ProviderStatRepo
  27. const geoStats = yield* GeoStatRepo
  28. yield* logRuntimeCheck()
  29. const [modelAggregates, providerAggregates, geoAggregates, geoModelAggregates] = yield* Effect.all(
  30. [
  31. athena
  32. .query(buildStatsQuery(periodStart, periodEnd, "model"))
  33. .pipe(Effect.map((rows) => rows.flatMap(toModelAggregate))),
  34. athena
  35. .query(buildStatsQuery(periodStart, periodEnd, "provider"))
  36. .pipe(Effect.map((rows) => rows.flatMap(toProviderAggregate))),
  37. athena
  38. .query(buildStatsQuery(periodStart, periodEnd, "geo"))
  39. .pipe(Effect.map((rows) => rows.flatMap(toGeoAggregate))),
  40. athena
  41. .query(buildStatsQuery(periodStart, periodEnd, "geo_model"))
  42. .pipe(Effect.map((rows) => rows.flatMap(toGeoAggregate))),
  43. ],
  44. { concurrency: "unbounded" },
  45. )
  46. const modelRows = modelRowsFromAggregates(modelAggregates)
  47. const providerRows = providerRowsFromAggregates(providerAggregates)
  48. const geoRows = geoRowsFromAggregates([...geoAggregates, ...geoModelAggregates])
  49. yield* Effect.all([modelStats.upsert(modelRows), providerStats.upsert(providerRows), geoStats.upsert(geoRows)], {
  50. concurrency: "unbounded",
  51. discard: true,
  52. })
  53. yield* Effect.all(
  54. [
  55. modelStats.deleteRetiredDimensions(modelRows),
  56. providerStats.deleteRetiredDimensions(providerRows),
  57. geoStats.deleteRetiredDimensions(geoRows),
  58. ],
  59. { concurrency: "unbounded", discard: true },
  60. )
  61. yield* Effect.logInfo(
  62. `stats sync complete ${JSON.stringify({
  63. startedAt: startedAt.toISOString(),
  64. periodStart: periodStart.toISOString(),
  65. periodEnd: periodEnd.toISOString(),
  66. rows: modelRows.length,
  67. providerRows: providerRows.length,
  68. geoRows: geoRows.length,
  69. stage: Resource.App.stage,
  70. })}`,
  71. )
  72. return {
  73. ok: true,
  74. rows: modelRows.length,
  75. startedAt: startedAt.toISOString(),
  76. periodStart: periodStart.toISOString(),
  77. periodEnd: periodEnd.toISOString(),
  78. }
  79. })
  80. function logRuntimeCheck() {
  81. return Effect.logInfo(
  82. `athena stats runtime check ${JSON.stringify({
  83. catalog: Resource.InferenceEvent.catalog,
  84. database: Resource.InferenceEvent.database,
  85. dataset: Resource.StatsSyncConfig.dataset,
  86. table: Resource.InferenceEvent.table,
  87. workgroup: Resource.InferenceEvent.workgroup,
  88. region: Resource.InferenceEvent.region,
  89. stage: Resource.App.stage,
  90. })}`,
  91. )
  92. }