Coverage for kv4p/__init__.py: 73%

263 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-13 13:27 +0000

1"""KV4P radio integration.""" 

2 

3from __future__ import annotations 

4 

5__version__ = "0.0.1" 

6 

7import logging 

8from collections.abc import Callable 

9 

10from kv4p.constants.kiss import KISS_CMD_DATA, KISS_CMD_SETHARDWARE 

11from kv4p.constants.messages import ( 

12 HOST_STATE_ENABLE_STATUS_REPORTS, 

13 HOST_STATE_FILTER_HIGH, 

14 HOST_STATE_FILTER_LOW, 

15 HOST_STATE_FILTER_PRE, 

16 HOST_STATE_HIGH_POWER, 

17 HOST_STATE_RSSI_ENABLED, 

18 HOST_STATE_RX_AUDIO_OPEN, 

19 HOST_STATE_TX_ALLOWED, 

20) 

21from kv4p.constants.vendor import ( 

22 COMMAND_AUDIO_ADPCM, 

23 COMMAND_AUDIO_OPUS, 

24 COMMAND_DEBUG_DEBUG, 

25 COMMAND_DEBUG_ERROR, 

26 COMMAND_DEBUG_INFO, 

27 COMMAND_DEBUG_TRACE, 

28 COMMAND_DEBUG_WARN, 

29 COMMAND_DEVICE_STATE, 

30 COMMAND_HELLO, 

31 COMMAND_HOST_DESIRED_STATE, 

32 COMMAND_WINDOW_UPDATE, 

33 KV4P_PROTOCOL_VERSION, 

34 KV4P_VENDOR_HEADER_LEN, 

35 KV4P_VENDOR_PREFIX, 

36) 

37from kv4p.flow_control import FlowControlWindow 

38from kv4p.messages.desired_state import HostDesiredState, bandwidth_to_dra818 

39from kv4p.messages.device_state import DeviceState, RadioMode 

40from kv4p.messages.hello import Hello 

41from kv4p.messages.version import RadioFeatures 

42from kv4p.messages.window_update import WindowUpdate 

43from kv4p.protocol.kiss import encode_kiss_frame 

44from kv4p.state_tracker import DeviceStateTracker 

45from kv4p.transports import Kv4pTransport 

46 

47logger = logging.getLogger(__name__) 

48 

49 

50def encode_vendor_payload(command: int, payload: bytes = b"") -> bytes: 

51 """Build KV4P vendor payload for a KISS SETHARDWARE frame.""" 

52 return KV4P_VENDOR_PREFIX + bytes([KV4P_PROTOCOL_VERSION, command]) + payload 

53 

54 

55def decode_vendor_payload(payload: bytes) -> tuple[int, bytes] | None: 

56 """Parse KV4P vendor payload.""" 

57 if len(payload) < KV4P_VENDOR_HEADER_LEN: 

58 return None 

59 if payload[:4] != KV4P_VENDOR_PREFIX: 

60 return None 

61 if payload[4] != KV4P_PROTOCOL_VERSION: 

62 return None 

63 return payload[5], payload[6:] 

64 

65 

66class RadioNotReadyError(RuntimeError): 

67 """Raised when an operation requires HELLO to have been received but it hasn't.""" 

68 

69 

70class RadioTransportError(RadioNotReadyError): 

71 """Raised when an operation is attempted after the transport failed unexpectedly.""" 

72 

73 

74class Kv4pRadio: 

75 """KV4P-HT radio side of the bridge.""" 

76 

77 def __init__(self, transport: Kv4pTransport) -> None: 

78 self._transport = transport 

79 self._flow = FlowControlWindow() 

80 self._on_rx_audio: Callable[[bytes], None] | None = None 

81 self._on_sql: Callable[[bool], None] | None = None 

82 self._on_phy_ptt: Callable[[bool], None] | None = None 

83 self._on_tx_active: Callable[[bool], None] | None = None 

84 self._on_ax25_frame: Callable[[bytes], None] | None = None 

