diff --git a/a_backend.py b/a_backend.py new file mode 100644 index 0000000..6dbecf0 --- /dev/null +++ b/a_backend.py @@ -0,0 +1,4 @@ +from abc import ABC + +class ABackend(ABC): + pass \ No newline at end of file diff --git a/msg_montainer.py b/msg_container.py similarity index 100% rename from msg_montainer.py rename to msg_container.py diff --git a/msg_receiver.py b/msg_receiver.py new file mode 100644 index 0000000..2951f1e --- /dev/null +++ b/msg_receiver.py @@ -0,0 +1,20 @@ +from msg_container import MsgContainer +from msg_sink import MsgSink +import re + +class MsgReceiver(MsgSink): + def __init__(self, name: str, filter: str): + super().__init__() + self.name = name + self.filter = filter + + def is_msg(self, msg: MsgContainer): + try: + data = msg.data.decode(encoding="utf-8") + return re.match(self.filter, data) + except UnicodeDecodeError: + pass + + def on_recv(self, msg: MsgContainer): + print(f"{self.name}: {msg}") + diff --git a/msg_sink.py b/msg_sink.py index d63cb76..ef721b7 100644 --- a/msg_sink.py +++ b/msg_sink.py @@ -1,5 +1,5 @@ from abc import ABC, abstractmethod -from msg_montainer import MsgContainer +from msg_container import MsgContainer class MsgSink(ABC): def __init__(self): @@ -8,3 +8,6 @@ class MsgSink(ABC): @abstractmethod def on_recv(self, msg: MsgContainer): pass + + def is_msg(self, msg: MsgContainer): + return False diff --git a/msg_source.py b/msg_source.py index 5569b87..e1b72f8 100644 --- a/msg_source.py +++ b/msg_source.py @@ -1,15 +1,15 @@ from abc import ABC -from msg_montainer import MsgContainer +from msg_container import MsgContainer from msg_sink import MsgSink class MsgSource(ABC): def __init__(self): - self.sinks = {} + self.sinks = [] - def register_recv(self, name: str, sink: MsgSink): - self.sinks[name] = {'sink': sink} + def register_recv(self, sink: MsgSink): + self.sinks.append(sink) def call_sinks(self, msg: MsgContainer): - if msg.name in self.sinks: - sink = self.sinks[msg.name]['sink'] - sink.on_recv(msg) + for sink in self.sinks: + if sink.is_msg(msg): + sink.on_recv(msg) diff --git a/multi_client.py b/network_backend.py similarity index 50% rename from multi_client.py rename to network_backend.py index 34a207b..7136459 100644 --- a/multi_client.py +++ b/network_backend.py @@ -2,34 +2,53 @@ import socket import selectors import time -from msg_montainer import MsgContainer -from nmea import NmeaReceiver -from msg_sink import MsgSink -from msg_source import MsgSource +from msg_container import MsgContainer from packetizer import Packetizer +from transceiver import Transceiver +from a_backend import ABackend +from msg_receiver import MsgReceiver - -class MultiClient(MsgSource): +class NetworkBackend(ABackend): def __init__(self, servers: list[dict]): - MsgSource.__init__(self) self.sel = None - self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port'])} for s in servers] - self.receivers = {} - for server in self.sock_list: - sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - sock.setblocking(False) - server['sock'] = sock + self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port']), 'xcvr': None} for s in servers] - def register_recv(self, name: str, rx: MsgSink): - self.receivers[name] = {'rx': rx} + def create_xcvr(self, name: str) -> Transceiver: + obj = Transceiver(self, name) + if not self.register_xcvr(obj): + raise Exception(f"No socket found for \"{name}\"") + return obj + self.sel.register(name, selectors.EVENT_READ, data=None) + def register_xcvr(self, xcvr: Transceiver): + found = False + for sock in self.sock_list: + if xcvr.name in sock['name']: + found = True + sock['xcvr'] = xcvr + break + + return found def find_by_addr(self, addr): for sock in self.sock_list: if sock['addr'] == addr: return sock return None + def find_by_name(self, name: str): + for sock in self.sock_list: + if sock['name'] == name: + return sock + return None + + def _prepare(self): + for server in self.sock_list: + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.setblocking(False) + server['sock'] = sock + def connect(self): + self._prepare() self.sel = selectors.DefaultSelector() events = selectors.EVENT_READ for sock in self.sock_list: @@ -44,30 +63,22 @@ class MultiClient(MsgSource): def event_loop(self): try: while True: - events = self.sel.select(timeout=1) - rx_data = [] + events = self.sel.select(timeout=100) for key, mask in events: descr = self.find_by_addr(key.fileobj.getpeername()) name = descr['name'] - - if name in self.receivers: - rx = self.receivers[name]['rx'] - if mask & selectors.EVENT_READ: - timestamp = time.time() - data = key.fileobj.recv(16) - rx.on_recv(MsgContainer(name, timestamp, data)) + sock = self.find_by_name(name) + rx = sock['xcvr'] + if mask & selectors.EVENT_READ: + timestamp = time.time() + data = key.fileobj.recv(16) + rx.on_recv(MsgContainer(name, timestamp, data)) except KeyboardInterrupt: print("Caught keyboard interrupt, exiting") finally: self.sel.close() -def GPGSV_callback(msg: MsgContainer): - print(f"GPGSV_callback: {msg}") - -def GNGLL_callback(msg: MsgContainer): - print(f"GNGLL_callback: {msg}") - if __name__ == "__main__": server_list = [ {'name': "neo-f9p", 'host': "192.168.22.93", 'port': 8721}, @@ -75,25 +86,26 @@ if __name__ == "__main__": ] # TCPIP multi source receiver - mc = MultiClient(server_list) - - # NMEA multi source receiver - receiver = NmeaReceiver() -# receiver.filter_by_msg_type('.*') - receiver.register_msg('^.?GPGSV.*', GPGSV_callback) - receiver.register_msg('^.?GNGLL.*', GNGLL_callback) + mc = NetworkBackend(server_list) + xcvr_1 = mc.create_xcvr("neo-f9p") + xcvr_2 = mc.create_xcvr("zed-x20p") # Packetizer for "zed-x20p" packetizer_0 = Packetizer() - packetizer_0.register_recv("zed-x20p", receiver) # Packetizer for "neo-f9p" packetizer_1 = Packetizer() - packetizer_1.register_recv("neo-f9p", receiver) - # Register both packetizers at multi source receiver - mc.register_recv("zed-x20p", packetizer_0) - mc.register_recv("neo-f9p", packetizer_1) + xcvr_1.register_recv(packetizer_0) + xcvr_2.register_recv(packetizer_1) + rx_1 = MsgReceiver('Rx1', '^.?GPGSV.*') + rx_2 = MsgReceiver('Rx2', '^.?GNGLL.*') + + # Register all Rx to all sources + packetizer_0.register_recv(rx_1) + packetizer_0.register_recv(rx_2) + packetizer_1.register_recv(rx_1) + packetizer_1.register_recv(rx_2) # Connect all sources mc.connect() diff --git a/nmea.py b/nmea.py deleted file mode 100644 index fced200..0000000 --- a/nmea.py +++ /dev/null @@ -1,23 +0,0 @@ -from msg_montainer import MsgContainer -from msg_sink import MsgSink -from typing import Callable -import re - -class NmeaReceiver(MsgSink): - def __init__(self): - MsgSink.__init__(self) - self.filters = [] - - def on_recv(self, msg: MsgContainer): - data = msg.data.decode(encoding="utf-8") - for filter in self.filters: - if re.match(filter['re'], data): - if filter['listener'] is not None: - filter['listener'](msg) - - def register_msg(self, msg_type: str, listener: Callable): - self.filters.append({'re': msg_type, 'listener': listener}) - - - def _filter_by_msg_type(self, msg_type: str): - self.filters.append(msg_type) diff --git a/packetizer.py b/packetizer.py index 8645fe5..91de594 100644 --- a/packetizer.py +++ b/packetizer.py @@ -1,4 +1,4 @@ -from msg_montainer import MsgContainer +from msg_container import MsgContainer from msg_sink import MsgSink from msg_source import MsgSource from struct import pack @@ -27,3 +27,5 @@ class Packetizer(MsgSink, MsgSource): + def is_msg(self, msg: MsgContainer): + return True \ No newline at end of file diff --git a/transceiver.py b/transceiver.py new file mode 100644 index 0000000..684d38a --- /dev/null +++ b/transceiver.py @@ -0,0 +1,12 @@ +from msg_source import MsgSource +from a_backend import ABackend +from msg_container import * + +class Transceiver(MsgSource): + def __init__(self, backend: ABackend, name: str): + MsgSource.__init__(self) + self.backend = backend + self.name = name + + def on_recv(self, msg: MsgContainer): + self.call_sinks(msg)