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 | ) |
| 50 | self.gearman_worker.registerFunction( |
Joshua Hesketh | 363d004 | 2013-07-26 11:44:07 +1000 | [diff] [blame] | 51 | 'stop:turbo-hipster-manager-%s' % hostname) |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 52 | |
| 53 | def stop(self): |
| 54 | self._stop.set() |
| 55 | # Unblock gearman |
| 56 | self.log.debug("Telling gearman to stop waiting for jobs") |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 57 | self.gearman_worker.shutdown() |
| 58 | |
| 59 | def stopped(self): |
| 60 | return self._stop.isSet() |
| 61 | |
| 62 | def run(self): |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 63 | while not self.stopped(): |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 64 | try: |
| 65 | # gearman_worker.getJob() blocks until a job is available |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 66 | self.log.debug("Waiting for server") |
| 67 | self.gearman_worker.waitForServer() |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 68 | if (not self.stopped() and self.gearman_worker.running and |
| 69 | self.gearman_worker.active_connections): |
| 70 | logging.debug("Waiting for job") |
| 71 | self.current_step = 0 |
| 72 | job = self.gearman_worker.getJob() |
| 73 | self._handle_job(job) |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 74 | except: |
| 75 | logging.exception('Exception retrieving log event.') |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 76 | self.log.debug("Finished manager thread") |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 77 | |
| 78 | def _handle_job(self, job): |
| 79 | """ Handle the requested job """ |
| 80 | try: |
| 81 | job_arguments = json.loads(job.arguments.decode('utf-8')) |
| 82 | self.tasks[job_arguments['name']].stop_worker( |
| 83 | job_arguments['number']) |
Joshua Hesketh | 4a30d27 | 2013-08-12 11:17:25 +1000 | [diff] [blame] | 84 | job.sendWorkComplete() |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 85 | except Exception as e: |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 86 | self.log.exception('Exception waiting for management job.') |
Joshua Hesketh | 0ddd638 | 2013-07-26 10:33:36 +1000 | [diff] [blame] | 87 | job.sendWorkException(str(e).encode('utf-8')) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 88 | |
| 89 | |
| 90 | class ZuulClient(threading.Thread): |
| 91 | |
| 92 | """ ...""" |
| 93 | |
| 94 | log = logging.getLogger("worker_manager.ZuulClient") |
| 95 | |
| 96 | def __init__(self, global_config, worker_name): |
| 97 | super(ZuulClient, self).__init__() |
| 98 | self._stop = threading.Event() |
| 99 | self.global_config = global_config |
| 100 | |
| 101 | self.worker_name = worker_name |
| 102 | |
| 103 | # Set up the runner worker |
| 104 | self.gearman_worker = None |
| 105 | self.functions = {} |
| 106 | |
| 107 | self.job = None |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 108 | |
| 109 | self.setup_gearman() |
| 110 | |
| 111 | def setup_gearman(self): |
| 112 | self.log.debug("Set up gearman worker") |
| 113 | self.gearman_worker = gear.Worker(self.worker_name) |
| 114 | self.gearman_worker.addServer( |
| 115 | self.global_config['zuul_server']['gearman_host'], |
| 116 | self.global_config['zuul_server']['gearman_port'] |
| 117 | ) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 118 | |
| 119 | def register_functions(self): |
Joshua Hesketh | fc44906 | 2013-11-20 16:06:06 +1100 | [diff] [blame] | 120 | self.log.debug("Register functions with gearman") |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 121 | for function_name, plugin in self.functions.items(): |
| 122 | self.gearman_worker.registerFunction(function_name) |
Joshua Hesketh | b39f884 | 2013-11-21 13:12:00 +1100 | [diff] [blame] | 123 | self.log.debug(self.gearman_worker.functions) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 124 | |
| 125 | def add_function(self, function_name, plugin): |
Joshua Hesketh | fc44906 | 2013-11-20 16:06:06 +1100 | [diff] [blame] | 126 | self.log.debug("Add function, %s, to list" % function_name) |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 127 | self.functions[function_name] = plugin |
| 128 | |
| 129 | def stop(self): |
| 130 | self._stop.set() |
| 131 | # Unblock gearman |
| 132 | self.log.debug("Telling gearman to stop waiting for jobs") |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 133 | self.gearman_worker.shutdown() |
| 134 | |
| 135 | def stopped(self): |
| 136 | return self._stop.isSet() |
| 137 | |
| 138 | def run(self): |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 139 | while not self.stopped(): |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 140 | try: |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 141 | # gearman_worker.getJob() blocks until a job is available |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 142 | self.log.debug("Waiting for server") |
| 143 | self.gearman_worker.waitForServer() |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 144 | if (not self.stopped() and self.gearman_worker.running and |
| 145 | self.gearman_worker.active_connections): |
| 146 | self.log.debug("Waiting for job") |
| 147 | self.job = self.gearman_worker.getJob() |
| 148 | self._handle_job() |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 149 | except: |
Joshua Hesketh | dd50d27 | 2014-01-05 15:35:52 +1100 | [diff] [blame] | 150 | self.log.exception('Exception waiting for job.') |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 151 | self.log.debug("Finished client thread") |
Joshua Hesketh | 4343b95 | 2013-11-20 12:11:55 +1100 | [diff] [blame] | 152 | |
| 153 | def _handle_job(self): |
| 154 | """ We have a job, give it to the right plugin """ |
Joshua Hesketh | 81b5fb8 | 2014-03-05 13:06:08 +1100 | [diff] [blame^] | 155 | if self.job: |
| 156 | self.log.debug("We have a job, we'll launch the task now.") |
| 157 | self.functions[self.job.name].start_job(self.job) |