Skip to content

Commit 4a9546a

Browse files
authored
feat: add safe package for unified goroutine recovery (#5528)
1 parent ff83990 commit 4a9546a

23 files changed

Lines changed: 449 additions & 101 deletions

File tree

pkg/api/pin.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import (
1212

1313
"github.com/ethersphere/bee/v2/pkg/file/redundancy"
1414
"github.com/ethersphere/bee/v2/pkg/jsonhttp"
15+
"github.com/ethersphere/bee/v2/pkg/safe"
1516
"github.com/ethersphere/bee/v2/pkg/storage"
1617
"github.com/ethersphere/bee/v2/pkg/storer"
1718
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -238,7 +239,9 @@ func (s *Service) pinIntegrityHandler(w http.ResponseWriter, r *http.Request) {
238239

239240
out := make(chan storer.PinStat)
240241

241-
go s.pinIntegrity.Check(r.Context(), logger, querie.Ref.String(), out)
242+
safe.Go(logger, "pin-integrity-check", func() {
243+
s.pinIntegrity.Check(r.Context(), logger, querie.Ref.String(), out)
244+
})
242245

243246
flusher, ok := w.(http.Flusher)
244247
if !ok {

pkg/file/joiner/joiner.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import (
1919
"github.com/ethersphere/bee/v2/pkg/file/redundancy"
2020
"github.com/ethersphere/bee/v2/pkg/file/redundancy/getter"
2121
"github.com/ethersphere/bee/v2/pkg/replicas"
22+
"github.com/ethersphere/bee/v2/pkg/safe"
2223
"github.com/ethersphere/bee/v2/pkg/storage"
2324
"github.com/ethersphere/bee/v2/pkg/swarm"
2425
"golang.org/x/sync/errgroup"
@@ -296,7 +297,7 @@ func (j *joiner) readAtOffset(
296297
currentReadSize = min(currentReadSize, subtrieSpan)
297298

298299
func(address swarm.Address, b []byte, cur, subTrieSize, off, bufferOffset, bytesToRead, subtrieSpanLimit int64) {
299-
eg.Go(func() error {
300+
eg.Go(safe.RunFunc(nil, "joiner-read-at-offset", func() error {
300301
ch, err := g.Get(j.ctx, addr)
301302
if err != nil {
302303
return err
@@ -312,7 +313,7 @@ func (j *joiner) readAtOffset(
312313

313314
j.readAtOffset(b, chunkData, cur, subtrieSpan, off, bufferOffset, currentReadSize, bytesRead, subtrieParity, eg)
314315
return nil
315-
})
316+
}))
316317
}(addr, b, cur, subtrieSpan, off, bufferOffset, currentReadSize, subtrieSpanLimit)
317318

318319
bufferOffset += currentReadSize

pkg/hive/hive.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
"github.com/ethersphere/bee/v2/pkg/p2p"
2626
"github.com/ethersphere/bee/v2/pkg/p2p/protobuf"
2727
"github.com/ethersphere/bee/v2/pkg/ratelimit"
28+
"github.com/ethersphere/bee/v2/pkg/safe"
2829
"github.com/ethersphere/bee/v2/pkg/settlement/swap/chequebook"
2930
"github.com/ethersphere/bee/v2/pkg/swarm"
3031
ma "github.com/multiformats/go-multiaddr"
@@ -313,7 +314,9 @@ func (s *Service) startCheckPeersHandler() {
313314
return
314315
case newPeers := <-s.peersChan:
315316
s.wg.Go(func() {
316-
s.checkAndAddPeers(ctx, newPeers)
317+
safe.Run(s.logger, "hive-check-and-add-peers", func() {
318+
s.checkAndAddPeers(ctx, newPeers)
319+
})
317320
})
318321
}
319322
}

pkg/postage/listener/listener.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import (
1919
"github.com/ethersphere/bee/v2/pkg/log"
2020
"github.com/ethersphere/bee/v2/pkg/postage"
2121
"github.com/ethersphere/bee/v2/pkg/postage/batchservice"
22+
"github.com/ethersphere/bee/v2/pkg/safe"
2223
"github.com/ethersphere/bee/v2/pkg/transaction"
2324
"github.com/ethersphere/bee/v2/pkg/util/syncutil"
2425
"github.com/prometheus/client_golang/prometheus"
@@ -250,7 +251,7 @@ func (l *listener) Listen(ctx context.Context, from uint64, updater postage.Even
250251
lastConfirmedBlock := uint64(0)
251252

252253
l.wg.Add(1)
253-
listenf := func() error {
254+
listenf := safe.RunFunc(l.logger, "postage-listener-func", func() error {
254255
defer l.wg.Done()
255256
for {
256257
// if for whatever reason we are stuck for too long we terminate
@@ -350,7 +351,7 @@ func (l *listener) Listen(ctx context.Context, from uint64, updater postage.Even
350351
totalTimeMetric(l.metrics.PageProcessDuration, start)
351352
l.metrics.PagesProcessed.Inc()
352353
}
353-
}
354+
})
354355

355356
go func() {
356357
err := listenf()

pkg/pss/pss.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import (
2020
"github.com/ethersphere/bee/v2/pkg/log"
2121
"github.com/ethersphere/bee/v2/pkg/postage"
2222
"github.com/ethersphere/bee/v2/pkg/pushsync"
23+
"github.com/ethersphere/bee/v2/pkg/safe"
2324
"github.com/ethersphere/bee/v2/pkg/swarm"
2425
"github.com/ethersphere/bee/v2/pkg/topology"
2526
)
@@ -180,7 +181,9 @@ func (p *pss) TryUnwrap(c swarm.Chunk) {
180181
wg.Add(1)
181182
go func(hh Handler) {
182183
defer wg.Done()
183-
hh(ctx, msg)
184+
safe.Run(p.logger, "pss-handler", func() {
185+
hh(ctx, msg)
186+
})
184187
}(*hh)
185188
}
186189
go func() {

pkg/puller/puller.go

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import (
2121
"github.com/ethersphere/bee/v2/pkg/puller/intervalstore"
2222
"github.com/ethersphere/bee/v2/pkg/pullsync"
2323
"github.com/ethersphere/bee/v2/pkg/rate"
24+
"github.com/ethersphere/bee/v2/pkg/safe"
2425
"github.com/ethersphere/bee/v2/pkg/storage"
2526
"github.com/ethersphere/bee/v2/pkg/storer"
2627
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -409,12 +410,16 @@ func (p *Puller) syncPeerBin(parentCtx context.Context, peer *syncPeer, bin uint
409410
if cursor > 0 {
410411
peer.wg.Add(1)
411412
p.wg.Add(1)
412-
go sync(true, peer.address, cursor)
413+
safe.Go(p.logger, "puller-sync-historical", func() {
414+
sync(true, peer.address, cursor)
415+
})
413416
}
414417

415418
peer.wg.Add(1)
416419
p.wg.Add(1)
417-
go sync(false, peer.address, cursor+1)
420+
safe.Go(p.logger, "puller-sync-live", func() {
421+
sync(false, peer.address, cursor+1)
422+
})
418423
}
419424

420425
func (p *Puller) Close() error {

pkg/pusher/pusher.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import (
1818
"github.com/ethersphere/bee/v2/pkg/log"
1919
"github.com/ethersphere/bee/v2/pkg/postage"
2020
"github.com/ethersphere/bee/v2/pkg/pushsync"
21+
"github.com/ethersphere/bee/v2/pkg/safe"
2122
"github.com/ethersphere/bee/v2/pkg/stabilization"
2223
storage "github.com/ethersphere/bee/v2/pkg/storage"
2324
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -236,7 +237,9 @@ func (s *Service) chunksWorker(startupStabilizer stabilization.Subscriber) {
236237
select {
237238
case sem <- struct{}{}:
238239
wg.Add(1)
239-
go push(op)
240+
safe.Go(s.logger, "pusher-push-worker", func() {
241+
push(op)
242+
})
240243
case <-s.quit:
241244
return
242245
}

pkg/pushsync/pushsync.go

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
"github.com/ethersphere/bee/v2/pkg/postage"
2323
"github.com/ethersphere/bee/v2/pkg/pricer"
2424
"github.com/ethersphere/bee/v2/pkg/pushsync/pb"
25+
"github.com/ethersphere/bee/v2/pkg/safe"
2526
"github.com/ethersphere/bee/v2/pkg/skippeers"
2627
"github.com/ethersphere/bee/v2/pkg/soc"
2728
"github.com/ethersphere/bee/v2/pkg/stabilization"
@@ -242,7 +243,9 @@ func (ps *PushSync) handler(ctx context.Context, p p2p.Peer, stream p2p.Stream)
242243
chunk.WithStamp(stamp)
243244

244245
if cac.Valid(chunk) {
245-
go ps.unwrap(chunk)
246+
safe.Go(ps.logger, "pushsync-unwrap-chunk", func() {
247+
ps.unwrap(chunk)
248+
})
246249
} else if chunk, err := soc.FromChunk(chunk); err == nil {
247250
addr, err := chunk.Address()
248251
if err != nil {
@@ -424,7 +427,9 @@ func (ps *PushSync) pushToClosest(ctx context.Context, ch swarm.Chunk, origin bo
424427
if inflight == 0 {
425428
if ps.fullNode {
426429
if cac.Valid(ch) {
427-
go ps.unwrap(ch)
430+
safe.Go(ps.logger, "pushsync-unwrap-ch", func() {
431+
ps.unwrap(ch)
432+
})
428433
}
429434
return nil, topology.ErrWantSelf
430435
}
@@ -477,7 +482,9 @@ func (ps *PushSync) pushToClosest(ctx context.Context, ch swarm.Chunk, origin bo
477482
ps.metrics.TotalSendAttempts.Inc()
478483
inflight++
479484

480-
go ps.push(ctx, resultChan, peer, ch, action)
485+
safe.Go(ps.logger, "pushsync-push", func() {
486+
ps.push(ctx, resultChan, peer, ch, action)
487+
})
481488

482489
case result := <-resultChan:
483490
inflight--

pkg/replicas/getter.go

Lines changed: 30 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import (
1313
"time"
1414

1515
"github.com/ethersphere/bee/v2/pkg/file/redundancy"
16+
"github.com/ethersphere/bee/v2/pkg/safe"
1617
"github.com/ethersphere/bee/v2/pkg/soc"
1718
"github.com/ethersphere/bee/v2/pkg/storage"
1819
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -70,15 +71,20 @@ func (g *getter) Get(ctx context.Context, addr swarm.Address) (ch swarm.Chunk, e
7071

7172
// concurrently call to retrieve chunk using original CAC address
7273
g.wg.Go(func() {
73-
ch, err := g.Getter.Get(ctx, addr)
74+
err := safe.RunFunc(nil, "replicas-get-original", func() error {
75+
ch, err := g.Getter.Get(ctx, addr)
76+
if err != nil {
77+
return err
78+
}
79+
80+
select {
81+
case resultC <- ch:
82+
case <-ctx.Done():
83+
}
84+
return nil
85+
})()
7486
if err != nil {
7587
errc <- err
76-
return
77-
}
78-
79-
select {
80-
case resultC <- ch:
81-
case <-ctx.Done():
8288
}
8389
})
8490
// counters
@@ -129,21 +135,25 @@ func (g *getter) Get(ctx context.Context, addr swarm.Address) (ch swarm.Chunk, e
129135
}
130136

131137
g.wg.Go(func() {
132-
ch, err := g.Getter.Get(ctx, swarm.NewAddress(so.addr))
138+
err := safe.RunFunc(nil, "replicas-get-replica", func() error {
139+
ch, err := g.Getter.Get(ctx, swarm.NewAddress(so.addr))
140+
if err != nil {
141+
return err
142+
}
143+
144+
soc, err := soc.FromChunk(ch)
145+
if err != nil {
146+
return err
147+
}
148+
149+
select {
150+
case resultC <- soc.WrappedChunk():
151+
case <-ctx.Done():
152+
}
153+
return nil
154+
})()
133155
if err != nil {
134156
errc <- err
135-
return
136-
}
137-
138-
soc, err := soc.FromChunk(ch)
139-
if err != nil {
140-
errc <- err
141-
return
142-
}
143-
144-
select {
145-
case resultC <- soc.WrappedChunk():
146-
case <-ctx.Done():
147157
}
148158
})
149159
n++

pkg/replicas/putter.go

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import (
1212
"sync"
1313

1414
"github.com/ethersphere/bee/v2/pkg/file/redundancy"
15+
"github.com/ethersphere/bee/v2/pkg/safe"
1516
"github.com/ethersphere/bee/v2/pkg/soc"
1617
"github.com/ethersphere/bee/v2/pkg/storage"
1718
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -44,10 +45,13 @@ func (p *putter) Put(ctx context.Context, ch swarm.Chunk) (err error) {
4445
wg := sync.WaitGroup{}
4546
for r := range rr.c {
4647
wg.Go(func() {
47-
sch, err := soc.New(r.id, ch).Sign(signer)
48-
if err == nil {
49-
err = p.putter.Put(ctx, sch)
50-
}
48+
err := safe.RunFunc(nil, "replicas-put", func() error {
49+
sch, err := soc.New(r.id, ch).Sign(signer)
50+
if err != nil {
51+
return err
52+
}
53+
return p.putter.Put(ctx, sch)
54+
})()
5155
errc <- err
5256
})
5357
}

0 commit comments

Comments
 (0)