85 self._on_device_state: Callable[[DeviceState], None] | None = None 

86 self._tracker = DeviceStateTracker( 

87 send_desired_state=self._send_desired_state, 

88 on_rx_audio=lambda payload: self._on_rx_audio(payload) if self._on_rx_audio else None, 

89 on_sql=lambda state: self._on_sql(state) if self._on_sql else None, 

90 on_phy_ptt=lambda state: self._on_phy_ptt(state) if self._on_phy_ptt else None, 

91 on_tx_active=lambda state: self._on_tx_active(state) if self._on_tx_active else None, 

92 ) 

93 self._open = False 

94 self._transport_error: Exception | None = None 

95 

96 self._dispatch: dict[int, Callable[[bytes], None]] = { 

97 COMMAND_DEBUG_INFO: self._handle_debug, 

98 COMMAND_DEBUG_ERROR: self._handle_debug, 

99 COMMAND_DEBUG_WARN: self._handle_debug, 

100 COMMAND_DEBUG_DEBUG: self._handle_debug, 

101 COMMAND_DEBUG_TRACE: self._handle_debug, 

102 COMMAND_HELLO: self._handle_hello, 

103 COMMAND_DEVICE_STATE: self._handle_device_state, 

104 COMMAND_AUDIO_OPUS: self._handle_rx_audio, 

105 COMMAND_AUDIO_ADPCM: self._handle_rx_audio, 

106 COMMAND_WINDOW_UPDATE: self._handle_window_update, 

107 } 

108 

109 # -- lifecycle ------------------------------------------------------------- 

110 

111 def __enter__(self) -> Kv4pRadio: 

112 try: 

113 self.connect() 

114 except BaseException: 

115 self.disconnect() 

116 raise 

117 return self 

118 

119 def __exit__(self, _exc_type: object, _exc: object, _tb: object) -> None: 

120 self.disconnect() 

121 

122 def connect(self, hello_timeout: float = 5.0) -> None: 

123 """Open the transport, reset the radio, and wait for HELLO.""" 

124 if self._open: 

125 return 

126 self._transport_error = None 

127 self._transport.open(self._on_kiss_frame, self._on_transport_error) 

128 self._open = True 

129 logger.info("radio connect") 

130 self.reset(hello_timeout=hello_timeout) 

131 

132 def disconnect(self) -> None: 

133 """Close radio transport.""" 

134 if not self._open: 

135 return 

136 logger.info("radio disconnect") 

137 try: 

138 if self.is_ready: 

139 self.set_ptt(False) 

140 except Exception: 

141 logger.exception("failed to clear PTT during disconnect") 

142 

143 self._transport.close() 

144 self._open = False 

145 logger.info("radio disconnected") 

146 

147 def reset(self, hello_timeout: float = 5.0) -> None: 

148 """Hardware-reset the radio and wait for it to re-announce itself via HELLO. 

149 

150 Public so callers can recover a radio that has hung, not just at connect(). 

151 """ 

152 if not self._open: 

153 raise RuntimeError("radio transport is not connected") 

154 self._transport.reset() 

155 if not self._tracker.wait_for_hello(timeout=hello_timeout): 

156 raise TimeoutError("timed out waiting for HELLO after reset") 

157 

158 # -- configuration ----------------------------------------------------------- 

159 # 

160 # All settings below are seeded from the DeviceState carried in HELLO — 

161 # the firmware always reports its actual tuned state there, right after 

162 # open()/reset() forces a reboot. Reading a property never involves I/O; 

163 # each set_*() call sends a full HostDesiredState snapshot (the protocol 

164 # always wants the complete state, not a delta). 

165 

166 def set_frequency(self, freq: float, txfreq: float | None = None) -> None: 

167 """Update the radio's frequency, validated against the range reported in HELLO. 

168 

169 `freq` sets both RX and TX (simplex). Pass `txfreq` too for 

170 split/repeater operation, where TX differs from RX. 

171 """ 

