* rtc: record a room's departure when the participant leaves it
RemoveParticipant removed the participant from the room, closed it, and
only then recorded leftAt. CloseIfEmpty runs every second, so a tick
landing during the close saw an empty room with a stale or zero
departure time and closed it at once: with no departure recorded it
compares against the room's creation time and empty_timeout, and with an
earlier departure it measures the grace period from that older
timestamp. Either way departure_timeout is skipped.
Record the departure where the participant is removed, under the lock
that already covers it.
* rtc: drop the removal wait from the departure timeout test
The close signal already pins the interleaving, and the room's close
state is decided by the CloseIfEmpty call itself.
* rtc: split EndSession out of MoveToRoom
A participant's session in a room can end without the participant
closing: today on a move, and next on a room re-key, where a room ended
by the sweeper while participants were joining restarts under a new room
id with the same participants in it.
EndSession is the room-independent half of MoveToRoom: report the
published tracks unpublished, report the session end, run the caller's
leave under the ending session's telemetry guard, then take a fresh
guard and reset the deferred resolvers for the next session. The guard
matters: ParticipantLeft releases it and the telemetry service refuses
to open a stats worker for a released guard, so without a fresh one the
next session's join would silently get no worker.
MoveToRoom keeps the subscriber teardown and the onClose fan-out, passed
in as the leave.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* rtc: rename EndSession to LeaveSession, add JoinSession
EndSession read as terminal, but the participant survives it and joins
another session, so name the pair for what happens to the participant:
LeaveSession and JoinSession.
JoinSession is the other half that review pointed out was missing from
the API: the caller reports the join, then the published tracks are
reported published again under the new session, pairing the unpublish
LeaveSession reported. Without it the re-publish lived only in the
caller, with nothing in OSS proving the two halves balance.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
participants.
While migrating, the old node will have the recently disconnected
participants in the migrating out participant's cache and needs to sent
out to the participant from the migrating in node via the new WebSocket.
This method can be used to exchange that data between the two nodes.
Keeping it as `recently` disconnected and leaving the semantics upto how
recent for implementation.
There are reports of disconnected participant update missing -
https://github.com/livekit/client-sdk-js/pull/2096. That results in
clients retaining disconnected participant(s) in its roster.
Two changes
1. Look for participant disconnects a little farther back than last
signal time as client app might have lost visibility earlier.
2. Look at all participants within the relaxed window above. The update
cache is LRU. So, in theory can stop when an entry is before the time
threshold above. But, theoretically it is possible that `updatedAt`
time and LRU queuing time can flip ordering although very unlikely.
I think we should do one of the following to make it more robust, but
both require client side changes
1. Send full participant list in `ReconnectResponse`. Clients can use
that to drop participants not in the list.
2. Add a sequence number to all signal responses (maybe signal requests
as well) and have client send up the last response it saw in
`SyncState` message. Server can then replay all the messages after
that. This also requires a message cache on server side.
* dynacast: re-send subscribed qualities after every publisher answer
The SFU answers publisher offers without the simulcast pause markers
(SetIgnoreRidPauseForRecv), so applying the answer re-enables in the
browser the layers that dynacast had paused. Dynacast only notifies
changes, so those layers stayed on until the subscribed quality of the
track changed again.
Send the publisher its committed subscribed qualities again 1 s after
each answer, once the client has applied it: livekit-client handles the
next signal message without waiting for the answer to be applied.
Fixes https://github.com/livekit/livekit/issues/4906
* dynacast: resend once per burst of publisher answers, only for paused layers
In SPC, every subscription change of a participant goes through a new publisher offer and answer, so answers come in bursts when a room is busy. Each answer scheduled its own resend of the subscribed qualities, sending the same update several times in a row.
Thist change keeps one pending resend per participant instead. It fires 2 s after the last answer, and never later than 8 s after the first answer of the burst, so a steady stream of answers cannot postpone it forever. It also skips the resend for tracks whose committed qualities are all HIGH: with every layer enabled, an answer has nothing paused to re-enable. We can save those messages.
* agent: release JobTerminate handler when job ends before registration
JobRequest registered the job's JobTerminate handler after AssignJob
returned, so an ended UpdateJobStatus or worker disconnect handled in
between deregistered nothing and the handler leaked. AssignJob also sent
the assignment before recording the job as running, so an immediate
ended update was dropped.
Record the job before sending the assignment, and after registering the
handler recheck under h.mu that the job is still running and its worker
still registered, deregistering otherwise.
Fixes#4901
* replace utils.CloneProto with proto.CloneOf
* agent: wait on worker readiness and poll for released handlers in test
* Increase receiver loadbalance threshold when batch io enabled
With batch IO the send syscall moves to the BatchConn flush goroutine, so WriteRTP only
translates the packet and enqueues it. Increash the threshold to reduce
* lint
* Do not request media sections for transceivers that cannot send
getNumUnmatchedTransceivers counted every local transceiver without a mid.
A transceiver whose track was removed before its section was negotiated
(direction recvonly, no sender) or that was stopped (inactive) can never be
matched to the recvonly section the remote adds in response to a
MediaSectionsRequirement, so it stayed "outstanding" forever: every answer
asked the remote for a section for it again, the remote appended another
batch of transceivers each round, and its offer grew by a fixed amount per
negotiation until it exceeded limit.signal_message_size_limit and the
signalling connection was closed (full reconnect, then the cycle restarts).
Count only transceivers that have something to send. A released transceiver
that AddTrack reuses later becomes sendrecv and is counted again.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
* Address review: blank lines between test functions, assert transceiver reuse
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
* test: stop registering the global DefaultConfig with the logger
InitLoggerFromConfig hands the pointer to zaputil.ComponentLeveler as its
level resolver, and protocol a879e94 gave logger.Config a
ResolveComponentLevel that takes c.lock. The two test packages passed
&config.DefaultConfig.Logging, so every first-time component level
resolution locked a mutex inside the global DefaultConfig while
config.NewConfig marshalled that same global from another goroutine:
Read at 0x39baa08 yaml.Marshal -> pkg/config/config.go:633 (NewConfig)
Write at 0x39baa08 ComponentLeveler.resolve -> zaputil/leveler.go:117
That tripped the race detector in pkg/rtc TestPreferMediaCodecForPublisher
and broke CI on master.
Register a test-local LoggingConfig instead, so the logger never touches
the global that NewConfig reads. DefaultConfig.Logging only ever sets
PionLevel, so behavior is unchanged. A plain struct copy is not an option
here: logger.Config holds a mutex and vet's copylocks would reject it.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* test: make RTCClient callbacks safe to set while the client runs
The scenario tests assign c2.OnDataReceived after waitUntilConnected, so
the assignment raced the data channel goroutine already reading the field
in handleDataMessage:
Read test/client.(*RTCClient).handleDataMessage client.go:1146
Write test.scenarioDataPublish.func1 scenarios.go:160
This is what failed TestMultinodeDataPublishing on master.
Replace the three exported callback fields with atomic pointers behind
SetOnConnected/SetOnDataReceived/SetOnDataUnlabeledReceived. OnConnected
and OnDataUnlabeledReceived have no writers today, but they are read from
the same background goroutines and would race the moment one appeared, so
all three move together rather than leaving a split API on one struct.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
SubscriptionManager tracks data-track subscriptions but exposed no accessor
for the active ones, so callers could see subscribed and published media tracks
and published data tracks, but not subscribed data tracks. Add
GetSubscribedDataTracks(), returning the bound data down-tracks, mirroring
GetSubscribedTracks() for media, so a participant's full track set can be
accounted for. Regenerated the LocalParticipant fake.
Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
* rtc: send connection quality to participants that subscribe after their first update
Added tracking for participants' connection quality updates to ensure all subscribed participants receive necessary updates, even if they were added after the last quality announcement
* rtc: record sent connection quality only after a successful send
A failed SendConnectionQualityUpdate must not mark the participant as informed, otherwise the update is never retried while qualities stay stable. Also simplify the untold-subscription check.
* 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>
* report a reason on room-ended telemetry
Room.Close already took a ParticipantCloseReason, but OnClose was func()
with no arguments, so the reason was structurally dropped before it could
reach RoomEnded.
Add types.RoomCloseReason and carry it through Close -> OnClose -> RoomEnded,
where it lands on both the analytics event and the room_finished webhook.
Room.Close now takes the room reason and derives the participant reason from
it, so the two can never disagree. ToParticipantCloseReason maps each reason
to exactly what its call site passed before, so no participant-facing
behaviour changes; a table test pins that mapping.
* fix(test): update webhook test for RoomCloseReason
test/ was outside the packages checked before pushing, so this Close call
site was missed. Also assert the reason reaches the room_finished webhook,
which the unit test cannot cover since it runs with a nil notifier.
In migration cases, dummy receiver trackInfo is used by relay tracks to
set up the receivers and those need the proper track info.
Also check for proper receiver when adding a migrated track.
TestUnsubscribe checked that the changed-notifier observer was gone as
soon as the unsubscribe had settled, but setDesired leaves the
RemoveObserver call to a goroutine of its own and nothing the test waits
on orders against it, so CI caught the assertion running first.
TestSubscribe has the same defect on the unsubscribed callback, which
unmarkSubscribedTo delivers with a bare go while the subscribed one is
called inline.
Wait for both, each with a message of its own, so that a leak that is
real still says which of them broke.
Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
pion no longer starts the repair stream reader when a custom BufferFactory
is set, so the mid/rid/rsid extensions were never observed and simulcast RTX
streams were never paired with their primary streams.
Extract the extensions on the buffer write path instead. Migrated publishers
send no extensions at all, so pair those from SimTracks.
Adds an integration test covering both paths, and moves the vnet setup it
shares with the downtrack test into pkg/testutils/vnettest.
Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
* Reduce locking in media track + telemetry listener on move participant.
Telemetry listener was not getting it from new room on room move.
* update comment
* fix signal bytes reporting
* room aware telemetry listener in media track
* data track telemetry listener
SendData API messages do not have sender ID or sequence number. So, they
were not cached and hence excluded from the reliable caching feature
which is meant to provide reliability of data channel messages between
the time other participants see participant as ACTIVE (which happens on
ICE connected) and data channel being open (DTLS done + data channels
opened).
Add a small cache for that and flush those messages on data channel
establishment.
* Check for ICE connection before closing participant on signal close.
Only close the participant if ICE has not connected by the time signal
source is closed. If ICE had connected, candidates have been exchanged
and link was established. So, it should be resumable.
Waiting for DTLS closed the participant sometimes in the windowa after
ICE connection, but before DTLS finishes and that unnecessarily closed
the participant forcing a full reconnect.
* variable name
* Limit number of pending tracks per participant.
Prevents just a signalling connection adding tracks without actually
publishing them growing a large number.
* add to supervisor only if pending track is accepted
With https://github.com/livekit/livekit/pull/4706, there was a case of
some downstream component taking a long time while lock was held. While
the underlying cause of holding a lock while doing callback was removed
in that PR, to catch such cases, some publish side metric anomaly would
be useful to monitor and alert on.
Adding a publish time record for pending tracks on participant close.
That would inflate the publish time for participants not being able to
publish and can be alerted on as it will spike up the value at node
level and at cluster level if multiple nodes have the issue.
* Add configurable read-message size limit on signalling WebSockets
Set a read limit on both the client-facing (/rtc) and agent worker
WebSocket connections so an oversized frame is rejected by the transport
before being buffered. The limits are operator-tunable via
signal_message_size_limit and agent_signal_message_size_limit, both
defaulting to 2 MiB (0 disables).
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Add tests for signalling WebSocket read-message size limit
Cover the configurable signal_message_size_limit added in the prior
commit:
- config: assert both limits default to 2 MiB and that a YAML override
(including 0 to disable) is parsed correctly.
- full-path integration: a real client connects to /rtc on a single-node
server and an oversized frame is rejected by the transport with a 1009
close; a 0 limit leaves the connection unbounded and signalling
proceeds.
Adds setupSingleNodeTestWithConfig so a single-node server can be started
with config overrides.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Bound decompressed size of signalling WebSocket messages
conn.SetReadLimit only accounts for the compressed bytes read off the
wire, and the client-facing /rtc upgrader negotiates permessage-deflate,
so a small compressed frame could still expand into a much larger buffer
once inflated. Enforce the same limit on the decompressed message by
reading through NextReader + io.LimitReader in WSSignalConnection instead
of the unbounded ReadMessage.
The transport-level SetReadLimit is kept as the cheap wire-level guard;
the new check is the decompressed-size backstop.
Adds NextReader to the WebsocketClient interface (regenerated fake) and a
unit test plus permessage-deflate integration tests covering the
amplification case.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Cover a couple of more cases on data track runt packet handling.
* Guard data track header parser against extensions-size integer wraparound.
Widen the extensions-size arithmetic to int so a 0xFFFF wire value no
longer wraps in uint16, and reject any packet whose computed hdrSize
exceeds the buffer before slicing the payload.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Fix publish track count on migration in.
https://github.com/livekit/livekit/pull/4707 addressed the case of
publish tracks overcounting due to synthesised track publish on migrate
in. But, it introduced an issue where published tracks count could go
negative because unpublish subtracted the counter irrespective of the
track actually migrated in or not.
Fix it by keeping track of local publish.
Also, the older code was skipping publisher track count increase if the
synthesised publish was handled first. Address it by checking if the
track is actually new (i. e. fresh local publish) when the track was
already created in the migrate in path.
* fix pub time for tracks published after migration
* test
* prevent multiple track egresses
* pass room proto directly to telemetry events
* Keep telemetry analytics events on a minimal room, gate full room in webhooks
---------
Co-authored-by: Simon Beeli <simon.beeli@gmx.ch>
onMediaLossUpdate notified the participant handler directly, which only
sends a leave request with resume action. handleConnectionFailed that
actually switches the ICE preference to TCP/TLS was never called on
this path, so the client reconnected over UDP again and the fallback
kept firing every 30-60s without ever migrating.
Fixeslivekit/livekit#4702
* log high stream start latency.
There is something wrong in measurement as audio is showing high p99
latency. Must be misattributing samples. So, logging for high latency to
understand this better.
* use correct variable
* time since create
* Register h264 main profile if enabled explicitly
We don't support the h264 main profile for compatibility,
user can enabled it by set fmtp explicitly in codec config
to enable it if want to use it in special scenario.
* go mod
- Count a publish attempt on a migrating in tarck as there is no
AddTrack for that.
- Add cancel publish only if the participant connection is canceled
- Do not add publish counter for synthetic publish attempts which
happens for migrating in tracks. It will be counted on migrating in
node when the track is actually published, i. e. negotiated/packets
flowing.
* Do not call telemetry listener under pending track lock.
Fix the TrackPublishRequested call of telemetry listener.
Audited other callbacks to ensure that it is not under lock.
* missed some paths of recording it, thanks Devin
* Record subscribe stream start time in prometheus.
Adjust for mutes, i. e. take the last unmute time as the start point and
calculate time till the first byte is sent.
* close the tiny window of race
* Prevent long tail sample when publisher glitches.
Thanks to @milos-lk for this.
Publisher restarting would have reset the layer and would have caused a
sample with very high stream start time. We only need to capture when we
do a dummy start or when the state is seeded to a different node upon
migration.
* reduce a diff
* test
* changed the wrong thing, thank you Devin