Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
101 changes: 101 additions & 0 deletions pkg/av1/rtp.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
package av1

import (
"time"

"github.com/AlexxIT/go2rtc/pkg/core"
"github.com/AlexxIT/go2rtc/pkg/h264"
"github.com/pion/rtp"
"github.com/pion/rtp/codecs"
)

// RTPDepay - depacketize AV1 RTP packets (https://aomediacodec.github.io/av1-rtp-spec/)
// into temporal units in the AV1 low overhead bitstream format (OBUs with obu_size_field).
// One output packet = one temporal unit (same convention as h264/h265 access units).
func RTPDepay(handler core.HandlerFunc) core.HandlerFunc {
depack := &codecs.AV1Depacketizer{}

buf := make([]byte, 0, 512*1024) // 512K
var seqNum uint16

return func(packet *rtp.Packet) {
// when we collect data into one buffer, we need to make sure
// that all of it falls into the same sequence
if len(buf) > 0 && packet.SequenceNumber-seqNum != 1 {
//log.Printf("broken AV1 sequence")
buf = buf[:0] // drop data
depack = &codecs.AV1Depacketizer{} // drop pending OBU fragment
return
}

seqNum = packet.SequenceNumber

obus, err := depack.Unmarshal(packet.Payload)
if err != nil {
buf = buf[:0]
depack = &codecs.AV1Depacketizer{}
return
}

buf = append(buf, obus...)

// collect all OBUs for temporal unit
if !packet.Marker {
return
}

if len(buf) == 0 {
return
}

clone := *packet
clone.Version = h264.RTPPacketVersionAVC
clone.Payload = buf

buf = buf[:0]

handler(&clone)
}
}

// RTPPay - packetize AV1 temporal units (low overhead bitstream) into RTP packets
// sized for the given MTU, per the AV1 RTP specification.
func RTPPay(mtu uint16, handler core.HandlerFunc) core.HandlerFunc {
if mtu == 0 {
mtu = 1472
}

payloader := &codecs.AV1Payloader{}
sequencer := rtp.NewRandomSequencer()
mtu -= 12 // rtp.Header size

return func(packet *rtp.Packet) {
if packet.Version != h264.RTPPacketVersionAVC {
handler(packet)
return
}

payloads := payloader.Payload(mtu, packet.Payload)
last := len(payloads) - 1
for i, payload := range payloads {
// 4K AV1 keyframe temporal units run to hundreds of packets;
// blasting them in one tight loop overflows UDP socket buffers
// (~50% observed loss on a gigabit LAN). Pace large bursts —
// each sender runs on its own goroutine, so sleeping here only
// delays this consumer, never the producer.
if i > 0 && i%16 == 0 {
time.Sleep(2 * time.Millisecond)
}
clone := rtp.Packet{
Header: rtp.Header{
Version: 2,
Marker: i == last,
SequenceNumber: sequencer.NextSequenceNumber(),
Timestamp: packet.Timestamp,
},
Payload: payload,
}
handler(&clone)
}
}
}
16 changes: 16 additions & 0 deletions pkg/webrtc/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,13 @@ func newUDPMux(address string, filters *Filters) ice.UDPMux {
var muxes []ice.UDPMux
for _, addr := range addrs {
if ln, _ := net.ListenPacket(networkUDP, addr); ln != nil {
// large keyframes (4K H265/AV1) are written in bursts of hundreds
// of packets; the default socket buffer (net.core.wmem_default,
// often ~208KB) drops much of such a burst
if udpConn, ok := ln.(*net.UDPConn); ok {
_ = udpConn.SetWriteBuffer(4 << 20)
_ = udpConn.SetReadBuffer(1 << 20)
}
OnNewListener(ln)
mux := ice.NewUDPMuxDefault(ice.UDPMuxParams{UDPConn: ln})
muxes = append(muxes, mux)
Expand Down Expand Up @@ -306,6 +313,15 @@ func RegisterDefaultCodecs(m *webrtc.MediaEngine) error {
},
PayloadType: 100,
},
// Chrome 113+, Firefox 136+, Safari 18.4+
{
RTPCodecCapability: webrtc.RTPCodecCapability{
MimeType: webrtc.MimeTypeAV1,
ClockRate: 90000,
RTCPFeedback: videoRTCPFeedback,
},
PayloadType: 105,
},
} {
if err := m.RegisterCodec(codec, webrtc.RTPCodecTypeVideo); err != nil {
return err
Expand Down
9 changes: 9 additions & 0 deletions pkg/webrtc/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package webrtc
import (
"errors"

"github.com/AlexxIT/go2rtc/pkg/av1"
"github.com/AlexxIT/go2rtc/pkg/core"
"github.com/AlexxIT/go2rtc/pkg/h264"
"github.com/AlexxIT/go2rtc/pkg/h265"
Expand Down Expand Up @@ -63,6 +64,14 @@ func (c *Conn) AddTrack(media *core.Media, codec *core.Codec, track *core.Receiv
sender.Handler = h265.RepairAVCC(track.Codec, sender.Handler)
}

case core.CodecAV1:
// repacketize to browser-safe MTU (RTSP sources are often TCP-interleaved
// and can carry RTP packets larger than the WebRTC path MTU)
sender.Handler = av1.RTPPay(1200, sender.Handler)
if track.Codec.IsRTP() {
sender.Handler = av1.RTPDepay(sender.Handler)
}

case core.CodecPCMA, core.CodecPCMU, core.CodecPCM, core.CodecPCML:
// Fix audio quality https://github.com/AlexxIT/WebRTC/issues/500
// should be before ResampleToG711, because it will be called last
Expand Down