Browse Source

etcd naming registry

master
tangmingyou 3 years ago
parent
commit
e35b9facab
  1. 6
      cmd/auth/main.go
  2. 12
      cmd/chat/main.go
  3. 28
      cmd/gateway_http/main.go
  4. 26
      cmd/gateway_ws/main.go
  5. 16
      cmd/mahjong/main.go
  6. 121
      deploy_k8s/prometheus.yaml
  7. 23
      internal/auth/logic/auth_server.go
  8. 23
      internal/chat/logic/chat_server.go
  9. 6
      internal/gateway_http/logic/postal_balancer.go
  10. 23
      internal/mahjong/logic/mahjong_server.go
  11. 26
      internal/postal/logic/postal_server.go
  12. 27
      pkg/grpc/balancer/consistent_hash.go
  13. 33
      pkg/grpc/discovery/etcd.go
  14. 80
      pkg/grpc/discovery/etcd_registry.go
  15. 11
      pkg/grpc/generic/generic_client_factory.go
  16. 6
      pkg/protocol/deliver/deliver.go

6
cmd/auth/main.go

@ -5,7 +5,6 @@ import (
"sonet/internal/auth/data" "sonet/internal/auth/data"
"sonet/internal/auth/logic" "sonet/internal/auth/logic"
"sonet/pkg/config" "sonet/pkg/config"
"sonet/pkg/grpc/discovery"
"sonet/pkg/utils/shutdown" "sonet/pkg/utils/shutdown"
) )
@ -23,14 +22,11 @@ func main() {
panic(err) panic(err)
} }
registry := discovery.NewRegister(etcdClient)
shutdown.AddShutdownHook(registry.Stop)
// user dao // user dao
userDao := data.NewUserDao(config.NewGorm(conf.Gorm)) userDao := data.NewUserDao(config.NewGorm(conf.Gorm))
authServer := logic.NewAuthServer(appConf.AesTokenKey, userDao) authServer := logic.NewAuthServer(appConf.AesTokenKey, userDao)
go func() { go func() {
err = authServer.Run(conf.Grpc, registry) err = authServer.Run(conf.Grpc, etcdClient)
if err != nil { if err != nil {
panic(err) panic(err)
} }

12
cmd/chat/main.go

@ -3,10 +3,10 @@ package main
import ( import (
"context" "context"
clientv3 "go.etcd.io/etcd/client/v3" clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/naming/resolver"
"sonet/api/gen/chat" "sonet/api/gen/chat"
"sonet/internal/chat/logic" "sonet/internal/chat/logic"
"sonet/pkg/config" "sonet/pkg/config"
"sonet/pkg/grpc/discovery"
"sonet/pkg/protocol/deliver" "sonet/pkg/protocol/deliver"
"sonet/pkg/utils/shutdown" "sonet/pkg/utils/shutdown"
) )
@ -20,17 +20,17 @@ func main() {
panic(err) panic(err)
} }
registry := discovery.NewRegister(etcdClient) etcdResolver, err := resolver.NewBuilder(etcdClient)
shutdown.AddShutdownHook(registry.Stop) if err != nil {
panic(err)
etcdResolver := discovery.NewResolver(etcdClient) }
deli := deliver.NewDeliver(chat.Chat_ServiceDesc.ServiceName) deli := deliver.NewDeliver(chat.Chat_ServiceDesc.ServiceName)
if err := deli.InitWithResolver(context.Background(), etcdResolver); err != nil { if err := deli.InitWithResolver(context.Background(), etcdResolver); err != nil {
panic(err) panic(err)
} }
chatServer := logic.NewChatServer(deli) chatServer := logic.NewChatServer(deli)
go func() { go func() {
err = chatServer.Run(conf.Grpc, registry) err = chatServer.Run(conf.Grpc, etcdClient)
if err != nil { if err != nil {
panic(err) panic(err)
} }

28
cmd/gateway_http/main.go

@ -5,15 +5,14 @@ import (
"fmt" "fmt"
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
clientv3 "go.etcd.io/etcd/client/v3" clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/naming/resolver"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/resolver"
"net/http" "net/http"
"runtime/debug" "runtime/debug"
config2 "sonet/internal/gateway_http/config" config2 "sonet/internal/gateway_http/config"
"sonet/internal/gateway_http/logic" "sonet/internal/gateway_http/logic"
"sonet/pkg/config" "sonet/pkg/config"
"sonet/pkg/grpc/discovery"
"sonet/pkg/grpc/generic" "sonet/pkg/grpc/generic"
"sonet/pkg/utils/logger" "sonet/pkg/utils/logger"
"sonet/pkg/utils/resp" "sonet/pkg/utils/resp"
@ -35,9 +34,9 @@ func main() {
if err != nil { if err != nil {
panic(err) panic(err)
} }
etcdResolver := discovery.NewResolver(etcdClient) //etcdResolver := discovery.NewResolver(etcdClient)
shutdown.AddShutdownHook(etcdResolver.Close) //shutdown.AddShutdownHook(etcdResolver.Close)
resolver.Register(etcdResolver) //resolver.Register(etcdResolver)
// gin http server // gin http server
authFilter, err := config2.NewAuthFilter(appConf.AesTokenKey, appConf.IgnoreUrls) authFilter, err := config2.NewAuthFilter(appConf.AesTokenKey, appConf.IgnoreUrls)
@ -61,18 +60,25 @@ func main() {
}) })
server.Use(authFilter.Filter) server.Use(authFilter.Filter)
// postal loadBalancer handler
postalBalancer := logic.NewPostalBalancer()
postalBalancer.Init(context.Background())
server.GET("/api/lb/ws", postalBalancer.Endpoint)
// grpc services // grpc services
grpcGroup := server.Group("/api/svc") 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() grpcFactory.Init()
grpcGenericHandler := logic.NewGrpcGenericHandler(grpcFactory) grpcGenericHandler := logic.NewGrpcGenericHandler(grpcFactory)
grpcGenericHandler.Route(grpcGroup) grpcGenericHandler.Route(grpcGroup)
// postal loadBalancer handler
postalBalancer := logic.NewPostalBalancer()
postalBalancer.Init(context.Background(), etcdResolver)
server.GET("/api/lb/ws", postalBalancer.Endpoint)
go func() { go func() {
err := server.Run(fmt.Sprintf(":%d", appConf.Port)) err := server.Run(fmt.Sprintf(":%d", appConf.Port))
if err != nil { if err != nil {

26
cmd/gateway_ws/main.go

@ -4,9 +4,9 @@ import (
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
"github.com/redis/go-redis/v9" "github.com/redis/go-redis/v9"
clientv3 "go.etcd.io/etcd/client/v3" clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/naming/resolver"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/resolver"
"sonet/internal/gateway_ws/server" "sonet/internal/gateway_ws/server"
"sonet/internal/postal/logic" "sonet/internal/postal/logic"
"sonet/pkg/config" "sonet/pkg/config"
@ -16,6 +16,7 @@ import (
"sonet/pkg/plugins/cache" "sonet/pkg/plugins/cache"
"sonet/pkg/plugins/mq" "sonet/pkg/plugins/mq"
"sonet/pkg/utils/conver" "sonet/pkg/utils/conver"
"sonet/pkg/utils/logger"
"sonet/pkg/utils/shutdown" "sonet/pkg/utils/shutdown"
"sync" "sync"
) )
@ -44,12 +45,15 @@ func main() {
if err != nil { if err != nil {
panic(err) 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() grpcFactory.Init()
postalAddr := discovery.MustRegisterAddress(conf.Grpc.Address) postalAddr := discovery.MustRegisterAddress(conf.Grpc.Address)
@ -63,11 +67,15 @@ func main() {
}() }()
// run postal server // run postal server
registry := discovery.NewRegister(etcdClient) shutdown.AddShutdownHook(func() {
shutdown.AddShutdownHook(registry.Stop) 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) postalServer := logic.NewPostalServer(appConf.EndpointAddress, sessionStore, subjectStore, clientFactory)
go func() { go func() {
err = postalServer.Run(conf.Grpc, registry) err = postalServer.Run(conf.Grpc, etcdClient)
if err != nil { if err != nil {
panic(err) panic(err)
} }

16
cmd/mahjong/main.go

@ -3,6 +3,7 @@ package main
import ( import (
"context" "context"
clientv3 "go.etcd.io/etcd/client/v3" clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/naming/resolver"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
"sonet/api/gen/auth" "sonet/api/gen/auth"
@ -28,7 +29,11 @@ func main() {
shutdown.AddShutdownHook(registry.Stop) shutdown.AddShutdownHook(registry.Stop)
// init deliver // 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) deli := deliver.NewDeliver(mahjong.Mahjong_ServiceDesc.ServiceName)
if err := deli.InitWithResolver(context.Background(), etcdResolver); err != nil { if err := deli.InitWithResolver(context.Background(), etcdResolver); err != nil {
panic(err) panic(err)
@ -36,14 +41,17 @@ func main() {
mjStore := store.NewStore(deli) mjStore := store.NewStore(deli)
// auth conn // auth conn
url := discovery.BuildResolverUrl(auth.Auth_ServiceDesc.ServiceName) url := discovery.EtcdDialUrl(auth.Auth_ServiceDesc.ServiceName)
authConn, err := grpc.DialContext(context.Background(), url, grpc.WithTransportCredentials(insecure.NewCredentials())) authConn, err := grpc.DialContext(context.Background(), url,
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithResolvers(etcdResolver),
)
if err != nil { if err != nil {
panic(err) panic(err)
} }
mahjongServer := logic.NewMahjongServer(mjStore, deli, auth.NewAuthClient(authConn)) mahjongServer := logic.NewMahjongServer(mjStore, deli, auth.NewAuthClient(authConn))
go func() { go func() {
err = mahjongServer.Run(conf.Grpc, registry) err = mahjongServer.Run(conf.Grpc, etcdClient)
if err != nil { if err != nil {
panic(err) panic(err)
} }

121
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

23
internal/auth/logic/auth_server.go

@ -5,6 +5,7 @@ import (
"encoding/base64" "encoding/base64"
"errors" "errors"
"github.com/bytedance/sonic" "github.com/bytedance/sonic"
clientv3 "go.etcd.io/etcd/client/v3"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/reflection" "google.golang.org/grpc/reflection"
"net" "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( server := grpc.NewServer(
config.GetGrpcOptions( config.GetGrpcOptions(
conf, conf,
@ -53,22 +54,20 @@ func (s *AuthServer) Run(conf config.GrpcConfig, register *discovery.Register) (
} }
// registry discovery // registry discovery
reg := conf.Register register := &(conf.Register)
if reg.Name == "" { if register.Name == "" {
reg.Name = auth.Auth_ServiceDesc.ServiceName register.Name = auth.Auth_ServiceDesc.ServiceName
} }
if reg.Addr == "" { if register.Addr == "" {
reg.Addr, err = discovery.RegisterAddress(conf.Address) register.Addr = conf.Address
if err != nil {
return
}
} }
if err = register.Register(reg); err != nil { err = discovery.EtcdRegistry(etcd, register)
return if err != nil {
panic(err)
} }
// run serve // 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) err = server.Serve(listen)
return return
} }

23
internal/chat/logic/chat_server.go

@ -2,6 +2,7 @@ package logic
import ( import (
"context" "context"
clientv3 "go.etcd.io/etcd/client/v3"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/reflection" "google.golang.org/grpc/reflection"
"google.golang.org/protobuf/types/known/emptypb" "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( server := grpc.NewServer(
config.GetGrpcOptions( config.GetGrpcOptions(
conf, conf,
@ -44,22 +45,20 @@ func (s *ChatServer) Run(conf config.GrpcConfig, register *discovery.Register) (
} }
// registry discovery // registry discovery
reg := conf.Register register := &(conf.Register)
if reg.Name == "" { if register.Name == "" {
reg.Name = chat.Chat_ServiceDesc.ServiceName register.Name = chat.Chat_ServiceDesc.ServiceName
} }
if reg.Addr == "" { if register.Addr == "" {
reg.Addr, err = discovery.RegisterAddress(conf.Address) register.Addr = conf.Address
if err != nil {
return
}
} }
if err = register.Register(reg); err != nil { err = discovery.EtcdRegistry(etcd, register)
return if err != nil {
panic(err)
} }
// run serve // 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) err = server.Serve(listen)
return return
} }

6
internal/gateway_http/logic/postal_balancer.go

@ -6,6 +6,7 @@ import (
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/resolver"
"google.golang.org/protobuf/types/known/emptypb" "google.golang.org/protobuf/types/known/emptypb"
"net/http" "net/http"
"sonet/api/gen/postal" "sonet/api/gen/postal"
@ -23,14 +24,15 @@ func NewPostalBalancer() *PostalBalancer {
return &PostalBalancer{} return &PostalBalancer{}
} }
func (h *PostalBalancer) Init(ctx context.Context) { func (h *PostalBalancer) Init(ctx context.Context, resolver resolver.Builder) {
balancer.InitConsistentHashBuilder() balancer.InitConsistentHashBuilder()
postalUrl := discovery.BuildResolverUrl(postal.Postal_ServiceDesc.ServiceName) postalUrl := discovery.EtcdDialUrl(postal.Postal_ServiceDesc.ServiceName)
conn, err := grpc.DialContext(ctx, postalUrl, conn, err := grpc.DialContext(ctx, postalUrl,
grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithTransportCredentials(insecure.NewCredentials()),
// consistent hash lb // consistent hash lb
grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)), grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)),
grpc.WithResolvers(resolver),
) )
if err != nil { if err != nil {
return return

23
internal/mahjong/logic/mahjong_server.go

@ -3,6 +3,7 @@ package logic
import ( import (
"context" "context"
"errors" "errors"
clientv3 "go.etcd.io/etcd/client/v3"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/reflection" "google.golang.org/grpc/reflection"
"google.golang.org/protobuf/types/known/emptypb" "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( server := grpc.NewServer(
config.GetGrpcOptions( config.GetGrpcOptions(
conf, conf,
@ -55,22 +56,20 @@ func (mj *MahjongServer) Run(conf config.GrpcConfig, register *discovery.Registe
} }
// registry discovery // registry discovery
reg := conf.Register register := &(conf.Register)
if reg.Name == "" { if register.Name == "" {
reg.Name = mahjong.Mahjong_ServiceDesc.ServiceName register.Name = mahjong.Mahjong_ServiceDesc.ServiceName
} }
if reg.Addr == "" { if register.Addr == "" {
reg.Addr, err = discovery.RegisterAddress(conf.Address) register.Addr = conf.Address
if err != nil {
return
}
} }
if err = register.Register(reg); err != nil { err = discovery.EtcdRegistry(etcd, register)
return if err != nil {
panic(err)
} }
// run serve // 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) err = server.Serve(listen)
return return
} }

26
internal/postal/logic/postal_server.go

@ -2,6 +2,7 @@ package logic
import ( import (
"context" "context"
clientv3 "go.etcd.io/etcd/client/v3"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/reflection" "google.golang.org/grpc/reflection"
"google.golang.org/protobuf/types/known/emptypb" "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( server := grpc.NewServer(
config.GetGrpcOptions( config.GetGrpcOptions(
conf, conf,
@ -58,24 +59,23 @@ func (s *PostalServer) Run(conf config.GrpcConfig, postalRegister *discovery.Reg
} }
// registry discovery // registry discovery
reg := conf.Register register := &(conf.Register)
if reg.Name == "" { if register.Name == "" {
reg.Name = postal.Postal_ServiceDesc.ServiceName register.Name = postal.Postal_ServiceDesc.ServiceName
} }
if reg.Addr == "" { if register.Addr == "" {
reg.Addr, err = discovery.RegisterAddress(conf.Address) register.Addr = conf.Address
if err != nil {
return
}
} }
if err = postalRegister.Register(reg); err != nil { err = discovery.EtcdRegistry(etcd, register)
return if err != nil {
panic(err)
} }
// 其他服务直连地址 // 其他服务直连地址
s.broadcastAddress = reg.Addr s.broadcastAddress = register.Addr
// run serve // 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) err = server.Serve(listen)
return return
} }

27
pkg/grpc/balancer/consistent_hash.go

@ -1,12 +1,14 @@
package balancer package balancer
import ( import (
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"google.golang.org/grpc/balancer" "google.golang.org/grpc/balancer"
"google.golang.org/grpc/balancer/base" "google.golang.org/grpc/balancer/base"
"google.golang.org/grpc/grpclog" "google.golang.org/grpc/grpclog"
"google.golang.org/grpc/resolver" "google.golang.org/grpc/resolver"
"sonet/pkg/utils/logger"
"strconv" "strconv"
) )
@ -76,11 +78,28 @@ func wrapAddr(addr string, idx int) string {
func GetWeight(addr resolver.Address) (weight int) { func GetWeight(addr resolver.Address) (weight int) {
weight = DefaultWeight 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 return
} }
val := addr.Attributes.Value(WeightKey)
switch val.(type) { switch val.(type) {
case int: case int:
weight = val.(int) weight = val.(int)
@ -91,8 +110,8 @@ func GetWeight(addr resolver.Address) (weight int) {
return return
} }
weight = w weight = w
default: //default:
grpclog.Errorf("instance weight value type not string: %v\n", val) // grpclog.Errorf("instance weight value type not string: %v\n", val)
} }
return return
} }

33
pkg/grpc/discovery/etcd.go

@ -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
}

80
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
}

11
pkg/grpc/generic/generic_client_factory.go

@ -4,31 +4,28 @@ import (
"context" "context"
"fmt" "fmt"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/resolver"
"sonet/pkg/grpc/discovery"
"sync" "sync"
) )
type GrpcGenericClientFactory struct { type GrpcGenericClientFactory struct {
resolver *discovery.Resolver scheme string
defaultOpts []grpc.DialOption defaultOpts []grpc.DialOption
clientCache *sync.Map clientCache *sync.Map
} }
func NewGpcGenericClientFactory(resolver *discovery.Resolver, defaultOpts ...grpc.DialOption) *GrpcGenericClientFactory { func NewGpcGenericClientFactory(scheme string, defaultOpts ...grpc.DialOption) *GrpcGenericClientFactory {
return &GrpcGenericClientFactory{ return &GrpcGenericClientFactory{
resolver: resolver, scheme: scheme,
defaultOpts: defaultOpts, defaultOpts: defaultOpts,
} }
} }
func (f *GrpcGenericClientFactory) Init() { func (f *GrpcGenericClientFactory) Init() {
resolver.Register(f.resolver)
f.clientCache = &sync.Map{} f.clientCache = &sync.Map{}
} }
func (f *GrpcGenericClientFactory) NewClient(ctx context.Context, serviceName string, opts ...grpc.DialOption) (client *GrpcGenericClient, err error) { 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 := make([]grpc.DialOption, 0, len(f.defaultOpts)+len(opts))
dialOpts = append(dialOpts, f.defaultOpts...) dialOpts = append(dialOpts, f.defaultOpts...)
dialOpts = append(dialOpts, opts...) dialOpts = append(dialOpts, opts...)

6
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) { func (d *Deliver) InitWithResolver(ctx context.Context, builder resolver.Builder) (err error) {
resolver.Register(postalResolver)
balancer.InitConsistentHashBuilder() balancer.InitConsistentHashBuilder()
postalUrl := discovery.BuildResolverUrl(postal.Postal_ServiceDesc.ServiceName) postalUrl := discovery.EtcdDialUrl(postal.Postal_ServiceDesc.ServiceName)
conn, err := grpc.DialContext(ctx, postalUrl, conn, err := grpc.DialContext(ctx, postalUrl,
grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithTransportCredentials(insecure.NewCredentials()),
// consistent hash lb // consistent hash lb
grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)), grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)),
grpc.WithResolvers(builder),
) )
if err != nil { if err != nil {
return return

Loading…
Cancel
Save