monitor-hue-lb 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294
  1. #!/usr/bin/env python
  2. # Copyright 2014 Cloudera, Inc.
  3. #
  4. # Licensed under the Apache License, Version 2.0 (the "License");
  5. # you may not use this file except in compliance with the License.
  6. # You may obtain a copy of the License at
  7. #
  8. # http://www.apache.org/licenses/LICENSE-2.0
  9. #
  10. # Unless required by applicable law or agreed to in writing, software
  11. # distributed under the License is distributed on an "AS IS" BASIS,
  12. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  13. # See the License for the specific language governing permissions and
  14. # limitations under the License.
  15. import getpass
  16. import optparse
  17. import os
  18. import signal
  19. import sys
  20. import toml
  21. from cm_api.api_client import ApiResource
  22. from supervisor import childutils
  23. # ------------------------------------------------------------------------------
  24. class HueServer(object):
  25. def __init__(self, address, state):
  26. self.address = address
  27. self.state = state
  28. def __hash__(self):
  29. return hash((self.address, self.state))
  30. def __str__(self):
  31. return '%s (%s)' % (self.address, self.state)
  32. def __eq__(self, other):
  33. return isinstance(other, HueServer) and \
  34. self.address == other.address and \
  35. self.state == other.state
  36. def __cmp__(self, other):
  37. return cmp(
  38. (self.address, self.state),
  39. (other.address, other.state))
  40. # ------------------------------------------------------------------------------
  41. class MonitorHue(object):
  42. def __init__(
  43. self,
  44. cm_client,
  45. stdin=sys.stdin,
  46. stdout=sys.stdout,
  47. stderr=sys.stderr):
  48. self.cm_client = cm_client
  49. self.stdin = stdin
  50. self.stdout = stdout
  51. self.stderr = stderr
  52. self.hue_servers = None
  53. self.listeners = []
  54. def run_forever(self):
  55. self.tick()
  56. while True:
  57. headers, payload = childutils.listener.wait(self.stdin, self.stdout)
  58. # ignore non-tick events.
  59. if not headers['eventname'].startswith('TICK'):
  60. childutils.listener.ok(self.stdout)
  61. continue
  62. try:
  63. self.tick()
  64. finally:
  65. childutils.listener.ok(self.stdout)
  66. def tick(self):
  67. """Update the load balancer if any Hue servers have been added or removed"""
  68. hue_servers = self.get_hue_servers()
  69. # Don't do anything if we've already processed this set of hue servers.
  70. if self.hue_servers == hue_servers:
  71. return
  72. self.hue_servers = hue_servers
  73. print >> sys.stderr, 'updating server list:'
  74. for server in sorted(hue_servers):
  75. print >> sys.stderr, ' %s' % server
  76. for listener in self.listeners:
  77. listener.update(hue_servers)
  78. def get_hue_servers(self):
  79. """Fetch all the known hue servers"""
  80. hue_servers = set()
  81. for cluster in self.cm_client.get_all_clusters():
  82. for service in cluster.get_all_services():
  83. if service.type == 'HUE':
  84. for role in service.get_all_roles():
  85. host = self.cm_client.get_host(role.hostRef.hostId)
  86. hostname = host.hostname
  87. config = role.get_config(view='full')
  88. port = config['hue_http_port']
  89. port = (port.value is None and port.default) or port.value
  90. address = '%s:%s' % (hostname, port)
  91. hue_servers.add(HueServer(address, role.roleState))
  92. return hue_servers
  93. # ------------------------------------------------------------------------------
  94. class ConfigListener(object):
  95. def __init__(self, rpc, config, process_names):
  96. self.rpc = rpc
  97. self.config = config
  98. self.process_names = process_names
  99. self.config_file = self.config['config_file']
  100. with open(self.config['config_template']) as f:
  101. self.config_template = f.read()
  102. with open(self.config['server_template']) as f:
  103. self.server_template = f.read()
  104. def update(self, hue_servers):
  105. processes = self.rpc.supervisor.getAllProcessInfo()
  106. config = self.expand_config_template(hue_servers)
  107. with open(self.config_file, 'w') as f:
  108. print >> f, config
  109. for process in processes:
  110. # Ignore down processes.
  111. if not process['pid']:
  112. continue
  113. # Ignore processes we aren't monitoring.
  114. if process['name'] not in self.process_names:
  115. continue
  116. self.reload_process_config(process)
  117. def expand_config_template(self, hue_servers):
  118. config_servers = []
  119. for index, server in enumerate(hue_servers):
  120. config_servers.append(self.expand_server_template(index, server))
  121. return self.config_template % {
  122. 'servers': '\n'.join(config_servers)
  123. }
  124. def expand_server_template(self, index, server):
  125. raise NotImplementedError
  126. def reload_process(self, process):
  127. raise NotImplementedError
  128. # ------------------------------------------------------------------------------
  129. class HAProxyListener(ConfigListener):
  130. def expand_server_template(self, index, server):
  131. """Expand the HAProxy server line"""
  132. return self.server_template % {
  133. 'index': index,
  134. 'address': server.address,
  135. }
  136. def reload_process_config(self, process):
  137. """Let haproxy-wrapper know the config file has changed."""
  138. os.kill(process['pid'], signal.SIGUSR2)
  139. # ------------------------------------------------------------------------------
  140. class NginxListener(ConfigListener):
  141. def expand_server_template(self, index, server):
  142. """Expand the Nginx server line"""
  143. config_server = self.server_template % {
  144. 'address': server.address,
  145. }
  146. # Let nginx know if this backend is offline.
  147. if server.state == 'STOPPED':
  148. config_server += ' down'
  149. return config_server + ';'
  150. def reload_process_config(self, process):
  151. """Let nginx know the config file has changed."""
  152. os.kill(process['pid'], signal.SIGHUP)
  153. # ------------------------------------------------------------------------------
  154. def main(argv):
  155. parser = optparse.OptionParser()
  156. parser.add_option('-c', '--config_file',
  157. help='config file')
  158. parser.add_option('--haproxy',
  159. default=[],
  160. type='str',
  161. action='append',
  162. help='monitor this haproxy process')
  163. parser.add_option('--nginx',
  164. default=[],
  165. type='str',
  166. action='append',
  167. help='monitor this haproxy process')
  168. options, args = parser.parse_args(argv[1:])
  169. if options.config_file is None:
  170. options.config_file = 'etc/hue-lb.toml'
  171. try:
  172. with open(options.config_file) as f:
  173. config = toml.load(f)
  174. except IOError, e:
  175. print >> sys.stderr, 'error opening %s: %s' % (options.config_file, e)
  176. return 1
  177. try:
  178. hostname = config['cloudera-manager']['hostname']
  179. except KeyError:
  180. hostname = 'localhost'
  181. try:
  182. username = config['cloudera-manager']['username']
  183. except KeyError:
  184. username = 'admin'
  185. try:
  186. password = config['cloudera-manager']['password']
  187. except KeyError:
  188. password = 'admin'
  189. if len(options.haproxy) == 0 and len(options.nginx) == 0:
  190. print >> sys.stderr, 'no processes specified'
  191. return 1
  192. client = ApiResource(hostname, username=username, password=password)
  193. monitor = MonitorHue(client)
  194. rpc = childutils.getRPCInterface(os.environ)
  195. if len(options.haproxy) != 0:
  196. listener = HAProxyListener(rpc, config['haproxy'], options.haproxy)
  197. monitor.listeners.append(listener)
  198. if len(options.nginx) != 0:
  199. listener = NginxListener(rpc, config['nginx'], options.nginx)
  200. monitor.listeners.append(listener)
  201. monitor.run_forever()
  202. return 0
  203. # ------------------------------------------------------------------------------
  204. if __name__ == '__main__':
  205. sys.exit(main(sys.argv))