mirror of
https://github.com/slackhq/nebula.git
synced 2026-08-15 09:36:58 +02:00
overlay/batch: don't seal the open slot on a pure ACK
Every non-coalesceable in-flow packet evicted the flow's open slot, so a bidirectional connection's inbound data run was broken by each peer ACK interleaved into it, largely defeating coalescing on concurrent upload+download. A bare acknowledgment (zero payload, nothing beyond ACK|PSH|ECE) carries no ordering obligation toward the flow's data -- delivered late it is just a stale ACK the receiver ignores -- so it can ride the lane as a passthrough without the evict, same as kernel GRO, which doesn't flush held data on pure ACKs. SYN/FIN/RST/CWR keep sealing. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014ugV2edVqoz3tBvq9J6yWp
This commit is contained in:
@@ -171,6 +171,17 @@ func (p parsedTCP) coalesceable() bool {
|
||||
return p.payLen > 0
|
||||
}
|
||||
|
||||
// pureAck reports whether a parsed segment is a bare acknowledgment: no
|
||||
// payload and nothing beyond ACK|PSH|ECE in the flags. These are the only
|
||||
// non-coalesceable shape that may safely pass through WITHOUT sealing the
|
||||
// flow's open slot — a late-delivered stale ACK is ignored by the receiver,
|
||||
// whereas SYN/FIN/RST/CWR mark transitions the flow must observe in order.
|
||||
func (p parsedTCP) pureAck() bool {
|
||||
return p.payLen == 0 &&
|
||||
p.flags&tcpFlagAck != 0 &&
|
||||
p.flags&^(tcpFlagAck|tcpFlagPsh|tcpFlagEce) == 0
|
||||
}
|
||||
|
||||
func (c *TCPCoalescer) Commit(pkt []byte) error {
|
||||
info, ok := parseTCPBase(pkt)
|
||||
if !ok {
|
||||
@@ -186,11 +197,21 @@ func (c *TCPCoalescer) Commit(pkt []byte) error {
|
||||
// after the dispatcher has already done so.
|
||||
func (c *TCPCoalescer) commitParsed(pkt []byte, info parsedTCP) error {
|
||||
if !info.coalesceable() {
|
||||
// TCP but not admissible (SYN/FIN/RST/URG/CWR or zero-payload).
|
||||
// Seal this flow's open slot so later in-flow packets don't extend
|
||||
// it and accidentally reorder past this passthrough. The len guard
|
||||
// skips hashing the 38-byte key on ack-dominant queues, where the
|
||||
// map is almost always empty.
|
||||
if info.pureAck() {
|
||||
// A bare window/ack update carries no ordering obligation toward
|
||||
// the flow's data: delivering it after later-arriving data only
|
||||
// makes it a stale ACK, which receivers ignore. Skipping the
|
||||
// evict keeps a bidirectional flow's inbound data run coalescing
|
||||
// across the peer ACKs interleaved into it — kernel GRO likewise
|
||||
// doesn't flush held data on a pure ACK.
|
||||
c.addPassthrough(pkt)
|
||||
return nil
|
||||
}
|
||||
// TCP but not admissible (SYN/FIN/RST/URG/CWR or a shape the flow
|
||||
// must observe in sequence). Seal this flow's open slot so later
|
||||
// in-flow packets don't extend it and accidentally reorder past this
|
||||
// passthrough. The len guard skips hashing the 38-byte key on
|
||||
// ack-dominant queues, where the map is almost always empty.
|
||||
if len(c.openSlots) != 0 {
|
||||
if last := c.lastSlot; last != nil && last.fk == info.fk {
|
||||
c.lastSlot = nil
|
||||
|
||||
@@ -229,6 +229,70 @@ func TestCoalescerSeedThenFlushAlone(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestCoalescerPureAckDoesNotSealRun pins the pure-ACK fast path: a bare
|
||||
// acknowledgment (zero payload, nothing beyond ACK|PSH|ECE) rides its lane
|
||||
// as a passthrough WITHOUT sealing the flow's open slot, so an inbound data
|
||||
// run on a bidirectional connection keeps coalescing across the peer ACKs
|
||||
// interleaved into it. The ACK is emitted after the superpacket (stale ACKs
|
||||
// are ignored by receivers, so the reorder is harmless by design).
|
||||
func TestCoalescerPureAckDoesNotSealRun(t *testing.T) {
|
||||
w := &fakeTunWriter{gsoEnabled: true}
|
||||
c := newTestTCPCoalescer(t, w)
|
||||
pay := make([]byte, 1200)
|
||||
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ack := buildTCPv4(2200, tcpAck, nil)
|
||||
if err := c.Commit(ack); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := c.Commit(buildTCPv4(2200, tcpAck, pay)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := c.Flush(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(w.gsoWrites) != 1 || len(w.writes) != 1 {
|
||||
t.Fatalf("ACK sealed the run: writes=%d gso=%d, want 1 gso (2 pays) + 1 plain", len(w.writes), len(w.gsoWrites))
|
||||
}
|
||||
if got := len(w.gsoWrites[0].pays); got != 2 {
|
||||
t.Errorf("pay count=%d want 2 (data kept coalescing across the ACK)", got)
|
||||
}
|
||||
if !bytes.Equal(w.writes[0], ack) {
|
||||
t.Errorf("plain write is not the ACK packet: got %d bytes want %d", len(w.writes[0]), len(ack))
|
||||
}
|
||||
if got, want := w.order, []string{"gso", "write"}; !stringSliceEq(got, want) {
|
||||
t.Errorf("flush order=%v want %v (slot order: data run seeded first)", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCoalescerFinStillSealsRun is the guard rail for the pure-ACK fast
|
||||
// path: control flags (here FIN|ACK, zero payload) must keep sealing the
|
||||
// open slot so data never reorders across a flow-state transition.
|
||||
func TestCoalescerFinStillSealsRun(t *testing.T) {
|
||||
w := &fakeTunWriter{gsoEnabled: true}
|
||||
c := newTestTCPCoalescer(t, w)
|
||||
pay := make([]byte, 1200)
|
||||
if err := c.Commit(buildTCPv4(1000, tcpAck, pay)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := c.Commit(buildTCPv4(2200, tcpFin|tcpAck, nil)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := c.Commit(buildTCPv4(2200, tcpAck, pay)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := c.Flush(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// FIN evicts the open slot; the third packet seeds a fresh one. All
|
||||
// three stay single-segment, so all three emit as plain writes in
|
||||
// arrival order — any gso write would mean data coalesced across FIN.
|
||||
if len(w.writes) != 3 || len(w.gsoWrites) != 0 {
|
||||
t.Fatalf("FIN must seal the run: writes=%d gso=%d", len(w.writes), len(w.gsoWrites))
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoalescerCoalescesAdjacentACKs(t *testing.T) {
|
||||
w := &fakeTunWriter{gsoEnabled: true}
|
||||
c := newTestTCPCoalescer(t, w)
|
||||
|
||||
Reference in New Issue
Block a user