Files
calibration/tools/rscap_v2/pipeline_common.py

363 lines
13 KiB
Python

from __future__ import annotations
import binascii
import json
import math
import struct
from pathlib import Path
from typing import Iterable
from capture_format_v2 import CaptureFile, RawChunk, iter_contiguous_segments, read_capture
DOTNET_UNIX_EPOCH_TICKS = 621355968000000000
def ticks_to_unix_ns(ticks: int) -> int:
return (ticks - DOTNET_UNIX_EPOCH_TICKS) * 100
def safe_float(value: str, default=None):
try:
return float(value)
except (TypeError, ValueError):
return default
def safe_int(value: str, default=None):
try:
return int(value)
except (TypeError, ValueError):
return default
def nmea_checksum_valid(line: str) -> bool:
star = line.rfind("*")
if star < 0:
return False
try:
expected = int(line[star + 1:star + 3], 16)
except ValueError:
return False
value = 0
for char in line[1:star]:
value ^= ord(char)
return value == expected
def unicore_crc32(text: str) -> int:
crc = 0
for value in text.encode("ascii", "replace"):
crc ^= value
for _ in range(8):
crc = (crc >> 1) ^ (0xEDB88320 if crc & 1 else 0)
return crc & 0xFFFFFFFF
def unicore_checksum_valid(line: str) -> bool:
star = line.rfind("*")
if star < 0 or len(line) < star + 9:
return False
try:
expected = int(line[star + 1:star + 9], 16)
except ValueError:
return False
return unicore_crc32(line[1:star]) == expected
def parse_checksum(line: str) -> bool:
if line.startswith("$"):
return nmea_checksum_valid(line)
if line.startswith("#"):
return unicore_checksum_valid(line)
return False
def parse_nmea_latlon(value: str, hemisphere: str):
raw = safe_float(value)
if raw is None:
return None
degrees = math.floor(raw / 100.0)
result = degrees + (raw - degrees * 100.0) / 60.0
if hemisphere.upper() in ("S", "W"):
result = -result
return result
def parse_gga(line: str) -> dict:
fields = line[:line.rfind("*")].split(",")
if len(fields) < 10:
raise ValueError("GGA has too few fields")
return {
"type": "GGA",
"position_time_utc": fields[1],
"lat_deg": parse_nmea_latlon(fields[2], fields[3]),
"lon_deg": parse_nmea_latlon(fields[4], fields[5]),
"fix_quality": safe_int(fields[6], -1),
"satellites": safe_int(fields[7], -1),
"hdop": safe_float(fields[8]),
"altitude_m": safe_float(fields[9]),
"geoid_separation_m": safe_float(fields[11]) if len(fields) > 11 else None,
"differential_age_s": safe_float(fields[13]) if len(fields) > 13 else None,
"station_id": fields[14].strip('"') if len(fields) > 14 else "",
}
def parse_heading(line: str) -> dict:
before_crc = line[:line.rfind("*")]
header, payload = before_crc.split(";", 1)
header_fields = header.split(",")
fields = payload.split(",")
if len(fields) < 7:
raise ValueError("UNIHEADINGA has too few fields")
raw_heading = safe_float(fields[3])
return {
"type": "UNIHEADINGA",
"gnss_week": safe_int(header_fields[4]) if len(header_fields) > 4 else None,
"gnss_tow_ms": safe_int(header_fields[5]) if len(header_fields) > 5 else None,
"heading_status": fields[0],
"heading_solution": fields[1],
"baseline_length_m": safe_float(fields[2]),
"raw_heading_deg": raw_heading,
"pitch_deg": safe_float(fields[4]),
"roll_deg": 0.0,
"heading_stddev_deg": safe_float(fields[6]),
"pitch_stddev_deg": safe_float(fields[7]) if len(fields) > 7 else None,
"station_id": fields[8].strip('"') if len(fields) > 8 else "",
"satellites": safe_int(fields[9], -1) if len(fields) > 9 else -1,
"solution_satellites": safe_int(fields[10], -1) if len(fields) > 10 else -1,
"observations": safe_int(fields[11], -1) if len(fields) > 11 else -1,
"multi_count": safe_int(fields[12], -1) if len(fields) > 12 else -1,
"heading_valid": fields[0] == "SOL_COMPUTED" and fields[1] in {"NARROW_INT", "NARROW_FLOAT"},
}
def parse_gnhpr(line: str) -> dict:
"""Parse Wheeltec G90 ``$GNHPR`` heading/pitch output."""
fields = line[:line.rfind("*")].split(",")
if len(fields) < 7:
raise ValueError("GNHPR has too few fields")
quality = safe_int(fields[5], -1)
return {
"type": "GNHPR",
"position_time_utc": fields[1],
"raw_heading_deg": safe_float(fields[2]),
"pitch_deg": safe_float(fields[3]),
"roll_deg": safe_float(fields[4]),
"heading_quality": quality,
"satellites": safe_int(fields[6], -1),
"heading_solution": f"GNHPR_QUALITY_{quality}",
"baseline_length_m": None,
"heading_stddev_deg": None,
"pitch_stddev_deg": None,
"solution_satellites": safe_int(fields[6], -1),
"heading_valid": quality in {4, 5},
}
def parse_pvtslna(line: str) -> dict:
"""Parse Unicore/G90 ``#PVTSLNA`` into GGA-compatible position fields.
``fix_quality`` is synthesized as 4 when checksum-valid coordinates exist so
the existing prepare gate (accepted fixes {4,5}) keeps working. Position
stddevs are retained for audits.
"""
star = line.rfind("*")
fields = line[1:star if star >= 0 else None].split(",")
if len(fields) < 16:
raise ValueError("PVTSLNA has too few fields")
tow = safe_float(fields[5])
return {
"type": "PVTSLNA",
"gnss_week": safe_int(fields[4]),
"gnss_tow_ms": int(tow) if tow is not None else None,
"altitude_m": safe_float(fields[10]),
"lat_deg": safe_float(fields[11]),
"lon_deg": safe_float(fields[12]),
"height_std_m": safe_float(fields[13]),
"latitude_std_m": safe_float(fields[14]),
"longitude_std_m": safe_float(fields[15]),
# Downstream prepare still filters on NMEA-style fix quality.
"fix_quality": 4,
"satellites": -1,
"hdop": None,
"differential_age_s": None,
"position_time_utc": "",
"geoid_separation_m": None,
"station_id": "",
}
def chunk_source(chunks: list[RawChunk], offset: int, end: int) -> dict:
first = chunks[0]
last = chunks[-1]
cursor = 0
start_chunk = first
end_chunk = last
for chunk in chunks:
chunk_start = cursor
chunk_end = cursor + len(chunk.raw)
if chunk_start <= offset < chunk_end:
start_chunk = chunk
if chunk_start < end <= chunk_end:
end_chunk = chunk
break
cursor = chunk_end
return {
"source_segment_id": None,
"source_chunk_sequence_first": start_chunk.sequence,
"source_chunk_sequence_last": end_chunk.sequence,
"source_raw_file_offset": start_chunk.raw_file_offset + max(0, offset - sum(len(c.raw) for c in chunks if c.sequence < start_chunk.sequence)),
"source_raw_byte_length": max(0, end - offset),
}
def parse_rtk_capture(capture: CaptureFile) -> list[dict]:
rows = []
for segment_id, chunks in iter_contiguous_segments(capture.chunks):
stream = b"".join(chunk.raw for chunk in chunks)
cursor = 0
while cursor < len(stream):
newline = stream.find(b"\n", cursor)
if newline < 0:
break
end = newline + 1
raw_line = stream[cursor:end].rstrip(b"\r\n")
cursor = end
if not raw_line:
continue
line = raw_line.decode("ascii", "replace")
valid = parse_checksum(line)
row = {
"type": "UNKNOWN",
"raw_line": line,
"checksum_valid": valid,
"host_receive_utc_ns": ticks_to_unix_ns(chunks[-1].receive_utc_ticks),
"host_receive_monotonic_ticks": chunks[-1].receive_monotonic_ticks,
"source_segment_id": segment_id,
"source_byte_offset_in_segment": cursor - len(raw_line) - 1,
"source_byte_length": len(raw_line) + 1,
}
try:
if line.startswith("$GNGGA") or line.startswith("$GPGGA"):
row.update(parse_gga(line))
elif line.startswith("#PVTSLNA"):
row.update(parse_pvtslna(line))
elif line.startswith("#UNIHEADINGA"):
row.update(parse_heading(line))
elif line.startswith("$GNHPR"):
row.update(parse_gnhpr(line))
except ValueError as ex:
row["parse_error"] = str(ex)
rows.append(row)
return rows
def crc16_hi13(data: bytes) -> int:
crc = 0
for value in data:
crc ^= value << 8
for _ in range(8):
crc = ((crc << 1) ^ 0x1021) & 0xFFFF if crc & 0x8000 else (crc << 1) & 0xFFFF
return crc
def decode_hi91(frame: bytes) -> dict:
f32 = lambda i: struct.unpack_from("<f", frame, i)[0]
return {
"tag": 0x91,
"pps_sync_stamp_ms": int.from_bytes(frame[7:9], "little"),
"temperature_c": struct.unpack_from("<b", frame, 9)[0],
"air_pressure_pa": f32(10),
"device_timestamp_ms": int.from_bytes(frame[14:18], "little"),
"accel_x_mps2": f32(18) * 9.80665,
"accel_y_mps2": f32(22) * 9.80665,
"accel_z_mps2": f32(26) * 9.80665,
"gyro_x_radps": f32(30) * math.pi / 180.0,
"gyro_y_radps": f32(34) * math.pi / 180.0,
"gyro_z_radps": f32(38) * math.pi / 180.0,
"mag_x_ut": f32(42), "mag_y_ut": f32(46), "mag_z_ut": f32(50),
"roll_deg": f32(54), "pitch_deg": f32(58), "yaw_deg": f32(62),
"quaternion_w": f32(66), "quaternion_x": f32(70),
"quaternion_y": f32(74), "quaternion_z": f32(78),
}
def decode_hi92(frame: bytes) -> dict:
i16 = lambda i: struct.unpack_from("<h", frame, i)[0]
i32 = lambda i: struct.unpack_from("<i", frame, i)[0]
return {
"tag": 0x92,
"status": int.from_bytes(frame[7:9], "little"),
"temperature_c": struct.unpack_from("<b", frame, 9)[0],
"pps_sync_stamp_ms": int.from_bytes(frame[10:12], "little"),
"air_pressure_pa": i16(12) + 100000.0,
"heave_m": i16(14) * 0.001,
"gyro_x_radps": i16(16) * 0.001, "gyro_y_radps": i16(18) * 0.001, "gyro_z_radps": i16(20) * 0.001,
"accel_x_mps2": i16(22) * 0.0048828, "accel_y_mps2": i16(24) * 0.0048828, "accel_z_mps2": i16(26) * 0.0048828,
"mag_x_ut": i16(28) * 0.030517, "mag_y_ut": i16(30) * 0.030517, "mag_z_ut": i16(32) * 0.030517,
"roll_deg": i32(34) * 0.001, "pitch_deg": i32(38) * 0.001, "yaw_deg": i32(42) * 0.001,
"quaternion_w": i16(46) * 0.0001, "quaternion_x": i16(48) * 0.0001,
"quaternion_y": i16(50) * 0.0001, "quaternion_z": i16(52) * 0.0001,
}
def parse_imu_capture(capture: CaptureFile) -> list[dict]:
rows = []
for segment_id, chunks in iter_contiguous_segments(capture.chunks):
stream = b"".join(chunk.raw for chunk in chunks)
cursor = 0
while True:
start = stream.find(b"\x5a\xa5", cursor)
if start < 0 or start + 6 > len(stream):
break
payload_length = int.from_bytes(stream[start + 2:start + 4], "little")
frame_length = 6 + payload_length
if payload_length <= 0 or payload_length > 512:
cursor = start + 1
continue
if start + frame_length > len(stream):
break
frame = stream[start:start + frame_length]
expected = int.from_bytes(frame[4:6], "little")
actual = crc16_hi13(frame[:4] + frame[6:])
end = start + frame_length
source = chunk_source(chunks, start, end)
source["source_segment_id"] = segment_id
row = {
"type": "HI13",
"tag": frame[6],
"frame_length": frame_length,
"crc_valid": expected == actual,
"host_receive_utc_ns": ticks_to_unix_ns(chunks[-1].receive_utc_ticks),
"host_receive_monotonic_ticks": chunks[-1].receive_monotonic_ticks,
"source_segment_id": segment_id,
"source_byte_offset_in_segment": start,
"source_byte_length": frame_length,
"raw_frame_hex": frame.hex(),
}
if expected == actual:
try:
row.update(decode_hi91(frame) if frame[6] == 0x91 else decode_hi92(frame) if frame[6] == 0x92 else {})
except (IndexError, struct.error, ValueError) as ex:
row["parse_error"] = str(ex)
rows.append(row)
cursor = end
return rows
def write_jsonl(path: Path, rows: Iterable[dict]) -> None:
with path.open("w", encoding="utf-8", newline="\n") as stream:
for row in rows:
stream.write(json.dumps(row, ensure_ascii=False, separators=(",", ":")) + "\n")
def write_json(path: Path, value: dict) -> None:
path.write_text(json.dumps(value, ensure_ascii=False, indent=2), encoding="utf-8")
def load_jsonl(path: Path) -> list[dict]:
with path.open("r", encoding="utf-8") as stream:
return [json.loads(line) for line in stream if line.strip()]