diff --git a/pkg/api/pin.go b/pkg/api/pin.go index 9ce1bf053bd..9a381fe065f 100644 --- a/pkg/api/pin.go +++ b/pkg/api/pin.go @@ -12,6 +12,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/file/redundancy" "github.com/ethersphere/bee/v2/pkg/jsonhttp" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storer" "github.com/ethersphere/bee/v2/pkg/swarm" @@ -238,7 +239,9 @@ func (s *Service) pinIntegrityHandler(w http.ResponseWriter, r *http.Request) { out := make(chan storer.PinStat) - go s.pinIntegrity.Check(r.Context(), logger, querie.Ref.String(), out) + safe.Go(logger, "pin-integrity-check", func() { + s.pinIntegrity.Check(r.Context(), logger, querie.Ref.String(), out) + }) flusher, ok := w.(http.Flusher) if !ok { diff --git a/pkg/file/joiner/joiner.go b/pkg/file/joiner/joiner.go index cadd4d61604..c44cf0aa0e2 100644 --- a/pkg/file/joiner/joiner.go +++ b/pkg/file/joiner/joiner.go @@ -19,6 +19,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/file/redundancy" "github.com/ethersphere/bee/v2/pkg/file/redundancy/getter" "github.com/ethersphere/bee/v2/pkg/replicas" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/swarm" "golang.org/x/sync/errgroup" @@ -279,7 +280,7 @@ func (j *joiner) readAtOffset( currentReadSize = min(currentReadSize, subtrieSpan) func(address swarm.Address, b []byte, cur, subTrieSize, off, bufferOffset, bytesToRead, subtrieSpanLimit int64) { - eg.Go(func() error { + eg.Go(safe.RunFunc(nil, "joiner-read-at-offset", func() error { ch, err := g.Get(j.ctx, addr) if err != nil { return err @@ -295,7 +296,7 @@ func (j *joiner) readAtOffset( j.readAtOffset(b, chunkData, cur, subtrieSpan, off, bufferOffset, currentReadSize, bytesRead, subtrieParity, eg) return nil - }) + })) }(addr, b, cur, subtrieSpan, off, bufferOffset, currentReadSize, subtrieSpanLimit) bufferOffset += currentReadSize diff --git a/pkg/hive/hive.go b/pkg/hive/hive.go index a991dbdd58b..e15e59b279d 100644 --- a/pkg/hive/hive.go +++ b/pkg/hive/hive.go @@ -25,6 +25,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/p2p" "github.com/ethersphere/bee/v2/pkg/p2p/protobuf" "github.com/ethersphere/bee/v2/pkg/ratelimit" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/settlement/swap/chequebook" "github.com/ethersphere/bee/v2/pkg/swarm" ma "github.com/multiformats/go-multiaddr" @@ -313,7 +314,9 @@ func (s *Service) startCheckPeersHandler() { return case newPeers := <-s.peersChan: s.wg.Go(func() { - s.checkAndAddPeers(ctx, newPeers) + safe.Run(s.logger, "hive-check-and-add-peers", func() { + s.checkAndAddPeers(ctx, newPeers) + }) }) } } diff --git a/pkg/postage/listener/listener.go b/pkg/postage/listener/listener.go index 349b6a05a14..2658771ca5c 100644 --- a/pkg/postage/listener/listener.go +++ b/pkg/postage/listener/listener.go @@ -19,6 +19,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/log" "github.com/ethersphere/bee/v2/pkg/postage" "github.com/ethersphere/bee/v2/pkg/postage/batchservice" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/transaction" "github.com/ethersphere/bee/v2/pkg/util/syncutil" "github.com/prometheus/client_golang/prometheus" @@ -250,7 +251,7 @@ func (l *listener) Listen(ctx context.Context, from uint64, updater postage.Even lastConfirmedBlock := uint64(0) l.wg.Add(1) - listenf := func() error { + listenf := safe.RunFunc(l.logger, "postage-listener-func", func() error { defer l.wg.Done() for { // 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 totalTimeMetric(l.metrics.PageProcessDuration, start) l.metrics.PagesProcessed.Inc() } - } + }) go func() { err := listenf() diff --git a/pkg/pss/pss.go b/pkg/pss/pss.go index 7202b37e792..ba996f8651e 100644 --- a/pkg/pss/pss.go +++ b/pkg/pss/pss.go @@ -20,6 +20,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/log" "github.com/ethersphere/bee/v2/pkg/postage" "github.com/ethersphere/bee/v2/pkg/pushsync" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/swarm" "github.com/ethersphere/bee/v2/pkg/topology" ) @@ -180,7 +181,9 @@ func (p *pss) TryUnwrap(c swarm.Chunk) { wg.Add(1) go func(hh Handler) { defer wg.Done() - hh(ctx, msg) + safe.Run(p.logger, "pss-handler", func() { + hh(ctx, msg) + }) }(*hh) } go func() { diff --git a/pkg/puller/puller.go b/pkg/puller/puller.go index 9248f82ec7c..018f48c7124 100644 --- a/pkg/puller/puller.go +++ b/pkg/puller/puller.go @@ -21,6 +21,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/puller/intervalstore" "github.com/ethersphere/bee/v2/pkg/pullsync" "github.com/ethersphere/bee/v2/pkg/rate" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storer" "github.com/ethersphere/bee/v2/pkg/swarm" @@ -409,12 +410,16 @@ func (p *Puller) syncPeerBin(parentCtx context.Context, peer *syncPeer, bin uint if cursor > 0 { peer.wg.Add(1) p.wg.Add(1) - go sync(true, peer.address, cursor) + safe.Go(p.logger, "puller-sync-historical", func() { + sync(true, peer.address, cursor) + }) } peer.wg.Add(1) p.wg.Add(1) - go sync(false, peer.address, cursor+1) + safe.Go(p.logger, "puller-sync-live", func() { + sync(false, peer.address, cursor+1) + }) } func (p *Puller) Close() error { diff --git a/pkg/pusher/pusher.go b/pkg/pusher/pusher.go index b756bb088d6..468dec32712 100644 --- a/pkg/pusher/pusher.go +++ b/pkg/pusher/pusher.go @@ -18,6 +18,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/log" "github.com/ethersphere/bee/v2/pkg/postage" "github.com/ethersphere/bee/v2/pkg/pushsync" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/stabilization" storage "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/swarm" @@ -236,7 +237,9 @@ func (s *Service) chunksWorker(startupStabilizer stabilization.Subscriber) { select { case sem <- struct{}{}: wg.Add(1) - go push(op) + safe.Go(s.logger, "pusher-push-worker", func() { + push(op) + }) case <-s.quit: return } diff --git a/pkg/pushsync/pushsync.go b/pkg/pushsync/pushsync.go index f11dfd55f9c..cd11f966583 100644 --- a/pkg/pushsync/pushsync.go +++ b/pkg/pushsync/pushsync.go @@ -22,6 +22,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/postage" "github.com/ethersphere/bee/v2/pkg/pricer" "github.com/ethersphere/bee/v2/pkg/pushsync/pb" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/skippeers" "github.com/ethersphere/bee/v2/pkg/soc" "github.com/ethersphere/bee/v2/pkg/stabilization" @@ -242,7 +243,9 @@ func (ps *PushSync) handler(ctx context.Context, p p2p.Peer, stream p2p.Stream) chunk.WithStamp(stamp) if cac.Valid(chunk) { - go ps.unwrap(chunk) + safe.Go(ps.logger, "pushsync-unwrap-chunk", func() { + ps.unwrap(chunk) + }) } else if chunk, err := soc.FromChunk(chunk); err == nil { addr, err := chunk.Address() if err != nil { @@ -424,7 +427,9 @@ func (ps *PushSync) pushToClosest(ctx context.Context, ch swarm.Chunk, origin bo if inflight == 0 { if ps.fullNode { if cac.Valid(ch) { - go ps.unwrap(ch) + safe.Go(ps.logger, "pushsync-unwrap-ch", func() { + ps.unwrap(ch) + }) } return nil, topology.ErrWantSelf } @@ -477,7 +482,9 @@ func (ps *PushSync) pushToClosest(ctx context.Context, ch swarm.Chunk, origin bo ps.metrics.TotalSendAttempts.Inc() inflight++ - go ps.push(ctx, resultChan, peer, ch, action) + safe.Go(ps.logger, "pushsync-push", func() { + ps.push(ctx, resultChan, peer, ch, action) + }) case result := <-resultChan: inflight-- diff --git a/pkg/replicas/getter.go b/pkg/replicas/getter.go index 7f0c8aa4159..bfa46c2005a 100644 --- a/pkg/replicas/getter.go +++ b/pkg/replicas/getter.go @@ -13,6 +13,7 @@ import ( "time" "github.com/ethersphere/bee/v2/pkg/file/redundancy" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/soc" "github.com/ethersphere/bee/v2/pkg/storage" "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 // concurrently call to retrieve chunk using original CAC address g.wg.Go(func() { - ch, err := g.Getter.Get(ctx, addr) + err := safe.RunFunc(nil, "replicas-get-original", func() error { + ch, err := g.Getter.Get(ctx, addr) + if err != nil { + return err + } + + select { + case resultC <- ch: + case <-ctx.Done(): + } + return nil + })() if err != nil { errc <- err - return - } - - select { - case resultC <- ch: - case <-ctx.Done(): } }) // counters @@ -129,21 +135,25 @@ func (g *getter) Get(ctx context.Context, addr swarm.Address) (ch swarm.Chunk, e } g.wg.Go(func() { - ch, err := g.Getter.Get(ctx, swarm.NewAddress(so.addr)) + err := safe.RunFunc(nil, "replicas-get-replica", func() error { + ch, err := g.Getter.Get(ctx, swarm.NewAddress(so.addr)) + if err != nil { + return err + } + + soc, err := soc.FromChunk(ch) + if err != nil { + return err + } + + select { + case resultC <- soc.WrappedChunk(): + case <-ctx.Done(): + } + return nil + })() if err != nil { errc <- err - return - } - - soc, err := soc.FromChunk(ch) - if err != nil { - errc <- err - return - } - - select { - case resultC <- soc.WrappedChunk(): - case <-ctx.Done(): } }) n++ diff --git a/pkg/replicas/putter.go b/pkg/replicas/putter.go index 7614dee56d0..c5bc399abec 100644 --- a/pkg/replicas/putter.go +++ b/pkg/replicas/putter.go @@ -12,6 +12,7 @@ import ( "sync" "github.com/ethersphere/bee/v2/pkg/file/redundancy" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/soc" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/swarm" @@ -44,10 +45,13 @@ func (p *putter) Put(ctx context.Context, ch swarm.Chunk) (err error) { wg := sync.WaitGroup{} for r := range rr.c { wg.Go(func() { - sch, err := soc.New(r.id, ch).Sign(signer) - if err == nil { - err = p.putter.Put(ctx, sch) - } + err := safe.RunFunc(nil, "replicas-put", func() error { + sch, err := soc.New(r.id, ch).Sign(signer) + if err != nil { + return err + } + return p.putter.Put(ctx, sch) + })() errc <- err }) } diff --git a/pkg/retrieval/retrieval.go b/pkg/retrieval/retrieval.go index f4bc9844a33..b4c2e7ab7c5 100644 --- a/pkg/retrieval/retrieval.go +++ b/pkg/retrieval/retrieval.go @@ -21,6 +21,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/p2p/protobuf" "github.com/ethersphere/bee/v2/pkg/pricer" pb "github.com/ethersphere/bee/v2/pkg/retrieval/pb" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/skippeers" "github.com/ethersphere/bee/v2/pkg/soc" storage "github.com/ethersphere/bee/v2/pkg/storage" @@ -257,13 +258,13 @@ func (s *Service) RetrieveChunk(ctx context.Context, chunkAddr, sourcePeerAddr s inflight++ - go func() { + safe.Go(loggerV1, "retrieval-retrieve-chunk", func() { span, _, ctx := s.tracer.FollowSpanFromContext(spanCtx, "retrieve-chunk", s.logger, trace.WithAttributes( attribute.String("address", chunkAddr.String()), )) defer span.End() s.retrieveChunk(ctx, quit, chunkAddr, peer, resultC, action, span) - }() + }) case res := <-resultC: diff --git a/pkg/safe/safe.go b/pkg/safe/safe.go new file mode 100644 index 00000000000..37ec65920cc --- /dev/null +++ b/pkg/safe/safe.go @@ -0,0 +1,62 @@ +// 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 safe + +import ( + "fmt" + "runtime/debug" + + "github.com/ethersphere/bee/v2/pkg/log" +) + +// Go runs the given function in a new goroutine and recovers from panic in it. +// Panics are logged using the provided logger (if non-nil). +// If the logger is nil, the function is run without panic recovery. +func Go(logger log.Logger, name string, fn func()) { + go func() { + if logger != nil { + defer func() { + if r := recover(); r != nil { + logger.Error(nil, "goroutine panic recovered", "name", name, "panic", fmt.Sprintf("%v", r), "stack", string(debug.Stack())) + } + }() + } + fn() + }() +} + +// Run runs the given function synchronously and recovers from panic in it. +// Panics are logged using the provided logger (if non-nil). +// If the logger is nil, the function is run without panic recovery. +func Run(logger log.Logger, name string, fn func()) { + if logger != nil { + defer func() { + if r := recover(); r != nil { + logger.Error(nil, "panic recovered", "name", name, "panic", fmt.Sprintf("%v", r), "stack", string(debug.Stack())) + } + }() + } + fn() +} + +// RunFunc returns a function wrapped with panic recovery, suitable for use in errgroup.Go. +// Panics are logged using the provided logger (if non-nil), and the returned function returns a non-nil error. +// Do not try to "unwrap" r to an error, since it is a panic and we don't +// want to lose the fact that it was a panic. +func RunFunc(logger log.Logger, name string, fn func() error) func() error { + return func() (err error) { + defer func() { + if r := recover(); r != nil { + if logger != nil { + logger.Error(nil, "errgroup goroutine panic recovered", "name", name, "panic", fmt.Sprintf("%v", r), "stack", string(debug.Stack())) + } + // Do not try to "unwrap" r to an error, since it is a panic and we don't + // want to lose the fact that it was a panic. + err = fmt.Errorf("panic in %s: %v", name, r) + } + }() + return fn() + } +} diff --git a/pkg/safe/safe_test.go b/pkg/safe/safe_test.go new file mode 100644 index 00000000000..1ae81a31c55 --- /dev/null +++ b/pkg/safe/safe_test.go @@ -0,0 +1,223 @@ +// 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 safe_test + +import ( + "errors" + "strings" + "sync" + "testing" + + "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/safe" +) + +func TestGo(t *testing.T) { + logger := &mockLogger{logged: make(chan struct{})} + + safe.Go(logger, "test-panic-goroutine", func() { + panic("intentional panic async") + }) + + <-logger.logged + + msg, keyvals, err := logger.getLogged() + if err != nil { + t.Errorf("expected nil error, got %v", err) + } + if msg != "goroutine panic recovered" { + t.Errorf("expected message 'goroutine panic recovered', got %q", msg) + } + + keyvalsMap := make(map[string]any) + for i := 0; i < len(keyvals); i += 2 { + keyvalsMap[keyvals[i].(string)] = keyvals[i+1] + } + + if keyvalsMap["name"] != "test-panic-goroutine" { + t.Errorf("expected name 'test-panic-goroutine', got %v", keyvalsMap["name"]) + } + if keyvalsMap["panic"] != "intentional panic async" { + t.Errorf("expected panic 'intentional panic async', got %v", keyvalsMap["panic"]) + } + stack, ok := keyvalsMap["stack"].(string) + if !ok || !strings.Contains(stack, "safe_test.go") { + t.Errorf("expected stack trace containing safe_test.go, got %q", stack) + } +} + +func TestRun(t *testing.T) { + logger := &mockLogger{logged: make(chan struct{})} + + safe.Run(logger, "test-panic-sync", func() { + panic("intentional panic sync") + }) + + msg, keyvals, err := logger.getLogged() + if err != nil { + t.Errorf("expected nil error, got %v", err) + } + if msg != "panic recovered" { + t.Errorf("expected message 'panic recovered', got %q", msg) + } + + keyvalsMap := make(map[string]any) + for i := 0; i < len(keyvals); i += 2 { + keyvalsMap[keyvals[i].(string)] = keyvals[i+1] + } + + if keyvalsMap["name"] != "test-panic-sync" { + t.Errorf("expected name 'test-panic-sync', got %v", keyvalsMap["name"]) + } + if keyvalsMap["panic"] != "intentional panic sync" { + t.Errorf("expected panic 'intentional panic sync', got %v", keyvalsMap["panic"]) + } + stack, ok := keyvalsMap["stack"].(string) + if !ok || !strings.Contains(stack, "safe_test.go") { + t.Errorf("expected stack trace containing safe_test.go, got %q", stack) + } +} + +func TestRunFunc(t *testing.T) { + t.Run("no panic", func(t *testing.T) { + logger := &mockLogger{logged: make(chan struct{})} + + wrapped := safe.RunFunc(logger, "test-run-func-ok", func() error { + return nil + }) + + err := wrapped() + if err != nil { + t.Errorf("expected nil error, got %v", err) + } + }) + + t.Run("with panic", func(t *testing.T) { + logger := &mockLogger{logged: make(chan struct{})} + + wrapped := safe.RunFunc(logger, "test-run-func-panic", func() error { + panic("intentional panic runfunc") + }) + + err := wrapped() + if err == nil { + t.Error("expected non-nil error from panic recovery, got nil") + } else if !strings.Contains(err.Error(), "intentional panic runfunc") { + t.Errorf("expected error message to contain panic value, got %q", err.Error()) + } + + <-logger.logged + + msg, keyvals, _ := logger.getLogged() + if msg != "errgroup goroutine panic recovered" { + t.Errorf("expected message 'errgroup goroutine panic recovered', got %q", msg) + } + + keyvalsMap := make(map[string]any) + for i := 0; i < len(keyvals); i += 2 { + keyvalsMap[keyvals[i].(string)] = keyvals[i+1] + } + + if keyvalsMap["name"] != "test-run-func-panic" { + t.Errorf("expected name 'test-run-func-panic', got %v", keyvalsMap["name"]) + } + if keyvalsMap["panic"] != "intentional panic runfunc" { + t.Errorf("expected panic 'intentional panic runfunc', got %v", keyvalsMap["panic"]) + } + stack, ok := keyvalsMap["stack"].(string) + if !ok || !strings.Contains(stack, "safe_test.go") { + t.Errorf("expected stack trace containing safe_test.go, got %q", stack) + } + }) + + t.Run("with panic error no wrapping", func(t *testing.T) { + logger := &mockLogger{logged: make(chan struct{})} + + type customErr struct { + error + } + var sentinelErr = customErr{error: errors.New("sentinel panic")} + + wrapped := safe.RunFunc(logger, "test-run-func-panic-err", func() error { + panic(sentinelErr) + }) + + err := wrapped() + if err == nil { + t.Error("expected non-nil error from panic recovery, got nil") + } else if errors.Is(err, sentinelErr) { + t.Errorf("expected wrapped error to be errors.Is sentinelErr, got %v", err) + } + }) + + t.Run("nil logger", func(t *testing.T) { + wrapped := safe.RunFunc(nil, "test-run-func-nil-logger", func() error { + panic("panic without logger") + }) + + err := wrapped() + if err == nil { + t.Error("expected non-nil error from panic recovery, got nil") + } else if !strings.Contains(err.Error(), "panic without logger") { + t.Errorf("expected error message to contain panic value, got %q", err.Error()) + } + }) +} + +func TestGoNilLogger(t *testing.T) { + done := make(chan struct{}) + safe.Go(nil, "test-nil-logger-go", func() { + close(done) + }) + <-done +} + +func TestRunNilLogger(t *testing.T) { + called := false + safe.Run(nil, "test-nil-logger-run", func() { + called = true + }) + if !called { + t.Error("expected function to be called") + } +} + +func TestRunNilLoggerPanic(t *testing.T) { + defer func() { + if r := recover(); r == nil { + t.Error("expected panic to propagate when logger is nil, but none occurred") + } else if r != "panic-to-propagate" { + t.Errorf("expected panic 'panic-to-propagate', got %v", r) + } + }() + + safe.Run(nil, "test-nil-logger-run-panic", func() { + panic("panic-to-propagate") + }) +} + +type mockLogger struct { + log.Logger + loggedErr error + loggedMsg string + loggedKeyvals []any + mtx sync.Mutex + logged chan struct{} +} + +func (m *mockLogger) Error(err error, msg string, keyvals ...any) { + m.mtx.Lock() + m.loggedErr = err + m.loggedMsg = msg + m.loggedKeyvals = keyvals + m.mtx.Unlock() + close(m.logged) +} + +func (m *mockLogger) getLogged() (string, []any, error) { + m.mtx.Lock() + defer m.mtx.Unlock() + return m.loggedMsg, m.loggedKeyvals, m.loggedErr +} diff --git a/pkg/salud/salud.go b/pkg/salud/salud.go index 4a3bebf76bf..fd46398e617 100644 --- a/pkg/salud/salud.go +++ b/pkg/salud/salud.go @@ -13,6 +13,7 @@ import ( "time" "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/stabilization" "github.com/ethersphere/bee/v2/pkg/status" "github.com/ethersphere/bee/v2/pkg/storer" @@ -144,31 +145,33 @@ func (s *service) salud(mode string, durPercentile float64, connsPercentile floa err := s.topology.EachConnectedPeer(func(addr swarm.Address, bin uint8) (stop bool, jumpToNext bool, err error) { wg.Go(func() { - ctx, cancel := context.WithTimeout(context.Background(), requestTimeout) - defer cancel() - - start := time.Now() - snapshot, err := s.status.PeerSnapshot(ctx, addr) - dur := time.Since(start) - - if err != nil { - s.topology.UpdatePeerHealth(addr, false, dur) - return - } - - if snapshot.BeeMode != mode { - return - } - - mtx.Lock() - totaldur += dur.Seconds() - peer := peer{snapshot, dur, addr, bin, s.reserve.IsWithinStorageRadius(addr)} - peers = append(peers, peer) - if peer.neighbor { - neighborhoodPeers++ - neighborhoodTotalDur += dur.Seconds() - } - mtx.Unlock() + safe.Run(s.logger, "salud-peer-snapshot", func() { + ctx, cancel := context.WithTimeout(context.Background(), requestTimeout) + defer cancel() + + start := time.Now() + snapshot, err := s.status.PeerSnapshot(ctx, addr) + dur := time.Since(start) + + if err != nil { + s.topology.UpdatePeerHealth(addr, false, dur) + return + } + + if snapshot.BeeMode != mode { + return + } + + mtx.Lock() + defer mtx.Unlock() + totaldur += dur.Seconds() + peer := peer{snapshot, dur, addr, bin, s.reserve.IsWithinStorageRadius(addr)} + peers = append(peers, peer) + if peer.neighbor { + neighborhoodPeers++ + neighborhoodTotalDur += dur.Seconds() + } + }) }) return false, false, nil }, topology.Select{}) diff --git a/pkg/storageincentives/agent.go b/pkg/storageincentives/agent.go index 5142e97836e..4a1c0a6e994 100644 --- a/pkg/storageincentives/agent.go +++ b/pkg/storageincentives/agent.go @@ -20,6 +20,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/log" "github.com/ethersphere/bee/v2/pkg/postage" "github.com/ethersphere/bee/v2/pkg/postage/postagecontract" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/settlement/swap/erc20" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storageincentives/redistribution" @@ -222,7 +223,9 @@ func (a *Agent) start(blockTime time.Duration, blocksPerRound, blocksPerPhase ui a.state.SetCurrentEvent(currentPhase, round) a.state.SetFullySynced(a.fullSyncedFunc()) a.state.SetHealthy(a.health.IsHealthy()) - go a.state.purgeStaleRoundData() + safe.Go(a.logger, "storageincentives-purge-stale-round-data", func() { + a.state.purgeStaleRoundData() + }) // check if node is frozen starting from the next block isFrozen, err := a.redistributionStatuser.IsOverlayFrozen(ctx, block+1) diff --git a/pkg/storer/internal/cache/cache.go b/pkg/storer/internal/cache/cache.go index 6e30a56d6c3..3039c3ce181 100644 --- a/pkg/storer/internal/cache/cache.go +++ b/pkg/storer/internal/cache/cache.go @@ -14,6 +14,7 @@ import ( "sync/atomic" "time" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storer/internal/transaction" "github.com/ethersphere/bee/v2/pkg/swarm" @@ -215,7 +216,7 @@ func (c *Cache) RemoveOldest(ctx context.Context, st transaction.Storage, count for _, item := range evictItems { func(item *cacheEntry) { - eg.Go(func() error { + eg.Go(safe.RunFunc(nil, "cache-evict", func() error { c.glock.Lock(item.Address.ByteString()) defer c.glock.Unlock(item.Address.ByteString()) err := st.Run(ctx, func(s transaction.Store) error { @@ -233,7 +234,7 @@ func (c *Cache) RemoveOldest(ctx context.Context, st transaction.Storage, count } c.size.Add(-1) return nil - }) + })) }(item) } diff --git a/pkg/storer/internal/pinning/pinning.go b/pkg/storer/internal/pinning/pinning.go index 01abe264fc8..05b5552308c 100644 --- a/pkg/storer/internal/pinning/pinning.go +++ b/pkg/storer/internal/pinning/pinning.go @@ -13,10 +13,11 @@ import ( "runtime" "github.com/ethersphere/bee/v2/pkg/encryption" - storage "github.com/ethersphere/bee/v2/pkg/storage" + "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storer/internal/transaction" "golang.org/x/sync/errgroup" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage/storageutil" "github.com/ethersphere/bee/v2/pkg/storer/internal" "github.com/ethersphere/bee/v2/pkg/swarm" @@ -273,14 +274,14 @@ func deleteCollectionChunks(ctx context.Context, st transaction.Storage, collect for _, item := range chunksToDelete { func(item *pinChunkItem) { - eg.Go(func() error { + eg.Go(safe.RunFunc(nil, "pinning-delete-chunks", func() error { return st.Run(ctx, func(s transaction.Store) error { return errors.Join( s.IndexStore().Delete(item), s.ChunkStore().Delete(ctx, item.Addr), ) }) - }) + })) }(item) } diff --git a/pkg/storer/internal/reserve/reserve.go b/pkg/storer/internal/reserve/reserve.go index 5566946942e..28c28192e96 100644 --- a/pkg/storer/internal/reserve/reserve.go +++ b/pkg/storer/internal/reserve/reserve.go @@ -17,6 +17,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/log" "github.com/ethersphere/bee/v2/pkg/postage" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storer/internal/chunkstamp" pinstore "github.com/ethersphere/bee/v2/pkg/storer/internal/pinning" @@ -393,7 +394,7 @@ func (r *Reserve) EvictBatchBin( for _, item := range evictedItems { func(item *BatchRadiusItem) { - eg.Go(func() error { + eg.Go(safe.RunFunc(r.logger, "reserve-eviction-remove-chunk", func() error { err := r.st.Run(ctx, func(s transaction.Store) error { return RemoveChunkWithItem(ctx, s, item) }) @@ -402,13 +403,13 @@ func (r *Reserve) EvictBatchBin( } evicted.Add(1) return nil - }) + })) }(item) } for _, item := range pinnedEvictedItems { func(item *BatchRadiusItem) { - eg.Go(func() error { + eg.Go(safe.RunFunc(r.logger, "reserve-eviction-remove-metadata", func() error { err := r.st.Run(ctx, func(s transaction.Store) error { return RemoveChunkMetaData(ctx, s, item) }) @@ -417,7 +418,7 @@ func (r *Reserve) EvictBatchBin( } evicted.Add(1) return nil - }) + })) }(item) } @@ -626,7 +627,7 @@ func (r *Reserve) Reset(ctx context.Context) error { return err } for _, item := range bRitems { - eg.Go(func() error { + eg.Go(safe.RunFunc(r.logger, "reserve-cleanup-delete-chunk", func() error { return r.st.Run(ctx, func(s transaction.Store) error { return errors.Join( s.ChunkStore().Delete(ctx, item.Address), @@ -634,7 +635,7 @@ func (r *Reserve) Reset(ctx context.Context) error { s.IndexStore().Delete(&ChunkBinItem{Bin: item.Bin, BinID: item.BinID}), ) }) - }) + })) } err = eg.Wait() @@ -655,14 +656,14 @@ func (r *Reserve) Reset(ctx context.Context) error { return err } for _, item := range sitems { - eg.Go(func() error { + eg.Go(safe.RunFunc(r.logger, "reserve-cleanup-delete-stamp", func() error { return r.st.Run(ctx, func(s transaction.Store) error { return errors.Join( s.IndexStore().Delete(item), chunkstamp.DeleteWithStamp(s.IndexStore(), reserveScope, item.ChunkAddress, postage.NewStamp(item.BatchID, item.StampIndex, item.StampTimestamp, nil)), ) }) - }) + })) } err = eg.Wait() diff --git a/pkg/storer/internal/upload/uploadstore.go b/pkg/storer/internal/upload/uploadstore.go index 51e99fa16d3..0fe42aa69f9 100644 --- a/pkg/storer/internal/upload/uploadstore.go +++ b/pkg/storer/internal/upload/uploadstore.go @@ -14,6 +14,7 @@ import ( "time" "github.com/ethersphere/bee/v2/pkg/encryption" + "github.com/ethersphere/bee/v2/pkg/safe" storage "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storage/storageutil" "github.com/ethersphere/bee/v2/pkg/storer/internal" @@ -514,7 +515,7 @@ func (u *uploadPutter) Cleanup(st transaction.Storage) error { for _, item := range itemsToDelete { func(item *pushItem) { - eg.Go(func() error { + eg.Go(safe.RunFunc(nil, "uploadstore-delete-chunks", func() error { return st.Run(context.Background(), func(s transaction.Store) error { ui := &uploadItem{Address: item.Address, BatchID: item.BatchID} return errors.Join( @@ -524,7 +525,7 @@ func (u *uploadPutter) Cleanup(st transaction.Storage) error { s.IndexStore().Delete(item), ) }) - }) + })) }(item) } diff --git a/pkg/storer/netstore.go b/pkg/storer/netstore.go index f9e43ae59c4..ea33ddbdd6f 100644 --- a/pkg/storer/netstore.go +++ b/pkg/storer/netstore.go @@ -10,6 +10,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/pusher" "github.com/ethersphere/bee/v2/pkg/pushsync" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/swarm" "github.com/ethersphere/bee/v2/pkg/topology" @@ -28,7 +29,7 @@ func (db *DB) DirectUpload() PutterSession { Putter: putterWithMetrics{ storage.PutterFunc(func(ctx context.Context, ch swarm.Chunk) error { db.directUploadLimiter <- struct{}{} - eg.Go(func() (err error) { + eg.Go(safe.RunFunc(db.logger, "storer-netstore-direct-upload", func() (err error) { defer func() { <-db.directUploadLimiter }() span, logger, ctx := db.tracer.FollowSpanFromContext(ctx, "put-direct-upload", db.logger) @@ -68,7 +69,7 @@ func (db *DB) DirectUpload() PutterSession { } } } - }) + })) return nil }), db.metrics, diff --git a/pkg/storer/sample.go b/pkg/storer/sample.go index a19136b84de..f4b6f8139e3 100644 --- a/pkg/storer/sample.go +++ b/pkg/storer/sample.go @@ -19,6 +19,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/bmt" "github.com/ethersphere/bee/v2/pkg/cac" "github.com/ethersphere/bee/v2/pkg/postage" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/soc" chunk "github.com/ethersphere/bee/v2/pkg/storage/testing" "github.com/ethersphere/bee/v2/pkg/storer/internal/chunkstamp" @@ -93,7 +94,7 @@ func (db *DB) ReserveSample( chunkC := make(chan *reserve.ChunkBinItem, 3*workers) // Phase 1: Iterate chunk addresses - g.Go(func() error { + g.Go(safe.RunFunc(db.logger, "storer-sample-iterate-chunks", func() error { start := time.Now() stats := SampleStats{} defer func() { @@ -115,7 +116,7 @@ func (db *DB) ReserveSample( } }) return err - }) + })) // Phase 2: Get the chunk data and calculate transformed hash sampleItemChan := make(chan SampleItem, 3*workers) @@ -123,7 +124,7 @@ func (db *DB) ReserveSample( db.logger.Debug("reserve sampler workers", "count", workers) for range workers { - g.Go(func() error { + g.Go(safe.RunFunc(db.logger, "storer-sample-worker", func() error { wstat := SampleStats{} hasher := bmt.NewPrefixHasher(anchor) defer func() { @@ -177,7 +178,7 @@ func (db *DB) ReserveSample( } return nil - }) + })) } go func() { diff --git a/pkg/storer/validate.go b/pkg/storer/validate.go index d4d0958a12b..2cda26db586 100644 --- a/pkg/storer/validate.go +++ b/pkg/storer/validate.go @@ -15,6 +15,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/cac" "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/sharky" "github.com/ethersphere/bee/v2/pkg/soc" "github.com/ethersphere/bee/v2/pkg/storage" @@ -153,7 +154,9 @@ func validateWork(logger log.Logger, store storage.Store, readFn func(context.Co wg.Go(func() { buf := make([]byte, swarm.SocMaxChunkSize) for item := range iteratateItemsC { - validChunk(item, buf[:item.Location.Length]) + safe.Run(logger, "reserve-validation-worker", func() { + validChunk(item, buf[:item.Location.Length]) + }) } }) } @@ -330,7 +333,11 @@ func (p *PinIntegrity) Check(ctx context.Context, logger log.Logger, pin string, if ctx.Err() != nil { break } - if !validChunk(item, buf[:item.Location.Length]) { + var isValid bool + safe.Run(logger, "pin-integrity-worker", func() { + isValid = validChunk(item, buf[:item.Location.Length]) + }) + if !isValid { invalid.Add(1) } } diff --git a/pkg/transaction/transaction.go b/pkg/transaction/transaction.go index 68fc74d6f53..43508fd1f17 100644 --- a/pkg/transaction/transaction.go +++ b/pkg/transaction/transaction.go @@ -22,6 +22,7 @@ import ( "github.com/ethereum/go-ethereum/rpc" "github.com/ethersphere/bee/v2/pkg/crypto" "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/safe" "github.com/ethersphere/bee/v2/pkg/sctx" "github.com/ethersphere/bee/v2/pkg/storage" ) @@ -232,20 +233,22 @@ func (t *transactionService) Send(ctx context.Context, request *TxRequest, boost func (t *transactionService) waitForPendingTx(txHash common.Hash) { t.wg.Go(func() { - switch _, err := t.WaitForReceipt(t.ctx, txHash); err { - case nil: - t.logger.Info("pending transaction confirmed", "tx", txHash) - err = t.store.Delete(pendingTransactionKey(txHash)) - if err != nil { - t.logger.Error(err, "unregistering finished pending transaction failed", "tx", txHash) - } - default: - if errors.Is(err, ErrTransactionCancelled) { - t.logger.Warning("pending transaction cancelled", "tx", txHash) - } else { - t.logger.Error(err, "waiting for pending transaction failed", "tx", txHash) + safe.Run(t.logger, "transaction-wait-pending", func() { + switch _, err := t.WaitForReceipt(t.ctx, txHash); err { + case nil: + t.logger.Info("pending transaction confirmed", "tx", txHash) + err = t.store.Delete(pendingTransactionKey(txHash)) + if err != nil { + t.logger.Error(err, "unregistering finished pending transaction failed", "tx", txHash) + } + default: + if errors.Is(err, ErrTransactionCancelled) { + t.logger.Warning("pending transaction cancelled", "tx", txHash) + } else { + t.logger.Error(err, "waiting for pending transaction failed", "tx", txHash) + } } - } + }) }) }