summaryrefslogtreecommitdiff
path: root/internal/telemetry/metrics.go
blob: 0631de437be504f4debdfcd37756bf2474cf9ae9 (plain)
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())
}