batch: move shared-arena Reset ownership from lanes to their owner

This commit is contained in:
JackDoan
2026-07-14 11:51:32 -05:00
parent b83893bad7
commit 327c8218e3
9 changed files with 163 additions and 90 deletions
+1 -1
View File
@@ -307,7 +307,7 @@ func (f *Interface) activate() error {
f.batchers[i] = batch.NewMultiCoalescer(f.queues[i], f.l, arena, caps.TSO, caps.USO) f.batchers[i] = batch.NewMultiCoalescer(f.queues[i], f.l, arena, caps.TSO, caps.USO)
} else { } else {
arena := batch.NewArena(batch.DefaultPassthroughArenaCap) arena := batch.NewArena(batch.DefaultPassthroughArenaCap)
f.batchers[i] = batch.NewPassthrough(f.queues[i], arena) f.batchers[i] = batch.NewPassthrough(f.queues[i], arena.Reserve, arena.Reset)
} }
} }
+7
View File
@@ -179,3 +179,10 @@ func (a *Arena) Reserve(sz int) []byte {
func (a *Arena) Reset() { func (a *Arena) Reset() {
a.buf = a.buf[:0] a.buf = a.buf[:0]
} }
// Reserver hands out an sz-byte slice valid until its Resetter runs.
type Reserver func(sz int) []byte
// Resetter clears all reservations held by a Reserver. Only the arena's
// owner holds one; lanes inside a MultiCoalescer get nil.
type Resetter func()
+9 -15
View File
@@ -27,14 +27,8 @@ type MultiCoalescer struct {
tcp *TCPCoalescer tcp *TCPCoalescer
udp *UDPCoalescer udp *UDPCoalescer
pt *Passthrough pt *Passthrough
// arena is owned by the Multi: lanes get only its Reserve (nil Resetter)
// arena is shared across every lane (constructor hands the same // and Flush resets it exactly once after every lane has drained.
// *Arena to TCP, UDP, and Passthrough), so there's exactly one
// backing slab per MultiCoalescer instance. Each lane's Flush calls
// Reset; the resets are idempotent because Multi.Flush drains lanes
// sequentially and never Reserves in between, so a later lane's
// slots stay readable across an earlier lane's Reset (the underlying
// bytes are still alive — Reset only re-slices len to 0).
arena *Arena arena *Arena
} }
@@ -51,14 +45,14 @@ const DefaultMultiArenaCap = initialSlots * 65535
// pre-sizes it via NewArena so the hot path never allocates. // pre-sizes it via NewArena so the hot path never allocates.
func NewMultiCoalescer(w io.Writer, l *slog.Logger, arena *Arena, tcpEnabled, udpEnabled bool) *MultiCoalescer { func NewMultiCoalescer(w io.Writer, l *slog.Logger, arena *Arena, tcpEnabled, udpEnabled bool) *MultiCoalescer {
m := &MultiCoalescer{ m := &MultiCoalescer{
pt: NewPassthrough(w, arena), pt: NewPassthrough(w, arena.Reserve, nil),
arena: arena, arena: arena,
} }
if tcpEnabled { if tcpEnabled {
m.tcp = NewTCPCoalescer(w, l, arena) m.tcp = NewTCPCoalescer(w, l, arena.Reserve, nil)
} }
if udpEnabled { if udpEnabled {
m.udp = NewUDPCoalescer(w, arena) m.udp = NewUDPCoalescer(w, arena.Reserve, nil)
} }
return m return m
} }
@@ -116,10 +110,9 @@ func (m *MultiCoalescer) Commit(pkt []byte) error {
return m.pt.Commit(pkt) return m.pt.Commit(pkt)
} }
// Flush drains every lane in a fixed order: TCP, UDP, passthrough. Errors // Flush drains every lane in a fixed order TCP, UDP, passthrough — then
// from a lane do not stop subsequent lanes from flushing, we keep // resets the shared arena once. A lane error doesn't stop the remaining
// draining and return the first observed error so a single bad packet // lanes; the joined errors are returned.
// doesn't strand the others.
func (m *MultiCoalescer) Flush() error { func (m *MultiCoalescer) Flush() error {
var errs []error var errs []error
if m.tcp != nil { if m.tcp != nil {
@@ -135,5 +128,6 @@ func (m *MultiCoalescer) Flush() error {
if err := m.pt.Flush(); err != nil { if err := m.pt.Flush(); err != nil {
errs = append(errs, err) errs = append(errs, err)
} }
m.arena.Reset()
return errors.Join(errs...) return errors.Join(errs...)
} }
+21 -11
View File
@@ -8,11 +8,11 @@ import (
// Passthrough is a RxBatcher that doesn't batch anything, it just accumulates and then sends packets. // Passthrough is a RxBatcher that doesn't batch anything, it just accumulates and then sends packets.
type Passthrough struct { type Passthrough struct {
out io.Writer out io.Writer
slots [][]byte slots [][]byte
// arena is injected; see TCPCoalescer.arena for the contract. reserver Reserver
arena *Arena resetter Resetter
cursor int cursor int
} }
const passthroughBaseNumSlots = 128 const passthroughBaseNumSlots = 128
@@ -21,16 +21,17 @@ const passthroughBaseNumSlots = 128
// standalone Passthrough batcher: 128 slots × udp.MTU ≈ 1.1 MiB. // standalone Passthrough batcher: 128 slots × udp.MTU ≈ 1.1 MiB.
const DefaultPassthroughArenaCap = passthroughBaseNumSlots * udp.MTU const DefaultPassthroughArenaCap = passthroughBaseNumSlots * udp.MTU
func NewPassthrough(w io.Writer, arena *Arena) *Passthrough { func NewPassthrough(w io.Writer, reserver Reserver, resetter Resetter) *Passthrough {
return &Passthrough{ return &Passthrough{
out: w, out: w,
slots: make([][]byte, 0, passthroughBaseNumSlots), slots: make([][]byte, 0, passthroughBaseNumSlots),
arena: arena, reserver: reserver,
resetter: resetter,
} }
} }
func (p *Passthrough) Reserve(sz int) []byte { func (p *Passthrough) Reserve(sz int) []byte {
return p.arena.Reserve(sz) return p.reserver(sz)
} }
func (p *Passthrough) Commit(pkt []byte) error { func (p *Passthrough) Commit(pkt []byte) error {
@@ -38,7 +39,17 @@ func (p *Passthrough) Commit(pkt []byte) error {
return nil return nil
} }
// Flush drains every queued packet and calls the configured Resetter
func (p *Passthrough) Flush() error { func (p *Passthrough) Flush() error {
firstErr := p.drain()
if p.resetter != nil {
p.resetter()
}
return firstErr
}
// drain writes out every queued packet and clears the slot list.
func (p *Passthrough) drain() error {
var firstErr error var firstErr error
for _, s := range p.slots { for _, s := range p.slots {
_, err := p.out.Write(s) _, err := p.out.Write(s)
@@ -48,6 +59,5 @@ func (p *Passthrough) Flush() error {
} }
clear(p.slots) clear(p.slots)
p.slots = p.slots[:0] p.slots = p.slots[:0]
p.arena.Reset()
return firstErr return firstErr
} }
+17 -11
View File
@@ -83,22 +83,19 @@ type TCPCoalescer struct {
// at is removed/sealed. // at is removed/sealed.
lastSlot *coalesceSlot lastSlot *coalesceSlot
pool []*coalesceSlot // free list for reuse pool []*coalesceSlot // free list for reuse
reserver Reserver
// arena is injected; the coalescer borrows slices from it via Reserve resetter Resetter
// and tells it to release them via Reset on Flush. When wrapped in l *slog.Logger
// MultiCoalescer the same *Arena is shared with the other lanes so
// there's exactly one backing slab per Multi instance.
arena *Arena
l *slog.Logger
} }
func NewTCPCoalescer(w io.Writer, l *slog.Logger, arena *Arena) *TCPCoalescer { func NewTCPCoalescer(w io.Writer, l *slog.Logger, reserver Reserver, resetter Resetter) *TCPCoalescer {
c := &TCPCoalescer{ c := &TCPCoalescer{
plainW: w, plainW: w,
slots: make([]*coalesceSlot, 0, initialSlots), slots: make([]*coalesceSlot, 0, initialSlots),
openSlots: make(map[flowKey]*coalesceSlot, initialSlots), openSlots: make(map[flowKey]*coalesceSlot, initialSlots),
pool: make([]*coalesceSlot, 0, initialSlots), pool: make([]*coalesceSlot, 0, initialSlots),
arena: arena, reserver: reserver,
resetter: resetter,
l: l, l: l,
} }
if gw, ok := tio.SupportsGSO(w, tio.GSOProtoTCP); ok { if gw, ok := tio.SupportsGSO(w, tio.GSOProtoTCP); ok {
@@ -178,7 +175,7 @@ func (p parsedTCP) coalesceable() bool {
} }
func (c *TCPCoalescer) Reserve(sz int) []byte { func (c *TCPCoalescer) Reserve(sz int) []byte {
return c.arena.Reserve(sz) return c.reserver(sz)
} }
// Commit borrows pkt. The caller must keep pkt valid until the next Flush, // Commit borrows pkt. The caller must keep pkt valid until the next Flush,
@@ -259,6 +256,16 @@ func (c *TCPCoalescer) commitParsed(pkt []byte, info parsedTCP) error {
// doesn't hold up the rest. After Flush returns, borrowed payload slices // doesn't hold up the rest. After Flush returns, borrowed payload slices
// may be recycled. // may be recycled.
func (c *TCPCoalescer) Flush() error { func (c *TCPCoalescer) Flush() error {
first := c.drain()
if c.resetter != nil {
c.resetter()
}
return first
}
// drain emits every queued slot (reordering/merging coalesced runs first)
// and clears the slot state.
func (c *TCPCoalescer) drain() error {
c.reorderForFlush() c.reorderForFlush()
var first error var first error
for _, s := range c.slots { for _, s := range c.slots {
@@ -278,7 +285,6 @@ func (c *TCPCoalescer) Flush() error {
clear(c.openSlots) clear(c.openSlots)
c.lastSlot = nil c.lastSlot = nil
c.arena.Reset()
return first return first
} }
+2 -1
View File
@@ -71,7 +71,8 @@ func buildICMPv4() []byte {
// between batches, and reports per-packet cost. // between batches, and reports per-packet cost.
func runCommitBench(b *testing.B, pkts [][]byte, batchSize int) { func runCommitBench(b *testing.B, pkts [][]byte, batchSize int) {
b.Helper() b.Helper()
c := NewTCPCoalescer(nopTunWriter{}, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(nopTunWriter{}, test.NewLogger(), arena.Reserve, arena.Reset)
b.ReportAllocs() b.ReportAllocs()
b.SetBytes(int64(len(pkts[0]))) b.SetBytes(int64(len(pkts[0])))
b.ResetTimer() b.ResetTimer()
+56 -28
View File
@@ -128,7 +128,8 @@ const (
func TestCoalescerPassthroughWhenGSOUnavailable(t *testing.T) { func TestCoalescerPassthroughWhenGSOUnavailable(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: false} w := &fakeTunWriter{gsoEnabled: false}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pkt := buildTCPv4(1000, tcpAck, []byte("hello")) pkt := buildTCPv4(1000, tcpAck, []byte("hello"))
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -147,7 +148,8 @@ func TestCoalescerPassthroughWhenGSOUnavailable(t *testing.T) {
func TestCoalescerNonTCPPassthrough(t *testing.T) { func TestCoalescerNonTCPPassthrough(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pkt := make([]byte, 28) pkt := make([]byte, 28)
pkt[0] = 0x45 pkt[0] = 0x45
binary.BigEndian.PutUint16(pkt[2:4], 28) binary.BigEndian.PutUint16(pkt[2:4], 28)
@@ -167,7 +169,8 @@ func TestCoalescerNonTCPPassthrough(t *testing.T) {
func TestCoalescerSeedThenFlushAlone(t *testing.T) { func TestCoalescerSeedThenFlushAlone(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pkt := buildTCPv4(1000, tcpAck, make([]byte, 1000)) pkt := buildTCPv4(1000, tcpAck, make([]byte, 1000))
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -194,7 +197,8 @@ func TestCoalescerSeedThenFlushAlone(t *testing.T) {
func TestCoalescerCoalescesAdjacentACKs(t *testing.T) { func TestCoalescerCoalescesAdjacentACKs(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -234,7 +238,8 @@ func TestCoalescerCoalescesAdjacentACKs(t *testing.T) {
func TestCoalescerRejectsSeqGap(t *testing.T) { func TestCoalescerRejectsSeqGap(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -253,7 +258,8 @@ func TestCoalescerRejectsSeqGap(t *testing.T) {
func TestCoalescerRejectsFlagMismatch(t *testing.T) { func TestCoalescerRejectsFlagMismatch(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -274,7 +280,8 @@ func TestCoalescerRejectsFlagMismatch(t *testing.T) {
func TestCoalescerRejectsFIN(t *testing.T) { func TestCoalescerRejectsFIN(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
fin := buildTCPv4(1000, tcpAck|tcpFin, []byte("x")) fin := buildTCPv4(1000, tcpAck|tcpFin, []byte("x"))
if err := c.Commit(fin); err != nil { if err := c.Commit(fin); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -290,7 +297,8 @@ func TestCoalescerRejectsFIN(t *testing.T) {
func TestCoalescerShortLastSegmentClosesChain(t *testing.T) { func TestCoalescerShortLastSegmentClosesChain(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
full := make([]byte, 1200) full := make([]byte, 1200)
half := make([]byte, 500) half := make([]byte, 500)
if err := c.Commit(buildTCPv4(1000, tcpAck, full)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, full)); err != nil {
@@ -325,7 +333,8 @@ func TestCoalescerShortLastSegmentClosesChain(t *testing.T) {
func TestCoalescerPSHFinalizesChain(t *testing.T) { func TestCoalescerPSHFinalizesChain(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -355,7 +364,8 @@ func TestCoalescerPSHFinalizesChain(t *testing.T) {
// coalescer drops it the sender's push signal never reaches the receiver. // coalescer drops it the sender's push signal never reaches the receiver.
func TestCoalescerPropagatesPSHFromAppended(t *testing.T) { func TestCoalescerPropagatesPSHFromAppended(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Seed has no PSH; second segment carries PSH and seals the chain. // Seed has no PSH; second segment carries PSH and seals the chain.
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
@@ -383,7 +393,8 @@ func TestCoalescerPropagatesPSHFromAppended(t *testing.T) {
func TestCoalescerRejectsDifferentFlow(t *testing.T) { func TestCoalescerRejectsDifferentFlow(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
p1 := buildTCPv4(1000, tcpAck, pay) p1 := buildTCPv4(1000, tcpAck, pay)
p2 := buildTCPv4(2200, tcpAck, pay) p2 := buildTCPv4(2200, tcpAck, pay)
@@ -405,7 +416,8 @@ func TestCoalescerRejectsDifferentFlow(t *testing.T) {
func TestCoalescerRejectsIPOptions(t *testing.T) { func TestCoalescerRejectsIPOptions(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 500) pay := make([]byte, 500)
pkt := buildTCPv4(1000, tcpAck, pay) pkt := buildTCPv4(1000, tcpAck, pay)
// Bump IHL to 6 to simulate 4 bytes of IP options. Don't actually add // Bump IHL to 6 to simulate 4 bytes of IP options. Don't actually add
@@ -425,7 +437,8 @@ func TestCoalescerRejectsIPOptions(t *testing.T) {
func TestCoalescerCapBySegments(t *testing.T) { func TestCoalescerCapBySegments(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 512) pay := make([]byte, 512)
seq := uint32(1000) seq := uint32(1000)
for i := 0; i < tcpCoalesceMaxSegs+5; i++ { for i := 0; i < tcpCoalesceMaxSegs+5; i++ {
@@ -449,7 +462,8 @@ func TestCoalescerCapBySegments(t *testing.T) {
// flows coalesce independently in a single Flush. // flows coalesce independently in a single Flush.
func TestCoalescerMultipleFlowsInSameBatch(t *testing.T) { func TestCoalescerMultipleFlowsInSameBatch(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Flow A: sport 1000. Flow B: sport 3000. // Flow A: sport 1000. Flow B: sport 3000.
@@ -506,7 +520,8 @@ func TestCoalescerMultipleFlowsInSameBatch(t *testing.T) {
// writing passthrough packets synchronously. // writing passthrough packets synchronously.
func TestCoalescerPreservesArrivalOrder(t *testing.T) { func TestCoalescerPreservesArrivalOrder(t *testing.T) {
w := &orderedFakeWriter{gsoEnabled: true} w := &orderedFakeWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
// Sequence: coalesceable TCP, ICMP (passthrough), coalesceable TCP on // Sequence: coalesceable TCP, ICMP (passthrough), coalesceable TCP on
// a different flow. Expected emit order: gso(X), plain(ICMP), gso(Y). // a different flow. Expected emit order: gso(X), plain(ICMP), gso(Y).
pay := make([]byte, 1200) pay := make([]byte, 1200)
@@ -574,7 +589,8 @@ func stringSliceEq(a, b []string) bool {
// packet (SYN) mid-flow only flushes its own flow, not others. // packet (SYN) mid-flow only flushes its own flow, not others.
func TestCoalescerInterleavedFlowsPreserveOrdering(t *testing.T) { func TestCoalescerInterleavedFlowsPreserveOrdering(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Flow A two segments. // Flow A two segments.
@@ -679,7 +695,8 @@ func buildTCPv6(tcLow byte, seq uint32, flags byte, payload []byte) []byte {
// retains ECE on the wire. // retains ECE on the wire.
func TestCoalescerCoalescesEceFlow(t *testing.T) { func TestCoalescerCoalescesEceFlow(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
flags := byte(tcpAck | tcpEce) flags := byte(tcpAck | tcpEce)
if err := c.Commit(buildTCPv4(1000, flags, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, flags, pay)); err != nil {
@@ -708,7 +725,8 @@ func TestCoalescerCoalescesEceFlow(t *testing.T) {
// in-flow segment seeds a new slot rather than extending the prior burst. // in-flow segment seeds a new slot rather than extending the prior burst.
func TestCoalescerCwrSealsFlow(t *testing.T) { func TestCoalescerCwrSealsFlow(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -741,7 +759,8 @@ func TestCoalescerCwrSealsFlow(t *testing.T) {
// a CE-echoing window or none. // a CE-echoing window or none.
func TestCoalescerEceMismatchReseeds(t *testing.T) { func TestCoalescerEceMismatchReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4(1000, tcpAck|tcpEce, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck|tcpEce, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -771,7 +790,8 @@ func TestCoalescerEceMismatchReseeds(t *testing.T) {
// across the whole burst. // across the whole burst.
func TestCoalescerDifferingECNReseeds(t *testing.T) { func TestCoalescerDifferingECNReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4WithToS(ecnECT0, 1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4WithToS(ecnECT0, 1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -816,7 +836,8 @@ func TestCoalescerDifferingECNReseeds(t *testing.T) {
// codepoint, and neither may end up CE-marked. // codepoint, and neither may end up CE-marked.
func TestCoalescerECT0ThenECT1NoCE(t *testing.T) { func TestCoalescerECT0ThenECT1NoCE(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
if err := c.Commit(buildTCPv4WithToS(ecnECT0, 1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4WithToS(ecnECT0, 1000, tcpAck, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -846,7 +867,8 @@ func TestCoalescerECT0ThenECT1NoCE(t *testing.T) {
// six DSCP bits must match too. // six DSCP bits must match too.
func TestCoalescerDscpMismatchReseeds(t *testing.T) { func TestCoalescerDscpMismatchReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Same ECN (Not-ECT), different DSCP (0x10 vs 0x20 in upper 6 bits). // Same ECN (Not-ECT), different DSCP (0x10 vs 0x20 in upper 6 bits).
tosA := byte(0x10<<2) | ecnNotECT tosA := byte(0x10<<2) | ecnNotECT
@@ -869,7 +891,8 @@ func TestCoalescerDscpMismatchReseeds(t *testing.T) {
// TestCoalescerCoalescesEceFlow. // TestCoalescerCoalescesEceFlow.
func TestCoalescerIPv6CoalescesEceFlow(t *testing.T) { func TestCoalescerIPv6CoalescesEceFlow(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
flags := byte(tcpAck | tcpEce) flags := byte(tcpAck | tcpEce)
if err := c.Commit(buildTCPv6(0, 1000, flags, pay)); err != nil { if err := c.Commit(buildTCPv6(0, 1000, flags, pay)); err != nil {
@@ -900,7 +923,8 @@ func TestCoalescerIPv6CoalescesEceFlow(t *testing.T) {
// seen had the wire never reordered. // seen had the wire never reordered.
func TestCoalescerSortsReorderedSeedsAndMerges(t *testing.T) { func TestCoalescerSortsReorderedSeedsAndMerges(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Arrival order: seq 1000, 3400, 2200. The 3400 seeds a separate slot // Arrival order: seq 1000, 3400, 2200. The 3400 seeds a separate slot
// because 3400 != nextSeq=2200, then 2200 fails to extend the 3400 slot // because 3400 != nextSeq=2200, then 2200 fails to extend the 3400 slot
@@ -936,7 +960,8 @@ func TestCoalescerSortsReorderedSeedsAndMerges(t *testing.T) {
// without any cross-flow contamination. // without any cross-flow contamination.
func TestCoalescerSortAcrossFlowsMergesEachIndependently(t *testing.T) { func TestCoalescerSortAcrossFlowsMergesEachIndependently(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Flow A (sport 1000) seq 100, 1300; flow B (sport 3000) seq 500, 1700. // Flow A (sport 1000) seq 100, 1300; flow B (sport 3000) seq 500, 1700.
// Arrival: A.1300, B.1700, A.100, B.500 — every flow reordered. // Arrival: A.1300, B.1700, A.100, B.500 — every flow reordered.
@@ -987,7 +1012,8 @@ func TestCoalescerSortAcrossFlowsMergesEachIndependently(t *testing.T) {
// boundary by an arbitrary number of segments. // boundary by an arbitrary number of segments.
func TestCoalescerSortKeepsPSHBoundary(t *testing.T) { func TestCoalescerSortKeepsPSHBoundary(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// Seq 1000 (no PSH) + 2200 (PSH) → seal one slot with PSH set. // Seq 1000 (no PSH) + 2200 (PSH) → seal one slot with PSH set.
// Seq 3400 (no PSH) is contiguous to 3400 from seq 2200+1200; without // Seq 3400 (no PSH) is contiguous to 3400 from seq 2200+1200; without
@@ -1015,7 +1041,8 @@ func TestCoalescerSortKeepsPSHBoundary(t *testing.T) {
// is sorted/merged independently. // is sorted/merged independently.
func TestCoalescerSortKeepsPassthroughBarrier(t *testing.T) { func TestCoalescerSortKeepsPassthroughBarrier(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// First two segments seed S1 (then a 3400 reorder seeds S2). // First two segments seed S1 (then a 3400 reorder seeds S2).
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
@@ -1049,7 +1076,8 @@ func TestCoalescerSortKeepsPassthroughBarrier(t *testing.T) {
// 0x30, so ipHeadersMatch (comparing byte 1 fully) still splits them. // 0x30, so ipHeadersMatch (comparing byte 1 fully) still splits them.
func TestCoalescerIPv6DifferingECNReseeds(t *testing.T) { func TestCoalescerIPv6DifferingECNReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewTCPCoalescer(w, test.NewLogger(), NewArena(0)) arena := NewArena(0)
c := NewTCPCoalescer(w, test.NewLogger(), arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
// tcLow is the low 4 bits of TC; ECN occupies the bottom 2 of those. // tcLow is the low 4 bits of TC; ECN occupies the bottom 2 of those.
if err := c.Commit(buildTCPv6(ecnECT0, 1000, tcpAck, pay)); err != nil { if err := c.Commit(buildTCPv6(ecnECT0, 1000, tcpAck, pay)); err != nil {
+18 -7
View File
@@ -65,9 +65,8 @@ type UDPCoalescer struct {
slots []*udpSlot slots []*udpSlot
openSlots map[flowKey]*udpSlot openSlots map[flowKey]*udpSlot
pool []*udpSlot pool []*udpSlot
reserver Reserver
// arena is injected; see TCPCoalescer.arena for the contract. resetter Resetter
arena *Arena
} }
// NewUDPCoalescer wraps w. The caller is responsible for only constructing // NewUDPCoalescer wraps w. The caller is responsible for only constructing
@@ -75,13 +74,14 @@ type UDPCoalescer struct {
// the kernel may reject GSO_UDP_L4 writes. If w does not implement // the kernel may reject GSO_UDP_L4 writes. If w does not implement
// tio.GSOWriter at all (single-packet Queue), the coalescer degrades to // tio.GSOWriter at all (single-packet Queue), the coalescer degrades to
// plain Writes — same defensive shape as the TCP coalescer. // plain Writes — same defensive shape as the TCP coalescer.
func NewUDPCoalescer(w io.Writer, arena *Arena) *UDPCoalescer { func NewUDPCoalescer(w io.Writer, reserver Reserver, resetter Resetter) *UDPCoalescer {
c := &UDPCoalescer{ c := &UDPCoalescer{
plainW: w, plainW: w,
slots: make([]*udpSlot, 0, initialSlots), slots: make([]*udpSlot, 0, initialSlots),
openSlots: make(map[flowKey]*udpSlot, initialSlots), openSlots: make(map[flowKey]*udpSlot, initialSlots),
pool: make([]*udpSlot, 0, initialSlots), pool: make([]*udpSlot, 0, initialSlots),
arena: arena, reserver: reserver,
resetter: resetter,
} }
if gw, ok := tio.SupportsGSO(w, tio.GSOProtoUDP); ok { if gw, ok := tio.SupportsGSO(w, tio.GSOProtoUDP); ok {
c.gsoW = gw c.gsoW = gw
@@ -127,7 +127,7 @@ func parseUDP(pkt []byte) (parsedUDP, bool) {
} }
func (c *UDPCoalescer) Reserve(sz int) []byte { func (c *UDPCoalescer) Reserve(sz int) []byte {
return c.arena.Reserve(sz) return c.reserver(sz)
} }
// Commit borrows pkt. The caller must keep pkt valid until the next Flush. // Commit borrows pkt. The caller must keep pkt valid until the next Flush.
@@ -178,7 +178,19 @@ func (c *UDPCoalescer) commitParsed(pkt []byte, info parsedUDP) error {
return nil return nil
} }
// Flush drains every queued slot and calls the configured Resetter.
func (c *UDPCoalescer) Flush() error { func (c *UDPCoalescer) Flush() error {
first := c.drain()
if c.resetter != nil {
c.resetter()
}
return first
}
// drain emits every queued slot in arrival order and clears the slot state.
// It does NOT reset the arena: borrowed payload slices stay valid until the
// arena's owner recycles it.
func (c *UDPCoalescer) drain() error {
var first error var first error
for _, s := range c.slots { for _, s := range c.slots {
var err error var err error
@@ -195,7 +207,6 @@ func (c *UDPCoalescer) Flush() error {
clear(c.slots) clear(c.slots)
c.slots = c.slots[:0] c.slots = c.slots[:0]
clear(c.openSlots) clear(c.openSlots)
c.arena.Reset()
return first return first
} }
+32 -16
View File
@@ -60,7 +60,8 @@ func buildUDPv6(sport, dport uint16, payload []byte) []byte {
func TestUDPCoalescerPassthroughWhenGSOUnavailable(t *testing.T) { func TestUDPCoalescerPassthroughWhenGSOUnavailable(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: false} w := &fakeTunWriter{gsoEnabled: false}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pkt := buildUDPv4(1000, 53, make([]byte, 100)) pkt := buildUDPv4(1000, 53, make([]byte, 100))
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -78,7 +79,8 @@ func TestUDPCoalescerPassthroughWhenGSOUnavailable(t *testing.T) {
func TestUDPCoalescerNonUDPPassthrough(t *testing.T) { func TestUDPCoalescerNonUDPPassthrough(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
// ICMP packet // ICMP packet
pkt := make([]byte, 28) pkt := make([]byte, 28)
pkt[0] = 0x45 pkt[0] = 0x45
@@ -99,7 +101,8 @@ func TestUDPCoalescerNonUDPPassthrough(t *testing.T) {
func TestUDPCoalescerSeedThenFlushAlone(t *testing.T) { func TestUDPCoalescerSeedThenFlushAlone(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pkt := buildUDPv4(1000, 53, make([]byte, 800)) pkt := buildUDPv4(1000, 53, make([]byte, 800))
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -116,7 +119,8 @@ func TestUDPCoalescerSeedThenFlushAlone(t *testing.T) {
func TestUDPCoalescerCoalescesEqualSized(t *testing.T) { func TestUDPCoalescerCoalescesEqualSized(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
for i := 0; i < 3; i++ { for i := 0; i < 3; i++ {
if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil { if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil {
@@ -156,7 +160,8 @@ func TestUDPCoalescerCoalescesEqualSized(t *testing.T) {
// Last segment may be shorter, sealing the chain. // Last segment may be shorter, sealing the chain.
func TestUDPCoalescerShortLastSegmentSeals(t *testing.T) { func TestUDPCoalescerShortLastSegmentSeals(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
full := make([]byte, 1200) full := make([]byte, 1200)
tail := make([]byte, 600) tail := make([]byte, 600)
if err := c.Commit(buildUDPv4(1000, 53, full)); err != nil { if err := c.Commit(buildUDPv4(1000, 53, full)); err != nil {
@@ -189,7 +194,8 @@ func TestUDPCoalescerShortLastSegmentSeals(t *testing.T) {
// A larger-than-gsoSize packet cannot extend the slot — it reseeds. // A larger-than-gsoSize packet cannot extend the slot — it reseeds.
func TestUDPCoalescerLargerThanSeedReseeds(t *testing.T) { func TestUDPCoalescerLargerThanSeedReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
if err := c.Commit(buildUDPv4(1000, 53, make([]byte, 800))); err != nil { if err := c.Commit(buildUDPv4(1000, 53, make([]byte, 800))); err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -207,7 +213,8 @@ func TestUDPCoalescerLargerThanSeedReseeds(t *testing.T) {
// Different 5-tuples must not coalesce. // Different 5-tuples must not coalesce.
func TestUDPCoalescerDifferentFlowsKeepSeparate(t *testing.T) { func TestUDPCoalescerDifferentFlowsKeepSeparate(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pay := make([]byte, 800) pay := make([]byte, 800)
if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil { if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -238,7 +245,8 @@ func TestUDPCoalescerDifferentFlowsKeepSeparate(t *testing.T) {
// Caps at udpCoalesceMaxSegs. // Caps at udpCoalesceMaxSegs.
func TestUDPCoalescerCapsAtMaxSegs(t *testing.T) { func TestUDPCoalescerCapsAtMaxSegs(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pay := make([]byte, 100) pay := make([]byte, 100)
for i := 0; i < udpCoalesceMaxSegs+5; i++ { for i := 0; i < udpCoalesceMaxSegs+5; i++ {
if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil { if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil {
@@ -267,7 +275,8 @@ func TestUDPCoalescerCapsAtMaxSegs(t *testing.T) {
// trailing Not-ECT datagram seeds another. // trailing Not-ECT datagram seeds another.
func TestUDPCoalescerDifferingECNReseeds(t *testing.T) { func TestUDPCoalescerDifferingECNReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pay := make([]byte, 800) pay := make([]byte, 800)
pkt0 := buildUDPv4(1000, 53, pay) // ECN=00 (Not-ECT) pkt0 := buildUDPv4(1000, 53, pay) // ECN=00 (Not-ECT)
pkt1 := buildUDPv4(1000, 53, pay) pkt1 := buildUDPv4(1000, 53, pay)
@@ -298,7 +307,8 @@ func TestUDPCoalescerDifferingECNReseeds(t *testing.T) {
// IPv6 path: same flow, equal-sized → coalesced. // IPv6 path: same flow, equal-sized → coalesced.
func TestUDPCoalescerIPv6Coalesces(t *testing.T) { func TestUDPCoalescerIPv6Coalesces(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pay := make([]byte, 1200) pay := make([]byte, 1200)
for i := 0; i < 3; i++ { for i := 0; i < 3; i++ {
if err := c.Commit(buildUDPv6(1000, 53, pay)); err != nil { if err := c.Commit(buildUDPv6(1000, 53, pay)); err != nil {
@@ -334,7 +344,8 @@ func TestUDPCoalescerIPv6Coalesces(t *testing.T) {
// DSCP differences must reseed: udpHeadersMatch compares the full ToS byte. // DSCP differences must reseed: udpHeadersMatch compares the full ToS byte.
func TestUDPCoalescerDSCPMismatchReseeds(t *testing.T) { func TestUDPCoalescerDSCPMismatchReseeds(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pay := make([]byte, 800) pay := make([]byte, 800)
pkt0 := buildUDPv4(1000, 53, pay) pkt0 := buildUDPv4(1000, 53, pay)
pkt1 := buildUDPv4(1000, 53, pay) pkt1 := buildUDPv4(1000, 53, pay)
@@ -356,7 +367,8 @@ func TestUDPCoalescerDSCPMismatchReseeds(t *testing.T) {
// Fragmented IPv4 must not be coalesced. // Fragmented IPv4 must not be coalesced.
func TestUDPCoalescerFragmentedIPv4PassesThrough(t *testing.T) { func TestUDPCoalescerFragmentedIPv4PassesThrough(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pkt := buildUDPv4(1000, 53, make([]byte, 200)) pkt := buildUDPv4(1000, 53, make([]byte, 200))
binary.BigEndian.PutUint16(pkt[6:8], 0x2000) // MF=1 binary.BigEndian.PutUint16(pkt[6:8], 0x2000) // MF=1
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
@@ -377,7 +389,8 @@ func TestUDPCoalescerFragmentedIPv4PassesThrough(t *testing.T) {
// reach the GSO path. Regression: must not panic and must be written. // reach the GSO path. Regression: must not panic and must be written.
func TestUDPCoalescerZeroLengthPayloadPassesThrough(t *testing.T) { func TestUDPCoalescerZeroLengthPayloadPassesThrough(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pkt := buildUDPv4(1000, 53, nil) // UDP length 8, zero payload pkt := buildUDPv4(1000, 53, nil) // UDP length 8, zero payload
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -396,7 +409,8 @@ func TestUDPCoalescerZeroLengthPayloadPassesThrough(t *testing.T) {
// IPv6 zero-length UDP datagram: same passthrough contract as v4. // IPv6 zero-length UDP datagram: same passthrough contract as v4.
func TestUDPCoalescerZeroLengthPayloadIPv6PassesThrough(t *testing.T) { func TestUDPCoalescerZeroLengthPayloadIPv6PassesThrough(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pkt := buildUDPv6(1000, 53, nil) // UDP length 8, zero payload pkt := buildUDPv6(1000, 53, nil) // UDP length 8, zero payload
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -417,7 +431,8 @@ func TestUDPCoalescerZeroLengthPayloadIPv6PassesThrough(t *testing.T) {
// wire — per-flow arrival order (full, empty, full) must be preserved. // wire — per-flow arrival order (full, empty, full) must be preserved.
func TestUDPCoalescerZeroLengthMidFlowSealsAndPreservesOrder(t *testing.T) { func TestUDPCoalescerZeroLengthMidFlowSealsAndPreservesOrder(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
full := make([]byte, 800) full := make([]byte, 800)
if err := c.Commit(buildUDPv4(1000, 53, full)); err != nil { if err := c.Commit(buildUDPv4(1000, 53, full)); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -441,7 +456,8 @@ func TestUDPCoalescerZeroLengthMidFlowSealsAndPreservesOrder(t *testing.T) {
// IPv4 with options is not admissible (we require IHL=5). // IPv4 with options is not admissible (we require IHL=5).
func TestUDPCoalescerIPv4WithOptionsPassesThrough(t *testing.T) { func TestUDPCoalescerIPv4WithOptionsPassesThrough(t *testing.T) {
w := &fakeTunWriter{gsoEnabled: true} w := &fakeTunWriter{gsoEnabled: true}
c := NewUDPCoalescer(w, NewArena(0)) arena := NewArena(0)
c := NewUDPCoalescer(w, arena.Reserve, arena.Reset)
pkt := buildUDPv4(1000, 53, make([]byte, 200)) pkt := buildUDPv4(1000, 53, make([]byte, 200))
pkt[0] = 0x46 // IHL = 6 (24-byte IPv4 header — has options) pkt[0] = 0x46 // IHL = 6 (24-byte IPv4 header — has options)
if err := c.Commit(pkt); err != nil { if err := c.Commit(pkt); err != nil {