diff --git a/internal/metrics/relay_metrics_collector.go b/internal/metrics/relay_metrics_collector.go index 646f68b9..b18e10ba 100644 --- a/internal/metrics/relay_metrics_collector.go +++ b/internal/metrics/relay_metrics_collector.go @@ -63,19 +63,28 @@ type RelayMetricsCollector struct { pollingCounts map[connectionsKeyType]pollingCounts mu sync.Mutex closer chan struct{} + now func() time.Time } func newRelayMetricsCollector(relayID, envName string, publisher events.EventPublisher, flushInterval time.Duration, logger *slog.Logger) *RelayMetricsCollector { + return newRelayMetricsCollectorWithTimeSource(relayID, envName, publisher, flushInterval, logger, time.Now) +} + +// newRelayMetricsCollectorWithTimeSource allows tests to control the timestamps +// used for interval boundaries. The now function must never return a value earlier +// than one it previously returned. +func newRelayMetricsCollectorWithTimeSource(relayID, envName string, publisher events.EventPublisher, flushInterval time.Duration, logger *slog.Logger, now func() time.Time) *RelayMetricsCollector { c := &RelayMetricsCollector{ relayID: relayID, envName: envName, publisher: publisher, logger: logger, closer: make(chan struct{}), - intervalStartTime: time.Now(), + intervalStartTime: now(), pollingDataIsDirty: false, currentConnections: make(map[connectionsKeyType]int64), pollingCounts: make(map[connectionsKeyType]pollingCounts), + now: now, } flushTicker := time.NewTicker(flushInterval) @@ -135,7 +144,7 @@ func (c *RelayMetricsCollector) hasMetricDataToReport() bool { func (c *RelayMetricsCollector) flush() { c.mu.Lock() startTime := c.intervalStartTime - stopTime := time.Now() + stopTime := c.now() c.intervalStartTime = stopTime if !c.hasMetricDataToReport() { diff --git a/internal/metrics/relay_metrics_collector_test.go b/internal/metrics/relay_metrics_collector_test.go index 1805a3d7..c230666e 100644 --- a/internal/metrics/relay_metrics_collector_test.go +++ b/internal/metrics/relay_metrics_collector_test.go @@ -122,15 +122,24 @@ func TestRelayMetricsCollector(t *testing.T) { t.Run("the event start time still shifts when events are not sent", func(t *testing.T) { publisher := newTestEventsPublisher() - withCollector(publisher, func(c *RelayMetricsCollector, relayID string) { - time.Sleep(time.Millisecond * 10) - startTime := ldtime.UnixMillisNow() - time.Sleep(time.Millisecond * 1) - c.RecordConnectionChange(platformValue, userAgentValue, "", 1) + clk := newFakeClock() + c := newRelayMetricsCollectorWithTimeSource(uuid.New(), "envName", publisher, time.Hour, slog.Default(), clk.now) + defer c.close() - c.flush() - metricsEvent := publisher.expectMetricsEvent(t, time.Second) - assert.True(t, metricsEvent.StartDate >= startTime) - }) + // A flush with no data publishes nothing, but should still move the + // interval start forward. + clk.advance(time.Millisecond * 10) + c.flush() + publisher.expectNoMetricsEvent(t, time.Millisecond*50) + shiftedStart := clk.now() + + clk.advance(time.Millisecond) + c.RecordConnectionChange(platformValue, userAgentValue, "", 1) + clk.advance(time.Millisecond) + c.flush() + + metricsEvent := publisher.expectMetricsEvent(t, time.Second) + assert.Equal(t, ldtime.UnixMillisFromTime(shiftedStart), metricsEvent.StartDate) + assert.Equal(t, ldtime.UnixMillisFromTime(clk.now()), metricsEvent.EndDate) }) } diff --git a/internal/metrics/test_utils_test.go b/internal/metrics/test_utils_test.go index 3c6df636..155f3664 100644 --- a/internal/metrics/test_utils_test.go +++ b/internal/metrics/test_utils_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "log/slog" + "sync" "testing" "time" @@ -75,6 +76,30 @@ func testWithOTel(t *testing.T, action func(testWithOTelParams)) { }) } +// fakeClock is a controllable time source for tests that measure durations. +// Advancing it is the test's replacement for sleeping real time, which is +// unreliable under CI scheduling jitter. +type fakeClock struct { + mu sync.Mutex + t time.Time +} + +func newFakeClock() *fakeClock { + return &fakeClock{t: time.Date(2026, time.January, 1, 12, 0, 0, 0, time.UTC)} +} + +func (c *fakeClock) now() time.Time { + c.mu.Lock() + defer c.mu.Unlock() + return c.t +} + +func (c *fakeClock) advance(d time.Duration) { + c.mu.Lock() + defer c.mu.Unlock() + c.t = c.t.Add(d) +} + type testEventsPublisher struct { events chan json.RawMessage } diff --git a/internal/metrics/usage.go b/internal/metrics/usage.go index d0c313d3..971b8e19 100644 --- a/internal/metrics/usage.go +++ b/internal/metrics/usage.go @@ -48,11 +48,17 @@ type usageActivityMessage struct { platformCategory string instanceID string tagsHeader string + + // The time at which the activity was recorded. Stamped when the message is + // handed to an environmentMetricUsage, before it crosses into the + // processing goroutine, so that tests can substitute a controllable clock + // and observe deterministic durations. + timestamp time.Time } type ( - usageActivityFlush struct{} - usageActivityShutdown struct{} + usageActivityFlush struct{ timestamp time.Time } + usageActivityShutdown struct{ timestamp time.Time } ) // metricUsage is used to track usage information for a single @@ -137,6 +143,7 @@ type environmentMetricUsage struct { publisher events.EventPublisher flushInterval time.Duration usageCh chan interface{} + now func() time.Time // All of this data is expected to only be accessed from within a single go // routine @@ -144,11 +151,18 @@ type environmentMetricUsage struct { } func NewEnvironmentMetricUsage(relayID string, publisher events.EventPublisher, flushInterval time.Duration) *environmentMetricUsage { + return newEnvironmentMetricUsage(relayID, publisher, flushInterval, time.Now) +} + +// newEnvironmentMetricUsage allows tests to control the time source. The now +// function must never return a value earlier than one it previously returned. +func newEnvironmentMetricUsage(relayID string, publisher events.EventPublisher, flushInterval time.Duration, now func() time.Time) *environmentMetricUsage { e := &environmentMetricUsage{ relayID: relayID, publisher: publisher, flushInterval: flushInterval, usageCh: make(chan interface{}), + now: now, usages: make(map[usageKeyType]*metricUsage), } @@ -159,6 +173,7 @@ func NewEnvironmentMetricUsage(relayID string, publisher events.EventPublisher, } func (e *environmentMetricUsage) usageActivityMessage(usage *usageActivityMessage) { + usage.timestamp = e.now() e.usageCh <- usage } @@ -169,20 +184,20 @@ func (e *environmentMetricUsage) run() { for { select { case <-ticker.C: - e.flushInternal() + e.flushInternal(e.now()) case usage, ok := <-e.usageCh: if !ok { return } - now := time.Now() switch u := usage.(type) { case *usageActivityShutdown: - e.flushInternal() + e.flushInternal(u.timestamp) return case *usageActivityFlush: - e.flushInternal() + e.flushInternal(u.timestamp) case *usageActivityMessage: + now := u.timestamp key := usageKeyType{userAgent: u.userAgent, platformCategory: u.platformCategory, instanceID: u.instanceID, tagsHeader: u.tagsHeader} if e.publisher == nil { continue @@ -238,15 +253,15 @@ func (e *environmentMetricUsage) run() { } func (e *environmentMetricUsage) flush() { //nolint:unused // used only in tests - e.usageCh <- &usageActivityFlush{} + e.usageCh <- &usageActivityFlush{timestamp: e.now()} } func (e *environmentMetricUsage) close() { - e.usageCh <- &usageActivityShutdown{} + e.usageCh <- &usageActivityShutdown{timestamp: e.now()} close(e.usageCh) } -func (e *environmentMetricUsage) flushInternal() { +func (e *environmentMetricUsage) flushInternal(now time.Time) { if e.publisher == nil { return } @@ -255,8 +270,6 @@ func (e *environmentMetricUsage) flushInternal() { return } - now := time.Now() - for key, usage := range e.usages { // Refer back to the metricUsage comment for an explanation on this // calculation. diff --git a/internal/metrics/usage_test.go b/internal/metrics/usage_test.go index 7513e920..cb892029 100644 --- a/internal/metrics/usage_test.go +++ b/internal/metrics/usage_test.go @@ -28,7 +28,8 @@ func TestIndividualCountMessage(t *testing.T) { func TestCountsWithDelay(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindCount, @@ -37,7 +38,7 @@ func TestCountsWithDelay(t *testing.T) { instanceID: "instanceID", }) - time.Sleep(1 * time.Millisecond) + clk.advance(1 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindCount, @@ -142,7 +143,8 @@ func TestFlushingResetsCounts(t *testing.T) { func TestCapturesSimpleStreamDuration(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamConnected, @@ -151,7 +153,7 @@ func TestCapturesSimpleStreamDuration(t *testing.T) { instanceID: "instanceID", }) - time.Sleep(10 * time.Millisecond) + clk.advance(10 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamDisconnected, @@ -166,12 +168,13 @@ func TestCapturesSimpleStreamDuration(t *testing.T) { assert.Equal(t, "userAgent", event.UserAgent) assert.Equal(t, "platform", event.PlatformCategory) assert.Equal(t, "instanceID", event.InstanceID) - assert.InDeltaf(t, 10, event.TotalStreamMs, 5, "stream time should be approximately 10ms") + assert.Equal(t, int64(10), event.TotalStreamMs) } func TestStreamWithoutDisconnect(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) // Connect to stream but don't disconnect env.usageActivityMessage(&usageActivityMessage{ @@ -181,7 +184,7 @@ func TestStreamWithoutDisconnect(t *testing.T) { instanceID: "instanceID", }) - time.Sleep(10 * time.Millisecond) + clk.advance(10 * time.Millisecond) // Flush without disconnecting env.flush() @@ -189,16 +192,16 @@ func TestStreamWithoutDisconnect(t *testing.T) { event := publisher.expectUsageEvent(t, time.Second) assert.Equal(t, "userAgent", event.UserAgent) assert.NotEqual(t, event.FirstActive, event.LastActive) // Ensure timestamps differ - assert.InDeltaf(t, 10, event.TotalStreamMs, 5, "stream time should be approximately 10ms") + assert.Equal(t, int64(10), event.TotalStreamMs) - // Stream is still connected, so wait and try to force another event. - time.Sleep(30 * time.Millisecond) + // Stream is still connected, so let more time pass and force another event. + clk.advance(30 * time.Millisecond) env.flush() event = publisher.expectUsageEvent(t, time.Second) assert.Equal(t, "userAgent", event.UserAgent) assert.NotEqual(t, event.FirstActive, event.LastActive) // Ensure timestamps differ - assert.InDeltaf(t, 30, event.TotalStreamMs, 5, "stream time should be approximately 30ms") + assert.Equal(t, int64(30), event.TotalStreamMs) } func TestStreamDisconnectWithoutConnect(t *testing.T) { @@ -224,7 +227,8 @@ func TestStreamDisconnectWithoutConnect(t *testing.T) { func TestStreamDurationIsNotAffectedByActivityCounts(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindCount, @@ -232,7 +236,7 @@ func TestStreamDurationIsNotAffectedByActivityCounts(t *testing.T) { platformCategory: "platform", instanceID: "instanceID", }) - time.Sleep(20 * time.Millisecond) + clk.advance(20 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamConnected, @@ -240,7 +244,7 @@ func TestStreamDurationIsNotAffectedByActivityCounts(t *testing.T) { platformCategory: "platform", instanceID: "instanceID", }) - time.Sleep(10 * time.Millisecond) + clk.advance(10 * time.Millisecond) env.flush() @@ -248,21 +252,22 @@ func TestStreamDurationIsNotAffectedByActivityCounts(t *testing.T) { assert.Equal(t, "userAgent", event.UserAgent) assert.Equal(t, "platform", event.PlatformCategory) assert.Equal(t, "instanceID", event.InstanceID) - assert.InDeltaf(t, 10, event.TotalStreamMs, 5, "stream time should be approximately 10ms") + assert.Equal(t, int64(10), event.TotalStreamMs) } func TestNonoverlappingStreams(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) - // First stream session: ~10ms + // First stream session: 10ms env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamConnected, userAgent: "userAgent", platformCategory: "platform", instanceID: "instanceID", }) - time.Sleep(10 * time.Millisecond) + clk.advance(10 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamDisconnected, userAgent: "userAgent", @@ -270,14 +275,14 @@ func TestNonoverlappingStreams(t *testing.T) { instanceID: "instanceID", }) - // Second stream session: ~30ms + // Second stream session: 30ms env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamConnected, userAgent: "userAgent", platformCategory: "platform", instanceID: "instanceID", }) - time.Sleep(30 * time.Millisecond) + clk.advance(30 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamDisconnected, userAgent: "userAgent", @@ -289,12 +294,13 @@ func TestNonoverlappingStreams(t *testing.T) { event := publisher.expectUsageEvent(t, time.Second) assert.Equal(t, "userAgent", event.UserAgent) - assert.InDeltaf(t, 40, event.TotalStreamMs, 5, "stream time should be approximately 40ms") + assert.Equal(t, int64(40), event.TotalStreamMs) } func TestOverlappingStreams(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamConnected, @@ -303,7 +309,7 @@ func TestOverlappingStreams(t *testing.T) { instanceID: "instanceID", }) - time.Sleep(100 * time.Millisecond) + clk.advance(100 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamConnected, @@ -312,7 +318,7 @@ func TestOverlappingStreams(t *testing.T) { instanceID: "instanceID", }) - time.Sleep(100 * time.Millisecond) + clk.advance(100 * time.Millisecond) env.usageActivityMessage(&usageActivityMessage{ kind: UsageActivityKindStreamDisconnected, @@ -326,12 +332,14 @@ func TestOverlappingStreams(t *testing.T) { event := publisher.expectUsageEvent(t, time.Second) assert.Equal(t, "userAgent", event.UserAgent) - assert.InDeltaf(t, 300, event.TotalStreamMs, 50, "stream time should be approximately 300ms") + // First connection streams for the full 200ms; the second overlaps for 100ms. + assert.Equal(t, int64(300), event.TotalStreamMs) } func TestMultipleUserStreams(t *testing.T) { publisher := newTestEventsPublisher() - env := NewEnvironmentMetricUsage("relayID", publisher, 1*time.Hour) + clk := newFakeClock() + env := newEnvironmentMetricUsage("relayID", publisher, 1*time.Hour, clk.now) // First user connects env.usageActivityMessage(&usageActivityMessage{ @@ -341,7 +349,7 @@ func TestMultipleUserStreams(t *testing.T) { instanceID: "instanceID1", }) - time.Sleep(50 * time.Millisecond) + clk.advance(50 * time.Millisecond) // Second user connects env.usageActivityMessage(&usageActivityMessage{ @@ -351,7 +359,7 @@ func TestMultipleUserStreams(t *testing.T) { instanceID: "instanceID2", }) - time.Sleep(100 * time.Millisecond) + clk.advance(100 * time.Millisecond) // First user disconnects env.usageActivityMessage(&usageActivityMessage{ @@ -361,7 +369,7 @@ func TestMultipleUserStreams(t *testing.T) { instanceID: "instanceID1", }) - time.Sleep(200 * time.Millisecond) + clk.advance(200 * time.Millisecond) // Second user disconnects env.usageActivityMessage(&usageActivityMessage{ @@ -383,12 +391,12 @@ func TestMultipleUserStreams(t *testing.T) { assert.Equal(t, "platform1", event1.PlatformCategory) assert.Equal(t, "instanceID1", event1.InstanceID) - assert.InDeltaf(t, 150, event1.TotalStreamMs, 50, "stream time should be approximately 150ms") + assert.Equal(t, int64(150), event1.TotalStreamMs) assert.Equal(t, "userAgent2", event2.UserAgent) assert.Equal(t, "platform2", event2.PlatformCategory) assert.Equal(t, "instanceID2", event2.InstanceID) - assert.InDeltaf(t, 300, event2.TotalStreamMs, 50, "stream time should be approximately 300ms") + assert.Equal(t, int64(300), event2.TotalStreamMs) } func TestTagsHeaderIsIncludedInUsageEvent(t *testing.T) {