feat(micromate): the event chain -- walk, records, download, and a self-check
Step 4 of docs/micromate_client_spec.md. MicromateEventRef plus list_events(), iter_events(), download_event(), get_event() and decode_error(). 14 new tests; 112 micromate tests total. The load-bearing test replays THOR's captured six-event download session through iter_events() + get_event() and asserts EVERY BYTE WE EMIT MATCHES THOR'S, in order -- 74 frames -- while decoding all six events and cross-checking each waveform's peak vector sum against the float the device computed itself. Two design decisions worth recording: 1. iter_events() EXISTS BECAUSE THOR INTERLEAVES. Its captured order is 0x93 -> 1E -> 0C -> 5A*n -> 0x93 -> 1F -> 0C -> 5A*n, downloading each event before advancing the chain. list_events() walks to the end first, which is fine for browsing but leaves the device cursor parked past the event a later download addresses. 0x5A is key-addressed so it very probably does not care -- but nothing observed says either way, so the interleaved path is the one offered for downloads, and it is the one the replay test exercises. 2. get_event(verify=True) RECOMPUTES THE PEAK VECTOR SUM from the decoded samples and compares it against the device's own 0x0C float. Two independent computations over the same samples, so a disagreement means our decode is wrong. Agreement on the bench events is 0.000%. Cheap insurance in a codebase whose decode failures have historically been silent -- unhandled block tags shorten a channel and nothing raises. It is a decode-correctness check, NOT a truncation detector: a channel cut after its peak still yields the right PVS, and the docstring says so. A test corrupts a stored peak to prove the check actually fires. MicromateEventRef.uid is SERIAL:key, because the key alone is ambiguous across units, and .filename generates THOR's name (<serial>_<YYYYMMDDHHMMSS>.IDFW/H) -- returning None rather than guessing when the record type is unknown, since read_idf_file() dispatches on exactly that suffix. Recorded as a CANDIDATE, not used: 0x06 content[0:4] looks like the EVENT COUNT -- 6 on a unit holding 6 events, zeros on an empty one, and THOR reads it BEFORE the walk then downloads exactly six events without ever reading the chain sentinel. If it holds it lets a caller decide whether to walk at all, which over cellular is the useful part. Two samples on one unit, and content[4:8] reads 9 unexplained, so list_events() still walks to the sentinel: slower by one round trip and correct on evidence rather than inference. Full suite: 465 passed, 16 pre-existing failures unchanged. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Ru8Lg9HkkYvX9VWWo65SmL
This commit is contained in:
+274
-1
@@ -32,11 +32,14 @@ from __future__ import annotations
|
||||
|
||||
import datetime
|
||||
import logging
|
||||
import math
|
||||
import os
|
||||
import struct
|
||||
from typing import Optional
|
||||
|
||||
from minimateplus.transport import BaseTransport
|
||||
|
||||
from .models import MicromateDeviceInfo, MicromateState
|
||||
from .models import MicromateDeviceInfo, MicromateEventRef, MicromateState
|
||||
from .protocol import MicromateProtocol, ProtocolError
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
@@ -57,6 +60,27 @@ _MS_BATTERY = slice(34, 36) # uint16 BE, volts × 100
|
||||
_MS_MEM_TOTAL = slice(36, 40) # uint32 BE
|
||||
_MS_MEM_FREE = slice(40, 44) # uint32 BE
|
||||
|
||||
# `SUB 0x0C` record, relative to content. Established against 7 events across
|
||||
# both firmware lines.
|
||||
_REC_DAY, _REC_MONTH, _REC_YEAR = 0, 1, slice(2, 4)
|
||||
_REC_UNKNOWN_4 = 4 # ⚠ same shape as 0x1C's content[6]; undecoded
|
||||
_REC_HOUR, _REC_MIN, _REC_SEC = 5, 6, 7
|
||||
_REC_TYPE = 11 # 0x07 waveform, 0x08 histogram
|
||||
_REC_LOCATION = 12
|
||||
_REC_SETUP = 34
|
||||
_REC_SERIAL = 76
|
||||
_RECORD_TYPES = {0x07: "waveform", 0x08: "histogram"}
|
||||
|
||||
# The channel labels the record carries, in the order they appear. The peak
|
||||
# float sits `label + 6`; the peak vector sum sits 12 bytes BEFORE "Tran".
|
||||
_REC_CHANNELS = (b"Tran", b"Vert", b"Long", b"Mic")
|
||||
_REC_PVS_BACK = 12
|
||||
|
||||
# A chain walk terminates on an all-zero key. The ceilings below are guards
|
||||
# against a device cursor that never advances, not fleet limits.
|
||||
_NULL_KEY = bytes(4)
|
||||
_MAX_EVENTS = 4096
|
||||
|
||||
# A setup-list walk that does not terminate is a bug, not a big fleet. The
|
||||
# bench unit holds 22 setups; this is a generous ceiling, not a limit.
|
||||
_MAX_SETUPS = 512
|
||||
@@ -71,6 +95,10 @@ def _cstring(buf: bytes, offset: int = 0) -> str:
|
||||
return buf[offset:].split(b"\x00")[0].decode("ascii", "replace").strip()
|
||||
|
||||
|
||||
class DecodeMismatch(ProtocolError):
|
||||
"""Our decoded peak disagrees with the one the device computed itself."""
|
||||
|
||||
|
||||
class MicromateClient:
|
||||
"""High-level read-only client for one Micromate.
|
||||
|
||||
@@ -88,6 +116,7 @@ class MicromateClient:
|
||||
transport, recv_timeout=recv_timeout, strict_checksums=strict_checksums
|
||||
)
|
||||
self._firmware_line: Optional[str] = None
|
||||
self._serial: Optional[str] = None
|
||||
|
||||
# ── Lifecycle ─────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -140,6 +169,7 @@ class MicromateClient:
|
||||
|
||||
manufacturer, model = self._parse_poll(poll.data)
|
||||
serial = _cstring(_content(self._proto.read_serial()))
|
||||
self._serial = serial
|
||||
monitoring = self._parse_state(self._proto.read_state())
|
||||
|
||||
info = MicromateDeviceInfo(
|
||||
@@ -293,3 +323,246 @@ class MicromateClient:
|
||||
f"setup list did not terminate after {_MAX_SETUPS} entries — the "
|
||||
f"device cursor is not advancing"
|
||||
)
|
||||
|
||||
# ── Events ────────────────────────────────────────────────────────────────
|
||||
|
||||
def serial(self) -> str:
|
||||
"""The unit's serial, cached from `connect()` or read on demand.
|
||||
|
||||
Needed by anything that handles an event, because an event key is
|
||||
ambiguous without it — see `MicromateEventRef`.
|
||||
"""
|
||||
if self._serial is None:
|
||||
self._serial = _cstring(_content(self._proto.read_serial()))
|
||||
return self._serial
|
||||
|
||||
def list_events(self, *, with_records: bool = True) -> list[MicromateEventRef]:
|
||||
"""Walk the event chain. `0x93 → 1E`, then `0x93 → 1F` until the sentinel.
|
||||
|
||||
THOR sends `0x93` before **every** chain read, and this mirrors that.
|
||||
The chain ends on an all-zero key.
|
||||
|
||||
⚠ `with_records=True` costs **one extra round trip per event** for the
|
||||
`0x0C` read, and over cellular a round trip is ~0.65 s regardless of
|
||||
size. On a unit with 40 events that is the difference between ~52 s and
|
||||
~78 s. Pass False when you only need "what is here and how big" — but
|
||||
note the **record is where the type and timestamp live**, so without it
|
||||
`ref.filename` is None and `get_event()` cannot pick a suffix.
|
||||
|
||||
⚠ **To download, prefer `iter_events()`.** This walks the whole chain
|
||||
first; THOR interleaves, downloading each event before advancing.
|
||||
`0x5A` addresses an event by key, so downloading afterwards *should*
|
||||
work — but "should" is doing real work in that sentence, and the
|
||||
interleaved order is the one with captures behind it. See
|
||||
`iter_events()`.
|
||||
"""
|
||||
serial = self.serial()
|
||||
refs: list[MicromateEventRef] = []
|
||||
|
||||
for i in range(_MAX_EVENTS):
|
||||
self._proto.arm_event()
|
||||
raw = (self._proto.read_event_first() if i == 0
|
||||
else self._proto.read_event_next())
|
||||
c = _content(raw)
|
||||
if len(c) < 8:
|
||||
raise ProtocolError(
|
||||
f"chain entry {i}: {len(c)} B of content, need 8 (key + size)"
|
||||
)
|
||||
key, size = c[0:4], int.from_bytes(c[4:8], "big")
|
||||
if key == _NULL_KEY:
|
||||
return refs # the sentinel, not an error
|
||||
|
||||
ref = MicromateEventRef(index=i, key=key, size=size, serial=serial)
|
||||
if with_records:
|
||||
self._read_record_into(ref)
|
||||
refs.append(ref)
|
||||
|
||||
raise ProtocolError(
|
||||
f"event chain did not terminate after {_MAX_EVENTS} entries — the "
|
||||
f"device cursor is not advancing"
|
||||
)
|
||||
|
||||
def iter_events(self, *, with_records: bool = True):
|
||||
"""Walk the chain, yielding each event **at the cursor position THOR uses.**
|
||||
|
||||
for ref in mm.iter_events():
|
||||
if ref.is_histogram:
|
||||
continue
|
||||
data = mm.download_event(ref) # ← safe here
|
||||
|
||||
Why this exists alongside `list_events()`: THOR's captured order is
|
||||
|
||||
0x93 → 1E → 0C → 5A×n → 0x93 → 1F → 0C → 5A×n → …
|
||||
|
||||
— it downloads each event *before* advancing the chain. `list_events()`
|
||||
walks to the end first, which is fine for browsing (the browse walk is
|
||||
separately attested) but means a later download happens with the device
|
||||
cursor parked past the event. `0x5A` is key-addressed, so it very
|
||||
probably does not care; nothing observed says it does, and nothing
|
||||
observed says it does not.
|
||||
|
||||
Downloading inside this loop reproduces THOR's sequence exactly, so it
|
||||
is the path to use when it matters. ⚠ Do not advance the generator
|
||||
before finishing with the event it yielded.
|
||||
"""
|
||||
serial = self.serial()
|
||||
for i in range(_MAX_EVENTS):
|
||||
self._proto.arm_event()
|
||||
raw = (self._proto.read_event_first() if i == 0
|
||||
else self._proto.read_event_next())
|
||||
c = _content(raw)
|
||||
if len(c) < 8:
|
||||
raise ProtocolError(
|
||||
f"chain entry {i}: {len(c)} B of content, need 8 (key + size)"
|
||||
)
|
||||
key, size = c[0:4], int.from_bytes(c[4:8], "big")
|
||||
if key == _NULL_KEY:
|
||||
return
|
||||
ref = MicromateEventRef(index=i, key=key, size=size, serial=serial)
|
||||
if with_records:
|
||||
self._read_record_into(ref)
|
||||
yield ref
|
||||
|
||||
raise ProtocolError(
|
||||
f"event chain did not terminate after {_MAX_EVENTS} entries — the "
|
||||
f"device cursor is not advancing"
|
||||
)
|
||||
|
||||
def _read_record_into(self, ref: MicromateEventRef) -> None:
|
||||
"""`SUB 0x0C` — 210 B of content: timestamp, type, names, peaks."""
|
||||
raw = self._proto.read_event_record(ref.key)
|
||||
c = _content(raw)
|
||||
if len(c) <= _REC_SERIAL:
|
||||
raise ProtocolError(f"event record is {len(c)} B, too short to decode")
|
||||
ref.raw_record = raw
|
||||
|
||||
ref.record_type = _RECORD_TYPES.get(c[_REC_TYPE])
|
||||
if ref.record_type is None:
|
||||
# Worth saying out loud rather than filing the event as a waveform:
|
||||
# the suffix decides which codec runs.
|
||||
log.warning("event %s: unknown record type 0x%02x at content[%d]",
|
||||
ref.key_hex, c[_REC_TYPE], _REC_TYPE)
|
||||
|
||||
try:
|
||||
ref.timestamp = datetime.datetime(
|
||||
year=int.from_bytes(c[_REC_YEAR], "big"),
|
||||
month=c[_REC_MONTH], day=c[_REC_DAY],
|
||||
hour=c[_REC_HOUR], minute=c[_REC_MIN], second=c[_REC_SEC],
|
||||
)
|
||||
except ValueError as e:
|
||||
log.warning("event %s: bad timestamp (%s): %s",
|
||||
ref.key_hex, e, c[:8].hex(" "))
|
||||
|
||||
ref.sensor_location = _cstring(c, _REC_LOCATION) or None
|
||||
ref.setup = _cstring(c, _REC_SETUP) or None
|
||||
# Prefer the record's own serial over the cached one — they have always
|
||||
# agreed, but the record is the event's own account of where it came from.
|
||||
if rec_serial := _cstring(c, _REC_SERIAL):
|
||||
ref.serial = rec_serial
|
||||
|
||||
peaks: dict[str, float] = {}
|
||||
for label in _REC_CHANNELS:
|
||||
i = c.find(label)
|
||||
if i < 0 or i + len(label) + 10 > len(c):
|
||||
continue
|
||||
peaks[label.decode()] = struct.unpack(
|
||||
">f", c[i + len(label) + 2: i + len(label) + 6])[0]
|
||||
ref.peaks_ips = peaks or None
|
||||
|
||||
tran = c.find(b"Tran")
|
||||
if tran >= _REC_PVS_BACK:
|
||||
ref.peak_vector_sum_ips = struct.unpack(
|
||||
">f", c[tran - _REC_PVS_BACK: tran - _REC_PVS_BACK + 4])[0]
|
||||
|
||||
def download_event(self, ref: MicromateEventRef) -> bytes:
|
||||
"""The raw `.IDFW`/`.IDFH` bytes, exactly as THOR would have stored them.
|
||||
|
||||
Feeds `micromate.idf_file.read_idf_file()` and `/db/import/idf_file`
|
||||
unchanged — no new codec work is needed for a directly downloaded event.
|
||||
"""
|
||||
return self._proto.read_event_file(ref.key, ref.size)
|
||||
|
||||
def get_event(self, ref: MicromateEventRef, *, verify: bool = True,
|
||||
tolerance: float = 0.01):
|
||||
"""Download and decode one event.
|
||||
|
||||
Returns the codec's `IdfReadResult`. Needs `ref.record_type`, since
|
||||
`read_idf_file()` dispatches on the filename suffix and there is no
|
||||
filename on the wire — so call `list_events(with_records=True)` first.
|
||||
|
||||
⚠ `verify=True` re-computes the **peak vector sum** from the decoded
|
||||
samples and compares it against the float the *device* put in the `0x0C`
|
||||
record. Those are two independent computations over the same samples —
|
||||
the device's from its own firmware, ours from our codec — so a
|
||||
disagreement means our decode is wrong. Measured agreement on the bench
|
||||
events is **0.000%**.
|
||||
|
||||
This is cheap insurance in a codebase whose decode failures have
|
||||
historically been *silent*: unhandled block tags shorten a channel and
|
||||
nothing raises. ⚠ It is a decode-correctness check, **not** a
|
||||
truncation detector — a channel cut after its peak still yields the
|
||||
right PVS.
|
||||
|
||||
Histograms are not verified: `samples` is empty for them.
|
||||
"""
|
||||
if not ref.record_type:
|
||||
raise ValueError(
|
||||
f"event {ref.key_hex}: record_type is unknown, so the codec "
|
||||
f"cannot be dispatched. Use list_events(with_records=True)."
|
||||
)
|
||||
blob = self.download_event(ref)
|
||||
|
||||
import tempfile
|
||||
from .idf_file import read_idf_file
|
||||
|
||||
# read_idf_file dispatches on the suffix, so the bytes need a name.
|
||||
with tempfile.NamedTemporaryFile(suffix=ref.suffix, delete=False) as f:
|
||||
f.write(blob)
|
||||
tmp = f.name
|
||||
try:
|
||||
result = read_idf_file(tmp)
|
||||
finally:
|
||||
os.unlink(tmp)
|
||||
|
||||
if verify and not ref.is_histogram:
|
||||
err = self.decode_error(ref, result)
|
||||
if err is not None and abs(err) > tolerance:
|
||||
raise DecodeMismatch(
|
||||
f"event {ref.uid}: decoded peak vector sum differs from the "
|
||||
f"device's own by {100 * err:+.3f}% (tolerance "
|
||||
f"{100 * tolerance:.1f}%) — the decode is suspect, not the "
|
||||
f"device. Stored {ref.peak_vector_sum_ips:.5f} in/s."
|
||||
)
|
||||
if err is not None:
|
||||
log.debug("event %s: PVS agrees to %+.4f%%", ref.uid, 100 * err)
|
||||
return result
|
||||
|
||||
@staticmethod
|
||||
def decode_error(ref: MicromateEventRef, result) -> Optional[float]:
|
||||
"""Relative error between our decoded PVS and the device's stored one.
|
||||
|
||||
None when either side is unavailable. Positive means the device's
|
||||
figure is higher than ours.
|
||||
"""
|
||||
from .idf_file import geo_count_to_ips
|
||||
|
||||
stored = ref.peak_vector_sum_ips
|
||||
if not stored or not getattr(result, "samples", None):
|
||||
return None
|
||||
|
||||
ch = {k.lower(): v for k, v in result.samples.items()}
|
||||
try:
|
||||
t, v, l = ch["tran"], ch["vert"], ch["long"]
|
||||
except KeyError:
|
||||
return None
|
||||
n = min(len(t), len(v), len(l))
|
||||
if not n:
|
||||
return None
|
||||
|
||||
pvs = max(
|
||||
math.sqrt(geo_count_to_ips(t[i]) ** 2
|
||||
+ geo_count_to_ips(v[i]) ** 2
|
||||
+ geo_count_to_ips(l[i]) ** 2)
|
||||
for i in range(n)
|
||||
)
|
||||
return (stored - pvs) / pvs if pvs else None
|
||||
|
||||
Reference in New Issue
Block a user