From afd7a59c58a7cb7bfbecb0f5fb683ac1344c3d9a Mon Sep 17 00:00:00 2001 From: David Chen Date: Wed, 10 Jun 2026 14:34:06 -0700 Subject: [PATCH] add flexFEC publisher to SFU support & test scripts --- config-sample.yaml | 3 + pkg/config/config.go | 3 + pkg/rtc/config.go | 16 +- pkg/rtc/mediaengine.go | 24 +- pkg/rtc/participant.go | 12 + pkg/rtc/transport.go | 64 +++ pkg/rtc/transport_test.go | 120 +++++- pkg/sfu/buffer/buffer.go | 181 +++++++++ pkg/sfu/buffer/buffer_base.go | 35 ++ pkg/sfu/buffer/factory.go | 30 ++ pkg/sfu/buffer/flexfec.go | 516 ++++++++++++++++++++++++ pkg/sfu/buffer/flexfec_test.go | 591 ++++++++++++++++++++++++++++ pkg/telemetry/prometheus/packets.go | 30 ++ scripts/fec/.gitignore | 2 + scripts/fec/README.md | 135 +++++++ scripts/fec/aggregate_fec.py | 212 ++++++++++ scripts/fec/plot_fec.py | 373 ++++++++++++++++++ scripts/fec/prom_poll.sh | 25 ++ scripts/fec/run_fec_test.sh | 344 ++++++++++++++++ scripts/fec/shape_linux.sh | 112 ++++++ scripts/fec/shape_macos.sh | 126 ++++++ scripts/fec/sweep_fec.sh | 152 +++++++ 22 files changed, 3098 insertions(+), 8 deletions(-) create mode 100644 pkg/sfu/buffer/flexfec.go create mode 100644 pkg/sfu/buffer/flexfec_test.go create mode 100644 scripts/fec/.gitignore create mode 100644 scripts/fec/README.md create mode 100755 scripts/fec/aggregate_fec.py create mode 100755 scripts/fec/plot_fec.py create mode 100755 scripts/fec/prom_poll.sh create mode 100755 scripts/fec/run_fec_test.sh create mode 100755 scripts/fec/shape_linux.sh create mode 100755 scripts/fec/shape_macos.sh create mode 100755 scripts/fec/sweep_fec.sh diff --git a/config-sample.yaml b/config-sample.yaml index f264e721e..0159fe5a5 100644 --- a/config-sample.yaml +++ b/config-sample.yaml @@ -120,6 +120,9 @@ rtc: # tcp_fallback_rtt_threshold: 0 # # migrate an established-but-lossy UDP connection to ICE/TCP or TURN/TLS. requires tcp_fallback_rtt_threshold > 0, default false # allow_udp_unstable_fallback: false + # # accept FlexFEC-03 from publishers and use it to recover lost packets without + # # waiting for an RTT (useful for high-loss links, e.g. cellular). default false + # enable_flexfec: false # # number of packets to buffer in the SFU for video, defaults to 500 # packet_buffer_size_video: 500 # # number of packets to buffer in the SFU for audio, defaults to 200 diff --git a/pkg/config/config.go b/pkg/config/config.go index 68d3e8a2c..b00ed6082 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -151,6 +151,9 @@ type RTCConfig struct { // enable rtp stream restart detection for published tracks EnableRTPStreamRestartDetection bool `yaml:"enable_rtp_stream_restart_detection,omitempty"` + + // accept FlexFEC-03 from publishers and use it to recover lost packets + EnableFlexFEC bool `yaml:"enable_flexfec,omitempty"` } type TURNServer struct { diff --git a/pkg/rtc/config.go b/pkg/rtc/config.go index 2f30060fd..a161b7d8f 100644 --- a/pkg/rtc/config.go +++ b/pkg/rtc/config.go @@ -37,6 +37,8 @@ type WebRTCConfig struct { Receiver ReceiverConfig Publisher DirectionConfig Subscriber DirectionConfig + + enableFlexFEC bool } type ReceiverConfig struct { @@ -57,6 +59,8 @@ type RTCPFeedbackConfig struct { type DirectionConfig struct { RTPHeaderExtension RTPHeaderExtensionConfig RTCPFeedback RTCPFeedbackConfig + // accept FlexFEC-03 in negotiation (publisher/receive direction only) + EnableFlexFEC bool } func NewWebRTCConfig(conf *config.Config) (*WebRTCConfig, error) { @@ -80,19 +84,23 @@ func NewWebRTCConfig(conf *config.Config) (*WebRTCConfig, error) { rtcConf.PacketBufferSizeAudio = rtcConf.PacketBufferSize } - return &WebRTCConfig{ + c := &WebRTCConfig{ WebRTCConfig: *webRTCConfig, Receiver: ReceiverConfig{ PacketBufferSizeVideo: rtcConf.PacketBufferSizeVideo, PacketBufferSizeAudio: rtcConf.PacketBufferSizeAudio, }, - Publisher: getPublisherConfig(false), - Subscriber: getSubscriberConfig(rtcConf.CongestionControl.UseSendSideBWEInterceptor || rtcConf.CongestionControl.UseSendSideBWE), - }, nil + Publisher: getPublisherConfig(false), + Subscriber: getSubscriberConfig(rtcConf.CongestionControl.UseSendSideBWEInterceptor || rtcConf.CongestionControl.UseSendSideBWE), + enableFlexFEC: rtcConf.EnableFlexFEC, + } + c.Publisher.EnableFlexFEC = c.enableFlexFEC + return c, nil } func (c *WebRTCConfig) UpdatePublisherConfig(consolidated bool) { c.Publisher = getPublisherConfig(consolidated) + c.Publisher.EnableFlexFEC = c.enableFlexFEC } func (c *WebRTCConfig) UpdateSubscriberConfig(ccConf config.CongestionControlConfig) { diff --git a/pkg/rtc/mediaengine.go b/pkg/rtc/mediaengine.go index cfc53eaaf..087242338 100644 --- a/pkg/rtc/mediaengine.go +++ b/pkg/rtc/mediaengine.go @@ -25,7 +25,21 @@ import ( "github.com/livekit/protocol/livekit" ) -func registerCodecs(me *webrtc.MediaEngine, codecs []*livekit.Codec, rtcpFeedback RTCPFeedbackConfig, filterOutH264HighProfile bool) error { +// flexFEC03CodecParameters describes the FlexFEC-03 codec as sent by libwebrtc. +// It is not part of livekit/protocol codecs - it is a repair mechanism negotiated +// alongside a video codec (a=ssrc-group:FEC-FR), not a publishable codec. +// The payload type is a local preference only, pion matches by capability. +var flexFEC03CodecParameters = webrtc.RTPCodecParameters{ + RTPCodecCapability: webrtc.RTPCodecCapability{ + MimeType: webrtc.MimeTypeFlexFEC03, + ClockRate: 90000, + SDPFmtpLine: "repair-window=10000000", + }, + PayloadType: 49, +} + +func registerCodecs(me *webrtc.MediaEngine, codecs []*livekit.Codec, config DirectionConfig, filterOutH264HighProfile bool) error { + rtcpFeedback := config.RTCPFeedback // audio codecs if IsCodecEnabled(codecs, protoCodecs.OpusCodecParameters.RTPCodecCapability) { cp := protoCodecs.OpusCodecParameters @@ -83,6 +97,12 @@ func registerCodecs(me *webrtc.MediaEngine, codecs []*livekit.Codec, rtcpFeedbac return err } } + + if config.EnableFlexFEC { + if err := me.RegisterCodec(flexFEC03CodecParameters, webrtc.RTPCodecTypeVideo); err != nil { + return err + } + } return nil } @@ -104,7 +124,7 @@ func registerHeaderExtensions(me *webrtc.MediaEngine, rtpHeaderExtension RTPHead func createMediaEngine(codecs []*livekit.Codec, config DirectionConfig, filterOutH264HighProfile bool) (*webrtc.MediaEngine, error) { me := &webrtc.MediaEngine{} - if err := registerCodecs(me, codecs, config.RTCPFeedback, filterOutH264HighProfile); err != nil { + if err := registerCodecs(me, codecs, config, filterOutH264HighProfile); err != nil { return nil, err } diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index bf2ebdb5b..5cdb59360 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -2267,6 +2267,18 @@ func (p *ParticipantImpl) onMediaTrack(rtcTrack *webrtc.TrackRemote, rtpReceiver "parameters", rtpReceiver.GetParameters(), ) + // FlexFEC repair flows are consumed by the primary stream's buffer and must + // not surface as media tracks. pion does not fire tracks for declared FEC + // SSRCs, this is a safety net. The mime package does not know flexfec. + if strings.EqualFold(codec.MimeType, webrtc.MimeTypeFlexFEC03) || strings.EqualFold(codec.MimeType, webrtc.MimeTypeFlexFEC) { + p.pubLogger.Infow( + "ignoring flexfec track", + "trackID", rtcTrack.ID(), + "ssrc", rtcTrack.SSRC(), + ) + return + } + var track sfu.TrackRemote = sfu.NewTrackRemoteFromSdp(rtcTrack, codec) publishedTrack, isNewTrack, isReceiverAdded, sdpRids := p.mediaTrackReceived(track, rtpReceiver) if publishedTrack == nil { diff --git a/pkg/rtc/transport.go b/pkg/rtc/transport.go index 0848ae137..e64fd7565 100644 --- a/pkg/rtc/transport.go +++ b/pkg/rtc/transport.go @@ -1661,6 +1661,14 @@ func (t *PCTransport) HandleRemoteDescription(sd webrtc.SessionDescription, remo t.params.Config.BufferFactory.SetRTXPair(repair, base, "") } } + + fecPairs := flexFECPairsFromSDP(parsed, t.params.Logger) + if len(fecPairs) > 0 { + t.params.Logger.Debugw("flexfec pairs found from sdp", "ssrcs", fecPairs) + for fec, base := range fecPairs { + t.params.Config.BufferFactory.SetFECPair(fec, base) + } + } return nil } @@ -2899,6 +2907,14 @@ func (t *PCTransport) handleRemoteOfferReceived(sd *webrtc.SessionDescription, o } } + fecPairs := flexFECPairsFromSDP(parsed, t.params.Logger) + if len(fecPairs) > 0 { + t.params.Logger.Debugw("flexfec pairs found from sdp", "ssrcs", fecPairs) + for fec, base := range fecPairs { + t.params.Config.BufferFactory.SetFECPair(fec, base) + } + } + if t.currentOfferIceCredential == "" || offerRestartICE { t.currentOfferIceCredential = iceCredential } @@ -3173,6 +3189,15 @@ func (t *PCTransport) restrictReceiverCodecsToPublishList() { if len(filtered) == 0 { continue } + if t.params.DirectionConfig.EnableFlexFEC && tr.Kind() == webrtc.RTPCodecTypeVideo { + // flexfec is a repair mechanism negotiated alongside a video codec, + // it is never part of the publish codec list + for _, c := range receiver.GetParameters().Codecs { + if strings.EqualFold(c.MimeType, webrtc.MimeTypeFlexFEC03) { + filtered = append(filtered, c) + } + } + } if err := tr.SetCodecPreferences(filtered); err != nil { t.params.Logger.Warnw("failed to set recv codec preferences", err, "mid", tr.Mid()) } @@ -3270,6 +3295,45 @@ func nonSimulcastRTXRepairsFromSDP(s *sdp.SessionDescription, logger logger.Logg return rtxRepairFlows } +// flexFECPairsFromSDP extracts FlexFEC repair flows (a=ssrc-group:FEC-FR ) +// from non-simulcast media sections, returning a map of FEC SSRC -> media SSRC. +// libwebrtc only protects a single stream per FEC-FR group, simulcast (rid) sections are skipped. +func flexFECPairsFromSDP(s *sdp.SessionDescription, logger logger.Logger) map[uint32]uint32 { + fecFlows := map[uint32]uint32{} + for _, media := range s.MediaDescriptions { + var ridFound bool + fecPairs := make(map[uint32]uint32) + findFEC: + for _, attr := range media.Attributes { + switch attr.Key { + case "rid": + ridFound = true + break findFEC + case sdp.AttrKeySSRCGroup: + split := strings.Split(attr.Value, " ") + if split[0] == sdp.SemanticTokenForwardErrorCorrectionFramework && len(split) == 3 { + baseSsrc, err := strconv.ParseUint(split[1], 10, 32) + if err != nil { + logger.Warnw("Failed to parse SSRC", err, "ssrc", split[1]) + continue + } + fecSsrc, err := strconv.ParseUint(split[2], 10, 32) + if err != nil { + logger.Warnw("Failed to parse SSRC", err, "ssrc", split[2]) + continue + } + fecPairs[uint32(fecSsrc)] = uint32(baseSsrc) + } + } + } + if !ridFound { + maps.Copy(fecFlows, fecPairs) + } + } + + return fecFlows +} + // ---------------------- type iceCandidatePairStatsEncoder struct { diff --git a/pkg/rtc/transport_test.go b/pkg/rtc/transport_test.go index ab328d7c5..24c6120a8 100644 --- a/pkg/rtc/transport_test.go +++ b/pkg/rtc/transport_test.go @@ -28,9 +28,11 @@ import ( "github.com/livekit/livekit-server/pkg/rtc/transport" "github.com/livekit/livekit-server/pkg/rtc/transport/transportfakes" + "github.com/livekit/livekit-server/pkg/sfu/buffer" "github.com/livekit/livekit-server/pkg/testutils" "github.com/livekit/protocol/codecs/mime" "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/logger" ) func TestMissingAnswerDuringICERestart(t *testing.T) { @@ -610,7 +612,7 @@ func TestConfigureAudioTransceiver(t *testing.T) { } { t.Run(fmt.Sprintf("nack=%v,stereo=%v", testcase.nack, testcase.stereo), func(t *testing.T) { var me webrtc.MediaEngine - registerCodecs(&me, []*livekit.Codec{{Mime: mime.MimeTypeOpus.String()}}, RTCPFeedbackConfig{Audio: []webrtc.RTCPFeedback{{Type: webrtc.TypeRTCPFBNACK}}}, false) + registerCodecs(&me, []*livekit.Codec{{Mime: mime.MimeTypeOpus.String()}}, DirectionConfig{RTCPFeedback: RTCPFeedbackConfig{Audio: []webrtc.RTCPFeedback{{Type: webrtc.TypeRTCPFBNACK}}}}, false) pc, err := webrtc.NewAPI(webrtc.WithMediaEngine(&me)).NewPeerConnection(webrtc.Configuration{}) require.NoError(t, err) defer pc.Close() @@ -732,7 +734,7 @@ func TestSinglePCAnswerStripsSubscribeOnlyCodecsFromRecvSide(t *testing.T) { defer server.Close() var clientME webrtc.MediaEngine - require.NoError(t, registerCodecs(&clientME, subscribeCodecs, RTCPFeedbackConfig{}, false)) + require.NoError(t, registerCodecs(&clientME, subscribeCodecs, DirectionConfig{}, false)) client, err := webrtc.NewAPI(webrtc.WithMediaEngine(&clientME)).NewPeerConnection(webrtc.Configuration{}) require.NoError(t, err) defer client.Close() @@ -777,3 +779,117 @@ func TestSinglePCAnswerStripsSubscribeOnlyCodecsFromRecvSide(t *testing.T) { "answer must not advertise H.264 in recv-side m-section: %s", a.Value) } } + +// TestFlexFECPairsFromSDP verifies extraction of FEC-FR repair flows and that +// simulcast (rid) sections are skipped +func TestFlexFECPairsFromSDP(t *testing.T) { + offer := `v=0 +o=- 8423650423 2 IN IP4 127.0.0.1 +s=- +t=0 0 +m=video 9 UDP/TLS/RTP/SAVPF 96 49 +a=rtpmap:96 VP8/90000 +a=rtpmap:49 flexfec-03/90000 +a=fmtp:49 repair-window=10000000 +a=ssrc-group:FEC-FR 1234 5678 +a=ssrc:1234 cname:test +a=ssrc:5678 cname:test +` + parsed := &sdp.SessionDescription{} + require.NoError(t, parsed.Unmarshal([]byte(offer))) + + pairs := flexFECPairsFromSDP(parsed, logger.GetLogger()) + require.Equal(t, map[uint32]uint32{5678: 1234}, pairs) + + // rid section must be skipped + offerWithRid := offer + "a=rid:hi send\n" + parsed = &sdp.SessionDescription{} + require.NoError(t, parsed.Unmarshal([]byte(offerWithRid))) + require.Empty(t, flexFECPairsFromSDP(parsed, logger.GetLogger())) +} + +// TestFlexFECAnswer verifies the answer accepts flexfec-03 from an offer with a +// FEC-FR group only when FlexFEC is enabled for the direction +func TestFlexFECAnswer(t *testing.T) { + publishCodecs := []*livekit.Codec{ + {Mime: mime.MimeTypeVP8.String()}, + } + + for _, enabled := range []bool{true, false} { + t.Run(fmt.Sprintf("enabled=%v", enabled), func(t *testing.T) { + cfg := &WebRTCConfig{} + cfg.SetBufferFactory(buffer.NewFactoryOfBufferFactory(500, 200).CreateBufferFactory()) + + handler := &transportfakes.FakeHandler{} + server, err := NewPCTransport(TransportParams{ + Config: cfg, + EnabledPublishCodecs: publishCodecs, + DirectionConfig: DirectionConfig{EnableFlexFEC: enabled}, + Handler: handler, + }) + require.NoError(t, err) + defer server.Close() + + var clientME webrtc.MediaEngine + require.NoError(t, registerCodecs(&clientME, publishCodecs, DirectionConfig{}, false)) + client, err := webrtc.NewAPI(webrtc.WithMediaEngine(&clientME)).NewPeerConnection(webrtc.Configuration{}) + require.NoError(t, err) + defer client.Close() + + _, err = client.AddTransceiverFromKind(webrtc.RTPCodecTypeVideo, webrtc.RTPTransceiverInit{ + Direction: webrtc.RTPTransceiverDirectionSendonly, + }) + require.NoError(t, err) + offer, err := client.CreateOffer(nil) + require.NoError(t, err) + require.NoError(t, client.SetLocalDescription(offer)) + + // munge the offer the way a libwebrtc publisher with FlexFEC + // enabled generates it: flexfec-03 codec + FEC-FR ssrc group + mungedSDP := client.LocalDescription().SDP + mediaSSRC := "" + for _, line := range strings.Split(mungedSDP, "\n") { + if strings.HasPrefix(line, "a=ssrc:") { + mediaSSRC = strings.Split(strings.TrimPrefix(line, "a=ssrc:"), " ")[0] + break + } + } + require.NotEmpty(t, mediaSSRC, "offer missing media ssrc") + + lines := strings.Split(strings.TrimRight(mungedSDP, "\r\n"), "\n") + for i, line := range lines { + if strings.HasPrefix(line, "m=video ") { + lines[i] = strings.TrimRight(line, "\r") + " 49" + } + } + lines = append(lines, + "a=rtpmap:49 flexfec-03/90000", + "a=fmtp:49 repair-window=10000000", + fmt.Sprintf("a=ssrc-group:FEC-FR %s 99999", mediaSSRC), + "a=ssrc:99999 cname:fec-test", + ) + munged := webrtc.SessionDescription{ + Type: webrtc.SDPTypeOffer, + SDP: strings.Join(lines, "\n") + "\n", + } + + var answer atomic.Pointer[webrtc.SessionDescription] + handler.OnAnswerCalls(func(sd webrtc.SessionDescription, _ uint32, _ map[string]string) error { + answer.Store(&sd) + return nil + }) + require.NoError(t, server.HandleRemoteDescription(munged, 1)) + + require.Eventually(t, func() bool { + return answer.Load() != nil + }, 5*time.Second, 10*time.Millisecond, "server did not produce answer") + + if enabled { + require.Contains(t, answer.Load().SDP, "flexfec-03/90000", "answer must accept flexfec-03") + require.Contains(t, answer.Load().SDP, "a=rtpmap:49 flexfec-03/90000", "answer must use the offered payload type") + } else { + require.NotContains(t, answer.Load().SDP, "flexfec", "answer must not accept flexfec when disabled") + } + }) + } +} diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index 5b376a490..bc9f88d71 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -18,6 +18,9 @@ import ( "encoding/binary" "errors" "io" + "math/rand/v2" + "os" + "strconv" "github.com/pion/rtcp" "github.com/pion/rtp" @@ -35,12 +38,30 @@ const ( InitPacketBufferSizeVideo = 300 InitPacketBufferSizeAudio = 70 + + flexFECStatsLogInterval = int64(30 * 1e9) // 30s ) var ( errInvalidCodec = errors.New("invalid codec") ) +// debugRxDropPercent enables receive-side fault injection for testing loss +// recovery (NACK/RTX, FlexFEC) without OS-level traffic shaping. Packets are +// dropped before any processing, like wire loss. Test harnesses only, see +// scripts/fec/README.md. +var debugRxDropPercent = func() float64 { + v := os.Getenv("LIVEKIT_DEBUG_RX_DROP_PCT") + if v == "" { + return 0 + } + pct, err := strconv.ParseFloat(v, 64) + if err != nil || pct <= 0 { + return 0 + } + return pct +}() + var _ BufferProvider = (*Buffer)(nil) type pendingPacket struct { @@ -72,6 +93,12 @@ type Buffer struct { primaryBufferForRTX *Buffer rtxPktBuf []byte + + primaryBufferForFEC *Buffer + fecDecoder *flexFECDecoder + fecInjecting bool + lastFECStatsLogAt int64 + lastFECStatsLogged FECStreamStats } func NewBuffer(ssrc uint32, maxVideoPkts, maxAudioPkts int) *Buffer { @@ -145,6 +172,10 @@ func (b *Buffer) Write(pkt []byte) (n int, err error) { return } + if debugRxDropPercent > 0 && rand.Float64()*100 < debugRxDropPercent { + return len(pkt), nil + } + b.Lock() if b.BufferBase.IsClosed() { b.Unlock() @@ -179,6 +210,19 @@ func (b *Buffer) Write(pkt []byte) (n int, err error) { return } + // handle FlexFEC packet + if pb := b.primaryBufferForFEC; pb != nil { + b.Unlock() + + // skip padding only packets + if rtpPacket.Padding && len(rtpPacket.Payload) == 0 { + return + } + + pb.writeFEC(&rtpPacket, now) + return + } + if !b.isBound { packet := make([]byte, len(pkt)) copy(packet, pkt) @@ -239,6 +283,26 @@ func (b *Buffer) NotifyRTX(ssrc uint32, repairSSRC uint32, rsid string) { } } +func (b *Buffer) SetPrimaryBufferForFEC(primaryBuffer *Buffer) { + b.Lock() + b.primaryBufferForFEC = primaryBuffer + pkts := b.pPackets + b.pPackets = nil + b.Unlock() + + for _, pp := range pkts { + var rtpPacket rtp.Packet + err := rtpPacket.Unmarshal(pp.packet) + if err != nil { + continue + } + if rtpPacket.Padding && len(rtpPacket.Payload) == 0 { + continue + } + primaryBuffer.writeFEC(&rtpPacket, pp.arrivalTime) + } +} + func (b *Buffer) writeRTX(rtxPkt *rtp.Packet, arrivalTime int64) { b.Lock() defer b.Unlock() @@ -274,6 +338,107 @@ func (b *Buffer) writeRTX(rtxPkt *rtp.Packet, arrivalTime int64) { b.calc(b.rtxPktBuf[:n], &repairedPkt, arrivalTime, false, true) } +// writeFEC processes a packet of this buffer's FlexFEC repair flow, recovering +// and injecting lost media packets when possible +func (b *Buffer) writeFEC(fecPkt *rtp.Packet, arrivalTime int64) { + b.Lock() + defer b.Unlock() + if !b.isBound { + return + } + + // FEC packets carry transport-wide sequence numbers and count toward the + // publisher's send side bandwidth estimate. The FEC stream is not bound by + // pion, so it is reported on the primary's responder, the streams share a + // media section and with it the extension IDs. + if b.twcc != nil && b.twccExtID != 0 { + if ext := fecPkt.GetExtension(b.twccExtID); ext != nil { + b.twcc.Push(fecPkt.SSRC, binary.BigEndian.Uint16(ext[0:2]), arrivalTime, fecPkt.Marker) + } + } + + // the payload type is learned from negotiated parameters at bind time, when + // unknown the FEC-FR pairing from the SDP and the header validation of the + // decoder are relied upon instead + if b.fecPayloadType != 0 && fecPkt.PayloadType != b.fecPayloadType { + b.logger.Debugw("unexpected fec payload type", "expected", b.fecPayloadType, "actual", fecPkt.PayloadType) + return + } + + if b.fecDecoder == nil { + b.logger.Infow("first flexfec packet received", "fecSSRC", fecPkt.SSRC, "payloadType", fecPkt.PayloadType) + b.fecDecoder = newFlexFECDecoder(flexFECDecoderParams{ + Logger: b.logger, + SSRC: b.BufferBase.SSRC(), + GetPacket: b.BufferBase.getWirePacketLocked, + ExtHighestSN: b.BufferBase.extHighestWireSNLocked, + WindowSize: b.BufferBase.bucketCapacityLocked, + }) + } + + b.injectRecoveredLocked(b.fecDecoder.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, arrivalTime), arrivalTime) + b.maybeLogFECStatsLocked(arrivalTime) +} + +// injectRecoveredLocked feeds FEC recovered packets through the regular packet +// path. An injected packet can complete other buffered FEC packets, those +// chained recoveries are processed in the same loop. +func (b *Buffer) injectRecoveredLocked(recovered [][]byte, arrivalTime int64) { + if len(recovered) == 0 || b.fecInjecting { + return + } + + b.fecInjecting = true + for len(recovered) > 0 { + raw := recovered[0] + recovered = recovered[1:] + + var rtpPacket rtp.Packet + if err := rtpPacket.Unmarshal(raw); err != nil { + continue + } + if rtpPacket.Padding && len(rtpPacket.Payload) == 0 { + // a padding only packet would have been dropped on arrival as well + b.logger.Debugw("discarding recovered padding only packet", "sn", rtpPacket.SequenceNumber) + continue + } + + wireSN := rtpPacket.SequenceNumber + b.calc(raw, &rtpPacket, arrivalTime, false, true) + if b.fecDecoder.HasPending() { + recovered = append(recovered, b.fecDecoder.OnMediaPacket(wireSN)...) + } + } + b.fecInjecting = false +} + +func (b *Buffer) maybeLogFECStatsLocked(now int64) { + if b.fecDecoder == nil || now-b.lastFECStatsLogAt < flexFECStatsLogInterval { + return + } + + stats := b.fecDecoder.Stats() + if stats == b.lastFECStatsLogged { + return + } + + b.lastFECStatsLogAt = now + b.lastFECStatsLogged = stats + b.logger.Debugw("flexfec stats", "stats", stats.String()) +} + +// GetFECStreamStats returns FlexFEC counters for this buffer's repair flow, +// false if no FlexFEC packet has been received +func (b *Buffer) GetFECStreamStats() (FECStreamStats, bool) { + b.RLock() + defer b.RUnlock() + + if b.fecDecoder == nil { + return FECStreamStats{}, false + } + return b.fecDecoder.Stats(), true +} + func (b *Buffer) Read(buff []byte) (n int, err error) { b.Lock() for { @@ -303,6 +468,10 @@ func (b *Buffer) Close() error { return err } + if fecStats, ok := b.GetFECStreamStats(); ok { + b.logger.Infow("flexfec final stats", "stats", fecStats.String()) + } + if stats != nil { if cb := b.getOnFinalRtpStats(); cb != nil { cb(stats) @@ -349,6 +518,13 @@ func (b *Buffer) sendPLI() { } func (b *Buffer) calc(rawPkt []byte, rtpPacket *rtp.Packet, arrivalTime int64, isBuffered bool, isRTX bool) []rtcp.Packet { + // HandleIncomingPacketLocked rewrites the sequence number for bucket storage, + // rtpPacket is nil for packets queued before bind, no FEC state exists then + var wireSN uint16 + if rtpPacket != nil { + wireSN = rtpPacket.SequenceNumber + } + b.BufferBase.HandleIncomingPacketLocked( rawPkt, rtpPacket, @@ -359,6 +535,11 @@ func (b *Buffer) calc(rawPkt []byte, rtpPacket *rtp.Packet, arrivalTime int64, i 0, ) + // an arriving packet can complete buffered FEC packets waiting on it + if rtpPacket != nil && b.fecDecoder != nil && !b.fecInjecting && b.fecDecoder.HasPending() { + b.injectRecoveredLocked(b.fecDecoder.OnMediaPacket(wireSN), arrivalTime) + } + return b.getRTCPPackets(arrivalTime) } diff --git a/pkg/sfu/buffer/buffer_base.go b/pkg/sfu/buffer/buffer_base.go index b94326cff..80c8ec940 100644 --- a/pkg/sfu/buffer/buffer_base.go +++ b/pkg/sfu/buffer/buffer_base.go @@ -168,6 +168,7 @@ type BufferBase struct { rtpParameters webrtc.RTPParameters payloadType uint8 rtxPayloadType uint8 + fecPayloadType uint8 snRangeMap *utils.RangeMap[uint64, uint64] @@ -312,6 +313,14 @@ func (b *BufferBase) BindLocked(rtpParameters webrtc.RTPParameters, codec webrtc } } + // find FlexFEC payload type, the mime package does not know flexfec-03 + for _, codec := range rtpParameters.Codecs { + if strings.EqualFold(codec.MimeType, webrtc.MimeTypeFlexFEC03) { + b.fecPayloadType = uint8(codec.PayloadType) + break + } + } + for _, ext := range rtpParameters.HeaderExtensions { switch ext.URI { case dd.ExtensionURI: @@ -695,6 +704,32 @@ func (b *BufferBase) NotifyRead() { b.readCond.Broadcast() } +// getWirePacketLocked fetches the raw packet stored for the given extended wire +// sequence number, mapping it through the padding-only-drop sequence number adjustment. +func (b *BufferBase) getWirePacketLocked(buf []byte, wireESN uint64) (int, error) { + snAdjustment, err := b.snRangeMap.GetValue(wireESN) + if err != nil { + return 0, err + } + return b.bucket.GetPacket(buf, wireESN-snAdjustment) +} + +// extHighestWireSNLocked returns the extended highest sequence number as seen on the +// wire, i.e. before the padding-only-drop adjustment applied for bucket storage. +func (b *BufferBase) extHighestWireSNLocked() (uint64, bool) { + if b.rtpStats == nil || !b.rtpStats.IsActive() { + return 0, false + } + return b.rtpStats.ExtendedHighestSequenceNumber(), true +} + +func (b *BufferBase) bucketCapacityLocked() int { + if b.bucket == nil { + return 0 + } + return b.bucket.Capacity() +} + func (b *BufferBase) HandleIncomingPacket( rawPkt []byte, rtpPacket *rtp.Packet, diff --git a/pkg/sfu/buffer/factory.go b/pkg/sfu/buffer/factory.go index 4a73ab057..802909e8a 100644 --- a/pkg/sfu/buffer/factory.go +++ b/pkg/sfu/buffer/factory.go @@ -40,6 +40,7 @@ func (f *FactoryOfBufferFactory) CreateBufferFactory() *Factory { rtpBuffers: make(map[uint32]*Buffer), rtcpReaders: make(map[uint32]*RTCPReader), rtxPair: make(map[uint32]uint32), + fecPair: make(map[uint32]uint32), } } @@ -50,6 +51,7 @@ type Factory struct { rtpBuffers map[uint32]*Buffer rtcpReaders map[uint32]*RTCPReader rtxPair map[uint32]uint32 // repair -> base + fecPair map[uint32]uint32 // fec -> base } func (f *Factory) GetOrNew(packetType packetio.BufferPacketType, ssrc uint32) io.ReadWriteCloser { @@ -89,10 +91,26 @@ func (f *Factory) GetOrNew(packetType packetio.BufferPacketType, ssrc uint32) io break } } + for fec, base := range f.fecPair { + if fec == ssrc { + baseBuffer, ok := f.rtpBuffers[base] + if ok { + buffer.SetPrimaryBufferForFEC(baseBuffer) + } + break + } else if base == ssrc { + fecBuffer, ok := f.rtpBuffers[fec] + if ok { + fecBuffer.SetPrimaryBufferForFEC(buffer) + } + break + } + } buffer.OnClose(func() { f.Lock() delete(f.rtpBuffers, ssrc) delete(f.rtxPair, ssrc) + delete(f.fecPair, ssrc) f.Unlock() }) return buffer @@ -132,3 +150,15 @@ func (f *Factory) SetRTXPair(repair, base uint32, rsid string) { } } } + +func (f *Factory) SetFECPair(fec, base uint32) { + f.Lock() + fecBuffer, baseBuffer := f.rtpBuffers[fec], f.rtpBuffers[base] + if fecBuffer == nil || baseBuffer == nil { + f.fecPair[fec] = base + } + f.Unlock() + if fecBuffer != nil && baseBuffer != nil { + fecBuffer.SetPrimaryBufferForFEC(baseBuffer) + } +} diff --git a/pkg/sfu/buffer/flexfec.go b/pkg/sfu/buffer/flexfec.go new file mode 100644 index 000000000..d030c6c32 --- /dev/null +++ b/pkg/sfu/buffer/flexfec.go @@ -0,0 +1,516 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package buffer + +import ( + "encoding/binary" + "errors" + "fmt" + + "github.com/pion/rtp" + + "github.com/livekit/mediatransportutil/pkg/bucket" + "github.com/livekit/protocol/logger" + + "github.com/livekit/livekit-server/pkg/telemetry/prometheus" +) + +// FlexFEC-03 (draft-ietf-payload-flexible-fec-scheme-03) decoder as sent by libwebrtc. +// +// A FlexFEC packet is an RTP packet on its own SSRC (negotiated via +// a=ssrc-group:FEC-FR ) whose payload is: +// +// 0 1 2 3 +// 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// |R|F| P|X| CC |M| PT recovery | length recovery | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// | TS recovery | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// | SSRCCount | reserved | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// | SSRC_i | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// | SN base_i |k| Mask [0-14] | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// |k| Mask [15-45] (optional) | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// |k| | +// +-+ Mask [46-108] (optional) | +// | | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// | ... repair payload (XOR of protected) ... | +// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+ +// +// The first 10 bytes are the XOR of the protected packets' first 8 header +// bytes with their 16-bit (length - 12) in place of the sequence number, +// the repair payload is the XOR of the protected packets' bytes [12:], +// zero padded to the longest protected packet. A missing packet is recovered +// by XORing the FEC bit string and repair payload with all other protected +// packets, which therefore requires every other protected packet. + +var ( + errFECPacketTruncated = errors.New("flexfec packet truncated") + errFECRetransmissionBit = errors.New("flexfec retransmission bit set not supported") + errFECInflexibleMask = errors.New("flexfec inflexible generator matrix not supported") + errFECMultipleSSRC = errors.New("flexfec multiple ssrc protection not supported") + errFECLastMaskKBitNotSet = errors.New("flexfec k-bit of last optional mask not set") + errFECEmptyMask = errors.New("flexfec empty packet mask") + errFECRepairPayloadTooSmall = errors.New("flexfec repair payload smaller than recovered length") +) + +const ( + // FlexFEC-03 masks cover at most 109 packets from SN base + flexFECMaxCoverage = 109 + + defaultMaxPendingFEC = 64 +) + +// FECStreamStats are cumulative counters for one media stream's FlexFEC repair flow +type FECStreamStats struct { + PacketsReceived uint32 + PacketsInvalid uint32 + RecoveryAttempts uint32 + PacketsRecovered uint32 + RecoveryFailed uint32 + PacketsUnused uint32 + PacketsDiscardedOld uint32 +} + +func (f FECStreamStats) String() string { + return fmt.Sprintf( + "received: %d, invalid: %d, recoveryAttempts: %d, recovered: %d, recoveryFailed: %d, unused: %d, discardedOld: %d", + f.PacketsReceived, f.PacketsInvalid, f.RecoveryAttempts, f.PacketsRecovered, f.RecoveryFailed, f.PacketsUnused, f.PacketsDiscardedOld, + ) +} + +type flexFECHeader struct { + protectedSSRC uint32 + snBase uint16 + // offsets from snBase of protected packets, ascending + offsets []uint16 + // offset of the repair payload within the FEC RTP payload + payloadOffset int +} + +// parseFlexFEC03Header parses the FlexFEC-03 header out of a FEC RTP payload +func parseFlexFEC03Header(payload []byte) (flexFECHeader, error) { + if len(payload) < 20 { + return flexFECHeader{}, fmt.Errorf("%w: length %d", errFECPacketTruncated, len(payload)) + } + + if payload[0]&0x80 != 0 { + return flexFECHeader{}, errFECRetransmissionBit + } + if payload[0]&0x40 != 0 { + return flexFECHeader{}, errFECInflexibleMask + } + if ssrcCount := payload[8]; ssrcCount != 1 { + return flexFECHeader{}, fmt.Errorf("%w: count %d", errFECMultipleSSRC, ssrcCount) + } + + h := flexFECHeader{ + protectedSSRC: binary.BigEndian.Uint32(payload[12:]), + snBase: binary.BigEndian.Uint16(payload[16:]), + } + + appendMask := func(mask uint64, bitCount uint16, baseOffset uint16) { + for i := uint16(0); i < bitCount; i++ { + if (mask>>(bitCount-1-i))&1 == 1 { + h.offsets = append(h.offsets, baseOffset+i) + } + } + } + + mask0 := binary.BigEndian.Uint16(payload[18:]) & 0x7fff + appendMask(uint64(mask0), 15, 0) + if payload[18]&0x80 != 0 { + // k-bit 0 set, mask is 15 bits + h.payloadOffset = 20 + } else { + if len(payload) < 24 { + return flexFECHeader{}, fmt.Errorf("%w: length %d", errFECPacketTruncated, len(payload)) + } + appendMask(uint64(binary.BigEndian.Uint32(payload[20:])&0x7fffffff), 31, 15) + if payload[20]&0x80 != 0 { + // k-bit 1 set, masks are 15 + 31 bits + h.payloadOffset = 24 + } else { + // k-bit 2 must be set, masks are 15 + 31 + 63 bits + if len(payload) < 32 { + return flexFECHeader{}, fmt.Errorf("%w: length %d", errFECPacketTruncated, len(payload)) + } + if payload[24]&0x80 == 0 { + return flexFECHeader{}, errFECLastMaskKBitNotSet + } + appendMask(binary.BigEndian.Uint64(payload[24:])&0x7fffffffffffffff, 63, 46) + h.payloadOffset = 32 + } + } + + if len(h.offsets) == 0 { + return flexFECHeader{}, errFECEmptyMask + } + return h, nil +} + +type pendingFECPacket struct { + fecSN uint16 + header flexFECHeader + payload []byte // copied FEC RTP payload + arrivalTime int64 + // extended wire sequence numbers of protected packets, ascending, + // resolved on first evaluation with an active sequence number reference + protected []uint64 + // set once an evaluation saw no not-yet-arrived protected packets, after + // which only arrivals of covered packets can change the outcome + rangeSettled bool +} + +type flexFECDecoderParams struct { + Logger logger.Logger + // SSRC of the protected media stream + SSRC uint32 + // GetPacket fetches the stored raw packet for an extended wire sequence number + GetPacket func(buf []byte, wireESN uint64) (int, error) + // ExtHighestSN returns the extended highest wire sequence number received, false if none yet + ExtHighestSN func() (uint64, bool) + // WindowSize returns the number of packets retrievable via GetPacket + WindowSize func() int + // MaxPendingFEC bounds the number of buffered FEC packets, 0 for default + MaxPendingFEC int +} + +// flexFECDecoder recovers lost media packets from a FlexFEC-03 repair flow. +// Received media packets are not duplicated here, presence and XOR inputs +// are sourced from the primary buffer's packet bucket via GetPacket. +// All methods must be called holding the primary buffer's lock. +type flexFECDecoder struct { + params flexFECDecoderParams + + pending []*pendingFECPacket + stats FECStreamStats + + pktBuf []byte // scratch for GetPacket + payloadAcc []byte // scratch for repair payload XOR accumulation + + // packets recovered within the current call, recovery of the same + // packet by an overlapping FEC is deferred until it is injected + justRecovered map[uint64]struct{} +} + +func newFlexFECDecoder(params flexFECDecoderParams) *flexFECDecoder { + if params.MaxPendingFEC == 0 { + params.MaxPendingFEC = defaultMaxPendingFEC + } + return &flexFECDecoder{ + params: params, + pktBuf: make([]byte, bucket.RTPMaxPktSize), + payloadAcc: make([]byte, bucket.RTPMaxPktSize), + justRecovered: make(map[uint64]struct{}, 4), + } +} + +func (d *flexFECDecoder) HasPending() bool { + return len(d.pending) > 0 +} + +func (d *flexFECDecoder) Stats() FECStreamStats { + return d.stats +} + +// AddFEC processes one received FlexFEC packet and returns raw recovered media +// packets, if any. The payload is copied and may be reused by the caller. +func (d *flexFECDecoder) AddFEC(fecSN uint16, payload []byte, arrivalTime int64) [][]byte { + d.stats.PacketsReceived++ + prometheus.IncrementFEC(prometheus.FECStateReceived, 1) + + for _, p := range d.pending { + if p.fecSN == fecSN { + // duplicate of a buffered FEC packet + return nil + } + } + + header, err := parseFlexFEC03Header(payload) + if err == nil && header.protectedSSRC != d.params.SSRC { + err = fmt.Errorf("flexfec protecting unexpected ssrc, expected %d, got %d", d.params.SSRC, header.protectedSSRC) + } + if err != nil { + d.stats.PacketsInvalid++ + prometheus.IncrementFEC(prometheus.FECStateInvalid, 1) + d.params.Logger.Debugw("dropping invalid flexfec packet", "error", err, "fecSN", fecSN) + return nil + } + + defer clear(d.justRecovered) + recoveredPackets := d.evictStale() + + if len(d.pending) >= d.params.MaxPendingFEC { + d.pending = d.pending[1:] + d.stats.PacketsDiscardedOld++ + prometheus.IncrementFEC(prometheus.FECStateDiscardedOld, 1) + } + + p := &pendingFECPacket{ + fecSN: fecSN, + header: header, + payload: append([]byte(nil), payload...), + arrivalTime: arrivalTime, + } + d.pending = append(d.pending, p) + + if recovered := d.evaluate(p); recovered != nil { + recoveredPackets = append(recoveredPackets, recovered) + } + return recoveredPackets +} + +// OnMediaPacket re-evaluates buffered FEC packets covering the given wire +// sequence number and returns raw recovered media packets, if any. +// Callers should check HasPending first to keep the common path cheap. +func (d *flexFECDecoder) OnMediaPacket(wireSN uint16) [][]byte { + if len(d.pending) == 0 { + return nil + } + + ref, ok := d.params.ExtHighestSN() + if !ok { + return nil + } + esn := unwrapNearESN(wireSN, ref) + + var recoveredPackets [][]byte + defer clear(d.justRecovered) + for _, p := range append([]*pendingFECPacket(nil), d.pending...) { + if !p.affectedBy(esn) { + continue + } + if recovered := d.evaluate(p); recovered != nil { + recoveredPackets = append(recoveredPackets, recovered) + } + } + return recoveredPackets +} + +// affectedBy reports whether the arrival of the given packet can change the +// outcome of this FEC packet, either because it is protected by it or because +// it moves a previously not-yet-arrived part of the protected range into the +// loss-detectable past +func (p *pendingFECPacket) affectedBy(esn uint64) bool { + if p.protected == nil { + // not resolved yet, conservatively assume affected + return true + } + if !p.rangeSettled && esn > p.protected[len(p.protected)-1] { + return true + } + for _, protectedESN := range p.protected { + if protectedESN == esn { + return true + } + } + return false +} + +// evictStale drops buffered FEC packets whose protected range has fallen out +// of the packet window, giving each a final evaluation first. Returns any +// packets recovered by those final evaluations. +func (d *flexFECDecoder) evictStale() [][]byte { + ref, ok := d.params.ExtHighestSN() + if !ok { + return nil + } + + var recoveredPackets [][]byte + window := uint64(d.params.WindowSize()) + for _, p := range append([]*pendingFECPacket(nil), d.pending...) { + if p.protected == nil { + continue + } + if ref > p.protected[0]+window { + // oldest protected packet is leaving the window, final attempt; + // evaluate evicts the packet on any conclusive outcome + if recovered := d.evaluate(p); recovered != nil { + recoveredPackets = append(recoveredPackets, recovered) + } + d.remove(p, prometheus.FECStateRecoveryFailed) + } + } + return recoveredPackets +} + +// remove drops a pending FEC packet, counting the disposition if it is still buffered +func (d *flexFECDecoder) remove(p *pendingFECPacket, state prometheus.FECState) { + for i, pp := range d.pending { + if pp != p { + continue + } + d.pending = append(d.pending[:i], d.pending[i+1:]...) + switch state { + case prometheus.FECStateUnused: + d.stats.PacketsUnused++ + prometheus.IncrementFEC(state, 1) + case prometheus.FECStateRecoveryFailed: + d.stats.RecoveryFailed++ + prometheus.IncrementFEC(state, 1) + case prometheus.FECStateDiscardedOld: + d.stats.PacketsDiscardedOld++ + prometheus.IncrementFEC(state, 1) + } + return + } +} + +// evaluate classifies the protected packets of a buffered FEC packet and runs +// XOR recovery when exactly one is missing, returning the recovered raw packet. +// Conclusive outcomes (recovered, unused, unrecoverable) evict the FEC packet. +func (d *flexFECDecoder) evaluate(p *pendingFECPacket) []byte { + ref, ok := d.params.ExtHighestSN() + if !ok { + // no media received yet, keep pending + return nil + } + + if p.protected == nil { + p.protected = make([]uint64, 0, len(p.header.offsets)) + for _, offset := range p.header.offsets { + p.protected = append(p.protected, unwrapNearESN(p.header.snBase+offset, ref)) + } + } + + var ( + missing int + missingESN uint64 + hdrAcc [12]byte + repair = p.payload[p.header.payloadOffset:] + payloadAccSize = len(repair) + ) + if payloadAccSize > len(d.payloadAcc) { + payloadAccSize = len(d.payloadAcc) + } + copy(hdrAcc[:10], p.payload[:10]) + copy(d.payloadAcc[:payloadAccSize], repair[:payloadAccSize]) + + for _, esn := range p.protected { + if _, isJustRecovered := d.justRecovered[esn]; isJustRecovered { + // recovered by an overlapping FEC in this call but not yet in the + // bucket, defer to the re-evaluation triggered by its injection + return nil + } + + n, err := d.params.GetPacket(d.pktBuf, esn) + switch { + case err == nil: + if n < 12 { + // not a parseable RTP packet, cannot use as XOR input + d.remove(p, prometheus.FECStateRecoveryFailed) + return nil + } + raw := d.pktBuf[:n] + hdrAcc[0] ^= raw[0] + hdrAcc[1] ^= raw[1] + lengthRecovery := uint16(n - 12) + hdrAcc[2] ^= byte(lengthRecovery >> 8) + hdrAcc[3] ^= byte(lengthRecovery) + for i := 4; i < 8; i++ { + hdrAcc[i] ^= raw[i] + } + for i := 12; i < n && i-12 < payloadAccSize; i++ { + d.payloadAcc[i-12] ^= raw[i] + } + + case errors.Is(err, bucket.ErrPacketSizeInvalid): + // in-window hole, the packet was lost + missing++ + missingESN = esn + if missing > 1 { + // not recoverable yet, keep pending + return nil + } + + case errors.Is(err, bucket.ErrPacketTooNew): + // not arrived yet, keep pending and re-evaluate as the stream advances + return nil + + default: + // packet aged out of the window or was excluded (e.g. a received + // padding-only packet that is not stored), never recoverable + d.remove(p, prometheus.FECStateRecoveryFailed) + return nil + } + } + // the classification completed without not-yet-arrived packets + p.rangeSettled = true + + if missing == 0 { + d.remove(p, prometheus.FECStateUnused) + return nil + } + + recovered, err := d.finishRecovery(p, hdrAcc, d.payloadAcc[:payloadAccSize], missingESN) + d.stats.RecoveryAttempts++ + prometheus.IncrementFEC(prometheus.FECStateRecoveryAttempt, 1) + if err != nil { + d.params.Logger.Debugw("flexfec recovery failed", "error", err, "fecSN", p.fecSN, "missingESN", missingESN) + d.remove(p, prometheus.FECStateRecoveryFailed) + return nil + } + + d.stats.PacketsRecovered++ + prometheus.IncrementFEC(prometheus.FECStateRecovered, 1) + d.justRecovered[missingESN] = struct{}{} + // the FEC packet did its job, no eviction disposition to count + d.remove(p, "") + return recovered +} + +// finishRecovery turns the XOR accumulators into the recovered raw packet +func (d *flexFECDecoder) finishRecovery(p *pendingFECPacket, hdrAcc [12]byte, payloadAcc []byte, missingESN uint64) ([]byte, error) { + // force RTP version 2 + hdrAcc[0] = (hdrAcc[0] | 0x80) & 0xbf + + recoveredLen := int(binary.BigEndian.Uint16(hdrAcc[2:4])) + if recoveredLen > len(payloadAcc) { + return nil, fmt.Errorf("%w: recovered %d, repair %d", errFECRepairPayloadTooSmall, recoveredLen, len(payloadAcc)) + } + + binary.BigEndian.PutUint16(hdrAcc[2:4], uint16(missingESN)) + binary.BigEndian.PutUint32(hdrAcc[8:12], d.params.SSRC) + + recovered := make([]byte, 12+recoveredLen) + copy(recovered, hdrAcc[:]) + copy(recovered[12:], payloadAcc[:recoveredLen]) + + var pkt rtp.Packet + if err := pkt.Unmarshal(recovered); err != nil { + return nil, fmt.Errorf("recovered packet does not unmarshal: %w", err) + } + return recovered, nil +} + +// unwrapNearESN expands a 16-bit wire sequence number to the extended sequence +// number closest to the reference +func unwrapNearESN(sn uint16, ref uint64) uint64 { + candidate := (ref &^ uint64(0xffff)) | uint64(sn) + if candidate > ref { + if candidate-ref > 0x8000 && candidate >= (1<<16) { + candidate -= 1 << 16 + } + } else if ref-candidate > 0x8000 { + candidate += 1 << 16 + } + return candidate +} diff --git a/pkg/sfu/buffer/flexfec_test.go b/pkg/sfu/buffer/flexfec_test.go new file mode 100644 index 000000000..02c7ec8e4 --- /dev/null +++ b/pkg/sfu/buffer/flexfec_test.go @@ -0,0 +1,591 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package buffer + +import ( + "fmt" + "testing" + + pionflexfec "github.com/pion/interceptor/pkg/flexfec" + "github.com/pion/rtp" + "github.com/pion/transport/v4/packetio" + "github.com/pion/webrtc/v4" + "github.com/stretchr/testify/require" + + "github.com/livekit/mediatransportutil/pkg/bucket" + "github.com/livekit/mediatransportutil/pkg/twcc" + "github.com/livekit/protocol/logger" +) + +const ( + testFECMediaSSRC = uint32(0x12345678) + testFECSSRC = uint32(0x23456789) + testFECPT = uint8(49) +) + +// fecPacketStore mimics the primary buffer's packet bucket keyed by extended +// wire sequence number +type fecPacketStore struct { + packets map[uint64][]byte + highest uint64 + hasAny bool + window uint64 +} + +func newFECPacketStore(window int) *fecPacketStore { + return &fecPacketStore{ + packets: make(map[uint64][]byte), + window: uint64(window), + } +} + +func (s *fecPacketStore) add(esn uint64, raw []byte) { + s.packets[esn] = raw + if !s.hasAny || esn > s.highest { + s.highest = esn + } + s.hasAny = true +} + +func (s *fecPacketStore) getPacket(buf []byte, esn uint64) (int, error) { + if !s.hasAny || esn > s.highest { + return 0, bucket.ErrPacketTooNew + } + if s.highest-esn >= s.window { + return 0, bucket.ErrPacketTooOld + } + raw, ok := s.packets[esn] + if !ok { + return 0, bucket.ErrPacketSizeInvalid + } + return copy(buf, raw), nil +} + +func (s *fecPacketStore) decoderParams() flexFECDecoderParams { + return flexFECDecoderParams{ + Logger: logger.GetLogger(), + SSRC: testFECMediaSSRC, + GetPacket: s.getPacket, + ExtHighestSN: func() (uint64, bool) { return s.highest, s.hasAny }, + WindowSize: func() int { return int(s.window) }, + } +} + +// makeFECMediaPackets builds count consecutive media packets starting at baseSN +// with varying payload sizes, header extensions and marker bits +func makeFECMediaPackets(t *testing.T, baseSN uint16, count int) ([]rtp.Packet, [][]byte) { + t.Helper() + + packets := make([]rtp.Packet, 0, count) + raws := make([][]byte, 0, count) + for i := 0; i < count; i++ { + pkt := rtp.Packet{ + Header: rtp.Header{ + Version: 2, + PayloadType: 96, + SequenceNumber: baseSN + uint16(i), + Timestamp: 0x1000 + uint32(i/3)*3000, + SSRC: testFECMediaSSRC, + Marker: i%3 == 2, + }, + Payload: make([]byte, 20+(i*37)%600), + } + for j := range pkt.Payload { + pkt.Payload[j] = byte(i + j) + } + if i%3 == 0 { + require.NoError(t, pkt.Header.SetExtension(3, []byte{byte(i), byte(i >> 8)})) + } + + raw, err := pkt.Marshal() + require.NoError(t, err) + packets = append(packets, pkt) + raws = append(raws, raw) + } + return packets, raws +} + +func encodeFECPackets(t *testing.T, media []rtp.Packet, numFEC uint32) []rtp.Packet { + t.Helper() + + fecPackets := pionflexfec.NewFlexEncoder03(testFECPT, testFECSSRC).EncodeFec(media, numFEC) + require.NotEmpty(t, fecPackets) + return fecPackets +} + +func TestParseFlexFEC03Header(t *testing.T) { + buildPayload := func(k0, k1, k2 bool, size int) []byte { + payload := make([]byte, size) + payload[8] = 1 // SSRCCount + // protected SSRC + payload[12], payload[13], payload[14], payload[15] = 0x12, 0x34, 0x56, 0x78 + // SN base + payload[16], payload[17] = 0x10, 0x01 + if k0 { + payload[18] |= 0x80 + } + if k1 && size > 20 { + payload[20] |= 0x80 + } + if k2 && size > 24 { + payload[24] |= 0x80 + } + return payload + } + + t.Run("mask sizes", func(t *testing.T) { + // k-bit 0 set, 15 bit mask, offsets 0 and 14 + payload := buildPayload(true, false, false, 30) + payload[18] |= 0x40 // offset 0 + payload[19] |= 0x01 // offset 14 + h, err := parseFlexFEC03Header(payload) + require.NoError(t, err) + require.Equal(t, uint32(0x12345678), h.protectedSSRC) + require.Equal(t, uint16(0x1001), h.snBase) + require.Equal(t, []uint16{0, 14}, h.offsets) + require.Equal(t, 20, h.payloadOffset) + + // k-bit 1 set, 15+31 bit masks, offsets 15 and 45 + payload = buildPayload(false, true, false, 40) + payload[20] |= 0x40 // offset 15 + payload[23] |= 0x01 // offset 45 + h, err = parseFlexFEC03Header(payload) + require.NoError(t, err) + require.Equal(t, []uint16{15, 45}, h.offsets) + require.Equal(t, 24, h.payloadOffset) + + // k-bit 2 set, 15+31+63 bit masks, offsets 46 and 108 + payload = buildPayload(false, false, true, 40) + payload[24] |= 0x40 // offset 46 + payload[31] |= 0x01 // offset 108 + h, err = parseFlexFEC03Header(payload) + require.NoError(t, err) + require.Equal(t, []uint16{46, 108}, h.offsets) + require.Equal(t, 32, h.payloadOffset) + }) + + t.Run("rejects", func(t *testing.T) { + // truncated + _, err := parseFlexFEC03Header(make([]byte, 19)) + require.ErrorIs(t, err, errFECPacketTruncated) + + // truncated with k-bit 0 unset + _, err = parseFlexFEC03Header(buildPayload(false, true, false, 20)) + require.ErrorIs(t, err, errFECPacketTruncated) + + // truncated with k-bits 0 and 1 unset + _, err = parseFlexFEC03Header(buildPayload(false, false, true, 24)) + require.ErrorIs(t, err, errFECPacketTruncated) + + // retransmission bit + payload := buildPayload(true, false, false, 30) + payload[0] |= 0x80 + _, err = parseFlexFEC03Header(payload) + require.ErrorIs(t, err, errFECRetransmissionBit) + + // inflexible mask bit + payload = buildPayload(true, false, false, 30) + payload[0] |= 0x40 + _, err = parseFlexFEC03Header(payload) + require.ErrorIs(t, err, errFECInflexibleMask) + + // multiple SSRC + payload = buildPayload(true, false, false, 30) + payload[8] = 2 + _, err = parseFlexFEC03Header(payload) + require.ErrorIs(t, err, errFECMultipleSSRC) + + // k-bit of last mask not set + payload = buildPayload(false, false, false, 40) + _, err = parseFlexFEC03Header(payload) + require.ErrorIs(t, err, errFECLastMaskKBitNotSet) + + // empty mask + payload = buildPayload(true, false, false, 30) + _, err = parseFlexFEC03Header(payload) + require.ErrorIs(t, err, errFECEmptyMask) + }) +} + +func TestFlexFECDecoderRoundTrip(t *testing.T) { + for _, blockSize := range []int{10, 30, 100} { + for _, dropIdx := range []int{0, blockSize / 2, blockSize - 1} { + t.Run(fmt.Sprintf("block=%d,drop=%d", blockSize, dropIdx), func(t *testing.T) { + media, raws := makeFECMediaPackets(t, 100, blockSize) + fecPkt := encodeFECPackets(t, media, 1)[0] + + store := newFECPacketStore(1000) + for i, raw := range raws { + if i != dropIdx { + store.add(uint64(media[i].SequenceNumber), raw) + } + } + + d := newFlexFECDecoder(store.decoderParams()) + recovered := d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0) + + if dropIdx == blockSize-1 { + // the dropped packet is past the highest received one, it + // is not detectable as lost until the stream advances + require.Empty(t, recovered) + next, nextRaws := makeFECMediaPackets(t, 100+uint16(blockSize), 1) + store.add(uint64(next[0].SequenceNumber), nextRaws[0]) + recovered = d.OnMediaPacket(next[0].SequenceNumber) + } + + require.Len(t, recovered, 1) + require.Equal(t, raws[dropIdx], recovered[0]) + + stats := d.Stats() + require.Equal(t, uint32(1), stats.PacketsReceived) + require.Equal(t, uint32(1), stats.RecoveryAttempts) + require.Equal(t, uint32(1), stats.PacketsRecovered) + require.Equal(t, uint32(0), stats.RecoveryFailed) + require.False(t, d.HasPending()) + }) + } + } +} + +func TestFlexFECDecoderNoLoss(t *testing.T) { + media, raws := makeFECMediaPackets(t, 5000, 12) + fecPkt := encodeFECPackets(t, media, 1)[0] + + store := newFECPacketStore(1000) + for i, raw := range raws { + store.add(uint64(media[i].SequenceNumber), raw) + } + + d := newFlexFECDecoder(store.decoderParams()) + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0)) + + stats := d.Stats() + require.Equal(t, uint32(1), stats.PacketsReceived) + require.Equal(t, uint32(1), stats.PacketsUnused) + require.Equal(t, uint32(0), stats.RecoveryAttempts) + require.False(t, d.HasPending()) +} + +func TestFlexFECDecoderFECBeforeMedia(t *testing.T) { + media, raws := makeFECMediaPackets(t, 300, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + + store := newFECPacketStore(1000) + d := newFlexFECDecoder(store.decoderParams()) + + // FEC arrives before any media + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0)) + require.True(t, d.HasPending()) + + // media arrives, packet 4 is lost + var recovered [][]byte + for i, raw := range raws { + if i == 4 { + continue + } + store.add(uint64(media[i].SequenceNumber), raw) + recovered = append(recovered, d.OnMediaPacket(media[i].SequenceNumber)...) + } + + require.Len(t, recovered, 1) + require.Equal(t, raws[4], recovered[0]) + require.False(t, d.HasPending()) +} + +func TestFlexFECDecoderTwoMissing(t *testing.T) { + media, raws := makeFECMediaPackets(t, 700, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + + store := newFECPacketStore(1000) + for i, raw := range raws { + if i == 2 || i == 6 { + continue + } + store.add(uint64(media[i].SequenceNumber), raw) + } + + d := newFlexFECDecoder(store.decoderParams()) + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0)) + require.True(t, d.HasPending()) + require.Equal(t, uint32(0), d.Stats().RecoveryAttempts) + + // packet 2 arrives late, e.g. via RTX, packet 6 becomes recoverable + store.add(uint64(media[2].SequenceNumber), raws[2]) + recovered := d.OnMediaPacket(media[2].SequenceNumber) + require.Len(t, recovered, 1) + require.Equal(t, raws[6], recovered[0]) + require.False(t, d.HasPending()) +} + +func TestFlexFECDecoderDuplicateFEC(t *testing.T) { + media, raws := makeFECMediaPackets(t, 900, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + + store := newFECPacketStore(1000) + for i, raw := range raws { + if i == 3 || i == 5 { + continue + } + store.add(uint64(media[i].SequenceNumber), raw) + } + + d := newFlexFECDecoder(store.decoderParams()) + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0)) + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0)) + + stats := d.Stats() + require.Equal(t, uint32(2), stats.PacketsReceived) + require.Equal(t, 1, len(d.pending), "duplicate must not be buffered twice") +} + +func TestFlexFECDecoderInvalid(t *testing.T) { + store := newFECPacketStore(1000) + d := newFlexFECDecoder(store.decoderParams()) + + // garbage payload + require.Empty(t, d.AddFEC(1, make([]byte, 10), 0)) + + // valid header protecting another SSRC + media, _ := makeFECMediaPackets(t, 100, 5) + for i := range media { + media[i].SSRC = testFECMediaSSRC + 1 + } + fecPkt := encodeFECPackets(t, media, 1)[0] + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0)) + + stats := d.Stats() + require.Equal(t, uint32(2), stats.PacketsReceived) + require.Equal(t, uint32(2), stats.PacketsInvalid) + require.False(t, d.HasPending()) +} + +func TestFlexFECDecoderWraparound(t *testing.T) { + // media packets crossing the 16-bit sequence number wrap + media, raws := makeFECMediaPackets(t, 65530, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + + const cycle = uint64(7) << 16 + esn := func(i int) uint64 { + // 65530..65535 are in the cycle, 0..3 in the next + if media[i].SequenceNumber >= 65530 { + return cycle | uint64(media[i].SequenceNumber) + } + return cycle + (1 << 16) + uint64(media[i].SequenceNumber) + } + + for _, dropIdx := range []int{2, 8} { // one drop on each side of the wrap + t.Run(fmt.Sprintf("drop=%d", dropIdx), func(t *testing.T) { + store := newFECPacketStore(1000) + for i, raw := range raws { + if i != dropIdx { + store.add(esn(i), raw) + } + } + + d := newFlexFECDecoder(store.decoderParams()) + recovered := d.AddFEC(fecPkt.SequenceNumber, fecPkt.Payload, 0) + require.Len(t, recovered, 1) + require.Equal(t, raws[dropIdx], recovered[0]) + }) + } +} + +func TestFlexFECDecoderOverlappingChain(t *testing.T) { + // two FEC packets with overlapping coverage, recovering a packet via the + // first enables the second once the recovered packet is fed back + media, raws := makeFECMediaPackets(t, 2000, 12) + fecA := encodeFECPackets(t, media[:8], 1)[0] // protects 0..7 + fecB := encodeFECPackets(t, media[4:], 1)[0] // protects 4..11 + + store := newFECPacketStore(1000) + for i, raw := range raws { + if i == 6 || i == 10 { + continue + } + store.add(uint64(media[i].SequenceNumber), raw) + } + + d := newFlexFECDecoder(store.decoderParams()) + + // distinct FEC stream sequence numbers, each encoder starts at the same one + // B cannot recover, 6 and 10 missing + require.Empty(t, d.AddFEC(1, fecB.Payload, 0)) + // A recovers 6 + recovered := d.AddFEC(2, fecA.Payload, 0) + require.Len(t, recovered, 1) + require.Equal(t, raws[6], recovered[0]) + + // the buffer injects the recovered packet and notifies the decoder, + // which lets B recover 10 + store.add(uint64(media[6].SequenceNumber), raws[6]) + recovered = d.OnMediaPacket(media[6].SequenceNumber) + require.Len(t, recovered, 1) + require.Equal(t, raws[10], recovered[0]) + + require.False(t, d.HasPending()) + require.Equal(t, uint32(2), d.Stats().PacketsRecovered) +} + +func TestFlexFECDecoderStaleEviction(t *testing.T) { + media, raws := makeFECMediaPackets(t, 100, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + + store := newFECPacketStore(50) + for i, raw := range raws { + if i == 2 || i == 6 { + continue + } + store.add(uint64(media[i].SequenceNumber), raw) + } + + d := newFlexFECDecoder(store.decoderParams()) + require.Empty(t, d.AddFEC(1, fecPkt.Payload, 0)) + require.True(t, d.HasPending()) + + // the stream advances beyond the window, next FEC packet triggers a sweep + next, nextRaws := makeFECMediaPackets(t, 300, 10) + for i, raw := range nextRaws { + store.add(uint64(next[i].SequenceNumber), raw) + } + nextFEC := encodeFECPackets(t, next, 1)[0] + require.Empty(t, d.AddFEC(2, nextFEC.Payload, 0)) + + stats := d.Stats() + require.Equal(t, uint32(1), stats.RecoveryFailed, "stale FEC with missing packets evicted as failed") + require.Equal(t, uint32(1), stats.PacketsUnused, "fresh FEC with no losses evicted as unused") + require.False(t, d.HasPending()) +} + +// TestBufferFlexFECIntegration exercises the full path, a FEC packet written to +// the paired FEC buffer recovers a lost media packet into the primary's bucket +func TestBufferFlexFECIntegration(t *testing.T) { + flexfecCodec := webrtc.RTPCodecParameters{ + RTPCodecCapability: webrtc.RTPCodecCapability{ + MimeType: webrtc.MimeTypeFlexFEC03, + ClockRate: 90000, + SDPFmtpLine: "repair-window=10000000", + }, + PayloadType: 49, + } + + media, raws := makeFECMediaPackets(t, 100, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + fecRaw, err := fecPkt.Marshal() + require.NoError(t, err) + + for _, pairFirst := range []bool{true, false} { + t.Run(fmt.Sprintf("pairBeforeFECBuffer=%v", pairFirst), func(t *testing.T) { + factory := NewFactoryOfBufferFactory(500, 200).CreateBufferFactory() + primary := factory.GetOrNew(packetio.RTPBufferPacket, testFECMediaSSRC).(*Buffer) + require.NoError(t, primary.Bind(webrtc.RTPParameters{ + Codecs: []webrtc.RTPCodecParameters{vp8Codec, flexfecCodec}, + }, vp8Codec.RTPCodecCapability, 0)) + require.Equal(t, uint8(49), primary.fecPayloadType, "fec payload type learned at bind") + + twccResponder := twcc.NewTransportWideCCResponder() + primary.SetTWCCAndExtID(twccResponder, 5) + + // media packets arrive, packet 4 is lost + for i, raw := range raws { + if i == 4 { + continue + } + _, err := primary.Write(raw) + require.NoError(t, err) + } + + if pairFirst { + factory.SetFECPair(testFECSSRC, testFECMediaSSRC) + } + fecBuffer := factory.GetOrNew(packetio.RTPBufferPacket, testFECSSRC).(*Buffer) + if !pairFirst { + factory.SetFECPair(testFECSSRC, testFECMediaSSRC) + } + + _, err = fecBuffer.Write(fecRaw) + require.NoError(t, err) + + // the recovered packet must be byte equal in the primary's bucket + primary.Lock() + ref, ok := primary.extHighestWireSNLocked() + require.True(t, ok) + buf := make([]byte, bucket.RTPMaxPktSize) + n, getErr := primary.getWirePacketLocked(buf, unwrapNearESN(media[4].SequenceNumber, ref)) + primary.Unlock() + require.NoError(t, getErr) + require.Equal(t, raws[4], buf[:n]) + + stats, hasStats := primary.GetFECStreamStats() + require.True(t, hasStats) + require.Equal(t, uint32(1), stats.PacketsReceived) + require.Equal(t, uint32(1), stats.RecoveryAttempts) + require.Equal(t, uint32(1), stats.PacketsRecovered) + }) + } +} + +// TestBufferFlexFECBeforePairing verifies FEC packets queued in the unbound FEC +// buffer are processed once the pairing is established +func TestBufferFlexFECBeforePairing(t *testing.T) { + media, raws := makeFECMediaPackets(t, 100, 10) + fecPkt := encodeFECPackets(t, media, 1)[0] + fecRaw, err := fecPkt.Marshal() + require.NoError(t, err) + + factory := NewFactoryOfBufferFactory(500, 200).CreateBufferFactory() + primary := factory.GetOrNew(packetio.RTPBufferPacket, testFECMediaSSRC).(*Buffer) + require.NoError(t, primary.Bind(webrtc.RTPParameters{ + Codecs: []webrtc.RTPCodecParameters{vp8Codec}, + }, vp8Codec.RTPCodecCapability, 0)) + require.Equal(t, uint8(0), primary.fecPayloadType, "no flexfec codec negotiated") + + // FEC packet arrives before the pairing is known, queued in the FEC buffer + fecBuffer := factory.GetOrNew(packetio.RTPBufferPacket, testFECSSRC).(*Buffer) + _, err = fecBuffer.Write(fecRaw) + require.NoError(t, err) + + for i, raw := range raws { + if i == 4 { + continue + } + _, err := primary.Write(raw) + require.NoError(t, err) + } + + // pairing flushes the queued FEC packet into the primary + factory.SetFECPair(testFECSSRC, testFECMediaSSRC) + + stats, hasStats := primary.GetFECStreamStats() + require.True(t, hasStats) + require.Equal(t, uint32(1), stats.PacketsReceived) + require.Equal(t, uint32(1), stats.PacketsRecovered) +} + +func TestFlexFECDecoderPendingOverflow(t *testing.T) { + store := newFECPacketStore(100000) + params := store.decoderParams() + params.MaxPendingFEC = 4 + d := newFlexFECDecoder(params) + + // no media received, FEC packets pile up + for i := 0; i < 6; i++ { + media, _ := makeFECMediaPackets(t, uint16(1000+i*20), 5) + fecPkt := encodeFECPackets(t, media, 1)[0] + require.Empty(t, d.AddFEC(fecPkt.SequenceNumber+uint16(i), fecPkt.Payload, 0)) + } + + require.Equal(t, 4, len(d.pending)) + require.Equal(t, uint32(2), d.Stats().PacketsDiscardedOld) +} diff --git a/pkg/telemetry/prometheus/packets.go b/pkg/telemetry/prometheus/packets.go index a9f6684e1..0e0d6ba55 100644 --- a/pkg/telemetry/prometheus/packets.go +++ b/pkg/telemetry/prometheus/packets.go @@ -35,6 +35,19 @@ const ( TransmissionRetransmit TransmissionType = "retransmit" ) +// FECState describes the disposition of incoming FlexFEC packets and recovery outcomes +type FECState string + +const ( + FECStateReceived FECState = "received" // valid FlexFEC packet received + FECStateInvalid FECState = "invalid" // FlexFEC packet failed to parse or was unsupported + FECStateRecoveryAttempt FECState = "recovery_attempt" // XOR recovery executed + FECStateRecovered FECState = "recovered" // recovered media packet injected + FECStateRecoveryFailed FECState = "recovery_failed" // FEC packet discarded with protected packets still missing + FECStateUnused FECState = "unused" // FEC packet discarded with all protected packets received + FECStateDiscardedOld FECState = "discarded_old" // FEC packet discarded due to age/overflow before evaluation +) + var ( bytesIn atomic.Uint64 bytesOut atomic.Uint64 @@ -57,6 +70,7 @@ var ( promPacketTotal *prometheus.CounterVec promPacketBytes *prometheus.CounterVec promRTCPLabels = []string{"direction", "country"} + promFlexFECTotal *prometheus.CounterVec promStreamLabels = []string{"direction", "source", "type", "country"} promNackTotal *prometheus.CounterVec promPliTotal *prometheus.CounterVec @@ -105,6 +119,12 @@ func initPacketStats(nodeID string, nodeType livekit.NodeType) { Name: "total", ConstLabels: prometheus.Labels{"node_id": nodeID, "node_type": nodeType.String()}, }, promRTCPLabels) + promFlexFECTotal = prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: livekitNamespace, + Subsystem: "flexfec_packet", + Name: "total", + ConstLabels: prometheus.Labels{"node_id": nodeID, "node_type": nodeType.String()}, + }, []string{"state"}) promPacketLossTotal = prometheus.NewCounterVec(prometheus.CounterOpts{ Namespace: livekitNamespace, Subsystem: "packet_loss", @@ -195,6 +215,7 @@ func initPacketStats(nodeID string, nodeType livekit.NodeType) { prometheus.MustRegister(promNackTotal) prometheus.MustRegister(promPliTotal) prometheus.MustRegister(promFirTotal) + prometheus.MustRegister(promFlexFECTotal) prometheus.MustRegister(promPacketLossTotal) prometheus.MustRegister(promPacketLoss) prometheus.MustRegister(promPacketOutOfOrderTotal) @@ -208,6 +229,15 @@ func initPacketStats(nodeID string, nodeType livekit.NodeType) { prometheus.MustRegister(promForwardLatencyHist) } +// IncrementFEC counts incoming FlexFEC packet dispositions and recovery outcomes. +// Safe to call before Init (e.g. from unit tests), increments are dropped in that case. +func IncrementFEC(state FECState, count uint32) { + if promFlexFECTotal == nil || count == 0 { + return + } + promFlexFECTotal.WithLabelValues(string(state)).Add(float64(count)) +} + func IncrementPackets(country string, direction Direction, count uint64, retransmit bool) { var transmission TransmissionType if retransmit { diff --git a/scripts/fec/.gitignore b/scripts/fec/.gitignore new file mode 100644 index 000000000..8cb86768a --- /dev/null +++ b/scripts/fec/.gitignore @@ -0,0 +1,2 @@ +# FEC test run artifacts (logs, CSVs, prom samples, reports, built binaries) +out/ diff --git a/scripts/fec/README.md b/scripts/fec/README.md new file mode 100644 index 000000000..086890206 --- /dev/null +++ b/scripts/fec/README.md @@ -0,0 +1,135 @@ +# FlexFEC test harness (publisher → SFU) + +Validates that FlexFEC-03 sent by a publisher recovers lost packets at the SFU +under cellular-like loss (continuous base loss plus periodic bursts), the kind +of uplink a robot on LTE/5G sees. The harness runs everything on loopback: + +``` +publisher (rust-sdks local_video, --test-pattern --flex-fec) + │ UDP → 127.0.0.1:7882 ← traffic shaper drops packets here + ▼ +livekit-server (enable_flexfec: true, prometheus on :6789) + │ + ▼ +subscriber (rust-sdks local_video, --headless --log-frames) +``` + +The publisher attaches a wall-clock timestamp and a frame id to every frame +via the packet trailer feature; the subscriber logs them per received frame. +That gives ground truth for end-to-end frame latency and frame loss, while +the SFU's prometheus counters (`livekit_flexfec_packet_total{state=...}`, +`livekit_nack_total`, `livekit_packet_loss_total`) show what FEC did. + +## Requirements + +- `go`, `cargo`, `curl`, `python3` with `matplotlib` (`pip3 install matplotlib`) +- a `rust-sdks` checkout with the local_video example + (default `../rust-sdks` relative to this repo, override with `RUST_SDKS_DIR`) +- `sudo` for traffic shaping: + - **macOS**: dummynet (`dnctl` + `pfctl`). If dummynet is unavailable on + your macOS build, run with `--no-shaping` and use Network Link Conditioner + manually, or test on Linux. + - **Linux**: `tc` with the `netem` qdisc (`iproute2`). + +## Usage + +```bash +# A/B comparison: baseline (NACK only) vs FlexFEC, 2 minutes each +./run_fec_test.sh --mode ab --duration 120 + +# single FEC run with a harsher profile +./run_fec_test.sh --mode fec --duration 60 --base-loss 0.05 --burst-loss 0.4 + +# sanity check without shaping (expect ~0 loss, ~0 recoveries) +./run_fec_test.sh --mode fec --duration 30 --no-shaping + +# no sudo available: drop 4% of received packets inside the SFU instead of +# OS shaping (uniform loss only, also useful for CI) +./run_fec_test.sh --mode ab --duration 60 --debug-drop 4 +``` + +Outputs land in `scripts/fec/out//`: + +- `fec_report.png` — stacked time series: per-frame latency with frame-gap + markers and burst shading, FEC received/recovered/failed rates, NACK and + packet-loss rates, delivered fps +- `summary.txt` — per-run table (frames lost, latency percentiles, FEC + counters, NACK totals) and the A/B comparison +- per run (`baseline/`, `fec/`): `server.log`, `publisher.log`, + `subscriber.log`, `frames.csv`, `prom.tsv`, `events.csv`, `meta.env` + +## Loss profile + +Defaults simulate a robot uplink over cellular: 2% continuous loss +(Gilbert-Elliott on Linux for realistic correlation, uniform on macOS) with a +3 s burst of 25% loss every 15 s. Tune via `--base-loss`, `--burst-loss`, +`--burst-every`, `--burst-len`. The shaper matches **UDP destined to port +7882** on loopback, which is the publisher→SFU media leg; the subscriber's +upstream RTCP shares that port and is shaped too, which is acceptable for A/B +comparisons since both runs see identical conditions. + +The publisher defaults to a fixed 30% FEC protection rate with the bursty +mask (`--fec-rate`, `--fec-mask-type`); pass `--fec-rate 0` to let libwebrtc +adapt the rate to its loss estimate instead. + +`--debug-drop PCT` is an unprivileged alternative to OS shaping: the SFU +drops PCT% of received packets before processing (uniform, all SSRCs, +enabled via the `LIVEKIT_DEBUG_RX_DROP_PCT` env var). Use the OS shapers for +the cellular burst profile, this knob for quick checks and CI. Note that with +uniform loss and a fixed protection rate, blocks losing two or more packets +are not FEC-recoverable and fall back to NACK/RTX, so expect partial +recovery; bursty loss with the bursty mask is the scenario FlexFEC targets. + +## What to expect + +- Sanity run (no shaping): `state="received"` grows, `recovered` ≈ 0, no + frame gaps. +- FEC run under bursts: `recovered` spikes inside burst windows, the + subscriber sees few or no frame-id gaps, latency stays near baseline. +- Baseline under the same bursts: frame gaps and latency spikes during + bursts (NACK/RTX needs a round trip per loss; FEC repairs immediately). +- Publisher bitrate should not collapse in the FEC run — FEC packets carry + transport-wide CC sequence numbers and the SFU reports them, so the + publisher's bandwidth estimate stays intact. + +## Parameter sweep + +`sweep_fec.sh` runs `run_fec_test.sh` across a matrix of loss levels and FEC +configurations (one baseline plus one FEC run per `(rate, mask)` at each loss +level), builds the binaries once, then `aggregate_fec.py` combines all cells +into a single comparison: `sweep_report.png` (frame loss, p99 latency, FEC +recovery rate, and recovered-packet count, each vs loss level, baseline vs +every FEC config), `sweep_summary.csv`, and a markdown table. + +```bash +# uniform-loss sweep, no sudo: 4 loss levels x 2 FEC rates x 1 mask +./sweep_fec.sh --loss-mode debug --loss-list "2 5 10 15" --fec-rate-list "20 50" + +# cellular-burst sweep: vary burst intensity, compare the two masks at 30% +sudo -v && ./sweep_fec.sh --loss-mode shaped --loss-list "0.15 0.30 0.50" \ + --fec-rate-list "30" --mask-list "random bursty" --duration 90 + +# re-aggregate an existing sweep without re-running it +python3 aggregate_fec.py --sweep out/sweep_ +``` + +Runtime ≈ cells × (~20s setup + `--duration`). Each cell is a directory under +the sweep output; `manifest.tsv` maps cell → parameters. In `debug` mode +`--loss-list` is packet-drop percentages; in `shaped` mode it is the burst +loss fraction (base loss and cadence fixed by `--base-loss`/`--burst-every`/ +`--burst-len`). Note the caveat above: uniform `debug` loss with a fixed FEC +rate only recovers blocks that lose a single packet, so `shaped` bursts with +the `bursty` mask show FlexFEC at its best. + +## Direct script use + +```bash +sudo ./shape_macos.sh start --port 7882 --base-loss 0.02 --burst-loss 0.25 \ + --burst-every 15 --burst-len 3 --events /tmp/events.csv +sudo ./shape_macos.sh stop # force cleanup (also: shape_linux.sh) + +./prom_poll.sh 6789 /tmp/prom.tsv + +# single run / A/B pair +python3 plot_fec.py --run out//baseline --run out//fec --out out/ +``` diff --git a/scripts/fec/aggregate_fec.py b/scripts/fec/aggregate_fec.py new file mode 100755 index 000000000..55c4b10a0 --- /dev/null +++ b/scripts/fec/aggregate_fec.py @@ -0,0 +1,212 @@ +#!/usr/bin/env python3 +"""Aggregate a FlexFEC parameter sweep into a single comparison report. + +Reads the cells produced by sweep_fec.sh (a manifest.tsv plus one run directory +per cell) and emits: + sweep_report.png - loss/latency/recovery vs loss level, baseline vs each config + sweep_summary.csv - one row per cell with the key metrics + a markdown table on stdout + +Usage: + aggregate_fec.py --sweep +""" + +import argparse +import csv +import os +import sys + +import matplotlib + +matplotlib.use("Agg") +import matplotlib.pyplot as plt # noqa: E402 + +# Run lives in plot_fec.py next to this script +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) +from plot_fec import Run # noqa: E402 + + +def read_manifest(sweep_dir): + rows = [] + with open(os.path.join(sweep_dir, "manifest.tsv")) as f: + for row in csv.DictReader(f, delimiter="\t"): + rows.append(row) + return rows + + +def collect(sweep_dir): + records = [] + for entry in read_manifest(sweep_dir): + run_dir = os.path.join(sweep_dir, entry["leaf"]) + if not os.path.isdir(run_dir): + print(f"warning: missing run dir {run_dir}", file=sys.stderr) + continue + run = Run(run_dir) + if not run.frames: + print(f"warning: no frames for {entry['leaf']}, skipping", file=sys.stderr) + continue + s = run.summary() + try: + loss = float(entry["loss"]) + except ValueError: + loss = 0.0 + expected = s["frames_expected"] or 1 + rec = { + "leaf": entry["leaf"], + "mode": entry["mode"], + "loss": loss, + "loss_label": entry["loss"], + "fec_rate": entry["fec_rate"], + "fec_mask": entry["fec_mask"], + "config": "baseline" + if entry["mode"] == "baseline" + else f"FEC r{entry['fec_rate']}/{entry['fec_mask']}", + "frame_loss_pct": 100.0 * s["frames_lost"] / expected, + "recovery_rate_pct": (100.0 * s["fec_recovered"] / s["fec_received"]) if s["fec_received"] else 0.0, + "debug_drop": run.meta.get("DEBUG_DROP", "0"), + **s, + } + records.append(rec) + return records + + +def x_unit(records): + if records and records[0]["debug_drop"] not in ("0", ""): + return "injected packet drop (%)" + return "burst loss fraction" + + +def write_csv(records, path): + cols = [ + "config", "mode", "loss_label", "fec_rate", "fec_mask", + "frames_received", "frames_expected", "frames_lost", "frame_loss_pct", + "latency_p50_ms", "latency_p95_ms", "latency_p99_ms", + "fec_received", "fec_recovery_attempts", "fec_recovered", + "fec_recovery_failed", "recovery_rate_pct", "nack_total", + ] + with open(path, "w", newline="") as f: + w = csv.writer(f) + w.writerow(cols) + for r in sorted(records, key=lambda r: (r["loss"], r["config"])): + w.writerow([_fmt(r.get(c)) for c in cols]) + + +def _fmt(v): + if isinstance(v, float): + return f"{v:.2f}" + return v + + +def series_by_config(records): + """config label -> sorted [(loss, record)]""" + by_cfg = {} + for r in records: + by_cfg.setdefault(r["config"], []).append(r) + for cfg in by_cfg: + by_cfg[cfg].sort(key=lambda r: r["loss"]) + return by_cfg + + +def plot(records, out_path): + by_cfg = series_by_config(records) + unit = x_unit(records) + # baseline first (gray), FEC configs in a color cycle + order = sorted(by_cfg, key=lambda c: (c != "baseline", c)) + cmap = plt.get_cmap("viridis") + fec_cfgs = [c for c in order if c != "baseline"] + colors = {"baseline": "#888888"} + for i, c in enumerate(fec_cfgs): + colors[c] = cmap(0.15 + 0.7 * (i / max(1, len(fec_cfgs) - 1))) + + fig, axes = plt.subplots(2, 2, figsize=(15, 11)) + ax_loss, ax_lat, ax_rec, ax_recn = axes.flat + + def line(ax, key, cfgs): + for cfg in cfgs: + pts = by_cfg[cfg] + xs = [r["loss"] for r in pts] + ys = [r[key] for r in pts] + ax.plot(xs, ys, marker="o", label=cfg, color=colors[cfg], + lw=2 if cfg != "baseline" else 1.5, + ls="-" if cfg != "baseline" else "--") + + line(ax_loss, "frame_loss_pct", order) + ax_loss.set_title("Frame loss vs loss level") + ax_loss.set_ylabel("frames lost (%)") + ax_loss.set_xlabel(unit) + ax_loss.legend(fontsize=8) + ax_loss.grid(alpha=0.3) + + line(ax_lat, "latency_p99_ms", order) + ax_lat.set_title("Tail latency (p99) vs loss level") + ax_lat.set_ylabel("capture→receive p99 (ms)") + ax_lat.set_xlabel(unit) + ax_lat.legend(fontsize=8) + ax_lat.grid(alpha=0.3) + + line(ax_rec, "recovery_rate_pct", fec_cfgs) + ax_rec.set_title("FEC recovery rate (recovered / FEC packets received)") + ax_rec.set_ylabel("recovery rate (%)") + ax_rec.set_xlabel(unit) + ax_rec.legend(fontsize=8) + ax_rec.grid(alpha=0.3) + + line(ax_recn, "fec_recovered", fec_cfgs) + ax_recn.set_title("Packets recovered by FEC") + ax_recn.set_ylabel("recovered packets") + ax_recn.set_xlabel(unit) + ax_recn.legend(fontsize=8) + ax_recn.grid(alpha=0.3) + + fig.suptitle("FlexFEC parameter sweep — baseline (no FEC) vs FEC configs", fontsize=13) + fig.tight_layout(rect=(0, 0, 1, 0.97)) + fig.savefig(out_path, dpi=130) + + +def print_table(records): + cols = [ + ("config", "config", "{}"), + ("loss_label", "loss", "{}"), + ("frames_lost", "lost", "{:.0f}"), + ("frame_loss_pct", "lost%", "{:.1f}"), + ("latency_p95_ms", "p95ms", "{:.1f}"), + ("latency_p99_ms", "p99ms", "{:.1f}"), + ("fec_received", "fecRx", "{:.0f}"), + ("fec_recovered", "recov", "{:.0f}"), + ("recovery_rate_pct", "recov%", "{:.1f}"), + ("nack_total", "nacks", "{:.0f}"), + ] + widths = [max(len(h), 8) for _, h, _ in cols] + print("| " + " | ".join(h.ljust(w) for (_, h, _), w in zip(cols, widths)) + " |") + print("|" + "|".join("-" * (w + 2) for w in widths) + "|") + for r in sorted(records, key=lambda r: (r["loss"], r["config"] != "baseline", r["config"])): + cells = [] + for (key, _, fmt), w in zip(cols, widths): + v = r.get(key) + cells.append((fmt.format(v) if v is not None else "-").ljust(w)) + print("| " + " | ".join(cells) + " |") + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--sweep", required=True, help="sweep output directory (contains manifest.tsv)") + args = parser.parse_args() + + records = collect(args.sweep) + if not records: + print("no usable cells found", file=sys.stderr) + sys.exit(1) + + csv_path = os.path.join(args.sweep, "sweep_summary.csv") + png_path = os.path.join(args.sweep, "sweep_report.png") + write_csv(records, csv_path) + plot(records, png_path) + + print(f"report: {png_path}") + print(f"summary: {csv_path}") + print() + print_table(records) + + +if __name__ == "__main__": + main() diff --git a/scripts/fec/plot_fec.py b/scripts/fec/plot_fec.py new file mode 100755 index 000000000..728ed2a66 --- /dev/null +++ b/scripts/fec/plot_fec.py @@ -0,0 +1,373 @@ +#!/usr/bin/env python3 +"""Render time-series plots and a summary for FlexFEC test harness runs. + +Inputs are run directories produced by run_fec_test.sh, each containing: + frames.csv - per received frame: recv_wall_us,frame_id,user_timestamp_us,width,height + prom.tsv - 1s prometheus samples: wall_usmetric{labels}value + events.csv - shaper events: wall_us,event (burst_on / burst_off / ...) + meta.env - run parameters incl. T0_US/T1_US measurement window + +Usage: + plot_fec.py --run [--run ] --out + +With two runs the first is treated as the baseline and the second as the FEC +run, both are overlaid and a comparison table is printed. +""" + +import argparse +import csv +import os +import re +import statistics +import sys + +import matplotlib + +matplotlib.use("Agg") +import matplotlib.pyplot as plt # noqa: E402 + +FLEXFEC_METRIC = "livekit_flexfec_packet_total" +METRIC_RE = re.compile(r"^(?P[a-zA-Z0-9_]+)(?:\{(?P.*)\})?$") + + +class Run: + def __init__(self, path): + self.path = path + self.name = os.path.basename(os.path.normpath(path)) + self.meta = self._read_meta() + self.t0 = int(self.meta.get("T0_US", 0)) + self.t1 = int(self.meta.get("T1_US", 1 << 62)) + self.frames = self._read_frames() + self.prom = self._read_prom() + self.bursts = self._read_bursts() + + def _read_meta(self): + meta = {} + path = os.path.join(self.path, "meta.env") + if os.path.exists(path): + with open(path) as f: + for line in f: + if "=" in line: + k, v = line.strip().split("=", 1) + meta[k] = v + return meta + + def _in_window(self, ts): + return self.t0 <= ts <= self.t1 + + def rel(self, ts): + return (ts - self.t0) / 1e6 + + def _read_frames(self): + frames = [] + path = os.path.join(self.path, "frames.csv") + if not os.path.exists(path): + return frames + with open(path) as f: + for row in csv.DictReader(f): + try: + recv = int(row["recv_wall_us"]) + except (KeyError, ValueError): + continue + if not self._in_window(recv): + continue + frames.append( + { + "recv": recv, + "frame_id": int(row["frame_id"]) if row.get("frame_id") else None, + "user_ts": int(row["user_timestamp_us"]) if row.get("user_timestamp_us") else None, + } + ) + frames.sort(key=lambda fr: fr["recv"]) + return frames + + def _read_prom(self): + """metric family -> label key -> [(wall_us, value)]""" + series = {} + path = os.path.join(self.path, "prom.tsv") + if not os.path.exists(path): + return series + with open(path) as f: + for line in f: + parts = line.rstrip("\n").split("\t") + if len(parts) != 3: + continue + ts, metric, value = parts + m = METRIC_RE.match(metric) + if not m: + continue + try: + ts, value = int(ts), float(value) + except ValueError: + continue + family = m.group("name") + labels = m.group("labels") or "" + series.setdefault(family, {}).setdefault(labels, []).append((ts, value)) + return series + + def _read_bursts(self): + """[(start_rel_s, end_rel_s)] of shaper burst windows""" + bursts, start = [], None + path = os.path.join(self.path, "events.csv") + if not os.path.exists(path): + return bursts + with open(path) as f: + for line in f: + parts = line.strip().split(",", 1) + if len(parts) != 2: + continue + ts, event = int(parts[0]), parts[1] + if event == "burst_on": + start = self.rel(ts) + elif event == "burst_off" and start is not None: + bursts.append((start, self.rel(ts))) + start = None + if start is not None: + bursts.append((start, self.rel(self.t1 if self.t1 < (1 << 62) else start))) + return bursts + + # ---------- derived series ---------- + + def latency_series(self): + xs, ys = [], [] + for fr in self.frames: + if fr["user_ts"]: + xs.append(self.rel(fr["recv"])) + ys.append((fr["recv"] - fr["user_ts"]) / 1e3) + return xs, ys + + def frame_gaps(self): + """[(rel_s, missing_count)] where frame ids skipped between consecutive frames""" + gaps, prev = [], None + for fr in self.frames: + fid = fr["frame_id"] + if fid is None: + continue + if prev is not None and fid > prev + 1: + gaps.append((self.rel(fr["recv"]), fid - prev - 1)) + prev = fid + return gaps + + def fps_series(self): + buckets = {} + for fr in self.frames: + buckets[int(self.rel(fr["recv"]))] = buckets.get(int(self.rel(fr["recv"])), 0) + 1 + xs = sorted(buckets) + return xs, [buckets[x] for x in xs] + + def counter_rate(self, family, label_filter=None): + """summed per-second rate across label sets of a counter family""" + fam = self.prom.get(family, {}) + per_ts = {} + for labels, samples in fam.items(): + if label_filter and label_filter not in labels: + continue + samples = [s for s in samples if self._in_window(s[0])] + for (t_a, v_a), (t_b, v_b) in zip(samples, samples[1:]): + dt = (t_b - t_a) / 1e6 + if dt <= 0: + continue + key = int(self.rel(t_b)) + per_ts[key] = per_ts.get(key, 0.0) + max(0.0, v_b - v_a) / dt + xs = sorted(per_ts) + return xs, [per_ts[x] for x in xs] + + def counter_total(self, family, label_filter=None): + """delta of a summed counter family over the measurement window""" + total = 0.0 + for labels, samples in self.prom.get(family, {}).items(): + if label_filter and label_filter not in labels: + continue + samples = [s for s in samples if self._in_window(s[0])] + if len(samples) >= 2: + total += samples[-1][1] - samples[0][1] + return total + + def flexfec_state_rate(self, state): + return self.counter_rate(FLEXFEC_METRIC, f'state="{state}"') + + def flexfec_state_total(self, state): + return self.counter_total(FLEXFEC_METRIC, f'state="{state}"') + + # ---------- summary ---------- + + def summary(self): + _, lat = self.latency_series() + gaps = self.frame_gaps() + frame_ids = [fr["frame_id"] for fr in self.frames if fr["frame_id"] is not None] + expected = (max(frame_ids) - min(frame_ids) + 1) if frame_ids else 0 + lost = sum(n for _, n in gaps) + + def pct(p): + if not lat: + return float("nan") + data = sorted(lat) + return data[min(len(data) - 1, int(len(data) * p))] + + return { + "frames_received": len(self.frames), + "frames_expected": expected, + "frames_lost": lost, + "gap_events": len(gaps), + "latency_mean_ms": statistics.fmean(lat) if lat else float("nan"), + "latency_p50_ms": pct(0.50), + "latency_p95_ms": pct(0.95), + "latency_p99_ms": pct(0.99), + "fec_received": self.flexfec_state_total("received"), + "fec_recovery_attempts": self.flexfec_state_total("recovery_attempt"), + "fec_recovered": self.flexfec_state_total("recovered"), + "fec_recovery_failed": self.flexfec_state_total("recovery_failed"), + "fec_unused": self.flexfec_state_total("unused"), + "fec_invalid": self.flexfec_state_total("invalid"), + "nack_total": self.counter_total("livekit_nack_total"), + "packet_loss_total": self.counter_total("livekit_packet_loss_total"), + } + + +def shade_bursts(ax, bursts): + for start, end in bursts: + ax.axvspan(start, end, color="red", alpha=0.08, lw=0) + + +def plot(runs, out_dir): + colors = {"baseline": "#888888", "fec": "#1f77b4"} + fig, axes = plt.subplots(4, 1, figsize=(14, 16), sharex=True) + ax_lat, ax_fec, ax_loss, ax_fps = axes + bursts = runs[-1].bursts + + for run in runs: + color = colors.get(run.name, None) + xs, ys = run.latency_series() + ax_lat.plot(xs, ys, ".", markersize=2.5, label=f"{run.name} latency", color=color, alpha=0.7) + for x, n in run.frame_gaps(): + ax_lat.axvline(x, color=color or "red", alpha=0.5, lw=min(0.5 + n * 0.3, 3)) + shade_bursts(ax_lat, bursts) + ax_lat.set_ylabel("capture→receive latency (ms)") + ax_lat.set_title("Frame latency (vertical lines: frame-id gaps = lost frames, red shading: loss bursts)") + ax_lat.legend(loc="upper right", fontsize=8) + ax_lat.grid(alpha=0.3) + + fec_runs = [r for r in runs if r.flexfec_state_total("received") > 0] or runs[-1:] + for run in fec_runs: + for state, style in [ + ("received", dict(color="#1f77b4", lw=1)), + ("recovered", dict(color="#2ca02c", lw=1.8)), + ("recovery_failed", dict(color="#d62728", lw=1.2)), + ("unused", dict(color="#9467bd", lw=0.8, alpha=0.6)), + ]: + xs, ys = run.flexfec_state_rate(state) + ax_fec.plot(xs, ys, label=f"{run.name} {state}/s", **style) + shade_bursts(ax_fec, bursts) + ax_fec.set_ylabel("FEC packets/s") + ax_fec.set_title("FlexFEC at the SFU") + ax_fec.legend(loc="upper right", fontsize=8) + ax_fec.grid(alpha=0.3) + + for run in runs: + color = colors.get(run.name, None) + xs, ys = run.counter_rate("livekit_nack_total") + ax_loss.plot(xs, ys, label=f"{run.name} nack/s", color=color, lw=1.2) + xs, ys = run.counter_rate("livekit_packet_loss_total") + ax_loss.plot(xs, ys, label=f"{run.name} packet_loss/s", color=color, lw=1.2, ls="--", alpha=0.7) + shade_bursts(ax_loss, bursts) + ax_loss.set_ylabel("packets/s") + ax_loss.set_title("NACK and reported packet loss") + ax_loss.legend(loc="upper right", fontsize=8) + ax_loss.grid(alpha=0.3) + + for run in runs: + color = colors.get(run.name, None) + xs, ys = run.fps_series() + ax_fps.plot(xs, ys, label=f"{run.name} fps", color=color, lw=1.2) + shade_bursts(ax_fps, bursts) + ax_fps.set_ylabel("frames/s") + ax_fps.set_xlabel("seconds since measurement start") + ax_fps.set_title("Delivered frame rate at subscriber") + ax_fps.legend(loc="lower right", fontsize=8) + ax_fps.grid(alpha=0.3) + + meta = runs[-1].meta + if meta.get("DEBUG_DROP", "0") not in ("0", ""): + profile = "uniform {}% loss injected at SFU receive".format(meta["DEBUG_DROP"]) + elif meta.get("SHAPING") == "0": + profile = "no loss" + else: + profile = "base loss {} / burst {} for {}s every {}s".format( + meta.get("BASE_LOSS", "?"), + meta.get("BURST_LOSS", "?"), + meta.get("BURST_LEN", "?"), + meta.get("BURST_EVERY", "?"), + ) + fig.suptitle( + "FlexFEC publisher→SFU recovery — {}, codec {}".format(profile, meta.get("CODEC", "?")), + fontsize=12, + ) + fig.tight_layout(rect=(0, 0, 1, 0.985)) + out_path = os.path.join(out_dir, "fec_report.png") + fig.savefig(out_path, dpi=130) + return out_path + + +def print_summary(runs): + summaries = [(run.name, run.summary()) for run in runs] + keys = [ + ("frames_received", "{:.0f}"), + ("frames_expected", "{:.0f}"), + ("frames_lost", "{:.0f}"), + ("gap_events", "{:.0f}"), + ("latency_mean_ms", "{:.1f}"), + ("latency_p50_ms", "{:.1f}"), + ("latency_p95_ms", "{:.1f}"), + ("latency_p99_ms", "{:.1f}"), + ("fec_received", "{:.0f}"), + ("fec_recovery_attempts", "{:.0f}"), + ("fec_recovered", "{:.0f}"), + ("fec_recovery_failed", "{:.0f}"), + ("fec_unused", "{:.0f}"), + ("fec_invalid", "{:.0f}"), + ("nack_total", "{:.0f}"), + ("packet_loss_total", "{:.0f}"), + ] + + name_w = 24 + header = "metric".ljust(name_w) + "".join(name.rjust(16) for name, _ in summaries) + print(header) + print("-" * len(header)) + for key, fmt in keys: + row = key.ljust(name_w) + for _, summary in summaries: + row += fmt.format(summary[key]).rjust(16) + print(row) + + if len(summaries) == 2: + base, fec = summaries[0][1], summaries[1][1] + print() + if fec["frames_lost"] < base["frames_lost"]: + print( + "frames lost reduced {} -> {} with FlexFEC".format( + int(base["frames_lost"]), int(fec["frames_lost"]) + ) + ) + if fec["fec_recovered"] > 0: + print("SFU recovered {} packets via FlexFEC".format(int(fec["fec_recovered"]))) + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--run", action="append", required=True, help="run directory (repeatable)") + parser.add_argument("--out", required=True, help="output directory for the report") + args = parser.parse_args() + + runs = [Run(path) for path in args.run] + for run in runs: + if not run.frames: + print(f"warning: no frames recorded for {run.path}", file=sys.stderr) + + out_path = plot(runs, args.out) + print(f"report: {out_path}") + print() + print_summary(runs) + + +if __name__ == "__main__": + main() diff --git a/scripts/fec/prom_poll.sh b/scripts/fec/prom_poll.sh new file mode 100755 index 000000000..755e80239 --- /dev/null +++ b/scripts/fec/prom_poll.sh @@ -0,0 +1,25 @@ +#!/usr/bin/env bash +# Polls the SFU's prometheus endpoint once per second and appends the FEC and +# packet counters relevant to the FlexFEC test harness as tab-separated rows: +# wall_usmetric{labels}value +# +# Usage: ./prom_poll.sh + +set -u + +PORT="${1:?prometheus port required}" +OUT="${2:?output file required}" + +now_us() { python3 -c 'import time; print(int(time.time() * 1e6))'; } + +while true; do + TS=$(now_us) + curl -s --max-time 2 "http://127.0.0.1:${PORT}/metrics" | \ + awk -v ts="$TS" '/^livekit_(flexfec_packet|nack|packet|packet_loss|packet_out_of_order)_total/ { + value = $NF + metric = $0 + sub(/ [^ ]*$/, "", metric) + print ts "\t" metric "\t" value + }' >> "$OUT" + sleep 1 +done diff --git a/scripts/fec/run_fec_test.sh b/scripts/fec/run_fec_test.sh new file mode 100755 index 000000000..10cbbd1da --- /dev/null +++ b/scripts/fec/run_fec_test.sh @@ -0,0 +1,344 @@ +#!/usr/bin/env bash +# End-to-end FlexFEC test harness: publisher (rust-sdks local_video) -> SFU, +# with traffic shaping on the publisher->SFU leg simulating a cellular uplink +# (continuous base loss + periodic loss bursts), collecting frame metadata at +# the subscriber and FEC/NACK/loss counters from the SFU, then plotting time +# series and printing a summary. +# +# Usage: +# ./run_fec_test.sh --mode {fec|baseline|ab} [options] +# +# Modes: +# fec single run with --flex-fec on the publisher +# baseline single run without FlexFEC (NACK/RTX recovery only) +# ab baseline run followed by a fec run with the same loss profile, +# producing a comparison report +# +# Options (defaults in brackets): +# --duration S measurement duration per run after media starts [120] +# --out DIR output directory [scripts/fec/out/] +# --base-loss F continuous loss fraction [0.02] +# --burst-loss F loss fraction inside bursts [0.1] +# --burst-every S seconds between burst starts [15] +# --burst-len S burst duration in seconds [1] +# --fec-rate N publisher FEC protection rate percent, 0 = adaptive [30] +# --fec-mask-type T random|bursty [bursty] +# --codec C video codec [h264] +# --width/--height/--fps test pattern format [1280/720/30] +# --no-shaping skip traffic shaping (sanity run) +# --debug-drop PCT drop PCT% of received packets inside the SFU instead of +# OS traffic shaping (uniform loss, no sudo required; +# implies --no-shaping) +# --server-bin PATH use this prebuilt livekit-server instead of building +# --skip-build skip go/cargo builds (requires --server-bin and +# pre-built example binaries); used by sweep_fec.sh +# +# Requires: go, cargo, python3 with matplotlib, curl, and sudo (for shaping). +# The rust-sdks checkout is located via RUST_SDKS_DIR [../rust-sdks]. + +set -u + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" +RUST_SDKS_DIR="${RUST_SDKS_DIR:-$(cd "$REPO_ROOT/.." && pwd)/rust-sdks}" + +MODE="" +DURATION=120 +OUT_DIR="" +BASE_LOSS=0.02 +BURST_LOSS=0.1 +BURST_EVERY=15 +BURST_LEN=1 +FEC_RATE=30 +FEC_MASK_TYPE="bursty" +CODEC="h264" +WIDTH=1280 +HEIGHT=720 +FPS=30 +SHAPING=1 +DEBUG_DROP=0 +SERVER_BIN="" +SKIP_BUILD=0 + +SIGNAL_PORT=7880 +MEDIA_PORT=7882 +PROM_PORT=6789 +API_KEY="devkey" +API_SECRET="fec-test-secret-fec-test-secret-00" +ROOM_NAME="fec-test" + +log() { echo "[fec-test] $*"; } +die() { echo "[fec-test] ERROR: $*" >&2; exit 1; } +now_us() { python3 -c 'import time; print(int(time.time() * 1e6))'; } + +while [ $# -gt 0 ]; do + case "$1" in + --mode) MODE="$2"; shift 2 ;; + --duration) DURATION="$2"; shift 2 ;; + --out) OUT_DIR="$2"; shift 2 ;; + --base-loss) BASE_LOSS="$2"; shift 2 ;; + --burst-loss) BURST_LOSS="$2"; shift 2 ;; + --burst-every) BURST_EVERY="$2"; shift 2 ;; + --burst-len) BURST_LEN="$2"; shift 2 ;; + --fec-rate) FEC_RATE="$2"; shift 2 ;; + --fec-mask-type) FEC_MASK_TYPE="$2"; shift 2 ;; + --codec) CODEC="$2"; shift 2 ;; + --width) WIDTH="$2"; shift 2 ;; + --height) HEIGHT="$2"; shift 2 ;; + --fps) FPS="$2"; shift 2 ;; + --no-shaping) SHAPING=0; shift ;; + --debug-drop) DEBUG_DROP="$2"; SHAPING=0; shift 2 ;; + --server-bin) SERVER_BIN="$2"; shift 2 ;; + --skip-build) SKIP_BUILD=1; shift ;; + -h|--help) sed -n '2,30p' "$0" | sed 's/^# \{0,1\}//'; exit 0 ;; + *) die "unknown argument: $1" ;; + esac +done + +case "$MODE" in + fec|baseline|ab) ;; + *) die "--mode must be fec, baseline or ab" ;; +esac + +case "$(uname -s)" in + Darwin) SHAPER="$SCRIPT_DIR/shape_macos.sh"; LOOPBACK_IF="lo0" ;; + Linux) SHAPER="$SCRIPT_DIR/shape_linux.sh"; LOOPBACK_IF="lo" ;; + *) die "unsupported platform $(uname -s)" ;; +esac + +# ---------- preflight ---------- + +command -v go >/dev/null || die "go not found" +command -v cargo >/dev/null || die "cargo not found" +command -v curl >/dev/null || die "curl not found" +python3 -c 'import matplotlib' 2>/dev/null || die "python3 with matplotlib required (pip3 install matplotlib)" +[ -d "$RUST_SDKS_DIR/examples/local_video" ] || die "rust-sdks not found at $RUST_SDKS_DIR (set RUST_SDKS_DIR)" + +if [ "$SHAPING" = "1" ]; then + log "shaping requires sudo, validating credentials..." + sudo -v || die "sudo required for traffic shaping (or pass --no-shaping)" + # keep the sudo timestamp alive for long runs + ( while true; do sudo -n true 2>/dev/null; sleep 60; done ) & + SUDO_KEEPALIVE_PID=$! +fi + +if [ -z "$OUT_DIR" ]; then + OUT_DIR="$SCRIPT_DIR/out/$(date +%Y%m%d_%H%M%S)" +fi +mkdir -p "$OUT_DIR" +log "output directory: $OUT_DIR" + +# ---------- builds ---------- + +if [ "$SKIP_BUILD" = "1" ]; then + [ -n "$SERVER_BIN" ] || die "--skip-build requires --server-bin" + [ -x "$SERVER_BIN" ] || die "server binary not found at $SERVER_BIN" + log "skipping builds, using prebuilt $SERVER_BIN" +else + SERVER_BIN="${SERVER_BIN:-$OUT_DIR/livekit-server}" + log "building livekit-server..." + (cd "$REPO_ROOT" && go build -o "$SERVER_BIN" ./cmd/server) || die "server build failed" + + log "building local_video examples (release)..." + (cd "$RUST_SDKS_DIR" && cargo build --release -p local_video -F desktop --bin publisher --bin subscriber) \ + || die "example build failed" +fi +PUBLISHER_BIN="$RUST_SDKS_DIR/target/release/publisher" +SUBSCRIBER_BIN="$RUST_SDKS_DIR/target/release/subscriber" +[ -x "$PUBLISHER_BIN" ] || die "publisher binary not found, run once without --skip-build first" +[ -x "$SUBSCRIBER_BIN" ] || die "subscriber binary not found, run once without --skip-build first" + +# ---------- process management ---------- + +SERVER_PID="" +SUBSCRIBER_PID="" +PUBLISHER_PID="" +PROM_POLL_PID="" +SHAPER_PID="" + +stop_pid() { + local pid="$1" sig="${2:-TERM}" + if [ -n "$pid" ] && kill -0 "$pid" 2>/dev/null; then + kill -"$sig" "$pid" 2>/dev/null + for _ in 1 2 3 4 5 6 7 8 9 10; do + kill -0 "$pid" 2>/dev/null || return 0 + sleep 0.5 + done + kill -KILL "$pid" 2>/dev/null + fi +} + +stop_shaper() { + if [ -n "$SHAPER_PID" ] && kill -0 "$SHAPER_PID" 2>/dev/null; then + sudo kill -TERM "$SHAPER_PID" 2>/dev/null + sleep 2 + fi + SHAPER_PID="" + sudo "$SHAPER" stop >/dev/null 2>&1 || true +} + +cleanup() { + trap - EXIT INT TERM + log "cleaning up processes" + [ "$SHAPING" = "1" ] && stop_shaper + stop_pid "$PROM_POLL_PID" + stop_pid "$PUBLISHER_PID" INT + stop_pid "$SUBSCRIBER_PID" INT + stop_pid "$SERVER_PID" + [ -n "${SUDO_KEEPALIVE_PID:-}" ] && kill "$SUDO_KEEPALIVE_PID" 2>/dev/null + exit "${1:-1}" +} +trap 'cleanup 1' INT TERM +trap 'cleanup $?' EXIT + +# ---------- single run ---------- + +run_one() { + local mode="$1" + local run_dir="$OUT_DIR/$mode" + mkdir -p "$run_dir" + log "=== $mode run: ${DURATION}s ===" + + # media is pinned to the loopback interface: it keeps every packet on the + # shaped path and, on macOS, avoids the application firewall silently + # dropping inbound UDP for unsigned freshly-built binaries + cat > "$run_dir/server.yaml" < "$run_dir/server.log" 2>&1 & + SERVER_PID=$! + + for i in $(seq 1 40); do + curl -s -o /dev/null --max-time 1 "http://127.0.0.1:$SIGNAL_PORT" && break + kill -0 "$SERVER_PID" 2>/dev/null || die "server exited early, see $run_dir/server.log" + [ "$i" = "40" ] && die "server did not become reachable" + sleep 0.5 + done + log "server is up (pid $SERVER_PID)" + + local conn_args="--url ws://127.0.0.1:$SIGNAL_PORT --api-key $API_KEY --api-secret $API_SECRET --room-name $ROOM_NAME" + + log "starting subscriber (headless)" + RUST_LOG=info "$SUBSCRIBER_BIN" $conn_args \ + --identity fec-sub --headless --log-frames "$run_dir/frames.csv" \ + > "$run_dir/subscriber.log" 2>&1 & + SUBSCRIBER_PID=$! + sleep 2 + + local fec_args="" + if [ "$mode" = "fec" ]; then + fec_args="--flex-fec --fec-mask-type $FEC_MASK_TYPE" + if [ "$FEC_RATE" != "0" ]; then + fec_args="$fec_args --fec-protection-rate $FEC_RATE" + fi + fi + + log "starting publisher (test pattern ${WIDTH}x${HEIGHT}@${FPS} $CODEC${fec_args:+,$fec_args})" + RUST_LOG=info "$PUBLISHER_BIN" $conn_args \ + --identity fec-pub --test-pattern \ + --width "$WIDTH" --height "$HEIGHT" --fps "$FPS" --codec "$CODEC" \ + --attach-timestamp --attach-frame-id $fec_args \ + > "$run_dir/publisher.log" 2>&1 & + PUBLISHER_PID=$! + + log "waiting for media to flow..." + for i in $(seq 1 120); do + if [ -f "$run_dir/frames.csv" ] && [ "$(wc -l < "$run_dir/frames.csv")" -gt 30 ]; then + break + fi + kill -0 "$PUBLISHER_PID" 2>/dev/null || die "publisher exited early, see $run_dir/publisher.log" + kill -0 "$SUBSCRIBER_PID" 2>/dev/null || die "subscriber exited early, see $run_dir/subscriber.log" + [ "$i" = "120" ] && die "no frames received after 60s" + sleep 0.5 + done + + local t0 + t0=$(now_us) + log "media flowing, starting measurement (t0=$t0)" + cat > "$run_dir/meta.env" < "$run_dir/shaper.log" 2>&1 & + SHAPER_PID=$! + sleep 1 + kill -0 "$SHAPER_PID" 2>/dev/null || die "shaper failed to start, see $run_dir/shaper.log" + fi + + sleep "$DURATION" + + log "measurement done, stopping" + [ "$SHAPING" = "1" ] && stop_shaper + echo "T1_US=$(now_us)" >> "$run_dir/meta.env" + stop_pid "$PROM_POLL_PID"; PROM_POLL_PID="" + stop_pid "$PUBLISHER_PID" INT; PUBLISHER_PID="" + stop_pid "$SUBSCRIBER_PID" INT; SUBSCRIBER_PID="" + # the final flexfec stats log line is emitted when the publisher's buffers close + sleep 2 + grep -h "flexfec" "$run_dir/server.log" | tail -5 || true + stop_pid "$SERVER_PID"; SERVER_PID="" + sleep 1 +} + +# ---------- runs + report ---------- + +case "$MODE" in + fec) run_one fec ;; + baseline) run_one baseline ;; + ab) + run_one baseline + sleep 3 + run_one fec + ;; +esac + +log "generating report" +if [ "$MODE" = "ab" ]; then + python3 "$SCRIPT_DIR/plot_fec.py" --run "$OUT_DIR/baseline" --run "$OUT_DIR/fec" --out "$OUT_DIR" \ + | tee "$OUT_DIR/summary.txt" +else + python3 "$SCRIPT_DIR/plot_fec.py" --run "$OUT_DIR/$MODE" --out "$OUT_DIR" \ + | tee "$OUT_DIR/summary.txt" +fi + +log "done. results in $OUT_DIR" diff --git a/scripts/fec/shape_linux.sh b/scripts/fec/shape_linux.sh new file mode 100755 index 000000000..ba3700c02 --- /dev/null +++ b/scripts/fec/shape_linux.sh @@ -0,0 +1,112 @@ +#!/usr/bin/env bash +# Traffic shaper for the FlexFEC test harness (Linux, tc/netem). +# +# Applies packet loss to UDP traffic destined to the SFU's pinned media port on +# loopback, simulating a robot uplink over cellular: continuous low base loss +# (Gilbert-Elliott model for realistic loss correlation) plus periodic +# high-loss bursts. Burst on/off transitions are appended to an events file +# (wall-clock microseconds) so plots can shade the burst windows. +# +# Usage: +# sudo ./shape_linux.sh start --port 7882 --base-loss 0.02 \ +# --burst-loss 0.25 --burst-every 15 --burst-len 3 --events events.csv +# sudo ./shape_linux.sh stop +# +# `start` runs in the foreground until terminated, cleaning up on exit. +# `stop` force-cleans shaping state from a previous run. + +set -u + +DEV="lo" + +PORT=7882 +BASE_LOSS=0.02 +BURST_LOSS=0.25 +BURST_EVERY=15 +BURST_LEN=3 +EVENTS_FILE="" + +log() { echo "[shape_linux] $*" >&2; } + +now_us() { python3 -c 'import time; print(int(time.time() * 1e6))'; } + +record_event() { + if [ -n "$EVENTS_FILE" ]; then + echo "$(now_us),$1" >> "$EVENTS_FILE" + fi +} + +pct() { python3 -c "print($1 * 100)"; } + +apply_base_loss() { + # Gilbert-Elliott: p = chance of entering the bad state, r = chance of + # leaving it. p derived from the target average loss with r fixed at 30% + # gives short correlated loss runs typical for radio links. + local p + p=$(pct "$BASE_LOSS") + tc qdisc change dev $DEV parent 1:4 handle 40: netem loss gemodel "${p}%" 30% +} + +cleanup() { + trap - EXIT INT TERM + log "cleaning up" + tc qdisc del dev $DEV root 2>/dev/null + record_event "shaper_stopped" + log "done" +} + +start() { + trap cleanup EXIT INT TERM + + # 4-band prio qdisc: default TOS mapping never selects band 4, so only the + # filtered SFU-bound UDP flow passes through the netem child + tc qdisc add dev $DEV root handle 1: prio bands 4 priomap 1 2 2 2 1 2 0 0 1 1 1 1 1 1 1 1 || { + log "failed to add root qdisc (already shaped? try '$0 stop')" + exit 1 + } + tc qdisc add dev $DEV parent 1:4 handle 40: netem loss gemodel "$(pct "$BASE_LOSS")%" 30% + tc filter add dev $DEV parent 1: protocol ip prio 1 u32 \ + match ip protocol 17 0xff \ + match ip dport "$PORT" 0xffff \ + flowid 1:4 + + log "shaping active: udp dport $PORT, base loss $BASE_LOSS (gemodel), burst $BURST_LOSS for ${BURST_LEN}s every ${BURST_EVERY}s" + record_event "shaper_started base=$BASE_LOSS burst=$BURST_LOSS" + + # periodic burst loop + while true; do + sleep "$BURST_EVERY" + tc qdisc change dev $DEV parent 1:4 handle 40: netem loss "$(pct "$BURST_LOSS")%" + record_event "burst_on" + log "burst on ($BURST_LOSS)" + sleep "$BURST_LEN" + apply_base_loss + record_event "burst_off" + log "burst off ($BASE_LOSS)" + done +} + +CMD="${1:-}" +shift || true +while [ $# -gt 0 ]; do + case "$1" in + --port) PORT="$2"; shift 2 ;; + --base-loss) BASE_LOSS="$2"; shift 2 ;; + --burst-loss) BURST_LOSS="$2"; shift 2 ;; + --burst-every) BURST_EVERY="$2"; shift 2 ;; + --burst-len) BURST_LEN="$2"; shift 2 ;; + --events) EVENTS_FILE="$2"; shift 2 ;; + *) log "unknown argument: $1"; exit 1 ;; + esac +done + +if [ "$(id -u)" -ne 0 ]; then + log "ERROR: must run as root (sudo)" + exit 1 +fi + +case "$CMD" in + start) start ;; + stop) cleanup ;; + *) echo "usage: $0 {start|stop} [--port N] [--base-loss F] [--burst-loss F] [--burst-every S] [--burst-len S] [--events FILE]" >&2; exit 1 ;; +esac diff --git a/scripts/fec/shape_macos.sh b/scripts/fec/shape_macos.sh new file mode 100755 index 000000000..bc42433ee --- /dev/null +++ b/scripts/fec/shape_macos.sh @@ -0,0 +1,126 @@ +#!/usr/bin/env bash +# Traffic shaper for the FlexFEC test harness (macOS, dummynet via dnctl/pfctl). +# +# Applies packet loss to UDP traffic destined to the SFU's pinned media port on +# loopback, simulating a robot uplink over cellular: continuous low base loss +# plus periodic high-loss bursts. Burst on/off transitions are appended to an +# events file (wall-clock microseconds) so plots can shade the burst windows. +# +# Usage: +# sudo ./shape_macos.sh start --port 7882 --base-loss 0.02 \ +# --burst-loss 0.25 --burst-every 15 --burst-len 3 --events events.csv +# sudo ./shape_macos.sh stop +# +# `start` runs in the foreground until terminated, cleaning up on exit. +# `stop` force-cleans shaping state from a previous run. + +set -u + +ANCHOR="livekit_fec" +PIPE=1 +STATE_DIR="${TMPDIR:-/tmp}/livekit_fec_shaper" + +PORT=7882 +BASE_LOSS=0.02 +BURST_LOSS=0.25 +BURST_EVERY=15 +BURST_LEN=3 +EVENTS_FILE="" + +log() { echo "[shape_macos] $*" >&2; } + +now_us() { python3 -c 'import time; print(int(time.time() * 1e6))'; } + +record_event() { + if [ -n "$EVENTS_FILE" ]; then + echo "$(now_us),$1" >> "$EVENTS_FILE" + fi +} + +cleanup() { + trap - EXIT INT TERM + log "cleaning up" + pfctl -a "$ANCHOR" -F all 2>/dev/null + # restore the system ruleset, dropping our anchor attachment + pfctl -f /etc/pf.conf 2>/dev/null + dnctl -q flush 2>/dev/null + if [ -f "$STATE_DIR/pf_token" ]; then + pfctl -X "$(cat "$STATE_DIR/pf_token")" 2>/dev/null + rm -f "$STATE_DIR/pf_token" + fi + record_event "shaper_stopped" + log "done" +} + +start() { + mkdir -p "$STATE_DIR" + + if ! command -v dnctl >/dev/null; then + log "ERROR: dnctl not found. dummynet is unavailable on this system," + log "consider Network Link Conditioner or running the test on Linux." + exit 1 + fi + + trap cleanup EXIT INT TERM + + # configure the dummynet pipe with the base loss + dnctl pipe $PIPE config plr "$BASE_LOSS" || { log "dnctl failed"; exit 1; } + + # enable pf, keeping the reference token for clean disable + local token + token=$(pfctl -E 2>&1 | awk '/Token/ {print $NF}') + if [ -n "$token" ]; then + echo "$token" > "$STATE_DIR/pf_token" + fi + + # attach our dummynet anchor on top of the system ruleset + pfctl -q -f - < SFU leg) + echo "dummynet in quick proto udp from any to any port $PORT pipe $PIPE" | \ + pfctl -q -a "$ANCHOR" -f - + + log "shaping active: udp dport $PORT, base loss $BASE_LOSS, burst $BURST_LOSS for ${BURST_LEN}s every ${BURST_EVERY}s" + record_event "shaper_started base=$BASE_LOSS burst=$BURST_LOSS" + + # periodic burst loop + while true; do + sleep "$BURST_EVERY" + dnctl pipe $PIPE config plr "$BURST_LOSS" + record_event "burst_on" + log "burst on ($BURST_LOSS)" + sleep "$BURST_LEN" + dnctl pipe $PIPE config plr "$BASE_LOSS" + record_event "burst_off" + log "burst off ($BASE_LOSS)" + done +} + +CMD="${1:-}" +shift || true +while [ $# -gt 0 ]; do + case "$1" in + --port) PORT="$2"; shift 2 ;; + --base-loss) BASE_LOSS="$2"; shift 2 ;; + --burst-loss) BURST_LOSS="$2"; shift 2 ;; + --burst-every) BURST_EVERY="$2"; shift 2 ;; + --burst-len) BURST_LEN="$2"; shift 2 ;; + --events) EVENTS_FILE="$2"; shift 2 ;; + *) log "unknown argument: $1"; exit 1 ;; + esac +done + +if [ "$(id -u)" -ne 0 ]; then + log "ERROR: must run as root (sudo)" + exit 1 +fi + +case "$CMD" in + start) start ;; + stop) cleanup ;; + *) echo "usage: $0 {start|stop} [--port N] [--base-loss F] [--burst-loss F] [--burst-every S] [--burst-len S] [--events FILE]" >&2; exit 1 ;; +esac diff --git a/scripts/fec/sweep_fec.sh b/scripts/fec/sweep_fec.sh new file mode 100755 index 000000000..5b414de04 --- /dev/null +++ b/scripts/fec/sweep_fec.sh @@ -0,0 +1,152 @@ +#!/usr/bin/env bash +# Parameter sweep over the FlexFEC harness: runs run_fec_test.sh across a matrix +# of loss levels and FEC configurations, then aggregates all cells into a single +# comparison report (sweep_report.png + sweep_summary.csv + a markdown table). +# +# At each loss level it runs one baseline (no FEC) plus one FEC run per +# (rate, mask) combination, so every FEC point has a same-loss baseline to +# compare against. The server and example binaries are built once and reused. +# +# Usage: +# ./sweep_fec.sh [options] +# +# Options (defaults in brackets): +# --loss-mode {debug|shaped} loss mechanism [debug] +# debug = uniform packet drop inside the SFU, no sudo, --loss-list is % +# shaped = OS traffic shaping bursts (needs sudo), --loss-list is the +# burst loss fraction, base loss/cadence fixed by --base-loss etc. +# --loss-list "L1 L2 .." loss levels to sweep [debug: "2 5 10 15"] +# --fec-rate-list "R1 .." publisher FEC protection rates (percent) ["20 30 50"] +# --mask-list "M1 .." FEC mask types: random and/or bursty ["bursty"] +# --duration S measurement seconds per cell [60] +# --codec C video codec [h264] +# --base-loss F shaped mode: continuous base loss [0.01] +# --burst-every S shaped mode: seconds between bursts [15] +# --burst-len S shaped mode: burst duration [3] +# --out DIR output directory [scripts/fec/out/sweep_] +# +# Example matrices: +# # uniform-loss sweep (no sudo), 4 loss levels x 2 rates x 1 mask = 12 cells +# ./sweep_fec.sh --loss-list "2 5 10 15" --fec-rate-list "20 50" +# +# # shaped cellular-burst sweep, vary burst intensity and compare masks +# sudo -v && ./sweep_fec.sh --loss-mode shaped --loss-list "0.15 0.30 0.50" \ +# --fec-rate-list "30" --mask-list "random bursty" +# +# Runtime ~= cells * (~20s setup + duration). The example above is ~12 * 80s. + +set -u + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" +RUST_SDKS_DIR="${RUST_SDKS_DIR:-$(cd "$REPO_ROOT/.." && pwd)/rust-sdks}" +RUNNER="$SCRIPT_DIR/run_fec_test.sh" + +LOSS_MODE="debug" +LOSS_LIST="" +FEC_RATE_LIST="20 30 50" +MASK_LIST="bursty" +DURATION=60 +CODEC="h264" +BASE_LOSS=0.01 +BURST_EVERY=15 +BURST_LEN=3 +OUT_DIR="" + +log() { echo "[sweep] $*"; } +die() { echo "[sweep] ERROR: $*" >&2; exit 1; } + +while [ $# -gt 0 ]; do + case "$1" in + --loss-mode) LOSS_MODE="$2"; shift 2 ;; + --loss-list) LOSS_LIST="$2"; shift 2 ;; + --fec-rate-list) FEC_RATE_LIST="$2"; shift 2 ;; + --mask-list) MASK_LIST="$2"; shift 2 ;; + --duration) DURATION="$2"; shift 2 ;; + --codec) CODEC="$2"; shift 2 ;; + --base-loss) BASE_LOSS="$2"; shift 2 ;; + --burst-every) BURST_EVERY="$2"; shift 2 ;; + --burst-len) BURST_LEN="$2"; shift 2 ;; + --out) OUT_DIR="$2"; shift 2 ;; + -h|--help) sed -n '2,38p' "$0" | sed 's/^# \{0,1\}//'; exit 0 ;; + *) die "unknown argument: $1" ;; + esac +done + +case "$LOSS_MODE" in + debug) : "${LOSS_LIST:=2 5 10 15}" ;; + shaped) : "${LOSS_LIST:=0.15 0.30 0.50}" ;; + *) die "--loss-mode must be debug or shaped" ;; +esac + +command -v go >/dev/null || die "go not found" +command -v cargo >/dev/null || die "cargo not found" +python3 -c 'import matplotlib' 2>/dev/null || die "python3 with matplotlib required" +[ -d "$RUST_SDKS_DIR/examples/local_video" ] || die "rust-sdks not found at $RUST_SDKS_DIR" + +if [ "$LOSS_MODE" = "shaped" ]; then + sudo -v || die "shaped mode needs sudo (or use --loss-mode debug)" + ( while true; do sudo -n true 2>/dev/null; sleep 60; done ) & + SUDO_KEEPALIVE_PID=$! + trap '[ -n "${SUDO_KEEPALIVE_PID:-}" ] && kill "$SUDO_KEEPALIVE_PID" 2>/dev/null' EXIT +fi + +if [ -z "$OUT_DIR" ]; then + OUT_DIR="$SCRIPT_DIR/out/sweep_$(date +%Y%m%d_%H%M%S)" +fi +mkdir -p "$OUT_DIR" +log "output directory: $OUT_DIR" + +# build once, reuse across all cells +SERVER_BIN="$OUT_DIR/livekit-server" +log "building livekit-server (once)..." +(cd "$REPO_ROOT" && go build -o "$SERVER_BIN" ./cmd/server) || die "server build failed" +log "building local_video examples (once)..." +(cd "$RUST_SDKS_DIR" && cargo build --release -p local_video -F desktop --bin publisher --bin subscriber) \ + || die "example build failed" + +MANIFEST="$OUT_DIR/manifest.tsv" +printf 'leaf\tmode\tloss\tfec_rate\tfec_mask\n' > "$MANIFEST" + +run_cell() { + # run_cell + local mode="$1" cell="$2"; shift 2 + log "cell: $cell" + "$RUNNER" --mode "$mode" --skip-build --server-bin "$SERVER_BIN" \ + --duration "$DURATION" --codec "$CODEC" --out "$OUT_DIR/$cell" "$@" \ + > "$OUT_DIR/$cell.log" 2>&1 || { log "WARNING: cell $cell failed, see $OUT_DIR/$cell.log"; return 1; } +} + +loss_args() { + local loss="$1" + if [ "$LOSS_MODE" = "debug" ]; then + echo "--debug-drop $loss" + else + echo "--base-loss $BASE_LOSS --burst-loss $loss --burst-every $BURST_EVERY --burst-len $BURST_LEN" + fi +} + +CELL_COUNT=0 +for loss in $LOSS_LIST; do + la=$(loss_args "$loss") + + base_cell="cell_loss${loss}_baseline" + if run_cell baseline "$base_cell" $la; then + printf '%s\t%s\t%s\t%s\t%s\n' "$base_cell/baseline" baseline "$loss" "-" "-" >> "$MANIFEST" + fi + CELL_COUNT=$((CELL_COUNT + 1)) + + for rate in $FEC_RATE_LIST; do + for mask in $MASK_LIST; do + fec_cell="cell_loss${loss}_rate${rate}_${mask}" + if run_cell fec "$fec_cell" --fec-rate "$rate" --fec-mask-type "$mask" $la; then + printf '%s\t%s\t%s\t%s\t%s\n' "$fec_cell/fec" fec "$loss" "$rate" "$mask" >> "$MANIFEST" + fi + CELL_COUNT=$((CELL_COUNT + 1)) + done + done +done + +log "ran $CELL_COUNT cells, aggregating" +python3 "$SCRIPT_DIR/aggregate_fec.py" --sweep "$OUT_DIR" | tee "$OUT_DIR/sweep_summary.txt" +log "done. report in $OUT_DIR/sweep_report.png"