Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES
Original file line number Diff line number Diff line change
Expand Up @@ -3,3 +3,4 @@ Version 0.2.0 - released mm/dd/18
- Initial release
- Change Android Controller -> WebServiceSerial: Now use WebService to reuse WebService methods
- Defined Target: Initial target is Android, but is possible implements other targets for communication with other devices, as Arduino via USB Serial communication
- Try reconnect when connection is closed
4 changes: 4 additions & 0 deletions webservice_serial/target/target.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,10 @@ def __init__(self):
self.application = None
self.port = None

@property
def name(self):
return self.__class__.__name__

def init(self, application, port):
"""
Target initialization
Expand Down
23 changes: 17 additions & 6 deletions webservice_serial/webservice_serial.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
from webservice_serial.webservice_serial_client import WebServiceSerialClient
from webservice_serial.websocket_client import WebSocketClient

from time import sleep


class WebServiceSerial(Component):
port = 8888
Expand All @@ -38,15 +40,21 @@ def __init__(self, application, target, ws_port=3000):
def init(self):
self._client.connected_listener = self._on_connected
self._client.message_listener = self._process_message
self._client.disconnected_listener = lambda: self._try_connect(5)

self.request_message_processor.processed_listener = self._on_processed

self.target.init(self.application, WebServiceSerial.port)

self._websocket_client.token_defined_listener = self._on_token_defined
self._websocket_client.message_listener = self._on_event

self._websocket_client.connect()

self._try_connect()

def _try_connect(self, delay=0):
self._log('Trying to connect with {}', self.target.name)
self.target.init(self.application, WebServiceSerial.port)
sleep(delay)
self._client.connect()

def _on_token_defined(self, token):
Expand All @@ -59,22 +67,25 @@ def close(self):
self._client.close()

def _on_connected(self):
self.application.log('AndroidController - DisplayView connected')
self._log('{} connected', self.target.name)

def _process_message(self, message):
self.application.log('AndroidController - Message received: {}', message)
self._log('Message received: {}', message)

self.request_message_processor.process(message)

def _on_processed(self, request_message, response_message):
response_message = self.target.process(request_message, response_message)

self.application.log('AndroidController - Message sent: {}', response_message)
self._log('Message sent: {}', response_message)
self._client.send(response_message)

def _on_event(self, message):
response_message = self.request_message_processor.process_event(message)
response_message = self.target.process(None, response_message)

self.application.log('AndroidController - Message sent: {}', response_message)
self._log('Message sent: {}', response_message)
self._client.send(response_message)

def _log(self, message, *args, **kwargs):
self.application.log('{} - {}'.format(self.__class__.__name__, message), *args, **kwargs)
38 changes: 34 additions & 4 deletions webservice_serial/webservice_serial_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from tornado import gen
from tornado.ioloop import IOLoop
from tornado.tcpclient import TCPClient
from tornado.iostream import StreamClosedError

from webservice_serial.protocol.message_builder import MessageBuilder

Expand All @@ -30,26 +31,55 @@ def __init__(self, address, port, encoding="utf-8"):
self.message_listener = lambda message: ...
self.connected_listener = lambda: ...

self.disconnected_listener = lambda: print('Disconnected :(')

def connect(self):
IOLoop.current().spawn_callback(lambda: self._connect())

@gen.coroutine
def _connect(self):
self.stream = yield TCPClient().connect(self.address, self.port)
self.stream = yield self._try_connect()
if self.stream is None:
return

self.connected_listener()
yield self._start_read_data()

@gen.coroutine
def _try_connect(self):
try:
stream = yield TCPClient().connect(self.address, self.port)
return stream
except StreamClosedError as e:
self.disconnected_listener()
return None

@gen.coroutine
def _start_read_data(self):
while True:
data = yield self.stream.read_until('\n'.encode(self.encoding))
data = yield self._read_data()
if data is None:
break
data = data.decode(self.encoding).strip()

generated = MessageBuilder.generate(data)
if generated is not None:
self.message_listener(generated)

@gen.coroutine
def _read_data(self):
try:
data = yield self.stream.read_until('\n'.encode(self.encoding))
except StreamClosedError as e:
self.disconnected_listener()
return None

return data

def send(self, message):
text = str(message).encode(self.encoding)
self.stream.write(text)

def close(self):
if self.stream is not None:
self.stream.close()
if self.stream is not None and not self.stream.closed():
self.stream.close()