2015-11-18 08:50:45 +00:00
|
|
|
package client
|
|
|
|
|
|
|
|
import (
|
|
|
|
"fmt"
|
2015-11-18 12:34:23 +00:00
|
|
|
consul "github.com/hashicorp/consul/api"
|
2015-11-18 08:50:45 +00:00
|
|
|
"github.com/hashicorp/go-multierror"
|
|
|
|
"github.com/hashicorp/nomad/nomad/structs"
|
2015-11-18 10:14:07 +00:00
|
|
|
"log"
|
2015-11-18 12:34:23 +00:00
|
|
|
"time"
|
2015-11-18 08:50:45 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
const (
|
2015-11-18 12:34:23 +00:00
|
|
|
syncInterval = 5 * time.Second
|
2015-11-18 08:50:45 +00:00
|
|
|
)
|
|
|
|
|
2015-11-18 12:34:23 +00:00
|
|
|
type trackedService struct {
|
|
|
|
allocId string
|
|
|
|
task *structs.Task
|
|
|
|
service *structs.Service
|
|
|
|
}
|
|
|
|
|
2015-11-18 08:50:45 +00:00
|
|
|
type ConsulClient struct {
|
2015-11-18 12:34:23 +00:00
|
|
|
client *consul.Client
|
|
|
|
logger *log.Logger
|
|
|
|
shutdownCh chan struct{}
|
2015-11-18 10:14:07 +00:00
|
|
|
|
2015-11-18 12:34:23 +00:00
|
|
|
trackedServices map[string]*trackedService
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
|
|
|
|
2015-11-18 13:15:52 +00:00
|
|
|
func NewConsulClient(logger *log.Logger, consulAddr string) (*ConsulClient, error) {
|
2015-11-18 08:50:45 +00:00
|
|
|
var err error
|
2015-11-18 12:34:23 +00:00
|
|
|
var c *consul.Client
|
2015-11-18 13:15:52 +00:00
|
|
|
cfg := consul.DefaultConfig()
|
|
|
|
cfg.Address = consulAddr
|
|
|
|
if c, err = consul.NewClient(cfg); err != nil {
|
2015-11-18 08:50:45 +00:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
consulClient := ConsulClient{
|
2015-11-18 12:34:23 +00:00
|
|
|
client: c,
|
|
|
|
logger: logger,
|
2015-11-18 12:59:57 +00:00
|
|
|
trackedServices: make(map[string]*trackedService),
|
|
|
|
shutdownCh: make(chan struct{}),
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return &consulClient, nil
|
|
|
|
}
|
|
|
|
|
2015-11-18 09:18:29 +00:00
|
|
|
func (c *ConsulClient) Register(task *structs.Task, allocID string) error {
|
2015-11-18 08:50:45 +00:00
|
|
|
var mErr multierror.Error
|
2015-11-18 10:37:34 +00:00
|
|
|
for _, service := range task.Services {
|
2015-11-18 12:34:23 +00:00
|
|
|
if err := c.registerService(service, task, allocID); err != nil {
|
|
|
|
mErr.Errors = append(mErr.Errors, err)
|
2015-11-18 09:18:29 +00:00
|
|
|
}
|
2015-11-18 12:34:23 +00:00
|
|
|
ts := &trackedService{
|
|
|
|
allocId: allocID,
|
|
|
|
task: task,
|
2015-11-18 12:59:57 +00:00
|
|
|
service: service,
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
2015-11-18 12:34:23 +00:00
|
|
|
c.trackedServices[service.Id] = ts
|
2015-11-18 08:50:45 +00:00
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
return mErr.ErrorOrNil()
|
|
|
|
}
|
|
|
|
|
2015-11-18 10:37:34 +00:00
|
|
|
func (c *ConsulClient) Deregister(task *structs.Task) error {
|
2015-11-18 08:50:45 +00:00
|
|
|
var mErr multierror.Error
|
|
|
|
for _, service := range task.Services {
|
2015-11-18 10:37:34 +00:00
|
|
|
c.logger.Printf("[INFO] De-Registering service %v with Consul", service.Name)
|
2015-11-18 12:34:23 +00:00
|
|
|
if err := c.deregisterService(service.Id); err != nil {
|
2015-11-18 10:37:34 +00:00
|
|
|
c.logger.Printf("[ERROR] Error in de-registering service %v from Consul", service.Name)
|
2015-11-18 08:50:45 +00:00
|
|
|
mErr.Errors = append(mErr.Errors, err)
|
|
|
|
}
|
2015-11-18 12:34:23 +00:00
|
|
|
delete(c.trackedServices, service.Id)
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
|
|
|
return mErr.ErrorOrNil()
|
|
|
|
}
|
2015-11-18 09:18:29 +00:00
|
|
|
|
2015-11-18 12:59:57 +00:00
|
|
|
func (c *ConsulClient) ShutDown() {
|
|
|
|
close(c.shutdownCh)
|
|
|
|
}
|
|
|
|
|
2015-11-18 09:18:29 +00:00
|
|
|
func (c *ConsulClient) findPortAndHostForLabel(portLabel string, task *structs.Task) (string, int) {
|
|
|
|
for _, network := range task.Resources.Networks {
|
|
|
|
if p, ok := network.MapLabelToValues()[portLabel]; ok {
|
2015-11-18 09:20:53 +00:00
|
|
|
return network.IP, p
|
2015-11-18 09:18:29 +00:00
|
|
|
}
|
|
|
|
}
|
2015-11-18 09:20:53 +00:00
|
|
|
return "", 0
|
2015-11-18 09:18:29 +00:00
|
|
|
}
|
2015-11-18 11:08:53 +00:00
|
|
|
|
2015-11-18 12:34:23 +00:00
|
|
|
func (c *ConsulClient) SyncWithConsul() {
|
|
|
|
sync := time.After(syncInterval)
|
|
|
|
agent := c.client.Agent()
|
|
|
|
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-sync:
|
2015-11-18 12:59:57 +00:00
|
|
|
sync = time.After(syncInterval)
|
2015-11-18 12:34:23 +00:00
|
|
|
var consulServices map[string]*consul.AgentService
|
|
|
|
var err error
|
|
|
|
if consulServices, err = agent.Services(); err != nil {
|
|
|
|
c.logger.Printf("[DEBUG] Error while syncing services with Consul: %v", err)
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
for serviceId := range c.trackedServices {
|
|
|
|
if _, ok := consulServices[serviceId]; !ok {
|
|
|
|
ts := c.trackedServices[serviceId]
|
|
|
|
c.registerService(ts.service, ts.task, ts.allocId)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
for serviceId := range consulServices {
|
2015-11-18 12:59:57 +00:00
|
|
|
if serviceId == "consul" {
|
|
|
|
continue
|
|
|
|
}
|
2015-11-18 12:34:23 +00:00
|
|
|
if _, ok := c.trackedServices[serviceId]; !ok {
|
|
|
|
if err := c.deregisterService(serviceId); err != nil {
|
|
|
|
c.logger.Printf("[DEBUG] Error while de-registering service with ID: %s", serviceId)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
case <-c.shutdownCh:
|
2015-11-18 12:59:57 +00:00
|
|
|
c.logger.Printf("[INFO] Shutting down Consul Client")
|
2015-11-18 12:34:23 +00:00
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *ConsulClient) registerService(service *structs.Service, task *structs.Task, allocID string) error {
|
|
|
|
var mErr multierror.Error
|
|
|
|
service.Id = fmt.Sprintf("%s-%s", allocID, task.Name)
|
|
|
|
host, port := c.findPortAndHostForLabel(service.PortLabel, task)
|
|
|
|
if host == "" || port == 0 {
|
|
|
|
return fmt.Errorf("The port:%s marked for registration of service: %s couldn't be found", service.PortLabel, service.Name)
|
|
|
|
}
|
|
|
|
checks := c.makeChecks(service, host, port)
|
|
|
|
asr := &consul.AgentServiceRegistration{
|
|
|
|
ID: service.Id,
|
|
|
|
Name: service.Name,
|
|
|
|
Tags: service.Tags,
|
|
|
|
Port: port,
|
|
|
|
Address: host,
|
|
|
|
Checks: checks,
|
|
|
|
}
|
|
|
|
if err := c.client.Agent().ServiceRegister(asr); err != nil {
|
|
|
|
c.logger.Printf("[ERROR] Error while registering service %v with Consul: %v", service.Name, err)
|
|
|
|
mErr.Errors = append(mErr.Errors, err)
|
|
|
|
}
|
|
|
|
return mErr.ErrorOrNil()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *ConsulClient) deregisterService(serviceId string) error {
|
|
|
|
if err := c.client.Agent().ServiceDeregister(serviceId); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *ConsulClient) makeChecks(service *structs.Service, ip string, port int) []*consul.AgentServiceCheck {
|
|
|
|
var checks []*consul.AgentServiceCheck
|
2015-11-18 11:08:53 +00:00
|
|
|
for _, check := range service.Checks {
|
2015-11-18 12:34:23 +00:00
|
|
|
c := &consul.AgentServiceCheck{
|
2015-11-18 11:08:53 +00:00
|
|
|
Interval: check.Interval.String(),
|
|
|
|
Timeout: check.Timeout.String(),
|
|
|
|
}
|
|
|
|
switch check.Type {
|
|
|
|
case structs.ServiceCheckHTTP:
|
|
|
|
c.HTTP = fmt.Sprintf("%s://%s:%d/%s", check.Protocol, ip, port, check.Http)
|
|
|
|
case structs.ServiceCheckTCP:
|
|
|
|
c.TCP = fmt.Sprintf("%s:%d", ip, port)
|
|
|
|
case structs.ServiceCheckScript:
|
|
|
|
c.Script = check.Script // TODO This needs to include the path of the alloc dir and based on driver types
|
|
|
|
}
|
|
|
|
checks = append(checks, c)
|
|
|
|
}
|
|
|
|
return checks
|
|
|
|
}
|