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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 49 additions & 0 deletions sei-cosmos/storev2/rootmulti/flatkv_snapshot_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -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,
Expand Down
33 changes: 30 additions & 3 deletions sei-cosmos/storev2/rootmulti/restore_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down
7 changes: 5 additions & 2 deletions sei-cosmos/storev2/rootmulti/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -1302,9 +1302,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,
Expand Down
9 changes: 5 additions & 4 deletions sei-db/state_db/sc/composite/importer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
26 changes: 21 additions & 5 deletions sei-db/state_db/sc/composite/importer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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")
}
3 changes: 2 additions & 1 deletion sei-db/state_db/sc/composite/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down
78 changes: 43 additions & 35 deletions sei-db/state_db/sc/flatkv/import_export_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -598,45 +598,53 @@ 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) {
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 nil key are silently skipped.
imp.AddNode(&types.SnapshotNode{
Key: nil,
Value: []byte{0xAA},
Height: 0,
})

require.NoError(t, imp.Close())
require.Equal(t, int64(1), s.Version())
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) {
Expand Down
41 changes: 32 additions & 9 deletions sei-db/state_db/sc/flatkv/importer.go
Original file line number Diff line number Diff line change
Expand Up @@ -337,22 +337,45 @@ 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 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
}
if err := imp.checkNode(node); err != nil {
imp.setErr(err)
return err
}
Comment thread
manav2401 marked this conversation as resolved.
select {
case imp.ingestCh <- rawKVPair{Key: node.Key, Value: node.Value}:
case <-imp.done:
}
// 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.
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 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",
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.
Expand Down
Loading
Loading