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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 31 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ Congestion Control (FileCC) are not implemented.
| ✅ | Live Congestion Control (LiveCC) |
| ✅ | NAK and Peridoc NAK |
| ✅ | Encryption |
| ✅ | Packet Filtering (FEC) |
| ❌ | Buffer mode |
| ❌ | Rendezvous Handshake |
| ❌ | File Transfer Congestion Control (FileCC) |
Expand Down Expand Up @@ -78,6 +79,36 @@ conn.Close()

In the `contrib/client` directory you'll find a complete example of a SRT client.

## Forward Error Correction (FEC)

Forward Error Correction (FEC) adds redundancy to the data stream, allowing the receiver to rebuild lost packets without needing them to be retransmitted. This is especially useful in high-latency or extremely lossy networks (such as satellite or microwave links) where traditional ARQ (retransmission) would take too long or be inefficient.

GoSRT implements FEC based on the [SMPTE 2022-1-2007 standard](https://github.com/Haivision/srt/blob/master/docs/features/packet-filtering-and-fec.md) via the SRT Packet Filtering framework. Both parties (sender and receiver) must configure FEC for it to be successfully negotiated.

You can enable FEC by providing a `PacketFilter` configuration string. The syntax is:
`fec,cols:<N>,rows:<N>,layout:<even|staircase>`

- **`cols`**: The number of columns in the FEC matrix (e.g., number of data packets before an FEC packet).
- **`rows`**: The number of rows in the FEC matrix.
- **`layout`**: The arrangement of packets (`even` or `staircase`).

### Example

```go
import "github.com/datarhei/gosrt"

config := srt.DefaultConfig()

// Enable FEC with 5 columns, 1 row, and an even layout.
// This means for every 5 data packets, 1 redundant FEC packet is sent.
config.PacketFilter = "fec,cols:5,rows:1,layout:even"

conn, err := srt.Dial("srt", "golang.org:6000", config)
if err != nil {
// handle error
}
```

## Listener example

```go
Expand Down
5 changes: 4 additions & 1 deletion config.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"fmt"
"net/url"
"strconv"
"strings"
"time"
)

Expand Down Expand Up @@ -694,7 +695,9 @@ func (c *Config) Validate() error {
}

if len(c.PacketFilter) != 0 {
return fmt.Errorf("config: PacketFilter are not supported")
if !strings.HasPrefix(strings.ToLower(c.PacketFilter), "fec") {
return fmt.Errorf("config: PacketFilter %s is not supported", c.PacketFilter)
}
}

if len(c.Passphrase) != 0 {
Expand Down
6 changes: 5 additions & 1 deletion conn_request.go
Original file line number Diff line number Diff line change
Expand Up @@ -462,7 +462,11 @@ func (req *connRequest) Accept() (Conn, error) {
req.handshake.SRTHS.SRTFlags.PERIODICNAK = true
req.handshake.SRTHS.SRTFlags.REXMITFLG = true
req.handshake.SRTHS.SRTFlags.STREAM = false
req.handshake.SRTHS.SRTFlags.PACKET_FILTER = false
req.handshake.SRTHS.SRTFlags.PACKET_FILTER = len(req.config.PacketFilter) > 0
if len(req.config.PacketFilter) > 0 {
req.handshake.HasFilter = true
req.handshake.PacketFilter = req.config.PacketFilter
}
req.handshake.SRTHS.RecvTSBPDDelay = recvTsbpdDelay
req.handshake.SRTHS.SendTSBPDDelay = sendTsbpdDelay
}
Expand Down
62 changes: 55 additions & 7 deletions connection.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"github.com/datarhei/gosrt/congestion"
"github.com/datarhei/gosrt/congestion/live"
"github.com/datarhei/gosrt/crypto"
"github.com/datarhei/gosrt/fec"
"github.com/datarhei/gosrt/packet"
)

Expand Down Expand Up @@ -198,6 +199,10 @@ type srtConn struct {
// Congestion control
recv congestion.Receiver
snd congestion.Sender

// FEC
fecGen *fec.Generator
fecRec *fec.Reconstructor

// context of all channels and routines
ctx context.Context
Expand Down Expand Up @@ -328,6 +333,16 @@ func newSRTConn(config srtConnConfig) *srtConn {
OnDeliver: c.pop,
})

if len(c.config.PacketFilter) > 0 {
fecCfg, err := fec.ParseConfig(c.config.PacketFilter)
if err == nil {
c.fecGen = fec.NewGenerator(fecCfg)
c.fecRec = fec.NewReconstructor(fecCfg)
} else {
c.log("connection:error", func() string { return fmt.Sprintf("invalid FEC config: %s", err) })
}
}

c.ctx, c.cancelCtx = context.WithCancel(context.Background())

go c.networkQueueReader(c.ctx)
Expand Down Expand Up @@ -581,6 +596,14 @@ func (c *srtConn) pop(p packet.Packet) {

// Send the packet on the wire
c.onSend(p)

if c.fecGen != nil && !p.Header().IsControlPacket {
fecPkts := c.fecGen.AddPacket(p)
for _, fecPkt := range fecPkts {
c.log("fec:send:dump", func() string { return fecPkt.Dump() })
c.onSend(fecPkt)
}
}
}

// networkQueueReader reads the packets from the network queue in order to process them.
Expand Down Expand Up @@ -689,7 +712,24 @@ func (c *srtConn) handlePacket(p packet.Packet) {
// "An FEC control packet is distinguished from a regular data packet by having
// its message number equal to 0. This value isn't normally used in SRT (message
// numbers start from 1, increment to a maximum, and then roll back to 1)."
if header.MessageNumber == 0 {
// Process FEC filter
if c.fecRec != nil {
recoveredPkts := c.fecRec.AddPacket(p)

if header.MessageNumber == 0 {
c.log("connection:filter", func() string { return "processed FEC filter control packet" })
for _, recPkt := range recoveredPkts {
c.log("connection:filter", func() string { return fmt.Sprintf("recovered packet %d", recPkt.Header().PacketSequenceNumber.Val()) })
c.handlePacket(recPkt)
}
return
}

for _, recPkt := range recoveredPkts {
c.log("connection:filter", func() string { return fmt.Sprintf("recovered packet %d", recPkt.Header().PacketSequenceNumber.Val()) })
c.handlePacket(recPkt)
}
} else if header.MessageNumber == 0 {
c.log("connection:filter", func() string { return "dropped FEC filter control packet" })
return
}
Expand Down Expand Up @@ -936,8 +976,12 @@ func (c *srtConn) handleHSRequest(p packet.Packet) {
return
}

if cif.SRTFlags.PACKET_FILTER {
c.log("control:recv:HSReq:error", func() string { return "PACKET_FILTER flag is set" })
if cif.SRTFlags.PACKET_FILTER && len(c.config.PacketFilter) == 0 {
c.log("control:recv:HSReq:error", func() string { return "Peer set PACKET_FILTER but local does not support it" })
c.close()
return
} else if !cif.SRTFlags.PACKET_FILTER && len(c.config.PacketFilter) > 0 {
c.log("control:recv:HSReq:error", func() string { return "Local requires PACKET_FILTER but peer did not set it" })
c.close()
return
}
Expand Down Expand Up @@ -1020,8 +1064,12 @@ func (c *srtConn) handleHSResponse(p packet.Packet) {
return
}

if cif.SRTFlags.PACKET_FILTER {
c.log("control:recv:HSReq:error", func() string { return "PACKET_FILTER flag is set" })
if cif.SRTFlags.PACKET_FILTER && len(c.config.PacketFilter) == 0 {
c.log("control:recv:HSRes:error", func() string { return "Peer set PACKET_FILTER but local does not support it" })
c.close()
return
} else if !cif.SRTFlags.PACKET_FILTER && len(c.config.PacketFilter) > 0 {
c.log("control:recv:HSRes:error", func() string { return "Local requires PACKET_FILTER but peer did not set it" })
c.close()
return
}
Expand Down Expand Up @@ -1305,8 +1353,8 @@ func (c *srtConn) sendHSRequest() {
TLPKTDROP: true, // must be set in live mode
PERIODICNAK: false, // not relevant for us as sender
REXMITFLG: true, // must alwasy be set
STREAM: false, // has been introducet in HSv5
PACKET_FILTER: false, // has been introducet in HSv5
STREAM: false,
PACKET_FILTER: len(c.config.PacketFilter) > 0,
},
RecvTSBPDDelay: 0,
SendTSBPDDelay: uint16(c.config.ReceiverLatency.Milliseconds()),
Expand Down
7 changes: 6 additions & 1 deletion dial.go
Original file line number Diff line number Diff line change
Expand Up @@ -378,7 +378,7 @@ func (dl *dialer) handleHandshake(p packet.Packet) {
PERIODICNAK: true,
REXMITFLG: true,
STREAM: false,
PACKET_FILTER: false,
PACKET_FILTER: len(dl.config.PacketFilter) > 0,
},
RecvTSBPDDelay: uint16(dl.config.ReceiverLatency.Milliseconds()),
SendTSBPDDelay: uint16(dl.config.PeerLatency.Milliseconds()),
Expand All @@ -387,6 +387,11 @@ func (dl *dialer) handleHandshake(p packet.Packet) {
cif.HasSID = true
cif.StreamId = dl.config.StreamId

if len(dl.config.PacketFilter) > 0 {
cif.HasFilter = true
cif.PacketFilter = dl.config.PacketFilter
}

if dl.crypto != nil {
cif.HasKM = true
cif.SRTKM = &packet.CIFKeyMaterialExtension{}
Expand Down
71 changes: 71 additions & 0 deletions fec/fec.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
package fec

import (
"fmt"
"strconv"
"strings"
)

// Config represents the parsed FEC configuration.
type Config struct {
Cols int
Rows int
Layout string // "even" or "staircase"
ARQ string // "always", "onreq", "never"
}

// ParseConfig parses a packet filter string into an FEC config.
// Example: "fec,cols:10,rows:5,layout:even"
func ParseConfig(filterStr string) (Config, error) {
cfg := Config{
Cols: 0,
Rows: 1, // default
Layout: "even",
ARQ: "always",
}

parts := strings.Split(filterStr, ",")
if len(parts) == 0 || strings.ToLower(parts[0]) != "fec" {
return cfg, fmt.Errorf("not an fec config")
}

for _, part := range parts[1:] {
kv := strings.SplitN(part, ":", 2)
if len(kv) != 2 {
continue
}
key := strings.ToLower(kv[0])
val := kv[1]

switch key {
case "cols":
cols, err := strconv.Atoi(val)
if err != nil || cols < 2 {
return cfg, fmt.Errorf("invalid cols: %s", val)
}
cfg.Cols = cols
case "rows":
rows, err := strconv.Atoi(val)
if err != nil {
return cfg, fmt.Errorf("invalid rows: %s", val)
}
cfg.Rows = rows
case "layout":
if val != "even" && val != "staircase" {
return cfg, fmt.Errorf("invalid layout: %s", val)
}
cfg.Layout = val
case "arq":
if val != "always" && val != "onreq" && val != "never" {
return cfg, fmt.Errorf("invalid arq: %s", val)
}
cfg.ARQ = val
}
}

if cfg.Cols < 2 {
return cfg, fmt.Errorf("cols must be specified and >= 2")
}

return cfg, nil
}
90 changes: 90 additions & 0 deletions fec/generator.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
package fec

import (
"sync"

"github.com/datarhei/gosrt/packet"
)

// Generator is responsible for generating FEC control packets.
type Generator struct {
config Config

mu sync.Mutex
buffer []packet.Packet
}

// NewGenerator creates a new FEC Generator.
func NewGenerator(cfg Config) *Generator {
return &Generator{
config: cfg,
buffer: make([]packet.Packet, 0, cfg.Cols),
}
}

// AddPacket processes an outgoing data packet and potentially generates an FEC control packet.
// Returns a slice of FEC control packets that are ready to be sent.
func (g *Generator) AddPacket(p packet.Packet) []packet.Packet {
g.mu.Lock()
defer g.mu.Unlock()

// Simplest 1D row-based FEC generation for illustration (to be expanded)
// We only XOR the payload bytes here.
clone := p.Clone()
g.buffer = append(g.buffer, clone)

if len(g.buffer) >= g.config.Cols {
// Generate one FEC control packet for the row
ctrl := g.generateRowFEC()

// Reset buffer for the next row
g.buffer = g.buffer[:0]

// Create the gosrt packet
fecPkt := WriteControlPacket(p.Header().DestinationSocketId, ctrl)
return []packet.Packet{fecPkt}
}

return nil
}

func (g *Generator) generateRowFEC() *ControlPacket {
if len(g.buffer) == 0 {
return nil
}

ctrl := &ControlPacket{
SNBase: g.buffer[0].Header().PacketSequenceNumber.Val(),
GroupIndex: -1, // -1 means row
TimestampRecov: 0,
FlagsRecov: 0,
LengthRecov: 0,
}

// Determine max payload size
maxLen := 0
for _, p := range g.buffer {
if int(p.Len()) > maxLen {
maxLen = int(p.Len())
}
}

ctrl.PayloadRecov = make([]byte, maxLen)

for _, p := range g.buffer {
ctrl.TimestampRecov ^= p.Header().Timestamp
// Flags: we'd XOR KK flags here
flags := byte(p.Header().KeyBaseEncryptionFlag.Val() & 0x03)
ctrl.FlagsRecov ^= flags

length := uint16(p.Len())
ctrl.LengthRecov ^= length

data := p.Data()
for i := 0; i < len(data); i++ {
ctrl.PayloadRecov[i] ^= data[i]
}
}

return ctrl
}
Loading
Loading