summaryrefslogtreecommitdiff
path: root/internal/operations/operations.go
diff options
context:
space:
mode:
Diffstat (limited to 'internal/operations/operations.go')
-rw-r--r--internal/operations/operations.go164
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)