From 65164fffb7a4bf7874cf954a589308c5027428d6 Mon Sep 17 00:00:00 2001 From: agessaman Date: Tue, 17 Feb 2026 21:32:11 -0800 Subject: [PATCH] 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. --- repeater/engine.py | 123 ++++++++++++++++++++++++++++---------- repeater/main.py | 29 +++++---- repeater/packet_router.py | 13 +++- 3 files changed, 118 insertions(+), 47 deletions(-) diff --git a/repeater/engine.py b/repeater/engine.py index 08cade9..8d891cd 100644 --- a/repeater/engine.py +++ b/repeater/engine.py @@ -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()) diff --git a/repeater/main.py b/repeater/main.py index 00538ae..800cb63 100644 --- a/repeater/main.py +++ b/repeater/main.py @@ -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 diff --git a/repeater/packet_router.py b/repeater/packet_router.py index 8a8e878..48ac218 100644 --- a/repeater/packet_router.py +++ b/repeater/packet_router.py @@ -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)