◐ Off-By-One · answer catalog

go-quic-multipath-congestion-recovery

1 answer(s)gogo1.26

A single global "in-flight" counter plus a single loss timer is unsafe in a multipath sender. ACKs arrive reordered, and a path can disappear while it still owns unacknowledged bytes. The fix is a per-path congestion/RTT domain, a bounded late-ACK window, and a migration grace period that releases the dead path's congestion budget immediately but delays retransmission long enough for a racing ACK to cancel it.

📦 Source in repository (JSON)

Answer

The solution is verified and saved to ~/mpquic/SOLUTION.md. Here is the complete document.


Multipath QUIC Sender: Correct Congestion & Recovery under Reordering and Abrupt Migration

Summary

A single global "in-flight" counter plus a single loss timer is unsafe in a multipath sender. ACKs arrive reordered, and a path can disappear while it still owns unacknowledged bytes. The fix is a per-path congestion/RTT domain, a bounded late-ACK window, and a migration grace period that releases the dead path's congestion budget immediately but delays retransmission long enough for a racing ACK to cancel it.

Files created and verified:

~/mpquic/go.mod
~/mpquic/mpquic.go
~/mpquic/mpquic_test.go
~/mpquic/SOLUTION.md

Root-cause analysis

1. A shared congestion/RTT/loss state cannot be per-path safe

If bytes_in_flight, cwnd, srtt, and loss detection are shared across paths, then: - a fast path's ACK advances the loss threshold for a slow path's packets, declaring them lost spuriously; - a slow path's RTT sample contaminates the fast path's PTO; - a loss on one path halves the window for both, so a congested path throttles a healthy one.

Rule: each path owns cwnd, ssthresh, bytes_in_flight, srtt/rttvar, largest_acked, and its outstanding packet set. Loss detection is evaluated only among packets sent on the same path that produced the ACK.

2. Reordered ACKs cause spurious retransmits

The packet-threshold rule (pn <= largest_acked - kPacketThreshold) can fire for a packet whose ACK is still in flight. If the sender has already queued a retransmission, it emits a duplicate datagram.

Rule: when an ACK later proves a "lost" packet was delivered, look it up and cancel the queued retransmission. That requires keeping the lost packet around briefly. The window is a bounded FIFO so memory cannot grow with connection age.

3. Abrupt migration causes deadlock or spurious retransmits

When a path is invalidated, its unacknowledged packets will never be ACKed by that path: - Deadlock: if those bytes remain counted against a connection- or path-level in-flight limit, the sender is permanently blocked waiting for ACKs that never come. - Spurious retransmit: if the sender immediately retransmits everything, it duplicates frames the peer may already have received (invalidation and data cross on the wire).

Rule (key insight): on invalidation 1. move every unacked packet to a migrating set and release the path's congestion budget now (kills deadlock); 2. do not retransmit immediately; set a PTO deadline (kills spurious retransmits for racing ACKs); 3. when the deadline expires, mark frames lost and reschedule on a surviving path (kills deadlock if the ACK never comes).

A migration is not a congestion signal: do not reduce cwnd for it.

4. Unbounded reordering memory

Storing every ACKed packet number, or expanding arbitrary ACK ranges, lets a peer exhaust memory.

Rule: test outstanding packets for ACK membership instead of expanding ranges; keep only unacked packets plus a bounded recent-loss window.


Exact fix

Concern Mechanism
Per-path safety pathState with its own reno controller, rttEstimator, outstanding map, largestAcked
Reordering tolerance per-path packet-threshold loss + bounded lostPackets/lostOrder FIFO, cancel-on-late-ACK
Migration InvalidatePath moves packets to migrating, zeroes bytesInFlight, sets a PTO deadline
Bounded memory membership testing against AckFrame; MaxReorderingPackets; MaxSendBuffer
                 send                      ack
  app frame ───────────▶ outstanding ───────────▶ gone
                             │  │
              packet-threshold│  │InvalidatePath
                     loss     │  │  (release budget)
                             ▼  ▼
                        retx queued ◀── lost window ──┐
                             │        (bounded)       │ late ACK: cancel retx
                             │  migration PTO expiry  │
                             └────────────────────────┘

Full source (mpquic.go)

// Package mpquic implements a QUIC-like multipath sender that stripes
// datagrams over two paths with independent RTT, loss and congestion state.
//
// The sender is written to be correct under two adversarial conditions that
// the task calls out explicitly:
//
//  1. ACK reordering. Acknowledgments may arrive out of order, so the
//     packet-threshold loss rule can fire for a packet whose ACK is still in
//     flight. The sender therefore keeps a *bounded* record of recently lost
//     packets so a late ACK can cancel the queued retransmission instead of
//     producing a spurious retransmit.
//
//  2. Abrupt path migration. When a path is invalidated, every unacknowledged
//     packet on it is moved into a bounded "migrating" set and its congestion
//     budget is released immediately. The frames are only rescheduled after a
//     PTO grace period, so an ACK that races with the migration cancels the
//     retransmission, while the grace timer guarantees forward progress if the
//     ACK never comes (no deadlock).
package mpquic

import (
    "errors"
    "sort"
    "time"
)

// PathID identifies a network path.
type PathID uint8

