feat(store/v2): snapshot manager (#18458)
This commit is contained in:
@@ -0,0 +1,40 @@
|
||||
package iavl
|
||||
|
||||
import (
|
||||
"errors"
|
||||
|
||||
"github.com/cosmos/iavl"
|
||||
|
||||
"cosmossdk.io/store/v2/commitment"
|
||||
snapshotstypes "cosmossdk.io/store/v2/snapshots/types"
|
||||
)
|
||||
|
||||
// Exporter is a wrapper around iavl.Exporter.
|
||||
type Exporter struct {
|
||||
exporter *iavl.Exporter
|
||||
}
|
||||
|
||||
// Next returns the next item in the exporter.
|
||||
func (e *Exporter) Next() (*snapshotstypes.SnapshotIAVLItem, error) {
|
||||
item, err := e.exporter.Next()
|
||||
if err != nil {
|
||||
if errors.Is(err, iavl.ErrorExportDone) {
|
||||
return nil, commitment.ErrorExportDone
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &snapshotstypes.SnapshotIAVLItem{
|
||||
Key: item.Key,
|
||||
Value: item.Value,
|
||||
Version: item.Version,
|
||||
Height: int32(item.Height),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Close closes the exporter.
|
||||
func (e *Exporter) Close() error {
|
||||
e.exporter.Close()
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package iavl
|
||||
|
||||
import (
|
||||
"github.com/cosmos/iavl"
|
||||
|
||||
snapshotstypes "cosmossdk.io/store/v2/snapshots/types"
|
||||
)
|
||||
|
||||
// Importer is a wrapper around iavl.Importer.
|
||||
type Importer struct {
|
||||
importer *iavl.Importer
|
||||
}
|
||||
|
||||
// Add adds the given item to the importer.
|
||||
func (i *Importer) Add(item *snapshotstypes.SnapshotIAVLItem) error {
|
||||
return i.importer.Add(&iavl.ExportNode{
|
||||
Key: item.Key,
|
||||
Value: item.Value,
|
||||
Version: item.Version,
|
||||
Height: int8(item.Height),
|
||||
})
|
||||
}
|
||||
|
||||
// Commit commits the importer.
|
||||
func (i *Importer) Commit() error {
|
||||
return i.importer.Commit()
|
||||
}
|
||||
|
||||
// Close closes the importer.
|
||||
func (i *Importer) Close() error {
|
||||
i.importer.Close()
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -77,6 +77,34 @@ func (t *IavlTree) Prune(version uint64) error {
|
||||
return t.tree.DeleteVersionsTo(int64(version))
|
||||
}
|
||||
|
||||
// Export exports the tree exporter at the given version.
|
||||
func (t *IavlTree) Export(version uint64) (commitment.Exporter, error) {
|
||||
tree, err := t.tree.GetImmutable(int64(version))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
exporter, err := tree.Export()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Exporter{
|
||||
exporter: exporter,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Import imports the tree importer at the given version.
|
||||
func (t *IavlTree) Import(version uint64) (commitment.Importer, error) {
|
||||
importer, err := t.tree.Import(int64(version))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Importer{
|
||||
importer: importer,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Close closes the iavl tree.
|
||||
func (t *IavlTree) Close() error {
|
||||
return nil
|
||||
|
||||
@@ -5,11 +5,29 @@ import (
|
||||
|
||||
dbm "github.com/cosmos/cosmos-db"
|
||||
"github.com/stretchr/testify/require"
|
||||
"github.com/stretchr/testify/suite"
|
||||
|
||||
"cosmossdk.io/log"
|
||||
"cosmossdk.io/store/v2/commitment"
|
||||
)
|
||||
|
||||
func generateTree(treeType string) *IavlTree {
|
||||
func TestCommitterSuite(t *testing.T) {
|
||||
s := &commitment.CommitStoreTestSuite{
|
||||
NewStore: func(db dbm.DB, storeKeys []string, logger log.Logger) (*commitment.CommitStore, error) {
|
||||
multiTrees := make(map[string]commitment.Tree)
|
||||
cfg := DefaultConfig()
|
||||
for _, storeKey := range storeKeys {
|
||||
prefixDB := dbm.NewPrefixDB(db, []byte(storeKey))
|
||||
multiTrees[storeKey] = NewIavlTree(prefixDB, logger, cfg)
|
||||
}
|
||||
return commitment.NewCommitStore(multiTrees, logger)
|
||||
},
|
||||
}
|
||||
|
||||
suite.Run(t, s)
|
||||
}
|
||||
|
||||
func generateTree() *IavlTree {
|
||||
cfg := DefaultConfig()
|
||||
db := dbm.NewMemDB()
|
||||
return NewIavlTree(db, log.NewNopLogger(), cfg)
|
||||
@@ -17,7 +35,7 @@ func generateTree(treeType string) *IavlTree {
|
||||
|
||||
func TestIavlTree(t *testing.T) {
|
||||
// generate a new tree
|
||||
tree := generateTree("iavl")
|
||||
tree := generateTree()
|
||||
require.NotNil(t, tree)
|
||||
|
||||
initVersion := tree.GetLatestVersion()
|
||||
|
||||
+149
-1
@@ -3,14 +3,22 @@ package commitment
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"math"
|
||||
|
||||
protoio "github.com/cosmos/gogoproto/io"
|
||||
ics23 "github.com/cosmos/ics23/go"
|
||||
|
||||
"cosmossdk.io/log"
|
||||
"cosmossdk.io/store/v2"
|
||||
"cosmossdk.io/store/v2/snapshots"
|
||||
snapshotstypes "cosmossdk.io/store/v2/snapshots/types"
|
||||
)
|
||||
|
||||
var _ store.Committer = (*CommitStore)(nil)
|
||||
var (
|
||||
_ store.Committer = (*CommitStore)(nil)
|
||||
_ snapshots.CommitSnapshotter = (*CommitStore)(nil)
|
||||
)
|
||||
|
||||
// CommitStore is a wrapper around multiple Tree objects mapped by a unique store
|
||||
// key. Each store key reflects dedicated and unique usage within a module. A caller
|
||||
@@ -127,6 +135,146 @@ func (c *CommitStore) Prune(version uint64) (ferr error) {
|
||||
return ferr
|
||||
}
|
||||
|
||||
// Snapshot implements snapshotstypes.CommitSnapshotter.
|
||||
func (c *CommitStore) Snapshot(version uint64, protoWriter protoio.Writer) error {
|
||||
if version == 0 {
|
||||
return fmt.Errorf("the snapshot version must be greater than 0")
|
||||
}
|
||||
|
||||
latestVersion, err := c.GetLatestVersion()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if version > latestVersion {
|
||||
return fmt.Errorf("the snapshot version %d is greater than the latest version %d", version, latestVersion)
|
||||
}
|
||||
|
||||
for storeKey, tree := range c.multiTrees {
|
||||
// TODO: check the parallelism of this loop
|
||||
if err := func() error {
|
||||
exporter, err := tree.Export(version)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to export tree for version %d: %w", version, err)
|
||||
}
|
||||
defer exporter.Close()
|
||||
|
||||
err = protoWriter.WriteMsg(&snapshotstypes.SnapshotItem{
|
||||
Item: &snapshotstypes.SnapshotItem_Store{
|
||||
Store: &snapshotstypes.SnapshotStoreItem{
|
||||
Name: storeKey,
|
||||
},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to write store name: %w", err)
|
||||
}
|
||||
|
||||
for {
|
||||
item, err := exporter.Next()
|
||||
if errors.Is(err, ErrorExportDone) {
|
||||
break
|
||||
} else if err != nil {
|
||||
return fmt.Errorf("failed to get the next export node: %w", err)
|
||||
}
|
||||
|
||||
if err = protoWriter.WriteMsg(&snapshotstypes.SnapshotItem{
|
||||
Item: &snapshotstypes.SnapshotItem_IAVL{
|
||||
IAVL: item,
|
||||
},
|
||||
}); err != nil {
|
||||
return fmt.Errorf("failed to write iavl node: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}(); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Restore implements snapshotstypes.CommitSnapshotter.
|
||||
func (c *CommitStore) Restore(version uint64, format uint32, protoReader protoio.Reader, chStorage chan<- *store.KVPair) (snapshotstypes.SnapshotItem, error) {
|
||||
var (
|
||||
importer Importer
|
||||
snapshotItem snapshotstypes.SnapshotItem
|
||||
storeKey string
|
||||
)
|
||||
|
||||
loop:
|
||||
for {
|
||||
snapshotItem = snapshotstypes.SnapshotItem{}
|
||||
err := protoReader.ReadMsg(&snapshotItem)
|
||||
if errors.Is(err, io.EOF) {
|
||||
break
|
||||
} else if err != nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("invalid protobuf message: %w", err)
|
||||
}
|
||||
|
||||
switch item := snapshotItem.Item.(type) {
|
||||
case *snapshotstypes.SnapshotItem_Store:
|
||||
if importer != nil {
|
||||
if err := importer.Commit(); err != nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("failed to commit importer: %w", err)
|
||||
}
|
||||
importer.Close()
|
||||
}
|
||||
storeKey = item.Store.Name
|
||||
tree := c.multiTrees[storeKey]
|
||||
if tree == nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("store %s not found", storeKey)
|
||||
}
|
||||
importer, err = tree.Import(version)
|
||||
if err != nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("failed to import tree for version %d: %w", version, err)
|
||||
}
|
||||
defer importer.Close()
|
||||
|
||||
case *snapshotstypes.SnapshotItem_IAVL:
|
||||
if importer == nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("received IAVL node item before store item")
|
||||
}
|
||||
node := item.IAVL
|
||||
if node.Height > int32(math.MaxInt8) {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("node height %v cannot exceed %v",
|
||||
item.IAVL.Height, math.MaxInt8)
|
||||
}
|
||||
// Protobuf does not differentiate between []byte{} and nil, but fortunately IAVL does
|
||||
// not allow nil keys nor nil values for leaf nodes, so we can always set them to empty.
|
||||
if node.Key == nil {
|
||||
node.Key = []byte{}
|
||||
}
|
||||
if node.Height == 0 {
|
||||
if node.Value == nil {
|
||||
node.Value = []byte{}
|
||||
}
|
||||
// If the node is a leaf node, it will be written to the storage.
|
||||
chStorage <- &store.KVPair{
|
||||
Key: node.Key,
|
||||
Value: node.Value,
|
||||
StoreKey: storeKey,
|
||||
}
|
||||
}
|
||||
err := importer.Add(node)
|
||||
if err != nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("failed to add node to importer: %w", err)
|
||||
}
|
||||
default:
|
||||
break loop
|
||||
}
|
||||
}
|
||||
|
||||
if importer != nil {
|
||||
if err := importer.Commit(); err != nil {
|
||||
return snapshotstypes.SnapshotItem{}, fmt.Errorf("failed to commit importer: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
return snapshotItem, c.LoadVersion(version)
|
||||
}
|
||||
|
||||
func (c *CommitStore) Close() (ferr error) {
|
||||
for _, tree := range c.multiTrees {
|
||||
if err := tree.Close(); err != nil {
|
||||
|
||||
@@ -0,0 +1,121 @@
|
||||
package commitment
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"sync"
|
||||
|
||||
dbm "github.com/cosmos/cosmos-db"
|
||||
"github.com/stretchr/testify/suite"
|
||||
|
||||
"cosmossdk.io/log"
|
||||
"cosmossdk.io/store/v2"
|
||||
"cosmossdk.io/store/v2/snapshots"
|
||||
snapshotstypes "cosmossdk.io/store/v2/snapshots/types"
|
||||
)
|
||||
|
||||
const (
|
||||
storeKey1 = "store1"
|
||||
storeKey2 = "store2"
|
||||
)
|
||||
|
||||
// CommitStoreTestSuite is a test suite to be used for all tree backends.
|
||||
type CommitStoreTestSuite struct {
|
||||
suite.Suite
|
||||
|
||||
NewStore func(db dbm.DB, storeKeys []string, logger log.Logger) (*CommitStore, error)
|
||||
}
|
||||
|
||||
func (s *CommitStoreTestSuite) TestSnapshotter() {
|
||||
storeKeys := []string{storeKey1, storeKey2}
|
||||
commitStore, err := s.NewStore(dbm.NewMemDB(), storeKeys, log.NewNopLogger())
|
||||
s.Require().NoError(err)
|
||||
|
||||
latestVersion := uint64(10)
|
||||
kvCount := 10
|
||||
for i := uint64(1); i <= latestVersion; i++ {
|
||||
kvPairs := make(map[string]store.KVPairs)
|
||||
for _, storeKey := range storeKeys {
|
||||
kvPairs[storeKey] = store.KVPairs{}
|
||||
for j := 0; j < kvCount; j++ {
|
||||
key := []byte(fmt.Sprintf("key-%d-%d", i, j))
|
||||
value := []byte(fmt.Sprintf("value-%d-%d", i, j))
|
||||
kvPairs[storeKey] = append(kvPairs[storeKey], store.KVPair{Key: key, Value: value})
|
||||
}
|
||||
}
|
||||
s.Require().NoError(commitStore.WriteBatch(store.NewChangeset(kvPairs)))
|
||||
|
||||
_, err = commitStore.Commit()
|
||||
s.Require().NoError(err)
|
||||
}
|
||||
|
||||
latestStoreInfos := commitStore.WorkingStoreInfos(latestVersion)
|
||||
s.Require().Equal(len(storeKeys), len(latestStoreInfos))
|
||||
|
||||
// create a snapshot
|
||||
dummyExtensionItem := snapshotstypes.SnapshotItem{
|
||||
Item: &snapshotstypes.SnapshotItem_Extension{
|
||||
Extension: &snapshotstypes.SnapshotExtensionMeta{
|
||||
Name: "test",
|
||||
Format: 1,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
targetStore, err := s.NewStore(dbm.NewMemDB(), storeKeys, log.NewNopLogger())
|
||||
s.Require().NoError(err)
|
||||
|
||||
chunks := make(chan io.ReadCloser, kvCount*int(latestVersion))
|
||||
go func() {
|
||||
streamWriter := snapshots.NewStreamWriter(chunks)
|
||||
s.Require().NotNil(streamWriter)
|
||||
defer streamWriter.Close()
|
||||
err := commitStore.Snapshot(latestVersion, streamWriter)
|
||||
s.Require().NoError(err)
|
||||
// write an extension metadata
|
||||
err = streamWriter.WriteMsg(&dummyExtensionItem)
|
||||
s.Require().NoError(err)
|
||||
}()
|
||||
|
||||
streamReader, err := snapshots.NewStreamReader(chunks)
|
||||
s.Require().NoError(err)
|
||||
chStorage := make(chan *store.KVPair, 100)
|
||||
leaves := make(map[string]string)
|
||||
wg := sync.WaitGroup{}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
for kv := range chStorage {
|
||||
leaves[fmt.Sprintf("%s_%s", kv.StoreKey, kv.Key)] = string(kv.Value)
|
||||
}
|
||||
wg.Done()
|
||||
}()
|
||||
nextItem, err := targetStore.Restore(latestVersion, snapshotstypes.CurrentFormat, streamReader, chStorage)
|
||||
s.Require().NoError(err)
|
||||
s.Require().Equal(*dummyExtensionItem.GetExtension(), *nextItem.GetExtension())
|
||||
|
||||
close(chStorage)
|
||||
wg.Wait()
|
||||
s.Require().Equal(len(storeKeys)*kvCount*int(latestVersion), len(leaves))
|
||||
for _, storeKey := range storeKeys {
|
||||
for i := 1; i <= int(latestVersion); i++ {
|
||||
for j := 0; j < kvCount; j++ {
|
||||
key := fmt.Sprintf("%s_key-%d-%d", storeKey, i, j)
|
||||
s.Require().Equal(leaves[key], fmt.Sprintf("value-%d-%d", i, j))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// check the restored tree hash
|
||||
targetStoreInfos := targetStore.WorkingStoreInfos(latestVersion)
|
||||
s.Require().Equal(len(storeKeys), len(targetStoreInfos))
|
||||
for _, storeInfo := range targetStoreInfos {
|
||||
matched := false
|
||||
for _, latestStoreInfo := range latestStoreInfos {
|
||||
if storeInfo.Name == latestStoreInfo.Name {
|
||||
s.Require().Equal(latestStoreInfo.GetHash(), storeInfo.GetHash())
|
||||
matched = true
|
||||
}
|
||||
}
|
||||
s.Require().True(matched)
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,17 @@
|
||||
package commitment
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"io"
|
||||
|
||||
ics23 "github.com/cosmos/ics23/go"
|
||||
|
||||
snapshotstypes "cosmossdk.io/store/v2/snapshots/types"
|
||||
)
|
||||
|
||||
// ErrorExportDone is returned by Exporter.Next() when all items have been exported.
|
||||
var ErrorExportDone = errors.New("export is complete")
|
||||
|
||||
// Tree is the interface that wraps the basic Tree methods.
|
||||
type Tree interface {
|
||||
Set(key, value []byte) error
|
||||
@@ -16,6 +22,23 @@ type Tree interface {
|
||||
Commit() ([]byte, error)
|
||||
GetProof(version uint64, key []byte) (*ics23.CommitmentProof, error)
|
||||
Prune(version uint64) error
|
||||
Export(version uint64) (Exporter, error)
|
||||
Import(version uint64) (Importer, error)
|
||||
|
||||
io.Closer
|
||||
}
|
||||
|
||||
// Exporter is the interface that wraps the basic Export methods.
|
||||
type Exporter interface {
|
||||
Next() (*snapshotstypes.SnapshotIAVLItem, error)
|
||||
|
||||
io.Closer
|
||||
}
|
||||
|
||||
// Importer is the interface that wraps the basic Import methods.
|
||||
type Importer interface {
|
||||
Add(*snapshotstypes.SnapshotIAVLItem) error
|
||||
Commit() error
|
||||
|
||||
io.Closer
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user