- added stream talker

- added stream listener
This commit is contained in:
2026-03-28 12:03:54 +01:00
parent 2e01a70d53
commit 8208edfc46
6 changed files with 54 additions and 42 deletions
+4 -5
View File
@@ -5,11 +5,10 @@ import signal
from threading import Thread from threading import Thread
from queue import Queue from queue import Queue
from msg.container import MsgContainer
from transceiver import Transceiver from transceiver import Transceiver
from a_backend import ABackend from a_backend import ABackend
from nmea import NmeaPacketizer from nmea.packetizer import NmeaPacketizer
from nmea.messages import Gll, Gsa, Gsv, Gga, Rmc from nmea.messages import Gll, Gsa, Gsv, Gga, Rmc
from ubx.ubx import frame_create from ubx.ubx import frame_create
@@ -98,7 +97,7 @@ class NetworkBackend(ABackend):
if mask & selectors.EVENT_READ: if mask & selectors.EVENT_READ:
timestamp = time.time() timestamp = time.time()
data = key.fileobj.recv(64) data = key.fileobj.recv(64)
sock['xcvr'].on_recv(MsgContainer(name, timestamp, data)) sock['xcvr'].on_recv(data)
if mask & selectors.EVENT_WRITE: if mask & selectors.EVENT_WRITE:
data = sock['queue'].get() data = sock['queue'].get()
key.fileobj.send(data) key.fileobj.send(data)
@@ -165,8 +164,8 @@ if __name__ == "__main__":
mc.register_xcvr(transceiver[grm]) mc.register_xcvr(transceiver[grm])
# Create NMEA-Packetizer # Create NMEA-Packetizer
nmea_packetizer = NmeaPacketizer() nmea_packetizer = NmeaPacketizer(grm)
ubx_packetizer = UbxPacketizer() ubx_packetizer = UbxPacketizer(grm)
# Register NMEA-Packetizer to transceiver # Register NMEA-Packetizer to transceiver
transceiver[grm].register_listener(nmea_packetizer) transceiver[grm].register_listener(nmea_packetizer)
+9 -16
View File
@@ -1,31 +1,24 @@
from msg.container import MsgContainer from msg.container import MsgContainer
from msg.listener import MsgListener
from msg.talker import MsgTalker from msg.talker import MsgTalker
from stream.listener import StreamListener
from struct import pack from struct import pack
import time
class NmeaPacketizer(MsgListener, MsgTalker): class NmeaPacketizer(StreamListener, MsgTalker):
def __init__(self): def __init__(self, name: str):
MsgListener.__init__(self) StreamListener.__init__(self)
MsgTalker.__init__(self) MsgTalker.__init__(self)
self.name = name
self.packet = b'' self.packet = b''
self.wait_sync = True self.wait_sync = True
def on_recv(self, msg: MsgContainer): def on_recv(self, rx_data: bytes):
name = msg.name for d in rx_data:
data = msg.data
for n in range (len(data)):
d = data[n]
if d == 13 or d == 10: if d == 13 or d == 10:
timestamp = msg.timestamp
self.wait_sync = False self.wait_sync = False
if len(self.packet) > 0: if len(self.packet) > 0:
self.call_listener(MsgContainer(name, timestamp, self.packet)) self.call_listener(MsgContainer(self.name, time.time(), self.packet))
self.wait_sync = True self.wait_sync = True
self.packet = b'' self.packet = b''
elif not self.wait_sync: elif not self.wait_sync:
self.packet += pack('B', d) self.packet += pack('B', d)
def is_msg(self, msg: MsgContainer):
return True
+9
View File
@@ -0,0 +1,9 @@
from abc import ABC, abstractmethod
class StreamListener(ABC):
def __init__(self):
pass
@abstractmethod
def on_recv(self, data: bytes):
pass
+13
View File
@@ -0,0 +1,13 @@
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 -6
View File
@@ -1,15 +1,17 @@
from msg.talker import MsgTalker
from a_backend import ABackend from a_backend import ABackend
from msg.container import * from stream.talker import StreamTalker
from stream.listener import StreamListener
class Transceiver(MsgTalker):
class Transceiver(StreamListener, StreamTalker):
def __init__(self, name: str): def __init__(self, name: str):
MsgTalker.__init__(self) StreamListener.__init__(self)
StreamTalker.__init__(self)
self.name = name self.name = name
self.backend = None self.backend = None
def on_recv(self, msg: MsgContainer): def on_recv(self, data: bytes):
self.call_listener(msg) self.call_listener(data)
def on_register(self, backend: ABackend): def on_register(self, backend: ABackend):
self.backend = backend self.backend = backend
+11 -15
View File
@@ -1,7 +1,9 @@
import struct import struct
import time
from msg.container import MsgContainer from msg.container import MsgContainer
from msg.listener import MsgListener from msg.listener import MsgListener
from msg.talker import MsgTalker from msg.talker import MsgTalker
from stream.listener import StreamListener
from ubx.ubx import UBX_SYNC_WORD, frame_parse, frame_create from ubx.ubx import UBX_SYNC_WORD, frame_parse, frame_create
from ubx.listener import UbxListener from ubx.listener import UbxListener
@@ -10,15 +12,15 @@ def ubx_pkt_debug(s: str):
if UBX_PKT_DEBUG: if UBX_PKT_DEBUG:
print(s) print(s)
class UbxPacketizer(MsgListener, MsgTalker): class UbxPacketizer(StreamListener, MsgTalker):
def __init__(self, sync_bytes: bytes = UBX_SYNC_WORD): def __init__(self, name: str, sync_bytes: bytes = UBX_SYNC_WORD):
MsgListener.__init__(self) StreamListener.__init__(self)
MsgTalker.__init__(self) MsgTalker.__init__(self)
self.name = name
self.sync_bytes = sync_bytes self.sync_bytes = sync_bytes
self.packet = b'' self.packet = b''
self.wait_sync = True self.wait_sync = True
self.sync_count = 0 self.sync_count = 0
self.timestamp = 0
self.d_save = b'' self.d_save = b''
def reset(self): def reset(self):
@@ -26,21 +28,19 @@ class UbxPacketizer(MsgListener, MsgTalker):
self.wait_sync = True self.wait_sync = True
self.sync_count = 0 self.sync_count = 0
def on_recv(self, msg: MsgContainer): def on_recv(self, rx_data: bytes):
for n in range (0, len(msg.data)): for d in rx_data:
d = msg.data[n]
if d == self.sync_bytes[self.sync_count]: if d == self.sync_bytes[self.sync_count]:
self.d_save += struct.pack('B', d) self.d_save += struct.pack('B', d)
self.sync_count += 1 self.sync_count += 1
if self.sync_count == len(self.sync_bytes): if self.sync_count == len(self.sync_bytes):
ubx_pkt_debug(f"{msg.name}:UBX-sync") ubx_pkt_debug(f"{self.name}:UBX-sync")
if len(self.packet) > 0: if len(self.packet) > 0:
ubx_pkt_debug(f"{msg.name}:Remaining packet: {self.packet}") ubx_pkt_debug(f"{self.name}:Remaining packet: {self.packet}")
self.sync_count = 0 self.sync_count = 0
self.d_save = b'' self.d_save = b''
self.packet = b'' self.packet = b''
self.wait_sync = False self.wait_sync = False
self.timestamp = msg.timestamp
else: else:
if self.sync_count > 0: if self.sync_count > 0:
self.packet += self.d_save self.packet += self.d_save
@@ -51,14 +51,10 @@ class UbxPacketizer(MsgListener, MsgTalker):
data, hdr = frame_parse(self.packet) data, hdr = frame_parse(self.packet)
if data is not None: if data is not None:
if hdr.length == len(data): if hdr.length == len(data):
self.call_listener(MsgContainer(msg.name, self.timestamp, self.packet)) self.call_listener(MsgContainer(self.name, time.time(), self.packet))
self.packet = b'' self.packet = b''
self.wait_sync = True self.wait_sync = True
def is_msg(self, msg: MsgContainer):
return True
if __name__ == "__main__": if __name__ == "__main__":
UBX_PKT_DEBUG = True UBX_PKT_DEBUG = True
class MySink(MsgListener): class MySink(MsgListener):