From 5acbd7e1f88a4adb58a8edcbdacbaa14b4fba736 Mon Sep 17 00:00:00 2001 From: saltbo Date: Wed, 22 Jul 2026 10:40:18 -0400 Subject: [PATCH] fix(stats): backfill registration source rollups --- scripts/backfill-admin-stats.ts | 223 ++++++++++++++++-- .../repos/admin-stats-counter-query.ts | 7 +- server/adapters/repos/admin-stats-hourly.ts | 49 ++-- .../admin-stats-rollup.integration.test.ts | 53 +++++ server/adapters/repos/admin-stats-rollup.ts | 11 +- server/adapters/repos/admin-stats.ts | 24 +- server/domain/admin-stats-metrics.ts | 3 +- server/http/admin-stats.integration.test.ts | 6 +- server/scripts/backfill-admin-stats.test.ts | 48 +++- 9 files changed, 369 insertions(+), 55 deletions(-) diff --git a/scripts/backfill-admin-stats.ts b/scripts/backfill-admin-stats.ts index 3b3e670f..a3fc6deb 100644 --- a/scripts/backfill-admin-stats.ts +++ b/scripts/backfill-admin-stats.ts @@ -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(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(options.target, validationSql) diff --git a/server/adapters/repos/admin-stats-counter-query.ts b/server/adapters/repos/admin-stats-counter-query.ts index 656a7b5f..4c5bc80e 100644 --- a/server/adapters/repos/admin-stats-counter-query.ts +++ b/server/adapters/repos/admin-stats-counter-query.ts @@ -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)) diff --git a/server/adapters/repos/admin-stats-hourly.ts b/server/adapters/repos/admin-stats-hourly.ts index a72b4094..3dd87c41 100644 --- a/server/adapters/repos/admin-stats-hourly.ts +++ b/server/adapters/repos/admin-stats-hourly.ts @@ -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>() - private readonly markerRowsPromises = new Map>() + private readonly markerRowsPromises = new Map>() constructor( private readonly db: Database, @@ -131,10 +132,13 @@ export class AdminStatsHourlyReader { })) } - async coverage(requiredScope: AdminStatsRollupScope = 'full'): Promise { + async coverage( + requiredScope: AdminStatsRollupScope = 'full', + metric?: AdminStatsMetric, + ): Promise { 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> { - const markerBuckets = await this.markerBuckets(requiredScope) + async completeDayKeys(requiredScope: AdminStatsRollupScope, metric?: AdminStatsMetric): Promise> { + const markerBuckets = await this.markerBuckets(requiredScope, metric) const queryTo = this.queryTo(requiredScope) const expectedByDay = new Map() const completedByDay = new Map() @@ -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> { - return this.markerRows(requiredScope).then((rows) => new Set(rows.map((row) => row.bucketStart.getTime()))) + private markerBuckets(requiredScope: AdminStatsRollupScope, metric?: AdminStatsMetric): Promise> { + return this.markerRows(requiredScope, metric).then((rows) => new Set(rows.map((row) => row.bucketStart.getTime()))) } - private markerRows(requiredScope: AdminStatsRollupScope): Promise { - const cached = this.markerRowsPromises.get(requiredScope) + private markerRows(requiredScope: AdminStatsRollupScope, metric?: AdminStatsMetric): Promise { + 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 { + private async loadMarkerRows( + requiredScope: AdminStatsRollupScope, + metric?: AdminStatsMetric, + ): Promise { 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 { diff --git a/server/adapters/repos/admin-stats-rollup.integration.test.ts b/server/adapters/repos/admin-stats-rollup.integration.test.ts index 0f53245d..cac48a5e 100644 --- a/server/adapters/repos/admin-stats-rollup.integration.test.ts +++ b/server/adapters/repos/admin-stats-rollup.integration.test.ts @@ -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') diff --git a/server/adapters/repos/admin-stats-rollup.ts b/server/adapters/repos/admin-stats-rollup.ts index 946b4538..44cdce19 100644 --- a/server/adapters/repos/admin-stats-rollup.ts +++ b/server/adapters/repos/admin-stats-rollup.ts @@ -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)), + ), ), ), ] diff --git a/server/adapters/repos/admin-stats.ts b/server/adapters/repos/admin-stats.ts index 119003ff..71e6a7a8 100644 --- a/server/adapters/repos/admin-stats.ts +++ b/server/adapters/repos/admin-stats.ts @@ -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) => ({ diff --git a/server/domain/admin-stats-metrics.ts b/server/domain/admin-stats-metrics.ts index 273e3a64..13ea7232 100644 --- a/server/domain/admin-stats-metrics.ts +++ b/server/domain/admin-stats-metrics.ts @@ -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), diff --git a/server/http/admin-stats.integration.test.ts b/server/http/admin-stats.integration.test.ts index 83496705..ce990b48 100644 --- a/server/http/admin-stats.integration.test.ts +++ b/server/http/admin-stats.integration.test.ts @@ -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, diff --git a/server/scripts/backfill-admin-stats.test.ts b/server/scripts/backfill-admin-stats.test.ts index 4d9eaeb8..0da157e3 100644 --- a/server/scripts/backfill-admin-stats.test.ts +++ b/server/scripts/backfill-admin-stats.test.ts @@ -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({