Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
70f4d666e3 | ||
|
|
dc76f1b732 | ||
|
|
f374cb7e46 | ||
|
|
9481e8f921 | ||
|
|
c8d4080ef5 |
@@ -1,6 +1,6 @@
|
|||||||
from os.path import dirname, realpath, join
|
from os.path import dirname, realpath, join
|
||||||
|
|
||||||
from setuptools import setup
|
from setuptools import setup, find_packages
|
||||||
|
|
||||||
here = dirname(realpath(__file__))
|
here = dirname(realpath(__file__))
|
||||||
|
|
||||||
@@ -17,7 +17,7 @@ setup(
|
|||||||
author='Avatao.com Innovative Learning Kft.',
|
author='Avatao.com Innovative Learning Kft.',
|
||||||
author_email='support@avatao.com',
|
author_email='support@avatao.com',
|
||||||
license='custom',
|
license='custom',
|
||||||
packages=['tfw'],
|
packages=find_packages(),
|
||||||
package_dir={'tfw': 'tfw'},
|
package_dir={'tfw': 'tfw'},
|
||||||
install_requires=requirements,
|
install_requires=requirements,
|
||||||
extras_require={
|
extras_require={
|
||||||
|
|||||||
@@ -16,10 +16,13 @@ class ProcessLogHandler:
|
|||||||
self._initial_log_tail = log_tail
|
self._initial_log_tail = log_tail
|
||||||
|
|
||||||
self.command_handlers = {
|
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):
|
def start(self):
|
||||||
|
if not self._monitor:
|
||||||
self._monitor = LogInotifyObserver(
|
self._monitor = LogInotifyObserver(
|
||||||
connector=self.connector,
|
connector=self.connector,
|
||||||
process_name=self.process_name,
|
process_name=self.process_name,
|
||||||
@@ -40,5 +43,16 @@ class ProcessLogHandler:
|
|||||||
if data.get('tail'):
|
if data.get('tail'):
|
||||||
self._monitor.log_tail = data['tail']
|
self._monitor.log_tail = data['tail']
|
||||||
|
|
||||||
def cleanup(self):
|
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.stop()
|
||||||
|
self._monitor = None
|
||||||
|
|
||||||
|
def cleanup(self):
|
||||||
|
self._stop_monitor()
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ class TerminadoMiniServer:
|
|||||||
url,
|
url,
|
||||||
TerminadoMiniServer.ResetterTermSocket,
|
TerminadoMiniServer.ResetterTermSocket,
|
||||||
{'term_manager': self._term_manager}
|
{'term_manager': self._term_manager}
|
||||||
)])
|
)], websocket_ping_interval=30)
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def term_manager(self):
|
def term_manager(self):
|
||||||
|
|||||||
+3
-4
@@ -1,6 +1,7 @@
|
|||||||
import logging
|
import logging
|
||||||
from collections import defaultdict
|
from collections import defaultdict
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
|
from contextlib import suppress
|
||||||
|
|
||||||
from transitions import Machine, MachineError
|
from transitions import Machine, MachineError
|
||||||
|
|
||||||
@@ -62,14 +63,12 @@ class FSMBase(Machine, CallbackMixin):
|
|||||||
)
|
)
|
||||||
|
|
||||||
if all(predicate_results):
|
if all(predicate_results):
|
||||||
try:
|
with suppress(AttributeError, MachineError):
|
||||||
from_state = self.state
|
from_state = self.state
|
||||||
self.trigger(trigger)
|
if self.trigger(trigger):
|
||||||
self.update_event_log(from_state, trigger)
|
self.update_event_log(from_state, trigger)
|
||||||
return True
|
return True
|
||||||
except (AttributeError, MachineError):
|
|
||||||
LOG.debug('FSM failed to execute nonexistent trigger: "%s"', trigger)
|
LOG.debug('FSM failed to execute nonexistent trigger: "%s"', trigger)
|
||||||
return False
|
|
||||||
|
|
||||||
def update_event_log(self, from_state, trigger):
|
def update_event_log(self, from_state, trigger):
|
||||||
self.event_log.append({
|
self.event_log.append({
|
||||||
|
|||||||
@@ -42,7 +42,8 @@ def deserialize_tfw_msg(*args):
|
|||||||
"""
|
"""
|
||||||
Return message from TFW multipart data
|
Return message from TFW multipart data
|
||||||
"""
|
"""
|
||||||
return _deserialize_all(*args)[1]
|
envelope = _deserialize_all(*args)
|
||||||
|
return _repair_if_needed(envelope)
|
||||||
|
|
||||||
|
|
||||||
def _serialize_all(*args):
|
def _serialize_all(*args):
|
||||||
@@ -84,6 +85,18 @@ def _deserialize_single(data):
|
|||||||
return _decode_if_needed(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):
|
def _encode_if_needed(value):
|
||||||
"""
|
"""
|
||||||
Return input as bytes
|
Return input as bytes
|
||||||
|
|||||||
@@ -28,7 +28,7 @@ class TFWServer:
|
|||||||
r'/ws', ZMQWebSocketRouter, {
|
r'/ws', ZMQWebSocketRouter, {
|
||||||
'listener': self._listener,
|
'listener': self._listener,
|
||||||
}
|
}
|
||||||
)])
|
)], websocket_ping_interval=30)
|
||||||
|
|
||||||
def listen(self):
|
def listen(self):
|
||||||
self.application.listen(TFWENV.WEB_PORT)
|
self.application.listen(TFWENV.WEB_PORT)
|
||||||
|
|||||||
Reference in New Issue
Block a user