From c9197da8bea6ebf896ed94664e5eb984b30ca729 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Tue, 21 May 2019 15:59:27 +0200 Subject: [PATCH 01/10] Marshalling membership and saving it on disk - Marshalled a membership and saved it on disk - Introduced a general functions for read / write on disk with an interface - Renamed NewGroupRegistry to NewGroups --- .gitignore | 1 + config/config_test.go | 2 +- pkg/beacon/beacon.go | 2 +- pkg/beacon/relay/registry/groups.go | 31 ++++++++--- pkg/storage/file_storage.go | 37 +++++++++++++ pkg/storage/storage.go | 84 +++++++++++++++++++++++++++++ pkg/storage/storage_test.go | 49 +++++++++++++++++ 7 files changed, 197 insertions(+), 9 deletions(-) create mode 100644 pkg/storage/file_storage.go create mode 100644 pkg/storage/storage.go create mode 100644 pkg/storage/storage_test.go diff --git a/.gitignore b/.gitignore index 9114b2f37c..35b1e1332a 100644 --- a/.gitignore +++ b/.gitignore @@ -2,6 +2,7 @@ node_modules/ ./contracts/solidity/truffle.js build/ vendor/ +data_storage/ _local/ *.swp *.swo diff --git a/config/config_test.go b/config/config_test.go index eb72b7e4b3..aaca20e397 100644 --- a/config/config_test.go +++ b/config/config_test.go @@ -47,7 +47,7 @@ func TestReadConfig(t *testing.T) { "KeepGroup": "0xcf64c2a367341170cb4e09cf8c0ed137d8473ceb", }, }, - "StateManagementData": { + "Storage.DataDir": { readValueFunc: func(c *Config) interface{} { return c.Storage.DataDir }, expectedValue: "/my/secure/location", }, diff --git a/pkg/beacon/beacon.go b/pkg/beacon/beacon.go index 8e29f68e10..a986300053 100644 --- a/pkg/beacon/beacon.go +++ b/pkg/beacon/beacon.go @@ -35,7 +35,7 @@ func Initialize( return err } - groupRegistry := registry.NewGroupRegistry(relayChain) + groupRegistry := registry.NewGroups(relayChain) node := relay.NewNode( staker, diff --git a/pkg/beacon/relay/registry/groups.go b/pkg/beacon/relay/registry/groups.go index b4efb51b19..7a8e006a36 100644 --- a/pkg/beacon/relay/registry/groups.go +++ b/pkg/beacon/relay/registry/groups.go @@ -5,8 +5,11 @@ import ( "os" "sync" + "encoding/hex" + relaychain "github.com/keep-network/keep-core/pkg/beacon/relay/chain" "github.com/keep-network/keep-core/pkg/beacon/relay/dkg" + "github.com/keep-network/keep-core/pkg/storage" ) // Groups represents a collection of Keep groups in which the given @@ -25,8 +28,8 @@ type Membership struct { ChannelName string } -// NewGroupRegistry returns an empty GroupRegistry. -func NewGroupRegistry( +// NewGroups returns an empty GroupRegistry. +func NewGroups( relayChain relaychain.GroupRegistrationInterface, ) *Groups { return &Groups{ @@ -47,11 +50,25 @@ func (gr *Groups) RegisterGroup( groupPublicKey := string(signer.GroupPublicKeyBytes()) - gr.myGroups[groupPublicKey] = append(gr.myGroups[groupPublicKey], - &Membership{ - Signer: signer, - ChannelName: channelName, - }) + membership := &Membership{ + Signer: signer, + ChannelName: channelName, + } + + membershipBytes, err := membership.Marshal() + if err != nil { + fmt.Fprintf(os.Stderr, "Marshalling of the membership failed: [%v]\n", err) + } + + fileStorage := storage.NewFileStorage() + if fileStorage != nil { + hexGroupPublicKey := hex.EncodeToString(signer.GroupPublicKeyBytes()) + fileStorage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) + } else { + fmt.Fprintf(os.Stderr, "An error occured while retrieving a config path") + } + + gr.myGroups[groupPublicKey] = append(gr.myGroups[groupPublicKey], membership) } // GetGroup gets a group by a groupPublicKey diff --git a/pkg/storage/file_storage.go b/pkg/storage/file_storage.go new file mode 100644 index 0000000000..c31b9aeea0 --- /dev/null +++ b/pkg/storage/file_storage.go @@ -0,0 +1,37 @@ +package storage + +import ( + "github.com/keep-network/keep-core/config" +) + +// Storage is an interface to persist data on disk +type Storage interface { + Save(data []byte, name string) +} + +// FileStorage struct is an implementation of Storage +type FileStorage struct { + dataDir string +} + +// NewFileStorage creates a new FileStorage +func NewFileStorage() *FileStorage { + cfg, _ := config.ReadConfig("../../../../config.local.1.toml") // path is probably different on test / prod env. + + if cfg == nil { + return nil + } + + return &FileStorage{ + dataDir: cfg.Storage.DataDir, + } +} + +// Save - writes data in file +func (fs *FileStorage) Save(data []byte, suffix string) { + file := &File{ + FileName: fs.dataDir + suffix, + } + + file.Write(data) +} diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go new file mode 100644 index 0000000000..0deb95a2b1 --- /dev/null +++ b/pkg/storage/storage.go @@ -0,0 +1,84 @@ +package storage + +import ( + "fmt" + "io/ioutil" + "os" + "path/filepath" +) + +var ( + //ErrNoFileExists an error is shown when no file name was provided + ErrNoFileExists = fmt.Errorf("please provide a file name") +) + +// File represents a file on disk that a caller can use to read and write into. +type File struct { + // FileName is the file name of the main storage file. + FileName string +} + +// NewFile creates a new file at the target location on disk +func (f *File) NewFile(FileName string) *File { + + filepath.Dir(FileName) + + return &File{ + FileName: FileName, + } +} + +// Create and write data to a file +func (f *File) Write(data []byte) error { + if f.FileName == "" { + return ErrNoFileExists + } + + var err error + writeFile, err := os.Create(f.FileName) + check(err) + + defer writeFile.Close() + + _, err = writeFile.Write(data) + check(err) + + writeFile.Sync() + + return nil +} + +// Read a file from a file system +func (f *File) Read(FileName string) ([]byte, error) { + if f.FileName == "" { + return nil, ErrNoFileExists + } + + readFile, err := os.Open(FileName) + check(err) + + defer readFile.Close() + + data, err := ioutil.ReadAll(readFile) + check(err) + + return data, nil +} + +// Remove a file from a file syste +func (f *File) Remove(FileName string) error { + if f.FileName == "" { + return ErrNoFileExists + } + + err := os.Remove(FileName) + check(err) + + return nil +} + +func check(e error) { + if e != nil { + panic(e) + } +} diff --git a/pkg/storage/storage_test.go b/pkg/storage/storage_test.go new file mode 100644 index 0000000000..da04f7c3c5 --- /dev/null +++ b/pkg/storage/storage_test.go @@ -0,0 +1,49 @@ +package storage + +import ( + "bytes" + "os" + "testing" +) + +var ( + FileName = "foo" +) + +func TestMain(m *testing.M) { + code := m.Run() + if _, err := os.Stat(FileName); err == nil { + os.Remove(FileName) + } + os.Exit(code) +} + +func TestFile_WriteRead(t *testing.T) { + file := &File{ + FileName: FileName, + } + bytesToTest := []byte{115, 111, 109, 101, 10} + + file.Write(bytesToTest) + + actual, _ := file.Read(FileName) + + if !bytes.Equal(bytesToTest, actual) { + t.Fatalf("Bytes do not match. \nExpected: [%+v]\nActual: [%+v]", + bytesToTest, + actual) + } +} + +func TestFile_Remove(t *testing.T) { + if _, err := os.Stat(FileName); err == nil { + err = os.Remove(FileName) + if err != nil { + t.Fatalf("Was not able to remove a file [%+v]", FileName) + } + } + + if _, err := os.Stat(FileName); err == nil { + t.Fatalf("File [%+v] was supposed to be removed", FileName) + } +} From 5243961e2e459f9bffdd5526455da7942983d9be Mon Sep 17 00:00:00 2001 From: Dmitry Date: Wed, 22 May 2019 08:03:25 +0200 Subject: [PATCH 02/10] Modifying Storage interface - Modified storage interface to accept path to data dir - passed storage through beacon.go - reverted NewGroups back to NewGroupRegistry --- cmd/smoketest.go | 4 ++++ cmd/start.go | 6 ++++++ pkg/beacon/beacon.go | 4 +++- pkg/beacon/relay/registry/groups.go | 18 ++++++++---------- pkg/beacon/relay/registry/groups_test.go | 2 ++ pkg/storage/file_storage.go | 22 ++++++---------------- 6 files changed, 29 insertions(+), 27 deletions(-) diff --git a/cmd/smoketest.go b/cmd/smoketest.go index 82e5e0ca69..001c9e4a49 100644 --- a/cmd/smoketest.go +++ b/cmd/smoketest.go @@ -13,6 +13,7 @@ import ( "github.com/keep-network/keep-core/pkg/chain/local" netlocal "github.com/keep-network/keep-core/pkg/net/local" "github.com/keep-network/keep-core/pkg/operator" + "github.com/keep-network/keep-core/pkg/storage" "github.com/urfave/cli" ) @@ -140,6 +141,8 @@ func createNode( )) } + storage := storage.NewStorage("path_to_data_storage") + netProvider := netlocal.Connect() go func() { @@ -157,6 +160,7 @@ func createNode( chainCounter, stakeMonitor, netProvider, + storage, ) if err != nil { panic(fmt.Sprintf( diff --git a/cmd/start.go b/cmd/start.go index de00a24aee..a18659e2c2 100644 --- a/cmd/start.go +++ b/cmd/start.go @@ -11,6 +11,7 @@ import ( "github.com/keep-network/keep-core/pkg/net/key" "github.com/keep-network/keep-core/pkg/net/libp2p" "github.com/keep-network/keep-core/pkg/operator" + "github.com/keep-network/keep-core/pkg/storage" "github.com/urfave/cli" ) @@ -76,6 +77,8 @@ func Start(c *cli.Context) error { hasMinimumStake, err := stakeMonitor.HasMinimumStake( config.Ethereum.Account.Address, ) + + fmt.Println("hasMinimumStake: ", hasMinimumStake) if err != nil { return fmt.Errorf("could not check the stake [%v]", err) } @@ -100,6 +103,8 @@ func Start(c *cli.Context) error { isBootstrapNode := config.LibP2P.Seed != 0 nodeHeader(isBootstrapNode, netProvider.AddrStrings(), port) + storage := storage.NewStorage(config.Storage.DataDir) + err = beacon.Initialize( ctx, config.Ethereum.Account.Address, @@ -107,6 +112,7 @@ func Start(c *cli.Context) error { blockCounter, stakeMonitor, netProvider, + storage, ) if err != nil { return fmt.Errorf("error initializing beacon: [%v]", err) diff --git a/pkg/beacon/beacon.go b/pkg/beacon/beacon.go index a986300053..3106c55860 100644 --- a/pkg/beacon/beacon.go +++ b/pkg/beacon/beacon.go @@ -11,6 +11,7 @@ import ( "github.com/keep-network/keep-core/pkg/beacon/relay/registry" "github.com/keep-network/keep-core/pkg/chain" "github.com/keep-network/keep-core/pkg/net" + "github.com/keep-network/keep-core/pkg/storage" ) // Initialize kicks off the random beacon by initializing internal state, @@ -24,6 +25,7 @@ func Initialize( blockCounter chain.BlockCounter, stakeMonitor chain.StakeMonitor, netProvider net.Provider, + storage storage.Storage, ) error { chainConfig, err := relayChain.GetConfig() if err != nil { @@ -35,7 +37,7 @@ func Initialize( return err } - groupRegistry := registry.NewGroups(relayChain) + groupRegistry := registry.NewGroupRegistry(relayChain, storage) node := relay.NewNode( staker, diff --git a/pkg/beacon/relay/registry/groups.go b/pkg/beacon/relay/registry/groups.go index 7a8e006a36..aa07651042 100644 --- a/pkg/beacon/relay/registry/groups.go +++ b/pkg/beacon/relay/registry/groups.go @@ -20,6 +20,8 @@ type Groups struct { myGroups map[string][]*Membership relayChain relaychain.GroupRegistrationInterface + + storage storage.Storage } // Membership represents a member of a group @@ -28,13 +30,15 @@ type Membership struct { ChannelName string } -// NewGroups returns an empty GroupRegistry. -func NewGroups( +// NewGroupRegistry returns an empty GroupRegistry. +func NewGroupRegistry( relayChain relaychain.GroupRegistrationInterface, + storage storage.Storage, ) *Groups { return &Groups{ myGroups: make(map[string][]*Membership), relayChain: relayChain, + storage: storage, } } @@ -59,14 +63,8 @@ func (gr *Groups) RegisterGroup( if err != nil { fmt.Fprintf(os.Stderr, "Marshalling of the membership failed: [%v]\n", err) } - - fileStorage := storage.NewFileStorage() - if fileStorage != nil { - hexGroupPublicKey := hex.EncodeToString(signer.GroupPublicKeyBytes()) - fileStorage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) - } else { - fmt.Fprintf(os.Stderr, "An error occured while retrieving a config path") - } + hexGroupPublicKey := hex.EncodeToString(signer.GroupPublicKeyBytes()) + gr.storage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) gr.myGroups[groupPublicKey] = append(gr.myGroups[groupPublicKey], membership) } diff --git a/pkg/beacon/relay/registry/groups_test.go b/pkg/beacon/relay/registry/groups_test.go index cd8ce13965..ded94a19b6 100644 --- a/pkg/beacon/relay/registry/groups_test.go +++ b/pkg/beacon/relay/registry/groups_test.go @@ -11,6 +11,7 @@ import ( "github.com/keep-network/keep-core/pkg/beacon/relay/event" "github.com/keep-network/keep-core/pkg/beacon/relay/group" chainLocal "github.com/keep-network/keep-core/pkg/chain/local" + "github.com/keep-network/keep-core/pkg/storage" "github.com/keep-network/keep-core/pkg/subscription" ) @@ -25,6 +26,7 @@ func TestRegisterGroup(t *testing.T) { mutex: sync.Mutex{}, myGroups: make(map[string][]*Membership), relayChain: chainLocal.Connect(5, 3, big.NewInt(200)).ThresholdRelay(), + storage: storage.NewStorage("../../../../data_storage"), } gr.RegisterGroup(signer, "test_channel") diff --git a/pkg/storage/file_storage.go b/pkg/storage/file_storage.go index c31b9aeea0..505e5ff8f3 100644 --- a/pkg/storage/file_storage.go +++ b/pkg/storage/file_storage.go @@ -1,34 +1,24 @@ package storage -import ( - "github.com/keep-network/keep-core/config" -) - // Storage is an interface to persist data on disk type Storage interface { Save(data []byte, name string) } // FileStorage struct is an implementation of Storage -type FileStorage struct { +type fileStorage struct { dataDir string } -// NewFileStorage creates a new FileStorage -func NewFileStorage() *FileStorage { - cfg, _ := config.ReadConfig("../../../../config.local.1.toml") // path is probably different on test / prod env. - - if cfg == nil { - return nil - } - - return &FileStorage{ - dataDir: cfg.Storage.DataDir, +// NewStorage creates a new FileStorage +func NewStorage(path string) Storage { + return &fileStorage{ + dataDir: path, } } // Save - writes data in file -func (fs *FileStorage) Save(data []byte, suffix string) { +func (fs *fileStorage) Save(data []byte, suffix string) { file := &File{ FileName: fs.dataDir + suffix, } From fd355a2ca496a5f54ca5205c62139ab97ff738f3 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Wed, 22 May 2019 08:41:06 +0200 Subject: [PATCH 03/10] Adding .gitignore file inside data_storage dir --- data_storage/.gitignore | 4 ++++ 1 file changed, 4 insertions(+) create mode 100644 data_storage/.gitignore diff --git a/data_storage/.gitignore b/data_storage/.gitignore new file mode 100644 index 0000000000..86d0cb2726 --- /dev/null +++ b/data_storage/.gitignore @@ -0,0 +1,4 @@ +# Ignore everything in this directory +* +# Except this file +!.gitignore \ No newline at end of file From 358a0f62af6f23ae19bef26b0efb1168fcdf8652 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Wed, 22 May 2019 08:59:56 +0200 Subject: [PATCH 04/10] Adding data_storage path in TestUnregisterStaleGroups test --- pkg/beacon/relay/registry/groups_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/pkg/beacon/relay/registry/groups_test.go b/pkg/beacon/relay/registry/groups_test.go index ded94a19b6..2df6543c50 100644 --- a/pkg/beacon/relay/registry/groups_test.go +++ b/pkg/beacon/relay/registry/groups_test.go @@ -57,6 +57,7 @@ func TestUnregisterStaleGroups(t *testing.T) { mutex: sync.Mutex{}, myGroups: make(map[string][]*Membership), relayChain: mockChain, + storage: storage.NewStorage("../../../../data_storage"), } signer1 := dkg.NewThresholdSigner( From c5312d14fa64312d857b2f353d360719775c3220 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Wed, 22 May 2019 16:16:43 +0200 Subject: [PATCH 05/10] Removing accidental printf --- cmd/start.go | 2 -- 1 file changed, 2 deletions(-) diff --git a/cmd/start.go b/cmd/start.go index a18659e2c2..11f5478358 100644 --- a/cmd/start.go +++ b/cmd/start.go @@ -77,8 +77,6 @@ func Start(c *cli.Context) error { hasMinimumStake, err := stakeMonitor.HasMinimumStake( config.Ethereum.Account.Address, ) - - fmt.Println("hasMinimumStake: ", hasMinimumStake) if err != nil { return fmt.Errorf("could not check the stake [%v]", err) } From fcfe6cc98c88a7381fcdb261deefc28c4eac443b Mon Sep 17 00:00:00 2001 From: Dmitry Date: Thu, 23 May 2019 13:45:51 +0200 Subject: [PATCH 06/10] Adding a mock type for data storage and renaming files - Added a mock type for data storage object to avoid writing actual files on disk. Modified the tests. - Renamed storage.go <-> file_storage.go - Removed data_storage folder --- .gitignore | 1 - cmd/smoketest.go | 10 ++- data_storage/.gitignore | 4 - pkg/beacon/relay/registry/groups.go | 1 + pkg/beacon/relay/registry/groups_test.go | 27 ++++-- pkg/storage/file_storage.go | 87 +++++++++++++++---- .../{storage_test.go => file_storage_test.go} | 0 pkg/storage/storage.go | 87 ++++--------------- 8 files changed, 114 insertions(+), 103 deletions(-) delete mode 100644 data_storage/.gitignore rename pkg/storage/{storage_test.go => file_storage_test.go} (100%) diff --git a/.gitignore b/.gitignore index 35b1e1332a..9114b2f37c 100644 --- a/.gitignore +++ b/.gitignore @@ -2,7 +2,6 @@ node_modules/ ./contracts/solidity/truffle.js build/ vendor/ -data_storage/ _local/ *.swp *.swo diff --git a/cmd/smoketest.go b/cmd/smoketest.go index 001c9e4a49..720d8f8352 100644 --- a/cmd/smoketest.go +++ b/cmd/smoketest.go @@ -13,7 +13,6 @@ import ( "github.com/keep-network/keep-core/pkg/chain/local" netlocal "github.com/keep-network/keep-core/pkg/net/local" "github.com/keep-network/keep-core/pkg/operator" - "github.com/keep-network/keep-core/pkg/storage" "github.com/urfave/cli" ) @@ -43,6 +42,13 @@ const smokeTestDescription = `The smoke-test command creates a local threshold executed, once again with an in-process broadcast channel and chain, and the final signature is verified by each member of the group.` +type devNullDataStorage struct { +} + +func (dnds *devNullDataStorage) Save(data []byte, name string) { + // noop +} + func init() { SmokeTestCommand = cli.Command{ Name: "smoke-test", @@ -141,7 +147,7 @@ func createNode( )) } - storage := storage.NewStorage("path_to_data_storage") + storage := &devNullDataStorage{} netProvider := netlocal.Connect() diff --git a/data_storage/.gitignore b/data_storage/.gitignore deleted file mode 100644 index 86d0cb2726..0000000000 --- a/data_storage/.gitignore +++ /dev/null @@ -1,4 +0,0 @@ -# Ignore everything in this directory -* -# Except this file -!.gitignore \ No newline at end of file diff --git a/pkg/beacon/relay/registry/groups.go b/pkg/beacon/relay/registry/groups.go index aa07651042..61f50af52b 100644 --- a/pkg/beacon/relay/registry/groups.go +++ b/pkg/beacon/relay/registry/groups.go @@ -62,6 +62,7 @@ func (gr *Groups) RegisterGroup( membershipBytes, err := membership.Marshal() if err != nil { fmt.Fprintf(os.Stderr, "Marshalling of the membership failed: [%v]\n", err) + return } hexGroupPublicKey := hex.EncodeToString(signer.GroupPublicKeyBytes()) gr.storage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) diff --git a/pkg/beacon/relay/registry/groups_test.go b/pkg/beacon/relay/registry/groups_test.go index 2df6543c50..9ad957ecfa 100644 --- a/pkg/beacon/relay/registry/groups_test.go +++ b/pkg/beacon/relay/registry/groups_test.go @@ -11,24 +11,31 @@ import ( "github.com/keep-network/keep-core/pkg/beacon/relay/event" "github.com/keep-network/keep-core/pkg/beacon/relay/group" chainLocal "github.com/keep-network/keep-core/pkg/chain/local" - "github.com/keep-network/keep-core/pkg/storage" "github.com/keep-network/keep-core/pkg/subscription" ) -func TestRegisterGroup(t *testing.T) { - signer := dkg.NewThresholdSigner( - group.MemberIndex(2), - new(bn256.G2).ScalarBaseMult(big.NewInt(10)), - big.NewInt(1), - ) +type devNullDataStorage struct { +} + +func (dnds *devNullDataStorage) Save(data []byte, name string) { + // noop +} +func TestRegisterGroup(t *testing.T) { + noopStorage := &devNullDataStorage{} gr := &Groups{ mutex: sync.Mutex{}, myGroups: make(map[string][]*Membership), relayChain: chainLocal.Connect(5, 3, big.NewInt(200)).ThresholdRelay(), - storage: storage.NewStorage("../../../../data_storage"), + storage: noopStorage, } + signer := dkg.NewThresholdSigner( + group.MemberIndex(2), + new(bn256.G2).ScalarBaseMult(big.NewInt(10)), + big.NewInt(1), + ) + gr.RegisterGroup(signer, "test_channel") actual := gr.GetGroup(signer.GroupPublicKeyBytes()) @@ -53,11 +60,13 @@ func TestUnregisterStaleGroups(t *testing.T) { groupsToRemove: [][]byte{}, } + noopStorage := &devNullDataStorage{} + gr := &Groups{ mutex: sync.Mutex{}, myGroups: make(map[string][]*Membership), relayChain: mockChain, - storage: storage.NewStorage("../../../../data_storage"), + storage: noopStorage, } signer1 := dkg.NewThresholdSigner( diff --git a/pkg/storage/file_storage.go b/pkg/storage/file_storage.go index 505e5ff8f3..0deb95a2b1 100644 --- a/pkg/storage/file_storage.go +++ b/pkg/storage/file_storage.go @@ -1,27 +1,84 @@ package storage -// Storage is an interface to persist data on disk -type Storage interface { - Save(data []byte, name string) +import ( + "fmt" + "io/ioutil" + "os" + "path/filepath" +) + +var ( + //ErrNoFileExists an error is shown when no file name was provided + ErrNoFileExists = fmt.Errorf("please provide a file name") +) + +// File represents a file on disk that a caller can use to read and write into. +type File struct { + // FileName is the file name of the main storage file. + FileName string } -// FileStorage struct is an implementation of Storage -type fileStorage struct { - dataDir string +// NewFile creates a new file at the target location on disk +func (f *File) NewFile(FileName string) *File { + + filepath.Dir(FileName) + + return &File{ + FileName: FileName, + } } -// NewStorage creates a new FileStorage -func NewStorage(path string) Storage { - return &fileStorage{ - dataDir: path, +// Create and write data to a file +func (f *File) Write(data []byte) error { + if f.FileName == "" { + return ErrNoFileExists } + + var err error + writeFile, err := os.Create(f.FileName) + check(err) + + defer writeFile.Close() + + _, err = writeFile.Write(data) + check(err) + + writeFile.Sync() + + return nil } -// Save - writes data in file -func (fs *fileStorage) Save(data []byte, suffix string) { - file := &File{ - FileName: fs.dataDir + suffix, +// Read a file from a file system +func (f *File) Read(FileName string) ([]byte, error) { + if f.FileName == "" { + return nil, ErrNoFileExists } - file.Write(data) + readFile, err := os.Open(FileName) + check(err) + + defer readFile.Close() + + data, err := ioutil.ReadAll(readFile) + check(err) + + return data, nil +} + +// Remove a file from a file syste +func (f *File) Remove(FileName string) error { + if f.FileName == "" { + return ErrNoFileExists + } + + err := os.Remove(FileName) + check(err) + + return nil +} + +func check(e error) { + if e != nil { + panic(e) + } } diff --git a/pkg/storage/storage_test.go b/pkg/storage/file_storage_test.go similarity index 100% rename from pkg/storage/storage_test.go rename to pkg/storage/file_storage_test.go diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go index 0deb95a2b1..505e5ff8f3 100644 --- a/pkg/storage/storage.go +++ b/pkg/storage/storage.go @@ -1,84 +1,27 @@ package storage -import ( - "fmt" - "io/ioutil" - "os" - "path/filepath" -) - -var ( - //ErrNoFileExists an error is shown when no file name was provided - ErrNoFileExists = fmt.Errorf("please provide a file name") -) - -// File represents a file on disk that a caller can use to read and write into. -type File struct { - // FileName is the file name of the main storage file. - FileName string +// Storage is an interface to persist data on disk +type Storage interface { + Save(data []byte, name string) } -// NewFile creates a new file at the target location on disk -func (f *File) NewFile(FileName string) *File { - - filepath.Dir(FileName) - - return &File{ - FileName: FileName, - } +// FileStorage struct is an implementation of Storage +type fileStorage struct { + dataDir string } -// Create and write data to a file -func (f *File) Write(data []byte) error { - if f.FileName == "" { - return ErrNoFileExists +// NewStorage creates a new FileStorage +func NewStorage(path string) Storage { + return &fileStorage{ + dataDir: path, } - - var err error - writeFile, err := os.Create(f.FileName) - check(err) - - defer writeFile.Close() - - _, err = writeFile.Write(data) - check(err) - - writeFile.Sync() - - return nil } -// Read a file from a file system -func (f *File) Read(FileName string) ([]byte, error) { - if f.FileName == "" { - return nil, ErrNoFileExists +// Save - writes data in file +func (fs *fileStorage) Save(data []byte, suffix string) { + file := &File{ + FileName: fs.dataDir + suffix, } - readFile, err := os.Open(FileName) - check(err) - - defer readFile.Close() - - data, err := ioutil.ReadAll(readFile) - check(err) - - return data, nil -} - -// Remove a file from a file syste -func (f *File) Remove(FileName string) error { - if f.FileName == "" { - return ErrNoFileExists - } - - err := os.Remove(FileName) - check(err) - - return nil -} - -func check(e error) { - if e != nil { - panic(e) - } + file.Write(data) } From 7f6d0077745a78833503f1fa1c286dd178c701a4 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Wed, 29 May 2019 18:30:49 +0200 Subject: [PATCH 07/10] Intoroducing a membership data storage interface - Renamed disk storage files to disk*.go. - Introduced a layer for membership storage. - Refactored areas affected by above changes. --- cmd/smoketest.go | 3 +- cmd/start.go | 4 +-- pkg/beacon/beacon.go | 5 ++- pkg/beacon/relay/node.go | 5 ++- pkg/beacon/relay/registry/groups.go | 22 +++++------- pkg/beacon/relay/registry/groups_test.go | 3 +- pkg/beacon/relay/registry/storage.go | 36 +++++++++++++++++++ pkg/disk_storage/disk_handler.go | 26 ++++++++++++++ .../disk_storage.go} | 26 ++++++++------ .../disk_storage_test.go} | 0 pkg/storage/storage.go | 27 -------------- 11 files changed, 98 insertions(+), 59 deletions(-) create mode 100644 pkg/beacon/relay/registry/storage.go create mode 100644 pkg/disk_storage/disk_handler.go rename pkg/{storage/file_storage.go => disk_storage/disk_storage.go} (88%) rename pkg/{storage/file_storage_test.go => disk_storage/disk_storage_test.go} (100%) delete mode 100644 pkg/storage/storage.go diff --git a/cmd/smoketest.go b/cmd/smoketest.go index 720d8f8352..22152296a7 100644 --- a/cmd/smoketest.go +++ b/cmd/smoketest.go @@ -45,8 +45,9 @@ const smokeTestDescription = `The smoke-test command creates a local threshold type devNullDataStorage struct { } -func (dnds *devNullDataStorage) Save(data []byte, name string) { +func (dnds *devNullDataStorage) Save(data []byte, name string) error { // noop + return nil } func init() { diff --git a/cmd/start.go b/cmd/start.go index 11f5478358..943afaa5e5 100644 --- a/cmd/start.go +++ b/cmd/start.go @@ -8,10 +8,10 @@ import ( "github.com/keep-network/keep-core/pkg/beacon" "github.com/keep-network/keep-core/pkg/chain/ethereum" "github.com/keep-network/keep-core/pkg/chain/ethereum/ethutil" + storage "github.com/keep-network/keep-core/pkg/disk_storage" "github.com/keep-network/keep-core/pkg/net/key" "github.com/keep-network/keep-core/pkg/net/libp2p" "github.com/keep-network/keep-core/pkg/operator" - "github.com/keep-network/keep-core/pkg/storage" "github.com/urfave/cli" ) @@ -101,7 +101,7 @@ func Start(c *cli.Context) error { isBootstrapNode := config.LibP2P.Seed != 0 nodeHeader(isBootstrapNode, netProvider.AddrStrings(), port) - storage := storage.NewStorage(config.Storage.DataDir) + storage := storage.NewDiskHandler(config.Storage.DataDir) err = beacon.Initialize( ctx, diff --git a/pkg/beacon/beacon.go b/pkg/beacon/beacon.go index 3ec34329c2..38f296c961 100644 --- a/pkg/beacon/beacon.go +++ b/pkg/beacon/beacon.go @@ -10,8 +10,8 @@ import ( "github.com/keep-network/keep-core/pkg/beacon/relay/event" "github.com/keep-network/keep-core/pkg/beacon/relay/registry" "github.com/keep-network/keep-core/pkg/chain" + ds "github.com/keep-network/keep-core/pkg/disk_storage" "github.com/keep-network/keep-core/pkg/net" - "github.com/keep-network/keep-core/pkg/storage" ) // Initialize kicks off the random beacon by initializing internal state, @@ -25,8 +25,7 @@ func Initialize( blockCounter chain.BlockCounter, stakeMonitor chain.StakeMonitor, netProvider net.Provider, - storage storage.Storage, -) error { + storage ds.DiskHandler) error { chainConfig, err := relayChain.GetConfig() if err != nil { return err diff --git a/pkg/beacon/relay/node.go b/pkg/beacon/relay/node.go index ce55d1bb4b..6fbb0d855e 100644 --- a/pkg/beacon/relay/node.go +++ b/pkg/beacon/relay/node.go @@ -97,10 +97,13 @@ func (n *Node) JoinGroupIfEligible( return } - n.groupRegistry.RegisterGroup( + err = n.groupRegistry.RegisterGroup( signer, broadcastChannelName, ) + if err != nil { + fmt.Fprintf(os.Stderr, "Failed to register a group: [%v].\n", err) + } }() } } diff --git a/pkg/beacon/relay/registry/groups.go b/pkg/beacon/relay/registry/groups.go index 61f50af52b..aa40d0297a 100644 --- a/pkg/beacon/relay/registry/groups.go +++ b/pkg/beacon/relay/registry/groups.go @@ -5,11 +5,9 @@ import ( "os" "sync" - "encoding/hex" - relaychain "github.com/keep-network/keep-core/pkg/beacon/relay/chain" "github.com/keep-network/keep-core/pkg/beacon/relay/dkg" - "github.com/keep-network/keep-core/pkg/storage" + ds "github.com/keep-network/keep-core/pkg/disk_storage" ) // Groups represents a collection of Keep groups in which the given @@ -21,7 +19,7 @@ type Groups struct { relayChain relaychain.GroupRegistrationInterface - storage storage.Storage + storage Handle } // Membership represents a member of a group @@ -33,12 +31,11 @@ type Membership struct { // NewGroupRegistry returns an empty GroupRegistry. func NewGroupRegistry( relayChain relaychain.GroupRegistrationInterface, - storage storage.Storage, -) *Groups { + diskStorage ds.DiskHandler) *Groups { return &Groups{ myGroups: make(map[string][]*Membership), relayChain: relayChain, - storage: storage, + storage: NewStorage(diskStorage), } } @@ -47,7 +44,7 @@ func NewGroupRegistry( func (gr *Groups) RegisterGroup( signer *dkg.ThresholdSigner, channelName string, -) { +) error { gr.mutex.Lock() defer gr.mutex.Unlock() @@ -59,15 +56,14 @@ func (gr *Groups) RegisterGroup( ChannelName: channelName, } - membershipBytes, err := membership.Marshal() + err := gr.storage.Save(membership) if err != nil { - fmt.Fprintf(os.Stderr, "Marshalling of the membership failed: [%v]\n", err) - return + return fmt.Errorf("could not persist membership to the storage: [%v]", err) } - hexGroupPublicKey := hex.EncodeToString(signer.GroupPublicKeyBytes()) - gr.storage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) gr.myGroups[groupPublicKey] = append(gr.myGroups[groupPublicKey], membership) + + return nil } // GetGroup gets a group by a groupPublicKey diff --git a/pkg/beacon/relay/registry/groups_test.go b/pkg/beacon/relay/registry/groups_test.go index 9ad957ecfa..2a686f6646 100644 --- a/pkg/beacon/relay/registry/groups_test.go +++ b/pkg/beacon/relay/registry/groups_test.go @@ -17,8 +17,9 @@ import ( type devNullDataStorage struct { } -func (dnds *devNullDataStorage) Save(data []byte, name string) { +func (dnds *devNullDataStorage) Save(membership *Membership) error { // noop + return nil } func TestRegisterGroup(t *testing.T) { diff --git a/pkg/beacon/relay/registry/storage.go b/pkg/beacon/relay/registry/storage.go new file mode 100644 index 0000000000..c763cb02f8 --- /dev/null +++ b/pkg/beacon/relay/registry/storage.go @@ -0,0 +1,36 @@ +package registry + +import ( + "fmt" + + ds "github.com/keep-network/keep-core/pkg/disk_storage" + + "encoding/hex" +) + +// Handle is an interface to handle memberships on disk +type Handle interface { + Save(membership *Membership) error +} + +type storage struct { + diskStorage ds.DiskHandler +} + +// NewStorage creates a new storage. +func NewStorage(diskStorage ds.DiskHandler) Handle { + return &storage{ + diskStorage: diskStorage, + } +} + +// Save converts a membership suitable for disk storage. +func (s *storage) Save(membership *Membership) error { + membershipBytes, err := membership.Marshal() + if err != nil { + return fmt.Errorf("marshalling of the membership failed: [%v]", err) + } + hexGroupPublicKey := hex.EncodeToString(membership.Signer.GroupPublicKeyBytes()) + + return s.diskStorage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) +} diff --git a/pkg/disk_storage/disk_handler.go b/pkg/disk_storage/disk_handler.go new file mode 100644 index 0000000000..63123a065f --- /dev/null +++ b/pkg/disk_storage/disk_handler.go @@ -0,0 +1,26 @@ +package storage + +// DiskHandler is an interface for data persistence on disk +type DiskHandler interface { + Save(data []byte, name string) error +} + +type diskStorage struct { + dataDir string +} + +// NewDiskHandler creates a new diskStorage +func NewDiskHandler(path string) DiskHandler { + return &diskStorage{ + dataDir: path, + } +} + +// Save - writes data to file +func (ds *diskStorage) Save(data []byte, suffix string) error { + file := &File{ + FileName: ds.dataDir + suffix, + } + + return file.Write(data) +} diff --git a/pkg/storage/file_storage.go b/pkg/disk_storage/disk_storage.go similarity index 88% rename from pkg/storage/file_storage.go rename to pkg/disk_storage/disk_storage.go index 0deb95a2b1..3524d67f30 100644 --- a/pkg/storage/file_storage.go +++ b/pkg/disk_storage/disk_storage.go @@ -36,12 +36,16 @@ func (f *File) Write(data []byte) error { var err error writeFile, err := os.Create(f.FileName) - check(err) + if err != nil { + return err + } defer writeFile.Close() _, err = writeFile.Write(data) - check(err) + if err != nil { + return err + } writeFile.Sync() @@ -55,12 +59,16 @@ func (f *File) Read(FileName string) ([]byte, error) { } readFile, err := os.Open(FileName) - check(err) + if err != nil { + return nil, err + } defer readFile.Close() data, err := ioutil.ReadAll(readFile) - check(err) + if err != nil { + return nil, err + } return data, nil } @@ -72,13 +80,9 @@ func (f *File) Remove(FileName string) error { } err := os.Remove(FileName) - check(err) + if err != nil { + return err + } return nil } - -func check(e error) { - if e != nil { - panic(e) - } -} diff --git a/pkg/storage/file_storage_test.go b/pkg/disk_storage/disk_storage_test.go similarity index 100% rename from pkg/storage/file_storage_test.go rename to pkg/disk_storage/disk_storage_test.go diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go deleted file mode 100644 index 505e5ff8f3..0000000000 --- a/pkg/storage/storage.go +++ /dev/null @@ -1,27 +0,0 @@ -package storage - -// Storage is an interface to persist data on disk -type Storage interface { - Save(data []byte, name string) -} - -// FileStorage struct is an implementation of Storage -type fileStorage struct { - dataDir string -} - -// NewStorage creates a new FileStorage -func NewStorage(path string) Storage { - return &fileStorage{ - dataDir: path, - } -} - -// Save - writes data in file -func (fs *fileStorage) Save(data []byte, suffix string) { - file := &File{ - FileName: fs.dataDir + suffix, - } - - file.Write(data) -} From e2b176c57f4ac32d71a0f67fbd4d3e7c2f9ac975 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Thu, 30 May 2019 12:23:27 +0200 Subject: [PATCH 08/10] Changing an interface for saving memberships - Renamed files to be more generic - Changed code and test around it --- cmd/start.go | 6 +- pkg/beacon/beacon.go | 7 +- pkg/beacon/relay/registry/groups.go | 14 +-- pkg/beacon/relay/registry/groups_test.go | 26 ++--- pkg/beacon/relay/registry/storage.go | 22 ++--- pkg/disk_storage/disk_handler.go | 26 ----- pkg/disk_storage/disk_storage.go | 88 ----------------- pkg/persistence/disk_persistence.go | 97 +++++++++++++++++++ .../disk_persistence_test.go} | 24 ++--- pkg/persistence/persistence.go | 7 ++ 10 files changed, 149 insertions(+), 168 deletions(-) delete mode 100644 pkg/disk_storage/disk_handler.go delete mode 100644 pkg/disk_storage/disk_storage.go create mode 100644 pkg/persistence/disk_persistence.go rename pkg/{disk_storage/disk_storage_test.go => persistence/disk_persistence_test.go} (51%) create mode 100644 pkg/persistence/persistence.go diff --git a/cmd/start.go b/cmd/start.go index 943afaa5e5..d22d95e6e5 100644 --- a/cmd/start.go +++ b/cmd/start.go @@ -8,10 +8,10 @@ import ( "github.com/keep-network/keep-core/pkg/beacon" "github.com/keep-network/keep-core/pkg/chain/ethereum" "github.com/keep-network/keep-core/pkg/chain/ethereum/ethutil" - storage "github.com/keep-network/keep-core/pkg/disk_storage" "github.com/keep-network/keep-core/pkg/net/key" "github.com/keep-network/keep-core/pkg/net/libp2p" "github.com/keep-network/keep-core/pkg/operator" + "github.com/keep-network/keep-core/pkg/persistence" "github.com/urfave/cli" ) @@ -101,7 +101,7 @@ func Start(c *cli.Context) error { isBootstrapNode := config.LibP2P.Seed != 0 nodeHeader(isBootstrapNode, netProvider.AddrStrings(), port) - storage := storage.NewDiskHandler(config.Storage.DataDir) + persistence := persistence.NewDiskHandle(config.Storage.DataDir) err = beacon.Initialize( ctx, @@ -110,7 +110,7 @@ func Start(c *cli.Context) error { blockCounter, stakeMonitor, netProvider, - storage, + persistence, ) if err != nil { return fmt.Errorf("error initializing beacon: [%v]", err) diff --git a/pkg/beacon/beacon.go b/pkg/beacon/beacon.go index 38f296c961..9af570ec08 100644 --- a/pkg/beacon/beacon.go +++ b/pkg/beacon/beacon.go @@ -10,8 +10,8 @@ import ( "github.com/keep-network/keep-core/pkg/beacon/relay/event" "github.com/keep-network/keep-core/pkg/beacon/relay/registry" "github.com/keep-network/keep-core/pkg/chain" - ds "github.com/keep-network/keep-core/pkg/disk_storage" "github.com/keep-network/keep-core/pkg/net" + "github.com/keep-network/keep-core/pkg/persistence" ) // Initialize kicks off the random beacon by initializing internal state, @@ -25,7 +25,8 @@ func Initialize( blockCounter chain.BlockCounter, stakeMonitor chain.StakeMonitor, netProvider net.Provider, - storage ds.DiskHandler) error { + persistence persistence.Handle, +) error { chainConfig, err := relayChain.GetConfig() if err != nil { return err @@ -36,7 +37,7 @@ func Initialize( return err } - groupRegistry := registry.NewGroupRegistry(relayChain, storage) + groupRegistry := registry.NewGroupRegistry(relayChain, persistence) node := relay.NewNode( staker, diff --git a/pkg/beacon/relay/registry/groups.go b/pkg/beacon/relay/registry/groups.go index aa40d0297a..a7d81d6746 100644 --- a/pkg/beacon/relay/registry/groups.go +++ b/pkg/beacon/relay/registry/groups.go @@ -7,7 +7,8 @@ import ( relaychain "github.com/keep-network/keep-core/pkg/beacon/relay/chain" "github.com/keep-network/keep-core/pkg/beacon/relay/dkg" - ds "github.com/keep-network/keep-core/pkg/disk_storage" + + "github.com/keep-network/keep-core/pkg/persistence" ) // Groups represents a collection of Keep groups in which the given @@ -19,7 +20,7 @@ type Groups struct { relayChain relaychain.GroupRegistrationInterface - storage Handle + storage storage } // Membership represents a member of a group @@ -31,11 +32,13 @@ type Membership struct { // NewGroupRegistry returns an empty GroupRegistry. func NewGroupRegistry( relayChain relaychain.GroupRegistrationInterface, - diskStorage ds.DiskHandler) *Groups { + persistence persistence.Handle, +) *Groups { return &Groups{ myGroups: make(map[string][]*Membership), relayChain: relayChain, - storage: NewStorage(diskStorage), + storage: newStorage(persistence), + mutex: sync.Mutex{}, } } @@ -45,7 +48,6 @@ func (gr *Groups) RegisterGroup( signer *dkg.ThresholdSigner, channelName string, ) error { - gr.mutex.Lock() defer gr.mutex.Unlock() @@ -56,7 +58,7 @@ func (gr *Groups) RegisterGroup( ChannelName: channelName, } - err := gr.storage.Save(membership) + err := gr.storage.save(membership) if err != nil { return fmt.Errorf("could not persist membership to the storage: [%v]", err) } diff --git a/pkg/beacon/relay/registry/groups_test.go b/pkg/beacon/relay/registry/groups_test.go index 2a686f6646..4390aa36d9 100644 --- a/pkg/beacon/relay/registry/groups_test.go +++ b/pkg/beacon/relay/registry/groups_test.go @@ -3,7 +3,6 @@ package registry import ( "bytes" "math/big" - "sync" "testing" bn256 "github.com/ethereum/go-ethereum/crypto/bn256/cloudflare" @@ -14,22 +13,19 @@ import ( "github.com/keep-network/keep-core/pkg/subscription" ) -type devNullDataStorage struct { +type noopPersistence struct { } -func (dnds *devNullDataStorage) Save(membership *Membership) error { +func (np *noopPersistence) Save([]byte, string) error { // noop return nil } func TestRegisterGroup(t *testing.T) { - noopStorage := &devNullDataStorage{} - gr := &Groups{ - mutex: sync.Mutex{}, - myGroups: make(map[string][]*Membership), - relayChain: chainLocal.Connect(5, 3, big.NewInt(200)).ThresholdRelay(), - storage: noopStorage, - } + noopPersistence := &noopPersistence{} + chain := chainLocal.Connect(5, 3, big.NewInt(200)).ThresholdRelay() + + gr := NewGroupRegistry(chain, noopPersistence) signer := dkg.NewThresholdSigner( group.MemberIndex(2), @@ -60,15 +56,9 @@ func TestUnregisterStaleGroups(t *testing.T) { mockChain := &mockGroupRegistrationInterface{ groupsToRemove: [][]byte{}, } + noopPersistence := &noopPersistence{} - noopStorage := &devNullDataStorage{} - - gr := &Groups{ - mutex: sync.Mutex{}, - myGroups: make(map[string][]*Membership), - relayChain: mockChain, - storage: noopStorage, - } + gr := NewGroupRegistry(mockChain, noopPersistence) signer1 := dkg.NewThresholdSigner( group.MemberIndex(1), diff --git a/pkg/beacon/relay/registry/storage.go b/pkg/beacon/relay/registry/storage.go index c763cb02f8..daedbc969f 100644 --- a/pkg/beacon/relay/registry/storage.go +++ b/pkg/beacon/relay/registry/storage.go @@ -3,34 +3,32 @@ package registry import ( "fmt" - ds "github.com/keep-network/keep-core/pkg/disk_storage" + "github.com/keep-network/keep-core/pkg/persistence" "encoding/hex" ) -// Handle is an interface to handle memberships on disk -type Handle interface { - Save(membership *Membership) error +type storage interface { + save(membership *Membership) error } -type storage struct { - diskStorage ds.DiskHandler +type persistentStorage struct { + handle persistence.Handle } -// NewStorage creates a new storage. -func NewStorage(diskStorage ds.DiskHandler) Handle { - return &storage{ - diskStorage: diskStorage, +func newStorage(persistence persistence.Handle) storage { + return &persistentStorage{ + handle: persistence, } } // Save converts a membership suitable for disk storage. -func (s *storage) Save(membership *Membership) error { +func (ps *persistentStorage) save(membership *Membership) error { membershipBytes, err := membership.Marshal() if err != nil { return fmt.Errorf("marshalling of the membership failed: [%v]", err) } hexGroupPublicKey := hex.EncodeToString(membership.Signer.GroupPublicKeyBytes()) - return s.diskStorage.Save(membershipBytes, "/membership_"+hexGroupPublicKey) + return ps.handle.Save(membershipBytes, "/membership_"+hexGroupPublicKey) } diff --git a/pkg/disk_storage/disk_handler.go b/pkg/disk_storage/disk_handler.go deleted file mode 100644 index 63123a065f..0000000000 --- a/pkg/disk_storage/disk_handler.go +++ /dev/null @@ -1,26 +0,0 @@ -package storage - -// DiskHandler is an interface for data persistence on disk -type DiskHandler interface { - Save(data []byte, name string) error -} - -type diskStorage struct { - dataDir string -} - -// NewDiskHandler creates a new diskStorage -func NewDiskHandler(path string) DiskHandler { - return &diskStorage{ - dataDir: path, - } -} - -// Save - writes data to file -func (ds *diskStorage) Save(data []byte, suffix string) error { - file := &File{ - FileName: ds.dataDir + suffix, - } - - return file.Write(data) -} diff --git a/pkg/disk_storage/disk_storage.go b/pkg/disk_storage/disk_storage.go deleted file mode 100644 index 3524d67f30..0000000000 --- a/pkg/disk_storage/disk_storage.go +++ /dev/null @@ -1,88 +0,0 @@ -package storage - -import ( - "fmt" - "io/ioutil" - "os" - "path/filepath" -) - -var ( - //ErrNoFileExists an error is shown when no file name was provided - ErrNoFileExists = fmt.Errorf("please provide a file name") -) - -// File represents a file on disk that a caller can use to read and write into. -type File struct { - // FileName is the file name of the main storage file. - FileName string -} - -// NewFile creates a new file at the target location on disk -func (f *File) NewFile(FileName string) *File { - - filepath.Dir(FileName) - - return &File{ - FileName: FileName, - } -} - -// Create and write data to a file -func (f *File) Write(data []byte) error { - if f.FileName == "" { - return ErrNoFileExists - } - - var err error - writeFile, err := os.Create(f.FileName) - if err != nil { - return err - } - - defer writeFile.Close() - - _, err = writeFile.Write(data) - if err != nil { - return err - } - - writeFile.Sync() - - return nil -} - -// Read a file from a file system -func (f *File) Read(FileName string) ([]byte, error) { - if f.FileName == "" { - return nil, ErrNoFileExists - } - - readFile, err := os.Open(FileName) - if err != nil { - return nil, err - } - - defer readFile.Close() - - data, err := ioutil.ReadAll(readFile) - if err != nil { - return nil, err - } - - return data, nil -} - -// Remove a file from a file syste -func (f *File) Remove(FileName string) error { - if f.FileName == "" { - return ErrNoFileExists - } - - err := os.Remove(FileName) - if err != nil { - return err - } - - return nil -} diff --git a/pkg/persistence/disk_persistence.go b/pkg/persistence/disk_persistence.go new file mode 100644 index 0000000000..81031287df --- /dev/null +++ b/pkg/persistence/disk_persistence.go @@ -0,0 +1,97 @@ +package persistence + +import ( + "fmt" + "io/ioutil" + "os" +) + +// NewDiskHandle creates on-disk data persistence handle +func NewDiskHandle(path string) Handle { + return &diskPersistence{ + dataDir: path, + } +} + +type diskPersistence struct { + dataDir string +} + +// Save - writes data to file +func (ds *diskPersistence) Save(data []byte, suffix string) error { + file := &file{ + fileName: ds.dataDir + suffix, + } + + return file.Write(data) +} + +var ( + //ErrNoFileExists an error is shown when no file name was provided + errNoFileExists = fmt.Errorf("please provide a file name") +) + +// File represents a file on disk that a caller can use to read and write into. +type file struct { + // FileName is the file name of the main storage file. + fileName string +} + +// Create and write data to a file +func (f *file) Write(data []byte) error { + if f.fileName == "" { + return errNoFileExists + } + + var err error + writeFile, err := os.Create(f.fileName) + if err != nil { + return err + } + + defer writeFile.Close() + + _, err = writeFile.Write(data) + if err != nil { + return err + } + + writeFile.Sync() + + return nil +} + +// Read a file from a file system +func (f *file) Read(fileName string) ([]byte, error) { + if f.fileName == "" { + return nil, errNoFileExists + } + + readFile, err := os.Open(fileName) + if err != nil { + return nil, err + } + + defer readFile.Close() + + data, err := ioutil.ReadAll(readFile) + if err != nil { + return nil, err + } + + return data, nil +} + +// Remove a file from a file syste +func (f *file) Remove(fileName string) error { + if f.fileName == "" { + return errNoFileExists + } + + err := os.Remove(fileName) + if err != nil { + return err + } + + return nil +} diff --git a/pkg/disk_storage/disk_storage_test.go b/pkg/persistence/disk_persistence_test.go similarity index 51% rename from pkg/disk_storage/disk_storage_test.go rename to pkg/persistence/disk_persistence_test.go index da04f7c3c5..896f5138c5 100644 --- a/pkg/disk_storage/disk_storage_test.go +++ b/pkg/persistence/disk_persistence_test.go @@ -1,4 +1,4 @@ -package storage +package persistence import ( "bytes" @@ -7,26 +7,26 @@ import ( ) var ( - FileName = "foo" + fileName = "foo" ) func TestMain(m *testing.M) { code := m.Run() - if _, err := os.Stat(FileName); err == nil { - os.Remove(FileName) + if _, err := os.Stat(fileName); err == nil { + os.Remove(fileName) } os.Exit(code) } func TestFile_WriteRead(t *testing.T) { - file := &File{ - FileName: FileName, + file := &file{ + fileName: fileName, } bytesToTest := []byte{115, 111, 109, 101, 10} file.Write(bytesToTest) - actual, _ := file.Read(FileName) + actual, _ := file.Read(fileName) if !bytes.Equal(bytesToTest, actual) { t.Fatalf("Bytes do not match. \nExpected: [%+v]\nActual: [%+v]", @@ -36,14 +36,14 @@ func TestFile_WriteRead(t *testing.T) { } func TestFile_Remove(t *testing.T) { - if _, err := os.Stat(FileName); err == nil { - err = os.Remove(FileName) + if _, err := os.Stat(fileName); err == nil { + err = os.Remove(fileName) if err != nil { - t.Fatalf("Was not able to remove a file [%+v]", FileName) + t.Fatalf("Was not able to remove a file [%+v]", fileName) } } - if _, err := os.Stat(FileName); err == nil { - t.Fatalf("File [%+v] was supposed to be removed", FileName) + if _, err := os.Stat(fileName); err == nil { + t.Fatalf("File [%+v] was supposed to be removed", fileName) } } diff --git a/pkg/persistence/persistence.go b/pkg/persistence/persistence.go new file mode 100644 index 0000000000..f28e1c3c11 --- /dev/null +++ b/pkg/persistence/persistence.go @@ -0,0 +1,7 @@ +package persistence + +// Handle is an interface for data persistence. Underlying implementation +// can write data e.g. to disk, cache, or hardware module. +type Handle interface { + Save(data []byte, name string) error +} From 66526d505e17e280420abfa613c62ee7d22e5560 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Thu, 30 May 2019 13:34:18 +0200 Subject: [PATCH 09/10] Unexporting disk_persistence::Remove() function --- pkg/persistence/disk_persistence.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/persistence/disk_persistence.go b/pkg/persistence/disk_persistence.go index 81031287df..bcf369d94d 100644 --- a/pkg/persistence/disk_persistence.go +++ b/pkg/persistence/disk_persistence.go @@ -82,8 +82,8 @@ func (f *file) Read(fileName string) ([]byte, error) { return data, nil } -// Remove a file from a file syste -func (f *file) Remove(fileName string) error { +// Remove a file from a file system +func (f *file) remove(fileName string) error { if f.fileName == "" { return errNoFileExists } From b2259d717161863a68009d5f783b518b1d7e4386 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Thu, 30 May 2019 14:04:07 +0200 Subject: [PATCH 10/10] Adding memberID to a membership file name and unexporting functions --- cmd/smoketest.go | 6 +++--- pkg/beacon/relay/registry/storage.go | 2 +- pkg/persistence/disk_persistence.go | 6 +++--- pkg/persistence/disk_persistence_test.go | 4 ++-- 4 files changed, 9 insertions(+), 9 deletions(-) diff --git a/cmd/smoketest.go b/cmd/smoketest.go index 22152296a7..8d784cf36a 100644 --- a/cmd/smoketest.go +++ b/cmd/smoketest.go @@ -42,10 +42,10 @@ const smokeTestDescription = `The smoke-test command creates a local threshold executed, once again with an in-process broadcast channel and chain, and the final signature is verified by each member of the group.` -type devNullDataStorage struct { +type noopPersistence struct { } -func (dnds *devNullDataStorage) Save(data []byte, name string) error { +func (np *noopPersistence) Save(data []byte, name string) error { // noop return nil } @@ -148,7 +148,7 @@ func createNode( )) } - storage := &devNullDataStorage{} + storage := &noopPersistence{} netProvider := netlocal.Connect() diff --git a/pkg/beacon/relay/registry/storage.go b/pkg/beacon/relay/registry/storage.go index daedbc969f..d96d09c796 100644 --- a/pkg/beacon/relay/registry/storage.go +++ b/pkg/beacon/relay/registry/storage.go @@ -30,5 +30,5 @@ func (ps *persistentStorage) save(membership *Membership) error { } hexGroupPublicKey := hex.EncodeToString(membership.Signer.GroupPublicKeyBytes()) - return ps.handle.Save(membershipBytes, "/membership_"+hexGroupPublicKey) + return ps.handle.Save(membershipBytes, "/membership_"+hexGroupPublicKey+"_"+fmt.Sprint(membership.Signer.MemberID())) } diff --git a/pkg/persistence/disk_persistence.go b/pkg/persistence/disk_persistence.go index bcf369d94d..bdbcea4f35 100644 --- a/pkg/persistence/disk_persistence.go +++ b/pkg/persistence/disk_persistence.go @@ -23,7 +23,7 @@ func (ds *diskPersistence) Save(data []byte, suffix string) error { fileName: ds.dataDir + suffix, } - return file.Write(data) + return file.write(data) } var ( @@ -38,7 +38,7 @@ type file struct { } // Create and write data to a file -func (f *file) Write(data []byte) error { +func (f *file) write(data []byte) error { if f.fileName == "" { return errNoFileExists } @@ -62,7 +62,7 @@ func (f *file) Write(data []byte) error { } // Read a file from a file system -func (f *file) Read(fileName string) ([]byte, error) { +func (f *file) read(fileName string) ([]byte, error) { if f.fileName == "" { return nil, errNoFileExists } diff --git a/pkg/persistence/disk_persistence_test.go b/pkg/persistence/disk_persistence_test.go index 896f5138c5..cb3f6531c5 100644 --- a/pkg/persistence/disk_persistence_test.go +++ b/pkg/persistence/disk_persistence_test.go @@ -24,9 +24,9 @@ func TestFile_WriteRead(t *testing.T) { } bytesToTest := []byte{115, 111, 109, 101, 10} - file.Write(bytesToTest) + file.write(bytesToTest) - actual, _ := file.Read(fileName) + actual, _ := file.read(fileName) if !bytes.Equal(bytesToTest, actual) { t.Fatalf("Bytes do not match. \nExpected: [%+v]\nActual: [%+v]",