fix(channel-monitor-v2): correct aggregation privacy and backfill

This commit is contained in:
IanShaw027
2026-08-08 01:59:44 +08:00
parent 4a9459d607
commit 0d98176c59
10 changed files with 122 additions and 38 deletions
@@ -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 ---
@@ -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
@@ -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) {
@@ -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
}
@@ -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()
@@ -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) {
+3 -1
View File
@@ -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",
}
}
@@ -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;
@@ -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;
@@ -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;