const (
    PathA PathID = 0
    PathB PathID = 1
)

const (
    // MaxDatagramSize is the fixed payload size used by the sender.
    MaxDatagramSize = 1200
    // MinCwnd is the floor for a path's congestion window.
    MinCwnd = 2 * MaxDatagramSize
)

// Config tunes the sender.
type Config struct {
    InitialRTT      time.Duration
    InitialCwnd     int
    PacketThreshold uint64  // reordering tolerance before a lower PN is declared lost
    TimeThreshold   float64 // multiplier of max(SRTT, latest RTT) for time-based loss
    MaxAckDelay     time.Duration
    // MaxReorderingPackets bounds the number of recently-lost packets kept
    // around so late ACKs can cancel retransmissions. This is the "bounded
    // reordering memory" of the sender.
    MaxReorderingPackets int
    // MaxConnectionInFlight, when > 0, caps the total bytes in flight across
    // all *active* paths. Migrating bytes are deliberately not counted.
    MaxConnectionInFlight int
    // MaxSendBuffer bounds the number of application frames waiting to be sent.
    MaxSendBuffer int
}

// DefaultConfig returns a conservative, production-shaped default.
func DefaultConfig() Config {
    return Config{
        InitialRTT:            100 * time.Millisecond,
        InitialCwnd:           10 * MaxDatagramSize,
        PacketThreshold:       3,
        TimeThreshold:         9.0 / 8.0,
        MaxAckDelay:           25 * time.Millisecond,
        MaxReorderingPackets:  512,
        MaxConnectionInFlight: 0,
        MaxSendBuffer:         1024,
    }
}

// Frame is an application data unit. ID must be unique per frame so that a
// retransmission can be matched back to the original.
type Frame struct {
    ID   uint64
    Data []byte
}

// Range is an inclusive range of packet numbers.
type Range struct {
    Smallest uint64
    Largest  uint64
}

// AckFrame is a (possibly non-contiguous) acknowledgment. Ranges are strictly
// below Largest. The sender tests outstanding packets for membership rather
// than expanding the ranges, so a malicious ACK cannot blow up memory.
type AckFrame struct {
    Largest uint64
    Ranges  []Range
}

func (a AckFrame) contains(pn uint64) bool {
    if pn == a.Largest {
        return true
    }
    for i := range a.Ranges {
        r := a.Ranges[i]
        if pn >= r.Smallest && pn <= r.Largest {
            return true
        }
    }
    return false
}

// ---------------------------------------------------------------------------
// RTT estimation
// ---------------------------------------------------------------------------

type rttEstimator struct {
    inited bool
    srtt   time.Duration
    rttvar time.Duration
    latest time.Duration
}

func (r *rttEstimator) update(sample time.Duration) {
    if sample <= 0 {
        return
    }
    if !r.inited {
        r.srtt = sample
        r.rttvar = sample / 2
        r.inited = true
    } else {
        diff := r.srtt - sample
        if diff < 0 {
            diff = -diff
        }
        r.rttvar = (3*r.rttvar + diff) / 4
        r.srtt = (7*r.srtt + sample) / 8
    }
    r.latest = sample
}

// pto is the probe timeout: SRTT + 4*RTTVAR + max_ack_delay.
func (r *rttEstimator) pto(maxAckDelay time.Duration) time.Duration {
    if !r.inited {
        return 0
    }
    p := r.srtt + 4*r.rttvar + maxAckDelay
    if p < 10*time.Millisecond {
        p = 10 * time.Millisecond
    }
    return p
}

// lossDelay is the time-threshold for declaring a packet lost.
func (r *rttEstimator) lossDelay(mult float64) time.Duration {
    if !r.inited {
        return 0
    }
    base := r.srtt
    if r.latest > base {
        base = r.latest
    }
    d := time.Duration(float64(base) * mult)
    if d < 10*time.Millisecond {
        d = 10 * time.Millisecond
    }
    return d
}

func (r *rttEstimator) srttOr(d time.Duration) time.Duration {
    if r.inited {
        return r.srtt
    }
    return d
}

// ---------------------------------------------------------------------------
// Per-path congestion control (NewReno)
// ---------------------------------------------------------------------------

type reno struct {
    cwnd          int
    ssthresh      int
    bytesInFlight int
    recoveryEnd   uint64
    inRecovery    bool
}

func newReno(initial int) reno {
    return reno{cwnd: initial, ssthresh: 1 << 30}
}

func (c *reno) onAck(pn uint64, ackedBytes int) {
    if c.inRecovery {
        if pn <= c.recoveryEnd {
            return // no window growth for ACKs inside a recovery period
        }
        c.inRecovery = false
    }
    if c.cwnd < c.ssthresh {
        c.cwnd += ackedBytes // slow start
    } else {
        inc := MaxDatagramSize * ackedBytes / c.cwnd
        if inc < 1 {
            inc = 1
        }
        c.cwnd += inc // congestion avoidance
    }
}

// onCongestion enters a recovery period ending at the largest PN sent when the
// loss was detected. Losses reported inside the same period do not reduce the
// window a second time.
func (c *reno) onCongestion(recoveryEnd uint64) {
    if c.inRecovery && recoveryEnd <= c.recoveryEnd {
        return
    }
    c.inRecovery = true
    c.recoveryEnd = recoveryEnd
    c.ssthresh = c.cwnd / 2
    if c.ssthresh < MinCwnd {
        c.ssthresh = MinCwnd
    }
    c.cwnd = c.ssthresh
}

