diff --git a/cmd/beekeeper/cmd/cluster.go b/cmd/beekeeper/cmd/cluster.go index 58782c85..518ff0c2 100644 --- a/cmd/beekeeper/cmd/cluster.go +++ b/cmd/beekeeper/cmd/cluster.go @@ -396,16 +396,21 @@ func setupNodeOptions(node config.ClusterNode, bConfig *orchestration.Config) or } const ( - fundAttempts = 2 + fundAttempts = 3 fundRetryDelay = 15 * time.Second ) // fund tops the given addresses up to the configured minimum amounts, retrying -// once so that a transient RPC error does not abort the cluster setup. +// so that a transient RPC error does not abort the cluster setup. // // node-funder reads balances from the latest mined block and returns as soon as // a transfer is broadcast, so a retry that runs before the previous transfer is // mined funds the address twice. Keep the delay above the block time. +// +// A failed send can leave a later transfer queued behind a nonce gap. The retry +// that fills the gap then rebuilds that queued transfer byte for byte, and geth +// rejects it as "already known" although it gets mined. The extra attempt only +// sees the settled balances and sends nothing. func fund( ctx context.Context, fundAddresses []string, diff --git a/config/testnet-bee-playground.yaml b/config/testnet-bee-playground.yaml index f0a3bc80..5bf9fa1e 100644 --- a/config/testnet-bee-playground.yaml +++ b/config/testnet-bee-playground.yaml @@ -9,7 +9,7 @@ clusters: api-insecure-tls: true api-scheme: http funding: - eth: 1 + eth: 5 bzz: 100.0 node-groups: bootnode: @@ -25,15 +25,15 @@ clusters: mode: node bee-config: geth-playground config: ng-bee-playground - count: 8 + count: 16 # node-groups defines node groups that can be registered in the cluster # node-groups may inherit it's configuration from already defined node-group and override specific fields from it node-groups: ng-bee-playground: _inherit: default - persistence-enabled: true - image: ethersphere/bee:latest + persistence-enabled: false + image: akram951/bee:feat-per-bin-sync-rate-4000 ingress-class: "nginx-oss" ingress-annotations: nginx.org/client-max-body-size: "10g" @@ -57,6 +57,7 @@ bee-configs: network-id: 12345 p2p-addr: :1634 password: "beekeeper" + payment-threshold: 108000000 storage-incentives-enable: true swap-enable: true verbosity: 5 @@ -200,3 +201,19 @@ checks: duration: 12h timeout: 13h type: smoke + pg-storage-radius: + options: + reserve-capacity: 4000 + target-fill-percent: 1.03 + chunks-per-upload: 512 + postage-depth: 22 + postage-amount: 2073600000 + postage-label: storage-radius + upload-timeout: 30m + upload-wave-pause: 5s + poll-interval: 2s + min-radius-wait: 5m + dilute-depth: 32 + dilute-wait: 45m + timeout: 80m + type: storage-radius diff --git a/pkg/bee/api/api.go b/pkg/bee/api/api.go index f943fc33..6346e47a 100644 --- a/pkg/bee/api/api.go +++ b/pkg/bee/api/api.go @@ -58,6 +58,7 @@ type Client struct { Status *StatusService Stewardship *StewardshipService Tags *TagsService + DebugStore *DebugStoreService } // NewClient constructs a new Client. @@ -108,6 +109,7 @@ func newClient(apiURL *url.URL, httpClient *http.Client) (c *Client) { c.Status = (*StatusService)(&c.service) c.Stewardship = (*StewardshipService)(&c.service) c.Tags = (*TagsService)(&c.service) + c.DebugStore = (*DebugStoreService)(&c.service) return c } diff --git a/pkg/bee/api/debugstore.go b/pkg/bee/api/debugstore.go index 66f6750c..f81d58e0 100644 --- a/pkg/bee/api/debugstore.go +++ b/pkg/bee/api/debugstore.go @@ -9,11 +9,49 @@ import ( type DebugStoreService service // DebugStore represents DebugStore's response -type DebugStore map[string]int +type DebugStore struct { + Upload UploadStat `json:"upload"` + Pinning PinningStat `json:"pinning"` + Cache CacheStat `json:"cache"` + Reserve ReserveStat `json:"reserve"` + ChunkStore ChunkStoreStat `json:"chunkStore"` +} + +// UploadStat reports the upload store, which holds chunks the pusher has not yet +// delivered to the network. PendingUpload is that undelivered backlog. +type UploadStat struct { + TotalUploaded int `json:"totalUploaded"` + TotalSynced int `json:"totalSynced"` + PendingUpload int `json:"pendingUpload"` +} + +type PinningStat struct { + TotalCollections int `json:"totalCollections"` + TotalChunks int `json:"totalChunks"` +} + +type CacheStat struct { + Size int `json:"size"` + Capacity int `json:"capacity"` +} + +type ReserveStat struct { + SizeWithinRadius int `json:"sizeWithinRadius"` + TotalSize int `json:"totalSize"` + Capacity int `json:"capacity"` + LastBinIDs []uint64 `json:"lastBinIDs"` + Epoch uint64 `json:"epoch"` +} + +type ChunkStoreStat struct { + TotalChunks int `json:"totalChunks"` + SharedSlots int `json:"sharedSlots"` + ReferenceCount int `json:"referenceCount"` +} // GetDebugStore gets db indices func (d *DebugStoreService) GetDebugStore(ctx context.Context) (DebugStore, error) { - resp := make(DebugStore) + var resp DebugStore err := d.client.requestJSON(ctx, http.MethodGet, "/debugstore", nil, &resp) return resp, err } diff --git a/pkg/bee/api/status.go b/pkg/bee/api/status.go index 32e98a16..de0028f0 100644 --- a/pkg/bee/api/status.go +++ b/pkg/bee/api/status.go @@ -22,6 +22,7 @@ type StatusResponse struct { IsReachable bool `json:"isReachable"` LastSyncedBlock uint64 `json:"lastSyncedBlock"` CommittedDepth uint8 `json:"committedDepth"` + IsWarmingUp bool `json:"isWarmingUp"` } // Ping pings given node diff --git a/pkg/check/storageradius/storageradius.go b/pkg/check/storageradius/storageradius.go new file mode 100644 index 00000000..8c568b82 --- /dev/null +++ b/pkg/check/storageradius/storageradius.go @@ -0,0 +1,617 @@ +package storageradius + +import ( + "context" + crand "crypto/rand" + "errors" + "fmt" + "sync" + "time" + + "github.com/ethersphere/beekeeper/pkg/bee" + "github.com/ethersphere/beekeeper/pkg/bee/api" + "github.com/ethersphere/beekeeper/pkg/beekeeper" + "github.com/ethersphere/beekeeper/pkg/logging" + "github.com/ethersphere/beekeeper/pkg/orchestration" + "github.com/ethersphere/beekeeper/pkg/random" + "golang.org/x/sync/errgroup" +) + +const ( + stablePollsBeforeGivingUp = 5 +) + +type Options struct { + TargetFillPercent float64 // fraction of capacity to fill (>1 overshoots) + ReserveCapacity int // per-node reserve capacity + PollInterval time.Duration // how often to check node status + ChunksPerUpload int // chunks per upload request + MinRadiusWait time.Duration // min time watching for radius increase + DiluteDepth uint64 // depth to dilute to (32 is max) + DiluteWait time.Duration // timeout for radius decrease + UploadWavePause time.Duration // pause between upload dispatches, so the watcher can catch up + UploadTimeout time.Duration // timeout per upload request + PostageAmount int64 + PostageLabel string + PostageDepth uint64 +} + +func NewDefaultOptions() Options { + return Options{ + PollInterval: 2 * time.Second, + TargetFillPercent: 1.2, + ChunksPerUpload: 512, + ReserveCapacity: 4000, + PostageDepth: 22, + PostageAmount: 2073600000, + PostageLabel: "storage-radius-check", + MinRadiusWait: 5 * time.Minute, + DiluteDepth: 32, + DiluteWait: 20 * time.Minute, + UploadWavePause: 5 * time.Second, + UploadTimeout: 5 * time.Minute, + } +} + +var _ beekeeper.Action = (*Check)(nil) + +type Check struct { + logger logging.Logger +} + +func NewCheck(logger logging.Logger) beekeeper.Action { + return &Check{logger: logger} +} + +func (c *Check) Run(ctx context.Context, cluster orchestration.Cluster, opts any) error { + o, ok := opts.(Options) + if !ok { + return errors.New("invalid options type") + } + + startedAt := time.Now() + + if err := c.waitForWarmup(ctx, cluster, o); err != nil { + return fmt.Errorf("wait for warmup: %w", err) + } + + fullNodes, err := cluster.ShuffledFullNodeClients(ctx, random.PseudoGenerator(time.Now().UnixNano())) + if err != nil { + return fmt.Errorf("get shuffled full node clients: %w", err) + } + if len(fullNodes) == 0 { + return errors.New("no full nodes available, storage-radius check requires at least one full node") + } + + status, err := fullNodes[0].Status(ctx) + if err != nil { + return fmt.Errorf("status: %w", err) + } + + initialStorageRadius := status.StorageRadius + + uploadPlan := newUploadPlan(o) + + c.logger.Infof("cluster: %d full nodes, sizing upload to trigger radius increase", len(fullNodes)) + c.logger.Infof("target %d chunks (%.0f%% of %d), %d chunks per upload", uploadPlan.totalChunks, o.TargetFillPercent*100, o.ReserveCapacity, uploadPlan.chunksPerUpload) + + chunksBefore, err := c.logClusterState(ctx, cluster, "before") + if err != nil { + return err + } + + batches, err := c.prepareBatches(ctx, fullNodes, len(fullNodes), o) + if err != nil { + return err + } + + uploadedChunks := 0 + if node, radius := c.radiusAlreadyRaised(ctx, batches); radius > 0 { + c.logger.Infof("%s already reports storage radius %d, skipping uploads", node.Name(), radius) + } else { + uploadedChunks, err = c.upload(ctx, batches, uploadPlan, o) + if err != nil { + return err + } + } + + risenNode, storageRadius, err := c.waitForStorageRadiusIncrease(ctx, fullNodes, o) + if err != nil { + return err + } + + chunksAfter, err := c.logClusterState(ctx, cluster, "after") + if err != nil { + return err + } + + c.logger.Infof("uploaded %d chunks in %s, cluster reserves hold %d", + uploadedChunks, time.Since(startedAt).Round(time.Second), chunksAfter) + + if storageRadius == 0 { + return c.radiusUnchangedError(chunksBefore, chunksAfter, uploadedChunks, o) + } + + c.logger.Infof("storage radius is %d (started at %d)", storageRadius, initialStorageRadius) + + if err := c.dilute(ctx, risenNode, batches, storageRadius, o); err != nil { + return err + } + + c.logger.Infof("storage-radius check finished in %s", time.Since(startedAt).Round(time.Second)) + + return nil +} + +func (c *Check) waitForWarmup(ctx context.Context, cluster orchestration.Cluster, o Options) error { + ticker := time.NewTicker(o.PollInterval) + defer ticker.Stop() + + for { + clients, err := cluster.ShuffledFullNodeClients(ctx, random.PseudoGenerator(0)) + if err != nil { + return fmt.Errorf("get full node clients: %w", err) + } + if len(clients) == 0 { + return errors.New("no full nodes available") + } + + warmingUp := 0 + for _, client := range clients { + status, err := client.Status(ctx) + if err != nil { + return fmt.Errorf("node %s: status: %w", client.Name(), err) + } + if status.IsWarmingUp { + warmingUp++ + } + } + + if warmingUp == 0 { + c.logger.Infof("all %d full nodes finished warming up", len(clients)) + return nil + } + c.logger.Infof("waiting for %d/%d full nodes to finish warming up", warmingUp, len(clients)) + + select { + case <-ctx.Done(): + return fmt.Errorf("timed out waiting for full nodes to finish warming up: %w", ctx.Err()) + case <-ticker.C: + } + } +} + +// uploadPlan is the computed upload sizing for a cluster. +type uploadPlan struct { + totalChunks int + chunksPerUpload int +} + +// newUploadPlan calculates upload size needed to trigger a radius increase. +// We size by reserve capacity with target fill percent; no need to multiply by +// neighborhoods since we're just overflowing reserves to trigger eviction/radius logic. +func newUploadPlan(o Options) uploadPlan { + totalChunks := int(o.TargetFillPercent * float64(o.ReserveCapacity)) + + return uploadPlan{ + totalChunks: totalChunks, + chunksPerUpload: min(o.ChunksPerUpload, totalChunks), + } +} + +// logClusterState reports each node's radius and reserve size, returning the +// cluster-wide chunk total. +func (c *Check) logClusterState(ctx context.Context, cluster orchestration.Cluster, label string) (int, error) { + clients, err := cluster.NodesClients(ctx) + if err != nil { + return 0, fmt.Errorf("get nodes clients: %w", err) + } + + reserveTotal := 0 + c.logger.Infof("cluster state (%s):", label) + for name, client := range clients { + status, err := client.Status(ctx) + if err != nil { + c.logger.Infof(" %s: status unavailable: %v", name, err) + continue + } + reserveTotal += int(status.ReserveSize) + c.logger.Infof(" %s: radius %d, reserve %d (within radius %d), committed depth %d", + name, status.StorageRadius, status.ReserveSize, status.ReserveSizeWithinRadius, status.CommittedDepth) + } + + return reserveTotal, nil +} + +// nodeBatch pairs a postage batch with the node that owns it. +type nodeBatch struct { + batchID string + node *bee.Client +} + +// prepareBatches buys postage batches for each node, reusing existing ones to save time and tokens. +// Uses WaitGroup rather than errgroup because a failed purchase must not fail the whole check: +// batch creation can revert on-chain per node, and the check only needs enough batches to fill +// one reserve, so failures are collected and reported while the usable batches are returned. +func (c *Check) prepareBatches(ctx context.Context, nodes orchestration.ClientList, batchCount int, o Options) ([]nodeBatch, error) { + c.logger.Infof("preparing %d postage batches in parallel (depth %d, amount %d)", + batchCount, o.PostageDepth, o.PostageAmount) + + startedAt := time.Now() + + var ( + mutex sync.Mutex + batches []nodeBatch + failedNodes []string + waitGroup sync.WaitGroup + ) + + for i := range batchCount { + node := nodes[i] + waitGroup.Add(1) + go func() { + defer waitGroup.Done() + + batchID, reused, err := c.batchForNode(ctx, node, o) + + mutex.Lock() + defer mutex.Unlock() + if err != nil { + c.logger.Infof("%s: no usable batch: %v", node.Name(), err) + failedNodes = append(failedNodes, node.Name()) + return + } + if reused { + c.logger.Infof("%s: reusing batch %s", node.Name(), batchID) + } else { + c.logger.Infof("%s: bought batch %s", node.Name(), batchID) + } + batches = append(batches, nodeBatch{batchID: batchID, node: node}) + }() + } + waitGroup.Wait() + + if len(batches) == 0 { + return nil, fmt.Errorf("no usable postage batches: all %d nodes failed", batchCount) + } + if len(failedNodes) > 0 { + c.logger.Infof("continuing with %d/%d batches, failed on %v", len(batches), batchCount, failedNodes) + } + c.logger.Infof("%d batches ready in %s", len(batches), time.Since(startedAt).Round(time.Second)) + + return batches, nil +} + +// batchForNode finds an existing usable batch or creates a new one. +func (c *Check) batchForNode(ctx context.Context, node *bee.Client, o Options) (batchID string, reused bool, err error) { + label := fmt.Sprintf("%s-%s", o.PostageLabel, node.Name()) + + existingBatches, err := node.PostageBatches(ctx) + if err != nil { + return "", false, fmt.Errorf("list batches: %w", err) + } + + for _, batch := range existingBatches { + if !batch.Exists || batch.ImmutableFlag || !batch.Usable || batch.Label != label { + continue + } + if batch.BatchTTL == 0 { + continue // expired + } + if batch.Utilization >= 1<<(batch.Depth-batch.BucketDepth) { + continue // buckets full, cannot issue more stamps + } + return batch.BatchID, true, nil + } + + batchID, err = node.CreatePostageBatch(ctx, o.PostageAmount, o.PostageDepth, label, false) + if err != nil { + return "", false, err + } + return batchID, false, nil +} + +// waitForStorageRadiusIncrease waits for any node's radius to climb above zero, returning +// that node and its radius as soon as one does, so the caller can later check the same +// node's radius for the decrease rather than a different node that never moved. It gives +// up and returns a nil node with radius 0 only once the reserves have stopped growing AND +// MinRadiusWait has elapsed, since pullsync replicates chunks after the uploads return. +func (c *Check) waitForStorageRadiusIncrease(ctx context.Context, nodes orchestration.ClientList, o Options) (*bee.Client, uint8, error) { + ticker := time.NewTicker(o.PollInterval) + defer ticker.Stop() + + c.logger.Infof("waiting up to %s for the reserves to fill and the radius to rise", o.MinRadiusWait) + + startedAt := time.Now() + prevReserveSize, stablePolls := -1, 0 + + for { + select { + case <-ctx.Done(): + return nil, 0, fmt.Errorf("timed out waiting for the storage radius to rise above 0: %w", ctx.Err()) + case <-ticker.C: + } + + reserveTotal, risenNode, highestRadius := c.reserveState(ctx, nodes) + if highestRadius > 0 { + c.logger.Infof("storage radius is %d after %s (reserves at %d chunks)", + highestRadius, time.Since(startedAt).Round(time.Second), reserveTotal) + return risenNode, highestRadius, nil + } + + elapsed := time.Since(startedAt) + if reserveTotal == prevReserveSize { + stablePolls++ + if stablePolls >= stablePollsBeforeGivingUp && elapsed >= o.MinRadiusWait { + c.logger.Infof("reserves settled at %d chunks and radius still 0 after %s", + reserveTotal, elapsed.Round(time.Second)) + return nil, 0, nil + } + } else { + if prevReserveSize >= 0 { + c.logger.Infof("reserves at %d chunks (+%d), radius 0 (%s elapsed)", + reserveTotal, reserveTotal-prevReserveSize, elapsed.Round(time.Second)) + } + stablePolls = 0 + } + prevReserveSize = reserveTotal + } +} + +// radiusAlreadyRaised returns the first batch node whose storage radius is already +// above 0, so a rerun against a filled cluster can skip the uploads. +func (c *Check) radiusAlreadyRaised(ctx context.Context, batches []nodeBatch) (*bee.Client, uint8) { + for _, batch := range batches { + status, err := batch.node.Status(ctx) + if err != nil { + continue + } + if status.StorageRadius > 0 { + return batch.node, status.StorageRadius + } + } + return nil, 0 +} + +// upload sends random data in parallel, stopping once the radius rises or a reserve is over capacity. +func (c *Check) upload(ctx context.Context, batches []nodeBatch, plan uploadPlan, options Options) (int, error) { + totalUploads := plan.uploadCount() + + c.logger.Infof("uploading %d chunks in %d requests across %d nodes", + plan.totalChunks, totalUploads, len(batches)) + + var ( + mutex sync.Mutex + uploadedChunks int + completedCount int + ) + + group, groupCtx := errgroup.WithContext(ctx) + group.SetLimit(len(batches)) + + // Once the radius rises, in-flight uploads are no longer needed: cancel them + // so a slow one cannot outlive the goal and fail the check on its timeout. + uploadsCtx, cancelUploads := context.WithCancel(groupCtx) + defer cancelUploads() + + enough := make(chan struct{}) + var stopOnce sync.Once + stopUploading := func() { + stopOnce.Do(func() { + close(enough) + cancelUploads() + }) + } + + watchCtx, cancelWatch := context.WithCancel(groupCtx) + defer cancelWatch() + go c.stopWhenRadiusRises(watchCtx, batches, options, stopUploading) + + for i := range totalUploads { + if i > 0 && i%len(batches) == 0 { + // Pause between waves so the watcher's independent poll loop gets a chance to observe a radius rise and call stopUploading before the next wave fires off. Without this, a fast cluster can spin up all uploads within a single PollInterval and the watcher never sees them. + select { + case <-enough: + case <-groupCtx.Done(): + case <-time.After(options.UploadWavePause): + } + } + + batch := batches[i%len(batches)] + group.Go(func() error { + select { + case <-enough: + return nil + default: + } + + data := make([]byte, int64(plan.chunksPerUpload)*bee.MaxChunkSize) + if _, err := crand.Read(data); err != nil { + return fmt.Errorf("generate random data: %w", err) + } + + select { + case <-enough: + return nil + default: + } + + uploadCtx, cancel := context.WithTimeout(uploadsCtx, options.UploadTimeout) + address, err := batch.node.UploadBytes(uploadCtx, data, api.UploadOptions{BatchID: batch.batchID, Direct: true}) + cancel() + if err != nil { + select { + case <-enough: + c.logger.Infof("upload to %s abandoned, no longer needed: %v", batch.node.Name(), err) + return nil + default: + } + return fmt.Errorf("upload to %s: %w", batch.node.Name(), err) + } + + mutex.Lock() + uploadedChunks += plan.chunksPerUpload + completedCount++ + completed, chunks := completedCount, uploadedChunks + mutex.Unlock() + + c.logger.Infof("upload %d/%d to %s: %s (%d/%d chunks)", + completed, totalUploads, batch.node.Name(), address, chunks, plan.totalChunks) + return nil + }) + } + + if err := group.Wait(); err != nil { + return uploadedChunks, err + } + + return uploadedChunks, nil +} + +// radiusUnchangedError explains why the radius stayed at zero. +func (c *Check) radiusUnchangedError(chunksBefore, chunksAfter, uploadedChunks int, options Options) error { + if chunksAfter <= chunksBefore { + return fmt.Errorf("storage radius is still 0 and the reserves did not grow (%d chunks before, %d after, %d uploaded): "+ + "bee accepted the uploads but the chunks are not reaching the reserves", + chunksBefore, chunksAfter, uploadedChunks) + } + return fmt.Errorf("storage radius is still 0: reserves grew from %d to %d chunks, "+ + "but no node exceeded its %d-chunk capacity long enough to force an increase", + chunksBefore, chunksAfter, options.ReserveCapacity) +} + +// dilute increases batch depths to push chunks outside the storage radius. +func (c *Check) dilute(ctx context.Context, risenNode *bee.Client, batches []nodeBatch, startRadius uint8, options Options) error { + c.logger.Infof("diluting %d batches to depth %d to push chunks outside the storage radius", + len(batches), options.DiluteDepth) + + dilutedCount := 0 + for _, batch := range batches { + stamp, err := batch.node.PostageStamp(ctx, batch.batchID) + if err != nil { + c.logger.Infof("%s: cannot read batch %s: %v", batch.node.Name(), batch.batchID, err) + continue + } + if uint64(stamp.Depth) >= options.DiluteDepth { + c.logger.Infof("%s: batch already at depth %d, skipping", batch.node.Name(), stamp.Depth) + continue + } + + if err := batch.node.DilutePostageBatch(ctx, batch.batchID, options.DiluteDepth, ""); err != nil { + c.logger.Infof("%s: dilute to depth %d failed: %v", batch.node.Name(), options.DiluteDepth, err) + continue + } + dilutedCount++ + c.logger.Infof("%s: diluted batch from depth %d to %d", batch.node.Name(), stamp.Depth, options.DiluteDepth) + } + + if dilutedCount == 0 { + return errors.New("no batches were diluted, cannot provoke a radius decrease") + } + + return c.waitForStorageRadiusDecrease(ctx, risenNode, startRadius, options) +} + +// waitForStorageRadiusDecrease waits for the node whose radius rose to drop back below +// startRadius. It checks the same node that triggered the increase rather than a +// cluster-wide extremum, since radius is decided locally per node and other nodes may +// never have left radius 0. +func (c *Check) waitForStorageRadiusDecrease(ctx context.Context, risenNode *bee.Client, startRadius uint8, options Options) error { + ticker := time.NewTicker(options.PollInterval) + defer ticker.Stop() + + c.logger.Infof("waiting up to %s for %s's storage radius to fall below %d", options.DiluteWait, risenNode.Name(), startRadius) + + startedAt := time.Now() + + for { + select { + case <-ctx.Done(): + return fmt.Errorf("timed out waiting for the storage radius to decrease: %w", ctx.Err()) + case <-ticker.C: + } + + elapsed := time.Since(startedAt) + + status, err := risenNode.Status(ctx) + if err != nil { + if elapsed >= options.DiluteWait { + return fmt.Errorf("storage radius stayed at %d after %s: %s did not answer: %w", + startRadius, options.DiluteWait, risenNode.Name(), err) + } + c.logger.Infof("%s did not answer, retrying (%s elapsed)", risenNode.Name(), elapsed.Round(time.Second)) + continue + } + + if status.StorageRadius < startRadius { + c.logger.Infof("%s's storage radius decreased %d -> %d after %s", + risenNode.Name(), startRadius, status.StorageRadius, elapsed.Round(time.Second)) + return nil + } + + if elapsed >= options.DiluteWait { + return fmt.Errorf("%s's storage radius stayed at %d after %s: %d chunks within radius, pullsync rate %.2f", + risenNode.Name(), startRadius, options.DiluteWait, status.ReserveSizeWithinRadius, status.PullsyncRate) + } + + c.logger.Infof("%s's radius still %d, %d chunks within radius, pullsync %.2f (%s elapsed)", + risenNode.Name(), status.StorageRadius, status.ReserveSizeWithinRadius, status.PullsyncRate, elapsed.Round(time.Second)) + } +} + +// reserveState returns the total reserve size, the highest radius in the cluster, +// and the node that reported it. +func (c *Check) reserveState(ctx context.Context, nodes orchestration.ClientList) (reserveTotal int, highestRadiusNode *bee.Client, highestRadius uint8) { + for _, node := range nodes { + status, err := node.Status(ctx) + if err != nil { + continue + } + reserveTotal += int(status.ReserveSize) + if status.StorageRadius > highestRadius || highestRadiusNode == nil { + highestRadius = status.StorageRadius + highestRadiusNode = node + } + } + return reserveTotal, highestRadiusNode, highestRadius +} + +// uploadCount returns the total number of upload requests needed. +func (p uploadPlan) uploadCount() int { + return max((p.totalChunks+p.chunksPerUpload-1)/p.chunksPerUpload, 1) +} + +// stopWhenRadiusRises halts uploads once any batch node's radius rises or its reserve +// exceeds capacity, since the remaining uploads are then no longer needed. +func (c *Check) stopWhenRadiusRises(ctx context.Context, batches []nodeBatch, options Options, stopUploading func()) { + ticker := time.NewTicker(options.PollInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + + for _, batch := range batches { + status, err := batch.node.Status(ctx) + if err != nil { + continue + } + + if status.StorageRadius > 0 { + c.logger.Infof("%s reports storage radius %d, stopping further uploads", + batch.node.Name(), status.StorageRadius) + stopUploading() + return + } + if int(status.ReserveSize) > options.ReserveCapacity { + c.logger.Infof("%s reached %d/%d chunks, stopping further uploads", + batch.node.Name(), status.ReserveSize, options.ReserveCapacity) + stopUploading() + return + } + } + } +} diff --git a/pkg/config/check.go b/pkg/config/check.go index 541afb7e..1f9090ee 100644 --- a/pkg/config/check.go +++ b/pkg/config/check.go @@ -35,6 +35,7 @@ import ( "github.com/ethersphere/beekeeper/pkg/check/smoke" "github.com/ethersphere/beekeeper/pkg/check/soc" "github.com/ethersphere/beekeeper/pkg/check/stake" + "github.com/ethersphere/beekeeper/pkg/check/storageradius" "github.com/ethersphere/beekeeper/pkg/check/withdraw" "github.com/ethersphere/beekeeper/pkg/logging" "github.com/ethersphere/beekeeper/pkg/random" @@ -734,6 +735,35 @@ var Checks = map[string]CheckType{ return nil, fmt.Errorf("applying options: %w", err) } + return opts, nil + }, + }, + "storage-radius": { + NewAction: storageradius.NewCheck, + NewOptions: func(checkGlobalConfig CheckGlobalConfig, check Check) (any, error) { + checkOpts := new(struct { + PollInterval *time.Duration `yaml:"poll-interval"` + TargetFillPercent *float64 `yaml:"target-fill-percent"` + ChunksPerUpload *int `yaml:"chunks-per-upload"` + ReserveCapacity *int `yaml:"reserve-capacity"` + PostageDepth *uint64 `yaml:"postage-depth"` + PostageAmount *int64 `yaml:"postage-amount"` + PostageLabel *string `yaml:"postage-label"` + MinRadiusWait *time.Duration `yaml:"min-radius-wait"` + DiluteDepth *uint64 `yaml:"dilute-depth"` + DiluteWait *time.Duration `yaml:"dilute-wait"` + UploadWavePause *time.Duration `yaml:"upload-wave-pause"` + UploadTimeout *time.Duration `yaml:"upload-timeout"` + }) + if err := check.Options.Decode(checkOpts); err != nil { + return nil, fmt.Errorf("decoding check %s options: %w", check.Type, err) + } + opts := storageradius.NewDefaultOptions() + + if err := applyCheckConfig(checkGlobalConfig, checkOpts, &opts); err != nil { + return nil, fmt.Errorf("applying options: %w", err) + } + return opts, nil }, },