diff --git a/pkg/hive/export_test.go b/pkg/hive/export_test.go index 3d190772f18..016c2424e1d 100644 --- a/pkg/hive/export_test.go +++ b/pkg/hive/export_test.go @@ -11,8 +11,11 @@ import ( "github.com/ethersphere/bee/v2/pkg/hive/pb" ) -var MaxBatchSize = maxBatchSize -var LimitBurst = limitBurst +var ( + MaxBatchSize = maxBatchSize + LimitBurst = limitBurst + CoalesceThreshold = coalesceThreshold +) func (s *Service) SetTimeFunc(f func() time.Time) { s.now = f diff --git a/pkg/hive/gossip_buffer.go b/pkg/hive/gossip_buffer.go new file mode 100644 index 00000000000..68e1ff0a59a --- /dev/null +++ b/pkg/hive/gossip_buffer.go @@ -0,0 +1,104 @@ +// Copyright 2026 The Swarm Authors. All rights reserved. +// Use of this source code is governed by a BSD-style +// license that can be found in the LICENSE file. + +package hive + +import ( + "maps" + "slices" + "sync" + "time" + + "github.com/ethersphere/bee/v2/pkg/swarm" +) + +const ( + defaultGossipCoalesceInterval = 5 * time.Second + // coalesceThreshold: gossips with fewer peers are buffered; larger + // (already-batched) messages are dispatched immediately. + coalesceThreshold = 5 +) + +// gossipBuffer accumulates single-peer outbound gossip per addressee so it can be +// flushed as one batched message. +type gossipBuffer struct { + mu sync.Mutex + pending map[string]map[string]swarm.Address // addressee key -> peer key -> peer + interval time.Duration + maxBatch int +} + +type gossipBatch struct { + addressee swarm.Address + peers []swarm.Address +} + +func newGossipBuffer(interval time.Duration, maxBatch int) *gossipBuffer { + if interval == 0 { + interval = defaultGossipCoalesceInterval + } + return &gossipBuffer{ + pending: make(map[string]map[string]swarm.Address), + interval: interval, + maxBatch: maxBatch, + } +} + +// stagePeers buffers peers for the addressee. If the buffer reaches maxBatch it is +// removed and returned so the caller can flush it immediately. A new addressee is +// not written to pending if the merged set already meets maxBatch. +func (b *gossipBuffer) stagePeers(addressee swarm.Address, peers ...swarm.Address) (flushPeers []swarm.Address, flush bool) { + b.mu.Lock() + defer b.mu.Unlock() + + key := addressee.ByteString() + peerSet, exist := b.pending[key] + if !exist { + peerSet = make(map[string]swarm.Address) + } + for _, p := range peers { + peerSet[p.ByteString()] = p + } + if len(peerSet) >= b.maxBatch { + if exist { + delete(b.pending, key) + } + return slices.Collect(maps.Values(peerSet)), true + } + + b.pending[key] = peerSet + return nil, false +} + +// takeAll removes and returns all buffered entries. +func (b *gossipBuffer) takeAll() []gossipBatch { + b.mu.Lock() + defer b.mu.Unlock() + + if len(b.pending) == 0 { + return nil + } + + out := make([]gossipBatch, 0, len(b.pending)) + for key, peerSet := range b.pending { + out = append(out, gossipBatch{ + addressee: swarm.NewAddress([]byte(key)), + peers: slices.Collect(maps.Values(peerSet)), + }) + } + clear(b.pending) + return out +} + +func (b *gossipBuffer) clearAddressee(addressee swarm.Address) { + b.mu.Lock() + defer b.mu.Unlock() + delete(b.pending, addressee.ByteString()) +} + +func (b *gossipBuffer) pendingAddressees() int { + b.mu.Lock() + defer b.mu.Unlock() + return len(b.pending) +} diff --git a/pkg/hive/gossip_buffer_test.go b/pkg/hive/gossip_buffer_test.go new file mode 100644 index 00000000000..365cf85746a --- /dev/null +++ b/pkg/hive/gossip_buffer_test.go @@ -0,0 +1,87 @@ +// Copyright 2026 The Swarm Authors. All rights reserved. +// Use of this source code is governed by a BSD-style +// license that can be found in the LICENSE file. + +package hive + +import ( + "testing" + "time" + + "github.com/ethersphere/bee/v2/pkg/swarm" +) + +func TestGossipBufferAddAndTakeAll(t *testing.T) { + t.Parallel() + + b := newGossipBuffer(time.Second, maxBatchSize) + addressee := swarm.RandAddress(t) + peer1 := swarm.RandAddress(t) + peer2 := swarm.RandAddress(t) + + if pending := b.takeAll(); len(pending) != 0 { + t.Fatalf("want no pending entries, got %d", len(pending)) + } + + if _, flush := b.stagePeers(addressee, peer1); flush { + t.Fatal("unexpected immediate flush") + } + + if _, flush := b.stagePeers(addressee, peer2); flush { + t.Fatal("unexpected immediate flush") + } + + pending := b.takeAll() + if len(pending) != 1 { + t.Fatalf("want 1 pending entry, got %d", len(pending)) + } + if got := len(pending[0].peers); got != 2 { + t.Fatalf("want 2 coalesced peers, got %d", got) + } + if !pending[0].addressee.Equal(addressee) { + t.Fatal("unexpected addressee in pending batch") + } + + if pending := b.takeAll(); len(pending) != 0 { + t.Fatalf("want empty buffer after takeAll, got %d pending", len(pending)) + } +} + +func TestGossipBufferMaxBatchFlush(t *testing.T) { + t.Parallel() + + b := newGossipBuffer(time.Second, 2) + addressee := swarm.RandAddress(t) + + b.stagePeers(addressee, swarm.RandAddress(t)) + flushPeers, flush := b.stagePeers(addressee, swarm.RandAddress(t)) + if !flush { + t.Fatal("want immediate flush at maxBatch") + } + if got := len(flushPeers); got != 2 { + t.Fatalf("want 2 peers in full batch, got %d", got) + } + if pending := b.takeAll(); len(pending) != 0 { + t.Fatalf("want empty buffer after maxBatch flush, got %d pending", len(pending)) + } +} + +func TestGossipBufferMaxBatchFlushWithoutPendingInsert(t *testing.T) { + t.Parallel() + + b := newGossipBuffer(time.Second, 2) + addressee := swarm.RandAddress(t) + peer1 := swarm.RandAddress(t) + peer2 := swarm.RandAddress(t) + + flushPeers, flush := b.stagePeers(addressee, peer1, peer2) + if !flush { + t.Fatal("want immediate flush when first insert meets maxBatch") + } + if got := len(flushPeers); got != 2 { + t.Fatalf("want 2 peers in full batch, got %d", got) + } + if pending := b.takeAll(); len(pending) != 0 { + t.Fatalf("want no pending insert after maxBatch flush, got %d pending", len(pending)) + } +} diff --git a/pkg/hive/hive.go b/pkg/hive/hive.go index 4f4fb7d1bcb..d60a195bd17 100644 --- a/pkg/hive/hive.go +++ b/pkg/hive/hive.go @@ -18,6 +18,10 @@ import ( "time" "github.com/ethereum/go-ethereum/common" + ma "github.com/multiformats/go-multiaddr" + manet "github.com/multiformats/go-multiaddr/net" + "golang.org/x/sync/semaphore" + "github.com/ethersphere/bee/v2/pkg/addressbook" "github.com/ethersphere/bee/v2/pkg/bzz" "github.com/ethersphere/bee/v2/pkg/hive/pb" @@ -28,9 +32,6 @@ import ( "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/settlement/swap/chequebook" "github.com/ethersphere/bee/v2/pkg/swarm" - ma "github.com/multiformats/go-multiaddr" - manet "github.com/multiformats/go-multiaddr/net" - "golang.org/x/sync/semaphore" ) // ChequebookStorer persists the overlay→chequebook mapping. Put holds its @@ -58,6 +59,11 @@ var ( ErrRateLimitExceeded = errors.New("rate limit exceeded") ) +const ( + coalesceFlushReasonTimer = "timer" + coalesceFlushReasonMaxBatch = "max_batch" +) + // Options configures hive.Service at construction. Chequebook fields are // optional: a nil ChequebookVerifier disables the verification gate (and // records without a chequebook are accepted); a nil ChequebookStorer means @@ -67,6 +73,8 @@ type Options struct { AllowPrivateCIDRs bool ChequebookVerifier chequebook.Verifier ChequebookStorer ChequebookStorer + + GossipCoalesceInterval time.Duration } type Service struct { @@ -91,6 +99,7 @@ type Service struct { // chequebook are dropped. chequebookVerifier chequebook.Verifier chequebookStorer ChequebookStorer + gossipBuf *gossipBuffer } func New(streamer p2p.Streamer, addressbook addressbook.GetPutSeener, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service { @@ -113,9 +122,12 @@ func New(streamer p2p.Streamer, addressbook addressbook.GetPutSeener, networkID chequebookStorer: o.ChequebookStorer, } + svc.gossipBuf = newGossipBuffer(o.GossipCoalesceInterval, maxBatchSize) + if !o.BootnodeMode { svc.startCheckPeersHandler() } + svc.startGossipCoalescer() return svc } @@ -137,34 +149,77 @@ func (s *Service) Protocol() p2p.ProtocolSpec { var ErrShutdownInProgress = errors.New("shutdown in progress") +// BroadcastPeers sends peer gossip to the addressee. Calls with fewer than +// coalesceThreshold peers are buffered and flushed asynchronously; errors +// during deferred dispatch are logged but not returned to the caller. +// Calls with coalesceThreshold or more peers are sent immediately. func (s *Service) BroadcastPeers(ctx context.Context, addressee swarm.Address, peers ...swarm.Address) error { - maxSize := maxBatchSize + if len(peers) == 0 { + return nil + } + s.metrics.BroadcastPeers.Inc() s.metrics.BroadcastPeersPeers.Add(float64(len(peers))) + // Already-batched messages go out immediately; single-peer gossips are coalesced. + if len(peers) >= coalesceThreshold { + s.metrics.GossipCoalesceImmediatePeers.Add(float64(len(peers))) + s.logger.Debug("gossip immediate send", "addressee", addressee, "peer_count", len(peers)) + return s.broadcastNow(ctx, addressee, false, peers...) + } + + select { + case <-s.quit: + return ErrShutdownInProgress + default: + } + + s.metrics.GossipCoalesceBufferedPeers.Add(float64(len(peers))) + s.logger.Debug("gossip buffered", "addressee", addressee, "peer_count", len(peers)) + + // Buffer; if it just filled up, flush it synchronously while still in the call + if flushPeers, flush := s.gossipBuf.stagePeers(addressee, peers...); flush { + s.recordCoalesceFlush(coalesceFlushReasonMaxBatch, addressee, flushPeers) + s.setCoalesceBufferGauge() + return s.broadcastNow(ctx, addressee, true, flushPeers...) + } + s.setCoalesceBufferGauge() + return nil +} + +// broadcastNow performs the synchronous, rate-limited, batched send. +func (s *Service) broadcastNow(ctx context.Context, addressee swarm.Address, coalesced bool, peers ...swarm.Address) error { + maxSize := maxBatchSize + for len(peers) > 0 { if maxSize > len(peers) { maxSize = len(peers) } - // If broadcasting limit is exceeded, return early if !s.outLimiter.Allow(addressee.ByteString(), maxSize) { + if coalesced { + s.metrics.GossipCoalesceDropped.Add(float64(len(peers))) + } return nil } select { + case <-ctx.Done(): + return ctx.Err() case <-s.quit: return ErrShutdownInProgress default: } if err := s.sendPeers(ctx, addressee, peers[:maxSize]); err != nil { + if coalesced { + s.metrics.GossipCoalesceDropped.Add(float64(len(peers))) + } return err } peers = peers[maxSize:] } - return nil } @@ -297,9 +352,59 @@ func (s *Service) peersHandler(ctx context.Context, peer p2p.Peer, stream p2p.St func (s *Service) disconnect(peer p2p.Peer) error { s.inLimiter.Clear(peer.Address.ByteString()) s.outLimiter.Clear(peer.Address.ByteString()) + s.gossipBuf.clearAddressee(peer.Address) + s.setCoalesceBufferGauge() return nil } +func (s *Service) startGossipCoalescer() { + s.wg.Go(func() { + ticker := time.NewTicker(s.gossipBuf.interval) + defer ticker.Stop() + for { + select { + case <-ticker.C: + for _, batch := range s.gossipBuf.takeAll() { + go func(batch gossipBatch) { + s.flushGossipBatch(batch.addressee, batch.peers, coalesceFlushReasonTimer) + }(batch) + } + + case <-s.quit: + return + } + } + }) +} + +func (s *Service) flushGossipBatch(addressee swarm.Address, peers []swarm.Address, reason string) { + s.recordCoalesceFlush(reason, addressee, peers) + + ctx, cancel := context.WithTimeout(context.Background(), messageTimeout) + defer cancel() + + err := s.broadcastNow(ctx, addressee, true, peers...) + if err != nil { + s.logger.Debug("coalesced gossip flush failed", "addressee", addressee, "reason", reason, "batch_size", len(peers), "error", err) + } + s.setCoalesceBufferGauge() +} + +func (s *Service) recordCoalesceFlush(reason string, addressee swarm.Address, peers []swarm.Address) { + batchSize := len(peers) + if batchSize == 0 { + return + } + + s.metrics.GossipCoalesceFlushTotal.WithLabelValues(reason).Inc() + s.metrics.GossipCoalesceFlushPeers.Add(float64(batchSize)) + s.logger.Debug("coalesced gossip flush", "addressee", addressee, "reason", reason, "batch_size", batchSize) +} + +func (s *Service) setCoalesceBufferGauge() { + s.metrics.GossipCoalesceBufferSize.Set(float64(s.gossipBuf.pendingAddressees())) +} + func (s *Service) startCheckPeersHandler() { ctx, cancel := context.WithCancel(context.Background()) s.wg.Go(func() { diff --git a/pkg/hive/hive_test.go b/pkg/hive/hive_test.go index 9f6f599da43..ac1df519540 100644 --- a/pkg/hive/hive_test.go +++ b/pkg/hive/hive_test.go @@ -17,12 +17,16 @@ import ( "time" "github.com/ethereum/go-ethereum/common" + ma "github.com/multiformats/go-multiaddr" + "github.com/multiformats/go-varint" + ab "github.com/ethersphere/bee/v2/pkg/addressbook" "github.com/ethersphere/bee/v2/pkg/bzz" "github.com/ethersphere/bee/v2/pkg/crypto" "github.com/ethersphere/bee/v2/pkg/hive" "github.com/ethersphere/bee/v2/pkg/hive/pb" "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/p2p" "github.com/ethersphere/bee/v2/pkg/p2p/protobuf" "github.com/ethersphere/bee/v2/pkg/p2p/streamtest" "github.com/ethersphere/bee/v2/pkg/settlement/swap/chequebook" @@ -31,13 +35,89 @@ import ( "github.com/ethersphere/bee/v2/pkg/statestore/mock" "github.com/ethersphere/bee/v2/pkg/swarm" "github.com/ethersphere/bee/v2/pkg/util/testutil" - ma "github.com/multiformats/go-multiaddr" - "github.com/multiformats/go-varint" ) var nonce = common.HexToHash("0x2").Bytes() -const spinTimeout = time.Second * 5 +const ( + spinTimeout = time.Second * 5 + testCoalesceInterval = 100 * time.Millisecond + // Must exceed one coalesce interval before the ticker flushes. + testCoalesceWait = 150 * time.Millisecond +) + +func waitForCoalesceFlush(t *testing.T) { + t.Helper() + time.Sleep(testCoalesceWait) + synctest.Wait() +} + +func newCoalescingClient( + t *testing.T, + recorder p2p.Streamer, + addressbook ab.GetPutSeener, + networkID uint64, + overlay swarm.Address, + logger log.Logger, + opts hive.Options, +) *hive.Service { + t.Helper() + + opts.GossipCoalesceInterval = testCoalesceInterval + client := hive.New(recorder, addressbook, networkID, overlay, logger, opts) + testutil.CleanupCloser(t, client) + return client +} + +// addTestOverlays puts n random peers into the addressbook and returns their overlays. +func addTestOverlays(t *testing.T, book ab.Putter, networkID uint64, n, portBase int) []swarm.Address { + t.Helper() + + if n <= 0 { + return nil + } + + overlays := make([]swarm.Address, n) + for i := range n { + underlay, err := ma.NewMultiaddr("/ip4/127.0.0.1/udp/" + strconv.Itoa(portBase+i)) + if err != nil { + t.Fatal(err) + } + pk, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + signer := crypto.NewDefaultSigner(pk) + overlay, err := crypto.NewOverlayAddress(pk.PublicKey, networkID, nonce) + if err != nil { + t.Fatal(err) + } + bzzAddr, err := bzz.NewAddress(signer, []ma.Multiaddr{underlay}, overlay, networkID, nonce, 1, common.Address{}) + if err != nil { + t.Fatal(err) + } + if err := book.Put(bzzAddr.Overlay, *bzzAddr, true); err != nil { + t.Fatal(err) + } + overlays[i] = bzzAddr.Overlay + } + return overlays +} + +func assertNoGossipRecords(t *testing.T, recorder *streamtest.Recorder, addressee swarm.Address) { + t.Helper() + + records, err := recorder.Records(addressee, "hive", "2.0.0", "peers") + if err == nil { + if len(records) != 0 { + t.Fatalf("got %d gossip records, want none before coalesce flush", len(records)) + } + return + } + if !errors.Is(err, streamtest.ErrRecordsNotFound) { + t.Fatal(err) + } +} func TestHandlerRateLimit(t *testing.T) { t.Parallel() @@ -200,14 +280,6 @@ func TestBroadcastPeers(t *testing.T) { wantBzzAddresses []bzz.Address allowPrivateCIDRs bool }{ - "OK - single record": { - addresee: swarm.MustParseHexAddress("ca1e9f3938cc1425c6061b96ad9eb93e134dfe8734ad490164ef20af9d1cf59c"), - peers: []swarm.Address{overlays[0]}, - wantMsgs: []pb.Peers{{Peers: wantMsgs[0].Peers[:1]}}, - wantOverlays: []swarm.Address{overlays[0]}, - wantBzzAddresses: []bzz.Address{bzzAddresses[0]}, - allowPrivateCIDRs: true, - }, "OK - single batch - multiple records": { addresee: swarm.MustParseHexAddress("ca1e9f3938cc1425c6061b96ad9eb93e134dfe8734ad490164ef20af9d1cf59c"), peers: overlays[:15], @@ -250,7 +322,9 @@ func TestBroadcastPeers(t *testing.T) { }, "Ok - don't advertise private CIDRs only (but include one public peer)": { addresee: overlays[0], - peers: overlays[58:], + // Last CoalesceThreshold peers so the call is sent immediately; + // only the final overlay is public. + peers: overlays[len(overlays)-hive.CoalesceThreshold:], wantMsgs: []pb.Peers{{Peers: func() []*pb.BzzAddress { ub, err := bzz.SerializeUnderlays(bzzAddresses[len(bzzAddresses)-1].Underlays) if err != nil { @@ -299,11 +373,11 @@ func TestBroadcastPeers(t *testing.T) { // create a hive client that will do broadcast clientAddress := swarm.RandAddress(t) client := hive.New(recorder, addressbook, networkID, clientAddress, logger, hive.Options{AllowPrivateCIDRs: tc.allowPrivateCIDRs}) + testutil.CleanupCloser(t, client) if err := client.BroadcastPeers(context.Background(), tc.addresee, tc.peers...); err != nil { t.Fatal(err) } - testutil.CleanupCloser(t, client) // get a record for this stream records, err := recorder.Records(tc.addresee, "hive", "2.0.0", "peers") @@ -329,6 +403,86 @@ func TestBroadcastPeers(t *testing.T) { } } +func TestBroadcastPeersSingleCoalesced(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + logger := log.Noop + statestore := mock.NewStateStore() + addressbook := ab.New(statestore) + networkID := uint64(1) + + underlay, err := ma.NewMultiaddr("/ip4/127.0.0.1/udp/2000") + if err != nil { + t.Fatal(err) + } + pk, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + signer := crypto.NewDefaultSigner(pk) + overlay, err := crypto.NewOverlayAddress(pk.PublicKey, networkID, nonce) + if err != nil { + t.Fatal(err) + } + bzzAddr, err := bzz.NewAddress(signer, []ma.Multiaddr{underlay}, overlay, networkID, nonce, 1, common.Address{}) + if err != nil { + t.Fatal(err) + } + if err := addressbook.Put(bzzAddr.Overlay, *bzzAddr, true); err != nil { + t.Fatal(err) + } + + underlayBytes, err := bzz.SerializeUnderlays(bzzAddr.Underlays) + if err != nil { + t.Fatal(err) + } + wantMsg := pb.Peers{Peers: []*pb.BzzAddress{{ + Overlay: bzzAddr.Overlay.Bytes(), + Underlay: underlayBytes, + Signature: bzzAddr.Signature, + Nonce: nonce, + Timestamp: bzzAddr.Timestamp, + }}} + + addresee := swarm.MustParseHexAddress("ca1e9f3938cc1425c6061b96ad9eb93e134dfe8734ad490164ef20af9d1cf59c") + + addressbookclean := ab.New(mock.NewStateStore()) + + streamer := streamtest.New() + serverAddress := swarm.RandAddress(t) + server := hive.New(streamer, addressbookclean, networkID, serverAddress, logger, hive.Options{AllowPrivateCIDRs: true}) + testutil.CleanupCloser(t, server) + + recorder := streamtest.New(streamtest.WithProtocols(server.Protocol())) + + clientAddress := swarm.RandAddress(t) + client := newCoalescingClient(t, recorder, addressbook, networkID, clientAddress, logger, hive.Options{AllowPrivateCIDRs: true}) + + if err := client.BroadcastPeers(context.Background(), addresee, bzzAddr.Overlay); err != nil { + t.Fatal(err) + } + + assertNoGossipRecords(t, recorder, addresee) + waitForCoalesceFlush(t) + + records, err := recorder.Records(addresee, "hive", "2.0.0", "peers") + if err != nil { + t.Fatal(err) + } + if l := len(records); l != 1 { + t.Fatalf("got %v records, want 1", l) + } + + messages, err := readAndAssertPeersMsgs(records[0].In(), 1) + if err != nil { + t.Fatal(err) + } + comparePeerMsgs(t, messages[0].Peers, wantMsg.Peers) + + expectOverlaysEventually(t, addressbookclean, []swarm.Address{bzzAddr.Overlay}) + expectBzzAddresessEventually(t, addressbookclean, []bzz.Address{*bzzAddr}) + }) +} + func expectOverlaysEventually(t *testing.T, exporter ab.Interface, wantOverlays []swarm.Address) { t.Helper() @@ -536,8 +690,11 @@ func TestBroadcastPeersSkipsSelf(t *testing.T) { client := hive.New(serverRecorder, addressbook, networkID, clientAddress, logger, hive.Options{AllowPrivateCIDRs: true}) testutil.CleanupCloser(t, client) - // Try to broadcast: peer1, clientAddress (self), and another peer - peersIncludingSelf := []swarm.Address{bzzAddr1.Overlay, clientAddress, peer1} + // Try to broadcast: peer1, clientAddress (self), and another peer. + // Pad to CoalesceThreshold so the message is sent immediately. + peersIncludingSelf := make([]swarm.Address, 0, hive.CoalesceThreshold) + peersIncludingSelf = append(peersIncludingSelf, bzzAddr1.Overlay, clientAddress, peer1) + peersIncludingSelf = append(peersIncludingSelf, addTestOverlays(t, addressbook, networkID, hive.CoalesceThreshold-len(peersIncludingSelf), 5000)...) err = client.BroadcastPeers(context.Background(), serverAddress, peersIncludingSelf...) if err != nil { @@ -659,8 +816,11 @@ func TestReceivePeersSkipsSelf(t *testing.T) { client := hive.New(serverRecorder, addressbook, networkID, clientAddress, logger, hive.Options{AllowPrivateCIDRs: true}) testutil.CleanupCloser(t, client) - // Client broadcasts: valid peer and server's own address - peersIncludingSelf := []swarm.Address{bzzAddr1.Overlay, serverAddress} + // Client broadcasts: valid peer and server's own address. + // Pad to CoalesceThreshold so the message is sent immediately. + peersIncludingSelf := make([]swarm.Address, 0, hive.CoalesceThreshold) + peersIncludingSelf = append(peersIncludingSelf, bzzAddr1.Overlay, serverAddress) + peersIncludingSelf = append(peersIncludingSelf, addTestOverlays(t, addressbook, networkID, hive.CoalesceThreshold-len(peersIncludingSelf), 6000)...) err = client.BroadcastPeers(context.Background(), serverAddress, peersIncludingSelf...) if err != nil { @@ -1235,7 +1395,10 @@ func TestBroadcastSkipsLegacyZeroRecord(t *testing.T) { t.Cleanup(func() { _ = client.Close() }) addressee := swarm.RandAddress(t) - if err := client.BroadcastPeers(context.Background(), addressee, legacyPeer.overlay, modernPeer.overlay); err != nil { + peers := make([]swarm.Address, 0, hive.CoalesceThreshold) + peers = append(peers, legacyPeer.overlay, modernPeer.overlay) + peers = append(peers, addTestOverlays(t, clientBook, networkID, hive.CoalesceThreshold-len(peers), 7000)...) + if err := client.BroadcastPeers(context.Background(), addressee, peers...); err != nil { t.Fatalf("BroadcastPeers: %v", err) } @@ -1253,11 +1416,21 @@ func TestBroadcastSkipsLegacyZeroRecord(t *testing.T) { } sent := messages[0].Peers - if len(sent) != 1 { - t.Fatalf("expected 1 peer sent (legacy dropped), got %d", len(sent)) + if len(sent) != hive.CoalesceThreshold-1 { + t.Fatalf("expected %d peers sent (legacy dropped), got %d", hive.CoalesceThreshold-1, len(sent)) + } + foundModern := false + for _, p := range sent { + overlay := swarm.NewAddress(p.Overlay) + if overlay.Equal(legacyPeer.overlay) { + t.Fatal("legacy overlay should have been dropped") + } + if overlay.Equal(modernPeer.overlay) { + foundModern = true + } } - if !swarm.NewAddress(sent[0].Overlay).Equal(modernPeer.overlay) { - t.Fatalf("expected modern overlay %s, got %s", modernPeer.overlay, swarm.NewAddress(sent[0].Overlay)) + if !foundModern { + t.Fatalf("expected modern overlay %s in sent peers", modernPeer.overlay) } } @@ -1344,3 +1517,249 @@ func TestHiveGossipUnderlayCaps(t *testing.T) { } }) } + +func TestBroadcastPeersCoalesce(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + logger := log.Noop + statestore := mock.NewStateStore() + addressbook := ab.New(statestore) + networkID := uint64(1) + + overlays := make([]swarm.Address, 3) + for i := range overlays { + underlay, err := ma.NewMultiaddr("/ip4/127.0.0.1/udp/" + strconv.Itoa(2000+i)) + if err != nil { + t.Fatal(err) + } + pk, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + signer := crypto.NewDefaultSigner(pk) + overlay, err := crypto.NewOverlayAddress(pk.PublicKey, networkID, nonce) + if err != nil { + t.Fatal(err) + } + bzzAddr, err := bzz.NewAddress(signer, []ma.Multiaddr{underlay}, overlay, networkID, nonce, 1, common.Address{}) + if err != nil { + t.Fatal(err) + } + if err := addressbook.Put(bzzAddr.Overlay, *bzzAddr, true); err != nil { + t.Fatal(err) + } + overlays[i] = bzzAddr.Overlay + } + + streamer := streamtest.New() + serverAddress := swarm.RandAddress(t) + server := hive.New(streamer, ab.New(mock.NewStateStore()), networkID, serverAddress, logger, hive.Options{AllowPrivateCIDRs: true}) + testutil.CleanupCloser(t, server) + + recorder := streamtest.New(streamtest.WithProtocols(server.Protocol())) + clientAddress := swarm.RandAddress(t) + client := newCoalescingClient(t, recorder, addressbook, networkID, clientAddress, logger, hive.Options{ + AllowPrivateCIDRs: true, + }) + + ctx := context.Background() + for _, overlay := range overlays { + if err := client.BroadcastPeers(ctx, serverAddress, overlay); err != nil { + t.Fatal(err) + } + } + + assertNoGossipRecords(t, recorder, serverAddress) + waitForCoalesceFlush(t) + + records, err := recorder.Records(serverAddress, "hive", "2.0.0", "peers") + if err != nil { + t.Fatal(err) + } + if got, want := len(records), 1; got != want { + t.Fatalf("after flush got %d gossip messages, want %d", got, want) + } + + messages, err := readAndAssertPeersMsgs(records[0].In(), 1) + if err != nil { + t.Fatal(err) + } + if got, want := len(messages[0].Peers), len(overlays); got != want { + t.Fatalf("coalesced peer count: got %d, want %d", got, want) + } + + // Batched gossip is sent immediately without coalescing (use fresh peers). + batchedOverlays := make([]swarm.Address, hive.CoalesceThreshold) + for i := range batchedOverlays { + underlay, err := ma.NewMultiaddr("/ip4/127.0.0.1/udp/" + strconv.Itoa(3000+i)) + if err != nil { + t.Fatal(err) + } + pk, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + signer := crypto.NewDefaultSigner(pk) + overlay, err := crypto.NewOverlayAddress(pk.PublicKey, networkID, nonce) + if err != nil { + t.Fatal(err) + } + bzzAddr, err := bzz.NewAddress(signer, []ma.Multiaddr{underlay}, overlay, networkID, nonce, 1, common.Address{}) + if err != nil { + t.Fatal(err) + } + if err := addressbook.Put(bzzAddr.Overlay, *bzzAddr, true); err != nil { + t.Fatal(err) + } + batchedOverlays[i] = bzzAddr.Overlay + } + + if err := client.BroadcastPeers(ctx, serverAddress, batchedOverlays...); err != nil { + t.Fatal(err) + } + records, err = recorder.Records(serverAddress, "hive", "2.0.0", "peers") + if err != nil { + t.Fatal(err) + } + if got, want := len(records), 2; got != want { + t.Fatalf("after batched broadcast got %d gossip messages, want %d", got, want) + } + }) +} + +const hiveGossipBufferingInterval = time.Second + +func TestHiveGossipBuffering(t *testing.T) { + t.Parallel() + + makeAddressbookWithPeers := func(t *testing.T, n int) (ab.GetPutSeener, []swarm.Address) { + t.Helper() + + addressbook := ab.New(mock.NewStateStore()) + networkID := uint64(1) + overlays := make([]swarm.Address, n) + + for i := range n { + underlay, err := ma.NewMultiaddr("/ip4/127.0.0.1/udp/" + strconv.Itoa(4000+i)) + if err != nil { + t.Fatal(err) + } + pk, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + signer := crypto.NewDefaultSigner(pk) + overlay, err := crypto.NewOverlayAddress(pk.PublicKey, networkID, nonce) + if err != nil { + t.Fatal(err) + } + bzzAddr, err := bzz.NewAddress(signer, []ma.Multiaddr{underlay}, overlay, networkID, nonce, 1, common.Address{}) + if err != nil { + t.Fatal(err) + } + if err := addressbook.Put(bzzAddr.Overlay, *bzzAddr, true); err != nil { + t.Fatal(err) + } + overlays[i] = bzzAddr.Overlay + } + + return addressbook, overlays + } + + setupClient := func(t *testing.T, addressbook ab.GetPutSeener, coalesceInterval time.Duration) (*hive.Service, *streamtest.Recorder, swarm.Address) { + t.Helper() + + logger := log.Noop + networkID := uint64(1) + + streamer := streamtest.New() + serverAddress := swarm.RandAddress(t) + server := hive.New(streamer, ab.New(mock.NewStateStore()), networkID, serverAddress, logger, hive.Options{AllowPrivateCIDRs: true}) + testutil.CleanupCloser(t, server) + + recorder := streamtest.New(streamtest.WithProtocols(server.Protocol())) + clientAddress := swarm.RandAddress(t) + client := hive.New(recorder, addressbook, networkID, clientAddress, logger, hive.Options{ + AllowPrivateCIDRs: true, + GossipCoalesceInterval: coalesceInterval, + }) + testutil.CleanupCloser(t, client) + + return client, recorder, serverAddress + } + + t.Run("waits for interval before flush", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + const peerCount = 5 + + addressbook, overlays := makeAddressbookWithPeers(t, peerCount) + client, recorder, serverAddress := setupClient(t, addressbook, hiveGossipBufferingInterval) + ctx := context.Background() + + for _, overlay := range overlays { + if err := client.BroadcastPeers(ctx, serverAddress, overlay); err != nil { + t.Fatal(err) + } + } + + assertNoGossipRecords(t, recorder, serverAddress) + + // One coalesce interval plus a small margin for the ticker to flush. + time.Sleep(hiveGossipBufferingInterval + 200*time.Millisecond) + synctest.Wait() + + records, err := recorder.Records(serverAddress, "hive", "2.0.0", "peers") + if err != nil { + t.Fatal(err) + } + if got, want := len(records), 1; got != want { + t.Fatalf("got %d gossip messages, want %d", got, want) + } + + messages, err := readAndAssertPeersMsgs(records[0].In(), 1) + if err != nil { + t.Fatal(err) + } + if got, want := len(messages[0].Peers), peerCount; got != want { + t.Fatalf("coalesced peer count: got %d, want %d", got, want) + } + }) + }) + + t.Run("flushes immediately when buffer is full", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + // Long interval so the coalescer ticker cannot fire before maxBatch flush. + const coalesceInterval = time.Hour + + peerCount := hive.MaxBatchSize + + addressbook, overlays := makeAddressbookWithPeers(t, peerCount) + client, recorder, serverAddress := setupClient(t, addressbook, coalesceInterval) + ctx := context.Background() + + for i, overlay := range overlays { + if err := client.BroadcastPeers(ctx, serverAddress, overlay); err != nil { + t.Fatal(err) + } + if i == peerCount-2 { + assertNoGossipRecords(t, recorder, serverAddress) + } + } + + records, err := recorder.Records(serverAddress, "hive", "2.0.0", "peers") + if err != nil { + t.Fatal(err) + } + if got, want := len(records), 1; got != want { + t.Fatalf("got %d gossip messages after max batch, want %d", got, want) + } + + messages, err := readAndAssertPeersMsgs(records[0].In(), 1) + if err != nil { + t.Fatal(err) + } + if got, want := len(messages[0].Peers), peerCount; got != want { + t.Fatalf("coalesced peer count: got %d, want %d", got, want) + } + }) + }) +} diff --git a/pkg/hive/metrics.go b/pkg/hive/metrics.go index 849c45abec9..6f9bf5e5b5a 100644 --- a/pkg/hive/metrics.go +++ b/pkg/hive/metrics.go @@ -32,6 +32,13 @@ type metrics struct { TimestampRejected *prometheus.CounterVec LegacyRecordSkipped prometheus.Counter + + GossipCoalesceImmediatePeers prometheus.Counter + GossipCoalesceBufferedPeers prometheus.Counter + GossipCoalesceFlushTotal *prometheus.CounterVec + GossipCoalesceFlushPeers prometheus.Counter + GossipCoalesceDropped prometheus.Counter + GossipCoalesceBufferSize prometheus.Gauge } func newMetrics() metrics { @@ -137,6 +144,45 @@ func newMetrics() metrics { }, []string{"reason"}, ), + GossipCoalesceImmediatePeers: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "gossip_coalesce_immediate_peers_total", + Help: "Number of peer gossip entries sent immediately without coalescing.", + }), + GossipCoalesceBufferedPeers: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "gossip_coalesce_buffered_peers_total", + Help: "Number of peer gossip entries enqueued into the coalesce buffer.", + }), + GossipCoalesceFlushTotal: prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "gossip_coalesce_flush_total", + Help: "Number of coalesced gossip flushes dispatched. The reason label is one of: timer, max_batch.", + }, + []string{"reason"}, + ), + GossipCoalesceFlushPeers: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "gossip_coalesce_flush_peers_total", + Help: "Number of peer gossip entries dispatched by coalesced flushes.", + }), + GossipCoalesceDropped: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "gossip_coalesce_dropped_total", + Help: "Number of peer gossip entries dropped during coalesced flush (e.g. outbound rate limiting or send failure).", + }), + GossipCoalesceBufferSize: prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "gossip_coalesce_buffer_size", + Help: "Number of addressees with outbound gossip buffered awaiting coalesced flush.", + }), ChequebookVerification: prometheus.NewCounterVec( prometheus.CounterOpts{ Namespace: m.Namespace, diff --git a/pkg/topology/kademlia/kademlia.go b/pkg/topology/kademlia/kademlia.go index d66b669b786..33a669de804 100644 --- a/pkg/topology/kademlia/kademlia.go +++ b/pkg/topology/kademlia/kademlia.go @@ -63,7 +63,7 @@ const ( defaultShortRetry = 10 * time.Second defaultTimeToRetry = 2 * defaultShortRetry defaultPruneWakeup = 5 * time.Minute - defaultBroadcastBinSize = 2 + defaultBroadcastBinSize = 6 ) var ( @@ -1081,6 +1081,12 @@ func (k *Kad) Announce(ctx context.Context, peer swarm.Address, fullnode bool) e depth := k.neighborhoodDepth() isNeighbor := swarm.Proximity(peer.Bytes(), k.base.Bytes()) >= depth + if isNeighbor { + k.metrics.AnnounceIsNeighborTotal.WithLabelValues("true").Inc() + } else { + k.metrics.AnnounceIsNeighborTotal.WithLabelValues("false").Inc() + } + outer: for bin := range swarm.MaxBins { @@ -1091,11 +1097,15 @@ outer: if bin >= depth && isNeighbor { connectedPeers = k.binPeers(bin, false) // broadcast all neighborhood peers + k.recordAnnounceBinSelection("full", len(connectedPeers), len(connectedPeers)) } else { - connectedPeers, err = randomSubset(k.binPeers(bin, true), k.opt.BroadcastBinSize) + binPeers := k.binPeers(bin, true) + connectedPeers, err = randomSubset(binPeers, k.opt.BroadcastBinSize) if err != nil { + k.metrics.AnnounceErrorsTotal.WithLabelValues("random_subset").Inc() return err } + k.recordAnnounceBinSelection("subset", len(binPeers), len(connectedPeers)) } for _, connectedPeer := range connectedPeers { @@ -1140,8 +1150,11 @@ outer: default: } + k.metrics.AnnouncePeersSentToNewPeer.Observe(float64(len(addrs))) + err := k.discovery.BroadcastPeers(ctx, peer, addrs...) if err != nil { + k.metrics.AnnounceErrorsTotal.WithLabelValues("broadcast_to_new").Inc() k.logger.Error(err, "could not broadcast to peer", "peer_address", peer) _ = k.p2p.Disconnect(peer, "failed broadcasting to peer") } @@ -1149,6 +1162,14 @@ outer: return err } +func (k *Kad) recordAnnounceBinSelection(mode string, available, selected int) { + if available == 0 { + return + } + k.metrics.AnnounceBinPeersAvailable.WithLabelValues(mode).Observe(float64(available)) + k.metrics.AnnounceBinPeersSelected.WithLabelValues(mode).Observe(float64(selected)) +} + // AnnounceTo announces a selected peer to another. func (k *Kad) AnnounceTo(ctx context.Context, addressee, peer swarm.Address, fullnode bool) error { if !fullnode { diff --git a/pkg/topology/kademlia/metrics.go b/pkg/topology/kademlia/metrics.go index 7fc9a53751e..60f17de3449 100644 --- a/pkg/topology/kademlia/metrics.go +++ b/pkg/topology/kademlia/metrics.go @@ -31,6 +31,12 @@ type metrics struct { Blocklist prometheus.Counter ReachabilityStatus *prometheus.GaugeVec PeersReachabilityStatus *prometheus.GaugeVec + + AnnounceIsNeighborTotal *prometheus.CounterVec + AnnounceBinPeersAvailable *prometheus.HistogramVec + AnnounceBinPeersSelected *prometheus.HistogramVec + AnnouncePeersSentToNewPeer prometheus.Histogram + AnnounceErrorsTotal *prometheus.CounterVec } // newMetrics is a convenient constructor for creating new metrics. @@ -164,6 +170,51 @@ func newMetrics() metrics { }, []string{"peers_reachability_status"}, ), + AnnounceIsNeighborTotal: prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "announce_is_neighbor_total", + Help: "Number of peer announce operations. The is_neighbor label is one of: true, false.", + }, + []string{"is_neighbor"}, + ), + AnnounceBinPeersAvailable: prometheus.NewHistogramVec( + prometheus.HistogramOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "announce_bin_peers_available", + Help: "Number of connected peers available in a bin before announce selection. The mode label is one of: full, subset.", + Buckets: []float64{1, 2, 3, 4, 5, 6, 8, 10, 12, 15, 18, 25, 32}, + }, + []string{"mode"}, + ), + AnnounceBinPeersSelected: prometheus.NewHistogramVec( + prometheus.HistogramOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "announce_bin_peers_selected", + Help: "Number of peers selected from a bin during announce. The mode label is one of: full, subset.", + Buckets: []float64{1, 2, 3, 4, 5, 6, 8, 10, 12, 15, 18, 25, 32}, + }, + []string{"mode"}, + ), + AnnouncePeersSentToNewPeer: prometheus.NewHistogram(prometheus.HistogramOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "announce_peers_sent_to_new_peer", + Help: "Number of existing peers sent to a newly connected peer in a single announce.", + Buckets: []float64{1, 2, 5, 10, 15, 20, 30, 40, 50, 75, 100, 150, 200}, + }), + AnnounceErrorsTotal: prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: m.Namespace, + Subsystem: subsystem, + Name: "announce_errors_total", + Help: "Number of announce errors. The reason label is one of: random_subset, broadcast_to_new.", + }, + []string{"reason"}, + ), } }