// ---------------------------------------------------------------------------
// Sender state
// ---------------------------------------------------------------------------

type packet struct {
    pn        uint64
    path      PathID
    sentAt    time.Time
    size      int
    frame     Frame
    acked     bool
    lost      bool
    migrating bool
    deadline  time.Time
}

type pathState struct {
    id              PathID
    active          bool
    cc              reno
    rtt             rttEstimator
    outstanding     map[uint64]*packet
    largestSent     uint64
    largestAcked    uint64
    hasLargestAcked bool
}

// Outgoing is a datagram produced by Send.
type Outgoing struct {
    Path  PathID
    PN    uint64
    Frame Frame
    Retx  bool
}

// Stats are diagnostic counters, useful for verification.
type Stats struct {
    Sent                 int
    Retransmits          int
    SpuriousLosses       int
    MigrationRetransmits int
    CongestionEvents     int
    Acked                int
    AckFrames            int
}

// Sender is a multipath QUIC-like sender. All methods take an explicit
// timestamp so behavior is deterministic and testable.
type Sender struct {
    cfg    Config
    nextPN uint64

    paths map[PathID]*pathState

    byPN map[uint64]*packet

    // lostPackets is a bounded window of recently-lost packets. A late ACK
    // for one of them cancels the queued retransmission (spurious-loss
    // suppression). Eviction is FIFO, so memory is bounded.
    lostPackets map[uint64]*packet
    lostOrder   []uint64

    // migrating holds packets from an invalidated path. They are still in
    // byPN (so ACKs can find them) but no longer consume any path's budget.
    migrating []*packet

    newFrames []Frame
    retxOrder []uint64
    scheduled map[uint64]bool
    frames    map[uint64]Frame

    stats Stats
}

// New creates a Sender. Zero-valued config fields fall back to DefaultConfig.
func New(cfg Config) *Sender {
    d := DefaultConfig()
    if cfg.InitialRTT <= 0 {
        cfg.InitialRTT = d.InitialRTT
    }
    if cfg.InitialCwnd <= 0 {
        cfg.InitialCwnd = d.InitialCwnd
    }
    if cfg.PacketThreshold == 0 {
        cfg.PacketThreshold = d.PacketThreshold
    }
    if cfg.TimeThreshold <= 0 {
        cfg.TimeThreshold = d.TimeThreshold
    }
    if cfg.MaxAckDelay <= 0 {
        cfg.MaxAckDelay = d.MaxAckDelay
    }
    if cfg.MaxReorderingPackets <= 0 {
        cfg.MaxReorderingPackets = d.MaxReorderingPackets
    }
    if cfg.MaxSendBuffer <= 0 {
        cfg.MaxSendBuffer = d.MaxSendBuffer
    }
    return &Sender{
        cfg:         cfg,
        paths:       map[PathID]*pathState{},
        byPN:        map[uint64]*packet{},
        lostPackets: map[uint64]*packet{},
        scheduled:   map[uint64]bool{},
        frames:      map[uint64]Frame{},
    }
}

// AddPath activates a path with a fresh, independent congestion controller.
// Re-adding an already-active path is a no-op; re-adding a previously
// invalidated path installs a clean controller.
func (s *Sender) AddPath(id PathID) {
    if p, ok := s.paths[id]; ok && p.active {
        return
    }
    s.paths[id] = &pathState{
        id:          id,
        active:      true,
        cc:          newReno(s.cfg.InitialCwnd),
        outstanding: map[uint64]*packet{},
    }
}

// Enqueue queues an application frame for transmission.
func (s *Sender) Enqueue(f Frame) error {
    if len(s.newFrames) >= s.cfg.MaxSendBuffer {
        return errors.New("mpquic: send buffer full")
    }
    cp := make([]byte, len(f.Data))
    copy(cp, f.Data)
    f.Data = cp
    s.newFrames = append(s.newFrames, f)
    return nil
}

// InvalidatePath handles an abrupt migration: the path is retired and every
// unacknowledged packet is moved to the migrating set with a PTO deadline.
//
// Crucially the path's congestion budget is released at once. If those bytes
// remained "in flight", a connection-level or path-level window could be
// exhausted forever because no ACK will ever arrive for the dead path. The
// frames are not retransmitted yet: a PTO grace period lets a racing ACK
// cancel the retransmission.
func (s *Sender) InvalidatePath(now time.Time, id PathID) {
    p := s.paths[id]
    if p == nil || !p.active {
        return
    }
    p.active = false

    pto := p.rtt.pto(s.cfg.MaxAckDelay)
    if pto <= 0 {
        pto = 2*p.rtt.srttOr(s.cfg.InitialRTT) + s.cfg.MaxAckDelay
    }
    deadline := now.Add(pto)

    for pn, pkt := range p.outstanding {
        pkt.migrating = true
        pkt.deadline = deadline
        s.migrating = append(s.migrating, pkt)
        delete(p.outstanding, pn)
    }
    p.cc.bytesInFlight = 0
}

