fix(stats): backfill registration source rollups

This commit is contained in:
saltbo
2026-07-22 10:40:18 -04:00
parent df0c5f3fbd
commit 5acbd7e1f8
9 changed files with 369 additions and 55 deletions
+207 -16
View File
@@ -87,6 +87,10 @@ interface ValidationSummary {
counterExpectedBuckets: number
counterCompletedBuckets: number
counterMissingBuckets: number
signupExpectedBuckets: number
signupCompletedBuckets: number
signupMissingBuckets: number
userSignupProviderMismatchGroups: number
openCounterMarkers: number
}
@@ -266,11 +270,46 @@ ${buildHourlyBackfillSql(now)}
function buildHourlyBackfillSql(now: Date): string {
const currentHour = Math.floor(now.getTime() / 3_600_000) * 3_600_000
const fromMs = `MAX(${MIN_VALID_TIMESTAMP_MS}, COALESCE(${statisticsFirstFullHourMsSql}, ${MIN_VALID_TIMESTAMP_MS}))`
const signupFromMs = userSignupHistoryStartSql(currentHour)
return [
...buildAdminStatsCounterRollupInsertSqlStatements({ fromMs, toMs: currentHour }),
...buildAdminStatsCounterRollupInsertSqlStatements({
fromMs,
toMs: currentHour,
metrics: ADMIN_STATS_FACT_COUNTER_METRICS.filter((metric) => metric !== M.userSignup),
}),
...buildAdminStatsCounterRollupInsertSqlStatements({
fromMs: signupFromMs,
toMs: currentHour,
metrics: [M.userSignup],
}),
rollupMarkerBackfillSql(now),
metricRollupMarkerBackfillSql(now, M.userSignup, signupFromMs),
].join(';\n\n')
}
function userSignupHistoryStartSql(currentHour: number): string {
return `MAX(
${MIN_VALID_TIMESTAMP_MS},
COALESCE((
SELECT CAST(MIN(occurred_at) / 3600000 AS INTEGER) * 3600000
FROM (
SELECT created_at * 1000 AS occurred_at
FROM audit_events
WHERE action = 'user_register'
AND target_id IS NOT NULL
AND user_id = target_id
AND id = 'event:user_register:' || target_id
AND json_valid(metadata) = 1
AND json_type(metadata, '$.provider') = 'text'
AND length(json_extract(metadata, '$.provider')) > 0
UNION ALL
SELECT created_at AS occurred_at FROM user
) registration_facts
WHERE occurred_at >= ${MIN_VALID_TIMESTAMP_MS}
AND occurred_at < ${currentHour}
), ${currentHour})
)`
}
function purgeLegacyRollupsSql(): string {
return `DELETE FROM stats_rollups_hourly
WHERE CASE WHEN json_valid(metadata) = 1 THEN
@@ -295,16 +334,28 @@ WHERE result.metric_key <> 'stats.rollup_run'
WHERE marker.bucket_start = result.bucket_start
AND marker.org_id = ''
AND marker.metric_key = 'stats.rollup_run'
AND marker.dimension_key = ''
AND marker.dimension_value = ''
AND json_valid(marker.metadata) = 1
AND json_extract(marker.metadata, '$.version') = 3
AND json_extract(marker.metadata, '$.quality') = 'exact'
AND (
(json_extract(result.metadata, '$.scope') = 'counters'
AND json_extract(marker.metadata, '$.scope') IN ('counters', 'full'))
OR (json_extract(result.metadata, '$.scope') = 'snapshots'
AND json_extract(marker.metadata, '$.scope') IN ('snapshots', 'full'))
OR (json_extract(result.metadata, '$.scope') = 'full'
AND json_extract(marker.metadata, '$.scope') = 'full')
(
marker.dimension_key = ''
AND marker.dimension_value = ''
AND (
(json_extract(result.metadata, '$.scope') = 'counters'
AND json_extract(marker.metadata, '$.scope') IN ('counters', 'full'))
OR (json_extract(result.metadata, '$.scope') = 'snapshots'
AND json_extract(marker.metadata, '$.scope') IN ('snapshots', 'full'))
OR (json_extract(result.metadata, '$.scope') = 'full'
AND json_extract(marker.metadata, '$.scope') = 'full')
)
)
OR (
json_extract(result.metadata, '$.scope') = 'counters'
AND marker.dimension_key = 'metric_key'
AND marker.dimension_value = result.metric_key
AND json_extract(marker.metadata, '$.scope') IN ('counters', 'full')
)
)
);`
}
@@ -329,10 +380,53 @@ WHERE metric_key IN (
)
OR (
metric_key = 'stats.rollup_run'
AND CASE WHEN json_valid(metadata) = 1 THEN json_extract(metadata, '$.scope') = 'counters' ELSE 0 END = 1
AND (
dimension_key = 'metric_key'
OR CASE WHEN json_valid(metadata) = 1 THEN json_extract(metadata, '$.scope') = 'counters' ELSE 0 END = 1
)
);`
}
function metricRollupMarkerBackfillSql(now: Date, metric: string, startAt: string): string {
const currentHour = Math.floor(now.getTime() / 3_600_000) * 3_600_000
const latestClosedHour = currentHour - 3_600_000
return `WITH
digits(n) AS (VALUES (0),(1),(2),(3),(4),(5),(6),(7),(8),(9)),
numbers(n) AS (
SELECT ones.n + tens.n * 10 + hundreds.n * 100 + thousands.n * 1000 + ten_thousands.n * 10000
FROM digits ones
CROSS JOIN digits tens
CROSS JOIN digits hundreds
CROSS JOIN digits thousands
CROSS JOIN digits ten_thousands
),
bounds AS (
SELECT ${startAt} AS start_at, ${latestClosedHour} AS end_at
),
buckets AS (
SELECT start_at + numbers.n * 3600000 AS bucket_start
FROM bounds
JOIN numbers ON start_at + numbers.n * 3600000 <= end_at
WHERE start_at <= end_at
)
INSERT INTO stats_rollups_hourly (
id, bucket_start, org_id, metric_key, dimension_key, dimension_value,
count, bytes, unique_count, metadata, updated_at
)
SELECT
CAST(bucket_start AS TEXT) || ':global:stats.rollup_run:metric_key:${metric}',
bucket_start, '', 'stats.rollup_run', 'metric_key', '${metric}', 1, 0, 0,
json_object('version', 3, 'scope', 'counters', 'quality', 'exact'),
bucket_start + 3600000
FROM buckets
WHERE true
ON CONFLICT(bucket_start, org_id, metric_key, dimension_key, dimension_value)
DO UPDATE SET count = excluded.count, bytes = excluded.bytes, unique_count = excluded.unique_count,
metadata = excluded.metadata, updated_at = excluded.updated_at
WHERE count <> excluded.count OR bytes <> excluded.bytes OR unique_count <> excluded.unique_count
OR metadata <> excluded.metadata OR updated_at <> excluded.updated_at;`
}
function rollupMarkerBackfillSql(now: Date): string {
const currentHour = Math.floor(now.getTime() / 3_600_000) * 3_600_000
const latestClosedHour = currentHour - 3_600_000
@@ -533,7 +627,7 @@ export function buildValidationSql(now = new Date()): string {
AND json_valid(metadata) = 1
AND json_type(metadata, '$.provider') = 'text'
AND length(json_extract(metadata, '$.provider')) > 0
AND created_at * 1000 >= MAX(${MIN_VALID_TIMESTAMP_MS}, COALESCE(${statisticsFirstFullHourMsSql}, ${MIN_VALID_TIMESTAMP_MS}))
AND created_at * 1000 >= ${MIN_VALID_TIMESTAMP_MS}
AND created_at * 1000 < ${currentHour}
),
'rawSharesCreated', (
@@ -690,11 +784,25 @@ counter_markers AS MATERIALIZED (
FROM completion_markers
WHERE scope IN ('counters', 'full')
),
metric_counter_markers AS MATERIALIZED (
SELECT bucket_start, dimension_value AS metric_key
FROM valid_rollups
WHERE metric_key = 'stats.rollup_run'
AND org_id = ''
AND dimension_key = 'metric_key'
AND json_extract(metadata, '$.scope') IN ('counters', 'full')
),
counter_rows AS MATERIALIZED (
SELECT result.*
FROM valid_rollups result
WHERE result.metric_key <> 'stats.rollup_run'
AND EXISTS (SELECT 1 FROM counter_markers marker WHERE marker.bucket_start = result.bucket_start)
AND (
EXISTS (SELECT 1 FROM counter_markers marker WHERE marker.bucket_start = result.bucket_start)
OR EXISTS (
SELECT 1 FROM metric_counter_markers marker
WHERE marker.bucket_start = result.bucket_start AND marker.metric_key = result.metric_key
)
)
),
required_dimensions(metric_key, dimension_key, check_bytes) AS (
VALUES
@@ -738,6 +846,31 @@ required_dimension_keys AS MATERIALIZED (
FROM required_dimensions required
JOIN dimension_rows dimension
ON dimension.metric_key = required.metric_key AND dimension.dimension_key = required.dimension_key
),
authoritative_registration_sources AS MATERIALIZED (
SELECT json_extract(metadata, '$.provider') AS provider, COUNT(*) AS count
FROM audit_events
WHERE action = 'user_register'
AND target_id IS NOT NULL
AND user_id = target_id
AND id = 'event:user_register:' || target_id
AND json_valid(metadata) = 1
AND json_type(metadata, '$.provider') = 'text'
AND length(json_extract(metadata, '$.provider')) > 0
AND created_at * 1000 >= ${MIN_VALID_TIMESTAMP_MS}
AND created_at * 1000 < ${currentHour}
GROUP BY 1
),
rollup_registration_sources AS MATERIALIZED (
SELECT dimension_value AS provider, SUM(count) AS count
FROM counter_rows
WHERE metric_key = 'user.signup' AND dimension_key = 'provider'
GROUP BY dimension_value
),
registration_providers AS MATERIALIZED (
SELECT provider FROM authoritative_registration_sources
UNION
SELECT provider FROM rollup_registration_sources
)
SELECT json_object(
'rollupUploadEvents', (
@@ -778,6 +911,20 @@ SELECT json_object(
OR (json_extract(r.metadata, '$.scope') = 'full' AND marker.scope = 'full')
)
)
AND NOT EXISTS (
SELECT 1
FROM metric_counter_markers marker
WHERE marker.bucket_start = r.bucket_start
AND marker.metric_key = r.metric_key
AND json_extract(r.metadata, '$.scope') = 'counters'
)
),
'userSignupProviderMismatchGroups', (
SELECT COUNT(*)
FROM registration_providers provider
LEFT JOIN authoritative_registration_sources authoritative USING (provider)
LEFT JOIN rollup_registration_sources rollup USING (provider)
WHERE COALESCE(authoritative.count, 0) <> COALESCE(rollup.count, 0)
),
'requiredDimensionMismatchGroups', (
SELECT COUNT(*)
@@ -833,10 +980,21 @@ counter_markers AS MATERIALIZED (
WHERE metric_key = 'stats.rollup_run' AND org_id = '' AND dimension_key = '' AND dimension_value = ''
AND json_extract(metadata, '$.scope') IN ('counters', 'full')
),
metric_counter_markers AS MATERIALIZED (
SELECT bucket_start, dimension_value AS metric_key FROM valid_rollups
WHERE metric_key = 'stats.rollup_run' AND org_id = '' AND dimension_key = 'metric_key'
AND json_extract(metadata, '$.scope') IN ('counters', 'full')
),
counter_rows AS MATERIALIZED (
SELECT result.* FROM valid_rollups result
WHERE result.metric_key <> 'stats.rollup_run'
AND EXISTS (SELECT 1 FROM counter_markers marker WHERE marker.bucket_start = result.bucket_start)
AND (
EXISTS (SELECT 1 FROM counter_markers marker WHERE marker.bucket_start = result.bucket_start)
OR EXISTS (
SELECT 1 FROM metric_counter_markers marker
WHERE marker.bucket_start = result.bucket_start AND marker.metric_key = result.metric_key
)
)
)
SELECT json_object(
'rollupUploadAttempts', (
@@ -888,9 +1046,23 @@ SELECT json_object(
ELSE 0 END = 1
AND json_extract(metadata, '$.scope') IN ('counters', 'full')
),
signup_markers AS MATERIALIZED (
SELECT bucket_start
FROM stats_rollups_hourly
WHERE metric_key = 'stats.rollup_run' AND org_id = ''
AND dimension_key = 'metric_key' AND dimension_value = '${M.userSignup}'
AND CASE WHEN json_valid(metadata) = 1 THEN
json_extract(metadata, '$.version') = 3
AND json_extract(metadata, '$.scope') IN ('counters', 'full')
AND json_extract(metadata, '$.quality') = 'exact'
ELSE 0 END = 1
),
coverage AS MATERIALIZED (
SELECT ${statsHistoryStartSql()} AS start_at, ${latestClosedHour} AS end_at
),
signup_coverage AS MATERIALIZED (
SELECT ${userSignupHistoryStartSql(currentHour)} AS start_at, ${latestClosedHour} AS end_at
),
coverage_counts AS MATERIALIZED (
SELECT
CASE WHEN start_at IS NULL OR start_at > end_at THEN 0
@@ -898,13 +1070,27 @@ coverage_counts AS MATERIALIZED (
CASE WHEN start_at IS NULL OR start_at > end_at THEN 0
ELSE (SELECT COUNT(*) FROM counter_markers WHERE bucket_start BETWEEN start_at AND end_at) END AS completed
FROM coverage
),
signup_coverage_counts AS MATERIALIZED (
SELECT
CASE WHEN start_at IS NULL OR start_at > end_at THEN 0
ELSE CAST((end_at - start_at) / 3600000 AS INTEGER) + 1 END AS expected,
CASE WHEN start_at IS NULL OR start_at > end_at THEN 0
ELSE (SELECT COUNT(*) FROM signup_markers WHERE bucket_start BETWEEN start_at AND end_at) END AS completed
FROM signup_coverage
)
SELECT json_object(
'hourlyRollups', (SELECT COUNT(*) FROM counter_markers),
'counterExpectedBuckets', (SELECT expected FROM coverage_counts),
'counterCompletedBuckets', (SELECT completed FROM coverage_counts),
'counterMissingBuckets', (SELECT expected - completed FROM coverage_counts),
'openCounterMarkers', (SELECT COUNT(*) FROM counter_markers WHERE bucket_start >= ${currentHour})
'signupExpectedBuckets', (SELECT expected FROM signup_coverage_counts),
'signupCompletedBuckets', (SELECT completed FROM signup_coverage_counts),
'signupMissingBuckets', (SELECT expected - completed FROM signup_coverage_counts),
'openCounterMarkers', (
(SELECT COUNT(*) FROM counter_markers WHERE bucket_start >= ${currentHour})
+ (SELECT COUNT(*) FROM signup_markers WHERE bucket_start >= ${currentHour})
)
) AS summary;
`
@@ -1106,11 +1292,13 @@ export function assertBackfillValidation(summary: ValidationSummary): void {
summary.missingUserRegistrationEvents > 0 ||
summary.invalidDownloadTaskEvents > 0 ||
summary.orphanRollupBuckets > 0 ||
summary.userSignupProviderMismatchGroups > 0 ||
summary.requiredDimensionMismatchGroups > 0 ||
summary.lowerBoundRollups > 0 ||
summary.legacyRollupRows > 0 ||
summary.incompatibleUserSnapshotRows > 0 ||
summary.counterMissingBuckets > 0 ||
summary.signupMissingBuckets > 0 ||
summary.openCounterMarkers > 0
) {
throw new Error(
@@ -1120,11 +1308,13 @@ export function assertBackfillValidation(summary: ValidationSummary): void {
missingUserRegistrationEvents: summary.missingUserRegistrationEvents,
invalidDownloadTaskEvents: summary.invalidDownloadTaskEvents,
orphanRollupBuckets: summary.orphanRollupBuckets,
userSignupProviderMismatchGroups: summary.userSignupProviderMismatchGroups,
requiredDimensionMismatchGroups: summary.requiredDimensionMismatchGroups,
lowerBoundRollups: summary.lowerBoundRollups,
legacyRollupRows: summary.legacyRollupRows,
incompatibleUserSnapshotRows: summary.incompatibleUserSnapshotRows,
counterMissingBuckets: summary.counterMissingBuckets,
signupMissingBuckets: summary.signupMissingBuckets,
openCounterMarkers: summary.openCounterMarkers,
})}`,
)
@@ -1139,8 +1329,9 @@ function main(): void {
const plan = querySummary<BackfillPlan>(options.target, BACKFILL_PLAN_SQL)
console.log(JSON.stringify({ mode: options.apply ? 'apply' : 'dry-run', before, plan }, null, 2))
if (!options.apply) return
if (before.counterExpectedBuckets > MAX_BACKFILL_HOURS) {
throw new Error(`admin_stats_backfill_range_too_large:${before.counterExpectedBuckets}`)
const maxExpectedBuckets = Math.max(before.counterExpectedBuckets, before.signupExpectedBuckets)
if (maxExpectedBuckets > MAX_BACKFILL_HOURS) {
throw new Error(`admin_stats_backfill_range_too_large:${maxExpectedBuckets}`)
}
apply(options.target, buildBackfillSql(now))
const after = querySummary<ValidationSummary>(options.target, validationSql)
@@ -7,6 +7,7 @@ const D1_MAX_COMPOUND_SELECT_TERMS = 5
export interface AdminStatsCounterQueryRange {
fromMs: number | string
toMs: number | string
metrics?: readonly AdminStatsMetric[]
}
type HourlySource = {
@@ -240,7 +241,11 @@ WHERE count <> excluded.count
}
function buildCounterSqlStatements(range: AdminStatsCounterQueryRange, statement: string): string[] {
const sources = [...HOURLY_SOURCES.map(hourlySourceRow), storageBalanceRow()]
const selectedMetrics = range.metrics ? new Set(range.metrics) : null
const sources = [
...HOURLY_SOURCES.filter((source) => !selectedMetrics || selectedMetrics.has(source.metric)).map(hourlySourceRow),
...(!selectedMetrics || selectedMetrics.has(M.storageLedgerBalance) ? [storageBalanceRow()] : []),
]
const statements: string[] = []
for (let index = 0; index < sources.length; index += D1_MAX_COMPOUND_SELECT_TERMS) {
statements.push(buildCounterSql(range, sources.slice(index, index + D1_MAX_COMPOUND_SELECT_TERMS), statement))
+33 -16
View File
@@ -1,7 +1,8 @@
import type { AdminStatsCoverage } from '@shared/types'
import { and, eq, gte, inArray, lt, sql } from 'drizzle-orm'
import { and, eq, gte, inArray, lt, or, sql } from 'drizzle-orm'
import { statsRollupsHourly } from '../../db/schema'
import {
ADMIN_STATS_DIMENSIONS,
ADMIN_STATS_METRICS,
type AdminStatsDimension,
type AdminStatsMetric,
@@ -35,7 +36,7 @@ export class AdminStatsHourlyReader {
private readonly counterQueryTo: Date
private readonly snapshotQueryTo: Date
private readonly metricRows = new Map<string, Promise<HourlyMetricRow[]>>()
private readonly markerRowsPromises = new Map<AdminStatsRollupScope, Promise<CompatibleMarkerRow[]>>()
private readonly markerRowsPromises = new Map<string, Promise<CompatibleMarkerRow[]>>()
constructor(
private readonly db: Database,
@@ -131,10 +132,13 @@ export class AdminStatsHourlyReader {
}))
}
async coverage(requiredScope: AdminStatsRollupScope = 'full'): Promise<AdminStatsCoverage> {
async coverage(
requiredScope: AdminStatsRollupScope = 'full',
metric?: AdminStatsMetric,
): Promise<AdminStatsCoverage> {
const queryTo = this.queryTo(requiredScope)
const expectedBuckets = Math.max(0, Math.floor((queryTo.getTime() - this.queryFrom.getTime()) / HOUR_MS))
const markerRows = await this.markerRows(requiredScope)
const markerRows = await this.markerRows(requiredScope, metric)
const completedBuckets = markerRows.length
const lowerBoundBuckets = markerRows.filter(
(row) => qualityForScope(row.metadata, requiredScope) === 'lower_bound',
@@ -153,8 +157,8 @@ export class AdminStatsHourlyReader {
}
}
async completeDayKeys(requiredScope: AdminStatsRollupScope): Promise<Set<string>> {
const markerBuckets = await this.markerBuckets(requiredScope)
async completeDayKeys(requiredScope: AdminStatsRollupScope, metric?: AdminStatsMetric): Promise<Set<string>> {
const markerBuckets = await this.markerBuckets(requiredScope, metric)
const queryTo = this.queryTo(requiredScope)
const expectedByDay = new Map<string, number>()
const completedByDay = new Map<string, number>()
@@ -185,7 +189,7 @@ export class AdminStatsHourlyReader {
const requiredScope = metricDefinition(metric).kind === 'gauge' ? 'snapshots' : 'counters'
const queryTo = this.queryTo(requiredScope)
if (this.queryFrom >= queryTo) return []
const markerBuckets = await this.markerBuckets(requiredScope)
const markerBuckets = await this.markerBuckets(requiredScope, metric)
if (markerBuckets.size === 0) return []
const rows = await this.db
.select({
@@ -229,19 +233,23 @@ export class AdminStatsHourlyReader {
}))
}
private markerBuckets(requiredScope: AdminStatsRollupScope): Promise<Set<number>> {
return this.markerRows(requiredScope).then((rows) => new Set(rows.map((row) => row.bucketStart.getTime())))
private markerBuckets(requiredScope: AdminStatsRollupScope, metric?: AdminStatsMetric): Promise<Set<number>> {
return this.markerRows(requiredScope, metric).then((rows) => new Set(rows.map((row) => row.bucketStart.getTime())))
}
private markerRows(requiredScope: AdminStatsRollupScope): Promise<CompatibleMarkerRow[]> {
const cached = this.markerRowsPromises.get(requiredScope)
private markerRows(requiredScope: AdminStatsRollupScope, metric?: AdminStatsMetric): Promise<CompatibleMarkerRow[]> {
const cacheKey = `${requiredScope}\u0000${metric ?? ''}`
const cached = this.markerRowsPromises.get(cacheKey)
if (cached) return cached
const rows = this.loadMarkerRows(requiredScope)
this.markerRowsPromises.set(requiredScope, rows)
const rows = this.loadMarkerRows(requiredScope, metric)
this.markerRowsPromises.set(cacheKey, rows)
return rows
}
private async loadMarkerRows(requiredScope: AdminStatsRollupScope): Promise<CompatibleMarkerRow[]> {
private async loadMarkerRows(
requiredScope: AdminStatsRollupScope,
metric?: AdminStatsMetric,
): Promise<CompatibleMarkerRow[]> {
const queryTo = this.queryTo(requiredScope)
if (this.queryFrom >= queryTo) return []
const rows = await this.db
@@ -251,12 +259,20 @@ export class AdminStatsHourlyReader {
and(
eq(statsRollupsHourly.metricKey, ADMIN_STATS_METRICS.statsRollupRun),
eq(statsRollupsHourly.orgId, ''),
eq(statsRollupsHourly.dimensionKey, ''),
metric
? or(
and(eq(statsRollupsHourly.dimensionKey, ''), eq(statsRollupsHourly.dimensionValue, '')),
and(
eq(statsRollupsHourly.dimensionKey, ADMIN_STATS_DIMENSIONS.metric),
eq(statsRollupsHourly.dimensionValue, metric),
),
)
: and(eq(statsRollupsHourly.dimensionKey, ''), eq(statsRollupsHourly.dimensionValue, '')),
gte(statsRollupsHourly.bucketStart, this.queryFrom),
lt(statsRollupsHourly.bucketStart, queryTo),
),
)
return rows.flatMap((row) => {
const compatibleRows = rows.flatMap((row) => {
const metadata = parseAdminStatsRollupMetadata(row.metadata)
return metadata &&
supportsScope(metadata.scope, requiredScope) &&
@@ -264,6 +280,7 @@ export class AdminStatsHourlyReader {
? [{ bucketStart: row.bucketStart, metadata }]
: []
})
return [...new Map(compatibleRows.map((row) => [row.bucketStart.getTime(), row])).values()]
}
private queryTo(requiredScope: AdminStatsRollupScope): Date {
@@ -314,6 +314,7 @@ describe('admin hourly stats rollup', () => {
expect(row(M.trafficReportSnapshot, '', '', '')).toMatchObject({ count: 5, bytes: 290 })
expect(row(M.webhookSnapshot, 'status', 'processed', '')).toMatchObject({ count: 1 })
expect(row(M.statsRollupRun, '', '', '')).toMatchObject({ count: 1 })
expect(row(M.statsRollupRun, 'metric_key', M.userSignup, '')).toMatchObject({ count: 1 })
expect(JSON.parse(row(M.transferUpload)?.metadata ?? '{}')).toMatchObject({ version: 3, quality: 'exact' })
expect(JSON.parse(row(M.statsRollupRun, '', '', '')?.metadata ?? '{}')).toMatchObject({
version: 3,
@@ -323,6 +324,11 @@ describe('admin hourly stats rollup', () => {
snapshotQuality: 'exact',
snapshotObservedAt: generatedAt.toISOString(),
})
expect(JSON.parse(row(M.statsRollupRun, 'metric_key', M.userSignup, '')?.metadata ?? '{}')).toMatchObject({
version: 3,
scope: 'counters',
quality: 'exact',
})
const reader = new AdminStatsHourlyReader(
db,
@@ -356,6 +362,17 @@ describe('admin hourly stats rollup', () => {
`)
expect(second.rows).toBe(first.rows)
expect(storedRows).toBeGreaterThan(first.rows)
await captureAdminStatsSnapshot(db, bucketStart, generatedAt)
const [{ count: signupMarkers }] = await db.all<{ count: number }>(sql`
SELECT COUNT(*) AS count
FROM stats_rollups_hourly
WHERE bucket_start = ${bucketStart.getTime()}
AND metric_key = ${M.statsRollupRun}
AND dimension_key = 'metric_key'
AND dimension_value = ${M.userSignup}
`)
expect(signupMarkers).toBe(1)
})
it('rejects buckets that are not aligned to a UTC hour', async () => {
@@ -527,6 +544,42 @@ describe('admin hourly stats rollup', () => {
expect(await reader.coverage()).toMatchObject({ status: 'empty', completedBuckets: 0 })
})
it('uses metric completion markers without claiming unrelated counters are complete', async () => {
const { db } = await createTestApp()
const bucketStart = Date.parse('2026-06-10T09:00:00.000Z')
const metadata = '{"version":3,"scope":"counters","quality":"exact"}'
await db.run(sql`
INSERT INTO stats_rollups_hourly
(id, bucket_start, org_id, metric_key, dimension_key, dimension_value,
count, bytes, unique_count, metadata, updated_at)
VALUES
('signup-marker', ${bucketStart}, '', 'stats.rollup_run', 'metric_key', 'user.signup', 1, 0, 0,
${metadata}, ${bucketStart}),
('signup-result', ${bucketStart}, '', 'user.signup', '', '', 2, 0, 0, ${metadata}, ${bucketStart}),
('signup-provider', ${bucketStart}, '', 'user.signup', 'provider', 'github', 2, 0, 0,
${metadata}, ${bucketStart}),
('unmarked-upload', ${bucketStart}, '', 'transfer.upload', '', '', 1, 42, 0, ${metadata}, ${bucketStart})
`)
const reader = new AdminStatsHourlyReader(
db,
{
from: new Date(bucketStart),
to: new Date(bucketStart + 3_600_000 - 1),
timeZone: 'UTC',
},
new Date(bucketStart + 7_200_000),
)
expect(await reader.rows(M.userSignup, ['', 'provider'])).toHaveLength(2)
expect(await reader.rows(M.transferUpload)).toEqual([])
expect(await reader.coverage('counters')).toMatchObject({ status: 'empty', completedBuckets: 0 })
expect(await reader.coverage('counters', M.userSignup)).toMatchObject({
status: 'complete',
completedBuckets: 1,
})
expect(await reader.completeDayKeys('counters', M.userSignup)).toEqual(new Set(['2026-06-10']))
})
it('marks a day complete only when every requested hour has an exact marker', async () => {
const { db } = await createTestApp()
const firstHour = Date.parse('2026-07-10T10:00:00.000Z')
+8 -3
View File
@@ -1,4 +1,4 @@
import { and, eq, inArray, isNull, sql } from 'drizzle-orm'
import { and, eq, inArray, isNull, ne, or, sql } from 'drizzle-orm'
import { organization } from '../../db/auth-schema'
import {
backgroundJobs,
@@ -17,6 +17,7 @@ import {
type AdminStatsDimension,
type AdminStatsMetric,
assertMetricDimension,
ADMIN_STATS_DIMENSIONS as D,
ADMIN_STATS_METRICS as M,
parseAdminStatsRollupMetadata,
ROLLUP_VERSION,
@@ -120,6 +121,7 @@ export async function rebuildAdminStatsHour(
}
const lowerBoundRows = 0
rollups.add(M.statsRollupRun, '', 1, 0, { outcome: 'success' })
rollups.incrementValue(M.statsRollupRun, '', 'metric_key', M.userSignup, 1, 0, 0)
const completionScope = capturedSnapshot ? 'full' : 'counters'
const counterQuality: 'exact' = 'exact'
@@ -135,7 +137,7 @@ export async function rebuildAdminStatsHour(
bytes: row.bytes,
uniqueCount: row.uniqueCount,
metadata: JSON.stringify(
row.metric === M.statsRollupRun
row.metric === M.statsRollupRun && row.dimensionKey === ''
? {
version: ROLLUP_VERSION,
scope: completionScope,
@@ -215,7 +217,10 @@ export async function captureAdminStatsSnapshot(
.where(
and(
eq(statsRollupsHourly.bucketStart, bucketStart),
inArray(statsRollupsHourly.metricKey, [...GAUGE_METRICS, M.statsRollupRun]),
or(
inArray(statsRollupsHourly.metricKey, GAUGE_METRICS),
and(eq(statsRollupsHourly.metricKey, M.statsRollupRun), ne(statsRollupsHourly.dimensionKey, D.metric)),
),
),
),
]
+16 -8
View File
@@ -76,6 +76,7 @@ async function getOverviewStatistics(
topUsage,
storageDataQuality,
counterDays,
signupDays,
] = await Promise.all([
getUserInventory(reader),
getActiveUserSnapshot(reader),
@@ -89,6 +90,7 @@ async function getOverviewStatistics(
getTopPersonalUsage(db, now),
getStorageDataQuality(reader),
reader.completeDayKeys('counters'),
reader.completeDayKeys('counters', ADMIN_STATS_METRICS.userSignup),
])
const dates = createDateBuckets(effective)
const activeByDate = new Map(activeByDay.map((row) => [row.date, row]))
@@ -115,7 +117,7 @@ async function getOverviewStatistics(
date,
totalUsers: totalUsersByDay.get(date) ?? null,
activeUsers: activeByDate.get(date)?.mau ?? null,
newUsers: counterDays.has(date) ? (newUsersByDay.get(date) ?? 0) : null,
newUsers: signupDays.has(date) ? (newUsersByDay.get(date) ?? 0) : null,
})),
topUsage: exactUsage ? topUsage : [],
},
@@ -272,6 +274,9 @@ async function getDashboardOverviewStats(
trafficLedgerAvailable,
previousTrafficLedgerAvailable,
counterDays,
signupDays,
signupCoverage,
previousSignupCoverage,
trafficFirstCompleteDay,
] = await Promise.all([
getUserInventory(reader),
@@ -292,6 +297,9 @@ async function getDashboardOverviewStats(
trafficLedgerAvailableInRange(db, effective),
trafficLedgerAvailableInRange(db, previous),
reader.completeDayKeys('counters'),
reader.completeDayKeys('counters', ADMIN_STATS_METRICS.userSignup),
reader.coverage('counters', ADMIN_STATS_METRICS.userSignup),
previousReader.coverage('counters', ADMIN_STATS_METRICS.userSignup),
trafficLedgerFirstCompleteDay(db),
])
const [trendNewUsers, activeByDay, storageUsedByDay, uploadByDay, downloadByDay, missingBytesByDay] =
@@ -307,7 +315,7 @@ async function getDashboardOverviewStats(
const missingBytes = missingBytesByDay.get(date)
return {
date,
newUsers: counterDays.has(date) ? (trendNewUsers.get(date) ?? 0) : null,
newUsers: signupDays.has(date) ? (trendNewUsers.get(date) ?? 0) : null,
activeUsers: activeByDay.get(date) ?? null,
storageUsedBytes: storageUsedByDay.get(date) ?? null,
uploadBytes: counterDays.has(date) && !missingBytes?.upload ? (uploadByDay.get(date) ?? 0) : null,
@@ -343,9 +351,9 @@ async function getDashboardOverviewStats(
totals: {
users: users?.total ?? null,
newUsers: delta(
currentCountersAvailable ? newUsers : null,
previousCountersAvailable ? previousNewUsers : null,
countersComparable,
hasExactCoverage(signupCoverage) ? newUsers : null,
hasExactCoverage(previousSignupCoverage) ? previousNewUsers : null,
comparable(signupCoverage, previousSignupCoverage),
),
activeUsers: delta(
activeUsers?.mau ?? null,
@@ -489,11 +497,11 @@ async function getDashboardGrowthStats(
getRollingActiveUserTrend(reader, effective),
getRegistrationSources(reader),
getUserTotalsByDay(reader),
reader.coverage('counters'),
previousReader.coverage('counters'),
reader.coverage('counters', ADMIN_STATS_METRICS.userSignup),
previousReader.coverage('counters', ADMIN_STATS_METRICS.userSignup),
reader.coverage('snapshots'),
previousReader.coverage('snapshots'),
reader.completeDayKeys('counters'),
reader.completeDayKeys('counters', ADMIN_STATS_METRICS.userSignup),
])
const newUsersByDay = await getSignupsByDay(reader)
const userScaleTrend = createDateBuckets(effective).map((date) => ({
+2 -1
View File
@@ -23,6 +23,7 @@ export const ADMIN_STATS_DIMENSIONS = {
jobType: 'job_type',
kind: 'kind',
lifecycle: 'lifecycle',
metric: 'metric_key',
orgType: 'org_type',
outcome: 'outcome',
provider: 'provider',
@@ -90,7 +91,7 @@ export const ADMIN_STATS_METRIC_REGISTRY = {
[M.shareSaved]: counter(['actor_type', 'share_id'], true),
[M.statsDataQualitySnapshot]: gauge(['kind'], 'events', true),
[M.statsMissingBytes]: counter(['direction', 'source']),
[M.statsRollupRun]: counter(['outcome']),
[M.statsRollupRun]: counter(['metric_key', 'outcome']),
[M.storageInventory]: gauge(['age_bucket', 'file_type_group', 'size_bucket', 'storage_id'], 'entities', true),
[M.storageLedgerBalance]: counter(['storage_id'], true),
[M.storageLedgerChange]: counter(['direction', 'reason', 'storage_id'], true),
+4 -2
View File
@@ -532,8 +532,10 @@ describe('site stats routes', () => {
(id, bucket_start, org_id, metric_key, dimension_key, dimension_value,
count, bytes, unique_count, metadata, updated_at)
VALUES
('growth-rollup-marker', ${at}, '', 'stats.rollup_run', '', '', 1, 0, 0,
'{"version":3,"scope":"full","quality":"exact"}', ${at + 3_600_000}),
('growth-snapshot-marker', ${at}, '', 'stats.rollup_run', '', '', 1, 0, 0,
'{"version":3,"scope":"snapshots","quality":"exact"}', ${at + 3_600_000}),
('growth-signup-marker', ${at}, '', 'stats.rollup_run', 'metric_key', 'user.signup', 1, 0, 0,
'{"version":3,"scope":"counters","quality":"exact"}', ${at + 3_600_000}),
('growth-rollup-total', ${at}, '', 'user.signup', '', '', 2, 0, 0,
'{"version":3,"scope":"counters","quality":"exact"}', ${at + 3_600_000}),
('growth-inventory-total', ${at}, '', 'user.inventory', '', '', 2, 0, 0,
+40 -8
View File
@@ -40,13 +40,24 @@ describe('admin stats backfill', () => {
expect(sql.match(/FROM audit_events ae/g)).toHaveLength(1)
expect(sql.match(/FROM cloud_traffic_reports traffic_report/g)).toHaveLength(1)
expect(sql.match(/FROM storage_usage_ledger storage_change/g)).toHaveLength(1)
expect(sql).toContain('audit_events registered_user')
expect(statements.every((statement) => (statement.match(/UNION ALL/g) ?? []).length <= 5)).toBe(true)
const signupStatements = buildAdminStatsCounterRowsSqlStatements({
fromMs: 0,
toMs: 3_600_000,
metrics: ['user.signup'],
})
expect(signupStatements).toHaveLength(1)
expect(signupStatements[0]).toContain('audit_events registered_user')
expect(signupStatements[0]).not.toContain('cloud_traffic_reports')
})
it('recovers exact available facts and is idempotent', () => {
const db = new Database(':memory:')
const now = new Date('2026-07-10T12:00:00.000Z')
const historyStartMs = Date.parse('2026-04-01T00:10:00.000Z')
const preExactSignupMs = Date.parse('2026-03-31T22:10:00.000Z')
const eventMs = Date.parse('2026-07-10T09:10:00.000Z')
const eventHourMs = Date.parse('2026-07-10T09:00:00.000Z')
const sessionCreatedMs = Date.parse('2026-07-10T08:00:00.000Z')
@@ -60,6 +71,8 @@ describe('admin stats backfill', () => {
const latestClosedHour = Date.parse('2026-07-10T11:00:00.000Z')
const firstExactHour = Math.ceil(historyStartMs / 3_600_000) * 3_600_000
const expectedBuckets = (latestClosedHour - firstExactHour) / 3_600_000 + 1
const signupFirstHour = Math.floor(preExactSignupMs / 3_600_000) * 3_600_000
const expectedSignupBuckets = (latestClosedHour - signupFirstHour) / 3_600_000 + 1
db.exec(`
CREATE TABLE user (id TEXT PRIMARY KEY, created_at INTEGER NOT NULL DEFAULT 0, last_active_at INTEGER);
CREATE TABLE account (id TEXT PRIMARY KEY, user_id TEXT NOT NULL, provider_id TEXT NOT NULL, created_at INTEGER NOT NULL);
@@ -127,9 +140,11 @@ describe('admin stats backfill', () => {
INSERT INTO user VALUES
('u0', 0, ${newerExistingActivityMs}),
('u3', ${preExactSignupMs}, NULL),
('u1', ${firstExactHour + 600_000}, NULL),
('u2', ${firstExactHour + 601_000}, NULL);
INSERT INTO account VALUES
('a3', 'u3', 'google', ${preExactSignupMs}),
('a1', 'u1', 'github', ${firstExactHour + 600_000}),
('a2', 'u2', 'github', ${firstExactHour + 601_000});
INSERT INTO session VALUES
@@ -224,12 +239,18 @@ describe('admin stats backfill', () => {
const firstSummary = readValidationSummary()
const firstRows = db.prepare('SELECT * FROM stats_rollups_hourly ORDER BY id').all()
const metricList = ADMIN_STATS_FACT_COUNTER_METRICS.map((metric) => `'${metric}'`).join(', ')
const calculatedFactRows = buildAdminStatsCounterRowsSqlStatements({
fromMs: firstExactHour,
toMs: currentHourMs,
})
.flatMap((statement) => db.prepare(statement).all())
.sort(compareCounterRows)
const calculatedFactRows = [
...buildAdminStatsCounterRowsSqlStatements({
fromMs: firstExactHour,
toMs: currentHourMs,
metrics: ADMIN_STATS_FACT_COUNTER_METRICS.filter((metric) => metric !== 'user.signup'),
}).flatMap((statement) => db.prepare(statement).all()),
...buildAdminStatsCounterRowsSqlStatements({
fromMs: signupFirstHour,
toMs: currentHourMs,
metrics: ['user.signup'],
}).flatMap((statement) => db.prepare(statement).all()),
].sort(compareCounterRows)
const backfilledFactRows = db
.prepare(
`SELECT
@@ -271,14 +292,18 @@ describe('admin stats backfill', () => {
counterExpectedBuckets: expectedBuckets,
counterCompletedBuckets: expectedBuckets,
counterMissingBuckets: 0,
signupExpectedBuckets: expectedSignupBuckets,
signupCompletedBuckets: expectedSignupBuckets,
signupMissingBuckets: 0,
openCounterMarkers: 0,
requiredDimensionMismatchGroups: 0,
userSignupProviderMismatchGroups: 0,
orphanRollupBuckets: 0,
lowerBoundRollups: 0,
rawUploadAttempts: 1,
rollupUploadAttempts: 1,
rawUserSignups: 2,
rollupUserSignups: 2,
rawUserSignups: 3,
rollupUserSignups: 3,
rawSharesCreated: 1,
rollupSharesCreated: 1,
rawFailedDownloads: 1,
@@ -362,6 +387,13 @@ describe('admin stats backfill', () => {
)
.get(),
).toEqual({ value: 2 })
expect(
db
.prepare(
"SELECT count AS value FROM stats_rollups_hourly WHERE metric_key = 'user.signup' AND dimension_key = 'provider' AND dimension_value = 'google'",
)
.get(),
).toEqual({ value: 1 })
expect(
db.prepare("SELECT COUNT(*) AS value FROM stats_rollups_hourly WHERE metric_key = 'traffic.report_sync'").get(),
).toEqual({