reverted introducing stream

This commit is contained in:
2026-03-29 10:37:49 +02:00
parent 8208edfc46
commit af94c39b3e
8 changed files with 50 additions and 55 deletions
+20
View File
@@ -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)
+3 -3
View File
@@ -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}"
+2 -1
View File
@@ -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)
+9 -7
View File
@@ -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:
-9
View File
@@ -1,9 +0,0 @@
from abc import ABC, abstractmethod
class StreamListener(ABC):
def __init__(self):
pass
@abstractmethod
def on_recv(self, data: bytes):
pass
-13
View File
@@ -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)
+8 -16
View File
@@ -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)
+8 -6
View File
@@ -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