diff --git a/a_transceiver.py b/a_transceiver.py new file mode 100644 index 0000000..9c304f3 --- /dev/null +++ b/a_transceiver.py @@ -0,0 +1,20 @@ +from a_backend import ABackend +from msg.container import MsgContainer +from msg.listener import MsgListener + + +class ATransceiver(MsgListener): + def __init__(self, name: str): + MsgListener.__init__(self) + self.name = name + self.backend = None + + def on_recv(self, msg: MsgContainer): + pass + + def on_register(self, backend: ABackend): + self.backend = backend + + def send(self, data: bytes): + self.backend.send(self.name, data) + diff --git a/msg/container.py b/msg/container.py index 30e9ed1..fd968a5 100644 --- a/msg/container.py +++ b/msg/container.py @@ -1,10 +1,10 @@ - +import time class MsgContainer: - def __init__(self, name: str, timestamp: float = 0, data = b''): + def __init__(self, data: bytes, name: str = "default", timestamp=time.time()): self.timestamp : float = timestamp - self.name : str= name self.data: bytes = data + self.name : str = name def __str__(self): return f"{self.timestamp}:{self.name}:{self.data}" diff --git a/network_backend.py b/network_backend.py index c2d9f07..eeca698 100644 --- a/network_backend.py +++ b/network_backend.py @@ -5,6 +5,7 @@ import signal from threading import Thread from queue import Queue +from msg.container import MsgContainer from transceiver import Transceiver from a_backend import ABackend @@ -97,7 +98,7 @@ class NetworkBackend(ABackend): if mask & selectors.EVENT_READ: timestamp = time.time() data = key.fileobj.recv(64) - sock['xcvr'].on_recv(data) + sock['xcvr'].on_recv(MsgContainer(data)) if mask & selectors.EVENT_WRITE: data = sock['queue'].get() key.fileobj.send(data) diff --git a/nmea/packetizer.py b/nmea/packetizer.py index bd7dffd..04694a7 100644 --- a/nmea/packetizer.py +++ b/nmea/packetizer.py @@ -1,23 +1,25 @@ from msg.container import MsgContainer from msg.talker import MsgTalker -from stream.listener import StreamListener +from msg.listener import MsgListener from struct import pack -import time -class NmeaPacketizer(StreamListener, MsgTalker): +class NmeaPacketizer(MsgListener, MsgTalker): def __init__(self, name: str): - StreamListener.__init__(self) + MsgListener.__init__(self) MsgTalker.__init__(self) self.name = name self.packet = b'' self.wait_sync = True - def on_recv(self, rx_data: bytes): - for d in rx_data: + def is_msg(self, msg: MsgContainer): + return True + + def on_recv(self, msg: MsgContainer): + for d in msg.data: if d == 13 or d == 10: self.wait_sync = False if len(self.packet) > 0: - self.call_listener(MsgContainer(self.name, time.time(), self.packet)) + self.call_listener(MsgContainer(self.packet, self.name)) self.wait_sync = True self.packet = b'' elif not self.wait_sync: diff --git a/stream/listener.py b/stream/listener.py deleted file mode 100644 index b8b4268..0000000 --- a/stream/listener.py +++ /dev/null @@ -1,9 +0,0 @@ -from abc import ABC, abstractmethod - -class StreamListener(ABC): - def __init__(self): - pass - - @abstractmethod - def on_recv(self, data: bytes): - pass diff --git a/stream/talker.py b/stream/talker.py deleted file mode 100644 index ab555d6..0000000 --- a/stream/talker.py +++ /dev/null @@ -1,13 +0,0 @@ -from abc import ABC -from stream.listener import StreamListener - -class StreamTalker(ABC): - def __init__(self): - self.sinks = [] - - def register_listener(self, sink: StreamListener): - self.sinks.append(sink) - - def call_listener(self, data: bytes): - for sink in self.sinks: - sink.on_recv(data) diff --git a/transceiver.py b/transceiver.py index 0ebc77a..35e0204 100644 --- a/transceiver.py +++ b/transceiver.py @@ -1,21 +1,13 @@ -from a_backend import ABackend -from stream.talker import StreamTalker -from stream.listener import StreamListener +from a_transceiver import ATransceiver +from msg.container import MsgContainer +from msg.talker import MsgTalker -class Transceiver(StreamListener, StreamTalker): +class Transceiver(ATransceiver, MsgTalker): def __init__(self, name: str): - StreamListener.__init__(self) - StreamTalker.__init__(self) - self.name = name - self.backend = None + ATransceiver.__init__(self, name) + MsgTalker.__init__(self) - def on_recv(self, data: bytes): - self.call_listener(data) - - def on_register(self, backend: ABackend): - self.backend = backend - - def send(self, data: bytes): - self.backend.send(self.name, data) + def on_recv(self, msg: MsgContainer): + self.call_listener(msg) diff --git a/ubx/packetizer.py b/ubx/packetizer.py index 37ff40c..3663df9 100644 --- a/ubx/packetizer.py +++ b/ubx/packetizer.py @@ -3,7 +3,6 @@ import time from msg.container import MsgContainer from msg.listener import MsgListener from msg.talker import MsgTalker -from stream.listener import StreamListener from ubx.ubx import UBX_SYNC_WORD, frame_parse, frame_create from ubx.listener import UbxListener @@ -12,9 +11,9 @@ def ubx_pkt_debug(s: str): if UBX_PKT_DEBUG: print(s) -class UbxPacketizer(StreamListener, MsgTalker): +class UbxPacketizer(MsgListener, MsgTalker): def __init__(self, name: str, sync_bytes: bytes = UBX_SYNC_WORD): - StreamListener.__init__(self) + MsgListener.__init__(self) MsgTalker.__init__(self) self.name = name self.sync_bytes = sync_bytes @@ -28,8 +27,11 @@ class UbxPacketizer(StreamListener, MsgTalker): self.wait_sync = True self.sync_count = 0 - def on_recv(self, rx_data: bytes): - for d in rx_data: + def is_msg(self, msg: MsgContainer): + return True + + def on_recv(self, rx_data: MsgContainer): + for d in rx_data.data: if d == self.sync_bytes[self.sync_count]: self.d_save += struct.pack('B', d) self.sync_count += 1 @@ -51,7 +53,7 @@ class UbxPacketizer(StreamListener, MsgTalker): data, hdr = frame_parse(self.packet) if data is not None: if hdr.length == len(data): - self.call_listener(MsgContainer(self.name, time.time(), self.packet)) + self.call_listener(MsgContainer(self.packet, self.name)) self.packet = b'' self.wait_sync = True