// CheckTimeouts drives the two timer-based transitions:
//   - migration grace expiry -> reschedule the frame on a surviving path;
//   - time-threshold loss on an active path -> reschedule and reduce cwnd.
func (s *Sender) CheckTimeouts(now time.Time) {
    // Migration deadlines.
    if len(s.migrating) > 0 {
        due := make([]*packet, 0, len(s.migrating))
        for _, pkt := range s.migrating {
            if !pkt.acked && !pkt.lost && !now.Before(pkt.deadline) {
                due = append(due, pkt)
            }
        }
        for _, pkt := range due {
            if s.declareLost(pkt, true) {
                s.stats.MigrationRetransmits++
            }
        }
    }

    // Time-threshold loss detection on active paths.
    for _, p := range s.paths {
        if !p.active {
            continue
        }
        delay := p.rtt.lossDelay(s.cfg.TimeThreshold)
        if delay <= 0 {
            delay = s.cfg.InitialRTT
        }
        var due []*packet
        for _, pkt := range p.outstanding {
            if now.Sub(pkt.sentAt) >= delay {
                due = append(due, pkt)
            }
        }
        for _, pkt := range due {
            s.declareLost(pkt, false)
        }
    }
}

func (s *Sender) removeMigrating(pkt *packet) {
    for i, p := range s.migrating {
        if p == pkt {
            s.migrating = append(s.migrating[:i], s.migrating[i+1:]...)
            return
        }
    }
}

// declareLost moves a packet out of flight, optionally applies a congestion
// event, remembers it in the bounded lost window, and schedules its frame for
// retransmission (once). It returns true if this call changed the packet's
// state.
func (s *Sender) declareLost(pkt *packet, fromMigration bool) bool {
    if pkt.acked || pkt.lost {
        return false
    }
    pkt.lost = true

    if pkt.migrating {
        s.removeMigrating(pkt)
    } else if p, ok := s.paths[pkt.path]; ok {
        delete(p.outstanding, pkt.pn)
        p.cc.bytesInFlight -= pkt.size
        if p.cc.bytesInFlight < 0 {
            p.cc.bytesInFlight = 0
        }
        if !fromMigration {
            // A migration is not a congestion signal: the path vanished,
            // it did not become congested.
            p.cc.onCongestion(p.largestSent)
            s.stats.CongestionEvents++
        }
    }

    delete(s.byPN, pkt.pn)
    s.rememberLost(pkt)

    if _, ok := s.scheduled[pkt.frame.ID]; !ok {
        s.scheduled[pkt.frame.ID] = true
        s.frames[pkt.frame.ID] = pkt.frame
        s.retxOrder = append(s.retxOrder, pkt.frame.ID)
    }
    return true
}

func (s *Sender) rememberLost(pkt *packet) {
    s.lostPackets[pkt.pn] = pkt
    s.lostOrder = append(s.lostOrder, pkt.pn)
    for len(s.lostOrder) > s.cfg.MaxReorderingPackets {
        old := s.lostOrder[0]
        s.lostOrder = s.lostOrder[1:]
        delete(s.lostPackets, old)
    }
}

// detectLostPath applies the packet-threshold rule within a single path. Only
// packets sent on the path that produced the ACK are considered, so
// cross-path reordering cannot cause spurious losses. Late ACKs are handled
// by the bounded lost window in OnAck.
func (s *Sender) detectLostPath(p *pathState) {
    if !p.hasLargestAcked || p.largestAcked < s.cfg.PacketThreshold {
        return
    }
    cutoff := p.largestAcked - s.cfg.PacketThreshold
    var due []*packet
    for pn, pkt := range p.outstanding {
        if pn <= cutoff {
            due = append(due, pkt)
        }
    }
    for _, pkt := range due {
        s.declareLost(pkt, false)
    }
}

// OnAck processes an ACK frame. It is safe against out-of-order and duplicate
// ACKs, and against ACKs that acknowledge packets already declared lost.
func (s *Sender) OnAck(now time.Time, ack AckFrame) {
    s.stats.AckFrames++

    // Collect live packets. Membership testing keeps memory proportional to
    // the send window regardless of how many ranges the ACK carries.
    var hits []*packet
    for pn, pkt := range s.byPN {
        if ack.contains(pn) {
            hits = append(hits, pkt)
        }
    }

    // Collect late ACKs for packets already declared lost. These are what
    // would otherwise become spurious retransmits.
    var lostHits []*packet
    for pn, pkt := range s.lostPackets {
        if ack.contains(pn) {
            lostHits = append(lostHits, pkt)
        }
    }
    for _, pkt := range lostHits {
        delete(s.scheduled, pkt.frame.ID)
        delete(s.lostPackets, pkt.pn)
        s.stats.SpuriousLosses++
    }

    // Process live ACKs in ascending PN order so largestAcked only grows.
    sort.Slice(hits, func(i, j int) bool { return hits[i].pn < hits[j].pn })

    prevLargest := map[PathID]uint64{}
    newLargest := map[PathID]*packet{}

    for _, pkt := range hits {
        if pkt.acked {
            continue
        }
        pkt.acked = true
        delete(s.byPN, pkt.pn)
        // A frame that was scheduled for retransmission is now known to
        // have been delivered: cancel the pending retransmit.
        delete(s.scheduled, pkt.frame.ID)

        p, ok := s.paths[pkt.path]
        if !ok {
            continue
        }
        if pkt.migrating {
            s.removeMigrating(pkt)
            // No congestion credit for an abandoned path.
        } else {
            delete(p.outstanding, pkt.pn)
            p.cc.bytesInFlight -= pkt.size
            if p.cc.bytesInFlight < 0 {
                p.cc.bytesInFlight = 0
            }
            p.cc.onAck(pkt.pn, pkt.size)
        }
        s.stats.Acked++

        if _, seen := prevLargest[pkt.path]; !seen {
            prevLargest[pkt.path] = p.largestAcked
        }
        if pkt.pn > p.largestAcked {
            p.largestAcked = pkt.pn
            p.hasLargestAcked = true
        }
        if m := newLargest[pkt.path]; m == nil || pkt.pn > m.pn {
            newLargest[pkt.path] = pkt
        }
    }

    // RTT sampling (largest newly acked PN per path) and loss detection.
    for pid, pkt := range newLargest {
        p := s.paths[pid]
        if p == nil {
            continue
        }
        if pkt.pn > prevLargest[pid] {
            p.rtt.update(now.Sub(pkt.sentAt))
        }
        s.detectLostPath(p)
    }
}

