diff --git a/internal/transport/vp8channel/transport.go b/internal/transport/vp8channel/transport.go index 9f4e455..dc3ffdc 100644 --- a/internal/transport/vp8channel/transport.go +++ b/internal/transport/vp8channel/transport.go @@ -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 { diff --git a/internal/transport/vp8channel/transport_unit_test.go b/internal/transport/vp8channel/transport_unit_test.go index ffb8213..2a37211 100644 --- a/internal/transport/vp8channel/transport_unit_test.go +++ b/internal/transport/vp8channel/transport_unit_test.go @@ -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) + } + } +}