Skip to content
This repository was archived by the owner on Dec 13, 2018. It is now read-only.

Commit 3b3356b

Browse files
author
Feng Honglin
committed
Merge pull request #11 from docker/support_compose_v2
Support compose v2 TUT-691
2 parents 1166382 + 4ce7047 commit 3b3356b

28 files changed

Lines changed: 1644 additions & 862 deletions

README.md

Lines changed: 133 additions & 61 deletions
Large diffs are not rendered by default.

haproxy/config.py

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,17 @@
11
import os
22
import re
33

4-
from parser import parse_extra_bind_settings
4+
5+
def parse_extra_bind_settings(extra_bind_settings):
6+
bind_dict = {}
7+
if extra_bind_settings:
8+
settings = re.split(r'(?<!\\),', extra_bind_settings)
9+
for setting in settings:
10+
term = setting.split(":", 1)
11+
if len(term) == 2:
12+
bind_dict[term[0].strip().replace("\,", ",")] = term[1].strip().replace("\,", ",")
13+
return bind_dict
14+
515

616
# envvar
717
DEFAULT_SSL_CERT = os.getenv("DEFAULT_SSL_CERT") or os.getenv("SSL_CERT")
@@ -28,6 +38,7 @@
2838
HAPROXY_SERVICE_URI = os.getenv("DOCKERCLOUD_SERVICE_API_URI")
2939
API_AUTH = os.getenv("DOCKERCLOUD_AUTH")
3040
DEBUG = os.getenv("DEBUG", False)
41+
LINK_MODE = ""
3142

3243
# const
3344
CERT_DIR = "/certs/"
@@ -37,6 +48,7 @@
3748
API_RETRY = 10 # seconds
3849
PID_FILE = "/tmp/dockercloud-haproxy.pid"
3950

51+
# regular expressions
4052
SERVICE_NAME_MATCH = re.compile(r"(.+)_\d+$")
4153
BACKEND_MATCH = re.compile(r"(?P<proto>tcp|udp):\/\/(?P<addr>[^:]*):(?P<port>.*)")
4254
SERVICE_ALIAS_MATCH = re.compile(r"_PORT_\d{1,5}_(TCP|UDP)$")

haproxy/eventhandler.py

Lines changed: 36 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,12 @@
1-
import logging
21
import json
2+
import logging
3+
4+
import dockercloud
5+
from compose.cli.docker_client import docker_client
6+
from docker.errors import APIError
37

48
import config
9+
import helper.cloud_link_helper
510
from haproxycfg import run_haproxy, Haproxy
611
from utils import get_uuid_from_resource_uri
712

