diff options
Diffstat (limited to '')
| -rw-r--r-- | internal/telemetry/usage_sink.go | 57 |
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) + } +} |
