refactored
This commit is contained in:
@@ -0,0 +1,4 @@
|
|||||||
|
from abc import ABC
|
||||||
|
|
||||||
|
class ABackend(ABC):
|
||||||
|
pass
|
||||||
@@ -0,0 +1,20 @@
|
|||||||
|
from msg_container import MsgContainer
|
||||||
|
from msg_sink import MsgSink
|
||||||
|
import re
|
||||||
|
|
||||||
|
class MsgReceiver(MsgSink):
|
||||||
|
def __init__(self, name: str, filter: str):
|
||||||
|
super().__init__()
|
||||||
|
self.name = name
|
||||||
|
self.filter = filter
|
||||||
|
|
||||||
|
def is_msg(self, msg: MsgContainer):
|
||||||
|
try:
|
||||||
|
data = msg.data.decode(encoding="utf-8")
|
||||||
|
return re.match(self.filter, data)
|
||||||
|
except UnicodeDecodeError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
def on_recv(self, msg: MsgContainer):
|
||||||
|
print(f"{self.name}: {msg}")
|
||||||
|
|
||||||
+4
-1
@@ -1,5 +1,5 @@
|
|||||||
from abc import ABC, abstractmethod
|
from abc import ABC, abstractmethod
|
||||||
from msg_montainer import MsgContainer
|
from msg_container import MsgContainer
|
||||||
|
|
||||||
class MsgSink(ABC):
|
class MsgSink(ABC):
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
@@ -8,3 +8,6 @@ class MsgSink(ABC):
|
|||||||
@abstractmethod
|
@abstractmethod
|
||||||
def on_recv(self, msg: MsgContainer):
|
def on_recv(self, msg: MsgContainer):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
|
def is_msg(self, msg: MsgContainer):
|
||||||
|
return False
|
||||||
|
|||||||
+7
-7
@@ -1,15 +1,15 @@
|
|||||||
from abc import ABC
|
from abc import ABC
|
||||||
from msg_montainer import MsgContainer
|
from msg_container import MsgContainer
|
||||||
from msg_sink import MsgSink
|
from msg_sink import MsgSink
|
||||||
|
|
||||||
class MsgSource(ABC):
|
class MsgSource(ABC):
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.sinks = {}
|
self.sinks = []
|
||||||
|
|
||||||
def register_recv(self, name: str, sink: MsgSink):
|
def register_recv(self, sink: MsgSink):
|
||||||
self.sinks[name] = {'sink': sink}
|
self.sinks.append(sink)
|
||||||
|
|
||||||
def call_sinks(self, msg: MsgContainer):
|
def call_sinks(self, msg: MsgContainer):
|
||||||
if msg.name in self.sinks:
|
for sink in self.sinks:
|
||||||
sink = self.sinks[msg.name]['sink']
|
if sink.is_msg(msg):
|
||||||
sink.on_recv(msg)
|
sink.on_recv(msg)
|
||||||
|
|||||||
@@ -2,34 +2,53 @@ import socket
|
|||||||
import selectors
|
import selectors
|
||||||
import time
|
import time
|
||||||
|
|
||||||
from msg_montainer import MsgContainer
|
from msg_container import MsgContainer
|
||||||
from nmea import NmeaReceiver
|
|
||||||
from msg_sink import MsgSink
|
|
||||||
from msg_source import MsgSource
|
|
||||||
from packetizer import Packetizer
|
from packetizer import Packetizer
|
||||||
|
from transceiver import Transceiver
|
||||||
|
from a_backend import ABackend
|
||||||
|
from msg_receiver import MsgReceiver
|
||||||
|
|
||||||
|
class NetworkBackend(ABackend):
|
||||||
class MultiClient(MsgSource):
|
|
||||||
def __init__(self, servers: list[dict]):
|
def __init__(self, servers: list[dict]):
|
||||||
MsgSource.__init__(self)
|
|
||||||
self.sel = None
|
self.sel = None
|
||||||
self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port'])} for s in servers]
|
self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port']), 'xcvr': None} for s in servers]
|
||||||
self.receivers = {}
|
|
||||||
for server in self.sock_list:
|
|
||||||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
||||||
sock.setblocking(False)
|
|
||||||
server['sock'] = sock
|
|
||||||
|
|
||||||
def register_recv(self, name: str, rx: MsgSink):
|
def create_xcvr(self, name: str) -> Transceiver:
|
||||||
self.receivers[name] = {'rx': rx}
|
obj = Transceiver(self, name)
|
||||||
|
if not self.register_xcvr(obj):
|
||||||
|
raise Exception(f"No socket found for \"{name}\"")
|
||||||
|
return obj
|
||||||
|
|
||||||
|
self.sel.register(name, selectors.EVENT_READ, data=None)
|
||||||
|
def register_xcvr(self, xcvr: Transceiver):
|
||||||
|
found = False
|
||||||
|
for sock in self.sock_list:
|
||||||
|
if xcvr.name in sock['name']:
|
||||||
|
found = True
|
||||||
|
sock['xcvr'] = xcvr
|
||||||
|
break
|
||||||
|
|
||||||
|
return found
|
||||||
def find_by_addr(self, addr):
|
def find_by_addr(self, addr):
|
||||||
for sock in self.sock_list:
|
for sock in self.sock_list:
|
||||||
if sock['addr'] == addr:
|
if sock['addr'] == addr:
|
||||||
return sock
|
return sock
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
def find_by_name(self, name: str):
|
||||||
|
for sock in self.sock_list:
|
||||||
|
if sock['name'] == name:
|
||||||
|
return sock
|
||||||
|
return None
|
||||||
|
|
||||||
|
def _prepare(self):
|
||||||
|
for server in self.sock_list:
|
||||||
|
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||||
|
sock.setblocking(False)
|
||||||
|
server['sock'] = sock
|
||||||
|
|
||||||
def connect(self):
|
def connect(self):
|
||||||
|
self._prepare()
|
||||||
self.sel = selectors.DefaultSelector()
|
self.sel = selectors.DefaultSelector()
|
||||||
events = selectors.EVENT_READ
|
events = selectors.EVENT_READ
|
||||||
for sock in self.sock_list:
|
for sock in self.sock_list:
|
||||||
@@ -44,30 +63,22 @@ class MultiClient(MsgSource):
|
|||||||
def event_loop(self):
|
def event_loop(self):
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
events = self.sel.select(timeout=1)
|
events = self.sel.select(timeout=100)
|
||||||
rx_data = []
|
|
||||||
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)
|
||||||
if name in self.receivers:
|
rx = sock['xcvr']
|
||||||
rx = self.receivers[name]['rx']
|
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))
|
||||||
rx.on_recv(MsgContainer(name, timestamp, data))
|
|
||||||
|
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
print("Caught keyboard interrupt, exiting")
|
print("Caught keyboard interrupt, exiting")
|
||||||
finally:
|
finally:
|
||||||
self.sel.close()
|
self.sel.close()
|
||||||
|
|
||||||
def GPGSV_callback(msg: MsgContainer):
|
|
||||||
print(f"GPGSV_callback: {msg}")
|
|
||||||
|
|
||||||
def GNGLL_callback(msg: MsgContainer):
|
|
||||||
print(f"GNGLL_callback: {msg}")
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
server_list = [
|
server_list = [
|
||||||
{'name': "neo-f9p", 'host': "192.168.22.93", 'port': 8721},
|
{'name': "neo-f9p", 'host': "192.168.22.93", 'port': 8721},
|
||||||
@@ -75,25 +86,26 @@ if __name__ == "__main__":
|
|||||||
]
|
]
|
||||||
|
|
||||||
# TCPIP multi source receiver
|
# TCPIP multi source receiver
|
||||||
mc = MultiClient(server_list)
|
mc = NetworkBackend(server_list)
|
||||||
|
xcvr_1 = mc.create_xcvr("neo-f9p")
|
||||||
# NMEA multi source receiver
|
xcvr_2 = mc.create_xcvr("zed-x20p")
|
||||||
receiver = NmeaReceiver()
|
|
||||||
# receiver.filter_by_msg_type('.*')
|
|
||||||
receiver.register_msg('^.?GPGSV.*', GPGSV_callback)
|
|
||||||
receiver.register_msg('^.?GNGLL.*', GNGLL_callback)
|
|
||||||
|
|
||||||
# Packetizer for "zed-x20p"
|
# Packetizer for "zed-x20p"
|
||||||
packetizer_0 = Packetizer()
|
packetizer_0 = Packetizer()
|
||||||
packetizer_0.register_recv("zed-x20p", receiver)
|
|
||||||
|
|
||||||
# Packetizer for "neo-f9p"
|
# Packetizer for "neo-f9p"
|
||||||
packetizer_1 = Packetizer()
|
packetizer_1 = Packetizer()
|
||||||
packetizer_1.register_recv("neo-f9p", receiver)
|
|
||||||
|
|
||||||
# Register both packetizers at multi source receiver
|
xcvr_1.register_recv(packetizer_0)
|
||||||
mc.register_recv("zed-x20p", packetizer_0)
|
xcvr_2.register_recv(packetizer_1)
|
||||||
mc.register_recv("neo-f9p", packetizer_1)
|
rx_1 = MsgReceiver('Rx1', '^.?GPGSV.*')
|
||||||
|
rx_2 = MsgReceiver('Rx2', '^.?GNGLL.*')
|
||||||
|
|
||||||
|
# Register all Rx to all sources
|
||||||
|
packetizer_0.register_recv(rx_1)
|
||||||
|
packetizer_0.register_recv(rx_2)
|
||||||
|
packetizer_1.register_recv(rx_1)
|
||||||
|
packetizer_1.register_recv(rx_2)
|
||||||
|
|
||||||
# Connect all sources
|
# Connect all sources
|
||||||
mc.connect()
|
mc.connect()
|
||||||
@@ -1,23 +0,0 @@
|
|||||||
from msg_montainer import MsgContainer
|
|
||||||
from msg_sink import MsgSink
|
|
||||||
from typing import Callable
|
|
||||||
import re
|
|
||||||
|
|
||||||
class NmeaReceiver(MsgSink):
|
|
||||||
def __init__(self):
|
|
||||||
MsgSink.__init__(self)
|
|
||||||
self.filters = []
|
|
||||||
|
|
||||||
def on_recv(self, msg: MsgContainer):
|
|
||||||
data = msg.data.decode(encoding="utf-8")
|
|
||||||
for filter in self.filters:
|
|
||||||
if re.match(filter['re'], data):
|
|
||||||
if filter['listener'] is not None:
|
|
||||||
filter['listener'](msg)
|
|
||||||
|
|
||||||
def register_msg(self, msg_type: str, listener: Callable):
|
|
||||||
self.filters.append({'re': msg_type, 'listener': listener})
|
|
||||||
|
|
||||||
|
|
||||||
def _filter_by_msg_type(self, msg_type: str):
|
|
||||||
self.filters.append(msg_type)
|
|
||||||
+3
-1
@@ -1,4 +1,4 @@
|
|||||||
from msg_montainer import MsgContainer
|
from msg_container import MsgContainer
|
||||||
from msg_sink import MsgSink
|
from msg_sink import MsgSink
|
||||||
from msg_source import MsgSource
|
from msg_source import MsgSource
|
||||||
from struct import pack
|
from struct import pack
|
||||||
@@ -27,3 +27,5 @@ class Packetizer(MsgSink, MsgSource):
|
|||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
def is_msg(self, msg: MsgContainer):
|
||||||
|
return True
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
from msg_source import MsgSource
|
||||||
|
from a_backend import ABackend
|
||||||
|
from msg_container import *
|
||||||
|
|
||||||
|
class Transceiver(MsgSource):
|
||||||
|
def __init__(self, backend: ABackend, name: str):
|
||||||
|
MsgSource.__init__(self)
|
||||||
|
self.backend = backend
|
||||||
|
self.name = name
|
||||||
|
|
||||||
|
def on_recv(self, msg: MsgContainer):
|
||||||
|
self.call_sinks(msg)
|
||||||
Reference in New Issue
Block a user