summaryrefslogtreecommitdiff
path: root/internal/telemetry/usage_sink.go
blob: aefe6f58b5dc8154755d98a912397b3395dfbdac (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
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)
	}
}