Skip to content

Commit bbdf8bc

Browse files
author
auxten
committed
remove periodicPingBlockProducer, BP will init nodes during consistent.InitConsistent()
1 parent d4ebe40 commit bbdf8bc

4 files changed

Lines changed: 40 additions & 39 deletions

File tree

cmd/thunderdbd/adapter.go

Lines changed: 32 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,10 @@ func (s *LocalStorage) Commit(ctx context.Context, wb twopc.WriteBatch) (err err
9898
log.Errorf("decode log failed: %s", err)
9999
return
100100
}
101+
return s.commit(ctx, payload)
102+
}
103+
104+
func (s *LocalStorage) commit(ctx context.Context, payload *KayakPayload) (err error) {
101105
var nodeToSet proto.Node
102106
err = utils.DecodeMsgPack(payload.Data, &nodeToSet)
103107
if err != nil {
@@ -117,10 +121,14 @@ func (s *LocalStorage) Commit(ctx context.Context, wb twopc.WriteBatch) (err err
117121
if err != nil {
118122
log.Errorf("kms set node failed: %v", err)
119123
}
120-
err = s.consistent.AddCache(nodeToSet)
121-
if err != nil {
122-
//TODO(auxten) even no error will be returned, there may be some inconsistency and needs sync periodically
123-
log.Errorf("add consistent cache failed: %s", err)
124+
125+
// if s.consistent == nil, it is called during Init. and AddCache will be called by consistent.InitConsistent
126+
if s.consistent != nil {
127+
err = s.consistent.AddCache(nodeToSet)
128+
if err != nil {
129+
//TODO(auxten) even no error will be returned, there may be some inconsistency and needs sync periodically
130+
log.Errorf("add consistent cache failed: %s", err)
131+
}
124132
}
125133

126134
return s.Storage.Commit(ctx, execLog)
@@ -234,7 +242,24 @@ type KayakKVServer struct {
234242

235243
// Init implements consistent.Persistence
236244
func (s *KayakKVServer) Init(storePath string, initNodes []proto.Node) (err error) {
237-
//FIXME(auxten) implements KayakKVServer.Init
245+
for _, n := range initNodes {
246+
var nodeBuf *bytes.Buffer
247+
nodeBuf, err = utils.EncodeMsgPack(n)
248+
if err != nil {
249+
log.Errorf("marshal node failed: %v", err)
250+
return
251+
}
252+
payload := &KayakPayload{
253+
Command: CmdSet,
254+
Data: nodeBuf.Bytes(),
255+
}
256+
257+
err = s.Storage.commit(context.Background(), payload)
258+
if err != nil {
259+
log.Errorf("init kayak KV node failed: %v", err)
260+
return
261+
}
262+
}
238263
return
239264
}
240265

@@ -248,7 +273,7 @@ type KayakPayload struct {
248273
func (s *KayakKVServer) SetNode(node *proto.Node) (err error) {
249274
nodeBuf, err := utils.EncodeMsgPack(node)
250275
if err != nil {
251-
log.Errorf("marshal node failed: %s", err)
276+
log.Errorf("marshal node failed: %v", err)
252277
return
253278
}
254279
payload := &KayakPayload{
@@ -258,7 +283,7 @@ func (s *KayakKVServer) SetNode(node *proto.Node) (err error) {
258283

259284
writeData, err := utils.EncodeMsgPack(payload)
260285
if err != nil {
261-
log.Errorf("marshal payload failed: %s", err)
286+
log.Errorf("marshal payload failed: %v", err)
262287
return err
263288
}
264289

cmd/thunderdbd/bootstrap.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -146,7 +146,7 @@ func runNode(nodeID proto.NodeID, listenAddr string) (err error) {
146146
}
147147

148148
log.Info(conf.StartSucceedMessage)
149-
go periodicPingBlockProducer()
149+
//go periodicPingBlockProducer()
150150

151151
// start server
152152
server.Serve()

cmd/thunderdbd/initconf.go

Lines changed: 0 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -60,23 +60,6 @@ func initNodePeers(nodeID proto.NodeID, publicKeystorePath string) (nodes *[]pro
6060
PubKey: publicKey,
6161
})
6262
}
63-
//if n.Role == proto.Client {
64-
// var publicKeyBytes []byte
65-
// var clientPublicKey *asymmetric.PublicKey
66-
// //02ec784ca599f21ef93fe7abdc68d78817ab6c9b31f2324d15ea174d9da498b4c4
67-
// publicKeyBytes, err = hex.DecodeString("02ec784ca599f21ef93fe7abdc68d78817ab6c9b31f2324d15ea174d9da498b4c4")
68-
// if err != nil {
69-
// log.Errorf("hex decode clientPublicKey error: %s", err)
70-
// return
71-
// }
72-
// clientPublicKey, err = asymmetric.ParsePubKey(publicKeyBytes)
73-
// if err != nil {
74-
// log.Errorf("parse clientPublicKey error: %s", err)
75-
// return
76-
// }
77-
// //FIXME(auxten): read public key from conf
78-
// (*conf.GConf.KnownNodes)[i].PublicKey = clientPublicKey
79-
//}
8063
}
8164
}
8265

consistent/consistent.go

Lines changed: 7 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -77,8 +77,6 @@ type Consistent struct {
7777
}
7878

7979
// InitConsistent creates a new Consistent object with a default setting of 20 replicas for each entry.
80-
//
81-
// To change the number of replicas, set NumberOfReplicas before adding entries.
8280
func InitConsistent(storePath string, persistImpl Persistence, initBP bool) (c *Consistent, err error) {
8381
var BPNodes []proto.Node
8482
if initBP {
@@ -90,24 +88,19 @@ func InitConsistent(storePath string, persistImpl Persistence, initBP bool) (c *
9088
BPNodes = conf.GConf.SeedBPNodes
9189
}
9290

93-
// Create new public key store
94-
//err = kms.InitPublicKeyStore(storePath, BPNode)
95-
//if err != nil {
96-
// log.Errorf("init public keystore failed: %s", err)
97-
// return
98-
//}
99-
//IDs, err := kms.GetAllNodeID()
100-
//if err != nil {
101-
// log.Errorf("get all node id failed: %s", err)
102-
// return
103-
//}
10491
c = &Consistent{
10592
//TODO(auxten): reduce NumberOfReplicas
10693
NumberOfReplicas: 20,
10794
circle: make(map[proto.NodeKey]*proto.Node),
10895
persist: persistImpl,
10996
}
110-
c.persist.Init(storePath, BPNodes)
97+
98+
err = c.persist.Init(storePath, BPNodes)
99+
if err != nil {
100+
log.Errorf("init persist BP nodes failed: %v", err)
101+
return
102+
}
103+
111104
nodes, err := c.persist.GetAllNodeInfo()
112105
if err != nil {
113106
log.Errorf("get all node id failed: %s", err)

0 commit comments

Comments
 (0)