Skip to content

Commit ea9a08f

Browse files
committed
fix: simplify implementation
1 parent a7a5a15 commit ea9a08f

5 files changed

Lines changed: 96 additions & 125 deletions

File tree

pkg/hive/gossip_buffer.go

Lines changed: 33 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55
package hive
66

77
import (
8-
"math/rand/v2"
8+
"maps"
9+
"slices"
910
"sync"
1011
"time"
1112

@@ -14,7 +15,6 @@ import (
1415

1516
const (
1617
defaultGossipCoalesceInterval = time.Second
17-
defaultGossipCoalesceJitter = 100 * time.Millisecond
1818
// coalesceThreshold: gossips with fewer peers are buffered; larger
1919
// (already-batched) messages are dispatched immediately.
2020
coalesceThreshold = 2
@@ -24,98 +24,78 @@ const (
2424
// flushed as one batched message.
2525
type gossipBuffer struct {
2626
mu sync.Mutex
27-
pending map[string]*pendingGossip // addressee bytestring -> buffered peers
27+
pending map[string]map[string]swarm.Address // addressee key -> peer key -> peer
2828
interval time.Duration
29-
jitter time.Duration
3029
maxBatch int
3130
}
3231

33-
type pendingGossip struct {
32+
type gossipBatch struct {
3433
addressee swarm.Address
35-
peers map[string]swarm.Address // peer bytestring -> address (set semantics)
36-
deadline time.Time
34+
peers []swarm.Address
3735
}
3836

3937
func newGossipBuffer(interval time.Duration, maxBatch int) *gossipBuffer {
4038
if interval == 0 {
4139
interval = defaultGossipCoalesceInterval
4240
}
4341
return &gossipBuffer{
44-
pending: make(map[string]*pendingGossip),
42+
pending: make(map[string]map[string]swarm.Address),
4543
interval: interval,
46-
jitter: defaultGossipCoalesceJitter,
4744
maxBatch: maxBatch,
4845
}
4946
}
5047

51-
// add buffers peers for the addressee. If the buffer reaches maxBatch it is
48+
// stagePeers buffers peers for the addressee. If the buffer reaches maxBatch it is
5249
// removed and returned so the caller can flush it immediately.
53-
func (b *gossipBuffer) add(now time.Time, addressee swarm.Address, peers ...swarm.Address) *pendingGossip {
50+
func (b *gossipBuffer) stagePeers(addressee swarm.Address, peers ...swarm.Address) (flushPeers []swarm.Address, flush bool) {
5451
b.mu.Lock()
5552
defer b.mu.Unlock()
5653

5754
key := addressee.ByteString()
58-
e, ok := b.pending[key]
55+
peerSet, ok := b.pending[key]
5956
if !ok {
60-
var jitter time.Duration
61-
if b.jitter > 0 {
62-
jitter = time.Duration(rand.Int64N(int64(b.jitter)))
63-
}
64-
e = &pendingGossip{
65-
addressee: addressee,
66-
peers: make(map[string]swarm.Address),
67-
deadline: now.Add(b.interval + jitter),
68-
}
69-
b.pending[key] = e
57+
peerSet = make(map[string]swarm.Address)
58+
b.pending[key] = peerSet
7059
}
7160
for _, p := range peers {
72-
e.peers[p.ByteString()] = p
61+
peerSet[p.ByteString()] = p
7362
}
7463

75-
if len(e.peers) >= b.maxBatch {
64+
if len(peerSet) >= b.maxBatch {
7665
delete(b.pending, key)
77-
return e
66+
return slices.Collect(maps.Values(peerSet)), true
7867
}
79-
return nil
68+
return nil, false
8069
}
8170

82-
// takeDue removes and returns all entries whose deadline has passed.
83-
func (b *gossipBuffer) takeDue(now time.Time) []*pendingGossip {
84-
return b.take(func(e *pendingGossip) bool { return !e.deadline.After(now) })
71+
// takeAll removes and returns all buffered entries.
72+
func (b *gossipBuffer) takeAll() []gossipBatch {
73+
b.mu.Lock()
74+
defer b.mu.Unlock()
75+
76+
if len(b.pending) == 0 {
77+
return nil
78+
}
79+
80+
out := make([]gossipBatch, 0, len(b.pending))
81+
for key, peerSet := range b.pending {
82+
out = append(out, gossipBatch{
83+
addressee: swarm.NewAddress([]byte(key)),
84+
peers: slices.Collect(maps.Values(peerSet)),
85+
})
86+
}
87+
b.pending = make(map[string]map[string]swarm.Address)
88+
return out
8589
}
8690

8791
func (b *gossipBuffer) clearAddressee(addressee swarm.Address) {
8892
b.mu.Lock()
8993
defer b.mu.Unlock()
90-
9194
delete(b.pending, addressee.ByteString())
9295
}
9396

9497
func (b *gossipBuffer) pendingAddressees() int {
9598
b.mu.Lock()
9699
defer b.mu.Unlock()
97-
98100
return len(b.pending)
99101
}
100-
101-
func (b *gossipBuffer) take(match func(*pendingGossip) bool) []*pendingGossip {
102-
b.mu.Lock()
103-
defer b.mu.Unlock()
104-
105-
var out []*pendingGossip
106-
for key, e := range b.pending {
107-
if match(e) {
108-
out = append(out, e)
109-
delete(b.pending, key)
110-
}
111-
}
112-
return out
113-
}
114-
115-
func (e *pendingGossip) addresses() []swarm.Address {
116-
out := make([]swarm.Address, 0, len(e.peers))
117-
for _, p := range e.peers {
118-
out = append(out, p)
119-
}
120-
return out
121-
}

pkg/hive/gossip_buffer_test.go

Lines changed: 24 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -11,55 +11,57 @@ import (
1111
"github.com/ethersphere/bee/v2/pkg/swarm"
1212
)
1313

14-
func TestGossipBufferAddAndDue(t *testing.T) {
14+
func TestGossipBufferAddAndTakeAll(t *testing.T) {
1515
t.Parallel()
1616

17-
const interval = 100 * time.Millisecond
18-
19-
b := newGossipBuffer(interval, maxBatchSize)
17+
b := newGossipBuffer(time.Second, maxBatchSize)
2018
addressee := swarm.RandAddress(t)
2119
peer1 := swarm.RandAddress(t)
2220
peer2 := swarm.RandAddress(t)
2321

24-
now := time.Now()
25-
if full := b.add(now, addressee, peer1); full != nil {
26-
t.Fatal("unexpected immediate flush")
22+
if pending := b.takeAll(); len(pending) != 0 {
23+
t.Fatalf("want no pending entries, got %d", len(pending))
2724
}
2825

29-
if due := b.takeDue(now); len(due) != 0 {
30-
t.Fatalf("want no due entries, got %d", len(due))
26+
if _, flush := b.stagePeers(addressee, peer1); flush {
27+
t.Fatal("unexpected immediate flush")
3128
}
3229

33-
if full := b.add(now, addressee, peer2); full != nil {
30+
if _, flush := b.stagePeers(addressee, peer2); flush {
3431
t.Fatal("unexpected immediate flush")
3532
}
3633

37-
afterDeadline := now.Add(interval + defaultGossipCoalesceJitter + time.Millisecond)
38-
due := b.takeDue(afterDeadline)
39-
if len(due) != 1 {
40-
t.Fatalf("want 1 due entry, got %d", len(due))
34+
pending := b.takeAll()
35+
if len(pending) != 1 {
36+
t.Fatalf("want 1 pending entry, got %d", len(pending))
4137
}
42-
if got := len(due[0].addresses()); got != 2 {
38+
if got := len(pending[0].peers); got != 2 {
4339
t.Fatalf("want 2 coalesced peers, got %d", got)
4440
}
41+
if !pending[0].addressee.Equal(addressee) {
42+
t.Fatal("unexpected addressee in pending batch")
43+
}
44+
45+
if pending := b.takeAll(); len(pending) != 0 {
46+
t.Fatalf("want empty buffer after takeAll, got %d pending", len(pending))
47+
}
4548
}
4649

4750
func TestGossipBufferMaxBatchFlush(t *testing.T) {
4851
t.Parallel()
4952

5053
b := newGossipBuffer(time.Second, 2)
5154
addressee := swarm.RandAddress(t)
52-
now := time.Now()
5355

54-
b.add(now, addressee, swarm.RandAddress(t))
55-
full := b.add(now, addressee, swarm.RandAddress(t))
56-
if full == nil {
56+
b.stagePeers(addressee, swarm.RandAddress(t))
57+
flushPeers, flush := b.stagePeers(addressee, swarm.RandAddress(t))
58+
if !flush {
5759
t.Fatal("want immediate flush at maxBatch")
5860
}
59-
if got := len(full.addresses()); got != 2 {
61+
if got := len(flushPeers); got != 2 {
6062
t.Fatalf("want 2 peers in full batch, got %d", got)
6163
}
62-
if due := b.takeDue(now.Add(time.Second)); len(due) != 0 {
63-
t.Fatalf("want empty buffer after maxBatch flush, got %d due", len(due))
64+
if pending := b.takeAll(); len(pending) != 0 {
65+
t.Fatalf("want empty buffer after maxBatch flush, got %d pending", len(pending))
6466
}
6567
}

pkg/hive/hive.go

Lines changed: 27 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -163,8 +163,7 @@ func (s *Service) BroadcastPeers(ctx context.Context, addressee swarm.Address, p
163163
if len(peers) >= coalesceThreshold {
164164
s.metrics.GossipCoalesceImmediatePeers.Add(float64(len(peers)))
165165
s.logger.Debug("gossip immediate send", "addressee", addressee, "peer_count", len(peers))
166-
_, err := s.broadcastNow(ctx, addressee, peers...)
167-
return err
166+
return s.broadcastNow(ctx, addressee, false, peers...)
168167
}
169168

170169
select {
@@ -177,52 +176,49 @@ func (s *Service) BroadcastPeers(ctx context.Context, addressee swarm.Address, p
177176
s.logger.Debug("gossip buffered", "addressee", addressee, "peer_count", len(peers))
178177

179178
// Buffer; if it just filled up, flush it synchronously while still in the call
180-
if full := s.gossipBuf.add(s.now(), addressee, peers...); full != nil {
181-
flushPeers := full.addresses()
179+
if flushPeers, flush := s.gossipBuf.stagePeers(addressee, peers...); flush {
182180
s.recordCoalesceFlush(coalesceFlushReasonMaxBatch, addressee, flushPeers)
183181
s.setCoalesceBufferGauge()
184-
sent, err := s.broadcastNow(ctx, addressee, flushPeers...)
185-
if dropped := len(flushPeers) - sent; dropped > 0 {
186-
s.metrics.GossipCoalesceDropped.Add(float64(dropped))
187-
}
188-
return err
182+
return s.broadcastNow(ctx, addressee, true, flushPeers...)
189183
}
190184
s.setCoalesceBufferGauge()
191185
return nil
192186
}
193187

194188
// broadcastNow performs the synchronous, rate-limited, batched send.
195-
// It returns the number of peers successfully sent.
196-
func (s *Service) broadcastNow(ctx context.Context, addressee swarm.Address, peers ...swarm.Address) (sent int, err error) {
189+
func (s *Service) broadcastNow(ctx context.Context, addressee swarm.Address, coalesced bool, peers ...swarm.Address) error {
197190
maxSize := maxBatchSize
198-
total := len(peers)
199191

200192
for len(peers) > 0 {
201193
if maxSize > len(peers) {
202194
maxSize = len(peers)
203195
}
204196

205-
// If broadcasting limit is exceeded, return early
206197
if !s.outLimiter.Allow(addressee.ByteString(), maxSize) {
207-
return total - len(peers), nil
198+
if coalesced {
199+
s.metrics.GossipCoalesceDropped.Add(float64(len(peers)))
200+
}
201+
return nil
208202
}
209203

210204
select {
211205
case <-ctx.Done():
212-
return total - len(peers), ctx.Err()
206+
return ctx.Err()
213207
case <-s.quit:
214-
return total - len(peers), ErrShutdownInProgress
208+
return ErrShutdownInProgress
215209
default:
216210
}
217211

218212
if err := s.sendPeers(ctx, addressee, peers[:maxSize]); err != nil {
219-
return total - len(peers), err
213+
if coalesced {
214+
s.metrics.GossipCoalesceDropped.Add(float64(len(peers)))
215+
}
216+
return err
220217
}
221218

222219
peers = peers[maxSize:]
223220
}
224-
225-
return total, nil
221+
return nil
226222
}
227223

228224
func (s *Service) SetAddPeersHandler(h func(addr ...swarm.Address)) {
@@ -360,42 +356,32 @@ func (s *Service) disconnect(peer p2p.Peer) error {
360356
}
361357

362358
func (s *Service) startGossipCoalescer() {
363-
tick := s.gossipBuf.interval / 2
364-
if tick <= 0 {
365-
tick = s.gossipBuf.interval
366-
}
367-
368359
s.wg.Go(func() {
369-
ticker := time.NewTicker(tick)
360+
ticker := time.NewTicker(s.gossipBuf.interval)
370361
defer ticker.Stop()
371362
for {
372363
select {
373364
case <-ticker.C:
374-
s.flushGossipEntries(s.gossipBuf.takeDue(s.now()), coalesceFlushReasonTimer)
365+
for _, batch := range s.gossipBuf.takeAll() {
366+
s.flushGossipBatch(batch.addressee, batch.peers, coalesceFlushReasonTimer)
367+
}
368+
s.setCoalesceBufferGauge()
375369
case <-s.quit:
376370
return
377371
}
378372
}
379373
})
380374
}
381375

382-
func (s *Service) flushGossipEntries(entries []*pendingGossip, reason string) {
383-
s.setCoalesceBufferGauge()
384-
385-
for _, e := range entries {
386-
peers := e.addresses()
387-
s.recordCoalesceFlush(reason, e.addressee, peers)
376+
func (s *Service) flushGossipBatch(addressee swarm.Address, peers []swarm.Address, reason string) {
377+
s.recordCoalesceFlush(reason, addressee, peers)
388378

389-
ctx, cancel := context.WithTimeout(context.Background(), messageTimeout)
390-
sent, err := s.broadcastNow(ctx, e.addressee, peers...)
391-
if dropped := len(peers) - sent; dropped > 0 {
392-
s.metrics.GossipCoalesceDropped.Add(float64(dropped))
393-
}
394-
if err != nil {
395-
s.logger.Debug("coalesced gossip flush failed", "addressee", e.addressee, "reason", reason, "batch_size", len(peers), "error", err)
396-
}
397-
cancel()
379+
ctx, cancel := context.WithTimeout(context.Background(), messageTimeout)
380+
err := s.broadcastNow(ctx, addressee, true, peers...)
381+
if err != nil {
382+
s.logger.Debug("coalesced gossip flush failed", "addressee", addressee, "reason", reason, "batch_size", len(peers), "error", err)
398383
}
384+
cancel()
399385
}
400386

401387
func (s *Service) recordCoalesceFlush(reason string, addressee swarm.Address, peers []swarm.Address) {

0 commit comments

Comments
 (0)