"""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(" bytes: return struct.pack(" bytes: return struct.pack(" bytes: return struct.pack(" 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(" 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()