package controlplane import ( "context" "encoding/json" "fmt" "time" ) func (s *Store) Overview(ctx context.Context) (Overview, error) { var result Overview err := s.db.QueryRow(ctx, ` SELECT (SELECT generation FROM control_state WHERE singleton = TRUE), (SELECT count(*) FROM tenants WHERE status = 'active'), (SELECT count(*) FROM projects WHERE status = 'active'), (SELECT count(*) FROM api_keys WHERE status = 'active'), (SELECT count(*) FROM providers WHERE enabled = TRUE), (SELECT count(*) FROM models WHERE enabled = TRUE)`, ).Scan(&result.Generation, &result.Tenants, &result.Projects, &result.APIKeys, &result.Providers, &result.Models) if err != nil { return Overview{}, fmt.Errorf("query control-plane overview: %w", err) } return result, nil } func (s *Store) ListTenants(ctx context.Context) ([]Tenant, error) { return s.ListTenantsFor(ctx, "") } func (s *Store) ListTenantsFor(ctx context.Context, tenantID string) ([]Tenant, error) { query := `SELECT id::text, slug, name, status, created_at FROM tenants` args := []any{} if tenantID != "" { query += ` WHERE id=$1` args = append(args, tenantID) } query += ` ORDER BY created_at DESC` rows, err := s.db.Query(ctx, query, args...) if err != nil { return nil, fmt.Errorf("query tenants: %w", err) } defer rows.Close() result := make([]Tenant, 0) for rows.Next() { var item Tenant if err := rows.Scan(&item.ID, &item.Slug, &item.Name, &item.Status, &item.CreatedAt); err != nil { return nil, fmt.Errorf("scan tenant: %w", err) } result = append(result, item) } return result, rows.Err() } func (s *Store) ListProjects(ctx context.Context) ([]Project, error) { return s.ListProjectsFor(ctx, "") } func (s *Store) ListProjectsFor(ctx context.Context, tenantID string) ([]Project, error) { query := `SELECT id::text, tenant_id::text, slug, name, status, created_at FROM projects` args := []any{} if tenantID != "" { query += ` WHERE tenant_id=$1` args = append(args, tenantID) } query += ` ORDER BY created_at DESC` rows, err := s.db.Query(ctx, query, args...) if err != nil { return nil, fmt.Errorf("query projects: %w", err) } defer rows.Close() result := make([]Project, 0) for rows.Next() { var item Project if err := rows.Scan(&item.ID, &item.TenantID, &item.Slug, &item.Name, &item.Status, &item.CreatedAt); err != nil { return nil, fmt.Errorf("scan project: %w", err) } result = append(result, item) } return result, rows.Err() } func (s *Store) ListAPIKeys(ctx context.Context) ([]APIKey, error) { return s.ListAPIKeysFor(ctx, "") } func (s *Store) ListAPIKeysFor(ctx context.Context, tenantID string) ([]APIKey, error) { now := time.Now().UTC() periodStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC) periodEnd := periodStart.AddDate(0, 1, 0) dayStart := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, time.UTC) dayEnd := dayStart.AddDate(0, 0, 1) query := ` SELECT k.id::text, k.tenant_id::text, k.project_id::text, k.name, k.key_prefix, k.key_suffix, k.scopes, k.tags, k.monthly_spend_micros, k.daily_spend_micros, k.requests_per_minute, k.tokens_per_minute, k.status, k.expires_at, k.last_used_at, k.created_at, usage.month_spend, usage.month_requests, pending.month_reserved, daily.day_spend, daily.day_requests, daily_pending.day_reserved, COALESCE(( SELECT jsonb_agg(m.public_id ORDER BY m.public_id) FROM api_key_model_restrictions r JOIN models m ON m.id = r.model_id WHERE r.api_key_id = k.id ), '[]'::jsonb) FROM api_keys k CROSS JOIN LATERAL ( SELECT COALESCE(SUM(u.cost_micros), 0)::bigint AS month_spend, COUNT(*)::bigint AS month_requests FROM usage_events u WHERE u.key_id = k.id AND u.started_at >= $1 AND u.started_at < $2 ) usage CROSS JOIN LATERAL ( SELECT COALESCE(SUM(b.reserved_micros), 0)::bigint AS month_reserved FROM billing_reservations b WHERE b.key_id = k.id AND b.status IN ('pending', 'metering_failed') AND b.created_at >= $1 AND b.created_at < $2 ) pending CROSS JOIN LATERAL ( SELECT COALESCE(SUM(u.cost_micros), 0)::bigint AS day_spend, COUNT(*)::bigint AS day_requests FROM usage_events u WHERE u.key_id = k.id AND u.started_at >= $3 AND u.started_at < $4 ) daily CROSS JOIN LATERAL ( SELECT COALESCE(SUM(b.reserved_micros), 0)::bigint AS day_reserved FROM billing_reservations b WHERE b.key_id = k.id AND b.status IN ('pending', 'metering_failed') AND b.created_at >= $3 AND b.created_at < $4 ) daily_pending` args := []any{periodStart, periodEnd, dayStart, dayEnd} if tenantID != "" { query += ` WHERE k.tenant_id=$5` args = append(args, tenantID) } query += ` ORDER BY k.created_at DESC` rows, err := s.db.Query(ctx, query, args...) if err != nil { return nil, fmt.Errorf("query API keys: %w", err) } defer rows.Close() result := make([]APIKey, 0) for rows.Next() { var item APIKey var scopesJSON, tagsJSON, allowedModelsJSON []byte if err := rows.Scan(&item.ID, &item.TenantID, &item.ProjectID, &item.Name, &item.KeyPrefix, &item.KeySuffix, &scopesJSON, &tagsJSON, &item.MonthlySpendMicros, &item.DailySpendMicros, &item.RequestsPerMinute, &item.TokensPerMinute, &item.Status, &item.ExpiresAt, &item.LastUsedAt, &item.CreatedAt, &item.CurrentMonthSpendMicros, &item.CurrentMonthRequests, &item.CurrentMonthReservedMicros, &item.CurrentDaySpendMicros, &item.CurrentDayRequests, &item.CurrentDayReservedMicros, &allowedModelsJSON); err != nil { return nil, fmt.Errorf("scan API key: %w", err) } if err := json.Unmarshal(scopesJSON, &item.Scopes); err != nil { return nil, fmt.Errorf("decode API key scopes: %w", err) } if err := json.Unmarshal(tagsJSON, &item.Tags); err != nil { return nil, fmt.Errorf("decode API key tags: %w", err) } if err := json.Unmarshal(allowedModelsJSON, &item.AllowedModels); err != nil { return nil, fmt.Errorf("decode API key model restrictions: %w", err) } result = append(result, item) } return result, rows.Err() } func (s *Store) OverviewFor(ctx context.Context, tenantID string) (Overview, error) { if tenantID == "" { return s.Overview(ctx) } var result Overview err := s.db.QueryRow(ctx, `SELECT (SELECT generation FROM control_state WHERE singleton=TRUE), (SELECT count(*) FROM tenants WHERE id=$1 AND status='active'), (SELECT count(*) FROM projects WHERE tenant_id=$1 AND status='active'), (SELECT count(*) FROM api_keys WHERE tenant_id=$1 AND status='active'), 0, (SELECT count(*) FROM models WHERE enabled=TRUE)`, tenantID, ).Scan(&result.Generation, &result.Tenants, &result.Projects, &result.APIKeys, &result.Providers, &result.Models) if err != nil { return Overview{}, fmt.Errorf("query tenant overview: %w", err) } return result, nil } func (s *Store) ResourceTenantID(ctx context.Context, resource, id string) (string, error) { var query string switch resource { case "project": query = `SELECT tenant_id::text FROM projects WHERE id=$1` case "api_key": query = `SELECT tenant_id::text FROM api_keys WHERE id=$1` case "console_user": query = `SELECT COALESCE(tenant_id::text,'') FROM console_users WHERE id=$1` default: return "", ErrNotFound } var tenantID string if err := s.db.QueryRow(ctx, query, id).Scan(&tenantID); err != nil { return "", ErrNotFound } return tenantID, nil } func (s *Store) ListProviders(ctx context.Context) ([]Provider, error) { rows, err := s.db.Query(ctx, ` SELECT p.id::text, p.slug, p.name, p.protocol, p.wire_api, p.base_url, p.enabled, count(r.id), p.created_at FROM providers p LEFT JOIN model_routes r ON r.provider_id = p.id GROUP BY p.id ORDER BY p.created_at DESC`) if err != nil { return nil, fmt.Errorf("query providers: %w", err) } defer rows.Close() result := make([]Provider, 0) for rows.Next() { var item Provider if err := rows.Scan(&item.ID, &item.Slug, &item.Name, &item.Protocol, &item.WireAPI, &item.BaseURL, &item.Enabled, &item.RouteCount, &item.CreatedAt); err != nil { return nil, fmt.Errorf("scan provider: %w", err) } result = append(result, item) } return result, rows.Err() } func (s *Store) ListModels(ctx context.Context) ([]Model, error) { rows, err := s.db.Query(ctx, ` SELECT m.id::text, m.public_id, m.display_name, m.description, m.owned_by, m.input_modalities, m.output_modalities, m.context_window, m.max_output_tokens, m.capabilities, m.regions, m.lifecycle, m.released_at, m.deprecated_at, m.retired_at, COALESCE(m.replacement_model,''), pv.id::text, pv.version, pv.currency, pv.effective_from, pv.input_price_micros_per_million, pv.output_price_micros_per_million, pv.cache_read_price_micros_per_million, pv.cache_write_price_micros_per_million, m.enabled, m.created_at FROM models m JOIN LATERAL ( SELECT * FROM model_price_versions v WHERE v.model_id=m.id AND v.effective_from <= now() AND (v.effective_to IS NULL OR v.effective_to > now()) ORDER BY v.effective_from DESC LIMIT 1 ) pv ON TRUE ORDER BY m.public_id`) if err != nil { return nil, fmt.Errorf("query models: %w", err) } models := make([]Model, 0) positions := make(map[string]int) for rows.Next() { var item Model var inputModalitiesJSON, outputModalitiesJSON, capabilitiesJSON, regionsJSON []byte if err := rows.Scan(&item.ID, &item.PublicID, &item.DisplayName, &item.Description, &item.OwnedBy, &inputModalitiesJSON, &outputModalitiesJSON, &item.ContextWindow, &item.MaxOutputTokens, &capabilitiesJSON, ®ionsJSON, &item.Lifecycle, &item.ReleasedAt, &item.DeprecatedAt, &item.RetiredAt, &item.ReplacementModel, &item.PriceVersionID, &item.PriceVersion, &item.PriceCurrency, &item.PriceEffectiveFrom, &item.InputPriceMicrosPerMillion, &item.OutputPriceMicrosPerMillion, &item.CacheReadPriceMicrosPerMillion, &item.CacheWritePriceMicrosPerMillion, &item.Enabled, &item.CreatedAt); err != nil { rows.Close() return nil, fmt.Errorf("scan model: %w", err) } _ = json.Unmarshal(inputModalitiesJSON, &item.InputModalities) _ = json.Unmarshal(outputModalitiesJSON, &item.OutputModalities) _ = json.Unmarshal(capabilitiesJSON, &item.Capabilities) _ = json.Unmarshal(regionsJSON, &item.Regions) item.Routes = []Route{} item.Aliases, item.AllowedTenantIDs, item.AllowedKeyIDs = []string{}, []string{}, []string{} positions[item.ID] = len(models) models = append(models, item) } if err := rows.Err(); err != nil { rows.Close() return nil, err } rows.Close() aliasRows, err := s.db.Query(ctx, `SELECT model_id::text, alias FROM model_aliases ORDER BY alias`) if err != nil { return nil, fmt.Errorf("query model aliases: %w", err) } for aliasRows.Next() { var modelID, alias string if err := aliasRows.Scan(&modelID, &alias); err != nil { aliasRows.Close() return nil, err } if p, ok := positions[modelID]; ok { models[p].Aliases = append(models[p].Aliases, alias) } } aliasRows.Close() tenantRows, err := s.db.Query(ctx, `SELECT model_id::text, tenant_id::text FROM model_tenant_allowlist`) if err != nil { return nil, fmt.Errorf("query tenant model allowlist: %w", err) } for tenantRows.Next() { var modelID, tenantID string if err := tenantRows.Scan(&modelID, &tenantID); err != nil { tenantRows.Close() return nil, err } if p, ok := positions[modelID]; ok { models[p].AllowedTenantIDs = append(models[p].AllowedTenantIDs, tenantID) } } tenantRows.Close() keyRows, err := s.db.Query(ctx, `SELECT model_id::text, api_key_id::text FROM api_key_model_allowlist`) if err != nil { return nil, fmt.Errorf("query key model allowlist: %w", err) } for keyRows.Next() { var modelID, keyID string if err := keyRows.Scan(&modelID, &keyID); err != nil { keyRows.Close() return nil, err } if p, ok := positions[modelID]; ok { models[p].AllowedKeyIDs = append(models[p].AllowedKeyIDs, keyID) } } keyRows.Close() routeRows, err := s.db.Query(ctx, ` SELECT r.id::text, r.model_id::text, r.provider_id::text, p.name, p.protocol, p.wire_api, p.enabled, r.upstream_model, r.priority, r.weight, r.enabled FROM model_routes r JOIN providers p ON p.id = r.provider_id ORDER BY r.priority, r.created_at`) if err != nil { return nil, fmt.Errorf("query model routes: %w", err) } defer routeRows.Close() for routeRows.Next() { var route Route var modelID string if err := routeRows.Scan(&route.ID, &modelID, &route.ProviderID, &route.ProviderName, &route.Protocol, &route.WireAPI, &route.ProviderEnabled, &route.UpstreamModel, &route.Priority, &route.Weight, &route.Enabled); err != nil { return nil, fmt.Errorf("scan model route: %w", err) } if position, ok := positions[modelID]; ok { models[position].Routes = append(models[position].Routes, route) } } return models, routeRows.Err() } func (s *Store) ListDeveloperModels(ctx context.Context, tenantID string) ([]DeveloperModel, error) { models, err := s.ListModels(ctx) if err != nil { return nil, err } keyIDs := map[string]struct{}{} if tenantID != "" { rows, err := s.db.Query(ctx, `SELECT id::text FROM api_keys WHERE tenant_id=$1 AND status='active'`, tenantID) if err != nil { return nil, fmt.Errorf("query developer API keys: %w", err) } for rows.Next() { var id string if err := rows.Scan(&id); err != nil { rows.Close() return nil, err } keyIDs[id] = struct{}{} } if err := rows.Err(); err != nil { rows.Close() return nil, err } rows.Close() } return developerModelsFor(models, tenantID, keyIDs), nil } func (s *Store) ListPublicModels(ctx context.Context) ([]PublicModel, error) { models, err := s.ListModels(ctx) if err != nil { return nil, err } return publicModelsFor(models), nil } func publicModelsFor(models []Model) []PublicModel { result := make([]PublicModel, 0, len(models)) for _, model := range models { if !model.Enabled || model.Lifecycle == "retired" || len(model.AllowedTenantIDs) != 0 || len(model.AllowedKeyIDs) != 0 { continue } wireSet := make(map[string]struct{}) providerSet := make(map[string]struct{}) wireAPIs := make([]string, 0, len(model.Routes)) for _, route := range model.Routes { if !route.Enabled || !route.ProviderEnabled { continue } wireAPI := route.WireAPI if wireAPI == "" && route.Protocol == "anthropic" { wireAPI = "messages" } else if wireAPI == "" { wireAPI = "chat_completions" } if _, exists := wireSet[wireAPI]; !exists { wireSet[wireAPI] = struct{}{} wireAPIs = append(wireAPIs, wireAPI) } providerSet[route.ProviderID] = struct{}{} } if len(wireAPIs) == 0 { continue } result = append(result, PublicModel{PublicID: model.PublicID, DisplayName: model.DisplayName, Description: model.Description, OwnedBy: model.OwnedBy, InputModalities: model.InputModalities, OutputModalities: model.OutputModalities, ContextWindow: model.ContextWindow, MaxOutputTokens: model.MaxOutputTokens, Capabilities: model.Capabilities, Regions: model.Regions, Lifecycle: model.Lifecycle, ReleasedAt: model.ReleasedAt, ReplacementModel: model.ReplacementModel, Aliases: model.Aliases, PriceCurrency: model.PriceCurrency, InputPriceMicrosPerMillion: model.InputPriceMicrosPerMillion, OutputPriceMicrosPerMillion: model.OutputPriceMicrosPerMillion, CacheReadPriceMicrosPerMillion: model.CacheReadPriceMicrosPerMillion, CacheWritePriceMicrosPerMillion: model.CacheWritePriceMicrosPerMillion, SupportedWireAPIs: wireAPIs, ProviderCount: len(providerSet), AvailableProviderCount: len(providerSet), HealthStatus: "available"}) } return result } func developerModelsFor(models []Model, tenantID string, keyIDs map[string]struct{}) []DeveloperModel { result := make([]DeveloperModel, 0, len(models)) for _, model := range models { if !model.Enabled || model.Lifecycle == "retired" || !stringAllowed(model.AllowedTenantIDs, tenantID) || !keyAllowed(model.AllowedKeyIDs, keyIDs) { continue } wireSet := map[string]struct{}{} wireAPIs := make([]string, 0, len(model.Routes)) for _, route := range model.Routes { if !route.Enabled || !route.ProviderEnabled { continue } wireAPI := route.WireAPI if wireAPI == "" && route.Protocol == "anthropic" { wireAPI = "messages" } else if wireAPI == "" { wireAPI = "chat_completions" } if _, exists := wireSet[wireAPI]; !exists { wireSet[wireAPI] = struct{}{} wireAPIs = append(wireAPIs, wireAPI) } } if len(wireAPIs) == 0 { continue } result = append(result, DeveloperModel{ID: model.ID, PublicID: model.PublicID, DisplayName: model.DisplayName, Description: model.Description, OwnedBy: model.OwnedBy, InputModalities: model.InputModalities, OutputModalities: model.OutputModalities, ContextWindow: model.ContextWindow, MaxOutputTokens: model.MaxOutputTokens, Capabilities: model.Capabilities, Regions: model.Regions, Lifecycle: model.Lifecycle, ReleasedAt: model.ReleasedAt, ReplacementModel: model.ReplacementModel, Aliases: model.Aliases, PriceCurrency: model.PriceCurrency, InputPriceMicrosPerMillion: model.InputPriceMicrosPerMillion, OutputPriceMicrosPerMillion: model.OutputPriceMicrosPerMillion, CacheReadPriceMicrosPerMillion: model.CacheReadPriceMicrosPerMillion, CacheWritePriceMicrosPerMillion: model.CacheWritePriceMicrosPerMillion, SupportedWireAPIs: wireAPIs}) } return result } func stringAllowed(allowed []string, value string) bool { if len(allowed) == 0 || value == "" { return true } for _, item := range allowed { if item == value { return true } } return false } func keyAllowed(allowed []string, keys map[string]struct{}) bool { if len(allowed) == 0 { return true } for _, id := range allowed { if _, ok := keys[id]; ok { return true } } return false }