Compare commits

...
6 changed files with 47 additions and 21 deletions
+2 -2
View File
@@ -1,6 +1,6 @@
from os.path import dirname, realpath, join
from setuptools import setup
from setuptools import setup, find_packages
here = dirname(realpath(__file__))
@@ -17,7 +17,7 @@ setup(
author='Avatao.com Innovative Learning Kft.',
author_email='support@avatao.com',
license='custom',
packages=['tfw'],
packages=find_packages(),
package_dir={'tfw': 'tfw'},
install_requires=requirements,
extras_require={
@@ -16,17 +16,20 @@ class ProcessLogHandler:
self._initial_log_tail = log_tail
self.command_handlers = {
'process.log.set': self.handle_set
'process.log.set': self.handle_set,
'process.log.start': self.handle_start,
'process.log.stop': self.handle_stop
}
def start(self):
self._monitor = LogInotifyObserver(
connector=self.connector,
process_name=self.process_name,
supervisor_uri=self._supervisor_uri,
log_tail=self._initial_log_tail
)
self._monitor.start()
if not self._monitor:
self._monitor = LogInotifyObserver(
connector=self.connector,
process_name=self.process_name,
supervisor_uri=self._supervisor_uri,
log_tail=self._initial_log_tail
)
self._monitor.start()
def handle_event(self, message, _):
try:
@@ -40,5 +43,16 @@ class ProcessLogHandler:
if data.get('tail'):
self._monitor.log_tail = data['tail']
def handle_start(self, _):
self.start()
def handle_stop(self, _):
self._stop_monitor()
def _stop_monitor(self):
if self._monitor:
self._monitor.stop()
self._monitor = None
def cleanup(self):
self._monitor.stop()
self._stop_monitor()
@@ -17,7 +17,7 @@ class TerminadoMiniServer:
url,
TerminadoMiniServer.ResetterTermSocket,
{'term_manager': self._term_manager}
)])
)], websocket_ping_interval=30)
@property
def term_manager(self):
+6 -7
View File
@@ -1,6 +1,7 @@
import logging
from collections import defaultdict
from datetime import datetime
from contextlib import suppress
from transitions import Machine, MachineError
@@ -62,14 +63,12 @@ class FSMBase(Machine, CallbackMixin):
)
if all(predicate_results):
try:
with suppress(AttributeError, MachineError):
from_state = self.state
self.trigger(trigger)
self.update_event_log(from_state, trigger)
return True
except (AttributeError, MachineError):
LOG.debug('FSM failed to execute nonexistent trigger: "%s"', trigger)
return False
if self.trigger(trigger):
self.update_event_log(from_state, trigger)
return True
LOG.debug('FSM failed to execute nonexistent trigger: "%s"', trigger)
def update_event_log(self, from_state, trigger):
self.event_log.append({
+14 -1
View File
@@ -42,7 +42,8 @@ def deserialize_tfw_msg(*args):
"""
Return message from TFW multipart data
"""
return _deserialize_all(*args)[1]
envelope = _deserialize_all(*args)
return _repair_if_needed(envelope)
def _serialize_all(*args):
@@ -84,6 +85,18 @@ def _deserialize_single(data):
return _decode_if_needed(data)
def _repair_if_needed(envelope):
"""
Quick fix for broken messages received from separate processes.
"""
if len(envelope) == 2:
return envelope[1]
for part in envelope:
if isinstance(part, dict):
return part
return {}
def _encode_if_needed(value):
"""
Return input as bytes
+1 -1
View File
@@ -28,7 +28,7 @@ class TFWServer:
r'/ws', ZMQWebSocketRouter, {
'listener': self._listener,
}
)])
)], websocket_ping_interval=30)
def listen(self):
self.application.listen(TFWENV.WEB_PORT)