Skip to content
Merged
11 changes: 11 additions & 0 deletions cmd/smoketest.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,14 @@ 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 noopPersistence struct {
}

func (np *noopPersistence) Save(data []byte, name string) error {
// noop
return nil
}

func init() {
SmokeTestCommand = cli.Command{
Name: "smoke-test",
Expand Down Expand Up @@ -140,6 +148,8 @@ func createNode(
))
}

storage := &noopPersistence{}

netProvider := netlocal.Connect()

go func() {
Expand All @@ -157,6 +167,7 @@ func createNode(
chainCounter,
stakeMonitor,
netProvider,
storage,
)
if err != nil {
panic(fmt.Sprintf(
Expand Down
4 changes: 4 additions & 0 deletions cmd/start.go
Original file line number Diff line number Diff line change
Expand Up @@ -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/persistence"
"github.com/urfave/cli"
)

Expand Down Expand Up @@ -100,13 +101,16 @@ func Start(c *cli.Context) error {
isBootstrapNode := config.LibP2P.Seed != 0
nodeHeader(isBootstrapNode, netProvider.AddrStrings(), port)

persistence := persistence.NewDiskHandle(config.Storage.DataDir)

err = beacon.Initialize(
ctx,
config.Ethereum.Account.Address,
chainProvider.ThresholdRelay(),
blockCounter,
stakeMonitor,
netProvider,
persistence,
)
if err != nil {
return fmt.Errorf("error initializing beacon: [%v]", err)
Expand Down
2 changes: 1 addition & 1 deletion config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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",
},
Expand Down
4 changes: 3 additions & 1 deletion pkg/beacon/beacon.go
Original file line number Diff line number Diff line change
Expand Up @@ -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/persistence"
)

// Initialize kicks off the random beacon by initializing internal state,
Expand All @@ -24,6 +25,7 @@ func Initialize(
blockCounter chain.BlockCounter,
stakeMonitor chain.StakeMonitor,
netProvider net.Provider,
persistence persistence.Handle,
) error {
chainConfig, err := relayChain.GetConfig()
if err != nil {
Expand All @@ -35,7 +37,7 @@ func Initialize(
return err
}

groupRegistry := registry.NewGroupRegistry(relayChain)
groupRegistry := registry.NewGroupRegistry(relayChain, persistence)

node := relay.NewNode(
staker,
Expand Down
5 changes: 4 additions & 1 deletion pkg/beacon/relay/node.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}()
}
}
Expand Down
28 changes: 21 additions & 7 deletions pkg/beacon/relay/registry/groups.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import (

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/persistence"
)

// Groups represents a collection of Keep groups in which the given
Expand All @@ -17,6 +19,8 @@ type Groups struct {
myGroups map[string][]*Membership

relayChain relaychain.GroupRegistrationInterface

storage storage
}

// Membership represents a member of a group
Expand All @@ -28,10 +32,13 @@ type Membership struct {
// NewGroupRegistry returns an empty GroupRegistry.
func NewGroupRegistry(
relayChain relaychain.GroupRegistrationInterface,
persistence persistence.Handle,
) *Groups {
return &Groups{
myGroups: make(map[string][]*Membership),
relayChain: relayChain,
storage: newStorage(persistence),
mutex: sync.Mutex{},
}
}

Expand All @@ -40,18 +47,25 @@ func NewGroupRegistry(
func (gr *Groups) RegisterGroup(
signer *dkg.ThresholdSigner,
channelName string,
) {

) error {
gr.mutex.Lock()
defer gr.mutex.Unlock()

groupPublicKey := string(signer.GroupPublicKeyBytes())

gr.myGroups[groupPublicKey] = append(gr.myGroups[groupPublicKey],
&Membership{
Signer: signer,
ChannelName: channelName,
})
membership := &Membership{
Signer: signer,
ChannelName: channelName,
}

err := gr.storage.save(membership)
if err != nil {
return fmt.Errorf("could not persist membership to the storage: [%v]", err)
}

gr.myGroups[groupPublicKey] = append(gr.myGroups[groupPublicKey], membership)

return nil
}

// GetGroup gets a group by a groupPublicKey
Expand Down
27 changes: 15 additions & 12 deletions pkg/beacon/relay/registry/groups_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@ package registry
import (
"bytes"
"math/big"
"sync"
"testing"

bn256 "github.com/ethereum/go-ethereum/crypto/bn256/cloudflare"
Expand All @@ -14,19 +13,26 @@ import (
"github.com/keep-network/keep-core/pkg/subscription"
)

type noopPersistence struct {
}

func (np *noopPersistence) Save([]byte, string) error {
// noop
return nil
}

func TestRegisterGroup(t *testing.T) {
noopPersistence := &noopPersistence{}
chain := chainLocal.Connect(5, 3, big.NewInt(200)).ThresholdRelay()

gr := NewGroupRegistry(chain, noopPersistence)

signer := dkg.NewThresholdSigner(
group.MemberIndex(2),
new(bn256.G2).ScalarBaseMult(big.NewInt(10)),
big.NewInt(1),
)

gr := &Groups{
mutex: sync.Mutex{},
myGroups: make(map[string][]*Membership),
relayChain: chainLocal.Connect(5, 3, big.NewInt(200)).ThresholdRelay(),
}

gr.RegisterGroup(signer, "test_channel")

actual := gr.GetGroup(signer.GroupPublicKeyBytes())
Expand All @@ -50,12 +56,9 @@ func TestUnregisterStaleGroups(t *testing.T) {
mockChain := &mockGroupRegistrationInterface{
groupsToRemove: [][]byte{},
}
noopPersistence := &noopPersistence{}

gr := &Groups{
mutex: sync.Mutex{},
myGroups: make(map[string][]*Membership),
relayChain: mockChain,
}
gr := NewGroupRegistry(mockChain, noopPersistence)

signer1 := dkg.NewThresholdSigner(
group.MemberIndex(1),
Expand Down
34 changes: 34 additions & 0 deletions pkg/beacon/relay/registry/storage.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package registry

import (
"fmt"

"github.com/keep-network/keep-core/pkg/persistence"

"encoding/hex"
)

type storage interface {
save(membership *Membership) error
}

type persistentStorage struct {
handle persistence.Handle
}

func newStorage(persistence persistence.Handle) storage {
return &persistentStorage{
handle: persistence,
}
}

// Save converts a membership suitable for disk storage.
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 ps.handle.Save(membershipBytes, "/membership_"+hexGroupPublicKey+"_"+fmt.Sprint(membership.Signer.MemberID()))
}
97 changes: 97 additions & 0 deletions pkg/persistence/disk_persistence.go
Original file line number Diff line number Diff line change
@@ -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 system
func (f *file) remove(fileName string) error {
if f.fileName == "" {
return errNoFileExists
}

err := os.Remove(fileName)
if err != nil {
return err
}

return nil
}
Loading