fix(vp8channel): reorder RTP and inject server keyframes (issue #95)

This commit is contained in:
neuronori
2026-06-25 16:07:19 +00:00
parent 2a9b1ec867
commit cd3d568fd5
2 changed files with 185 additions and 6 deletions
+111 -6
View File
@@ -885,6 +885,90 @@ func (p *streamTransport) drainTrack(track *webrtc.TrackRemote) {
}
}
// reorderWindow bounds how many out-of-order RTP packets the reorder buffer
// holds while waiting for a gap to fill. Real SFUs reorder within a handful of
// packets; once this many newer packets pile up behind a hole, the missing
// sequence is treated as genuinely lost and we advance, so a truly dropped
// packet cannot stall delivery indefinitely.
const reorderWindow = 256
// seqLess reports whether RTP sequence a precedes b using wrap-around aware
// comparison (RFC 1982 serial arithmetic on uint16).
func seqLess(a, b uint16) bool {
// bit15 of the wrap-around difference is the serial sign bit: set means a
// precedes b. Avoids a signed conversion gosec flags as overflow.
return (a-b)&0x8000 != 0
}
// reorderBuffer restores RTP sequence order before frame assembly. The SFU may
// deliver packets out of order or drop them; feeding that stream straight into
// the strict contiguity check in processRTPPacket made every reorder look like
// loss and discarded whole frames (issue #95: ~80-90% of VP8 frames dropped on
// a live SFU). Buffering by sequence number and draining in order means only
// genuine loss produces a gap.
type reorderBuffer struct {
pkts map[uint16]*rtp.Packet
nextSeq uint16
started bool
}
func newReorderBuffer() *reorderBuffer {
return &reorderBuffer{pkts: make(map[uint16]*rtp.Packet, reorderWindow)}
}
// push adds pkt and returns any packets now deliverable in strict sequence
// order. The caller reuses its read buffer across packets, so the payload is
// copied before buffering.
func (b *reorderBuffer) push(pkt *rtp.Packet) []*rtp.Packet {
if !b.started {
b.started = true
b.nextSeq = pkt.SequenceNumber
}
// Drop packets older than our current position: already delivered, or
// skipped past as lost.
if seqLess(pkt.SequenceNumber, b.nextSeq) {
return nil
}
cp := &rtp.Packet{Header: pkt.Header}
cp.Payload = append([]byte(nil), pkt.Payload...)
b.pkts[pkt.SequenceNumber] = cp
// Holding a full window behind a hole means the head sequence is
// genuinely lost: skip forward to the oldest buffered packet.
if len(b.pkts) > reorderWindow {
b.skipToOldest()
}
return b.drain()
}
// drain pops contiguous packets starting at nextSeq.
func (b *reorderBuffer) drain() []*rtp.Packet {
var out []*rtp.Packet
for {
pkt, ok := b.pkts[b.nextSeq]
if !ok {
return out
}
out = append(out, pkt)
delete(b.pkts, b.nextSeq)
b.nextSeq++
}
}
// skipToOldest advances nextSeq to the lowest buffered sequence, abandoning a
// lost packet so drain can make progress.
func (b *reorderBuffer) skipToOldest() {
first := true
var oldest uint16
for seq := range b.pkts {
if first || seqLess(seq, oldest) {
oldest = seq
first = false
}
}
b.nextSeq = oldest
}
type vp8FrameState struct {
vp8Pkt codecs.VP8Packet
frameBuf []byte
@@ -941,6 +1025,7 @@ func (s *vp8FrameState) processRTPPacket(pkt *rtp.Packet) []byte {
func (p *streamTransport) readVP8Track(track *webrtc.TrackRemote) {
var state vp8FrameState
reorder := newReorderBuffer()
buf := make([]byte, rtpBufSize)
var rtpCount, frameCount int
@@ -958,13 +1043,16 @@ func (p *streamTransport) readVP8Track(track *webrtc.TrackRemote) {
continue
}
frame := state.processRTPPacket(pkt)
if frame == nil {
continue
// Restore sequence order before assembly so SFU reordering is not
// mistaken for loss.
for _, ordered := range reorder.push(pkt) {
frame := state.processRTPPacket(ordered)
if frame == nil {
continue
}
frameCount++
p.handleIncomingFrame(frame)
}
frameCount++
p.handleIncomingFrame(frame)
}
}
@@ -1106,11 +1194,28 @@ func (p *streamTransport) peerWriterPump(_ uint32, out chan []byte) {
ticker := time.NewTicker(p.frameInterval)
defer ticker.Stop()
// Inject a decodable VP8 keyframe on the same cadence writerLoop uses for
// the client->server path. The server's per-peer bulk path previously
// emitted only opaque KCP data frames, which never form a decodable VP8
// keyframe: the SFU's decoder times out (~40s without a keyframe) and stops
// forwarding the server's track to the subscriber. The client side was
// kept alive by writerLoop.forceKeepalive; the server side had no
// equivalent, so the server->client direction collapsed first while the
// client->server direction kept flowing (issue #95).
keyframeEvery := max(int((2*time.Second)/p.frameInterval), 1)
ticksSinceKeyframe := 0
for {
select {
case <-p.closeCh:
return
case <-ticker.C:
ticksSinceKeyframe++
if ticksSinceKeyframe >= keyframeEvery {
ticksSinceKeyframe = 0
hdr := p.epochHeader()
_ = p.writeSampleLocked(hdr[:])
}
select {
case frame, ok := <-out:
if !ok {
@@ -5,6 +5,7 @@ import (
"context"
"encoding/binary"
"errors"
"reflect"
"testing"
"time"
@@ -393,3 +394,76 @@ func TestHandleIncomingFrameEpochFilteringAndReconnect(t *testing.T) {
t.Fatalf("peer epoch not re-latched: got %d want 2", tr.peerEpoch.Load())
}
}
func seqList(pkts []*rtp.Packet) []uint16 {
out := make([]uint16, len(pkts))
for i, p := range pkts {
out[i] = p.SequenceNumber
}
return out
}
func TestReorderBufferRestoresOrderAndSurvivesLoss(t *testing.T) {
// In-order packets pass straight through.
b := newReorderBuffer()
got := make([]uint16, 0, 3)
for _, seq := range []uint16{100, 101, 102} {
got = append(got, seqList(b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: seq}}))...)
}
if !reflect.DeepEqual(got, []uint16{100, 101, 102}) {
t.Fatalf("in-order drain = %v, want [100 101 102]", got)
}
// A reordered packet is held until the gap fills, then both drain in order.
b = newReorderBuffer()
if out := b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: 10}}); !reflect.DeepEqual(seqList(out), []uint16{10}) {
t.Fatalf("first packet = %v, want [10]", seqList(out))
}
// 12 arrives before 11: must be buffered, nothing delivered yet.
if out := b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: 12}}); out != nil {
t.Fatalf("out-of-order packet drained early = %v, want nil", seqList(out))
}
// 11 fills the hole: 11 and 12 drain in order.
out := b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: 11}})
if !reflect.DeepEqual(seqList(out), []uint16{11, 12}) {
t.Fatalf("gap fill drain = %v, want [11 12]", seqList(out))
}
// Genuine loss: a full window piles up behind a hole, buffer skips the
// lost sequence rather than stalling forever.
b = newReorderBuffer()
_ = b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: 0}})
var delivered int
for i := 2; i <= reorderWindow+2; i++ {
seq := uint16(i & 0xffff)
delivered += len(b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: seq}}))
}
if delivered == 0 {
t.Fatal("buffer stalled on lost packet: nothing delivered after window overflow")
}
// Stale packets older than the current position are dropped.
b = newReorderBuffer()
_ = b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: 50}})
if out := b.push(&rtp.Packet{Header: rtp.Header{SequenceNumber: 49}}); out != nil {
t.Fatalf("stale packet delivered = %v, want nil", seqList(out))
}
}
func TestSeqLessWrapAround(t *testing.T) {
cases := []struct {
a, b uint16
want bool
}{
{1, 2, true},
{2, 1, false},
{65535, 0, true}, // wrap: 65535 precedes 0
{0, 65535, false}, // wrap: 0 follows 65535
{10, 10, false},
}
for _, c := range cases {
if got := seqLess(c.a, c.b); got != c.want {
t.Fatalf("seqLess(%d, %d) = %v, want %v", c.a, c.b, got, c.want)
}
}
}