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-02 12:54:55 +00:00
|
|
|
|
2019-09-05 14:10:46 +00:00
|
|
|
from pipe_io_server import PipeIOServer
|
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-07-24 13:17:16 +00:00
|
|
|
class PipeIOHandlerBase:
|
2019-07-12 21:25:16 +00:00
|
|
|
keys = ['']
|
|
|
|
|
|
|
|
def __init__(self, in_pipe_path, out_pipe_path, permissions=DEFAULT_PERMISSIONS):
|
2019-07-30 13:17:29 +00:00
|
|
|
self.connector = None
|
2019-09-05 13:52:34 +00:00
|
|
|
self.in_pipe = in_pipe_path
|
|
|
|
self.out_pipe = out_pipe_path
|
2019-09-04 16:07:36 +00:00
|
|
|
self.pipe_io = PipeIOServer(
|
2019-05-04 19:10:05 +00:00
|
|
|
in_pipe_path,
|
|
|
|
out_pipe_path,
|
|
|
|
permissions
|
|
|
|
)
|
2019-09-04 16:07:36 +00:00
|
|
|
self.pipe_io.handle_message = self._server_handle_message
|
2019-05-06 13:23:21 +00:00
|
|
|
self.pipe_io.start()
|
|
|
|
|
2019-09-04 16:07:36 +00:00
|
|
|
def _server_handle_message(self, message):
|
|
|
|
try:
|
|
|
|
self.handle_pipe_event(message)
|
|
|
|
except: # pylint: disable=bare-except
|
2019-09-05 13:52:34 +00:00
|
|
|
LOG.exception('Failed to handle message %s from pipe %s!', message, self.in_pipe)
|
2019-09-04 16:07:36 +00:00
|
|
|
|
2019-05-06 13:23:21 +00:00
|
|
|
@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()
|
|
|
|
|
|
|
|
|
2019-07-24 13:17:16 +00:00
|
|
|
class PipeIOHandler(PipeIOHandlerBase):
|
2019-07-12 21:25:16 +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-07-30 13:17:29 +00:00
|
|
|
self.connector.send_message(json)
|