diff --git a/scratch/mm_stream_probe.py b/scratch/mm_stream_probe.py new file mode 100644 index 0000000..a233d82 --- /dev/null +++ b/scratch/mm_stream_probe.py @@ -0,0 +1,262 @@ +#!/usr/bin/env python3 +""" +mm_stream_probe.py — is there a one-request streaming mode for `SUB 0x5A`? + +The question +------------ +THOR downloads an event as a **chunk loop**: `ceil(size / 1024)` requests, each +asking for `min(1024, remaining)` bytes. Verified byte-for-byte against its own +frames, and confirmed on hardware up to 71 chunks. + +But our own 2026-09-23 probes recorded something different — a **single** request +with `offset_hi = 0x10` that appeared to return an entire 11 KB event: + + offset_word = 0x1000 + 2 * ceil(size / 512) + +`0x10` in `offset_hi` is exactly the bulk-stream marker Series III's +`build_5a_frame()` writes raw, so "`0x10XX` means stream until done, and the +device sends several frames" is a plausible reading. Those captures never landed +in the repo, so it cannot be re-derived from bytes on disk. + +Why it matters +-------------- +A round trip over cellular costs ~0.65 s regardless of payload. UM20147's +72,560-byte event is **71 chunks ≈ 46 seconds**. If one request can fetch it, +that becomes under a second. On a fleet of units called daily, that is the +difference between a workable receiver and an unworkable one. + +What this does +-------------- +Downloads **the same event twice** — once with the known-good chunk loop, once +with a single `0x10XX` request — and **diffs the bytes**. A differential test +rather than a suggestive one: if the streaming form returns byte-identical +output, it is safe to adopt; if it returns anything else, we learn exactly what. + +Then it re-POLLs, because the honest risk here is leaving the session in an odd +state, and the script should say so rather than leave you guessing. + +⚠ **Read-only.** `0x5A` is a read we have already sent thousands of times; the +only thing new is the value in its offset field. Nothing here writes, erases or +changes monitoring state. The worst realistic outcome is an unanswered frame or +a session that needs reconnecting. + +Usage +----- + python3 scratch/mm_stream_probe.py /dev/ttyACM1 + python3 scratch/mm_stream_probe.py /dev/ttyACM1 --event 055d4a82 + python3 scratch/mm_stream_probe.py 63.45.161.30:9034 --event smallest +""" + +from __future__ import annotations + +import argparse +import math +import sys +import time +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) +sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "bridges")) + +from micromate.client import MicromateClient, _content # noqa: E402 +from micromate.framing import MicromateFrameParser, build_request # noqa: E402 +from micromate.protocol import SUB_BULK_DOWNLOAD # noqa: E402 +from mm_client_check import StdlibSerial # noqa: E402 +from minimateplus.transport import TcpTransport # noqa: E402 + +_CHUNK_PREFIX = 11 + + +def collect(transport, parser, *, idle_gap: float, deadline: float) -> list: + """Read until `idle_gap` seconds pass with no new bytes, or `deadline`.""" + frames, last = [], time.monotonic() + while time.monotonic() < deadline: + chunk = transport.read(4096) + if chunk: + frames += parser.feed(chunk) + last = time.monotonic() + continue + if time.monotonic() - last > idle_gap: + break + time.sleep(0.005) + return frames + + +def main() -> int: + ap = argparse.ArgumentParser( + description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter + ) + ap.add_argument("target", help="host:port, or a serial device path") + ap.add_argument("--baud", type=int, default=115200) + ap.add_argument("--event", default="smallest", + help="a key in hex, 'smallest' (default — the gentlest " + "first test) or 'largest' (the one that matters)") + ap.add_argument("--idle-gap", type=float, default=2.0, + help="seconds of silence that end a streaming read; 2.0 is " + "generous for a modem, which buffers ~1 s") + ap.add_argument("--timeout", type=float, default=90.0) + a = ap.parse_args() + + if ":" in a.target and not Path(a.target).exists(): + host, _, port = a.target.rpartition(":") + inner = TcpTransport(host, int(port), connect_timeout=10.0) + label = f"TCP {host}:{port}" + else: + inner = StdlibSerial(a.target, baud=a.baud) + label = f"serial {a.target} @ {a.baud}" + + mm = MicromateClient(inner, recv_timeout=20.0) + print(f"\n{label} READ-ONLY: chain walk + two downloads of one event\n") + mm.open() + try: + info = mm.connect(with_active_setup=False) + print(f" {info}\n") + + refs = mm.list_events() + if not refs: + print(" no events stored — nothing to download. Record one first.") + return 1 + for r in refs: + print(f" {r}") + + which = a.event.strip().lower() + if which == "smallest": + ref = min(refs, key=lambda r: r.size) + elif which == "largest": + ref = max(refs, key=lambda r: r.size) + else: + matches = [r for r in refs if r.key_hex.lower() == which] + if not matches: + print(f"\n --event {a.event!r} matched nothing") + return 2 + ref = matches[0] + + n_chunks = math.ceil(ref.size / 1024) + print(f"\n target: {ref.key_hex} {ref.size} B {ref.record_type} " + f"({n_chunks} chunks the known-good way)") + + # ── 1. the known-good chunk loop ────────────────────────────────────── + t0 = time.monotonic() + chunked = mm.protocol.read_event_file(ref.key, ref.size) + dt_chunked = time.monotonic() - t0 + print(f"\n [1] chunk loop ...... {len(chunked)} B in {dt_chunked:.2f} s " + f"({n_chunks} requests)") + + # ── 2. one request, offset_hi = 0x10 ────────────────────────────────── + # The form our 2026-09-23 probes recorded. `pages` is 512-byte pages; + # the +0x1000 is the bulk-stream marker. + pages = math.ceil(ref.size / 512) + offset = 0x1000 + 2 * pages + params = ref.key + bytes(6) + + # ⚠ TRY BOTH ESCAPINGS, or a null result means nothing. + # + # On Series III, `offset_hi = 0x10` in a 5A frame must be written RAW -- + # doubled to `10 10`, the device SILENTLY IGNORES the frame. That is a + # documented, hard-won rule for this exact command. + # + # The Micromate escapes `0x04` in offset_hi (218/218 captured THOR + # frames), which argues the uniform escape set applies to 0x10 too. But + # Series III is a direct counter-example in the same command, so testing + # only one form risks concluding "no streaming mode" when the real + # finding is "that frame was malformed". + escaped = build_request(SUB_BULK_DOWNLOAD, offset, params) + candidates = [("escaped offset_hi (uniform rule)", escaped)] + if (offset >> 8) == 0x10: + # Hand-build the raw form: the uniform builder cannot express it. + payload = bytes([0x10, 0x00, SUB_BULK_DOWNLOAD, 0x00]) + \ + bytes([(offset >> 8) & 0xFF, offset & 0xFF]) + params + body = payload + bytes([sum(payload) & 0xFF]) + out = bytearray([0x41, 0x02]) + for i, b in enumerate(body): + # escape everything EXCEPT the offset_hi at payload index 4 + if b in (0x02, 0x03, 0x04, 0x10) and i != 4: + out.append(0x10) + out.append(b) + out.append(0x03) + candidates.append(("RAW offset_hi (series-III rule)", bytes(out))) + else: + print(f"\n note: offset_hi is 0x{offset >> 8:02x}, not 0x10, so the " + f"escaping question does not arise for this event") + + print(f"\n [2] streaming ....... one request, offset=0x{offset:04x} " + f"(0x1000 + 2 x {pages} pages)") + + frames, dt_stream, used = [], 0.0, None + for name, frame in candidates: + print(f"\n trying {name}") + print(f" wire: {frame.hex(' ')}") + parser = MicromateFrameParser() + t0 = time.monotonic() + mm.protocol._send(frame) + got = collect(inner, parser, idle_gap=a.idle_gap, + deadline=t0 + a.timeout) + dt = time.monotonic() - t0 + print(f" -> {len(got)} frame(s), {parser.bytes_fed} raw bytes, " + f"{dt:.2f} s") + if got: + frames, dt_stream, used = got, dt, name + break + + if used: + print(f"\n answered by: {used}") + if not frames: + print("\n VERDICT: no answer to EITHER escaping. The 0x10XX form") + print(" is not a streaming mode — or not with these params. Since") + print(" both escapings were tried, this is not a framing artefact.") + print(" The chunk loop stands as the only way to download an event,") + print(" and the 2026-09-23 note should be retracted.") + else: + bad = [f for f in frames if not f.checksum_valid] + subs = sorted({f"0x{f.sub:02x}" for f in frames}) + print(f" SUBs {subs}, {len(bad)} bad checksum") + + # Assemble the same way a chunk response is assembled. + streamed = b"".join(f.data[_CHUNK_PREFIX:] for f in frames) + print(f" assembled {len(streamed)} B " + f"(event is {ref.size} B)") + + print() + if streamed == chunked: + print(f" VERDICT: ** BYTE-IDENTICAL ** in {len(frames)} frame(s) " + f"against {n_chunks}.") + print(f" {dt_chunked:.2f} s -> {dt_stream:.2f} s here; over " + f"cellular that is ~{n_chunks * 0.65:.0f} s -> ~0.7 s.") + print(" The streaming mode is real. Worth adopting.") + elif len(streamed) == ref.size: + print(" VERDICT: right LENGTH, wrong BYTES. So it streams, but") + print(" the assembly differs — likely a different per-frame") + print(" prefix than the chunk form's 11 bytes. Compare below.") + for i in range(min(len(streamed), len(chunked))): + if streamed[i] != chunked[i]: + print(f" first difference at byte {i}") + print(f" chunked {chunked[max(0,i-4):i+8].hex(' ')}") + print(f" streamed {streamed[max(0,i-4):i+8].hex(' ')}") + break + else: + print(f" VERDICT: answered, but {len(streamed)} B against " + f"{ref.size} B expected.") + print(" Inconclusive — dump the frame sizes and look again:") + for i, f in enumerate(frames[:12]): + print(f" frame {i}: data {len(f.data)} B, " + f"page=0x{f.page_key:04x}") + + # ── 3. is the unit still healthy? ───────────────────────────────────── + # The real risk of this experiment is leaving the session wedged, so + # check rather than assume. + print() + try: + p = mm.protocol.poll() + print(f" [3] unit still answering POLL (SUB 0x{p.sub:02x}) — " + f"session is healthy") + except Exception as e: + print(f" [3] ⚠ POLL FAILED after the probe: {type(e).__name__}: {e}") + print(" Reconnect; if that does not help, power-cycle the unit.") + return 4 + finally: + mm.close() + return 0 + + +if __name__ == "__main__": + raise SystemExit(main())