Wired coalescing transmit buffer and optimized HDLC deframer into Backbone and Local interfaces
parent
0ec84c7491
commit
89e73db94a
|
|
@ -29,6 +29,8 @@
|
|||
# SOFTWARE.
|
||||
|
||||
from RNS.Interfaces.Interface import Interface
|
||||
from RNS.Interfaces.util.TransmitBuffer import TransmitBuffer
|
||||
from RNS.Interfaces.util.HDLC import HDLC, ReceiveBuffer
|
||||
import threading
|
||||
import socket
|
||||
import select
|
||||
|
|
@ -37,22 +39,11 @@ import sys
|
|||
import os
|
||||
import RNS
|
||||
|
||||
class HDLC():
|
||||
FLAG = 0x7E
|
||||
ESC = 0x7D
|
||||
ESC_MASK = 0x20
|
||||
|
||||
@staticmethod
|
||||
def escape(data):
|
||||
data = data.replace(bytes([HDLC.ESC]), bytes([HDLC.ESC, HDLC.ESC^HDLC.ESC_MASK]))
|
||||
data = data.replace(bytes([HDLC.FLAG]), bytes([HDLC.ESC, HDLC.FLAG^HDLC.ESC_MASK]))
|
||||
return data
|
||||
|
||||
class BackboneInterface(Interface):
|
||||
HW_MTU = 1048576
|
||||
BITRATE_GUESS = 100_000_000
|
||||
DEFAULT_IFAC_SIZE = 16
|
||||
AUTOCONFIGURE_MTU = True
|
||||
HW_MTU = 1048576
|
||||
BITRATE_GUESS = 100_000_000
|
||||
DEFAULT_IFAC_SIZE = 16
|
||||
AUTOCONFIGURE_MTU = True
|
||||
|
||||
BLOCK_FAST_FLAPPING = True
|
||||
FAST_FLAP_THRESHOLD = 20
|
||||
|
|
@ -352,7 +343,7 @@ class BackboneInterface(Interface):
|
|||
interface.dp_ingress_gated = True
|
||||
interface.dp_ingress_tcount += 1
|
||||
events = select.EPOLLHUP
|
||||
if len(interface.transmit_buffer) > 0: events |= select.EPOLLOUT
|
||||
if interface.transmit_buffer.sendable > 0: events |= select.EPOLLOUT
|
||||
BackboneInterface.epoll.modify(fileno, events)
|
||||
if hold: interface.dp_ingress_hold = time.time() + hold
|
||||
RNS.log(f"Ingress throttled on {interface}", RNS.LOG_NOTICE) if RNS.sl(RNS.LOG_NOTICE) else None
|
||||
|
|
@ -368,7 +359,7 @@ class BackboneInterface(Interface):
|
|||
fileno = interface.socket.fileno()
|
||||
if fileno in BackboneInterface.spawned_interface_filenos:
|
||||
events = select.EPOLLIN | select.EPOLLHUP
|
||||
if len(interface.transmit_buffer) > 0: events |= select.EPOLLOUT
|
||||
if interface.transmit_buffer.sendable > 0: events |= select.EPOLLOUT
|
||||
BackboneInterface.epoll.modify(fileno, events)
|
||||
interface.dp_ingress_gated = False
|
||||
RNS.log(f"Released ingress throttle on {interface}", RNS.LOG_NOTICE) if RNS.sl(RNS.LOG_NOTICE) else None
|
||||
|
|
@ -511,7 +502,7 @@ class BackboneInterface(Interface):
|
|||
|
||||
socket_valid_after_read = socket_valid and fileno in BackboneInterface.spawned_interface_filenos
|
||||
if socket_valid_after_read and (event & select.EPOLLOUT):
|
||||
try: written = client_socket.send(spawned_interface.transmit_buffer)
|
||||
try: written = spawned_interface.transmit_buffer.drain_to(client_socket)
|
||||
except Exception as e:
|
||||
written = 0
|
||||
if not spawned_interface.detached:
|
||||
|
|
@ -538,9 +529,8 @@ class BackboneInterface(Interface):
|
|||
except Exception as e: RNS.log(f"Error while closing socket for {spawned_interface}: {e}", RNS.LOG_WARNING)
|
||||
spawned_interface.receive(b"")
|
||||
|
||||
spawned_interface.transmit_buffer = spawned_interface.transmit_buffer[written:]
|
||||
try:
|
||||
if len(spawned_interface.transmit_buffer) == 0:
|
||||
if spawned_interface.transmit_buffer.sendable == 0:
|
||||
events = select.EPOLLHUP
|
||||
if not spawned_interface.dp_ingress_gated: events |= select.EPOLLIN
|
||||
BackboneInterface.epoll.modify(fileno, events)
|
||||
|
|
@ -792,8 +782,10 @@ class BackboneClientInterface(Interface):
|
|||
self.i2p_tunneled = i2p_tunneled
|
||||
self.mode = RNS.Interfaces.Interface.Interface.MODE_FULL
|
||||
self.bitrate = BackboneClientInterface.BITRATE_GUESS
|
||||
self.frame_buffer = b""
|
||||
self.transmit_buffer = b""
|
||||
self.transmit_buffer = TransmitBuffer()
|
||||
self.receive_buffer = ReceiveBuffer(mtu=lambda: self.HW_MTU, min_frame_len=RNS.Reticulum.HEADER_MINSIZE,
|
||||
max_frame_len=lambda: self.HW_MTU+(getattr(self, "ifac_size", None) or 0),
|
||||
on_frame=self.process_incoming, on_invalid=self.invalid_frame)
|
||||
|
||||
if max_reconnect_tries == None:
|
||||
self.max_reconnect_tries = BackboneClientInterface.RECONNECT_MAX_TRIES
|
||||
|
|
@ -949,7 +941,7 @@ class BackboneClientInterface(Interface):
|
|||
def process_outgoing(self, data):
|
||||
if self.online and not self.detached:
|
||||
try:
|
||||
self.transmit_buffer += bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG])
|
||||
self.transmit_buffer.append(bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]))
|
||||
BackboneInterface.tx_ready(self)
|
||||
|
||||
except Exception as e:
|
||||
|
|
@ -957,41 +949,12 @@ class BackboneClientInterface(Interface):
|
|||
RNS.log("The contained exception was: "+str(e), RNS.LOG_ERROR)
|
||||
self.teardown()
|
||||
|
||||
def check_frame_len(self, frame_len):
|
||||
if frame_len <= RNS.Reticulum.HEADER_MINSIZE: return False
|
||||
elif frame_len > self.HW_MTU + (self.ifac_size or 0): return False
|
||||
else: return True
|
||||
|
||||
def invalid_frame(self, frame_len):
|
||||
RNS.log(f"Invalid HDLC frame of {RNS.prettysize(frame_len)} received on {self}, dropping frame", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
|
||||
|
||||
def receive(self, data_in):
|
||||
try:
|
||||
if len(data_in) > 0:
|
||||
self.frame_buffer += data_in
|
||||
flags_remaining = True
|
||||
while flags_remaining:
|
||||
frame_start = self.frame_buffer.find(HDLC.FLAG)
|
||||
if frame_start != -1:
|
||||
frame_end = self.frame_buffer.find(HDLC.FLAG, frame_start+1)
|
||||
if frame_end != -1:
|
||||
frame = self.frame_buffer[frame_start+1:frame_end]
|
||||
frame = frame.replace(bytes([HDLC.ESC, HDLC.FLAG ^ HDLC.ESC_MASK]), bytes([HDLC.FLAG]))
|
||||
frame = frame.replace(bytes([HDLC.ESC, HDLC.ESC ^ HDLC.ESC_MASK]), bytes([HDLC.ESC]))
|
||||
frame_len = len(frame)
|
||||
if frame_len != 0:
|
||||
if self.check_frame_len(frame_len): self.process_incoming(frame)
|
||||
else: self.invalid_frame(len(frame))
|
||||
|
||||
self.frame_buffer = self.frame_buffer[frame_end:]
|
||||
|
||||
else:
|
||||
if len(self.frame_buffer) > self.HW_MTU*2: self.frame_buffer = b""
|
||||
flags_remaining = False
|
||||
else:
|
||||
self.frame_buffer = b""
|
||||
flags_remaining = False
|
||||
|
||||
if len(data_in) > 0: self.receive_buffer.feed(data_in)
|
||||
else:
|
||||
self.online = False
|
||||
if self.initiator and not self.detached:
|
||||
|
|
@ -1010,8 +973,7 @@ class BackboneClientInterface(Interface):
|
|||
RNS.log("Attempting to reconnect...", RNS.LOG_WARNING)
|
||||
def job(): self.reconnect()
|
||||
threading.Thread(target=job, daemon=True).start()
|
||||
else:
|
||||
self.teardown()
|
||||
else: self.teardown()
|
||||
|
||||
def teardown(self):
|
||||
if self.initiator and not self.detached:
|
||||
|
|
|
|||
|
|
@ -30,6 +30,8 @@
|
|||
|
||||
from RNS.Interfaces.Interface import Interface
|
||||
from RNS.Interfaces.BackboneInterface import BackboneInterface
|
||||
from RNS.Interfaces.util.TransmitBuffer import TransmitBuffer
|
||||
from RNS.Interfaces.util.HDLC import HDLC, ReceiveBuffer
|
||||
import socketserver
|
||||
import threading
|
||||
import socket
|
||||
|
|
@ -39,17 +41,6 @@ import os
|
|||
import RNS
|
||||
from threading import Lock
|
||||
|
||||
class HDLC():
|
||||
FLAG = 0x7E
|
||||
ESC = 0x7D
|
||||
ESC_MASK = 0x20
|
||||
|
||||
@staticmethod
|
||||
def escape(data):
|
||||
data = data.replace(bytes([HDLC.ESC]), bytes([HDLC.ESC, HDLC.ESC^HDLC.ESC_MASK]))
|
||||
data = data.replace(bytes([HDLC.FLAG]), bytes([HDLC.ESC, HDLC.FLAG^HDLC.ESC_MASK]))
|
||||
return data
|
||||
|
||||
class ThreadingTCPServer(socketserver.ThreadingMixIn, socketserver.TCPServer):
|
||||
def server_bind(self):
|
||||
if RNS.vendor.platformutils.is_windows():
|
||||
|
|
@ -84,8 +75,9 @@ class LocalClientInterface(Interface):
|
|||
self.detached = False
|
||||
self.name = name
|
||||
self.mode = RNS.Interfaces.Interface.Interface.MODE_FULL
|
||||
self.frame_buffer = b""
|
||||
self.transmit_buffer = b""
|
||||
self.transmit_buffer = TransmitBuffer()
|
||||
self.receive_buffer = ReceiveBuffer(mtu=lambda: self.HW_MTU, min_frame_len=RNS.Reticulum.HEADER_MINSIZE,
|
||||
max_frame_len=lambda: self.HW_MTU, on_frame=self.process_incoming)
|
||||
|
||||
if RNS.vendor.platformutils.use_epoll(): self.epoll_backend = True
|
||||
|
||||
|
|
@ -198,7 +190,7 @@ class LocalClientInterface(Interface):
|
|||
RNS.log(f"Sending keepalive on {self}", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None
|
||||
try:
|
||||
if self.epoll_backend:
|
||||
self.transmit_buffer += bytes([HDLC.FLAG])+bytes([HDLC.FLAG])
|
||||
self.transmit_buffer.append(bytes([HDLC.FLAG])+bytes([HDLC.FLAG]))
|
||||
BackboneInterface.tx_ready(self)
|
||||
|
||||
else:
|
||||
|
|
@ -226,16 +218,14 @@ class LocalClientInterface(Interface):
|
|||
if self.online:
|
||||
try:
|
||||
if self.epoll_backend:
|
||||
self.transmit_buffer += bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG])
|
||||
self.transmit_buffer.append(bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]))
|
||||
BackboneInterface.tx_ready(self)
|
||||
|
||||
else:
|
||||
self.writing = True
|
||||
|
||||
if self._force_bitrate:
|
||||
if not hasattr(self, "send_lock"):
|
||||
self.send_lock = Lock()
|
||||
|
||||
if not hasattr(self, "send_lock"): self.send_lock = Lock()
|
||||
with self.send_lock:
|
||||
# RNS.log(f"Simulating latency of {RNS.prettytime(s)} for {len(data)} bytes", RNS.LOG_EXTREME)
|
||||
s = len(data) / self.bitrate * 8
|
||||
|
|
@ -255,22 +245,7 @@ class LocalClientInterface(Interface):
|
|||
self.teardown()
|
||||
|
||||
def handle_hdlc(self, data_in):
|
||||
self.frame_buffer += data_in
|
||||
flags_remaining = True
|
||||
while flags_remaining:
|
||||
frame_start = self.frame_buffer.find(HDLC.FLAG)
|
||||
if frame_start != -1:
|
||||
frame_end = self.frame_buffer.find(HDLC.FLAG, frame_start+1)
|
||||
if frame_end != -1:
|
||||
frame = self.frame_buffer[frame_start+1:frame_end]
|
||||
frame = frame.replace(bytes([HDLC.ESC, HDLC.FLAG ^ HDLC.ESC_MASK]), bytes([HDLC.FLAG]))
|
||||
frame = frame.replace(bytes([HDLC.ESC, HDLC.ESC ^ HDLC.ESC_MASK]), bytes([HDLC.ESC]))
|
||||
if len(frame) > RNS.Reticulum.HEADER_MINSIZE: self.process_incoming(frame)
|
||||
self.frame_buffer = self.frame_buffer[frame_end:]
|
||||
|
||||
else: flags_remaining = False
|
||||
|
||||
else: flags_remaining = False
|
||||
self.receive_buffer.feed(data_in)
|
||||
|
||||
def receive(self, data_in):
|
||||
try:
|
||||
|
|
@ -284,8 +259,7 @@ class LocalClientInterface(Interface):
|
|||
# there's no other connectivity left to block anyway, it might be
|
||||
# unnecessary.
|
||||
self.reconnect()
|
||||
else:
|
||||
self.teardown(nowarning=True)
|
||||
else: self.teardown(nowarning=True)
|
||||
|
||||
except Exception as e:
|
||||
self.online = False
|
||||
|
|
@ -297,7 +271,7 @@ class LocalClientInterface(Interface):
|
|||
|
||||
def read_loop(self):
|
||||
try:
|
||||
self.frame_buffer = b""
|
||||
self.receive_buffer.reset()
|
||||
data_in = b""
|
||||
while True:
|
||||
data_in = self.socket.recv(4096)
|
||||
|
|
@ -513,5 +487,4 @@ class LocalInterfaceHandler(socketserver.BaseRequestHandler):
|
|||
self.callback = callback
|
||||
socketserver.BaseRequestHandler.__init__(self, *args, **keys)
|
||||
|
||||
def handle(self):
|
||||
self.callback(handler=self)
|
||||
def handle(self): self.callback(handler=self)
|
||||
Loading…
Reference in New Issue