Skip to content
Merged
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
45 changes: 45 additions & 0 deletions internal/router/linkqueue_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package router

import (
"log/slog"
"testing"

"github.com/packethacking/net-sim/internal/audio"
"github.com/packethacking/net-sim/internal/config"
)

// A transmission far longer than the old fixed 3 s buffer must be delivered
// gap-free: net-sim models a continuous carrier and must never drop audio
// within a transmission. Dropping mid-transmission gaps the receiver's audio,
// collapses its carrier sense, and makes the far end key up on top of the
// in-progress transmission (a collision) — the bug this non-dropping FIFO
// fixes. Regression for the 3 s drop-on-overflow `linkQueue`.
func TestLinkQueueDeliversLongBurstGapFree(t *testing.T) {
q := newLinkQueue(
config.PortRef{NodeID: "a", PortID: "vhf"},
config.PortRef{NodeID: "b", PortID: "vhf"},
0, 0,
)

// ~10 s of audio — well past the old 3 s cap that used to drop.
n := 10 * audio.SampleRate / audio.BlockSamples
for i := 0; i < n; i++ {
blk := make(audio.Block, audio.BlockBytes)
blk[0] = byte(i)
blk[1] = byte(i >> 8)
q.push(blk, slog.Default())
}

for i := 0; i < n; i++ {
blk, ok := q.pop()
if !ok {
t.Fatalf("block %d of %d was dropped — the queue must never drop within a transmission", i, n)
}
if got := int(blk[0]) | int(blk[1])<<8; got != i {
t.Fatalf("FIFO order broken at block %d: got marker %d", i, got)
}
}
if _, ok := q.pop(); ok {
t.Fatal("queue should be empty after draining the whole transmission")
}
}
89 changes: 51 additions & 38 deletions internal/router/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -146,57 +146,70 @@ type txTracker struct {
lastBusyNanos atomic.Int64
}

// linkQueue is one source→destination link's audio buffer. The capacity
// is sized for ~3 s of audio at SampleRate / BlockSamples blocks per
// second — large enough to absorb a full samoyed TX burst (preamble +
// frame + postamble for typical AX.25) without dropping.
// maxQueueBlocks is a safety cap on a single link's audio backlog (~60 s).
// The channel is modelled as a CONTINUOUS CARRIER: blocks are never dropped
// within a transmission (see linkQueue.push). 60 s only triggers on a genuine
// runaway (a wedged rxFeeder), where dropping the head is the lesser evil — a
// real keyup is at most a few seconds.
const maxQueueBlocks = 60 * audio.SampleRate / audio.BlockSamples

// linkQueue is one source→destination link's audio buffer: a non-dropping FIFO
// of PCM blocks the source has transmitted, drained by the destination's
// rxFeeder at the channel sample rate (so the receiver hears the carrier in
// real time).
//
// A TNC bursts a whole keyup's worth of audio with no real-time pacing on its
// end (see txReader), so the buffer must hold an entire transmission and the
// rxFeeder plays it out at real time. It is non-dropping and grows as needed:
// dropping mid-transmission would gap the receiver's audio, collapse its DCD /
// carrier sense, and make the far end key up on top of an in-progress
// transmission — a collision that cannot happen on a real continuous-carrier
// channel. Memory is bounded in practice (a keyup is finite; the backing array
// is released once drained); only a pathological runaway past maxQueueBlocks
// ever drops.
type linkQueue struct {
src, dst config.PortRef
loss float64
noise float64
ch chan audio.Block

mu sync.Mutex
buf []audio.Block // FIFO; index 0 = oldest
}

func newLinkQueue(src, dst config.PortRef, loss, noise float64) *linkQueue {
const capBlocks = 3 * audio.SampleRate / audio.BlockSamples
return &linkQueue{
src: src,
dst: dst,
loss: loss,
noise: noise,
ch: make(chan audio.Block, capBlocks),
}
return &linkQueue{src: src, dst: dst, loss: loss, noise: noise}
}

// pushNonBlocking enqueues a block; drops oldest if full. Logs the drop.
func (q *linkQueue) pushNonBlocking(blk audio.Block, logger *slog.Logger) {
select {
case q.ch <- blk:
return
default:
}
// Full: drop the oldest block to make room. This indicates the
// downstream rxFeeder is falling behind — usually a sign something
// has stalled rather than a normal-operations event, so log it.
select {
case <-q.ch:
default:
}
select {
case q.ch <- blk:
default:
}
logger.Warn("audio queue overflow", "from", q.src, "to", q.dst)
// push appends a block to the link's buffer. It never blocks and — modelling a
// continuous carrier — never drops within a transmission, no matter how far the
// source TNC has run ahead of real time. The only drop is the maxQueueBlocks
// safety valve, which signals a stalled rxFeeder rather than normal operation,
// so it is logged.
func (q *linkQueue) push(blk audio.Block, logger *slog.Logger) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.buf) >= maxQueueBlocks {
q.buf[0] = nil
q.buf = q.buf[1:]
logger.Warn("audio queue safety cap reached (rxFeeder stalled?)", "from", q.src, "to", q.dst)
}
q.buf = append(q.buf, blk)
}

// pop returns one block if immediately available.
// pop returns one block if immediately available (FIFO, oldest first).
func (q *linkQueue) pop() (audio.Block, bool) {
select {
case b := <-q.ch:
return b, true
default:
q.mu.Lock()
defer q.mu.Unlock()
if len(q.buf) == 0 {
return nil, false
}
blk := q.buf[0]
q.buf[0] = nil // release the reference
q.buf = q.buf[1:]
if len(q.buf) == 0 {
q.buf = nil // release the backing array once drained
}
return blk, true
}

// Start spawns all samoyed children and begins routing audio.
Expand Down Expand Up @@ -481,7 +494,7 @@ func (r *Router) txReader(ctx context.Context, ref config.PortRef, c *tnc.Child)
}
}
for _, q := range outgoing {
q.pushNonBlocking(blk, r.logger)
q.push(blk, r.logger)
}
}
}
Expand Down
Loading