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
47 changes: 47 additions & 0 deletions cmd/util/cmd/checkpoint-collect-stats/cmd.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package checkpoint_collect_stats
import (
"cmp"
"encoding/hex"
"fmt"
"math"
"slices"
"strings"
Expand Down Expand Up @@ -315,6 +316,15 @@ func getPayloadStatsFromCheckpoint(payloadCallBack func(payload *ledger.Payload)
memAllocBefore := debug.GetHeapAllocsBytes()
log.Info().Msgf("loading checkpoint(s) from %v", flagCheckpointDir)

// checkpoint-collect-stats analyzes payload contents (register types, sizes,
// account info). V7 (payloadless) checkpoints store only leaf hashes and contain
// no payloads, so they cannot be processed here. The WAL replay below loads only
// V6 checkpoints and silently ignores V7 files, which would otherwise produce
// misleading (stale or empty) stats. Fail fast with a clear error instead.
if err := requireV6Checkpoint(flagCheckpointDir); err != nil {
log.Fatal().Err(err).Msg("cannot collect stats from checkpoint")
}

diskWal, err := wal.NewDiskWAL(zerolog.Nop(), nil, &metrics.NoopCollector{}, flagCheckpointDir, complete.DefaultCacheSize, pathfinder.PathByteSize, wal.SegmentSize)
if err != nil {
log.Fatal().Err(err).Msg("cannot create WAL")
Expand Down Expand Up @@ -369,6 +379,43 @@ func getPayloadStatsFromCheckpoint(payloadCallBack func(payload *ledger.Payload)
return ledgerStats
}

// requireV6Checkpoint returns an error if the directory's newest checkpoint is a V7
// (payloadless) checkpoint, i.e. if the newest V7 number is greater than the newest
// V6 number. checkpoint-collect-stats requires full payloads, which V7 checkpoints
// do not contain.
//
// Only numbered checkpoints are considered (the WAL bootstrap loads the latest
// numbered V6 checkpoint). The two versions are compared per version rather than
// via the combined latest, because a payloadless triedir produced by
// checkpoint-convert-v7 holds both checkpoint.N (V6) and checkpoint.N.v7 for the
// same number. Such a directory is accepted: the WAL replay loads the V6 checkpoint
// and the stats are correct. Only a strictly newer V7 checkpoint would make the
// replay silently fall back to an older V6 checkpoint or an empty state, reporting
// misleading stats.
//
// Expected error returns during normal operation:
// - an error when the newest checkpoint in dir is a V7 (payloadless) checkpoint
func requireV6Checkpoint(dir string) error {
_, latestV6, err := wal.ListV6Checkpoints(dir)
if err != nil {
return fmt.Errorf("cannot list V6 checkpoints in %s: %w", dir, err)
}

_, latestV7, err := wal.ListV7Checkpoints(dir)
if err != nil {
return fmt.Errorf("cannot list V7 checkpoints in %s: %w", dir, err)
}

if latestV7 > latestV6 {
return fmt.Errorf(
"checkpoint %d in %s is a V7 (payloadless) checkpoint, which contains no payloads; "+
"checkpoint-collect-stats requires a V6 checkpoint",
latestV7, dir)
}

return nil
}

func getRegisterStats(valueSizesByType sizesByType) []RegisterStatsByTypes {
domainStats := make([]RegisterStatsByTypes, 0, len(common.AllStorageDomains))
var allDomainSizes []float64
Expand Down
89 changes: 89 additions & 0 deletions cmd/util/cmd/checkpoint-collect-stats/cmd_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
package checkpoint_collect_stats

import (
"testing"

"github.com/rs/zerolog"
"github.com/stretchr/testify/require"

"github.com/onflow/flow-go/ledger"
"github.com/onflow/flow-go/ledger/common/testutils"
"github.com/onflow/flow-go/ledger/complete/mtrie/trie"
"github.com/onflow/flow-go/ledger/complete/payloadless"
"github.com/onflow/flow-go/ledger/complete/wal"
)

// TestRequireV6Checkpoint_EmptyDir verifies that a directory without any numbered
// checkpoint is accepted (the caller proceeds with WAL replay / root checkpoint).
func TestRequireV6Checkpoint_EmptyDir(t *testing.T) {
require.NoError(t, requireV6Checkpoint(t.TempDir()))
}

// TestRequireV6Checkpoint_V6 verifies that a directory whose latest checkpoint is
// V6 is accepted.
func TestRequireV6Checkpoint_V6(t *testing.T) {
dir := t.TempDir()
storeV6Checkpoint(t, dir, 1)

require.NoError(t, requireV6Checkpoint(dir))
}

// TestRequireV6Checkpoint_V7 verifies that a directory whose latest checkpoint is
// V7 (payloadless) is rejected, since this command requires full payloads.
func TestRequireV6Checkpoint_V7(t *testing.T) {
dir := t.TempDir()
storeV7Checkpoint(t, dir, 1)

err := requireV6Checkpoint(dir)
require.Error(t, err)
require.Contains(t, err.Error(), "V7")
}

// TestRequireV6Checkpoint_V6AndV7SameNumber verifies that a payloadless triedir
// holding both checkpoint.N (V6) and checkpoint.N.v7 for the same number is
// accepted: the WAL replay loads the V6 checkpoint, so the stats are correct.
func TestRequireV6Checkpoint_V6AndV7SameNumber(t *testing.T) {
dir := t.TempDir()
storeV6Checkpoint(t, dir, 1)
storeV7Checkpoint(t, dir, 1)

require.NoError(t, requireV6Checkpoint(dir))
}

// TestRequireV6Checkpoint_V7NewerThanV6 verifies that a strictly newer V7
// checkpoint is rejected even when older V6 checkpoints exist, since the WAL replay
// would silently fall back to an older V6 checkpoint and report stale stats.
func TestRequireV6Checkpoint_V7NewerThanV6(t *testing.T) {
dir := t.TempDir()
storeV6Checkpoint(t, dir, 1)
storeV7Checkpoint(t, dir, 2)

err := requireV6Checkpoint(dir)
require.Error(t, err)
require.Contains(t, err.Error(), "V7")
}

// storeV6Checkpoint writes a single-trie V6 checkpoint numbered `number` into dir.
func storeV6Checkpoint(t *testing.T, dir string, number int) {
p := testutils.PathByUint8(0)
v := testutils.LightPayload8('A', 'a')
tr, _, err := trie.NewTrieWithUpdatedRegisters(
trie.NewEmptyMTrie(), []ledger.Path{p}, []ledger.Payload{*v}, true)
require.NoError(t, err)

require.NoError(t, wal.StoreCheckpointV6Concurrently(
[]*trie.MTrie{tr}, dir, wal.NumberToFilename(number), zerolog.Nop()))
}

// storeV7Checkpoint writes a single-trie V7 (payloadless) checkpoint numbered
// `number` into dir.
func storeV7Checkpoint(t *testing.T, dir string, number int) {
p := testutils.PathByUint8(0)
v := testutils.LightPayload8('A', 'a')
tr, _, err := payloadless.NewTrieWithUpdatedRegisters(
payloadless.NewEmptyMTrie(), []ledger.Path{p}, [][]byte{v.Value()}, true)
require.NoError(t, err)

require.NoError(t, wal.StoreCheckpointV7Concurrently(
[]*payloadless.MTrie{tr}, dir, wal.NumberToFilenameV7(number), zerolog.Nop()))
}
30 changes: 24 additions & 6 deletions cmd/util/cmd/checkpoint-list-tries/cmd.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,13 @@ package checkpoint_list_tries

import (
"fmt"
"path/filepath"

"github.com/rs/zerolog"
"github.com/rs/zerolog/log"
"github.com/spf13/cobra"

"github.com/onflow/flow-go/ledger"
"github.com/onflow/flow-go/ledger/complete/wal"
)

Expand All @@ -28,14 +31,29 @@ func init() {

func run(*cobra.Command, []string) {

log.Info().Msgf("loading checkpoint %v", flagCheckpoint)
tries, err := wal.LoadCheckpoint(flagCheckpoint, log.Logger)
log.Info().Msgf("reading trie root hashes from checkpoint %v", flagCheckpoint)

hashes, err := readTrieRootHashes(log.Logger, flagCheckpoint)
if err != nil {
log.Fatal().Err(err).Msg("error while loading checkpoint")
log.Fatal().Err(err).Msg("error while reading trie root hashes from checkpoint")
}
log.Info().Msgf("checkpoint loaded, total tries: %v", len(tries))
log.Info().Msgf("checkpoint read, total tries: %v", len(hashes))

for _, trie := range tries {
fmt.Printf("trie root hash: %s\n", trie.RootHash())
for _, h := range hashes {
fmt.Printf("trie root hash: %s\n", h)
}
}

// readTrieRootHashes reads only the trie root hashes from the checkpoint file at
// the given path, without materializing the full trie forest. Only the top-trie
// part file (containing the trie root records) is read.
//
// Both V6 and V7 (payloadless) checkpoints are supported; the version is
// determined by the V7 filename suffix ([wal.V7FileSuffix]). The root hashes are
// returned in the order they are stored in the checkpoint.
//
// No error returns are expected during normal operation.
func readTrieRootHashes(logger zerolog.Logger, checkpointFilePath string) ([]ledger.RootHash, error) {
dir, fileName := filepath.Split(checkpointFilePath)
return wal.ReadCheckpointTriesRootHash(logger, dir, fileName)
}
94 changes: 94 additions & 0 deletions cmd/util/cmd/checkpoint-list-tries/cmd_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
package checkpoint_list_tries

import (
"path/filepath"
"testing"

"github.com/rs/zerolog"
"github.com/stretchr/testify/require"

"github.com/onflow/flow-go/ledger"
"github.com/onflow/flow-go/ledger/common/testutils"
"github.com/onflow/flow-go/ledger/complete/mtrie/trie"
"github.com/onflow/flow-go/ledger/complete/payloadless"
"github.com/onflow/flow-go/ledger/complete/wal"
)

// TestReadTrieRootHashesV6 verifies that the trie root hashes are read from a V6
// checkpoint in the order they were stored, without loading the full forest.
func TestReadTrieRootHashesV6(t *testing.T) {
dir := t.TempDir()
const fileName = "checkpoint"

tries := createV6Tries(t)

err := wal.StoreCheckpointV6Concurrently(tries, dir, fileName, zerolog.Nop())
require.NoError(t, err)

hashes, err := readTrieRootHashes(zerolog.Nop(), filepath.Join(dir, fileName))
require.NoError(t, err)

expected := make([]ledger.RootHash, len(tries))
for i, tr := range tries {
expected[i] = tr.RootHash()
}
require.Equal(t, expected, hashes)
}

// TestReadTrieRootHashesV7 verifies that the trie root hashes are read from a V7
// (payloadless) checkpoint in the order they were stored, by dispatching on the
// V7 filename suffix.
func TestReadTrieRootHashesV7(t *testing.T) {
dir := t.TempDir()
fileName := "checkpoint" + wal.V7FileSuffix

tries := createV7Tries(t)

err := wal.StoreCheckpointV7Concurrently(tries, dir, fileName, zerolog.Nop())
require.NoError(t, err)

hashes, err := readTrieRootHashes(zerolog.Nop(), filepath.Join(dir, fileName))
require.NoError(t, err)

expected := make([]ledger.RootHash, len(tries))
for i, tr := range tries {
expected[i] = tr.RootHash()
}
require.Equal(t, expected, hashes)
}

// createV6Tries builds a chain of two distinct full-payload tries for use as V6
// checkpoint content.
func createV6Tries(t *testing.T) []*trie.MTrie {
p1 := testutils.PathByUint8(0)
v1 := testutils.LightPayload8('A', 'a')
trie1, _, err := trie.NewTrieWithUpdatedRegisters(
trie.NewEmptyMTrie(), []ledger.Path{p1}, []ledger.Payload{*v1}, true)
require.NoError(t, err)

p2 := testutils.PathByUint8(1)
v2 := testutils.LightPayload8('B', 'b')
trie2, _, err := trie.NewTrieWithUpdatedRegisters(
trie1, []ledger.Path{p2}, []ledger.Payload{*v2}, true)
require.NoError(t, err)

return []*trie.MTrie{trie1, trie2}
}

// createV7Tries builds a chain of two distinct payloadless tries for use as V7
// checkpoint content.
func createV7Tries(t *testing.T) []*payloadless.MTrie {
p1 := testutils.PathByUint8(0)
v1 := testutils.LightPayload8('A', 'a')
trie1, _, err := payloadless.NewTrieWithUpdatedRegisters(
payloadless.NewEmptyMTrie(), []ledger.Path{p1}, [][]byte{v1.Value()}, true)
require.NoError(t, err)

p2 := testutils.PathByUint8(1)
v2 := testutils.LightPayload8('B', 'b')
trie2, _, err := payloadless.NewTrieWithUpdatedRegisters(
trie1, []ledger.Path{p2}, [][]byte{v2.Value()}, true)
require.NoError(t, err)

return []*payloadless.MTrie{trie1, trie2}
}
Loading
Loading