Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions block/components.go
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ func NewSyncComponents(
if p, ok := exec.(coreexecutor.ExecPruner); ok {
execPruner = p
}
pruner := pruner.New(logger, store, execPruner, config.Pruning, config.Node.BlockTime.Duration, config.DA.Address)
pruner := pruner.New(logger, store, execPruner, config.Pruning, config.Node.BlockTime.Duration)

var submitter *submitting.Submitter
if daClient != nil {
Expand Down Expand Up @@ -278,7 +278,7 @@ func newAggregatorComponents(
if p, ok := exec.(coreexecutor.ExecPruner); ok {
execPruner = p
}
pruner := pruner.New(logger, store, execPruner, config.Pruning, config.Node.BlockTime.Duration, config.DA.Address)
pruner := pruner.New(logger, store, execPruner, config.Pruning, config.Node.BlockTime.Duration)

reaper, err := reaping.NewReaper(
exec,
Expand Down
157 changes: 84 additions & 73 deletions block/internal/pruner/pruner.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,13 +16,12 @@ import (
"github.com/evstack/ev-node/pkg/store"
)

// Pruner periodically removes old state and execution metadata entries.
// Pruner periodically removes old blocks, state snapshots, and execution metadata.
type Pruner struct {
store store.Store
execPruner coreexecutor.ExecPruner
cfg config.PruningConfig
blockTime time.Duration
daEnabled bool
logger zerolog.Logger

// Lifecycle
Expand All @@ -38,14 +37,12 @@ func New(
execPruner coreexecutor.ExecPruner,
cfg config.PruningConfig,
blockTime time.Duration,
daAddress string,
) *Pruner {
return &Pruner{
store: store,
execPruner: execPruner,
cfg: cfg,
blockTime: blockTime,
daEnabled: daAddress != "", // DA is enabled if address is provided
logger: logger.With().Str("component", "pruner").Logger(),
}
}
Expand Down Expand Up @@ -108,75 +105,71 @@ func (p *Pruner) pruneLoop() {
}
}

// pruneBlocks prunes blocks and their metadatas.
// pruneBlocks prunes blocks and state snapshots outside local retention.
func (p *Pruner) pruneBlocks() error {
storeHeight, err := p.store.Height(p.ctx)
if err != nil {
return fmt.Errorf("failed to get store height for pruning: %w", err)
}

upperBound := storeHeight

// If DA is enabled, only prune blocks that are DA included
if p.daEnabled {
var currentDAIncluded uint64
currentDAIncludedBz, err := p.store.GetMetadata(p.ctx, store.DAIncludedHeightKey)
if err == nil && len(currentDAIncludedBz) == 8 {
currentDAIncluded = binary.LittleEndian.Uint64(currentDAIncludedBz)
} else {
p.logger.Debug().Msg("skipping pruning: DA is enabled but DA included height is not available yet")
return nil
}

// Never prune blocks that are not DA included
upperBound = min(storeHeight, currentDAIncluded)
}

if upperBound <= p.cfg.KeepRecent {
// Not enough fully included blocks to prune
// Retention follows the locally committed head. DA inclusion can remain
// stalled when historical DA heights are no longer served, and must not
// prevent operators from reclaiming local storage.
if storeHeight <= p.cfg.KeepRecent {
return nil
}
targetHeight := storeHeight - p.cfg.KeepRecent

targetHeight := upperBound - p.cfg.KeepRecent

// Get the last pruned height to determine batch size
lastPruned, err := p.getLastPrunedBlockHeight(p.ctx)
lastBlock, err := p.getLastPrunedBlockHeight(p.ctx)
if err != nil {
return fmt.Errorf("failed to get last pruned block height: %w", err)
}

// prune in batches to avoid overwhelming the system
batchSize := p.calculateBatchSize()
batchEnd := min(lastPruned+batchSize, targetHeight)

if err := p.store.PruneBlocks(p.ctx, batchEnd); err != nil {
p.logger.Error().Err(err).Uint64("target_height", batchEnd).Msg("failed to prune old block data")
return err
lastState, err := p.getLastPrunedStateHeight(p.ctx)
if err != nil {
return fmt.Errorf("failed to get last pruned state height: %w", err)
}

if p.execPruner != nil {
if err := p.execPruner.PruneExec(p.ctx, batchEnd); err != nil && !errors.Is(err, ds.ErrNotFound) {
// State may lag behind blocks after upgrading from a version that only
// deleted blocks in all mode, or after a failed pruning operation.
start := min(lastBlock, lastState)
if start >= targetHeight {
return nil
}
end := start + min(p.calculateBatchSize(), targetHeight-start)
for start < end {
if err := p.ctx.Err(); err != nil {
return err
}
batchEnd := start + min(maxPruningBatchSize, end-start)
if lastBlock < batchEnd {
if err := p.store.PruneBlocks(p.ctx, batchEnd); err != nil {
return fmt.Errorf("failed to prune blocks through height %d: %w", batchEnd, err)
}
lastBlock = batchEnd
}
if err := p.pruneState(lastState, batchEnd); err != nil {
return err
}
lastState = max(lastState, batchEnd)
start = batchEnd
}

p.logger.Debug().Uint64("pruned_up_to_height", batchEnd).Bool("da_enabled", p.daEnabled).Msg("pruned blocks up to height")
p.logger.Debug().Uint64("pruned_up_to_height", end).Msg("pruned blocks and state snapshots")
return nil
}

// calculateBatchSize returns the appropriate batch size for pruning operations.
// The batch size is based on the pruning interval and block time to ensure reasonable progress
// without overwhelming the node.
// maxPruningBatchSize bounds each datastore batch, rather than the total
// progress per interval, so fast chains can still catch up with old history.
const maxPruningBatchSize uint64 = 10000

// calculateBatchSize returns a per-interval work budget of four times the
// expected block production. Each datastore batch is bounded independently.
func (p *Pruner) calculateBatchSize() uint64 {
// Calculate batch size based on pruning interval and block time.
// We use 4x the blocks produced during one pruning interval as the batch size.
// This ensures we catch up at 3x the block production rate when there's a backlog.
// Example: With 100ms blocks and 15min interval: 15*60/0.1 = 9000 blocks/interval
// - Batch size: 36,000 blocks (prunes 36k, chain grows 9k = net 27k catch-up per interval)
blocksPerInterval := uint64(p.cfg.Interval.Duration / p.blockTime)

// Ensure reasonable minimum
return min(blocksPerInterval*4, 10000)
if p.blockTime <= 0 || p.cfg.Interval.Duration <= 0 {
return 1
}
blocksPerInterval := max(uint64(p.cfg.Interval.Duration/p.blockTime), 1)
return min(blocksPerInterval, ^uint64(0)/4) * 4
}

// pruneMetadata prunes old state and execution metadata entries based on the configured retention depth.
Expand All @@ -197,37 +190,49 @@ func (p *Pruner) pruneMetadata() error {
return fmt.Errorf("failed to get last pruned state height: %w", err)
}

if lastPrunedBlock, err := p.getLastPrunedBlockHeight(p.ctx); err == nil && lastPrunedBlock > lastPrunedState {
lastPrunedState = lastPrunedBlock
}

target := height - p.cfg.KeepRecent
if target <= lastPrunedState {
return nil
}

batchSize := p.calculateBatchSize()
end := lastPrunedState + min(p.calculateBatchSize(), target-lastPrunedState)
for lastPrunedState < end {
if err := p.ctx.Err(); err != nil {
return err
}
batchEnd := lastPrunedState + min(maxPruningBatchSize, end-lastPrunedState)
if err := p.pruneState(lastPrunedState, batchEnd); err != nil {
return err
}
lastPrunedState = batchEnd
}

// prune in batches to avoid overwhelming the system
batchEnd := min(lastPrunedState+batchSize, target)
p.logger.Debug().Uint64("pruned_to", end).Msg("pruned state height metadata up to height")
return nil
}

for h := lastPrunedState + 1; h <= batchEnd; h++ {
if err := p.store.DeleteStateAtHeight(p.ctx, h); err != nil && !errors.Is(err, ds.ErrNotFound) {
// pruneState retries execution pruning before advancing the state cursor, so
// failures remain recoverable even if block deletion has already committed.
func (p *Pruner) pruneState(lastPruned, end uint64) error {
for h := lastPruned; h < end; {
h++
if err := p.ctx.Err(); err != nil {
return err
}
if err := p.store.DeleteStateAtHeight(p.ctx, h); err != nil && !errors.Is(err, ds.ErrNotFound) {
return fmt.Errorf("failed to prune state at height %d: %w", h, err)
}
}

if p.execPruner != nil {
if err := p.execPruner.PruneExec(p.ctx, batchEnd); err != nil && !errors.Is(err, ds.ErrNotFound) {
return err
if err := p.execPruner.PruneExec(p.ctx, end); err != nil && !errors.Is(err, ds.ErrNotFound) {
return fmt.Errorf("failed to prune execution metadata through height %d: %w", end, err)
}
}

if err := p.setLastPrunedStateHeight(p.ctx, batchEnd); err != nil {
return fmt.Errorf("failed to set last pruned block height: %w", err)
if end > lastPruned {
if err := p.setLastPrunedStateHeight(p.ctx, end); err != nil {
return fmt.Errorf("failed to set last pruned state height: %w", err)
}
}

p.logger.Debug().Uint64("pruned_to", batchEnd).Msg("pruned state height metadata up to height")
return nil
}

Expand All @@ -239,8 +244,11 @@ func (p *Pruner) getLastPrunedBlockHeight(ctx context.Context) (uint64, error) {
return 0, nil
}

if err != nil || len(lastPrunedBlockHeightBz) != 8 {
return 0, fmt.Errorf("failed to get last pruned block height or invalid format: %w", err)
if err != nil {
return 0, fmt.Errorf("failed to get last pruned block height: %w", err)
}
if len(lastPrunedBlockHeightBz) != 8 {
return 0, errors.New("invalid last pruned block height format")
}

lastPrunedBlockHeight := binary.LittleEndian.Uint64(lastPrunedBlockHeightBz)
Expand All @@ -259,13 +267,16 @@ func (p *Pruner) getLastPrunedStateHeight(ctx context.Context) (uint64, error) {
return 0, nil
}

if err != nil || len(lastPrunedStateHeightBz) != 8 {
return 0, fmt.Errorf("failed to get last pruned block height or invalid format: %w", err)
if err != nil {
return 0, fmt.Errorf("failed to get last pruned state height: %w", err)
}
if len(lastPrunedStateHeightBz) != 8 {
return 0, errors.New("invalid last pruned state height format")
}

lastPrunedStateHeight := binary.LittleEndian.Uint64(lastPrunedStateHeightBz)
if lastPrunedStateHeight == 0 {
return 0, fmt.Errorf("invalid last pruned block height")
return 0, fmt.Errorf("invalid last pruned state height")
}

return lastPrunedStateHeight, nil
Expand Down
Loading
Loading