Compare commits
91
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
06e17a6591 | ||
|
|
8396c10a7b | ||
|
|
17c5ade3ce | ||
|
|
f4310c855f | ||
|
|
c831ba5ca4 | ||
|
|
3b30b333e5 | ||
|
|
9b85c19d40 | ||
|
|
0b651abb7d | ||
|
|
b1592e8ebb | ||
|
|
3f36826597 | ||
|
|
db1b008fc9 | ||
|
|
9d93d055b8 | ||
|
|
addd517ba7 | ||
|
|
690f9bb190 | ||
|
|
c909317ca0 | ||
|
|
25d039e84e | ||
|
|
117e7ffc18 | ||
|
|
8e8963b0f8 | ||
|
|
b592a1625a | ||
|
|
13e90dde41 | ||
|
|
63c235bb41 | ||
|
|
1061220f82 | ||
|
|
ca1278e0b4 | ||
|
|
137d1dbce2 | ||
|
|
3b1f42858d | ||
|
|
18124bf8c3 | ||
|
|
d8988ea097 | ||
|
|
986a67f5db | ||
|
|
39243811dc | ||
|
|
22cdbaac2c | ||
|
|
743fab366b | ||
|
|
efb8af4cd1 | ||
|
|
a65d3127ba | ||
|
|
1c29b700c2 | ||
|
|
e2bb126e6f | ||
|
|
586c435cc2 | ||
|
|
b082279916 | ||
|
|
ecec1cde59 | ||
|
|
0cd11f9d64 | ||
|
|
fe7bd4a678 | ||
|
|
3c12b5503e | ||
|
|
499b575adf | ||
|
|
9eb6717972 | ||
|
|
aabaadc4c4 | ||
|
|
cd36b6c22c | ||
|
|
67a23166c6 | ||
|
|
18e9ee5f52 | ||
|
|
d04479157d | ||
|
|
b8cdb30318 | ||
|
|
512a965328 | ||
|
|
c33a2f8103 | ||
|
|
9f4364ec10 | ||
|
|
7877e977ab | ||
|
|
5d8d541394 | ||
|
|
7e991426e7 | ||
|
|
16c1fe7b50 | ||
|
|
8caf879ffd | ||
|
|
5f48de037d | ||
|
|
73fa65c6e8 | ||
|
|
4842574cf3 | ||
|
|
6113149c58 | ||
|
|
87fa86d314 | ||
|
|
ba66188d80 | ||
|
|
dbe6fe0f92 | ||
|
|
8a255053aa | ||
|
|
29da0a7ca6 | ||
|
|
1922050e87 | ||
|
|
8bd6005b5d | ||
|
|
4a216d8875 | ||
|
|
8043338e58 | ||
|
|
de8c2ae6e2 | ||
|
|
b785b88af6 | ||
|
|
96c322c17c | ||
|
|
c05f9f6a28 | ||
|
|
11743e830a | ||
|
|
0464acbdc4 | ||
|
|
64bd2f1ba0 | ||
|
|
1f2e1f0489 | ||
|
|
b942c1cf53 | ||
|
|
bd84e4fe06 | ||
|
|
bc340e2e19 | ||
|
|
b73b7307bd | ||
|
|
b74ff39438 | ||
|
|
35421649c9 | ||
|
|
b54c91848b | ||
|
|
19f819c142 | ||
|
|
68fc4ca050 | ||
|
|
ade0936c6f | ||
|
|
d072f9ee9c | ||
|
|
eaa0eee4ca | ||
|
|
b8bcda134c |
+13
@@ -0,0 +1,13 @@
|
|||||||
|
pipeline:
|
||||||
|
build:
|
||||||
|
image: eu.gcr.io/avatao-public/docker:el7
|
||||||
|
volumes:
|
||||||
|
- /etc/docker:/etc/docker:ro
|
||||||
|
- /root/.docker:/root/.docker:ro
|
||||||
|
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||||
|
commands:
|
||||||
|
- docker build --pull -t eu.gcr.io/avatao-challengestore/tutorial-framework:${DRONE_TAG} .
|
||||||
|
- docker push eu.gcr.io/avatao-challengestore/tutorial-framework:${DRONE_TAG}
|
||||||
|
when:
|
||||||
|
event: 'tag'
|
||||||
|
branch: refs/tags/egyptianmau-20*
|
||||||
@@ -1,4 +0,0 @@
|
|||||||
export PYENV_ROOT="$HOME/.pyenv"
|
|
||||||
export PATH="$PYENV_ROOT/bin:$PATH"
|
|
||||||
eval "$(pyenv init -)"
|
|
||||||
eval "$(pyenv virtualenv-init -)"
|
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
[TYPECHECK]
|
||||||
|
|
||||||
|
ignored-modules = zmq
|
||||||
|
max-line-length = 150
|
||||||
|
disable = missing-docstring, too-few-public-methods
|
||||||
+17
-16
@@ -11,7 +11,8 @@ RUN curl -sL https://deb.nodesource.com/setup_8.x | sudo -E bash -
|
|||||||
libzmq5 \
|
libzmq5 \
|
||||||
nginx \
|
nginx \
|
||||||
gettext-base &&\
|
gettext-base &&\
|
||||||
rm -rf /var/lib/apt/lists/*
|
rm -rf /var/lib/apt/lists/* &&\
|
||||||
|
ln -sf /bin/bash /bin/sh
|
||||||
|
|
||||||
COPY requirements.txt /tmp
|
COPY requirements.txt /tmp
|
||||||
RUN pip3 install -r /tmp/requirements.txt
|
RUN pip3 install -r /tmp/requirements.txt
|
||||||
@@ -26,7 +27,7 @@ ENV TFW_PUBLIC_PORT=8888 \
|
|||||||
|
|
||||||
EXPOSE ${TFW_PUBLIC_PORT}
|
EXPOSE ${TFW_PUBLIC_PORT}
|
||||||
|
|
||||||
ENV PYTHONPATH="/usr/local/lib/" \
|
ENV PYTHONPATH="/usr/local/lib" \
|
||||||
TFW_SUPERVISOR_HTTP_URI="http://localhost:${TFW_SUPERVISOR_HTTP_PORT}" \
|
TFW_SUPERVISOR_HTTP_URI="http://localhost:${TFW_SUPERVISOR_HTTP_PORT}" \
|
||||||
TFW_SUPERVISORD_CONF="/etc/supervisor/supervisord.conf" \
|
TFW_SUPERVISORD_CONF="/etc/supervisor/supervisord.conf" \
|
||||||
TFW_SUPERVISORD_COMPONENTS="/etc/supervisor/conf" \
|
TFW_SUPERVISORD_COMPONENTS="/etc/supervisor/conf" \
|
||||||
@@ -39,30 +40,30 @@ ENV PYTHONPATH="/usr/local/lib/" \
|
|||||||
TFW_HISTFILE="/home/${AVATAO_USER}/.bash_history" \
|
TFW_HISTFILE="/home/${AVATAO_USER}/.bash_history" \
|
||||||
PROMPT_COMMAND="history -a"
|
PROMPT_COMMAND="history -a"
|
||||||
|
|
||||||
RUN echo "shopt -s cmdhist\n" \
|
COPY bashrc /tmp
|
||||||
"shopt -s histappend\n" \
|
RUN echo "export HISTFILE=${TFW_HISTFILE}" >> /tmp/bashrc &&\
|
||||||
"unset HISTCONTROL\n" \
|
cat /tmp/bashrc >> /home/${AVATAO_USER}/.bashrc
|
||||||
"export HISTFILE=$TFW_HISTFILE\n" \
|
|
||||||
"export HISTFILESIZE=1000\n" \
|
|
||||||
"export HISTSIZE=1000\n" \
|
|
||||||
'PROMPT_COMMAND="history -a"\n' \
|
|
||||||
>> /home/${AVATAO_USER}/.bashrc
|
|
||||||
|
|
||||||
COPY supervisor/supervisord.conf ${TFW_SUPERVISORD_CONF}
|
COPY supervisor/supervisord.conf ${TFW_SUPERVISORD_CONF}
|
||||||
COPY nginx/nginx.conf ${TFW_NGINX_CONF}
|
COPY nginx/nginx.conf ${TFW_NGINX_CONF}
|
||||||
COPY nginx/default.conf ${TFW_NGINX_DEFAULT}
|
COPY nginx/default.conf ${TFW_NGINX_DEFAULT}
|
||||||
COPY lib ${TFW_LIB_DIR}
|
COPY nginx/components/ ${TFW_NGINX_COMPONENTS}
|
||||||
|
COPY lib LICENSE ${TFW_LIB_DIR}
|
||||||
|
|
||||||
ONBUILD ARG BUILD_CONTEXT="."
|
RUN for dir in "${TFW_LIB_DIR}"/{tfw,tao,envvars.py} "/etc/nginx" "/etc/supervisor"; do \
|
||||||
|
chown -R root:root "$dir" && chmod -R 700 "$dir"; \
|
||||||
|
done
|
||||||
|
|
||||||
|
ONBUILD ARG BUILD_CONTEXT="solvable"
|
||||||
ONBUILD ARG NOFRONTEND=""
|
ONBUILD ARG NOFRONTEND=""
|
||||||
|
|
||||||
ONBUILD COPY ${BUILD_CONTEXT}/nginx/components/ ${TFW_NGINX_COMPONENTS}
|
ONBUILD COPY ${BUILD_CONTEXT}/nginx/ ${TFW_NGINX_COMPONENTS}
|
||||||
ONBUILD COPY ${BUILD_CONTEXT}/supervisor/components/ ${TFW_SUPERVISORD_COMPONENTS}
|
ONBUILD COPY ${BUILD_CONTEXT}/supervisor/ ${TFW_SUPERVISORD_COMPONENTS}
|
||||||
|
|
||||||
ONBUILD RUN chown -R ${AVATAO_USER} /var/log/nginx /var/lib/nginx &&\
|
ONBUILD RUN for f in "${TFW_NGINX_DEFAULT}" ${TFW_NGINX_COMPONENTS}/*.conf; do \
|
||||||
for f in "${TFW_NGINX_DEFAULT}" ${TFW_NGINX_COMPONENTS}/*.conf; do \
|
|
||||||
envsubst "$(printenv | cut -d= -f1 | grep TFW_ | sed -e 's/^/$/g')" < $f > $f~ && mv $f~ $f ;\
|
envsubst "$(printenv | cut -d= -f1 | grep TFW_ | sed -e 's/^/$/g')" < $f > $f~ && mv $f~ $f ;\
|
||||||
done
|
done
|
||||||
|
ONBUILD VOLUME ["/etc/nginx", "/var/lib/nginx", "/var/log/nginx"]
|
||||||
|
|
||||||
ONBUILD COPY ${BUILD_CONTEXT}/frontend /data/
|
ONBUILD COPY ${BUILD_CONTEXT}/frontend /data/
|
||||||
ONBUILD RUN test -z "${NOFRONTEND}" && cd /data && yarn install --frozen-lockfile || :
|
ONBUILD RUN test -z "${NOFRONTEND}" && cd /data && yarn install --frozen-lockfile || :
|
||||||
|
|||||||
@@ -0,0 +1,12 @@
|
|||||||
|
AVATAO CONFIDENTIAL
|
||||||
|
|
||||||
|
Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
All Rights Reserved.
|
||||||
|
|
||||||
|
All source code, configuration files and documentation contained herein
|
||||||
|
is, and remains the exclusive property of Avatao.com Innovative Learning Kft.
|
||||||
|
The intellectual and technical concepts contained herein are proprietary
|
||||||
|
to Avatao.com Innovative Learning Kft. and are protected by trade secret
|
||||||
|
or copyright law. Dissemination of this information or reproduction of
|
||||||
|
this material is strictly forbidden unless prior written permission is
|
||||||
|
obtained from Avatao.com Innovative Learning Kft.
|
||||||
@@ -1,16 +1,83 @@
|
|||||||
Tutorial Framework
|
# baseimage-tutorial-framework
|
||||||
==================
|
|
||||||
|
|
||||||
## Building and running with magic starter
|
This is the beating heart of TFW – the Docker baseimage containing the internals of the framework.
|
||||||
|
|
||||||
Execute the `tfw_magic_start.sh` script and it will automagically run a backend instance in Docker and serve the Angular frontend locally.
|
Every tutorial-framework based challenge has a `solvable` Docker image based on this one: their `Dockerfile`s begin with `FROM eu.gcr.io/avatao-challengestore/tutorial-framework`.
|
||||||
|
Note that TFW is not avaliable on Docker Hub due to legal reasons and is only accessible through local builds (don't worry, we've got you covered with build scripts in the [test-tutorial-framework](https://github.com/avatao-content/test-tutorial-framework) repo).
|
||||||
|
|
||||||
This requires the `test-tutorial-framework` (Docker image `FROM baseimage-tutorial-framework`) and `frontend-tutorial-framework` repositories to have a common parent folder.
|
This document explains the general concepts of TFW and should be the first thing you read before getting started with development.
|
||||||
|
|
||||||
## Building with Docker
|
For more on building and running you should check the [test-tutorial-framework](https://github.com/avatao-content/test-tutorial-framework) repo.
|
||||||
|
|
||||||
Simply issue `docker build -t baseimage-tutorial-framework .` in the project root.
|
## The framework
|
||||||
|
|
||||||
## Creating child images
|
The goal of the tutorial-framework is to help content developers in creating interactive tutorials for the Avatao platform.
|
||||||
|
|
||||||
Documentation in progress...
|
To make this possible TFW implements a robust messaging system and provides several pre-written components built upon it, such as a file editor and a terminal (both running in your browser).
|
||||||
|
|
||||||
|
The foundation of the whole framework is the messaging system connecting the frontend with the backend.
|
||||||
|
Frontend components use websockets to connect to the TFW server, to which you can hook several *event handlers* defining how to handle specific messages.
|
||||||
|
|
||||||
|

|
||||||
|
|
||||||
|
### Event handlers
|
||||||
|
|
||||||
|
Imagine event handlers as callbacks that are invoked when TFW receives a specific type of message. For instance, you could send a message to the framework when the user does something of note.
|
||||||
|
|
||||||
|
Event handler allow you to define actions triggered on the backend when the user presses a button on the frontend or moves the cursor to a specific area, etc.
|
||||||
|
|
||||||
|
Event handlers use ZeroMQ to connect to the framework. Due to this they are as loosely-coupled as possible: usually they are running in separate processes and only communicate with TFW through ZMQ.
|
||||||
|
|
||||||
|
Our pre-made event handlers are written in Python3, but you can write event handlers in any language that has ZeroMQ bindings (this means virtually any language).
|
||||||
|
|
||||||
|
This makes the framework really flexible: you can demonstrate the concepts you want to in any language while using the same set of tools provided by TFW.
|
||||||
|
Inside Avatao this means that any of the content teams can use the framework with ease.
|
||||||
|
|
||||||
|
### FSM
|
||||||
|
|
||||||
|
Another unique feature of the framework is the FSM – finite state machine – representing the state of your challenge.
|
||||||
|
This allows you to track users progressing with the tasks you've defined for them to complete.
|
||||||
|
|
||||||
|
For instance, you could represent whether the user managed to create a malicious user with a state called `user_registered` and subscribe callbacks to events regarding that state (like entering or leaving).
|
||||||
|
|
||||||
|
You could create challenges that can be completed in several different ways: imagine a state called `challenge_complete`, which indicates if the challenge is completed. Several series of actions (triggers) could lead to this state.
|
||||||
|
|
||||||
|
This enables you to guide your users through the experience you've envisioned with your tutorial.
|
||||||
|
We can provide a whole new level of interactivity in our challenges because we know what the user is doing.
|
||||||
|
This includes context-dependent hints and the automatic typing of commands to a terminal.
|
||||||
|
|
||||||
|
### Frontend
|
||||||
|
|
||||||
|
Note that our frontend implementation is written in Angular. It is maintained and documented in the [frontend-tutorial-framework](https://github.com/avatao-content/frontend-tutorial-framework) repository.
|
||||||
|
|
||||||
|
### Messaging format
|
||||||
|
|
||||||
|
The framework uses JSON messages internally and in exposed APIs as well.
|
||||||
|
These messages must comply with some rules.
|
||||||
|
Don't worry, we are not too fond of rules around these parts.
|
||||||
|
|
||||||
|
The TFW message format:
|
||||||
|
|
||||||
|
```text
|
||||||
|
{
|
||||||
|
"key: "some identifier used for addressing",
|
||||||
|
"data":
|
||||||
|
{
|
||||||
|
...
|
||||||
|
JSON object carrying anything, preferably cats
|
||||||
|
...
|
||||||
|
},
|
||||||
|
"trigger": "FSM action"
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
- The `key` field is used by TFW for addressing and every message must have one (it can be an empty string though)
|
||||||
|
- The `data` object can contain anything you might want to send
|
||||||
|
- The `trigger` key is an optional field that triggers an FSM action with that name from the current state (whatever that might be)
|
||||||
|
|
||||||
|
|
||||||
|
## Where to go next
|
||||||
|
|
||||||
|
Most of the components you need have docstrings included (hang on tight, this is work in progress) – refer to them for usage info.
|
||||||
|
|
||||||
|
To get started you should take a look at the [test-tutorial-framework](https://github.com/avatao-content/test-tutorial-framework) repository, which serves as an example project as well.
|
||||||
|
|||||||
@@ -0,0 +1,20 @@
|
|||||||
|
_postcmd()
|
||||||
|
{
|
||||||
|
AT_PROMPT=1
|
||||||
|
history -a
|
||||||
|
}
|
||||||
|
PROMPT_COMMAND="_postcmd"
|
||||||
|
|
||||||
|
_precmd()
|
||||||
|
{
|
||||||
|
[ -z "$AT_PROMPT" ] && return
|
||||||
|
unset AT_PROMPT
|
||||||
|
history -a
|
||||||
|
}
|
||||||
|
trap _precmd DEBUG
|
||||||
|
|
||||||
|
shopt -s cmdhist
|
||||||
|
shopt -s histappend
|
||||||
|
unset HISTCONTROL
|
||||||
|
export HISTFILESIZE=1000
|
||||||
|
export HISTSIZE=1000
|
||||||
Binary file not shown.
|
After Width: | Height: | Size: 40 KiB |
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from collections import namedtuple
|
from collections import namedtuple
|
||||||
from os import environ
|
from os import environ
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,2 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|||||||
@@ -1 +1,4 @@
|
|||||||
from .envvars import * # pylint: disable=wildcard-import
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .envvars import TAOENV
|
||||||
|
|||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from envvars import prefixed_envvars_to_namedtuple
|
from envvars import prefixed_envvars_to_namedtuple
|
||||||
|
|
||||||
TAOENV = prefixed_envvars_to_namedtuple('AVATAO_', 'taoenvtuple')
|
TAOENV = prefixed_envvars_to_namedtuple('AVATAO_', 'taoenvtuple')
|
||||||
|
|||||||
@@ -0,0 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .event_handler_base import EventHandlerBase, TriggeredEventHandler
|
||||||
|
from .fsm_base import FSMBase
|
||||||
|
from .linear_fsm import LinearFSM
|
||||||
|
|||||||
@@ -0,0 +1,9 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .directory_monitoring_event_handler import DirectoryMonitoringEventHandler
|
||||||
|
from .process_managing_event_handler import ProcessManagingEventHandler
|
||||||
|
from .terminal_event_handler import TerminalEventHandler
|
||||||
|
from .ide_event_handler import IdeEventHandler
|
||||||
|
from .history_monitor import HistoryMonitor, BashMonitor, GDBMonitor
|
||||||
|
from .terminal_commands import TerminalCommands
|
||||||
@@ -1 +0,0 @@
|
|||||||
from .rate_limiter import RateLimiter
|
|
||||||
@@ -1,30 +1,26 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from functools import wraps
|
from functools import wraps
|
||||||
|
|
||||||
from watchdog.observers import Observer
|
|
||||||
from watchdog.events import FileSystemEventHandler as FileSystemWatchdogEventHandler
|
from watchdog.events import FileSystemEventHandler as FileSystemWatchdogEventHandler
|
||||||
|
|
||||||
from tfw.networking.event_handlers.server_connector import ServerUplinkConnector
|
from tfw.networking.event_handlers import ServerUplinkConnector
|
||||||
from tfw.components.decorators import RateLimiter
|
from tfw.decorators import RateLimiter
|
||||||
|
from tfw.mixins import ObserverMixin
|
||||||
|
|
||||||
from tfw.config.logs import logging
|
from tfw.config.logs import logging
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class DirectoryMonitor:
|
class DirectoryMonitor(ObserverMixin):
|
||||||
def __init__(self, directory):
|
def __init__(self, directory):
|
||||||
self.observer = Observer()
|
ObserverMixin.__init__(self)
|
||||||
self.eventhandler = WebideReloadWatchdogEventHandler()
|
self.eventhandler = IdeReloadWatchdogEventHandler()
|
||||||
self.observer.schedule(self.eventhandler, directory, recursive=True)
|
self.observer.schedule(self.eventhandler, directory, recursive=True)
|
||||||
self.pause, self.resume = self.eventhandler.pause, self.eventhandler.resume
|
self.pause, self.resume = self.eventhandler.pause, self.eventhandler.resume
|
||||||
|
|
||||||
def watch(self):
|
|
||||||
self.observer.start()
|
|
||||||
|
|
||||||
def stop(self):
|
|
||||||
self.observer.stop()
|
|
||||||
self.observer.join()
|
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def ignore(self):
|
def ignore(self):
|
||||||
return self.eventhandler.ignore
|
return self.eventhandler.ignore
|
||||||
@@ -46,7 +42,7 @@ class DirectoryMonitor:
|
|||||||
self.directorymonitor.resume()
|
self.directorymonitor.resume()
|
||||||
|
|
||||||
|
|
||||||
class WebideReloadWatchdogEventHandler(FileSystemWatchdogEventHandler):
|
class IdeReloadWatchdogEventHandler(FileSystemWatchdogEventHandler):
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
super().__init__()
|
super().__init__()
|
||||||
self.uplink = ServerUplinkConnector()
|
self.uplink = ServerUplinkConnector()
|
||||||
@@ -67,8 +63,8 @@ class WebideReloadWatchdogEventHandler(FileSystemWatchdogEventHandler):
|
|||||||
self.ignore = self.ignore - 1
|
self.ignore = self.ignore - 1
|
||||||
return
|
return
|
||||||
LOG.debug(event)
|
LOG.debug(event)
|
||||||
key = 'webide'
|
self.uplink.send({'key': 'ide',
|
||||||
self.uplink.send(key, {'data': {'command': 'reload'}})
|
'data': {'command': 'reload'}})
|
||||||
|
|
||||||
|
|
||||||
def with_monitor_paused(fun):
|
def with_monitor_paused(fun):
|
||||||
|
|||||||
@@ -1,18 +1,22 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from os.path import isdir, exists
|
from os.path import isdir, exists
|
||||||
|
|
||||||
from tfw.components.directory_monitor import DirectoryMonitor
|
from tfw import EventHandlerBase
|
||||||
from tfw.event_handler_base import TriggerlessEventHandler
|
|
||||||
from tfw.config.logs import logging
|
from tfw.config.logs import logging
|
||||||
|
from tfw.mixins import MonitorManagerMixin
|
||||||
|
from .directory_monitor import DirectoryMonitor
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class DirectoryMonitoringEventHandler(TriggerlessEventHandler):
|
class DirectoryMonitoringEventHandler(EventHandlerBase, MonitorManagerMixin):
|
||||||
def __init__(self, key, directory):
|
def __init__(self, key, directory):
|
||||||
super().__init__(key)
|
super().__init__(key)
|
||||||
self._directory = directory
|
self._directory = directory
|
||||||
self._monitor = None
|
MonitorManagerMixin.__init__(self, DirectoryMonitor, self._directory)
|
||||||
self.reload_monitor()
|
|
||||||
self.commands = {'pause': self.pause,
|
self.commands = {'pause': self.pause,
|
||||||
'resume': self.resume,
|
'resume': self.resume,
|
||||||
'ignore': self.ignore,
|
'ignore': self.ignore,
|
||||||
@@ -28,20 +32,7 @@ class DirectoryMonitoringEventHandler(TriggerlessEventHandler):
|
|||||||
raise EnvironmentError('No such directory!')
|
raise EnvironmentError('No such directory!')
|
||||||
self._directory = directory
|
self._directory = directory
|
||||||
|
|
||||||
@property
|
def handle_event(self, message):
|
||||||
def monitor(self):
|
|
||||||
return self._monitor
|
|
||||||
|
|
||||||
def reload_monitor(self):
|
|
||||||
if self._monitor:
|
|
||||||
try:
|
|
||||||
self._monitor.stop()
|
|
||||||
except KeyError:
|
|
||||||
logging.debug('Working directory was removed – ignoring...')
|
|
||||||
self._monitor = DirectoryMonitor(self._directory)
|
|
||||||
self._monitor.watch() # This runs on a separate thread
|
|
||||||
|
|
||||||
def handle_event(self, key, message):
|
|
||||||
try:
|
try:
|
||||||
message['data'] = self.commands[message['data']['command']](message['data'])
|
message['data'] = self.commands[message['data']['command']](message['data'])
|
||||||
return message
|
return message
|
||||||
|
|||||||
@@ -1,13 +1,15 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from os.path import dirname
|
from os.path import dirname
|
||||||
from re import findall
|
from re import findall
|
||||||
from re import compile as compileregex
|
from re import compile as compileregex
|
||||||
from abc import ABC, abstractmethod
|
from abc import ABC, abstractmethod
|
||||||
|
|
||||||
from watchdog.observers import Observer
|
|
||||||
from watchdog.events import PatternMatchingEventHandler
|
from watchdog.events import PatternMatchingEventHandler
|
||||||
|
|
||||||
from tfw.components.mixins.callback_mixin import CallbackMixin
|
from tfw.mixins import CallbackMixin, ObserverMixin
|
||||||
from tfw.components.decorators.rate_limiter import RateLimiter
|
from tfw.decorators import RateLimiter
|
||||||
|
|
||||||
|
|
||||||
class CallbackEventHandler(PatternMatchingEventHandler, ABC):
|
class CallbackEventHandler(PatternMatchingEventHandler, ABC):
|
||||||
@@ -21,13 +23,24 @@ class CallbackEventHandler(PatternMatchingEventHandler, ABC):
|
|||||||
callback()
|
callback()
|
||||||
|
|
||||||
|
|
||||||
class HistoryMonitor(CallbackMixin, ABC):
|
class HistoryMonitor(CallbackMixin, ObserverMixin, ABC):
|
||||||
|
"""
|
||||||
|
Abstract class capable of monitoring and parsing a history file such as
|
||||||
|
bash HISTFILEs. Monitoring means detecting when the file was changed and
|
||||||
|
notifying subscribers about new content in the file.
|
||||||
|
|
||||||
|
This is useful for monitoring CLI sessions.
|
||||||
|
|
||||||
|
To specify a custom HistoryMonitor inherit from this class and override the
|
||||||
|
command pattern property and optionally the sanitize_command method.
|
||||||
|
See examples below.
|
||||||
|
"""
|
||||||
def __init__(self, histfile):
|
def __init__(self, histfile):
|
||||||
CallbackMixin.__init__(self)
|
CallbackMixin.__init__(self)
|
||||||
|
ObserverMixin.__init__(self)
|
||||||
self.histfile = histfile
|
self.histfile = histfile
|
||||||
self._history = []
|
self._history = []
|
||||||
self._last_length = len(self._history)
|
self._last_length = len(self._history)
|
||||||
self.observer = Observer()
|
|
||||||
self.observer.schedule(CallbackEventHandler([self.histfile],
|
self.observer.schedule(CallbackEventHandler([self.histfile],
|
||||||
self._fetch_history,
|
self._fetch_history,
|
||||||
self._invoke_callbacks),
|
self._invoke_callbacks),
|
||||||
@@ -57,15 +70,17 @@ class HistoryMonitor(CallbackMixin, ABC):
|
|||||||
if self._last_length < len(self._history):
|
if self._last_length < len(self._history):
|
||||||
self._execute_callbacks(self.history)
|
self._execute_callbacks(self.history)
|
||||||
|
|
||||||
def watch(self):
|
|
||||||
self.observer.start()
|
|
||||||
|
|
||||||
def stop(self):
|
|
||||||
self.observer.stop()
|
|
||||||
self.observer.join()
|
|
||||||
|
|
||||||
|
|
||||||
class BashMonitor(HistoryMonitor):
|
class BashMonitor(HistoryMonitor):
|
||||||
|
"""
|
||||||
|
HistoryMonitor for monitoring bash CLI sessions.
|
||||||
|
This requires the following to be set in bash
|
||||||
|
(note that this is done automatically by TFW):
|
||||||
|
PROMPT_COMMAND="history -a"
|
||||||
|
shopt -s cmdhist
|
||||||
|
shopt -s histappend
|
||||||
|
unset HISTCONTROL
|
||||||
|
"""
|
||||||
@property
|
@property
|
||||||
def command_pattern(self):
|
def command_pattern(self):
|
||||||
return r'.+'
|
return r'.+'
|
||||||
@@ -75,6 +90,10 @@ class BashMonitor(HistoryMonitor):
|
|||||||
|
|
||||||
|
|
||||||
class GDBMonitor(HistoryMonitor):
|
class GDBMonitor(HistoryMonitor):
|
||||||
|
"""
|
||||||
|
HistoryMonitor to monitor GDB sessions.
|
||||||
|
For this to work "set trace-commands on" must be set in GDB.
|
||||||
|
"""
|
||||||
@property
|
@property
|
||||||
def command_pattern(self):
|
def command_pattern(self):
|
||||||
return r'(?<=\n)\+(.+)\n'
|
return r'(?<=\n)\+(.+)\n'
|
||||||
|
|||||||
@@ -0,0 +1,223 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from os.path import isfile, join, relpath, exists, isdir, realpath
|
||||||
|
from glob import glob
|
||||||
|
from fnmatch import fnmatchcase
|
||||||
|
from collections import Iterable
|
||||||
|
|
||||||
|
from tfw import EventHandlerBase
|
||||||
|
from tfw.mixins import MonitorManagerMixin
|
||||||
|
from tfw.config.logs import logging
|
||||||
|
from .directory_monitor import DirectoryMonitor
|
||||||
|
|
||||||
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class FileManager: # pylint: disable=too-many-instance-attributes
|
||||||
|
def __init__(self, working_directory, allowed_directories, selected_file=None, exclude=None):
|
||||||
|
self._exclude, self.exclude = None, exclude
|
||||||
|
self._allowed_directories, self.allowed_directories = None, allowed_directories
|
||||||
|
self._workdir, self.workdir = None, working_directory
|
||||||
|
self._filename, self.filename = None, selected_file or self.files[0]
|
||||||
|
|
||||||
|
@property
|
||||||
|
def exclude(self):
|
||||||
|
return self._exclude
|
||||||
|
|
||||||
|
@exclude.setter
|
||||||
|
def exclude(self, exclude):
|
||||||
|
if exclude is None:
|
||||||
|
return
|
||||||
|
if not isinstance(exclude, Iterable):
|
||||||
|
raise TypeError('Exclude must be Iterable!')
|
||||||
|
self._exclude = exclude
|
||||||
|
|
||||||
|
@property
|
||||||
|
def workdir(self):
|
||||||
|
return self._workdir
|
||||||
|
|
||||||
|
@workdir.setter
|
||||||
|
def workdir(self, directory):
|
||||||
|
if not exists(directory) or not isdir(directory):
|
||||||
|
raise EnvironmentError(f'"{directory}" is not a directory!')
|
||||||
|
if not self._is_in_whitelisted_dir(directory):
|
||||||
|
raise EnvironmentError(f'Directory "{directory}" is not in whitelist!')
|
||||||
|
self._workdir = directory
|
||||||
|
|
||||||
|
@property
|
||||||
|
def allowed_directories(self):
|
||||||
|
return self._allowed_directories
|
||||||
|
|
||||||
|
@allowed_directories.setter
|
||||||
|
def allowed_directories(self, directories):
|
||||||
|
self._allowed_directories = directories
|
||||||
|
|
||||||
|
@property
|
||||||
|
def filename(self):
|
||||||
|
return self._filename
|
||||||
|
|
||||||
|
@filename.setter
|
||||||
|
def filename(self, filename):
|
||||||
|
if filename not in self.files:
|
||||||
|
raise EnvironmentError('No such file in workdir!')
|
||||||
|
self._filename = filename
|
||||||
|
|
||||||
|
@property
|
||||||
|
def files(self):
|
||||||
|
return [self._relpath(file) for file in glob(join(self._workdir, '**/*'), recursive=True)
|
||||||
|
if isfile(file) and self._is_in_whitelisted_dir(file) and not self._is_blacklisted(file)]
|
||||||
|
|
||||||
|
@property
|
||||||
|
def file_contents(self):
|
||||||
|
with open(self._filepath(self.filename), 'r', errors='surrogateescape') as ifile:
|
||||||
|
return ifile.read()
|
||||||
|
|
||||||
|
@file_contents.setter
|
||||||
|
def file_contents(self, value):
|
||||||
|
with open(self._filepath(self.filename), 'w', errors='surrogateescape') as ofile:
|
||||||
|
ofile.write(value)
|
||||||
|
|
||||||
|
def _is_in_whitelisted_dir(self, path):
|
||||||
|
return any(realpath(path).startswith(allowed_dir) for allowed_dir in self.allowed_directories)
|
||||||
|
|
||||||
|
def _is_blacklisted(self, file):
|
||||||
|
return any(fnmatchcase(file, blacklisted) for blacklisted in self.exclude)
|
||||||
|
|
||||||
|
def _filepath(self, filename):
|
||||||
|
return join(self._workdir, filename)
|
||||||
|
|
||||||
|
def _relpath(self, filename):
|
||||||
|
return relpath(self._filepath(filename), start=self._workdir)
|
||||||
|
|
||||||
|
|
||||||
|
class IdeEventHandler(EventHandlerBase, MonitorManagerMixin):
|
||||||
|
# pylint: disable=too-many-arguments
|
||||||
|
"""
|
||||||
|
Event handler implementing the backend of our browser based IDE.
|
||||||
|
By default all files in the directory specified in __init__ are displayed
|
||||||
|
on the fontend. Note that this is a stateful component.
|
||||||
|
|
||||||
|
This EventHandler accepts messages that have a data["command"] key specifying
|
||||||
|
a command to be executed.
|
||||||
|
The API of each command is documented in their respective handlers.
|
||||||
|
"""
|
||||||
|
def __init__(self, key, directory, allowed_directories, selected_file=None, exclude=None):
|
||||||
|
"""
|
||||||
|
:param key: the key this instance should listen to
|
||||||
|
:param directory: working directory which the EventHandler should serve files from
|
||||||
|
:param allowed_directories: list of directories that can be switched to using the selectdir command
|
||||||
|
:param selected_file: file that is selected by default
|
||||||
|
:param exclude: list of filenames that should not appear between files (for *.o, *.pyc, etc.)
|
||||||
|
"""
|
||||||
|
super().__init__(key)
|
||||||
|
try:
|
||||||
|
self.filemanager = FileManager(allowed_directories=allowed_directories, working_directory=directory,
|
||||||
|
selected_file=selected_file, exclude=exclude)
|
||||||
|
except IndexError:
|
||||||
|
raise EnvironmentError(f'No file(s) in IdeEventHandler working_directory "{directory}"!')
|
||||||
|
MonitorManagerMixin.__init__(self, DirectoryMonitor, self.filemanager.workdir)
|
||||||
|
|
||||||
|
self.commands = {'read': self.read,
|
||||||
|
'write': self.write,
|
||||||
|
'select': self.select,
|
||||||
|
'selectdir': self.select_dir,
|
||||||
|
'exclude': self.exclude}
|
||||||
|
|
||||||
|
def read(self, data):
|
||||||
|
"""
|
||||||
|
Read the currently selected file.
|
||||||
|
|
||||||
|
:return: message with the contents of the file in data['content']
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
data['content'] = self.filemanager.file_contents
|
||||||
|
except PermissionError:
|
||||||
|
data['content'] = 'You have no permission to open that file :('
|
||||||
|
except FileNotFoundError:
|
||||||
|
data['content'] = 'This file was removed :('
|
||||||
|
except Exception: # pylint: disable=broad-except
|
||||||
|
data['content'] = 'Failed to read file :('
|
||||||
|
return data
|
||||||
|
|
||||||
|
def write(self, data):
|
||||||
|
"""
|
||||||
|
Overwrites a file with the desired string.
|
||||||
|
|
||||||
|
:param data: TFW message data containing keys:
|
||||||
|
|-string: containing the desired file contents
|
||||||
|
"""
|
||||||
|
self.monitor.ignore = self.monitor.ignore + 1
|
||||||
|
try:
|
||||||
|
self.filemanager.file_contents = data['content']
|
||||||
|
except Exception: # pylint: disable=broad-except
|
||||||
|
LOG.exception('Error writing file!')
|
||||||
|
del data['content']
|
||||||
|
return data
|
||||||
|
|
||||||
|
def select(self, data):
|
||||||
|
"""
|
||||||
|
Selects a file from the current directory.
|
||||||
|
|
||||||
|
:param data: TFW message data containing keys:
|
||||||
|
|-filename: name of file to select relative to the current directory
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
self.filemanager.filename = data['filename']
|
||||||
|
except EnvironmentError:
|
||||||
|
LOG.exception('Failed to select file "%s"', data['filename'])
|
||||||
|
return data
|
||||||
|
|
||||||
|
def select_dir(self, data):
|
||||||
|
"""
|
||||||
|
Select a new working directory to display files from.
|
||||||
|
|
||||||
|
:param data: TFW message data containing keys:
|
||||||
|
|-directory: absolute path of diretory to select.
|
||||||
|
must be a path whitelisted in
|
||||||
|
self.allowed_directories
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
self.filemanager.workdir = data['directory']
|
||||||
|
self.reload_monitor()
|
||||||
|
try:
|
||||||
|
self.filemanager.filename = self.filemanager.files[0]
|
||||||
|
self.read(data)
|
||||||
|
except IndexError:
|
||||||
|
data['content'] = 'No files in this directory :('
|
||||||
|
except EnvironmentError as err:
|
||||||
|
LOG.error('Failed to select directory "%s". Reason: %s', data['directory'], str(err))
|
||||||
|
return data
|
||||||
|
|
||||||
|
def exclude(self, data):
|
||||||
|
"""
|
||||||
|
Overwrite list of excluded files
|
||||||
|
|
||||||
|
:param data: TFW message data containing keys:
|
||||||
|
|-exclude: list of filename patterns to be excluded, e.g.: ["*.pyc", "*.o"]
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
self.filemanager.exclude = list(data['exclude'])
|
||||||
|
except TypeError:
|
||||||
|
LOG.error('Exclude must be Iterable!')
|
||||||
|
return data
|
||||||
|
|
||||||
|
def attach_fileinfo(self, data):
|
||||||
|
"""
|
||||||
|
Basic information included in every response to the frontend.
|
||||||
|
"""
|
||||||
|
data['filename'] = self.filemanager.filename
|
||||||
|
data['files'] = self.filemanager.files
|
||||||
|
data['directory'] = self.filemanager.workdir
|
||||||
|
|
||||||
|
def handle_event(self, message):
|
||||||
|
try:
|
||||||
|
data = message['data']
|
||||||
|
message['data'] = self.commands[data['command']](data)
|
||||||
|
self.attach_fileinfo(data)
|
||||||
|
return message
|
||||||
|
except KeyError:
|
||||||
|
LOG.error('IGNORING MESSAGE: Invalid message received: %s', message)
|
||||||
|
|
||||||
|
def cleanup(self):
|
||||||
|
self.monitor.stop()
|
||||||
@@ -1 +0,0 @@
|
|||||||
from .supervisor_mixin import SupervisorMixin
|
|
||||||
@@ -1,9 +1,12 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from xmlrpc.client import Fault as SupervisorFault
|
from xmlrpc.client import Fault as SupervisorFault
|
||||||
|
|
||||||
from tfw.event_handler_base import TriggerlessEventHandler
|
from tfw import EventHandlerBase
|
||||||
from tfw.components.mixins import SupervisorMixin
|
from tfw.mixins import SupervisorMixin
|
||||||
from tfw.components.directory_monitor import with_monitor_paused
|
|
||||||
from tfw.config.logs import logging
|
from tfw.config.logs import logging
|
||||||
|
from .directory_monitor import with_monitor_paused
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -19,7 +22,19 @@ class ProcessManager(SupervisorMixin):
|
|||||||
return self.commands[command](process_name)
|
return self.commands[command](process_name)
|
||||||
|
|
||||||
|
|
||||||
class ProcessManagingEventHandler(TriggerlessEventHandler):
|
class ProcessManagingEventHandler(EventHandlerBase):
|
||||||
|
"""
|
||||||
|
Event handler that can manage processes managed by supervisor.
|
||||||
|
|
||||||
|
This EventHandler accepts messages that have a data["command"] key specifying
|
||||||
|
a command to be executed.
|
||||||
|
Every message must contain a data["process_name"] field with the name of the
|
||||||
|
process to manage. This is the name specified in supervisor config files like so:
|
||||||
|
[program:someprogram]
|
||||||
|
|
||||||
|
Commands available: start, stop, restart, readlog
|
||||||
|
(the names are as self-documenting as it gets)
|
||||||
|
"""
|
||||||
def __init__(self, key, dirmonitor=None):
|
def __init__(self, key, dirmonitor=None):
|
||||||
super().__init__(key)
|
super().__init__(key)
|
||||||
self.key = key
|
self.key = key
|
||||||
@@ -27,7 +42,7 @@ class ProcessManagingEventHandler(TriggerlessEventHandler):
|
|||||||
self.processmanager = ProcessManager()
|
self.processmanager = ProcessManager()
|
||||||
|
|
||||||
@with_monitor_paused
|
@with_monitor_paused
|
||||||
def handle_event(self, key, message):
|
def handle_event(self, message):
|
||||||
try:
|
try:
|
||||||
data = message['data']
|
data = message['data']
|
||||||
self.processmanager(data['command'], data['process_name'])
|
self.processmanager(data['command'], data['process_name'])
|
||||||
|
|||||||
@@ -1,163 +0,0 @@
|
|||||||
from os.path import isfile, join, relpath, exists, isdir
|
|
||||||
from glob import glob
|
|
||||||
from fnmatch import fnmatchcase
|
|
||||||
from collections import Iterable
|
|
||||||
|
|
||||||
from tfw.event_handler_base import TriggerlessEventHandler
|
|
||||||
from tfw.components.directory_monitor import DirectoryMonitor
|
|
||||||
from tfw.config.logs import logging
|
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
class FileManager:
|
|
||||||
def __init__(self, working_directory, selected_file=None, exclude=None):
|
|
||||||
self._exclude, self.exclude = None, exclude
|
|
||||||
self._workdir, self.workdir = None, working_directory
|
|
||||||
self._filename, self.filename = None, selected_file or self.files[0]
|
|
||||||
|
|
||||||
@property
|
|
||||||
def exclude(self):
|
|
||||||
return self._exclude
|
|
||||||
|
|
||||||
@exclude.setter
|
|
||||||
def exclude(self, exclude):
|
|
||||||
if exclude is None:
|
|
||||||
return
|
|
||||||
if not isinstance(exclude, Iterable):
|
|
||||||
raise TypeError('Exclude must be Iterable!')
|
|
||||||
self._exclude = exclude
|
|
||||||
|
|
||||||
@property
|
|
||||||
def workdir(self):
|
|
||||||
return self._workdir
|
|
||||||
|
|
||||||
@workdir.setter
|
|
||||||
def workdir(self, directory):
|
|
||||||
if not exists(directory) or not isdir(directory):
|
|
||||||
raise EnvironmentError('"{}" is not a directory!'.format(directory))
|
|
||||||
self._workdir = directory
|
|
||||||
|
|
||||||
@property
|
|
||||||
def filename(self):
|
|
||||||
return self._filename
|
|
||||||
|
|
||||||
@filename.setter
|
|
||||||
def filename(self, filename):
|
|
||||||
if not filename in self.files:
|
|
||||||
raise EnvironmentError('No such file in workdir!')
|
|
||||||
self._filename = filename
|
|
||||||
|
|
||||||
@property
|
|
||||||
def files(self):
|
|
||||||
return [self._relpath(file) for file in glob(join(self._workdir, '**/*'), recursive=True)
|
|
||||||
if isfile(file) and
|
|
||||||
not any(fnmatchcase(file, blacklisted) for blacklisted in self.exclude)]
|
|
||||||
|
|
||||||
@property
|
|
||||||
def file_contents(self):
|
|
||||||
with open(self._filepath(self.filename), 'r', errors='surrogateescape') as ifile:
|
|
||||||
return ifile.read()
|
|
||||||
|
|
||||||
@file_contents.setter
|
|
||||||
def file_contents(self, value):
|
|
||||||
with open(self._filepath(self.filename), 'w', errors='surrogateescape') as ofile:
|
|
||||||
ofile.write(value)
|
|
||||||
|
|
||||||
def _filepath(self, filename):
|
|
||||||
return join(self._workdir, filename)
|
|
||||||
|
|
||||||
def _relpath(self, filename):
|
|
||||||
return relpath(self._filepath(filename), start=self._workdir)
|
|
||||||
|
|
||||||
|
|
||||||
class SourceCodeEventHandler(TriggerlessEventHandler):
|
|
||||||
def __init__(self, key, directory, selected_file=None, exclude=None):
|
|
||||||
super().__init__(key)
|
|
||||||
self.filemanager = FileManager(directory, selected_file=selected_file, exclude=exclude)
|
|
||||||
|
|
||||||
self.commands = {'read': self.read,
|
|
||||||
'write': self.write,
|
|
||||||
'select': self.select,
|
|
||||||
'selectdir': self.select_dir,
|
|
||||||
'exclude': self.exclude}
|
|
||||||
|
|
||||||
self._monitor = None
|
|
||||||
self.reload_monitor()
|
|
||||||
|
|
||||||
@property
|
|
||||||
def monitor(self):
|
|
||||||
return self._monitor
|
|
||||||
|
|
||||||
def reload_monitor(self):
|
|
||||||
if self._monitor:
|
|
||||||
try:
|
|
||||||
self._monitor.stop()
|
|
||||||
except KeyError:
|
|
||||||
logging.debug('Working directory was removed – ignoring...')
|
|
||||||
self._monitor = DirectoryMonitor(self.filemanager.workdir)
|
|
||||||
self._monitor.watch() # This runs on a separate thread
|
|
||||||
|
|
||||||
def read(self, data):
|
|
||||||
try:
|
|
||||||
data['content'] = self.filemanager.file_contents
|
|
||||||
except PermissionError:
|
|
||||||
data['content'] = 'You have no permission to open that file :('
|
|
||||||
except FileNotFoundError:
|
|
||||||
data['content'] = 'This file was removed :('
|
|
||||||
except Exception: # pylint: disable=broad-except
|
|
||||||
data['content'] = 'Failed to read file :('
|
|
||||||
return data
|
|
||||||
|
|
||||||
def write(self, data):
|
|
||||||
self.monitor.ignore = self.monitor.ignore + 1
|
|
||||||
try:
|
|
||||||
self.filemanager.file_contents = data['content']
|
|
||||||
except Exception: # pylint: disable=broad-except
|
|
||||||
LOG.exception('Error writing file!')
|
|
||||||
del data['content']
|
|
||||||
return data
|
|
||||||
|
|
||||||
def select(self, data):
|
|
||||||
try:
|
|
||||||
self.filemanager.filename = data['filename']
|
|
||||||
except EnvironmentError:
|
|
||||||
LOG.exception('Failed to select file "%s"', data['filename'])
|
|
||||||
return data
|
|
||||||
|
|
||||||
def select_dir(self, data):
|
|
||||||
try:
|
|
||||||
self.filemanager.workdir = data['directory']
|
|
||||||
self.reload_monitor()
|
|
||||||
try:
|
|
||||||
self.filemanager.filename = self.filemanager.files[0]
|
|
||||||
self.read(data)
|
|
||||||
except IndexError:
|
|
||||||
data['content'] = 'No files in this directory :('
|
|
||||||
except EnvironmentError:
|
|
||||||
LOG.error('Failed to select directory "%s"', data['directory'])
|
|
||||||
return data
|
|
||||||
|
|
||||||
def exclude(self, data):
|
|
||||||
try:
|
|
||||||
self.filemanager.exclude = list(data['exclude'])
|
|
||||||
except TypeError:
|
|
||||||
LOG.error('Exclude must be Iterable!')
|
|
||||||
return data
|
|
||||||
|
|
||||||
def attach_fileinfo(self, data):
|
|
||||||
data['filename'] = self.filemanager.filename
|
|
||||||
data['files'] = self.filemanager.files
|
|
||||||
data['directory'] = self.filemanager.workdir
|
|
||||||
|
|
||||||
def handle_event(self, key, message):
|
|
||||||
try:
|
|
||||||
data = message['data']
|
|
||||||
message['data'] = self.commands[data['command']](data)
|
|
||||||
self.attach_fileinfo(data)
|
|
||||||
return message
|
|
||||||
except KeyError:
|
|
||||||
LOG.error('IGNORING MESSAGE: Invalid message received: %s', message)
|
|
||||||
|
|
||||||
def cleanup(self):
|
|
||||||
self.monitor.stop()
|
|
||||||
@@ -1,44 +0,0 @@
|
|||||||
from tfw.components.terminado_mini_server import TerminadoMiniServer
|
|
||||||
from tfw.event_handler_base import TriggerlessEventHandler
|
|
||||||
from tfw.config import TFWENV
|
|
||||||
from tfw.config.logs import logging
|
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
class TerminadoEventHandler(TriggerlessEventHandler):
|
|
||||||
def __init__(self, key, monitor):
|
|
||||||
super().__init__(key)
|
|
||||||
self.working_directory = TFWENV.TERMINADO_DIR
|
|
||||||
self._historymonitor = monitor
|
|
||||||
self.terminado_server = TerminadoMiniServer('/terminal', TFWENV.TERMINADO_PORT, TFWENV.TERMINADO_WD, ['bash'])
|
|
||||||
self.commands = {'write': self.write,
|
|
||||||
'read': self.read}
|
|
||||||
if self._historymonitor:
|
|
||||||
self._historymonitor.watch()
|
|
||||||
self.terminado_server.listen()
|
|
||||||
|
|
||||||
@property
|
|
||||||
def historymonitor(self):
|
|
||||||
return self._historymonitor
|
|
||||||
|
|
||||||
def handle_event(self, key, message):
|
|
||||||
LOG.debug('TerminadoEventHandler received event: %s', message)
|
|
||||||
try:
|
|
||||||
message['data'] = self.commands[message['data']['command']](message['data'])
|
|
||||||
return message
|
|
||||||
except KeyError:
|
|
||||||
LOG.error('IGNORING MESSAGE: Invalid message received: %s', message)
|
|
||||||
|
|
||||||
def write(self, data):
|
|
||||||
self.terminado_server.pty.write(data['shellcmd'])
|
|
||||||
|
|
||||||
def read(self, data):
|
|
||||||
data['count'] = int(data.get('count', 1))
|
|
||||||
if self.historymonitor:
|
|
||||||
data['history'] = self.historymonitor.history[-data['count']:]
|
|
||||||
return data
|
|
||||||
|
|
||||||
def cleanup(self):
|
|
||||||
if self.historymonitor:
|
|
||||||
self.historymonitor.stop()
|
|
||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from tornado.ioloop import IOLoop
|
from tornado.ioloop import IOLoop
|
||||||
from tornado.web import Application
|
from tornado.web import Application
|
||||||
from terminado import TermSocket, SingleTermManager
|
from terminado import TermSocket, SingleTermManager
|
||||||
|
|||||||
@@ -0,0 +1,66 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from abc import ABC
|
||||||
|
from re import match
|
||||||
|
|
||||||
|
from tfw.config.logs import logging
|
||||||
|
|
||||||
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class TerminalCommands(ABC):
|
||||||
|
"""
|
||||||
|
A class you can use to define hooks for terminal commands. This means that you can
|
||||||
|
have python code executed when the user enters a specific command to the terminal on
|
||||||
|
our frontend.
|
||||||
|
|
||||||
|
To receive events you need to subscribe TerminalCommand.callback to a HistoryMonitor
|
||||||
|
instance.
|
||||||
|
|
||||||
|
Inherit from this class and define methods which start with "command_". When the user
|
||||||
|
executes the command specified after the underscore, your method will be invoked. All
|
||||||
|
such commands must expect the parameter *args which will contain the arguments of the
|
||||||
|
command.
|
||||||
|
|
||||||
|
For example to define a method that runs when someone starts vim in the terminal
|
||||||
|
you have to define a method like: "def command_vim(self, *args)"
|
||||||
|
|
||||||
|
You can also use this class to create new commands similarly.
|
||||||
|
"""
|
||||||
|
def __init__(self, bashrc=None):
|
||||||
|
self._command_method_regex = r'^command_(.+)$'
|
||||||
|
self.command_implemetations = self._build_command_to_implementation_dict()
|
||||||
|
if bashrc is not None:
|
||||||
|
self._setup_bashrc_aliases(bashrc)
|
||||||
|
|
||||||
|
def _build_command_to_implementation_dict(self):
|
||||||
|
return {self._parse_command_name(fun): getattr(self, fun) for fun in dir(self)
|
||||||
|
if callable(getattr(self, fun)) and self._is_command_implementation(fun)}
|
||||||
|
|
||||||
|
def _setup_bashrc_aliases(self, bashrc):
|
||||||
|
with open(bashrc, 'a') as ofile:
|
||||||
|
alias_template = 'type {0} &> /dev/null || alias {0}="{0} &> /dev/null"\n'
|
||||||
|
for command in self.command_implemetations.keys():
|
||||||
|
ofile.write(alias_template.format(command))
|
||||||
|
|
||||||
|
def _is_command_implementation(self, method_name):
|
||||||
|
return bool(self._match_command_regex(method_name))
|
||||||
|
|
||||||
|
def _parse_command_name(self, method_name):
|
||||||
|
try:
|
||||||
|
return self._match_command_regex(method_name).groups()[0]
|
||||||
|
except AttributeError:
|
||||||
|
return ''
|
||||||
|
|
||||||
|
def _match_command_regex(self, string):
|
||||||
|
return match(self._command_method_regex, string)
|
||||||
|
|
||||||
|
def callback(self, history):
|
||||||
|
parts = history[-1].split()
|
||||||
|
command = parts[0]
|
||||||
|
if command in self.command_implemetations.keys():
|
||||||
|
try:
|
||||||
|
self.command_implemetations[command](*parts[1:])
|
||||||
|
except IndexError:
|
||||||
|
LOG.debug('Command "%s" failed!', command)
|
||||||
@@ -0,0 +1,77 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from tfw import EventHandlerBase
|
||||||
|
from tfw.config import TFWENV
|
||||||
|
from tfw.config.logs import logging
|
||||||
|
from tao.config import TAOENV
|
||||||
|
from .terminado_mini_server import TerminadoMiniServer
|
||||||
|
|
||||||
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class TerminalEventHandler(EventHandlerBase):
|
||||||
|
"""
|
||||||
|
Event handler responsible for managing terminal sessions for frontend xterm
|
||||||
|
sessions to connect to. You need to instanciate this in order for frontend
|
||||||
|
terminals to work.
|
||||||
|
|
||||||
|
This EventHandler accepts messages that have a data["command"] key specifying
|
||||||
|
a command to be executed.
|
||||||
|
The API of each command is documented in their respective handlers.
|
||||||
|
"""
|
||||||
|
def __init__(self, key, monitor):
|
||||||
|
"""
|
||||||
|
:param key: key this EventHandler listens to
|
||||||
|
:param monitor: tfw.components.HistoryMonitor instance to read command history from
|
||||||
|
"""
|
||||||
|
super().__init__(key)
|
||||||
|
self.working_directory = TFWENV.TERMINADO_DIR
|
||||||
|
self._historymonitor = monitor
|
||||||
|
bash_as_user_cmd = ['sudo', '-u', TAOENV.USER, 'bash']
|
||||||
|
self.terminado_server = TerminadoMiniServer('/terminal', TFWENV.TERMINADO_PORT, TFWENV.TERMINADO_WD, bash_as_user_cmd)
|
||||||
|
self.commands = {'write': self.write,
|
||||||
|
'read': self.read}
|
||||||
|
if self._historymonitor:
|
||||||
|
self._historymonitor.watch()
|
||||||
|
self.terminado_server.listen()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def historymonitor(self):
|
||||||
|
return self._historymonitor
|
||||||
|
|
||||||
|
def handle_event(self, message):
|
||||||
|
LOG.debug('TerminadoEventHandler received event: %s', message)
|
||||||
|
try:
|
||||||
|
data = message['data']
|
||||||
|
message['data'] = self.commands[data['command']](data)
|
||||||
|
return message
|
||||||
|
except KeyError:
|
||||||
|
LOG.error('IGNORING MESSAGE: Invalid message received: %s', message)
|
||||||
|
|
||||||
|
def write(self, data):
|
||||||
|
"""
|
||||||
|
Writes a string to the terminal session (on the pty level).
|
||||||
|
Useful for pre-typing and executing commands for the user.
|
||||||
|
|
||||||
|
:param data: TFW message data containing keys:
|
||||||
|
|-shellcmd: command to be written to the pty
|
||||||
|
"""
|
||||||
|
self.terminado_server.pty.write(data['shellcmd'])
|
||||||
|
|
||||||
|
def read(self, data):
|
||||||
|
"""
|
||||||
|
Reads the history of commands executed.
|
||||||
|
|
||||||
|
:param data: TFW message data containing keys:
|
||||||
|
|-count: the number of history elements to return
|
||||||
|
:return: message with list of commands in data['history']
|
||||||
|
"""
|
||||||
|
data['count'] = int(data.get('count', 1))
|
||||||
|
if self.historymonitor:
|
||||||
|
data['history'] = self.historymonitor.history[-data['count']:]
|
||||||
|
return data
|
||||||
|
|
||||||
|
def cleanup(self):
|
||||||
|
if self.historymonitor:
|
||||||
|
self.historymonitor.stop()
|
||||||
@@ -1 +1,4 @@
|
|||||||
from .envvars import *
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .envvars import TFWENV
|
||||||
|
|||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from envvars import prefixed_envvars_to_namedtuple
|
from envvars import prefixed_envvars_to_namedtuple
|
||||||
|
|
||||||
TFWENV = prefixed_envvars_to_namedtuple('TFW_', 'tfwenvtuple')
|
TFWENV = prefixed_envvars_to_namedtuple('TFW_', 'tfwenvtuple')
|
||||||
|
|||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
|
||||||
logging.basicConfig(level=logging.DEBUG)
|
logging.basicConfig(level=logging.DEBUG)
|
||||||
|
|||||||
@@ -0,0 +1,4 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .rate_limiter import RateLimiter
|
||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from functools import wraps
|
from functools import wraps
|
||||||
from time import time, sleep
|
from time import time, sleep
|
||||||
|
|
||||||
@@ -1,85 +1,104 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from abc import ABC, abstractmethod
|
from abc import ABC, abstractmethod
|
||||||
|
|
||||||
from tfw.networking.serialization import deserialize_all
|
from tfw.networking import deserialize_tfw_msg
|
||||||
from tfw.networking.event_handlers.server_connector import ServerConnector
|
from tfw.networking.event_handlers import ServerConnector
|
||||||
|
|
||||||
|
|
||||||
class EventHandlerBase(ABC):
|
class EventHandlerBase(ABC):
|
||||||
|
"""
|
||||||
|
Abstract base class for all Python based EventHandlers. Useful implementation template
|
||||||
|
for other languages.
|
||||||
|
|
||||||
|
Derived classes must implement the handle_event() method
|
||||||
|
"""
|
||||||
def __init__(self, key):
|
def __init__(self, key):
|
||||||
self.server_connector = ServerConnector()
|
self.server_connector = ServerConnector()
|
||||||
self.key = key
|
self.key = key
|
||||||
self.subscriptions = set()
|
self.subscribe(self.key, 'reset')
|
||||||
self.subscribe(self.key)
|
|
||||||
self.subscribe('reset')
|
|
||||||
self.server_connector.register_callback(self.event_handler_callback)
|
self.server_connector.register_callback(self.event_handler_callback)
|
||||||
|
|
||||||
def event_handler_callback(self, msg_parts):
|
def event_handler_callback(self, msg_parts):
|
||||||
key, message = deserialize_all(*msg_parts)
|
"""
|
||||||
response = self.dispatch_handling(key, message)
|
Callback that is invoked when receiving a message.
|
||||||
if response is None:
|
Dispatches messages to handler methods and sends
|
||||||
return
|
a response back in case the handler returned something.
|
||||||
self.server_connector.send(key, response)
|
This is subscribed in __init__().
|
||||||
|
"""
|
||||||
|
message = deserialize_tfw_msg(*msg_parts)
|
||||||
|
response = self.dispatch_handling(message)
|
||||||
|
if response:
|
||||||
|
response['key'] = message['key']
|
||||||
|
self.server_connector.send(response)
|
||||||
|
|
||||||
|
def dispatch_handling(self, message):
|
||||||
|
"""
|
||||||
|
Used to dispatch messages to their specific handlers.
|
||||||
|
"""
|
||||||
|
if message['key'] != 'reset':
|
||||||
|
return self.handle_event(message)
|
||||||
|
return self.handle_reset(message)
|
||||||
|
|
||||||
@abstractmethod
|
@abstractmethod
|
||||||
def dispatch_handling(self, key, message):
|
def handle_event(self, message):
|
||||||
raise NotImplementedError
|
"""
|
||||||
|
Abstract method that implements the handling of messages.
|
||||||
|
|
||||||
def _dispatch_handling(self, key, message):
|
:param message: the message received
|
||||||
# pylint: disable=no-else-return
|
"""
|
||||||
if key != 'reset':
|
|
||||||
return self.handle_event(key, message)
|
|
||||||
else:
|
|
||||||
return self.handle_reset(message)
|
|
||||||
|
|
||||||
@abstractmethod
|
|
||||||
def handle_event(self, key, message):
|
|
||||||
raise NotImplementedError
|
raise NotImplementedError
|
||||||
|
|
||||||
def handle_reset(self, message):
|
def handle_reset(self, message):
|
||||||
# pylint: disable=unused-argument,no-self-use
|
# pylint: disable=unused-argument,no-self-use
|
||||||
|
"""
|
||||||
|
Usually 'reset' events receive some sort of special treatment.
|
||||||
|
|
||||||
|
:param message: the message received
|
||||||
|
"""
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def cleanup(self):
|
def subscribe(self, *keys):
|
||||||
pass
|
"""
|
||||||
|
Subscribe this EventHandler to receive events for given keys.
|
||||||
|
Note that you can subscribe to the same key several times in which
|
||||||
|
case you will need to unsubscribe multiple times in order to stop
|
||||||
|
receiving events.
|
||||||
|
|
||||||
def message_other(self, key, data):
|
:param keys: list of keys to subscribe to
|
||||||
message = {
|
"""
|
||||||
'key': key,
|
for key in keys:
|
||||||
'data': data
|
|
||||||
}
|
|
||||||
self.server_connector.send(key, message)
|
|
||||||
|
|
||||||
def subscribe(self, key):
|
|
||||||
if key not in self.subscriptions:
|
|
||||||
self.subscriptions.add(key)
|
|
||||||
self.server_connector.subscribe(key)
|
self.server_connector.subscribe(key)
|
||||||
|
|
||||||
def unsubscribe(self, key):
|
def unsubscribe(self, *keys):
|
||||||
try:
|
"""
|
||||||
self.subscriptions.remove(key)
|
Unsubscribe this eventhandler from the given keys.
|
||||||
|
|
||||||
|
:param keys: list of keys to unsubscribe from
|
||||||
|
"""
|
||||||
|
for key in keys:
|
||||||
self.server_connector.unsubscribe(key)
|
self.server_connector.unsubscribe(key)
|
||||||
except KeyError:
|
|
||||||
pass
|
|
||||||
|
|
||||||
def unsubscribe_all(self):
|
def cleanup(self):
|
||||||
for sub in self.subscriptions:
|
"""
|
||||||
self.server_connector.unsubscribe(key=sub)
|
Perform cleanup actions such as releasing database
|
||||||
self.subscriptions.clear()
|
connections and stuff like that.
|
||||||
|
"""
|
||||||
|
pass
|
||||||
class TriggerlessEventHandler(EventHandlerBase, ABC):
|
|
||||||
# pylint: disable=abstract-method
|
|
||||||
def dispatch_handling(self, key, message):
|
|
||||||
return self._dispatch_handling(key, message)
|
|
||||||
|
|
||||||
|
|
||||||
class TriggeredEventHandler(EventHandlerBase, ABC):
|
class TriggeredEventHandler(EventHandlerBase, ABC):
|
||||||
# pylint: disable=abstract-method
|
# pylint: disable=abstract-method
|
||||||
|
"""
|
||||||
|
Abstract base class for EventHandlers which are only triggered in case
|
||||||
|
TFWServer has successfully triggered an FSM step defined in __init__.
|
||||||
|
"""
|
||||||
def __init__(self, key, trigger):
|
def __init__(self, key, trigger):
|
||||||
super().__init__(key)
|
super().__init__(key)
|
||||||
self.trigger = trigger
|
self.trigger = trigger
|
||||||
|
|
||||||
def dispatch_handling(self, key, message):
|
def dispatch_handling(self, message):
|
||||||
if message.get('trigger') == self.trigger:
|
if message.get('trigger') == self.trigger:
|
||||||
return self._dispatch_handling(key, message)
|
return super().dispatch_handling(message)
|
||||||
return None
|
return None
|
||||||
|
|||||||
+11
-1
@@ -1,11 +1,21 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from typing import List
|
from typing import List
|
||||||
|
|
||||||
from transitions import Machine
|
from transitions import Machine
|
||||||
|
|
||||||
from tfw.components.mixins.callback_mixin import CallbackMixin
|
from tfw.mixins import CallbackMixin
|
||||||
|
|
||||||
|
|
||||||
class FSMBase(CallbackMixin):
|
class FSMBase(CallbackMixin):
|
||||||
|
"""
|
||||||
|
A general FSM base class you can inherit from to track user progress.
|
||||||
|
See linear_fsm.py for an example use-case.
|
||||||
|
TFW the transitions library for state machines, please refer to their
|
||||||
|
documentation for more information on creating your own machines:
|
||||||
|
https://github.com/pytransitions/transitions
|
||||||
|
"""
|
||||||
states, transitions = [], []
|
states, transitions = [], []
|
||||||
|
|
||||||
def __init__(self, initial: str = None, accepted_states: List[str] = None):
|
def __init__(self, initial: str = None, accepted_states: List[str] = None):
|
||||||
|
|||||||
@@ -0,0 +1,18 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .fsm_base import FSMBase
|
||||||
|
|
||||||
|
|
||||||
|
class LinearFSM(FSMBase):
|
||||||
|
"""
|
||||||
|
This is a state machine for challenges with linear progression, consisting of
|
||||||
|
a number of steps specified in the constructor. It automatically sets up a single
|
||||||
|
action between states as such:
|
||||||
|
0 ==step_1==> 1 ==step_2==> 2 ==step_3==> 3 ... and so on
|
||||||
|
"""
|
||||||
|
def __init__(self, number_of_steps):
|
||||||
|
self.states = list(map(str, range(number_of_steps)))
|
||||||
|
self.transitions = [{'trigger': f'step_{int(index)+1}', 'source': index, 'dest': str(int(index)+1)}
|
||||||
|
for index in self.states[:-1]]
|
||||||
|
super(LinearFSM, self).__init__()
|
||||||
@@ -1,18 +0,0 @@
|
|||||||
from datetime import datetime
|
|
||||||
|
|
||||||
from tfw.networking.event_handlers.server_connector import ServerUplinkConnector
|
|
||||||
|
|
||||||
|
|
||||||
class MessageSender:
|
|
||||||
def __init__(self, custom_key: str = None):
|
|
||||||
self.server_connector = ServerUplinkConnector()
|
|
||||||
self.key = custom_key or 'message'
|
|
||||||
|
|
||||||
def send(self, originator, message):
|
|
||||||
data = {
|
|
||||||
'originator': originator,
|
|
||||||
'timestamp': datetime.now().isoformat(),
|
|
||||||
'message': message
|
|
||||||
}
|
|
||||||
response = {'data': data}
|
|
||||||
self.server_connector.send(self.key, response)
|
|
||||||
@@ -0,0 +1,7 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .supervisor_mixin import SupervisorMixin
|
||||||
|
from .callback_mixin import CallbackMixin
|
||||||
|
from .observer_mixin import ObserverMixin
|
||||||
|
from .monitor_manager_mixin import MonitorManagerMixin
|
||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from functools import partial
|
from functools import partial
|
||||||
|
|
||||||
|
|
||||||
@@ -6,6 +9,12 @@ class CallbackMixin:
|
|||||||
self._callbacks = []
|
self._callbacks = []
|
||||||
|
|
||||||
def subscribe_callback(self, callback, *args, **kwargs):
|
def subscribe_callback(self, callback, *args, **kwargs):
|
||||||
|
"""
|
||||||
|
Subscribe a callable to invoke once an event is triggered.
|
||||||
|
:param callback: callable to be executed on events
|
||||||
|
:param *args: arguments passed to callable
|
||||||
|
:param **kwargs: kwargs passed to callable
|
||||||
|
"""
|
||||||
fun = partial(callback, *args, **kwargs)
|
fun = partial(callback, *args, **kwargs)
|
||||||
self._callbacks.append(fun)
|
self._callbacks.append(fun)
|
||||||
|
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from tfw.config.logs import logging
|
||||||
|
|
||||||
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class MonitorManagerMixin:
|
||||||
|
def __init__(self, monitor_type, directory):
|
||||||
|
self._monitor_type = monitor_type
|
||||||
|
self._monitor = None
|
||||||
|
self._monitored_directory = directory
|
||||||
|
self.reload_monitor()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def monitor(self):
|
||||||
|
return self._monitor
|
||||||
|
|
||||||
|
def reload_monitor(self):
|
||||||
|
if self._monitor:
|
||||||
|
try:
|
||||||
|
self._monitor.stop()
|
||||||
|
except KeyError:
|
||||||
|
LOG.debug('Working directory was removed – ignoring...')
|
||||||
|
self._monitor = self._monitor_type(self._monitored_directory)
|
||||||
|
self._monitor.watch() # This runs on a separate thread
|
||||||
@@ -0,0 +1,16 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from watchdog.observers import Observer
|
||||||
|
|
||||||
|
|
||||||
|
class ObserverMixin:
|
||||||
|
def __init__(self):
|
||||||
|
self.observer = Observer()
|
||||||
|
|
||||||
|
def watch(self):
|
||||||
|
self.observer.start()
|
||||||
|
|
||||||
|
def stop(self):
|
||||||
|
self.observer.stop()
|
||||||
|
self.observer.join()
|
||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
import xmlrpc.client
|
import xmlrpc.client
|
||||||
from xmlrpc.client import Fault as SupervisorFault
|
from xmlrpc.client import Fault as SupervisorFault
|
||||||
from contextlib import suppress
|
from contextlib import suppress
|
||||||
@@ -0,0 +1,9 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .serialization import serialize_tfw_msg, deserialize_tfw_msg, validate_message
|
||||||
|
from .zmq_connector_base import ZMQConnectorBase
|
||||||
|
# from .controller_connector import ControllerConnector # TODO: readd once controller stuff is resolved
|
||||||
|
from .message_sender import MessageSender
|
||||||
|
from .event_handlers.server_connector import ServerUplinkConnector as TFWServerConnector
|
||||||
|
from .server.tfw_server import TFWServer
|
||||||
|
|||||||
@@ -1,15 +1,18 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
import zmq
|
import zmq
|
||||||
from zmq.eventloop.zmqstream import ZMQStream
|
from zmq.eventloop.zmqstream import ZMQStream
|
||||||
|
|
||||||
from tfw.config import TFWENV
|
from tfw.config import TFWENV
|
||||||
from tfw.networking.zmq_connector_base import ZMQConnectorBase
|
from tfw.networking import ZMQConnectorBase
|
||||||
|
|
||||||
|
|
||||||
class ControllerConnector(ZMQConnectorBase):
|
class ControllerConnector(ZMQConnectorBase):
|
||||||
def __init__(self, zmq_context=None):
|
def __init__(self, zmq_context=None):
|
||||||
super(ControllerConnector, self).__init__(zmq_context)
|
super(ControllerConnector, self).__init__(zmq_context)
|
||||||
self._zmq_rep_socket = self._zmq_context.socket(zmq.REP)
|
self._zmq_rep_socket = self._zmq_context.socket(zmq.REP)
|
||||||
self._zmq_rep_socket.connect('tcp://localhost:{}'.format(TFWENV.CONTROLLER_PORT))
|
self._zmq_rep_socket.connect(f'tcp://localhost:{TFWENV.CONTROLLER_PORT}')
|
||||||
self._zmq_rep_stream = ZMQStream(self._zmq_rep_socket)
|
self._zmq_rep_stream = ZMQStream(self._zmq_rep_socket)
|
||||||
|
|
||||||
self.register_callback = self._zmq_rep_stream.on_recv_stream
|
self.register_callback = self._zmq_rep_stream.on_recv_stream
|
||||||
|
|||||||
@@ -0,0 +1,4 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .server_connector import ServerConnector, ServerUplinkConnector, ServerDownlinkConnector
|
||||||
|
|||||||
@@ -1,18 +1,21 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from functools import partial
|
from functools import partial
|
||||||
|
|
||||||
import zmq
|
import zmq
|
||||||
from zmq.eventloop.zmqstream import ZMQStream
|
from zmq.eventloop.zmqstream import ZMQStream
|
||||||
|
|
||||||
from tfw.networking.serialization import serialize_all
|
from tfw.networking import serialize_tfw_msg
|
||||||
|
from tfw.networking import ZMQConnectorBase
|
||||||
from tfw.config import TFWENV
|
from tfw.config import TFWENV
|
||||||
from tfw.networking.zmq_connector_base import ZMQConnectorBase
|
|
||||||
|
|
||||||
|
|
||||||
class ServerDownlinkConnector(ZMQConnectorBase):
|
class ServerDownlinkConnector(ZMQConnectorBase):
|
||||||
def __init__(self, zmq_context=None):
|
def __init__(self, zmq_context=None):
|
||||||
super(ServerDownlinkConnector, self).__init__(zmq_context)
|
super(ServerDownlinkConnector, self).__init__(zmq_context)
|
||||||
self._zmq_sub_socket = self._zmq_context.socket(zmq.SUB)
|
self._zmq_sub_socket = self._zmq_context.socket(zmq.SUB)
|
||||||
self._zmq_sub_socket.connect('tcp://localhost:{}'.format(TFWENV.PUBLISHER_PORT))
|
self._zmq_sub_socket.connect(f'tcp://localhost:{TFWENV.PUBLISHER_PORT}')
|
||||||
self._zmq_sub_stream = ZMQStream(self._zmq_sub_socket)
|
self._zmq_sub_stream = ZMQStream(self._zmq_sub_socket)
|
||||||
|
|
||||||
self.subscribe = partial(self._zmq_sub_socket.setsockopt_string, zmq.SUBSCRIBE)
|
self.subscribe = partial(self._zmq_sub_socket.setsockopt_string, zmq.SUBSCRIBE)
|
||||||
@@ -21,18 +24,31 @@ class ServerDownlinkConnector(ZMQConnectorBase):
|
|||||||
|
|
||||||
|
|
||||||
class ServerUplinkConnector(ZMQConnectorBase):
|
class ServerUplinkConnector(ZMQConnectorBase):
|
||||||
|
"""
|
||||||
|
Class capable of sending messages to the TFW server and event handlers.
|
||||||
|
"""
|
||||||
def __init__(self, zmq_context=None):
|
def __init__(self, zmq_context=None):
|
||||||
super(ServerUplinkConnector, self).__init__(zmq_context)
|
super(ServerUplinkConnector, self).__init__(zmq_context)
|
||||||
self._zmq_push_socket = self._zmq_context.socket(zmq.PUSH)
|
self._zmq_push_socket = self._zmq_context.socket(zmq.PUSH)
|
||||||
self._zmq_push_socket.connect('tcp://localhost:{}'.format(TFWENV.RECEIVER_PORT))
|
self._zmq_push_socket.connect(f'tcp://localhost:{TFWENV.RECEIVER_PORT}')
|
||||||
|
|
||||||
def send_to_eventhandler(self, key, message):
|
def send_to_eventhandler(self, message):
|
||||||
message['data']['key'] = key
|
"""
|
||||||
self.send('mirror', message)
|
Send a message to an event handler.
|
||||||
|
:param message: JSON message you want to send
|
||||||
|
:param message['key']: key of event handler you want to address
|
||||||
|
"""
|
||||||
|
nested_message = {'key': message['key'], 'data': message.pop('data')}
|
||||||
|
message['key'] = 'mirror'
|
||||||
|
message['data'] = nested_message
|
||||||
|
self.send(message)
|
||||||
|
|
||||||
def send(self, key, message):
|
def send(self, message):
|
||||||
message['key'] = key
|
"""
|
||||||
self._zmq_push_socket.send_multipart(serialize_all(key, message))
|
Send a message to the TFW server
|
||||||
|
:param message: JSON message you want to send
|
||||||
|
"""
|
||||||
|
self._zmq_push_socket.send_multipart(serialize_tfw_msg(message))
|
||||||
|
|
||||||
|
|
||||||
class ServerConnector(ServerUplinkConnector, ServerDownlinkConnector):
|
class ServerConnector(ServerUplinkConnector, ServerDownlinkConnector):
|
||||||
|
|||||||
@@ -0,0 +1,30 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from datetime import datetime
|
||||||
|
|
||||||
|
from tfw.networking.event_handlers import ServerUplinkConnector
|
||||||
|
|
||||||
|
|
||||||
|
class MessageSender:
|
||||||
|
"""
|
||||||
|
Provides a mechanism to send messages to our frontend messaging component which
|
||||||
|
displays messages with the key "message".
|
||||||
|
"""
|
||||||
|
def __init__(self, custom_key: str = None):
|
||||||
|
self.server_connector = ServerUplinkConnector()
|
||||||
|
self.key = custom_key or 'message'
|
||||||
|
|
||||||
|
def send(self, originator, message):
|
||||||
|
"""
|
||||||
|
Sends a message to the key specified in __init__.
|
||||||
|
:param originator: name of sender to be displayed on the frontend
|
||||||
|
:param message: message to send
|
||||||
|
"""
|
||||||
|
data = {
|
||||||
|
'originator': originator,
|
||||||
|
'timestamp': datetime.now().isoformat(),
|
||||||
|
'message': message
|
||||||
|
}
|
||||||
|
self.server_connector.send({'key': self.key,
|
||||||
|
'data': data})
|
||||||
@@ -1,21 +1,69 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
"""
|
||||||
|
TFW JSON message format
|
||||||
|
|
||||||
|
message:
|
||||||
|
{
|
||||||
|
"key": string, # addressing
|
||||||
|
"data": {...}, # payload
|
||||||
|
"trigger": string # FSM trigger
|
||||||
|
}
|
||||||
|
|
||||||
|
ZeroMQ's sub-pub sockets use enveloped messages
|
||||||
|
(http://zguide.zeromq.org/page:all#Pub-Sub-Message-Envelopes)
|
||||||
|
and TFW also uses them internally. This means that on ZMQ sockets
|
||||||
|
we always send the messages key separately and then the actual
|
||||||
|
message (which contains the key as well) like so:
|
||||||
|
|
||||||
|
socket.send_multipart([message['key'], message])
|
||||||
|
|
||||||
|
The purpose of this module is abstracting away this low level behaviour.
|
||||||
|
"""
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
|
||||||
|
|
||||||
def encode_if_needed(value):
|
def validate_message(message):
|
||||||
|
return 'key' in message
|
||||||
|
|
||||||
|
|
||||||
|
def serialize_tfw_msg(message):
|
||||||
|
return _serialize_all(message['key'], message)
|
||||||
|
|
||||||
|
|
||||||
|
def deserialize_tfw_msg(*args):
|
||||||
|
return _deserialize_all(*args)[1]
|
||||||
|
|
||||||
|
|
||||||
|
def _serialize_all(*args):
|
||||||
|
return tuple(_serialize_single(arg) for arg in args)
|
||||||
|
|
||||||
|
|
||||||
|
def _deserialize_all(*args):
|
||||||
|
return tuple(_deserialize_single(arg) for arg in args)
|
||||||
|
|
||||||
|
|
||||||
|
def _serialize_single(data):
|
||||||
|
if not isinstance(data, str):
|
||||||
|
data = json.dumps(data)
|
||||||
|
return _encode_if_needed(data)
|
||||||
|
|
||||||
|
|
||||||
|
def _deserialize_single(data):
|
||||||
|
try:
|
||||||
|
return json.loads(data)
|
||||||
|
except ValueError:
|
||||||
|
return _decode_if_needed(data)
|
||||||
|
|
||||||
|
|
||||||
|
def _encode_if_needed(value):
|
||||||
if isinstance(value, str):
|
if isinstance(value, str):
|
||||||
value = value.encode('utf-8')
|
value = value.encode('utf-8')
|
||||||
return value
|
return value
|
||||||
|
|
||||||
|
|
||||||
def decode_if_needed(value):
|
def _decode_if_needed(value):
|
||||||
if isinstance(value, (bytes, bytearray)):
|
if isinstance(value, (bytes, bytearray)):
|
||||||
value = value.decode('utf-8')
|
value = value.decode('utf-8')
|
||||||
return value
|
return value
|
||||||
|
|
||||||
|
|
||||||
def serialize_all(key, data):
|
|
||||||
return [encode_if_needed(frame) for frame in (key, json.dumps(data))]
|
|
||||||
|
|
||||||
|
|
||||||
def deserialize_all(key, data):
|
|
||||||
return decode_if_needed(key), json.loads(data)
|
|
||||||
|
|||||||
@@ -0,0 +1,7 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from .event_handler_connector import EventHandlerConnector, EventHandlerUplinkConnector, EventHandlerDownlinkConnector
|
||||||
|
from .tfw_server import TFWServer
|
||||||
|
from .zmq_websocket_handler import ZMQWebSocketProxy
|
||||||
|
# from .controller_responder import ControllerResponder # TODO: readd once controller stuff is resolved
|
||||||
|
|||||||
@@ -1,5 +1,7 @@
|
|||||||
from tfw.networking.controller_connector import ControllerConnector
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
from tfw.networking.serialization import deserialize_all, serialize_all
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
|
from tfw.networking import deserialize_all, serialize_all, ControllerConnector
|
||||||
|
|
||||||
|
|
||||||
class ControllerResponder:
|
class ControllerResponder:
|
||||||
|
|||||||
@@ -1,8 +1,10 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
import zmq
|
import zmq
|
||||||
from zmq.eventloop.zmqstream import ZMQStream
|
from zmq.eventloop.zmqstream import ZMQStream
|
||||||
|
|
||||||
from tfw.networking.zmq_connector_base import ZMQConnectorBase
|
from tfw.networking import ZMQConnectorBase, serialize_tfw_msg
|
||||||
from tfw.networking.serialization import serialize_all
|
|
||||||
from tfw.config import TFWENV
|
from tfw.config import TFWENV
|
||||||
from tfw.config.logs import logging
|
from tfw.config.logs import logging
|
||||||
|
|
||||||
@@ -14,7 +16,7 @@ class EventHandlerDownlinkConnector(ZMQConnectorBase):
|
|||||||
super(EventHandlerDownlinkConnector, self).__init__(zmq_context)
|
super(EventHandlerDownlinkConnector, self).__init__(zmq_context)
|
||||||
self._zmq_pull_socket = self._zmq_context.socket(zmq.PULL)
|
self._zmq_pull_socket = self._zmq_context.socket(zmq.PULL)
|
||||||
self._zmq_pull_stream = ZMQStream(self._zmq_pull_socket)
|
self._zmq_pull_stream = ZMQStream(self._zmq_pull_socket)
|
||||||
address = 'tcp://*:{}'.format(TFWENV.RECEIVER_PORT)
|
address = f'tcp://*:{TFWENV.RECEIVER_PORT}'
|
||||||
self._zmq_pull_socket.bind(address)
|
self._zmq_pull_socket.bind(address)
|
||||||
LOG.debug('Pull socket bound to %s', address)
|
LOG.debug('Pull socket bound to %s', address)
|
||||||
|
|
||||||
@@ -23,7 +25,7 @@ class EventHandlerUplinkConnector(ZMQConnectorBase):
|
|||||||
def __init__(self, zmq_context=None):
|
def __init__(self, zmq_context=None):
|
||||||
super(EventHandlerUplinkConnector, self).__init__(zmq_context)
|
super(EventHandlerUplinkConnector, self).__init__(zmq_context)
|
||||||
self._zmq_pub_socket = self._zmq_context.socket(zmq.PUB)
|
self._zmq_pub_socket = self._zmq_context.socket(zmq.PUB)
|
||||||
address = 'tcp://*:{}'.format(TFWENV.PUBLISHER_PORT)
|
address = f'tcp://*:{TFWENV.PUBLISHER_PORT}'
|
||||||
self._zmq_pub_socket.bind(address)
|
self._zmq_pub_socket.bind(address)
|
||||||
LOG.debug('Pub socket bound to %s', address)
|
LOG.debug('Pub socket bound to %s', address)
|
||||||
|
|
||||||
@@ -32,7 +34,5 @@ class EventHandlerConnector(EventHandlerDownlinkConnector, EventHandlerUplinkCon
|
|||||||
def register_callback(self, callback):
|
def register_callback(self, callback):
|
||||||
self._zmq_pull_stream.on_recv(callback)
|
self._zmq_pull_stream.on_recv(callback)
|
||||||
|
|
||||||
def send_message(self, message: dict, key: str = None):
|
def send_message(self, message: dict):
|
||||||
if not key:
|
self._zmq_pub_socket.send_multipart(serialize_tfw_msg(message))
|
||||||
key = message.get('key', '')
|
|
||||||
self._zmq_pub_socket.send_multipart(serialize_all(key, message))
|
|
||||||
|
|||||||
@@ -1,25 +1,40 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
from collections import defaultdict
|
from collections import defaultdict
|
||||||
|
|
||||||
from tornado.web import Application
|
from tornado.web import Application
|
||||||
|
|
||||||
from tfw.networking.server.zmq_websocket_handler import ZMQWebSocketProxy
|
from tfw.networking import MessageSender
|
||||||
from tfw.networking.event_handlers.server_connector import ServerUplinkConnector
|
from tfw.networking.event_handlers import ServerUplinkConnector
|
||||||
from tfw.message_sender import MessageSender
|
from tfw.networking.server import EventHandlerConnector
|
||||||
from tfw.config.logs import logging
|
from tfw.config.logs import logging
|
||||||
|
from .zmq_websocket_handler import ZMQWebSocketProxy
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class TFWServer:
|
class TFWServer:
|
||||||
|
"""
|
||||||
|
This class handles the proxying of messages between the frontend and event handers.
|
||||||
|
It proxies messages from the "/ws" route to all event handlers subscribed to a ZMQ
|
||||||
|
SUB socket. It also manages an FSM you can define as a constructor argument.
|
||||||
|
"""
|
||||||
def __init__(self, fsm_type):
|
def __init__(self, fsm_type):
|
||||||
|
"""
|
||||||
|
:param fsm_type: the type of FSM you want TFW to use
|
||||||
|
"""
|
||||||
self._fsm = fsm_type()
|
self._fsm = fsm_type()
|
||||||
self._fsm_updater = FSMUpdater(self._fsm)
|
self._fsm_updater = FSMUpdater(self._fsm)
|
||||||
self._fsm_manager = FSMManager(self._fsm)
|
self._fsm_manager = FSMManager(self._fsm)
|
||||||
self._fsm.subscribe_callback(self._fsm_updater.update)
|
self._fsm.subscribe_callback(self._fsm_updater.update)
|
||||||
|
self._event_handler_connector = EventHandlerConnector()
|
||||||
|
|
||||||
self.application = Application(
|
self.application = Application(
|
||||||
[(r'/ws', ZMQWebSocketProxy, {'make_response': self.make_response,
|
[(r'/ws', ZMQWebSocketProxy, {'make_eventhandler_message': self.make_eventhandler_message,
|
||||||
'proxy_filter': self.proxy_filter})]
|
'proxy_filter': self.proxy_filter,
|
||||||
|
'handle_trigger': self.handle_trigger,
|
||||||
|
'event_handler_connector': self._event_handler_connector})]
|
||||||
)
|
)
|
||||||
#self.controller_responder = ControllerResponder(self.fsm) TODO: add this once controller stuff is resolved
|
#self.controller_responder = ControllerResponder(self.fsm) TODO: add this once controller stuff is resolved
|
||||||
|
|
||||||
@@ -31,13 +46,21 @@ class TFWServer:
|
|||||||
def fsm_manager(self):
|
def fsm_manager(self):
|
||||||
return self._fsm_manager
|
return self._fsm_manager
|
||||||
|
|
||||||
def make_response(self, message):
|
def make_eventhandler_message(self, message):
|
||||||
|
self.trigger_fsm(message)
|
||||||
|
message['FSMUpdate'] = self._fsm_updater.get_fsm_state_and_transitions()
|
||||||
|
return message
|
||||||
|
|
||||||
|
def handle_trigger(self, message):
|
||||||
|
LOG.debug('Executing handler for trigger "%s"', message.get('trigger', ''))
|
||||||
|
self.trigger_fsm(message)
|
||||||
|
|
||||||
|
def trigger_fsm(self, message):
|
||||||
trigger = message.get('trigger', '')
|
trigger = message.get('trigger', '')
|
||||||
try:
|
try:
|
||||||
self._fsm_manager.trigger(trigger, message)
|
self._fsm_manager.trigger(trigger, message)
|
||||||
except AttributeError:
|
except AttributeError:
|
||||||
LOG.debug('FSM failed to execute nonexistent trigger: "%s"', trigger)
|
LOG.debug('FSM failed to execute nonexistent trigger: "%s"', trigger)
|
||||||
return message
|
|
||||||
|
|
||||||
def proxy_filter(self, message):
|
def proxy_filter(self, message):
|
||||||
# pylint: disable=unused-argument,no-self-use
|
# pylint: disable=unused-argument,no-self-use
|
||||||
@@ -60,7 +83,7 @@ class FSMManager:
|
|||||||
def trigger(self, trigger, message):
|
def trigger(self, trigger, message):
|
||||||
predicate_results = []
|
predicate_results = []
|
||||||
for predicate in self.trigger_predicates[trigger]:
|
for predicate in self.trigger_predicates[trigger]:
|
||||||
success, message = predicate()
|
success, message = predicate(message)
|
||||||
predicate_results.append(success)
|
predicate_results.append(success)
|
||||||
self.messenge_sender.send('FSM', message)
|
self.messenge_sender.send('FSM', message)
|
||||||
|
|
||||||
@@ -85,11 +108,13 @@ class FSMUpdater:
|
|||||||
|
|
||||||
def update(self, kwargs_dict):
|
def update(self, kwargs_dict):
|
||||||
# pylint: disable=unused-argument
|
# pylint: disable=unused-argument
|
||||||
self.uplink.send(*self.generate_fsm_update())
|
self.uplink.send(self.generate_fsm_update())
|
||||||
|
|
||||||
def generate_fsm_update(self):
|
def generate_fsm_update(self):
|
||||||
key = 'FSMUpdate'
|
return {'key': 'FSMUpdate',
|
||||||
response = {'data': {'current_state': self.fsm.state,
|
'data': self.get_fsm_state_and_transitions()}
|
||||||
'valid_transitions':
|
|
||||||
[{'trigger': trigger} for trigger in self.fsm.machine.get_triggers(self.fsm.state)]}}
|
def get_fsm_state_and_transitions(self):
|
||||||
return key, response
|
state = self.fsm.state
|
||||||
|
valid_transitions = [{'trigger': trigger} for trigger in self.fsm.machine.get_triggers(self.fsm.state)]
|
||||||
|
return {'current_state': state, 'valid_transitions': valid_transitions}
|
||||||
|
|||||||
@@ -1,19 +1,23 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
import json
|
import json
|
||||||
from abc import ABC, abstractmethod
|
from abc import ABC, abstractmethod
|
||||||
|
|
||||||
from tornado.websocket import WebSocketHandler
|
from tornado.websocket import WebSocketHandler
|
||||||
|
|
||||||
from tfw.networking.serialization import deserialize_all
|
from tfw.networking import deserialize_tfw_msg, validate_message
|
||||||
from tfw.networking.server.event_handler_connector import EventHandlerConnector
|
|
||||||
from tfw.config.logs import logging
|
from tfw.config.logs import logging
|
||||||
|
|
||||||
LOG = logging.getLogger(__name__)
|
LOG = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class ZMQWebSocketHandler(WebSocketHandler, ABC):
|
class ZMQWebSocketHandler(WebSocketHandler, ABC):
|
||||||
_event_handler_connector = EventHandlerConnector()
|
|
||||||
instances = set()
|
instances = set()
|
||||||
|
|
||||||
|
def initialize(self, **kwargs): # pylint: disable=arguments-differ
|
||||||
|
self._event_handler_connector = kwargs['event_handler_connector']
|
||||||
|
|
||||||
def prepare(self):
|
def prepare(self):
|
||||||
ZMQWebSocketHandler.instances.add(self)
|
ZMQWebSocketHandler.instances.add(self)
|
||||||
|
|
||||||
@@ -24,36 +28,43 @@ class ZMQWebSocketHandler(WebSocketHandler, ABC):
|
|||||||
LOG.debug('WebSocket connection initiated')
|
LOG.debug('WebSocket connection initiated')
|
||||||
self._event_handler_connector.register_callback(self.zmq_callback)
|
self._event_handler_connector.register_callback(self.zmq_callback)
|
||||||
|
|
||||||
@classmethod
|
def zmq_callback(self, msg_parts):
|
||||||
def zmq_callback(cls, msg_parts):
|
keyhandlers = {'mirror': self.mirror}
|
||||||
keyhandlers = {'mirror': cls.mirror}
|
|
||||||
|
|
||||||
key, data = deserialize_all(*msg_parts)
|
message = deserialize_tfw_msg(*msg_parts)
|
||||||
LOG.debug('Received on pull socket: %s', data)
|
LOG.debug('Received on pull socket: %s', message)
|
||||||
if key not in keyhandlers:
|
if not validate_message(message):
|
||||||
for instance in cls.instances:
|
return
|
||||||
instance.write_message(data)
|
|
||||||
|
self.handle_trigger(message)
|
||||||
|
if message['key'] not in keyhandlers:
|
||||||
|
for instance in ZMQWebSocketHandler.instances:
|
||||||
|
instance.write_message(message)
|
||||||
else:
|
else:
|
||||||
try:
|
try:
|
||||||
keyhandlers[key](data['data'])
|
keyhandlers[message['key']](message)
|
||||||
except KeyError:
|
except KeyError:
|
||||||
LOG.error('Invalid mirror message format! Ignoring.')
|
LOG.error('Invalid mirror message format! Ignoring.')
|
||||||
|
|
||||||
@classmethod
|
def mirror(self, message):
|
||||||
def mirror(cls, data):
|
message = message['data']
|
||||||
key = data['key']
|
self._event_handler_connector.send_message(message)
|
||||||
cls._event_handler_connector.send_message({'data': data}, key)
|
|
||||||
|
|
||||||
def on_message(self, message):
|
def on_message(self, message):
|
||||||
LOG.debug('Received on WebSocket: %s', message)
|
LOG.debug('Received on WebSocket: %s', message)
|
||||||
self.send_message(self.make_response(message))
|
if validate_message(message):
|
||||||
|
self.send_message(self.make_eventhandler_message(message))
|
||||||
|
|
||||||
@abstractmethod
|
@abstractmethod
|
||||||
def make_response(self, message):
|
def make_eventhandler_message(self, message):
|
||||||
raise NotImplementedError
|
raise NotImplementedError
|
||||||
|
|
||||||
def send_message(self, message: dict, key: str = None):
|
def send_message(self, message: dict):
|
||||||
self._event_handler_connector.send_message(message, key)
|
self._event_handler_connector.send_message(message)
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
def handle_trigger(self, message):
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
# much secure, very cors, wow
|
# much secure, very cors, wow
|
||||||
def check_origin(self, origin):
|
def check_origin(self, origin):
|
||||||
@@ -62,15 +73,19 @@ class ZMQWebSocketHandler(WebSocketHandler, ABC):
|
|||||||
|
|
||||||
class ZMQWebSocketProxy(ZMQWebSocketHandler):
|
class ZMQWebSocketProxy(ZMQWebSocketHandler):
|
||||||
# pylint: disable=abstract-method
|
# pylint: disable=abstract-method
|
||||||
def initialize(self, make_response, proxy_filter):
|
def initialize(self, **kwargs): # pylint: disable=arguments-differ
|
||||||
# pylint: disable=arguments-differ
|
super(ZMQWebSocketProxy, self).initialize(**kwargs)
|
||||||
self._make_response = make_response
|
self._make_eventhandler_message = kwargs['make_eventhandler_message']
|
||||||
self._proxy_filter = proxy_filter
|
self._proxy_filter = kwargs['proxy_filter']
|
||||||
|
self._handle_trigger = kwargs['handle_trigger']
|
||||||
|
|
||||||
def on_message(self, message):
|
def on_message(self, message):
|
||||||
message = json.loads(message)
|
message = json.loads(message)
|
||||||
if self._proxy_filter(message):
|
if self._proxy_filter(message):
|
||||||
super().on_message(message)
|
super().on_message(message)
|
||||||
|
|
||||||
def make_response(self, message):
|
def make_eventhandler_message(self, message):
|
||||||
return self._make_response(message)
|
return self._make_eventhandler_message(message)
|
||||||
|
|
||||||
|
def handle_trigger(self, message):
|
||||||
|
self._handle_trigger(message)
|
||||||
|
|||||||
@@ -1,3 +1,6 @@
|
|||||||
|
# Copyright (C) 2018 Avatao.com Innovative Learning Kft.
|
||||||
|
# All Rights Reserved. See LICENSE file for details.
|
||||||
|
|
||||||
import zmq
|
import zmq
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,3 +0,0 @@
|
|||||||
location = /login {
|
|
||||||
proxy_pass http://127.0.0.1:${TFW_LOGIN_APP_PORT};
|
|
||||||
}
|
|
||||||
@@ -3,4 +3,4 @@ location = /terminal {
|
|||||||
proxy_http_version 1.1;
|
proxy_http_version 1.1;
|
||||||
proxy_set_header Upgrade $http_upgrade;
|
proxy_set_header Upgrade $http_upgrade;
|
||||||
proxy_set_header Connection "upgrade";
|
proxy_set_header Connection "upgrade";
|
||||||
}
|
}
|
||||||
|
|||||||
+2
-4
@@ -1,13 +1,11 @@
|
|||||||
worker_processes auto;
|
worker_processes auto;
|
||||||
pid /tmp/nginx.pid;
|
pid /tmp/nginx.pid;
|
||||||
|
|
||||||
events
|
events {
|
||||||
{
|
|
||||||
worker_connections 1024;
|
worker_connections 1024;
|
||||||
}
|
}
|
||||||
|
|
||||||
http
|
http {
|
||||||
{
|
|
||||||
sendfile on;
|
sendfile on;
|
||||||
tcp_nopush on;
|
tcp_nopush on;
|
||||||
tcp_nodelay on;
|
tcp_nodelay on;
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
[supervisord]
|
[supervisord]
|
||||||
user=user
|
user=root
|
||||||
logfile = /tmp/supervisord.log
|
logfile = /tmp/supervisord.log
|
||||||
loglevel = debug
|
loglevel = debug
|
||||||
pidfile = /tmp/supervisord.pid
|
pidfile = /tmp/supervisord.pid
|
||||||
@@ -18,9 +18,5 @@ command=/usr/sbin/nginx -g 'daemon off;'
|
|||||||
autostart=true
|
autostart=true
|
||||||
autorestart=true
|
autorestart=true
|
||||||
|
|
||||||
[program:app]
|
|
||||||
directory=%(ENV_TFW_APP_DIR)s
|
|
||||||
command=python3 app.py
|
|
||||||
|
|
||||||
[include]
|
[include]
|
||||||
files=%(ENV_TFW_SUPERVISORD_COMPONENTS)s/*.conf
|
files=%(ENV_TFW_SUPERVISORD_COMPONENTS)s/*.conf
|
||||||
|
|||||||
Reference in New Issue
Block a user