package telemetry import ( "fmt" "net/http" "sync/atomic" "time" "aigw/internal/billing" ) type Metrics struct { requests atomic.Uint64 failed atomic.Uint64 inFlight atomic.Int64 attempts atomic.Uint64 upstreamTTFTCount atomic.Uint64 upstreamTTFTMSSum atomic.Uint64 providerProbes atomic.Uint64 providerProbeFailed atomic.Uint64 providerSharedPublished atomic.Uint64 providerSharedImported atomic.Uint64 providerSharedDropped atomic.Uint64 providerSharedFailures atomic.Uint64 providerSharedConnected atomic.Int64 droppedUsage atomic.Uint64 settlementBacklog atomic.Int64 settlementSpool atomic.Int64 stripeRefundBacklog atomic.Int64 stripeUncollected atomic.Int64 stripeMismatches atomic.Int64 stripeWebhooks atomic.Int64 unmeteredSuccesses atomic.Int64 ready atomic.Int64 } func (m *Metrics) SetSettlementQueue(backlog int64, spool int) { m.settlementBacklog.Store(backlog) m.settlementSpool.Store(int64(spool)) } func (m *Metrics) SetReady(ready bool) { if ready { m.ready.Store(1) } else { m.ready.Store(0) } } func (m *Metrics) SetStripeOperations(status billing.OperationalStatus) { m.stripeRefundBacklog.Store(status.RefundBacklog) m.stripeUncollected.Store(status.UncollectedMicros) m.stripeMismatches.Store(status.ReconciliationMismatches) m.stripeWebhooks.Store(status.UnprocessedWebhooks) m.unmeteredSuccesses.Store(status.UnmeteredSuccesses) } func (m *Metrics) RequestStarted() { m.requests.Add(1) m.inFlight.Add(1) } func (m *Metrics) RequestFinished(success bool) { m.inFlight.Add(-1) if !success { m.failed.Add(1) } } func (m *Metrics) UpstreamAttempt() { m.attempts.Add(1) } func (m *Metrics) UpstreamTTFT(latency time.Duration) { if latency <= 0 { return } milliseconds := latency.Milliseconds() if milliseconds < 1 { milliseconds = 1 } m.upstreamTTFTCount.Add(1) m.upstreamTTFTMSSum.Add(uint64(milliseconds)) } func (m *Metrics) ProviderProbe(success bool) { m.providerProbes.Add(1) if !success { m.providerProbeFailed.Add(1) } } func (m *Metrics) ProviderHealthSharedPublished() { m.providerSharedPublished.Add(1) } func (m *Metrics) ProviderHealthSharedImported() { m.providerSharedImported.Add(1) } func (m *Metrics) ProviderHealthSharedDropped() { m.providerSharedDropped.Add(1) } func (m *Metrics) ProviderHealthSharedRedisFailure() { m.providerSharedFailures.Add(1) } func (m *Metrics) ProviderHealthSharedConnected(connected bool) { if connected { m.providerSharedConnected.Store(1) return } m.providerSharedConnected.Store(0) } func (m *Metrics) UsageDropped() { m.droppedUsage.Add(1) } func (m *Metrics) ServeHTTP(w http.ResponseWriter, _ *http.Request) { w.Header().Set("Content-Type", "text/plain; version=0.0.4") fmt.Fprintf(w, "# TYPE aigw_requests_total counter\naigw_requests_total %d\n", m.requests.Load()) fmt.Fprintf(w, "# TYPE aigw_requests_failed_total counter\naigw_requests_failed_total %d\n", m.failed.Load()) fmt.Fprintf(w, "# TYPE aigw_requests_in_flight gauge\naigw_requests_in_flight %d\n", m.inFlight.Load()) fmt.Fprintf(w, "# TYPE aigw_upstream_attempts_total counter\naigw_upstream_attempts_total %d\n", m.attempts.Load()) fmt.Fprintf(w, "# TYPE aigw_upstream_ttft_ms_count counter\naigw_upstream_ttft_ms_count %d\n", m.upstreamTTFTCount.Load()) fmt.Fprintf(w, "# TYPE aigw_upstream_ttft_ms_sum counter\naigw_upstream_ttft_ms_sum %d\n", m.upstreamTTFTMSSum.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_probes_total counter\naigw_provider_probes_total %d\n", m.providerProbes.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_probe_failures_total counter\naigw_provider_probe_failures_total %d\n", m.providerProbeFailed.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_health_shared_published_total counter\naigw_provider_health_shared_published_total %d\n", m.providerSharedPublished.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_health_shared_imported_total counter\naigw_provider_health_shared_imported_total %d\n", m.providerSharedImported.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_health_shared_dropped_total counter\naigw_provider_health_shared_dropped_total %d\n", m.providerSharedDropped.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_health_shared_redis_failures_total counter\naigw_provider_health_shared_redis_failures_total %d\n", m.providerSharedFailures.Load()) fmt.Fprintf(w, "# TYPE aigw_provider_health_shared_connected gauge\naigw_provider_health_shared_connected %d\n", m.providerSharedConnected.Load()) fmt.Fprintf(w, "# TYPE aigw_usage_events_dropped_total counter\naigw_usage_events_dropped_total %d\n", m.droppedUsage.Load()) fmt.Fprintf(w, "# TYPE aigw_billing_settlement_backlog gauge\naigw_billing_settlement_backlog %d\n", m.settlementBacklog.Load()) fmt.Fprintf(w, "# TYPE aigw_billing_settlement_spool_records gauge\naigw_billing_settlement_spool_records %d\n", m.settlementSpool.Load()) fmt.Fprintf(w, "# TYPE aigw_stripe_refund_backlog gauge\naigw_stripe_refund_backlog %d\n", m.stripeRefundBacklog.Load()) fmt.Fprintf(w, "# TYPE aigw_billing_uncollected_micros gauge\naigw_billing_uncollected_micros %d\n", m.stripeUncollected.Load()) fmt.Fprintf(w, "# TYPE aigw_stripe_reconciliation_mismatches gauge\naigw_stripe_reconciliation_mismatches %d\n", m.stripeMismatches.Load()) fmt.Fprintf(w, "# TYPE aigw_stripe_webhook_backlog gauge\naigw_stripe_webhook_backlog %d\n", m.stripeWebhooks.Load()) fmt.Fprintf(w, "# TYPE aigw_billing_unmetered_successes gauge\naigw_billing_unmetered_successes %d\n", m.unmeteredSuccesses.Load()) fmt.Fprintf(w, "# TYPE aigw_ready gauge\naigw_ready %d\n", m.ready.Load()) }