"""CAN transport adapters.""" from __future__ import annotations import threading import time from typing import Protocol from .config import Settings from .protocol import CanFrame class CanTransportError(RuntimeError): """Raised when the CANalyst-II interface cannot be opened or written.""" class CanTransport(Protocol): @property def is_open(self) -> bool: ... def send(self, frame: CanFrame, *, tx_queue_timeout_s: float | None = None) -> None: ... def receive(self, *, timeout_s: float = 0.0) -> CanFrame | None: ... def shutdown(self) -> None: ... class DisabledCanTransport: """Fail-closed transport used until hardware is explicitly enabled.""" @property def is_open(self) -> bool: return False def send(self, frame: CanFrame, *, tx_queue_timeout_s: float | None = None) -> None: del frame del tx_queue_timeout_s raise CanTransportError("hardware access is disabled by MOTOR_HARDWARE_ENABLED") def receive(self, *, timeout_s: float = 0.0) -> CanFrame | None: del timeout_s raise CanTransportError("hardware access is disabled by MOTOR_HARDWARE_ENABLED") def shutdown(self) -> None: return None class PythonCanTransport: """python-can adapter for the CANalyst-II USB backend.""" def __init__(self, settings: Settings) -> None: try: import can except ImportError as exc: # pragma: no cover - packaging failure raise CanTransportError("python-can is not installed; run `uv sync`") from exc try: self._bus = can.Bus( interface="canalystii", channel=settings.can_channel, device=settings.can_device, bitrate=settings.can_bitrate, rx_queue_size=256, receive_own_messages=False, ) except Exception as exc: raise CanTransportError( f"cannot open CANalyst-II device {settings.can_device}, " f"channel {settings.can_channel}: {exc}" ) from exc self._can = can self._lock = threading.Lock() self._is_open = True @property def is_open(self) -> bool: with self._lock: return self._is_open def send(self, frame: CanFrame, *, tx_queue_timeout_s: float | None = None) -> None: message = self._can.Message( arbitration_id=frame.arbitration_id, data=frame.data, is_extended_id=frame.is_extended_id, check=True, ) try: with self._lock: if not self._is_open: raise CanTransportError("CAN transport is closed") # Continuous control uses None so each 2 ms cycle is not # blocked. Safety-critical Stop frames use a finite timeout to # wait until CANalyst-II has processed its hardware TX queue. self._bus.send(message, timeout=tx_queue_timeout_s) except CanTransportError: raise except Exception as exc: raise CanTransportError(f"CAN send failed: {exc}") from exc def receive(self, *, timeout_s: float = 0.0) -> CanFrame | None: deadline = time.monotonic() + max(0.0, timeout_s) try: with self._lock: if not self._is_open: raise CanTransportError("CAN transport is closed") while True: remaining_s = max(0.0, deadline - time.monotonic()) message = self._bus.recv(timeout=remaining_s) if message is None: return None if len(message.data) == 8: return CanFrame( arbitration_id=message.arbitration_id, data=bytes(message.data), is_extended_id=message.is_extended_id, ) if remaining_s == 0.0: return None except CanTransportError: raise except Exception as exc: raise CanTransportError(f"CAN receive failed: {exc}") from exc def shutdown(self) -> None: with self._lock: if not self._is_open: return self._is_open = False try: self._bus.shutdown() except Exception as exc: raise CanTransportError(f"CAN shutdown failed: {exc}") from exc def create_transport(settings: Settings) -> CanTransport: if not settings.hardware_enabled: return DisabledCanTransport() return PythonCanTransport(settings)