diff --git a/pkg/hive/gossip_buffer.go b/pkg/hive/gossip_buffer.go new file mode 100644 index 00000000000..7d43eda2f76 --- /dev/null +++ b/pkg/hive/gossip_buffer.go @@ -0,0 +1,101 @@ +// 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 = time.Second + // coalesceThreshold: gossips with fewer peers are buffered; larger + // (already-batched) messages are dispatched immediately. + coalesceThreshold = 2 +) + +// 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. +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, ok := b.pending[key] + if !ok { + peerSet = make(map[string]swarm.Address) + b.pending[key] = peerSet + } + for _, p := range peers { + peerSet[p.ByteString()] = p + } + + if len(peerSet) >= b.maxBatch { + delete(b.pending, key) + return slices.Collect(maps.Values(peerSet)), true + } + 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)), + }) + } + b.pending = make(map[string]map[string]swarm.Address) + 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..b95539faed9 --- /dev/null +++ b/pkg/hive/gossip_buffer_test.go @@ -0,0 +1,67 @@ +// 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)) + } +} diff --git a/pkg/hive/hive.go b/pkg/hive/hive.go index a991dbdd58b..cbcd86e23b8 100644 --- a/pkg/hive/hive.go +++ b/pkg/hive/hive.go @@ -57,6 +57,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 @@ -66,6 +71,8 @@ type Options struct { AllowPrivateCIDRs bool ChequebookVerifier chequebook.Verifier ChequebookStorer ChequebookStorer + + GossipCoalesceInterval time.Duration } type Service struct { @@ -90,6 +97,7 @@ type Service struct { // chequebook are dropped. chequebookVerifier chequebook.Verifier chequebookStorer ChequebookStorer + gossipBuf *gossipBuffer } func New(streamer p2p.Streamer, addressbook addressbook.GetPutter, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service { @@ -112,9 +120,12 @@ func New(streamer p2p.Streamer, addressbook addressbook.GetPutter, networkID uin chequebookStorer: o.ChequebookStorer, } + svc.gossipBuf = newGossipBuffer(o.GossipCoalesceInterval, maxBatchSize) + if !o.BootnodeMode { svc.startCheckPeersHandler() } + svc.startGossipCoalescer() return svc } @@ -136,34 +147,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 } @@ -296,9 +350,55 @@ 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() { + s.flushGossipBatch(batch.addressee, batch.peers, coalesceFlushReasonTimer) + } + s.setCoalesceBufferGauge() + 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) + 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) + } + cancel() +} + +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..4c3f90f4596 100644 --- a/pkg/hive/hive_test.go +++ b/pkg/hive/hive_test.go @@ -23,6 +23,7 @@ import ( "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" @@ -37,7 +38,50 @@ import ( 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.GetPutter, + 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 +} + +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 +244,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], @@ -299,11 +335,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 +365,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() @@ -1344,3 +1460,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, 2) + 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.GetPutter, []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.GetPutter, 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,