diff --git a/cmd/cqld/adapter.go b/cmd/cqld/adapter.go index 8b41a2769..96437dd25 100644 --- a/cmd/cqld/adapter.go +++ b/cmd/cqld/adapter.go @@ -17,46 +17,29 @@ package main import ( - "bytes" "context" "database/sql" "os" + "sync" + "time" - bp "github.com/CovenantSQL/CovenantSQL/blockproducer" "github.com/CovenantSQL/CovenantSQL/consistent" "github.com/CovenantSQL/CovenantSQL/crypto/kms" - "github.com/CovenantSQL/CovenantSQL/kayak" - kt "github.com/CovenantSQL/CovenantSQL/kayak/types" "github.com/CovenantSQL/CovenantSQL/proto" "github.com/CovenantSQL/CovenantSQL/route" + "github.com/CovenantSQL/CovenantSQL/rpc" "github.com/CovenantSQL/CovenantSQL/storage" - "github.com/CovenantSQL/CovenantSQL/types" "github.com/CovenantSQL/CovenantSQL/utils" "github.com/CovenantSQL/CovenantSQL/utils/log" "github.com/pkg/errors" ) -const ( - // CmdSet is the command to set node - CmdSet = "set" - // CmdSetDatabase is the command to set database - CmdSetDatabase = "set_database" - // CmdDeleteDatabase is the command to del database - CmdDeleteDatabase = "delete_database" -) - // LocalStorage holds consistent and storage struct type LocalStorage struct { consistent *consistent.Consistent *storage.Storage } -type compiledLog struct { - cmdType string - queries []storage.Query - nodeToSet *proto.Node -} - func initStorage(dbFile string) (stor *LocalStorage, err error) { var st *storage.Storage if st, err = storage.New(dbFile); err != nil { @@ -67,9 +50,6 @@ func initStorage(dbFile string) (stor *LocalStorage, err error) { { Pattern: "CREATE TABLE IF NOT EXISTS `dht` (`id` TEXT NOT NULL PRIMARY KEY, `node` BLOB);", }, - { - Pattern: "CREATE TABLE IF NOT EXISTS `databases` (`id` TEXT NOT NULL PRIMARY KEY, `meta` BLOB);", - }, }) if err != nil { wd, _ := os.Getwd() @@ -84,339 +64,129 @@ func initStorage(dbFile string) (stor *LocalStorage, err error) { return } -// EncodePayload implements kayak.types.Handler.EncodePayload. -func (s *LocalStorage) EncodePayload(request interface{}) (data []byte, err error) { - var buf *bytes.Buffer - if buf, err = utils.EncodeMsgPack(request); err != nil { - err = errors.Wrap(err, "encode kayak payload failed") - return - } - - data = buf.Bytes() - return -} - -// DecodePayload implements kayak.types.Handler.DecodePayload. -func (s *LocalStorage) DecodePayload(data []byte) (request interface{}, err error) { - var kp *KayakPayload +// SetNode handles dht storage node update. +func (s *LocalStorage) SetNode(node *proto.Node) (err error) { + query := "INSERT OR REPLACE INTO `dht` (`id`, `node`) VALUES (?, ?);" + log.Debugf("sql: %#v", query) - if err = utils.DecodeMsgPack(data, &kp); err != nil { - err = errors.Wrap(err, "decode kayak payload failed") + nodeBuf, err := utils.EncodeMsgPack(node) + if err != nil { + err = errors.Wrap(err, "encode node failed") return } - request = kp - return -} - -// Check implements kayak.types.Handler.Check. -func (s *LocalStorage) Check(req interface{}) (err error) { - return nil -} - -// Commit implements kayak.types.Handler.Commit. -func (s *LocalStorage) Commit(req interface{}, isLeader bool) (_ interface{}, err error) { - var kp *KayakPayload - var cl *compiledLog - var ok bool - - if kp, ok = req.(*KayakPayload); !ok || kp == nil { - err = errors.Wrapf(kt.ErrInvalidLog, "invalid kayak payload %#v", req) - return + err = route.SetNodeAddrCache(node.ID.ToRawNodeID(), node.Addr) + if err != nil { + log.WithFields(log.Fields{ + "id": node.ID, + "addr": node.Addr, + }).WithError(err).Error("set node addr cache failed") } - - if cl, err = s.compileLog(kp); err != nil { - err = errors.Wrap(err, "compile log failed") - return + err = kms.SetNode(node) + if err != nil { + log.WithField("node", node).WithError(err).Error("kms set node failed") } - if cl.nodeToSet != nil { - err = route.SetNodeAddrCache(cl.nodeToSet.ID.ToRawNodeID(), cl.nodeToSet.Addr) - if err != nil { - log.WithFields(log.Fields{ - "id": cl.nodeToSet.ID, - "addr": cl.nodeToSet.Addr, - }).WithError(err).Error("set node addr cache failed") - } - err = kms.SetNode(cl.nodeToSet) - if err != nil { - log.WithField("node", cl.nodeToSet).WithError(err).Error("kms set node failed") - } - - // if s.consistent == nil, it is called during Init. and AddCache will be called by consistent.InitConsistent - if s.consistent != nil { - s.consistent.AddCache(*cl.nodeToSet) - } + // if s.consistent == nil, it is called during Init. and AddCache will be called by consistent.InitConsistent + if s.consistent != nil { + s.consistent.AddCache(*node) } // execute query - if _, err = s.Storage.Exec(context.Background(), cl.queries); err != nil { + if _, err = s.Storage.Exec(context.Background(), []storage.Query{ + { + Pattern: query, + Args: []sql.NamedArg{ + sql.Named("", node.ID), + sql.Named("", nodeBuf.Bytes()), + }, + }, + }); err != nil { err = errors.Wrap(err, "execute query in dht database failed") } return } -func (s *LocalStorage) compileLog(payload *KayakPayload) (result *compiledLog, err error) { - switch payload.Command { - case CmdSet: - var nodeToSet proto.Node - err = utils.DecodeMsgPack(payload.Data, &nodeToSet) - if err != nil { - log.WithError(err).Error("compileLog: unmarshal node from payload failed") - return - } - query := "INSERT OR REPLACE INTO `dht` (`id`, `node`) VALUES (?, ?);" - log.Debugf("sql: %#v", query) - result = &compiledLog{ - cmdType: payload.Command, - queries: []storage.Query{ - { - Pattern: query, - Args: []sql.NamedArg{ - sql.Named("", nodeToSet.ID), - sql.Named("", payload.Data), - }, - }, - }, - nodeToSet: &nodeToSet, - } - case CmdSetDatabase: - var instance types.ServiceInstance - if err = utils.DecodeMsgPack(payload.Data, &instance); err != nil { - log.WithError(err).Error("compileLog: unmarshal instance meta failed") - return - } - query := "INSERT OR REPLACE INTO `databases` (`id`, `meta`) VALUES (? ,?);" - result = &compiledLog{ - cmdType: payload.Command, - queries: []storage.Query{ - { - Pattern: query, - Args: []sql.NamedArg{ - sql.Named("", string(instance.DatabaseID)), - sql.Named("", payload.Data), - }, - }, - }, - } - case CmdDeleteDatabase: - var instance types.ServiceInstance - if err = utils.DecodeMsgPack(payload.Data, &instance); err != nil { - log.WithError(err).Error("compileLog: unmarshal instance id failed") - return - } - // TODO(xq262144), should add additional limit 1 after delete clause - // however, currently the go-sqlite3 - query := "DELETE FROM `databases` WHERE `id` = ?" - result = &compiledLog{ - cmdType: payload.Command, - queries: []storage.Query{ - { - Pattern: query, - Args: []sql.NamedArg{ - sql.Named("", string(instance.DatabaseID)), - }, - }, - }, - } - default: - err = errors.Errorf("undefined command: %v", payload.Command) - log.WithError(err).Error("compile log failed") - } - return +// KVServer holds LocalStorage instance and implements consistent persistence interface. +type KVServer struct { + current proto.NodeID + peers *proto.Peers + storage *LocalStorage + ctx context.Context + cancelCtx context.CancelFunc + timeout time.Duration + wg sync.WaitGroup } -// KayakKVServer holds kayak.Runtime and LocalStorage -type KayakKVServer struct { - Runtime *kayak.Runtime - KVStorage *LocalStorage +// NewKVServer returns the kv server instance. +func NewKVServer(currentNode proto.NodeID, peers *proto.Peers, storage *LocalStorage, timeout time.Duration) (s *KVServer) { + ctx, cancelCtx := context.WithCancel(context.Background()) + + return &KVServer{ + current: currentNode, + peers: peers, + storage: storage, + ctx: ctx, + cancelCtx: cancelCtx, + timeout: timeout, + } } // Init implements consistent.Persistence -func (s *KayakKVServer) Init(storePath string, initNodes []proto.Node) (err error) { +func (s *KVServer) Init(storePath string, initNodes []proto.Node) (err error) { for _, n := range initNodes { - var nodeBuf *bytes.Buffer - nodeBuf, err = utils.EncodeMsgPack(n) - if err != nil { - log.WithError(err).Error("marshal node failed") - return - } - payload := &KayakPayload{ - Command: CmdSet, - Data: nodeBuf.Bytes(), - } - _, err = s.KVStorage.Commit(payload, true) + err = s.storage.SetNode(&n) if err != nil { - log.WithError(err).Error("init kayak KV commit node failed") + log.WithError(err).Error("init dht kv server failed") return } } return } -// KayakPayload is the payload used in kayak Leader and Follower -type KayakPayload struct { - Command string - Data []byte +// SetNode implements consistent.Persistence +func (s *KVServer) SetNode(node *proto.Node) (err error) { + return s.SetNodeEx(node, 1, s.current) } -// SetNode implements consistent.Persistence -func (s *KayakKVServer) SetNode(node *proto.Node) (err error) { - nodeBuf, err := utils.EncodeMsgPack(node) - if err != nil { - log.WithError(err).Error("marshal node failed") +// SetNodeEx is used by gossip service to broadcast to other nodes. +func (s *KVServer) SetNodeEx(node *proto.Node, ttl uint32, origin proto.NodeID) (err error) { + log.WithFields(log.Fields{ + "node": node.ID, + "ttl": ttl, + "origin": origin, + }).Debug("update node to kv storage") + + // set local + if err = s.storage.SetNode(node); err != nil { + err = errors.Wrap(err, "set node failed") return } - payload := &KayakPayload{ - Command: CmdSet, - Data: nodeBuf.Bytes(), - } - _, _, err = s.Runtime.Apply(context.Background(), payload) - if err != nil { - log.Errorf("apply set node failed: %#v\nPayload:\n %#v", err, payload) + if ttl > 0 { + s.nonBlockingSync(node, origin, ttl-1) } return } // DelNode implements consistent.Persistence -func (s *KayakKVServer) DelNode(nodeID proto.NodeID) (err error) { +func (s *KVServer) DelNode(nodeID proto.NodeID) (err error) { // no need to del node currently return } // Reset implements consistent.Persistence -func (s *KayakKVServer) Reset() (err error) { - // no need to reset for kayak - return -} - -// GetDatabase implements blockproducer.DBMetaPersistence. -func (s *KayakKVServer) GetDatabase(dbID proto.DatabaseID) (instance types.ServiceInstance, err error) { - var result [][]interface{} - query := "SELECT `meta` FROM `databases` WHERE `id` = ? LIMIT 1" - _, _, result, err = s.KVStorage.Query(context.Background(), []storage.Query{ - { - Pattern: query, - Args: []sql.NamedArg{ - sql.Named("", string(dbID)), - }, - }, - }) - if err != nil { - log.WithField("db", dbID).WithError(err).Error("query database instance meta failed") - return - } - - if len(result) <= 0 || len(result[0]) <= 0 { - err = bp.ErrNoSuchDatabase - return - } - - var rawInstanceMeta []byte - var ok bool - if rawInstanceMeta, ok = result[0][0].([]byte); !ok { - err = bp.ErrNoSuchDatabase - return - } - - err = utils.DecodeMsgPack(rawInstanceMeta, &instance) - return -} - -// SetDatabase implements blockproducer.DBMetaPersistence. -func (s *KayakKVServer) SetDatabase(meta types.ServiceInstance) (err error) { - var metaBuf *bytes.Buffer - if metaBuf, err = utils.EncodeMsgPack(meta); err != nil { - return - } - - payload := &KayakPayload{ - Command: CmdSetDatabase, - Data: metaBuf.Bytes(), - } - - _, _, err = s.Runtime.Apply(context.Background(), payload) - if err != nil { - log.Errorf("apply set database failed: %#v\nPayload:\n %#v", err, payload) - } - - return -} - -// DeleteDatabase implements blockproducer.DBMetaPersistence. -func (s *KayakKVServer) DeleteDatabase(dbID proto.DatabaseID) (err error) { - meta := types.ServiceInstance{ - DatabaseID: dbID, - } - - var metaBuf *bytes.Buffer - if metaBuf, err = utils.EncodeMsgPack(meta); err != nil { - return - } - payload := &KayakPayload{ - Command: CmdDeleteDatabase, - Data: metaBuf.Bytes(), - } - - _, _, err = s.Runtime.Apply(context.Background(), payload) - if err != nil { - log.Errorf("apply set database failed: %#v\nPayload:\n %#v", err, payload) - } - - return -} - -// GetAllDatabases implements blockproducer.DBMetaPersistence. -func (s *KayakKVServer) GetAllDatabases() (instances []types.ServiceInstance, err error) { - var result [][]interface{} - query := "SELECT `meta` FROM `databases`" - _, _, result, err = s.KVStorage.Query(context.Background(), []storage.Query{ - { - Pattern: query, - }, - }) - if err != nil { - log.WithError(err).Error("query all database instance meta failed") - return - } - - instances = make([]types.ServiceInstance, 0, len(result)) - - for _, row := range result { - if len(row) <= 0 { - continue - } - - var instance types.ServiceInstance - var rawInstanceMeta []byte - var ok bool - if rawInstanceMeta, ok = row[0].([]byte); !ok { - err = bp.ErrNoSuchDatabase - continue - } - - if err = utils.DecodeMsgPack(rawInstanceMeta, &instance); err != nil { - continue - } - - instances = append(instances, instance) - } - - if len(instances) > 0 { - err = nil - } - +func (s *KVServer) Reset() (err error) { return } // GetAllNodeInfo implements consistent.Persistence -func (s *KayakKVServer) GetAllNodeInfo() (nodes []proto.Node, err error) { +func (s *KVServer) GetAllNodeInfo() (nodes []proto.Node, err error) { var result [][]interface{} query := "SELECT `node` FROM `dht`;" - _, _, result, err = s.KVStorage.Query(context.Background(), []storage.Query{ + _, _, result, err = s.storage.Query(context.Background(), []storage.Query{ { Pattern: query, }, @@ -451,5 +221,38 @@ func (s *KayakKVServer) GetAllNodeInfo() (nodes []proto.Node, err error) { if len(nodes) > 0 { err = nil } + return } + +// Stop stops the dht node server and wait for inflight gossip requests. +func (s *KVServer) Stop() { + if s.cancelCtx != nil { + s.cancelCtx() + } + s.wg.Wait() +} + +func (s *KVServer) nonBlockingSync(node *proto.Node, origin proto.NodeID, ttl uint32) { + if s.peers == nil { + return + } + + c, cancel := context.WithTimeout(s.ctx, s.timeout) + defer cancel() + for _, n := range s.peers.Servers { + if n != s.current && n != origin { + // sync + req := &GossipRequest{ + Node: node, + TTL: ttl, + } + + s.wg.Add(1) + go func(node proto.NodeID) { + defer s.wg.Done() + _ = rpc.NewCaller().CallNodeWithContext(c, node, route.DHTGSetNode.String(), req, nil) + }(n) + } + } +} diff --git a/cmd/cqld/bench_test.go b/cmd/cqld/bench_test.go index 37771d128..3bbd21fd0 100644 --- a/cmd/cqld/bench_test.go +++ b/cmd/cqld/bench_test.go @@ -235,7 +235,7 @@ func TestStartBP_CallRPC(t *testing.T) { } } -func BenchmarkKayakKVServer_GetAllNodeInfo(b *testing.B) { +func BenchmarkKVServer_GetAllNodeInfo(b *testing.B) { log.SetLevel(log.DebugLevel) start3BPs() diff --git a/cmd/cqld/bootstrap.go b/cmd/cqld/bootstrap.go index 80bf2d423..72dca2816 100644 --- a/cmd/cqld/bootstrap.go +++ b/cmd/cqld/bootstrap.go @@ -20,18 +20,13 @@ import ( "fmt" "os" "os/signal" - "path/filepath" "syscall" "time" "github.com/CovenantSQL/CovenantSQL/api" - bp "github.com/CovenantSQL/CovenantSQL/blockproducer" "github.com/CovenantSQL/CovenantSQL/conf" "github.com/CovenantSQL/CovenantSQL/crypto/kms" - "github.com/CovenantSQL/CovenantSQL/kayak" - kt "github.com/CovenantSQL/CovenantSQL/kayak/types" - kl "github.com/CovenantSQL/CovenantSQL/kayak/wal" "github.com/CovenantSQL/CovenantSQL/proto" "github.com/CovenantSQL/CovenantSQL/route" "github.com/CovenantSQL/CovenantSQL/rpc" @@ -42,18 +37,11 @@ import ( ) const ( - kayakServiceName = "Kayak" - kayakApplyMethodName = "Apply" - kayakFetchMethodName = "Fetch" - kayakWalFileName = "kayak.ldb" - kayakPrepareTimeout = 5 * time.Second - kayakCommitTimeout = time.Minute - kayakLogWaitTimeout = 10 * time.Second + dhtGossipServiceName = "DHTG" + dhtGossipTimeout = time.Second * 20 ) func runNode(nodeID proto.NodeID, listenAddr string) (err error) { - rootPath := conf.GConf.WorkingRoot - genesis, err := loadGenesis() if err != nil { return @@ -124,30 +112,29 @@ func runNode(nodeID proto.NodeID, listenAddr string) (err error) { return err } - // init kayak - log.Info("init kayak runtime") - var kayakRuntime *kayak.Runtime - if kayakRuntime, err = initKayakTwoPC(rootPath, thisNode, peers, st, server); err != nil { - log.WithError(err).Error("init kayak runtime failed") - return err - } - - // init kayak and consistent - log.Info("init kayak and consistent runtime") - kvServer := &KayakKVServer{ - Runtime: kayakRuntime, - KVStorage: st, - } + // init dht node server + log.Info("init consistent runtime") + kvServer := NewKVServer(thisNode.ID, peers, st, dhtGossipTimeout) dht, err := route.NewDHTService(conf.GConf.DHTFileName, kvServer, true) if err != nil { log.WithError(err).Error("init consistent hash failed") return err } + defer kvServer.Stop() - // set consistent handler to kayak storage - kvServer.KVStorage.consistent = dht.Consistent + // set consistent handler to local storage + kvServer.storage.consistent = dht.Consistent - // register service rpc + // register gossip service rpc + gossipService := NewGossipService(kvServer) + log.Info("register dht gossip service rpc") + err = server.RegisterService(route.DHTGossipRPCName, gossipService) + if err != nil { + log.WithError(err).Error("register dht gossip service failed") + return err + } + + // register dht service rpc log.Info("register dht service rpc") err = server.RegisterService(route.DHTRPCName, dht) if err != nil { @@ -211,49 +198,8 @@ func createServer(privateKeyPath, pubKeyStorePath string, masterKey []byte, list return } -func initKayakTwoPC(rootDir string, node *proto.Node, peers *proto.Peers, h kt.Handler, server *rpc.Server) (runtime *kayak.Runtime, err error) { - // create kayak config - log.Info("create kayak config") - - walPath := filepath.Join(rootDir, kayakWalFileName) - - var logWal kt.Wal - if logWal, err = kl.NewLevelDBWal(walPath); err != nil { - err = errors.Wrap(err, "init kayak log pool failed") - return - } - - config := &kt.RuntimeConfig{ - Handler: h, - PrepareThreshold: 1.0, - CommitThreshold: 1.0, - PrepareTimeout: kayakPrepareTimeout, - CommitTimeout: kayakCommitTimeout, - LogWaitTimeout: kayakLogWaitTimeout, - Peers: peers, - Wal: logWal, - NodeID: node.ID, - ServiceName: kayakServiceName, - ApplyMethodName: kayakApplyMethodName, - FetchMethodName: kayakFetchMethodName, - } - - // create kayak runtime - log.Info("init kayak runtime") - if runtime, err = kayak.NewRuntime(config); err != nil { - err = errors.Wrap(err, "init kayak runtime failed") - return - } - - // register rpc service - if _, err = NewKayakService(server, kayakServiceName, runtime); err != nil { - err = errors.Wrap(err, "init kayak rpc service failed") - return - } - - // init runtime - log.Info("start kayak runtime") - runtime.Start() +func initDHTGossip() (err error) { + log.Info("init gossip service") return } diff --git a/cmd/cqld/gossip.go b/cmd/cqld/gossip.go new file mode 100644 index 000000000..46125332c --- /dev/null +++ b/cmd/cqld/gossip.go @@ -0,0 +1,43 @@ +/* + * Copyright 2019 The CovenantSQL Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package main + +import "github.com/CovenantSQL/CovenantSQL/proto" + +// GossipRequest defines the gossip request payload. +type GossipRequest struct { + proto.Envelope + Node *proto.Node + TTL uint32 +} + +// GossipService defines the gossip service instance. +type GossipService struct { + s *KVServer +} + +// NewGossipService returns new gossip service. +func NewGossipService(s *KVServer) *GossipService { + return &GossipService{ + s: s, + } +} + +// SetNode update current node info and broadcast node update request. +func (s *GossipService) SetNode(req *GossipRequest, resp *interface{}) (err error) { + return s.s.SetNodeEx(req.Node, req.TTL, req.GetNodeID().ToNodeID()) +} diff --git a/cmd/cqld/kayak.go b/cmd/cqld/kayak.go deleted file mode 100644 index ee36f886a..000000000 --- a/cmd/cqld/kayak.go +++ /dev/null @@ -1,55 +0,0 @@ -/* - * Copyright 2018 The CovenantSQL Authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package main - -import ( - "github.com/CovenantSQL/CovenantSQL/kayak" - kt "github.com/CovenantSQL/CovenantSQL/kayak/types" - "github.com/CovenantSQL/CovenantSQL/rpc" -) - -// KayakService defines the leader service kayak. -type KayakService struct { - serviceName string - rt *kayak.Runtime -} - -// NewKayakService returns new kayak service instance for block producer consensus. -func NewKayakService(server *rpc.Server, serviceName string, rt *kayak.Runtime) (s *KayakService, err error) { - s = &KayakService{ - serviceName: serviceName, - rt: rt, - } - err = server.RegisterService(serviceName, s) - return -} - -// Apply handles kayak apply call. -func (s *KayakService) Apply(req *kt.ApplyRequest, _ *interface{}) (err error) { - return s.rt.FollowerApply(req.Log) -} - -// Fetch handles kayak log fetch call. -func (s *KayakService) Fetch(req *kt.FetchRequest, resp *kt.FetchResponse) (err error) { - var l *kt.Log - if l, err = s.rt.Fetch(req.GetContext(), req.Index); err != nil { - return - } - - resp.Log = l - return -} diff --git a/cmd/cqld/main.go b/cmd/cqld/main.go index 723dee0f0..00e7fd1bb 100644 --- a/cmd/cqld/main.go +++ b/cmd/cqld/main.go @@ -143,7 +143,7 @@ func main() { defer utils.StopProfile() if err := runNode(conf.GConf.ThisNodeID, conf.GConf.ListenAddr); err != nil { - log.WithError(err).Fatal("run kayak failed") + log.WithError(err).Fatal("run block producer node failed") } log.Info("server stopped") diff --git a/route/acl.go b/route/acl.go index 032aff54e..ae7a465ee 100644 --- a/route/acl.go +++ b/route/acl.go @@ -67,8 +67,8 @@ const ( DHTFindNeighbor // DHTFindNode gets node info DHTFindNode - // KayakCall is used by BP for data consistency - KayakCall + // DHTGSetNode is used by BP for dht data gossip + DHTGSetNode // MetricUploadMetrics uploads node metrics MetricUploadMetrics // DBSQuery is used by client to read/write database @@ -121,6 +121,8 @@ const ( // DHTRPCName defines the block producer dh-rpc service name DHTRPCName = "DHT" + // DHTGossipRPCName defines the block producer dh-rpc gossip service name + DHTGossipRPCName = "DHTG" // BlockProducerRPCName defines main chain rpc name BlockProducerRPCName = "MCC" // SQLChainRPCName defines the sql chain rpc name @@ -138,10 +140,10 @@ func (s RemoteFunc) String() string { return "DHT.FindNeighbor" case DHTFindNode: return "DHT.FindNode" + case DHTGSetNode: + return "DHTG.SetNode" case MetricUploadMetrics: return "Metric.UploadMetrics" - case KayakCall: - return "Kayak.Call" case DBSQuery: return "DBS.Query" case DBSAck: @@ -213,8 +215,8 @@ func IsPermitted(callerEnvelope *proto.Envelope, funcName RemoteFunc) (ok bool) // DHT related case DHTPing, DHTFindNode, DHTFindNeighbor, MetricUploadMetrics: return true - // Kayak related - case KayakCall: + // DHTGSetNode is for block producer to update node info + case DHTGSetNode: return false // DBSDeploy case DBSDeploy: diff --git a/route/acl_test.go b/route/acl_test.go index 60732f031..630e649fb 100644 --- a/route/acl_test.go +++ b/route/acl_test.go @@ -51,8 +51,8 @@ func TestIsPermitted(t *testing.T) { nodeID := proto.NodeID("0000") testEnv := &proto.Envelope{NodeID: nodeID.ToRawNodeID()} testAnonymous := &proto.Envelope{NodeID: kms.AnonymousRawNodeID} - So(IsPermitted(&proto.Envelope{NodeID: &conf.GConf.BP.RawNodeID}, KayakCall), ShouldBeTrue) - So(IsPermitted(testEnv, KayakCall), ShouldBeFalse) + So(IsPermitted(&proto.Envelope{NodeID: &conf.GConf.BP.RawNodeID}, DHTGSetNode), ShouldBeTrue) + So(IsPermitted(testEnv, DHTGSetNode), ShouldBeFalse) So(IsPermitted(testEnv, DHTFindNode), ShouldBeTrue) So(IsPermitted(testEnv, RemoteFunc(9999)), ShouldBeFalse) So(IsPermitted(testAnonymous, DHTFindNode), ShouldBeFalse) diff --git a/rpc/rpcutil.go b/rpc/rpcutil.go index 3588ca5a5..c2bacb5f4 100644 --- a/rpc/rpcutil.go +++ b/rpc/rpcutil.go @@ -28,7 +28,6 @@ import ( "time" "github.com/CovenantSQL/CovenantSQL/crypto/kms" - kt "github.com/CovenantSQL/CovenantSQL/kayak/types" "github.com/CovenantSQL/CovenantSQL/proto" "github.com/CovenantSQL/CovenantSQL/route" "github.com/CovenantSQL/CovenantSQL/utils/log" @@ -250,26 +249,13 @@ func GetNodeAddr(id *proto.RawNodeID) (addr string, err error) { if err != nil { //log.WithField("target", id.String()).WithError(err).Debug("get node addr from cache failed") if err == route.ErrUnknownNodeID { - BPs := route.GetBPs() - if len(BPs) == 0 { - err = errors.New("no available BP") - return - } - client := NewCaller() - reqFN := &proto.FindNodeReq{ - ID: proto.NodeID(id.String()), - } - respFN := new(proto.FindNodeResp) - - bp := BPs[rand.Intn(len(BPs))] - method := "DHT.FindNode" - err = client.CallNode(bp, method, reqFN, respFN) + var node *proto.Node + node, err = FindNodeInBP(id) if err != nil { - err = errors.Wrapf(err, "call dht rpc %s to %s failed", method, bp) return } - route.SetNodeAddrCache(id, respFN.Node.Addr) - addr = respFN.Node.Addr + route.SetNodeAddrCache(id, node.Addr) + addr = node.Addr } } return @@ -281,24 +267,10 @@ func GetNodeInfo(id *proto.RawNodeID) (nodeInfo *proto.Node, err error) { if err != nil { //log.WithField("target", id.String()).WithError(err).Info("get node info from KMS failed") if errors.Cause(err) == kms.ErrKeyNotFound { - BPs := route.GetBPs() - if len(BPs) == 0 { - err = errors.New("no available BP") - return - } - client := NewCaller() - reqFN := &proto.FindNodeReq{ - ID: proto.NodeID(id.String()), - } - respFN := new(proto.FindNodeResp) - bp := BPs[rand.Intn(len(BPs))] - method := "DHT.FindNode" - err = client.CallNode(bp, method, reqFN, respFN) + nodeInfo, err = FindNodeInBP(id) if err != nil { - err = errors.Wrapf(err, "call dht rpc %s to %s failed", method, bp) return } - nodeInfo = respFN.Node errSet := route.SetNodeAddrCache(id, nodeInfo.Addr) if errSet != nil { log.WithError(errSet).Warning("set node addr cache failed") @@ -312,6 +284,40 @@ func GetNodeInfo(id *proto.RawNodeID) (nodeInfo *proto.Node, err error) { return } +// FindNodeInBP find node in block producer dht service. +func FindNodeInBP(id *proto.RawNodeID) (node *proto.Node, err error) { + bps := route.GetBPs() + if len(bps) == 0 { + err = errors.New("no available BP") + return + } + client := NewCaller() + req := &proto.FindNodeReq{ + ID: proto.NodeID(id.String()), + } + resp := new(proto.FindNodeResp) + bpCount := len(bps) + offset := rand.Intn(bpCount) + method := route.DHTFindNode.String() + + for i := 0; i != bpCount; i++ { + bp := bps[(offset+i)%bpCount] + err = client.CallNode(bp, method, req, resp) + if err == nil { + node = resp.Node + return + } + + log.WithFields(log.Fields{ + "method": method, + "bp": bp, + }).WithError(err).Warning("call dht rpc failed") + } + + err = errors.Wrapf(err, "could not find node in all block producers") + return +} + // PingBP Send DHT.Ping Request with Anonymous ETLS session. func PingBP(node *proto.Node, BPNodeID proto.NodeID) (err error) { client := NewCaller() @@ -430,11 +436,10 @@ func RegisterNodeToBP(timeout time.Duration) (err error) { err := PingBP(localNodeInfo, id) if err == nil { log.Infof("ping BP succeed: %v", localNodeInfo) - ch <- id - return - } - if strings.Contains(err.Error(), kt.ErrNotLeader.Error()) { - log.Debug("stop ping non leader BP node") + select { + case ch <- id: + default: + } return } @@ -446,7 +451,6 @@ func RegisterNodeToBP(timeout time.Duration) (err error) { select { case bp := <-pingWaitCh: - close(pingWaitCh) log.WithField("BP", bp).Infof("ping BP succeed") case <-time.After(timeout): return errors.New("ping BP timeout")