// pickPath selects an active path with room in its congestion window. The
// connection-level cap counts only active paths; migrating bytes are excluded
// precisely so an invalidated path cannot deadlock the connection.
func (s *Sender) pickPath() *pathState {
    if s.cfg.MaxConnectionInFlight > 0 &&
        s.connInFlight()+MaxDatagramSize > s.cfg.MaxConnectionInFlight {
        return nil
    }
    var best *pathState
    bestAvail := 0
    for _, p := range s.paths {
        if !p.active {
            continue
        }
        avail := p.cc.cwnd - p.cc.bytesInFlight
        if avail < MaxDatagramSize {
            continue
        }
        if best == nil || avail > bestAvail || (avail == bestAvail && p.id < best.id) {
            best = p
            bestAvail = avail
        }
    }
    return best
}

func (s *Sender) connInFlight() int {
    total := 0
    for _, p := range s.paths {
        if p.active {
            total += p.cc.bytesInFlight
        }
    }
    return total
}

// Send produces the next datagram, retransmissions first. It returns false
// when there is no frame to send or no congestion capacity on an active path.
func (s *Sender) Send(now time.Time) (*Outgoing, bool) {
    s.CheckTimeouts(now)

    var fr Frame
    isRetx := false
    for len(s.retxOrder) > 0 {
        id := s.retxOrder[0]
        s.retxOrder = s.retxOrder[1:]
        if _, ok := s.scheduled[id]; !ok {
            continue // canceled by a late ACK
        }
        delete(s.scheduled, id)
        f, ok := s.frames[id]
        if !ok {
            continue
        }
        fr = f
        isRetx = true
        break
    }
    if !isRetx {
        if len(s.newFrames) == 0 {
            return nil, false
        }
        fr = s.newFrames[0]
        s.newFrames = s.newFrames[1:]
    }

    p := s.pickPath()
    if p == nil {
        if isRetx {
            s.scheduled[fr.ID] = true
            s.retxOrder = append([]uint64{fr.ID}, s.retxOrder...)
        } else {
            s.newFrames = append([]Frame{fr}, s.newFrames...)
        }
        return nil, false
    }

    pn := s.nextPN
    s.nextPN++
    pkt := &packet{pn: pn, path: p.id, sentAt: now, size: MaxDatagramSize, frame: fr}
    p.outstanding[pn] = pkt
    p.cc.bytesInFlight += pkt.size
    p.largestSent = pn
    s.byPN[pn] = pkt

    s.stats.Sent++
    if isRetx {
        s.stats.Retransmits++
    }
    return &Outgoing{Path: p.id, PN: pn, Frame: fr, Retx: isRetx}, true
}

// Stats returns a snapshot of the sender counters.
func (s *Sender) Stats() Stats { return s.stats }

// BytesInFlight returns the congestion-controlled bytes in flight on a path.
func (s *Sender) BytesInFlight(id PathID) int {
    if p, ok := s.paths[id]; ok {
        return p.cc.bytesInFlight
    }
    return 0
}

// MigratingCount returns the number of packets awaiting the migration grace.
func (s *Sender) MigratingCount() int { return len(s.migrating) }

// PendingFrames returns the number of unsent application frames.
func (s *Sender) PendingFrames() int { return len(s.newFrames) }

// ScheduledRetransmits returns the number of retransmissions queued but not
// yet sent (including those already canceled but not lazily removed).
func (s *Sender) ScheduledRetransmits() int { return len(s.scheduled) }

Test suite (mpquic_test.go)

package mpquic

import (
    "math/rand"
    "sort"
    "testing"
    "time"
)

func mkFrame(id uint64) Frame { return Frame{ID: id, Data: []byte{byte(id)}} }

// sendAll drains the sender until it refuses to produce more datagrams.
func sendAll(s *Sender, now time.Time) []*Outgoing {
    var out []*Outgoing
    for {
        o, ok := s.Send(now)
        if !ok {
            return out
        }
        out = append(out, o)
    }
}

