diff --git a/openapi/SwarmCommon.yaml b/openapi/SwarmCommon.yaml index a1b9718eb3b..502907e52a3 100644 --- a/openapi/SwarmCommon.yaml +++ b/openapi/SwarmCommon.yaml @@ -503,6 +503,25 @@ components: type: boolean swapEnabled: type: boolean + beeStatus: + type: string + description: > + Current node startup or runtime phase. The zero/unknown value means + the phase has not been reported yet. + * `starting` + * `waiting_chain_sync` + * `starting_chequebook` + * `opening_localstore` + * `resetting_reserve` + * `loading_postage_snapshot` + * `syncing_postage` + * `updating_stake_height` + * `warming_up` + * `counting_reserve` + * `evicting_reserve` + * `syncing_reserve` + * `ready` + * `unknown` HealthStatus: type: object @@ -521,6 +540,9 @@ components: type: string default: "0.0.0" description: The default value is set in case the bee binary was not build correctly. + beeStatus: + type: string + description: Current node startup or runtime phase. See Node.beeStatus. PostageBatch: type: object @@ -971,6 +993,9 @@ components: type: integer isWarmingUp: type: boolean + beeStatus: + type: string + description: Current node startup or runtime phase. See Node.beeStatus. StatusPeersResponse: type: object diff --git a/pkg/api/api.go b/pkg/api/api.go index 63a04c390ff..51bc959eefa 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -226,6 +226,12 @@ type Service struct { statusService *status.Service isWarmingUp bool + beeStatus BeeStatus +} + +type BeeStatus interface { + StatusCode() int32 + StatusString() string } func (s *Service) SetP2P(p2p p2p.DebugService) { @@ -396,6 +402,26 @@ func (s *Service) SetIsWarmingUp(v bool) { } } +func (s *Service) SetBeeStatus(status BeeStatus) { + if s != nil { + s.beeStatus = status + } +} + +func (s *Service) BeeStatusCode() int32 { + if s == nil || s.beeStatus == nil { + return 0 + } + return s.beeStatus.StatusCode() +} + +func (s *Service) BeeStatusString() string { + if s == nil || s.beeStatus == nil { + return "unknown" + } + return s.beeStatus.StatusString() +} + // Close hangs up running websockets on shutdown. func (s *Service) Close() error { s.logger.Info("api shutting down") diff --git a/pkg/api/api_test.go b/pkg/api/api_test.go index 23782bcbc6b..01271f4d5ac 100644 --- a/pkg/api/api_test.go +++ b/pkg/api/api_test.go @@ -137,6 +137,7 @@ type testServerOptions struct { ChequebookDisabled bool SwapDisabled bool Erc20ServiceNil bool + BeeStatus api.BeeStatus } func newTestServer(t *testing.T, o testServerOptions) (*http.Client, *websocket.Conn, string, *chanStorer) { @@ -231,6 +232,7 @@ func newTestServer(t *testing.T, o testServerOptions) (*http.Client, *websocket. s.SetSwarmAddress(&o.Overlay) s.SetProbe(o.Probe) + s.SetBeeStatus(o.BeeStatus) tracer := o.Tracer if tracer == nil { diff --git a/pkg/api/bee_status_test.go b/pkg/api/bee_status_test.go new file mode 100644 index 00000000000..3dfc88504b0 --- /dev/null +++ b/pkg/api/bee_status_test.go @@ -0,0 +1,131 @@ +// 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 api_test + +import ( + "crypto/ecdsa" + "net/http" + "testing" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethersphere/bee/v2" + "github.com/ethersphere/bee/v2/pkg/api" + "github.com/ethersphere/bee/v2/pkg/jsonhttp" + "github.com/ethersphere/bee/v2/pkg/jsonhttp/jsonhttptest" + "github.com/ethersphere/bee/v2/pkg/log" +) + +type testBeeStatus struct { + code int32 + name string +} + +func (t testBeeStatus) StatusCode() int32 { return t.code } +func (t testBeeStatus) StatusString() string { return t.name } + +func TestSetBeeStatus(t *testing.T) { + t.Parallel() + + var unset *api.Service + if unset.BeeStatusCode() != 0 || unset.BeeStatusString() != "unknown" { + t.Fatal("nil service must report unknown status") + } + + s := api.New( + ecdsa.PublicKey{}, + ecdsa.PublicKey{}, + common.Address{}, + nil, + log.Noop, + nil, + nil, + api.FullMode, + false, + false, + nil, + nil, + nil, + ) + + if s.BeeStatusString() != "unknown" { + t.Fatalf("got %q, want unknown", s.BeeStatusString()) + } + + s.SetBeeStatus(testBeeStatus{code: 4, name: "opening_localstore"}) + if s.BeeStatusCode() != 4 { + t.Fatalf("got code %d, want 4", s.BeeStatusCode()) + } + if s.BeeStatusString() != "opening_localstore" { + t.Fatalf("got %q, want opening_localstore", s.BeeStatusString()) + } +} + +func TestBeeStatusEndpoints(t *testing.T) { + t.Parallel() + + status := testBeeStatus{code: 4, name: "opening_localstore"} + probe := api.NewProbe() + probe.SetHealthy(api.ProbeStatusOK) + + t.Run("available endpoints", func(t *testing.T) { + t.Parallel() + + testServer, _, _, _ := newTestServer(t, testServerOptions{ + BeeStatus: status, + Probe: probe, + }) + + jsonhttptest.Request(t, testServer, http.MethodGet, "/node", http.StatusOK, + jsonhttptest.WithExpectedJSONResponse(api.NodeResponse{ + BeeMode: api.FullMode.String(), + ChequebookEnabled: true, + SwapEnabled: true, + BeeStatus: "opening_localstore", + }), + ) + jsonhttptest.Request(t, testServer, http.MethodGet, "/health", http.StatusOK, + jsonhttptest.WithExpectedJSONResponse(api.HealthStatusResponse{ + Status: "ok", + Version: bee.Version, + APIVersion: api.Version, + BeeStatus: "opening_localstore", + }), + ) + }) + + t.Run("status and reservestate before full api", func(t *testing.T) { + t.Parallel() + + testServer, _, _, _ := newTestServer(t, testServerOptions{ + BeeStatus: status, + FullAPIDisabled: true, + }) + + jsonhttptest.Request(t, testServer, http.MethodGet, "/status", http.StatusOK, + jsonhttptest.WithExpectedJSONResponse(api.StatusSnapshotResponse{ + Proximity: 256, + BeeMode: api.FullMode.String(), + BeeStatus: "opening_localstore", + }), + ) + jsonhttptest.Request(t, testServer, http.MethodGet, "/reservestate", http.StatusOK) + }) + + t.Run("unavailable message", func(t *testing.T) { + t.Parallel() + + testServer, _, _, _ := newTestServer(t, testServerOptions{ + BeeStatus: status, + FullAPIDisabled: true, + }) + + jsonhttptest.Request(t, testServer, http.MethodGet, "/bytes/00", http.StatusServiceUnavailable, + jsonhttptest.WithExpectedJSONResponse(jsonhttp.StatusResponse{ + Code: http.StatusServiceUnavailable, + Message: "node is not ready: opening_localstore", + }), + ) + }) +} diff --git a/pkg/api/health.go b/pkg/api/health.go index e153b9114f0..5942475501d 100644 --- a/pkg/api/health.go +++ b/pkg/api/health.go @@ -15,6 +15,7 @@ type healthStatusResponse struct { Status string `json:"status"` Version string `json:"version"` APIVersion string `json:"apiVersion"` + BeeStatus string `json:"beeStatus"` } func (s *Service) healthHandler(w http.ResponseWriter, _ *http.Request) { @@ -23,5 +24,6 @@ func (s *Service) healthHandler(w http.ResponseWriter, _ *http.Request) { Status: status.String(), Version: bee.Version, APIVersion: Version, + BeeStatus: s.BeeStatusString(), }) } diff --git a/pkg/api/health_test.go b/pkg/api/health_test.go index 11c79f5366c..de7ac734f2b 100644 --- a/pkg/api/health_test.go +++ b/pkg/api/health_test.go @@ -26,6 +26,7 @@ func TestHealth(t *testing.T) { Status: "nok", Version: bee.Version, APIVersion: api.Version, + BeeStatus: "unknown", })) }) @@ -42,6 +43,7 @@ func TestHealth(t *testing.T) { Status: "nok", Version: bee.Version, APIVersion: api.Version, + BeeStatus: "unknown", })) // When we set health probe to OK it should indicate that node is healthy @@ -50,6 +52,7 @@ func TestHealth(t *testing.T) { Status: "ok", Version: bee.Version, APIVersion: api.Version, + BeeStatus: "unknown", })) // When we set health probe to NOK it should indicate that node is not healthy @@ -58,6 +61,7 @@ func TestHealth(t *testing.T) { Status: "nok", Version: bee.Version, APIVersion: api.Version, + BeeStatus: "unknown", })) }) } diff --git a/pkg/api/node.go b/pkg/api/node.go index b3a26cacb03..ea98e074567 100644 --- a/pkg/api/node.go +++ b/pkg/api/node.go @@ -23,6 +23,7 @@ type nodeResponse struct { BeeMode string `json:"beeMode"` ChequebookEnabled bool `json:"chequebookEnabled"` SwapEnabled bool `json:"swapEnabled"` + BeeStatus string `json:"beeStatus"` } func (b BeeNodeMode) String() string { @@ -43,5 +44,6 @@ func (s *Service) nodeGetHandler(w http.ResponseWriter, _ *http.Request) { BeeMode: s.beeMode.String(), ChequebookEnabled: s.chequebookEnabled, SwapEnabled: s.swapEnabled, + BeeStatus: s.BeeStatusString(), }) } diff --git a/pkg/api/p2p_test.go b/pkg/api/p2p_test.go index b4e60973ee9..775a402b8b7 100644 --- a/pkg/api/p2p_test.go +++ b/pkg/api/p2p_test.go @@ -69,6 +69,7 @@ func TestAddresses(t *testing.T) { BeeMode: api.FullMode.String(), ChequebookEnabled: true, SwapEnabled: true, + BeeStatus: "unknown", }), ) } diff --git a/pkg/api/postage.go b/pkg/api/postage.go index 278eb27b9d2..22555652cbd 100644 --- a/pkg/api/postage.go +++ b/pkg/api/postage.go @@ -431,20 +431,26 @@ type chainStateResponse struct { func (s *Service) reserveStateHandler(w http.ResponseWriter, _ *http.Request) { logger := s.logger.WithName("get_reservestate").Build() - commitment, err := s.batchStore.Commitment() - if err != nil { - logger.Debug("batch store commitment calculation failed", "error", err) - logger.Error(nil, "batch store commitment calculation failed") - jsonhttp.InternalServerError(w, "unable to calculate commitment") - return + resp := reserveStateResponse{} + + if s.batchStore != nil { + commitment, err := s.batchStore.Commitment() + if err != nil { + logger.Debug("batch store commitment calculation failed", "error", err) + logger.Error(nil, "batch store commitment calculation failed") + jsonhttp.InternalServerError(w, "unable to calculate commitment") + return + } + resp.Radius = s.batchStore.Radius() + resp.Commitment = commitment } - jsonhttp.OK(w, reserveStateResponse{ - Radius: s.batchStore.Radius(), - StorageRadius: s.storer.StorageRadius(), - Commitment: commitment, - ReserveCapacityDoubling: s.storer.CapacityDoubling(), - }) + if s.storer != nil { + resp.StorageRadius = s.storer.StorageRadius() + resp.ReserveCapacityDoubling = s.storer.CapacityDoubling() + } + + jsonhttp.OK(w, resp) } // chainStateHandler returns the current chain state. diff --git a/pkg/api/readiness.go b/pkg/api/readiness.go index 35685483892..023b88fe663 100644 --- a/pkg/api/readiness.go +++ b/pkg/api/readiness.go @@ -19,12 +19,14 @@ func (s *Service) readinessHandler(w http.ResponseWriter, _ *http.Request) { Status: "ready", Version: bee.Version, APIVersion: Version, + BeeStatus: s.BeeStatusString(), }) } else { jsonhttp.BadRequest(w, ReadyStatusResponse{ Status: "notReady", Version: bee.Version, APIVersion: Version, + BeeStatus: s.BeeStatusString(), }) } } diff --git a/pkg/api/readiness_test.go b/pkg/api/readiness_test.go index 0454f541c2f..6287b592a7a 100644 --- a/pkg/api/readiness_test.go +++ b/pkg/api/readiness_test.go @@ -42,6 +42,7 @@ func TestReadiness(t *testing.T) { Status: "ready", Version: "-dev", APIVersion: "0.0.0", + BeeStatus: "unknown", })) // When we set readiness probe to NOK it should indicate that API is not ready @@ -51,6 +52,7 @@ func TestReadiness(t *testing.T) { Status: "notReady", Version: "-dev", APIVersion: "0.0.0", + BeeStatus: "unknown", })) }) } diff --git a/pkg/api/router.go b/pkg/api/router.go index 941c63fc89e..de2e50d4c2b 100644 --- a/pkg/api/router.go +++ b/pkg/api/router.go @@ -182,7 +182,7 @@ func (s *Service) mountTechnicalDebug() { func (s *Service) checkRouteAvailability(handler http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if !s.fullAPIEnabled { - jsonhttp.ServiceUnavailable(w, "Node is syncing. This endpoint is unavailable. Try again later.") + jsonhttp.ServiceUnavailable(w, fmt.Sprintf("node is not ready: %s", s.BeeStatusString())) return } handler.ServeHTTP(w, r) @@ -421,6 +421,12 @@ func (s *Service) mountBusinessDebug() { s.router.Handle(rootPath+path, routeHandler) } + // handleAlways registers diagnostic GET routes that must work before EnableFullAPI. + handleAlways := func(path string, handler http.Handler) { + s.router.Handle(path, handler) + s.router.Handle(rootPath+path, handler) + } + if s.transaction != nil { handle("/transactions", jsonhttp.MethodHandler{ "GET": http.HandlerFunc(s.transactionListHandler), @@ -441,7 +447,7 @@ func (s *Service) mountBusinessDebug() { "POST": http.HandlerFunc(s.pingpongHandler), }) - handle("/reservestate", jsonhttp.MethodHandler{ + handleAlways("/reservestate", jsonhttp.MethodHandler{ "GET": http.HandlerFunc(s.reserveStateHandler), }) @@ -679,7 +685,7 @@ func (s *Service) mountBusinessDebug() { })), ) - handle("/status", jsonhttp.MethodHandler{ + handleAlways("/status", jsonhttp.MethodHandler{ "GET": web.ChainHandlers( httpaccess.NewHTTPAccessSuppressLogHandler(), web.FinalHandlerFunc(s.statusGetHandler), diff --git a/pkg/api/router_test.go b/pkg/api/router_test.go index db70b645446..eeec5dd71bc 100644 --- a/pkg/api/router_test.go +++ b/pkg/api/router_test.go @@ -175,7 +175,7 @@ func TestEndpointOptions(t *testing.T) { {"/transactions/{hash}", nil, http.StatusServiceUnavailable}, {"/peers", nil, http.StatusServiceUnavailable}, {"/pingpong/{address}", nil, http.StatusServiceUnavailable}, - {"/reservestate", nil, http.StatusServiceUnavailable}, + {"/reservestate", []string{"GET"}, http.StatusNoContent}, {"/connect/{multi-address:.+}", nil, http.StatusServiceUnavailable}, {"/blocklist", nil, http.StatusServiceUnavailable}, {"/peers/{address}", nil, http.StatusServiceUnavailable}, @@ -209,7 +209,7 @@ func TestEndpointOptions(t *testing.T) { {"/stake/{amount}", nil, http.StatusServiceUnavailable}, {"/stake", nil, http.StatusServiceUnavailable}, {"/redistributionstate", nil, http.StatusServiceUnavailable}, - {"/status", nil, http.StatusServiceUnavailable}, + {"/status", []string{"GET"}, http.StatusNoContent}, {"/status/peers", nil, http.StatusServiceUnavailable}, {"/status/neighborhoods", nil, http.StatusServiceUnavailable}, {"/rchash/{depth}/{anchor1}/{anchor2}", nil, http.StatusServiceUnavailable}, diff --git a/pkg/api/status.go b/pkg/api/status.go index c1eaacb49c0..26d5b1f00ed 100644 --- a/pkg/api/status.go +++ b/pkg/api/status.go @@ -32,6 +32,7 @@ type statusSnapshotResponse struct { LastSyncedBlock uint64 `json:"lastSyncedBlock"` CommittedDepth uint8 `json:"committedDepth"` IsWarmingUp bool `json:"isWarmingUp"` + BeeStatus string `json:"beeStatus,omitempty"` } type statusResponse struct { @@ -69,6 +70,33 @@ func (s *Service) statusAccessHandler(h http.Handler) http.Handler { func (s *Service) statusGetHandler(w http.ResponseWriter, _ *http.Request) { logger := s.logger.WithName("get_status").Build() + overlay := "" + if s.overlay != nil { + overlay = s.overlay.String() + } + + resp := statusSnapshotResponse{ + Proximity: 256, + Overlay: overlay, + BeeMode: s.beeMode.String(), + IsWarmingUp: s.isWarmingUp, + BeeStatus: s.BeeStatusString(), + } + + if s.batchStore != nil { + if commitment, err := s.batchStore.Commitment(); err == nil { + resp.BatchCommitment = commitment + } + if cs := s.batchStore.GetChainState(); cs != nil { + resp.LastSyncedBlock = cs.Block + } + } + + if s.statusService == nil { + jsonhttp.OK(w, resp) + return + } + ss, err := s.statusService.LocalSnapshot() if err != nil { logger.Debug("status snapshot", "error", err) @@ -77,22 +105,19 @@ func (s *Service) statusGetHandler(w http.ResponseWriter, _ *http.Request) { return } - jsonhttp.OK(w, statusSnapshotResponse{ - Proximity: 256, - Overlay: s.overlay.String(), - BeeMode: ss.BeeMode, - ReserveSize: ss.ReserveSize, - ReserveSizeWithinRadius: ss.ReserveSizeWithinRadius, - PullsyncRate: ss.PullsyncRate, - StorageRadius: uint8(ss.StorageRadius), - ConnectedPeers: ss.ConnectedPeers, - NeighborhoodSize: ss.NeighborhoodSize, - BatchCommitment: ss.BatchCommitment, - IsReachable: ss.IsReachable, - LastSyncedBlock: ss.LastSyncedBlock, - CommittedDepth: uint8(ss.CommittedDepth), - IsWarmingUp: s.isWarmingUp, - }) + resp.BeeMode = ss.BeeMode + resp.ReserveSize = ss.ReserveSize + resp.ReserveSizeWithinRadius = ss.ReserveSizeWithinRadius + resp.PullsyncRate = ss.PullsyncRate + resp.StorageRadius = uint8(ss.StorageRadius) + resp.ConnectedPeers = ss.ConnectedPeers + resp.NeighborhoodSize = ss.NeighborhoodSize + resp.BatchCommitment = ss.BatchCommitment + resp.IsReachable = ss.IsReachable + resp.LastSyncedBlock = ss.LastSyncedBlock + resp.CommittedDepth = uint8(ss.CommittedDepth) + + jsonhttp.OK(w, resp) } // statusGetPeersHandler returns the status of currently connected peers. diff --git a/pkg/api/status_test.go b/pkg/api/status_test.go index 9da5797caf8..b08e0dc2e3d 100644 --- a/pkg/api/status_test.go +++ b/pkg/api/status_test.go @@ -41,6 +41,7 @@ func TestGetStatus(t *testing.T) { IsReachable: true, LastSyncedBlock: 6092500, CommittedDepth: 1, + BeeStatus: "unknown", } ssMock := &statusSnapshotMock{ @@ -75,6 +76,23 @@ func TestGetStatus(t *testing.T) { ) }) + t.Run("available before full api", func(t *testing.T) { + t.Parallel() + + client, _, _, _ := newTestServer(t, testServerOptions{ + BeeMode: api.FullMode, + FullAPIDisabled: true, + BeeStatus: testBeeStatus{code: 4, name: "opening_localstore"}, + }) + + jsonhttptest.Request(t, client, http.MethodGet, url, http.StatusOK, + jsonhttptest.WithExpectedJSONResponse(api.StatusSnapshotResponse{ + Proximity: 256, + BeeMode: api.FullMode.String(), + BeeStatus: "opening_localstore", + }), + ) + }) } // TestGetStatusPeersIncludesBootnodes is a regression test for diff --git a/pkg/api/warmup_test.go b/pkg/api/warmup_test.go index 91f9d79eee0..3773004bd09 100644 --- a/pkg/api/warmup_test.go +++ b/pkg/api/warmup_test.go @@ -22,4 +22,5 @@ func TestSetIsWarmingUpNilReceiver(t *testing.T) { // Must not panic. s.SetIsWarmingUp(false) s.SetIsWarmingUp(true) + s.SetBeeStatus(nil) } diff --git a/pkg/node/export_test.go b/pkg/node/export_test.go index b67e259a588..22225e9c828 100644 --- a/pkg/node/export_test.go +++ b/pkg/node/export_test.go @@ -4,13 +4,35 @@ package node -import "io" +import ( + "io" + + "github.com/ethersphere/bee/v2/pkg/log" +) var ( ValidatePublicAddress = validatePublicAddress UseEmbeddedSnapshot = useEmbeddedSnapshot ) +// NewTestBeeWithStatus returns a Bee with a status store for phase tests. +func NewTestBeeWithStatus() *Bee { + return &Bee{ + logger: log.Noop, + status: NewStatusStore(), + } +} + +func (b *Bee) ApplyReservePhase(phase string, syncRate func() float64, stabilized func() bool) { + b.applyReservePhase(phase, syncRate, stabilized) +} + +func (b *Bee) SetReadyFromWarmup() { b.setReadyFromWarmup() } + +func (b *Bee) CurrentStatus() Status { + return b.status.Status() +} + // NewTestBeeWithClosers builds a Bee with only the push-sync and retrieval // closer fields set, for exercising the Shutdown closer registration. func NewTestBeeWithClosers(pushSync, retrieval io.Closer) *Bee { diff --git a/pkg/node/node.go b/pkg/node/node.go index 9ef7704b9dd..37b1afae0d5 100644 --- a/pkg/node/node.go +++ b/pkg/node/node.go @@ -129,6 +129,7 @@ type Bee struct { syncingStopped *syncutil.Signaler accesscontrolCloser io.Closer ethClientCloser func() + status *StatusStore } type Options struct { @@ -290,7 +291,9 @@ func NewBee( ctxCancel: ctxCancel, errorLogWriter: sink, syncingStopped: syncutil.NewSignaler(), + status: NewStatusStore(), } + b.setStatus(StatusStarting) defer func(b *Bee) { if err != nil { @@ -539,6 +542,7 @@ func NewBee( apiService.SetProbe(probe) apiService.SetIsWarmingUp(true) apiService.SetSwarmAddress(&swarmAddress) + apiService.SetBeeStatus(b.status) apiServer := &http.Server{ IdleTimeout: 30 * time.Second, @@ -568,6 +572,7 @@ func NewBee( } if !isSynced { logger.Info("waiting to sync with the blockchain backend") + b.setStatus(StatusWaitingChainSync) err := transaction.WaitSynced(ctx, logger, chainBackend, maxDelay) if err != nil { @@ -598,6 +603,7 @@ func NewBee( } if o.ChequebookEnable { + b.setStatus(StatusStartingChequebook) chequebookService, err = InitChequebookService( ctx, logger, @@ -877,6 +883,7 @@ func NewBee( lo.ReserveCapacityDoubling = o.ReserveCapacityDoubling } + b.setStatus(StatusOpeningLocalstore) localStore, err := storer.New(ctx, path, lo) if err != nil { return nil, fmt.Errorf("localstore: %w", err) @@ -886,6 +893,7 @@ func NewBee( if resetReserve { logger.Warning("resetting the reserve") + b.setStatus(StatusResettingReserve) err := localStore.ResetReserve(ctx) if err != nil { return nil, fmt.Errorf("reset reserve: %w", err) @@ -912,6 +920,7 @@ func NewBee( var batchSnapshot *batchservice.Snapshot if useEmbeddedSnapshot(o.SkipPostageSnapshot, batchStoreExists, o.Resync, networkID, beeNodeMode) { + b.setStatus(StatusLoadingPostageSnapshot) batchSnapshot, err = snapshot.New(ctx, logger, archive.Getter{}, b.syncingStopped, postageStampContractAddress, postageStampContractABI, o.BlockTime, postageSyncingStallingTimeout, postageSyncingBackoffTimeout, postageSyncStart) if err != nil { // A corrupt snapshot is not fatal: rebuild from the chain instead. @@ -984,6 +993,7 @@ func NewBee( } if o.FullNodeMode { + b.setStatus(StatusSyncingPostage) err = batchSvc.Start(ctx, postageSyncStart) syncStatus.Store(true) if err != nil { @@ -1215,6 +1225,7 @@ func NewBee( defer unsubscribe() <-sub logger.Info("node warmup stabilization complete, updating API status") + b.setReadyFromWarmup() // apiService is nil when the API is disabled (empty --api-addr). if apiService != nil { apiService.SetIsWarmingUp(false) @@ -1256,6 +1267,7 @@ func NewBee( logger.Warning("staked amount does not sufficiently cover the additional reserve capacity. On-chain height update will be skipped. Node will start, but storage incentives may not function for this capacity.", "missing_stake", new(big.Int).Sub(minStake, stake)) } else { // make sure that the staking contract has the up to date height + b.setStatus(StatusUpdatingStakeHeight) tx, updated, err := stakingContract.UpdateHeight(ctx) if err != nil { return nil, fmt.Errorf("update height in staking contract: %w", err) @@ -1276,6 +1288,10 @@ func NewBee( pullerService = puller.New(swarmAddress, stateStore, kad, localStore, pullSyncProtocol, p2ps, logger, puller.Options{}) b.pullerCloser = pullerService + localStore.SetReservePhaseFunc(func(phase string) { + b.applyReservePhase(phase, pullerService.SyncRate, detector.IsStabilized) + }) + // we pass an empty channel since startup synchronization is not needed for production code, only tests. localStore.StartReserveWorker(ctx, pullerService, waitNetworkRFunc, nil) nodeStatus.SetSync(pullerService) @@ -1303,6 +1319,7 @@ func NewBee( fullSyncTime := pullSyncStartTime.Sub(t) logger.Info("full sync done", "duration", fullSyncTime) nodeMetrics.FullSyncDuration.Observe(fullSyncTime.Minutes()) + b.setReadyAfterSync() syncCheckTicker.Stop() return } @@ -1441,6 +1458,12 @@ func NewBee( WsPingPeriod: 60 * time.Second, }, extraOpts, chainID, erc20Service) + if detector.IsStabilized() { + b.setReadyFromWarmup() + } else { + b.setStatusIfIdle(StatusWarmingUp) + } + apiService.EnableFullAPI() apiService.SetRedistributionAgent(agent) @@ -1448,6 +1471,10 @@ func NewBee( // api metrics are constructed on api.Service.Configure apiService.MustRegisterMetrics(apiService.Metrics()...) statusMetricsRegistry.MustRegister(apiService.StatusMetrics()...) + } else if detector.IsStabilized() { + b.setReadyFromWarmup() + } else { + b.setStatusIfIdle(StatusWarmingUp) } if err := kad.Start(ctx); err != nil { diff --git a/pkg/node/status.go b/pkg/node/status.go new file mode 100644 index 00000000000..b7298347407 --- /dev/null +++ b/pkg/node/status.go @@ -0,0 +1,191 @@ +// 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 node + +import ( + "sync/atomic" + + "github.com/ethersphere/bee/v2/pkg/api" + "github.com/ethersphere/bee/v2/pkg/storer" +) + +// Status is a numeric node startup/runtime phase. +// The zero value is unset and String reports "unknown". +type Status int32 + +const ( + StatusStarting Status = iota + 1 + StatusWaitingChainSync + StatusStartingChequebook + StatusOpeningLocalstore + StatusResettingReserve + StatusLoadingPostageSnapshot + StatusSyncingPostage + StatusUpdatingStakeHeight + StatusWarmingUp + StatusCountingReserve + StatusEvictingReserve + StatusSyncingReserve + StatusReady +) + +// String returns the operator-facing name of the status. +func (s Status) String() string { + switch s { + case StatusStarting: + return "starting" + case StatusWaitingChainSync: + return "waiting_chain_sync" + case StatusStartingChequebook: + return "starting_chequebook" + case StatusOpeningLocalstore: + return "opening_localstore" + case StatusResettingReserve: + return "resetting_reserve" + case StatusLoadingPostageSnapshot: + return "loading_postage_snapshot" + case StatusSyncingPostage: + return "syncing_postage" + case StatusUpdatingStakeHeight: + return "updating_stake_height" + case StatusWarmingUp: + return "warming_up" + case StatusCountingReserve: + return "counting_reserve" + case StatusEvictingReserve: + return "evicting_reserve" + case StatusSyncingReserve: + return "syncing_reserve" + case StatusReady: + return "ready" + default: + return "unknown" + } +} + +// StatusProvider is the read side of StatusStore. +// The HTTP API depends on api.BeeStatus (same method set, no node import). +type StatusProvider interface { + Status() Status + StatusCode() int32 + StatusString() string +} + +// StatusStore holds the current Status in an atomic integer. +type StatusStore struct { + current atomic.Int32 +} + +var ( + _ StatusProvider = (*StatusStore)(nil) + _ api.BeeStatus = (*StatusStore)(nil) +) + +// NewStatusStore returns a store with an unset (unknown) status. +func NewStatusStore() *StatusStore { + return &StatusStore{} +} + +// Set stores the given status. +func (s *StatusStore) Set(st Status) { + if s == nil { + return + } + s.current.Store(int32(st)) +} + +// Status returns the current status. A nil store is unknown. +func (s *StatusStore) Status() Status { + if s == nil { + return 0 + } + return Status(s.current.Load()) +} + +// StatusCode returns the current status as a number. +func (s *StatusStore) StatusCode() int32 { + return int32(s.Status()) +} + +// StatusString returns the current status name. +func (s *StatusStore) StatusString() string { + return s.Status().String() +} + +func (b *Bee) setStatus(st Status) { + if b == nil || b.status == nil { + return + } + b.status.Set(st) + if b.logger != nil { + b.logger.Info("node status", "status", st.String()) + } +} + +func isReserveBlocking(s Status) bool { + return s == StatusCountingReserve || s == StatusEvictingReserve +} + +// setStatusIfIdle writes st unless the reserve worker is in a blocking phase +// or pullsync has already started. +func (b *Bee) setStatusIfIdle(st Status) { + if b == nil || b.status == nil { + return + } + cur := b.status.Status() + if isReserveBlocking(cur) || cur == StatusSyncingReserve { + return + } + b.setStatus(st) +} + +// setReadyFromWarmup marks the node ready after warmup, but does not interrupt +// reserve counting, eviction, or historical pullsync. +func (b *Bee) setReadyFromWarmup() { + if b == nil || b.status == nil { + return + } + cur := b.status.Status() + if isReserveBlocking(cur) || cur == StatusSyncingReserve { + return + } + b.setStatus(StatusReady) +} + +// setReadyAfterSync marks the node ready when pullsync has caught up. +func (b *Bee) setReadyAfterSync() { + if isReserveBlocking(b.status.Status()) { + return + } + b.setStatus(StatusReady) +} + +func (b *Bee) applyReservePhase(phase string, syncRate func() float64, stabilized func() bool) { + switch phase { + case storer.ReservePhaseCounting: + b.setStatus(StatusCountingReserve) + case storer.ReservePhaseEvicting: + b.setStatus(StatusEvictingReserve) + case storer.ReservePhaseSyncing: + if isReserveBlocking(b.status.Status()) { + return + } + b.setStatus(StatusSyncingReserve) + case storer.ReservePhaseIdle: + b.setRuntimeIdle(syncRate, stabilized) + } +} + +func (b *Bee) setRuntimeIdle(syncRate func() float64, stabilized func() bool) { + if stabilized != nil && !stabilized() { + b.setStatus(StatusWarmingUp) + return + } + if syncRate != nil && syncRate() > 0 { + b.setStatus(StatusSyncingReserve) + return + } + b.setStatus(StatusReady) +} diff --git a/pkg/node/status_test.go b/pkg/node/status_test.go new file mode 100644 index 00000000000..53cac6f12e8 --- /dev/null +++ b/pkg/node/status_test.go @@ -0,0 +1,161 @@ +// 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 node_test + +import ( + "sync" + "testing" + + "github.com/ethersphere/bee/v2/pkg/api" + "github.com/ethersphere/bee/v2/pkg/node" +) + +func TestStatusString(t *testing.T) { + t.Parallel() + + tests := []struct { + status node.Status + want string + }{ + {0, "unknown"}, + {node.StatusStarting, "starting"}, + {node.StatusWaitingChainSync, "waiting_chain_sync"}, + {node.StatusStartingChequebook, "starting_chequebook"}, + {node.StatusOpeningLocalstore, "opening_localstore"}, + {node.StatusResettingReserve, "resetting_reserve"}, + {node.StatusLoadingPostageSnapshot, "loading_postage_snapshot"}, + {node.StatusSyncingPostage, "syncing_postage"}, + {node.StatusUpdatingStakeHeight, "updating_stake_height"}, + {node.StatusWarmingUp, "warming_up"}, + {node.StatusCountingReserve, "counting_reserve"}, + {node.StatusEvictingReserve, "evicting_reserve"}, + {node.StatusSyncingReserve, "syncing_reserve"}, + {node.StatusReady, "ready"}, + {node.Status(99), "unknown"}, + } + + for _, tc := range tests { + if got := tc.status.String(); got != tc.want { + t.Fatalf("Status(%d).String() = %q, want %q", tc.status, got, tc.want) + } + } +} + +func TestStatusStore(t *testing.T) { + t.Parallel() + + var unset *node.StatusStore + if unset.Status() != 0 || unset.StatusCode() != 0 || unset.StatusString() != "unknown" { + t.Fatal("nil store must report unknown") + } + unset.Set(node.StatusReady) // must not panic + + store := node.NewStatusStore() + if store.StatusString() != "unknown" { + t.Fatalf("new store = %q, want unknown", store.StatusString()) + } + + store.Set(node.StatusOpeningLocalstore) + if store.Status() != node.StatusOpeningLocalstore { + t.Fatalf("got status %d, want %d", store.Status(), node.StatusOpeningLocalstore) + } + if store.StatusCode() != int32(node.StatusOpeningLocalstore) { + t.Fatalf("got code %d, want %d", store.StatusCode(), node.StatusOpeningLocalstore) + } + if store.StatusString() != "opening_localstore" { + t.Fatalf("got name %q, want opening_localstore", store.StatusString()) + } + + var _ api.BeeStatus = store +} + +func TestStatusStoreConcurrent(t *testing.T) { + t.Parallel() + + store := node.NewStatusStore() + var wg sync.WaitGroup + for range 8 { + wg.Add(1) + go func() { + defer wg.Done() + for range 1000 { + store.Set(node.StatusReady) + _ = store.Status() + _ = store.StatusCode() + _ = store.StatusString() + } + }() + } + wg.Wait() +} + +func TestApplyReservePhase(t *testing.T) { + t.Parallel() + + syncRate := func() float64 { return 0 } + notStabilized := func() bool { return false } + stabilized := func() bool { return true } + + t.Run("counting and evicting", func(t *testing.T) { + t.Parallel() + b := node.NewTestBeeWithStatus() + b.ApplyReservePhase("counting_reserve", syncRate, notStabilized) + if b.CurrentStatus() != node.StatusCountingReserve { + t.Fatalf("got %s, want counting_reserve", b.CurrentStatus()) + } + b.ApplyReservePhase("evicting_reserve", syncRate, notStabilized) + if b.CurrentStatus() != node.StatusEvictingReserve { + t.Fatalf("got %s, want evicting_reserve", b.CurrentStatus()) + } + }) + + t.Run("idle while warming up", func(t *testing.T) { + t.Parallel() + b := node.NewTestBeeWithStatus() + b.ApplyReservePhase("counting_reserve", syncRate, notStabilized) + b.ApplyReservePhase("idle", syncRate, notStabilized) + if b.CurrentStatus() != node.StatusWarmingUp { + t.Fatalf("got %s, want warming_up", b.CurrentStatus()) + } + }) + + t.Run("idle after warmup with pullsync", func(t *testing.T) { + t.Parallel() + b := node.NewTestBeeWithStatus() + b.ApplyReservePhase("idle", func() float64 { return 1.5 }, stabilized) + if b.CurrentStatus() != node.StatusSyncingReserve { + t.Fatalf("got %s, want syncing_reserve", b.CurrentStatus()) + } + }) + + t.Run("idle after warmup ready", func(t *testing.T) { + t.Parallel() + b := node.NewTestBeeWithStatus() + b.ApplyReservePhase("idle", syncRate, stabilized) + if b.CurrentStatus() != node.StatusReady { + t.Fatalf("got %s, want ready", b.CurrentStatus()) + } + }) + + t.Run("syncing does not override evicting", func(t *testing.T) { + t.Parallel() + b := node.NewTestBeeWithStatus() + b.ApplyReservePhase("evicting_reserve", syncRate, stabilized) + b.ApplyReservePhase("syncing_reserve", syncRate, stabilized) + if b.CurrentStatus() != node.StatusEvictingReserve { + t.Fatalf("got %s, want evicting_reserve", b.CurrentStatus()) + } + }) + + t.Run("warmup ready skipped while evicting", func(t *testing.T) { + t.Parallel() + b := node.NewTestBeeWithStatus() + b.ApplyReservePhase("evicting_reserve", syncRate, stabilized) + b.SetReadyFromWarmup() + if b.CurrentStatus() != node.StatusEvictingReserve { + t.Fatalf("got %s, want evicting_reserve", b.CurrentStatus()) + } + }) +} diff --git a/pkg/storer/reserve.go b/pkg/storer/reserve.go index b6799ea1cee..da0ebf9403d 100644 --- a/pkg/storer/reserve.go +++ b/pkg/storer/reserve.go @@ -29,6 +29,12 @@ const ( reserveUnreserved = "reserveUnreserved" batchExpiry = "batchExpiry" batchExpiryDone = "batchExpiryDone" + + // ReservePhase* values are reported through ReservePhaseFunc / SetReservePhaseFunc. + ReservePhaseCounting = "counting_reserve" + ReservePhaseEvicting = "evicting_reserve" + ReservePhaseSyncing = "syncing_reserve" + ReservePhaseIdle = "idle" ) var ( @@ -82,6 +88,7 @@ func (db *DB) startReserveWorkers( } // syncing can now begin now that the reserver worker is running + db.reportPhase(ReservePhaseSyncing) db.syncer.Start(ctx) } @@ -130,7 +137,9 @@ func (db *DB) reserveWorker(ctx context.Context, ready chan<- struct{}) { thresholdTicker := time.NewTicker(db.reserveOptions.wakeupDuration) defer thresholdTicker.Stop() + db.reportPhase(ReservePhaseCounting) _, _ = db.countWithinRadius(ctx) + db.reportPhase(ReservePhaseIdle) if !db.reserve.IsWithinCapacity() { db.events.Trigger(reserveOverCapacity) @@ -347,6 +356,9 @@ func (db *DB) unreserve(ctx context.Context) (err error) { return nil } + db.reportPhase(ReservePhaseEvicting) + defer db.reportPhase(ReservePhaseIdle) + db.logger.Info("unreserve start", "target", target, "radius", radius) batchExpiry, unsub := db.events.Subscribe(batchExpiry) diff --git a/pkg/storer/storer.go b/pkg/storer/storer.go index 43563683a9a..fabc5bb0d22 100644 --- a/pkg/storer/storer.go +++ b/pkg/storer/storer.go @@ -397,6 +397,9 @@ type Options struct { CacheMinEvictCount uint64 MinimumStorageRadius uint + + // ReservePhaseFunc reports reserve worker activity to the node status store. + ReservePhaseFunc func(phase string) } func defaultOptions() *Options { @@ -459,6 +462,7 @@ type reserveOpts struct { cacheMinEvictCount uint64 minimumRadius uint8 capacityDoubling int + phaseFunc func(string) } // New returns a newly constructed DB object which implements all the above @@ -553,6 +557,7 @@ func New(ctx context.Context, dirPath string, opts *Options) (*DB, error) { cacheMinEvictCount: opts.CacheMinEvictCount, minimumRadius: uint8(opts.MinimumStorageRadius), capacityDoubling: opts.ReserveCapacityDoubling, + phaseFunc: opts.ReservePhaseFunc, }, directUploadLimiter: make(chan struct{}, pusher.ConcurrentPushes), pinIntegrity: pinIntegrity, @@ -687,6 +692,21 @@ func (db *DB) StartReserveWorker(ctx context.Context, s Syncer, radius func() (u }) } +// SetReservePhaseFunc sets the callback used by the reserve worker to report +// counting, eviction and pullsync phases. Must be called before StartReserveWorker. +func (db *DB) SetReservePhaseFunc(fn func(string)) { + if db == nil { + return + } + db.reserveOptions.phaseFunc = fn +} + +func (db *DB) reportPhase(phase string) { + if db.reserveOptions.phaseFunc != nil { + db.reserveOptions.phaseFunc(phase) + } +} + type noopRetrieval struct{} func (noopRetrieval) RetrieveChunk(_ context.Context, _ swarm.Address, _ swarm.Address) (swarm.Chunk, error) {