diff options
Diffstat (limited to 'internal/operations/operations.go')
| -rw-r--r-- | internal/operations/operations.go | 164 |
1 files changed, 112 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) |
