diff --git a/RNS/Interfaces/BackboneInterface.py b/RNS/Interfaces/BackboneInterface.py index c41483f8..ce476540 100644 --- a/RNS/Interfaces/BackboneInterface.py +++ b/RNS/Interfaces/BackboneInterface.py @@ -60,6 +60,13 @@ class BackboneInterface(Interface): DP_IC_RCVBUF = 32768 DP_IC_IF_HEADROOM = 32 DP_IC_PENALTY = 1.5 + DP_EC_INTERVAL = 1.0 + DP_EC_MID_WM = 128*1024 + DP_EC_HIGH_WM = 4*1024*1024 + DP_EC_STALL_TICKS = 3 + DP_EC_MAX_ETA = 10.0 + DP_EC_RELEASE_ETA = 5.0 # Gate release hysterisis + DP_EC_DEAD_TIME = 12.0 epoll = None listener_filenos = {} @@ -67,7 +74,9 @@ class BackboneInterface(Interface): _job_active = False _ic_job_active = False + _ec_job_active = False _job_lock = threading.Lock() + _ec_job_lock = threading.Lock() _ic_job_lock = threading.Lock() _dp_ic_lock = threading.Lock() _dp_ic_snapshot = 0 @@ -253,6 +262,7 @@ class BackboneInterface(Interface): def start(): if not BackboneInterface._job_active: threading.Thread(target=BackboneInterface.__job, daemon=True).start() if not BackboneInterface._ic_job_active: threading.Thread(target=BackboneInterface.__dp_ic_job, daemon=True).start() + if not BackboneInterface._ec_job_active: threading.Thread(target=BackboneInterface.__dp_ec_job, daemon=True).start() @staticmethod def ensure_epoll(): @@ -398,6 +408,86 @@ class BackboneInterface(Interface): RNS.log(f"Available ingress budget {RNS.prettyspeed(avail*8)}, hard-allocated {round(alloc*100.0, 2)}% ({RNS.prettyspeed(avail*alloc*8)}), handled in {RNS.prettyshorttime(taken, compact=True, tight=True)}", RNS.LOG_DEBUG) RNS.log(f"Holding for {RNS.prettyshorttime(hold, compact=True, tight=True)}, throttling handled in {RNS.prettyshorttime(taken, compact=True, tight=True)} at depth {q_depth}", RNS.LOG_DEBUG) + @staticmethod + def _dp_ec_evaluate(interface, now): + tb = interface.transmit_buffer + drained = tb._tx_sent - interface._dp_ec_prev_sent + interface._dp_ec_prev_sent = tb._tx_sent + + sendable = tb.sendable + buffered = len(tb) + + if buffered == 0 or sendable == 0: + interface._dp_ec_zero_ticks = 0 + interface._dp_ec_last_drain = now + interface.tx_stalled = False + return False + + if now - interface._dp_ec_last_drain >= BackboneInterface.DP_EC_DEAD_TIME: + RNS.log(f"No egress control drain progress for {RNS.prettyshorttime(BackboneInterface.DP_EC_DEAD_TIME, compact=True)} on {interface}, tearing down", RNS.LOG_WARNING) + try: + if hasattr(interface, "socket") and interface.socket: + fileno = interface.socket.fileno() + BackboneInterface.deregister_fileno(fileno) + if fileno in BackboneInterface.spawned_interface_filenos: BackboneInterface.spawned_interface_filenos.pop(fileno) + try: interface.socket.close() + except Exception as e: RNS.log(f"Egress control could not close socket for {interface}: {e}", RNS.LOG_ERROR) + except Exception as e: RNS.log(f"Egress control cleanup error for {interface}: {e}", RNS.LOG_ERROR) + + interface.receive(b"") + return True + + if drained > 0: + interface._dp_ec_last_drain = now + interface._dp_ec_zero_ticks = 0 + drain_rate = drained / BackboneInterface.DP_EC_INTERVAL + clear_eta = buffered / drain_rate if drain_rate > 0 else float("inf") + + if buffered > BackboneInterface.DP_EC_MID_WM and clear_eta > BackboneInterface.DP_EC_MAX_ETA: + if not interface.tx_stalled: RNS.log(f"Egress control drain ETA of {RNS.prettyshorttime(clear_eta, compact=True)} exceeds maximum on {interface}, gating outbound", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + interface.tx_stalled = True + + elif clear_eta < BackboneInterface.DP_EC_RELEASE_ETA or buffered <= BackboneInterface.DP_EC_MID_WM: + if interface.tx_stalled: RNS.log(f"Egress control drain recovered on {interface}, resuming outbound", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + interface.tx_stalled = False + + else: + # No drain progress on this tick + if buffered > BackboneInterface.DP_EC_MID_WM: + interface._dp_ec_zero_ticks += 1 + if interface._dp_ec_zero_ticks >= BackboneInterface.DP_EC_STALL_TICKS: + if not interface.tx_stalled: RNS.log(f"No egress control drain progress on {interface}, gating outbound", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + interface.tx_stalled = True + else: + interface._dp_ec_zero_ticks = 0 + interface.tx_stalled = False + + return False + + @staticmethod + def __dp_ec_job(): + with BackboneInterface._ec_job_lock: + if BackboneInterface._ec_job_active: return + else: + BackboneInterface._ec_job_active = True + RNS.log(f"Started dataplane egress control", RNS.LOG_DEBUG) + try: + while True: + time.sleep(BackboneInterface.DP_EC_INTERVAL) + now = time.time() + try: interfaces = list(BackboneInterface.spawned_interface_filenos.values()) + except RuntimeError: + RNS.log(f"Deferring egress control evaluation due to error: {e}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + continue + + for interface in interfaces: + if interface.detached: continue + BackboneInterface._dp_ec_evaluate(interface, now) + + except Exception as e: + RNS.log(f"BackboneInterface egress control error: {e}", RNS.LOG_ERROR) + RNS.trace_exception(e) + @staticmethod def __dp_ic_job(): with BackboneInterface._ic_job_lock: @@ -416,7 +506,6 @@ class BackboneInterface(Interface): with BackboneInterface._dp_ic_lock: st = time.time() interfaces = list(BackboneInterface.spawned_interface_filenos.values()) - # active = [iface for iface in interfaces if iface.dp_ingress_bytes > 0] if q_depth > BackboneInterface.DP_IC_MID_WM: producers = [iface for iface in interfaces if iface.dp_ingress_packets > 0] @@ -941,8 +1030,12 @@ class BackboneClientInterface(Interface): def process_outgoing(self, data): if self.online and not self.detached: try: - self.transmit_buffer.append(bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG])) - BackboneInterface.tx_ready(self) + frame = bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]) + if self.tx_stalled or not self.transmit_buffer.append(frame, self.tx_hwm): + self.tx_drops += 1 + self.tx_dropped_bytes += len(frame) + RNS.log(f"Egress control dropping outbound frame of {RNS.prettysize(len(frame))} on {self}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + else: BackboneInterface.tx_ready(self) except Exception as e: RNS.log("Exception occurred while transmitting via "+str(self)+", tearing down interface", RNS.LOG_ERROR) diff --git a/RNS/Interfaces/Interface.py b/RNS/Interfaces/Interface.py index 7823cbe9..ca301f86 100755 --- a/RNS/Interfaces/Interface.py +++ b/RNS/Interfaces/Interface.py @@ -163,6 +163,14 @@ class Interface: self.dp_ingress_hold = None self.dp_ingress_gated = False + self.tx_hwm = 4*1024*1024 + self.tx_stalled = False + self.tx_drops = 0 + self.tx_dropped_bytes = 0 + self._dp_ec_prev_sent = 0 + self._dp_ec_zero_ticks = 0 + self._dp_ec_last_drain = time.time() + self.reports_phy_stats = False self.r_stat_rssi = None self.r_stat_snr = None diff --git a/RNS/Interfaces/LocalInterface.py b/RNS/Interfaces/LocalInterface.py index 330bd345..cfc5eb2b 100644 --- a/RNS/Interfaces/LocalInterface.py +++ b/RNS/Interfaces/LocalInterface.py @@ -218,8 +218,12 @@ class LocalClientInterface(Interface): if self.online: try: if self.epoll_backend: - self.transmit_buffer.append(bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG])) - BackboneInterface.tx_ready(self) + frame = bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]) + if self.tx_stalled or not self.transmit_buffer.append(frame, self.tx_hwm): + self.tx_drops += 1 + self.tx_dropped_bytes += len(frame) + RNS.log(f"Egress control dropping outbound frame of {RNS.prettysize(len(frame))} on {self}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + else: BackboneInterface.tx_ready(self) else: self.writing = True