From 0d98176c595a0b0b7e152107dc96ab8769759fac Mon Sep 17 00:00:00 2001 From: IanShaw027 Date: Sat, 8 Aug 2026 01:59:44 +0800 Subject: [PATCH] fix(channel-monitor-v2): correct aggregation privacy and backfill --- .../handler/channel_monitor_user_handler.go | 3 +- .../channel_monitor_v2_aggregation.go | 53 +++++++++++++++---- .../channel_monitor_v2_repo_test.go | 11 ++++ .../internal/service/channel_monitor_v2.go | 7 ++- .../service/channel_monitor_v2_aggregator.go | 49 +++++++++++++---- .../service/channel_monitor_v2_test.go | 2 +- backend/internal/service/setting_public.go | 4 +- ...el_monitor_v2_default_ignore_and_cache.sql | 14 ++--- .../204_channel_monitor_hide_throughput.sql | 4 +- ...itor_v2_reset_factory_cache_thresholds.sql | 13 +++++ 10 files changed, 122 insertions(+), 38 deletions(-) create mode 100644 backend/migrations/205_channel_monitor_v2_reset_factory_cache_thresholds.sql diff --git a/backend/internal/handler/channel_monitor_user_handler.go b/backend/internal/handler/channel_monitor_user_handler.go index cc36b33461..7e42c5a74d 100644 --- a/backend/internal/handler/channel_monitor_user_handler.go +++ b/backend/internal/handler/channel_monitor_user_handler.go @@ -35,7 +35,8 @@ func (h *ChannelMonitorUserHandler) featureEnabled(c *gin.Context) bool { if h.settingService == nil { return true } - return h.settingService.GetChannelMonitorRuntime(c.Request.Context()).Enabled + runtime := h.settingService.GetChannelMonitorRuntime(c.Request.Context()) + return runtime.Enabled && runtime.Mode == service.ChannelMonitorModeV1 } // --- Response --- diff --git a/backend/internal/repository/channel_monitor_v2_aggregation.go b/backend/internal/repository/channel_monitor_v2_aggregation.go index 9d908e955e..cd98d0f1af 100644 --- a/backend/internal/repository/channel_monitor_v2_aggregation.go +++ b/backend/internal/repository/channel_monitor_v2_aggregation.go @@ -8,7 +8,7 @@ import ( ) // Platform is derived from group/account (usage_logs has no provider column on upstream schema). -const channelMonitorV2PlatformSQL = `lower(COALESCE(NULLIF(TRIM(g.platform), ''), NULLIF(TRIM(a.platform), ''), 'unknown'))` +const channelMonitorV2PlatformSQL = `lower(` + usageLogEffectivePlatformExpr + `)` const channelMonitorV2ModelSQL = `COALESCE(NULLIF(TRIM(ul.requested_model), ''), NULLIF(TRIM(ul.model), ''), 'unknown')` // Tiered retention balances UI windows against storage: @@ -71,6 +71,13 @@ var channelMonitorV2RetentionRules = []channelMonitorV2RetentionRule{ } func (r *channelMonitorV2Repository) pruneChannelMonitorV2Retention(ctx context.Context, tx *sql.Tx, now time.Time) error { + // During historical bootstrap, retain all 1m facts until the cursor reaches + // the oldest rollup boundary. Otherwise adjacent chunks would rebuild the + // same daily bucket from source rows already pruned by the prior chunk. + var backfillCursor time.Time + if err := tx.QueryRowContext(ctx, `SELECT backfill_cursor FROM channel_monitor_v2_watermarks WHERE id = 1`).Scan(&backfillCursor); err == nil && backfillCursor.After(channelMonitorV2RetentionCutoff(now, channelMonitorV2RetentionMax)) { + return nil + } for _, rule := range channelMonitorV2RetentionRules { cutoff := channelMonitorV2RetentionCutoff(now, rule.retention) var err error @@ -164,11 +171,15 @@ INSERT INTO channel_monitor_v2_metrics_1m ( ) SELECT date_trunc('minute', ul.created_at), %s, COALESCE(ul.group_id, 0), %s, COUNT(DISTINCT COALESCE(NULLIF(ul.request_id, ''), 'usage:' || ul.id::text)) - FILTER (WHERE COALESCE(ul.request_type, 0) NOT IN (4, 6)), - SUM(COALESCE(ul.input_tokens, 0)), SUM(COALESCE(ul.output_tokens, 0)), - SUM(COALESCE(ul.cache_creation_tokens, 0)), SUM(COALESCE(ul.cache_read_tokens, 0)), - COALESCE(SUM(ul.first_token_ms) FILTER (WHERE ul.first_token_ms IS NOT NULL), 0), COUNT(ul.first_token_ms), - COALESCE(SUM(ul.duration_ms) FILTER (WHERE ul.duration_ms IS NOT NULL), 0), COUNT(ul.duration_ms), NOW() + FILTER (WHERE COALESCE(ul.request_type, 0) NOT IN (4, 6) AND ` + usageLogSuccessFilterUL + `), + COALESCE(SUM(ul.input_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.output_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.cache_creation_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.cache_read_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.first_token_ms) FILTER (WHERE ul.first_token_ms IS NOT NULL AND ` + usageLogSuccessFilterUL + `), 0), + COUNT(ul.first_token_ms) FILTER (WHERE ` + usageLogSuccessFilterUL + `), + COALESCE(SUM(ul.duration_ms) FILTER (WHERE ul.duration_ms IS NOT NULL AND ` + usageLogSuccessFilterUL + `), 0), + COUNT(ul.duration_ms) FILTER (WHERE ` + usageLogSuccessFilterUL + `), NOW() FROM usage_logs ul LEFT JOIN groups g ON g.id = ul.group_id LEFT JOIN accounts a ON a.id = ul.account_id @@ -183,11 +194,15 @@ INSERT INTO channel_monitor_v2_user_metrics_1m ( ) SELECT date_trunc('minute', ul.created_at), %s, COALESCE(ul.group_id, 0), %s, ul.user_id, COUNT(DISTINCT COALESCE(NULLIF(ul.request_id, ''), 'usage:' || ul.id::text)) - FILTER (WHERE COALESCE(ul.request_type, 0) NOT IN (4, 6)), - SUM(COALESCE(ul.input_tokens, 0)), SUM(COALESCE(ul.output_tokens, 0)), - SUM(COALESCE(ul.cache_creation_tokens, 0)), SUM(COALESCE(ul.cache_read_tokens, 0)), - COALESCE(SUM(ul.first_token_ms) FILTER (WHERE ul.first_token_ms IS NOT NULL), 0), COUNT(ul.first_token_ms), - COALESCE(SUM(ul.duration_ms) FILTER (WHERE ul.duration_ms IS NOT NULL), 0), COUNT(ul.duration_ms), NOW() + FILTER (WHERE COALESCE(ul.request_type, 0) NOT IN (4, 6) AND ` + usageLogSuccessFilterUL + `), + COALESCE(SUM(ul.input_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.output_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.cache_creation_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.cache_read_tokens) FILTER (WHERE ` + usageLogSuccessFilterUL + `), 0), + COALESCE(SUM(ul.first_token_ms) FILTER (WHERE ul.first_token_ms IS NOT NULL AND ` + usageLogSuccessFilterUL + `), 0), + COUNT(ul.first_token_ms) FILTER (WHERE ` + usageLogSuccessFilterUL + `), + COALESCE(SUM(ul.duration_ms) FILTER (WHERE ul.duration_ms IS NOT NULL AND ` + usageLogSuccessFilterUL + `), 0), + COUNT(ul.duration_ms) FILTER (WHERE ` + usageLogSuccessFilterUL + `), NOW() FROM usage_logs ul LEFT JOIN groups g ON g.id = ul.group_id LEFT JOIN accounts a ON a.id = ul.account_id @@ -207,6 +222,7 @@ CROSS JOIN LATERAL (VALUES (0::bigint), (ul.user_id)) audience(user_id) CROSS JOIN LATERAL (VALUES ('ttft'::text, ul.first_token_ms), ('duration'::text, ul.duration_ms)) latency(metric, value_ms) WHERE ul.created_at >= $1 AND ul.created_at < $2 AND audience.user_id IS NOT NULL AND latency.value_ms IS NOT NULL AND latency.value_ms >= 0 + AND ` + usageLogSuccessFilterUL + ` GROUP BY 1, 2, 3, 4, 5, 6, 7` func channelMonitorV2HistogramBoundSQL(column string) string { @@ -319,6 +335,13 @@ var channelMonitorV2FixedRollupSeconds = []int{300, 3600, 43200, 86400} func (r *channelMonitorV2Repository) recomputeFixedRollups(ctx context.Context, tx *sql.Tx, start, end time.Time) error { for _, seconds := range channelMonitorV2FixedRollupSeconds { + // Coarse buckets are immutable between boundaries during the normal + // trailing refresh. Historical backfills and boundary-crossing windows + // still rebuild them; this avoids repeatedly regrouping the full current + // day/user table every few minutes. + if seconds >= 43200 && sameFixedRollupBucket(start, end, seconds) { + continue + } interval := fmt.Sprintf("%d seconds", seconds) for _, table := range []string{ "channel_monitor_v2_latency_histograms_rollup", @@ -346,6 +369,14 @@ func (r *channelMonitorV2Repository) recomputeFixedRollups(ctx context.Context, return nil } +func sameFixedRollupBucket(start, end time.Time, seconds int) bool { + if !end.After(start) { + return true + } + interval := time.Duration(seconds) * time.Second + return start.Truncate(interval).Equal(end.Add(-time.Nanosecond).Truncate(interval)) +} + const channelMonitorV2FixedRollupBoundsSQL = ` WITH bounds AS ( SELECT diff --git a/backend/internal/repository/channel_monitor_v2_repo_test.go b/backend/internal/repository/channel_monitor_v2_repo_test.go index 7582b4c5aa..3dff580fef 100644 --- a/backend/internal/repository/channel_monitor_v2_repo_test.go +++ b/backend/internal/repository/channel_monitor_v2_repo_test.go @@ -106,7 +106,11 @@ func TestChannelMonitorV2ErrorAggregationCountsFinalUserErrorsOnly(t *testing.T) func TestChannelMonitorV2UsageSuccessExcludesCyberBillingRows(t *testing.T) { for _, query := range []string{channelMonitorV2UsageMetricsSQL, channelMonitorV2UserMetricsSQL} { require.Contains(t, query, "COALESCE(ul.request_type, 0) NOT IN (4, 6)") + require.Contains(t, query, "ul.actual_cost > 0") } + require.Contains(t, channelMonitorV2PlatformSQL, "g.platform = 'composite'") + require.Contains(t, channelMonitorV2PlatformSQL, "a.platform") + require.Contains(t, channelMonitorV2HistogramSQL, "ul.actual_cost > 0") } func TestChannelMonitorV2RatesUseCoveredWindow(t *testing.T) { @@ -170,6 +174,13 @@ func TestChannelMonitorV2TierRetentionPolicy(t *testing.T) { require.Equal(t, now.Add(-90*24*time.Hour), channelMonitorV2RetentionCutoff(now, channelMonitorV2MaxRetention())) } +func TestSameFixedRollupBucket(t *testing.T) { + start := time.Date(2026, 8, 7, 10, 0, 0, 0, time.UTC) + require.True(t, sameFixedRollupBucket(start, start.Add(10*time.Minute), 86400)) + require.False(t, sameFixedRollupBucket(start, start.Add(15*time.Hour), 43200)) + require.False(t, sameFixedRollupBucket(start, start.Add(24*time.Hour), 86400)) +} + // Needles present in service.ClassifyChannelMonitorV2Error must appear in the // aggregation SQL CASE so rollup categories match drilldown classification. func TestChannelMonitorV2SQLTaxonomyContainsGoNeedles(t *testing.T) { diff --git a/backend/internal/service/channel_monitor_v2.go b/backend/internal/service/channel_monitor_v2.go index 3f6b566500..df0cd71bc6 100644 --- a/backend/internal/service/channel_monitor_v2.go +++ b/backend/internal/service/channel_monitor_v2.go @@ -385,9 +385,14 @@ func (s *ChannelMonitorV2Service) SetRuntimeReader(r channelMonitorRuntimeReader } func (s *ChannelMonitorV2Service) hideThroughputForViewer(ctx context.Context, admin bool) bool { - if admin || s == nil || s.settings == nil { + if admin { return false } + // Privacy is fail-closed: an absent or unavailable settings reader must not + // expose fleet-scale rates to ordinary users. + if s == nil || s.settings == nil { + return true + } return s.settings.GetChannelMonitorRuntime(ctx).HideThroughput } diff --git a/backend/internal/service/channel_monitor_v2_aggregator.go b/backend/internal/service/channel_monitor_v2_aggregator.go index a70a869a50..8a2384f0cf 100644 --- a/backend/internal/service/channel_monitor_v2_aggregator.go +++ b/backend/internal/service/channel_monitor_v2_aggregator.go @@ -22,8 +22,9 @@ const ( channelMonitorV2BootstrapFirst = 2 * time.Hour // Subsequent bootstrap chunks grow so 24h/7d/30d fill without waiting full 90d pace. // Order after first tick: 22h → full 1d chunks until 30d, then 1d toward 90d. - channelMonitorV2RecentOverlap = 10 * time.Minute - channelMonitorV2BackfillChunk = 24 * time.Hour + channelMonitorV2RecentOverlap = 10 * time.Minute + channelMonitorV2BackfillChunk = 24 * time.Hour + channelMonitorV2MinBackfillChunk = time.Hour ) // channelMonitorRuntimeSubscriber is the optional settings hook that lets the @@ -45,7 +46,9 @@ type ChannelMonitorV2Aggregator struct { mu sync.Mutex // backfillAt is the earliest minute already recomputed (mirrors DB cursor). // Zero means "not yet loaded from durable watermark this process". - backfillAt time.Time + backfillAt time.Time + backfillChunk time.Duration + backfillFailures int // cursorLoaded is true after the first successful watermark read (or init). cursorLoaded bool // hasAggregated is true once any recompute in this process (or durable data) exists. @@ -57,12 +60,13 @@ type ChannelMonitorV2Aggregator struct { func NewChannelMonitorV2Aggregator(repo ChannelMonitorV2Repository, db *sql.DB, settings channelMonitorRuntimeReader) *ChannelMonitorV2Aggregator { return &ChannelMonitorV2Aggregator{ - repo: repo, - db: db, - settings: settings, - instanceID: uuid.NewString(), - stopCh: make(chan struct{}), - kickCh: make(chan struct{}, 1), + repo: repo, + db: db, + settings: settings, + instanceID: uuid.NewString(), + stopCh: make(chan struct{}), + kickCh: make(chan struct{}, 1), + backfillChunk: channelMonitorV2BackfillChunk, } } @@ -179,6 +183,9 @@ func (s *ChannelMonitorV2Aggregator) bootstrapActive() bool { if !s.hasAggregated { return true } + if s.backfillFailures >= 3 { + return false + } now := time.Now().UTC().Truncate(time.Minute) target := now.Add(-channelMonitorV2BootstrapWindow) return s.backfillAt.IsZero() || s.backfillAt.After(target) @@ -270,7 +277,18 @@ func (s *ChannelMonitorV2Aggregator) runOnce() { // past 24h so the default + daily UI ranges fill within a few ticks; 7d/30d // follow as more chunks complete (banner tracks progress against 30d). end := cursor - start := end.Add(-channelMonitorV2BackfillChunk) + s.mu.Lock() + chunk := s.backfillChunk + s.mu.Unlock() + if chunk <= 0 { + chunk = channelMonitorV2BackfillChunk + } + start := end.Add(-chunk) + // Once bootstrap reaches historical data, keep chunks on day boundaries so + // daily rollups never depend on 1m rows from two independently pruned chunks. + if end.Before(now.Add(-7 * 24 * time.Hour)) { + start = end.Add(-chunk).Truncate(24 * time.Hour) + } if start.Before(retentionCutoff) { start = retentionCutoff } @@ -279,9 +297,20 @@ func (s *ChannelMonitorV2Aggregator) runOnce() { } if err := s.repo.RecomputeRange(ctx, start, end); err != nil { logger.LegacyPrintf("service.channel_monitor_v2", "[ChannelMonitorV2] backfill failed %s..%s: %v", start, end, err) + s.mu.Lock() + s.backfillFailures++ + if s.backfillChunk > channelMonitorV2MinBackfillChunk { + s.backfillChunk /= 2 + if s.backfillChunk < channelMonitorV2MinBackfillChunk { + s.backfillChunk = channelMonitorV2MinBackfillChunk + } + } + s.mu.Unlock() return } s.mu.Lock() + s.backfillChunk = channelMonitorV2BackfillChunk + s.backfillFailures = 0 s.backfillAt = start s.hasAggregated = true s.mu.Unlock() diff --git a/backend/internal/service/channel_monitor_v2_test.go b/backend/internal/service/channel_monitor_v2_test.go index 8407946aed..9172b333df 100644 --- a/backend/internal/service/channel_monitor_v2_test.go +++ b/backend/internal/service/channel_monitor_v2_test.go @@ -422,7 +422,7 @@ func TestSnapshotRedactsPublicConfigPolicyFields(t *testing.T) { require.Equal(t, 300, snap.Config.RefreshIntervalSeconds) require.Zero(t, snap.Metrics.RequestCount) require.InDelta(t, 0.1, snap.Metrics.ErrorRate, 0.0001) - require.InDelta(t, 5.0, snap.Metrics.RPM, 0.0001) + require.Zero(t, snap.Metrics.RPM) } func TestRedactChannelMonitorV2MetricKeepsRates(t *testing.T) { diff --git a/backend/internal/service/setting_public.go b/backend/internal/service/setting_public.go index 53853d9be5..f3c86a8842 100644 --- a/backend/internal/service/setting_public.go +++ b/backend/internal/service/setting_public.go @@ -440,6 +440,7 @@ func (s *SettingService) GetChannelMonitorRuntime(ctx context.Context) ChannelMo Enabled: true, Mode: defaultChannelMonitorMode, DefaultIntervalSeconds: channelMonitorIntervalFallback, + HideThroughput: true, } } vals, err := s.settingRepo.GetMultiple(ctx, []string{ @@ -453,13 +454,14 @@ func (s *SettingService) GetChannelMonitorRuntime(ctx context.Context) ChannelMo Enabled: true, Mode: defaultChannelMonitorMode, DefaultIntervalSeconds: channelMonitorIntervalFallback, + HideThroughput: true, } } return ChannelMonitorRuntime{ Enabled: !isFalseSettingValue(vals[SettingKeyChannelMonitorEnabled]), Mode: normalizeChannelMonitorMode(vals[SettingKeyChannelMonitorMode]), DefaultIntervalSeconds: parseChannelMonitorInterval(vals[SettingKeyChannelMonitorDefaultIntervalSeconds]), - HideThroughput: vals[SettingKeyChannelMonitorHideThroughput] == "true", + HideThroughput: vals[SettingKeyChannelMonitorHideThroughput] != "false", } } diff --git a/backend/migrations/203_channel_monitor_v2_default_ignore_and_cache.sql b/backend/migrations/203_channel_monitor_v2_default_ignore_and_cache.sql index 4a03f0054f..48a80d5d33 100644 --- a/backend/migrations/203_channel_monitor_v2_default_ignore_and_cache.sql +++ b/backend/migrations/203_channel_monitor_v2_default_ignore_and_cache.sql @@ -1,7 +1,9 @@ -- Factory presets for Channel Monitor V2 config: -- 1) ignored_error_categories: non-ops client/policy failures that should not -- dominate health error_rate by default (still shown greyed in breakdown). --- 2) health_thresholds cache floors: 85% watch / 60% critical (higher is better). +-- 2) health_thresholds cache floors are intentionally left at the tolerant +-- zero/zero defaults from migration 198. Platforms without prompt caching +-- must not be marked unhealthy by factory configuration. -- -- Only apply when the row still looks like factory empty ignore list and/or -- zero cache thresholds (operator customizations are left alone). @@ -19,13 +21,3 @@ SET ignored_error_categories = ARRAY[ ]::text[] WHERE id = 1 AND COALESCE(cardinality(ignored_error_categories), 0) = 0; - -UPDATE channel_monitor_v2_config -SET health_thresholds = health_thresholds - || jsonb_build_object( - 'warning_cache_rate', 0.85, - 'critical_cache_rate', 0.60 - ) -WHERE id = 1 - AND COALESCE((health_thresholds->>'warning_cache_rate')::float8, 0) = 0 - AND COALESCE((health_thresholds->>'critical_cache_rate')::float8, 0) = 0; diff --git a/backend/migrations/204_channel_monitor_hide_throughput.sql b/backend/migrations/204_channel_monitor_hide_throughput.sql index 194eb7ce60..1c706bb2ab 100644 --- a/backend/migrations/204_channel_monitor_hide_throughput.sql +++ b/backend/migrations/204_channel_monitor_hide_throughput.sql @@ -1,5 +1,5 @@ -- Soft switch: hide RPM/TPM throughput rates on user-facing Channel Monitor V2. --- Default false (rates visible). Admins always see full metrics. +-- Default true (rates hidden). Admins always see full metrics. INSERT INTO settings (key, value) -VALUES ('channel_monitor_hide_throughput', 'false') +VALUES ('channel_monitor_hide_throughput', 'true') ON CONFLICT (key) DO NOTHING; diff --git a/backend/migrations/205_channel_monitor_v2_reset_factory_cache_thresholds.sql b/backend/migrations/205_channel_monitor_v2_reset_factory_cache_thresholds.sql new file mode 100644 index 0000000000..f30e4e71f0 --- /dev/null +++ b/backend/migrations/205_channel_monitor_v2_reset_factory_cache_thresholds.sql @@ -0,0 +1,13 @@ +-- Keep factory cache scoring tolerant. Migration 203 briefly wrote 0.85/0.60; +-- reset only that exact factory pair, preserving operator-customized values. +UPDATE channel_monitor_v2_config +SET health_thresholds = health_thresholds + || jsonb_build_object( + 'warning_cache_rate', 0, + 'critical_cache_rate', 0 + ) +WHERE id = 1 + AND COALESCE((health_thresholds->>'warning_cache_rate')::float8, 0) = 0.85 + AND COALESCE((health_thresholds->>'critical_cache_rate')::float8, 0) = 0.60 + AND version = 1 + AND updated_by IS NULL;