diff options
| author | Chia <Chia@93.nz> | 2026-08-06 09:29:41 +1200 |
|---|---|---|
| committer | Chia <Chia@93.nz> | 2026-08-06 09:32:46 +1200 |
| commit | 41e322c53d7b4b796eb377d0df9c29ecd10ba431 (patch) | |
| tree | c730526150e55e39b822d5197e4a20318ecaa449 /internal/operations | |
| parent | eadb2ffe85c43cf6fc741c9823cd28eedb4a844c (diff) | |
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
Diffstat (limited to '')
| -rw-r--r-- | internal/operations/operations.go | 164 | ||||
| -rw-r--r-- | internal/operations/operations_test.go | 89 |
2 files changed, 201 insertions, 52 deletions
diff --git a/internal/operations/operations.go b/internal/operations/operations.go index e566488..502c7f8 100644 --- a/internal/operations/operations.go +++ b/internal/operations/operations.go @@ -12,11 +12,29 @@ import ( ) type Handler struct { - Store *controlplane.Store - Manager *controlplane.Manager - Billing *billing.Service - Metrics *telemetry.Metrics - MaxSnapshotAge time.Duration + Store readinessStore + Manager *controlplane.Manager + Billing readinessBilling + Metrics *telemetry.Metrics + MaxSnapshotAge time.Duration + ReadinessTimeout time.Duration +} + +type readinessStore interface { + Ping(context.Context) error + MailQueueStatus(context.Context) (controlplane.MailQueueStatus, error) +} + +type readinessBilling interface { + Ping(context.Context) error + SettlementQueueStatus(context.Context) (billing.SettlementQueueStatus, error) + OperationalStatus(context.Context) (billing.OperationalStatus, error) +} + +type readinessResult struct { + name string + value map[string]any + ready bool } func (h Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { @@ -32,17 +50,47 @@ func (h Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { http.NotFound(w, r) return } - ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second) + + timeout := h.ReadinessTimeout + if timeout <= 0 { + timeout = 1500 * time.Millisecond + } + // Readiness probes are independently bounded below. Some container health + // clients close their request side aggressively after sending the GET; do + // not let that client lifecycle make every backend look unavailable. + ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() checks := map[string]any{} ready := true + results := make(chan readinessResult, 5) + expected := map[string]struct{}{} + launch := func(name string, check func(context.Context) readinessResult) { + expected[name] = struct{}{} + go func() { + result := check(ctx) + result.name = name + results <- result + }() + } + if h.Store != nil { - if err := h.Store.Ping(ctx); err != nil { - checks["postgres"] = map[string]any{"status": "failed", "error": err.Error()} - ready = false - } else { - checks["postgres"] = map[string]any{"status": "ok"} - } + launch("postgres", func(ctx context.Context) readinessResult { + if err := h.Store.Ping(ctx); err != nil { + return failedResult(err) + } + return okResult() + }) + launch("mail_queue", func(ctx context.Context) readinessResult { + mail, err := h.Store.MailQueueStatus(ctx) + if err != nil { + return failedResult(err) + } + mailOK := mail.Failed == 0 + if mail.OldestPending != nil && time.Since(*mail.OldestPending) > 10*time.Minute { + mailOK = false + } + return readinessResult{ready: mailOK, value: map[string]any{"status": status(mailOK), "details": mail}} + }) } if h.Manager != nil { healthyAt := h.Manager.LastHealthyAt() @@ -65,61 +113,60 @@ func (h Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { checks["redis"] = map[string]any{"status": redisStatus, "required": false} } if h.Billing != nil { - if err := h.Billing.Ping(ctx); err != nil { - checks["billing_postgres"] = map[string]any{"status": "failed", "error": err.Error()} - ready = false - } else { - checks["billing_postgres"] = map[string]any{"status": "ok"} - } - queue, err := h.Billing.SettlementQueueStatus(ctx) - if err != nil { - checks["settlement_queue"] = map[string]any{"status": "failed", "error": err.Error()} - ready = false - } else { + launch("billing_postgres", func(ctx context.Context) readinessResult { + if err := h.Billing.Ping(ctx); err != nil { + return failedResult(err) + } + return okResult() + }) + launch("settlement_queue", func(ctx context.Context) readinessResult { + queue, err := h.Billing.SettlementQueueStatus(ctx) + if err != nil { + return failedResult(err) + } backlog := queue.AwaitingEvent + queue.Pending + queue.Processing + queue.Retrying queueOK := queue.SpoolRecords == 0 if queue.OldestPending != nil && time.Since(*queue.OldestPending) > 15*time.Minute { queueOK = false } - if !queueOK { - ready = false - } - checks["settlement_queue"] = map[string]any{"status": status(queueOK), "backlog": backlog, "spool_records": queue.SpoolRecords, "oldest_pending": queue.OldestPending} if h.Metrics != nil { h.Metrics.SetSettlementQueue(backlog, queue.SpoolRecords) } + return readinessResult{ready: queueOK, value: map[string]any{"status": status(queueOK), "backlog": backlog, "spool_records": queue.SpoolRecords, "oldest_pending": queue.OldestPending}} + }) + launch("billing_operations", func(ctx context.Context) readinessResult { billingHealth, err := h.Billing.OperationalStatus(ctx) if err != nil { - checks["billing_operations"] = map[string]any{"status": "failed", "error": err.Error()} - ready = false - } else { - billingOK := billingHealth.Ready(time.Now().UTC()) - checks["billing_operations"] = map[string]any{"status": status(billingOK), "details": billingHealth} - if !billingOK { - ready = false - } - if h.Metrics != nil { - h.Metrics.SetStripeOperations(billingHealth) - } + return failedResult(err) } - } - if h.Store != nil { - mail, err := h.Store.MailQueueStatus(ctx) - if err != nil { - checks["mail_queue"] = map[string]any{"status": "failed", "error": err.Error()} + billingOK := billingHealth.Ready(time.Now().UTC()) + if h.Metrics != nil { + h.Metrics.SetStripeOperations(billingHealth) + } + return readinessResult{ready: billingOK, value: map[string]any{"status": status(billingOK), "details": billingHealth}} + }) + } + + for len(expected) > 0 { + select { + case result := <-results: + if _, ok := expected[result.name]; !ok { + continue + } + delete(expected, result.name) + checks[result.name] = result.value + if !result.ready { ready = false - } else { - mailOK := mail.Failed == 0 - if mail.OldestPending != nil && time.Since(*mail.OldestPending) > 10*time.Minute { - mailOK = false - } - checks["mail_queue"] = map[string]any{"status": status(mailOK), "details": mail} - if !mailOK { - ready = false - } } + case <-ctx.Done(): + for name := range expected { + checks[name] = map[string]any{"status": "failed", "error": "check timed out"} + } + ready = false + expected = map[string]struct{}{} } } + if h.Metrics != nil { h.Metrics.SetReady(ready) } @@ -130,12 +177,25 @@ func (h Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { write(w, code, map[string]any{"status": status(ready), "checks": checks}) } +func okResult() readinessResult { + return readinessResult{ready: true, value: map[string]any{"status": "ok"}} +} + +func failedResult(err error) readinessResult { + message := err.Error() + if err == context.Canceled || err == context.DeadlineExceeded { + message = "check timed out" + } + return readinessResult{value: map[string]any{"status": "failed", "error": message}} +} + func status(ok bool) string { if ok { return "ok" } return "failed" } + func write(w http.ResponseWriter, code int, value any) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(code) diff --git a/internal/operations/operations_test.go b/internal/operations/operations_test.go new file mode 100644 index 0000000..6936b1a --- /dev/null +++ b/internal/operations/operations_test.go @@ -0,0 +1,89 @@ +package operations + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "aigw/internal/billing" + "aigw/internal/controlplane" +) + +type healthyStore struct{} + +func (healthyStore) Ping(context.Context) error { return nil } +func (healthyStore) MailQueueStatus(context.Context) (controlplane.MailQueueStatus, error) { + return controlplane.MailQueueStatus{}, nil +} + +type slowBilling struct{} + +func (slowBilling) Ping(context.Context) error { return nil } +func (slowBilling) SettlementQueueStatus(context.Context) (billing.SettlementQueueStatus, error) { + return billing.SettlementQueueStatus{}, nil +} +func (slowBilling) OperationalStatus(ctx context.Context) (billing.OperationalStatus, error) { + <-ctx.Done() + return billing.OperationalStatus{}, ctx.Err() +} + +func TestReadinessTimeoutDoesNotMislabelCompletedChecks(t *testing.T) { + handler := Handler{Store: healthyStore{}, Billing: slowBilling{}, ReadinessTimeout: 20 * time.Millisecond} + request := httptest.NewRequest(http.MethodGet, "/readyz", nil) + recorder := httptest.NewRecorder() + handler.ServeHTTP(recorder, request) + if recorder.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d, want 503", recorder.Code) + } + var payload struct { + Checks map[string]struct { + Status string `json:"status"` + Error string `json:"error"` + } `json:"checks"` + } + if err := json.NewDecoder(recorder.Body).Decode(&payload); err != nil { + t.Fatal(err) + } + for _, name := range []string{"postgres", "mail_queue", "billing_postgres", "settlement_queue"} { + if payload.Checks[name].Status != "ok" { + t.Fatalf("%s status = %q, want ok; payload=%s", name, payload.Checks[name].Status, recorder.Body.String()) + } + } + if payload.Checks["billing_operations"].Status != "failed" || payload.Checks["billing_operations"].Error != "check timed out" { + t.Fatalf("billing operations = %+v", payload.Checks["billing_operations"]) + } +} + +func TestMailQueueIsCheckedWithoutBilling(t *testing.T) { + handler := Handler{Store: healthyStore{}, ReadinessTimeout: time.Second} + recorder := httptest.NewRecorder() + handler.ServeHTTP(recorder, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + if recorder.Code != http.StatusOK { + t.Fatalf("status = %d, body=%s", recorder.Code, recorder.Body.String()) + } + var payload struct { + Checks map[string]any `json:"checks"` + } + if err := json.NewDecoder(recorder.Body).Decode(&payload); err != nil { + t.Fatal(err) + } + if _, ok := payload.Checks["mail_queue"]; !ok { + t.Fatal("mail_queue check missing when billing is disabled") + } +} + +func TestReadinessChecksDoNotInheritCanceledClientContext(t *testing.T) { + handler := Handler{Store: healthyStore{}, ReadinessTimeout: time.Second} + request := httptest.NewRequest(http.MethodGet, "/readyz", nil) + ctx, cancel := context.WithCancel(request.Context()) + cancel() + request = request.WithContext(ctx) + recorder := httptest.NewRecorder() + handler.ServeHTTP(recorder, request) + if recorder.Code != http.StatusOK { + t.Fatalf("status = %d, body=%s", recorder.Code, recorder.Body.String()) + } +} |
