2019-06-10 13:32:45 +00:00
|
|
|
import logging
|
2019-05-06 13:23:21 +00:00
|
|
|
from abc import abstractmethod
|
2019-05-09 14:56:16 +00:00
|
|
|
from json import loads, dumps
|
2019-05-09 16:56:30 +00:00
|
|
|
from subprocess import run, PIPE, Popen
|
2019-05-09 16:22:33 +00:00
|
|
|
from functools import partial
|
2019-05-09 16:56:30 +00:00
|
|
|
from os import getpgid, killpg
|
|
|
|
from os.path import join
|
|
|
|
from signal import SIGTERM
|
|
|
|
from secrets import token_urlsafe
|
2019-05-12 20:30:49 +00:00
|
|
|
from threading import Thread
|
|
|
|
from contextlib import suppress
|
2019-05-02 12:54:55 +00:00
|
|
|
|
2019-05-27 12:09:13 +00:00
|
|
|
from tfw.event_handlers import EventHandlerBase
|
2019-05-02 12:54:55 +00:00
|
|
|
|
2019-05-12 20:30:49 +00:00
|
|
|
from .pipe_io_server import PipeIOServer, terminate_process_on_failure
|
2019-05-02 12:54:55 +00:00
|
|
|
|
|
|
|
LOG = logging.getLogger(__name__)
|
2019-05-09 14:56:16 +00:00
|
|
|
DEFAULT_PERMISSIONS = 0o600
|
2019-05-02 12:54:55 +00:00
|
|
|
|
|
|
|
|
2019-05-06 13:23:21 +00:00
|
|
|
class PipeIOEventHandlerBase(EventHandlerBase):
|
2019-05-09 14:56:16 +00:00
|
|
|
def __init__(self, key, in_pipe_path, out_pipe_path, permissions=DEFAULT_PERMISSIONS):
|
2019-05-02 12:54:55 +00:00
|
|
|
super().__init__(key)
|
2019-05-06 13:23:21 +00:00
|
|
|
self.pipe_io = CallbackPipeIOServer(
|
2019-05-04 19:10:05 +00:00
|
|
|
in_pipe_path,
|
|
|
|
out_pipe_path,
|
2019-05-06 13:23:21 +00:00
|
|
|
self.handle_pipe_event,
|
2019-05-04 19:10:05 +00:00
|
|
|
permissions
|
|
|
|
)
|
2019-05-06 13:23:21 +00:00
|
|
|
self.pipe_io.start()
|
|
|
|
|
|
|
|
@abstractmethod
|
2019-05-09 13:14:47 +00:00
|
|
|
def handle_pipe_event(self, message_bytes):
|
2019-05-06 13:23:21 +00:00
|
|
|
raise NotImplementedError()
|
2019-05-02 12:54:55 +00:00
|
|
|
|
|
|
|
def cleanup(self):
|
2019-05-06 13:23:21 +00:00
|
|
|
self.pipe_io.stop()
|
|
|
|
|
|
|
|
|
|
|
|
class CallbackPipeIOServer(PipeIOServer):
|
|
|
|
def __init__(self, in_pipe_path, out_pipe_path, callback, permissions):
|
|
|
|
super().__init__(in_pipe_path, out_pipe_path, permissions)
|
|
|
|
self.callback = callback
|
|
|
|
|
|
|
|
def handle_message(self, message):
|
2019-05-12 20:28:31 +00:00
|
|
|
try:
|
|
|
|
self.callback(message)
|
|
|
|
except: # pylint: disable=bare-except
|
|
|
|
LOG.exception('Failed to handle message %s from pipe %s!', message, self.in_pipe)
|
2019-05-02 12:54:55 +00:00
|
|
|
|
2019-05-06 13:23:21 +00:00
|
|
|
|
|
|
|
class PipeIOEventHandler(PipeIOEventHandlerBase):
|
2019-05-02 12:54:55 +00:00
|
|
|
def handle_event(self, message):
|
2019-05-09 13:12:42 +00:00
|
|
|
json_bytes = dumps(message).encode()
|
|
|
|
self.pipe_io.send_message(json_bytes)
|
2019-05-04 19:13:58 +00:00
|
|
|
|
2019-05-09 13:14:47 +00:00
|
|
|
def handle_pipe_event(self, message_bytes):
|
|
|
|
json = loads(message_bytes)
|
2019-05-26 16:26:33 +00:00
|
|
|
self.send_message(json)
|
2019-05-09 14:56:16 +00:00
|
|
|
|
|
|
|
|
|
|
|
class TransformerPipeIOEventHandler(PipeIOEventHandlerBase):
|
|
|
|
# pylint: disable=too-many-arguments
|
|
|
|
def __init__(
|
|
|
|
self, key, in_pipe_path, out_pipe_path,
|
2019-05-09 16:22:33 +00:00
|
|
|
transform_in_cmd, transform_out_cmd,
|
2019-05-09 14:56:16 +00:00
|
|
|
permissions=DEFAULT_PERMISSIONS
|
|
|
|
):
|
2019-05-09 16:22:33 +00:00
|
|
|
self._transform_in = partial(self._transform_message, transform_in_cmd)
|
|
|
|
self._transform_out = partial(self._transform_message, transform_out_cmd)
|
2019-05-09 14:56:16 +00:00
|
|
|
super().__init__(key, in_pipe_path, out_pipe_path, permissions)
|
|
|
|
|
|
|
|
@staticmethod
|
|
|
|
def _transform_message(transform_cmd, message):
|
|
|
|
proc = run(
|
|
|
|
transform_cmd,
|
|
|
|
input=message,
|
|
|
|
stdout=PIPE,
|
|
|
|
stderr=PIPE,
|
|
|
|
shell=True
|
|
|
|
)
|
|
|
|
if proc.returncode == 0:
|
|
|
|
return proc.stdout
|
|
|
|
raise ValueError(f'Transforming message {message} failed!')
|
|
|
|
|
|
|
|
def handle_event(self, message):
|
|
|
|
json_bytes = dumps(message).encode()
|
2019-05-09 16:22:33 +00:00
|
|
|
transformed_bytes = self._transform_out(json_bytes)
|
2019-05-09 14:56:16 +00:00
|
|
|
if transformed_bytes:
|
|
|
|
self.pipe_io.send_message(transformed_bytes)
|
|
|
|
|
|
|
|
def handle_pipe_event(self, message_bytes):
|
2019-05-09 16:22:33 +00:00
|
|
|
transformed_bytes = self._transform_in(message_bytes)
|
2019-05-09 14:56:16 +00:00
|
|
|
if transformed_bytes:
|
|
|
|
json_message = loads(transformed_bytes)
|
2019-05-26 16:26:33 +00:00
|
|
|
self.send_message(json_message)
|
2019-05-09 16:56:30 +00:00
|
|
|
|
|
|
|
|
|
|
|
class CommandEventHandler(PipeIOEventHandler):
|
|
|
|
def __init__(self, key, command, permissions=DEFAULT_PERMISSIONS):
|
|
|
|
super().__init__(
|
|
|
|
key,
|
|
|
|
self._generate_tempfilename(),
|
|
|
|
self._generate_tempfilename(),
|
|
|
|
permissions
|
|
|
|
)
|
|
|
|
|
|
|
|
self._proc_stdin = open(self.pipe_io.out_pipe, 'rb')
|
|
|
|
self._proc_stdout = open(self.pipe_io.in_pipe, 'wb')
|
|
|
|
self._proc = Popen(
|
|
|
|
command, shell=True, executable='/bin/bash',
|
2019-05-12 20:30:49 +00:00
|
|
|
stdin=self._proc_stdin, stdout=self._proc_stdout, stderr=PIPE,
|
2019-05-09 16:56:30 +00:00
|
|
|
start_new_session=True
|
|
|
|
)
|
2019-05-13 09:17:30 +00:00
|
|
|
self._monitor_proc_thread = self._start_monitor_proc()
|
2019-05-09 16:56:30 +00:00
|
|
|
|
|
|
|
def _generate_tempfilename(self):
|
|
|
|
# pylint: disable=no-self-use
|
|
|
|
random_filename = partial(token_urlsafe, 10)
|
|
|
|
return join('/tmp', f'{type(self).__name__}.{random_filename()}')
|
|
|
|
|
2019-05-13 09:17:30 +00:00
|
|
|
def _start_monitor_proc(self):
|
|
|
|
thread = Thread(target=self._monitor_proc, daemon=True)
|
2019-05-12 20:30:49 +00:00
|
|
|
thread.start()
|
|
|
|
return thread
|
|
|
|
|
|
|
|
@terminate_process_on_failure
|
2019-05-13 09:17:30 +00:00
|
|
|
def _monitor_proc(self):
|
|
|
|
return_code = self._proc.wait()
|
2019-05-13 12:53:31 +00:00
|
|
|
if return_code == -int(SIGTERM):
|
|
|
|
# supervisord asked the program to terminate, this is fine
|
|
|
|
return
|
2019-05-13 09:17:30 +00:00
|
|
|
if return_code != 0:
|
2019-05-12 20:30:49 +00:00
|
|
|
_, stderr = self._proc.communicate()
|
2019-05-13 12:52:17 +00:00
|
|
|
raise RuntimeError(f'Subprocess failed ({return_code})! Stderr:\n{stderr.decode()}')
|
2019-05-12 20:30:49 +00:00
|
|
|
|
2019-05-09 16:56:30 +00:00
|
|
|
def cleanup(self):
|
2019-05-12 20:30:49 +00:00
|
|
|
with suppress(ProcessLookupError):
|
|
|
|
process_group_id = getpgid(self._proc.pid)
|
|
|
|
killpg(process_group_id, SIGTERM)
|
|
|
|
self._proc_stdin.close()
|
|
|
|
self._proc_stdout.close()
|
2019-05-09 16:56:30 +00:00
|
|
|
super().cleanup()
|