Skip to content

Commit ce36b8f

Browse files
authored
fix: speed up node shutdown (#5408)
1 parent 65df151 commit ce36b8f

3 files changed

Lines changed: 21 additions & 5 deletions

File tree

pkg/node/node.go

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,7 @@ import (
9090
const LoggerName = "node"
9191

9292
type Bee struct {
93+
logger log.Logger
9394
p2pService io.Closer
9495
p2pHalter p2p.Halter
9596
ctxCancel context.CancelFunc
@@ -261,6 +262,7 @@ func NewBee(
261262
})
262263

263264
b = &Bee{
265+
logger: logger,
264266
ctxCancel: ctxCancel,
265267
errorLogWriter: sink,
266268
tracerCloser: tracerCloser,
@@ -1351,12 +1353,18 @@ func (b *Bee) Shutdown() error {
13511353
}
13521354
// tryClose is a convenient closure which decrease
13531355
// repetitive io.Closer tryClose procedure.
1354-
tryClose := func(c io.Closer, errMsg string) {
1356+
tryClose := func(c io.Closer, component string) {
13551357
if c == nil {
13561358
return
13571359
}
1360+
1361+
start := time.Now()
1362+
b.logger.Debug("starting shutdown", "component", component)
1363+
defer func() {
1364+
b.logger.Debug("finished shutdown", "component", component, "elapsed", time.Since(start))
1365+
}()
13581366
if err := c.Close(); err != nil {
1359-
mErr = multierror.Append(mErr, fmt.Errorf("%s: %w", errMsg, err))
1367+
mErr = multierror.Append(mErr, fmt.Errorf("%s: %w", component, err))
13601368
}
13611369
}
13621370

@@ -1431,9 +1439,10 @@ func (b *Bee) Shutdown() error {
14311439
tryClose(b.topologyCloser, "topology driver")
14321440
tryClose(b.storageIncetivesCloser, "storage incentives agent")
14331441
tryClose(b.stabilizationDetector, "stabilization detector")
1442+
// close localstore before StateStore to avoid ErrClosed / incomplete flush.
1443+
tryClose(b.localstoreCloser, "localstore")
14341444
tryClose(b.stateStoreCloser, "statestore")
14351445
tryClose(b.stamperStoreCloser, "stamperstore")
1436-
tryClose(b.localstoreCloser, "localstore")
14371446
tryClose(b.resolverCloser, "resolver service")
14381447

14391448
return mErr

pkg/p2p/libp2p/libp2p.go

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1019,7 +1019,14 @@ func (s *Service) Connect(ctx context.Context, addrs []ma.Multiaddr) (address *b
10191019
s.metrics.ConnectBreakerCount.Inc()
10201020
return nil, p2p.NewConnectionBackoffError(err, s.connectionBreaker.ClosedUntil())
10211021
}
1022-
s.logger.Warning("libp2p connect", "peer_id", peerID, "underlay", info.Addrs, "error", err)
1022+
if !errors.Is(err, context.Canceled) {
1023+
select {
1024+
case <-s.halt:
1025+
s.logger.Debug("libp2p connect", "peer_id", peerID, "underlay", info.Addrs, "error", err)
1026+
default:
1027+
s.logger.Warning("libp2p connect", "peer_id", peerID, "underlay", info.Addrs, "error", err)
1028+
}
1029+
}
10231030
connectErr = err
10241031
continue
10251032
}

pkg/pullsync/pullsync.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -193,7 +193,7 @@ func (s *Syncer) handler(streamCtx context.Context, p p2p.Peer, stream p2p.Strea
193193
}
194194

195195
// slow down future requests
196-
waitDur, err := s.limiter.Wait(streamCtx, p.Address.ByteString(), max(1, len(chs)))
196+
waitDur, err := s.limiter.Wait(ctx, p.Address.ByteString(), max(1, len(chs)))
197197
if err != nil {
198198
return fmt.Errorf("rate limiter: %w", err)
199199
}

0 commit comments

Comments
 (0)