172 self._require_ready() 

173 version = self._tracker.hello.version 

174 for value in (freq, txfreq) if txfreq is not None else (freq,): 

175 if not (version.min_radio_freq <= value <= version.max_radio_freq): 

176 raise ValueError( 

177 f"frequency {value} outside radio range " 

178 f"{version.min_radio_freq}-{version.max_radio_freq}" 

179 ) 

180 self._tracker.set_frequency(rx=freq, tx=txfreq if txfreq is not None else freq) 

181 

182 def set_bandwidth(self, bandwidth: str) -> None: 

183 """Update bandwidth ("12.5k" or "25k").""" 

184 self._require_ready() 

185 self._tracker.set_bandwidth(bandwidth_to_dra818(bandwidth)) 

186 

187 def set_squelch(self, squelch: int) -> None: 

188 """Update squelch level.""" 

189 self._require_ready() 

190 self._tracker.set_squelch(squelch) 

191 

192 def set_ctcss(self, *, rx: int | None = None, tx: int | None = None) -> None: 

193 """Update RX/TX CTCSS tone.""" 

194 self._require_ready() 

195 self._tracker.set_ctcss(rx=rx, tx=tx) 

196 

197 def set_high_power(self, enabled: bool) -> None: 

198 """Enable/disable high power output.""" 

199 self._require_ready() 

200 self._tracker.set_flag(HOST_STATE_HIGH_POWER, enabled) 

201 

202 def set_tx_allowed(self, enabled: bool) -> None: 

203 """Enable/disable TX capability.""" 

204 self._require_ready() 

205 self._tracker.set_flag(HOST_STATE_TX_ALLOWED, enabled) 

206 

207 def set_rssi(self, enabled: bool) -> None: 

208 """Enable/disable RSSI reporting.""" 

209 self._require_ready() 

210 self._tracker.set_flag(HOST_STATE_RSSI_ENABLED, enabled) 

211 

212 def set_rx_audio_open(self, enabled: bool) -> None: 

213 """Enable/disable receiving RX audio from the firmware.""" 

214 self._require_ready() 

215 self._tracker.set_flag(HOST_STATE_RX_AUDIO_OPEN, enabled) 

216 

217 def set_status_reports(self, enabled: bool) -> None: 

218 """Enable/disable periodic DEVICE_STATE reports from the firmware.""" 

219 self._require_ready() 

220 self._tracker.set_flag(HOST_STATE_ENABLE_STATUS_REPORTS, enabled) 

221 

222 def set_filters(self, *, pre: bool | None = None, high: bool | None = None, low: bool | None = None) -> None: 

223 """Enable/disable the pre-emphasis/high-pass/low-pass audio filters.""" 

224 self._require_ready() 

225 for flag, value in ( 

226 (HOST_STATE_FILTER_PRE, pre), 

227 (HOST_STATE_FILTER_HIGH, high), 

228 (HOST_STATE_FILTER_LOW, low), 

229 ): 

230 if value is not None: 

231 self._tracker.set_flag(flag, value) 

232 

233 def set_ptt(self, enabled: bool) -> None: 

234 """Set PTT requested state.""" 

235 self._require_ready() 

236 self._tracker.request_ptt(enabled) 

237 

238 # -- data path --------------------------------------------------------------- 

239 

240 def send_tx_audio(self, payload: bytes) -> None: 

241 """Send KV4P-native TX audio payload.""" 

242 self._require_ready() 

243 self._send_vendor(self._tracker.tx_audio_command, payload, flow_controlled=True) 

244 

245 def send_ax25_frame(self, payload: bytes) -> None: 

246 """Send a raw AX.25 frame over the KISS data port.""" 

247 logger.debug("ax25 frame tx bytes=%d", len(payload)) 

248 self._claim_and_write(KISS_CMD_DATA, payload, flow_controlled=True) 

249 

250 def flush(self) -> None: 

251 """Flush pending serial writes.""" 

