Joshua Hesketh | 39a0fee | 2013-07-31 12:00:53 +1000 | [diff] [blame] | 1 | # Copyright 2013 Rackspace Australia |
| 2 | # |
| 3 | # Licensed under the Apache License, Version 2.0 (the "License"); you may |
| 4 | # not use this file except in compliance with the License. You may obtain |
| 5 | # a copy of the License at |
| 6 | # |
| 7 | # http://www.apache.org/licenses/LICENSE-2.0 |
| 8 | # |
| 9 | # Unless required by applicable law or agreed to in writing, software |
| 10 | # distributed under the License is distributed on an "AS IS" BASIS, WITHOUT |
| 11 | # WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the |
| 12 | # License for the specific language governing permissions and limitations |
| 13 | # under the License. |
| 14 | |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 15 | |
| 16 | import gear |
| 17 | import json |
| 18 | import logging |
| 19 | import os |
| 20 | import threading |
| 21 | |
| 22 | |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 23 | class ZuulManager(threading.Thread): |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 24 | |
| 25 | """ This thread manages all of the launched gearman workers. |
| 26 | As required by the zuul protocol it handles stopping builds when they |
Joshua Hesketh | 363d004 | 2013-07-26 11:44:07 +1000 | [diff] [blame] | 27 | are cancelled through stop:turbo-hipster-manager-%hostname. |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 28 | To do this it implements its own gearman worker waiting for events on |
| 29 | that manager. """ |
| 30 | |
Joshua Hesketh | 6eb5fdc | 2013-11-20 12:36:06 +1100 | [diff] [blame] | 31 | log = logging.getLogger("worker_manager.ZuulManager") |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 32 | |
| 33 | def __init__(self, config, tasks): |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 34 | super(ZuulManager, self).__init__() |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 35 | self._stop = threading.Event() |
| 36 | self.config = config |
| 37 | self.tasks = tasks |
| 38 | |
| 39 | self.gearman_worker = None |
| 40 | self.setup_gearman() |
| 41 | |
| 42 | def setup_gearman(self): |
| 43 | hostname = os.uname()[1] |
Joshua Hesketh | 363d004 | 2013-07-26 11:44:07 +1000 | [diff] [blame] | 44 | self.gearman_worker = gear.Worker('turbo-hipster-manager-%s' |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 45 | % hostname) |
| 46 | self.gearman_worker.addServer( |
| 47 | self.config['zuul_server']['gearman_host'], |
| 48 | self.config['zuul_server']['gearman_port'] |
| 49 | ) |
Joshua Hesketh | 9cd2f93 | 2014-03-05 16:49:22 +1100 | [diff] [blame] | 50 | |
| 51 | def register_functions(self): |
| 52 | hostname = os.uname()[1] |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 53 | self.gearman_worker.registerFunction( |
Joshua Hesketh | 363d004 | 2013-07-26 11:44:07 +1000 | [diff] [blame] | 54 | 'stop:turbo-hipster-manager-%s' % hostname) |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 55 | |
| 56 | def stop(self): |
| 57 | self._stop.set() |
| 58 | # Unblock gearman |
| 59 | self.log.debug("Telling gearman to stop waiting for jobs") |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 60 | self.gearman_worker.shutdown() |
| 61 | |
| 62 | def stopped(self): |
| 63 | return self._stop.isSet() |
| 64 | |
| 65 | def run(self): |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 66 | while not self.stopped(): |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 67 | try: |
| 68 | # gearman_worker.getJob() blocks until a job is available |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 69 | self.log.debug("Waiting for server") |
| 70 | self.gearman_worker.waitForServer() |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 71 | if (not self.stopped() and self.gearman_worker.running and |
| 72 | self.gearman_worker.active_connections): |
Joshua Hesketh | 9cd2f93 | 2014-03-05 16:49:22 +1100 | [diff] [blame] | 73 | self.register_functions() |
| 74 | self.gearman_worker.waitForServer() |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 75 | logging.debug("Waiting for job") |
| 76 | self.current_step = 0 |
| 77 | job = self.gearman_worker.getJob() |
| 78 | self._handle_job(job) |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 79 | except: |
| 80 | logging.exception('Exception retrieving log event.') |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 81 | self.log.debug("Finished manager thread") |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 82 | |
| 83 | def _handle_job(self, job): |
| 84 | """ Handle the requested job """ |
| 85 | try: |
| 86 | job_arguments = json.loads(job.arguments.decode('utf-8')) |
Joshua Hesketh | 38a1718 | 2014-03-05 14:19:38 +1100 | [diff] [blame] | 87 | self.tasks[job_arguments['name']].stop_working( |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 88 | job_arguments['number']) |
Joshua Hesketh | 4a30d27 | 2013-08-12 11:17:25 +1000 | [diff] [blame] | 89 | job.sendWorkComplete() |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 90 | except Exception as e: |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 91 | self.log.exception('Exception waiting for management job.') |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 92 | job.sendWorkException(str(e).encode('utf-8')) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 93 | |
| 94 | |
| 95 | class ZuulClient(threading.Thread): |
| 96 | |
| 97 | """ ...""" |
| 98 | |
| 99 | log = logging.getLogger("worker_manager.ZuulClient") |
| 100 | |
| 101 | def __init__(self, global_config, worker_name): |
| 102 | super(ZuulClient, self).__init__() |
| 103 | self._stop = threading.Event() |
| 104 | self.global_config = global_config |
| 105 | |
| 106 | self.worker_name = worker_name |
| 107 | |
| 108 | # Set up the runner worker |
| 109 | self.gearman_worker = None |
| 110 | self.functions = {} |
| 111 | |
| 112 | self.job = None |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 113 | |
| 114 | self.setup_gearman() |
| 115 | |
| 116 | def setup_gearman(self): |
| 117 | self.log.debug("Set up gearman worker") |
| 118 | self.gearman_worker = gear.Worker(self.worker_name) |
| 119 | self.gearman_worker.addServer( |
| 120 | self.global_config['zuul_server']['gearman_host'], |
| 121 | self.global_config['zuul_server']['gearman_port'] |
| 122 | ) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 123 | |
| 124 | def register_functions(self): |
Joshua Hesketh | fc44906 | 2013-11-20 16:06:06 +1100 | [diff] [blame] | 125 | self.log.debug("Register functions with gearman") |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 126 | for function_name, plugin in self.functions.items(): |
| 127 | self.gearman_worker.registerFunction(function_name) |
Joshua Hesketh | b39f884 | 2013-11-21 13:12:00 +1100 | [diff] [blame] | 128 | self.log.debug(self.gearman_worker.functions) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 129 | |
| 130 | def add_function(self, function_name, plugin): |
Joshua Hesketh | fc44906 | 2013-11-20 16:06:06 +1100 | [diff] [blame] | 131 | self.log.debug("Add function, %s, to list" % function_name) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 132 | self.functions[function_name] = plugin |
| 133 | |
| 134 | def stop(self): |
| 135 | self._stop.set() |
Joshua Hesketh | 38a1718 | 2014-03-05 14:19:38 +1100 | [diff] [blame] | 136 | for task in self.functions.values(): |
| 137 | task.stop_working() |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 138 | # Unblock gearman |
| 139 | self.log.debug("Telling gearman to stop waiting for jobs") |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 140 | self.gearman_worker.shutdown() |
| 141 | |
| 142 | def stopped(self): |
| 143 | return self._stop.isSet() |
| 144 | |
| 145 | def run(self): |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 146 | while not self.stopped(): |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 147 | try: |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 148 | # gearman_worker.getJob() blocks until a job is available |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 149 | self.log.debug("Waiting for server") |
| 150 | self.gearman_worker.waitForServer() |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 151 | if (not self.stopped() and self.gearman_worker.running and |
| 152 | self.gearman_worker.active_connections): |
Joshua Hesketh | 9cd2f93 | 2014-03-05 16:49:22 +1100 | [diff] [blame] | 153 | self.register_functions() |
| 154 | self.gearman_worker.waitForServer() |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 155 | self.log.debug("Waiting for job") |
| 156 | self.job = self.gearman_worker.getJob() |
| 157 | self._handle_job() |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 158 | except: |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 159 | self.log.exception('Exception waiting for job.') |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 160 | self.log.debug("Finished client thread") |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 161 | |
| 162 | def _handle_job(self): |
| 163 | """ We have a job, give it to the right plugin """ |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame] | 164 | if self.job: |
| 165 | self.log.debug("We have a job, we'll launch the task now.") |
| 166 | self.functions[self.job.name].start_job(self.job) |