Reverts #4816. Since livekit/cloud#4747 a room re-key closes every
participant's stats worker under the old room id and opens a new one
under the new id (LeaveSession/JoinSession, #4917), so RoomIDChanged
and reKeyRoom have nothing left to move, and with them go the sealed
stats batches, SetRoom, ForceClose and the reference hand-over that
only the re-key collision case needed. The worker map goes back to a
flat (room, participant) key.
The guard fixes that landed after #4816 (#4817, #4856, #4860) are kept.
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
* Deregister agent job when the server terminates it
JobTerminate ended the job on the worker but left its JobTerminate
handler and jobToWorker entry registered. Those were only cleaned up
when the worker later sent an ended job status, or on worker disconnect
for jobs still running. Workers that do not report job status after a
termination (agents-js never sends UpdateJobStatus) leaked a psrpc
handler for every job terminated by the server, e.g. each time an agent
participant left its room.
* Wait for worker registration in JobTerminate test
The worker's registration response is sent before its job request topic
is registered, so wait for the WorkerRegistered event before requesting
jobs. Also cover a worker reporting an ended status after the server
terminated the job, and clarify the JobTerminate comment.
* 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
* sfu: don't hold bindLock across the blank-frame flush in CloseWithFlush
CloseWithFlush held bindLock across the up-to-flushTimeout (1s) wait for
the closing blank-frame flush. bindLock is a control-plane lock, also
taken by Bind/SetConnected/ReceiverRestart, so every close serialized
those behind its flush wait. This was the single largest accumulator in
mutex contention profiling under track churn, aggregated across all
closing downtracks.
The flush runs in its own goroutine (writeBlankFrameRTP) and is cancelled
via blankFramesGeneration, so the lock guards nothing during the wait;
release it and re-acquire for the unbind/delete teardown.
Guard the control-plane mutators on isClosed so none can act on a
downtrack through the newly opened window. Bind in particular was
reachable: its bindState != unbound check does not catch a closed track,
since close leaves the state unbound.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
* sfu: drop redundant isClosed guard in Bind
Bind already handles a closed downtrack after codec matching (returns the
codec with no error, preserving pion's non-error contract and never
rebuilding state). The early guard returned errDownTrackClosed, which
handleRemoteAnswerReceived treats as a negotiation failure, so a
close/negotiation race could abort the remote answer. Remove it and the
error type.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
* 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>
* ccutils: use millisecond resolution for probe interval backoff
ProbeSignal grew probeInterval via Seconds(), which truncates the float
to the nearest nanosecond before scaling by time.Second:
time.Duration(probeInterval.Seconds() * BackoffFactor) * time.Second
With BackoffFactor=1.5 and the default BaseInterval of 3 s, the first
congestion signal produced 4 s instead of 4.5 s, and the error compounded
on every subsequent signal (4 s → 6 s → 9 s → 13 s ... instead of
4.5 s → 6.75 s → 10.125 s ...). The interval converges to MaxInterval
later than intended, so the server probes more frequently during sustained
congestion than the configuration requests.
probeDuration in the same function already uses the correct idiom:
time.Duration(float64(probeDuration.Milliseconds()) * Factor) * time.Millisecond
Apply the same pattern to probeInterval so both fields are computed
consistently with millisecond precision.
Found by a defect-hunting pipeline I build and run ([Dev-next-gen](https://github.com/Dev-next-gen)), using Claude Code with Anthropic's Claude Opus 5.
* ccutils: remove probe interval backoff test
Requested by reviewer.
* sfu: add opt-in verbatim abs-capture-time forwarding
Downtrack ACT forwarding gates on an RTCP sender report and
unconditionally rewrites the capture timestamp into SFU clock domain,
which suits cross-stream A/V sync but breaks consumers that need the
original publisher capture time (e.g. end-to-end latency measurement).
Add ForwardAbsCaptureTimeVerbatim (RTC config: forward_abs_capture_time_verbatim),
off by default, which skips both the gate and the rewrite.
* fix: use DisableSenderReportPassThrough setting to also gate abs-capture-time forwarding
* deps: update psrpc to v0.8.0
psrpc v0.8.0 splits the brokers into leaf packages so importing psrpc no
longer pulls in nats.go and go-redis, and drops the root constructors:
- return psrpc.NewRedisMessageBus(rc, opts...)
+ return redisbus.New(rc, opts...)
getMessageBus keeps its NewLocalMessageBus branch -- that shim stays at the
root because localbus costs no dependency. wire_gen.go carries a verbatim
copy of the provider, so it gets the identical edit rather than a wire
rerun, which would churn unrelated output.
redisbus.New takes ...bus.BusOption and psrpc.BusOption is an alias for it,
so PSRPCConfig.BusOptions() spreads in unchanged.
The nats-io requirements drop out of go.mod; this module only ever used the
Redis bus, and nats came in indirectly through psrpc's root package.
* Reorder import statements in ioservice_sip_test.go
* 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>
helpers:pinGitHubActionDigestsToSemver extends the preset we already used, so
actions stay pinned by digest; it adds an extractVersion/versioning pair that
makes Renovate follow the full vX.Y.Z tag behind the digest instead of the
mutable major, so the comment a reviewer reads carries the exact version and
updates arrive typed as patch/minor with a changelog range.
minimumReleaseAge exists to let a third-party release sit before we adopt it.
Our own actions have nothing to wait out, and the quarantine actively hurts any
referenced by branch: Renovate ages a branch ref against its head commit, so
slack-notifier-action only moves once that repo goes two weeks without a push.
This mirrors the exemption github.com/livekit/** already has under gomod.
Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
* sfu: broadcast RTP to down tracks without a per-packet closure
Every forwarded packet allocated a closure capturing the packet and layer,
and an atomic write counter that escaped with it, even for tracks with a
single subscriber. Add BroadcastRTP, which carries the packet and layer in
the worker state and sums the writes per worker. The serial path no longer
allocates; the parallel path allocates the shared state and the worker
funcval only.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* sfu: reuse the forwarded packet copy in the RED receivers
Both RED receivers copied the ExtPacket (and the opus one the rtp.Packet)
per forwarded packet, and the copies moved to the heap because the
broadcast takes their address. Keep them on the receiver like the existing
redPayloadBuf: ForwardRTP runs on one goroutine and down tracks do not keep
the packet past WriteRTP.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* sfu: skip BroadcastRTP allocation assertion under the race detector
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* sfu: group BroadcastRTP and its interfaces after the spreader methods
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
sync.Pool drops a quarter of returned items when built with -race, so the
pooled send path averages about one allocation per packet in CI and the
zero-allocation assertion flips between passing and failing. Keep the
extension checks and skip only the allocation count under race.
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
SendPacket allocated four times per down track write: a 3 byte slice for
abs-send-time, a 2 byte slice for transport-cc, and two appends to a nil
Extensions slice because the pooled header was reset to a zero value.
Give pacer.Packet fixed scratch arrays for the two extensions and marshal
into them. The Packet is pooled and owned by one send until SendPacket
returns, so nothing is shared across down tracks. Keep the pooled header's
Extensions capacity across reuse, both when returning it to the pool and
when a down track initializes it.
Adds a test that asserts zero allocations per SendPacket.
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
CodecMunger.UpdateAndGet returned a freshly allocated slice for every
forwarded packet on every down track. Return the header in a fixed array
with a length instead and marshal into it with MarshalTo. TranslationParams
carries the array; WriteRTP slices it locally.
UpdateAndGet now allocates zero times per call (was one). Escape analysis
confirms TranslationParams stays on the stack.
Co-authored-by: Claude Fable 5.1 <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>
* 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>
* telemetry: add visibility into stats worker reference underflow
Log the paths that can leave a stats worker with no references, so the
`-1` never-closed cases seen in production can be traced to their origin.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* telemetry: guard StatsWorker.MarshalLogObject against a nil receiver
The new worker-created log passes the existing worker, which is a typed
nil when there is none. Mirror ReferenceGuard.MarshalLogObject.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
* turn: accept PROXY protocol on the TCP listener
Behind a TLS-terminating or reverse proxy that dials from its own address,
the embedded TURN server reports the proxy's address to the client as
XOR-MAPPED-ADDRESS. Firefox rejects a loopback or wildcard mapped address
and abandons the allocation, so relay-only clients never get a relay
candidate (#4851).
Add turn.proxy_protocol. When set, the TCP listener requires a PROXY
protocol v1/v2 header on every connection and uses the client address it
carries; connections without the header are rejected. The header is read
before TLS, so it works with both the built-in TLS listener and
external_tls.
* turn: only trust PROXY headers from configured proxies
A PROXY header from any peer that can reach the port would let a direct
client claim an arbitrary source address. Add
turn.proxy_protocol_trusted_cidrs, defaulting to loopback, and close
connections from any other address before reading the header.
getCPUStats has no callers: node CPU load comes from hwstats.CPUStats in
GetNodeStats, and nothing in the tree reads getCPUStats. Only getLoadAvg
is still used.
Remove the function from both the windows and non-windows variants along
with the state it kept. getLoadAvg is untouched, and go-osstat stays a
direct dependency through it.
Co-authored-by: XiaoShao <26596822+xiaoshao9704@users.noreply.github.com>
getMessageBus built the bus from the redis client alone, so there was no
way to reach the gzip settings psrpc v0.7.6 added at the bus boundary.
Take rpc.PSRPCConfig, which the wire graph already provides, and pass its
bus options to both the redis and the local bus.
Compression is off by default. A peer on an older psrpc cannot decode a
compressed payload and drops it silently, so egress, ingress, SIP and
agent workers all have to be upgraded before quality is raised.
* 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.