James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 1 | # Copyright 2014 OpenStack Foundation |
| 2 | # Copyright 2014 Hewlett-Packard Development Company, L.P. |
| 3 | # Copyright 2016 Red Hat |
| 4 | # |
| 5 | # Licensed under the Apache License, Version 2.0 (the "License"); you may |
| 6 | # not use this file except in compliance with the License. You may obtain |
| 7 | # a copy of the License at |
| 8 | # |
| 9 | # http://www.apache.org/licenses/LICENSE-2.0 |
| 10 | # |
| 11 | # Unless required by applicable law or agreed to in writing, software |
| 12 | # distributed under the License is distributed on an "AS IS" BASIS, WITHOUT |
| 13 | # WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the |
| 14 | # License for the specific language governing permissions and limitations |
| 15 | # under the License. |
| 16 | |
| 17 | import logging |
| 18 | import os |
| 19 | import socket |
| 20 | import threading |
Clint Byrum | ff97edb | 2017-05-10 21:30:14 -0700 | [diff] [blame] | 21 | from six.moves import queue |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 22 | |
| 23 | |
| 24 | class CommandSocket(object): |
| 25 | log = logging.getLogger("zuul.CommandSocket") |
| 26 | |
| 27 | def __init__(self, path): |
| 28 | self.running = False |
| 29 | self.path = path |
Clint Byrum | ff97edb | 2017-05-10 21:30:14 -0700 | [diff] [blame] | 30 | self.queue = queue.Queue() |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 31 | |
| 32 | def start(self): |
| 33 | self.running = True |
| 34 | if os.path.exists(self.path): |
| 35 | os.unlink(self.path) |
| 36 | self.socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) |
| 37 | self.socket.bind(self.path) |
| 38 | self.socket.listen(1) |
| 39 | self.socket_thread = threading.Thread(target=self._socketListener) |
| 40 | self.socket_thread.daemon = True |
| 41 | self.socket_thread.start() |
| 42 | |
| 43 | def stop(self): |
| 44 | # First, wake up our listener thread with a connection and |
| 45 | # tell it to stop running. |
| 46 | self.running = False |
| 47 | s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) |
| 48 | s.connect(self.path) |
Clint Byrum | f322fe2 | 2017-05-10 20:53:12 -0700 | [diff] [blame] | 49 | s.sendall(b'_stop\n') |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 50 | # The command '_stop' will be ignored by our listener, so |
| 51 | # directly inject it into the queue so that consumers of this |
| 52 | # class which are waiting in .get() are awakened. They can |
| 53 | # either handle '_stop' or just ignore the unknown command and |
| 54 | # then check to see if they should continue to run before |
| 55 | # re-entering their loop. |
Clint Byrum | f322fe2 | 2017-05-10 20:53:12 -0700 | [diff] [blame] | 56 | self.queue.put(b'_stop') |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 57 | self.socket_thread.join() |
| 58 | |
| 59 | def _socketListener(self): |
| 60 | while self.running: |
| 61 | try: |
| 62 | s, addr = self.socket.accept() |
| 63 | self.log.debug("Accepted socket connection %s" % (s,)) |
Clint Byrum | f322fe2 | 2017-05-10 20:53:12 -0700 | [diff] [blame] | 64 | buf = b'' |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 65 | while True: |
| 66 | buf += s.recv(1) |
Clint Byrum | f322fe2 | 2017-05-10 20:53:12 -0700 | [diff] [blame] | 67 | if buf[-1:] == b'\n': |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 68 | break |
| 69 | buf = buf.strip() |
| 70 | self.log.debug("Received %s from socket" % (buf,)) |
| 71 | s.close() |
| 72 | # Because we use '_stop' internally to wake up a |
| 73 | # waiting thread, don't allow it to actually be |
| 74 | # injected externally. |
Clint Byrum | f322fe2 | 2017-05-10 20:53:12 -0700 | [diff] [blame] | 75 | if buf != b'_stop': |
James E. Blair | c4b2041 | 2016-05-27 16:45:26 -0700 | [diff] [blame] | 76 | self.queue.put(buf) |
| 77 | except Exception: |
| 78 | self.log.exception("Exception in socket handler") |
| 79 | |
| 80 | def get(self): |
| 81 | if not self.running: |
| 82 | raise Exception("CommandSocket.get called while stopped") |
| 83 | return self.queue.get() |