From eaec8d8691e58cff5a7fabc8fc81985e4af4f4f8 Mon Sep 17 00:00:00 2001 From: Jens Ahrensfeld Date: Tue, 17 Mar 2026 19:09:48 +0100 Subject: [PATCH] implemented send --- a_backend.py | 4 +++- network_backend.py | 54 ++++++++++++++++++++++++++++++++++++++-------- transceiver.py | 4 ++++ 3 files changed, 52 insertions(+), 10 deletions(-) diff --git a/a_backend.py b/a_backend.py index 6dbecf0..ac3eb76 100644 --- a/a_backend.py +++ b/a_backend.py @@ -1,4 +1,6 @@ from abc import ABC +from msg_container import MsgContainer class ABackend(ABC): - pass \ No newline at end of file + def send(self, name: str, msg: MsgContainer): + pass diff --git a/network_backend.py b/network_backend.py index a36d9da..37d232d 100644 --- a/network_backend.py +++ b/network_backend.py @@ -13,12 +13,29 @@ from nmea.GSA import Gsa from nmea.GSV import Gsv from nmea.GGA import Gga from nmea.RMC import Rmc +from threading import Thread +from queue import Queue class NetworkBackend(ABackend): def __init__(self, servers: list[dict]): self.sel = None - self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port']), 'xcvr': None} for s in servers] + self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port']), 'xcvr': None, 'queue': Queue()} for s in servers] + self.thread = Thread(target=self.event_loop) + self.loop_enable = True + + def send(self, name: str, msg: MsgContainer): + sock = self.find_by_name(name) + sock['queue'].put(msg) + if sock is not None: + self.sel.modify(sock['sock'], selectors.EVENT_READ + selectors.EVENT_WRITE, None) + + def start(self): + self.thread.start() + + def stop(self): + self.loop_enable = False + self.thread.join() def create_xcvr(self, name: str) -> Transceiver: obj = Transceiver(self, name) @@ -67,19 +84,29 @@ class NetworkBackend(ABackend): for sock in self.sock_list: sock['sock'].close() + def on_send(self, name: str, data: bytes): + sock = self.find_by_name(name) + if sock is not None: + self.sel.modify(sock['sock'], selectors.EVENT_READ + selectors.EVENT_WRITE, data) + def event_loop(self): try: - while True: + while self.loop_enable: events = self.sel.select(timeout=100) for key, mask in events: descr = self.find_by_addr(key.fileobj.getpeername()) name = descr['name'] 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)) + sock['xcvr'].on_recv(MsgContainer(name, timestamp, data)) + if mask & selectors.EVENT_WRITE: + data = sock['queue'].get() + key.fileobj.send(data) + print(f"Writing {key.data} to {sock['name']}") + if sock['queue'].empty(): + self.sel.modify(sock['sock'], selectors.EVENT_READ, None) except KeyboardInterrupt: print("Caught keyboard interrupt, exiting") @@ -92,6 +119,7 @@ if __name__ == "__main__": {'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731} ] + transceiver = {} # TCPIP multi source receiver mc = NetworkBackend(server_list) for grm in [" neo-f9p", "zed-x20p"]: @@ -117,15 +145,15 @@ if __name__ == "__main__": # Receiver-1.1 / # Create transceiver - transceiver = mc.create_xcvr(grm) + transceiver[grm] = mc.create_xcvr(grm) # Create NMEA-Packetizer nmea_packetizer = NmeaPacketizer() ubx_packetizer = UbxPacketizer() # Register NMEA-Packetizer to transceiver - transceiver.register_recv(nmea_packetizer) - transceiver.register_recv(ubx_packetizer) + transceiver[grm].register_recv(nmea_packetizer) + transceiver[grm].register_recv(ubx_packetizer) # Register Receiver to NMEA-Packetizer for rx in nmea_rx_list: @@ -139,8 +167,16 @@ if __name__ == "__main__": # Connect all sources mc.connect() - # receive and block - mc.event_loop() + # Start thread + mc.start() + + for n in range(10): + for grm in [" neo-f9p", "zed-x20p"]: + transceiver[grm].send(b'Hallo') + time.sleep(1) + + # Stop thread + mc.stop() # code never reached mc.disconnect() diff --git a/transceiver.py b/transceiver.py index 684d38a..dd63128 100644 --- a/transceiver.py +++ b/transceiver.py @@ -10,3 +10,7 @@ class Transceiver(MsgSource): def on_recv(self, msg: MsgContainer): self.call_sinks(msg) + + def send(self, msg: MsgContainer): + self.backend.send(self.name, msg) +