Added timestamp in milliseconds to messages

This commit is contained in:
Your Name
2026-09-29 15:00:06 +03:00
parent 92271ef376
commit 12f05d8f10
4 changed files with 260 additions and 28 deletions
+243 -22
View File
@@ -1,5 +1,18 @@
#!/usr/bin/env python3
"""Send binary commands to a device over a serial port, optionally receiving and/or requiring an ACK."""
"""Send binary commands to a device over a serial port, optionally receiving and/or requiring an ACK.
Packet layout matches:
struct command_message_t {
uint8_t prefix;
uint8_t length;
uint8_t id;
uint8_t command;
int64_t tick; // milliseconds, little-endian
uint8_t crc;
uint8_t data[COMMAND_DATA_SIZE];
} __attribute__((packed));
"""
import argparse
import struct
@@ -59,6 +72,18 @@ EXIT_ERROR = 1
EXIT_NO_ACK = 2
EXIT_NACK = 3
# Fixed header layout: prefix, length, id, command, tick(int64), crc
TICK_FORMAT = "<q" # signed 64-bit little-endian, milliseconds
TICK_SIZE = struct.calcsize(TICK_FORMAT) # 8
CRC_OFFSET = 4 + TICK_SIZE # index of the crc byte -> 12
HEADER_SIZE = CRC_OFFSET + 1 # bytes before data -> 13
# Safety valve for the ASCII-log branch: a real firmware log line shouldn't
# be longer than this. If we've buffered more than this with no newline,
# we're desynced (likely sitting inside binary packet data), so drop a byte
# and retry instead of blocking forever waiting for a '\n' that won't come.
MAX_LOG_LINE = 256
# --------------------------------------------------------------------------
# Packet helpers
@@ -69,27 +94,35 @@ def calculate_crc(msg: bytes) -> int:
return (-s) & 0xFF
def make_packet(command: int, data: bytes = b"", device_id: int = DEFAULT_DEVICE_ID) -> bytes:
def make_packet(command: int, data: bytes = b"", device_id: int = DEFAULT_DEVICE_ID,
tick: int = 0) -> bytes:
pkt = bytearray()
pkt.append(COMMAND_PREFIX)
pkt.append(len(data))
pkt.append(device_id)
pkt.append(command)
pkt.extend(struct.pack(TICK_FORMAT, tick))
pkt.append(0) # CRC placeholder
pkt.extend(data)
pkt[4] = calculate_crc(pkt[:4] + pkt[5:])
pkt[CRC_OFFSET] = calculate_crc(pkt[:CRC_OFFSET] + pkt[CRC_OFFSET + 1:])
return bytes(pkt)
def verify_crc(packet: bytes) -> bool:
return packet[4] == calculate_crc(packet[:4] + packet[5:])
crc = packet[CRC_OFFSET]
calc = calculate_crc(packet[:CRC_OFFSET] + packet[CRC_OFFSET + 1:])
return crc == calc
def packet_size(buf: bytes):
if len(buf) < 2:
return None
return 5 + buf[1]
return HEADER_SIZE + buf[1]
def unpack_tick(packet: bytes) -> int:
return struct.unpack(TICK_FORMAT, packet[4:4 + TICK_SIZE])[0]
# --------------------------------------------------------------------------
@@ -97,21 +130,72 @@ def packet_size(buf: bytes):
# --------------------------------------------------------------------------
class Receiver(threading.Thread):
"""Parses incoming packets / ASCII logs. Prints them only if verbose."""
"""Parses incoming packets / ASCII logs.
def __init__(self, ser, verbose):
`verbose` gates whether anything is printed at all. `mode` picks the
print format for received command packets:
normal - name, tick, length, then the raw hex line (default)
debug - the full raw hex bitstream of every packet/line, nothing else
data - only the data payload bytes
timestamp - tick + data payload bytes ("-l")
Stats for --telemetry are always collected, regardless of `verbose`.
"""
def __init__(self, ser, verbose=False, mode="normal"):
super().__init__(daemon=True)
self.ser = ser
self.verbose = verbose
self.mode = mode
self.stop_event = threading.Event()
self.ack_event = threading.Event()
self.result = None # "ACK" or "NACK" once received
# telemetry
self.start_time = None
self.end_time = None
self.command_counts = {}
self.ack_count = 0
self.nack_count = 0
self.bad_crc_count = 0
self.log_line_count = 0
self.data_byte_count = 0
self.packet_count = 0
# device-clock (tick, ms) timestamps of every CRC-valid packet, in
# arrival order, used to compute real RX frequency. Kept per name
# and combined across all types.
self.tick_history = {}
self.all_ticks = []
def _print(self, *args):
if self.verbose:
print(*args)
def _print_packet(self, pkt, length, dev_id, cmd, tick, payload):
name = COMMAND_NAMES.get(cmd, str(cmd))
if self.mode == "debug":
self._print(f"<-- RAW: {pkt.hex(' ')}")
elif self.mode == "data":
self._print(f"<-- DATA: {payload.hex(' ')}")
elif self.mode == "timestamp":
self._print(f"<-- t={tick}ms DATA: {payload.hex(' ')}")
else: # normal
self._print(f"<-- {name} (device={dev_id}, tick={tick}ms) len={length}")
self._print(f"<-- RX: {pkt.hex(' ')}")
def _print_ack_nack(self, pkt, dev_id, tick, name):
if self.mode == "debug":
self._print(f"<-- RAW: {pkt.hex(' ')}")
elif self.mode == "data":
pass # no data payload to show
elif self.mode == "timestamp":
self._print(f"<-- t={tick}ms {name}")
else: # normal
self._print(f"<-- {name} (device={dev_id}, tick={tick}ms)")
def run(self):
self.start_time = time.time()
rx = bytearray()
while not self.stop_event.is_set():
@@ -128,41 +212,125 @@ class Receiver(threading.Thread):
break
pkt = bytes(rx[:size])
del rx[:size]
if not verify_crc(pkt):
self._print("RX: Bad CRC:", pkt.hex(" "))
self.bad_crc_count += 1
# The CRC is what's wrong, not necessarily the tick
# field's position, so still surface it - useful for
# spotting exactly this kind of firmware bug.
bad_tick = unpack_tick(pkt)
if self.mode == "debug":
self._print(f"RX: Bad CRC: {pkt.hex(' ')}")
else:
self._print(f"RX: Bad CRC (tick={bad_tick}ms): {pkt.hex(' ')}")
# The guessed size was based on a length byte that
# might itself be garbage (stale/misaligned bytes).
# Drop just the leading prefix byte and re-scan,
# rather than eating `size` bytes and risking a
# cascade of misalignment.
del rx[:1]
continue
del rx[:size]
length = pkt[1]
dev_id = pkt[2]
cmd = pkt[3]
tick = unpack_tick(pkt)
payload = pkt[HEADER_SIZE:HEADER_SIZE + length]
self.packet_count += 1
if cmd == COMMAND_ACK:
self._print(f"<-- ACK (device={dev_id})")
self.ack_count += 1
self.result = "ACK"
self._print_ack_nack(pkt, dev_id, tick, "ACK")
self.tick_history.setdefault("ACK", []).append(tick)
self.all_ticks.append(tick)
self.ack_event.set()
elif cmd == COMMAND_NACK:
self._print(f"<-- NACK (device={dev_id})")
self.nack_count += 1
self.result = "NACK"
self._print_ack_nack(pkt, dev_id, tick, "NACK")
self.tick_history.setdefault("NACK", []).append(tick)
self.all_ticks.append(tick)
self.ack_event.set()
else:
name = COMMAND_NAMES.get(cmd, str(cmd))
self._print(f"<-- Command {name} len={length}")
self._print(f"<-- RX: {pkt.hex(' ')}")
self.command_counts[name] = self.command_counts.get(name, 0) + 1
self.data_byte_count += length
self._print_packet(pkt, length, dev_id, cmd, tick, payload)
self.tick_history.setdefault(name, []).append(tick)
self.all_ticks.append(tick)
else:
idx = rx.find(b"\n")
if idx == -1:
if len(rx) > MAX_LOG_LINE:
# Not really a log line: desynced, likely sitting
# inside binary packet data. Shed a byte and
# keep looking for the next prefix instead of
# stalling forever waiting for '\n'.
del rx[:1]
continue
break
line = rx[:idx + 1]
del rx[:idx + 1]
self.log_line_count += 1
if self.mode == "data":
continue # data-only stream: skip firmware log text
try:
self._print("[LOG]", line.decode().rstrip())
decoded = line.decode().rstrip()
if self.mode == "debug":
self._print(f"[RAW] {line.hex(' ')}")
else:
self._print("[LOG]", decoded)
except UnicodeDecodeError:
self._print("[RAW]", line.hex())
self.end_time = time.time()
def print_telemetry(self):
start = self.start_time or time.time()
end = self.end_time or time.time()
elapsed = max(end - start, 1e-9)
total = self.packet_count + self.log_line_count
def hz(ticks):
"""Frequency in Hz from device-clock (tick, ms) deltas between
the first and last sample. None if there aren't enough points
or the ticks didn't advance (e.g. all identical / rolled over)."""
if len(ticks) < 2:
return None
span_ms = ticks[-1] - ticks[0]
if span_ms <= 0:
return None
return (len(ticks) - 1) * 1000.0 / span_ms
print("--- Telemetry ---")
print(f"elapsed: {elapsed:.3f} s")
print(f"packets: {self.packet_count} ({self.packet_count / elapsed:.2f} pkt/s host-side)")
print(f"data bytes: {self.data_byte_count} ({self.data_byte_count / elapsed:.2f} B/s)")
print(f"log lines: {self.log_line_count}")
print(f"ACK / NACK: {self.ack_count} / {self.nack_count}")
print(f"bad CRC: {self.bad_crc_count}")
print(f"total frames: {total} ({total / elapsed:.2f} frames/s)")
overall_hz = hz(self.all_ticks)
if overall_hz is not None:
print(f"RX frequency: {overall_hz:.2f} Hz (from device tick deltas, all valid packets)")
else:
print("RX frequency: n/a (need >=2 CRC-valid packets with advancing tick)")
if self.command_counts:
print("by command (count, Hz from tick deltas):")
for name, count in sorted(self.command_counts.items()):
f = hz(self.tick_history.get(name, []))
f_str = f"{f:.2f} Hz" if f is not None else "n/a"
print(f" {name:<18} {count:>6} {f_str}")
# --------------------------------------------------------------------------
# Argument parsing
@@ -211,6 +379,15 @@ def parse_data_token(tok: str) -> bytes:
return struct.pack(fmt, num)
def parse_tick(text: str) -> int:
if text.lower() == "now":
return int(time.time() * 1000)
try:
return int(text, 0)
except ValueError:
raise argparse.ArgumentTypeError("tick must be an integer (ms) or 'now'")
def build_parser():
commands_help = ", ".join(f"{name}={cid}" for name, cid in COMMANDS.items())
optional_help = ", ".join(sorted(COMMANDS_DATA_OPTIONAL)) or "(none)"
@@ -230,11 +407,18 @@ data tokens (all multi-byte values are little-endian):
u16:500 also: u8 i8 u16 i16 u32 i32 f32
hex:0a0b raw bytes
RX display modes (pick at most one; require --rx, --debug, --data or -l to receive at all):
(none) name, tick, length, and raw hex per packet
--debug raw hex bitstream of every packet/line, nothing decoded
--data only the data payload bytes of received commands
-l tick timestamp + data payload bytes
examples:
%(prog)s LED_TOGGLE
%(prog)s ADC_SET_READ u32:1 --require-ack
%(prog)s SERVO_SET 0 u16:1500 --rx --rx-time 5
%(prog)s -p /dev/ttyUSB0 -b 9600 ADC_READ_ALL --rx
%(prog)s ADC_READ_ALL -l --telemetry
%(prog)s -p /dev/ttyUSB0 -b 9600 ADC_READ_ALL --debug
""",
)
@@ -248,17 +432,34 @@ examples:
help="baud rate (default: %(default)s)")
p.add_argument("-d", "--device-id", type=lambda s: int(s, 0), default=DEFAULT_DEVICE_ID,
help="device id in the packet header (default: %(default)s)")
p.add_argument("-t", "--tick", type=parse_tick, default=0,
help="tick value in ms for the outgoing packet header, or 'now' (default: %(default)s)")
p.add_argument("--rx", action="store_true",
help="receive and print incoming packets/logs (default: off)")
help="receive and print incoming packets/logs in the default format (default: off)")
p.add_argument("--rx-time", type=float, default=10.0,
help="seconds to keep listening after sending, with --rx (default: %(default)s)")
help="seconds to keep listening after sending, when receiving (default: %(default)s)")
rx_format = p.add_mutually_exclusive_group()
rx_format.add_argument("--debug", action="store_true",
help="RX mode: print the raw hex bitstream of every packet/line (default: off)")
rx_format.add_argument("--data", dest="data_only", action="store_true",
help="RX mode: print only the data payload of received commands (default: off)")
rx_format.add_argument("-l", "--timestamp", action="store_true",
help="RX mode: print data payload with its tick timestamp (default: off)")
p.add_argument("--require-ack", action="store_true",
help="wait for an ACK; exit 2 on timeout, 3 on NACK (default: off)")
p.add_argument("--ack-timeout", type=float, default=1.0,
help="seconds to wait for ACK (default: %(default)s)")
p.add_argument("--boot-delay", type=float, default=0.0,
help="seconds to wait after opening the port before flushing/sending, "
"for boards that reset on connect (default: %(default)s, off)")
p.add_argument("--telemetry", action="store_true",
help="print an RX telemetry summary (packet/byte rates, counts) at exit (default: off)")
p.add_argument("--list-commands", action="store_true",
help="list known commands and exit")
return p
@@ -296,10 +497,27 @@ def main():
if len(data) > 255:
parser.error(f"data too long ({len(data)} bytes, max 255)")
packet = make_packet(cmd_id, bytes(data), args.device_id)
packet = make_packet(cmd_id, bytes(data), args.device_id, args.tick)
# --debug/--data/-l/--telemetry imply listening for the rx-time window,
# same as --rx, even if --rx itself wasn't passed.
listen = args.rx or args.debug or args.data_only or args.timestamp or args.telemetry
need_receiver = listen or args.require_ack
if args.debug:
mode = "debug"
elif args.data_only:
mode = "data"
elif args.timestamp:
mode = "timestamp"
else:
mode = "normal"
try:
ser = serial.Serial(args.port, args.baud, timeout=0.05)
if args.boot_delay > 0:
time.sleep(args.boot_delay) # let a board that resets on connect finish booting
ser.reset_input_buffer() # discard stale/boot-garbage bytes before we start
except serial.SerialException as e:
print(f"ERROR: could not open {args.port}: {e}", file=sys.stderr)
return EXIT_ERROR
@@ -309,11 +527,11 @@ def main():
try:
# Start the receiver BEFORE sending so a fast ACK isn't missed.
if args.rx or args.require_ack:
receiver = Receiver(ser, verbose=args.rx)
if need_receiver:
receiver = Receiver(ser, verbose=listen, mode=mode)
receiver.start()
print(f"--> TX {cmd_name}: {packet.hex(' ')}")
print(f"--> TX {cmd_name} (tick={args.tick}ms): {packet.hex(' ')}")
ser.write(packet)
ser.flush()
@@ -327,7 +545,7 @@ def main():
else:
print("ACK received.")
if args.rx:
if listen:
time.sleep(args.rx_time)
except KeyboardInterrupt:
@@ -338,6 +556,9 @@ def main():
receiver.join(timeout=1)
ser.close()
if args.telemetry and receiver:
receiver.print_telemetry()
print("Done.")
return exit_code
+4 -1
View File
@@ -2,6 +2,7 @@
#include <string.h>
#include <zephyr/logging/log.h>
#include <zephyr/kernel.h>
LOG_MODULE_REGISTER(command_message, LOG_LEVEL_INF);
@@ -9,6 +10,7 @@ void command_message_init(struct command_message_t *msg) {
memset(msg, 0, sizeof(struct command_message_t));
msg->prefix = COMMAND_PREFIX;
msg->id = COMMAND_ID;
msg->tick = 0;
}
void command_create_message(struct command_message_t *msg, uint8_t length, commands_e command, uint8_t data[]) {
@@ -20,6 +22,7 @@ void command_create_message(struct command_message_t *msg, uint8_t length, comma
command_message_init(msg);
msg->length = length;
msg->command = command;
msg->tick = k_uptime_get();
// Copy the data
if (data != NULL) {
@@ -38,7 +41,7 @@ uint8_t command_calculate_crc(struct command_message_t *msg) {
int loop_length = (sizeof(struct command_message_t) - sizeof(msg->data) + msg->length);
for (int i = 0; i < loop_length; i++) {
if (i == 4) { continue; }
if (i == (COMMAND_HEADER_SIZE - 1)) { continue; }
sum += byte_ptr[i];
}
+2 -1
View File
@@ -8,7 +8,7 @@
#define COMMAND_PREFIX 0x69
#define COMMAND_ID 0x00
#define COMMAND_DATA_SIZE 160
#define COMMAND_HEADER_SIZE 5
#define COMMAND_HEADER_SIZE 13
typedef enum {
COMMAND_ACK,
@@ -32,6 +32,7 @@ struct command_message_t {
uint8_t length;
uint8_t id;
uint8_t command;
int64_t tick;
uint8_t crc;
uint8_t data[COMMAND_DATA_SIZE];
} __attribute__((packed));
+11 -4
View File
@@ -104,13 +104,14 @@ static void usb_rx_thread(void *p1, void *p2, void *p3) {
len = ring_buf_get(&ringbuf, &buf_prefix, 1);
if (len && (buf_prefix == COMMAND_PREFIX)) {
uint8_t buf_header[4];
len = ring_buf_get(&ringbuf, buf_header, 4);
uint8_t buf_header[COMMAND_HEADER_SIZE];
len = ring_buf_get(&ringbuf, buf_header, (COMMAND_HEADER_SIZE - 1));
if ((len == 4) && (buf_header[1] == COMMAND_ID) && (buf_header[0] <= COMMAND_DATA_SIZE)) {
if ((len == (COMMAND_HEADER_SIZE - 1)) && (buf_header[1] == COMMAND_ID) && (buf_header[0] <= COMMAND_DATA_SIZE)) {
msg.length = buf_header[0];
msg.command = buf_header[2];
msg.crc = buf_header[3];
// Ignore the tick for now
msg.crc = buf_header[COMMAND_HEADER_SIZE - 2];
if (msg.length) {
len = ring_buf_get(&ringbuf, msg.data, msg.length);
@@ -120,6 +121,8 @@ static void usb_rx_thread(void *p1, void *p2, void *p3) {
if (calculated_crc != msg.crc) {
if (RETURN_ACK) {
// Send NACK
usb_tx_buffer[TX_BUFFER_SIZE + 1].tick = k_uptime_get();
usb_tx_buffer[TX_BUFFER_SIZE + 1].crc = command_calculate_crc(&usb_tx_buffer[TX_BUFFER_SIZE + 1]);
usb_send_command(&usb_tx_buffer[TX_BUFFER_SIZE + 1]);
}
continue;
@@ -129,12 +132,16 @@ static void usb_rx_thread(void *p1, void *p2, void *p3) {
if (ret == 0) {
if (RETURN_ACK) {
// Send ACK
usb_tx_buffer[TX_BUFFER_SIZE].tick = k_uptime_get();
usb_tx_buffer[TX_BUFFER_SIZE].crc = command_calculate_crc(&usb_tx_buffer[TX_BUFFER_SIZE]);
usb_send_command(&usb_tx_buffer[TX_BUFFER_SIZE]);
}
}
else {
if (RETURN_ACK) {
// Send NACK
usb_tx_buffer[TX_BUFFER_SIZE + 1].tick = k_uptime_get();
usb_tx_buffer[TX_BUFFER_SIZE + 1].crc = command_calculate_crc(&usb_tx_buffer[TX_BUFFER_SIZE + 1]);
usb_send_command(&usb_tx_buffer[TX_BUFFER_SIZE + 1]);
}
}