// TestBasicTwoPath checks the happy path: independent paths, all ACKed, no
// retransmits, no spurious losses.
func TestBasicTwoPath(t *testing.T) {
    s := New(DefaultConfig())
    s.AddPath(PathA)
    s.AddPath(PathB)

    now := time.Unix(0, 0)
    for i := uint64(0); i < 10; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }
    out := sendAll(s, now)
    if len(out) != 10 {
        t.Fatalf("want 10 datagrams, got %d", len(out))
    }

    s.OnAck(now.Add(time.Second), AckFrame{Largest: 9, Ranges: []Range{{Smallest: 0, Largest: 8}}})

    if s.Stats().Acked != 10 {
        t.Fatalf("acked = %d, want 10", s.Stats().Acked)
    }
    if s.Stats().Retransmits != 0 {
        t.Fatalf("retransmits = %d, want 0", s.Stats().Retransmits)
    }
    if s.Stats().SpuriousLosses != 0 {
        t.Fatalf("spurious losses = %d, want 0", s.Stats().SpuriousLosses)
    }
}

// TestReorderedAckCancelsSpuriousRetransmit exercises the packet-threshold
// rule firing for PN0..2 because only PN5 is ACKed, followed by a delayed ACK
// for PN0..4. The queued retransmissions must be canceled; no retransmit is
// allowed to leave the sender.
func TestReorderedAckCancelsSpuriousRetransmit(t *testing.T) {
    s := New(DefaultConfig())
    s.AddPath(PathA)

    now := time.Unix(0, 0)
    for i := uint64(0); i < 6; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }
    if got := len(sendAll(s, now)); got != 6 {
        t.Fatalf("want 6 datagrams, got %d", got)
    }

    // Only PN5 is acknowledged. PN0..2 fall below the packet threshold and
    // are (prematurely) declared lost, queueing retransmissions.
    s.OnAck(now.Add(10*time.Millisecond), AckFrame{Largest: 5})
    if s.ScheduledRetransmits() != 3 {
        t.Fatalf("scheduled retransmits = %d, want 3", s.ScheduledRetransmits())
    }

    // The delayed ACK for PN0..4 races in before any retransmit is sent.
    s.OnAck(now.Add(20*time.Millisecond), AckFrame{Largest: 4, Ranges: []Range{{Smallest: 0, Largest: 3}}})

    if s.Stats().SpuriousLosses != 3 {
        t.Fatalf("spurious losses = %d, want 3", s.Stats().SpuriousLosses)
    }
    if len(s.byPN) != 0 {
        t.Fatalf("%d packets still in flight, want 0", len(s.byPN))
    }
    if s.ScheduledRetransmits() != 0 {
        t.Fatalf("scheduled retransmits = %d, want 0", s.ScheduledRetransmits())
    }

    o, ok := s.Send(now.Add(30 * time.Millisecond))
    if ok && o.Retx {
        t.Fatalf("spurious retransmit produced: %+v", o)
    }
    if s.Stats().Retransmits != 0 {
        t.Fatalf("retransmits = %d, want 0", s.Stats().Retransmits)
    }
}

// TestMigrationAckRaceNoSpuriousRetransmit checks that an ACK which arrives
// after a path is invalidated but before the migration grace period expires
// cancels the retransmission.
func TestMigrationAckRaceNoSpuriousRetransmit(t *testing.T) {
    s := New(DefaultConfig())
    s.AddPath(PathA)

    now := time.Unix(0, 0)
    for i := uint64(0); i < 4; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }
    if got := len(sendAll(s, now)); got != 4 {
        t.Fatalf("want 4 datagrams on path A, got %d", got)
    }

    s.InvalidatePath(now, PathA)
    s.AddPath(PathB) // the surviving path

    if s.MigratingCount() != 4 {
        t.Fatalf("migrating = %d, want 4", s.MigratingCount())
    }
    if s.BytesInFlight(PathA) != 0 {
        t.Fatalf("bytes in flight on dead path = %d, want 0", s.BytesInFlight(PathA))
    }

    // The peer's ACK arrives just after migration.
    s.OnAck(now.Add(10*time.Millisecond), AckFrame{Largest: 3, Ranges: []Range{{Smallest: 0, Largest: 2}}})

    if s.Stats().MigrationRetransmits != 0 {
        t.Fatalf("migration retransmits = %d, want 0", s.Stats().MigrationRetransmits)
    }
    if s.MigratingCount() != 0 {
        t.Fatalf("migrating = %d, want 0", s.MigratingCount())
    }
    if o, ok := s.Send(now.Add(20 * time.Millisecond)); ok && o.Retx {
        t.Fatalf("spurious retransmit after migration race: %+v", o)
    }
}

