Skip to content

Commit 54b228d

Browse files
committed
Add support for symmetric RTP
1 parent 271ef71 commit 54b228d

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
@@ -190,15 +190,16 @@ type UDPConn interface {
190190
WriteToUDPAddrPort(b []byte, addr netip.AddrPort) (int, error)
191191
}
192192

193-
func newUDPConn(log logger.Logger, conn UDPConn) *udpConn {
194-
return &udpConn{UDPConn: conn, log: log, stopped: make(chan struct{})}
193+
func newUDPConn(log logger.Logger, conn UDPConn, symmetricRTP bool) *udpConn {
194+
return &udpConn{UDPConn: conn, log: log, stopped: make(chan struct{}), symmetricRTP: symmetricRTP}
195195
}
196196

197197
type udpConn struct {
198198
UDPConn
199199
stopping core.Fuse
200200
stopped chan struct{}
201201
log logger.Logger
202+
symmetricRTP bool
202203
src atomic.Pointer[netip.AddrPort]
203204
dst atomic.Pointer[netip.AddrPort]
204205
}
@@ -231,6 +232,12 @@ func (c *udpConn) Read(b []byte) (n int, err error) {
231232
} else if *prev != addr {
232233
c.log.Infow("changing media source", "addr", addr.String())
233234
}
235+
if c.symmetricRTP {
236+
dst := c.dst.Load()
237+
if dst == nil || !dst.IsValid() || *dst != addr {
238+
c.SetDst(addr)
239+
}
240+
}
234241
return n, err
235242
}
236243

@@ -289,6 +296,7 @@ type MediaOptions struct {
289296
Ports rtcconfig.PortRange
290297
MediaTimeoutInitial time.Duration
291298
MediaTimeout time.Duration
299+
SymmetricRTP bool
292300
Stats *PortStats
293301
EnableJitterBuffer bool
294302
NoInputResample bool
@@ -335,7 +343,7 @@ func NewMediaPortWith(tid traceid.ID, log logger.Logger, mon *stats.CallMonitor,
335343
timeoutResetTick: make(chan time.Duration, 1),
336344
jitterEnabled: opts.EnableJitterBuffer,
337345
logSignalChanges: opts.LogSignalChanges,
338-
port: newUDPConn(log, conn),
346+
port: newUDPConn(log, conn, opts.SymmetricRTP),
339347
audioOut: msdk.NewSwitchWriter(sampleRate),
340348
audioIn: msdk.NewSwitchWriter(inSampleRate),
341349
stats: opts.Stats,

pkg/sip/media_port_test.go

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

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)