From 2ef8d9edf0cc0f8083240888404c4b64ebf565ce Mon Sep 17 00:00:00 2001 From: Wago Networking Agent Date: Sun, 26 Jul 2026 12:40:03 +0000 Subject: [PATCH 1/3] perf: bound idle TCP reuse storage --- agent-todo.md | 6 +- docs/architecture.md | 14 +- internal/backend/lneto/tcp/tcp.go | 61 ++++----- internal/backend/lneto/tcp/tcp_test.go | 173 +++++++++++++++++++++++-- 4 files changed, 205 insertions(+), 49 deletions(-) diff --git a/agent-todo.md b/agent-todo.md index 652580b..9ca09ac 100644 --- a/agent-todo.md +++ b/agent-todo.md @@ -205,7 +205,11 @@ but host-facing operations call only immediate `tcp.Handler` state/buffer method under the namespace lock; `tcp.Conn.Read`, `Write`, and `Flush` remain absent. Connect/accept, partial I/O, EOF/reset semantics, half-close, policy, exact resource/retained-storage quota, bounded readiness, port reuse, and abort cleanup -are covered. Accepted-stream close releases resource quota immediately; lneto's +are covered. Closed listener/outbound bytes are zeroed, and idle reuse retains +at most one listener pool capped at 256 slots and 1 MiB plus one outbound buffer +capped at 1 MiB; excess high-water storage is dropped after quota release. +Adapter creation no longer allocates reuse-index arrays proportional to maximum +listener/outbound counts. Accepted-stream close releases resource quota immediately; lneto's private accepted list is preserved until the next bounded egress service probe, which reclaims the pool slot and now reports one charged maintenance operation even when it emits no frame. DNS uses adapter-owned immediate IPv4 UDP packets plus lneto DNS codecs, diff --git a/docs/architecture.md b/docs/architecture.md index 7031e21..a530727 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -174,10 +174,16 @@ It uses only immediate `tcp.Handler` buffer/state primitives and never calls pools and outbound streams have bounded receive/transmit storage, partial I/O, connect and accept progress, half-close, level readiness, endpoint policy, quota ownership, -port reuse, and deterministic abort cleanup. Adapter creation seeds only a -small stream-registry capacity hint and grows that registry as streams are -actually created rather than preallocating for the full theoretical -`MaxOutboundStreams + MaxListeners*AcceptBacklog` population. Closing an +port reuse, and deterministic abort cleanup. Closed listener and outbound +buffers are zeroed before reuse or release. The adapter retains at most one idle +listener pool of no more than 256 slots and 1 MiB of storage plus one idle +outbound buffer of no more than 1 MiB; concurrent high-water buffers and larger +configurations are dropped after close rather than remaining as uncharged cache. +Adapter creation seeds only a small stream-registry capacity hint and no longer +allocates reuse-index arrays proportional to configured listener or outbound +limits. It grows the registry as streams are actually created rather than +preallocating for the full theoretical `MaxOutboundStreams + +MaxListeners*AcceptBacklog` population. Closing an accepted stream releases its resource quota immediately. lneto retains the closed pool entry until its listener performs maintenance; the next bounded egress service probe reclaims that entry and now reports one charged service diff --git a/internal/backend/lneto/tcp/tcp.go b/internal/backend/lneto/tcp/tcp.go index 84ab8e0..f34c5d2 100644 --- a/internal/backend/lneto/tcp/tcp.go +++ b/internal/backend/lneto/tcp/tcp.go @@ -24,6 +24,8 @@ const ( ingressOrder = 15 closeOrder = 20 maxEagerTCPListenerStorageBytes = 256 << 20 + maxIdleTCPReuseStorageBytes = 1 << 20 + maxIdleTCPListenerSlots = 256 maxTCPStreamCapacityHint = 16 ) @@ -42,10 +44,10 @@ type Adapter struct { quotas *quota.Account config Config listeners []*tcpListener - freeListenerPools []tcpPool + freeListenerPool tcpPool streams []*tcpStream outboundStreams int - freeOutboundStorage [][]byte + freeOutboundStorage []byte portOwner *lnetocore.TCPPortOwner nextISS lnetotcp.Value } @@ -97,7 +99,6 @@ func New(common *lnetocore.Namespace, config Config) (*Adapter, error) { } n.listeners = make([]*tcpListener, 0, config.MaxListeners) n.streams = make([]*tcpStream, 0, streamCapacityHint(config)) - n.prepareReusePools() common.Unlock() if err := common.Install(lnetocore.Participant{IngressOrder: ingressOrder, Ingress: n.ingressLocked, CloseOrder: closeOrder, Close: n.CloseLocked}); err != nil { return nil, err @@ -175,32 +176,19 @@ func streamCapacityHint(config Config) int { return int(hint) } -func (n *Adapter) prepareReusePools() { - if n == nil { - return - } - if n.config.MaxListeners > 0 { - n.freeListenerPools = make([]tcpPool, 0, n.config.MaxListeners) - } - if n.config.MaxOutboundStreams > 0 { - n.freeOutboundStorage = make([][]byte, 0, n.config.MaxOutboundStreams) - } -} - func (n *Adapter) acquireListenerLocked(local nscore.Endpoint) (*tcpListener, error) { if n == nil { return nil, lneto.ErrInvalidConfig } - var pool tcpPool - if len(n.freeListenerPools) == 0 { + pool := n.freeListenerPool + n.freeListenerPool = tcpPool{} + if pool.slots == nil { created, err := newTCPPool(n, n.config.AcceptBacklog, n.config) if err != nil { return nil, err } pool = created } else { - pool = n.freeListenerPools[len(n.freeListenerPools)-1] - n.freeListenerPools = n.freeListenerPools[:len(n.freeListenerPools)-1] pool.resetLocked(n) } return &tcpListener{owner: n, local: local, pool: pool}, nil @@ -211,7 +199,11 @@ func (n *Adapter) recycleListenerLocked(listener *tcpListener) { return } listener.pool.releaseLocked() - n.freeListenerPools = append(n.freeListenerPools, listener.pool) + if n.freeListenerPool.slots == nil && len(listener.pool.slots) <= maxIdleTCPListenerSlots && len(listener.pool.storage) <= maxIdleTCPReuseStorageBytes { + n.freeListenerPool = listener.pool + } else { + listener.pool.destroyLocked() + } listener.pool = tcpPool{} listener.listener = lnetotcp.Listener{} listener.local = nscore.Endpoint{} @@ -228,11 +220,11 @@ func (n *Adapter) prepareOutboundStreamLocked(stream *tcpStream, retained uint64 return err } stream.allocation = &stream.retained - if len(n.freeOutboundStorage) == 0 { + if n.freeOutboundStorage == nil { stream.storage = make([]byte, int(retained)) } else { - stream.storage = n.freeOutboundStorage[len(n.freeOutboundStorage)-1] - n.freeOutboundStorage = n.freeOutboundStorage[:len(n.freeOutboundStorage)-1] + stream.storage = n.freeOutboundStorage + n.freeOutboundStorage = nil } return nil } @@ -242,7 +234,9 @@ func (n *Adapter) recycleOutboundStreamLocked(stream *tcpStream) { return } clear(stream.storage) - n.freeOutboundStorage = append(n.freeOutboundStorage, stream.storage) + if n.freeOutboundStorage == nil && len(stream.storage) <= maxIdleTCPReuseStorageBytes { + n.freeOutboundStorage = stream.storage + } stream.conn = nil stream.connValue = lnetotcp.Conn{} stream.storage = nil @@ -292,6 +286,7 @@ type tcpStream struct { type tcpPool struct { owner *Adapter slots []tcpPoolSlot + storage []byte nextISS lnetotcp.Value } @@ -311,14 +306,14 @@ func newTCPPool(owner *Adapter, count uint16, config Config) (tcpPool, error) { return pool, nil } stride, _ := tcpStreamStorageBytes(config) - storage := make([]byte, int(uint64(count)*stride)) + pool.storage = make([]byte, int(uint64(count)*stride)) strideBytes := int(stride) for i := range pool.slots { start := i * strideBytes rxEnd := start + config.ReceiveBytes if err := pool.slots[i].conn.Configure(lnetotcp.ConnConfig{ - RxBuf: storage[start:rxEnd], - TxBuf: storage[rxEnd : start+strideBytes], + RxBuf: pool.storage[start:rxEnd], + TxBuf: pool.storage[rxEnd : start+strideBytes], TxPacketQueueSize: config.TransmitPackets, RWBackoff: immediateBackoff, }); err != nil { @@ -412,6 +407,7 @@ func (p *tcpPool) releaseLocked() { slot.inUse = false slot.quotaOwned = false } + clear(p.storage) } func (p *tcpPool) destroyLocked() { @@ -420,6 +416,7 @@ func (p *tcpPool) destroyLocked() { } p.releaseLocked() p.slots = nil + p.storage = nil p.owner = nil } @@ -1114,14 +1111,10 @@ func (n *Adapter) CloseLocked() { for len(n.streams) > 0 { n.streams[len(n.streams)-1].closeLocked() } - for i := range n.freeListenerPools { - n.freeListenerPools[i].destroyLocked() - } - for i := range n.freeOutboundStorage { - clear(n.freeOutboundStorage[i]) - } + n.freeListenerPool.destroyLocked() + clear(n.freeOutboundStorage) n.portOwner = nil - n.freeListenerPools = nil + n.freeListenerPool = tcpPool{} n.listeners = nil n.freeOutboundStorage = nil n.streams = nil diff --git a/internal/backend/lneto/tcp/tcp_test.go b/internal/backend/lneto/tcp/tcp_test.go index da29b57..7806a9f 100644 --- a/internal/backend/lneto/tcp/tcp_test.go +++ b/internal/backend/lneto/tcp/tcp_test.go @@ -227,7 +227,7 @@ func TestQuotaDeniedConnectDoesNotAllocateOrRetainOutboundStorage(t *testing.T) common.Lock() leaseCount := common.TCPPortLeaseCountLocked() common.Unlock() - if len(adapter.freeOutboundStorage) != 0 || len(adapter.streams) != 0 || leaseCount != 0 || adapter.outboundStreams != 0 { + if adapter.freeOutboundStorage != nil || len(adapter.streams) != 0 || leaseCount != 0 || adapter.outboundStreams != 0 { t.Fatalf("quota denied connect retained state: storage=%d streams=%d leases=%d outbound=%d", len(adapter.freeOutboundStorage), len(adapter.streams), leaseCount, adapter.outboundStreams) } if usage, closed := adapter.quotas.Snapshot(); closed || usage != (quota.Usage{QueuedBytes: 16 << 10}) { @@ -305,10 +305,10 @@ func TestOutboundStorageReuseClearsDataAndIsolatesStaleStream(t *testing.T) { if stale.storage != nil || stale.conn != nil || stale.allocation != nil || !stale.closed { t.Fatalf("closed stale stream retained state: storage=%v conn=%p allocation=%p closed=%v", stale.storage, stale.conn, stale.allocation, stale.closed) } - if len(adapter.freeOutboundStorage) != 1 || &adapter.freeOutboundStorage[0][0] != storageAddress { - t.Fatalf("outbound storage was not retained for bounded reuse: pools=%d", len(adapter.freeOutboundStorage)) + if len(adapter.freeOutboundStorage) != len(storage) || &adapter.freeOutboundStorage[0] != storageAddress { + t.Fatalf("outbound storage was not retained for bounded reuse: bytes=%d", len(adapter.freeOutboundStorage)) } - for i, value := range adapter.freeOutboundStorage[0] { + for i, value := range adapter.freeOutboundStorage { if value != 0 { t.Fatalf("recycled storage byte %d = %d", i, value) } @@ -347,6 +347,159 @@ func TestOutboundStorageReuseClearsDataAndIsolatesStaleStream(t *testing.T) { } } +func TestOutboundStorageReuseRetainsAtMostOneIdleBuffer(t *testing.T) { + core, adapter := newTestAdapter(t, 7, 0, 3) + remote := nscore.Endpoint{Address: netip.MustParseAddr("192.0.2.8"), Port: 4208} + streams := make([]*tcpStream, 0, 3) + storages := make([][]byte, 0, 3) + for i := 0; i < 3; i++ { + value, progress, err := adapter.TryConnect(remote) + if err != nil || progress != nscore.ProgressInProgress { + t.Fatalf("connect %d = %T, %v, %v", i, value, progress, err) + } + stream := value.(*tcpStream) + for j := range stream.storage { + stream.storage[j] = byte(i + 1) + } + streams = append(streams, stream) + storages = append(storages, stream.storage) + } + wantCached := &storages[0][0] + for i, stream := range streams { + if err := stream.Close(); err != nil { + t.Fatalf("close %d: %v", i, err) + } + } + if len(adapter.freeOutboundStorage) != 512 || &adapter.freeOutboundStorage[0] != wantCached { + t.Fatalf("idle outbound cache = bytes:%d address:%p want:%p", len(adapter.freeOutboundStorage), firstByte(adapter.freeOutboundStorage), wantCached) + } + for i, storage := range storages { + for j, value := range storage { + if value != 0 { + t.Fatalf("released storage %d byte %d = %d", i, j, value) + } + } + } + if len(adapter.streams) != 0 || adapter.outboundStreams != 0 { + t.Fatalf("closed outbound state = streams:%d outbound:%d", len(adapter.streams), adapter.outboundStreams) + } + if usage, _ := adapter.quotas.Snapshot(); usage != (quota.Usage{}) { + t.Fatalf("closed outbound quota = %+v", usage) + } + if err := core.Close(); err != nil { + t.Fatal(err) + } + if adapter.freeOutboundStorage != nil { + t.Fatalf("namespace close retained %d outbound bytes", len(adapter.freeOutboundStorage)) + } +} + +func TestListenerPoolReuseRetainsAtMostOneClearedIdlePool(t *testing.T) { + core, adapter := newTestAdapterWithBacklog(t, 9, 3, 0, 2) + listeners := make([]*tcpListener, 0, 3) + storages := make([][]byte, 0, 3) + for i := 0; i < 3; i++ { + local := nscore.Endpoint{Address: netip.MustParseAddr("192.0.2.9"), Port: uint16(4210 + i)} + value, progress, err := adapter.TryListen(local) + if err != nil || progress != nscore.ProgressDone { + t.Fatalf("listen %d = %T, %v, %v", i, value, progress, err) + } + listener := value.(*tcpListener) + for j := range listener.pool.storage { + listener.pool.storage[j] = byte(i + 1) + } + listeners = append(listeners, listener) + storages = append(storages, listener.pool.storage) + } + wantCached := &storages[0][0] + for i, listener := range listeners { + if err := listener.Close(); err != nil { + t.Fatalf("close %d: %v", i, err) + } + if listener.pool.slots != nil || listener.pool.storage != nil { + t.Fatalf("closed listener %d retained pool", i) + } + } + if len(adapter.freeListenerPool.slots) != 2 || len(adapter.freeListenerPool.storage) != 1024 || &adapter.freeListenerPool.storage[0] != wantCached { + t.Fatalf("idle listener cache = slots:%d bytes:%d address:%p want:%p", len(adapter.freeListenerPool.slots), len(adapter.freeListenerPool.storage), firstByte(adapter.freeListenerPool.storage), wantCached) + } + for i, storage := range storages { + for j, value := range storage { + if value != 0 { + t.Fatalf("released listener storage %d byte %d = %d", i, j, value) + } + } + } + if len(adapter.listeners) != 0 { + t.Fatalf("closed listener state = listeners:%d", len(adapter.listeners)) + } + if usage, _ := adapter.quotas.Snapshot(); usage != (quota.Usage{}) { + t.Fatalf("closed listener quota = %+v", usage) + } + if err := core.Close(); err != nil { + t.Fatal(err) + } + if adapter.freeListenerPool.slots != nil || adapter.freeListenerPool.storage != nil { + t.Fatal("namespace close retained listener pool") + } +} + +func TestIdleReuseDropsOversizedTCPStorageAndListenerMetadata(t *testing.T) { + adapter := &Adapter{} + outboundStorage := make([]byte, maxIdleTCPReuseStorageBytes+1) + for i := range outboundStorage { + outboundStorage[i] = 0xff + } + stream := &tcpStream{owner: adapter, storage: outboundStorage} + adapter.recycleOutboundStreamLocked(stream) + if adapter.freeOutboundStorage != nil || stream.storage != nil { + t.Fatalf("oversized outbound storage retained: cache=%d stream=%d", len(adapter.freeOutboundStorage), len(stream.storage)) + } + for i, value := range outboundStorage { + if value != 0 { + t.Fatalf("oversized outbound byte %d = %d", i, value) + } + } + + oversizedConfig := Config{AcceptBacklog: 1, ReceiveBytes: maxIdleTCPReuseStorageBytes, TransmitBytes: 256, TransmitPackets: 4} + oversizedPool, err := newTCPPool(adapter, 1, oversizedConfig) + if err != nil { + t.Fatal(err) + } + oversizedListenerStorage := oversizedPool.storage + for i := range oversizedListenerStorage { + oversizedListenerStorage[i] = 0xff + } + oversizedListener := &tcpListener{pool: oversizedPool} + adapter.recycleListenerLocked(oversizedListener) + if adapter.freeListenerPool.slots != nil || adapter.freeListenerPool.storage != nil || oversizedListener.pool.slots != nil || oversizedListener.pool.storage != nil { + t.Fatal("oversized listener pool retained") + } + for i, value := range oversizedListenerStorage { + if value != 0 { + t.Fatalf("oversized listener byte %d = %d", i, value) + } + } + + metadataConfig := Config{AcceptBacklog: maxIdleTCPListenerSlots + 1, ReceiveBytes: 256, TransmitBytes: 256, TransmitPackets: 4} + metadataPool, err := newTCPPool(adapter, maxIdleTCPListenerSlots+1, metadataConfig) + if err != nil { + t.Fatal(err) + } + metadataListener := &tcpListener{pool: metadataPool} + adapter.recycleListenerLocked(metadataListener) + if adapter.freeListenerPool.slots != nil || adapter.freeListenerPool.storage != nil || metadataListener.pool.slots != nil || metadataListener.pool.storage != nil { + t.Fatal("oversized listener metadata retained") + } +} + +func firstByte(storage []byte) *byte { + if len(storage) == 0 { + return nil + } + return &storage[0] +} + func TestAcceptedCloseRetainsSlotUntilChargedMaintenance(t *testing.T) { clientCore, client := newTestAdapter(t, 1, 0, 2) serverCore, server := newTestAdapter(t, 2, 1, 0) @@ -1688,11 +1841,11 @@ func TestListenerBacklogCloseDetachesAllSlotsBeforePoolReuse(t *testing.T) { if usage, _ := server.quotas.Snapshot(); usage != (quota.Usage{}) { t.Fatalf("listener close quota = %+v", usage) } - if len(server.freeListenerPools) != 1 { - t.Fatalf("free listener pools = %d", len(server.freeListenerPools)) + if server.freeListenerPool.slots == nil { + t.Fatal("listener pool was not retained for bounded reuse") } - for i := range server.freeListenerPools[0].slots { - slot := &server.freeListenerPools[0].slots[i] + for i := range server.freeListenerPool.slots { + slot := &server.freeListenerPool.slots[i] if slot.inUse || slot.stream != nil || slot.quotaOwned || slot.resource.ResetReleased() { t.Fatalf("released slot %d = in_use=%v stream=%p quota_owned=%v", i, slot.inUse, slot.stream, slot.quotaOwned) } @@ -1706,8 +1859,8 @@ func TestListenerBacklogCloseDetachesAllSlotsBeforePoolReuse(t *testing.T) { t.Fatalf("replacement listen = %T, %v, %v", replacementValue, progress, err) } replacement := replacementValue.(*tcpListener) - if replacement == listener || len(server.freeListenerPools) != 0 { - t.Fatalf("replacement wrapper/pool reuse = same_wrapper=%v free_pools=%d", replacement == listener, len(server.freeListenerPools)) + if replacement == listener || server.freeListenerPool.slots != nil { + t.Fatalf("replacement wrapper/pool reuse = same_wrapper=%v free_pool=%v", replacement == listener, server.freeListenerPool.slots != nil) } connect(clientCore1, client1, clientMAC1) thirdServerValue, progress, err := replacement.TryAccept() From 68d2c95a1a5f8a107e7ca493677d785f84e4c253 Mon Sep 17 00:00:00 2001 From: Wago Networking Agent Date: Sun, 26 Jul 2026 12:43:47 +0000 Subject: [PATCH 2/3] perf: cap TCP registry initialization --- agent-todo.md | 5 +++-- docs/architecture.md | 10 +++++----- internal/backend/lneto/tcp/benchmark_test.go | 19 +++++++++++++++++++ internal/backend/lneto/tcp/tcp.go | 10 +++++++++- internal/backend/lneto/tcp/tcp_test.go | 7 +++++-- 5 files changed, 41 insertions(+), 10 deletions(-) diff --git a/agent-todo.md b/agent-todo.md index 9ca09ac..052bb75 100644 --- a/agent-todo.md +++ b/agent-todo.md @@ -208,8 +208,9 @@ resource/retained-storage quota, bounded readiness, port reuse, and abort cleanu are covered. Closed listener/outbound bytes are zeroed, and idle reuse retains at most one listener pool capped at 256 slots and 1 MiB plus one outbound buffer capped at 1 MiB; excess high-water storage is dropped after quota release. -Adapter creation no longer allocates reuse-index arrays proportional to maximum -listener/outbound counts. Accepted-stream close releases resource quota immediately; lneto's +Adapter creation now uses small listener/stream registry hints and no longer +allocates registries or reuse-index arrays proportional to maximum listener or +outbound counts. Accepted-stream close releases resource quota immediately; lneto's private accepted list is preserved until the next bounded egress service probe, which reclaims the pool slot and now reports one charged maintenance operation even when it emits no frame. DNS uses adapter-owned immediate IPv4 UDP packets plus lneto DNS codecs, diff --git a/docs/architecture.md b/docs/architecture.md index a530727..3ad4857 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -179,11 +179,11 @@ buffers are zeroed before reuse or release. The adapter retains at most one idle listener pool of no more than 256 slots and 1 MiB of storage plus one idle outbound buffer of no more than 1 MiB; concurrent high-water buffers and larger configurations are dropped after close rather than remaining as uncharged cache. -Adapter creation seeds only a small stream-registry capacity hint and no longer -allocates reuse-index arrays proportional to configured listener or outbound -limits. It grows the registry as streams are actually created rather than -preallocating for the full theoretical `MaxOutboundStreams + -MaxListeners*AcceptBacklog` population. Closing an +Adapter creation seeds only small listener/stream registry capacity hints and no +longer allocates registries or reuse-index arrays proportional to configured +listener or outbound limits. It grows the registries as resources are actually +created rather than preallocating for the full theoretical +`MaxOutboundStreams + MaxListeners*AcceptBacklog` population. Closing an accepted stream releases its resource quota immediately. lneto retains the closed pool entry until its listener performs maintenance; the next bounded egress service probe reclaims that entry and now reports one charged service diff --git a/internal/backend/lneto/tcp/benchmark_test.go b/internal/backend/lneto/tcp/benchmark_test.go index cd7224c..de8d6a5 100644 --- a/internal/backend/lneto/tcp/benchmark_test.go +++ b/internal/backend/lneto/tcp/benchmark_test.go @@ -38,6 +38,25 @@ func BenchmarkAdapterNew(b *testing.B) { } } +func BenchmarkAdapterNewMaximumListenerConfig(b *testing.B) { + config := Config{MaxListeners: ^uint16(0), AcceptBacklog: 1, ReceiveBytes: 256, TransmitBytes: 256, TransmitPackets: 4} + b.ReportAllocs() + for b.Loop() { + common := newConfigTestCore(b, config.MaxListeners) + adapter, err := New(common, config) + if err != nil { + common.Close() + b.Fatal(err) + } + if err := common.Close(); err != nil { + b.Fatal(err) + } + if adapter == nil { + b.Fatal("nil adapter") + } + } +} + func BenchmarkAdapterTryListenClose(b *testing.B) { _, adapter := newTestAdapter(b, 111, 1, 0) local := nscore.Endpoint{Address: netip.MustParseAddr("192.0.2.111"), Port: 4211} diff --git a/internal/backend/lneto/tcp/tcp.go b/internal/backend/lneto/tcp/tcp.go index f34c5d2..15e5e27 100644 --- a/internal/backend/lneto/tcp/tcp.go +++ b/internal/backend/lneto/tcp/tcp.go @@ -97,7 +97,7 @@ func New(common *lnetocore.Namespace, config Config) (*Adapter, error) { common.Unlock() return n, nil } - n.listeners = make([]*tcpListener, 0, config.MaxListeners) + n.listeners = make([]*tcpListener, 0, listenerCapacityHint(config)) n.streams = make([]*tcpStream, 0, streamCapacityHint(config)) common.Unlock() if err := common.Install(lnetocore.Participant{IngressOrder: ingressOrder, Ingress: n.ingressLocked, CloseOrder: closeOrder, Close: n.CloseLocked}); err != nil { @@ -168,6 +168,14 @@ func tcpStreamStorageBytes(config Config) (uint64, bool) { return checked.AddUint64(uint64(config.ReceiveBytes), uint64(config.TransmitBytes)) } +func listenerCapacityHint(config Config) int { + hint := uint64(config.MaxListeners) + if hint > maxTCPStreamCapacityHint { + hint = maxTCPStreamCapacityHint + } + return int(hint) +} + func streamCapacityHint(config Config) int { hint := uint64(config.MaxListeners) + uint64(config.MaxOutboundStreams) if hint > maxTCPStreamCapacityHint { diff --git a/internal/backend/lneto/tcp/tcp_test.go b/internal/backend/lneto/tcp/tcp_test.go index 7806a9f..18fef47 100644 --- a/internal/backend/lneto/tcp/tcp_test.go +++ b/internal/backend/lneto/tcp/tcp_test.go @@ -94,11 +94,14 @@ func TestValidConfigRejectsOverflowAndKeepsAdapterCreationBounded(t *testing.T) t.Fatalf("New error = %v", err) } }() + if cap(adapter.listeners) != maxTCPStreamCapacityHint { + t.Fatalf("listener capacity hint = %d, want %d", cap(adapter.listeners), maxTCPStreamCapacityHint) + } if cap(adapter.streams) != maxTCPStreamCapacityHint { t.Fatalf("stream capacity hint = %d, want %d", cap(adapter.streams), maxTCPStreamCapacityHint) } - if len(adapter.streams) != 0 { - t.Fatalf("new adapter eagerly populated streams = %d", len(adapter.streams)) + if len(adapter.listeners) != 0 || len(adapter.streams) != 0 || adapter.freeListenerPool.slots != nil || adapter.freeOutboundStorage != nil { + t.Fatalf("new adapter eagerly populated state: listeners=%d streams=%d listener-cache=%v outbound-cache=%d", len(adapter.listeners), len(adapter.streams), adapter.freeListenerPool.slots != nil, len(adapter.freeOutboundStorage)) } } From 8dd761f02d132ec71812d36c4dd7a30f93e08e79 Mon Sep 17 00:00:00 2001 From: Wago Networking Agent Date: Sun, 26 Jul 2026 12:51:01 +0000 Subject: [PATCH 3/3] perf: reuse TCP quota accounting --- agent-todo.md | 7 +- docs/architecture.md | 9 ++- internal/backend/lneto/tcp/tcp.go | 85 ++++++++++++++++-------- internal/backend/lneto/tcp/tcp_test.go | 92 +++++++++++++++++++++----- 4 files changed, 144 insertions(+), 49 deletions(-) diff --git a/agent-todo.md b/agent-todo.md index 052bb75..dda20c2 100644 --- a/agent-todo.md +++ b/agent-todo.md @@ -206,8 +206,11 @@ under the namespace lock; `tcp.Conn.Read`, `Write`, and `Flush` remain absent. Connect/accept, partial I/O, EOF/reset semantics, half-close, policy, exact resource/retained-storage quota, bounded readiness, port reuse, and abort cleanup are covered. Closed listener/outbound bytes are zeroed, and idle reuse retains -at most one listener pool capped at 256 slots and 1 MiB plus one outbound buffer -capped at 1 MiB; excess high-water storage is dropped after quota release. +at most one listener pool capped at 256 slots and 1 MiB plus one outbound +buffer/released-accounting slot capped at 1 MiB; excess high-water storage is +dropped after quota release. The accounting reuse lowers steady connect/close +allocation from 1,128 to 936 B/op (17.0%) while `tcp.Conn` remains deliberately +generation-distinct because lneto may retain its pointer after abort. Adapter creation now uses small listener/stream registry hints and no longer allocates registries or reuse-index arrays proportional to maximum listener or outbound counts. Accepted-stream close releases resource quota immediately; lneto's diff --git a/docs/architecture.md b/docs/architecture.md index 3ad4857..6d1c95c 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -175,9 +175,12 @@ pools and outbound streams have bounded receive/transmit storage, partial I/O, connect and accept progress, half-close, level readiness, endpoint policy, quota ownership, port reuse, and deterministic abort cleanup. Closed listener and outbound -buffers are zeroed before reuse or release. The adapter retains at most one idle -listener pool of no more than 256 slots and 1 MiB of storage plus one idle -outbound buffer of no more than 1 MiB; concurrent high-water buffers and larger +buffers are zeroed before reuse or release. The outbound buffer cache also owns +one released/reset quota charge, reducing steady wrapper garbage while each +`tcp.Conn` remains generation-distinct because lneto may retain its registration +pointer after abort. The adapter retains at most one idle listener pool of no +more than 256 slots and 1 MiB of storage plus one idle outbound buffer/accounting +slot of no more than 1 MiB; concurrent high-water buffers and larger configurations are dropped after close rather than remaining as uncharged cache. Adapter creation seeds only small listener/stream registry capacity hints and no longer allocates registries or reuse-index arrays proportional to configured diff --git a/internal/backend/lneto/tcp/tcp.go b/internal/backend/lneto/tcp/tcp.go index 15e5e27..ec213a3 100644 --- a/internal/backend/lneto/tcp/tcp.go +++ b/internal/backend/lneto/tcp/tcp.go @@ -47,7 +47,7 @@ type Adapter struct { freeListenerPool tcpPool streams []*tcpStream outboundStreams int - freeOutboundStorage []byte + freeOutboundStorage *tcpOutboundStorage portOwner *lnetocore.TCPPortOwner nextISS lnetotcp.Value } @@ -224,36 +224,61 @@ func (n *Adapter) prepareOutboundStreamLocked(stream *tcpStream, retained uint64 if n == nil || stream == nil || stream.owner != n || stream.portLease.TCPPort() == 0 { return lneto.ErrInvalidConfig } - if err := n.quotas.AcquireResourceAndQueuedBytes(&stream.retained, quota.ResourceTCP, 1, retained); err != nil { + storage := n.freeOutboundStorage + reused := storage != nil + n.freeOutboundStorage = nil + if storage == nil { + storage = new(tcpOutboundStorage) + } + stream.outboundStorage = storage + if err := n.quotas.AcquireResourceAndQueuedBytes(&storage.retained, quota.ResourceTCP, 1, retained); err != nil { + storage.retained = quota.Charge{} + if reused { + n.freeOutboundStorage = storage + } + stream.outboundStorage = nil return err } - stream.allocation = &stream.retained - if n.freeOutboundStorage == nil { - stream.storage = make([]byte, int(retained)) - } else { - stream.storage = n.freeOutboundStorage - n.freeOutboundStorage = nil + if storage.bytes == nil { + storage.bytes = make([]byte, int(retained)) } + stream.conn = &stream.connValue + stream.storage = storage.bytes + stream.allocation = &storage.retained return nil } +func (n *Adapter) recycleOutboundStorageLocked(storage *tcpOutboundStorage) { + if n == nil || storage == nil { + return + } + storage.retained = quota.Charge{} + clear(storage.bytes) + if n.freeOutboundStorage == nil && len(storage.bytes) <= maxIdleTCPReuseStorageBytes { + n.freeOutboundStorage = storage + } else { + storage.bytes = nil + } +} + func (n *Adapter) recycleOutboundStreamLocked(stream *tcpStream) { if n == nil || stream == nil { return } - clear(stream.storage) - if n.freeOutboundStorage == nil && len(stream.storage) <= maxIdleTCPReuseStorageBytes { - n.freeOutboundStorage = stream.storage + if stream.outboundStorage != nil { + n.recycleOutboundStorageLocked(stream.outboundStorage) + } else { + clear(stream.storage) } stream.conn = nil stream.connValue = lnetotcp.Conn{} stream.storage = nil + stream.outboundStorage = nil stream.local = nscore.Endpoint{} stream.remote = nscore.Endpoint{} stream.portLease = lnetocore.TCPPortLease{} stream.slot = nil stream.allocation = nil - stream.retained = quota.Charge{} stream.connected = false stream.shutdown = false stream.terminal = false @@ -272,17 +297,22 @@ type tcpListener struct { closed bool } +type tcpOutboundStorage struct { + retained quota.Charge + bytes []byte +} + type tcpStream struct { - owner *Adapter - conn *lnetotcp.Conn - connValue lnetotcp.Conn - storage []byte - local nscore.Endpoint - remote nscore.Endpoint - slot *tcpPoolSlot + owner *Adapter + conn *lnetotcp.Conn + connValue lnetotcp.Conn + storage []byte + outboundStorage *tcpOutboundStorage + local nscore.Endpoint + remote nscore.Endpoint + slot *tcpPoolSlot allocation *quota.Charge - retained quota.Charge portLease lnetocore.TCPPortLease connected bool shutdown bool @@ -690,7 +720,7 @@ func (n *Adapter) TryConnectAuthorized(remote nscore.Endpoint, authorize Connect stream.portLease.ReleaseLocked() return nil, 0, lnetocore.MapError(err) } - conn := &stream.connValue + conn := stream.conn if err := conn.Configure(lnetotcp.ConnConfig{ RxBuf: stream.storage[:n.config.ReceiveBytes], TxBuf: stream.storage[n.config.ReceiveBytes:], @@ -698,7 +728,7 @@ func (n *Adapter) TryConnectAuthorized(remote nscore.Endpoint, authorize Connect RWBackoff: immediateBackoff, }); err != nil { stream.allocation.Release() - stream.retained.ResetReleased() + stream.outboundStorage.retained.ResetReleased() stream.portLease.ReleaseLocked() n.recycleOutboundStreamLocked(stream) return nil, 0, lnetocore.MapError(err) @@ -706,7 +736,7 @@ func (n *Adapter) TryConnectAuthorized(remote nscore.Endpoint, authorize Connect if err := n.stack.DialTCP(conn, localPort, netip.AddrPortFrom(remote.Address, remote.Port)); err != nil { conn.Abort() stream.allocation.Release() - stream.retained.ResetReleased() + stream.outboundStorage.retained.ResetReleased() stream.portLease.ReleaseLocked() n.recycleOutboundStreamLocked(stream) return nil, 0, lnetocore.MapError(err) @@ -1050,8 +1080,8 @@ func (s *tcpStream) closeLocked() error { } if s.allocation != nil { s.allocation.Release() - if s.allocation == &s.retained { - s.retained.ResetReleased() + if s.outboundStorage != nil && s.allocation == &s.outboundStorage.retained { + s.outboundStorage.retained.ResetReleased() } s.allocation = nil } @@ -1120,7 +1150,10 @@ func (n *Adapter) CloseLocked() { n.streams[len(n.streams)-1].closeLocked() } n.freeListenerPool.destroyLocked() - clear(n.freeOutboundStorage) + if n.freeOutboundStorage != nil { + n.recycleOutboundStorageLocked(n.freeOutboundStorage) + n.freeOutboundStorage = nil + } n.portOwner = nil n.freeListenerPool = tcpPool{} n.listeners = nil diff --git a/internal/backend/lneto/tcp/tcp_test.go b/internal/backend/lneto/tcp/tcp_test.go index 18fef47..964d185 100644 --- a/internal/backend/lneto/tcp/tcp_test.go +++ b/internal/backend/lneto/tcp/tcp_test.go @@ -101,7 +101,7 @@ func TestValidConfigRejectsOverflowAndKeepsAdapterCreationBounded(t *testing.T) t.Fatalf("stream capacity hint = %d, want %d", cap(adapter.streams), maxTCPStreamCapacityHint) } if len(adapter.listeners) != 0 || len(adapter.streams) != 0 || adapter.freeListenerPool.slots != nil || adapter.freeOutboundStorage != nil { - t.Fatalf("new adapter eagerly populated state: listeners=%d streams=%d listener-cache=%v outbound-cache=%d", len(adapter.listeners), len(adapter.streams), adapter.freeListenerPool.slots != nil, len(adapter.freeOutboundStorage)) + t.Fatalf("new adapter eagerly populated state: listeners=%d streams=%d listener-cache=%v outbound-cache=%d", len(adapter.listeners), len(adapter.streams), adapter.freeListenerPool.slots != nil, len(cachedOutboundBytes(adapter))) } } @@ -231,7 +231,7 @@ func TestQuotaDeniedConnectDoesNotAllocateOrRetainOutboundStorage(t *testing.T) leaseCount := common.TCPPortLeaseCountLocked() common.Unlock() if adapter.freeOutboundStorage != nil || len(adapter.streams) != 0 || leaseCount != 0 || adapter.outboundStreams != 0 { - t.Fatalf("quota denied connect retained state: storage=%d streams=%d leases=%d outbound=%d", len(adapter.freeOutboundStorage), len(adapter.streams), leaseCount, adapter.outboundStreams) + t.Fatalf("quota denied connect retained state: storage=%d streams=%d leases=%d outbound=%d", len(cachedOutboundBytes(adapter)), len(adapter.streams), leaseCount, adapter.outboundStreams) } if usage, closed := adapter.quotas.Snapshot(); closed || usage != (quota.Usage{QueuedBytes: 16 << 10}) { t.Fatalf("quota denied usage = %+v, closed=%v", usage, closed) @@ -241,6 +241,50 @@ func TestQuotaDeniedConnectDoesNotAllocateOrRetainOutboundStorage(t *testing.T) } } +func TestQuotaDeniedConnectPreservesClearedIdleOutboundStorage(t *testing.T) { + common, adapter := newTestAdapter(t, 36, 0, 1) + remote := nscore.Endpoint{Address: netip.MustParseAddr("192.0.2.37"), Port: 4037} + value, progress, err := adapter.TryConnect(remote) + if err != nil || progress != nscore.ProgressInProgress { + t.Fatalf("initial connect = %T, %v, %v", value, progress, err) + } + if err := value.Close(); err != nil { + t.Fatal(err) + } + cached := adapter.freeOutboundStorage + if cached == nil || len(cached.bytes) != 512 { + t.Fatalf("initial cache = %p bytes=%d", cached, len(cachedOutboundBytes(adapter))) + } + cachedAddress := &cached.bytes[0] + var occupied quota.Charge + if err := adapter.quotas.AcquireQueuedBytes(&occupied, 16<<10); err != nil { + t.Fatal(err) + } + if value, progress, err := adapter.TryConnect(remote); value != nil || progress != 0 || failureOf(t, err) != nscore.FailureResourceLimit { + t.Fatalf("quota denied reuse = %T, %v, %v", value, progress, err) + } + if adapter.freeOutboundStorage != cached || len(cached.bytes) != 512 || &cached.bytes[0] != cachedAddress { + t.Fatalf("quota denial replaced cache = got:%p want:%p bytes:%d", adapter.freeOutboundStorage, cached, len(cached.bytes)) + } + for i, value := range cached.bytes { + if value != 0 { + t.Fatalf("cache byte %d = %d", i, value) + } + } + common.Lock() + leaseCount := common.TCPPortLeaseCountLocked() + common.Unlock() + if leaseCount != 0 || len(adapter.streams) != 0 || adapter.outboundStreams != 0 { + t.Fatalf("quota denial retained state: leases=%d streams=%d outbound=%d", leaseCount, len(adapter.streams), adapter.outboundStreams) + } + if usage, _ := adapter.quotas.Snapshot(); usage != (quota.Usage{QueuedBytes: 16 << 10}) { + t.Fatalf("quota denied usage = %+v", usage) + } + if !occupied.Release() || !occupied.ResetReleased() { + t.Fatal("release occupied quota") + } +} + func TestOutboundStreamCountTracksReuseAndCoreClose(t *testing.T) { core, adapter := newTestAdapter(t, 4, 0, 2) remote := nscore.Endpoint{Address: netip.MustParseAddr("192.0.2.5"), Port: 4205} @@ -297,6 +341,8 @@ func TestOutboundStorageReuseClearsDataAndIsolatesStaleStream(t *testing.T) { } stale := firstValue.(*tcpStream) stalePort := stale.local.Port + storageOwner := stale.outboundStorage + connAddress := stale.conn storage := stale.storage for i := range storage { storage[i] = byte(i | 1) @@ -305,13 +351,14 @@ func TestOutboundStorageReuseClearsDataAndIsolatesStaleStream(t *testing.T) { if err := stale.Close(); err != nil { t.Fatal(err) } - if stale.storage != nil || stale.conn != nil || stale.allocation != nil || !stale.closed { - t.Fatalf("closed stale stream retained state: storage=%v conn=%p allocation=%p closed=%v", stale.storage, stale.conn, stale.allocation, stale.closed) + if stale.storage != nil || stale.conn != nil || stale.outboundStorage != nil || stale.allocation != nil || !stale.closed { + t.Fatalf("closed stale stream retained state: storage=%v owner=%p conn=%p allocation=%p closed=%v", stale.storage, stale.outboundStorage, stale.conn, stale.allocation, stale.closed) } - if len(adapter.freeOutboundStorage) != len(storage) || &adapter.freeOutboundStorage[0] != storageAddress { - t.Fatalf("outbound storage was not retained for bounded reuse: bytes=%d", len(adapter.freeOutboundStorage)) + cached := cachedOutboundBytes(adapter) + if adapter.freeOutboundStorage != storageOwner || len(cached) != len(storage) || &cached[0] != storageAddress { + t.Fatalf("outbound storage was not retained for bounded reuse: owner=%p want:%p bytes=%d", adapter.freeOutboundStorage, storageOwner, len(cached)) } - for i, value := range adapter.freeOutboundStorage { + for i, value := range cached { if value != 0 { t.Fatalf("recycled storage byte %d = %d", i, value) } @@ -321,8 +368,8 @@ func TestOutboundStorageReuseClearsDataAndIsolatesStaleStream(t *testing.T) { t.Fatalf("fresh connect = %T, %v, %v", freshValue, progress, err) } fresh := freshValue.(*tcpStream) - if fresh == stale || fresh.local.Port == stalePort || &fresh.storage[0] != storageAddress { - t.Fatalf("fresh reuse = same wrapper %v, port %d stale %d, storage reused %v", fresh == stale, fresh.local.Port, stalePort, &fresh.storage[0] == storageAddress) + if fresh == stale || fresh.conn == connAddress || fresh.local.Port == stalePort || fresh.outboundStorage != storageOwner || &fresh.storage[0] != storageAddress { + t.Fatalf("fresh reuse = same wrapper %v, conn reused %v, port %d stale %d, owner reused %v, storage reused %v", fresh == stale, fresh.conn == connAddress, fresh.local.Port, stalePort, fresh.outboundStorage == storageOwner, &fresh.storage[0] == storageAddress) } if progress, err := stale.TryFinishConnect(); progress != 0 || failureOf(t, err) != nscore.FailureClosed { t.Fatalf("stale finish connect = %v, %v", progress, err) @@ -373,8 +420,9 @@ func TestOutboundStorageReuseRetainsAtMostOneIdleBuffer(t *testing.T) { t.Fatalf("close %d: %v", i, err) } } - if len(adapter.freeOutboundStorage) != 512 || &adapter.freeOutboundStorage[0] != wantCached { - t.Fatalf("idle outbound cache = bytes:%d address:%p want:%p", len(adapter.freeOutboundStorage), firstByte(adapter.freeOutboundStorage), wantCached) + cached := cachedOutboundBytes(adapter) + if len(cached) != 512 || &cached[0] != wantCached { + t.Fatalf("idle outbound cache = bytes:%d address:%p want:%p", len(cached), firstByte(cached), wantCached) } for i, storage := range storages { for j, value := range storage { @@ -393,7 +441,7 @@ func TestOutboundStorageReuseRetainsAtMostOneIdleBuffer(t *testing.T) { t.Fatal(err) } if adapter.freeOutboundStorage != nil { - t.Fatalf("namespace close retained %d outbound bytes", len(adapter.freeOutboundStorage)) + t.Fatalf("namespace close retained %d outbound bytes", len(cachedOutboundBytes(adapter))) } } @@ -453,10 +501,10 @@ func TestIdleReuseDropsOversizedTCPStorageAndListenerMetadata(t *testing.T) { for i := range outboundStorage { outboundStorage[i] = 0xff } - stream := &tcpStream{owner: adapter, storage: outboundStorage} + stream := &tcpStream{owner: adapter, storage: outboundStorage, outboundStorage: &tcpOutboundStorage{bytes: outboundStorage}} adapter.recycleOutboundStreamLocked(stream) - if adapter.freeOutboundStorage != nil || stream.storage != nil { - t.Fatalf("oversized outbound storage retained: cache=%d stream=%d", len(adapter.freeOutboundStorage), len(stream.storage)) + if adapter.freeOutboundStorage != nil || stream.storage != nil || stream.outboundStorage != nil { + t.Fatalf("oversized outbound storage retained: cache=%d stream=%d owner=%p", len(cachedOutboundBytes(adapter)), len(stream.storage), stream.outboundStorage) } for i, value := range outboundStorage { if value != 0 { @@ -496,6 +544,13 @@ func TestIdleReuseDropsOversizedTCPStorageAndListenerMetadata(t *testing.T) { } } +func cachedOutboundBytes(adapter *Adapter) []byte { + if adapter == nil || adapter.freeOutboundStorage == nil { + return nil + } + return adapter.freeOutboundStorage.bytes +} + func firstByte(storage []byte) *byte { if len(storage) == 0 { return nil @@ -1396,6 +1451,7 @@ func TestConnectResetBeforeEstablishment(t *testing.T) { t.Fatalf("connect = %T, %v, %v", resource, progress, err) } stream := resource.(*tcpStream) + storageOwner := stream.outboundStorage if usage, closed := adapter.quotas.Snapshot(); closed || usage != (quota.Usage{Resources: 1, TCPResources: 1, QueuedBytes: 512}) { t.Fatalf("outbound quota = %+v, closed=%v", usage, closed) } @@ -1411,9 +1467,9 @@ func TestConnectResetBeforeEstablishment(t *testing.T) { if err := stream.Close(); err != nil { t.Fatal(err) } - streamRetainedReset := stream.retained.ResetReleased() - if stream.conn != nil || streamRetainedReset { - t.Fatalf("closed outbound stream retained graph state: conn=%p retained_reset=%v", stream.conn, streamRetainedReset) + streamRetainedReset := storageOwner.retained.ResetReleased() + if stream.conn != nil || stream.outboundStorage != nil || streamRetainedReset { + t.Fatalf("closed outbound stream retained graph state: conn=%p owner=%p retained_reset=%v", stream.conn, stream.outboundStorage, streamRetainedReset) } if usage, _ := adapter.quotas.Snapshot(); usage != (quota.Usage{}) { t.Fatalf("outbound close quota = %+v", usage)