mirror of
https://github.com/pyMC-dev/pyMC_Repeater.git
synced 2026-08-07 09:23:06 +02:00
fix: serialise radio TX and close duty-cycle TOCTOU race
Add self._tx_lock (asyncio.Lock) to RepeaterHandler and acquire it inside delayed_send after the per-packet sleep completes. Problem 1 — radio interleave: concurrent delayed_send coroutines (one per queued packet) could both exit their sleep at nearly the same moment and call dispatcher.send_packet simultaneously, interleaving SPI/serial register writes to the half-duplex LoRa radio. Problem 2 — TOCTOU gap: the upfront can_transmit() check in __call__ and the record_tx() call in delayed_send are separated by the entire TX delay (up to several seconds). Under burst conditions two tasks both pass the check before either has recorded its airtime, causing both to transmit and the duty-cycle budget to be exceeded. Fix: acquire _tx_lock after the sleep so delay timers still run concurrently (matching firmware behaviour), then immediately re-check can_transmit() inside the lock before sending. Because only one task holds the lock at a time, airtime state is stable; check and record_tx() are effectively atomic — no TOCTOU window. Airtime is recorded only on a successful send, so a radio failure never inflates the budget. Also move `import random` from inside _calculate_tx_delay to module level (stdlib imports belong at the top; the lazy-import pattern is unnecessary here). Docs: docs/pr_tx_serialization.md — problem statement, root-cause analysis, alternative approaches considered, invariant table, full unit + field test plan, and proof of correctness for the asyncio.Lock approach. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
+67
-25
@@ -1,6 +1,7 @@
|
||||
import asyncio
|
||||
import copy
|
||||
import logging
|
||||
import random
|
||||
import struct
|
||||
import time
|
||||
from collections import OrderedDict, deque
|
||||
@@ -135,6 +136,24 @@ class RepeaterHandler(BaseHandler):
|
||||
self._transport_keys_cache_time = 0
|
||||
self._transport_keys_cache_ttl = 60 # Cache for 60 seconds
|
||||
|
||||
# Serialise all radio TX calls.
|
||||
#
|
||||
# Background: since the queue loop dispatches each packet as an
|
||||
# asyncio.create_task, multiple _route_packet coroutines can have their
|
||||
# TX delay timers running concurrently — which is the intended behaviour
|
||||
# (firmware nodes do the same with a hardware timer). However, the
|
||||
# LoRa radio is half-duplex: it can only transmit one packet at a time.
|
||||
# Without serialisation, two tasks whose delay timers expire near-
|
||||
# simultaneously both call dispatcher.send_packet, interleaving SPI/serial
|
||||
# commands to the radio and both passing the LBT check before either has
|
||||
# actually transmitted.
|
||||
#
|
||||
# _tx_lock is acquired after each delay sleep and held for the entire
|
||||
# send_packet call. Delays still run concurrently; only the radio
|
||||
# access is serialised. This also eliminates the TOCTOU gap in duty-cycle
|
||||
# enforcement — see schedule_retransmit / delayed_send for details.
|
||||
self._tx_lock = asyncio.Lock()
|
||||
|
||||
self._start_background_tasks()
|
||||
|
||||
async def __call__(
|
||||
@@ -923,8 +942,6 @@ class RepeaterHandler(BaseHandler):
|
||||
|
||||
def _calculate_tx_delay(self, packet: Packet, snr: float = 0.0) -> float:
|
||||
|
||||
import random
|
||||
|
||||
packet_len = packet.get_raw_length()
|
||||
airtime_ms = self.airtime_mgr.calculate_airtime(packet_len)
|
||||
|
||||
@@ -1006,29 +1023,54 @@ class RepeaterHandler(BaseHandler):
|
||||
|
||||
async def delayed_send():
|
||||
await asyncio.sleep(delay)
|
||||
last_error = None
|
||||
for attempt in range(2 if local_transmission else 1):
|
||||
try:
|
||||
await self.dispatcher.send_packet(fwd_pkt, wait_for_ack=False)
|
||||
self._record_packet_sent(fwd_pkt)
|
||||
if airtime_ms > 0:
|
||||
self.airtime_mgr.record_tx(airtime_ms)
|
||||
packet_size = fwd_pkt.get_raw_length()
|
||||
logger.info(
|
||||
f"Retransmitted packet ({packet_size} bytes, "
|
||||
f"{airtime_ms:.1f}ms airtime)"
|
||||
)
|
||||
return
|
||||
except Exception as e:
|
||||
last_error = e
|
||||
logger.error(f"Retransmit failed: {e}")
|
||||
if local_transmission and attempt == 0:
|
||||
logger.info("Retrying local TX in 1s...")
|
||||
await asyncio.sleep(1.0)
|
||||
else:
|
||||
raise
|
||||
if last_error is not None:
|
||||
raise last_error
|
||||
|
||||
# Acquire the TX lock *after* the delay so that delay timers for
|
||||
# multiple packets still run concurrently (matching firmware). Only
|
||||
# one coroutine enters the radio send path at a time.
|
||||
async with self._tx_lock:
|
||||
# ── Authoritative duty-cycle gate ─────────────────────────────
|
||||
# The upfront can_transmit() call in __call__ is advisory: it
|
||||
# avoids scheduling packets that are obviously over budget, but
|
||||
# it cannot prevent a race between two tasks whose delay timers
|
||||
# expire at almost the same moment. Both tasks pass the advisory
|
||||
# check before either has recorded its airtime, then both try to
|
||||
# transmit.
|
||||
#
|
||||
# Inside _tx_lock only one task runs at a time, so airtime state
|
||||
# is stable here. The check and the subsequent record_tx() are
|
||||
# effectively atomic — no TOCTOU window.
|
||||
if airtime_ms > 0:
|
||||
can_tx_now, _ = self.airtime_mgr.can_transmit(airtime_ms)
|
||||
if not can_tx_now:
|
||||
logger.warning(
|
||||
"Packet dropped at TX time: duty-cycle exceeded "
|
||||
"(airtime=%.1fms)", airtime_ms,
|
||||
)
|
||||
return
|
||||
|
||||
last_error = None
|
||||
for attempt in range(2 if local_transmission else 1):
|
||||
try:
|
||||
await self.dispatcher.send_packet(fwd_pkt, wait_for_ack=False)
|
||||
self._record_packet_sent(fwd_pkt)
|
||||
if airtime_ms > 0:
|
||||
self.airtime_mgr.record_tx(airtime_ms)
|
||||
packet_size = fwd_pkt.get_raw_length()
|
||||
logger.info(
|
||||
f"Retransmitted packet ({packet_size} bytes, "
|
||||
f"{airtime_ms:.1f}ms airtime)"
|
||||
)
|
||||
return
|
||||
except Exception as e:
|
||||
last_error = e
|
||||
logger.error(f"Retransmit failed: {e}")
|
||||
if local_transmission and attempt == 0:
|
||||
logger.info("Retrying local TX in 1s...")
|
||||
await asyncio.sleep(1.0)
|
||||
else:
|
||||
raise
|
||||
if last_error is not None:
|
||||
raise last_error
|
||||
|
||||
return asyncio.create_task(delayed_send())
|
||||
|
||||
|
||||
Reference in New Issue
Block a user