diff --git a/cmd/auth/main.go b/cmd/auth/main.go index 9553d12..0b7237e 100644 --- a/cmd/auth/main.go +++ b/cmd/auth/main.go @@ -5,7 +5,6 @@ import ( "sonet/internal/auth/data" "sonet/internal/auth/logic" "sonet/pkg/config" - "sonet/pkg/grpc/discovery" "sonet/pkg/utils/shutdown" ) @@ -23,14 +22,11 @@ func main() { panic(err) } - registry := discovery.NewRegister(etcdClient) - shutdown.AddShutdownHook(registry.Stop) - // user dao userDao := data.NewUserDao(config.NewGorm(conf.Gorm)) authServer := logic.NewAuthServer(appConf.AesTokenKey, userDao) go func() { - err = authServer.Run(conf.Grpc, registry) + err = authServer.Run(conf.Grpc, etcdClient) if err != nil { panic(err) } diff --git a/cmd/chat/main.go b/cmd/chat/main.go index b3a0739..a303aa3 100644 --- a/cmd/chat/main.go +++ b/cmd/chat/main.go @@ -3,10 +3,10 @@ package main import ( "context" clientv3 "go.etcd.io/etcd/client/v3" + "go.etcd.io/etcd/client/v3/naming/resolver" "sonet/api/gen/chat" "sonet/internal/chat/logic" "sonet/pkg/config" - "sonet/pkg/grpc/discovery" "sonet/pkg/protocol/deliver" "sonet/pkg/utils/shutdown" ) @@ -20,17 +20,17 @@ func main() { panic(err) } - registry := discovery.NewRegister(etcdClient) - shutdown.AddShutdownHook(registry.Stop) - - etcdResolver := discovery.NewResolver(etcdClient) + etcdResolver, err := resolver.NewBuilder(etcdClient) + if err != nil { + panic(err) + } deli := deliver.NewDeliver(chat.Chat_ServiceDesc.ServiceName) if err := deli.InitWithResolver(context.Background(), etcdResolver); err != nil { panic(err) } chatServer := logic.NewChatServer(deli) go func() { - err = chatServer.Run(conf.Grpc, registry) + err = chatServer.Run(conf.Grpc, etcdClient) if err != nil { panic(err) } diff --git a/cmd/gateway_http/main.go b/cmd/gateway_http/main.go index 8262688..b5025d5 100644 --- a/cmd/gateway_http/main.go +++ b/cmd/gateway_http/main.go @@ -5,15 +5,14 @@ import ( "fmt" "github.com/gin-gonic/gin" clientv3 "go.etcd.io/etcd/client/v3" + "go.etcd.io/etcd/client/v3/naming/resolver" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" - "google.golang.org/grpc/resolver" "net/http" "runtime/debug" config2 "sonet/internal/gateway_http/config" "sonet/internal/gateway_http/logic" "sonet/pkg/config" - "sonet/pkg/grpc/discovery" "sonet/pkg/grpc/generic" "sonet/pkg/utils/logger" "sonet/pkg/utils/resp" @@ -35,9 +34,9 @@ func main() { if err != nil { panic(err) } - etcdResolver := discovery.NewResolver(etcdClient) - shutdown.AddShutdownHook(etcdResolver.Close) - resolver.Register(etcdResolver) + //etcdResolver := discovery.NewResolver(etcdClient) + //shutdown.AddShutdownHook(etcdResolver.Close) + //resolver.Register(etcdResolver) // gin http server authFilter, err := config2.NewAuthFilter(appConf.AesTokenKey, appConf.IgnoreUrls) @@ -61,18 +60,25 @@ func main() { }) server.Use(authFilter.Filter) - // postal loadBalancer handler - postalBalancer := logic.NewPostalBalancer() - postalBalancer.Init(context.Background()) - server.GET("/api/lb/ws", postalBalancer.Endpoint) - // grpc services grpcGroup := server.Group("/api/svc") - grpcFactory := generic.NewGpcGenericClientFactory(etcdResolver, grpc.WithTransportCredentials(insecure.NewCredentials())) + etcdResolver, err := resolver.NewBuilder(etcdClient) + if err != nil { + panic(err) + } + grpcFactory := generic.NewGpcGenericClientFactory("etcd", + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithResolvers(etcdResolver), + ) grpcFactory.Init() grpcGenericHandler := logic.NewGrpcGenericHandler(grpcFactory) grpcGenericHandler.Route(grpcGroup) + // postal loadBalancer handler + postalBalancer := logic.NewPostalBalancer() + postalBalancer.Init(context.Background(), etcdResolver) + server.GET("/api/lb/ws", postalBalancer.Endpoint) + go func() { err := server.Run(fmt.Sprintf(":%d", appConf.Port)) if err != nil { diff --git a/cmd/gateway_ws/main.go b/cmd/gateway_ws/main.go index 8f60400..71437f7 100644 --- a/cmd/gateway_ws/main.go +++ b/cmd/gateway_ws/main.go @@ -4,9 +4,9 @@ import ( "github.com/nats-io/nats.go" "github.com/redis/go-redis/v9" clientv3 "go.etcd.io/etcd/client/v3" + "go.etcd.io/etcd/client/v3/naming/resolver" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" - "google.golang.org/grpc/resolver" "sonet/internal/gateway_ws/server" "sonet/internal/postal/logic" "sonet/pkg/config" @@ -16,6 +16,7 @@ import ( "sonet/pkg/plugins/cache" "sonet/pkg/plugins/mq" "sonet/pkg/utils/conver" + "sonet/pkg/utils/logger" "sonet/pkg/utils/shutdown" "sync" ) @@ -44,12 +45,15 @@ func main() { if err != nil { panic(err) } - clientFactory := client.NewGrpcDirectClientFactory(grpc.WithTransportCredentials(insecure.NewCredentials())) - - etcdResolver := discovery.NewResolver(etcdClient) - resolver.Register(etcdResolver) - grpcFactory := generic.NewGpcGenericClientFactory(etcdResolver, grpc.WithTransportCredentials(insecure.NewCredentials())) + etcdResolver, err := resolver.NewBuilder(etcdClient) + if err != nil { + panic(err) + } + grpcFactory := generic.NewGpcGenericClientFactory("etcd", + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithResolvers(etcdResolver), + ) grpcFactory.Init() postalAddr := discovery.MustRegisterAddress(conf.Grpc.Address) @@ -63,11 +67,15 @@ func main() { }() // run postal server - registry := discovery.NewRegister(etcdClient) - shutdown.AddShutdownHook(registry.Stop) + shutdown.AddShutdownHook(func() { + if err := etcdClient.Close(); err != nil { + logger.Error("etcd close error: ", err) + } + }) + clientFactory := client.NewGrpcDirectClientFactory(grpc.WithTransportCredentials(insecure.NewCredentials())) postalServer := logic.NewPostalServer(appConf.EndpointAddress, sessionStore, subjectStore, clientFactory) go func() { - err = postalServer.Run(conf.Grpc, registry) + err = postalServer.Run(conf.Grpc, etcdClient) if err != nil { panic(err) } diff --git a/cmd/mahjong/main.go b/cmd/mahjong/main.go index 0f2ccb0..1dae610 100644 --- a/cmd/mahjong/main.go +++ b/cmd/mahjong/main.go @@ -3,6 +3,7 @@ package main import ( "context" clientv3 "go.etcd.io/etcd/client/v3" + "go.etcd.io/etcd/client/v3/naming/resolver" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" "sonet/api/gen/auth" @@ -28,7 +29,11 @@ func main() { shutdown.AddShutdownHook(registry.Stop) // init deliver - etcdResolver := discovery.NewResolver(etcdClient) + // etcdResolver := discovery.NewResolver(etcdClient) + etcdResolver, err := resolver.NewBuilder(etcdClient) + if err != nil { + panic(err) + } deli := deliver.NewDeliver(mahjong.Mahjong_ServiceDesc.ServiceName) if err := deli.InitWithResolver(context.Background(), etcdResolver); err != nil { panic(err) @@ -36,14 +41,17 @@ func main() { mjStore := store.NewStore(deli) // auth conn - url := discovery.BuildResolverUrl(auth.Auth_ServiceDesc.ServiceName) - authConn, err := grpc.DialContext(context.Background(), url, grpc.WithTransportCredentials(insecure.NewCredentials())) + url := discovery.EtcdDialUrl(auth.Auth_ServiceDesc.ServiceName) + authConn, err := grpc.DialContext(context.Background(), url, + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithResolvers(etcdResolver), + ) if err != nil { panic(err) } mahjongServer := logic.NewMahjongServer(mjStore, deli, auth.NewAuthClient(authConn)) go func() { - err = mahjongServer.Run(conf.Grpc, registry) + err = mahjongServer.Run(conf.Grpc, etcdClient) if err != nil { panic(err) } diff --git a/deploy_k8s/prometheus.yaml b/deploy_k8s/prometheus.yaml new file mode 100644 index 0000000..828cdd0 --- /dev/null +++ b/deploy_k8s/prometheus.yaml @@ -0,0 +1,121 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: prometheus-map + namespace: sopod + labels: + app: prometheus-map +data: + prometheus.yml: |- + global: + scrape_interval: 15s # Set the scrape interval to every 15 seconds. Default is every 1 minute. + evaluation_interval: 15s # Evaluate rules every 15 seconds. The default is every 1 minute. + # scrape_timeout is set to the global default (10s). + + scrape_configs: + - job_name: "prometheus" + static_configs: + - targets: ["localhost:9090"] + # - job_name: "gateway-http-svc" + # static_configs: + # - targets: ["gateway-http-svc:9100"] + # - job_name: "svr-gateway-ws1-svc" + # static_configs: + # - targets: ["svr-gateway-ws1-svc:9100"] + # - job_name: "svr-gateway-ws2-svc" + # static_configs: + # - targets: ["svr-gateway-ws2-svc:9100"] + +--- + +apiVersion: v1 +kind: PersistentVolume +metadata: + namespace: sopod + name: prometheus-pv + labels: + pv: prometheus-pv +spec: + capacity: + storage: 20Gi + accessModes: + - ReadWriteMany + persistentVolumeReclaimPolicy: Retain + storageClassName: local-storage + local: + path: /data/k8s_data/prometheus/data + nodeAffinity: + required: + nodeSelectorTerms: + - matchExpressions: + - key: node-role.kubernetes.io/master + operator: In + values: + - "true" +--- +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + name: prometheus-pvc + namespace: sopod +spec: + storageClassName: local-storage + accessModes: + - ReadWriteMany + resources: + requests: + storage: 20Gi + selector: + matchLabels: + pv: prometheus-pv + +--- +apiVersion: v1 +kind: Pod +metadata: + name: prometheus + namespace: sopod + labels: + app: prometheus +spec: + containers: + - name: prometheus + securityContext: + runAsUser: 0 # open /opt/bitnami/prometheus/data/queries.active: permission denied + image: bitnami/prometheus:2.47.0 + imagePullPolicy: IfNotPresent + args: + - "--config.file=/opt/bitnami/prometheus/conf/prometheus.yml" + - "--web.listen-address=:9090" + - "--storage.tsdb.path=/opt/bitnami/prometheus/data/" + ports: + - containerPort: 9090 + volumeMounts: + - name: prometheus-data + mountPath: /opt/bitnami/prometheus/data/ + - name: prometheus-map-vol + mountPath: /opt/bitnami/prometheus/conf/ + volumes: + - name: prometheus-data + persistentVolumeClaim: + claimName: prometheus-pvc + - name: prometheus-map-vol + configMap: + name: prometheus-map + +--- +apiVersion: v1 +kind: Service +metadata: + namespace: sopod + name: prometheus-svc + labels: + app: prometheus-svc +spec: + type: NodePort + ports: + - name: prometheus + port: 9090 + nodePort: 30090 + selector: + app: prometheus diff --git a/internal/auth/logic/auth_server.go b/internal/auth/logic/auth_server.go index 54e2455..3d30d87 100644 --- a/internal/auth/logic/auth_server.go +++ b/internal/auth/logic/auth_server.go @@ -5,6 +5,7 @@ import ( "encoding/base64" "errors" "github.com/bytedance/sonic" + clientv3 "go.etcd.io/etcd/client/v3" "google.golang.org/grpc" "google.golang.org/grpc/reflection" "net" @@ -35,7 +36,7 @@ func NewAuthServer(aesTokenKey string, userDao *data.AuthUserDao) *AuthServer { } } -func (s *AuthServer) Run(conf config.GrpcConfig, register *discovery.Register) (err error) { +func (s *AuthServer) Run(conf config.GrpcConfig, etcd *clientv3.Client) (err error) { server := grpc.NewServer( config.GetGrpcOptions( conf, @@ -53,22 +54,20 @@ func (s *AuthServer) Run(conf config.GrpcConfig, register *discovery.Register) ( } // registry discovery - reg := conf.Register - if reg.Name == "" { - reg.Name = auth.Auth_ServiceDesc.ServiceName + register := &(conf.Register) + if register.Name == "" { + register.Name = auth.Auth_ServiceDesc.ServiceName } - if reg.Addr == "" { - reg.Addr, err = discovery.RegisterAddress(conf.Address) - if err != nil { - return - } + if register.Addr == "" { + register.Addr = conf.Address } - if err = register.Register(reg); err != nil { - return + err = discovery.EtcdRegistry(etcd, register) + if err != nil { + panic(err) } // run serve - logger.Infof("%s grpc server running %s\n", reg.Name, listen.Addr().String()) + logger.Infof("%s grpc server running %s\n", register.Name, listen.Addr().String()) err = server.Serve(listen) return } diff --git a/internal/chat/logic/chat_server.go b/internal/chat/logic/chat_server.go index 633a59c..6e23468 100644 --- a/internal/chat/logic/chat_server.go +++ b/internal/chat/logic/chat_server.go @@ -2,6 +2,7 @@ package logic import ( "context" + clientv3 "go.etcd.io/etcd/client/v3" "google.golang.org/grpc" "google.golang.org/grpc/reflection" "google.golang.org/protobuf/types/known/emptypb" @@ -26,7 +27,7 @@ func NewChatServer(deliver *deliver.Deliver) *ChatServer { } } -func (s *ChatServer) Run(conf config.GrpcConfig, register *discovery.Register) (err error) { +func (s *ChatServer) Run(conf config.GrpcConfig, etcd *clientv3.Client) (err error) { server := grpc.NewServer( config.GetGrpcOptions( conf, @@ -44,22 +45,20 @@ func (s *ChatServer) Run(conf config.GrpcConfig, register *discovery.Register) ( } // registry discovery - reg := conf.Register - if reg.Name == "" { - reg.Name = chat.Chat_ServiceDesc.ServiceName + register := &(conf.Register) + if register.Name == "" { + register.Name = chat.Chat_ServiceDesc.ServiceName } - if reg.Addr == "" { - reg.Addr, err = discovery.RegisterAddress(conf.Address) - if err != nil { - return - } + if register.Addr == "" { + register.Addr = conf.Address } - if err = register.Register(reg); err != nil { - return + err = discovery.EtcdRegistry(etcd, register) + if err != nil { + panic(err) } // run serve - logger.Infof("%s grpc server running %s\n", reg.Name, listen.Addr().String()) + logger.Infof("%s grpc server running %s\n", register.Name, listen.Addr().String()) err = server.Serve(listen) return } diff --git a/internal/gateway_http/logic/postal_balancer.go b/internal/gateway_http/logic/postal_balancer.go index 1b9cba8..760bc8d 100644 --- a/internal/gateway_http/logic/postal_balancer.go +++ b/internal/gateway_http/logic/postal_balancer.go @@ -6,6 +6,7 @@ import ( "github.com/gin-gonic/gin" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/resolver" "google.golang.org/protobuf/types/known/emptypb" "net/http" "sonet/api/gen/postal" @@ -23,14 +24,15 @@ func NewPostalBalancer() *PostalBalancer { return &PostalBalancer{} } -func (h *PostalBalancer) Init(ctx context.Context) { +func (h *PostalBalancer) Init(ctx context.Context, resolver resolver.Builder) { balancer.InitConsistentHashBuilder() - postalUrl := discovery.BuildResolverUrl(postal.Postal_ServiceDesc.ServiceName) + postalUrl := discovery.EtcdDialUrl(postal.Postal_ServiceDesc.ServiceName) conn, err := grpc.DialContext(ctx, postalUrl, grpc.WithTransportCredentials(insecure.NewCredentials()), // consistent hash lb grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)), + grpc.WithResolvers(resolver), ) if err != nil { return diff --git a/internal/mahjong/logic/mahjong_server.go b/internal/mahjong/logic/mahjong_server.go index 67859e8..15bf9b5 100644 --- a/internal/mahjong/logic/mahjong_server.go +++ b/internal/mahjong/logic/mahjong_server.go @@ -3,6 +3,7 @@ package logic import ( "context" "errors" + clientv3 "go.etcd.io/etcd/client/v3" "google.golang.org/grpc" "google.golang.org/grpc/reflection" "google.golang.org/protobuf/types/known/emptypb" @@ -37,7 +38,7 @@ func NewMahjongServer(store *store.Store, deliver *deliver.Deliver, authClient a } } -func (mj *MahjongServer) Run(conf config.GrpcConfig, register *discovery.Register) (err error) { +func (mj *MahjongServer) Run(conf config.GrpcConfig, etcd *clientv3.Client) (err error) { server := grpc.NewServer( config.GetGrpcOptions( conf, @@ -55,22 +56,20 @@ func (mj *MahjongServer) Run(conf config.GrpcConfig, register *discovery.Registe } // registry discovery - reg := conf.Register - if reg.Name == "" { - reg.Name = mahjong.Mahjong_ServiceDesc.ServiceName + register := &(conf.Register) + if register.Name == "" { + register.Name = mahjong.Mahjong_ServiceDesc.ServiceName } - if reg.Addr == "" { - reg.Addr, err = discovery.RegisterAddress(conf.Address) - if err != nil { - return - } + if register.Addr == "" { + register.Addr = conf.Address } - if err = register.Register(reg); err != nil { - return + err = discovery.EtcdRegistry(etcd, register) + if err != nil { + panic(err) } // run serve - logger.Infof("%s grpc server running %s\n", reg.Name, listen.Addr().String()) + logger.Infof("%s grpc server running %s\n", register.Name, listen.Addr().String()) err = server.Serve(listen) return } diff --git a/internal/postal/logic/postal_server.go b/internal/postal/logic/postal_server.go index 9982245..4d62253 100644 --- a/internal/postal/logic/postal_server.go +++ b/internal/postal/logic/postal_server.go @@ -2,6 +2,7 @@ package logic import ( "context" + clientv3 "go.etcd.io/etcd/client/v3" "google.golang.org/grpc" "google.golang.org/grpc/reflection" "google.golang.org/protobuf/types/known/emptypb" @@ -40,7 +41,7 @@ func NewPostalServer( } } -func (s *PostalServer) Run(conf config.GrpcConfig, postalRegister *discovery.Register) (err error) { +func (s *PostalServer) Run(conf config.GrpcConfig, etcd *clientv3.Client) (err error) { server := grpc.NewServer( config.GetGrpcOptions( conf, @@ -58,24 +59,23 @@ func (s *PostalServer) Run(conf config.GrpcConfig, postalRegister *discovery.Reg } // registry discovery - reg := conf.Register - if reg.Name == "" { - reg.Name = postal.Postal_ServiceDesc.ServiceName + register := &(conf.Register) + if register.Name == "" { + register.Name = postal.Postal_ServiceDesc.ServiceName } - if reg.Addr == "" { - reg.Addr, err = discovery.RegisterAddress(conf.Address) - if err != nil { - return - } + if register.Addr == "" { + register.Addr = conf.Address } - if err = postalRegister.Register(reg); err != nil { - return + err = discovery.EtcdRegistry(etcd, register) + if err != nil { + panic(err) } + // 其他服务直连地址 - s.broadcastAddress = reg.Addr + s.broadcastAddress = register.Addr // run serve - logger.Infof("%s grpc server running %s\n", reg.Name, listen.Addr().String()) + logger.Infof("%s grpc server running %s\n", register.Name, listen.Addr().String()) err = server.Serve(listen) return } diff --git a/pkg/grpc/balancer/consistent_hash.go b/pkg/grpc/balancer/consistent_hash.go index 478daf9..93eaf02 100644 --- a/pkg/grpc/balancer/consistent_hash.go +++ b/pkg/grpc/balancer/consistent_hash.go @@ -1,12 +1,14 @@ package balancer import ( + "encoding/json" "errors" "fmt" "google.golang.org/grpc/balancer" "google.golang.org/grpc/balancer/base" "google.golang.org/grpc/grpclog" "google.golang.org/grpc/resolver" + "sonet/pkg/utils/logger" "strconv" ) @@ -76,11 +78,28 @@ func wrapAddr(addr string, idx int) string { func GetWeight(addr resolver.Address) (weight int) { weight = DefaultWeight - if addr.Attributes == nil { + var val any = nil + + // from metadata... + if marshal, ok := addr.Metadata.(string); ok { + m := make(map[string]any) + err := json.Unmarshal([]byte(marshal), &m) + if err != nil { + logger.Error("unmarshal metadata error: ", err) + } else { + val = m[WeightKey] + } + } + + // from attributes... + if addr.Attributes != nil { + val = addr.Attributes.Value(WeightKey) + } + + if val == nil { return } - val := addr.Attributes.Value(WeightKey) switch val.(type) { case int: weight = val.(int) @@ -91,8 +110,8 @@ func GetWeight(addr resolver.Address) (weight int) { return } weight = w - default: - grpclog.Errorf("instance weight value type not string: %v\n", val) + //default: + // grpclog.Errorf("instance weight value type not string: %v\n", val) } return } diff --git a/pkg/grpc/discovery/etcd.go b/pkg/grpc/discovery/etcd.go deleted file mode 100644 index ae33e26..0000000 --- a/pkg/grpc/discovery/etcd.go +++ /dev/null @@ -1,33 +0,0 @@ -package discovery - -import ( - "context" - "fmt" - clientv3 "go.etcd.io/etcd/client/v3" - "go.etcd.io/etcd/client/v3/naming/endpoints" -) - -func EtcdDialUrl(serviceName string) string { - return fmt.Sprintf("etcd:///%s", serviceName) -} - -func EtcdRegistry(client *clientv3.Client, serviceName, addr string, attrs map[string]string) (err error) { - em, err := endpoints.NewManager(client, "Hello") - if err != nil { - return - } - ip, port, err := RegisterIpPort(addr) - if err != nil { - return - } - fmt.Println() - - err = em.AddEndpoint(context.Background(), - fmt.Sprintf("%s/%s", serviceName, ip), - endpoints.Endpoint{ - Addr: fmt.Sprintf("%s:%d", ip, port), - Metadata: attrs, - }, - ) - return -} diff --git a/pkg/grpc/discovery/etcd_registry.go b/pkg/grpc/discovery/etcd_registry.go new file mode 100644 index 0000000..71b779f --- /dev/null +++ b/pkg/grpc/discovery/etcd_registry.go @@ -0,0 +1,80 @@ +package discovery + +import ( + "context" + "encoding/json" + "fmt" + clientv3 "go.etcd.io/etcd/client/v3" + "go.etcd.io/etcd/client/v3/naming/endpoints" + "sonet/pkg/utils/logger" + "sonet/pkg/utils/shutdown" + "time" +) + +func EtcdDialUrl(serviceName string) string { + return fmt.Sprintf("etcd:///%s", serviceName) +} + +func EtcdRegistry(client *clientv3.Client, server *Server) (err error) { + em, err := endpoints.NewManager(client, server.Name) + if err != nil { + return + } + ip, port, err := RegisterIpPort(server.Addr) + if err != nil { + return + } + + server.Addr = fmt.Sprintf("%s:%d", ip, port) + // 序列化 metadata 信息 + meta := "{}" + if server.Attrs != nil { + bytes, e := json.Marshal(server.Attrs) + if e != nil { + err = e + return + } + meta = string(bytes) + } + + ctx, cancel := context.WithTimeout(context.Background(), time.Second*10) + defer cancel() + + lease, err := client.Grant(ctx, DefaultRegisterTTL) // 使用租约注册端点确保如果主机无法维持保活心跳, 从服务中删除 + if err != nil { + return + } + endpointKey := fmt.Sprintf("%s/%s", server.Name, ip) + err = em.AddEndpoint(ctx, + endpointKey, + endpoints.Endpoint{ + Addr: server.Addr, + Metadata: meta, + }, + clientv3.WithLease(lease.ID), + ) + + // keepalive lease + keepAliveCh, err := client.KeepAlive(context.Background(), lease.ID) + doneCh := make(chan bool) + c := func() { + doneCh <- true + <-doneCh + } + shutdown.AddShutdownHook(c) + go func() { + for { + select { + case <-doneCh: // async... not done + logger.Info("keepalive done") + ctx, c := context.WithTimeout(context.Background(), time.Second*5) + _, _ = client.Revoke(ctx, lease.ID) + c() + doneCh <- true + return + case _ = <-keepAliveCh: + } + } + }() + return +} diff --git a/pkg/grpc/generic/generic_client_factory.go b/pkg/grpc/generic/generic_client_factory.go index 85f8f53..d6974a2 100644 --- a/pkg/grpc/generic/generic_client_factory.go +++ b/pkg/grpc/generic/generic_client_factory.go @@ -4,31 +4,28 @@ import ( "context" "fmt" "google.golang.org/grpc" - "google.golang.org/grpc/resolver" - "sonet/pkg/grpc/discovery" "sync" ) type GrpcGenericClientFactory struct { - resolver *discovery.Resolver + scheme string defaultOpts []grpc.DialOption clientCache *sync.Map } -func NewGpcGenericClientFactory(resolver *discovery.Resolver, defaultOpts ...grpc.DialOption) *GrpcGenericClientFactory { +func NewGpcGenericClientFactory(scheme string, defaultOpts ...grpc.DialOption) *GrpcGenericClientFactory { return &GrpcGenericClientFactory{ - resolver: resolver, + scheme: scheme, defaultOpts: defaultOpts, } } func (f *GrpcGenericClientFactory) Init() { - resolver.Register(f.resolver) f.clientCache = &sync.Map{} } func (f *GrpcGenericClientFactory) NewClient(ctx context.Context, serviceName string, opts ...grpc.DialOption) (client *GrpcGenericClient, err error) { - addr := fmt.Sprintf("%s:///%s", f.resolver.Scheme(), serviceName) + addr := fmt.Sprintf("%s:///%s", f.scheme, serviceName) dialOpts := make([]grpc.DialOption, 0, len(f.defaultOpts)+len(opts)) dialOpts = append(dialOpts, f.defaultOpts...) dialOpts = append(dialOpts, opts...) diff --git a/pkg/protocol/deliver/deliver.go b/pkg/protocol/deliver/deliver.go index 8888d05..667ff50 100644 --- a/pkg/protocol/deliver/deliver.go +++ b/pkg/protocol/deliver/deliver.go @@ -35,15 +35,15 @@ func NewDeliver(msgInServiceName string) *Deliver { } } -func (d *Deliver) InitWithResolver(ctx context.Context, postalResolver *discovery.Resolver) (err error) { - resolver.Register(postalResolver) +func (d *Deliver) InitWithResolver(ctx context.Context, builder resolver.Builder) (err error) { balancer.InitConsistentHashBuilder() - postalUrl := discovery.BuildResolverUrl(postal.Postal_ServiceDesc.ServiceName) + postalUrl := discovery.EtcdDialUrl(postal.Postal_ServiceDesc.ServiceName) conn, err := grpc.DialContext(ctx, postalUrl, grpc.WithTransportCredentials(insecure.NewCredentials()), // consistent hash lb grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)), + grpc.WithResolvers(builder), ) if err != nil { return