// TestMigrationTimeoutPreventsDeadlock checks the other side of the race: if
// the ACK never arrives, the grace timer must reschedule every frame on the
// surviving path. A connection-level in-flight cap makes the deadlock that
// would occur if migrating bytes were still counted explicit.
func TestMigrationTimeoutPreventsDeadlock(t *testing.T) {
    cfg := DefaultConfig()
    cfg.MaxConnectionInFlight = 4 * MaxDatagramSize // exactly the 4 frames we send
    s := New(cfg)
    s.AddPath(PathA)

    now := time.Unix(0, 0)
    for i := uint64(0); i < 4; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }
    if got := len(sendAll(s, now)); got != 4 {
        t.Fatalf("want 4 datagrams, got %d", got)
    }

    s.InvalidatePath(now, PathA)
    s.AddPath(PathB)

    // If the dead path's bytes were still counted, the connection cap would
    // be full (4*1200) and Send could never make progress.
    if s.connInFlight() != 0 {
        t.Fatalf("connection in-flight after migration = %d, want 0", s.connInFlight())
    }

    later := now.Add(time.Second) // well past the PTO grace period
    s.CheckTimeouts(later)
    if s.Stats().MigrationRetransmits != 4 {
        t.Fatalf("migration retransmits = %d, want 4", s.Stats().MigrationRetransmits)
    }

    var retx int
    for {
        o, ok := s.Send(later)
        if !ok {
            break
        }
        if o.Path != PathB {
            t.Fatalf("retransmit went to path %d, want path B", o.Path)
        }
        if o.Retx {
            retx++
        }
        later = later.Add(time.Millisecond)
    }
    if retx != 4 {
        t.Fatalf("retransmits on path B = %d, want 4", retx)
    }
}

// TestCrossPathReorderingDoesNotCauseLoss confirms that an ACK for a packet on
// one path cannot declare packets on the other path lost, even when the global
// packet numbers are interleaved.
func TestCrossPathReorderingDoesNotCauseLoss(t *testing.T) {
    s := New(DefaultConfig())
    s.AddPath(PathA)
    s.AddPath(PathB)

    now := time.Unix(0, 0)
    for i := uint64(0); i < 8; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }
    outs := sendAll(s, now)
    if len(outs) != 8 {
        t.Fatalf("want 8 datagrams, got %d", len(outs))
    }

    // Acknowledge all path-B packets as one ACK (PNs 1, 3, 5, largest 7).
    // Path-A outstanding packets must be untouched: loss detection is
    // per-path, so interleaved global PNs cannot cause cross-path losses.
    beforeA := len(s.paths[PathA].outstanding)
    s.OnAck(now.Add(10*time.Millisecond), AckFrame{
        Largest: 7,
        Ranges:  []Range{{Smallest: 1, Largest: 1}, {Smallest: 3, Largest: 3}, {Smallest: 5, Largest: 5}},
    })
    if s.ScheduledRetransmits() != 0 {
        t.Fatalf("cross-path ACK scheduled retransmits: %d", s.ScheduledRetransmits())
    }
    if got := len(s.paths[PathA].outstanding); got != beforeA {
        t.Fatalf("path A outstanding changed: %d -> %d", beforeA, got)
    }
}

// TestBoundedReorderingMemory confirms the recently-lost window never grows
// past its configured bound even under a long run of losses.
func TestBoundedReorderingMemory(t *testing.T) {
    cfg := DefaultConfig()
    cfg.MaxReorderingPackets = 8
    cfg.InitialCwnd = 100 * MaxDatagramSize
    s := New(cfg)
    s.AddPath(PathA)

    now := time.Unix(0, 0)
    for i := uint64(0); i < 100; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }
    sendAll(s, now)

    // Force the whole window to be declared lost via the packet threshold.
    s.OnAck(now.Add(10*time.Millisecond), AckFrame{Largest: 99})
    if len(s.lostPackets) > cfg.MaxReorderingPackets {
        t.Fatalf("lostPackets = %d, exceeds bound %d", len(s.lostPackets), cfg.MaxReorderingPackets)
    }
    if len(s.lostOrder) > cfg.MaxReorderingPackets {
        t.Fatalf("lostOrder = %d, exceeds bound %d", len(s.lostOrder), cfg.MaxReorderingPackets)
    }
}

// TestRandomizedStress drives the sender with random ACKs, random path
// invalidations and random timeouts, then checks the safety invariants and
// convergence.
func TestRandomizedStress(t *testing.T) {
    rng := rand.New(rand.NewSource(7))
    s := New(DefaultConfig())
    s.AddPath(PathA)
    s.AddPath(PathB)
    now := time.Unix(0, 0)

    const total = 300
    for i := uint64(0); i < total; i++ {
        if err := s.Enqueue(mkFrame(i)); err != nil {
            t.Fatal(err)
        }
    }

    ackAll := func(at time.Time) {
        if len(s.byPN) == 0 {
            return
        }
        var pns []uint64
        for pn := range s.byPN {
            pns = append(pns, pn)
        }
        sort.Slice(pns, func(i, j int) bool { return pns[i] < pns[j] })
        ack := AckFrame{Largest: pns[len(pns)-1]}
        for _, pn := range pns[:len(pns)-1] {
            ack.Ranges = append(ack.Ranges, Range{Smallest: pn, Largest: pn})
        }
        s.OnAck(at, ack)
    }

    for round := 0; round < 3000; round++ {
        now = now.Add(time.Millisecond)
        for {
            if _, ok := s.Send(now); !ok {
                break
            }
        }
        if rng.Intn(40) == 0 {
            id := PathID(rng.Intn(2))
            s.InvalidatePath(now, id)
            s.AddPath(id) // revive cleanly
        }
        if len(s.byPN) > 0 && rng.Intn(4) != 0 {
            ackAll(now)
        }
        for _, p := range s.paths {
            if p.cc.bytesInFlight < 0 {
                t.Fatalf("negative bytesInFlight on path %d: %d", p.id, p.cc.bytesInFlight)
            }
        }
        if len(s.lostPackets) > s.cfg.MaxReorderingPackets {
            t.Fatalf("lostPackets unbounded: %d", len(s.lostPackets))
        }
    }

    for round := 0; round < 5000; round++ {
        if len(s.byPN) == 0 && len(s.newFrames) == 0 && len(s.migrating) == 0 && len(s.scheduled) == 0 {
            break
        }
        now = now.Add(10 * time.Millisecond)
        for {
            if _, ok := s.Send(now); !ok {
                break
            }
        }
        ackAll(now)
    }
    if len(s.byPN) != 0 || len(s.newFrames) != 0 || len(s.migrating) != 0 {
        t.Fatalf("stress did not converge: inFlight=%d new=%d migrating=%d scheduled=%d",
            len(s.byPN), len(s.newFrames), len(s.migrating), len(s.scheduled))
    }
}

