micro/registry/consul_registry.go

222 lines
4.2 KiB
Go
Raw Normal View History

2015-01-13 23:31:27 +00:00
package registry
import (
"encoding/json"
2015-01-13 23:31:27 +00:00
"errors"
"fmt"
"net"
2015-02-14 23:00:47 +00:00
"sync"
2015-01-13 23:31:27 +00:00
consul "github.com/hashicorp/consul/api"
2015-01-13 23:31:27 +00:00
)
type consulRegistry struct {
2015-02-14 23:00:47 +00:00
Address string
Client *consul.Client
2015-01-13 23:31:27 +00:00
2015-02-14 23:00:47 +00:00
mtx sync.RWMutex
2015-11-08 01:48:48 +00:00
services map[string][]*Service
}
2015-10-11 12:05:20 +01:00
func encodeEndpoints(en []*Endpoint) []string {
var tags []string
for _, e := range en {
if b, err := json.Marshal(e); err == nil {
tags = append(tags, "e="+string(b))
}
}
return tags
}
func decodeEndpoints(tags []string) []*Endpoint {
var en []*Endpoint
for _, tag := range tags {
if len(tag) == 0 || tag[0] != 'e' {
continue
}
var e *Endpoint
if err := json.Unmarshal([]byte(tag[2:]), &e); err == nil {
en = append(en, e)
}
}
return en
}
2015-05-26 22:39:48 +01:00
func encodeMetadata(md map[string]string) []string {
var tags []string
for k, v := range md {
if b, err := json.Marshal(map[string]string{
k: v,
}); err == nil {
2015-10-11 12:05:20 +01:00
tags = append(tags, "t="+string(b))
}
}
return tags
}
2015-05-26 22:39:48 +01:00
func decodeMetadata(tags []string) map[string]string {
md := make(map[string]string)
for _, tag := range tags {
2015-10-11 12:05:20 +01:00
if len(tag) == 0 || tag[0] != 't' {
continue
}
var kv map[string]string
2015-10-11 12:05:20 +01:00
if err := json.Unmarshal([]byte(tag[2:]), &kv); err == nil {
for k, v := range kv {
md[k] = v
}
}
}
return md
2015-02-14 23:00:47 +00:00
}
2015-01-13 23:31:27 +00:00
func newConsulRegistry(addrs []string, opts ...Option) Registry {
config := consul.DefaultConfig()
if len(addrs) > 0 {
addr, port, err := net.SplitHostPort(addrs[0])
if ae, ok := err.(*net.AddrError); ok && ae.Err == "missing port in address" {
port = "8500"
config.Address = fmt.Sprintf("%s:%s", addr, port)
} else if err == nil {
config.Address = fmt.Sprintf("%s:%s", addr, port)
}
}
2015-10-26 20:23:57 +00:00
client, _ := consul.NewClient(config)
cr := &consulRegistry{
Address: config.Address,
Client: client,
2015-11-08 01:48:48 +00:00
services: make(map[string][]*Service),
}
return cr
}
func (c *consulRegistry) Deregister(s *Service) error {
if len(s.Nodes) == 0 {
2015-01-13 23:31:27 +00:00
return errors.New("Require at least one node")
}
node := s.Nodes[0]
2015-01-13 23:31:27 +00:00
_, err := c.Client.Catalog().Deregister(&consul.CatalogDeregistration{
Node: node.Id,
Address: node.Address,
2015-01-13 23:31:27 +00:00
}, nil)
return err
}
func (c *consulRegistry) Register(s *Service) error {
if len(s.Nodes) == 0 {
2015-01-13 23:31:27 +00:00
return errors.New("Require at least one node")
}
node := s.Nodes[0]
2015-05-26 22:39:48 +01:00
tags := encodeMetadata(node.Metadata)
2015-10-11 12:05:20 +01:00
tags = append(tags, encodeEndpoints(s.Endpoints)...)
2015-01-13 23:31:27 +00:00
_, err := c.Client.Catalog().Register(&consul.CatalogRegistration{
Node: node.Id,
Address: node.Address,
2015-01-13 23:31:27 +00:00
Service: &consul.AgentService{
2015-11-08 01:48:48 +00:00
ID: s.Version,
Service: s.Name,
Port: node.Port,
Tags: tags,
2015-01-13 23:31:27 +00:00
},
}, nil)
return err
}
2015-11-08 01:48:48 +00:00
func (c *consulRegistry) GetService(name string) ([]*Service, error) {
2015-02-14 23:00:47 +00:00
c.mtx.RLock()
service, ok := c.services[name]
c.mtx.RUnlock()
if ok {
return service, nil
}
2015-01-13 23:31:27 +00:00
rsp, _, err := c.Client.Catalog().Service(name, "", nil)
if err != nil {
return nil, err
}
2015-11-08 01:48:48 +00:00
serviceMap := map[string]*Service{}
2015-01-13 23:31:27 +00:00
for _, s := range rsp {
if s.ServiceName != name {
continue
}
2015-11-08 01:48:48 +00:00
id := s.Node
key := s.ServiceID
version := s.ServiceID
// We're adding service version but
// don't want to break backwards compatibility
if id == version {
key = "default"
version = ""
}
svc, ok := serviceMap[key]
if !ok {
svc = &Service{
Endpoints: decodeEndpoints(s.ServiceTags),
Name: s.ServiceName,
Version: version,
}
serviceMap[key] = svc
}
svc.Nodes = append(svc.Nodes, &Node{
Id: id,
Address: s.Address,
Port: s.ServicePort,
2015-05-26 22:39:48 +01:00
Metadata: decodeMetadata(s.ServiceTags),
2015-01-13 23:31:27 +00:00
})
}
2015-11-08 01:48:48 +00:00
var services []*Service
for _, service := range serviceMap {
services = append(services, service)
}
return services, nil
2015-01-13 23:31:27 +00:00
}
func (c *consulRegistry) ListServices() ([]*Service, error) {
c.mtx.RLock()
serviceMap := c.services
c.mtx.RUnlock()
var services []*Service
if len(serviceMap) > 0 {
2015-11-08 01:48:48 +00:00
for _, service := range serviceMap {
services = append(services, service...)
}
return services, nil
}
rsp, _, err := c.Client.Catalog().Services(&consul.QueryOptions{})
if err != nil {
return nil, err
}
for service, _ := range rsp {
services = append(services, &Service{Name: service})
}
return services, nil
}
func (c *consulRegistry) Watch() (Watcher, error) {
return newConsulWatcher(c)
2015-01-13 23:31:27 +00:00
}