monitor-hue-lb 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321
  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_roles_by_type("HUE_SERVER"):
  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 __init__(self, *args, **kwargs):
  142. super(NginxListener, self).__init__(*args, **kwargs)
  143. self.alias_file = self.config['alias_file']
  144. self.generate_alias_file()
  145. def find_static_files(self):
  146. paths = (
  147. '/opt/cloudera/parcels/CDH/lib/hue/build/static/',
  148. '/usr/lib/hue/build/static/',
  149. os.path.join(os.path.dirname(__file__), '../../../build/static/'),
  150. )
  151. for path in paths:
  152. if os.path.exists(path):
  153. return path
  154. else:
  155. raise Exception('could not find the Hue static files')
  156. def generate_alias_file(self):
  157. static_files = self.find_static_files()
  158. with open(self.alias_file, 'w') as f:
  159. print >> f, 'alias %s;' % static_files
  160. def expand_server_template(self, index, server):
  161. """Expand the Nginx server line"""
  162. config_server = self.server_template % {
  163. 'address': server.address,
  164. }
  165. # Let nginx know if this backend is offline.
  166. if server.state == 'STOPPED':
  167. config_server += ' down'
  168. return config_server + ';'
  169. def reload_process_config(self, process):
  170. """Let nginx know the config file has changed."""
  171. os.kill(process['pid'], signal.SIGHUP)
  172. # ------------------------------------------------------------------------------
  173. def main(argv):
  174. parser = optparse.OptionParser()
  175. parser.add_option('-c', '--config_file',
  176. help='config file')
  177. parser.add_option('--haproxy',
  178. default=[],
  179. type='str',
  180. action='append',
  181. help='monitor this haproxy process')
  182. parser.add_option('--nginx',
  183. default=[],
  184. type='str',
  185. action='append',
  186. help='monitor this haproxy process')
  187. options, args = parser.parse_args(argv[1:])
  188. if options.config_file is None:
  189. options.config_file = 'etc/hue-lb.toml'
  190. try:
  191. with open(options.config_file) as f:
  192. config = toml.load(f)
  193. except IOError, e:
  194. print >> sys.stderr, 'error opening %s: %s' % (options.config_file, e)
  195. return 1
  196. try:
  197. hostname = config['cloudera-manager']['hostname']
  198. except KeyError:
  199. hostname = 'localhost'
  200. try:
  201. username = config['cloudera-manager']['username']
  202. except KeyError:
  203. username = 'admin'
  204. try:
  205. password = config['cloudera-manager']['password']
  206. except KeyError:
  207. password = 'admin'
  208. if len(options.haproxy) == 0 and len(options.nginx) == 0:
  209. print >> sys.stderr, 'no processes specified'
  210. return 1
  211. client = ApiResource(hostname, username=username, password=password)
  212. monitor = MonitorHue(client)
  213. rpc = childutils.getRPCInterface(os.environ)
  214. if len(options.haproxy) != 0:
  215. listener = HAProxyListener(rpc, config['haproxy'], options.haproxy)
  216. monitor.listeners.append(listener)
  217. if len(options.nginx) != 0:
  218. listener = NginxListener(rpc, config['nginx'], options.nginx)
  219. monitor.listeners.append(listener)
  220. monitor.run_forever()
  221. return 0
  222. # ------------------------------------------------------------------------------
  223. if __name__ == '__main__':
  224. sys.exit(main(sys.argv))