Files
livekit/pkg/rtc/signalling/signallerasyncbase_test.go
Raja SubramanianandClaude Opus 5 0aae7a4894 fix: hold signal messages until the ReconnectResponse goes out (#4827)
* fix: hold signal messages until the ReconnectResponse goes out

Clients take the ReconnectResponse as the first message on a resumed or
migrated in signal connection, anything ahead of it is dropped. Hold
messages back until it has been written.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: keep queued participant updates on a path that always flushes

Queue only while the participant is not ready, that queue is always
drained by the join or reconnect response. Log a dropped SDP, it leaves
the negotiation waiting until the state machine recovers it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: flush queued updates whenever the connection opens

Queue participant updates while the handshake is pending again, and give
the signaller a hook that fires when the connection opens, on an explicit
open and on the handshake window expiring, so the queue always drains.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: drive the handshake window off a timer and read the gate atomically

The window ran only when something asked whether the handshake was
pending, so a connection with nothing else to send held its queue.
Reading the gate under the lock the flush takes closes the race where an
update queued just after a flush.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: ignore a handshake timeout from a wait that has ended

Stop cannot cancel a timeout that is already running, so tag each wait
and let a timeout act only on its own. Count opens atomically in the
test, the timer fires on its own goroutine.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
2026-09-01 13:59:04 +05:30

158 lines
4.9 KiB
Go

package signalling
import (
"testing"
"time"
"github.com/stretchr/testify/require"
"go.uber.org/atomic"
"github.com/livekit/protocol/logger"
"github.com/livekit/livekit-server/pkg/routing/routingfakes"
"github.com/livekit/livekit-server/pkg/rtc/types"
)
func newTestSignallerBase(onHandshakeOpened func()) *signallerAsyncBase {
return newSignallerAsyncBase(signallerAsyncBaseParams{
Logger: logger.GetLogger(),
OnHandshakeOpened: onHandshakeOpened,
})
}
// currentHandshakeGeneration is what a timeout armed right now would be tagged with
func currentHandshakeGeneration(s *signallerAsyncBase) uint32 {
s.resSinkMu.Lock()
defer s.resSinkMu.Unlock()
return s.handshakeGeneration
}
func TestHandshakeGate(t *testing.T) {
t.Run("a resumed connection holds messages back", func(t *testing.T) {
s := newTestSignallerBase(nil)
require.False(t, s.HandshakePending())
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
require.True(t, s.HandshakePending())
s.OpenHandshake()
require.False(t, s.HandshakePending())
})
t.Run("other sink swaps do not hold messages back", func(t *testing.T) {
s := newTestSignallerBase(nil)
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonUnknown)
require.False(t, s.HandshakePending())
})
t.Run("closing the connection clears the gate", func(t *testing.T) {
s := newTestSignallerBase(nil)
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
require.True(t, s.HandshakePending())
s.CloseSignalConnection(types.SignallingCloseReasonParticipantClose)
require.False(t, s.HandshakePending())
})
t.Run("the handshake window opens the gate", func(t *testing.T) {
s := newTestSignallerBase(nil)
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
require.True(t, s.HandshakePending())
s.onHandshakeTimeout(currentHandshakeGeneration(s))
require.False(t, s.HandshakePending())
})
t.Run("a swap to a fresh connection opens the gate", func(t *testing.T) {
s := newTestSignallerBase(nil)
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
require.True(t, s.HandshakePending())
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonUnknown)
require.False(t, s.HandshakePending())
})
t.Run("a timeout of a wait that has ended does nothing", func(t *testing.T) {
var opened atomic.Int32
s := newTestSignallerBase(func() { opened.Inc() })
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
stale := currentHandshakeGeneration(s)
// the client resumes again while the first timeout is running, Stop cannot
// cancel a callback that has already started
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
s.onHandshakeTimeout(stale)
require.True(t, s.HandshakePending())
require.Zero(t, opened.Load())
// and the wait in progress still opens on its own timeout
s.onHandshakeTimeout(currentHandshakeGeneration(s))
require.False(t, s.HandshakePending())
require.EqualValues(t, 1, opened.Load())
})
}
func TestHandshakeGateNotifiesOnOpen(t *testing.T) {
t.Run("on an explicit open", func(t *testing.T) {
var opened atomic.Int32
s := newTestSignallerBase(func() { opened.Inc() })
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
require.Zero(t, opened.Load())
s.OpenHandshake()
require.EqualValues(t, 1, opened.Load())
// only the transition notifies
s.OpenHandshake()
require.EqualValues(t, 1, opened.Load())
})
t.Run("on the handshake window expiring", func(t *testing.T) {
var opened atomic.Int32
s := newTestSignallerBase(func() { opened.Inc() })
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
s.onHandshakeTimeout(currentHandshakeGeneration(s))
require.False(t, s.HandshakePending())
require.EqualValues(t, 1, opened.Load())
})
t.Run("the window fires without anything else touching the gate", func(t *testing.T) {
var opened atomic.Int32
s := newTestSignallerBase(func() { opened.Inc() })
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
require.Eventually(t, func() bool {
return !s.HandshakePending() && opened.Load() == 1
}, 2*handshakeWindow, 100*time.Millisecond)
})
t.Run("not on a connection close", func(t *testing.T) {
var opened atomic.Int32
s := newTestSignallerBase(func() { opened.Inc() })
s.SwapResponseSink(&routingfakes.FakeMessageSink{}, types.SignallingCloseReasonResume)
s.CloseSignalConnection(types.SignallingCloseReasonParticipantClose)
require.False(t, s.HandshakePending())
require.Zero(t, opened.Load())
})
t.Run("not when the gate was never armed", func(t *testing.T) {
var opened atomic.Int32
s := newTestSignallerBase(func() { opened.Inc() })
s.OpenHandshake()
require.Zero(t, opened.Load())
})
}