1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
|
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())
}
|