2015-11-18 08:50:45 +00:00
|
|
|
package client
|
|
|
|
|
|
|
|
import (
|
2015-11-25 21:39:16 +00:00
|
|
|
"crypto/tls"
|
2015-11-18 08:50:45 +00:00
|
|
|
"fmt"
|
2015-11-18 10:14:07 +00:00
|
|
|
"log"
|
2015-11-25 21:39:16 +00:00
|
|
|
"net/http"
|
2015-11-18 22:19:58 +00:00
|
|
|
"net/url"
|
2015-11-25 21:39:16 +00:00
|
|
|
"strings"
|
2015-11-18 17:36:37 +00:00
|
|
|
"sync"
|
2015-11-18 12:34:23 +00:00
|
|
|
"time"
|
2015-11-20 06:18:19 +00:00
|
|
|
|
|
|
|
consul "github.com/hashicorp/consul/api"
|
|
|
|
"github.com/hashicorp/go-multierror"
|
|
|
|
"github.com/hashicorp/nomad/nomad/structs"
|
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-26 09:03:16 +00:00
|
|
|
type consulApi interface {
|
|
|
|
CheckRegister(check *consul.AgentCheckRegistration) error
|
|
|
|
CheckDeregister(checkID string) error
|
|
|
|
ServiceRegister(service *consul.AgentServiceRegistration) error
|
|
|
|
ServiceDeregister(ServiceID string) error
|
|
|
|
Services() (map[string]*consul.AgentService, error)
|
|
|
|
Checks() (map[string]*consul.AgentCheck, error)
|
|
|
|
}
|
|
|
|
|
|
|
|
type consulApiClient struct {
|
|
|
|
client *consul.Client
|
|
|
|
}
|
|
|
|
|
|
|
|
func (a *consulApiClient) CheckRegister(check *consul.AgentCheckRegistration) error {
|
|
|
|
return a.client.Agent().CheckRegister(check)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (a *consulApiClient) CheckDeregister(checkID string) error {
|
|
|
|
return a.client.Agent().CheckDeregister(checkID)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (a *consulApiClient) ServiceRegister(service *consul.AgentServiceRegistration) error {
|
|
|
|
return a.client.Agent().ServiceRegister(service)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (a *consulApiClient) ServiceDeregister(serviceId string) error {
|
|
|
|
return a.client.Agent().ServiceDeregister(serviceId)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (a *consulApiClient) Services() (map[string]*consul.AgentService, error) {
|
|
|
|
return a.client.Agent().Services()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (a *consulApiClient) Checks() (map[string]*consul.AgentCheck, error) {
|
|
|
|
return a.client.Agent().Checks()
|
|
|
|
}
|
|
|
|
|
2015-11-25 01:26:30 +00:00
|
|
|
type trackedTask struct {
|
|
|
|
allocID string
|
2015-11-18 12:34:23 +00:00
|
|
|
task *structs.Task
|
|
|
|
}
|
|
|
|
|
2015-11-24 20:34:26 +00:00
|
|
|
type ConsulService struct {
|
2015-11-26 09:03:16 +00:00
|
|
|
client consulApi
|
2015-11-18 12:34:23 +00:00
|
|
|
logger *log.Logger
|
|
|
|
shutdownCh chan struct{}
|
2015-11-18 10:14:07 +00:00
|
|
|
|
2015-11-26 01:28:44 +00:00
|
|
|
trackedTasks map[string]*trackedTask
|
|
|
|
serviceStates map[string]string
|
|
|
|
trackedTskLock sync.Mutex
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
|
|
|
|
2015-11-26 02:31:11 +00:00
|
|
|
// A factory method to create new consul service
|
2015-11-25 21:39:16 +00:00
|
|
|
func NewConsulService(logger *log.Logger, consulAddr string, token string,
|
|
|
|
auth string, enableSSL bool, verifySSL bool) (*ConsulService, 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
|
2015-11-25 21:39:16 +00:00
|
|
|
if token != "" {
|
|
|
|
cfg.Token = token
|
|
|
|
}
|
|
|
|
|
|
|
|
if auth != "" {
|
|
|
|
var username, password string
|
|
|
|
if strings.Contains(auth, ":") {
|
|
|
|
split := strings.SplitN(auth, ":", 2)
|
|
|
|
username = split[0]
|
|
|
|
password = split[1]
|
|
|
|
} else {
|
|
|
|
username = auth
|
|
|
|
}
|
|
|
|
|
|
|
|
cfg.HttpAuth = &consul.HttpBasicAuth{
|
|
|
|
Username: username,
|
|
|
|
Password: password,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
if enableSSL {
|
|
|
|
cfg.Scheme = "https"
|
|
|
|
}
|
|
|
|
if enableSSL && !verifySSL {
|
|
|
|
cfg.HttpClient.Transport = &http.Transport{
|
|
|
|
TLSClientConfig: &tls.Config{
|
|
|
|
InsecureSkipVerify: true,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|
2015-11-18 13:15:52 +00:00
|
|
|
if c, err = consul.NewClient(cfg); err != nil {
|
2015-11-18 08:50:45 +00:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2015-11-24 20:34:26 +00:00
|
|
|
consulService := ConsulService{
|
2015-11-26 09:03:16 +00:00
|
|
|
client: &consulApiClient{client: c},
|
2015-11-26 01:28:44 +00:00
|
|
|
logger: logger,
|
|
|
|
trackedTasks: make(map[string]*trackedTask),
|
|
|
|
serviceStates: make(map[string]string),
|
|
|
|
shutdownCh: make(chan struct{}),
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
|
|
|
|
2015-11-24 20:34:26 +00:00
|
|
|
return &consulService, nil
|
2015-11-18 08:50:45 +00:00
|
|
|
}
|
|
|
|
|
2015-11-26 02:31:11 +00:00
|
|
|
// Starts tracking a task for changes to it's services and tasks
|
2015-11-24 20:34:26 +00:00
|
|
|
func (c *ConsulService) Register(task *structs.Task, allocID string) error {
|
2015-11-18 08:50:45 +00:00
|
|
|
var mErr multierror.Error
|
2015-11-25 01:26:30 +00:00
|
|
|
c.trackedTskLock.Lock()
|
|
|
|
tt := &trackedTask{allocID: allocID, task: task}
|
|
|
|
c.trackedTasks[fmt.Sprintf("%s-%s", allocID, task.Name)] = tt
|
|
|
|
c.trackedTskLock.Unlock()
|
2015-11-18 10:37:34 +00:00
|
|
|
for _, service := range task.Services {
|
2015-11-19 01:33:29 +00:00
|
|
|
c.logger.Printf("[INFO] consul: Registering service %s with Consul.", service.Name)
|
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 08:50:45 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return mErr.ErrorOrNil()
|
|
|
|
}
|
|
|
|
|
2015-11-26 02:31:11 +00:00
|
|
|
// Stops tracking a task for changes to it's services and checks
|
2015-11-25 01:26:30 +00:00
|
|
|
func (c *ConsulService) Deregister(task *structs.Task, allocID string) error {
|
2015-11-18 08:50:45 +00:00
|
|
|
var mErr multierror.Error
|
2015-11-25 01:26:30 +00:00
|
|
|
c.trackedTskLock.Lock()
|
|
|
|
delete(c.trackedTasks, fmt.Sprintf("%s-%s", allocID, task.Name))
|
|
|
|
c.trackedTskLock.Unlock()
|
2015-11-18 08:50:45 +00:00
|
|
|
for _, service := range task.Services {
|
2015-11-19 03:31:29 +00:00
|
|
|
if service.Id == "" {
|
|
|
|
continue
|
|
|
|
}
|
2015-11-19 01:33:29 +00:00
|
|
|
c.logger.Printf("[INFO] consul: 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-24 20:34:26 +00:00
|
|
|
c.logger.Printf("[DEBUG] consul: Error in de-registering service %v from Consul", service.Name)
|
2015-11-18 08:50:45 +00:00
|
|
|
mErr.Errors = append(mErr.Errors, err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return mErr.ErrorOrNil()
|
|
|
|
}
|
2015-11-18 09:18:29 +00:00
|
|
|
|
2015-11-24 20:34:26 +00:00
|
|
|
func (c *ConsulService) ShutDown() {
|
2015-11-18 12:59:57 +00:00
|
|
|
close(c.shutdownCh)
|
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// SyncWithConsul is a long lived function that performs calls to sync
|
|
|
|
// checks and services periodically with Consul Agent
|
2015-11-24 20:34:26 +00:00
|
|
|
func (c *ConsulService) SyncWithConsul() {
|
2015-11-18 12:34:23 +00:00
|
|
|
sync := time.After(syncInterval)
|
|
|
|
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-sync:
|
2015-11-26 09:03:16 +00:00
|
|
|
c.performSync()
|
2015-11-19 02:47:12 +00:00
|
|
|
sync = time.After(syncInterval)
|
2015-11-18 12:34:23 +00:00
|
|
|
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
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// performSync syncs checks and services with Consul and removed tracked
|
|
|
|
// services which are no longer present in tasks
|
2015-11-26 09:03:16 +00:00
|
|
|
func (c *ConsulService) performSync() {
|
2015-11-26 01:28:44 +00:00
|
|
|
// Get the list of the services and that Consul knows about
|
2015-11-26 09:03:16 +00:00
|
|
|
consulServices, _ := c.client.Services()
|
|
|
|
consulChecks, _ := c.client.Checks()
|
2015-11-26 01:28:44 +00:00
|
|
|
delete(consulServices, "consul")
|
2015-11-24 22:37:14 +00:00
|
|
|
|
2015-11-26 01:28:44 +00:00
|
|
|
knownChecks := make(map[string]struct{})
|
2015-11-26 02:23:47 +00:00
|
|
|
knownServices := make(map[string]struct{})
|
|
|
|
|
|
|
|
// Add services and checks which Consul doesn't know about
|
2015-11-25 01:26:30 +00:00
|
|
|
for _, trackedTask := range c.trackedTasks {
|
|
|
|
for _, service := range trackedTask.task.Services {
|
2015-11-26 08:21:25 +00:00
|
|
|
|
|
|
|
// Add new services which Consul agent isn't aware of
|
2015-11-26 02:23:47 +00:00
|
|
|
knownServices[service.Id] = struct{}{}
|
2015-11-26 01:28:44 +00:00
|
|
|
if _, ok := consulServices[service.Id]; !ok {
|
2015-11-25 01:26:30 +00:00
|
|
|
c.registerService(service, trackedTask.task, trackedTask.allocID)
|
2015-11-26 01:28:44 +00:00
|
|
|
continue
|
2015-11-25 01:26:30 +00:00
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// If a service has changed, re-register it with Consul agent
|
2015-11-26 01:28:44 +00:00
|
|
|
if service.Hash() != c.serviceStates[service.Id] {
|
|
|
|
c.registerService(service, trackedTask.task, trackedTask.allocID)
|
|
|
|
continue
|
|
|
|
}
|
2015-11-26 08:21:25 +00:00
|
|
|
|
|
|
|
// Add new checks that Consul isn't aware of
|
2015-11-26 01:28:44 +00:00
|
|
|
for _, check := range service.Checks {
|
2015-11-26 01:44:57 +00:00
|
|
|
knownChecks[check.Id] = struct{}{}
|
|
|
|
if _, ok := consulChecks[check.Id]; !ok {
|
2015-11-26 01:28:44 +00:00
|
|
|
host, port := trackedTask.task.FindHostAndPortFor(service.PortLabel)
|
2015-11-26 01:44:57 +00:00
|
|
|
cr := c.makeCheck(service, check, host, port)
|
2015-11-26 01:28:44 +00:00
|
|
|
c.registerCheck(cr)
|
|
|
|
}
|
2015-11-24 22:37:14 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// Remove services from the service tracker which no longer exists
|
|
|
|
for serviceId := range c.serviceStates {
|
|
|
|
if _, ok := knownServices[serviceId]; !ok {
|
|
|
|
delete(c.serviceStates, serviceId)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-11-26 02:23:47 +00:00
|
|
|
// Remove services that are not present anymore
|
2015-11-26 01:28:44 +00:00
|
|
|
for _, consulService := range consulServices {
|
2015-11-26 02:23:47 +00:00
|
|
|
if _, ok := knownServices[consulService.ID]; !ok {
|
|
|
|
delete(c.serviceStates, consulService.ID)
|
2015-11-26 01:28:44 +00:00
|
|
|
c.deregisterService(consulService.ID)
|
2015-11-25 02:39:38 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-11-26 02:23:47 +00:00
|
|
|
// Remove checks that are not present anymore
|
2015-11-26 01:28:44 +00:00
|
|
|
for _, consulCheck := range consulChecks {
|
|
|
|
if _, ok := knownChecks[consulCheck.CheckID]; !ok {
|
|
|
|
c.deregisterCheck(consulCheck.CheckID)
|
2015-11-25 02:39:38 +00:00
|
|
|
}
|
|
|
|
}
|
2015-11-24 22:37:14 +00:00
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// registerService registers a Service with Consul
|
2015-11-24 20:34:26 +00:00
|
|
|
func (c *ConsulService) registerService(service *structs.Service, task *structs.Task, allocID string) error {
|
2015-11-18 12:34:23 +00:00
|
|
|
var mErr multierror.Error
|
2015-11-24 20:34:26 +00:00
|
|
|
service.Id = fmt.Sprintf("%s-%s", allocID, service.Name)
|
2015-11-25 02:39:38 +00:00
|
|
|
host, port := task.FindHostAndPortFor(service.PortLabel)
|
2015-11-18 12:34:23 +00:00
|
|
|
if host == "" || port == 0 {
|
2015-11-19 01:33:29 +00:00
|
|
|
return fmt.Errorf("consul: The port:%s marked for registration of service: %s couldn't be found", service.PortLabel, service.Name)
|
2015-11-18 12:34:23 +00:00
|
|
|
}
|
2015-11-26 01:28:44 +00:00
|
|
|
c.serviceStates[service.Id] = service.Hash()
|
2015-11-19 02:35:22 +00:00
|
|
|
|
2015-11-25 02:58:53 +00:00
|
|
|
asr := &consul.AgentServiceRegistration{
|
|
|
|
ID: service.Id,
|
|
|
|
Name: service.Name,
|
|
|
|
Tags: service.Tags,
|
|
|
|
Port: port,
|
|
|
|
Address: host,
|
|
|
|
}
|
|
|
|
|
2015-11-26 09:03:16 +00:00
|
|
|
if err := c.client.ServiceRegister(asr); err != nil {
|
2015-11-24 20:34:26 +00:00
|
|
|
c.logger.Printf("[DEBUG] consul: Error while registering service %v with Consul: %v", service.Name, err)
|
2015-11-18 12:34:23 +00:00
|
|
|
mErr.Errors = append(mErr.Errors, err)
|
|
|
|
}
|
2015-11-26 01:28:44 +00:00
|
|
|
for _, check := range service.Checks {
|
2015-11-26 01:44:57 +00:00
|
|
|
cr := c.makeCheck(service, check, host, port)
|
2015-11-26 01:28:44 +00:00
|
|
|
if err := c.registerCheck(cr); err != nil {
|
2015-11-23 07:27:59 +00:00
|
|
|
c.logger.Printf("[ERROR] consul: Error while registerting check %v with Consul: %v", check.Name, err)
|
|
|
|
mErr.Errors = append(mErr.Errors, err)
|
|
|
|
}
|
2015-11-26 01:28:44 +00:00
|
|
|
|
2015-11-23 07:27:59 +00:00
|
|
|
}
|
2015-11-18 12:34:23 +00:00
|
|
|
return mErr.ErrorOrNil()
|
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// registerCheck registers a check with Consul
|
2015-11-25 02:39:38 +00:00
|
|
|
func (c *ConsulService) registerCheck(check *consul.AgentCheckRegistration) error {
|
|
|
|
c.logger.Printf("[DEBUG] Registering Check with ID: %v for Service: %v", check.ID, check.ServiceID)
|
2015-11-26 09:03:16 +00:00
|
|
|
return c.client.CheckRegister(check)
|
2015-11-25 02:39:38 +00:00
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// deregisterCheck de-registers a check with a specific ID from Consul
|
2015-11-25 02:39:38 +00:00
|
|
|
func (c *ConsulService) deregisterCheck(checkID string) error {
|
2015-11-25 02:43:23 +00:00
|
|
|
c.logger.Printf("[DEBUG] Removing check with ID: %v", checkID)
|
2015-11-26 09:03:16 +00:00
|
|
|
return c.client.CheckDeregister(checkID)
|
2015-11-25 02:39:38 +00:00
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// deregisterService de-registers a Service with a specific id from Consul
|
2015-11-24 20:34:26 +00:00
|
|
|
func (c *ConsulService) deregisterService(serviceId string) error {
|
2015-11-26 01:28:44 +00:00
|
|
|
delete(c.serviceStates, serviceId)
|
2015-11-26 09:03:16 +00:00
|
|
|
if err := c.client.ServiceDeregister(serviceId); err != nil {
|
2015-11-18 12:34:23 +00:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2015-11-26 08:21:25 +00:00
|
|
|
// makeCheck creates a Consul Check Registration struct
|
2015-11-26 01:28:44 +00:00
|
|
|
func (c *ConsulService) makeCheck(service *structs.Service, check *structs.ServiceCheck, ip string, port int) *consul.AgentCheckRegistration {
|
|
|
|
if check.Name == "" {
|
2015-11-26 02:31:11 +00:00
|
|
|
check.Name = fmt.Sprintf("service: %q%s%q check", service.Name)
|
2015-11-26 01:28:44 +00:00
|
|
|
}
|
2015-11-26 01:44:57 +00:00
|
|
|
check.Id = check.Hash(service.Id)
|
2015-11-26 02:31:11 +00:00
|
|
|
|
2015-11-26 01:28:44 +00:00
|
|
|
cr := &consul.AgentCheckRegistration{
|
2015-11-26 01:44:57 +00:00
|
|
|
ID: check.Id,
|
2015-11-26 01:28:44 +00:00
|
|
|
Name: check.Name,
|
|
|
|
ServiceID: service.Id,
|
|
|
|
}
|
|
|
|
cr.Interval = check.Interval.String()
|
|
|
|
cr.Timeout = check.Timeout.String()
|
2015-11-26 02:31:11 +00:00
|
|
|
|
2015-11-26 01:28:44 +00:00
|
|
|
switch check.Type {
|
|
|
|
case structs.ServiceCheckHTTP:
|
|
|
|
if check.Protocol == "" {
|
|
|
|
check.Protocol = "http"
|
|
|
|
}
|
|
|
|
url := url.URL{
|
|
|
|
Scheme: check.Protocol,
|
|
|
|
Host: fmt.Sprintf("%s:%d", ip, port),
|
|
|
|
Path: check.Path,
|
|
|
|
}
|
|
|
|
cr.HTTP = url.String()
|
|
|
|
case structs.ServiceCheckTCP:
|
|
|
|
cr.TCP = fmt.Sprintf("%s:%d", ip, port)
|
|
|
|
case structs.ServiceCheckScript:
|
|
|
|
cr.Script = check.Script // TODO This needs to include the path of the alloc dir and based on driver types
|
|
|
|
}
|
|
|
|
return cr
|
2015-11-18 11:08:53 +00:00
|
|
|
}
|