c34449aac5
First Upload to Gitea
204 lines
5.9 KiB
Python
204 lines
5.9 KiB
Python
import socket
|
|
import struct
|
|
import time
|
|
import numpy as np
|
|
from threading import Thread, Event
|
|
from queue import Queue
|
|
|
|
import config as cfg
|
|
import DataTypes as dt
|
|
|
|
|
|
# ============================================================
|
|
# Protocol definitions (SHARED with ESP32)
|
|
# ============================================================
|
|
|
|
PACKET_VERSION = 1
|
|
|
|
# Packet IDs
|
|
PKT_COMMAND = 1
|
|
PKT_STATUS = 2
|
|
PKT_DEBUG = 3
|
|
|
|
# Field Types
|
|
FT_ARRAY_F32 = 1 # float32 array
|
|
FT_STRING = 2 # utf-8 string
|
|
FT_UINT8 = 3
|
|
FT_FLOAT32 = 4
|
|
|
|
# Header: packet_id | version | payload_len | timestamp_ms
|
|
HEADER_FMT = "<B B H I"
|
|
HEADER_SIZE = struct.calcsize(HEADER_FMT)
|
|
|
|
# TLV field: type | length | value
|
|
TLV_FMT = "<B H"
|
|
TLV_SIZE = struct.calcsize(TLV_FMT)
|
|
|
|
|
|
# ============================================================
|
|
# ESP32 Communication Thread
|
|
# ============================================================
|
|
|
|
class ESP32Communication(Thread):
|
|
def __init__(self, host=cfg.esp32_ip, port=cfg.esp32_port, timeout=cfg.timeout):
|
|
super().__init__(daemon=True)
|
|
|
|
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
self.sock.settimeout(timeout)
|
|
|
|
print(f"Connecting to ESP32 at {host}:{port} ...")
|
|
self.sock.connect((host, port))
|
|
print("Connected to ESP32!")
|
|
|
|
self.tx_queue = Queue()
|
|
self.rx_queue = Queue()
|
|
self.running = Event()
|
|
self.running.set()
|
|
|
|
# --------------------------------------------------------
|
|
# Thread loop
|
|
# --------------------------------------------------------
|
|
def run(self):
|
|
self.sock.setblocking(False)
|
|
|
|
while self.running.is_set():
|
|
self._handle_tx()
|
|
self._handle_rx()
|
|
time.sleep(0.002)
|
|
|
|
# --------------------------------------------------------
|
|
# Public API
|
|
# --------------------------------------------------------
|
|
def send_command(self, radial_array :dt.RadArray, state=0, name=""):
|
|
"""
|
|
Example command packet:
|
|
- 3x6 float array
|
|
- uint8 state
|
|
- optional string
|
|
"""
|
|
data = np.degrees(radial_array).astype(np.uint8)
|
|
print(data)
|
|
fields = [
|
|
self._tlv_array(data),
|
|
self._tlv_uint8(state)
|
|
]
|
|
|
|
if name:
|
|
fields.append(self._tlv_string(name))
|
|
|
|
packet = self._build_packet(PKT_COMMAND, fields)
|
|
self.tx_queue.put(packet)
|
|
|
|
def receive(self):
|
|
"""
|
|
Non-blocking receive.
|
|
Returns (packet_id, timestamp, fields) or None
|
|
"""
|
|
try:
|
|
return self.rx_queue.get_nowait()
|
|
except Exception:
|
|
return None
|
|
|
|
def stop(self):
|
|
self.running.clear()
|
|
self.join()
|
|
self.sock.close()
|
|
|
|
# --------------------------------------------------------
|
|
# TX / RX handlers
|
|
# --------------------------------------------------------
|
|
def _handle_tx(self):
|
|
if not self.tx_queue.empty():
|
|
data = self.tx_queue.get()
|
|
self.sock.sendall(data)
|
|
|
|
def _handle_rx(self):
|
|
try:
|
|
header = self._recv_exact(HEADER_SIZE)
|
|
if not header:
|
|
return
|
|
|
|
packet_id, version, payload_len, timestamp = struct.unpack(HEADER_FMT, header)
|
|
|
|
if payload_len == 0:
|
|
# Nothing to parse
|
|
self.rx_queue.put((packet_id, timestamp, []))
|
|
return
|
|
|
|
payload = self._recv_exact(payload_len)
|
|
if payload is None:
|
|
return # incomplete, skip
|
|
|
|
fields = self._parse_tlv(payload)
|
|
self.rx_queue.put((packet_id, timestamp, fields))
|
|
|
|
except BlockingIOError:
|
|
pass
|
|
|
|
# --------------------------------------------------------
|
|
# Packet construction
|
|
# --------------------------------------------------------
|
|
def _build_packet(self, packet_id, fields):
|
|
payload = b"".join(fields)
|
|
|
|
header = struct.pack(
|
|
HEADER_FMT,
|
|
packet_id,
|
|
PACKET_VERSION,
|
|
len(payload),
|
|
int(time.time() * 1000) & 0xFFFFFFFF
|
|
)
|
|
|
|
return header + payload
|
|
|
|
# --------------------------------------------------------
|
|
# TLV helpers
|
|
# --------------------------------------------------------
|
|
def _tlv(self, ftype, data):
|
|
return struct.pack(TLV_FMT, ftype, len(data)) + data
|
|
|
|
def _tlv_array(self, array):
|
|
# Convert to uint8 degrees before sending
|
|
arr_uint8 = array.astype(np.uint8)
|
|
return self._tlv(FT_ARRAY_F32, arr_uint8.tobytes())
|
|
|
|
def _tlv_string(self, text):
|
|
return self._tlv(FT_STRING, text.encode("utf-8"))
|
|
|
|
def _tlv_uint8(self, value):
|
|
return self._tlv(FT_UINT8, struct.pack("<B", value))
|
|
|
|
# --------------------------------------------------------
|
|
# TLV parsing
|
|
# --------------------------------------------------------
|
|
def _parse_tlv(self, payload):
|
|
pos = 0
|
|
fields = []
|
|
|
|
while pos < len(payload):
|
|
ftype, flen = struct.unpack(
|
|
TLV_FMT, payload[pos:pos + TLV_SIZE]
|
|
)
|
|
pos += TLV_SIZE
|
|
|
|
data = payload[pos:pos + flen]
|
|
pos += flen
|
|
|
|
fields.append((ftype, data))
|
|
|
|
return fields
|
|
|
|
# --------------------------------------------------------
|
|
# Socket helper
|
|
# --------------------------------------------------------
|
|
def _recv_exact(self, size):
|
|
data = b""
|
|
while len(data) < size:
|
|
try:
|
|
chunk = self.sock.recv(size - len(data))
|
|
if not chunk:
|
|
return None
|
|
data += chunk
|
|
except BlockingIOError:
|
|
return None
|
|
return data |