diff --git a/.gitignore b/.gitignore index 9f11b75..2483976 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,2 @@ .idea/ +__pycache__/ diff --git a/msg_montainer.py b/msg_montainer.py new file mode 100644 index 0000000..c6e9428 --- /dev/null +++ b/msg_montainer.py @@ -0,0 +1,13 @@ + + +class MsgContainer: + def __init__(self, name: str, timestamp: float = 0, data: bytes = b''): + self.name : str= name + self.timestamp : float = timestamp + self.data: bytes = data + + def __str__(self): + return f"{self.name}: {self.timestamp}: {self.data}" + + def __repr__(self): + return f"{self.name}: {self.timestamp}: {self.data}" diff --git a/msg_sink.py b/msg_sink.py new file mode 100644 index 0000000..d63cb76 --- /dev/null +++ b/msg_sink.py @@ -0,0 +1,10 @@ +from abc import ABC, abstractmethod +from msg_montainer import MsgContainer + +class MsgSink(ABC): + def __init__(self): + pass + + @abstractmethod + def on_recv(self, msg: MsgContainer): + pass diff --git a/msg_source.py b/msg_source.py new file mode 100644 index 0000000..5569b87 --- /dev/null +++ b/msg_source.py @@ -0,0 +1,15 @@ +from abc import ABC +from msg_montainer import MsgContainer +from msg_sink import MsgSink + +class MsgSource(ABC): + def __init__(self): + self.sinks = {} + + def register_recv(self, name: str, sink: MsgSink): + self.sinks[name] = {'sink': sink} + + def call_sinks(self, msg: MsgContainer): + if msg.name in self.sinks: + sink = self.sinks[msg.name]['sink'] + sink.on_recv(msg) diff --git a/multi_client.py b/multi_client.py index 4bd3bc4..16841ba 100644 --- a/multi_client.py +++ b/multi_client.py @@ -2,15 +2,27 @@ import socket import selectors import time -class MultiClient(object): +from msg_montainer import MsgContainer +from nmea import NmeaReceiver +from msg_sink import MsgSink +from msg_source import MsgSource +from packetizer import Packetizer + + +class MultiClient(MsgSource): def __init__(self, servers: list[dict]): - self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port'])} for s in servers] + 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 + def register_recv(self, name: str, rx: MsgSink): + self.receivers[name] = {'rx': rx} + def find_by_addr(self, addr): for sock in self.sock_list: if sock['addr'] == addr: @@ -29,16 +41,6 @@ class MultiClient(object): for sock in self.sock_list: sock['sock'].close() - def on_receive(self, rx_list: list): - for data in rx_list: - print(f"{self}: on_receive: 'Time': {data['timestamp'], data['name']}") - - def send_request(self): - events = selectors.EVENT_READ | selectors.EVENT_WRITE - for sock in self.sock_list: - self.sel.register(sock['sock'], events, data=None) - print("receive") - def event_loop(self): try: while True: @@ -46,12 +48,14 @@ class MultiClient(object): rx_data = [] for key, mask in events: descr = self.find_by_addr(key.fileobj.getpeername()) - if mask & selectors.EVENT_READ: - timestamp = time.time() - data = key.fileobj.recv(1024) - rx_data.append({'timestamp': timestamp, 'name': descr['name'], 'data': data}) + name = descr['name'] - self.on_receive(rx_data) + 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)) except KeyboardInterrupt: print("Caught keyboard interrupt, exiting") @@ -65,8 +69,32 @@ if __name__ == "__main__": {'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731} ] + # TCPIP multi source receiver mc = MultiClient(server_list) + + # NMEA multi source receiver + receiver = NmeaReceiver() +# receiver.filter_by_msg_type('.*') + receiver.filter_by_msg_type('^.?GPGSV.*') + receiver.filter_by_msg_type('^.?GNGLL.*') + + # 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) + + # Connect all sources mc.connect() + + # receive and block mc.event_loop() - time.sleep(1) + + # code never reached mc.disconnect() diff --git a/nmea.py b/nmea.py new file mode 100644 index 0000000..99cf433 --- /dev/null +++ b/nmea.py @@ -0,0 +1,20 @@ +from urllib.parse import to_bytes + +from msg_montainer import MsgContainer +from msg_sink import MsgSink +import re + +class NmeaReceiver(MsgSink): + def __init__(self): + MsgSink.__init__(self) + self.filters_msg_type = [] + + def filter_by_msg_type(self, msg_type: str): + self.filters_msg_type.append(msg_type) + + def on_recv(self, msg: MsgContainer): + data = msg.data.decode(encoding="utf-8") + for pattern in self.filters_msg_type: + if re.match(pattern, data): + print(f"{self}: on_receive: {msg}") + diff --git a/packetizer.py b/packetizer.py new file mode 100644 index 0000000..8645fe5 --- /dev/null +++ b/packetizer.py @@ -0,0 +1,29 @@ +from msg_montainer import MsgContainer +from msg_sink import MsgSink +from msg_source import MsgSource +from struct import pack + +class Packetizer(MsgSink, MsgSource): + def __init__(self): + MsgSink.__init__(self) + MsgSource.__init__(self) + self.packet = b'' + self.wait_sync = True + + def on_recv(self, msg: MsgContainer): + name = msg.name + data = msg.data + for n in range (len(data)): + d = data[n] + if d == 13 or d == 10: + timestamp = msg.timestamp + self.wait_sync = False + if len(self.packet) > 0: + self.call_sinks(MsgContainer(name, timestamp, self.packet)) + self.wait_sync = True + self.packet = b'' + elif not self.wait_sync: + self.packet += pack('B', d) + + +