diff --git a/app/config_fuzz_test.go b/app/config_fuzz_test.go index fd735504b9..1b8c7617ce 100644 --- a/app/config_fuzz_test.go +++ b/app/config_fuzz_test.go @@ -23,11 +23,12 @@ import ( // // - parseSCConfigs guards almost every read with `if v := opts.Get(k); v != nil`, // so a key absent from an older app.toml keeps its non-zero in-code default. -// - parseSSConfigs guards nothing. Every read is a bare cast of a possibly-nil -// value, so an absent key resolves to the zero value and overwrites the -// default. ss-keep-recent becomes 0 (keep everything, unbounded disk growth), +// - parseSSConfigs leaves every legacy read unguarded. Each is a bare cast of a +// possibly-nil value, so an absent key resolves to the zero value and overwrites +// the default. ss-keep-recent becomes 0 (keep everything, unbounded disk growth), // ss-async-write-buffer becomes 0 (synchronous writes), ss-backend becomes "" -// and ss-enable becomes false. +// and ss-enable becomes false. ss-snapshot-enable is the one guarded read, so an +// app.toml written before SS snapshots existed keeps the in-code default. // // Neither reader returns an error, so nothing about the second case is visible at // boot. It is recorded here as behavior rather than reported as a defect: the @@ -84,8 +85,8 @@ var scKeys = []configtest.KeySpec{ }, } -// ssKeys is the [state-store] read-site manifest. Every row is unguarded and -// unchecked — the section has no presence checks at all. +// ssKeys is the [state-store] read-site manifest. Every row is unchecked; +// SnapshotEnable is guarded while the legacy rows remain unguarded. // // StateStoreConfig also carries KeepLastVersion and UseDefaultComparer, which are // absent here because parseSSConfigs reads neither: they hold their in-code @@ -112,6 +113,10 @@ var ssKeys = []configtest.KeySpec{ {Key: FlagSSImportNumWorkers, Path: "ImportNumWorkers", Cast: configtest.CastInt, Unguarded: true}, {Key: FlagSSDirectory, Path: "DBDirectory", Cast: configtest.CastString, Unguarded: true}, {Key: FlagSSReadWriteMetrics, Path: "EnableReadWriteMetrics", Cast: configtest.CastBool, Unguarded: true}, + { + Key: FlagSSSnapshotEnable, Path: "SnapshotEnable", Cast: configtest.CastBool, + Why: "guarded so app.toml files created before SS snapshots keep the default-off rollout", + }, {Key: FlagEVMSSDirectory, Path: "EVMDBDirectory", Cast: configtest.CastString, Unguarded: true}, {Key: FlagEVMSSSeparateDBs, Path: "SeparateEVMSubDBs", Cast: configtest.CastBool, Unguarded: true}, {Key: FlagEVMSSSplit, Path: "EVMSplit", Cast: configtest.CastBool, Unguarded: true}, @@ -253,9 +258,8 @@ func FuzzParseSCConfigs(f *testing.F) { } // FuzzParseSSConfigs drives every [state-store] key through arbitrary raw values. -// Because the whole section is unguarded, the property being pinned for a nil -// value is the clobber itself: the resolved field must equal the cast's zero, not -// the in-code default. +// For legacy unguarded rows, a nil value must clobber the field to the cast's +// zero. Guarded rows such as SnapshotEnable must retain their in-code default. func FuzzParseSSConfigs(f *testing.F) { seeds := configtest.NewSeeds(f, fuzzing.ConfigValue) @@ -277,7 +281,8 @@ func FuzzParseSSConfigs(f *testing.F) { seeds.AddRow(uint(3), fuzzing.KindInt64, "", int64(200000), false) seeds.AddRow(uint(3), fuzzing.KindNil, "", int64(0), false) // nil clobbers KeepRecent to 0 seeds.AddRow(uint(6), fuzzing.KindString, "/var/lib/sei/ss", int64(0), false) - seeds.AddRow(uint(10), fuzzing.KindBoolString, "", int64(0), true) + seeds.AddRow(uint(8), fuzzing.KindBool, "", int64(0), true) // explicit snapshot opt-in; the default is off + seeds.AddRow(uint(11), fuzzing.KindBoolString, "", int64(0), true) // The clobber cuts both ways for the four rows below. Because the section is unguarded, // an absent key resolves them to their cast's zero, and so does the malformed seed on an @@ -287,7 +292,7 @@ func FuzzParseSSConfigs(f *testing.F) { seeds.AddRow(uint(4), fuzzing.KindInt64, "", int64(1800), false) // prune every 30 min rather than the default 600s seeds.AddRow(uint(5), fuzzing.KindInt64, "", int64(4), false) // four import workers rather than the default 1 seeds.AddRow(uint(7), fuzzing.KindBool, "", int64(0), true) // pebbledb read/write metrics on; the default is off - seeds.AddRow(uint(9), fuzzing.KindBool, "", int64(0), true) // EVM state in its own sub-DBs; the default is shared + seeds.AddRow(uint(10), fuzzing.KindBool, "", int64(0), true) // EVM state in its own sub-DBs; the default is shared configtest.CheckEveryRowHasADiscriminatingSeed(f, "state-store", readSS, ssKeys, seeds) @@ -520,8 +525,9 @@ func TestParseSCConfigsAbsentBaseline(t *testing.T) { // TestParseSSConfigsAbsentBaselineIsZeroClobbered records the clobber in full: an // app.toml with no [state-store] section resolves to a config in which every -// operator-visible knob has been overwritten with a zero value, including the two -// that change the node's disk behavior without any log line. +// unguarded operator-visible knob has been overwritten with a zero value, including +// the two that change the node's disk behavior without any log line. SnapshotEnable +// is the one field that survives, because its read is guarded. func TestParseSSConfigsAbsentBaselineIsZeroClobbered(t *testing.T) { got := parseSSConfigs(configtest.AppOpts{}) @@ -673,6 +679,13 @@ func TestManifestNamesEveryField(t *testing.T) { // manager would otherwise try to map a key onto. "KeepLastVersion", "UseDefaultComparer", + // The three below are unreachable for a different reason, and the distinction is the + // point: they are tagged mapstructure:"-" so no key can bind them even in principle, + // and AlignSSSnapshotWithSC derives all three at runtime from the state-commit cadence. + // ss-snapshot-enable is the only SS-side knob, and it has a row of its own above. + "SnapshotInterval", + "SnapshotKeepRecent", + "SnapshotMinTimeInterval", ) }) t.Run("light_invariance", func(t *testing.T) { diff --git a/app/seidb.go b/app/seidb.go index 9307af2bb0..716583de37 100644 --- a/app/seidb.go +++ b/app/seidb.go @@ -48,6 +48,7 @@ const ( FlagSSPruneInterval = "state-store.ss-prune-interval" FlagSSImportNumWorkers = "state-store.ss-import-num-workers" FlagSSReadWriteMetrics = "state-store.ss-enable-read-write-metrics" + FlagSSSnapshotEnable = "state-store.ss-snapshot-enable" // EVM SS optimization (embedded in SS config, controlled via write/read mode) FlagEVMSSDirectory = "state-store.evm-ss-db-directory" @@ -204,6 +205,12 @@ func parseSSConfigs(appOpts servertypes.AppOptions) config.StateStoreConfig { ssConfig.DBDirectory = cast.ToString(appOpts.Get(FlagSSDirectory)) ssConfig.EnableReadWriteMetrics = cast.ToBool(appOpts.Get(FlagSSReadWriteMetrics)) + // An absent key is an app.toml rendered before SS snapshots existed. Keep + // the in-code default (off) rather than relying on a nil cast. + if v := appOpts.Get(FlagSSSnapshotEnable); v != nil { + ssConfig.SnapshotEnable = cast.ToBool(v) + } + // EVM optimization fields (embedded in SS config) ssConfig.EVMDBDirectory = cast.ToString(appOpts.Get(FlagEVMSSDirectory)) ssConfig.SeparateEVMSubDBs = cast.ToBool(appOpts.Get(FlagEVMSSSeparateDBs)) diff --git a/app/testdata/state-store.golden b/app/testdata/state-store.golden index c60c8b49e2..57d01bd92f 100644 --- a/app/testdata/state-store.golden +++ b/app/testdata/state-store.golden @@ -8,6 +8,10 @@ ImportNumWorkers = int(1) EnableReadWriteMetrics = bool(false) KeepLastVersion = bool(true) UseDefaultComparer = bool(false) +SnapshotEnable = bool(false) +SnapshotInterval = int64(0) +SnapshotKeepRecent = int(0) +SnapshotMinTimeInterval = time.Duration(0s) EVMSplit = bool(false) EVMDBDirectory = string("") SeparateEVMSubDBs = bool(false) diff --git a/app/testdata/state-store.keys.golden b/app/testdata/state-store.keys.golden index fc808ad460..a96612730c 100644 --- a/app/testdata/state-store.keys.golden +++ b/app/testdata/state-store.keys.golden @@ -6,6 +6,7 @@ "state-store.ss-import-num-workers" "state-store.ss-db-directory" "state-store.ss-enable-read-write-metrics" +"state-store.ss-snapshot-enable" "state-store.evm-ss-db-directory" "state-store.evm-ss-separate-dbs" "state-store.evm-ss-split" diff --git a/sei-cosmos/server/config/config.go b/sei-cosmos/server/config/config.go index c46a296125..7bf8b4d70a 100644 --- a/sei-cosmos/server/config/config.go +++ b/sei-cosmos/server/config/config.go @@ -508,6 +508,13 @@ func GetConfig(v *viper.Viper) (Config, error) { memIAVLConfig.SnapshotPrefetchThreshold = v.GetFloat64("state-commit.sc-snapshot-prefetch-threshold") } + // Absent key means an app.toml rendered before SS snapshots existed, which + // should keep the in-code default (off) rather than rely on viper's zero. + ssSnapshotEnable := config.DefaultStateStoreConfig().SnapshotEnable + if v.IsSet("state-store.ss-snapshot-enable") { + ssSnapshotEnable = v.GetBool("state-store.ss-snapshot-enable") + } + // Apply the in-code default when the key is absent so that nodes upgrading // with an older app.toml (which lacks this key) are still bounded rather // than running with unlimited connections. @@ -636,6 +643,7 @@ func GetConfig(v *viper.Viper) (Config, error) { EnableReadWriteMetrics: v.GetBool( "state-store.ss-enable-read-write-metrics", ), + SnapshotEnable: ssSnapshotEnable, EVMSplit: v.GetBool("state-store.evm-ss-split"), EVMDBDirectory: v.GetString("state-store.evm-ss-db-directory"), SeparateEVMSubDBs: v.GetBool("state-store.evm-ss-separate-dbs"), diff --git a/sei-cosmos/server/config/testdata/server_config.golden b/sei-cosmos/server/config/testdata/server_config.golden index a8b4777fa6..442b08b5d9 100644 --- a/sei-cosmos/server/config/testdata/server_config.golden +++ b/sei-cosmos/server/config/testdata/server_config.golden @@ -140,6 +140,10 @@ StateStore.ImportNumWorkers = int(1) StateStore.EnableReadWriteMetrics = bool(false) StateStore.KeepLastVersion = bool(true) StateStore.UseDefaultComparer = bool(false) +StateStore.SnapshotEnable = bool(false) +StateStore.SnapshotInterval = int64(0) +StateStore.SnapshotKeepRecent = int(0) +StateStore.SnapshotMinTimeInterval = time.Duration(0s) StateStore.EVMSplit = bool(false) StateStore.EVMDBDirectory = string("") StateStore.SeparateEVMSubDBs = bool(false) diff --git a/sei-cosmos/storev2/rootmulti/store.go b/sei-cosmos/storev2/rootmulti/store.go index 43b8387347..f7a23cd0a1 100644 --- a/sei-cosmos/storev2/rootmulti/store.go +++ b/sei-cosmos/storev2/rootmulti/store.go @@ -37,6 +37,7 @@ import ( "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/hashlog" sctypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" "github.com/sei-protocol/sei-chain/sei-db/state_db/ss" + sscomposite "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/composite" abci "github.com/sei-protocol/sei-chain/sei-tendermint/abci/types" dbm "github.com/tendermint/tm-db" ) @@ -48,10 +49,24 @@ var ( _ types.Queryable = (*Store)(nil) ) +// stateStoreSnapshotScheduler is the commit path's half of the SS snapshot +// contract: flush tells the state store which version it has just finished +// enqueueing, and the store decides whether that version is a boundary. +type stateStoreSnapshotScheduler interface { + ScheduleSnapshot(version int64) +} + +// ss.NewStateStore returns the interface, so the capability is resolved by type +// assertion at startup. This pins the only implementation, so wrapping the state +// store without carrying the method through fails the build here rather than +// silently ending SS snapshots at runtime. +var _ stateStoreSnapshotScheduler = (*sscomposite.CompositeStateStore)(nil) + type Store struct { mtx sync.RWMutex scStore sctypes.Committer ssStore seidbtypes.StateStore + ssSnapshots stateStoreSnapshotScheduler lastCommitInfo *types.CommitInfo storesParams map[types.StoreKey]storeParams storeKeys map[string]types.StoreKey @@ -136,6 +151,7 @@ func NewStore( scDir: scDir, } if ssConfig.Enable { + config.AlignSSSnapshotWithSC(scConfig, &ssConfig) ssStore, err := ss.NewStateStore(homeDir, ssConfig) if err != nil { panic(err) @@ -150,6 +166,16 @@ func NewStore( panic("Enabling SS store without state sync could cause data corruption") } store.ssStore = ssStore + scheduler, ok := ssStore.(stateStoreSnapshotScheduler) + if !ok { + // Unreachable while CompositeStateStore is the only implementation, + // which the assertion above pins. Log rather than drop silently, so + // a wrapper that loses the method is visible as a boot line instead + // of as snapshots that never appear. + logger.Error("state store does not schedule snapshots; SS snapshots are disabled", + "type", fmt.Sprintf("%T", ssStore)) + } + store.ssSnapshots = scheduler } return store @@ -255,6 +281,14 @@ func (rs *Store) flush() error { telemetry.SetGauge(float32(currentVersion), "storeV2", "ss", "version") } } + // Both branches above have finished handing currentVersion to SS and have + // enqueued nothing above it, which is what makes an SS snapshot label exact. + // Triggering here rather than inside either branch keeps populated and empty + // blocks on one path. A repeat within the same block (flush runs twice, and + // the second pass sees an empty changeset) is ignored by the state store. + if rs.ssSnapshots != nil { + rs.ssSnapshots.ScheduleSnapshot(currentVersion) + } return rs.scStore.ApplyChangeSets(changeSets) } @@ -346,6 +380,9 @@ func (rs *Store) CacheMultiStoreWithVersion(version int64) (types.CacheMultiStor if version <= 0 { version = rs.ssStore.GetLatestVersion() } + if err := rs.validateSSReadVersion(version); err != nil { + return nil, err + } // add the transient/mem stores registered in current app. for k, store := range rs.ckvStores { if store.GetStoreType() != types.StoreTypeIAVL { @@ -364,6 +401,16 @@ func (rs *Store) CacheMultiStoreWithVersion(version int64) (types.CacheMultiStor return cachemulti.NewStore(nil, stores, rs.storeKeys, nil, nil, nil), nil } +// validateSSReadVersion rejects a historical query below the common SS floor +// before constructing stores that would otherwise return partial state. +func (rs *Store) validateSSReadVersion(version int64) error { + earliest := rs.ssStore.GetEarliestVersion() + if version < earliest { + return fmt.Errorf("state store version %d is below earliest available version %d", version, earliest) + } + return nil +} + func (rs *Store) CacheMultiStoreForExport(version int64) (types.CacheMultiStore, error) { if version <= 0 || (rs.lastCommitInfo != nil && version == rs.lastCommitInfo.Version) { return rs.CacheMultiStore(), nil @@ -848,6 +895,9 @@ func (rs *Store) Query(ctx context.Context, req abci.RequestQuery) abci.Response // Fast path: no proof + SS enabled if !needProof && rs.ssStore != nil { + if err := rs.validateSSReadVersion(version); err != nil { + return sdkerrors.QueryResult(errors.Wrap(sdkerrors.ErrInvalidHeight, err.Error())) + } store := types.Queryable(state.NewStore(rs.ssStore, types.NewKVStoreKey(storeName), version)) return store.Query(ctx, req) } diff --git a/sei-cosmos/storev2/rootmulti/store_test.go b/sei-cosmos/storev2/rootmulti/store_test.go index 90b7aa98c9..e026add773 100644 --- a/sei-cosmos/storev2/rootmulti/store_test.go +++ b/sei-cosmos/storev2/rootmulti/store_test.go @@ -3,6 +3,7 @@ package rootmulti import ( "context" "fmt" + "path/filepath" "sync" "testing" "time" @@ -11,6 +12,7 @@ import ( "github.com/sei-protocol/sei-chain/sei-cosmos/store/types" "github.com/sei-protocol/sei-chain/sei-cosmos/storev2/state" "github.com/sei-protocol/sei-chain/sei-db/config" + sscomposite "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/composite" abci "github.com/sei-protocol/sei-chain/sei-tendermint/abci/types" "github.com/stretchr/testify/require" "golang.org/x/time/rate" @@ -133,6 +135,70 @@ func TestSCSS_WriteAndHistoricalRead(t *testing.T) { }) require.EqualValues(t, 0, resp.Code) require.Equal(t, valV1, resp.Value) + + // Once SS reports a higher floor, historical SS queries below it must fail + // before a cache store can mix available Cosmos data with unavailable routed + // data. + require.NoError(t, store.ssStore.SetEarliestVersion(c2.Version, false)) + _, err = store.CacheMultiStoreWithVersion(c1.Version) + require.ErrorContains(t, err, "below earliest available version 2") + + resp = store.Query(context.Background(), abci.RequestQuery{ + Path: "/bank/key", + Data: keyBytes, + Height: c1.Version, + Prove: false, + }) + require.NotEqualValues(t, 0, resp.Code) +} + +// flush owns the SS snapshot trigger for every block, so a boundary must be +// scheduled whether or not the block carried changesets. The composite package +// cannot pin this: its tests call ScheduleSnapshot themselves, so a regression +// in either branch of flush is invisible there. +func TestFlushSchedulesSSSnapshotAtABoundary(t *testing.T) { + for _, tc := range []struct { + name string + writeAtBlock int64 + }{ + {name: "boundary block is populated", writeAtBlock: 2}, + {name: "boundary block is empty", writeAtBlock: 1}, + } { + t.Run(tc.name, func(t *testing.T) { + home := t.TempDir() + scCfg := config.DefaultStateCommitConfig() + scCfg.Enable = true + scCfg.MemIAVLConfig.AsyncCommitBuffer = 0 + // SS mirrors the SC cadence, so this is what puts the SS boundary at 2. + scCfg.MemIAVLConfig.SnapshotInterval = 2 + scCfg.MemIAVLConfig.SnapshotKeepRecent = 1 + + ssCfg := config.DefaultStateStoreConfig() + ssCfg.Enable = true + ssCfg.SnapshotEnable = true + + store := NewStore(home, scCfg, ssCfg, []string{}) + defer func() { _ = store.Close() }() + require.NotNil(t, store.ssSnapshots, "SS snapshot capability was not resolved") + + key := types.NewKVStoreKey("bank") + store.MountStoreWithDB(key, types.StoreTypeIAVL, nil) + require.NoError(t, store.LoadLatestVersion()) + + for block := int64(1); block <= 2; block++ { + if block == tc.writeAtBlock { + store.GetStoreByName("bank").(types.KVStore).Set([]byte("k"), []byte("v")) + } + require.Equal(t, block, store.Commit(true).Version) + } + + root := filepath.Join(home, "data", "state_store", sscomposite.SnapshotsDirName) + require.Eventually(t, func() bool { + versions, err := sscomposite.ListSnapshotVersions(root) + return err == nil && len(versions) == 1 && versions[0] == 2 + }, 10*time.Second, 20*time.Millisecond, "boundary did not produce an SS snapshot") + }) + } } // TestCacheMultiStoreWithVersion_OnlyUsesSSStores verifies that CacheMultiStoreWithVersion diff --git a/sei-db/common/utils/path.go b/sei-db/common/utils/path.go index d8a03fe841..d073091bad 100644 --- a/sei-db/common/utils/path.go +++ b/sei-db/common/utils/path.go @@ -7,6 +7,8 @@ import ( "strings" ) +const StateStoreSnapshotsDirName = "snapshots" + // DirExists returns true if path exists and is a directory. func DirExists(path string) bool { info, err := os.Stat(path) @@ -63,6 +65,11 @@ func GetEVMStateStorePath(homePath string, backend string) string { return filepath.Join(homePath, "data", "state_store", "evm", backend) } +// GetStateStoreSnapshotsPath returns the path for online state-store snapshots. +func GetStateStoreSnapshotsPath(homePath string) string { + return filepath.Join(homePath, "data", "state_store", StateStoreSnapshotsDirName) +} + // GetReceiptStorePath returns the path for the receipt store. // New nodes use data/ledger/receipt/{backend}; existing nodes with // data/receipt.db continue using the legacy path for backward compatibility. diff --git a/sei-db/config/sc_config.go b/sei-db/config/sc_config.go index 48ec3635ba..6b3bda0a54 100644 --- a/sei-db/config/sc_config.go +++ b/sei-db/config/sc_config.go @@ -2,6 +2,7 @@ package config import ( "fmt" + "time" "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/config" "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/memiavl" @@ -16,6 +17,34 @@ const ( legacySCWriteModeCosmosOnly = "cosmos_only" ) +// EffectiveMemIAVLSnapshotCadence resolves memIAVL's snapshot cadence the way +// Options.FillDefaults resolves it at OpenDB, so a caller mirroring the cadence +// onto another backend sees the values memIAVL will actually run with rather +// than the raw config. A zero means "unset" here, not "disabled": memIAVL heals +// it to the default, so mirroring the raw zero would silently disable snapshots +// on the mirroring backend. +func EffectiveMemIAVLSnapshotCadence(cfg memiavl.Config) (interval, keepRecent uint32) { + interval = cfg.SnapshotInterval + if interval == 0 { + interval = memiavl.DefaultSnapshotInterval + } + keepRecent = cfg.SnapshotKeepRecent + if keepRecent == 0 { + keepRecent = memiavl.DefaultSnapshotKeepRecent + } + return interval, keepRecent +} + +// EffectiveMemIAVLSnapshotMinTimeInterval resolves the minimum wall-clock +// interval the same way memIAVL Options.FillDefaults does. +func EffectiveMemIAVLSnapshotMinTimeInterval(cfg memiavl.Config) time.Duration { + seconds := cfg.SnapshotMinTimeInterval + if seconds == 0 { + seconds = memiavl.DefaultSnapshotMinTimeInterval + } + return time.Duration(seconds) * time.Second +} + // StateCommitConfig defines configuration for the state commit (SC) layer. type StateCommitConfig struct { // Enable defines if the state-commit (SeiDB) should be enabled. diff --git a/sei-db/config/ss_config.go b/sei-db/config/ss_config.go index a3a89b898e..567faf9138 100644 --- a/sei-db/config/ss_config.go +++ b/sei-db/config/ss_config.go @@ -1,5 +1,7 @@ package config +import "time" + // DBBackend defines the SS DB backend. type DBBackend string @@ -63,6 +65,38 @@ type StateStoreConfig struct { // defaults to false (use MVCCComparer for backwards compatibility) UseDefaultComparer bool `mapstructure:"use-default-comparer"` + // SnapshotEnable controls whether the state store takes periodic online + // snapshots. Snapshots are Pebble checkpoints (hardlink trees), so the + // backend must be pebbledb and every SS database must be able to hardlink + // into the snapshot root. Startup fails on either rather than running + // without snapshots. A custom Cosmos SS directory moves the snapshot root + // beside that directory, which keeps the link inside one filesystem. + // + // Taking a snapshot occupies each backend's SS apply goroutine for the WAL + // flush, the filesystem sync, and the checkpoint. No data is copied up + // front, but a queue that fills during that window applies write + // backpressure. + // + // Each retained snapshot pins the SSTs it references and prevents compaction + // from reclaiming them. Steady-state disk overhead is therefore the + // compaction churn accumulated over SnapshotInterval blocks, per retained + // snapshot — significant on a multi-TB state store. Managed snapshots are + // rollback restore points, not an archive format. They have no lease in this + // release, so node-external tools must not resolve a snapshot path and open it + // later without first adding a hold mechanism. Attempts, skips, outcomes, + // duration, in-flight state, height, count, and apparent bytes are exported as + // ss_snapshot_* metrics. + // defaults to false + SnapshotEnable bool `mapstructure:"snapshot-enable"` + + // SnapshotInterval, SnapshotKeepRecent, and SnapshotMinTimeInterval are + // mirrored from the state-commit snapshot settings at runtime by + // AlignSSSnapshotWithSC. They are intentionally not exposed in app.toml; + // SnapshotEnable is the only SS-side knob. + SnapshotInterval int64 `mapstructure:"-"` + SnapshotKeepRecent int `mapstructure:"-"` + SnapshotMinTimeInterval time.Duration `mapstructure:"-"` + // --- EVM optimization fields --- // EVMSplit controls whether EVM data is routed to a dedicated SS backend. @@ -93,7 +127,25 @@ func DefaultStateStoreConfig() StateStoreConfig { ImportNumWorkers: DefaultSSImportWorkers, KeepLastVersion: true, UseDefaultComparer: false, + SnapshotEnable: false, EVMSplit: false, SeparateEVMSubDBs: false, } } + +// AlignSSSnapshotWithSC mirrors the state-commit interval, minimum time +// interval, and retention settings onto the state store. SC and SS apply their +// in-flight gates independently, so this does not promise identical retained +// heights. When SS snapshots are disabled the cadence is zeroed. +func AlignSSSnapshotWithSC(scConfig StateCommitConfig, ssConfig *StateStoreConfig) { + if !ssConfig.SnapshotEnable { + ssConfig.SnapshotInterval = 0 + ssConfig.SnapshotKeepRecent = 0 + ssConfig.SnapshotMinTimeInterval = 0 + return + } + interval, keepRecent := EffectiveMemIAVLSnapshotCadence(scConfig.MemIAVLConfig) + ssConfig.SnapshotInterval = int64(interval) + ssConfig.SnapshotKeepRecent = int(keepRecent) + ssConfig.SnapshotMinTimeInterval = EffectiveMemIAVLSnapshotMinTimeInterval(scConfig.MemIAVLConfig) +} diff --git a/sei-db/config/ss_config_test.go b/sei-db/config/ss_config_test.go new file mode 100644 index 0000000000..543f50ba0c --- /dev/null +++ b/sei-db/config/ss_config_test.go @@ -0,0 +1,106 @@ +package config + +import ( + "testing" + "time" + + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/memiavl" + "github.com/stretchr/testify/require" +) + +func TestAlignSSSnapshotWithSC(t *testing.T) { + scConfig := DefaultStateCommitConfig() + ssConfig := DefaultStateStoreConfig() + ssConfig.SnapshotEnable = true + + scConfig.MemIAVLConfig.SnapshotInterval = 123 + scConfig.MemIAVLConfig.SnapshotKeepRecent = 4 + scConfig.MemIAVLConfig.SnapshotMinTimeInterval = 17 + + AlignSSSnapshotWithSC(scConfig, &ssConfig) + + require.Equal(t, int64(123), ssConfig.SnapshotInterval) + require.Equal(t, 4, ssConfig.SnapshotKeepRecent) + require.Equal(t, 17*time.Second, ssConfig.SnapshotMinTimeInterval) +} + +func TestAlignSSSnapshotWithSCHealsZeroToSCDefaults(t *testing.T) { + scConfig := DefaultStateCommitConfig() + ssConfig := DefaultStateStoreConfig() + ssConfig.SnapshotEnable = true + + scConfig.MemIAVLConfig.SnapshotInterval = 0 + scConfig.MemIAVLConfig.SnapshotKeepRecent = 0 + + AlignSSSnapshotWithSC(scConfig, &ssConfig) + + require.Equal(t, int64(memiavl.DefaultSnapshotInterval), ssConfig.SnapshotInterval) + require.Equal(t, memiavl.DefaultSnapshotKeepRecent, ssConfig.SnapshotKeepRecent) + require.Equal( + t, + time.Duration(memiavl.DefaultSnapshotMinTimeInterval)*time.Second, + ssConfig.SnapshotMinTimeInterval, + ) +} + +func TestDefaultStateStoreConfigDisablesSnapshots(t *testing.T) { + require.False(t, DefaultStateStoreConfig().SnapshotEnable, + "snapshots require an explicit ss-snapshot-enable opt-in") +} + +// A zero cadence is what the snapshot manager reads as "do not run", so the +// off switch has to zero it rather than mirror SC's. +func TestAlignSSSnapshotWithSCZeroesCadenceWhenDisabled(t *testing.T) { + scConfig := DefaultStateCommitConfig() + scConfig.MemIAVLConfig.SnapshotInterval = 123 + scConfig.MemIAVLConfig.SnapshotKeepRecent = 4 + scConfig.MemIAVLConfig.SnapshotMinTimeInterval = 17 + + ssConfig := DefaultStateStoreConfig() + ssConfig.SnapshotEnable = false + + AlignSSSnapshotWithSC(scConfig, &ssConfig) + + require.Zero(t, ssConfig.SnapshotInterval) + require.Zero(t, ssConfig.SnapshotKeepRecent) + require.Zero(t, ssConfig.SnapshotMinTimeInterval) +} + +// FlatKV and SS both mirror memIAVL's cadence, and they must resolve it +// identically or the two backends drift onto different snapshot heights. +func TestAlignSSSnapshotMatchesEffectiveMemIAVLCadence(t *testing.T) { + for _, tc := range []struct { + name string + interval, keepRecent uint32 + minTime uint32 + wantInterval int64 + wantMinTime time.Duration + }{ + { + name: "explicit", interval: 500, keepRecent: 3, minTime: 45, + wantInterval: 500, wantMinTime: 45 * time.Second, + }, + { + name: "zero heals to default", interval: 0, keepRecent: 0, + wantInterval: memiavl.DefaultSnapshotInterval, + wantMinTime: time.Duration(memiavl.DefaultSnapshotMinTimeInterval) * time.Second, + }, + } { + t.Run(tc.name, func(t *testing.T) { + scConfig := DefaultStateCommitConfig() + scConfig.MemIAVLConfig.SnapshotInterval = tc.interval + scConfig.MemIAVLConfig.SnapshotKeepRecent = tc.keepRecent + scConfig.MemIAVLConfig.SnapshotMinTimeInterval = tc.minTime + + ssConfig := DefaultStateStoreConfig() + ssConfig.SnapshotEnable = true + AlignSSSnapshotWithSC(scConfig, &ssConfig) + + wantInterval, wantKeepRecent := EffectiveMemIAVLSnapshotCadence(scConfig.MemIAVLConfig) + require.Equal(t, int64(wantInterval), ssConfig.SnapshotInterval) + require.Equal(t, int(wantKeepRecent), ssConfig.SnapshotKeepRecent) + require.Equal(t, tc.wantInterval, ssConfig.SnapshotInterval) + require.Equal(t, tc.wantMinTime, ssConfig.SnapshotMinTimeInterval) + }) + } +} diff --git a/sei-db/config/toml.go b/sei-db/config/toml.go index d7e782bb89..6847fd3650 100644 --- a/sei-db/config/toml.go +++ b/sei-db/config/toml.go @@ -140,6 +140,20 @@ ss-import-num-workers = {{ .StateStore.ImportNumWorkers }} # Applies when ss-backend = "pebbledb". Default: false. ss-enable-read-write-metrics = {{ .StateStore.EnableReadWriteMetrics }} +# SnapshotEnable turns on periodic online state-store snapshots. The cadence is +# not configurable here: it mirrors the state-commit snapshot settings. +# Two configurations fail startup rather than run without snapshots: an +# ss-backend other than "pebbledb", and SS databases that cannot hardlink into +# the snapshot root, which needs every SS database and that root on one +# filesystem. +# Each retained snapshot pins the SST files it references against compaction, so +# budget the write churn of one snapshot interval per retained snapshot. This is +# substantial on a multi-TB state store. +# Snapshot directories are rollback restore points, not an archive format. +# They have no lease, so do not build tools that resolve one and open it later. +# Cost and progress are exported as ss_snapshot_* metrics. Default: false. +ss-snapshot-enable = {{ .StateStore.SnapshotEnable }} + # EVMDBDirectory defines the directory for the optional EVM state-store DB(s). # If unset, defaults to /data/evm_ss when EVM SS is enabled. evm-ss-db-directory = "{{ .StateStore.EVMDBDirectory }}" diff --git a/sei-db/config/toml_test.go b/sei-db/config/toml_test.go index 5ffb274af6..a07fbb05fd 100644 --- a/sei-db/config/toml_test.go +++ b/sei-db/config/toml_test.go @@ -97,6 +97,7 @@ func TestStateStoreConfigTemplate(t *testing.T) { require.Contains(t, output, "ss-prune-interval =", "Missing ss-prune-interval") require.Contains(t, output, "ss-import-num-workers =", "Missing ss-import-num-workers") require.Contains(t, output, "ss-enable-read-write-metrics = false", "Missing state-store read/write metrics flag") + require.Contains(t, output, "ss-snapshot-enable = false", "Missing or incorrect ss-snapshot-enable") require.Contains(t, output, `evm-ss-db-directory = ""`, "Missing evm-ss-db-directory") require.Contains(t, output, `evm-ss-split = false`, "Missing or incorrect evm-ss-split") require.Contains(t, output, "evm-ss-separate-dbs = false", "Missing or incorrect evm-ss-separate-dbs") diff --git a/sei-db/db_engine/pebbledb/mvcc/db.go b/sei-db/db_engine/pebbledb/mvcc/db.go index 114659be57..d4ea4c2f94 100644 --- a/sei-db/db_engine/pebbledb/mvcc/db.go +++ b/sei-db/db_engine/pebbledb/mvcc/db.go @@ -71,7 +71,8 @@ type Database struct { asyncWriteWG sync.WaitGroup config config.StateStoreConfig // Earliest version for db after pruning - earliestVersion atomic.Int64 + earliestVersion atomic.Int64 + earliestVersionMu sync.Mutex // Latest version for db latestVersion atomic.Int64 // descending indicates whether this DB uses descending-version MVCC @@ -85,6 +86,12 @@ type Database struct { // Used in pruning to skip over stores that have not been updated recently storeKeyDirty sync.Map + // pruneIncomplete records that a pass raised the earliest-version marker and + // then failed before it finished deleting. The next pass cannot use that + // marker as its skip baseline: rows below it are still on disk, and a store + // that has gone idle since would be skipped for as long as it stays idle. + pruneIncomplete atomic.Bool + // Changelog used to support async write streamHandler wal.ChangelogWAL @@ -101,12 +108,12 @@ type VersionedChangesets struct { Version int64 Changesets []*proto.NamedChangeSet Done chan struct{} // non-nil for barrier: closed when this entry is processed + // AtDrain, when non-nil, is run by the apply goroutine in queue order + // instead of applying a changeset. See ScheduleAtDrain. + AtDrain func() } -func OpenDB(dataDir string, config config.StateStoreConfig) (types.StateStore, error) { - cache := pebble.NewCache(1024 * 1024 * 32) - defer cache.Unref() - +func newPebbleOptions(config config.StateStoreConfig, cache *pebble.Cache) *pebble.Options { // Select comparer based on config. Note: UseDefaultComparer is NOT backwards compatible // with existing databases created with MVCCComparer - Pebble will refuse to open due to // comparer name mismatch. Only use UseDefaultComparer for NEW databases. @@ -162,6 +169,14 @@ func OpenDB(dataDir string, config config.StateStoreConfig) (types.StateStore, e //TODO: add a new config and check if readonly = true to support readonly mode + return opts +} + +func OpenDB(dataDir string, config config.StateStoreConfig) (types.StateStore, error) { + cache := pebble.NewCache(1024 * 1024 * 32) + defer cache.Unref() + + opts := newPebbleOptions(config, cache) db, err := pebble.Open(dataDir, opts) if err != nil { return nil, fmt.Errorf("failed to open PebbleDB: %w", err) @@ -278,6 +293,51 @@ func (db *Database) PebbleMetrics() *pebble.Metrics { return db.storage.Metrics() } +// Checkpoint writes a point-in-time snapshot of the database into destDir +// (which must not exist yet). Pebble implements this with hardlinks to +// already-fsynced SSTs plus a flushed WAL. SS schedules it on the apply +// goroutine at an ordered queue boundary, so that backend cannot apply more +// changes until the WAL flush, filesystem sync, and checkpoint creation finish. +// Satisfies types.Checkpointable. +func (db *Database) Checkpoint(destDir string) error { + if err := db.storage.Checkpoint(destDir, pebble.WithFlushedWAL()); err != nil { + return fmt.Errorf("pebble checkpoint to %q: %w", destDir, err) + } + return nil +} + +// SetCheckpointVersion writes the logical latest version into a completed +// checkpoint without changing the live database marker. +func (db *Database) SetCheckpointVersion(destDir string, version int64) error { + if version < 0 { + return fmt.Errorf("version must be non-negative") + } + + opts := newPebbleOptions(db.config, nil) + opts.DisableAutomaticCompactions = true + checkpoint, err := pebble.Open(destDir, opts) + if err != nil { + return fmt.Errorf("open checkpoint %q to set markers: %w", destDir, err) + } + + // Converted here, where the non-negative check above is in view. + setErr := setCheckpointMarker(checkpoint, latestVersionKey, uint64(version)) + closeErr := checkpoint.Close() + if setErr != nil { + setErr = fmt.Errorf("set checkpoint version %d: %w", version, setErr) + } + if closeErr != nil { + closeErr = fmt.Errorf("close checkpoint after setting version: %w", closeErr) + } + return errors.Join(setErr, closeErr) +} + +func setCheckpointMarker(checkpoint *pebble.DB, key string, version uint64) error { + var marker [VersionSize]byte + binary.LittleEndian.PutUint64(marker[:], version) + return checkpoint.Set([]byte(key), marker[:], pebble.Sync) +} + func (db *Database) SetLatestVersion(version int64) error { if version < 0 { return fmt.Errorf("version must be non-negative") @@ -337,21 +397,21 @@ func (db *Database) SetEarliestVersion(version int64, ignoreVersion bool) error if version < 0 { return fmt.Errorf("version must be non-negative") } + db.earliestVersionMu.Lock() + defer db.earliestVersionMu.Unlock() + earliestVersion := db.earliestVersion.Load() - if version > earliestVersion || ignoreVersion { - swapped := db.earliestVersion.CompareAndSwap(earliestVersion, version) - if swapped { - var ts [VersionSize]byte - binary.LittleEndian.PutUint64(ts[:], uint64(version)) - err := db.storage.Set([]byte(earliestVersionKey), ts[:], defaultWriteOpts) - if err == nil { - db.operationMetrics.AddWrite(1) - } - return err - } else { - return fmt.Errorf("failed to set earliest version to: %d", version) - } + if version <= earliestVersion && !ignoreVersion { + return nil + } + + var ts [VersionSize]byte + binary.LittleEndian.PutUint64(ts[:], uint64(version)) + if err := db.storage.Set([]byte(earliestVersionKey), ts[:], defaultWriteOpts); err != nil { + return err } + db.earliestVersion.Store(version) + db.operationMetrics.AddWrite(1) return nil } @@ -359,6 +419,18 @@ func (db *Database) GetEarliestVersion() int64 { return db.earliestVersion.Load() } +// advanceEarliestVersion raises the earliest-version marker to target for a +// prune pass that has not deleted anything yet. +// +// SetEarliestVersion serializes competing writers and changes the in-memory +// marker only after Pebble accepts the metadata write. A persistence failure is +// therefore returned with both markers unchanged. Deleting history under a +// marker that only moved in memory would advertise, after a restart, versions +// the pass has already dropped. +func (db *Database) advanceEarliestVersion(target int64) error { + return db.SetEarliestVersion(target, false) +} + // Retrieves earliest version from db, if not found, return 0 func retrieveEarliestVersion(db *pebble.DB) (int64, error) { return retrieveVersionKey(db, earliestVersionKey) @@ -490,6 +562,10 @@ func (db *Database) ApplyChangesetAsync(version int64, changesets []*proto.Named func (db *Database) writeAsyncInBackground() { defer db.asyncWriteWG.Done() for nextChange := range db.pendingChanges { + if nextChange.AtDrain != nil { + nextChange.AtDrain() + continue + } if nextChange.Done != nil { close(nextChange.Done) continue @@ -508,6 +584,19 @@ func (db *Database) WaitForPendingWrites() { <-done } +// ScheduleAtDrain runs fn on the apply goroutine at the point in the queue +// where every changeset enqueued before this call has been applied and none +// enqueued after it has. Unlike WaitForPendingWrites it does not block the +// caller, which is what lets a caller capture the DB at an exact version +// without stalling the block it is committing: the version is pinned by fn's +// position in the queue rather than by when it runs. +// +// fn runs on the writer, so it must not enqueue more work on this DB (that +// deadlocks once the buffer fills) and must not panic. +func (db *Database) ScheduleAtDrain(fn func()) { + db.pendingChanges <- VersionedChangesets{AtDrain: fn} +} + // Prune dispatches between descending- and ascending-mode implementations // depending on the on-disk encoding detected at open time. func (db *Database) Prune(version int64) error { @@ -612,6 +701,11 @@ func (db *Database) getDescending(storeKey string, targetVersion int64, key []by // NOTE: There is a rare case when a module's keys are skipped during pruning even though // it has been updated. This occurs when that module's keys are updated in between pruning runs, the node after is restarted. // This is not a large issue given the next time that module is updated, it will be properly pruned thereafter. +// NOTE: the marker is raised before the deletes, so a pass that fails partway +// leaves rows below the marker on disk. pruneIncomplete makes the next pass in +// the same process rescan every store to reach them. A crash inside that window +// loses the flag, and those rows stay on disk — unreachable by any read, since +// the marker bounds reads too — until the store is written to again. func (db *Database) pruneDescending(version int64) (_err error) { // Defensive check: ensure database is not closed if db.storage == nil { @@ -630,6 +724,22 @@ func (db *Database) pruneDescending(version int64) (_err error) { }() earliestVersion := version + 1 // we increment by 1 to include the provided version + skipBelow := db.GetEarliestVersion() + if err := db.advanceEarliestVersion(earliestVersion); err != nil { + return err + } + if db.pruneIncomplete.Load() { + // A previous pass raised the marker and then stopped short of its + // deletes, so the marker no longer bounds what is on disk. Scan every + // store to reach the rows it left behind. + skipBelow = 0 + } + db.pruneIncomplete.Store(true) + defer func() { + if _err == nil { + db.pruneIncomplete.Store(false) + } + }() itr, err := db.storage.NewIter(nil) if err != nil { @@ -676,8 +786,12 @@ func (db *Database) pruneDescending(version int64) (_err error) { prevStore = storeKey updated, ok := db.storeKeyDirty.Load(storeKey) versionUpdated, typeOk := updated.(int64) - // Skip a store's keys if version it was last updated is less than last prune height - if !ok || (typeOk && versionUpdated < db.GetEarliestVersion()) { + // The marker is advanced before deletes so checkpoints never claim + // history that the prune has already dropped. skipBelow is the marker + // as it stood before this pass raised it; comparing against the raised + // value would skip every store whose latest update is at or below the + // prune height. + if !ok || (typeOk && versionUpdated < skipBelow) { itr.SeekGE(storePrefix(storeKey + "0")) continue } @@ -748,9 +862,6 @@ func (db *Database) pruneDescending(version int64) (_err error) { } db.operationMetrics.AddRead(scanReads) - if err := db.SetEarliestVersion(earliestVersion, false); err != nil { - return err - } return db.compactPrunedRange(firstDeletedKey, lastDeletedKey) } diff --git a/sei-db/db_engine/pebbledb/mvcc/db_ascending.go b/sei-db/db_engine/pebbledb/mvcc/db_ascending.go index 4075f9eea1..6932d06fd2 100644 --- a/sei-db/db_engine/pebbledb/mvcc/db_ascending.go +++ b/sei-db/db_engine/pebbledb/mvcc/db_ascending.go @@ -108,6 +108,22 @@ func (db *Database) pruneAscending(version int64) (_err error) { }() earliestVersion := version + 1 // we increment by 1 to include the provided version + skipBelow := db.GetEarliestVersion() + if err := db.advanceEarliestVersion(earliestVersion); err != nil { + return err + } + if db.pruneIncomplete.Load() { + // A previous pass raised the marker and then stopped short of its + // deletes, so the marker no longer bounds what is on disk. Scan every + // store to reach the rows it left behind. + skipBelow = 0 + } + db.pruneIncomplete.Store(true) + defer func() { + if _err == nil { + db.pruneIncomplete.Store(false) + } + }() itr, err := db.storage.NewIter(nil) if err != nil { @@ -154,8 +170,12 @@ func (db *Database) pruneAscending(version int64) (_err error) { prevStore = storeKey updated, ok := db.storeKeyDirty.Load(storeKey) versionUpdated, typeOk := updated.(int64) - // Skip a store's keys if version it was last updated is less than last prune height - if !ok || (typeOk && versionUpdated < db.GetEarliestVersion()) { + // The marker is advanced before deletes so checkpoints never claim + // history that the prune has already dropped. skipBelow is the marker + // as it stood before this pass raised it; comparing against the raised + // value would skip every store whose latest update is at or below the + // prune height. + if !ok || (typeOk && versionUpdated < skipBelow) { itr.SeekGE(storePrefix(storeKey + "0")) continue } @@ -224,9 +244,6 @@ func (db *Database) pruneAscending(version int64) (_err error) { } db.operationMetrics.AddRead(scanReads) - if err := db.SetEarliestVersion(earliestVersion, false); err != nil { - return err - } return db.compactPrunedRange(firstDeletedKey, lastDeletedKey) } diff --git a/sei-db/db_engine/pebbledb/mvcc/db_test.go b/sei-db/db_engine/pebbledb/mvcc/db_test.go index b923188718..ea914203e5 100644 --- a/sei-db/db_engine/pebbledb/mvcc/db_test.go +++ b/sei-db/db_engine/pebbledb/mvcc/db_test.go @@ -1,13 +1,17 @@ package mvcc import ( + "encoding/binary" + "path/filepath" "testing" + "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" "github.com/sei-protocol/sei-chain/sei-db/config" - "github.com/sei-protocol/sei-chain/sei-db/db_engine/test" + sstest "github.com/sei-protocol/sei-chain/sei-db/db_engine/test" "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" + "github.com/sei-protocol/sei-chain/sei-db/management" ) func TestStorageTestSuite(t *testing.T) { @@ -46,3 +50,57 @@ func TestStorageTestSuiteDefaultComparer(t *testing.T) { suite.Run(t, s) } + +func TestVersionedCheckpointPreservesFutureLiveMarker(t *testing.T) { + cfg := config.DefaultStateStoreConfig() + cfg.Backend = config.PebbleDBBackend + + store, err := OpenDB(filepath.Join(t.TempDir(), "live"), cfg) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + require.NoError(t, store.SetLatestVersion(10)) + require.NoError(t, store.SetEarliestVersion(4, false)) + + dest := filepath.Join(t.TempDir(), "snapshot") + done := make(chan error, 1) + management.ScheduleCheckpoint(store, dest, nil, func(err error) { + done <- err + }) + require.NoError(t, <-done) + // The caller stamps only the label. Earliest is inherited from the + // checkpointed DB because prune advances it before deleting history. + require.NoError(t, management.SetCheckpointVersion(store, dest, 5)) + + require.Equal(t, int64(10), store.GetLatestVersion()) + require.Equal(t, int64(4), store.GetEarliestVersion()) + for key, want := range map[string]uint64{latestVersionKey: 10, earliestVersionKey: 4} { + marker, closer, err := store.(*Database).storage.Get([]byte(key)) + require.NoError(t, err) + require.Equal(t, want, binary.LittleEndian.Uint64(marker), "live %s changed", key) + require.NoError(t, closer.Close()) + } + + checkpoint, err := OpenDB(dest, cfg) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, checkpoint.Close()) }) + require.Equal(t, int64(5), checkpoint.GetLatestVersion()) + require.Equal(t, int64(4), checkpoint.GetEarliestVersion()) +} + +func TestScheduledCheckpointCanBeCanceledAtBarrier(t *testing.T) { + cfg := config.DefaultStateStoreConfig() + cfg.Backend = config.PebbleDBBackend + + store, err := OpenDB(filepath.Join(t.TempDir(), "live"), cfg) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + + dest := filepath.Join(t.TempDir(), "snapshot") + done := make(chan error, 1) + management.ScheduleCheckpoint(store, dest, func() bool { return false }, func(err error) { + done <- err + }) + + require.ErrorIs(t, <-done, management.ErrCheckpointCanceled) + require.NoDirExists(t, dest) +} diff --git a/sei-db/db_engine/pebbledb/mvcc/prune_test.go b/sei-db/db_engine/pebbledb/mvcc/prune_test.go index 6334760534..c7909c426c 100644 --- a/sei-db/db_engine/pebbledb/mvcc/prune_test.go +++ b/sei-db/db_engine/pebbledb/mvcc/prune_test.go @@ -4,6 +4,7 @@ import ( "testing" "github.com/cockroachdb/pebble/v2" + "github.com/cockroachdb/pebble/v2/vfs" "github.com/stretchr/testify/require" "github.com/sei-protocol/sei-chain/sei-db/config" @@ -119,4 +120,91 @@ func TestPruneDescendingOrder_DeletesOldVersions(t *testing.T) { require.ElementsMatch(t, []int64{140}, rawVersionsForKey(t, db, store, k2)) }) + t.Run("idle store still prunes against previous earliest marker", func(t *testing.T) { + db := newTestDB(t, true) + + applyVersion(t, db, store, 50, key, []byte("v50")) + applyVersion(t, db, store, 100, key, []byte("v100")) + + require.NoError(t, db.Prune(150)) + + versions := rawVersionsForKey(t, db, store, key) + require.ElementsMatch(t, []int64{100}, versions, + "prune must not use the just-advanced marker to skip this store") + }) + +} + +func TestPruneAdvancesEarliestBeforeDeletingHistory(t *testing.T) { + db := newTestDB(t, true) + + require.NoError(t, db.storage.Set([]byte("invalid-mvcc-key"), []byte("value"), defaultWriteOpts)) + + err := db.Prune(10) + require.Error(t, err) + require.Equal(t, int64(11), db.GetEarliestVersion(), + "earliest marker must advance before a later prune failure") +} + +// TestPruneAfterFailedPassRescansIdleStores covers the other half of raising the +// marker first: the pass that follows a failure cannot use that marker as its +// skip baseline. store1 goes idle at version 100, below the raised marker, so +// skipping it would leave v50 on disk with no read able to reach it. +func TestPruneAfterFailedPassRescansIdleStores(t *testing.T) { + const store = "store1" + key := []byte("k") + db := newTestDB(t, true) + + applyVersion(t, db, store, 50, key, []byte("v50")) + applyVersion(t, db, store, 100, key, []byte("v100")) + + // "invalid-mvcc-key" sorts ahead of every "s/k:" store key, so the pass + // fails after raising the marker and before deleting anything. + badKey := []byte("invalid-mvcc-key") + require.NoError(t, db.storage.Set(badKey, []byte("value"), defaultWriteOpts)) + require.Error(t, db.Prune(150)) + require.Equal(t, int64(151), db.GetEarliestVersion()) + require.ElementsMatch(t, []int64{50, 100}, rawVersionsForKey(t, db, store, key), + "the failed pass must not have deleted anything") + + require.NoError(t, db.storage.Delete(badKey, defaultWriteOpts)) + require.NoError(t, db.Prune(150)) + + require.ElementsMatch(t, []int64{100}, rawVersionsForKey(t, db, store, key), + "the pass after a failure must rescan a store the raised marker would skip") +} + +// TestAdvanceEarliestVersionAcceptsAHigherMarker pins the outcome a prune pass +// sees when another writer moves the marker past its target. Raising the marker +// now runs ahead of the deletes, so reporting that as a failure would cost the +// whole pass rather than just the marker write. +func TestAdvanceEarliestVersionAcceptsAHigherMarker(t *testing.T) { + db := newTestDB(t, true) + + require.NoError(t, db.SetEarliestVersion(200, false)) + require.NoError(t, db.advanceEarliestVersion(151)) + require.Equal(t, int64(200), db.GetEarliestVersion(), + "the target must not lower a marker another writer raised past it") +} + +// TestAdvanceEarliestVersionReturnsPersistenceFailure pins that Pebble must +// accept the metadata write before the in-memory marker moves. Otherwise a +// later call with the same target would see the target in memory, return nil, +// and let pruning delete history under a marker that was never persisted. +func TestAdvanceEarliestVersionReturnsPersistenceFailure(t *testing.T) { + fs := vfs.NewMem() + storage, err := pebble.Open("db", &pebble.Options{FS: fs}) + require.NoError(t, err) + require.NoError(t, storage.Close()) + + storage, err = pebble.Open("db", &pebble.Options{FS: fs, ReadOnly: true}) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, storage.Close()) }) + + db := &Database{storage: storage} + err = db.advanceEarliestVersion(151) + + require.ErrorIs(t, err, pebble.ErrReadOnly) + require.Zero(t, db.GetEarliestVersion(), + "a failed metadata write must not move the in-memory marker") } diff --git a/sei-db/db_engine/types/types.go b/sei-db/db_engine/types/types.go index 00096bf691..24856dc539 100644 --- a/sei-db/db_engine/types/types.go +++ b/sei-db/db_engine/types/types.go @@ -120,6 +120,30 @@ type Checkpointable interface { Checkpoint(destDir string) error } +// CheckpointVersionSetter writes the logical latest version into a completed +// checkpoint without changing the live database. +// +// The latest marker has to be stamped because a checkpoint is a copy of a +// database whose marker may already have moved on. The earliest marker is +// inherited from the checkpoint; pruning advances it before deleting history, so +// every checkpoint boundary either sees the old marker with old data or the new +// marker with data that is at least as deep as advertised. +type CheckpointVersionSetter interface { + SetCheckpointVersion(destDir string, version int64) error +} + +// DrainBarrier is an optional capability for engines that apply changesets from +// an async queue. It lets a caller place work at an exact point in the write +// order without waiting for the queue to drain. +type DrainBarrier interface { + ScheduleAtDrain(fn func()) +} + +// The three interfaces above are engine capabilities. Deciding when a checkpoint +// runs, and what version it is labeled with, is coordination rather than engine +// behavior and lives in sei-db/management: CheckpointScheduler, +// ScheduleCheckpoint, SetCheckpointVersion and ErrCheckpointCanceled. + // --------------------------------------------------------------------------- // SS DB layer // --------------------------------------------------------------------------- diff --git a/sei-db/management/checkpoint_scheduler.go b/sei-db/management/checkpoint_scheduler.go new file mode 100644 index 0000000000..8e68bb92b9 --- /dev/null +++ b/sei-db/management/checkpoint_scheduler.go @@ -0,0 +1,59 @@ +// Package management holds the coordination layer above the DB engines: work +// that decides when an engine-level operation runs, rather than how the engine +// performs it. +package management + +import ( + "errors" + "fmt" + + "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" +) + +// CheckpointScheduler coordinates checkpoints for stores with in-flight writes. +// +// The engine-side capabilities this builds on — types.Checkpointable, +// types.DrainBarrier and types.CheckpointVersionSetter — stay with the engines +// that implement them. What lives here is the decision of when a checkpoint runs +// and what version it is labeled with. +type CheckpointScheduler interface { + SupportsCheckpoint() bool + ScheduleCheckpoint(destDir string, shouldRun func() bool, done func(error)) + SetCheckpointVersion(destDir string, version int64) error +} + +// ErrCheckpointCanceled reports that a queued checkpoint was canceled before +// it started. +var ErrCheckpointCanceled = errors.New("state store checkpoint canceled") + +// ScheduleCheckpoint checkpoints an engine after all writes already enqueued +// on it have been applied. +func ScheduleCheckpoint(db types.StateStore, destDir string, shouldRun func() bool, done func(error)) { + cp, ok := db.(types.Checkpointable) + if !ok { + done(fmt.Errorf("state store backend %T does not support checkpoints", db)) + return + } + barrier, ok := db.(types.DrainBarrier) + if !ok { + done(fmt.Errorf("state store backend %T does not support ordered checkpoint barriers", db)) + return + } + barrier.ScheduleAtDrain(func() { + if shouldRun != nil && !shouldRun() { + done(ErrCheckpointCanceled) + return + } + done(cp.Checkpoint(destDir)) + }) +} + +// SetCheckpointVersion makes a completed checkpoint self-describing without +// changing the live database. +func SetCheckpointVersion(db types.StateStore, destDir string, version int64) error { + setter, ok := db.(types.CheckpointVersionSetter) + if !ok { + return fmt.Errorf("state store backend %T cannot set checkpoint version", db) + } + return setter.SetCheckpointVersion(destDir, version) +} diff --git a/sei-db/state_db/sc/composite/store.go b/sei-db/state_db/sc/composite/store.go index 0bca128e2d..c82d7a0ad6 100644 --- a/sei-db/state_db/sc/composite/store.go +++ b/sei-db/state_db/sc/composite/store.go @@ -247,14 +247,7 @@ func NewCompositeCommitStore( // different default. Note that mirroring a raw 0 is never correct here (0 means // "disable auto-snapshots" for FlatKV), which is why the zero is resolved first. func alignFlatKVSnapshotWithMemIAVL(cfg *config.StateCommitConfig) { - interval := cfg.MemIAVLConfig.SnapshotInterval - if interval == 0 { - interval = memiavl.DefaultSnapshotInterval - } - keepRecent := cfg.MemIAVLConfig.SnapshotKeepRecent - if keepRecent == 0 { - keepRecent = memiavl.DefaultSnapshotKeepRecent - } + interval, keepRecent := config.EffectiveMemIAVLSnapshotCadence(cfg.MemIAVLConfig) cfg.FlatKVConfig.SnapshotInterval = interval cfg.FlatKVConfig.SnapshotKeepRecent = keepRecent } diff --git a/sei-db/state_db/ss/composite/recovery_test.go b/sei-db/state_db/ss/composite/recovery_test.go index 9688b3afac..c6b0b14fe3 100644 --- a/sei-db/state_db/ss/composite/recovery_test.go +++ b/sei-db/state_db/ss/composite/recovery_test.go @@ -15,6 +15,7 @@ import ( "github.com/sei-protocol/sei-chain/sei-db/wal" evmtypes "github.com/sei-protocol/sei-chain/x/evm/types" "github.com/stretchr/testify/require" + dbm "github.com/tendermint/tm-db" ) func newCompositeStateStoreWithStores( @@ -78,34 +79,97 @@ func TestEVMSSPreRecoveryAfterStateSync(t *testing.T) { require.Contains(t, err.Error(), "EVM SS is empty") } -// TestEVMSSPostRecoveryEarliestMismatch: diverging earliest versions must abort startup. +// TestEVMSSPostRecoveryEarliestMismatch: diverging earliest versions are allowed +// because the composite reports the highest member floor. func TestEVMSSPostRecoveryEarliestMismatch(t *testing.T) { cosmos := &fakeStateStore{latest: 100, earliest: 50} evm := &fakeStateStore{latest: 100, earliest: 75} cs := newCompositeStateStoreWithStores(cosmos, evm, config.StateStoreConfig{EVMSplit: true}) - err := cs.validateEVMSSPostRecovery() - require.Error(t, err) - require.Contains(t, err.Error(), "earliest version") + cs.validateEVMSSPostRecovery() + require.Equal(t, int64(75), cs.GetEarliestVersion()) // Matching earliest → pass. evm.earliest = 50 - require.NoError(t, cs.validateEVMSSPostRecovery()) + cs.validateEVMSSPostRecovery() + require.Equal(t, int64(50), cs.GetEarliestVersion()) // Both zero → pass (fresh DBs). cosmos.earliest = 0 evm.earliest = 0 - require.NoError(t, cs.validateEVMSSPostRecovery()) + cs.validateEVMSSPostRecovery() + require.Zero(t, cs.GetEarliestVersion()) +} + +func TestCompositeGetEarliestVersionReportsHighestMemberFloor(t *testing.T) { + cosmos := &fakeStateStore{latest: 100, earliest: 50} + evm := &fakeStateStore{latest: 100, earliest: 75} + cs := newCompositeStateStoreWithStores(cosmos, evm, config.StateStoreConfig{EVMSplit: true}) + require.Equal(t, int64(75), cs.GetEarliestVersion()) + + cosmos.earliest = 90 + require.Equal(t, int64(90), cs.GetEarliestVersion()) + + cs.evmStore = nil + require.Equal(t, int64(90), cs.GetEarliestVersion()) +} + +// TestCompositeReadsBelowFloorDoNotError pins the read contract that keeps a +// prune racing an in-flight query from crashing the node: the cosmos KVStore +// wrapper panics on any read error, so a version below the reported floor must +// route through and report absence instead of returning an error. +func TestCompositeReadsBelowFloorDoNotError(t *testing.T) { + cosmos := &fakeStateStore{latest: 100, earliest: 50} + evmStore := &fakeStateStore{latest: 100, earliest: 75} + cs := newCompositeStateStoreWithStores(cosmos, evmStore, config.StateStoreConfig{EVMSplit: true}) + require.Equal(t, int64(75), cs.GetEarliestVersion()) + + value, err := cs.Get("bank", 74, []byte("key")) + require.NoError(t, err) + require.Nil(t, value) + + has, err := cs.Has(evm.EVMStoreKey, 74, []byte("key")) + require.NoError(t, err) + require.False(t, has) + + _, err = cs.Iterator("bank", 74, nil, nil) + require.NoError(t, err) + + _, err = cs.ReverseIterator(evm.EVMStoreKey, 74, nil, nil) + require.NoError(t, err) + + require.Equal(t, 4, cosmos.reads+evmStore.reads, "every read must reach its routed member") } -// fakeStateStore stubs latest/earliest for validator tests. +// fakeStateStore stubs latest/earliest and absent reads for validator tests. type fakeStateStore struct { types.StateStore latest, earliest int64 + reads int } func (f *fakeStateStore) GetLatestVersion() int64 { return f.latest } func (f *fakeStateStore) GetEarliestVersion() int64 { return f.earliest } +func (f *fakeStateStore) Get(string, int64, []byte) ([]byte, error) { + f.reads++ + return nil, nil +} + +func (f *fakeStateStore) Has(string, int64, []byte) (bool, error) { + f.reads++ + return false, nil +} + +func (f *fakeStateStore) Iterator(string, int64, []byte, []byte) (dbm.Iterator, error) { + f.reads++ + return nil, nil +} + +func (f *fakeStateStore) ReverseIterator(string, int64, []byte, []byte) (dbm.Iterator, error) { + f.reads++ + return nil, nil +} + func TestRecoverCompositeStateStore(t *testing.T) { dir, err := os.MkdirTemp("", "composite_recovery_test") require.NoError(t, err) diff --git a/sei-db/state_db/ss/composite/snapshot.go b/sei-db/state_db/ss/composite/snapshot.go new file mode 100644 index 0000000000..f51e4712f5 --- /dev/null +++ b/sei-db/state_db/ss/composite/snapshot.go @@ -0,0 +1,761 @@ +package composite + +import ( + "context" + "errors" + "fmt" + "io/fs" + "os" + "path/filepath" + "slices" + "strconv" + "strings" + "sync" + "time" + + "github.com/sei-protocol/sei-chain/sei-db/common/utils" + "github.com/sei-protocol/sei-chain/sei-db/management" +) + +// Online state-store snapshots. Every SnapshotInterval blocks the store takes a +// Pebble checkpoint of each backend while the node keeps producing blocks. +// Checkpoints are hardlink trees, so they do not copy database contents, but +// each one occupies its backend's apply goroutine for the full checkpoint +// operation. Writes continue to enter the bounded queue, but a full queue +// applies backpressure until the checkpoint finishes. The result is an +// immutable, crash-consistent image of the query store. +// +// These snapshots are an input to SS rollback, not an export format. The +// intended restore model matches SC FlatKV: restore from an SS snapshot, then +// replay the state WAL forward to the target height. State sync imports the SC +// snapshot stream and rebuilds SS from that stream; it does not consume these +// SS snapshot directories. +// +// On-disk layout under the snapshot root. By default the root is +// /data/state_store/snapshots. A custom Cosmos SS directory moves it +// to the sibling -snapshots directory so Pebble can use hardlinks. +// +// snapshots/ +// current -> snapshot-NNNNN (symlink to newest snapshot) +// snapshot-NNNNN/ (immutable; NNNNN = label version) +// cosmos// (Pebble checkpoint of Cosmos SS) +// evm// (Pebble checkpoint of EVM SS, if split) +// / (when EVM sub-DBs are separate) +// +// Snapshots are eligible at the same interval boundaries and minimum time +// cadence as state commit. This composite implementation uses one trigger and +// one current link for all member stores, so its member snapshots share a label. +// That same label is a property of this layout, not a rollback requirement: +// rollback can replay the state WAL from each store's own nearest snapshot. For +// every accepted SS snapshot, the label is exact: it is the version the write +// path had just handed to the backends when the snapshot was requested. Placing +// a barrier in each backend's apply queue — rather than sampling what the +// backends had applied — makes that label exact without the request having to +// wait. See requestSnapshot. +// +// The barrier orders only the async block-commit queues. Import, recovery, +// pruning, and direct version-marker writes bypass those queues and must not +// call ScheduleSnapshot. The rootmulti commit path owns the trigger for every +// block, populated or empty, and is the only caller of ScheduleSnapshot. +// +// Pruning is the one writer nothing orders a snapshot against. A checkpoint can +// capture a partially applied prune — the same state a crash mid-prune leaves on +// the live DB. This is safe for a snapshot because pruning advances each DB's +// earliest marker before deleting history, so the checkpoint never claims a +// range the DB has already dropped. Reopening a snapshot with different member +// floors is allowed; the composite reports the highest floor any member carries. +// +// SS rollback is not implemented in this feature. When it is added, it should +// use these snapshots the same way SC FlatKV does: restore from a snapshot +// boundary, then replay the state WAL forward. Until then, rolling back or +// state-syncing to a lower height in a reused home directory leaves two stale +// facts behind: lastRequested still carries the old high-water mark, so repeated +// boundaries can be skipped, and already published snapshot-NNNNN directories +// keep labels from the abandoned chain. Clear the snapshot root by hand in that +// case. +// +// Managed snapshot directories have no lease because they are not a node-external +// consumption API. Retention may remove any snapshot that rollback does not need. +// If a future tool opens or copies these directories directly, it must first add +// a lease or other hold mechanism. +// +// This file is the layer the planned per-SS restructure has to move. The +// lifecycle here — layout, retention, the current symlink, staging and +// publication, restart recovery — is reachable only as a method on +// *CompositeStateStore, and startSnapshotManager requires a checkpointable +// Cosmos store, so an EVM-only store cannot use it as written. The agreed +// direction is for each SS to own its own snapshot root, current link, creation, +// and retention behind gc.PrunableStore, mirroring SC FlatKV. Composite mode can +// then fan out to Cosmos SS and EVM SS, while Giga can use EVM SS directly. That +// also removes the second retention path this file adds: prune here is +// count-based and has no ExternalPruning stand-down, so pointing +// StorageGarbageCollector at SS before then would give a store two independent +// pruners. The shape waits on rollback not because GetRollbackFloor is unknown; +// SC FlatKV already defines that floor. It waits because gc.PrunableStore must +// not report a floor above what the store can actually restore to. The current +// link semantics should change with that work too: this implementation points to +// the newest published snapshot, while FlatKV's current link points to the +// active snapshot that open/rollback clones and replays from. + +const ( + // SnapshotsDirName is the directory under data/state_store that holds + // online snapshots. + SnapshotsDirName = utils.StateStoreSnapshotsDirName + + snapshotPrefix = "snapshot-" + // snapshotDirLen is "snapshot-" + 20-digit zero-padded version. + snapshotDirLen = len(snapshotPrefix) + 20 + + snapshotCurrentLink = "current" + snapshotCurrentTmpLink = "current-tmp" + snapshotTmpPrefix = "tmp-" + snapshotSizeFile = ".apparent-size" +) + +// SnapshotDirName returns the directory name for a snapshot labeled with the +// given version. +func SnapshotDirName(version int64) string { + return fmt.Sprintf("%s%020d", snapshotPrefix, version) +} + +// ParseSnapshotVersion parses a snapshot directory name; ok is false for +// anything that is not a snapshot-<20 digits> name. +func ParseSnapshotVersion(name string) (version int64, ok bool) { + if !strings.HasPrefix(name, snapshotPrefix) || len(name) != snapshotDirLen { + return 0, false + } + v, err := strconv.ParseInt(name[len(snapshotPrefix):], 10, 64) + if err != nil || v < 0 { + return 0, false + } + return v, true +} + +// ListSnapshotVersions returns the labels of all snapshots under root in +// ascending order. A missing root is not an error (no snapshots yet). +func ListSnapshotVersions(root string) ([]int64, error) { + entries, err := os.ReadDir(root) + if err != nil { + if os.IsNotExist(err) { + return nil, nil + } + return nil, fmt.Errorf("read snapshots dir %q: %w", root, err) + } + var versions []int64 + for _, entry := range entries { + if !entry.IsDir() { + continue + } + if v, ok := ParseSnapshotVersion(entry.Name()); ok { + versions = append(versions, v) + } + } + slices.Sort(versions) + return versions, nil +} + +// snapshotManager owns the snapshots directory and the one-at-a-time discipline +// for filling it. It has no goroutine of its own: snapshots are requested from +// the write path and completed on the backends' apply goroutines. +type snapshotManager struct { + root string + backend string + interval int64 + keepRecent int + minTime time.Duration + + cosmosScheduler management.CheckpointScheduler + evmScheduler management.CheckpointScheduler + snapshotSizes map[int64]int64 + + mu sync.Mutex + // lastRequested is the newest label already requested or on disk, so a + // boundary is not snapshotted twice across a restart or a re-sent version. + lastRequested int64 + lastRequestAt time.Time + inFlight bool + stopped bool + // scheduling closes the gap between accepting a request and enqueueing its + // barriers. Close waits for it before closing backend queues. + scheduling sync.WaitGroup + // publishing tracks the goroutine finishing the accepted snapshot off. + publishing sync.WaitGroup + + // publishMu serializes the publish step, which reads and rewrites the + // shared directory (the current link, and pruning). + publishMu sync.Mutex + lastPublished int64 +} + +type checkpointTarget struct { + store management.CheckpointScheduler + dest string +} + +// startSnapshotManager wires the manager into the composite store. Snapshot +// enablement is fail-closed: every backend must support checkpoints, and every +// live DB must be able to hardlink into root. Pebble otherwise silently falls +// back to copying SSTs across filesystems while its apply worker is blocked. +func (s *CompositeStateStore) startSnapshotManager(root string, sourceDirs []string) error { + if s.config.SnapshotInterval <= 0 { + return nil + } + cosmosScheduler, ok := s.cosmosStore.(management.CheckpointScheduler) + if !ok || !cosmosScheduler.SupportsCheckpoint() { + return fmt.Errorf("cosmos backend %q does not support checkpoints", s.config.Backend) + } + var evmScheduler management.CheckpointScheduler + if s.evmStore != nil { + evmScheduler, ok = s.evmStore.(management.CheckpointScheduler) + if !ok || !evmScheduler.SupportsCheckpoint() { + return fmt.Errorf("EVM backend %q does not support checkpoints", s.config.Backend) + } + } + if err := verifySnapshotHardlinks(root, sourceDirs); err != nil { + return err + } + m := &snapshotManager{ + root: root, + backend: s.config.Backend, + interval: s.config.SnapshotInterval, + keepRecent: s.config.SnapshotKeepRecent, + minTime: s.config.SnapshotMinTimeInterval, + cosmosScheduler: cosmosScheduler, + evmScheduler: evmScheduler, + snapshotSizes: map[int64]int64{}, + } + m.lastRequested = m.newestSnapshotVersion() + m.lastPublished = m.lastRequested + m.lastRequestAt = m.snapshotModTime(m.lastRequested) + m.removeStaleTmpDirs() + m.prune() + if m.lastPublished > 0 { + if err := m.updateCurrentLink(SnapshotDirName(m.lastPublished)); err != nil { + logger.Error("failed to restore state store snapshot current link", + "version", m.lastPublished, "error", err) + } + snapshotMetrics.CurrentHeight.Record(context.Background(), m.lastPublished) + } + s.snapshotMgr = m + logger.Info("state store snapshotting enabled", + "root", root, + "interval", m.interval, + "minTimeInterval", m.minTime, + "keepRecent", m.keepRecent, + ) + return nil +} + +func verifySnapshotHardlinks(root string, sourceDirs []string) error { + if err := os.MkdirAll(root, 0o750); err != nil { + return fmt.Errorf("create snapshot root %q: %w", root, err) + } + for _, sourceDir := range sourceDirs { + probe, err := os.CreateTemp(sourceDir, ".ss-snapshot-link-probe-*") + if err != nil { + return fmt.Errorf("create hardlink probe in state store %q: %w", sourceDir, err) + } + source := probe.Name() + if err := probe.Close(); err != nil { + _ = os.Remove(source) + return fmt.Errorf("close hardlink probe in state store %q: %w", sourceDir, err) + } + target := filepath.Join(root, filepath.Base(source)) + if err := os.Link(source, target); err != nil { + _ = os.Remove(source) + return fmt.Errorf( + "state store %q cannot hardlink snapshots into %q; place all SS databases and the snapshot root on one filesystem: %w", + sourceDir, + root, + err, + ) + } + if err := os.Remove(source); err != nil { + _ = os.Remove(target) + return fmt.Errorf("remove hardlink probe %q: %w", source, err) + } + if err := os.Remove(target); err != nil { + return fmt.Errorf("remove hardlink probe %q: %w", target, err) + } + } + return nil +} + +// stop prevents further snapshots, waits for accepted requests to enqueue their +// barriers, and then waits for active publication. Queued barriers are canceled +// before they start when backend close drains their queues. +func (m *snapshotManager) stop() { + m.mu.Lock() + m.stopped = true + m.mu.Unlock() + m.scheduling.Wait() + m.publishing.Wait() +} + +func (m *snapshotManager) isRunning() bool { + m.mu.Lock() + defer m.mu.Unlock() + return !m.stopped +} + +// maybeSnapshot takes a snapshot when version lands on an interval boundary. +// It is called from the write path for every version, so the common case is the +// modulo test and nothing else. +func (m *snapshotManager) maybeSnapshot(version int64) { + if m == nil || version <= 0 || m.interval <= 0 || version%m.interval != 0 { + return + } + now := time.Now() + m.mu.Lock() + previous := m.lastRequested + previousRequestAt := m.lastRequestAt + var skipReason string + accepted := false + switch { + case m.stopped || version <= m.lastRequested: + // A repeated commit-path call is expected and is not a skipped attempt. + case m.inFlight: + skipReason = "in_flight" + case !m.lastRequestAt.IsZero() && now.Sub(m.lastRequestAt) < m.minTime: + skipReason = "minimum_time_interval" + default: + m.lastRequested = version + m.lastRequestAt = now + m.inFlight = true + m.scheduling.Add(1) + recordSnapshotInFlight(1) + accepted = true + } + m.mu.Unlock() + if !accepted { + if skipReason != "" { + recordSnapshotSkipped(skipReason) + // A skipped boundary is the reason a snapshot an operator expected is + // not on disk, so name the gate rather than leaving only a metric. + logger.Info("skipping state store snapshot", "version", version, "reason", skipReason) + } + return + } + defer m.scheduling.Done() + start := time.Now() + recordSnapshotAttempt() + if err := m.requestSnapshot(version, start); err != nil { + recordSnapshotCompletion(start, "failure") + m.mu.Lock() + if m.lastRequested == version { + m.lastRequested = previous + m.lastRequestAt = previousRequestAt + m.inFlight = false + recordSnapshotInFlight(0) + } + m.mu.Unlock() + logger.Error("state store snapshot failed", "version", version, "error", err) + } +} + +func (m *snapshotManager) finishSnapshot() { + m.mu.Lock() + m.inFlight = false + recordSnapshotInFlight(0) + m.mu.Unlock() +} + +// requestSnapshot asks every backend to checkpoint itself into a staging +// directory and publishes the result once they all have. +// +// The label is exact because of when this runs: the caller has just enqueued +// version on the backends and has not enqueued anything above it, so a barrier +// placed in each apply queue now captures that backend with everything up to +// version applied and nothing after it. The backends reach their barriers +// independently and at different wall-clock times, and the caller waits for none +// of the checkpointing — enqueueing a barrier costs what enqueueing a changeset +// costs. The caller does wait for the staging directories below: one Stat, one +// RemoveAll and one MkdirAll per target, on the commit path and ahead of the SC +// apply. +func (m *snapshotManager) requestSnapshot(version int64, start time.Time) error { + name := SnapshotDirName(version) + finalDir := filepath.Join(m.root, name) + if _, err := os.Stat(finalDir); err == nil { + return fmt.Errorf("snapshot dir %q already exists", finalDir) + } else if !os.IsNotExist(err) { + return fmt.Errorf("inspect snapshot dir %q: %w", finalDir, err) + } + tmpDir := filepath.Join(m.root, snapshotTmpPrefix+name) + if err := os.RemoveAll(tmpDir); err != nil { + return fmt.Errorf("clear stale snapshot tmp dir: %w", err) + } + + targets := []checkpointTarget{ + {m.cosmosScheduler, filepath.Join(tmpDir, "cosmos", m.backend)}, + } + if m.evmScheduler != nil { + targets = append(targets, checkpointTarget{ + m.evmScheduler, + filepath.Join(tmpDir, "evm", m.backend), + }) + } + for _, target := range targets { + if err := os.MkdirAll(filepath.Dir(target.dest), 0o750); err != nil { + _ = os.RemoveAll(tmpDir) + return fmt.Errorf("create snapshot dir: %w", err) + } + } + + var ( + mu sync.Mutex + remaining = len(targets) + firstErr error + ) + // Set up before scheduling because callbacks can complete while the loop is + // still scheduling the remaining targets. + for _, target := range targets { + target.store.ScheduleCheckpoint(target.dest, m.isRunning, func(err error) { + mu.Lock() + if err != nil && firstErr == nil { + firstErr = err + } + remaining-- + last, outcome := remaining == 0, firstErr + mu.Unlock() + if !last { + return + } + m.startPublish(version, tmpDir, finalDir, targets, outcome, start) + }) + } + return nil +} + +// startPublish hands a finished set of checkpoints off to a goroutine. It runs +// on whichever backend's apply goroutine finished last, so it must not do the +// work itself: publishing renames directories and prunes old snapshots, and a +// writer stalled on that is a writer not applying blocks. +func (m *snapshotManager) startPublish( + version int64, + tmpDir, finalDir string, + targets []checkpointTarget, + checkpointErr error, + start time.Time, +) { + // Taken under the same lock stop uses, so no goroutine is registered after + // stop has started waiting. + m.mu.Lock() + if m.stopped { + m.mu.Unlock() + _ = os.RemoveAll(tmpDir) + recordSnapshotCompletion(start, "canceled") + m.finishSnapshot() + return + } + m.publishing.Add(1) + m.mu.Unlock() + + go func() { + defer m.publishing.Done() + defer m.finishSnapshot() + if checkpointErr != nil { + if errors.Is(checkpointErr, management.ErrCheckpointCanceled) { + recordSnapshotCompletion(start, "canceled") + } else { + recordSnapshotCompletion(start, "failure") + logger.Error("state store snapshot failed", "version", version, "error", checkpointErr) + } + _ = os.RemoveAll(tmpDir) + return + } + for _, target := range targets { + if err := target.store.SetCheckpointVersion(target.dest, version); err != nil { + recordSnapshotCompletion(start, "failure") + logger.Error("failed to set state store snapshot version", + "version", version, "dir", target.dest, "error", err) + _ = os.RemoveAll(tmpDir) + return + } + } + if m.publish(version, tmpDir, finalDir, start) { + recordSnapshotCompletion(start, "success") + } else { + recordSnapshotCompletion(start, "failure") + } + }() +} + +// publish moves a finished checkpoint into place and reports whether the whole +// publication succeeded. Retention runs either way. +// +// A boundary that fails anywhere past the barrier is given up on, and this is +// deliberate. maybeSnapshot restores lastRequested when requestSnapshot fails, +// because that failure happens before any barrier is enqueued and the boundary +// was never claimed. Once the barriers are out, the version they captured is +// the only image of that boundary there will ever be: the write path has moved +// on, so re-running the attempt would checkpoint a later state under the older +// label, which is the one thing the label is supposed to rule out. Recovery is +// therefore the next boundary rather than a retry of this one, at the cost of +// one snapshot interval of coverage. The error log and the outcome="failure" +// counter are the signal. +func (m *snapshotManager) publish(version int64, tmpDir, finalDir string, start time.Time) bool { + apparentBytes, sizeErr := snapshotDirApparentBytes(tmpDir) + if sizeErr != nil { + logger.Error("failed to measure state store snapshot", "dir", tmpDir, "error", sizeErr) + } else if err := writeSnapshotSize(tmpDir, apparentBytes); err != nil { + logger.Error("failed to persist state store snapshot size", "dir", tmpDir, "error", err) + sizeErr = err + } + + m.publishMu.Lock() + defer m.publishMu.Unlock() + defer m.prune() + + if err := os.Rename(tmpDir, finalDir); err != nil { + logger.Error("failed to finalize state store snapshot", "version", version, "error", err) + _ = os.RemoveAll(tmpDir) + return false + } + if err := syncDir(m.root); err != nil { + logger.Error("failed to persist state store snapshot publication", + "version", version, "dir", finalDir, "error", err) + return false + } + if sizeErr == nil { + if m.snapshotSizes == nil { + m.snapshotSizes = map[int64]int64{} + } + m.snapshotSizes[version] = apparentBytes + } + logger.Info("state store snapshot created", + "version", version, "dir", finalDir, "took", time.Since(start).String()) + + // Snapshots can finish out of order, so only move the link forward. + if version > m.lastPublished { + if err := m.updateCurrentLink(SnapshotDirName(version)); err != nil { + // The snapshot itself is intact and discoverable by name; only the + // convenience symlink is stale. The link is part of the publication + // contract, so record this attempt as a failure. + logger.Error("failed to update state store snapshot current link", + "version", version, "error", err) + return false + } + m.lastPublished = version + } + snapshotMetrics.CurrentHeight.Record(context.Background(), m.lastPublished) + return true +} + +func (m *snapshotManager) newestSnapshotVersion() int64 { + versions, err := ListSnapshotVersions(m.root) + if err != nil { + logger.Error("failed to list state store snapshots", "error", err) + return 0 + } + if len(versions) == 0 { + return 0 + } + return versions[len(versions)-1] +} + +func (m *snapshotManager) snapshotModTime(version int64) time.Time { + if version <= 0 { + return time.Time{} + } + info, err := os.Stat(filepath.Join(m.root, SnapshotDirName(version))) + if err != nil { + logger.Error("failed to read state store snapshot modification time", + "version", version, "error", err) + return time.Time{} + } + return info.ModTime() +} + +// removeStaleTmpDirs clears staging directories left behind by a crash or a +// shutdown that landed mid-snapshot. They are named after the snapshot they +// were staging, so they would otherwise sit there until that exact boundary +// came round again. +func (m *snapshotManager) removeStaleTmpDirs() { + tmpLink := filepath.Join(m.root, snapshotCurrentTmpLink) + if err := os.Remove(tmpLink); err != nil && !os.IsNotExist(err) { + logger.Error("failed to remove stale state store snapshot link", "path", tmpLink, "error", err) + } + + entries, err := os.ReadDir(m.root) + if err != nil { + if !os.IsNotExist(err) { + logger.Error("failed to scan state store snapshots dir", "error", err) + } + return + } + for _, entry := range entries { + if !entry.IsDir() || !strings.HasPrefix(entry.Name(), snapshotTmpPrefix) { + continue + } + dir := filepath.Join(m.root, entry.Name()) + if err := os.RemoveAll(dir); err != nil { + logger.Error("failed to remove stale snapshot tmp dir", "dir", dir, "error", err) + continue + } + logger.Info("removed stale state store snapshot tmp dir", "dir", dir) + } +} + +// updateCurrentLink atomically points the current symlink at name. +func (m *snapshotManager) updateCurrentLink(name string) error { + tmpLink := filepath.Join(m.root, snapshotCurrentTmpLink) + _ = os.Remove(tmpLink) + if err := os.Symlink(name, tmpLink); err != nil { + return fmt.Errorf("create snapshot current symlink: %w", err) + } + if err := os.Rename(tmpLink, filepath.Join(m.root, snapshotCurrentLink)); err != nil { + return fmt.Errorf("swap snapshot current symlink: %w", err) + } + return syncDir(m.root) +} + +func syncDir(path string) error { + // #nosec G304 -- path is an internal database or snapshot directory, not request input. + dir, err := os.Open(path) + if err != nil { + return fmt.Errorf("open directory %q for sync: %w", path, err) + } + syncErr := dir.Sync() + closeErr := dir.Close() + if syncErr != nil { + syncErr = fmt.Errorf("sync directory %q: %w", path, syncErr) + } + if closeErr != nil { + closeErr = fmt.Errorf("close directory %q after sync: %w", path, closeErr) + } + return errors.Join(syncErr, closeErr) +} + +// prune removes all but the newest 1+keepRecent snapshots. +func (m *snapshotManager) prune() { + versions, err := ListSnapshotVersions(m.root) + if err != nil { + logger.Error("failed to list state store snapshots for pruning", "error", err) + return + } + defer m.recordRetentionMetrics() + currentVersion, hasCurrent, err := m.currentSnapshotVersion() + if err != nil { + logger.Error("failed to resolve current state store snapshot before pruning", "error", err) + return + } + keep := 1 + m.keepRecent + if len(versions) <= keep { + return + } + for _, v := range versions[:len(versions)-keep] { + if hasCurrent && v == currentVersion { + continue + } + dir := filepath.Join(m.root, SnapshotDirName(v)) + if err := os.RemoveAll(dir); err != nil { + logger.Error("failed to prune state store snapshot", "dir", dir, "error", err) + continue + } + logger.Info("pruned state store snapshot", "dir", dir) + } +} + +func (m *snapshotManager) currentSnapshotVersion() (version int64, exists bool, err error) { + target, err := os.Readlink(filepath.Join(m.root, snapshotCurrentLink)) + if err != nil { + if os.IsNotExist(err) { + return 0, false, nil + } + return 0, false, fmt.Errorf("read current snapshot link: %w", err) + } + version, ok := ParseSnapshotVersion(filepath.Base(target)) + if !ok { + return 0, false, fmt.Errorf("current snapshot link has invalid target %q", target) + } + return version, true, nil +} + +func (m *snapshotManager) recordRetentionMetrics() { + versions, err := ListSnapshotVersions(m.root) + if err != nil { + logger.Error("failed to list state store snapshots for metrics", "error", err) + return + } + snapshotMetrics.RetainedCount.Record(context.Background(), int64(len(versions))) + + if m.snapshotSizes == nil { + m.snapshotSizes = map[int64]int64{} + } + retained := make(map[int64]struct{}, len(versions)) + var apparentBytes int64 + for _, version := range versions { + retained[version] = struct{}{} + if size, ok := m.snapshotSizes[version]; ok { + apparentBytes += size + continue + } + dir := filepath.Join(m.root, SnapshotDirName(version)) + size, err := readSnapshotSize(dir) + if err != nil { + if !os.IsNotExist(err) { + logger.Error("failed to read state store snapshot size", "dir", dir, "error", err) + } + continue + } + m.snapshotSizes[version] = size + apparentBytes += size + } + for version := range m.snapshotSizes { + if _, ok := retained[version]; !ok { + delete(m.snapshotSizes, version) + } + } + snapshotMetrics.ApparentBytes.Record(context.Background(), apparentBytes) +} + +func snapshotDirApparentBytes(dir string) (int64, error) { + var apparentBytes int64 + err := filepath.WalkDir(dir, func(_ string, entry fs.DirEntry, err error) error { + if err != nil { + return err + } + if !entry.Type().IsRegular() { + return nil + } + info, err := entry.Info() + if err != nil { + return err + } + apparentBytes += info.Size() + return nil + }) + return apparentBytes, err +} + +func writeSnapshotSize(dir string, size int64) error { + path := filepath.Join(dir, snapshotSizeFile) + // #nosec G304 -- dir is a managed snapshot directory and the file name is fixed. + file, err := os.OpenFile(path, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600) + if err != nil { + return err + } + _, writeErr := fmt.Fprintf(file, "%d\n", size) + syncErr := file.Sync() + closeErr := file.Close() + if err := errors.Join(writeErr, syncErr, closeErr); err != nil { + return err + } + return syncDir(dir) +} + +func readSnapshotSize(dir string) (int64, error) { + // #nosec G304 -- dir is a managed snapshot directory and the file name is fixed. + data, err := os.ReadFile(filepath.Join(dir, snapshotSizeFile)) + if err != nil { + return 0, err + } + size, err := strconv.ParseInt(strings.TrimSpace(string(data)), 10, 64) + if err != nil { + return 0, fmt.Errorf("parse snapshot size in %q: %w", dir, err) + } + if size < 0 { + return 0, fmt.Errorf("snapshot size in %q must be non-negative", dir) + } + return size, nil +} diff --git a/sei-db/state_db/ss/composite/snapshot_metrics.go b/sei-db/state_db/ss/composite/snapshot_metrics.go new file mode 100644 index 0000000000..d908055f01 --- /dev/null +++ b/sei-db/state_db/ss/composite/snapshot_metrics.go @@ -0,0 +1,94 @@ +package composite + +import ( + "context" + "time" + + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" + + commonmetrics "github.com/sei-protocol/sei-chain/sei-db/common/metrics" +) + +var snapshotMeter = otel.Meter("seidb_ss_snapshot") + +var snapshotMetrics = struct { + Attempts metric.Int64Counter + Skipped metric.Int64Counter + Completions metric.Int64Counter + Duration metric.Float64Histogram + InFlight metric.Int64Gauge + CurrentHeight metric.Int64Gauge + RetainedCount metric.Int64Gauge + ApparentBytes metric.Int64Gauge +}{ + Attempts: must(snapshotMeter.Int64Counter( + "ss_snapshot_attempts", + metric.WithDescription("Number of state-store snapshot attempts"), + metric.WithUnit("{count}"), + )), + Skipped: must(snapshotMeter.Int64Counter( + "ss_snapshot_skipped", + metric.WithDescription("Number of state-store snapshot boundaries skipped by a scheduling gate"), + metric.WithUnit("{count}"), + )), + Completions: must(snapshotMeter.Int64Counter( + "ss_snapshot_completions", + metric.WithDescription("Number of completed state-store snapshot attempts"), + metric.WithUnit("{count}"), + )), + Duration: must(snapshotMeter.Float64Histogram( + "ss_snapshot_duration", + metric.WithDescription("Time from a state-store snapshot request to completion"), + metric.WithUnit("s"), + metric.WithExplicitBucketBoundaries(commonmetrics.LongLatencyBuckets...), + )), + InFlight: must(snapshotMeter.Int64Gauge( + "ss_snapshot_in_flight", + metric.WithDescription("Whether one state-store snapshot is currently in flight"), + )), + CurrentHeight: must(snapshotMeter.Int64Gauge( + "ss_snapshot_current_height", + metric.WithDescription("Height of the newest published state-store snapshot"), + )), + RetainedCount: must(snapshotMeter.Int64Gauge( + "ss_snapshot_retained_count", + metric.WithDescription("Number of retained state-store snapshots"), + metric.WithUnit("{count}"), + )), + ApparentBytes: must(snapshotMeter.Int64Gauge( + "ss_snapshot_retained_apparent_bytes", + metric.WithDescription("Apparent bytes referenced by retained state-store snapshots; hardlinks can share physical blocks"), + metric.WithUnit("By"), + )), +} + +func must[V any](instrument V, err error) V { + if err != nil { + panic(err) + } + return instrument +} + +func recordSnapshotAttempt() { + snapshotMetrics.Attempts.Add(context.Background(), 1) +} + +func recordSnapshotSkipped(reason string) { + snapshotMetrics.Skipped.Add( + context.Background(), + 1, + metric.WithAttributes(attribute.String("reason", reason)), + ) +} + +func recordSnapshotInFlight(value int64) { + snapshotMetrics.InFlight.Record(context.Background(), value) +} + +func recordSnapshotCompletion(start time.Time, outcome string) { + attrs := metric.WithAttributes(attribute.String("outcome", outcome)) + snapshotMetrics.Completions.Add(context.Background(), 1, attrs) + snapshotMetrics.Duration.Record(context.Background(), time.Since(start).Seconds(), attrs) +} diff --git a/sei-db/state_db/ss/composite/snapshot_test.go b/sei-db/state_db/ss/composite/snapshot_test.go new file mode 100644 index 0000000000..c0967ae875 --- /dev/null +++ b/sei-db/state_db/ss/composite/snapshot_test.go @@ -0,0 +1,968 @@ +package composite + +import ( + "os" + "path/filepath" + "testing" + "time" + + "github.com/sei-protocol/sei-chain/sei-db/config" + "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" + "github.com/sei-protocol/sei-chain/sei-db/management" + "github.com/sei-protocol/sei-chain/sei-db/proto" + "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/cosmos" + "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/evm" + "github.com/stretchr/testify/require" +) + +type noCheckpointStateStore struct { + types.StateStore +} + +type noBarrierStateStore struct { + types.StateStore +} + +func (*noBarrierStateStore) Checkpoint(string) error { + return nil +} + +func (*noBarrierStateStore) SetCheckpointVersion(string, int64) error { + return nil +} + +type controlledSnapshotScheduler struct { + pending chan func() + entered chan struct{} + checkpointCalls int +} + +func (*controlledSnapshotScheduler) SupportsCheckpoint() bool { + return true +} + +func (s *controlledSnapshotScheduler) ScheduleCheckpoint( + destDir string, + shouldRun func() bool, + done func(error), +) { + if s.entered != nil { + close(s.entered) + } + s.pending <- func() { + if !shouldRun() { + done(management.ErrCheckpointCanceled) + return + } + s.checkpointCalls++ + _ = os.MkdirAll(destDir, 0o750) + done(nil) + } +} + +func (*controlledSnapshotScheduler) SetCheckpointVersion(string, int64) error { + return nil +} + +func bankChangeset(key, value string) []*proto.NamedChangeSet { + return []*proto.NamedChangeSet{ + { + Name: "bank", + Changeset: proto.ChangeSet{ + Pairs: []*proto.KVPair{{Key: []byte(key), Value: []byte(value)}}, + }, + }, + } +} + +// evmStorageKey builds a key in the EVM storage family (0x03 prefix), which +// routes to the storage sub-DB when sub-DBs are separate. +func evmStorageKey() []byte { + return append([]byte{0x03}, make([]byte, 20+32)...) +} + +// setupSnapshotStore opens a store with snapshotting on at a small interval so +// tests can cross boundaries cheaply. It returns the store and its snapshots +// root. +func setupSnapshotStore(t *testing.T, interval int64, keepRecent int, separateEVMSubDBs bool) (*CompositeStateStore, string) { + t.Helper() + dir := t.TempDir() + store, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 100, + KeepRecent: 100000, + EVMSplit: true, + SeparateEVMSubDBs: separateEVMSubDBs, + EVMDBDirectory: filepath.Join(dir, "evm_ss"), + SnapshotEnable: true, + SnapshotInterval: interval, + SnapshotKeepRecent: keepRecent, + }, dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + require.NotNil(t, store.snapshotMgr) + return store, filepath.Join(dir, "data", "state_store", SnapshotsDirName) +} + +type pendingWaiter interface { + WaitForPendingWrites() +} + +// settle waits until every snapshot requested so far has been published and its +// pruning finished. Snapshot barriers sit in the backends' apply queues, so +// draining those queues is what guarantees the barriers ran. +func settle(t *testing.T, store *CompositeStateStore) { + t.Helper() + if w, ok := store.cosmosStore.(pendingWaiter); ok { + w.WaitForPendingWrites() + } + if w, ok := store.evmStore.(pendingWaiter); ok { + w.WaitForPendingWrites() + } + store.snapshotMgr.publishing.Wait() +} + +// commitBlock is what rootmulti.flush does for a populated block: enqueue the +// changesets, then hand the version to the snapshot manager. ApplyChangesetAsync +// alone schedules nothing, so tests that expect a snapshot must come through +// here. +func commitBlock(t *testing.T, store *CompositeStateStore, version int64, changesets []*proto.NamedChangeSet) { + t.Helper() + require.NoError(t, store.ApplyChangesetAsync(version, changesets)) + store.ScheduleSnapshot(version) +} + +func writeBlock(t *testing.T, store *CompositeStateStore, version int64) { + t.Helper() + commitBlock(t, store, version, []*proto.NamedChangeSet{ + { + Name: "bank", + Changeset: proto.ChangeSet{ + Pairs: []*proto.KVPair{{Key: []byte("balance"), Value: []byte{byte(version)}}}, + }, + }, + { + Name: evm.EVMStoreKey, + Changeset: proto.ChangeSet{ + Pairs: []*proto.KVPair{{Key: evmStorageKey(), Value: []byte{byte(version)}}}, + }, + }, + }) +} + +// The snapshot manager keys off the mirrored cadence, so the ss-snapshot-enable +// switch has to reach it as a zero interval and leave no manager running. +func TestSnapshotManagerRespectsSnapshotEnable(t *testing.T) { + for _, tc := range []struct { + name string + enable bool + wantRunning bool + }{ + {name: "enabled", enable: true, wantRunning: true}, + {name: "disabled", enable: false, wantRunning: false}, + } { + t.Run(tc.name, func(t *testing.T) { + dir := t.TempDir() + ssConfig := config.DefaultStateStoreConfig() + ssConfig.SnapshotEnable = tc.enable + config.AlignSSSnapshotWithSC(config.DefaultStateCommitConfig(), &ssConfig) + + store, err := NewCompositeStateStore(ssConfig, dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + + if tc.wantRunning { + require.NotNil(t, store.snapshotMgr, "explicit opt-in starts snapshotting") + require.Positive(t, ssConfig.SnapshotInterval) + } else { + require.Nil(t, store.snapshotMgr) + require.Zero(t, ssConfig.SnapshotInterval) + } + }) + } +} + +func TestCustomStateStoreDirectoryMovesSnapshotRootBesideDatabase(t *testing.T) { + home := t.TempDir() + customDB := filepath.Join(t.TempDir(), "cosmos-state") + cfg := config.DefaultStateStoreConfig() + cfg.Backend = config.PebbleDBBackend + cfg.DBDirectory = customDB + cfg.SnapshotEnable = true + cfg.SnapshotInterval = 5 + cfg.SnapshotKeepRecent = 1 + + store, err := NewCompositeStateStore(cfg, home) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + + require.Equal(t, customDB+"-"+SnapshotsDirName, store.snapshotMgr.root) +} + +func TestSnapshotHardlinkPreflightCleansProbeFiles(t *testing.T) { + source := t.TempDir() + root := t.TempDir() + require.NoError(t, verifySnapshotHardlinks(root, []string{source})) + + sourceEntries, err := os.ReadDir(source) + require.NoError(t, err) + require.Empty(t, sourceEntries) + rootEntries, err := os.ReadDir(root) + require.NoError(t, err) + require.Empty(t, rootEntries) +} + +func TestSnapshotHardlinkPreflightRejectsCrossFilesystem(t *testing.T) { + root, err := os.MkdirTemp("/dev/shm", "ss-snapshot-test-*") + if err != nil { + t.Skipf("no separate /dev/shm filesystem: %v", err) + } + t.Cleanup(func() { require.NoError(t, os.RemoveAll(root)) }) + + err = verifySnapshotHardlinks(root, []string{t.TempDir()}) + if err == nil { + t.Skip("temporary directory and /dev/shm use the same filesystem") + } + require.ErrorContains(t, err, "cannot hardlink snapshots") +} + +func TestSnapshotManagerRejectsUnsupportedBackend(t *testing.T) { + store := &CompositeStateStore{ + cosmosStore: cosmos.NewCosmosStateStore(&noCheckpointStateStore{}), + config: config.StateStoreConfig{ + Backend: config.RocksDBBackend, + SnapshotInterval: 10, + }, + } + + err := store.startSnapshotManager(t.TempDir(), nil) + require.ErrorContains(t, err, "does not support checkpoints") + require.Nil(t, store.snapshotMgr) +} + +func TestSnapshotManagerRejectsBackendWithoutBarrier(t *testing.T) { + store := &CompositeStateStore{ + cosmosStore: cosmos.NewCosmosStateStore(&noBarrierStateStore{}), + config: config.StateStoreConfig{ + Backend: config.PebbleDBBackend, + SnapshotInterval: 10, + }, + } + + err := store.startSnapshotManager(t.TempDir(), nil) + require.ErrorContains(t, err, "does not support checkpoints") + require.Nil(t, store.snapshotMgr) +} + +func TestSnapshotStopCancelsQueuedCheckpoint(t *testing.T) { + scheduler := &controlledSnapshotScheduler{pending: make(chan func(), 1)} + manager := &snapshotManager{ + root: t.TempDir(), + backend: config.PebbleDBBackend, + interval: 5, + keepRecent: 1, + cosmosScheduler: scheduler, + } + + manager.maybeSnapshot(5) + manager.stop() + (<-scheduler.pending)() + + require.Zero(t, scheduler.checkpointCalls) + versions, err := ListSnapshotVersions(manager.root) + require.NoError(t, err) + require.Empty(t, versions) +} + +func TestSnapshotManagerAllowsOnlyOneInFlightSnapshot(t *testing.T) { + scheduler := &controlledSnapshotScheduler{pending: make(chan func(), 2)} + manager := &snapshotManager{ + root: t.TempDir(), + backend: config.PebbleDBBackend, + interval: 5, + keepRecent: 1, + cosmosScheduler: scheduler, + } + + manager.maybeSnapshot(5) + manager.maybeSnapshot(10) + require.Len(t, scheduler.pending, 1, "a second boundary must not enqueue while one snapshot is active") + + (<-scheduler.pending)() + manager.publishing.Wait() + require.False(t, manager.inFlight) + require.Equal(t, int64(5), manager.lastRequested) +} + +func TestSnapshotManagerAppliesMinimumTimeInterval(t *testing.T) { + scheduler := &controlledSnapshotScheduler{pending: make(chan func(), 2)} + manager := &snapshotManager{ + root: t.TempDir(), + backend: config.PebbleDBBackend, + interval: 5, + minTime: time.Hour, + keepRecent: 1, + cosmosScheduler: scheduler, + } + + manager.maybeSnapshot(5) + (<-scheduler.pending)() + manager.publishing.Wait() + + manager.maybeSnapshot(10) + require.Empty(t, scheduler.pending, "a rapid boundary must be skipped") + + manager.mu.Lock() + manager.lastRequestAt = time.Now().Add(-2 * time.Hour) + manager.mu.Unlock() + manager.maybeSnapshot(10) + require.Len(t, scheduler.pending, 1) + (<-scheduler.pending)() + manager.publishing.Wait() +} + +func TestSnapshotStopWaitsForBarrierScheduling(t *testing.T) { + scheduler := &controlledSnapshotScheduler{ + pending: make(chan func()), + entered: make(chan struct{}), + } + manager := &snapshotManager{ + root: t.TempDir(), + backend: config.PebbleDBBackend, + interval: 5, + keepRecent: 1, + cosmosScheduler: scheduler, + } + + requestDone := make(chan struct{}) + go func() { + manager.maybeSnapshot(5) + close(requestDone) + }() + <-scheduler.entered + + stopDone := make(chan struct{}) + go func() { + manager.stop() + close(stopDone) + }() + require.Never(t, func() bool { + select { + case <-stopDone: + return true + default: + return false + } + }, 50*time.Millisecond, 5*time.Millisecond) + + callback := <-scheduler.pending + <-requestDone + <-stopDone + callback() + require.False(t, manager.inFlight) +} + +// Snapshot labels are the interval boundaries themselves, not whatever version +// the store happened to be at when some background pass noticed. That is the +// property the in-queue barrier buys. It keeps each accepted SS snapshot's +// contents aligned with its own label even when SC independently skips that +// boundary. +func TestSnapshotTakenAtExactIntervalBoundaries(t *testing.T) { + store, root := setupSnapshotStore(t, 5, 5, false) + + for v := int64(1); v <= 12; v++ { + writeBlock(t, store, v) + if v%5 == 0 { + settle(t, store) + } + } + settle(t, store) + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{5, 10}, versions, + "snapshots must land on interval boundaries and nowhere else") + + target, err := os.Readlink(filepath.Join(root, snapshotCurrentLink)) + require.NoError(t, err) + require.Equal(t, SnapshotDirName(10), target) + + snapDir := filepath.Join(root, SnapshotDirName(10)) + apparentBytes, err := readSnapshotSize(snapDir) + require.NoError(t, err) + require.Positive(t, apparentBytes) + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: config.PebbleDBBackend, + AsyncWriteBuffer: 0, + KeepRecent: 100000, + EVMSplit: true, + DBDirectory: filepath.Join(snapDir, "cosmos", config.PebbleDBBackend), + EVMDBDirectory: filepath.Join(snapDir, "evm", config.PebbleDBBackend), + }, t.TempDir()) + require.NoError(t, err) + defer reopened.Close() + + require.Equal(t, int64(10), reopened.GetLatestVersion()) + cosmosValue, err := reopened.Get("bank", 12, []byte("balance")) + require.NoError(t, err) + require.Equal(t, []byte{10}, cosmosValue, "snapshot 10 must exclude Cosmos writes 11 and 12") + evmValue, err := reopened.Get(evm.EVMStoreKey, 12, evmStorageKey()) + require.NoError(t, err) + require.Equal(t, []byte{10}, evmValue, "snapshot 10 must exclude EVM writes 11 and 12") +} + +func TestSnapshotTakenAtExactIntervalBoundaryWithoutEVMSplit(t *testing.T) { + dir := t.TempDir() + ssConfig := config.DefaultStateStoreConfig() + ssConfig.Backend = config.PebbleDBBackend + ssConfig.AsyncWriteBuffer = 100 + ssConfig.KeepRecent = 100000 + ssConfig.EVMSplit = false + ssConfig.SnapshotEnable = true + ssConfig.SnapshotInterval = 5 + ssConfig.SnapshotKeepRecent = 1 + + store, err := NewCompositeStateStore(ssConfig, dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + require.NotNil(t, store.snapshotMgr) + + for version := int64(1); version <= 5; version++ { + commitBlock(t, store, version, bankChangeset("balance", "value")) + } + settle(t, store) + + root := filepath.Join(dir, "data", "state_store", SnapshotsDirName) + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{5}, versions) +} + +// A snapshot must be a complete image of every version at or below its label, +// reopenable as a store in its own right. +func TestSnapshotReopensWithEveryVersionBelowLabel(t *testing.T) { + store, root := setupSnapshotStore(t, 10, 5, false) + + for v := int64(1); v <= 10; v++ { + writeBlock(t, store, v) + } + settle(t, store) + + const label = int64(10) + snapDir := filepath.Join(root, SnapshotDirName(label)) + require.DirExists(t, snapDir) + + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + EVMSplit: true, + DBDirectory: filepath.Join(snapDir, "cosmos", "pebbledb"), + EVMDBDirectory: filepath.Join(snapDir, "evm", "pebbledb"), + }, t.TempDir()) + require.NoError(t, err) + defer reopened.Close() + + require.Equal(t, label, reopened.GetLatestVersion(), + "the label is the version the snapshot was requested at") + for v := int64(1); v <= label; v++ { + val, err := reopened.Get("bank", v, []byte("balance")) + require.NoError(t, err) + require.Equal(t, []byte{byte(v)}, val, "cosmos version %d missing from snapshot", v) + val, err = reopened.Get(evm.EVMStoreKey, v, evmStorageKey()) + require.NoError(t, err) + require.Equal(t, []byte{byte(v)}, val, "evm version %d missing from snapshot", v) + } +} + +// The property the barrier exists for: the label stays exact while the write +// path keeps going. Nothing is drained between block 10 and blocks 11 and 12, so +// the checkpoint runs with later versions already queued behind the barrier — the +// case a post-hoc "snapshot what has been applied" scheme would get wrong. +func TestSnapshotExcludesVersionsWrittenAfterTheBoundary(t *testing.T) { + store, root := setupSnapshotStore(t, 10, 5, false) + + const label = int64(10) + for v := int64(1); v <= label; v++ { + writeBlock(t, store, v) + } + for v := label + 1; v <= label+2; v++ { + writeBlock(t, store, v) + } + settle(t, store) + + snapDir := filepath.Join(root, SnapshotDirName(label)) + require.DirExists(t, snapDir) + + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + EVMSplit: true, + DBDirectory: filepath.Join(snapDir, "cosmos", "pebbledb"), + EVMDBDirectory: filepath.Join(snapDir, "evm", "pebbledb"), + }, t.TempDir()) + require.NoError(t, err) + defer reopened.Close() + + require.Equal(t, label, reopened.GetLatestVersion()) + // Reading above the label returns the label's value rather than 11 or 12, + // which is what "excluded" means for an MVCC store: the later writes are not + // in this image at any version. + for _, above := range []int64{label + 1, label + 2} { + val, err := reopened.Get("bank", above, []byte("balance")) + require.NoError(t, err) + require.Equal(t, []byte{byte(label)}, val, + "cosmos read at %d saw a write from after the boundary", above) + val, err = reopened.Get(evm.EVMStoreKey, above, evmStorageKey()) + require.NoError(t, err) + require.Equal(t, []byte{byte(label)}, val, + "evm read at %d saw a write from after the boundary", above) + } + + // The live store keeps them, so the snapshot dropped them rather than the + // writes never landing. + val, err := store.Get("bank", label+2, []byte("balance")) + require.NoError(t, err) + require.Equal(t, []byte{byte(label + 2)}, val) +} + +// A snapshot inherits each database's earliest marker. The composite is allowed +// to reopen with different member floors because it reports the highest one. +func TestSnapshotInheritsPerStoreEarliestMarkers(t *testing.T) { + store, root := setupSnapshotStore(t, 10, 5, false) + + for v := int64(1); v <= 9; v++ { + writeBlock(t, store, v) + } + settle(t, store) + require.NoError(t, store.cosmosStore.SetEarliestVersion(2, false)) + require.NoError(t, store.evmStore.SetEarliestVersion(5, false)) + require.Equal(t, int64(2), store.cosmosStore.GetEarliestVersion()) + require.Equal(t, int64(5), store.evmStore.GetEarliestVersion()) + + writeBlock(t, store, 10) + settle(t, store) + + snapDir := filepath.Join(root, SnapshotDirName(10)) + require.DirExists(t, snapDir) + + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + EVMSplit: true, + DBDirectory: filepath.Join(snapDir, "cosmos", "pebbledb"), + EVMDBDirectory: filepath.Join(snapDir, "evm", "pebbledb"), + }, t.TempDir()) + require.NoError(t, err, "a snapshot with different member floors must reopen") + defer reopened.Close() + + require.Equal(t, int64(2), reopened.cosmosStore.GetEarliestVersion()) + require.Equal(t, int64(5), reopened.GetEarliestVersion()) + require.Equal(t, int64(5), reopened.evmStore.GetEarliestVersion()) +} + +func TestSnapshotInheritsEarliestMarkerAfterPrune(t *testing.T) { + store, root := setupSnapshotStore(t, 10, 5, false) + + for v := int64(1); v <= 9; v++ { + writeBlock(t, store, v) + } + settle(t, store) + require.NoError(t, store.Prune(4)) + + writeBlock(t, store, 10) + settle(t, store) + + snapDir := filepath.Join(root, SnapshotDirName(10)) + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + EVMSplit: true, + DBDirectory: filepath.Join(snapDir, "cosmos", "pebbledb"), + EVMDBDirectory: filepath.Join(snapDir, "evm", "pebbledb"), + }, t.TempDir()) + require.NoError(t, err) + defer reopened.Close() + + require.Equal(t, int64(5), reopened.GetEarliestVersion()) + require.Equal(t, int64(5), reopened.cosmosStore.GetEarliestVersion()) + require.Equal(t, int64(5), reopened.evmStore.GetEarliestVersion()) +} + +// With separate sub-DBs the latest label has to reach every one of them, not +// just the sub-DBs that took writes, or the snapshot is not self-describing at +// its exact boundary. +func TestSnapshotSetsLatestVersionEveryEVMSubDB(t *testing.T) { + store, root := setupSnapshotStore(t, 10, 5, true) + + for v := int64(1); v <= 9; v++ { + writeBlock(t, store, v) + } + + writeBlock(t, store, 10) + settle(t, store) + + snapDir := filepath.Join(root, SnapshotDirName(10)) + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + EVMSplit: true, + SeparateEVMSubDBs: true, + DBDirectory: filepath.Join(snapDir, "cosmos", "pebbledb"), + EVMDBDirectory: filepath.Join(snapDir, "evm", "pebbledb"), + }, t.TempDir()) + require.NoError(t, err) + + require.Equal(t, int64(10), reopened.evmStore.GetLatestVersion()) + require.NoError(t, reopened.Close()) + + evmRoot := filepath.Join(snapDir, "evm", "pebbledb") + for _, storeType := range evm.AllEVMStoreTypes() { + subDir := filepath.Join(evmRoot, evm.StoreTypeName(storeType)) + subDB, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + UseDefaultComparer: true, + DBDirectory: subDir, + }, t.TempDir()) + require.NoError(t, err) + require.Equal(t, int64(10), subDB.GetLatestVersion(), "sub-DB %s latest marker", evm.StoreTypeName(storeType)) + require.NoError(t, subDB.Close()) + } +} + +// The reason the barrier has to be a message in every queue rather than a wait: +// a block that only touches storage keys is enqueued only on the storage sub-DB, +// so the idle sub-DBs never observe that version and no amount of waiting would +// tell them it passed. Every sub-DB must still be captured. +func TestSnapshotCapturesIdleEVMSubDBs(t *testing.T) { + store, root := setupSnapshotStore(t, 5, 5, true) + + // Storage keys only: codehash, code and misc sub-DBs stay idle throughout. + for v := int64(1); v <= 5; v++ { + commitBlock(t, store, v, []*proto.NamedChangeSet{ + { + Name: evm.EVMStoreKey, + Changeset: proto.ChangeSet{ + Pairs: []*proto.KVPair{{Key: evmStorageKey(), Value: []byte{byte(v)}}}, + }, + }, + }) + } + settle(t, store) + + evmRoot := filepath.Join(root, SnapshotDirName(5), "evm", "pebbledb") + for _, storeType := range evm.AllEVMStoreTypes() { + name := evm.StoreTypeName(storeType) + subDir := filepath.Join(evmRoot, name) + require.DirExists(t, subDir, "sub-DB %s missing from snapshot", name) + // A checkpoint always carries a manifest. An empty directory would mean + // the barrier never reached that sub-DB. + manifests, err := filepath.Glob(filepath.Join(subDir, "MANIFEST-*")) + require.NoError(t, err) + require.NotEmpty(t, manifests, "sub-DB %s was not checkpointed", name) + } + + // The storage sub-DB is the one that actually took writes, and it must be + // readable at every version up to the label. + storage, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + // EVM sub-DBs are opened with the plain byte comparer. + UseDefaultComparer: true, + DBDirectory: filepath.Join(evmRoot, evm.StoreTypeName(evm.StoreStorage)), + }, t.TempDir()) + require.NoError(t, err) + defer storage.Close() + + for v := int64(1); v <= 5; v++ { + val, err := storage.Get(evm.EVMStoreKey, v, evmStorageKey()) + require.NoError(t, err) + require.Equal(t, []byte{byte(v)}, val, "evm storage version %d missing from snapshot", v) + } +} + +// An interval boundary that happens to be an empty block arrives through +// SetLatestVersion rather than the changeset path, and must still snapshot — +// otherwise a quiet chain skips whole intervals. +func TestSnapshotTakenOnEmptyBoundaryBlock(t *testing.T) { + store, root := setupSnapshotStore(t, 5, 5, false) + + for v := int64(1); v <= 4; v++ { + writeBlock(t, store, v) + } + // Block 5 is empty: marker only, nothing enqueued. + require.NoError(t, store.SetLatestVersion(5)) + store.ScheduleSnapshot(5) + settle(t, store) + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{5}, versions) + + // Every data version below the label is inside the snapshot, and the + // checkpoint marker advances to the empty block's version. + snapDir := filepath.Join(root, SnapshotDirName(5)) + reopened, err := NewCompositeStateStore(config.StateStoreConfig{ + Backend: "pebbledb", + AsyncWriteBuffer: 0, + KeepRecent: 100000, + DBDirectory: filepath.Join(snapDir, "cosmos", "pebbledb"), + }, t.TempDir()) + require.NoError(t, err) + defer reopened.Close() + + require.Equal(t, int64(5), reopened.GetLatestVersion()) + val, err := reopened.Get("bank", 4, []byte("balance")) + require.NoError(t, err) + require.Equal(t, []byte{4}, val) +} + +func TestSetLatestVersionDoesNotSnapshotDuringImport(t *testing.T) { + store, root := setupSnapshotStore(t, 5, 5, false) + nodes := make(chan types.SnapshotNode) + importDone := make(chan error, 1) + go func() { + importDone <- store.Import(5, nodes) + }() + closed := false + t.Cleanup(func() { + if !closed { + close(nodes) + <-importDone + } + }) + + nodes <- types.SnapshotNode{StoreKey: "bank", Key: []byte("balance"), Value: []byte{5}} + require.NoError(t, store.SetLatestVersion(5)) + settle(t, store) + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Empty(t, versions, "direct restore metadata writes must not trigger a snapshot") + + close(nodes) + closed = true + require.NoError(t, <-importDone) +} + +// TestSnapshotPrune verifies retention: with keepRecent=1, only the newest two +// snapshots survive and current tracks the newest. +func TestSnapshotPrune(t *testing.T) { + store, root := setupSnapshotStore(t, 5, 1, false) + + for v := int64(1); v <= 15; v++ { + writeBlock(t, store, v) + if v%5 == 0 { + settle(t, store) + } + } + settle(t, store) + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{10, 15}, versions) + + target, err := os.Readlink(filepath.Join(root, snapshotCurrentLink)) + require.NoError(t, err) + require.Equal(t, SnapshotDirName(15), target) +} + +func TestSnapshotPruneKeepsCurrentTarget(t *testing.T) { + root := t.TempDir() + for _, version := range []int64{5, 10, 15, 20} { + require.NoError(t, os.MkdirAll(filepath.Join(root, SnapshotDirName(version)), 0o750)) + } + require.NoError(t, os.Symlink(SnapshotDirName(5), filepath.Join(root, snapshotCurrentLink))) + + manager := &snapshotManager{root: root, keepRecent: 1} + manager.prune() + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{5, 15, 20}, versions, + "retention may keep one extra snapshot but must not dangle current") + target, err := os.Readlink(filepath.Join(root, snapshotCurrentLink)) + require.NoError(t, err) + require.Equal(t, SnapshotDirName(5), target) +} + +func TestSnapshotManagerResumesFromNewestSnapshot(t *testing.T) { + dir := t.TempDir() + cfg := config.DefaultStateStoreConfig() + cfg.Backend = config.PebbleDBBackend + cfg.AsyncWriteBuffer = 100 + cfg.KeepRecent = 100000 + cfg.SnapshotEnable = true + cfg.SnapshotInterval = 5 + cfg.SnapshotKeepRecent = 1 + cfg.SnapshotMinTimeInterval = time.Hour + + store, err := NewCompositeStateStore(cfg, dir) + require.NoError(t, err) + for version := int64(1); version <= 5; version++ { + commitBlock(t, store, version, bankChangeset("balance", "value")) + } + settle(t, store) + + root := filepath.Join(dir, "data", "state_store", SnapshotsDirName) + snapshotDir := filepath.Join(root, SnapshotDirName(5)) + before, err := os.Stat(snapshotDir) + require.NoError(t, err) + require.NoError(t, store.Close()) + require.NoError(t, os.Remove(filepath.Join(root, snapshotCurrentLink))) + + reopened, err := NewCompositeStateStore(cfg, dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, reopened.Close()) }) + require.Equal(t, int64(5), reopened.snapshotMgr.lastRequested) + require.WithinDuration(t, before.ModTime(), reopened.snapshotMgr.lastRequestAt, time.Second) + target, err := os.Readlink(filepath.Join(root, snapshotCurrentLink)) + require.NoError(t, err) + require.Equal(t, SnapshotDirName(5), target) + + reopened.ScheduleSnapshot(5) + settle(t, reopened) + after, err := os.Stat(snapshotDir) + require.NoError(t, err) + require.True(t, os.SameFile(before, after), "restart must not replace an existing boundary snapshot") + + for version := int64(6); version <= 10; version++ { + commitBlock(t, reopened, version, bankChangeset("balance", "value")) + } + settle(t, reopened) + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{5}, versions, "restart must preserve the minimum-time gate") +} + +func TestOutOfOrderPublishDoesNotMoveCurrentBackward(t *testing.T) { + root := t.TempDir() + manager := &snapshotManager{root: root, keepRecent: 5} + + publish := func(version int64) { + tmpDir := filepath.Join(root, snapshotTmpPrefix+SnapshotDirName(version)) + require.NoError(t, os.MkdirAll(tmpDir, 0o750)) + manager.publish(version, tmpDir, filepath.Join(root, SnapshotDirName(version)), time.Now()) + } + publish(10) + publish(5) + + target, err := os.Readlink(filepath.Join(root, snapshotCurrentLink)) + require.NoError(t, err) + require.Equal(t, SnapshotDirName(10), target) +} + +func TestSnapshotManagerPrunesExistingSnapshotsAtStartup(t *testing.T) { + dir := t.TempDir() + root := filepath.Join(dir, "data", "state_store", SnapshotsDirName) + for _, version := range []int64{5, 10, 15} { + require.NoError(t, os.MkdirAll(filepath.Join(root, SnapshotDirName(version)), 0o750)) + } + + cfg := config.DefaultStateStoreConfig() + cfg.Backend = config.PebbleDBBackend + cfg.SnapshotEnable = true + cfg.SnapshotInterval = 5 + cfg.SnapshotKeepRecent = 1 + store, err := NewCompositeStateStore(cfg, dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{10, 15}, versions) +} + +func TestFailedPublishStillEnforcesRetention(t *testing.T) { + root := t.TempDir() + for _, version := range []int64{5, 10, 15} { + require.NoError(t, os.MkdirAll(filepath.Join(root, SnapshotDirName(version)), 0o750)) + } + manager := &snapshotManager{root: root, keepRecent: 1} + + published := manager.publish( + 20, + filepath.Join(root, "missing-staging-dir"), + filepath.Join(root, SnapshotDirName(20)), + time.Now(), + ) + require.False(t, published) + + versions, err := ListSnapshotVersions(root) + require.NoError(t, err) + require.Equal(t, []int64{10, 15}, versions) +} + +func TestRetentionMetricsCacheSnapshotSizes(t *testing.T) { + root := t.TempDir() + snapshotDir := filepath.Join(root, SnapshotDirName(5)) + require.NoError(t, os.MkdirAll(snapshotDir, 0o750)) + dataFile := filepath.Join(snapshotDir, "data.sst") + require.NoError(t, os.WriteFile(dataFile, []byte("one"), 0o600)) + require.NoError(t, writeSnapshotSize(snapshotDir, 3)) + + manager := &snapshotManager{root: root} + manager.recordRetentionMetrics() + require.Equal(t, int64(3), manager.snapshotSizes[5]) + + // Published snapshots are immutable, so later metric records reuse the + // cached total rather than walking every retained hardlink tree again. + require.NoError(t, os.WriteFile(dataFile, []byte("a longer value"), 0o600)) + manager.recordRetentionMetrics() + require.Equal(t, int64(3), manager.snapshotSizes[5]) + + require.NoError(t, os.RemoveAll(snapshotDir)) + manager.recordRetentionMetrics() + require.NotContains(t, manager.snapshotSizes, int64(5)) +} + +func TestSnapshotRequestReturnsUnexpectedStatError(t *testing.T) { + root := t.TempDir() + name := SnapshotDirName(5) + require.NoError(t, os.Symlink(name, filepath.Join(root, name))) + + manager := &snapshotManager{root: root, backend: config.PebbleDBBackend} + err := manager.requestSnapshot(5, time.Now()) + require.ErrorContains(t, err, "inspect snapshot dir") +} + +// A crash mid-snapshot leaves a staging directory named after the boundary it +// was staging, which would otherwise sit there until that exact boundary came +// round again. +func TestStaleSnapshotTmpDirRemovedAtStartup(t *testing.T) { + dir := t.TempDir() + root := filepath.Join(dir, "data", "state_store", SnapshotsDirName) + stale := filepath.Join(root, snapshotTmpPrefix+SnapshotDirName(40)) + require.NoError(t, os.MkdirAll(filepath.Join(stale, "cosmos"), 0o750)) + tmpLink := filepath.Join(root, snapshotCurrentTmpLink) + require.NoError(t, os.Symlink(filepath.Base(stale), tmpLink)) + + ssConfig := config.DefaultStateStoreConfig() + ssConfig.SnapshotEnable = true + config.AlignSSSnapshotWithSC(config.DefaultStateCommitConfig(), &ssConfig) + store, err := NewCompositeStateStore(ssConfig, dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, store.Close()) }) + + require.NoDirExists(t, stale) + _, err = os.Lstat(tmpLink) + require.ErrorIs(t, err, os.ErrNotExist) +} + +func TestParseSnapshotVersion(t *testing.T) { + v, ok := ParseSnapshotVersion(SnapshotDirName(219140000)) + require.True(t, ok) + require.Equal(t, int64(219140000), v) + + for _, bad := range []string{"snapshot-", "snapshot-123", "current", "tmp-snapshot-00000000000000000010", "snapshot-0000000000000000001x"} { + _, ok := ParseSnapshotVersion(bad) + require.False(t, ok, "expected %q to be rejected", bad) + } +} diff --git a/sei-db/state_db/ss/composite/store.go b/sei-db/state_db/ss/composite/store.go index d88c77647e..dbaf727921 100644 --- a/sei-db/state_db/ss/composite/store.go +++ b/sei-db/state_db/ss/composite/store.go @@ -4,6 +4,7 @@ import ( "encoding/binary" "fmt" "os" + "path/filepath" "sync" dbm "github.com/tendermint/tm-db" @@ -34,6 +35,7 @@ type CompositeStateStore struct { cosmosStore types.StateStore // CosmosStateStore wrapping MVCC DB evmStore types.StateStore // EVMStateStore wrapping sub MVCC DBs (nil if disabled) pruningManager *pruning.Manager + snapshotMgr *snapshotManager config config.StateStoreConfig closeOnce sync.Once closeErr error @@ -60,6 +62,7 @@ func NewCompositeStateStore( cosmosStore: cosmosStore, config: ssConfig, } + snapshotSourceDirs := []string{dbHome} if ssConfig.EVMSplit { evmDir := ssConfig.EVMDBDirectory @@ -79,6 +82,16 @@ func NewCompositeStateStore( return nil, fmt.Errorf("failed to create EVM store: %w", err) } cs.evmStore = evmStore + if ssConfig.SeparateEVMSubDBs { + for _, storeType := range evm.AllEVMStoreTypes() { + snapshotSourceDirs = append( + snapshotSourceDirs, + filepath.Join(evmDir, evm.StoreTypeName(storeType)), + ) + } + } else { + snapshotSourceDirs = append(snapshotSourceDirs, evmDir) + } logger.Info("EVM state store enabled", "dir", evmDir, "separateDBs", ssConfig.SeparateEVMSubDBs, @@ -97,12 +110,22 @@ func NewCompositeStateStore( return nil, fmt.Errorf("failed to recover state store: %w", err) } - // Mismatched earliest versions = DBs from different snapshots; reads would diverge. - if err := cs.validateEVMSSPostRecovery(); err != nil { - _ = cs.Close() - return nil, err - } + cs.validateEVMSSPostRecovery() + if ssConfig.SnapshotInterval > 0 { + snapshotRoot := utils.GetStateStoreSnapshotsPath(homeDir) + if ssConfig.DBDirectory != "" { + cleanDBHome := filepath.Clean(dbHome) + snapshotRoot = filepath.Join( + filepath.Dir(cleanDBHome), + filepath.Base(cleanDBHome)+"-"+utils.StateStoreSnapshotsDirName, + ) + } + if err := cs.startSnapshotManager(snapshotRoot, snapshotSourceDirs); err != nil { + _ = cs.Close() + return nil, fmt.Errorf("start state store snapshot manager: %w", err) + } + } cs.StartPruning() return cs, nil @@ -150,20 +173,23 @@ func (s *CompositeStateStore) validateEVMSSPreRecovery() error { return nil } -// validateEVMSSPostRecovery rejects mismatched earliest versions between the two SS DBs. -func (s *CompositeStateStore) validateEVMSSPostRecovery() error { +// validateEVMSSPostRecovery reports mismatched earliest versions between SS DBs. +// Divergence is safe because GetEarliestVersion reports the highest member +// floor, which is the first version every routed store can serve. +func (s *CompositeStateStore) validateEVMSSPostRecovery() { if s.evmStore == nil { - return nil + return } cosmosEarliest := s.cosmosStore.GetEarliestVersion() evmEarliest := s.evmStore.GetEarliestVersion() if cosmosEarliest != evmEarliest && (cosmosEarliest > 0 || evmEarliest > 0) { - return fmt.Errorf( - "EVM SS earliest version %d does not match Cosmos SS earliest version %d: state sync the EVM SS DB, or set evm-ss-split=false", - evmEarliest, cosmosEarliest, + logger.Warn( + "EVM SS earliest version does not match Cosmos SS earliest version; serving the highest floor", + "evmEarliest", evmEarliest, + "cosmosEarliest", cosmosEarliest, + "reportedEarliest", max(cosmosEarliest, evmEarliest), ) } - return nil } func (s *CompositeStateStore) StartPruning() { @@ -179,6 +205,14 @@ func (s *CompositeStateStore) evmRouted(storeKey string) bool { return s.evmStore != nil && storeKey == evm.EVMStoreKey } +// The read methods below route by store key and do not re-check +// GetEarliestVersion. The cosmos KVStore wrapper over a StateStore panics on any +// read error, and pruning can raise the floor after a query store was built, so +// an error here would crash the process for a request that must merely fail. The +// floor is enforced where an error is representable: query-store construction +// and VersionExists. Below the floor a pruned member reports the key as absent, +// as its engine already does. + func (s *CompositeStateStore) Get(storeKey string, version int64, key []byte) ([]byte, error) { if s.evmRouted(storeKey) { return s.evmStore.Get(storeKey, version, key) @@ -216,11 +250,18 @@ func (s *CompositeStateStore) GetLatestVersion() int64 { } func (s *CompositeStateStore) GetEarliestVersion() int64 { - return s.cosmosStore.GetEarliestVersion() + earliest := s.cosmosStore.GetEarliestVersion() + if s.evmStore != nil { + earliest = max(earliest, s.evmStore.GetEarliestVersion()) + } + return earliest } func (s *CompositeStateStore) Close() error { s.closeOnce.Do(func() { + if s.snapshotMgr != nil { + s.snapshotMgr.stop() + } if s.pruningManager != nil { s.pruningManager.Stop() } @@ -306,6 +347,19 @@ func (s *CompositeStateStore) ApplyChangesetAsync(version int64, changesets []*p return nil } +// ScheduleSnapshot asks the snapshot manager to capture version once the caller +// has enqueued every state change for that version and nothing above it. +// +// This is deliberately not called from ApplyChangesetAsync. That method is part +// of the general StateStore interface and has callers outside the commit path, +// such as the benchmark wrappers, which would inherit a snapshot trigger they +// never asked for. The rootmulti commit path is the single choke point that +// sees both the populated and the empty block, so it owns the trigger. Direct +// writes such as import, recovery, and prune must not use this hook. +func (s *CompositeStateStore) ScheduleSnapshot(version int64) { + s.snapshotMgr.maybeSnapshot(version) +} + func filterEVMChangesets(changesets []*proto.NamedChangeSet) []*proto.NamedChangeSet { var evmCS []*proto.NamedChangeSet for _, cs := range changesets { diff --git a/sei-db/state_db/ss/cosmos/store.go b/sei-db/state_db/ss/cosmos/store.go index 5b02d8ed15..c512335e49 100644 --- a/sei-db/state_db/ss/cosmos/store.go +++ b/sei-db/state_db/ss/cosmos/store.go @@ -4,6 +4,7 @@ import ( dbm "github.com/tendermint/tm-db" "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" + "github.com/sei-protocol/sei-chain/sei-db/management" "github.com/sei-protocol/sei-chain/sei-db/proto" ) @@ -76,3 +77,24 @@ func (s *CosmosStateStore) Import(version int64, ch <-chan types.SnapshotNode) e func (s *CosmosStateStore) Close() error { return s.db.Close() } + +func (s *CosmosStateStore) SupportsCheckpoint() bool { + _, checkpointable := s.db.(types.Checkpointable) + _, barrier := s.db.(types.DrainBarrier) + _, versionSetter := s.db.(types.CheckpointVersionSetter) + return checkpointable && barrier && versionSetter +} + +func (s *CosmosStateStore) ScheduleCheckpoint(destDir string, shouldRun func() bool, done func(error)) { + management.ScheduleCheckpoint(s.db, destDir, shouldRun, done) +} + +func (s *CosmosStateStore) SetCheckpointVersion(destDir string, version int64) error { + return management.SetCheckpointVersion(s.db, destDir, version) +} + +func (s *CosmosStateStore) WaitForPendingWrites() { + if w, ok := s.db.(interface{ WaitForPendingWrites() }); ok { + w.WaitForPendingWrites() + } +} diff --git a/sei-db/state_db/ss/evm/db_test.go b/sei-db/state_db/ss/evm/db_test.go index 4b7682f302..69dc3ea2b6 100644 --- a/sei-db/state_db/ss/evm/db_test.go +++ b/sei-db/state_db/ss/evm/db_test.go @@ -31,6 +31,28 @@ func openTestStore(t *testing.T) types.StateStore { return store } +// GetEarliestVersion reports the highest sub-DB floor, which is the earliest +// version every routed sub-DB can serve. +func TestGetEarliestVersionReportsTheFurthestPrunedSubDB(t *testing.T) { + dir := t.TempDir() + cfg := testConfig() + cfg.SeparateEVMSubDBs = true + + store, err := NewEVMStateStore(dir, cfg) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + require.Greater(t, len(store.managedDBs), 1) + + require.Zero(t, store.GetEarliestVersion()) + + require.NoError(t, store.subDBs[StoreStorage].SetEarliestVersion(40, false)) + require.Equal(t, int64(40), store.GetEarliestVersion()) + + // Once the pass finishes, the reported floor stays the same. + require.NoError(t, store.SetEarliestVersion(40, false)) + require.Equal(t, int64(40), store.GetEarliestVersion()) +} + func TestEVMStateStoreDefaultUsesUnifiedDB(t *testing.T) { dir := t.TempDir() cfg := testConfig() diff --git a/sei-db/state_db/ss/evm/store.go b/sei-db/state_db/ss/evm/store.go index 01337940f3..4c20aa6ceb 100644 --- a/sei-db/state_db/ss/evm/store.go +++ b/sei-db/state_db/ss/evm/store.go @@ -1,7 +1,9 @@ package evm import ( + "errors" "fmt" + "os" "path/filepath" "sync" @@ -10,6 +12,7 @@ import ( commonevm "github.com/sei-protocol/sei-chain/sei-db/common/keys" "github.com/sei-protocol/sei-chain/sei-db/config" "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" + "github.com/sei-protocol/sei-chain/sei-db/management" "github.com/sei-protocol/sei-chain/sei-db/proto" "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/backend" ) @@ -152,16 +155,13 @@ func (s *EVMStateStore) SetLatestVersion(version int64) error { } func (s *EVMStateStore) GetEarliestVersion() int64 { - var minVersion int64 = -1 + var maxVersion int64 for _, db := range s.managedDBs { - if v := db.GetEarliestVersion(); minVersion < 0 || v < minVersion { - minVersion = v + if v := db.GetEarliestVersion(); v > maxVersion { + maxVersion = v } } - if minVersion < 0 { - return 0 - } - return minVersion + return maxVersion } func (s *EVMStateStore) SetEarliestVersion(version int64, ignoreVersion bool) error { @@ -370,6 +370,100 @@ func (s *EVMStateStore) Close() error { return lastErr } +func (s *EVMStateStore) SupportsCheckpoint() bool { + for _, db := range s.managedDBs { + if _, checkpointable := db.(types.Checkpointable); !checkpointable { + return false + } + if _, barrier := db.(types.DrainBarrier); !barrier { + return false + } + if _, versionSetter := db.(types.CheckpointVersionSetter); !versionSetter { + return false + } + } + return len(s.managedDBs) > 0 +} + +// ScheduleCheckpoint places one barrier on each managed apply queue. +// +// Every sub-DB checkpoints at the same block without the sub-DBs having to agree +// on anything. The caller runs this after it has enqueued the target block on +// every sub-DB and before it enqueues any later block, so each barrier lands at +// the same point in its own queue: after that block and before the next one. A +// sub-DB then checkpoints its own state as of that block. Wall-clock times +// differ, and no lock is shared. A sub-DB that received no change at the target +// block stays at its last version, which is that sub-DB's correct state for the +// block. SetCheckpointVersion afterwards labels every sub-DB with the same +// block, so a reopened snapshot reports one version rather than five. +func (s *EVMStateStore) ScheduleCheckpoint(destDir string, shouldRun func() bool, done func(error)) { + if !s.separateDBs { + db := s.primaryDB() + if db == nil { + // Unreachable: NewEVMStateStore either opens a managed DB or fails. + // Reporting success would publish a snapshot with no evm tree in it, + // which is only discovered by whoever tries to restore from it. + done(errors.New("EVM state store has no managed DB to checkpoint")) + return + } + management.ScheduleCheckpoint(db, destDir, shouldRun, done) + return + } + + if err := os.MkdirAll(destDir, 0o750); err != nil { + done(fmt.Errorf("create EVM checkpoint dir %q: %w", destDir, err)) + return + } + + storeTypes := AllEVMStoreTypes() + var ( + mu sync.Mutex + remaining = len(storeTypes) + firstErr error + ) + for _, storeType := range storeTypes { + name := StoreTypeName(storeType) + dest := filepath.Join(destDir, name) + management.ScheduleCheckpoint(s.subDBs[storeType], dest, shouldRun, func(err error) { + mu.Lock() + if err != nil && firstErr == nil { + firstErr = fmt.Errorf("checkpoint EVM sub-DB %s: %w", name, err) + } + remaining-- + last, outcome := remaining == 0, firstErr + mu.Unlock() + if last { + done(outcome) + } + }) + } +} + +func (s *EVMStateStore) SetCheckpointVersion(destDir string, version int64) error { + if !s.separateDBs { + db := s.primaryDB() + if db == nil { + return errors.New("EVM state store has no managed DB to stamp") + } + return management.SetCheckpointVersion(db, destDir, version) + } + for _, storeType := range AllEVMStoreTypes() { + dest := filepath.Join(destDir, StoreTypeName(storeType)) + if err := management.SetCheckpointVersion(s.subDBs[storeType], dest, version); err != nil { + return fmt.Errorf("set EVM sub-DB %s checkpoint version: %w", StoreTypeName(storeType), err) + } + } + return nil +} + +func (s *EVMStateStore) WaitForPendingWrites() { + for _, db := range s.managedDBs { + if w, ok := db.(interface{ WaitForPendingWrites() }); ok { + w.WaitForPendingWrites() + } + } +} + func filterEVMChangesets(changesets []*proto.NamedChangeSet) []*proto.NamedChangeSet { filtered := make([]*proto.NamedChangeSet, 0, len(changesets)) for _, cs := range changesets {