implemented send
This commit is contained in:
@@ -1,4 +1,6 @@
|
|||||||
from abc import ABC
|
from abc import ABC
|
||||||
|
from msg_container import MsgContainer
|
||||||
|
|
||||||
class ABackend(ABC):
|
class ABackend(ABC):
|
||||||
|
def send(self, name: str, msg: MsgContainer):
|
||||||
pass
|
pass
|
||||||
+45
-9
@@ -13,12 +13,29 @@ from nmea.GSA import Gsa
|
|||||||
from nmea.GSV import Gsv
|
from nmea.GSV import Gsv
|
||||||
from nmea.GGA import Gga
|
from nmea.GGA import Gga
|
||||||
from nmea.RMC import Rmc
|
from nmea.RMC import Rmc
|
||||||
|
from threading import Thread
|
||||||
|
from queue import Queue
|
||||||
|
|
||||||
|
|
||||||
class NetworkBackend(ABackend):
|
class NetworkBackend(ABackend):
|
||||||
def __init__(self, servers: list[dict]):
|
def __init__(self, servers: list[dict]):
|
||||||
self.sel = None
|
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:
|
def create_xcvr(self, name: str) -> Transceiver:
|
||||||
obj = Transceiver(self, name)
|
obj = Transceiver(self, name)
|
||||||
@@ -67,19 +84,29 @@ class NetworkBackend(ABackend):
|
|||||||
for sock in self.sock_list:
|
for sock in self.sock_list:
|
||||||
sock['sock'].close()
|
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):
|
def event_loop(self):
|
||||||
try:
|
try:
|
||||||
while True:
|
while self.loop_enable:
|
||||||
events = self.sel.select(timeout=100)
|
events = self.sel.select(timeout=100)
|
||||||
for key, mask in events:
|
for key, mask in events:
|
||||||
descr = self.find_by_addr(key.fileobj.getpeername())
|
descr = self.find_by_addr(key.fileobj.getpeername())
|
||||||
name = descr['name']
|
name = descr['name']
|
||||||
sock = self.find_by_name(name)
|
sock = self.find_by_name(name)
|
||||||
rx = sock['xcvr']
|
|
||||||
if mask & selectors.EVENT_READ:
|
if mask & selectors.EVENT_READ:
|
||||||
timestamp = time.time()
|
timestamp = time.time()
|
||||||
data = key.fileobj.recv(16)
|
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:
|
except KeyboardInterrupt:
|
||||||
print("Caught keyboard interrupt, exiting")
|
print("Caught keyboard interrupt, exiting")
|
||||||
@@ -92,6 +119,7 @@ if __name__ == "__main__":
|
|||||||
{'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731}
|
{'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731}
|
||||||
]
|
]
|
||||||
|
|
||||||
|
transceiver = {}
|
||||||
# TCPIP multi source receiver
|
# TCPIP multi source receiver
|
||||||
mc = NetworkBackend(server_list)
|
mc = NetworkBackend(server_list)
|
||||||
for grm in [" neo-f9p", "zed-x20p"]:
|
for grm in [" neo-f9p", "zed-x20p"]:
|
||||||
@@ -117,15 +145,15 @@ if __name__ == "__main__":
|
|||||||
# Receiver-1.1 /
|
# Receiver-1.1 /
|
||||||
|
|
||||||
# Create transceiver
|
# Create transceiver
|
||||||
transceiver = mc.create_xcvr(grm)
|
transceiver[grm] = mc.create_xcvr(grm)
|
||||||
|
|
||||||
# Create NMEA-Packetizer
|
# Create NMEA-Packetizer
|
||||||
nmea_packetizer = NmeaPacketizer()
|
nmea_packetizer = NmeaPacketizer()
|
||||||
ubx_packetizer = UbxPacketizer()
|
ubx_packetizer = UbxPacketizer()
|
||||||
|
|
||||||
# Register NMEA-Packetizer to transceiver
|
# Register NMEA-Packetizer to transceiver
|
||||||
transceiver.register_recv(nmea_packetizer)
|
transceiver[grm].register_recv(nmea_packetizer)
|
||||||
transceiver.register_recv(ubx_packetizer)
|
transceiver[grm].register_recv(ubx_packetizer)
|
||||||
|
|
||||||
# Register Receiver to NMEA-Packetizer
|
# Register Receiver to NMEA-Packetizer
|
||||||
for rx in nmea_rx_list:
|
for rx in nmea_rx_list:
|
||||||
@@ -139,8 +167,16 @@ if __name__ == "__main__":
|
|||||||
# Connect all sources
|
# Connect all sources
|
||||||
mc.connect()
|
mc.connect()
|
||||||
|
|
||||||
# receive and block
|
# Start thread
|
||||||
mc.event_loop()
|
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
|
# code never reached
|
||||||
mc.disconnect()
|
mc.disconnect()
|
||||||
|
|||||||
@@ -10,3 +10,7 @@ class Transceiver(MsgSource):
|
|||||||
|
|
||||||
def on_recv(self, msg: MsgContainer):
|
def on_recv(self, msg: MsgContainer):
|
||||||
self.call_sinks(msg)
|
self.call_sinks(msg)
|
||||||
|
|
||||||
|
def send(self, msg: MsgContainer):
|
||||||
|
self.backend.send(self.name, msg)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user