2021-01-26 06:46:54 +00:00
|
|
|
package main
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"log"
|
|
|
|
"os"
|
|
|
|
"os/signal"
|
|
|
|
"syscall"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
dn "github.com/zilliztech/milvus-distributed/internal/datanode"
|
|
|
|
dnc "github.com/zilliztech/milvus-distributed/internal/distributed/datanode"
|
|
|
|
dsc "github.com/zilliztech/milvus-distributed/internal/distributed/dataservice"
|
|
|
|
msc "github.com/zilliztech/milvus-distributed/internal/distributed/masterservice"
|
2021-01-26 09:47:38 +00:00
|
|
|
ms "github.com/zilliztech/milvus-distributed/internal/masterservice"
|
2021-01-26 06:46:54 +00:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/commonpb"
|
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/internalpb2"
|
|
|
|
)
|
|
|
|
|
|
|
|
const retry = 10
|
2021-01-26 08:56:36 +00:00
|
|
|
const interval = 200
|
2021-01-26 06:46:54 +00:00
|
|
|
|
|
|
|
func main() {
|
|
|
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
|
|
|
svr, err := dnc.New(ctx)
|
|
|
|
if err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
log.Println("Datanode is", dn.Params.NodeID)
|
|
|
|
|
|
|
|
// --- Master Service Client ---
|
2021-01-26 09:47:38 +00:00
|
|
|
ms.Params.Init()
|
2021-01-26 06:46:54 +00:00
|
|
|
log.Println("Master service address:", dn.Params.MasterAddress)
|
2021-01-26 11:24:09 +00:00
|
|
|
log.Println("Init master service client ...")
|
2021-01-26 06:46:54 +00:00
|
|
|
masterClient, err := msc.NewGrpcClient(dn.Params.MasterAddress, 20*time.Second)
|
|
|
|
if err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err = masterClient.Init(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err = masterClient.Start(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
var cnt int
|
|
|
|
for cnt = 0; cnt < retry; cnt++ {
|
2021-01-26 08:56:36 +00:00
|
|
|
time.Sleep(time.Duration(cnt*interval) * time.Millisecond)
|
|
|
|
if cnt != 0 {
|
|
|
|
log.Println("Master service isn't ready ...")
|
|
|
|
log.Printf("Retrying getting master service's states in ... %v ms", interval)
|
|
|
|
}
|
|
|
|
|
2021-01-26 06:46:54 +00:00
|
|
|
msStates, err := masterClient.GetComponentStates()
|
2021-01-26 08:56:36 +00:00
|
|
|
|
2021-01-26 06:46:54 +00:00
|
|
|
if err != nil {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
if msStates.Status.ErrorCode != commonpb.ErrorCode_SUCCESS {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
if msStates.State.StateCode != internalpb2.StateCode_HEALTHY {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
break
|
|
|
|
}
|
|
|
|
if cnt >= retry {
|
2021-01-26 08:56:36 +00:00
|
|
|
panic("Master service isn't ready")
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
if err := svr.SetMasterServiceInterface(masterClient); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
// --- Data Service Client ---
|
|
|
|
log.Println("Data service address: ", dn.Params.ServiceAddress)
|
2021-01-26 11:24:09 +00:00
|
|
|
log.Println("Init data service client ...")
|
2021-01-26 06:46:54 +00:00
|
|
|
dataService := dsc.NewClient(dn.Params.ServiceAddress)
|
|
|
|
if err = dataService.Init(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
if err = dataService.Start(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
for cnt = 0; cnt < retry; cnt++ {
|
|
|
|
dsStates, err := dataService.GetComponentStates()
|
|
|
|
if err != nil {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
if dsStates.Status.ErrorCode != commonpb.ErrorCode_SUCCESS {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
if dsStates.State.StateCode != internalpb2.StateCode_INITIALIZING && dsStates.State.StateCode != internalpb2.StateCode_HEALTHY {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
break
|
|
|
|
}
|
|
|
|
if cnt >= retry {
|
2021-01-26 08:56:36 +00:00
|
|
|
panic("Data service isn't ready")
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
if err := svr.SetDataServiceInterface(dataService); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err := svr.Init(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
2021-01-27 03:34:16 +00:00
|
|
|
if err := svr.Start(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
log.Println("Data node successfully started ...")
|
|
|
|
|
2021-01-26 06:46:54 +00:00
|
|
|
sc := make(chan os.Signal, 1)
|
|
|
|
signal.Notify(sc,
|
|
|
|
syscall.SIGHUP,
|
|
|
|
syscall.SIGINT,
|
|
|
|
syscall.SIGTERM,
|
|
|
|
syscall.SIGQUIT)
|
|
|
|
|
|
|
|
var sig os.Signal
|
|
|
|
go func() {
|
|
|
|
sig = <-sc
|
|
|
|
cancel()
|
|
|
|
}()
|
|
|
|
|
|
|
|
<-ctx.Done()
|
|
|
|
log.Println("Got signal to exit signal:", sig.String())
|
|
|
|
|
2021-01-27 03:34:16 +00:00
|
|
|
if err := svr.Stop(); err != nil {
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
2021-01-26 06:46:54 +00:00
|
|
|
switch sig {
|
|
|
|
case syscall.SIGTERM:
|
|
|
|
exit(0)
|
|
|
|
default:
|
|
|
|
exit(1)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func exit(code int) {
|
|
|
|
os.Exit(code)
|
|
|
|
}
|