* 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>
* 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.
* ingress: add opt-in support for udp:// URL pull ingress
URL pull ingress previously accepted only http, https and srt source
URLs. Add udp:// as an accepted scheme, gated behind a new
`ingress.enable_udp_url_pull` config option that defaults to false, so
unauthenticated UDP sources are not reachable unless the operator opts
in.
Also add a counterfeiter fake for IngressLauncher and a test covering
the scheme validation matrix.
The ingress handler binds a local socket on the caller supplied address
and port instead of connecting out like the http and srt sources do.
Spell out what that means for operators in both the config field doc
comment and config-sample.yaml: caller controlled local port binding,
unauthenticated and spoofable input, and multicast relaying of traffic
on the handler's local network.
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
* 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>
The link target was passed through url.PathEscape, which turns the '?' of
turn:host:3478?transport=udp into %3F. The target of a Link header field
is a URI reference and "?transport=" is part of TURN URI syntax
, not a query string, so a publisher that percent-decodes the
target ends up with an unparsable host.
Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
* telemetry: support roomID change for a participant
A room can get a new id while participants are connected. Key stats
workers as map[roomID]map[participantID] so moving a room is a single
map splice, and add reKeyRoom/RoomIDChanged to do the move.
Stats collected before the change are sealed off with the room they
were collected in so they stay attributed to it.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* telemetry: close superseded worker on re-key collision
Only one worker can be keyed at (room, participant). If a re-key lands
on a room that already has a worker for the same participant, keep the
one already filed there and close the superseded one so it drains and
is reaped instead of lingering in the flush list.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* telemetry: hand references to the successor on force close
A ReferenceGuard records that it activated some worker, not which one,
so a superseded worker cannot just drop its references - the survivor
would be left with references it never sees released and would never
close. Hand them over instead.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <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.
* Carry worker kind details into agent job tokens.
WorkerRegistration gains a server-controlled KindDetails field that is
passed to BuildAgentToken when assigning jobs, so kind details flow into
the participant join token for the whole session.
Co-authored-by: Cursor <cursoragent@cursor.com>
Signed-off-by: shishir gowda <shishir@livekit.io>
Co-authored-by: Cursor <cursoragent@cursor.com>
* utils: make Median generic, overflow-safe, and add comprehensive
* tests: fix staticcheck unused variable warning in changenotifier_test.go
* tests: switch from assert to require for consistency with existing tests
* trigger ci rerun
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>
Let the docker-backed service tests be skipped with a flag.
TestMain called log.Fatalf when it could not reach a docker daemon, so
the whole package refused to run without one, including every test in it
that needs no container at all.
Record why docker is unavailable instead, and gate the tests that want a
container on it. A run asks to go without them with -docker=false;
otherwise a missing daemon still fails the package, so an unreachable
daemon stays a broken build rather than a run that quietly covers less
than the last one did.
Co-Authored-By: Claude Opus 5 (1M context) <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>
VP9 SVC sends each spatial layer as its own encoded frame,
so a picture carries one trailer per layer, but only the
top layer's last packet has the RTP marker bit set.
Fix https://github.com/livekit/egress/issues/1347
* 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