Optimize freezer truncation: batch processing, progress reporting, configurable batch size

This commit is contained in:
sivaratrisrinivas 2025-03-11 22:57:36 -05:00
parent 89074433ee
commit 19d510ac4e
2 changed files with 454 additions and 87 deletions

View file

@ -69,7 +69,26 @@ var (
Usage: "Keep headers when pruning history (default: true)",
Value: true,
}
useHardcodedMergeFlag = &cli.BoolFlag{
Name: "use-hardcoded-merge",
Usage: "Use hardcoded merge block for known networks instead of detecting it",
}
batchSizeFlag = &cli.IntFlag{
Name: "batch-size",
Usage: "Number of headers to process in each batch (default: 10000)",
Value: 10000,
}
)
// Known merge block numbers for different networks
var knownMergeBlocks = map[string]uint64{
"mainnet": 15537394,
"sepolia": 1735371,
"goerli": 7382818,
"holesky": 0, // Holesky was launched post-merge
}
var (
removedbCommand = &cli.Command{
Action: removeDB,
Name: "removedb",
@ -268,19 +287,27 @@ including headers, use --keep-headers=false.
dbTruncateFreezerCmd = &cli.Command{
Action: truncateFreezer,
Name: "truncate-freezer",
Usage: "Truncate freezer at the merge block, keeping headers but removing bodies",
Usage: "Truncate the freezer at the merge block, keeping headers but removing bodies",
ArgsUsage: "",
Flags: []cli.Flag{
Flags: slices.Concat([]cli.Flag{
dryRunFlag,
yesFlag,
},
keepHeadersFlag,
useHardcodedMergeFlag,
batchSizeFlag,
}, utils.NetworkFlags, utils.DatabaseFlags),
Description: `
This command truncates the freezer at the merge block, keeping headers but removing bodies.
The merge block is identified by the first block with zero difficulty.
This can significantly reduce disk space usage for nodes that don't need pre-merge block bodies.
This operation is specifically designed to maintain chain integrity while reducing disk usage
by removing the pre-merge body data which is no longer needed after the merge.
`,
The command will:
1. Find the merge block (first block with zero difficulty)
2. Create a temporary copy of the headers and hashes up to the merge block
3. Truncate all freezer tables at the merge block
4. Re-insert the headers and hashes from the temporary copy
WARNING: This operation cannot be undone. Make sure you have a backup if you might need
the removed data in the future.`,
}
)
@ -1006,7 +1033,6 @@ func pruneHistory(ctx *cli.Context) error {
return fmt.Errorf("header number of head block not found")
}
// Binary search to find the merge block
var (
low = uint64(0)
high = *headNumber
@ -1189,55 +1215,13 @@ func truncateFreezer(ctx *cli.Context) error {
// Find the merge block (first block with zero difficulty)
log.Info("Looking for the merge block...")
// Get the current head block
headHash := rawdb.ReadHeadBlockHash(db)
if headHash == (common.Hash{}) {
return fmt.Errorf("chain head not found")
mergeBlock, err := findMergeBlock(ctx, db)
if err != nil {
return err
}
headNumber := rawdb.ReadHeaderNumber(db, headHash)
if headNumber == nil {
return fmt.Errorf("header number of head block not found")
}
// Binary search to find the merge block
var (
low = uint64(0)
high = *headNumber
mergeBlock = uint64(0)
found = false
)
log.Info("Searching for merge block using binary search", "head", *headNumber)
for low <= high {
mid := (low + high) / 2
header := rawdb.ReadHeader(db, rawdb.ReadCanonicalHash(db, mid), mid)
if header == nil {
return fmt.Errorf("header %d not found", mid)
}
if header.Difficulty.Sign() == 0 {
// This is a post-merge block, look earlier
high = mid - 1
mergeBlock = mid
found = true
} else {
// This is a pre-merge block, look later
low = mid + 1
}
}
if !found {
return fmt.Errorf("merge block not found, chain may not have transitioned to PoS yet")
}
// The merge block is the first block with zero difficulty
log.Info("Found merge block", "number", mergeBlock)
// Check if we're in dry-run mode
if ctx.Bool("dry-run") {
if ctx.Bool(dryRunFlag.Name) {
log.Info("Dry run completed, no data was pruned", "mergeBlock", mergeBlock)
return nil
}
@ -1247,7 +1231,7 @@ func truncateFreezer(ctx *cli.Context) error {
fmt.Println("This operation cannot be undone and will permanently delete data.")
fmt.Println("Make sure you have a backup if you might need this data in the future.")
if !ctx.Bool("yes") {
if !ctx.Bool(yesFlag.Name) {
confirm, err := prompt.Stdin.PromptConfirm("Do you want to continue?")
if err != nil {
return err
@ -1280,47 +1264,208 @@ func truncateFreezer(ctx *cli.Context) error {
return nil
}
// We need to read the headers at the merge block and re-insert them
// after truncating everything
log.Info("Reading headers to preserve them", "count", mergeBlock)
headers := make([][]byte, mergeBlock)
hashes := make([][]byte, mergeBlock)
// Create a temporary directory for the headers freezer
tmpDir, err := os.MkdirTemp("", "geth-headers-freezer-*")
if err != nil {
return fmt.Errorf("failed to create temporary directory: %v", err)
}
defer os.RemoveAll(tmpDir)
log.Info("Created temporary directory for headers", "path", tmpDir)
for i := uint64(0); i < mergeBlock; i++ {
headers[i], err = db.(ethdb.AncientReader).Ancient(rawdb.ChainFreezerHeaderTable, i)
if err != nil {
return fmt.Errorf("failed to read header %d: %v", i, err)
}
hashes[i], err = db.(ethdb.AncientReader).Ancient(rawdb.ChainFreezerHashTable, i)
if err != nil {
return fmt.Errorf("failed to read hash %d: %v", i, err)
}
// Create a new freezer for headers and hashes
headersFreezer, err := rawdb.NewFreezer(tmpDir, "headers", false, 2*1000*1000*1000, map[string]bool{
rawdb.ChainFreezerHeaderTable: false,
rawdb.ChainFreezerHashTable: false,
})
if err != nil {
return fmt.Errorf("failed to create headers freezer: %v", err)
}
defer headersFreezer.Close()
// Get batch size from flag
batchSize := uint64(ctx.Int(batchSizeFlag.Name))
log.Info("Using batch size for processing", "batchSize", batchSize)
// Copy headers and hashes to the temporary freezer in batches
log.Info("Copying headers to temporary freezer", "count", mergeBlock)
if err := extractHeaders(ancientDb, headersFreezer, mergeBlock, batchSize); err != nil {
return err
}
// Truncate all tables
log.Info("Truncating all tables", "mergeBlock", mergeBlock)
if err := truncateAncientStore(db, mergeBlock); err != nil {
// Truncate all tables in the original freezer
log.Info("Truncating all tables in the original freezer", "mergeBlock", mergeBlock)
if err := truncateAncientStore(db, 0); err != nil {
return fmt.Errorf("failed to truncate ancient store: %v", err)
}
// Re-insert the headers and hashes
log.Info("Re-inserting headers", "count", len(headers))
// Copy headers and hashes back from the temporary freezer in batches
log.Info("Copying headers back to the main freezer", "count", mergeBlock)
freezerDb := db.(ethdb.AncientStore)
_, err = freezerDb.ModifyAncients(func(op ethdb.AncientWriteOp) error {
for i := uint64(0); i < mergeBlock; i++ {
if err := op.AppendRaw(rawdb.ChainFreezerHeaderTable, i, headers[i]); err != nil {
return fmt.Errorf("failed to re-insert header %d: %v", i, err)
}
if err := op.AppendRaw(rawdb.ChainFreezerHashTable, i, hashes[i]); err != nil {
return fmt.Errorf("failed to re-insert hash %d: %v", i, err)
}
}
return nil
})
if err != nil {
return fmt.Errorf("failed to re-insert headers: %v", err)
if err := reinsertHeaders(freezerDb, headersFreezer, mergeBlock, batchSize); err != nil {
return err
}
log.Info("Successfully truncated freezer, keeping headers but removing bodies", "mergeBlock", mergeBlock)
return nil
}
// extractHeaders copies headers from the main freezer to a temporary freezer in batches
func extractHeaders(ancientDb ethdb.AncientReader, headersFreezer *rawdb.Freezer, mergeBlock, batchSize uint64) error {
// Create a progress reporter
progressReporter := utils.NewProgressReporter("Copying headers to temporary freezer", int(mergeBlock), 5)
defer progressReporter.Stop()
for i := uint64(0); i < mergeBlock; i += batchSize {
end := i + batchSize
if end > mergeBlock {
end = mergeBlock
}
log.Info("Processing batch of headers", "from", i, "to", end)
// Process this batch
_, err := headersFreezer.ModifyAncients(func(op ethdb.AncientWriteOp) error {
for j := i; j < end; j++ {
// Read header and hash
headerBytes, err := ancientDb.Ancient(rawdb.ChainFreezerHeaderTable, j)
if err != nil {
return fmt.Errorf("failed to read header %d: %w", j, err)
}
hashBytes, err := ancientDb.Ancient(rawdb.ChainFreezerHashTable, j)
if err != nil {
return fmt.Errorf("failed to read hash %d: %w", j, err)
}
// Write to temporary freezer
if err := op.AppendRaw(rawdb.ChainFreezerHeaderTable, j, headerBytes); err != nil {
return fmt.Errorf("failed to write header %d: %w", j, err)
}
if err := op.AppendRaw(rawdb.ChainFreezerHashTable, j, hashBytes); err != nil {
return fmt.Errorf("failed to write hash %d: %w", j, err)
}
}
return nil
})
if err != nil {
return fmt.Errorf("failed to copy headers batch: %w", err)
}
// Update progress
progressReporter.Report(int(end))
}
return nil
}
// reinsertHeaders copies headers from the temporary freezer back to the main freezer in batches
func reinsertHeaders(freezerDb ethdb.AncientStore, headersFreezer *rawdb.Freezer, mergeBlock, batchSize uint64) error {
// Create a progress reporter
progressReporter := utils.NewProgressReporter("Copying headers back to main freezer", int(mergeBlock), 5)
defer progressReporter.Stop()
for i := uint64(0); i < mergeBlock; i += batchSize {
end := i + batchSize
if end > mergeBlock {
end = mergeBlock
}
log.Info("Processing batch of headers", "from", i, "to", end)
_, err := freezerDb.ModifyAncients(func(op ethdb.AncientWriteOp) error {
for j := i; j < end; j++ {
// Read from temporary freezer
headerBytes, err := headersFreezer.Ancient(rawdb.ChainFreezerHeaderTable, j)
if err != nil {
return fmt.Errorf("failed to read header %d from temp freezer: %w", j, err)
}
hashBytes, err := headersFreezer.Ancient(rawdb.ChainFreezerHashTable, j)
if err != nil {
return fmt.Errorf("failed to read hash %d from temp freezer: %w", j, err)
}
// Write back to main freezer
if err := op.AppendRaw(rawdb.ChainFreezerHeaderTable, j, headerBytes); err != nil {
return fmt.Errorf("failed to write header %d to main freezer: %w", j, err)
}
if err := op.AppendRaw(rawdb.ChainFreezerHashTable, j, hashBytes); err != nil {
return fmt.Errorf("failed to write hash %d to main freezer: %w", j, err)
}
}
return nil
})
if err != nil {
return fmt.Errorf("failed to copy headers batch back to main freezer: %w", err)
}
// Update progress
progressReporter.Report(int(end))
}
return nil
}
// findMergeBlock detects the merge block (first block with zero difficulty)
func findMergeBlock(ctx *cli.Context, db ethdb.Database) (uint64, error) {
// Get the current head block
headHash := rawdb.ReadHeadBlockHash(db)
if headHash == (common.Hash{}) {
return 0, fmt.Errorf("chain head not found")
}
headNumber := rawdb.ReadHeaderNumber(db, headHash)
if headNumber == nil {
return 0, fmt.Errorf("header number of head block not found")
}
var mergeBlock uint64
var found bool
// Check if we should use hardcoded merge block values
if ctx.Bool(useHardcodedMergeFlag.Name) {
// Get the network ID to determine which hardcoded value to use
networkName := ctx.String(utils.NetworkIdFlag.Name)
if blockNum, ok := knownMergeBlocks[networkName]; ok {
mergeBlock = blockNum
found = true
log.Info("Using hardcoded merge block", "network", networkName, "block", mergeBlock)
} else {
log.Warn("No hardcoded merge block for network, detecting dynamically", "network", networkName)
}
}
// If no hardcoded value was used or found, detect the merge block
if !found {
// Binary search to find the merge block
var low uint64 = 0
var high uint64 = *headNumber
log.Info("Searching for merge block using binary search", "head", *headNumber)
for low <= high {
mid := (low + high) / 2
header := rawdb.ReadHeader(db, rawdb.ReadCanonicalHash(db, mid), mid)
if header == nil {
return 0, fmt.Errorf("header %d not found", mid)
}
if header.Difficulty.Sign() == 0 {
// This is a post-merge block, look earlier
high = mid - 1
mergeBlock = mid
found = true
} else {
// This is a pre-merge block, look later
low = mid + 1
}
}
if !found {
return 0, fmt.Errorf("merge block not found, chain may not have transitioned to PoS yet")
}
// The merge block is the first block with zero difficulty
log.Info("Found merge block", "number", mergeBlock)
}
return mergeBlock, nil
}

222
cmd/geth/dbcmd_test.go Normal file
View file

@ -0,0 +1,222 @@
// Copyright 2023 The go-ethereum Authors
// This file is part of go-ethereum.
//
// go-ethereum is free software: you can redistribute it and/or modify
// it under the terms of the GNU General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// go-ethereum is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU General Public License for more details.
//
// You should have received a copy of the GNU General Public License
// along with go-ethereum. If not, see <http://www.gnu.org/licenses/>.
package main
import (
"math/big"
"os"
"path/filepath"
"testing"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/rawdb"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/ethdb"
"github.com/ethereum/go-ethereum/params"
)
// TestTruncateFreezerBatching tests the batch processing approach for headers
// in the truncate-freezer command.
func TestTruncateFreezerBatching(t *testing.T) {
// Create a temporary directory for the test
datadir, err := os.MkdirTemp("", "geth-test-freezer-truncation-*")
if err != nil {
t.Fatalf("Failed to create temporary directory: %v", err)
}
defer os.RemoveAll(datadir)
// Create a freezer database
freezerDir := filepath.Join(datadir, "freezer")
db, err := rawdb.NewDatabaseWithFreezer(rawdb.NewMemoryDatabase(), freezerDir, "", false)
if err != nil {
t.Fatalf("Failed to create freezer database: %v", err)
}
defer db.Close()
// Create test data: 100 pre-merge blocks and 10 post-merge blocks
// The merge block will be at index 100
const (
preMergeBlocks = 100
postMergeBlocks = 10
mergeBlock = preMergeBlocks
totalBlocks = preMergeBlocks + postMergeBlocks
)
// Generate and insert blocks
var (
parentHash = common.Hash{}
genesis = types.Header{
Number: big.NewInt(0),
Difficulty: big.NewInt(params.GenesisDifficulty.Int64()),
ParentHash: parentHash,
}
)
// Insert genesis header
rawdb.WriteHeader(db, &genesis)
parentHash = genesis.Hash()
// Insert pre-merge blocks with non-zero difficulty
for i := 1; i <= preMergeBlocks; i++ {
header := types.Header{
Number: big.NewInt(int64(i)),
Difficulty: big.NewInt(params.GenesisDifficulty.Int64()),
ParentHash: parentHash,
}
hash := header.Hash()
rawdb.WriteHeader(db, &header)
rawdb.WriteCanonicalHash(db, hash, uint64(i))
parentHash = hash
}
// Insert post-merge blocks with zero difficulty
for i := preMergeBlocks + 1; i <= totalBlocks; i++ {
header := types.Header{
Number: big.NewInt(int64(i)),
Difficulty: big.NewInt(0), // Zero difficulty = post-merge
ParentHash: parentHash,
}
hash := header.Hash()
rawdb.WriteHeader(db, &header)
rawdb.WriteCanonicalHash(db, hash, uint64(i))
parentHash = hash
}
// Set the head block
rawdb.WriteHeadBlockHash(db, parentHash)
// Create a temporary directory for the headers freezer
tmpDir, err := os.MkdirTemp("", "geth-headers-freezer-*")
if err != nil {
t.Fatalf("Failed to create temporary directory: %v", err)
}
defer os.RemoveAll(tmpDir)
// Create a new freezer for headers and hashes
headersFreezer, err := rawdb.NewFreezer(tmpDir, "headers", false, 2*1000*1000*1000, map[string]bool{
rawdb.ChainFreezerHeaderTable: false,
rawdb.ChainFreezerHashTable: false,
})
if err != nil {
t.Fatalf("Failed to create headers freezer: %v", err)
}
defer headersFreezer.Close()
// Get the ancient reader
ancientDb, ok := db.(ethdb.AncientReader)
if !ok {
t.Fatal("Database doesn't support ancient storage")
}
// Get the number of items in the freezer
ancients, err := ancientDb.Ancients()
if err != nil {
t.Fatalf("Failed to get ancients count: %v", err)
}
t.Logf("Number of items in freezer: %d", ancients)
// Test batch processing by copying headers to the temporary freezer
const batchSize = 10 // Small batch size for testing
for i := uint64(0); i < mergeBlock; i += batchSize {
end := i + batchSize
if end > mergeBlock {
end = mergeBlock
}
// Process this batch
_, err = headersFreezer.ModifyAncients(func(op ethdb.AncientWriteOp) error {
for j := i; j < end; j++ {
// Skip if j is out of range
if j >= ancients {
continue
}
// Read header and hash
headerBytes, err := ancientDb.Ancient(rawdb.ChainFreezerHeaderTable, j)
if err != nil {
return err
}
hashBytes, err := ancientDb.Ancient(rawdb.ChainFreezerHashTable, j)
if err != nil {
return err
}
// Write to temporary freezer
if err := op.AppendRaw(rawdb.ChainFreezerHeaderTable, j, headerBytes); err != nil {
return err
}
if err := op.AppendRaw(rawdb.ChainFreezerHashTable, j, hashBytes); err != nil {
return err
}
}
return nil
})
if err != nil {
t.Fatalf("Failed to copy headers batch: %v", err)
}
}
// Verify that the headers were copied correctly
headersFreezerAncients, err := headersFreezer.Ancients()
if err != nil {
t.Fatalf("Failed to get ancients count from headers freezer: %v", err)
}
// The number of items in the headers freezer should be equal to the merge block
// or the number of items in the original freezer, whichever is smaller
expectedCount := mergeBlock
if ancients < mergeBlock {
expectedCount = ancients
}
if headersFreezerAncients != expectedCount {
t.Fatalf("Expected %d items in headers freezer, got %d", expectedCount, headersFreezerAncients)
}
// Verify that the headers in the temporary freezer match the originals
for i := uint64(0); i < expectedCount; i++ {
// Read from original freezer
originalHeader, err := ancientDb.Ancient(rawdb.ChainFreezerHeaderTable, i)
if err != nil {
t.Fatalf("Failed to read header %d from original freezer: %v", i, err)
}
originalHash, err := ancientDb.Ancient(rawdb.ChainFreezerHashTable, i)
if err != nil {
t.Fatalf("Failed to read hash %d from original freezer: %v", i, err)
}
// Read from temporary freezer
tempHeader, err := headersFreezer.Ancient(rawdb.ChainFreezerHeaderTable, i)
if err != nil {
t.Fatalf("Failed to read header %d from temporary freezer: %v", i, err)
}
tempHash, err := headersFreezer.Ancient(rawdb.ChainFreezerHashTable, i)
if err != nil {
t.Fatalf("Failed to read hash %d from temporary freezer: %v", i, err)
}
// Compare
if string(originalHeader) != string(tempHeader) {
t.Fatalf("Header %d mismatch", i)
}
if string(originalHash) != string(tempHash) {
t.Fatalf("Hash %d mismatch", i)
}
}
t.Log("Batch processing test passed successfully")
}