diff --git a/bridges/mm_ach_server.py b/bridges/mm_ach_server.py new file mode 100644 index 0000000..2e2c2b0 --- /dev/null +++ b/bridges/mm_ach_server.py @@ -0,0 +1,312 @@ +#!/usr/bin/env python3 +""" +mm_ach_server.py — inbound Auto Call Home server for Micromate (Series IV). + +**Read-only.** It identifies the unit, reads its state and setups, and walks the +event chain. It does not erase, does not write, does not start or stop +monitoring, and downloads nothing unless asked. Series III's `ach_server.py` +has "rescue actions" that quiet a runaway unit; those are deliberately absent +here, because no command has ever been originated against a unit by this project +and an inbound session is the worst place to break that. + +What this is for +---------------- +The inbound call-home session is the last unmapped part of the Series IV +protocol — and it may well not be a separate protocol at all. Series III's +answer, from 72 captured sessions, is that the unit **dials in and then waits**: +the server runs the ordinary handshake and drives with the ordinary command set. +If Series IV behaves the same way, this file is just `MicromateClient` attached +to an inbound socket, and the "unknown" evaporates. + +**The one question that matters is push vs pull**, and it cannot be answered by +sniffing — only by connecting and keeping quiet. So this server **says nothing +for `--grace` seconds** after accepting. If bytes arrive in that window the unit +speaks first (push) and the server reports exactly what it sent. If the window +passes in silence, it is pull, and the normal handshake follows. + +Setup +----- +1. Run this where the unit can reach it: + + python3 bridges/mm_ach_server.py --port 12345 + +2. Point the call-home destination at this machine. **For Series III that + lived in the MODEM's config (ACEmanager), not in the seismograph** — so try + AirLink OS first. If the RX55 can be told to open an outbound TCP session + to `:12345` on serial activity, nothing on the instrument needs to + change, which is much the better outcome. + + If it has to come from the unit instead, its `SUB 0x2C` call-home block + holds a 40-byte dial string — but **writing that is THOR's job, not ours.** + +3. Trigger a call. A scheduled time is the controllable option: repeatable, + needs no monitoring and no thumping. An event-triggered call is the more + interesting session, because it shows how an event gets offered. + +4. Watch. Every frame is decoded as it arrives. + +⚠ Use a bench unit. Changing a call-home destination means a unit stops +reporting to THOR until it is changed back. + +Output per session, in the layout `scratch/mm_frame_parse.py` reads directly: + + bridges/captures/mm_ach_/ + raw_bw_.bin us → unit + raw_s3_.bin unit → us + session_.log + session.json what we learned, including push-vs-pull +""" + +from __future__ import annotations + +import argparse +import datetime +import json +import logging +import socket +import sys +import threading +import time +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) + +from micromate.client import MicromateClient # noqa: E402 +from micromate.framing import MicromateFrameParser # noqa: E402 +from micromate.protocol import ProtocolError # noqa: E402 +from minimateplus.transport import SocketTransport # noqa: E402 + +log = logging.getLogger("mm_ach") + +DEPTHS = ("identify", "state", "events", "download") + + +class Session: + """One inbound call. Captures every byte; decides push vs pull first.""" + + def __init__(self, sock: socket.socket, peer: str, args) -> None: + self.sock, self.peer, self.args = sock, peer, args + self.ts = datetime.datetime.now().strftime("%Y%m%d_%H%M%S") + self.found: dict = {"peer": peer, "started": self.ts} + self._rx: list[bytes] = [] + self._tx: list[bytes] = [] + self._rx_fh = self._tx_fh = None + + # ── byte taps ───────────────────────────────────────────────────────────── + + def _tap(self, transport) -> None: + """Mirror every byte, buffered until we know this is a real unit.""" + orig_read, orig_write = transport.read, transport.write + + def read(n: int) -> bytes: + data = orig_read(n) + if data: + (self._rx_fh.write(data), self._rx_fh.flush()) if self._rx_fh \ + else self._rx.append(data) + return data + + def write(data: bytes) -> None: + orig_write(data) + if data: + (self._tx_fh.write(data), self._tx_fh.flush()) if self._tx_fh \ + else self._tx.append(data) + + transport.read, transport.write = read, write + + def _open_capture(self) -> Path: + """Create the session directory — only once we trust the peer. + + A port scanner should leave nothing on disk. + """ + d = Path(self.args.output) / f"mm_ach_{self.ts}" + d.mkdir(parents=True, exist_ok=True) + # Named raw_bw/raw_s3 so scratch/mm_frame_parse.py reads the pair as-is. + self._tx_fh = open(d / f"raw_bw_{self.ts}.bin", "wb") + self._rx_fh = open(d / f"raw_s3_{self.ts}.bin", "wb") + for c in self._tx: + self._tx_fh.write(c) + for c in self._rx: + self._rx_fh.write(c) + self._tx_fh.flush() + self._rx_fh.flush() + self._tx.clear() + self._rx.clear() + return d + + # ── the experiment ──────────────────────────────────────────────────────── + + def _listen_first(self, transport) -> bytes: + """Say nothing for `--grace` seconds. Does the unit speak first? + + This is the whole reason the server exists. A packet capture cannot + answer it: you have to connect, keep quiet, and see. + """ + print(f" listening in silence for {self.args.grace:.1f} s " + f"(push-vs-pull test)...") + deadline = time.monotonic() + self.args.grace + got = bytearray() + while time.monotonic() < deadline: + chunk = transport.read(4096) + if chunk: + got += chunk + else: + time.sleep(0.02) + return bytes(got) + + def run(self) -> None: + print(f"\n=== inbound from {self.peer} at {self.ts} ===") + transport = SocketTransport(self.sock, peer=self.peer) + self._tap(transport) + + unsolicited = self._listen_first(transport) + self.found["push"] = bool(unsolicited) + self.found["unsolicited_bytes"] = len(unsolicited) + + if unsolicited: + print(f" ** PUSH ** — the unit sent {len(unsolicited)} B unprompted") + print(f" hex : {unsolicited[:64].hex(' ')}") + txt = "".join(chr(c) if 32 <= c < 127 else "." for c in unsolicited[:64]) + print(f" ascii : |{txt}|") + frames = MicromateFrameParser().feed(unsolicited) + if frames: + for f in frames: + print(f" frame : SUB=0x{f.sub:02x} " + f"(answers 0x{f.request_sub:02x}) {len(f.data)} B data, " + f"flags=0x{f.flags:02x}, chk " + f"{'ok' if f.checksum_valid else 'BAD'}") + self.found["unsolicited_subs"] = [f"0x{f.sub:02x}" for f in frames] + else: + print(" frame : none parsed — not a response frame. Could be") + print(" a banner (series III sends 'Operating System'") + print(" on cold boot) or a greeting we have not seen.") + else: + print(" ** PULL ** — silence. The unit is waiting to be driven,") + print(" same as series III. Running the normal handshake.") + + # From here on it is the ordinary client, which is the point. + mm = MicromateClient(transport, recv_timeout=self.args.timeout) + try: + info = mm.connect() + except ProtocolError as e: + print(f" handshake FAILED: {type(e).__name__}: {e}") + if not unsolicited: + print(" Nothing unprompted AND no answer to POLL — so this is") + print(" neither push nor pull as we understand it. The raw") + print(" capture is the thing to look at; keep it.") + self._finish(self._open_capture()) + return + + d = self._open_capture() + print(f" {info}") + self.found["device"] = { + "serial": info.serial, "model": info.model, + "firmware_line": info.firmware_line, "monitoring": info.monitoring, + "active_setup": info.active_setup, + } + + try: + if DEPTHS.index(self.args.depth) >= DEPTHS.index("state"): + st = mm.get_state() + print(f" {st}") + self.found["state"] = { + "monitoring": st.monitoring, + "device_time": st.device_time.isoformat() if st.device_time else None, + "battery_volts": st.battery_volts, + "memory_free_bytes": st.memory_free_bytes, + } + self.found["setups"] = mm.list_setups() + print(f" {len(self.found['setups'])} setups") + + if DEPTHS.index(self.args.depth) >= DEPTHS.index("events"): + refs = mm.list_events() + print(f" {len(refs)} events") + self.found["events"] = [] + for r in refs: + print(f" {r}") + self.found["events"].append({ + "uid": r.uid, "key": r.key_hex, "size": r.size, + "type": r.record_type, "filename": r.filename, + "timestamp": r.timestamp.isoformat() if r.timestamp else None, + "pvs_ips": r.peak_vector_sum_ips, + }) + + if self.args.depth == "download": + # Bounded on purpose: an unknown session is a poor place to + # discover that a 70 KB pull takes 30 s and the unit hung up. + for r in refs[: self.args.max_events]: + blob = mm.download_event(r) + out = d / (r.filename or f"{r.key_hex}.bin") + out.write_bytes(blob) + print(f" saved {out.name} ({len(blob)} B)") + except ProtocolError as e: + print(f" session ended early: {type(e).__name__}: {e}") + self.found["error"] = f"{type(e).__name__}: {e}" + + self._finish(d) + + def _finish(self, d: Path) -> None: + for fh in (self._rx_fh, self._tx_fh): + if fh: + fh.close() + (d / "session.json").write_text(json.dumps(self.found, indent=2)) + try: + self.sock.close() + except OSError: + pass + print(f" saved {d}") + print(f" parse with: python3 scratch/mm_frame_parse.py {d}") + + +def main() -> int: + ap = argparse.ArgumentParser( + description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter + ) + ap.add_argument("--port", type=int, default=12345) + ap.add_argument("--bind", default="0.0.0.0") + ap.add_argument("--output", default="bridges/captures") + ap.add_argument("--timeout", type=float, default=30.0, + help="per-command receive timeout; generous, because an " + "inbound cellular session is the slow case") + ap.add_argument("--grace", type=float, default=5.0, + help="seconds to stay SILENT after accepting, to see whether " + "the unit speaks first. This is the experiment; do not " + "set it to 0 on a first run.") + ap.add_argument("--depth", choices=DEPTHS, default="events", + help="how far to take the session. 'events' walks the chain " + "without downloading (default); 'download' also fetches " + "events, which on cellular is ~0.21 s + bytes/2350 each") + ap.add_argument("--max-events", type=int, default=3, + help="with --depth download, how many to fetch") + ap.add_argument("--once", action="store_true", + help="exit after the first session instead of listening on") + a = ap.parse_args() + + logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(levelname)-7s %(message)s") + + srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + srv.bind((a.bind, a.port)) + srv.listen(5) + print(f"\nMicromate ACH server on {a.bind}:{a.port} (READ-ONLY)") + print(f" depth={a.depth} grace={a.grace}s timeout={a.timeout}s") + print(" waiting for a unit to call in. Ctrl-C to stop.\n") + + try: + while True: + sock, addr = srv.accept() + peer = f"{addr[0]}:{addr[1]}" + s = Session(sock, peer, a) + if a.once: + s.run() + return 0 + threading.Thread(target=s.run, daemon=False).start() + except KeyboardInterrupt: + print("\nstopped") + return 0 + finally: + srv.close() + + +if __name__ == "__main__": + raise SystemExit(main())