Files
livekit/pkg/telemetry/statsworker_test.go
Raja SubramanianandClaude Fable 5.1 7c82d3e103 telemetry: do not recreate a stats worker for a released guard (#4860)
* telemetry: do not recreate a stats worker for a released guard

A ParticipantActive overtaken by the participant's close arrives with a
guard ParticipantLeft already released and replaced the closed worker
with one nothing could release.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

* telemetry: handle a released guard independently of map presence

A released guard reaching getOrCreateWorker after the closed worker was
reaped still created a zero-reference worker. Return nil as found
instead, and make SetConnected nil-safe for ParticipantActive.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-11 21:23:13 +05:30

103 lines
3.6 KiB
Go

package telemetry
import (
"context"
"testing"
"github.com/stretchr/testify/require"
"go.uber.org/zap/zapcore"
"github.com/livekit/protocol/livekit"
)
func TestStatsWorker(t *testing.T) {
t.Run("reference counted close works", func(t *testing.T) {
var g0, g1 ReferenceGuard
w := newStatsWorker(t.Context(), nil, "", "", "", "", &g0)
require.False(t, w.Closed(&g1))
require.False(t, w.Close(&g0))
require.False(t, w.Closed(&g1))
require.True(t, w.Close(&g1))
require.True(t, w.Closed(&g1))
})
// a ReferenceGuard records that it activated some worker, not which one, so a
// superseded worker has to hand its references to the one reachable in its place
t.Run("force close hands references to the successor", func(t *testing.T) {
t.Run("a guard shared by both workers", func(t *testing.T) {
// the second worker never got a reference, the guard was already activated
var g ReferenceGuard
superseded := newStatsWorker(t.Context(), nil, "", "", "", "", &g)
survivor := newStatsWorker(t.Context(), nil, "", "", "", "", &g)
require.Equal(t, 1, superseded.refCount.count)
require.Equal(t, 0, survivor.refCount.count)
require.True(t, superseded.ForceClose(survivor))
require.Equal(t, 0, superseded.refCount.count)
require.Equal(t, 1, survivor.refCount.count)
// without the hand over this would leave the survivor at -1 and never closed
require.True(t, survivor.Close(&g))
require.True(t, survivor.Closed(&g))
})
t.Run("a guard per worker", func(t *testing.T) {
var gSuperseded, gSurvivor ReferenceGuard
superseded := newStatsWorker(t.Context(), nil, "", "", "", "", &gSuperseded)
survivor := newStatsWorker(t.Context(), nil, "", "", "", "", &gSurvivor)
require.True(t, superseded.ForceClose(survivor))
require.Equal(t, 2, survivor.refCount.count)
// the superseded worker's owner departs, it must not close the survivor early
require.False(t, survivor.Close(&gSuperseded))
require.True(t, survivor.Close(&gSurvivor))
})
t.Run("closing an already closed worker holds on to its references", func(t *testing.T) {
var g ReferenceGuard
superseded := newStatsWorker(t.Context(), nil, "", "", "", "", &g)
survivor := newStatsWorker(t.Context(), nil, "", "", "", "", nil)
require.True(t, superseded.ForceClose(nil))
require.False(t, superseded.ForceClose(survivor))
require.Equal(t, 0, survivor.refCount.count)
})
})
t.Run("logging a nil worker does not panic", func(t *testing.T) {
var w *StatsWorker
require.NoError(t, w.MarshalLogObject(zapcore.NewMapObjectEncoder()))
})
}
func TestGetOrCreateWorkerReleasedGuard(t *testing.T) {
// ParticipantActive overtaken by the participant's close arrives with a guard that
// ParticipantLeft already released. It must not replace the closed worker with one
// nothing can release.
ts := &telemetryService{workers: make(map[livekit.RoomID]map[livekit.ParticipantID]*StatsWorker)}
roomID, pID := livekit.RoomID("room"), livekit.ParticipantID("participant")
var g ReferenceGuard
w, found := ts.getOrCreateWorker(context.Background(), roomID, "", pID, "", &g)
require.False(t, found)
require.True(t, w.Close(&g))
t.Run("closed worker still in the map", func(t *testing.T) {
late, found := ts.getOrCreateWorker(context.Background(), roomID, "", pID, "", &g)
require.True(t, found)
require.Same(t, w, late)
require.Same(t, w, ts.workers[roomID][pID])
})
t.Run("closed worker already reaped", func(t *testing.T) {
delete(ts.workers[roomID], pID)
late, found := ts.getOrCreateWorker(context.Background(), roomID, "", pID, "", &g)
require.True(t, found)
require.Nil(t, late)
require.Empty(t, ts.workers[roomID])
late.SetConnected()
})
}