summaryrefslogtreecommitdiff
path: root/internal/telemetry/usage_sink.go
diff options
context:
space:
mode:
Diffstat (limited to 'internal/telemetry/usage_sink.go')
-rw-r--r--internal/telemetry/usage_sink.go57
1 files changed, 57 insertions, 0 deletions
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)
+ }
+}