Skip to content

Commit 7dffd3c

Browse files
committed
fix: address review comments
1 parent bea3e7d commit 7dffd3c

3 files changed

Lines changed: 97 additions & 11 deletions

File tree

pkg/hive/export_test.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ var (
1515
MaxBatchSize = maxBatchSize
1616
LimitBurst = limitBurst
1717
CoalesceThreshold = coalesceThreshold
18+
MessageTimeout = messageTimeout
1819
)
1920

2021
func (s *Service) SetTimeFunc(f func() time.Time) {

pkg/hive/hive.go

Lines changed: 12 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,8 @@ type Service struct {
8787
inLimiter *ratelimit.Limiter
8888
outLimiter *ratelimit.Limiter
8989
quit chan struct{}
90+
bgCtx context.Context
91+
bgCancel context.CancelFunc
9092
wg sync.WaitGroup
9193
peersChan chan pb.Peers
9294
sem *semaphore.Weighted
@@ -103,6 +105,7 @@ type Service struct {
103105
}
104106

105107
func New(streamer p2p.Streamer, addressbook addressbook.GetPutSeener, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service {
108+
bgCtx, bgCancel := context.WithCancel(context.Background())
106109
svc := &Service{
107110
streamer: streamer,
108111
logger: logger.WithName(loggerName).Register(),
@@ -112,6 +115,8 @@ func New(streamer p2p.Streamer, addressbook addressbook.GetPutSeener, networkID
112115
inLimiter: ratelimit.New(limitRate, limitBurst),
113116
outLimiter: ratelimit.New(limitRate, limitBurst),
114117
quit: make(chan struct{}),
118+
bgCtx: bgCtx,
119+
bgCancel: bgCancel,
115120
peersChan: make(chan pb.Peers),
116121
sem: semaphore.NewWeighted(int64(swarm.MaxBins)),
117122
bootnode: o.BootnodeMode,
@@ -200,6 +205,7 @@ func (s *Service) broadcastNow(ctx context.Context, addressee swarm.Address, coa
200205
if coalesced {
201206
s.metrics.GossipCoalesceDropped.Add(float64(len(peers)))
202207
}
208+
s.logger.Debug("gossip dropped by outbound rate limiter", "addressee", addressee, "dropped", len(peers), "coalesced", coalesced)
203209
return nil
204210
}
205211

@@ -229,6 +235,7 @@ func (s *Service) SetAddPeersHandler(h func(addr ...swarm.Address)) {
229235

230236
func (s *Service) Close() error {
231237
close(s.quit)
238+
s.bgCancel()
232239

233240
stopped := make(chan struct{})
234241
go func() {
@@ -365,9 +372,9 @@ func (s *Service) startGossipCoalescer() {
365372
select {
366373
case <-ticker.C:
367374
for _, batch := range s.gossipBuf.takeAll() {
368-
go func(batch gossipBatch) {
375+
s.wg.Go(func() {
369376
s.flushGossipBatch(batch.addressee, batch.peers, coalesceFlushReasonTimer)
370-
}(batch)
377+
})
371378
}
372379

373380
case <-s.quit:
@@ -380,7 +387,7 @@ func (s *Service) startGossipCoalescer() {
380387
func (s *Service) flushGossipBatch(addressee swarm.Address, peers []swarm.Address, reason string) {
381388
s.recordCoalesceFlush(reason, addressee, peers)
382389

383-
ctx, cancel := context.WithTimeout(context.Background(), messageTimeout)
390+
ctx, cancel := context.WithTimeout(s.bgCtx, messageTimeout)
384391
defer cancel()
385392

386393
err := s.broadcastNow(ctx, addressee, true, peers...)
@@ -406,21 +413,15 @@ func (s *Service) setCoalesceBufferGauge() {
406413
}
407414

408415
func (s *Service) startCheckPeersHandler() {
409-
ctx, cancel := context.WithCancel(context.Background())
410-
s.wg.Go(func() {
411-
<-s.quit
412-
cancel()
413-
})
414-
415416
s.wg.Go(func() {
416417
for {
417418
select {
418-
case <-ctx.Done():
419+
case <-s.bgCtx.Done():
419420
return
420421
case newPeers := <-s.peersChan:
421422
s.wg.Go(func() {
422423
safe.Run(s.logger, "hive-check-and-add-peers", func() {
423-
s.checkAndAddPeers(ctx, newPeers)
424+
s.checkAndAddPeers(s.bgCtx, newPeers)
424425
})
425426
})
426427
}

pkg/hive/hive_test.go

Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1763,3 +1763,87 @@ func TestHiveGossipBuffering(t *testing.T) {
17631763
})
17641764
})
17651765
}
1766+
1767+
// blockUntilCancelStreamer blocks NewStream until ctx is cancelled, then
1768+
// reports that the blocked call has returned.
1769+
type blockUntilCancelStreamer struct {
1770+
entered chan struct{}
1771+
finished chan struct{}
1772+
}
1773+
1774+
func (s *blockUntilCancelStreamer) NewStream(ctx context.Context, _ swarm.Address, _ p2p.Headers, _, _, _ string) (p2p.Stream, error) {
1775+
select {
1776+
case <-s.entered:
1777+
default:
1778+
close(s.entered)
1779+
}
1780+
<-ctx.Done()
1781+
select {
1782+
case <-s.finished:
1783+
default:
1784+
close(s.finished)
1785+
}
1786+
return nil, ctx.Err()
1787+
}
1788+
1789+
// TestCoalescedFlushCancelsOnClose verifies that an in-flight async flush does
1790+
// not keep running until messageTimeout after Service.Close. The flush must
1791+
// observe bgCtx cancellation and exit promptly.
1792+
func TestCoalescedFlushCancelsOnClose(t *testing.T) {
1793+
t.Parallel()
1794+
1795+
synctest.Test(t, func(t *testing.T) {
1796+
const peerCount = 3
1797+
1798+
addressbook := ab.New(mock.NewStateStore())
1799+
networkID := uint64(1)
1800+
overlays := addTestOverlays(t, addressbook, networkID, peerCount, 5000)
1801+
1802+
entered := make(chan struct{})
1803+
finished := make(chan struct{})
1804+
streamer := &blockUntilCancelStreamer{entered: entered, finished: finished}
1805+
1806+
clientAddress := swarm.RandAddress(t)
1807+
addressee := swarm.RandAddress(t)
1808+
client := hive.New(streamer, addressbook, networkID, clientAddress, log.Noop, hive.Options{
1809+
AllowPrivateCIDRs: true,
1810+
GossipCoalesceInterval: hiveGossipBufferingInterval,
1811+
})
1812+
1813+
ctx := context.Background()
1814+
for _, overlay := range overlays {
1815+
if err := client.BroadcastPeers(ctx, addressee, overlay); err != nil {
1816+
t.Fatal(err)
1817+
}
1818+
}
1819+
1820+
time.Sleep(hiveGossipBufferingInterval + 200*time.Millisecond)
1821+
synctest.Wait()
1822+
1823+
select {
1824+
case <-entered:
1825+
default:
1826+
t.Fatal("expected coalesced flush to block in NewStream before Close")
1827+
}
1828+
1829+
closeAt := time.Now()
1830+
if err := client.Close(); err != nil {
1831+
t.Fatalf("Close during coalesced flush: %v", err)
1832+
}
1833+
1834+
// synctest.Wait advances the fake clock only as far as needed for all
1835+
// goroutines to become durably blocked. If the flush were still waiting
1836+
// on messageTimeout alone, Wait would advance by a full minute.
1837+
synctest.Wait()
1838+
1839+
select {
1840+
case <-finished:
1841+
default:
1842+
t.Fatal("flush goroutine still blocked in NewStream after Close")
1843+
}
1844+
1845+
if elapsed := time.Since(closeAt); elapsed >= hive.MessageTimeout {
1846+
t.Fatalf("flush lasted %v after Close; want exit via bgCtx cancel before messageTimeout (%v)", elapsed, hive.MessageTimeout)
1847+
}
1848+
})
1849+
}

0 commit comments

Comments
 (0)