summaryrefslogtreecommitdiff
path: root/internal/telemetry
diff options
context:
space:
mode:
Diffstat (limited to '')
-rw-r--r--internal/telemetry/metrics.go44
-rw-r--r--internal/telemetry/usage_sink.go57
2 files changed, 101 insertions, 0 deletions
diff --git a/internal/telemetry/metrics.go b/internal/telemetry/metrics.go
new file mode 100644
index 0000000..4942d8d
--- /dev/null
+++ b/internal/telemetry/metrics.go
@@ -0,0 +1,44 @@
+package telemetry
+
+import (
+ "fmt"
+ "net/http"
+ "sync/atomic"
+)
+
+type Metrics struct {
+ requests atomic.Uint64
+ failed atomic.Uint64
+ inFlight atomic.Int64
+ attempts atomic.Uint64
+ droppedUsage atomic.Uint64
+}
+
+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) 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_usage_events_dropped_total counter\naigw_usage_events_dropped_total %d\n", m.droppedUsage.Load())
+}
diff --git a/internal/telemetry/usage_sink.go b/internal/telemetry/usage_sink.go
new file mode 100644
index 0000000..aefe6f5
--- /dev/null
+++ b/internal/telemetry/usage_sink.go
@@ -0,0 +1,57 @@
+package telemetry
+
+import (
+ "context"
+ "log/slog"
+ "sync"
+
+ "aigw/internal/domain"
+)
+
+type UsageSink interface {
+ Publish(domain.UsageEvent)
+}
+
+type AsyncUsageLogger struct {
+ logger *slog.Logger
+ metrics *Metrics
+ events chan domain.UsageEvent
+ done chan struct{}
+ once sync.Once
+}
+
+func NewAsyncUsageLogger(logger *slog.Logger, metrics *Metrics, buffer int) *AsyncUsageLogger {
+ sink := &AsyncUsageLogger{
+ logger: logger,
+ metrics: metrics,
+ events: make(chan domain.UsageEvent, buffer),
+ done: make(chan struct{}),
+ }
+ go sink.run()
+ return sink
+}
+
+func (s *AsyncUsageLogger) Publish(event domain.UsageEvent) {
+ select {
+ case s.events <- event:
+ default:
+ s.metrics.UsageDropped()
+ }
+}
+
+func (s *AsyncUsageLogger) Close(ctx context.Context) error {
+ s.once.Do(func() { close(s.events) })
+ select {
+ case <-s.done:
+ return nil
+ case <-ctx.Done():
+ return ctx.Err()
+ }
+}
+
+func (s *AsyncUsageLogger) run() {
+ defer close(s.done)
+ for event := range s.events {
+ s.logger.Info("usage_event", "event", event)
+ }
+}