You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
86 lines
2.3 KiB
86 lines
2.3 KiB
package discovery |
|
|
|
import ( |
|
"errors" |
|
"fmt" |
|
"sig-pub/pkg/grpc/discovery/consul" |
|
"sig-pub/pkg/zlog" |
|
"strings" |
|
|
|
"github.com/hashicorp/consul/api" |
|
"google.golang.org/grpc" |
|
"google.golang.org/grpc/health/grpc_health_v1" |
|
"google.golang.org/grpc/resolver" |
|
) |
|
|
|
const ( |
|
ConsulSchema = "consul" |
|
) |
|
|
|
func ConsulDialUrl(svrName string) string { |
|
url := fmt.Sprintf("%s:///%s", ConsulSchema, svrName) |
|
return url |
|
} |
|
|
|
type ConsulDiscovery struct { |
|
client *api.Client |
|
watcher *consul.Watcher |
|
} |
|
|
|
func NewConsulDiscovery(client *api.Client) *ConsulDiscovery { |
|
return &ConsulDiscovery{ |
|
client: client, |
|
} |
|
} |
|
|
|
func (r *ConsulDiscovery) Registry(grpcServer grpc.ServiceRegistrar, register Server) (deregister func(), err error) { |
|
if strings.Trim(register.Name, " ") == "" { |
|
err = errors.New("registry name is empty") |
|
return |
|
} |
|
// 健康检查服务 |
|
healthSrv := NewHealthServer() |
|
grpc_health_v1.RegisterHealthServer(grpcServer, healthSrv) |
|
|
|
// 服务注册 |
|
ip, port, err := RegisterIpPort(register.Addr) |
|
if err != nil { |
|
panic(err) |
|
} |
|
|
|
reg := &api.AgentServiceRegistration{ |
|
ID: fmt.Sprintf("%s-%d", register.Name, register.NodeId), // 唯一 ID |
|
Name: register.Name, // 服务名 |
|
Port: port, |
|
Tags: []string{"v1", "grpc"}, |
|
Address: ip, |
|
// Weights: &api.AgentWeights{}, |
|
Check: &api.AgentServiceCheck{ |
|
CheckID: fmt.Sprintf("%s-%d-check", register.Name, register.NodeId), |
|
GRPC: fmt.Sprintf("%s:%d/grpc.health.v1.Health", ip, port), // gRPC 健康检查 |
|
Interval: "10s", |
|
Timeout: "3s", |
|
DeregisterCriticalServiceAfter: "30s", |
|
}, |
|
} |
|
if err = r.client.Agent().ServiceRegister(reg); err != nil { |
|
zlog.Errorf("registry service %s error: %v", register.Name, err) |
|
return |
|
} |
|
deregister = func() { |
|
err := r.client.Agent().ServiceDeregister(reg.ID) |
|
zlog.Infof("service id %s deregister to consul, err=%v", reg.ID, err) |
|
} |
|
zlog.Infof("service %s registered to consul", register.Name) |
|
return |
|
} |
|
|
|
func (r *ConsulDiscovery) Resolver() (builder resolver.Builder) { |
|
builder = consul.NewConsulBuilder(r.client) |
|
return |
|
} |
|
|
|
func (r *ConsulDiscovery) WatchServices(handler func(service string, passing bool)) (err error) { |
|
r.watcher = consul.NewWatcher(r.client, handler) |
|
return r.watcher.WatchServices() |
|
}
|
|
|