Files
radar_system/python_app/tests/test_kamil_adc_service.py
2026-07-02 17:44:24 +03:00

320 lines
13 KiB
Python

"""Tests for the Kamil ADC TTY reader, config round-trip, and producer wiring."""
from __future__ import annotations
import json
import os
from pathlib import Path
import pty
import struct
import subprocess
import sys
import tempfile
import time
import tty
import unittest
from unittest import mock
from python_app.hardware_full.kamil_adc import KamilAdcService, KamilAdcTtyReader
from python_app.hardware_full.kamil_adc.protocol import (
COMBO_MARKER,
MAIN_MARKER,
REFERENCE_MARKER,
)
from python_app.models.run_config_model import RunConfigModel
from python_app.orchestration.process_supervisor import ProcessSupervisor
def _boundary() -> bytes:
return struct.pack("<HHHH", MAIN_MARKER, 0xFFFF, 0xFFFF, 0xFFFF)
def _main(step: int, real: int, imag: int) -> bytes:
return struct.pack("<HHhh", MAIN_MARKER, step, real, imag)
def _reference(step: int, real: int, imag: int) -> bytes:
return struct.pack("<HHhh", REFERENCE_MARKER, step, real, imag)
def _combo(input_pos: int, output_pos: int, dirty: int = 0) -> bytes:
return struct.pack("<HHhh", COMBO_MARKER, input_pos, output_pos, dirty)
class KamilAdcTtyReaderTest(unittest.TestCase):
"""End-to-end tests over a PTY exercising the background reader thread."""
def _open_pty_reader(self) -> tuple[int, int, KamilAdcTtyReader]:
master_fd, slave_fd = pty.openpty()
tty.setraw(slave_fd)
reader = KamilAdcTtyReader(os.ttyname(slave_fd))
reader.open()
return master_fd, slave_fd, reader
@staticmethod
def _close(master_fd: int, slave_fd: int, reader: KamilAdcTtyReader) -> None:
try:
reader.close()
finally:
os.close(master_fd)
os.close(slave_fd)
def test_publishes_sweep_with_aligned_main_and_reference(self) -> None:
master_fd, slave_fd, reader = self._open_pty_reader()
try:
os.write(
master_fd,
_boundary()
+ _main(1, 10, -1) + _reference(1, 100, 5)
+ _main(2, -20, 2) + _reference(2, 200, 6)
+ _boundary(),
)
sweep = reader.read_sweep(timeout_s=1.0)
self.assertEqual(sweep.steps.tolist(), [1, 2])
self.assertEqual(sweep.main.tolist(), [complex(10, -1), complex(-20, 2)])
self.assertEqual(sweep.reference.tolist(), [complex(100, 5), complex(200, 6)])
finally:
self._close(master_fd, slave_fd, reader)
def test_variable_length_sweeps_are_allowed(self) -> None:
"""Unlike the old format, sweep length may vary — no locking, just resample later."""
master_fd, slave_fd, reader = self._open_pty_reader()
try:
# One complete sweep per write, read between, so delivery is deterministic
# (the reader publishes only the latest, overwriting unread sweeps).
os.write(master_fd, _boundary() + _main(1, 1, 0) + _reference(1, 9, 0) + _boundary())
first = reader.read_sweep(timeout_s=1.0)
self.assertEqual(first.steps.tolist(), [1])
os.write(
master_fd,
_main(1, 2, 0) + _reference(1, 8, 0)
+ _main(2, 3, 0) + _reference(2, 7, 0)
+ _boundary(),
)
second = reader.read_sweep(timeout_s=1.0)
self.assertEqual(second.steps.tolist(), [1, 2])
finally:
self._close(master_fd, slave_fd, reader)
def test_only_latest_sweep_is_published(self) -> None:
master_fd, slave_fd, reader = self._open_pty_reader()
try:
payload = (
_boundary() + _main(1, 1, 0) + _reference(1, 1, 0)
+ _boundary() + _main(1, 2, 0) + _reference(1, 2, 0)
+ _boundary() + _main(1, 3, 0) + _reference(1, 3, 0)
+ _boundary()
)
os.write(master_fd, payload)
deadline = time.monotonic() + 1.0
while time.monotonic() < deadline and reader.published_count < 3:
time.sleep(0.005)
self.assertGreaterEqual(reader.published_count, 3)
sweep = reader.read_sweep(timeout_s=1.0)
self.assertEqual(sweep.main.tolist(), [complex(3, 0)])
finally:
self._close(master_fd, slave_fd, reader)
def test_corrupt_frame_fails_fast(self) -> None:
master_fd, slave_fd, reader = self._open_pty_reader()
try:
corrupt = struct.pack("<HHhh", 0x001A, 1, 5, 5) # unknown marker
os.write(master_fd, _boundary() + _main(1, 10, -1) + corrupt + _boundary())
with self.assertRaises((ValueError, RuntimeError)):
reader.read_sweep(timeout_s=1.0)
finally:
self._close(master_fd, slave_fd, reader)
def test_no_completed_sweep_times_out(self) -> None:
master_fd, slave_fd, reader = self._open_pty_reader()
try:
os.write(master_fd, _boundary() + _main(1, 10, -1) + _reference(1, 1, 0))
with self.assertRaisesRegex(TimeoutError, "Timed out waiting for Kamil ADC sweep"):
reader.read_sweep(timeout_s=0.1)
finally:
self._close(master_fd, slave_fd, reader)
def test_read_sweep_for_demuxes_by_combo(self) -> None:
master_fd, slave_fd, reader = self._open_pty_reader()
try:
os.write(
master_fd,
_boundary() + _combo(0, 0) + _main(1, 11, 0) + _reference(1, 1, 0)
+ _boundary() + _combo(0, 1) + _main(1, 22, 0) + _reference(1, 1, 0)
+ _boundary(),
)
# Each combination is served from its own slot, regardless of order.
second = reader.read_sweep_for((0, 1), timeout_s=1.0)
self.assertEqual(second.main.real.tolist(), [22])
self.assertEqual(second.combo, (0, 1))
first = reader.read_sweep_for((0, 0), timeout_s=1.0)
self.assertEqual(first.main.real.tolist(), [11])
finally:
self._close(master_fd, slave_fd, reader)
def test_read_sweep_for_drops_dirty_and_takes_retake(self) -> None:
master_fd, slave_fd, reader = self._open_pty_reader()
try:
os.write(
master_fd,
# A dirty combo (0,1) sweep, then its clean re-take of the same combo.
_boundary() + _combo(0, 1, dirty=1) + _main(1, 99, 0) + _reference(1, 1, 0)
+ _boundary() + _combo(0, 1) + _main(1, 42, 0) + _reference(1, 1, 0)
+ _boundary(),
)
sweep = reader.read_sweep_for((0, 1), timeout_s=1.0)
# The dirty sweep (99) is dropped; only the clean re-take (42) is served.
self.assertEqual(sweep.main.real.tolist(), [42])
self.assertFalse(sweep.dirty)
finally:
self._close(master_fd, slave_fd, reader)
class KamilAdcConfigTest(unittest.TestCase):
def test_config_round_trip_preserves_kamil_sections(self) -> None:
payload = {
"radar": {
"model": "kamil_adc",
"serial": "kamil_adc",
"driver_mode": "native",
"kamil_adc": {
"project_dir": "",
"executable_path": "build/bin/kamil_adc_collector",
"tty_path": "/tmp/ttyADC_data",
"args": ["profile:phase", "do8_freq_ref"],
"env": {"ADC_ENV": "1"},
"startup_timeout_s": 7.0,
"sweep_timeout_s": 8.0,
"stop_timeout_s": 3.0,
"phase_calibration": {
"phase0_rad": 1.5,
"freq0_hz": 2_046_000_000.0,
"phase1_rad": 301.0,
"freq1_hz": 5_612_000_000.0,
},
"band": {
"start_hz": 2_100_000_000.0,
"stop_hz": 5_500_000_000.0,
"points": 1024,
},
},
"laser_control": {"enabled": True, "port": "/dev/ttyUSB0", "mode": "manual"},
"sweep": {"start_hz": 1.0, "stop_hz": 2.0, "points": 2},
},
"switches": {"port1": {"positions": 1}, "port2": {"positions": 1}},
}
encoded = RunConfigModel.from_dict(payload).to_dict()
kamil = encoded["radar"]["kamil_adc"]
self.assertEqual(encoded["radar"]["model"], "kamil_adc")
self.assertEqual(kamil["executable_path"], "build/bin/kamil_adc_collector")
self.assertEqual(kamil["args"], ["profile:phase", "do8_freq_ref"])
self.assertEqual(kamil["phase_calibration"]["phase0_rad"], 1.5)
self.assertEqual(kamil["phase_calibration"]["freq1_hz"], 5_612_000_000.0)
self.assertEqual(kamil["band"], {"start_hz": 2_100_000_000.0, "stop_hz": 5_500_000_000.0, "points": 1024})
def test_close_never_raises_when_collector_refuses_to_die(self) -> None:
"""close() must stay exception-safe even if a SIGKILL'd collector is not
reaped within the grace window (e.g. wedged in USB D-state)."""
with tempfile.TemporaryDirectory() as tmp_dir:
config = RunConfigModel.from_dict(
{
"radar": {
"model": "kamil_adc",
"driver_mode": "native",
"kamil_adc": {
"project_dir": tmp_dir,
"executable_path": "/bin/sh", # any real executable
"tty_path": "/tmp/ttyADC_test",
},
},
"switches": {"port1": {"positions": 1}, "port2": {"positions": 1}},
}
)
service = KamilAdcService(config)
class _UnreapableProcess:
pid = 2_000_000_000 # implausible; killpg is patched out below anyway
def poll(self) -> None:
return None # always "alive"
def wait(self, timeout: float | None = None) -> int:
raise subprocess.TimeoutExpired(cmd="kamil_adc_collector", timeout=timeout)
service._process = _UnreapableProcess() # type: ignore[assignment]
with mock.patch("python_app.hardware_full.kamil_adc.service.os.killpg"):
service.close() # must not raise
def test_drain_after_switch_waits_for_fresh_sweeps(self) -> None:
"""After a switch change, drain must skip the configured number of freshly
published sweeps before returning, so the next capture is post-switch."""
import types
from python_app.hardware_full.kamil_adc import service as service_module
with tempfile.TemporaryDirectory() as tmp_dir:
config = RunConfigModel.from_dict(
{
"radar": {
"model": "kamil_adc",
"driver_mode": "native",
"kamil_adc": {
"project_dir": tmp_dir,
"executable_path": "/bin/sh",
"tty_path": "/tmp/ttyADC_test",
"sweep_timeout_s": 5.0,
},
},
"switches": {"port1": {"positions": 1}, "port2": {"positions": 1}},
}
)
service = KamilAdcService(config)
service._reader = types.SimpleNamespace(published_count=10) # type: ignore[assignment]
service._process = types.SimpleNamespace(poll=lambda: None) # type: ignore[assignment]
# Each poll-sleep advances the published count, as the reader thread would.
def _advance(_seconds: float) -> None:
service._reader.published_count += 1
with mock.patch.object(service_module.time, "sleep", _advance):
service.drain_after_switch(sweeps=3)
# Started at 10, must have waited for at least 3 more sweeps.
self.assertGreaterEqual(service._reader.published_count, 13)
def test_drain_after_switch_is_noop_when_not_open(self) -> None:
with tempfile.TemporaryDirectory() as tmp_dir:
config = RunConfigModel.from_dict(
{
"radar": {
"model": "kamil_adc",
"driver_mode": "native",
"kamil_adc": {
"project_dir": tmp_dir,
"executable_path": "/bin/sh",
"tty_path": "/tmp/ttyADC_test",
},
},
"switches": {"port1": {"positions": 1}, "port2": {"positions": 1}},
}
)
KamilAdcService(config).drain_after_switch() # no reader → must not raise
def test_supervisor_selects_kamil_adc_producer(self) -> None:
with tempfile.TemporaryDirectory() as tmp_dir:
config_path = Path(tmp_dir) / "run_config.json"
config_path.write_text(json.dumps({"radar": {"model": "kamil_adc"}}), encoding="utf-8")
command = ProcessSupervisor(Path("/repo"))._acquisition_command(config_path)
self.assertEqual(
command,
[sys.executable, "-m", "python_app.scripts.kamil_adc_raw_producer", "--config", str(config_path)],
)
if __name__ == "__main__":
unittest.main()