Improve retransmission logic and duty cycle handling in RepeaterHandler

- Improved local transmission handling by deferring local TX when duty cycle limits are exceeded, instead of dropping packets.
- Added LBT metadata extraction and logging for better monitoring of transmission attempts and delays.
- Refactored `schedule_retransmit` to support retrying local transmissions on failure, enhancing reliability.
- Introduced a lock in PacketRouter to serialize local TX operations, preventing race conditions during packet processing.
This commit is contained in:
agessaman
2026-02-17 21:32:11 -08:00
parent c2f8a2e3cd
commit 65164fffb7
3 changed files with 118 additions and 47 deletions
+92 -31
View File
@@ -164,28 +164,69 @@ class RepeaterHandler(BaseHandler):
can_tx, wait_time = self.airtime_mgr.can_transmit(airtime_ms)
# LBT metadata (set after any TX path that awaits send)
tx_metadata = None
lbt_attempts = 0
lbt_backoff_delays_ms = None
lbt_channel_busy = False
if not can_tx:
logger.warning(
f"Duty-cycle limit exceeded. Airtime={airtime_ms:.1f}ms, "
f"wait={wait_time:.1f}s before retry"
)
self.dropped_count += 1
drop_reason = "Duty cycle limit"
if local_transmission:
# Defer local TX until duty cycle allows instead of dropping
deferred_delay = delay + wait_time
logger.info(
f"Duty-cycle limit: deferring local TX by {wait_time:.1f}s "
f"(airtime={airtime_ms:.1f}ms)"
)
self.forwarded_count += 1
transmitted = True
tx_task = await self.schedule_retransmit(
fwd_pkt, deferred_delay, airtime_ms, local_transmission=True
)
try:
await tx_task
except Exception as e:
self.forwarded_count -= 1
transmitted = False
drop_reason = "TX failed (deferred)"
logger.warning(f"Deferred local TX failed: {e}")
raise
tx_metadata = getattr(fwd_pkt, "_tx_metadata", None)
if tx_metadata:
lbt_attempts = tx_metadata.get("lbt_attempts", 0)
lbt_backoff_delays_ms = tx_metadata.get(
"lbt_backoff_delays_ms", []
)
lbt_channel_busy = tx_metadata.get("lbt_channel_busy", False)
if lbt_attempts > 0:
total_lbt_delay = sum(lbt_backoff_delays_ms)
logger.info(
f"LBT: {lbt_attempts} attempts, "
f"{total_lbt_delay:.0f}ms delay, "
f"backoffs={lbt_backoff_delays_ms}"
)
else:
logger.warning(
f"Duty-cycle limit exceeded. Airtime={airtime_ms:.1f}ms, "
f"wait={wait_time:.1f}s before retry"
)
self.dropped_count += 1
drop_reason = "Duty cycle limit"
else:
self.forwarded_count += 1
transmitted = True
# Schedule retransmit with delay (returns task)
tx_task = await self.schedule_retransmit(fwd_pkt, delay, airtime_ms)
# Wait for transmission to complete to get LBT metadata
await tx_task
# Extract LBT metadata after transmission
tx_task = await self.schedule_retransmit(
fwd_pkt, delay, airtime_ms, local_transmission=local_transmission
)
try:
await tx_task
except Exception as e:
self.forwarded_count -= 1
transmitted = False
drop_reason = "TX failed"
logger.warning(f"Local TX failed: {e}")
raise
tx_metadata = getattr(fwd_pkt, "_tx_metadata", None)
lbt_attempts = 0
lbt_backoff_delays_ms = None
lbt_channel_busy = False
if tx_metadata:
lbt_attempts = tx_metadata.get("lbt_attempts", 0)
lbt_backoff_delays_ms = tx_metadata.get("lbt_backoff_delays_ms", [])
@@ -672,23 +713,43 @@ class RepeaterHandler(BaseHandler):
packet.drop_reason = f"Unknown route type: {route_type}"
return None
async def schedule_retransmit(self, fwd_pkt: Packet, delay: float, airtime_ms: float = 0.0):
"""Schedule a packet retransmission with delay and return the task."""
async def schedule_retransmit(
self,
fwd_pkt: Packet,
delay: float,
airtime_ms: float = 0.0,
local_transmission: bool = False,
):
"""Schedule a packet retransmission with delay and return the task.
If local_transmission is True and the first send fails, retry once after
a short delay (handles transient radio/LBT failures).
"""
async def delayed_send():
await asyncio.sleep(delay)
try:
await self.dispatcher.send_packet(fwd_pkt, wait_for_ack=False)
# Record airtime after successful TX
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, {airtime_ms:.1f}ms airtime)"
)
except Exception as e:
logger.error(f"Retransmit failed: {e}")
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)
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())
+16 -13
View File
@@ -429,21 +429,24 @@ class RepeaterDaemon:
records.append(d)
bridge.contacts.load_from_dicts(records)
# Load channels from SQLite
# Load channels from SQLite (normalize secret to 32 bytes to match
# CompanionBase.set_channel and GroupTextHandler/PacketBuilder)
channel_rows = sqlite_handler.companion_load_channels(companion_hash_str)
for row in channel_rows:
ch = Channel(
name=row.get("name", ""),
secret=(
row.get("secret", b"")
if isinstance(row.get("secret"), bytes)
else (
bytes.fromhex(row.get("secret", ""))
if row.get("secret")
else b""
)
),
)
s = row.get("secret", b"")
if isinstance(s, bytes):
raw = s
elif isinstance(s, (bytearray, memoryview)):
raw = bytes(s)
elif s:
raw = bytes.fromhex(s if isinstance(s, str) else str(s))
else:
raw = b""
if len(raw) < 32:
raw = raw + b"\x00" * (32 - len(raw))
elif len(raw) > 32:
raw = raw[:32]
ch = Channel(name=row.get("name", ""), secret=raw)
bridge.channels.set(row.get("channel_idx", 0), ch)
# Preload queued messages from SQLite into bridge
+10 -3
View File
@@ -23,7 +23,9 @@ class PacketRouter:
self.queue = asyncio.Queue()
self.running = False
self.router_task = None
# Serialize injects so one local TX completes before the next is processed
self._inject_lock = asyncio.Lock()
async def start(self):
self.running = True
self.router_task = asyncio.create_task(self._process_queue())
@@ -51,8 +53,13 @@ class PacketRouter:
"timestamp": getattr(packet, "timestamp", 0),
}
# Use local_transmission=True to bypass forwarding logic
await self.daemon.repeater_handler(packet, metadata, local_transmission=True)
# Serialize injects so one local TX completes before the next runs
# (avoids duty-cycle or dispatcher races where a later packet goes out first)
async with self._inject_lock:
# Use local_transmission=True to bypass forwarding logic
await self.daemon.repeater_handler(
packet, metadata, local_transmission=True
)
# Enqueue so router can deliver to companion(s): TXT_MSG -> dest bridge, ACK -> all bridges (sender sees ACK)
await self.enqueue(packet)