252 self._transport.flush() 

253 

254 def on_rx_audio(self, callback: Callable[[bytes], None] | None) -> None: 

255 """Register callback for incoming audio frames.""" 

256 self._on_rx_audio = callback 

257 

258 def on_sql(self, callback: Callable[[bool], None] | None) -> None: 

259 """Register callback for squelch open/close events.""" 

260 self._on_sql = callback 

261 

262 def on_phy_ptt(self, callback: Callable[[bool], None] | None) -> None: 

263 """Register callback for physical PTT down/up events.""" 

264 self._on_phy_ptt = callback 

265 

266 def on_tx_active(self, callback: Callable[[bool], None] | None) -> None: 

267 """Register callback for TX active start/stop events.""" 

268 self._on_tx_active = callback 

269 

270 def on_ax25_frame(self, callback: Callable[[bytes], None] | None) -> None: 

271 """Register callback for incoming AX.25 frames.""" 

272 self._on_ax25_frame = callback 

273 

274 def on_device_state(self, callback: Callable[[DeviceState], None] | None) -> None: 

275 """Register callback for device state updates.""" 

276 self._on_device_state = callback 

277 

278 # -- read-only state ----------------------------------------------------- 

279 

280 @property 

281 def hello(self) -> Hello | None: 

282 return self._tracker.hello 

283 

284 @property 

285 def is_ready(self) -> bool: 

286 return self._open and self._transport_error is None and self._tracker.hello is not None 

287 

288 @property 

289 def phy_ptt(self) -> bool: 

290 return self._tracker.phy_ptt 

291 

292 @property 

293 def tx_active(self) -> bool: 

294 return self._tracker.tx_active 

295 

296 @property 

297 def sql_open(self) -> bool: 

298 return self._tracker.sql_open 

299 

300 @property 

301 def mode(self) -> RadioMode | None: 

302 return self._tracker.mode 

303 

304 @property 

305 def codec(self) -> int: 

306 """Audio codec command in use: COMMAND_AUDIO_OPUS or COMMAND_AUDIO_ADPCM.""" 

307 return self._tracker.tx_audio_command 

308 

309 @property 

310 def freq_rx(self) -> float: 

311 return self._tracker.freq_rx 

312 

313 @property 

314 def freq_tx(self) -> float: 

315 return self._tracker.freq_tx 

316 

317 @property 

318 def bandwidth(self) -> str: 

319 return self._tracker.bandwidth 

320 

321 @property 

322 def squelch(self) -> int: 

323 return self._tracker.squelch 

324 

325 @property 

326 def ctcss_rx(self) -> int: 

327 return self._tracker.ctcss_rx 

328 

329 @property 

330 def ctcss_tx(self) -> int: 

331 return self._tracker.ctcss_tx 

332 

333 @property 

334 def high_power(self) -> bool: 

335 return bool(self._tracker.flags & HOST_STATE_HIGH_POWER) 

336 

337 @property 

338 def tx_allowed(self) -> bool: 

339 return bool(self._tracker.flags & HOST_STATE_TX_ALLOWED) 

340 

341 @property 

342 def rssi(self) -> bool: 

343 return bool(self._tracker.flags & HOST_STATE_RSSI_ENABLED) 

344 

345 @property 

346 def rx_audio_open(self) -> bool: 

347 return bool(self._tracker.flags & HOST_STATE_RX_AUDIO_OPEN) 

348 

349 @property 

350 def status_reports(self) -> bool: 

351 return bool(self._tracker.flags & HOST_STATE_ENABLE_STATUS_REPORTS) 

352 

353 @property 

354 def filter_pre(self) -> bool: 

355 return bool(self._tracker.flags & HOST_STATE_FILTER_PRE) 

356 

357 @property 

358 def filter_high(self) -> bool: 

359 return bool(self._tracker.flags & HOST_STATE_FILTER_HIGH) 

360 

361 @property 

362 def filter_low(self) -> bool: 