@@ -21,17 +26,17 @@ def on_cloud_event(message):
2126
if event.get("state", "") not in ["In progress", "Pending", "Terminating", "Starting", "Scaling", "Stopping"] and \
2227
event.get("type", "").lower() in ["container", "service"] and \
2328
len(set(Haproxy.cls_linked_services).intersection(set(event.get("parents", [])))) > 0:
24-
msg = "Event: %s %s is %s" % (
29+
msg = "Docker Cloud Event: %s %s is %s" % (
2530
event["type"], get_uuid_from_resource_uri(event.get("resource_uri", "")), event["state"].lower())
2631
run_haproxy(msg)
2732

2833
# Add/remove services linked to haproxy
2934
if event.get("state", "") == "Success" and config.HAPROXY_SERVICE_URI in event.get("parents", []):
30-
run_haproxy("Event: New action is executed on the Haproxy container")
35+
run_haproxy("Docker Cloud Event: New action is executed on the Haproxy container")
3136

3237

3338
def on_websocket_open():
34-
Haproxy.LINKED_CONTAINER_CACHE.clear()
39+
helper.cloud_link_helper.LINKED_CONTAINER_CACHE.clear()
3540
run_haproxy("Websocket open")
3641

3742

@@ -41,3 +46,30 @@ def on_websocket_close():
4146

4247
def on_user_reload(signum, frame):
4348
run_haproxy("User reload")
49+
50+
51+
def listen_dockercloud_events():
52+
events = dockercloud.Events()
53+
events.on_open(on_websocket_open)
54+
events.on_close(on_websocket_close)
55+
events.on_message(on_cloud_event)
56+
events.run_forever()
57+
58+
59+
def listen_docker_events():
60+
try:
61+
docker = docker_client()
62+
docker.ping()
63+
for event in docker.events(decode=True):
64+
logger.debug(event)
65+
attr = event.get("Actor", {}).get("Attributes")
66+
compose_project = attr.get("com.docker.compose.project", "")
67+
compose_service = attr.get("com.docker.compose.service", "")
68+
container_name = attr.get("name", "")
69+
event_action = event.get("Action", "")
70+
service = "%s_%s" % (compose_project, compose_service)
71+
if service in Haproxy.cls_linked_services and event_action in ["start", "die"]:
72+
msg = "Docker event: container %s %s" % (container_name, event_action)
73+
run_haproxy(msg)
74+
except APIError as e:
75+
logger.info("Docker API error: %s" % e)

haproxy/haproxycfg.py

Lines changed: 66 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -2,22 +2,26 @@
22
import logging
33
from collections import OrderedDict
44

5+
from compose.cli.docker_client import docker_client
6+
7+
import config
58
import helper.backend_helper as BackendHelper
9+
import helper.cloud_link_helper as CloudLinkHelper
610
import helper.config_helper as ConfigHelper
711
import helper.frontend_helper as FrontendHelper
8-
import helper.init_helper as InitHelper
12+
import helper.new_link_helper as NewLinkHelper
913
import helper.ssl_helper as SslHelper
1014
import helper.tcp_helper as TcpHelper
1115
import helper.update_helper as UpdateHelper
1216
from haproxy.config import *
13-
from parser import Specs
17+
from haproxy.parser import LegacyLinkSpecs, NewLinkSpecs
1418
from utils import fetch_remote_obj, prettify, save_to_file, get_service_attribute, get_bind_string
1519

1620
logger = logging.getLogger("haproxy")
1721

1822

1923
def run_haproxy(msg=None):
20-
haproxy = Haproxy(msg)
24+
haproxy = Haproxy(config.LINK_MODE, msg)
2125
haproxy.update()
2226

2327

@@ -27,62 +31,78 @@ class Haproxy(object):
2731
cls_process = None
2832
cls_certs = []
2933

30-
cls_service_name_match = re.compile(r"(.+)_\d+$")
31-
32-
LINKED_CONTAINER_CACHE = {}
33-
34-
def __init__(self, msg=""):
34+
def __init__(self, link_mode="", msg=""):
3535
logger.info("==========BEGIN==========")
3636
if msg:
3737
logger.info(msg)
3838

39+
self.link_mode = link_mode
3940
self.ssl_bind_string = None
4041
self.ssl_updated = False
4142
self.routes_added = []
4243
self.require_default_route = False
44+
self.specs = None
4345

44-
self._initialize()
45-
46-
def _initialize(self):
47-
if HAPROXY_CONTAINER_URI and HAPROXY_SERVICE_URI and API_AUTH:
48-
haproxy_container = fetch_remote_obj(HAPROXY_CONTAINER_URI)
46+
self.specs = self._initialize(self.link_mode)
4947

50-
haproxy_links = InitHelper.get_links_from_haproxy(haproxy_container.linked_to_container)
51-
new_added_container_uris = InitHelper.get_new_added_link_uri(Haproxy.LINKED_CONTAINER_CACHE, haproxy_links)
52-
new_added_containers = InitHelper.get_container_object_from_uri(new_added_container_uris)
53-
InitHelper.update_container_cache(Haproxy.LINKED_CONTAINER_CACHE, new_added_container_uris,
54-
new_added_containers)
55-
linked_containers = InitHelper.get_linked_containers(Haproxy.LINKED_CONTAINER_CACHE,
56-
haproxy_container.linked_to_container)
57-
InitHelper.update_haproxy_links(haproxy_links, linked_containers)
58-
59-
logger.info("Service links: %s", ", ".join(InitHelper.get_service_links_str(haproxy_links)))
60-
logger.info("Container links: %s", ", ".join(InitHelper.get_container_links_str(haproxy_links)))
61-
62-
Haproxy.cls_linked_services = InitHelper.get_linked_services(haproxy_links)
63-
self.specs = Specs(haproxy_links)
48+
@staticmethod
49+
def _initialize(link_mode):
50+
if link_mode == "cloud":
51+
links = Haproxy._init_cloud_links()
52+
specs = NewLinkSpecs(links)
53+
elif link_mode == "new":
54+
links = Haproxy._init_new_links()
55+
if links is None:
56+
specs = LegacyLinkSpecs()
57+
else:
58+
specs = NewLinkSpecs(links)
6459
else:
65-
logger.info("Loading HAProxy definition from environment variables")
66-
Haproxy.cls_linked_services = None
67-
Haproxy.specs = Specs()
60+
specs = LegacyLinkSpecs()
61+
return specs
6862

69-
def update(self):
70-
self._config_ssl()
63+
@staticmethod
64+
def _init_cloud_links():
65+
haproxy_container = fetch_remote_obj(HAPROXY_CONTAINER_URI)
66+
links = CloudLinkHelper.get_cloud_links(haproxy_container)
67+
Haproxy.cls_linked_services = CloudLinkHelper.get_linked_services(links)
68+
logger.info("Linked service: %s", ", ".join(CloudLinkHelper.get_service_links_str(links)))
69+
logger.info("Linked container: %s", ", ".join(CloudLinkHelper.get_container_links_str(links)))
70+
return links
7171

72-
cfg_dict = OrderedDict()
73-
cfg_dict.update(self._config_global_section())
74-
cfg_dict.update(self._config_defaults_section())
75-
cfg_dict.update(self._config_stats_section())
76-
cfg_dict.update(self._config_userlist_section(HTTP_BASIC_AUTH))
77-
cfg_dict.update(self._config_tcp_sections())
78-
cfg_dict.update(self._config_frontend_sections())
79-
cfg_dict.update(self._config_backend_sections())
72+
@staticmethod
73+
def _init_new_links():
74+
try:
75+
docker = docker_client()
76+
docker.ping()
77+
container_id = os.environ.get("HOSTNAME", "")
78+
haproxy_container = docker.inspect_container(container_id)
79+
except Exception as e:
80+
logger.info("Docker API error, regressing to legacy links mode: ", e)
81+
return None
82+
links, Haproxy.cls_linked_services = NewLinkHelper.get_new_links(docker, haproxy_container)
83+
logger.info("Linked service: %s", ", ".join(NewLinkHelper.get_service_links_str(links)))
84+
logger.info("Linked container: %s", ", ".join(NewLinkHelper.get_container_links_str(links)))
85+
return links
8086

81-
cfg = prettify(cfg_dict)
82-
self._update_haproxy(cfg)
87+
def update(self):
88+
if self.specs:
89+
self._config_ssl()
90+
cfg_dict = OrderedDict()
91+
cfg_dict.update(self._config_global_section())
92+
cfg_dict.update(self._config_defaults_section())
93+
cfg_dict.update(self._config_stats_section())
94+
cfg_dict.update(self._config_userlist_section(HTTP_BASIC_AUTH))
95+
cfg_dict.update(self._config_tcp_sections())
96+
cfg_dict.update(self._config_frontend_sections())
97+
cfg_dict.update(self._config_backend_sections())
98+
99+
cfg = prettify(cfg_dict)
100+
self._update_haproxy(cfg)
101+
else:
102+
logger.info("Internal error: Specs is not initialized")
83103

84104
def _update_haproxy(self, cfg):
85-
if HAPROXY_SERVICE_URI and HAPROXY_CONTAINER_URI and API_AUTH:
105+
if self.link_mode in ["cloud", "new"]:
86106
if Haproxy.cls_cfg != cfg:
87107
logger.info("HAProxy configuration:\n%s" % cfg)
88108
Haproxy.cls_cfg = cfg
@@ -94,10 +114,10 @@ def _update_haproxy(self, cfg):
94114
else:
95115
logger.info("HAProxy configuration remains unchanged")
96116
logger.info("===========END===========")
97-
else:
117+
elif self.link_mode in ["legacy"]:
98118
logger.info("HAProxy configuration:\n%s" % cfg)
99-
save_to_file(HAPROXY_CONFIG_FILE, cfg)
100-
UpdateHelper.run_once()
119+
if save_to_file(HAPROXY_CONFIG_FILE, cfg):
120+
UpdateHelper.run_once()
101121

102122
def _config_ssl(self):
103123
ssl_bind_string = ""
Lines changed: 23 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,20 @@
33
from haproxy.config import SERVICE_NAME_MATCH
44
from haproxy.utils import get_uuid_from_resource_uri, fetch_remote_obj
55

6+
LINKED_CONTAINER_CACHE = {}
67

7-
def get_links_from_haproxy(container_links):
8+
9+
def get_cloud_links(haproxy_container):
10+
links = _init_links(haproxy_container.linked_to_container)
11+
new_added_container_uris = _get_new_added_link_uri(LINKED_CONTAINER_CACHE, links)
12+
new_added_containers = _get_container_object_from_uri(new_added_container_uris)
13+
_update_container_cache(LINKED_CONTAINER_CACHE, new_added_container_uris, new_added_containers)
14+
linked_containers = _get_linked_containers(LINKED_CONTAINER_CACHE, haproxy_container.linked_to_container)
15+
_update_links(links, linked_containers)
16+
return links
17+
18+
19+
def _init_links(container_links):
820
links = {}
921
for link in container_links:
1022
linked_container_uri = link["to_container"]
@@ -24,37 +36,30 @@ def get_links_from_haproxy(container_links):
2436
return links
2537

2638

27-
def get_new_added_link_uri(container_object_cache, links):
39+
def _get_new_added_link_uri(container_object_cache, links):
2840
return filter(lambda x: x not in container_object_cache, links)
2941

3042

31-
def get_container_object_from_uri(container_uris):
43+
def _get_container_object_from_uri(container_uris):
3244
pool = ThreadPool(processes=10)
3345
container_objects = pool.map(fetch_remote_obj, container_uris)
3446
return container_objects
3547

3648

37-
def update_container_cache(cache, new_container_object_uris, new_container_objects):
49+
def _update_container_cache(cache, new_container_object_uris, new_container_objects):
3850
for i, container_uri in enumerate(new_container_object_uris):
3951
cache[container_uri] = new_container_objects[i]
4052

4153

42-
def update_haproxy_links(links, linked_containers):
54+
def _update_links(links, linked_containers):
4355
for linked_container in linked_containers:
4456
linked_container_uri = linked_container.resource_uri
4557
linked_container_service_uri = linked_container.service
46-
linked_container_name = links[linked_container_uri]["container_name"]
47-
48-
linked_container_envvars = {}
49-
for envvar in linked_container.container_envvars:
50-
if "_ENV_" not in envvar['key']:
51-
linked_container_envvars["%s_ENV_%s" % (linked_container_name, envvar['key'])] = envvar['value']
52-
5358
links[linked_container_uri]["service_uri"] = linked_container_service_uri
54-
links[linked_container_uri]["container_envvars"] = linked_container_envvars
59+
links[linked_container_uri]["container_envvars"] = linked_container.container_envvars
5560

5661

57-
def get_linked_containers(cache, container_links):
62+
def _get_linked_containers(cache, container_links):
5863
linked_containers = [cache[link["to_container"]] for link in container_links]
5964
return linked_containers
6065

@@ -68,10 +73,11 @@ def get_linked_services(haproxy_links):
6873

6974

7075
def get_service_links_str(haproxy_links):
71-
return sorted(set(["%s(%s)" % (link.get("service_name"), get_uuid_from_resource_uri(link.get("service_uri")))
76+
return sorted(set(["%s(%s)" % (link.get("service_name"), get_uuid_from_resource_uri(link.get("service_uri", "")))
7277
for link in haproxy_links.itervalues()]))
7378

7479

7580
def get_container_links_str(haproxy_links):
76-
return sorted(set(["%s(%s)" % (link.get("container_name"), get_uuid_from_resource_uri(link.get("container_uri")))
77-
for link in haproxy_links.itervalues()]))
81+
return sorted(
82+
set(["%s(%s)" % (link.get("container_name"), get_uuid_from_resource_uri(link.get("container_uri", "")))
83+
for link in haproxy_links.itervalues()]))

0 commit comments

Comments
 (0)