Verification

Reproduce from scratch

mkdir -p ~/mpquic && cd ~/mpquic
# create go.mod, mpquic.go and mpquic_test.go from the listings above
go vet ./...
gofmt -l .            # must print nothing
go test -race -count=1 -v ./...

What each test proves

Test Property verified
TestBasicTwoPath Two independent paths, all packets ACKed, 0 retransmits, 0 spurious losses
TestReorderedAckCancelsSpuriousRetransmit A late ACK cancels 3 queued retransmissions; SpuriousLosses == 3, Retransmits == 0, nothing left in flight
TestMigrationAckRaceNoSpuriousRetransmit ACK racing the invalidation → MigrationRetransmits == 0, dead path's BytesInFlight == 0, MigratingCount == 0
TestMigrationTimeoutPreventsDeadlock ACK never arrives → PTO reschedules 4 frames on the surviving path, all sent on Path B; connection in-flight after migration is 0
TestCrossPathReorderingDoesNotCauseLoss An ACK for Path B packets never touches Path A's outstanding set
TestBoundedReorderingMemory lostPackets/lostOrder never exceed MaxReorderingPackets (8)
TestRandomizedStress 300 frames × 3000 rounds of random ACKs, random invalidations, random timeouts: no negative in-flight, bounded lost window, full convergence

Observed result

$ cd ~/mpquic && go vet ./... && gofmt -l . && go test -race -count=1 -v ./...
=== RUN   TestBasicTwoPath
--- PASS: TestBasicTwoPath (0.00s)
=== RUN   TestReorderedAckCancelsSpuriousRetransmit
--- PASS: TestReorderedAckCancelsSpuriousRetransmit (0.00s)
=== RUN   TestMigrationAckRaceNoSpuriousRetransmit
--- PASS: TestMigrationAckRaceNoSpuriousRetransmit (0.00s)
=== RUN   TestMigrationTimeoutPreventsDeadlock
--- PASS: TestMigrationTimeoutPreventsDeadlock (0.00s)
=== RUN   TestCrossPathReorderingDoesNotCauseLoss
--- PASS: TestCrossPathReorderingDoesNotCauseLoss (0.00s)
=== RUN   TestBoundedReorderingMemory
--- PASS: TestBoundedReorderingMemory (0.00s)
=== RUN   TestRandomizedStress
--- PASS: TestRandomizedStress (0.00s)
PASS
ok      mpquic  1.015s

Coverage: go test -cover ./... → coverage: 86.9% of statements.

Key assertions in code form

Spurious-retransmit suppression:

s.OnAck(now.Add(10*time.Millisecond), AckFrame{Largest: 5})          // queues 3 retx
s.OnAck(now.Add(20*time.Millisecond), AckFrame{Largest: 4,
        Ranges: []Range{{Smallest: 0, Largest: 3}}})                 // late ACK cancels them
if s.Stats().Retransmits != 0 { /* fail */ }

Migration deadlock avoidance:

s.InvalidatePath(now, PathA)      // bytesInFlight(PathA) -> 0 immediately
s.AddPath(PathB)
if s.connInFlight() != 0 { /* fail */ }   // dead path no longer blocks the connection
s.CheckTimeouts(now.Add(time.Second))     // PTO grace expires
// exactly 4 frames are retransmitted on PathB

Implementation notes / production hardening

Evidence & signatures

# Evidence
- Problem class: go-quic-multipath-congestion-recovery
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-10-05T04:09:01.414Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement a QUIC-like sender that stripes datagrams over two paths with independent RTT, loss, and congestion state, while packet acknowledgments may arrive reordered and path migration can invalidate one path abruptly. Preserve per-path congestion safety and bounded reordering memory, and define recovery behavior that avoids both spurious retransmits and deadlock when acknowledgments race with migration.", "environment": "go1.26", "language": "go", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "go-quic-multipath-congestion-recovery", "provider": "openrouter", "solved_at": "2026-10-05T04:09:01.415Z", "version": "1.26"}
Generated from the verified corpus · MIT licensedBack to the catalog