Skip to content

Commit 7735fdf

Browse files
committed
Add direct RPC server for miner
1 parent ac184ee commit 7735fdf

6 files changed

Lines changed: 59 additions & 23 deletions

File tree

cmd/cql-minerd/dbms.go

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,13 +21,14 @@ import (
2121

2222
"github.com/CovenantSQL/CovenantSQL/conf"
2323
"github.com/CovenantSQL/CovenantSQL/crypto/hash"
24-
rpc "github.com/CovenantSQL/CovenantSQL/rpc/mux"
24+
"github.com/CovenantSQL/CovenantSQL/rpc"
25+
"github.com/CovenantSQL/CovenantSQL/rpc/mux"
2526
"github.com/CovenantSQL/CovenantSQL/worker"
2627
)
2728

2829
var rootHash = hash.Hash{}
2930

30-
func startDBMS(server *rpc.Server, onCreateDB func()) (dbms *worker.DBMS, err error) {
31+
func startDBMS(server *mux.Server, direct *rpc.Server, onCreateDB func()) (dbms *worker.DBMS, err error) {
3132
if conf.GConf.Miner == nil {
3233
err = errors.New("invalid database config")
3334
return
@@ -36,6 +37,7 @@ func startDBMS(server *rpc.Server, onCreateDB func()) (dbms *worker.DBMS, err er
3637
cfg := &worker.DBMSConfig{
3738
RootDir: conf.GConf.Miner.RootDir,
3839
Server: server,
40+
DirectServer: direct,
3941
MaxReqTimeGap: conf.GConf.Miner.MaxReqTimeGap,
4042
OnCreateDatabase: onCreateDB,
4143
}

cmd/cql-minerd/main.go

Lines changed: 17 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,8 @@ import (
3434
"github.com/CovenantSQL/CovenantSQL/crypto/asymmetric"
3535
"github.com/CovenantSQL/CovenantSQL/crypto/kms"
3636
"github.com/CovenantSQL/CovenantSQL/metric"
37-
rpc "github.com/CovenantSQL/CovenantSQL/rpc/mux"
37+
"github.com/CovenantSQL/CovenantSQL/rpc"
38+
"github.com/CovenantSQL/CovenantSQL/rpc/mux"
3839
"github.com/CovenantSQL/CovenantSQL/utils"
3940
"github.com/CovenantSQL/CovenantSQL/utils/log"
4041
_ "github.com/CovenantSQL/CovenantSQL/utils/log/debug"
@@ -169,8 +170,11 @@ func main() {
169170
conf.GConf.GenerateKeyPair = genKeyPair
170171

171172
// start rpc
172-
var server *rpc.Server
173-
if server, err = initNode(); err != nil {
173+
var (
174+
server *mux.Server
175+
direct *rpc.Server
176+
)
177+
if server, direct, err = initNode(); err != nil {
174178
log.WithError(err).Fatal("init node failed")
175179
}
176180

@@ -213,7 +217,7 @@ func main() {
213217

214218
// start dbms
215219
var dbms *worker.DBMS
216-
if dbms, err = startDBMS(server, func() {
220+
if dbms, err = startDBMS(server, direct, func() {
217221
sendProvideService(reg)
218222
}); err != nil {
219223
log.WithError(err).Fatal("start dbms failed")
@@ -225,10 +229,15 @@ func main() {
225229
go func() {
226230
server.Serve()
227231
}()
228-
defer func() {
229-
_ = server.Listener.Close()
230-
server.Stop()
231-
}()
232+
defer server.Stop()
233+
234+
// start direct rpc server
235+
if direct != nil {
236+
go func() {
237+
direct.Serve()
238+
}()
239+
defer direct.Stop()
240+
}
232241

233242
if metricLog {
234243
go metrics.Log(metrics.DefaultRegistry, 5*time.Second, log.StandardLogger())

cmd/cql-minerd/node.go

Lines changed: 23 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -26,12 +26,13 @@ import (
2626
"github.com/CovenantSQL/CovenantSQL/conf"
2727
"github.com/CovenantSQL/CovenantSQL/crypto/kms"
2828
"github.com/CovenantSQL/CovenantSQL/route"
29-
rpc "github.com/CovenantSQL/CovenantSQL/rpc/mux"
29+
"github.com/CovenantSQL/CovenantSQL/rpc"
30+
"github.com/CovenantSQL/CovenantSQL/rpc/mux"
3031
"github.com/CovenantSQL/CovenantSQL/utils"
3132
"github.com/CovenantSQL/CovenantSQL/utils/log"
3233
)
3334

34-
func initNode() (server *rpc.Server, err error) {
35+
func initNode() (server *mux.Server, direct *rpc.Server, err error) {
3536
var masterKey []byte
3637
if !conf.GConf.UseTestMasterKey {
3738
// read master key
@@ -53,30 +54,44 @@ func initNode() (server *rpc.Server, err error) {
5354
// init kms routing
5455
route.InitKMS(conf.GConf.PubKeyStoreFile)
5556

56-
err = rpc.RegisterNodeToBP(30 * time.Second)
57+
err = mux.RegisterNodeToBP(30 * time.Second)
5758
if err != nil {
5859
log.Fatalf("register node to BP failed: %v", err)
5960
}
6061

6162
// init server
63+
utils.RemoveAll(conf.GConf.PubKeyStoreFile + "*")
6264
if server, err = createServer(
63-
conf.GConf.PrivateKeyFile, conf.GConf.PubKeyStoreFile, masterKey, conf.GConf.ListenAddr); err != nil {
65+
conf.GConf.PrivateKeyFile, masterKey, conf.GConf.ListenAddr); err != nil {
6466
log.WithError(err).Error("create server failed")
6567
return
6668
}
69+
if direct, err = createDirectServer(
70+
conf.GConf.PrivateKeyFile, masterKey, conf.GConf.ListenDirectAddr); err != nil {
71+
log.WithError(err).Error("create direct server failed")
72+
return
73+
}
6774

6875
return
6976
}
7077

71-
func createServer(privateKeyPath, pubKeyStorePath string, masterKey []byte, listenAddr string) (server *rpc.Server, err error) {
72-
utils.RemoveAll(pubKeyStorePath + "*")
78+
func createServer(privateKeyPath string, masterKey []byte, listenAddr string) (server *mux.Server, err error) {
79+
server = mux.NewServer()
80+
if err != nil {
81+
return
82+
}
83+
err = server.InitRPCServer(listenAddr, privateKeyPath, masterKey)
84+
return
85+
}
7386

87+
func createDirectServer(privateKeyPath string, masterKey []byte, listenAddr string) (server *rpc.Server, err error) {
88+
if listenAddr == "" {
89+
return nil, nil
90+
}
7491
server = rpc.NewServer()
7592
if err != nil {
7693
return
7794
}
78-
7995
err = server.InitRPCServer(listenAddr, privateKeyPath, masterKey)
80-
8196
return
8297
}

worker/dbms.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,7 +109,7 @@ func NewDBMS(cfg *DBMSConfig) (dbms *DBMS, err error) {
109109
}
110110

111111
// init service
112-
dbms.rpc = NewDBMSRPCService(route.DBRPCName, cfg.Server, dbms)
112+
dbms.rpc = NewDBMSRPCService(route.DBRPCName, cfg.Server, cfg.DirectServer, dbms)
113113
return
114114
}
115115

worker/dbms_config.go

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,8 @@ package worker
1919
import (
2020
"time"
2121

22-
rpc "github.com/CovenantSQL/CovenantSQL/rpc/mux"
22+
"github.com/CovenantSQL/CovenantSQL/rpc"
23+
"github.com/CovenantSQL/CovenantSQL/rpc/mux"
2324
)
2425

2526
var (
@@ -30,7 +31,8 @@ var (
3031
// DBMSConfig defines the local multi-database management system config.
3132
type DBMSConfig struct {
3233
RootDir string
33-
Server *rpc.Server
34+
Server *mux.Server
35+
DirectServer *rpc.Server // optional server to provide DBMS service
3436
MaxReqTimeGap time.Duration
3537
OnCreateDatabase func()
3638
}

worker/dbms_rpc.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,8 @@ import (
2222

2323
"github.com/CovenantSQL/CovenantSQL/proto"
2424
"github.com/CovenantSQL/CovenantSQL/route"
25-
rpc "github.com/CovenantSQL/CovenantSQL/rpc/mux"
25+
"github.com/CovenantSQL/CovenantSQL/rpc"
26+
"github.com/CovenantSQL/CovenantSQL/rpc/mux"
2627
"github.com/CovenantSQL/CovenantSQL/types"
2728
)
2829

@@ -50,11 +51,18 @@ type DBMSRPCService struct {
5051
}
5152

5253
// NewDBMSRPCService returns new dbms rpc service endpoint.
53-
func NewDBMSRPCService(serviceName string, server *rpc.Server, dbms *DBMS) (service *DBMSRPCService) {
54+
func NewDBMSRPCService(
55+
serviceName string, server *mux.Server, direct *rpc.Server, dbms *DBMS,
56+
) (
57+
service *DBMSRPCService,
58+
) {
5459
service = &DBMSRPCService{
5560
dbms: dbms,
5661
}
5762
server.RegisterService(serviceName, service)
63+
if direct != nil {
64+
direct.RegisterService(serviceName, service)
65+
}
5866

5967
dbQuerySuccCounter = metrics.NewMeter()
6068
metrics.Register("db-query-succ", dbQuerySuccCounter)

0 commit comments

Comments
 (0)