- added Msg frame work with Msg container , Msg sources, Msg sinks
- added filter by message type to nmea receiver - show case new features in MultiClient.main()
This commit is contained in:
@@ -1 +1,2 @@
|
|||||||
.idea/
|
.idea/
|
||||||
|
__pycache__/
|
||||||
|
|||||||
@@ -0,0 +1,13 @@
|
|||||||
|
|
||||||
|
|
||||||
|
class MsgContainer:
|
||||||
|
def __init__(self, name: str, timestamp: float = 0, data: bytes = b''):
|
||||||
|
self.name : str= name
|
||||||
|
self.timestamp : float = timestamp
|
||||||
|
self.data: bytes = data
|
||||||
|
|
||||||
|
def __str__(self):
|
||||||
|
return f"{self.name}: {self.timestamp}: {self.data}"
|
||||||
|
|
||||||
|
def __repr__(self):
|
||||||
|
return f"{self.name}: {self.timestamp}: {self.data}"
|
||||||
+10
@@ -0,0 +1,10 @@
|
|||||||
|
from abc import ABC, abstractmethod
|
||||||
|
from msg_montainer import MsgContainer
|
||||||
|
|
||||||
|
class MsgSink(ABC):
|
||||||
|
def __init__(self):
|
||||||
|
pass
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
def on_recv(self, msg: MsgContainer):
|
||||||
|
pass
|
||||||
@@ -0,0 +1,15 @@
|
|||||||
|
from abc import ABC
|
||||||
|
from msg_montainer import MsgContainer
|
||||||
|
from msg_sink import MsgSink
|
||||||
|
|
||||||
|
class MsgSource(ABC):
|
||||||
|
def __init__(self):
|
||||||
|
self.sinks = {}
|
||||||
|
|
||||||
|
def register_recv(self, name: str, sink: MsgSink):
|
||||||
|
self.sinks[name] = {'sink': sink}
|
||||||
|
|
||||||
|
def call_sinks(self, msg: MsgContainer):
|
||||||
|
if msg.name in self.sinks:
|
||||||
|
sink = self.sinks[msg.name]['sink']
|
||||||
|
sink.on_recv(msg)
|
||||||
+46
-18
@@ -2,15 +2,27 @@ import socket
|
|||||||
import selectors
|
import selectors
|
||||||
import time
|
import time
|
||||||
|
|
||||||
class MultiClient(object):
|
from msg_montainer import MsgContainer
|
||||||
|
from nmea import NmeaReceiver
|
||||||
|
from msg_sink import MsgSink
|
||||||
|
from msg_source import MsgSource
|
||||||
|
from packetizer import Packetizer
|
||||||
|
|
||||||
|
|
||||||
|
class MultiClient(MsgSource):
|
||||||
def __init__(self, servers: list[dict]):
|
def __init__(self, servers: list[dict]):
|
||||||
self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port'])} for s in servers]
|
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.receivers = {}
|
||||||
for server in self.sock_list:
|
for server in self.sock_list:
|
||||||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||||
sock.setblocking(False)
|
sock.setblocking(False)
|
||||||
server['sock'] = sock
|
server['sock'] = sock
|
||||||
|
|
||||||
|
def register_recv(self, name: str, rx: MsgSink):
|
||||||
|
self.receivers[name] = {'rx': rx}
|
||||||
|
|
||||||
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:
|
||||||
@@ -29,16 +41,6 @@ class MultiClient(object):
|
|||||||
for sock in self.sock_list:
|
for sock in self.sock_list:
|
||||||
sock['sock'].close()
|
sock['sock'].close()
|
||||||
|
|
||||||
def on_receive(self, rx_list: list):
|
|
||||||
for data in rx_list:
|
|
||||||
print(f"{self}: on_receive: 'Time': {data['timestamp'], data['name']}")
|
|
||||||
|
|
||||||
def send_request(self):
|
|
||||||
events = selectors.EVENT_READ | selectors.EVENT_WRITE
|
|
||||||
for sock in self.sock_list:
|
|
||||||
self.sel.register(sock['sock'], events, data=None)
|
|
||||||
print("receive")
|
|
||||||
|
|
||||||
def event_loop(self):
|
def event_loop(self):
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
@@ -46,12 +48,14 @@ class MultiClient(object):
|
|||||||
rx_data = []
|
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())
|
||||||
if mask & selectors.EVENT_READ:
|
name = descr['name']
|
||||||
timestamp = time.time()
|
|
||||||
data = key.fileobj.recv(1024)
|
|
||||||
rx_data.append({'timestamp': timestamp, 'name': descr['name'], 'data': data})
|
|
||||||
|
|
||||||
self.on_receive(rx_data)
|
if name in self.receivers:
|
||||||
|
rx = self.receivers[name]['rx']
|
||||||
|
if mask & selectors.EVENT_READ:
|
||||||
|
timestamp = time.time()
|
||||||
|
data = key.fileobj.recv(16)
|
||||||
|
rx.on_recv(MsgContainer(name, timestamp, data))
|
||||||
|
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
print("Caught keyboard interrupt, exiting")
|
print("Caught keyboard interrupt, exiting")
|
||||||
@@ -65,8 +69,32 @@ if __name__ == "__main__":
|
|||||||
{'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731}
|
{'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731}
|
||||||
]
|
]
|
||||||
|
|
||||||
|
# TCPIP multi source receiver
|
||||||
mc = MultiClient(server_list)
|
mc = MultiClient(server_list)
|
||||||
|
|
||||||
|
# NMEA multi source receiver
|
||||||
|
receiver = NmeaReceiver()
|
||||||
|
# receiver.filter_by_msg_type('.*')
|
||||||
|
receiver.filter_by_msg_type('^.?GPGSV.*')
|
||||||
|
receiver.filter_by_msg_type('^.?GNGLL.*')
|
||||||
|
|
||||||
|
# Packetizer for "zed-x20p"
|
||||||
|
packetizer_0 = Packetizer()
|
||||||
|
packetizer_0.register_recv("zed-x20p", receiver)
|
||||||
|
|
||||||
|
# Packetizer for "neo-f9p"
|
||||||
|
packetizer_1 = Packetizer()
|
||||||
|
packetizer_1.register_recv("neo-f9p", receiver)
|
||||||
|
|
||||||
|
# Register both packetizers at multi source receiver
|
||||||
|
mc.register_recv("zed-x20p", packetizer_0)
|
||||||
|
mc.register_recv("neo-f9p", packetizer_1)
|
||||||
|
|
||||||
|
# Connect all sources
|
||||||
mc.connect()
|
mc.connect()
|
||||||
|
|
||||||
|
# receive and block
|
||||||
mc.event_loop()
|
mc.event_loop()
|
||||||
time.sleep(1)
|
|
||||||
|
# code never reached
|
||||||
mc.disconnect()
|
mc.disconnect()
|
||||||
|
|||||||
@@ -0,0 +1,20 @@
|
|||||||
|
from urllib.parse import to_bytes
|
||||||
|
|
||||||
|
from msg_montainer import MsgContainer
|
||||||
|
from msg_sink import MsgSink
|
||||||
|
import re
|
||||||
|
|
||||||
|
class NmeaReceiver(MsgSink):
|
||||||
|
def __init__(self):
|
||||||
|
MsgSink.__init__(self)
|
||||||
|
self.filters_msg_type = []
|
||||||
|
|
||||||
|
def filter_by_msg_type(self, msg_type: str):
|
||||||
|
self.filters_msg_type.append(msg_type)
|
||||||
|
|
||||||
|
def on_recv(self, msg: MsgContainer):
|
||||||
|
data = msg.data.decode(encoding="utf-8")
|
||||||
|
for pattern in self.filters_msg_type:
|
||||||
|
if re.match(pattern, data):
|
||||||
|
print(f"{self}: on_receive: {msg}")
|
||||||
|
|
||||||
@@ -0,0 +1,29 @@
|
|||||||
|
from msg_montainer import MsgContainer
|
||||||
|
from msg_sink import MsgSink
|
||||||
|
from msg_source import MsgSource
|
||||||
|
from struct import pack
|
||||||
|
|
||||||
|
class Packetizer(MsgSink, MsgSource):
|
||||||
|
def __init__(self):
|
||||||
|
MsgSink.__init__(self)
|
||||||
|
MsgSource.__init__(self)
|
||||||
|
self.packet = b''
|
||||||
|
self.wait_sync = True
|
||||||
|
|
||||||
|
def on_recv(self, msg: MsgContainer):
|
||||||
|
name = msg.name
|
||||||
|
data = msg.data
|
||||||
|
for n in range (len(data)):
|
||||||
|
d = data[n]
|
||||||
|
if d == 13 or d == 10:
|
||||||
|
timestamp = msg.timestamp
|
||||||
|
self.wait_sync = False
|
||||||
|
if len(self.packet) > 0:
|
||||||
|
self.call_sinks(MsgContainer(name, timestamp, self.packet))
|
||||||
|
self.wait_sync = True
|
||||||
|
self.packet = b''
|
||||||
|
elif not self.wait_sync:
|
||||||
|
self.packet += pack('B', d)
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
Reference in New Issue
Block a user