363 return bool(self._tracker.flags & HOST_STATE_FILTER_LOW) 

364 

365 # -- incoming frame routing ----------------------------------------------- 

366 

367 def _on_kiss_frame(self, kiss_command: int, payload: bytes) -> None: 

368 port_command = kiss_command & 0x0F 

369 if port_command == KISS_CMD_DATA: 

370 self._handle_ax25_frame(payload) 

371 return 

372 

373 if port_command != KISS_CMD_SETHARDWARE: 

374 logger.debug("ignore KISS command=0x%02x bytes=%d", kiss_command, len(payload)) 

375 return 

376 

377 decoded = decode_vendor_payload(payload) 

378 if decoded is None: 

379 logger.warning("ignore non-KV4P vendor frame bytes=%d", len(payload)) 

380 return 

381 

382 command, body = decoded 

383 handler = self._dispatch.get(command) 

384 if handler is None: 

385 logger.debug("ignore KV4P command=0x%02x bytes=%d", command, len(body)) 

386 return 

387 try: 

388 handler(body) 

389 except Exception: 

390 logger.exception("failed to handle KV4P command=0x%02x bytes=%d", command, len(body)) 

391 

392 def _handle_ax25_frame(self, payload: bytes) -> None: 

393 logger.debug("ax25 frame rx bytes=%d", len(payload)) 

394 if self._on_ax25_frame is not None: 

395 self._on_ax25_frame(payload) 

396 

397 def _handle_debug(self, payload: bytes) -> None: 

398 text = payload.decode("utf-8", errors="replace").strip() 

399 if text: 

400 logger.info("firmware: %s", text) 

401 

402 def _handle_hello(self, payload: bytes) -> None: 

403 hello = Hello.from_bytes(payload) 

404 self._flow.reset(hello.version.window_size) 

405 self._tracker.on_hello(hello) 

406 

407 def _handle_device_state(self, payload: bytes) -> None: 

408 state = DeviceState.from_bytes(payload) 

409 self._tracker.on_device_state(state) 

410 if self._on_device_state is not None: 

411 self._on_device_state(state) 

412 

413 def _handle_rx_audio(self, payload: bytes) -> None: 

414 self._tracker.on_rx_audio(payload) 

415 

416 def _handle_window_update(self, payload: bytes) -> None: 

417 size = WindowUpdate.from_bytes(payload).size 

418 self._flow.add(size) 

419 logger.debug("window update size=%d", size) 

420 

421 def _on_transport_error(self, exc: Exception) -> None: 

422 """Called from the transport's background thread when it dies unexpectedly.""" 

423 logger.error("transport error: %s", exc) 

424 self._transport_error = exc 

425 

426 # -- outgoing frames ------------------------------------------------------- 

427 

428 def _send_desired_state(self, state: HostDesiredState) -> None: 

429 self._send_vendor(COMMAND_HOST_DESIRED_STATE, state.to_bytes(), flow_controlled=False) 

430 

431 def _send_vendor(self, command: int, payload: bytes = b"", *, flow_controlled: bool) -> None: 

432 self._claim_and_write(KISS_CMD_SETHARDWARE, encode_vendor_payload(command, payload), flow_controlled=flow_controlled) 

433 

434 def _claim_and_write(self, kiss_command: int, payload: bytes, *, flow_controlled: bool) -> None: 

435 if flow_controlled: 

436 frame_size = len(encode_kiss_frame(kiss_command, payload)) 

437 if not self._flow.claim(frame_size): 

438 logger.warning("drop KISS command=0x%02x frame; flow-control window exhausted", kiss_command) 

439 return 

440 self._transport.write_frame(kiss_command, payload) 

441 

442 def _require_ready(self) -> None: 

443 if self._transport_error is not None: 

444 raise RadioTransportError(f"transport failed: {self._transport_error}") from self._transport_error 

445 if not self.is_ready: 

446 raise RadioNotReadyError("radio has not completed the HELLO handshake yet")