From 97eb3c635aca0dcfe8ab94167e340f5cf6e1cf0e Mon Sep 17 00:00:00 2001 From: JackDoan Date: Tue, 14 Jul 2026 11:51:32 -0500 Subject: [PATCH] batch: move shared-arena Reset ownership from lanes to their owner --- interface.go | 2 +- overlay/batch/coalesce_core.go | 7 ++ overlay/batch/multi_coalesce.go | 24 +++---- overlay/batch/passthrough.go | 32 +++++---- overlay/batch/tcp_coalesce.go | 28 ++++---- overlay/batch/tcp_coalesce_bench_test.go | 3 +- overlay/batch/tcp_coalesce_test.go | 84 ++++++++++++++++-------- overlay/batch/udp_coalesce.go | 25 +++++-- overlay/batch/udp_coalesce_test.go | 48 +++++++++----- 9 files changed, 163 insertions(+), 90 deletions(-) diff --git a/interface.go b/interface.go index ab7c19df..a23fd821 100644 --- a/interface.go +++ b/interface.go @@ -307,7 +307,7 @@ func (f *Interface) activate() error { f.batchers[i] = batch.NewMultiCoalescer(f.queues[i], f.l, arena, caps.TSO, caps.USO) } else { 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) } } diff --git a/overlay/batch/coalesce_core.go b/overlay/batch/coalesce_core.go index 80100a81..4cb54e8c 100644 --- a/overlay/batch/coalesce_core.go +++ b/overlay/batch/coalesce_core.go @@ -179,3 +179,10 @@ func (a *Arena) Reserve(sz int) []byte { func (a *Arena) Reset() { 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() diff --git a/overlay/batch/multi_coalesce.go b/overlay/batch/multi_coalesce.go index b6f76fe6..12328e65 100644 --- a/overlay/batch/multi_coalesce.go +++ b/overlay/batch/multi_coalesce.go @@ -27,14 +27,8 @@ type MultiCoalescer struct { tcp *TCPCoalescer udp *UDPCoalescer pt *Passthrough - - // arena is shared across every lane (constructor hands the same - // *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 is owned by the Multi: lanes get only its Reserve (nil Resetter) + // and Flush resets it exactly once after every lane has drained. arena *Arena } @@ -51,14 +45,14 @@ const DefaultMultiArenaCap = initialSlots * 65535 // 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 { m := &MultiCoalescer{ - pt: NewPassthrough(w, arena), + pt: NewPassthrough(w, arena.Reserve, nil), arena: arena, } if tcpEnabled { - m.tcp = NewTCPCoalescer(w, l, arena) + m.tcp = NewTCPCoalescer(w, l, arena.Reserve, nil) } if udpEnabled { - m.udp = NewUDPCoalescer(w, arena) + m.udp = NewUDPCoalescer(w, arena.Reserve, nil) } return m } @@ -116,10 +110,9 @@ func (m *MultiCoalescer) Commit(pkt []byte) error { return m.pt.Commit(pkt) } -// Flush drains every lane in a fixed order: TCP, UDP, passthrough. Errors -// from a lane do not stop subsequent lanes from flushing, we keep -// draining and return the first observed error so a single bad packet -// doesn't strand the others. +// Flush drains every lane in a fixed order — TCP, UDP, passthrough — then +// resets the shared arena once. A lane error doesn't stop the remaining +// lanes; the joined errors are returned. func (m *MultiCoalescer) Flush() error { var errs []error if m.tcp != nil { @@ -135,5 +128,6 @@ func (m *MultiCoalescer) Flush() error { if err := m.pt.Flush(); err != nil { errs = append(errs, err) } + m.arena.Reset() return errors.Join(errs...) } diff --git a/overlay/batch/passthrough.go b/overlay/batch/passthrough.go index 6b216005..781a1978 100644 --- a/overlay/batch/passthrough.go +++ b/overlay/batch/passthrough.go @@ -8,11 +8,11 @@ import ( // Passthrough is a RxBatcher that doesn't batch anything, it just accumulates and then sends packets. type Passthrough struct { - out io.Writer - slots [][]byte - // arena is injected; see TCPCoalescer.arena for the contract. - arena *Arena - cursor int + out io.Writer + slots [][]byte + reserver Reserver + resetter Resetter + cursor int } const passthroughBaseNumSlots = 128 @@ -21,16 +21,17 @@ const passthroughBaseNumSlots = 128 // standalone Passthrough batcher: 128 slots × udp.MTU ≈ 1.1 MiB. 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{ - out: w, - slots: make([][]byte, 0, passthroughBaseNumSlots), - arena: arena, + out: w, + slots: make([][]byte, 0, passthroughBaseNumSlots), + reserver: reserver, + resetter: resetter, } } func (p *Passthrough) Reserve(sz int) []byte { - return p.arena.Reserve(sz) + return p.reserver(sz) } func (p *Passthrough) Commit(pkt []byte) error { @@ -38,7 +39,17 @@ func (p *Passthrough) Commit(pkt []byte) error { return nil } +// Flush drains every queued packet and calls the configured Resetter 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 for _, s := range p.slots { _, err := p.out.Write(s) @@ -48,6 +59,5 @@ func (p *Passthrough) Flush() error { } clear(p.slots) p.slots = p.slots[:0] - p.arena.Reset() return firstErr } diff --git a/overlay/batch/tcp_coalesce.go b/overlay/batch/tcp_coalesce.go index f7630c12..51f0674d 100644 --- a/overlay/batch/tcp_coalesce.go +++ b/overlay/batch/tcp_coalesce.go @@ -83,22 +83,19 @@ type TCPCoalescer struct { // at is removed/sealed. lastSlot *coalesceSlot pool []*coalesceSlot // free list for reuse - - // arena is injected; the coalescer borrows slices from it via Reserve - // and tells it to release them via Reset on Flush. When wrapped in - // 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 + reserver Reserver + resetter Resetter + 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{ plainW: w, slots: make([]*coalesceSlot, 0, initialSlots), openSlots: make(map[flowKey]*coalesceSlot, initialSlots), pool: make([]*coalesceSlot, 0, initialSlots), - arena: arena, + reserver: reserver, + resetter: resetter, l: l, } 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 { - return c.arena.Reserve(sz) + return c.reserver(sz) } // 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 // may be recycled. 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() var first error for _, s := range c.slots { @@ -278,7 +285,6 @@ func (c *TCPCoalescer) Flush() error { clear(c.openSlots) c.lastSlot = nil - c.arena.Reset() return first } diff --git a/overlay/batch/tcp_coalesce_bench_test.go b/overlay/batch/tcp_coalesce_bench_test.go index f4ebfcdf..1d6232bf 100644 --- a/overlay/batch/tcp_coalesce_bench_test.go +++ b/overlay/batch/tcp_coalesce_bench_test.go @@ -71,7 +71,8 @@ func buildICMPv4() []byte { // between batches, and reports per-packet cost. func runCommitBench(b *testing.B, pkts [][]byte, batchSize int) { b.Helper() - c := NewTCPCoalescer(nopTunWriter{}, test.NewLogger(), NewArena(0)) + arena := NewArena(0) + c := NewTCPCoalescer(nopTunWriter{}, test.NewLogger(), arena.Reserve, arena.Reset) b.ReportAllocs() b.SetBytes(int64(len(pkts[0]))) b.ResetTimer() diff --git a/overlay/batch/tcp_coalesce_test.go b/overlay/batch/tcp_coalesce_test.go index 941161fa..99560c53 100644 --- a/overlay/batch/tcp_coalesce_test.go +++ b/overlay/batch/tcp_coalesce_test.go @@ -128,7 +128,8 @@ const ( func TestCoalescerPassthroughWhenGSOUnavailable(t *testing.T) { 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")) if err := c.Commit(pkt); err != nil { t.Fatal(err) @@ -147,7 +148,8 @@ func TestCoalescerPassthroughWhenGSOUnavailable(t *testing.T) { func TestCoalescerNonTCPPassthrough(t *testing.T) { 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[0] = 0x45 binary.BigEndian.PutUint16(pkt[2:4], 28) @@ -167,7 +169,8 @@ func TestCoalescerNonTCPPassthrough(t *testing.T) { func TestCoalescerSeedThenFlushAlone(t *testing.T) { 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)) if err := c.Commit(pkt); err != nil { t.Fatal(err) @@ -194,7 +197,8 @@ func TestCoalescerSeedThenFlushAlone(t *testing.T) { func TestCoalescerCoalescesAdjacentACKs(t *testing.T) { 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) if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -234,7 +238,8 @@ func TestCoalescerCoalescesAdjacentACKs(t *testing.T) { func TestCoalescerRejectsSeqGap(t *testing.T) { 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) if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -253,7 +258,8 @@ func TestCoalescerRejectsSeqGap(t *testing.T) { func TestCoalescerRejectsFlagMismatch(t *testing.T) { 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) if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -274,7 +280,8 @@ func TestCoalescerRejectsFlagMismatch(t *testing.T) { func TestCoalescerRejectsFIN(t *testing.T) { 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")) if err := c.Commit(fin); err != nil { t.Fatal(err) @@ -290,7 +297,8 @@ func TestCoalescerRejectsFIN(t *testing.T) { func TestCoalescerShortLastSegmentClosesChain(t *testing.T) { 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) half := make([]byte, 500) if err := c.Commit(buildTCPv4(1000, tcpAck, full)); err != nil { @@ -325,7 +333,8 @@ func TestCoalescerShortLastSegmentClosesChain(t *testing.T) { func TestCoalescerPSHFinalizesChain(t *testing.T) { 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) if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -355,7 +364,8 @@ func TestCoalescerPSHFinalizesChain(t *testing.T) { // coalescer drops it the sender's push signal never reaches the receiver. func TestCoalescerPropagatesPSHFromAppended(t *testing.T) { 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) // Seed has no PSH; second segment carries PSH and seals the chain. if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { @@ -383,7 +393,8 @@ func TestCoalescerPropagatesPSHFromAppended(t *testing.T) { func TestCoalescerRejectsDifferentFlow(t *testing.T) { 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) p1 := buildTCPv4(1000, tcpAck, pay) p2 := buildTCPv4(2200, tcpAck, pay) @@ -405,7 +416,8 @@ func TestCoalescerRejectsDifferentFlow(t *testing.T) { func TestCoalescerRejectsIPOptions(t *testing.T) { 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) pkt := buildTCPv4(1000, tcpAck, pay) // 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) { 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) seq := uint32(1000) for i := 0; i < tcpCoalesceMaxSegs+5; i++ { @@ -449,7 +462,8 @@ func TestCoalescerCapBySegments(t *testing.T) { // flows coalesce independently in a single Flush. func TestCoalescerMultipleFlowsInSameBatch(t *testing.T) { 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) // Flow A: sport 1000. Flow B: sport 3000. @@ -506,7 +520,8 @@ func TestCoalescerMultipleFlowsInSameBatch(t *testing.T) { // writing passthrough packets synchronously. func TestCoalescerPreservesArrivalOrder(t *testing.T) { 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 // a different flow. Expected emit order: gso(X), plain(ICMP), gso(Y). 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. func TestCoalescerInterleavedFlowsPreserveOrdering(t *testing.T) { 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) // Flow A two segments. @@ -679,7 +695,8 @@ func buildTCPv6(tcLow byte, seq uint32, flags byte, payload []byte) []byte { // retains ECE on the wire. func TestCoalescerCoalescesEceFlow(t *testing.T) { 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) flags := byte(tcpAck | tcpEce) 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. func TestCoalescerCwrSealsFlow(t *testing.T) { 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) if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -741,7 +759,8 @@ func TestCoalescerCwrSealsFlow(t *testing.T) { // a CE-echoing window or none. func TestCoalescerEceMismatchReseeds(t *testing.T) { 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) if err := c.Commit(buildTCPv4(1000, tcpAck|tcpEce, pay)); err != nil { t.Fatal(err) @@ -771,7 +790,8 @@ func TestCoalescerEceMismatchReseeds(t *testing.T) { // across the whole burst. func TestCoalescerDifferingECNReseeds(t *testing.T) { 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) if err := c.Commit(buildTCPv4WithToS(ecnECT0, 1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -816,7 +836,8 @@ func TestCoalescerDifferingECNReseeds(t *testing.T) { // codepoint, and neither may end up CE-marked. func TestCoalescerECT0ThenECT1NoCE(t *testing.T) { 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) if err := c.Commit(buildTCPv4WithToS(ecnECT0, 1000, tcpAck, pay)); err != nil { t.Fatal(err) @@ -846,7 +867,8 @@ func TestCoalescerECT0ThenECT1NoCE(t *testing.T) { // six DSCP bits must match too. func TestCoalescerDscpMismatchReseeds(t *testing.T) { 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) // Same ECN (Not-ECT), different DSCP (0x10 vs 0x20 in upper 6 bits). tosA := byte(0x10<<2) | ecnNotECT @@ -869,7 +891,8 @@ func TestCoalescerDscpMismatchReseeds(t *testing.T) { // TestCoalescerCoalescesEceFlow. func TestCoalescerIPv6CoalescesEceFlow(t *testing.T) { 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) flags := byte(tcpAck | tcpEce) 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. func TestCoalescerSortsReorderedSeedsAndMerges(t *testing.T) { 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) // Arrival order: seq 1000, 3400, 2200. The 3400 seeds a separate 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. func TestCoalescerSortAcrossFlowsMergesEachIndependently(t *testing.T) { 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) // 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. @@ -987,7 +1012,8 @@ func TestCoalescerSortAcrossFlowsMergesEachIndependently(t *testing.T) { // boundary by an arbitrary number of segments. func TestCoalescerSortKeepsPSHBoundary(t *testing.T) { 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) // 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 @@ -1015,7 +1041,8 @@ func TestCoalescerSortKeepsPSHBoundary(t *testing.T) { // is sorted/merged independently. func TestCoalescerSortKeepsPassthroughBarrier(t *testing.T) { 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) // First two segments seed S1 (then a 3400 reorder seeds S2). 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. func TestCoalescerIPv6DifferingECNReseeds(t *testing.T) { 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) // 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 { diff --git a/overlay/batch/udp_coalesce.go b/overlay/batch/udp_coalesce.go index 7bd7c89f..af5a3541 100644 --- a/overlay/batch/udp_coalesce.go +++ b/overlay/batch/udp_coalesce.go @@ -65,9 +65,8 @@ type UDPCoalescer struct { slots []*udpSlot openSlots map[flowKey]*udpSlot pool []*udpSlot - - // arena is injected; see TCPCoalescer.arena for the contract. - arena *Arena + reserver Reserver + resetter Resetter } // 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 // tio.GSOWriter at all (single-packet Queue), the coalescer degrades to // 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{ plainW: w, slots: make([]*udpSlot, 0, initialSlots), openSlots: make(map[flowKey]*udpSlot, initialSlots), pool: make([]*udpSlot, 0, initialSlots), - arena: arena, + reserver: reserver, + resetter: resetter, } if gw, ok := tio.SupportsGSO(w, tio.GSOProtoUDP); ok { c.gsoW = gw @@ -127,7 +127,7 @@ func parseUDP(pkt []byte) (parsedUDP, bool) { } 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. @@ -178,7 +178,19 @@ func (c *UDPCoalescer) commitParsed(pkt []byte, info parsedUDP) error { return nil } +// Flush drains every queued slot and calls the configured Resetter. 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 for _, s := range c.slots { var err error @@ -195,7 +207,6 @@ func (c *UDPCoalescer) Flush() error { clear(c.slots) c.slots = c.slots[:0] clear(c.openSlots) - c.arena.Reset() return first } diff --git a/overlay/batch/udp_coalesce_test.go b/overlay/batch/udp_coalesce_test.go index 273c7d69..bd88835a 100644 --- a/overlay/batch/udp_coalesce_test.go +++ b/overlay/batch/udp_coalesce_test.go @@ -60,7 +60,8 @@ func buildUDPv6(sport, dport uint16, payload []byte) []byte { func TestUDPCoalescerPassthroughWhenGSOUnavailable(t *testing.T) { 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)) if err := c.Commit(pkt); err != nil { t.Fatal(err) @@ -78,7 +79,8 @@ func TestUDPCoalescerPassthroughWhenGSOUnavailable(t *testing.T) { func TestUDPCoalescerNonUDPPassthrough(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) // ICMP packet pkt := make([]byte, 28) pkt[0] = 0x45 @@ -99,7 +101,8 @@ func TestUDPCoalescerNonUDPPassthrough(t *testing.T) { func TestUDPCoalescerSeedThenFlushAlone(t *testing.T) { 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)) if err := c.Commit(pkt); err != nil { t.Fatal(err) @@ -116,7 +119,8 @@ func TestUDPCoalescerSeedThenFlushAlone(t *testing.T) { func TestUDPCoalescerCoalescesEqualSized(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) pay := make([]byte, 1200) for i := 0; i < 3; i++ { 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. func TestUDPCoalescerShortLastSegmentSeals(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) full := make([]byte, 1200) tail := make([]byte, 600) 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. func TestUDPCoalescerLargerThanSeedReseeds(t *testing.T) { 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 { t.Fatal(err) } @@ -207,7 +213,8 @@ func TestUDPCoalescerLargerThanSeedReseeds(t *testing.T) { // Different 5-tuples must not coalesce. func TestUDPCoalescerDifferentFlowsKeepSeparate(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) pay := make([]byte, 800) if err := c.Commit(buildUDPv4(1000, 53, pay)); err != nil { t.Fatal(err) @@ -238,7 +245,8 @@ func TestUDPCoalescerDifferentFlowsKeepSeparate(t *testing.T) { // Caps at udpCoalesceMaxSegs. func TestUDPCoalescerCapsAtMaxSegs(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) pay := make([]byte, 100) for i := 0; i < udpCoalesceMaxSegs+5; i++ { 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. func TestUDPCoalescerDifferingECNReseeds(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) pay := make([]byte, 800) pkt0 := buildUDPv4(1000, 53, pay) // ECN=00 (Not-ECT) pkt1 := buildUDPv4(1000, 53, pay) @@ -298,7 +307,8 @@ func TestUDPCoalescerDifferingECNReseeds(t *testing.T) { // IPv6 path: same flow, equal-sized → coalesced. func TestUDPCoalescerIPv6Coalesces(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) pay := make([]byte, 1200) for i := 0; i < 3; i++ { 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. func TestUDPCoalescerDSCPMismatchReseeds(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) pay := make([]byte, 800) pkt0 := buildUDPv4(1000, 53, pay) pkt1 := buildUDPv4(1000, 53, pay) @@ -356,7 +367,8 @@ func TestUDPCoalescerDSCPMismatchReseeds(t *testing.T) { // Fragmented IPv4 must not be coalesced. func TestUDPCoalescerFragmentedIPv4PassesThrough(t *testing.T) { 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)) binary.BigEndian.PutUint16(pkt[6:8], 0x2000) // MF=1 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. func TestUDPCoalescerZeroLengthPayloadPassesThrough(t *testing.T) { 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 if err := c.Commit(pkt); err != nil { t.Fatal(err) @@ -396,7 +409,8 @@ func TestUDPCoalescerZeroLengthPayloadPassesThrough(t *testing.T) { // IPv6 zero-length UDP datagram: same passthrough contract as v4. func TestUDPCoalescerZeroLengthPayloadIPv6PassesThrough(t *testing.T) { 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 if err := c.Commit(pkt); err != nil { t.Fatal(err) @@ -417,7 +431,8 @@ func TestUDPCoalescerZeroLengthPayloadIPv6PassesThrough(t *testing.T) { // wire — per-flow arrival order (full, empty, full) must be preserved. func TestUDPCoalescerZeroLengthMidFlowSealsAndPreservesOrder(t *testing.T) { w := &fakeTunWriter{gsoEnabled: true} - c := NewUDPCoalescer(w, NewArena(0)) + arena := NewArena(0) + c := NewUDPCoalescer(w, arena.Reserve, arena.Reset) full := make([]byte, 800) if err := c.Commit(buildUDPv4(1000, 53, full)); err != nil { t.Fatal(err) @@ -441,7 +456,8 @@ func TestUDPCoalescerZeroLengthMidFlowSealsAndPreservesOrder(t *testing.T) { // IPv4 with options is not admissible (we require IHL=5). func TestUDPCoalescerIPv4WithOptionsPassesThrough(t *testing.T) { 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[0] = 0x46 // IHL = 6 (24-byte IPv4 header — has options) if err := c.Commit(pkt); err != nil {