From a05e76641a064a08676f2b2ecced1015a875b59f Mon Sep 17 00:00:00 2001 From: Manav Darji <36959497+manav2401@users.noreply.github.com> Date: Thu, 1 Oct 2026 21:44:19 +0530 Subject: [PATCH 1/2] fix(rootmulti): reject invalid snapshot restore --- .../storev2/rootmulti/flatkv_snapshot_test.go | 49 +++++++++++++ sei-cosmos/storev2/rootmulti/restore_test.go | 33 ++++++++- sei-cosmos/storev2/rootmulti/store.go | 7 +- sei-db/state_db/sc/composite/importer.go | 9 +-- sei-db/state_db/sc/composite/importer_test.go | 26 +++++-- sei-db/state_db/sc/composite/store_test.go | 3 +- .../state_db/sc/flatkv/import_export_test.go | 44 ++++++------ sei-db/state_db/sc/flatkv/importer.go | 40 ++++++++--- sei-db/state_db/sc/flatkv/importer_test.go | 70 +++++++++++++++++++ .../sc/flatkv/lthash_agreement_test.go | 5 +- sei-db/state_db/sc/memiavl/import.go | 6 +- sei-db/state_db/sc/types/types.go | 4 +- .../operations/import_flatkv_from_memiavl.go | 22 ++++-- 13 files changed, 261 insertions(+), 57 deletions(-) diff --git a/sei-cosmos/storev2/rootmulti/flatkv_snapshot_test.go b/sei-cosmos/storev2/rootmulti/flatkv_snapshot_test.go index d08df94dd4..8f6c331d2b 100644 --- a/sei-cosmos/storev2/rootmulti/flatkv_snapshot_test.go +++ b/sei-cosmos/storev2/rootmulti/flatkv_snapshot_test.go @@ -9,9 +9,12 @@ import ( "testing" protoio "github.com/gogo/protobuf/io" + snapshottypes "github.com/sei-protocol/sei-chain/sei-cosmos/snapshots/types" "github.com/sei-protocol/sei-chain/sei-cosmos/store/types" "github.com/sei-protocol/sei-chain/sei-db/common/keys" seidbconfig "github.com/sei-protocol/sei-chain/sei-db/config" + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/ktype" + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/vtype" abci "github.com/sei-protocol/sei-chain/sei-tendermint/abci/types" "github.com/stretchr/testify/require" ) @@ -355,6 +358,52 @@ func TestFlatKVOnlySnapshotRestorePopulatesSS(t *testing.T) { queryEqual("/evm/key", evmData.nonKey, makeNonce(uint64(snapHeight))) } +// TestFlatKVOnlySnapshotRestoreRejectsNodeAtOtherVersion restores a flatkv_only snapshot whose flatkv +// section ends with a node stamped with a version other than the snapshot height. The restore must fail, +// and the state store must not hold that node's value. +func TestFlatKVOnlySnapshotRestoreRejectsNodeAtOtherVersion(t *testing.T) { + cfg := flatKVOnlyConfig() + ssCfg := seidbconfig.DefaultStateStoreConfig() + ssCfg.Enable = true + ssCfg.AsyncWriteBuffer = 0 + evmData := newEVMTestData(0x42) + + const snapHeight = 8 + + srcStore, srcKeys := newTestRootMultiWithSS(t, t.TempDir(), cfg, ssCfg) + for block := 1; block <= snapHeight; block++ { + simulateFlatKVOnlyBlock(t, srcStore, srcKeys, block, evmData) + } + waitUntilSSVersion(t, srcStore, snapHeight) + + var buf bytes.Buffer + writer := protoio.NewDelimitedWriter(&buf) + require.NoError(t, srcStore.Snapshot(snapHeight, writer)) + require.NoError(t, srcStore.Close()) + + // The flatkv section is the only one in a flatkv_only snapshot, so an item written after the snapshot + // belongs to it. + value := []byte("v9") + require.NoError(t, writer.WriteMsg(&snapshottypes.SnapshotItem{ + Item: &snapshottypes.SnapshotItem_IAVL{IAVL: &snapshottypes.SnapshotIAVLItem{ + Key: ktype.ModulePhysicalKey("bank", []byte("supply")), + Value: vtype.SerializeMisc(snapHeight, value), + Version: snapHeight + 1, + }}, + })) + + dstStore, _ := newTestRootMultiWithSS(t, t.TempDir(), cfg, ssCfg) + // The import reopened the commitment store, so closing releases it along with the state store. + defer func() { _ = dstStore.Close() }() + reader := protoio.NewDelimitedReader(bytes.NewReader(buf.Bytes()), 1<<30) + _, err := dstStore.Restore(snapHeight, 1, reader) + require.ErrorContains(t, err, "the import is at version") + + got, err := dstStore.GetStateStore().Get("bank", snapHeight, []byte("supply")) + require.NoError(t, err) + require.NotEqual(t, value, got, "a node the commitment store rejected must not reach the state store") +} + func simulateFlatKVOnlyBlock( t *testing.T, store *Store, diff --git a/sei-cosmos/storev2/rootmulti/restore_test.go b/sei-cosmos/storev2/rootmulti/restore_test.go index e9b7ab1772..2ac6632b65 100644 --- a/sei-cosmos/storev2/rootmulti/restore_test.go +++ b/sei-cosmos/storev2/rootmulti/restore_test.go @@ -11,9 +11,11 @@ import ( protoio "github.com/gogo/protobuf/io" snapshottypes "github.com/sei-protocol/sei-chain/sei-cosmos/snapshots/types" + "github.com/sei-protocol/sei-chain/sei-db/common/keys" "github.com/sei-protocol/sei-chain/sei-db/common/utils" seidbconfig "github.com/sei-protocol/sei-chain/sei-db/config" seidbtypes "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/vtype" "github.com/stretchr/testify/require" ) @@ -118,18 +120,20 @@ func TestRestoreRejectsMalformedStream(t *testing.T) { } } -// fakeStateStore is a state store whose Import returns importErr, either at once or after reading every -// node, and which counts the version writes restore makes. +// fakeStateStore is a state store whose Import returns importErr, either at once or after reading and +// recording every node, and which counts the version writes restore makes. type fakeStateStore struct { seidbtypes.StateStore returnEarly bool importErr error + imported []seidbtypes.SnapshotNode versionWrites int } func (f *fakeStateStore) Import(_ int64, ch <-chan seidbtypes.SnapshotNode) error { if !f.returnEarly { - for range ch { + for node := range ch { + f.imported = append(f.imported, node) } } return f.importErr @@ -196,6 +200,29 @@ func TestRestoreFailurePublishesNothing(t *testing.T) { } } +// TestRestoreStopsWhenCommitStoreRejectsNode pins that a node the SC importer rejects fails the restore, and +// that the state store receives only the nodes accepted before it. +func TestRestoreStopsWhenCommitStoreRejectsNode(t *testing.T) { + store, _ := newTestRootMulti(t, t.TempDir(), flatKVOnlyConfig()) + ss := &fakeStateStore{} + store.ssStore = ss + + leaf := func(key string, version int64) snapshottypes.SnapshotItem { + return snapshottypes.SnapshotItem{Item: &snapshottypes.SnapshotItem_IAVL{IAVL: &snapshottypes.SnapshotIAVLItem{ + Key: []byte(key), + Value: vtype.SerializeMisc(1, []byte("v")), + Version: version, + }}} + } + items := []snapshottypes.SnapshotItem{storeItem(keys.FlatKVStoreKey), leaf("bank/a", 1), leaf("bank/b", 2)} + + err := restoreWithin(t, store, snapshotStream(t, items)) + require.ErrorContains(t, err, "the import is at version") + require.Len(t, ss.imported, 1) + require.Equal(t, []byte("bank/a"), ss.imported[0].Key) + require.Zero(t, ss.versionWrites, "a failed restore must not record the snapshot height") +} + // TestRestoreSuccessPublishes pins that a successful restore publishes the memIAVL snapshot and records // the snapshot height on the state store. func TestRestoreSuccessPublishes(t *testing.T) { diff --git a/sei-cosmos/storev2/rootmulti/store.go b/sei-cosmos/storev2/rootmulti/store.go index b28ade6276..c77a44178b 100644 --- a/sei-cosmos/storev2/rootmulti/store.go +++ b/sei-cosmos/storev2/rootmulti/store.go @@ -1281,9 +1281,12 @@ loop: if node.Height == 0 && node.Value == nil { node.Value = []byte{} } - scImporter.AddNode(node) + if err = scImporter.AddNode(node); err != nil { + restoreErr = err + break loop + } - // Check if we should also import to SS store + // Only leaves the SC importer accepted reach the state store. if ssImport != nil && node.Height == 0 { if err = ssImport.send(seidbtypes.SnapshotNode{ StoreKey: storeKey, diff --git a/sei-db/state_db/sc/composite/importer.go b/sei-db/state_db/sc/composite/importer.go index bb4174e313..ed59ecb6de 100644 --- a/sei-db/state_db/sc/composite/importer.go +++ b/sei-db/state_db/sc/composite/importer.go @@ -61,16 +61,17 @@ func (si *SnapshotImporter) AddModule(name string) error { return nil } -func (si *SnapshotImporter) AddNode(node *types.SnapshotNode) { +func (si *SnapshotImporter) AddNode(node *types.SnapshotNode) error { if si.currentModule == keys.FlatKVStoreKey { if si.flatkvImporter != nil { - si.flatkvImporter.AddNode(node) + return si.flatkvImporter.AddNode(node) } - return + return nil } if si.cosmosImporter != nil { - si.cosmosImporter.AddNode(node) + return si.cosmosImporter.AddNode(node) } + return nil } // Close publishes both backends' imports. When the cosmos import fails, the flatkv import is discarded diff --git a/sei-db/state_db/sc/composite/importer_test.go b/sei-db/state_db/sc/composite/importer_test.go index 53a83cdcf8..9cbb2a1a11 100644 --- a/sei-db/state_db/sc/composite/importer_test.go +++ b/sei-db/state_db/sc/composite/importer_test.go @@ -4,20 +4,22 @@ import ( "errors" "testing" + "github.com/sei-protocol/sei-chain/sei-db/common/keys" "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" "github.com/stretchr/testify/require" ) -// endRecordingImporter is an importer whose Close returns closeErr, and which records whether Close or -// Abort ended it. +// endRecordingImporter is an importer whose AddNode returns addNodeErr and whose Close returns closeErr, and +// which records whether Close or Abort ended it. type endRecordingImporter struct { - closeErr error - endedBy string + addNodeErr error + closeErr error + endedBy string } func (e *endRecordingImporter) AddModule(string) error { return nil } -func (e *endRecordingImporter) AddNode(*types.SnapshotNode) {} +func (e *endRecordingImporter) AddNode(*types.SnapshotNode) error { return e.addNodeErr } func (e *endRecordingImporter) Close() error { e.endedBy = "close" @@ -51,3 +53,17 @@ func TestSnapshotImporterCloseDiscardsFlatKVWhenCosmosFails(t *testing.T) { require.Equal(t, "abort", flatkv.endedBy) }) } + +// TestSnapshotImporterAddNodeReturnsBackendError pins that AddNode returns the error of the backend that the +// current section routes to. +func TestSnapshotImporterAddNodeReturnsBackendError(t *testing.T) { + cosmos := &endRecordingImporter{addNodeErr: errors.New("cosmos rejected node")} + flatkv := &endRecordingImporter{addNodeErr: errors.New("flatkv rejected node")} + imp := NewImporter(cosmos, flatkv, nil) + node := &types.SnapshotNode{Key: []byte("k"), Value: []byte("v")} + + require.NoError(t, imp.AddModule("bank")) + require.ErrorContains(t, imp.AddNode(node), "cosmos rejected node") + require.NoError(t, imp.AddModule(keys.FlatKVStoreKey)) + require.ErrorContains(t, imp.AddNode(node), "flatkv rejected node") +} diff --git a/sei-db/state_db/sc/composite/store_test.go b/sei-db/state_db/sc/composite/store_test.go index f120e10e0f..ed877ffeb8 100644 --- a/sei-db/state_db/sc/composite/store_test.go +++ b/sei-db/state_db/sc/composite/store_test.go @@ -1281,8 +1281,9 @@ func (ti *trackingImporter) AddModule(name string) error { return nil } -func (ti *trackingImporter) AddNode(node *types.SnapshotNode) { +func (ti *trackingImporter) AddNode(node *types.SnapshotNode) error { *ti.nodes = append(*ti.nodes, node) + return nil } func (ti *trackingImporter) Close() error { return nil } diff --git a/sei-db/state_db/sc/flatkv/import_export_test.go b/sei-db/state_db/sc/flatkv/import_export_test.go index af34441958..752f624eea 100644 --- a/sei-db/state_db/sc/flatkv/import_export_test.go +++ b/sei-db/state_db/sc/flatkv/import_export_test.go @@ -585,7 +585,7 @@ func TestImporterOnReadOnlyStore(t *testing.T) { require.NoError(t, s.Close()) } -func TestImporterHeightNonZeroSkipped(t *testing.T) { +func TestImporterHeightNonZeroRejected(t *testing.T) { dir := t.TempDir() cfg := config.DefaultTestConfig(t) cfg.DataDir = filepath.Join(dir, flatkvRootDir) @@ -598,23 +598,25 @@ func TestImporterHeightNonZeroSkipped(t *testing.T) { imp, err := s.Importer(1) require.NoError(t, err) - // Non-leaf nodes (Height != 0) are silently skipped. - imp.AddNode(&types.SnapshotNode{ - Key: keys.BuildEVMKey(keys.EVMKeyStorage, ktype.StorageKey(addrN(0x01), slotN(0x01))), - Value: padLeft32(0x11), - Height: 1, // non-leaf - }) + key := storagePhysKey(addrN(0x01), slotN(0x01)) - require.NoError(t, imp.Close()) + // Non-leaf nodes (Height != 0) fail the import. + addErr := imp.AddNode(&types.SnapshotNode{ + Key: key, + Value: padLeft32(0x11), + Version: 1, + Height: 1, // non-leaf + }) + require.ErrorContains(t, addErr, "only leaves can be imported") + require.ErrorIs(t, imp.Close(), addErr) - // Data should NOT have been imported. - key := keys.BuildEVMKey(keys.EVMKeyStorage, ktype.StorageKey(addrN(0x01), slotN(0x01))) - _, found := s.Get(keys.EVMStoreKey, key) - require.False(t, found, "height != 0 node should be skipped") + _, found := s.Get(keys.EVMStoreKey, keys.BuildEVMKey(keys.EVMKeyStorage, ktype.StorageKey(addrN(0x01), slotN(0x01)))) + require.False(t, found, "height != 0 node should be rejected") + require.Zero(t, s.Version(), "a rejected import must not finalize") require.NoError(t, s.Close()) } -func TestImporterNilKeySkipped(t *testing.T) { +func TestImporterNilKeyRejected(t *testing.T) { dir := t.TempDir() cfg := config.DefaultTestConfig(t) cfg.DataDir = filepath.Join(dir, flatkvRootDir) @@ -627,15 +629,17 @@ func TestImporterNilKeySkipped(t *testing.T) { imp, err := s.Importer(1) require.NoError(t, err) - // Nodes with nil key are silently skipped. - imp.AddNode(&types.SnapshotNode{ - Key: nil, - Value: []byte{0xAA}, - Height: 0, + // Nodes with a nil key fail the import. + addErr := imp.AddNode(&types.SnapshotNode{ + Key: nil, + Value: []byte{0xAA}, + Version: 1, + Height: 0, }) + require.ErrorContains(t, addErr, "node has no key") - require.NoError(t, imp.Close()) - require.Equal(t, int64(1), s.Version()) + require.ErrorIs(t, imp.Close(), addErr) + require.Zero(t, s.Version(), "a rejected import must not finalize") require.NoError(t, s.Close()) } diff --git a/sei-db/state_db/sc/flatkv/importer.go b/sei-db/state_db/sc/flatkv/importer.go index cf69f1dfcc..8c69fd3f6a 100644 --- a/sei-db/state_db/sc/flatkv/importer.go +++ b/sei-db/state_db/sc/flatkv/importer.go @@ -337,22 +337,44 @@ func (imp *KVImporter) AddModule(_ string) error { return nil } -func (imp *KVImporter) AddNode(node *types.SnapshotNode) { - if node.Height != 0 || node.Key == nil || node.Version != imp.version { - return +// AddNode queues node for import. It returns an error, and fails the import, when node is not a leaf with a +// key, a non-empty value and the import's version. Once the import has failed, it returns that failure. +func (imp *KVImporter) AddNode(node *types.SnapshotNode) error { + if err := imp.getErr(); err != nil { + return err + } + if err := imp.checkNode(node); err != nil { + imp.setErr(err) + return err + } + select { + case imp.ingestCh <- rawKVPair{Key: node.Key, Value: node.Value}: + return nil + case <-imp.done: + return imp.getErr() + } +} + +// checkNode returns an error unless node is a row this import can store. +func (imp *KVImporter) checkNode(node *types.SnapshotNode) error { + if node.Height != 0 { + return fmt.Errorf("flatkv import: node %x has height %d; only leaves can be imported", node.Key, node.Height) + } + if node.Key == nil { + return errors.New("flatkv import: node has no key") + } + if node.Version != imp.version { + return fmt.Errorf("flatkv import: node %x has version %d; the import is at version %d", + node.Key, node.Version, imp.version) } // FlatKV import nodes carry already-serialized physical values. Even an // empty logical misc value has a non-empty serialized header, so a // zero-length physical value is malformed. Reject it instead of writing a // Pebble row that LtHash and verification would both skip. if len(node.Value) == 0 { - imp.setErr(fmt.Errorf("flatkv import: empty physical value for key %x", node.Key)) - return - } - select { - case imp.ingestCh <- rawKVPair{Key: node.Key, Value: node.Value}: - case <-imp.done: + return fmt.Errorf("flatkv import: empty physical value for key %x", node.Key) } + return nil } // Abort tears down the worker pipeline without finalizing the import. diff --git a/sei-db/state_db/sc/flatkv/importer_test.go b/sei-db/state_db/sc/flatkv/importer_test.go index 51993b62fe..dd617a47eb 100644 --- a/sei-db/state_db/sc/flatkv/importer_test.go +++ b/sei-db/state_db/sc/flatkv/importer_test.go @@ -136,6 +136,76 @@ func TestKVImporter_EmptyPhysicalValueRejected(t *testing.T) { } } +// TestKVImporter_NodeVersionMustMatchImport verifies that Importer accepts a node only at the import's own +// version, and that a node at any other version fails the whole import, including the nodes before it. +func TestKVImporter_NodeVersionMustMatchImport(t *testing.T) { + const importVersion = 5 + nodeAt := func(key string, version int64) *types.SnapshotNode { + return &types.SnapshotNode{ + Key: ktype.ModulePhysicalKey("bank", []byte(key)), + Value: vtype.SerializeMisc(importVersion, []byte("v")), + Version: version, + } + } + + t.Run("same version", func(t *testing.T) { + s, imp := newKVImporterForTest(t, importVersion) + defer func() { require.NoError(t, s.Close()) }() + + require.NoError(t, imp.AddNode(nodeAt("a", importVersion))) + require.NoError(t, imp.Close()) + require.Equal(t, int64(importVersion), s.Version()) + got, found := s.Get("bank", []byte("a")) + require.True(t, found) + require.Equal(t, []byte("v"), got) + }) + + for name, version := range map[string]int64{ + "earlier version": importVersion - 1, + "later version": importVersion + 1, + } { + t.Run(name, func(t *testing.T) { + s, imp := newKVImporterForTest(t, importVersion) + defer func() { require.NoError(t, s.Close()) }() + + require.NoError(t, imp.AddNode(nodeAt("a", importVersion))) + err := imp.AddNode(nodeAt("b", version)) + require.ErrorContains(t, err, "the import is at version") + require.ErrorIs(t, imp.AddNode(nodeAt("c", importVersion)), err, + "once the import has failed, AddNode must keep returning the failure") + require.ErrorIs(t, imp.Close(), err) + require.Zero(t, s.Version(), "a rejected import must not finalize") + }) + } +} + +// TestKVImporter_ImportAfterFailedImport verifies that a failed import does not carry its failure into the +// next import on the same store. +func TestKVImporter_ImportAfterFailedImport(t *testing.T) { + s, failed := newKVImporterForTest(t, 1) + defer func() { require.NoError(t, s.Close()) }() + + node := &types.SnapshotNode{ + Key: ktype.ModulePhysicalKey("bank", []byte("k")), + Value: vtype.SerializeMisc(1, []byte("v")), + Version: 2, + } + require.Error(t, failed.AddNode(node)) + require.Error(t, failed.Close()) + require.Zero(t, s.Version()) + + next, err := s.Importer(1) + require.NoError(t, err) + node.Version = 1 + require.NoError(t, next.AddNode(node)) + require.NoError(t, next.Close()) + + require.Equal(t, int64(1), s.Version()) + got, found := s.Get("bank", []byte("k")) + require.True(t, found) + require.Equal(t, []byte("v"), got) +} + // TestKVImporter_ErrLifecycle locks in the contract that Err() returns the // first pipeline error as soon as it propagates, before Close is invoked. // This is the path the seidb tool relies on to short-circuit a failing import diff --git a/sei-db/state_db/sc/flatkv/lthash_agreement_test.go b/sei-db/state_db/sc/flatkv/lthash_agreement_test.go index a184d22e30..516fa1f0ba 100644 --- a/sei-db/state_db/sc/flatkv/lthash_agreement_test.go +++ b/sei-db/state_db/sc/flatkv/lthash_agreement_test.go @@ -638,9 +638,8 @@ func sortedKeys(m map[string][]byte) []string { // importFromInto exports src at version and imports the stream into dst, returning how many rows // crossed the boundary. // -// The count is returned because AddNode silently drops any node whose version does not match the -// importer's: without asserting it, an empty destination would compare equal to nothing and the whole -// suite would pass vacuously. +// The count is returned so callers can assert the stream was not empty: an empty destination would +// compare equal to an empty source and the whole suite would pass vacuously. func importFromInto(t *testing.T, src *CommitStore, dst *CommitStore, version int64) int { t.Helper() diff --git a/sei-db/state_db/sc/memiavl/import.go b/sei-db/state_db/sc/memiavl/import.go index 58d9432624..38d09cf189 100644 --- a/sei-db/state_db/sc/memiavl/import.go +++ b/sei-db/state_db/sc/memiavl/import.go @@ -74,8 +74,7 @@ func (mti *MultiTreeImporter) tmpDir() string { func (mti *MultiTreeImporter) Add(item interface{}) error { switch item := item.(type) { case *types.SnapshotNode: - mti.AddNode(item) - return nil + return mti.AddNode(item) case string: return mti.AddModule(item) default: @@ -93,8 +92,9 @@ func (mti *MultiTreeImporter) AddModule(name string) error { return nil } -func (mti *MultiTreeImporter) AddNode(node *types.SnapshotNode) { +func (mti *MultiTreeImporter) AddNode(node *types.SnapshotNode) error { mti.importer.Add(node) + return nil } func (mti *MultiTreeImporter) Close() (err error) { diff --git a/sei-db/state_db/sc/types/types.go b/sei-db/state_db/sc/types/types.go index c60cb8681e..916b1baba9 100644 --- a/sei-db/state_db/sc/types/types.go +++ b/sei-db/state_db/sc/types/types.go @@ -221,7 +221,9 @@ type CommitKVStore interface { type Importer interface { AddModule(name string) error - AddNode(node *SnapshotNode) + // AddNode queues node for import. A non-nil error means node was not accepted and the import cannot + // complete. + AddNode(node *SnapshotNode) error // Abort discards the import in place of Close, publishing none of it. The returned error may be // reason itself, or another error that already ended the import. diff --git a/sei-db/tools/cmd/seidb/operations/import_flatkv_from_memiavl.go b/sei-db/tools/cmd/seidb/operations/import_flatkv_from_memiavl.go index f1fe59a07c..45e8db685d 100644 --- a/sei-db/tools/cmd/seidb/operations/import_flatkv_from_memiavl.go +++ b/sei-db/tools/cmd/seidb/operations/import_flatkv_from_memiavl.go @@ -138,16 +138,18 @@ func importerErr(importer sctypes.Importer) error { // emitPairs forwards translator output to the FlatKV importer, returning the // number of pairs written. -func emitPairs(importer sctypes.Importer, pairs []flatkv.PhysicalKVPair, height int64) int64 { +func emitPairs(importer sctypes.Importer, pairs []flatkv.PhysicalKVPair, height int64) (int64, error) { for _, p := range pairs { - importer.AddNode(&sctypes.SnapshotNode{ + if err := importer.AddNode(&sctypes.SnapshotNode{ Key: p.Key, Value: p.Value, Version: height, Height: 0, - }) + }); err != nil { + return 0, err + } } - return int64(len(pairs)) + return int64(len(pairs)), nil } func importMemiavlModulesToFlatKV(ctx context.Context, homeDir string, modules []string, height int64, force bool) (err error) { @@ -265,7 +267,11 @@ func importMemiavlModulesToFlatKV(ctx context.Context, homeDir string, modules [ if err != nil { return fmt.Errorf("translate batch (module=%s): %w", batch.Name, err) } - written += emitPairs(importer, pairs, height) + n, err := emitPairs(importer, pairs, height) + if err != nil { + return fmt.Errorf("FlatKV import failed: %w", err) + } + written += n batch.Changeset.Pairs = batch.Changeset.Pairs[:0] return nil } @@ -352,7 +358,11 @@ func importMemiavlModulesToFlatKV(ctx context.Context, homeDir string, modules [ return fmt.Errorf("FlatKV import failed: %w", err) } - written += emitPairs(importer, translator.Finalize(), height) + n, err := emitPairs(importer, translator.Finalize(), height) + if err != nil { + return fmt.Errorf("FlatKV import failed: %w", err) + } + written += n if err := importer.Close(); err != nil { return fmt.Errorf("failed to finalize FlatKV import: %w", err) From d0b5d357ff6b6352e0daba64fbc437fe867cf7b4 Mon Sep 17 00:00:00 2001 From: Manav Darji <36959497+manav2401@users.noreply.github.com> Date: Mon, 5 Oct 2026 12:23:44 +0530 Subject: [PATCH 2/2] Address comments --- .../state_db/sc/flatkv/import_export_test.go | 54 ++++++++++--------- sei-db/state_db/sc/flatkv/importer.go | 13 ++--- 2 files changed, 36 insertions(+), 31 deletions(-) diff --git a/sei-db/state_db/sc/flatkv/import_export_test.go b/sei-db/state_db/sc/flatkv/import_export_test.go index 752f624eea..b17ef1d090 100644 --- a/sei-db/state_db/sc/flatkv/import_export_test.go +++ b/sei-db/state_db/sc/flatkv/import_export_test.go @@ -616,31 +616,35 @@ func TestImporterHeightNonZeroRejected(t *testing.T) { require.NoError(t, s.Close()) } -func TestImporterNilKeyRejected(t *testing.T) { - dir := t.TempDir() - cfg := config.DefaultTestConfig(t) - cfg.DataDir = filepath.Join(dir, flatkvRootDir) - - s, err := newCommitStoreWithWAL(t.Context(), cfg) - require.NoError(t, err) - err = s.LoadLatest() - require.NoError(t, err) - - imp, err := s.Importer(1) - require.NoError(t, err) - - // Nodes with a nil key fail the import. - addErr := imp.AddNode(&types.SnapshotNode{ - Key: nil, - Value: []byte{0xAA}, - Version: 1, - Height: 0, - }) - require.ErrorContains(t, addErr, "node has no key") - - require.ErrorIs(t, imp.Close(), addErr) - require.Zero(t, s.Version(), "a rejected import must not finalize") - require.NoError(t, s.Close()) +func TestImporterEmptyKeyRejected(t *testing.T) { + // A restore hands the importer an empty key where the snapshot carried none, so both forms must fail. + for name, key := range map[string][]byte{"nil": nil, "zero-length": {}} { + t.Run(name, func(t *testing.T) { + dir := t.TempDir() + cfg := config.DefaultTestConfig(t) + cfg.DataDir = filepath.Join(dir, flatkvRootDir) + + s, err := newCommitStoreWithWAL(t.Context(), cfg) + require.NoError(t, err) + err = s.LoadLatest() + require.NoError(t, err) + + imp, err := s.Importer(1) + require.NoError(t, err) + + addErr := imp.AddNode(&types.SnapshotNode{ + Key: key, + Value: []byte{0xAA}, + Version: 1, + Height: 0, + }) + require.ErrorContains(t, addErr, "node has an empty key") + + require.ErrorIs(t, imp.Close(), addErr) + require.Zero(t, s.Version(), "a rejected import must not finalize") + require.NoError(t, s.Close()) + }) + } } func TestImporterEmptyStore(t *testing.T) { diff --git a/sei-db/state_db/sc/flatkv/importer.go b/sei-db/state_db/sc/flatkv/importer.go index 8c69fd3f6a..51a9555898 100644 --- a/sei-db/state_db/sc/flatkv/importer.go +++ b/sei-db/state_db/sc/flatkv/importer.go @@ -337,8 +337,8 @@ func (imp *KVImporter) AddModule(_ string) error { return nil } -// AddNode queues node for import. It returns an error, and fails the import, when node is not a leaf with a -// key, a non-empty value and the import's version. Once the import has failed, it returns that failure. +// AddNode queues node for import. It fails the import and returns an error unless node is a leaf at the +// import's version with a non-empty key and value. Once the import has failed, it returns that failure. func (imp *KVImporter) AddNode(node *types.SnapshotNode) error { if err := imp.getErr(); err != nil { return err @@ -349,10 +349,11 @@ func (imp *KVImporter) AddNode(node *types.SnapshotNode) error { } select { case imp.ingestCh <- rawKVPair{Key: node.Key, Value: node.Value}: - return nil case <-imp.done: - return imp.getErr() } + // A worker can fail while the send is pending, and select may still pick the send, so the failure is + // read again here rather than reporting the node as accepted. + return imp.getErr() } // checkNode returns an error unless node is a row this import can store. @@ -360,8 +361,8 @@ func (imp *KVImporter) checkNode(node *types.SnapshotNode) error { if node.Height != 0 { return fmt.Errorf("flatkv import: node %x has height %d; only leaves can be imported", node.Key, node.Height) } - if node.Key == nil { - return errors.New("flatkv import: node has no key") + if len(node.Key) == 0 { + return errors.New("flatkv import: node has an empty key") } if node.Version != imp.version { return fmt.Errorf("flatkv import: node %x has version %d; the import is at version %d",