Skip to content

Commit afd7dd4

Browse files
authored
Add support for symmetric RTP (#624)
1 parent 338f4eb commit afd7dd4

5 files changed

Lines changed: 65 additions & 3 deletions

File tree

pkg/config/config.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,7 @@ type Config struct {
102102

103103
MediaTimeout time.Duration `yaml:"media_timeout"`
104104
MediaTimeoutInitial time.Duration `yaml:"media_timeout_initial"`
105+
SymmetricRTP bool `yaml:"symmetric_rtp"`
105106
Codecs map[string]bool `yaml:"codecs"`
106107

107108
// HideInboundPort controls how SIP endpoint responds to unverified inbound requests.

pkg/sip/inbound.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -980,6 +980,7 @@ func (c *inboundCall) runMediaConn(tid traceid.ID, offerData []byte, enc livekit
980980
Ports: conf.RTPPort,
981981
MediaTimeoutInitial: c.s.conf.MediaTimeoutInitial,
982982
MediaTimeout: c.s.conf.MediaTimeout,
983+
SymmetricRTP: conf.SymmetricRTP,
983984
EnableJitterBuffer: c.jitterBuf,
984985
LogSignalChanges: logSignalChanges,
985986
Stats: &c.stats.Port,

pkg/sip/media_port.go

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -198,15 +198,16 @@ type UDPConn interface {
198198
WriteToUDPAddrPort(b []byte, addr netip.AddrPort) (int, error)
199199
}
200200

201-
func newUDPConn(log logger.Logger, conn UDPConn) *udpConn {
202-
return &udpConn{UDPConn: conn, log: log, stopped: make(chan struct{})}
201+
func newUDPConn(log logger.Logger, conn UDPConn, symmetricRTP bool) *udpConn {
202+
return &udpConn{UDPConn: conn, log: log, stopped: make(chan struct{}), symmetricRTP: symmetricRTP}
203203
}
204204

205205
type udpConn struct {
206206
UDPConn
207207
stopping core.Fuse
208208
stopped chan struct{}
209209
log logger.Logger
210+
symmetricRTP bool
210211
src atomic.Pointer[netip.AddrPort]
211212
dst atomic.Pointer[netip.AddrPort]
212213
}
@@ -239,6 +240,12 @@ func (c *udpConn) Read(b []byte) (n int, err error) {
239240
} else if *prev != addr {
240241
c.log.Infow("changing media source", "addr", addr.String())
241242
}
243+
if c.symmetricRTP {
244+
dst := c.dst.Load()
245+
if dst == nil || !dst.IsValid() || *dst != addr {
246+
c.SetDst(addr)
247+
}
248+
}
242249
return n, err
243250
}
244251

@@ -302,6 +309,7 @@ type MediaOptions struct {
302309
Ports rtcconfig.PortRange
303310
MediaTimeoutInitial time.Duration
304311
MediaTimeout time.Duration
312+
SymmetricRTP bool
305313
Stats *PortStats
306314
EnableJitterBuffer bool
307315
NoInputResample bool
@@ -348,7 +356,7 @@ func NewMediaPortWith(tid traceid.ID, log logger.Logger, mon *stats.CallMonitor,
348356
timeoutResetTick: make(chan time.Duration, 1),
349357
jitterEnabled: opts.EnableJitterBuffer,
350358
logSignalChanges: opts.LogSignalChanges,
351-
port: newUDPConn(log, conn),
359+
port: newUDPConn(log, conn, opts.SymmetricRTP),
352360
audioOut: msdk.NewSwitchWriter(sampleRate),
353361
audioIn: msdk.NewSwitchWriter(inSampleRate),
354362
stats: opts.Stats,

pkg/sip/media_port_test.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -553,3 +553,54 @@ func TestMediaTimeout(t *testing.T) {
553553
}
554554
})
555555
}
556+
557+
func TestSymmetricRTP(t *testing.T) {
558+
t.Run("disabled", func(t *testing.T) {
559+
m1, m2 := newMediaPair(t, &MediaOptions{SymmetricRTP: false}, nil)
560+
dstPtr := m1.port.dst.Load()
561+
require.NotNil(t, dstPtr)
562+
dst := *dstPtr
563+
require.True(t, dst.IsValid())
564+
565+
c2 := m2.port.UDPConn.(*testUDPConn)
566+
newAddr := netip.AddrPortFrom(newIP("9.9.9.9"), 9999)
567+
c2.addr = newAddr
568+
569+
err := m2.GetAudioWriter().WriteSample(msdk.PCM16Sample{0, 0})
570+
require.NoError(t, err)
571+
572+
select {
573+
case <-m1.Received():
574+
case <-time.After(time.Second):
575+
t.Fatal("no media received")
576+
}
577+
578+
curDstPtr := m1.port.dst.Load()
579+
require.NotNil(t, curDstPtr)
580+
require.Equal(t, dst, *curDstPtr)
581+
})
582+
583+
t.Run("enabled", func(t *testing.T) {
584+
m1, m2 := newMediaPair(t, &MediaOptions{SymmetricRTP: true}, nil)
585+
dstPtr := m1.port.dst.Load()
586+
require.NotNil(t, dstPtr)
587+
require.True(t, dstPtr.IsValid())
588+
589+
c2 := m2.port.UDPConn.(*testUDPConn)
590+
newAddr := netip.AddrPortFrom(newIP("9.9.9.9"), 9999)
591+
c2.addr = newAddr
592+
593+
err := m2.GetAudioWriter().WriteSample(msdk.PCM16Sample{0, 0})
594+
require.NoError(t, err)
595+
596+
select {
597+
case <-m1.Received():
598+
case <-time.After(time.Second):
599+
t.Fatal("no media received")
600+
}
601+
602+
curDstPtr := m1.port.dst.Load()
603+
require.NotNil(t, curDstPtr)
604+
require.Equal(t, newAddr, *curDstPtr)
605+
})
606+
}

pkg/sip/outbound.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,7 @@ func (c *Client) newCall(ctx context.Context, tid traceid.ID, conf *config.Confi
144144
Ports: conf.RTPPort,
145145
MediaTimeoutInitial: c.conf.MediaTimeoutInitial,
146146
MediaTimeout: c.conf.MediaTimeout,
147+
SymmetricRTP: c.conf.SymmetricRTP,
147148
EnableJitterBuffer: call.jitterBuf,
148149
LogSignalChanges: signalLoggingEnabled,
149150
Stats: &call.stats.Port,

0 commit comments

Comments
 (0)