From f5c2dafe0ab4a22428f0d0f2fd85a521d1b43e71 Mon Sep 17 00:00:00 2001 From: Alex Date: Thu, 20 Aug 2026 16:44:43 -0700 Subject: [PATCH 1/8] Self-review --- pkg/sip/inbound.go | 9 +++--- pkg/sip/media_pipeline.go | 6 ++-- pkg/sip/media_pipeline_test.go | 52 +++++++++++++++++++++++++++++----- pkg/sip/room.go | 2 +- 4 files changed, 53 insertions(+), 16 deletions(-) diff --git a/pkg/sip/inbound.go b/pkg/sip/inbound.go index c2d42c23..5d07cf16 100644 --- a/pkg/sip/inbound.go +++ b/pkg/sip/inbound.go @@ -1107,11 +1107,11 @@ func (c *inboundCall) waitForCallEnd(ctx context.Context, ackReceived <-chan str // Today we seek to enforce all calls to be ACKed or dropped. // Sometimes, though, we do not see ACKs for invites (e.g due to possible // issues with load balancing). - // To accomodate this issue, instead of ending the call right here, we instead + // To accommodate this issue, instead of ending the call right here, we instead // set an aggressive timeout as a softer fallback. // If the issue really is a dropped ACK, media is expected to flow shortly, - // allowing us to accomodate this eventuality. If, however, there is no media - // obserrved, the call still ends quickly. + // allowing us to accommodate this eventuality. If, however, there is no media + // observed, the call still ends quickly. // Once ACKs are certain to be reliable, we will end the call here. c.media.SetTimeout(min(inviteOkAckLateTimeout, c.s.conf.MediaTimeoutInitial), mediaTimeout) } @@ -1657,8 +1657,7 @@ func (c *inboundCall) publishTrack(features []livekit.SIPFeature, featureFlags m old.Close() } if old := c.media.WriteInboundDTMFTo(c.lkRoom.GetInboundDTMFWriter()); old != nil { - c.log().Warnw("media port has unexpected inbound dtmf writer", nil) - old.Close() + old.Close() // Can be pinDTMFWriter } return nil } diff --git a/pkg/sip/media_pipeline.go b/pkg/sip/media_pipeline.go index 18bcec6f..2112f916 100644 --- a/pkg/sip/media_pipeline.go +++ b/pkg/sip/media_pipeline.go @@ -152,7 +152,7 @@ func (p *mediaPortPipeline) init( return nil } -// Construct the Audio and optionally DTM pipleine from SIP RTP to LK PCM, in reverse order. +// Construct the Audio and optionally DTMF pipline from SIP RTP to LK PCM, in reverse order. func (p *mediaPortPipeline) setupInput(mc *sdp.MediaConfig, audioToRoom msdk.PCM16Writer, dtmfToRoom msdk.WriteCloser[*livekit.SipDTMF]) error { var err error var inboundLatencyEntry atomic.Int64 @@ -231,8 +231,8 @@ func (p *mediaPortPipeline) handleEventRTP(h *rtp.Header, payload []byte) error }) } -// Construct the Audio and optionally DTM pipleine from LK PCM to SIP RTP -// Retuirns the insulated (nopCloser) connectors, and an error. +// Construct the Audio and optionally DTMF pipline from LK PCM to SIP RTP +// Returns the insulated (nopCloser) connectors, and an error. func (p *mediaPortPipeline) setupOutput(mc *sdp.MediaConfig, incomingSampleRate int) error { p.rtpLoopWG.Go(p.rtpLoop) w, err := p.sess.OpenWriteStream() diff --git a/pkg/sip/media_pipeline_test.go b/pkg/sip/media_pipeline_test.go index b2b3e21e..ed629b1c 100644 --- a/pkg/sip/media_pipeline_test.go +++ b/pkg/sip/media_pipeline_test.go @@ -123,6 +123,44 @@ func (c *dtmfCollector) snapshot() []*livekit.SipDTMF { return out } +// pcmCollector accumulates decoded room audio. The pipeline writes from the RTP +// read goroutine while the test reads, so every access is guarded. +type pcmCollector struct { + sampleRate int + + mu sync.Mutex + buf msdk.PCM16Sample +} + +func (c *pcmCollector) String() string { return fmt.Sprintf("pcmCollector(%d)", c.sampleRate) } + +func (c *pcmCollector) SampleRate() int { return c.sampleRate } + +func (c *pcmCollector) Close() error { return nil } + +func (c *pcmCollector) WriteSample(sample msdk.PCM16Sample) error { + c.mu.Lock() + defer c.mu.Unlock() + c.buf = append(c.buf, sample...) + return nil +} + +func (c *pcmCollector) len() int { + c.mu.Lock() + defer c.mu.Unlock() + return len(c.buf) +} + +// since returns a copy of everything written after the first n samples. +func (c *pcmCollector) since(n int) msdk.PCM16Sample { + c.mu.Lock() + defer c.mu.Unlock() + if n >= len(c.buf) { + return nil + } + return slices.Clone(c.buf[n:]) +} + // pipelineHarness is the durable side of a mediaPort: UDP pipe, pipeline config, // buffer anchors, and a synthesized MediaConfig. The pipeline itself is swapped // on configure / reconfigure. @@ -136,7 +174,7 @@ type pipelineHarness struct { audioOut *msdk.WriteCloserSwitch[msdk.PCM16Sample] dtmfIn *msdk.WriteCloserSwitch[*livekit.SipDTMF] dtmfOut *msdk.WriteCloserSwitch[*livekit.SipDTMF] - roomAudio *msdk.PCM16Sample + roomAudio *pcmCollector roomDTMF *dtmfCollector pipeline *mediaPortPipeline ssrcCount atomic.Uint64 @@ -159,10 +197,10 @@ func newPipelineHarness(t *testing.T, sampleRate int) *pipelineHarness { audioOut: msdk.NewWriteCloserSwitch[msdk.PCM16Sample](sampleRate), dtmfIn: msdk.NewWriteCloserSwitch[*livekit.SipDTMF](dtmf.SampleRate), dtmfOut: msdk.NewWriteCloserSwitch[*livekit.SipDTMF](dtmf.SampleRate), - roomAudio: new(msdk.PCM16Sample), + roomAudio: &pcmCollector{sampleRate: sampleRate}, roomDTMF: &dtmfCollector{}, } - h.audioIn.Swap(msdk.NewPCM16BufferWriter(h.roomAudio, sampleRate)) + h.audioIn.Swap(h.roomAudio) h.dtmfIn.Swap(h.roomDTMF) h.conf = &MediaPortPipelineConfig{ log: log, @@ -348,7 +386,7 @@ func (h *pipelineHarness) testAudioFromRoom(t *testing.T) { } func (h *pipelineHarness) testAudioFromPort(t *testing.T) { - before := len(*h.roomAudio) + before := h.roomAudio.len() packetsBefore := h.packetCount.Load() clock := h.codec.Info().RTPClockRate if clock == 0 { @@ -363,15 +401,15 @@ func (h *pipelineHarness) testAudioFromPort(t *testing.T) { return h.packetCount.Load() >= packetsBefore+5 }, time.Second, 5*time.Millisecond, "RTP should be accepted") require.Eventually(t, func() bool { - return len(*h.roomAudio) > before + return h.roomAudio.len() > before }, time.Second, 5*time.Millisecond, "decoded PCM should reach room (packets=%d input=%d failed=%d ignored=%d room=%d)", h.packetCount.Load(), h.pipeline.conf.stats.InputPackets.Load(), h.pipeline.conf.stats.FailedPackets.Load(), h.pipeline.conf.stats.IgnoredPackets.Load(), - len(*h.roomAudio), + h.roomAudio.len(), ) - require.Greater(t, pcmEnergy((*h.roomAudio)[before:]), int64(0), "decoded room audio should carry energy") + require.Greater(t, pcmEnergy(h.roomAudio.since(before)), int64(0), "decoded room audio should carry energy") } func (h *pipelineHarness) testDTMFFromRoom(t *testing.T) { diff --git a/pkg/sip/room.go b/pkg/sip/room.go index ac9fcb35..2d4abb32 100644 --- a/pkg/sip/room.go +++ b/pkg/sip/room.go @@ -203,7 +203,7 @@ type RoomInterface interface { // GetInboundAudioWriter returns a writer that, when written to, writes // audio to the room. GetInboundAudioWriter() (msdk.PCM16Writer, error) - // GetInboundDTMFWriter returns a writer that, when weritten to, writes DTMF + // GetInboundDTMFWriter returns a writer that, when written to, writes DTMF // to the room. GetInboundDTMFWriter() msdk.WriteCloser[*livekit.SipDTMF] } From 9c9660431ba15bd2e477c43e3c94e4fa9f71ad3b Mon Sep 17 00:00:00 2001 From: Alex Date: Thu, 20 Aug 2026 17:01:55 -0700 Subject: [PATCH 2/8] Self-review: dead ends without response to reINVITE --- pkg/sip/inbound.go | 22 +++++++++++++++++----- 1 file changed, 17 insertions(+), 5 deletions(-) diff --git a/pkg/sip/inbound.go b/pkg/sip/inbound.go index 5d07cf16..6262f732 100644 --- a/pkg/sip/inbound.go +++ b/pkg/sip/inbound.go @@ -72,7 +72,9 @@ const ( var allowHeader = sip.NewHeader("Allow", "INVITE, ACK, CANCEL, BYE, NOTIFY, REFER, MESSAGE, OPTIONS, INFO, SUBSCRIBE") var errNoACK = errors.New("no ACK received for 200 OK") -var errInternal = errors.New("internal error") + +// RFC 3261 §21.4.27 / §14.2 — glare: INVITE received while an INVITE we sent is in progress. +const statusRequestPending sip.StatusCode = 491 // hashPassword creates a SHA256 hash of the password for logging purposes func hashPassword(password string) string { @@ -408,7 +410,11 @@ func (s *Server) processInvite(req *sip.Request, tx sip.ServerTransaction) (retE existing.log().Infow("reinvite", "content-length", req.ContentLength(), "cseq", cc.InviteCSeq()) if err := existing.updateRemoteFromSDP(sdpBodyFromRequest(req)); err != nil { log.Errorw("failed to update inbound call SDP", err) - cc.RejectAsKeepAlive(sip.StatusBadRequest, "Bad Request") + if ok := errors.As(err, &SDPError{}); ok { + cc.RejectAsKeepAlive(sip.StatusBadRequest, "Bad Request") + } else { + cc.RejectAsKeepAlive(sip.StatusInternalServerError, "Internal Server Error") + } return nil } // TODO(alexfish): Reply with the new SDP. @@ -421,17 +427,23 @@ func (s *Server) processInvite(req *sip.Request, tx sip.ServerTransaction) (retE if oc != nil && oc.cc != nil && oc.cc.InviteCSeq() < newCSeq { if oc.media == nil { oc.log.Errorw("outbound call media has not been negotiated", nil) - return errInternal + cc.RejectAsKeepAlive(statusRequestPending, "Request Pending") + return nil } localSDP, err := oc.media.GetLocalSDP() if err != nil || len(localSDP) == 0 { oc.log.Errorw("outbound call does not have an SDP", nil) - return errInternal + cc.RejectAsKeepAlive(statusRequestPending, "Request Pending") + return nil } oc.log.Infow("accepting reinvite", "content-length", req.ContentLength(), "cseq", cc.InviteCSeq()) if err := oc.updateRemoteFromSDP(sdpBodyFromRequest(req)); err != nil { log.Errorw("failed to update outbound call SDP", err) - cc.RejectAsKeepAlive(sip.StatusBadRequest, "Bad Request") + if ok := errors.As(err, &SDPError{}); ok { + cc.RejectAsKeepAlive(sip.StatusBadRequest, "Bad Request") + } else { + cc.RejectAsKeepAlive(sip.StatusInternalServerError, "Internal Server Error") + } return nil } oc.cc.RecordInvite(newCSeq) From b11494e06bb6d10865108c8593ef78151b4e94b7 Mon Sep 17 00:00:00 2001 From: Alex Date: Fri, 21 Aug 2026 00:19:26 -0700 Subject: [PATCH 3/8] Self-review: negotiate before room join; clear locals on pipeline close --- pkg/sip/inbound.go | 6 +++--- pkg/sip/media_port.go | 2 ++ 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/pkg/sip/inbound.go b/pkg/sip/inbound.go index 6262f732..72ff576b 100644 --- a/pkg/sip/inbound.go +++ b/pkg/sip/inbound.go @@ -1035,13 +1035,13 @@ func (c *inboundCall) handleInvite(ctx context.Context, tid traceid.ID, req *sip if pinPrompt { status = CallActive } - if err := c.joinRoom(ctx, disp.Room, status); err != nil { - return fmt.Errorf("failed joining room: %w", err) - } answerData, err = c.negotiateMedia(rawSDP) if err != nil { return rejectMedia(err) } + if err := c.joinRoom(ctx, disp.Room, status); err != nil { + return fmt.Errorf("failed joining room: %w", err) + } // Publish our own track. if err := c.publishTrack(disp.EnabledFeatures, disp.FeatureFlags); err != nil { c.log().Errorw("Cannot publish track", err) diff --git a/pkg/sip/media_port.go b/pkg/sip/media_port.go index f1312ca0..9369acef 100644 --- a/pkg/sip/media_port.go +++ b/pkg/sip/media_port.go @@ -920,6 +920,8 @@ func (p *mediaPort) configure(c *sdp.MediaConfig, localSDP []byte) error { } p.closePipelineLocked() + audioToPort = nil + dtmfToPort = nil p.port.stopDiscarding() // Needs readDeadline. Must be ahead of Reopen() and NewMediaPortPipeline() p.port.Reopen() // Allow reads from socket again From 1c00fccc0a9b4a5f29396417caaea0a49392c582 Mon Sep 17 00:00:00 2001 From: Alex Date: Fri, 21 Aug 2026 00:40:13 -0700 Subject: [PATCH 4/8] Self-review: Warn Downgrade; Drop dead param; use const; reinvite tag --- pkg/sip/inbound.go | 4 ++-- pkg/sip/media_pipeline.go | 2 +- pkg/sip/media_port.go | 31 +++++++++++--------------- pkg/sip/media_port_negotiation_test.go | 16 ++++++------- pkg/sip/media_port_test.go | 24 +++++++++++--------- pkg/sip/outbound.go | 2 +- pkg/stats/monitor.go | 10 +++++---- 7 files changed, 44 insertions(+), 45 deletions(-) diff --git a/pkg/sip/inbound.go b/pkg/sip/inbound.go index 72ff576b..7b9e6fa8 100644 --- a/pkg/sip/inbound.go +++ b/pkg/sip/inbound.go @@ -1222,7 +1222,7 @@ func (c *inboundCall) negotiateMedia(offerData []byte) ([]byte, error) { c.mon.SDPSize(len(offerData), true) c.log().Debugw("SDP offer", "sdp", string(offerData)) - answerData, err := c.media.GenerateAnswer(offerData, false) + answerData, err := c.media.GenerateAnswer(offerData) if err != nil { return nil, err } @@ -1583,7 +1583,7 @@ func (c *inboundCall) updateRemoteFromSDP(body []byte) error { if mp == nil { return nil } - _, err := mp.GenerateAnswer(body, false) + _, err := mp.GenerateAnswer(body) return err } diff --git a/pkg/sip/media_pipeline.go b/pkg/sip/media_pipeline.go index 2112f916..b9373f96 100644 --- a/pkg/sip/media_pipeline.go +++ b/pkg/sip/media_pipeline.go @@ -445,7 +445,7 @@ func (w *dtmfOutWriter) WriteSample(sample *livekit.SipDTMF) error { digits = string([]byte{digit}) } else if sample.Code > 0 { // We can't distinguish between a code0 and no code, but better have something here - w.log.Warnw("code payload detected, ignored due to explicit digits", nil, "code", sample.Code, "digits", sample.Digit) + w.log.Debugw("code payload detected, ignored due to explicit digits", "code", sample.Code, "digits", sample.Digit) } w.mu.Lock() diff --git a/pkg/sip/media_port.go b/pkg/sip/media_port.go index 9369acef..801c8848 100644 --- a/pkg/sip/media_port.go +++ b/pkg/sip/media_port.go @@ -46,6 +46,7 @@ const ( defaultMediaTimeoutInitial = 30 * time.Second dstChangePrintInterval = 10 * 1000 * 1000 * 1000 // 10 seconds, in nanoseconds srcChangePrintInterval = dstChangePrintInterval + holdEnabled = false // Disabled in current code ) var ErrRenegotiationDisabled = errors.New("renegotiation is not supported") @@ -416,12 +417,10 @@ type MediaPort interface { GenerateOffer() ([]byte, error) // GenerateAnswer returns an encoded SDP answer for the given offer. - // activateTimeout - If set to true, resets media timeout to defaults, - // starting from time of function call. This is designed to - // be set to false when not immediately sending the answer back. + // This does not arm the media timeout, use SetTimeout to do so. // // SIDE EFFECT: May cause a rebuild of the pipeline. - GenerateAnswer(offer []byte, activateTimeout bool) ([]byte, error) + GenerateAnswer(offer []byte) ([]byte, error) // ProcessAnswer processes an encoded SDP answer from the remote client. Returns an // error if the answer is invalid, the offer has not yet been generated, or @@ -718,11 +717,11 @@ func (p *mediaPort) RemoteAddr() netip.AddrPort { // Reported for inbound (SetOffer) only since outbound (SetAnswer) only contains the // codec picked by the end user, and not what they actually support -func (p *mediaPort) reportPeerCodecs(d sdp.MediaDesc) { +func (p *mediaPort) reportPeerCodecs(d sdp.MediaDesc, reinvite bool) { if p.mon == nil { return } - p.mon.PeerSDP(peerCodecNames(d)) + p.mon.PeerSDP(peerCodecNames(d), reinvite) } // Plumbing @@ -769,7 +768,7 @@ func (p *mediaPort) GenerateOffer() ([]byte, error) { return offer.SDP.Marshal() } -func (p *mediaPort) GenerateAnswer(offerData []byte, activateTimeout bool) ([]byte, error) { +func (p *mediaPort) GenerateAnswer(offerData []byte) ([]byte, error) { if len(offerData) == 0 { return p.GetLocalSDP() } @@ -778,11 +777,10 @@ func (p *mediaPort) GenerateAnswer(offerData []byte, activateTimeout bool) ([]by if err != nil { return nil, SDPError{Err: err} } - p.reportPeerCodecs(offer.MediaDesc) - p.mu.Lock() - p.offer = offer - p.mu.Unlock() - + p.mu.RLock() + isReinvite := p.offer != nil + p.mu.RUnlock() + p.reportPeerCodecs(offer.MediaDesc, isReinvite) answer, mc, err := offer.Answer(p.externalIP, p.Port(), p.encryption, sdp.WithLocalProfiles(p.localCrypto)) if err != nil { return nil, SDPError{Err: err} @@ -796,9 +794,6 @@ func (p *mediaPort) GenerateAnswer(offerData []byte, activateTimeout bool) ([]by if err != nil { return nil, err } - if activateTimeout { - p.SetTimeout(p.opts.MediaTimeoutInitial, p.opts.MediaTimeout) - } return answerData, nil } @@ -869,6 +864,8 @@ func (p *mediaPort) configure(c *sdp.MediaConfig, localSDP []byte) error { p.mu.Lock() // No concurrent rebuilding of the pipeline defer p.mu.Unlock() + p.offer = nil + if p.closed.IsBroken() { return errors.New("media is already closed") } @@ -901,7 +898,7 @@ func (p *mediaPort) configure(c *sdp.MediaConfig, localSDP []byte) error { // maybe gate these on timers being active on the session to prevent dud calls hold = c.PeerDirection == psdp.DirectionSendOnly } - if false && hold { // Disabled in current code + if holdEnabled && hold { audioToPort = nil dtmfToPort = nil zero := netip.IPv4Unspecified() @@ -915,7 +912,6 @@ func (p *mediaPort) configure(c *sdp.MediaConfig, localSDP []byte) error { if changeSetSummary != changeSetNew { // Explicitly disable renegotiation for now // Compatibility to todays behavior: return 200 OK, but don't reconfigure the pipeline - p.offer = nil return nil } @@ -950,7 +946,6 @@ func (p *mediaPort) configure(c *sdp.MediaConfig, localSDP []byte) error { p.localSDP = localSDP // TODO: Move to end of function when reconfiguring is supported } - p.offer = nil // Pipeline build done, can now proceed to offer anew p.negotiated = c return nil } diff --git a/pkg/sip/media_port_negotiation_test.go b/pkg/sip/media_port_negotiation_test.go index 45eb80b6..eeac93ff 100644 --- a/pkg/sip/media_port_negotiation_test.go +++ b/pkg/sip/media_port_negotiation_test.go @@ -171,7 +171,7 @@ func TestMediaPortCodecSet(t *testing.T) { // Peer offers both, only PCMA is enabled here. offer := sdpWithMedia("m=audio 5004 RTP/AVP 0 8", "a=rtpmap:0 PCMU/8000", "a=rtpmap:8 PCMA/8000") - answerData, err := m.GenerateAnswer(offer, true) + answerData, err := m.GenerateAnswer(offer) require.NoError(t, err) assert.Equal(t, g711.ALawSDPNameAndRate, answerCodec(t, answerData)) }) @@ -180,7 +180,7 @@ func TestMediaPortCodecSet(t *testing.T) { m := newLocked(t, g711.ALawSDPNameAndRate) offer := sdpWithMedia("m=audio 5004 RTP/AVP 0", "a=rtpmap:0 PCMU/8000") - _, err := m.GenerateAnswer(offer, true) + _, err := m.GenerateAnswer(offer) require.ErrorIs(t, err, sdp.ErrNoCommonMedia) }) } @@ -197,16 +197,16 @@ func TestMediaPortRejectsDifferentCodecOffer(t *testing.T) { sdpB := sdpWithMedia("m=audio 5004 RTP/AVP 9", "a=rtpmap:9 G722/8000") // Offer codec A - answer, err := m.GenerateAnswer(sdpA, true) + answer, err := m.GenerateAnswer(sdpA) require.NoError(t, err) require.Equal(t, g711.ULawSDPNameAndRate, answerCodec(t, answer)) // Attempt to offer only codec B, expect failure - answer, err = m.GenerateAnswer(sdpB, true) + answer, err = m.GenerateAnswer(sdpB) require.ErrorIs(t, err, ErrRenegotiationDisabled) // Offer codec A again, expect success - answer, err = m.GenerateAnswer(sdpA, true) + answer, err = m.GenerateAnswer(sdpA) require.NoError(t, err) require.Equal(t, g711.ULawSDPNameAndRate, answerCodec(t, answer)) } @@ -308,7 +308,7 @@ func TestMediaPortHold(t *testing.T) { // m2 re-INVITEs with the hold form of its offer. base, err := m2.GenerateOffer() require.NoError(t, err) - _, err = m1.GenerateAnswer([]byte(tc.hold(t, string(base))), true) + _, err = m1.GenerateAnswer([]byte(tc.hold(t, string(base)))) require.NoError(t, err) // m1 no longer sends: no destination to write to, and the room-facing @@ -329,7 +329,7 @@ func TestMediaPortHold(t *testing.T) { requireAudioFlows(t, m2, recv1) // Resume with the original offer. - _, err = m1.GenerateAnswer(base, true) + _, err = m1.GenerateAnswer(base) require.NoError(t, err) dst = m1.port.dst.Load() @@ -398,7 +398,7 @@ func TestMediaPortEncryptionPolicy(t *testing.T) { require.NoError(t, err) offerData, err := offer.SDP.Marshal() require.NoError(t, err) - answerData, err := mp.GenerateAnswer(offerData, true) + answerData, err := mp.GenerateAnswer(offerData) if err != nil { return nil, err } diff --git a/pkg/sip/media_port_test.go b/pkg/sip/media_port_test.go index 34cd1755..1cfea02b 100644 --- a/pkg/sip/media_port_test.go +++ b/pkg/sip/media_port_test.go @@ -261,23 +261,23 @@ func TestMediaPortUpdateRemote(t *testing.T) { require.False(t, mp.RemoteAddr().IsValid(), "RemoteAddr should be invalid before any offer") addr := netip.MustParseAddrPort("9.8.7.6:12345") - _, err := mp.GenerateAnswer(offerAt(t, addr), true) + _, err := mp.GenerateAnswer(offerAt(t, addr)) require.NoError(t, err) require.Equal(t, addr, mp.RemoteAddr(), "GenerateAnswer should set RemoteAddr from the offer") // Body-less re-INVITE: empty offer returns the local SDP and must not change dest. - _, err = mp.GenerateAnswer(nil, true) + _, err = mp.GenerateAnswer(nil) require.NoError(t, err) require.Equal(t, addr, mp.RemoteAddr(), "empty offer should not change RemoteAddr") // Hold form c=0.0.0.0 must not clobber dest once media is established. - _, err = mp.GenerateAnswer(offerAt(t, netip.MustParseAddrPort("0.0.0.0:12345")), true) + _, err = mp.GenerateAnswer(offerAt(t, netip.MustParseAddrPort("0.0.0.0:12345"))) require.NoError(t, err) require.Equal(t, addr, mp.RemoteAddr(), "offer with unspecified addr should not change RemoteAddr") // successful re-INVITE update addr = netip.MustParseAddrPort("10.10.10.10:54321") - _, err = mp.GenerateAnswer(offerAt(t, addr), true) + _, err = mp.GenerateAnswer(offerAt(t, addr)) require.NoError(t, err) require.Equal(t, addr, mp.RemoteAddr(), "re-INVITE offer should update RemoteAddr") } @@ -294,7 +294,7 @@ func TestMediaPortReinviteSameCrypto(t *testing.T) { addr := netip.MustParseAddrPort("9.8.7.6:12345") offer := offerAtEnc(t, addr, sdp.EncryptionRequire) - _, err := mp.GenerateAnswer(offer, true) + _, err := mp.GenerateAnswer(offer) require.NoError(t, err) require.Equal(t, addr, mp.RemoteAddr()) @@ -309,7 +309,7 @@ func TestMediaPortReinviteSameCrypto(t *testing.T) { require.NotEmpty(t, localSDP) // Same offer bytes: NewOfferWith would generate a new peer key. - _, err = mp.GenerateAnswer(offer, true) + _, err = mp.GenerateAnswer(offer) require.NoError(t, err, "re-INVITE with the same offer must be accepted") require.Equal(t, addr, mp.RemoteAddr(), "same offer must not change dest") require.Equal(t, localKey, mp.negotiated.Crypto.Keys.LocalMasterKey, "local master key must not change") @@ -359,10 +359,12 @@ func negotiate(t testing.TB, m1, m2 *mediaPort) []byte { offerData, err := m1.GenerateOffer() require.NoError(t, err) - answerData, err := m2.GenerateAnswer(offerData, true) + answerData, err := m2.GenerateAnswer(offerData) require.NoError(t, err) require.NoError(t, m1.ProcessAnswer(answerData)) + + m2.SetTimeout(m2.opts.MediaTimeoutInitial, m2.opts.MediaTimeout) return answerData } @@ -571,7 +573,7 @@ func TestPipelineChains(t *testing.T) { require.NoError(t, err) answerData, err := offer.SDP.Marshal() require.NoError(t, err) - _, err = mp.GenerateAnswer(answerData, true) + _, err = mp.GenerateAnswer(answerData) require.NoError(t, err) codecName := strings.Split(info.SDPName, "/")[0] @@ -881,7 +883,7 @@ func TestSetOfferReportsCodecsBeforeFailing(t *testing.T) { pcmuBefore := gatherCounter(t, offeredMetric, pcmu) offer := sdpWithMedia("m=audio 5004 RTP/AVP 96", "a=rtpmap:96 SPEEX/16000") - _, err := mp.GenerateAnswer(offer, true) + _, err := mp.GenerateAnswer(offer) require.ErrorIs(t, err, sdp.ErrNoCommonMedia) // Codecs that are not part of the internal set are classified as "other" @@ -903,7 +905,7 @@ func TestSetOfferReportsCodecsPerProvider(t *testing.T) { offer := sdpWithMedia("m=audio 5004 RTP/AVP 0 9", "a=rtpmap:0 PCMU/8000", "a=rtpmap:9 G722/8000") - _, err := mp.GenerateAnswer(offer, true) + _, err := mp.GenerateAnswer(offer) require.NoError(t, err) require.Equal(t, parsedBefore+1, gatherCounter(t, parsedMetric, parsed)) @@ -920,7 +922,7 @@ func TestSetOfferReportsUnknownProvider(t *testing.T) { before := gatherCounter(t, parsedMetric, parsed) offer := sdpWithMedia("m=audio 5004 RTP/AVP 0", "a=rtpmap:0 PCMU/8000") - _, err := mp.GenerateAnswer(offer, true) + _, err := mp.GenerateAnswer(offer) require.NoError(t, err) require.Equal(t, before+1, gatherCounter(t, parsedMetric, parsed)) diff --git a/pkg/sip/outbound.go b/pkg/sip/outbound.go index 6bca7722..2cf4f36a 100644 --- a/pkg/sip/outbound.go +++ b/pkg/sip/outbound.go @@ -537,7 +537,7 @@ func (c *outboundCall) updateRemoteFromSDP(body []byte) error { if mp == nil { return nil } - _, err := mp.GenerateAnswer(body, false) + _, err := mp.GenerateAnswer(body) return err } diff --git a/pkg/stats/monitor.go b/pkg/stats/monitor.go index d016c1b8..bfeb6a62 100644 --- a/pkg/stats/monitor.go +++ b/pkg/stats/monitor.go @@ -16,6 +16,7 @@ package stats import ( "errors" + "strconv" "sync/atomic" "time" @@ -250,7 +251,7 @@ func (m *Monitor) Start(conf *config.Config) error { Name: "sdp_parsed_total", Help: "Number of SDP bodies parsed successfully during SDP negotiation", ConstLabels: prometheus.Labels{"node_id": conf.NodeID}, - }, []string{"dir", "provider"})) + }, []string{"dir", "provider", "reinvite"})) m.codecOffered = mustRegister(m, prometheus.NewCounterVec(prometheus.CounterOpts{ Namespace: "livekit", @@ -258,7 +259,7 @@ func (m *Monitor) Start(conf *config.Config) error { Name: "codec_offered_total", Help: "Number of SDP bodies that advertised a given audio codec", ConstLabels: prometheus.Labels{"node_id": conf.NodeID}, - }, []string{"dir", "provider", "codec"})) + }, []string{"dir", "provider", "codec", "reinvite"})) m.nodeAvailable = mustRegister(m, prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Namespace: "livekit", @@ -526,14 +527,15 @@ func (c *CallMonitor) StageDurTimer(stage string) func() time.Duration { // PeerSDP increments SDP count and each individual codec from the SDP body. // Should be called before codec selection such that failed negotiations are still counted -func (c *CallMonitor) PeerSDP(names []string) { +func (c *CallMonitor) PeerSDP(names []string, reinvite bool) { provider := c.providerLabel() - c.m.sdpParsed.With(prometheus.Labels{"dir": c.dir, "provider": provider}).Inc() + c.m.sdpParsed.With(prometheus.Labels{"dir": c.dir, "provider": provider, "reinvite": strconv.FormatBool(reinvite)}).Inc() for _, name := range names { c.m.codecOffered.With(prometheus.Labels{ "dir": c.dir, "provider": provider, "codec": name, + "reinvite": strconv.FormatBool(reinvite), }).Inc() } } From e577bdd4f4f3e0447b558201a4d3d8cd8fc05f34 Mon Sep 17 00:00:00 2001 From: Alex Date: Fri, 21 Aug 2026 00:47:01 -0700 Subject: [PATCH 5/8] Self-review: Typos --- pkg/sip/media_pipeline.go | 12 ++++-------- pkg/sip/media_port.go | 10 +++++----- 2 files changed, 9 insertions(+), 13 deletions(-) diff --git a/pkg/sip/media_pipeline.go b/pkg/sip/media_pipeline.go index b9373f96..cc1c2f21 100644 --- a/pkg/sip/media_pipeline.go +++ b/pkg/sip/media_pipeline.go @@ -152,19 +152,15 @@ func (p *mediaPortPipeline) init( return nil } -// Construct the Audio and optionally DTMF pipline from SIP RTP to LK PCM, in reverse order. +// Construct the Audio and optionally DTMF pipeline from SIP RTP to LK PCM, in reverse order. func (p *mediaPortPipeline) setupInput(mc *sdp.MediaConfig, audioToRoom msdk.PCM16Writer, dtmfToRoom msdk.WriteCloser[*livekit.SipDTMF]) error { var err error var inboundLatencyEntry atomic.Int64 sink := msdk.NopCloser(audioToRoom) // Prevent pipeline close from closing room sink = newLatencyPCMExit(sink, &inboundLatencyEntry, &p.conf.stats.LatencyInE2E) - codecInfo := mc.Audio.Codec.Info() sink = msdk.ResampleWriter(sink, codecInfo.SampleRate) - - if p.conf.stats != nil { - sink = newMediaWriterCount(sink, &p.conf.stats.AudioInFrames, &p.conf.stats.AudioInSamples) - } + sink = newMediaWriterCount(sink, &p.conf.stats.AudioInFrames, &p.conf.stats.AudioInSamples) if p.conf.opts.LogSignalChanges { sink, err = NewSignalLogger(p.conf.log, "input", sink) @@ -231,7 +227,7 @@ func (p *mediaPortPipeline) handleEventRTP(h *rtp.Header, payload []byte) error }) } -// Construct the Audio and optionally DTMF pipline from LK PCM to SIP RTP +// Construct the Audio and optionally DTMF pipeline from LK PCM to SIP RTP // Returns the insulated (nopCloser) connectors, and an error. func (p *mediaPortPipeline) setupOutput(mc *sdp.MediaConfig, incomingSampleRate int) error { p.rtpLoopWG.Go(p.rtpLoop) @@ -440,7 +436,7 @@ func (w *dtmfOutWriter) WriteSample(sample *livekit.SipDTMF) error { if len(digits) == 0 { digit := dtmf.CodeToChar(byte(sample.Code)) if digit == 0 { - return fmt.Errorf("code %d not supoported", sample.Code) + return fmt.Errorf("code %d not supported", sample.Code) } digits = string([]byte{digit}) } else if sample.Code > 0 { diff --git a/pkg/sip/media_port.go b/pkg/sip/media_port.go index 801c8848..995126ce 100644 --- a/pkg/sip/media_port.go +++ b/pkg/sip/media_port.go @@ -406,9 +406,9 @@ type MediaPort interface { // GetOutboundDTMFWriter returns the LK room -> SIP DTMF writer. GetOutboundDTMFWriter() msdk.WriteCloser[*livekit.SipDTMF] - // WriteInboundAudioTo tells the MediaSegment where to write inbound SIP audio. + // WriteInboundAudioTo tells port where to write inbound SIP audio. WriteInboundAudioTo(w msdk.PCM16Writer) msdk.PCM16Writer - // WriteInboundDTMFTo tells the MediaSegment where to write inbound SIP DTMF. + // WriteInboundDTMFTo tells port where to write inbound SIP DTMF. WriteInboundDTMFTo(w msdk.WriteCloser[*livekit.SipDTMF]) msdk.WriteCloser[*livekit.SipDTMF] // If there is no offer, this generates an offer. @@ -692,8 +692,8 @@ func (p *mediaPort) Close() { } p.audioIn.Close() // Propagate Close() to onwards to room p.dtmfIn.Close() // Propagate Close() to onwards to room - p.audioOut.Close() // Pipeline insulated, but close switch - p.dtmfOut.Close() // Pipeline insulated, but close switch + p.audioOut.Close() // No-op, but do anyway + p.dtmfOut.Close() // No-op, but do anyway }) } @@ -983,7 +983,7 @@ func NewChangeSetSummary(current, new *sdp.MediaConfig) changeSetSummary { if a != b { changeSetSummary |= changeSetCrypto } - } else { // Prodile exists on both + } else { // Profile exists on both if a.Profile != b.Profile || !bytes.Equal(a.Keys.LocalMasterKey, b.Keys.LocalMasterKey) || !bytes.Equal(a.Keys.LocalMasterSalt, b.Keys.LocalMasterSalt) || From d2793dbddb123bc1e9db8562dbebe491fbf7b392 Mon Sep 17 00:00:00 2001 From: Alex Date: Fri, 21 Aug 2026 01:02:17 -0700 Subject: [PATCH 6/8] Bah --- pkg/sip/media_port.go | 4 +--- pkg/sip/media_port_negotiation_test.go | 2 +- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/pkg/sip/media_port.go b/pkg/sip/media_port.go index 995126ce..f9cc950b 100644 --- a/pkg/sip/media_port.go +++ b/pkg/sip/media_port.go @@ -49,8 +49,6 @@ const ( holdEnabled = false // Disabled in current code ) -var ErrRenegotiationDisabled = errors.New("renegotiation is not supported") - type PortStatsSnapshot struct { Streams uint64 `json:"streams"` Packets uint64 `json:"packets"` @@ -778,7 +776,7 @@ func (p *mediaPort) GenerateAnswer(offerData []byte) ([]byte, error) { return nil, SDPError{Err: err} } p.mu.RLock() - isReinvite := p.offer != nil + isReinvite := p.negotiated != nil p.mu.RUnlock() p.reportPeerCodecs(offer.MediaDesc, isReinvite) answer, mc, err := offer.Answer(p.externalIP, p.Port(), p.encryption, sdp.WithLocalProfiles(p.localCrypto)) diff --git a/pkg/sip/media_port_negotiation_test.go b/pkg/sip/media_port_negotiation_test.go index eeac93ff..a33292e1 100644 --- a/pkg/sip/media_port_negotiation_test.go +++ b/pkg/sip/media_port_negotiation_test.go @@ -203,7 +203,7 @@ func TestMediaPortRejectsDifferentCodecOffer(t *testing.T) { // Attempt to offer only codec B, expect failure answer, err = m.GenerateAnswer(sdpB) - require.ErrorIs(t, err, ErrRenegotiationDisabled) + require.ErrorIs(t, err, sdp.ErrNoCommonMedia) // Offer codec A again, expect success answer, err = m.GenerateAnswer(sdpA) From 5fcb4813d0166b6fc4a831ea04d43512f430a10d Mon Sep 17 00:00:00 2001 From: Alex Date: Fri, 21 Aug 2026 01:04:11 -0700 Subject: [PATCH 7/8] one more typo --- pkg/sip/media_port.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/sip/media_port.go b/pkg/sip/media_port.go index f9cc950b..6882030f 100644 --- a/pkg/sip/media_port.go +++ b/pkg/sip/media_port.go @@ -909,7 +909,7 @@ func (p *mediaPort) configure(c *sdp.MediaConfig, localSDP []byte) error { if changeSetSummary.shouldReconfigure() { if changeSetSummary != changeSetNew { // Explicitly disable renegotiation for now - // Compatibility to todays behavior: return 200 OK, but don't reconfigure the pipeline + // Compatibility to today's behavior: return 200 OK, but don't reconfigure the pipeline return nil } From fcc202107260c6c9649b199dd6d92923c892c60f Mon Sep 17 00:00:00 2001 From: Alex Date: Fri, 21 Aug 2026 10:42:36 -0700 Subject: [PATCH 8/8] PR comment - print close errors: --- pkg/sip/media_port.go | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/pkg/sip/media_port.go b/pkg/sip/media_port.go index 6882030f..10875886 100644 --- a/pkg/sip/media_port.go +++ b/pkg/sip/media_port.go @@ -678,20 +678,27 @@ func (p *mediaPort) Close() { p.closed.Once(func() { defer p.stats.Closed.Store(true) + logError := func(comp string, err error) { + if err != nil { + p.log.Errorw("error closing media port", err, "component", comp) + } + } + p.mu.Lock() defer p.mu.Unlock() p.closePipelineLocked() - p.port.Close() + logError("port", p.port.Close()) conn := p.port.unwrap() if uc, ok := conn.(*net.UDPConn); ok { go DrainPort(p.log, uc, p.opts.DrainingIdleTimeout, p.opts.DrainingDuration, nil) } else { - _ = conn.Close() + logError("conn", conn.Close()) } - p.audioIn.Close() // Propagate Close() to onwards to room - p.dtmfIn.Close() // Propagate Close() to onwards to room - p.audioOut.Close() // No-op, but do anyway - p.dtmfOut.Close() // No-op, but do anyway + + logError("audioIn", p.audioIn.Close()) // Propagate Close() to onwards to room + logError("dtmfIn", p.dtmfIn.Close()) // Propagate Close() to onwards to room + logError("audioOut", p.audioOut.Close()) // No-op, but do anyway + logError("dtmfOut", p.dtmfOut.Close()) // No-op, but do anyway }) }