diff --git a/Makefile b/Makefile index 164897929..8221f7c0b 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,4 @@ -.PHONY: all build run clean tidy test test-race test-dashboard test-e2e test-integration test-contract test-all lint lint-fix record-api swagger docs-openapi install-tools perf-check perf-bench infra image +.PHONY: all build run clean tidy test test-race test-dashboard test-e2e test-integration test-contract test-all lint lint-fix record-api swagger docs-openapi install-tools perf-check perf-bench infra image seed-demo-data all: build @@ -42,6 +42,11 @@ infra: image: docker compose --profile app up -d +# Seed rolling demo usage/audit data into SQLite. +# Usage: SQLITE_PATH=data/gomodel.db make seed-demo-data +seed-demo-data: + bash tools/seed-demo-data.sh + # Run unit tests only test: go test ./cmd/... ./internal/... ./config/... -v diff --git a/internal/auditlog/reader_postgresql.go b/internal/auditlog/reader_postgresql.go index b455e65cb..7dee35fe5 100644 --- a/internal/auditlog/reader_postgresql.go +++ b/internal/auditlog/reader_postgresql.go @@ -9,12 +9,18 @@ import ( "github.com/goccy/go-json" + "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) +type postgreSQLQueryer interface { + Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error) + QueryRow(ctx context.Context, sql string, args ...any) pgx.Row +} + // PostgreSQLReader implements Reader for PostgreSQL databases. type PostgreSQLReader struct { - pool *pgxpool.Pool + pool postgreSQLQueryer } // NewPostgreSQLReader creates a new PostgreSQL audit log reader. @@ -115,9 +121,10 @@ func (r *PostgreSQLReader) GetLogs(ctx context.Context, params LogQueryParams) ( var authKeyID *string var authMethod *string var userPath *string + var errorType *string if err := rows.Scan(&e.ID, &e.Timestamp, &e.DurationNs, &e.RequestedModel, &e.ResolvedModel, &e.Provider, &providerName, &e.AliasUsed, &workflowVersionID, &cacheType, &e.StatusCode, - &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &e.Stream, &e.ErrorType, &dataJSON); err != nil { + &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &e.Stream, &errorType, &dataJSON); err != nil { return nil, fmt.Errorf("failed to scan audit log row: %w", err) } if workflowVersionID != nil { @@ -140,6 +147,9 @@ func (r *PostgreSQLReader) GetLogs(ctx context.Context, params LogQueryParams) ( if userPath != nil { e.UserPath = *userPath } + if errorType != nil { + e.ErrorType = *errorType + } if dataJSON != nil && *dataJSON != "" { var data LogData @@ -257,9 +267,10 @@ func scanPostgreSQLLogEntry(rows interface { var authKeyID *string var authMethod *string var userPath *string + var errorType *string if err := rows.Scan(&e.ID, &e.Timestamp, &e.DurationNs, &e.RequestedModel, &e.ResolvedModel, &e.Provider, &providerName, &e.AliasUsed, &workflowVersionID, &cacheType, &e.StatusCode, - &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &e.Stream, &e.ErrorType, &dataJSON); err != nil { + &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &e.Stream, &errorType, &dataJSON); err != nil { return nil, fmt.Errorf("failed to scan audit log row: %w", err) } if workflowVersionID != nil { @@ -282,6 +293,9 @@ func scanPostgreSQLLogEntry(rows interface { if userPath != nil { e.UserPath = *userPath } + if errorType != nil { + e.ErrorType = *errorType + } if dataJSON != nil && *dataJSON != "" { var data LogData diff --git a/internal/auditlog/reader_postgresql_test.go b/internal/auditlog/reader_postgresql_test.go new file mode 100644 index 000000000..2255569b3 --- /dev/null +++ b/internal/auditlog/reader_postgresql_test.go @@ -0,0 +1,202 @@ +package auditlog + +import ( + "context" + "fmt" + "reflect" + "strings" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" +) + +type fakePostgreSQLRow struct { + values []any +} + +func (r fakePostgreSQLRow) Scan(dest ...any) error { + if len(dest) != len(r.values) { + return fmt.Errorf("scan destination count = %d, want %d", len(dest), len(r.values)) + } + for i, value := range r.values { + target := reflect.ValueOf(dest[i]) + if target.Kind() != reflect.Pointer || target.IsNil() { + return fmt.Errorf("scan destination %d is not a non-nil pointer", i) + } + elem := target.Elem() + if value == nil { + elem.Set(reflect.Zero(elem.Type())) + continue + } + if elem.Kind() == reflect.Pointer { + pointerValue := reflect.New(elem.Type().Elem()) + if err := assignScannedValue(pointerValue.Elem(), value); err != nil { + return fmt.Errorf("scan destination %d: %w", i, err) + } + elem.Set(pointerValue) + continue + } + if err := assignScannedValue(elem, value); err != nil { + return fmt.Errorf("scan destination %d: %w", i, err) + } + } + return nil +} + +func assignScannedValue(target reflect.Value, value any) error { + source := reflect.ValueOf(value) + if source.Type().AssignableTo(target.Type()) { + target.Set(source) + return nil + } + if source.Type().ConvertibleTo(target.Type()) { + target.Set(source.Convert(target.Type())) + return nil + } + return fmt.Errorf("cannot assign %T to %s", value, target.Type()) +} + +type fakePostgreSQLQueryer struct { + count int + rows pgx.Rows +} + +func (q fakePostgreSQLQueryer) QueryRow(_ context.Context, _ string, _ ...any) pgx.Row { + return fakePostgreSQLRow{values: []any{q.count}} +} + +func (q fakePostgreSQLQueryer) Query(_ context.Context, sql string, _ ...any) (pgx.Rows, error) { + if !strings.Contains(sql, "FROM audit_logs") { + return nil, fmt.Errorf("unexpected query: %s", sql) + } + return q.rows, nil +} + +type fakePostgreSQLRows struct { + values []any + read bool + closed bool + err error +} + +func (r *fakePostgreSQLRows) Close() { + r.closed = true +} + +func (r *fakePostgreSQLRows) Err() error { + return r.err +} + +func (r *fakePostgreSQLRows) CommandTag() pgconn.CommandTag { + return pgconn.CommandTag{} +} + +func (r *fakePostgreSQLRows) FieldDescriptions() []pgconn.FieldDescription { + return nil +} + +func (r *fakePostgreSQLRows) Next() bool { + if r.read { + r.Close() + return false + } + r.read = true + return true +} + +func (r *fakePostgreSQLRows) Scan(dest ...any) error { + return fakePostgreSQLRow{values: r.values}.Scan(dest...) +} + +func (r *fakePostgreSQLRows) Values() ([]any, error) { + return r.values, nil +} + +func (r *fakePostgreSQLRows) RawValues() [][]byte { + return nil +} + +func (r *fakePostgreSQLRows) Conn() *pgx.Conn { + return nil +} + +func postgreSQLAuditLogRowValues(errorType any) []any { + return []any{ + "entry-null-error-type", + time.Unix(1700000000, 0).UTC(), + int64(1234), + "gpt-4o-mini", + "gpt-4o-mini", + "openai", + "primary-openai", + false, + nil, + nil, + 200, + "req-1", + nil, + "master_key", + "127.0.0.1", + "POST", + "/v1/chat/completions", + "/", + false, + errorType, + `{"user_agent":"test-agent"}`, + } +} + +func TestPostgreSQLReaderGetLogsAllowsNullErrorType(t *testing.T) { + rows := &fakePostgreSQLRows{values: postgreSQLAuditLogRowValues(nil)} + reader := &PostgreSQLReader{ + pool: fakePostgreSQLQueryer{ + count: 1, + rows: rows, + }, + } + + result, err := reader.GetLogs(context.Background(), LogQueryParams{Limit: 10}) + if err != nil { + t.Fatalf("GetLogs failed: %v", err) + } + if result.Total != 1 { + t.Fatalf("Total = %d, want 1", result.Total) + } + if len(result.Entries) != 1 { + t.Fatalf("len(Entries) = %d, want 1", len(result.Entries)) + } + entry := result.Entries[0] + if entry.ErrorType != "" { + t.Fatalf("ErrorType = %q, want empty", entry.ErrorType) + } + if entry.ProviderName != "primary-openai" { + t.Fatalf("ProviderName = %q, want primary-openai", entry.ProviderName) + } + if entry.Data == nil || entry.Data.UserAgent != "test-agent" { + t.Fatalf("Data = %#v, want user_agent", entry.Data) + } + if !rows.closed { + t.Fatal("rows were not closed") + } +} + +func TestScanPostgreSQLLogEntryAllowsNullErrorType(t *testing.T) { + entry, err := scanPostgreSQLLogEntry(fakePostgreSQLRow{values: postgreSQLAuditLogRowValues(nil)}) + if err != nil { + t.Fatalf("scanPostgreSQLLogEntry failed: %v", err) + } + if entry.ErrorType != "" { + t.Fatalf("ErrorType = %q, want empty", entry.ErrorType) + } + if entry.ProviderName != "primary-openai" { + t.Fatalf("ProviderName = %q, want primary-openai", entry.ProviderName) + } + if entry.AuthMethod != "master_key" { + t.Fatalf("AuthMethod = %q, want master_key", entry.AuthMethod) + } + if entry.Data == nil || entry.Data.UserAgent != "test-agent" { + t.Fatalf("Data = %#v, want user_agent", entry.Data) + } +} diff --git a/internal/auditlog/reader_sqlite.go b/internal/auditlog/reader_sqlite.go index bdf06fbb2..96061c3db 100644 --- a/internal/auditlog/reader_sqlite.go +++ b/internal/auditlog/reader_sqlite.go @@ -113,9 +113,10 @@ func (r *SQLiteReader) GetLogs(ctx context.Context, params LogQueryParams) (*Log var authKeyID sql.NullString var authMethod sql.NullString var userPath sql.NullString + var errorType sql.NullString if err := rows.Scan(&e.ID, &ts, &e.DurationNs, &e.RequestedModel, &e.ResolvedModel, &e.Provider, &providerName, &aliasUsedInt, &workflowVersionID, &cacheType, &e.StatusCode, - &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &streamInt, &e.ErrorType, &dataJSON); err != nil { + &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &streamInt, &errorType, &dataJSON); err != nil { return nil, fmt.Errorf("failed to scan audit log row: %w", err) } @@ -142,6 +143,9 @@ func (r *SQLiteReader) GetLogs(ctx context.Context, params LogQueryParams) (*Log if userPath.Valid { e.UserPath = userPath.String } + if errorType.Valid { + e.ErrorType = errorType.String + } if dataJSON != nil && *dataJSON != "" { var data LogData @@ -343,9 +347,10 @@ func scanSQLiteLogEntry(rows *sql.Rows) (*LogEntry, error) { var authKeyID sql.NullString var authMethod sql.NullString var userPath sql.NullString + var errorType sql.NullString if err := rows.Scan(&e.ID, &ts, &e.DurationNs, &e.RequestedModel, &e.ResolvedModel, &e.Provider, &providerName, &aliasUsedInt, &workflowVersionID, &cacheType, &e.StatusCode, - &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &streamInt, &e.ErrorType, &dataJSON); err != nil { + &e.RequestID, &authKeyID, &authMethod, &e.ClientIP, &e.Method, &e.Path, &userPath, &streamInt, &errorType, &dataJSON); err != nil { return nil, fmt.Errorf("failed to scan audit log row: %w", err) } @@ -372,6 +377,9 @@ func scanSQLiteLogEntry(rows *sql.Rows) (*LogEntry, error) { if userPath.Valid { e.UserPath = userPath.String } + if errorType.Valid { + e.ErrorType = errorType.String + } if dataJSON != nil && *dataJSON != "" { var data LogData diff --git a/internal/auditlog/store_sqlite_test.go b/internal/auditlog/store_sqlite_test.go index 0f6d20baa..9e81d219a 100644 --- a/internal/auditlog/store_sqlite_test.go +++ b/internal/auditlog/store_sqlite_test.go @@ -279,7 +279,7 @@ func TestSQLiteStore_WriteBatch_PersistsAliasFields(t *testing.T) { } } -func TestSQLiteReader_AllowsNullWorkflowVersionID(t *testing.T) { +func TestSQLiteReader_AllowsNullWorkflowVersionIDAndErrorType(t *testing.T) { db := createTestDB(t) defer db.Close() @@ -310,7 +310,7 @@ func TestSQLiteReader_AllowsNullWorkflowVersionID(t *testing.T) { "POST", "/v1/chat/completions", 0, - "", + nil, nil, ); err != nil { t.Fatalf("failed to insert audit log row: %v", err) @@ -332,6 +332,9 @@ func TestSQLiteReader_AllowsNullWorkflowVersionID(t *testing.T) { if entry.WorkflowVersionID != "" { t.Fatalf("WorkflowVersionID = %q, want empty", entry.WorkflowVersionID) } + if entry.ErrorType != "" { + t.Fatalf("ErrorType = %q, want empty", entry.ErrorType) + } logs, err := reader.GetLogs(context.Background(), LogQueryParams{Limit: 10}) if err != nil { @@ -343,6 +346,9 @@ func TestSQLiteReader_AllowsNullWorkflowVersionID(t *testing.T) { if logs.Entries[0].WorkflowVersionID != "" { t.Fatalf("list WorkflowVersionID = %q, want empty", logs.Entries[0].WorkflowVersionID) } + if logs.Entries[0].ErrorType != "" { + t.Fatalf("list ErrorType = %q, want empty", logs.Entries[0].ErrorType) + } } func TestSQLiteReader_GetLogsFiltersByUserPathSubtree(t *testing.T) { diff --git a/tools/seed-demo-data.sh b/tools/seed-demo-data.sh new file mode 100755 index 000000000..0af92f78f --- /dev/null +++ b/tools/seed-demo-data.sh @@ -0,0 +1,713 @@ +#!/usr/bin/env bash +set -euo pipefail + +db_path="${SQLITE_PATH:-data/gomodel.db}" +days="${DEMO_DAYS:-90}" +end_date="${DEMO_END_DATE:-}" +avg_requests="${DEMO_AVG_REQUESTS_PER_DAY:-850}" +max_requests="${DEMO_MAX_REQUESTS_PER_DAY:-1600}" +exact_cache_pct="${DEMO_EXACT_CACHE_PCT:-12}" +semantic_cache_pct="${DEMO_SEMANTIC_CACHE_PCT:-7}" +prompt_cache_pct="${DEMO_PROMPT_CACHE_PCT:-28}" +prefix="${DEMO_SEED_PREFIX:-demo-generated}" + +usage() { + cat <&2 + exit 2 + fi +} + +require_int DEMO_DAYS "$days" +require_int DEMO_AVG_REQUESTS_PER_DAY "$avg_requests" +require_int DEMO_MAX_REQUESTS_PER_DAY "$max_requests" +require_int DEMO_EXACT_CACHE_PCT "$exact_cache_pct" +require_int DEMO_SEMANTIC_CACHE_PCT "$semantic_cache_pct" +require_int DEMO_PROMPT_CACHE_PCT "$prompt_cache_pct" + +if (( days < 1 )); then + echo "DEMO_DAYS must be at least 1" >&2 + exit 2 +fi +if (( max_requests < avg_requests )); then + echo "DEMO_MAX_REQUESTS_PER_DAY must be >= DEMO_AVG_REQUESTS_PER_DAY" >&2 + exit 2 +fi +if (( exact_cache_pct + semantic_cache_pct > 65 )); then + echo "Exact + semantic cache percentages should stay realistic and <= 65" >&2 + exit 2 +fi +if (( prompt_cache_pct > 85 )); then + echo "DEMO_PROMPT_CACHE_PCT must be <= 85" >&2 + exit 2 +fi +if [[ -n "$end_date" && ! "$end_date" =~ ^[0-9]{4}-[0-9]{2}-[0-9]{2}$ ]]; then + echo "DEMO_END_DATE must use YYYY-MM-DD, got: $end_date" >&2 + exit 2 +fi +if [[ ! "$prefix" =~ ^[A-Za-z0-9_.-]+$ ]]; then + echo "DEMO_SEED_PREFIX may only contain letters, numbers, dot, underscore, and dash" >&2 + exit 2 +fi + +command -v sqlite3 >/dev/null 2>&1 || { + echo "sqlite3 is required" >&2 + exit 127 +} + +mkdir -p "$(dirname "$db_path")" + +sqlite3 "$db_path" "PRAGMA journal_mode = WAL;" >/dev/null + +sqlite3 "$db_path" <= p.min_bucket AND b.path_bucket < p.max_bucket + JOIN demo_templates t ON b.template_bucket >= t.min_bucket AND b.template_bucket < t.max_bucket +), +tokens AS ( + SELECT + *, + input_min + (token_noise % input_span) AS input_tokens, + output_min + ((token_noise / 97) % output_span) AS output_tokens + FROM chosen +), +cache_decisions AS ( + SELECT + *, + CASE + WHEN local_cache_eligible = 1 AND cache_bucket < (${exact_cache_pct} * 100) THEN 'exact' + WHEN local_cache_eligible = 1 AND cache_bucket < ((${exact_cache_pct} + ${semantic_cache_pct}) * 100) THEN 'semantic' + ELSE NULL + END AS cache_type, + CASE + WHEN prompt_cache_eligible = 1 + AND NOT (local_cache_eligible = 1 AND cache_bucket < ((${exact_cache_pct} + ${semantic_cache_pct}) * 100)) + AND prompt_bucket < (${prompt_cache_pct} * 100) + THEN 1 + ELSE 0 + END AS prompt_cache_hit + FROM tokens +), +prompt_parts AS ( + SELECT + *, + CASE + WHEN prompt_cache_hit = 1 THEN CAST(input_tokens * (35 + (prompt_bucket % 46)) / 100 AS INTEGER) + ELSE 0 + END AS prompt_cached_tokens, + CASE + WHEN prompt_cache_hit = 1 AND provider = 'anthropic' THEN CAST(input_tokens * (8 + (prompt_bucket % 13)) / 100 AS INTEGER) + ELSE 0 + END AS prompt_cache_write_tokens + FROM cache_decisions +) +SELECT + *, + input_tokens + output_tokens AS total_tokens, + strftime('%Y-%m-%dT%H:%M:%fZ', day || ' 00:00:00', '+' || second_of_day || ' seconds') AS timestamp, + '${prefix}-usage-' || day_idx || '-' || slot_idx AS usage_id, + '${prefix}-audit-' || day_idx || '-' || slot_idx AS audit_id, + '${prefix}-req-' || day_idx || '-' || slot_idx AS request_id, + '${prefix}-provider-' || day_idx || '-' || slot_idx AS provider_id +FROM prompt_parts; + +INSERT INTO usage ( + id, request_id, provider_id, timestamp, model, provider, provider_name, + endpoint, user_path, cache_type, input_tokens, output_tokens, total_tokens, + raw_data, input_cost, output_cost, total_cost, cost_source, costs_calculation_caveat +) +SELECT + usage_id, + request_id, + provider_id, + timestamp, + model, + provider, + provider_name, + endpoint, + user_path, + cache_type, + input_tokens, + output_tokens, + total_tokens, + CASE + WHEN cache_type = 'exact' THEN json_object( + 'demo_seed', 1, + 'cache_story', 'exact local response cache hit', + 'locally_cached_tokens', total_tokens + ) + WHEN cache_type = 'semantic' THEN json_object( + 'demo_seed', 1, + 'cache_story', 'semantic local response cache hit', + 'semantic_similarity', 0.88 + ((prompt_bucket % 12) / 100.0), + 'locally_cached_tokens', total_tokens + ) + WHEN prompt_cache_hit = 1 AND provider = 'anthropic' THEN json_object( + 'demo_seed', 1, + 'cache_story', 'provider prompt cache read/write', + 'cache_read_input_tokens', prompt_cached_tokens, + 'cache_creation_input_tokens', prompt_cache_write_tokens + ) + WHEN prompt_cache_hit = 1 AND provider = 'gemini' THEN json_object( + 'demo_seed', 1, + 'cache_story', 'provider prompt cache read', + 'cached_tokens', prompt_cached_tokens + ) + WHEN prompt_cache_hit = 1 THEN json_object( + 'demo_seed', 1, + 'cache_story', 'provider prompt cache read', + 'prompt_cached_tokens', prompt_cached_tokens + ) + ELSE json_object('demo_seed', 1, 'cache_story', 'uncached provider request') + END AS raw_data, + round((CASE + WHEN cache_type IS NOT NULL THEN 0 + WHEN prompt_cache_hit = 1 THEN ((input_tokens - prompt_cached_tokens) * input_price + prompt_cached_tokens * input_price * 0.25) / 1000000.0 + ELSE input_tokens * input_price / 1000000.0 + END), 8) AS input_cost, + round(CASE + WHEN cache_type IS NOT NULL THEN 0 + ELSE output_tokens * output_price / 1000000.0 + END, 8) AS output_cost, + round((CASE + WHEN cache_type IS NOT NULL THEN 0 + WHEN prompt_cache_hit = 1 THEN (((input_tokens - prompt_cached_tokens) * input_price + prompt_cached_tokens * input_price * 0.25) / 1000000.0) + (output_tokens * output_price / 1000000.0) + ELSE (input_tokens * input_price / 1000000.0) + (output_tokens * output_price / 1000000.0) + END), 8) AS total_cost, + CASE + WHEN cache_type IS NOT NULL THEN 'demo_local_cache' + WHEN prompt_cache_hit = 1 THEN 'demo_prompt_cache' + ELSE 'demo_model_pricing' + END AS cost_source, + '' +FROM demo_generated; + +INSERT INTO audit_logs ( + id, timestamp, duration_ns, requested_model, resolved_model, provider, provider_name, + alias_used, workflow_version_id, cache_type, status_code, request_id, auth_key_id, + auth_method, client_ip, method, path, user_path, stream, error_type, data +) +SELECT + audit_id, + timestamp, + CASE WHEN cache_type IS NOT NULL THEN 8000000 + (token_noise % 12000000) ELSE 90000000 + (token_noise % 260000000) END, + provider_name || '/' || model, + provider_name || '/' || model, + provider, + provider_name, + 0, + NULL, + cache_type, + CASE WHEN abs(token_noise / 131) % 1000 < 994 THEN 200 ELSE 500 END, + request_id, + NULL, + 'master_key', + '127.0.0.1', + 'POST', + endpoint, + user_path, + 0, + CASE WHEN abs(token_noise / 131) % 1000 < 994 THEN '' ELSE 'provider_error' END, + json_object( + 'demo_seed', 1, + 'workflow_features', json_object( + 'cache', json('true'), + 'audit', json('true'), + 'usage', json('true'), + 'budget', json('true'), + 'guardrails', json('false'), + 'fallback', json('true') + ), + 'cache_type', cache_type, + 'cache_story', CASE + WHEN cache_type = 'exact' THEN 'Exact response cache hit' + WHEN cache_type = 'semantic' THEN 'Semantic response cache hit' + WHEN prompt_cache_hit = 1 THEN 'Provider prompt cache telemetry' + ELSE 'Uncached provider request' + END, + 'request_body', json(CASE + WHEN label IN ('chat-openai', 'chat-groq', 'chat-gemini', 'chat-bailian') THEN json_object( + 'model', provider_name || '/' || model, + 'messages', json_array( + json_object('role', 'system', 'content', 'You are a concise assistant for internal demo traffic. Respect the user path and return actionable JSON when useful.'), + json_object('role', 'user', 'content', 'Summarize daily gateway usage for ' || user_path || ' and call out cache savings, error spikes, and next actions.'), + json_object('role', 'assistant', 'content', 'I will compare current traffic against the recent baseline and identify cost or latency anomalies.'), + json_object('role', 'user', 'content', 'Use request id ' || request_id || ' and include provider ' || provider_name || '.') + ), + 'temperature', round(0.15 + ((token_noise % 70) / 100.0), 2), + 'max_tokens', output_tokens, + 'stream', CASE WHEN slot_idx % 5 = 0 THEN json('true') ELSE json('false') END, + 'metadata', json_object( + 'demo', json('true'), + 'user_path', user_path, + 'cache_expected', CASE WHEN cache_type IS NOT NULL OR prompt_cache_hit = 1 THEN json('true') ELSE json('false') END + ) + ) + WHEN label = 'responses' THEN json_object( + 'model', provider_name || '/' || model, + 'input', json_array( + json_object('role', 'system', 'content', 'You are GoModel demo analysis worker.'), + json_object('role', 'user', 'content', 'Create a short incident-style report for ' || user_path || ' using token totals and cache telemetry.') + ), + 'instructions', 'Return sections named summary, observations, and recommendation.', + 'previous_response_id', CASE WHEN slot_idx > 0 AND slot_idx % 7 = 0 THEN '${prefix}-response-' || day_idx || '-' || (slot_idx - 1) ELSE NULL END, + 'max_output_tokens', output_tokens, + 'metadata', json_object('demo', json('true'), 'request_id', request_id) + ) + WHEN label = 'messages' THEN json_object( + 'model', provider_name || '/' || model, + 'system', 'You help the engineering and sales teams reason about AI gateway telemetry.', + 'messages', json_array( + json_object('role', 'user', 'content', json_array( + json_object('type', 'text', 'text', 'Draft a weekly update for ' || user_path || ' with token volume, model mix, cache behavior, and budget risk.') + )) + ), + 'max_tokens', output_tokens, + 'temperature', round(0.10 + ((token_noise % 55) / 100.0), 2) + ) + WHEN label = 'embeddings' THEN json_object( + 'model', provider_name || '/' || model, + 'input', json_array( + 'gateway usage dashboard prompt cache overview for ' || user_path, + 'semantic cache hit investigation request ' || request_id, + 'budget variance notes for provider ' || provider_name + ), + 'encoding_format', 'float' + ) + WHEN label = 'stt' THEN json_object( + '__audio__', json('true'), + 'content_type', 'audio/mpeg', + 'bytes', 48000 + (token_noise % 180000), + 'stored', json('false'), + 'meta', json_object( + 'model', provider_name || '/' || model, + 'language', CASE WHEN token_noise % 4 = 0 THEN 'pl' ELSE 'en' END, + 'prompt', 'Demo meeting note for ' || user_path, + 'temperature', round((token_noise % 20) / 100.0, 2) + ) + ) + WHEN label = 'tts' THEN json_object( + 'model', provider_name || '/' || model, + 'input', 'Read a concise dashboard summary for ' || user_path || ': tokens are trending up, prompt caching is active, and budgets remain under review.', + 'voice', CASE token_noise % 4 WHEN 0 THEN 'alloy' WHEN 1 THEN 'verse' WHEN 2 THEN 'coral' ELSE 'sage' END, + 'format', CASE WHEN token_noise % 3 = 0 THEN 'wav' ELSE 'mp3' END, + 'speed', round(0.90 + ((token_noise % 30) / 100.0), 2) + ) + ELSE json_object('model', provider_name || '/' || model, 'input', 'Generated demo request') + END), + 'response_body', json(CASE + WHEN abs(token_noise / 131) % 1000 >= 994 THEN json_object( + 'error', json_object( + 'message', 'Synthetic upstream provider error for demo audit inspection.', + 'type', 'provider_error', + 'code', 'demo_provider_error', + 'param', 'model' + ) + ) + WHEN label IN ('chat-openai', 'chat-groq', 'chat-gemini', 'chat-bailian') THEN json_object( + 'id', '${prefix}-response-' || day_idx || '-' || slot_idx, + 'object', 'chat.completion', + 'created', strftime('%s', timestamp), + 'model', provider_name || '/' || model, + 'choices', json_array(json_object( + 'index', 0, + 'finish_reason', 'stop', + 'message', json_object( + 'role', 'assistant', + 'content', 'Usage for ' || user_path || ' is healthy. Total tokens were ' || total_tokens || ', with cache mode ' || coalesce(cache_type, CASE WHEN prompt_cache_hit = 1 THEN 'prompt-cache' ELSE 'uncached' END) || '.' + ) + )), + 'usage', json_object( + 'prompt_tokens', input_tokens, + 'completion_tokens', output_tokens, + 'total_tokens', total_tokens, + 'prompt_tokens_details', json_object('cached_tokens', prompt_cached_tokens) + ) + ) + WHEN label = 'responses' THEN json_object( + 'id', '${prefix}-response-' || day_idx || '-' || slot_idx, + 'object', 'response', + 'status', 'completed', + 'model', provider_name || '/' || model, + 'output', json_array(json_object( + 'id', '${prefix}-msg-' || day_idx || '-' || slot_idx, + 'type', 'message', + 'role', 'assistant', + 'content', json_array(json_object( + 'type', 'output_text', + 'text', 'Summary: ' || user_path || ' generated ' || total_tokens || ' tokens. Observation: cache savings were ' || CASE WHEN cache_type IS NOT NULL OR prompt_cache_hit = 1 THEN 'visible' ELSE 'not present' END || '. Recommendation: keep monitoring budget drift.' + )) + )), + 'usage', json_object( + 'input_tokens', input_tokens, + 'output_tokens', output_tokens, + 'total_tokens', total_tokens, + 'input_tokens_details', json_object('cached_tokens', prompt_cached_tokens) + ) + ) + WHEN label = 'messages' THEN json_object( + 'id', '${prefix}-response-' || day_idx || '-' || slot_idx, + 'type', 'message', + 'role', 'assistant', + 'model', provider_name || '/' || model, + 'content', json_array(json_object( + 'type', 'text', + 'text', 'Weekly update for ' || user_path || ': model usage is balanced, semantic cache checks are active, and budget burn is within demo limits.' + )), + 'stop_reason', 'end_turn', + 'usage', json_object( + 'input_tokens', input_tokens, + 'output_tokens', output_tokens, + 'cache_read_input_tokens', prompt_cached_tokens, + 'cache_creation_input_tokens', prompt_cache_write_tokens + ) + ) + WHEN label = 'embeddings' THEN json_object( + 'object', 'list', + 'model', provider_name || '/' || model, + 'data', json_array( + json_object('object', 'embedding', 'index', 0, 'embedding', json_array(0.012, -0.034, 0.087, 0.003)), + json_object('object', 'embedding', 'index', 1, 'embedding', json_array(-0.021, 0.045, 0.016, -0.008)), + json_object('object', 'embedding', 'index', 2, 'embedding', json_array(0.005, 0.019, -0.042, 0.071)) + ), + 'usage', json_object('prompt_tokens', input_tokens, 'total_tokens', total_tokens) + ) + WHEN label = 'stt' THEN json_object( + 'text', 'Synthetic transcript for ' || user_path || ': review gateway usage, cache hit rates, and budget status.', + 'duration_seconds', round(18.0 + ((token_noise % 2400) / 100.0), 2), + 'language', CASE WHEN token_noise % 4 = 0 THEN 'pl' ELSE 'en' END, + 'segments', json_array( + json_object('id', 0, 'start', 0.00, 'end', 6.20, 'text', 'Review gateway usage and token volume.'), + json_object('id', 1, 'start', 6.20, 'end', 12.80, 'text', 'Check prompt caching and semantic cache hits.'), + json_object('id', 2, 'start', 12.80, 'end', 18.00, 'text', 'Confirm budgets for the user path.') + ) + ) + WHEN label = 'tts' THEN json_object( + '__audio__', json('true'), + 'content_type', CASE WHEN token_noise % 3 = 0 THEN 'audio/wav' ELSE 'audio/mpeg' END, + 'bytes', 32000 + (token_noise % 160000), + 'stored', json('false'), + 'meta', json_object( + 'model', provider_name || '/' || model, + 'voice', CASE token_noise % 4 WHEN 0 THEN 'alloy' WHEN 1 THEN 'verse' WHEN 2 THEN 'coral' ELSE 'sage' END, + 'format', CASE WHEN token_noise % 3 = 0 THEN 'wav' ELSE 'mp3' END + ) + ) + ELSE json_object('id', '${prefix}-response-' || day_idx || '-' || slot_idx, 'object', label) + END) + ) +FROM demo_generated; + +DROP TABLE IF EXISTS temp.demo_budget_paths; +CREATE TEMP TABLE demo_budget_paths(user_path TEXT, daily_amount REAL, weekly_amount REAL, monthly_amount REAL); +INSERT INTO demo_budget_paths VALUES + ('/', 420.00, 2500.00, 9500.00), + ('/agents/team1', 82.00, 510.00, 1900.00), + ('/agents/team1/research', 38.00, 225.00, 850.00), + ('/agents/team2', 74.00, 455.00, 1700.00), + ('/agents/team2/ops', 30.00, 180.00, 690.00), + ('/engineering', 160.00, 980.00, 3700.00), + ('/engineering/ai', 140.00, 850.00, 3200.00), + ('/engineering/ai/mike', 92.00, 570.00, 2100.00), + ('/engineering/ai/mike/evals', 54.00, 320.00, 1200.00), + ('/engineering/ai/bot', 105.00, 650.00, 2450.00), + ('/engineering/ai/bot/batch', 68.00, 420.00, 1600.00), + ('/sales', 95.00, 585.00, 2200.00), + ('/sales/john', 58.00, 355.00, 1350.00), + ('/sales/john/prospects', 28.00, 165.00, 620.00); + +WITH budget_rows AS ( + SELECT user_path, 86400 AS period_seconds, daily_amount AS amount FROM demo_budget_paths + UNION ALL + SELECT user_path, 604800 AS period_seconds, weekly_amount AS amount FROM demo_budget_paths + UNION ALL + SELECT user_path, 2592000 AS period_seconds, monthly_amount AS amount FROM demo_budget_paths +) +INSERT INTO budgets (user_path, period_seconds, amount, source, last_reset_at, created_at, updated_at) +SELECT + user_path, + period_seconds, + amount, + '${prefix}', + strftime('%s', date(CASE WHEN '${end_date}' = '' THEN 'now' ELSE '${end_date}' END)), + strftime('%s', 'now'), + strftime('%s', 'now') +FROM budget_rows +WHERE true +ON CONFLICT(user_path, period_seconds) DO UPDATE SET + amount = excluded.amount, + source = excluded.source, + last_reset_at = excluded.last_reset_at, + updated_at = excluded.updated_at; + +INSERT INTO budget_settings (key, value, updated_at) +VALUES + ('daily_reset_hour', '0', strftime('%s', 'now')), + ('daily_reset_minute', '0', strftime('%s', 'now')), + ('weekly_reset_weekday', '1', strftime('%s', 'now')), + ('weekly_reset_hour', '0', strftime('%s', 'now')), + ('weekly_reset_minute', '0', strftime('%s', 'now')), + ('monthly_reset_day', '1', strftime('%s', 'now')), + ('monthly_reset_hour', '0', strftime('%s', 'now')), + ('monthly_reset_minute', '0', strftime('%s', 'now')) +ON CONFLICT(key) DO UPDATE SET + value = excluded.value, + updated_at = excluded.updated_at; + +COMMIT; + +SELECT 'seed_prefix', '${prefix}'; +SELECT 'date_range', min(date(REPLACE(timestamp, 'T', ' '))), max(date(REPLACE(timestamp, 'T', ' '))) FROM usage WHERE id GLOB '${prefix}-*'; +SELECT 'usage_rows', count(*), coalesce(sum(total_tokens), 0) FROM usage WHERE id GLOB '${prefix}-*'; +SELECT 'audit_rows', count(*) FROM audit_logs WHERE id GLOB '${prefix}-*'; +SELECT 'budget_rows', count(*) FROM budgets WHERE source = '${prefix}'; +SELECT 'cache_mix', coalesce(cache_type, CASE + WHEN coalesce(json_extract(raw_data, '$.prompt_cached_tokens'), 0) > 0 + OR coalesce(json_extract(raw_data, '$.cached_tokens'), 0) > 0 + OR coalesce(json_extract(raw_data, '$.cache_read_input_tokens'), 0) > 0 + THEN 'prompt-cache' + ELSE 'uncached' +END), count(*) +FROM usage +WHERE id GLOB '${prefix}-*' +GROUP BY 2 +ORDER BY 2; +SELECT 'user_paths', count(DISTINCT user_path) FROM usage WHERE id GLOB '${prefix}-*'; +SELECT 'daily_requests_min_max', min(rows), max(rows), round(avg(rows), 1) +FROM ( + SELECT date(REPLACE(timestamp, 'T', ' ')) AS day, count(*) AS rows + FROM usage + WHERE id GLOB '${prefix}-*' + GROUP BY day +); +SQL + +cat <