From 41e322c53d7b4b796eb377d0df9c29ecd10ba431 Mon Sep 17 00:00:00 2001 From: Chia Date: Thu, 6 Aug 2026 09:29:41 +1200 Subject: feat: complete commercial control plane, billing, auth, and model catalog - add PostgreSQL control-plane persistence with Redis-degraded hot reload - implement prepaid balance, usage ledger, Stripe top-up and reconciliation - add registration, email verification, password reset, invitations and RBAC - support TOTP, Passkey MFA, device sessions, quotas and rate limits - add tenant billing profiles, audit logs and operational readiness checks - build authenticated admin console, Quickstart, Playground and usage analytics - add public model catalog with pricing, filtering and cost estimation - support OpenAI Responses providers and provider health failover - validate real upstream usage reporting and balance settlement --- internal/controlplane/usage_analytics.go | 215 +++++++++++++++++++++++++++++++ 1 file changed, 215 insertions(+) create mode 100644 internal/controlplane/usage_analytics.go (limited to 'internal/controlplane/usage_analytics.go') diff --git a/internal/controlplane/usage_analytics.go b/internal/controlplane/usage_analytics.go new file mode 100644 index 0000000..7cc042f --- /dev/null +++ b/internal/controlplane/usage_analytics.go @@ -0,0 +1,215 @@ +package controlplane + +import ( + "context" + "fmt" + "strings" + "time" +) + +// UsageAnalytics aggregates the immutable usage ledger for the developer +// console. Queries run outside the inference path and are scoped by the +// caller's tenant before reaching this store. +func (s *Store) UsageAnalytics(ctx context.Context, query UsageQuery) (UsageAnalytics, error) { + to := query.To + if to.IsZero() { + to = time.Now().UTC() + } + from := query.From + if from.IsZero() { + from = to.Add(-30 * 24 * time.Hour) + } + if !to.After(from) { + return UsageAnalytics{}, fmt.Errorf("usage analytics range must be positive") + } + result := UsageAnalytics{RangeStart: from, RangeEnd: to, Models: make([]UsageModelAnalytics, 0), Providers: make([]UsageProviderAnalytics, 0)} + + modelPrevious, err := s.usageModelCharges(ctx, query, from.Add(-to.Sub(from)), from) + if err != nil { + return UsageAnalytics{}, err + } + models, err := s.usageModelAnalytics(ctx, query, from, to) + if err != nil { + return UsageAnalytics{}, err + } + for index := range models { + models[index].PreviousChargedMicros = modelPrevious[models[index].PublicModel] + models[index].ChargeChangePercent = chargeChange(models[index].ChargedMicros, models[index].PreviousChargedMicros) + } + + providerPrevious, err := s.usageProviderCharges(ctx, query, from.Add(-to.Sub(from)), from) + if err != nil { + return UsageAnalytics{}, err + } + providers, err := s.usageProviderAnalytics(ctx, query, from, to) + if err != nil { + return UsageAnalytics{}, err + } + for index := range providers { + providers[index].PreviousChargedMicros = providerPrevious[providers[index].ProviderID] + providers[index].ChargeChangePercent = chargeChange(providers[index].ChargedMicros, providers[index].PreviousChargedMicros) + } + result.Models = models + result.Providers = providers + return result, nil +} + +func chargeChange(current, previous int64) *float64 { + if previous == 0 { + return nil + } + value := (float64(current) - float64(previous)) / float64(previous) * 100 + return &value +} + +func analyticsUsageWhere(query UsageQuery, from, to time.Time) (string, []any) { + where := []string{"1=1"} + args := make([]any, 0, 12) + index := 1 + for _, item := range []struct { + value string + clause string + }{ + {query.TenantID, "e.tenant_id=$"}, + {query.ProjectID, "e.project_id=$"}, + {query.KeyID, "e.key_id=$"}, + {query.Model, "e.public_model=$"}, + {query.RequestID, "e.request_id=$"}, + } { + if strings.TrimSpace(item.value) != "" { + where = append(where, item.clause+fmt.Sprint(index)) + args = append(args, item.value) + index++ + } + } + for _, item := range []struct { + value string + clause string + }{ + {query.Protocol, "e.protocol=$"}, + {query.ErrorType, "e.error_type=$"}, + } { + if strings.TrimSpace(item.value) != "" { + where = append(where, item.clause+fmt.Sprint(index)) + args = append(args, item.value) + index++ + } + } + if strings.TrimSpace(query.Provider) != "" { + where = append(where, "EXISTS (SELECT 1 FROM providers filter_provider WHERE filter_provider.id::text=e.provider_id AND filter_provider.slug=$"+fmt.Sprint(index)+")") + args = append(args, query.Provider) + index++ + } + if query.Stream != nil { + where = append(where, "e.stream=$"+fmt.Sprint(index)) + args = append(args, *query.Stream) + index++ + } + if query.Status == "success" { + where = append(where, "e.success=TRUE") + } else if query.Status == "error" { + where = append(where, "e.success=FALSE") + } + if !from.IsZero() { + where = append(where, "e.started_at >= $"+fmt.Sprint(index)) + args = append(args, from) + index++ + } + if !to.IsZero() { + where = append(where, "e.started_at < $"+fmt.Sprint(index)) + args = append(args, to) + } + return strings.Join(where, " AND "), args +} + +func (s *Store) usageModelAnalytics(ctx context.Context, query UsageQuery, from, to time.Time) ([]UsageModelAnalytics, error) { + where, args := analyticsUsageWhere(query, from, to) + rows, err := s.db.Query(ctx, `SELECT e.public_model, count(*), count(*) FILTER (WHERE e.success), count(*) FILTER (WHERE NOT e.success), + count(DISTINCT NULLIF(e.provider_id,'')), COALESCE(sum(e.input_tokens),0), COALESCE(sum(e.output_tokens),0), + COALESCE(sum(e.total_tokens),0), COALESCE(sum(e.cache_read_input_tokens),0), COALESCE(sum(e.cache_creation_input_tokens),0), + COALESCE(sum(e.charged_micros),0), COALESCE(sum(e.uncollected_micros),0), count(*) FILTER (WHERE e.metering_status='missing'), + COALESCE(round(avg(e.duration_ms)),0)::bigint, COALESCE(round(percentile_cont(0.95) WITHIN GROUP (ORDER BY e.duration_ms)),0)::bigint + FROM usage_events e WHERE `+where+` GROUP BY e.public_model ORDER BY sum(e.charged_micros) DESC, e.public_model`, args...) + if err != nil { + return nil, fmt.Errorf("query usage model analytics: %w", err) + } + defer rows.Close() + result := make([]UsageModelAnalytics, 0) + for rows.Next() { + var item UsageModelAnalytics + if err := rows.Scan(&item.PublicModel, &item.RequestCount, &item.SuccessfulRequests, &item.ErrorCount, &item.ProviderCount, + &item.InputTokens, &item.OutputTokens, &item.TotalTokens, &item.CacheReadInputTokens, &item.CacheCreationInputTokens, + &item.ChargedMicros, &item.UncollectedMicros, &item.MissingUsageRequests, &item.AverageDurationMS, &item.P95DurationMS); err != nil { + return nil, fmt.Errorf("scan usage model analytics: %w", err) + } + result = append(result, item) + } + return result, rows.Err() +} + +func (s *Store) usageModelCharges(ctx context.Context, query UsageQuery, from, to time.Time) (map[string]int64, error) { + where, args := analyticsUsageWhere(query, from, to) + rows, err := s.db.Query(ctx, `SELECT e.public_model, COALESCE(sum(e.charged_micros),0) + FROM usage_events e WHERE `+where+` GROUP BY e.public_model`, args...) + if err != nil { + return nil, fmt.Errorf("query previous model charges: %w", err) + } + defer rows.Close() + result := make(map[string]int64) + for rows.Next() { + var model string + var charged int64 + if err := rows.Scan(&model, &charged); err != nil { + return nil, fmt.Errorf("scan previous model charges: %w", err) + } + result[model] = charged + } + return result, rows.Err() +} + +func (s *Store) usageProviderAnalytics(ctx context.Context, query UsageQuery, from, to time.Time) ([]UsageProviderAnalytics, error) { + where, args := analyticsUsageWhere(query, from, to) + rows, err := s.db.Query(ctx, `SELECT COALESCE(e.provider_id,''), COALESCE(NULLIF(p.name,''),'Unassigned'), COALESCE(p.wire_api,''), + count(*), count(*) FILTER (WHERE e.success), count(*) FILTER (WHERE NOT e.success), count(DISTINCT e.public_model), + COALESCE(sum(e.input_tokens),0), COALESCE(sum(e.output_tokens),0), COALESCE(sum(e.total_tokens),0), + COALESCE(sum(e.cache_read_input_tokens),0), COALESCE(sum(e.cache_creation_input_tokens),0), COALESCE(sum(e.charged_micros),0), + COALESCE(sum(e.uncollected_micros),0), count(*) FILTER (WHERE e.metering_status='missing'), COALESCE(round(avg(e.duration_ms)),0)::bigint, + COALESCE(round(percentile_cont(0.95) WITHIN GROUP (ORDER BY e.duration_ms)),0)::bigint + FROM usage_events e LEFT JOIN providers p ON p.id::text=e.provider_id WHERE `+where+` + GROUP BY e.provider_id, p.name, p.wire_api ORDER BY sum(e.charged_micros) DESC, COALESCE(NULLIF(p.name,''),'Unassigned')`, args...) + if err != nil { + return nil, fmt.Errorf("query usage provider analytics: %w", err) + } + defer rows.Close() + result := make([]UsageProviderAnalytics, 0) + for rows.Next() { + var item UsageProviderAnalytics + if err := rows.Scan(&item.ProviderID, &item.ProviderName, &item.WireAPI, &item.RequestCount, &item.SuccessfulRequests, &item.ErrorCount, + &item.ModelCount, &item.InputTokens, &item.OutputTokens, &item.TotalTokens, &item.CacheReadInputTokens, &item.CacheCreationInputTokens, + &item.ChargedMicros, &item.UncollectedMicros, &item.MissingUsageRequests, &item.AverageDurationMS, &item.P95DurationMS); err != nil { + return nil, fmt.Errorf("scan usage provider analytics: %w", err) + } + result = append(result, item) + } + return result, rows.Err() +} + +func (s *Store) usageProviderCharges(ctx context.Context, query UsageQuery, from, to time.Time) (map[string]int64, error) { + where, args := analyticsUsageWhere(query, from, to) + rows, err := s.db.Query(ctx, `SELECT COALESCE(e.provider_id,''), COALESCE(sum(e.charged_micros),0) + FROM usage_events e WHERE `+where+` GROUP BY e.provider_id`, args...) + if err != nil { + return nil, fmt.Errorf("query previous provider charges: %w", err) + } + defer rows.Close() + result := make(map[string]int64) + for rows.Next() { + var providerID string + var charged int64 + if err := rows.Scan(&providerID, &charged); err != nil { + return nil, fmt.Errorf("scan previous provider charges: %w", err) + } + result[providerID] = charged + } + return result, rows.Err() +} -- cgit v1.2.3