From 00d57a22eecb7adb775b0f11d48fde21ddc2dc21 Mon Sep 17 00:00:00 2001 From: Mark Qvist Date: Fri, 24 Jul 2026 02:30:34 +0200 Subject: [PATCH] Cleaned up I2P interface and control process code, fixed various i2plib bugs --- RNS/Interfaces/I2PInterface.py | 315 +++++++++----------------------- RNS/vendor/i2plib/aiosam.py | 65 +++---- RNS/vendor/i2plib/exceptions.py | 21 +-- RNS/vendor/i2plib/sam.py | 7 +- 4 files changed, 125 insertions(+), 283 deletions(-) diff --git a/RNS/Interfaces/I2PInterface.py b/RNS/Interfaces/I2PInterface.py index 18a5bf2d..381e50c5 100644 --- a/RNS/Interfaces/I2PInterface.py +++ b/RNS/Interfaces/I2PInterface.py @@ -88,8 +88,7 @@ class I2PController: self.ready = False self.storagepath = rns_storagepath+"/i2p" - if not os.path.isdir(self.storagepath): - os.makedirs(self.storagepath) + if not os.path.isdir(self.storagepath): os.makedirs(self.storagepath) def start(self): @@ -110,28 +109,22 @@ class I2PController: except Exception as e: self.ready = False RNS.log("Exception on event loop for "+str(self)+": "+str(e), RNS.LOG_ERROR) - finally: - self.loop.close() + finally: self.loop.close() def stop(self): for i2ptunnel in self.i2plib_tunnels: - if hasattr(i2ptunnel, "stop") and callable(i2ptunnel.stop): - i2ptunnel.stop() + if hasattr(i2ptunnel, "stop") and callable(i2ptunnel.stop): i2ptunnel.stop() if hasattr(asyncio.Task, "all_tasks") and callable(asyncio.Task.all_tasks): - for task in asyncio.Task.all_tasks(loop=self.loop): - task.cancel() + for task in asyncio.Task.all_tasks(loop=self.loop): task.cancel() time.sleep(0.2) - self.loop.stop() - def get_free_port(self): return self.i2plib.utils.get_free_port() - def stop_tunnel(self, i2ptunnel): if hasattr(i2ptunnel, "stop") and callable(i2ptunnel.stop): i2ptunnel.stop() @@ -159,10 +152,9 @@ class I2PController: tn = self.i2plib_tunnels[i2p_destination] if tn != None and hasattr(tn, "status"): - RNS.log("Waiting for status from I2P control process", RNS.LOG_EXTREME) - while not tn.status["setup_ran"]: - time.sleep(0.1) - RNS.log("Got status from I2P control process", RNS.LOG_EXTREME) + RNS.log("Waiting for status from I2P control process", RNS.LOG_DEBUG) + while not tn.status["setup_ran"]: time.sleep(0.1) + RNS.log("Got status from I2P control process", RNS.LOG_DEBUG) if tn.status["setup_failed"]: self.stop_tunnel(tn) @@ -172,29 +164,19 @@ class I2PController: if owner.socket != None: if hasattr(owner.socket, "close"): if callable(owner.socket.close): - try: - owner.socket.shutdown(socket.SHUT_RDWR) - except Exception as e: - RNS.log("Error while shutting down socket for "+str(owner)+": "+str(e)) + try: owner.socket.shutdown(socket.SHUT_RDWR) + except Exception as e: RNS.log("Error while shutting down socket for "+str(owner)+": "+str(e)) + try: owner.socket.close() + except Exception as e: RNS.log("Error while closing socket for "+str(owner)+": "+str(e)) - try: - owner.socket.close() - except Exception as e: - RNS.log("Error while closing socket for "+str(owner)+": "+str(e)) self.client_tunnels[i2p_destination] = True owner.awaiting_i2p_tunnel = False - RNS.log(str(owner)+" tunnel setup complete", RNS.LOG_VERBOSE) - else: - raise IOError("Got no status response from SAM API") - - except ConnectionRefusedError as e: - raise e - - except ConnectionAbortedError as e: - raise e + else: raise IOError("Got no status response from SAM API") + except ConnectionRefusedError as e: raise e + except ConnectionAbortedError as e: raise e except Exception as e: RNS.log("Unexpected error type from I2P SAM: "+str(e), RNS.LOG_ERROR) raise e @@ -206,54 +188,33 @@ class I2PController: if i2ptunnel.status["setup_ran"] == False: RNS.log(str(self)+" I2P tunnel setup did not complete", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False elif i2p_exception != None: RNS.log("An error ocurred while setting up I2P tunnel to "+str(i2p_destination), RNS.LOG_ERROR) - - if isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.CantReachPeer): - RNS.log("The I2P daemon can't reach peer "+str(i2p_destination), RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedDest): - RNS.log("The I2P daemon reported that the destination is already in use", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedId): - RNS.log("The I2P daemon reported that the ID is arleady in use", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidId): - RNS.log("The I2P daemon reported that the stream session ID doesn't exist", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidKey): - RNS.log("The I2P daemon reported that the key for "+str(i2p_destination)+" is invalid", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.KeyNotFound): - RNS.log("The I2P daemon could not find the key for "+str(i2p_destination), RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.PeerNotFound): - RNS.log("The I2P daemon mould not find the peer "+str(i2p_destination), RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.I2PError): - RNS.log("The I2P daemon experienced an unspecified error", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.Timeout): - RNS.log("I2P daemon timed out while setting up client tunnel to "+str(i2p_destination), RNS.LOG_ERROR) + if isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.CantReachPeer): RNS.log("The I2P daemon can't reach peer "+str(i2p_destination), RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedDest): RNS.log("The I2P daemon reported that the destination is already in use", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedId): RNS.log("The I2P daemon reported that the ID is arleady in use", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidId): RNS.log("The I2P daemon reported that the stream session ID doesn't exist", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidKey): RNS.log("The I2P daemon reported that the key for "+str(i2p_destination)+" is invalid", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.KeyNotFound): RNS.log("The I2P daemon could not find the key for "+str(i2p_destination), RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.PeerNotFound): RNS.log("The I2P daemon could not find the peer "+str(i2p_destination), RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.I2PError): RNS.log("The I2P daemon experienced an unspecified error", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.Timeout): RNS.log("I2P daemon timed out while setting up client tunnel to "+str(i2p_destination), RNS.LOG_ERROR) + else: RNS.log(f"Unspecified I2P daemon error: {i2p_exception}", RNS.LOG_ERROR) RNS.log("Resetting I2P tunnel and retrying later", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False elif i2ptunnel.status["setup_failed"] == True: RNS.log(str(self)+" Unspecified I2P tunnel setup error, resetting I2P tunnel", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False else: RNS.log(str(self)+" Got no status from SAM API, resetting I2P tunnel", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False @@ -274,10 +235,8 @@ class I2PController: i2p_keyfile_nf = self.storagepath+"/"+RNS.hexrep(i2p_dest_hash_nf, delimit=False)+".i2p" # Use old format if a key is already present - if os.path.isfile(i2p_keyfile_of): - i2p_keyfile = i2p_keyfile_of - else: - i2p_keyfile = i2p_keyfile_nf + if os.path.isfile(i2p_keyfile_of): i2p_keyfile = i2p_keyfile_of + else: i2p_keyfile = i2p_keyfile_nf i2p_dest = None if not os.path.isfile(i2p_keyfile): @@ -312,8 +271,7 @@ class I2PController: asyncio.run_coroutine_threadsafe(tunnel_up(), self.loop).result() self.server_tunnels[i2p_b32] = True - except Exception as e: - raise e + except Exception as e: raise e else: i2ptunnel = self.i2plib_tunnels[i2p_b32] @@ -322,66 +280,45 @@ class I2PController: if i2ptunnel.status["setup_ran"] == False: RNS.log(str(self)+" I2P tunnel setup did not complete", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False elif i2p_exception != None: RNS.log("An error ocurred while setting up I2P tunnel", RNS.LOG_ERROR) - - if isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.CantReachPeer): - RNS.log("The I2P daemon can't reach peer "+str(i2p_destination), RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedDest): - RNS.log("The I2P daemon reported that the destination is already in use", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedId): - RNS.log("The I2P daemon reported that the ID is arleady in use", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidId): - RNS.log("The I2P daemon reported that the stream session ID doesn't exist", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidKey): - RNS.log("The I2P daemon reported that the key for "+str(i2p_destination)+" is invalid", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.KeyNotFound): - RNS.log("The I2P daemon could not find the key for "+str(i2p_destination), RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.PeerNotFound): - RNS.log("The I2P daemon mould not find the peer "+str(i2p_destination), RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.I2PError): - RNS.log("The I2P daemon experienced an unspecified error", RNS.LOG_ERROR) - - elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.Timeout): - RNS.log("I2P daemon timed out while setting up client tunnel to "+str(i2p_destination), RNS.LOG_ERROR) + if isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.CantReachPeer): RNS.log("The I2P daemon can't reach peer "+str(i2p_destination), RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedDest): RNS.log("The I2P daemon reported that the destination is already in use", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.DuplicatedId): RNS.log("The I2P daemon reported that the ID is arleady in use", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidId): RNS.log("The I2P daemon reported that the stream session ID doesn't exist", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.InvalidKey): RNS.log("The I2P daemon reported that the key for "+str(i2p_destination)+" is invalid", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.KeyNotFound): RNS.log("The I2P daemon could not find the key for "+str(i2p_destination), RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.PeerNotFound): RNS.log("The I2P daemon could not find the peer "+str(i2p_destination), RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.I2PError): RNS.log("The I2P daemon experienced an unspecified error", RNS.LOG_ERROR) + elif isinstance(i2p_exception, RNS.vendor.i2plib.exceptions.Timeout): RNS.log("I2P daemon timed out while setting up client tunnel to "+str(i2p_destination), RNS.LOG_ERROR) + else: RNS.log(f"Unspecified I2P daemon error: {i2p_exception}", RNS.LOG_ERROR) RNS.log("Resetting I2P tunnel and retrying later", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False elif i2ptunnel.status["setup_failed"] == True: RNS.log(str(self)+" Unspecified I2P tunnel setup error, resetting I2P tunnel", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False else: RNS.log(str(self)+" Got no status from SAM API, resetting I2P tunnel", RNS.LOG_ERROR) - self.stop_tunnel(i2ptunnel) return False time.sleep(5) - def get_loop(self): - return asyncio.get_event_loop() + def get_loop(self): return asyncio.get_event_loop() class ThreadingI2PServer(socketserver.ThreadingMixIn, socketserver.TCPServer): pass + class I2PInterfacePeer(Interface): RECONNECT_WAIT = 15 RECONNECT_MAX_TRIES = None @@ -431,18 +368,11 @@ class I2PInterfacePeer(Interface): self.ifac_netkey = self.parent_interface.ifac_netkey if self.ifac_netname != None or self.ifac_netkey != None: ifac_origin = b"" - if self.ifac_netname != None: - ifac_origin += RNS.Identity.full_hash(self.ifac_netname.encode("utf-8")) - if self.ifac_netkey != None: - ifac_origin += RNS.Identity.full_hash(self.ifac_netkey.encode("utf-8")) + if self.ifac_netname != None: ifac_origin += RNS.Identity.full_hash(self.ifac_netname.encode("utf-8")) + if self.ifac_netkey != None: ifac_origin += RNS.Identity.full_hash(self.ifac_netkey.encode("utf-8")) ifac_origin_hash = RNS.Identity.full_hash(ifac_origin) - self.ifac_key = RNS.Cryptography.hkdf( - length=64, - derive_from=ifac_origin_hash, - salt=RNS.Reticulum.IFAC_SALT, - context=None - ) + self.ifac_key = RNS.Cryptography.hkdf(length=64, derive_from=ifac_origin_hash, salt=RNS.Reticulum.IFAC_SALT, context=None) self.ifac_identity = RNS.Identity.from_bytes(self.ifac_key) self.ifac_signature = self.ifac_identity.sign(RNS.Identity.full_hash(self.ifac_key)) @@ -450,10 +380,8 @@ class I2PInterfacePeer(Interface): self.announce_rate_grace = None self.announce_rate_penalty = None - if max_reconnect_tries == None: - self.max_reconnect_tries = I2PInterfacePeer.RECONNECT_MAX_TRIES - else: - self.max_reconnect_tries = max_reconnect_tries + if max_reconnect_tries == None: self.max_reconnect_tries = I2PInterfacePeer.RECONNECT_MAX_TRIES + else: self.max_reconnect_tries = max_reconnect_tries if connected_socket != None: self.receives = True @@ -461,15 +389,12 @@ class I2PInterfacePeer(Interface): self.target_port = None self.socket = connected_socket - if platform.system() == "Linux": - self.set_timeouts_linux() - elif platform.system() == "Darwin": - self.set_timeouts_osx() + if platform.system() == "Linux": self.set_timeouts_linux() + elif platform.system() == "Darwin": self.set_timeouts_osx() elif target_i2p_dest != None: self.receives = True self.initiator = True - self.bind_ip = "127.0.0.1" self.awaiting_i2p_tunnel = True @@ -497,12 +422,10 @@ class I2PInterfacePeer(Interface): thread.start() def wait_job(): - while self.awaiting_i2p_tunnel: - time.sleep(0.25) + while self.awaiting_i2p_tunnel: time.sleep(0.25) time.sleep(2) - if not self.kiss_framing: - self.wants_tunnel = True + if not self.kiss_framing: self.wants_tunnel = True if not self.connect(initial=True): thread = threading.Thread(target=self.reconnect) @@ -517,7 +440,6 @@ class I2PInterfacePeer(Interface): thread.daemon = True thread.start() - def set_timeouts_linux(self): self.socket.setsockopt(socket.IPPROTO_TCP, socket.TCP_USER_TIMEOUT, int(I2PInterfacePeer.I2P_USER_TIMEOUT * 1000)) self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1) @@ -526,10 +448,8 @@ class I2PInterfacePeer(Interface): self.socket.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, int(I2PInterfacePeer.I2P_PROBES)) def set_timeouts_osx(self): - if hasattr(socket, "TCP_KEEPALIVE"): - TCP_KEEPIDLE = socket.TCP_KEEPALIVE - else: - TCP_KEEPIDLE = 0x10 + if hasattr(socket, "TCP_KEEPALIVE"): TCP_KEEPIDLE = socket.TCP_KEEPALIVE + else: TCP_KEEPIDLE = 0x10 self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1) self.socket.setsockopt(socket.IPPROTO_TCP, TCP_KEEPIDLE, int(I2PInterfacePeer.I2P_PROBE_AFTER)) @@ -537,16 +457,12 @@ class I2PInterfacePeer(Interface): def shutdown_socket(self, target_socket): if callable(target_socket.close): try: - if socket != None: - target_socket.shutdown(socket.SHUT_RDWR) - except Exception as e: - RNS.log("Error while shutting down socket for "+str(self)+": "+str(e)) + if socket != None: target_socket.shutdown(socket.SHUT_RDWR) + except Exception as e: RNS.log("Error while shutting down socket for "+str(self)+": "+str(e)) try: - if socket != None: - target_socket.close() - except Exception as e: - RNS.log("Error while closing socket for "+str(self)+": "+str(e)) + if socket != None: target_socket.close() + except Exception as e: RNS.log("Error while closing socket for "+str(self)+": "+str(e)) def detach(self): RNS.log("Detaching "+str(self), RNS.LOG_DEBUG) @@ -555,15 +471,11 @@ class I2PInterfacePeer(Interface): if callable(self.socket.close): self.detached = True - try: - self.socket.shutdown(socket.SHUT_RDWR) - except Exception as e: - RNS.log("Error while shutting down socket for "+str(self)+": "+str(e)) + try: self.socket.shutdown(socket.SHUT_RDWR) + except Exception as e: RNS.log("Error while shutting down socket for "+str(self)+": "+str(e)) - try: - self.socket.close() - except Exception as e: - RNS.log("Error while closing socket for "+str(self)+": "+str(e)) + try: self.socket.close() + except Exception as e: RNS.log("Error while closing socket for "+str(self)+": "+str(e)) self.socket = None @@ -578,23 +490,17 @@ class I2PInterfacePeer(Interface): if not self.awaiting_i2p_tunnel: RNS.log("Initial connection for "+str(self)+" could not be established: "+str(e), RNS.LOG_ERROR) RNS.log("Leaving unconnected and retrying connection in "+str(I2PInterfacePeer.RECONNECT_WAIT)+" seconds.", RNS.LOG_ERROR) - return False - - else: - raise e + else: raise e - if platform.system() == "Linux": - self.set_timeouts_linux() - elif platform.system() == "Darwin": - self.set_timeouts_osx() + if platform.system() == "Linux": self.set_timeouts_linux() + elif platform.system() == "Darwin": self.set_timeouts_osx() self.online = True self.writing = False self.never_connected = False - if not self.kiss_framing and self.wants_tunnel: - RNS.Transport.synthesize_tunnel(self) + if not self.kiss_framing and self.wants_tunnel: RNS.Transport.synthesize_tunnel(self) return True @@ -612,14 +518,10 @@ class I2PInterfacePeer(Interface): self.teardown() break - try: - self.connect() - + try: self.connect() except Exception as e: - if not self.awaiting_i2p_tunnel: - RNS.log("Connection attempt for "+str(self)+" failed: "+str(e), RNS.LOG_DEBUG) - else: - RNS.log(str(self)+" still waiting for I2P tunnel to appear", RNS.LOG_VERBOSE) + if not self.awaiting_i2p_tunnel: RNS.log("Connection attempt for "+str(self)+" failed: "+str(e), RNS.LOG_DEBUG) + else: RNS.log(str(self)+" still waiting for I2P tunnel to appear", RNS.LOG_VERBOSE) if not self.never_connected: RNS.log(str(self)+" Re-established connection via I2P tunnel", RNS.LOG_INFO) @@ -628,8 +530,7 @@ class I2PInterfacePeer(Interface): thread = threading.Thread(target=self.read_loop) thread.daemon = True thread.start() - if not self.kiss_framing: - RNS.Transport.synthesize_tunnel(self) + if not self.kiss_framing: RNS.Transport.synthesize_tunnel(self) else: RNS.log("Attempt to reconnect on a non-initiator I2P interface. This should not happen.", RNS.LOG_ERROR) @@ -639,29 +540,24 @@ class I2PInterfacePeer(Interface): self.rxb += len(data) if hasattr(self, "parent_interface") and self.parent_interface != None and self.parent_count: self.parent_interface.rxb += len(data) - self.owner.inbound(data, self) def process_outgoing(self, data): if self.online: - while self.writing: - time.sleep(0.001) + while self.writing: time.sleep(0.001) try: self.writing = True - if self.kiss_framing: - data = bytes([KISS.FEND])+bytes([KISS.CMD_DATA])+KISS.escape(data)+bytes([KISS.FEND]) - else: - data = bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]) + if self.kiss_framing: data = bytes([KISS.FEND])+bytes([KISS.CMD_DATA])+KISS.escape(data)+bytes([KISS.FEND]) + else: data = bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]) self.socket.sendall(data) self.writing = False self.txb += len(data) self.last_write = time.time() - if hasattr(self, "parent_interface") and self.parent_interface != None and self.parent_count: - self.parent_interface.txb += len(data) + if hasattr(self, "parent_interface") and self.parent_interface != None and self.parent_count: self.parent_interface.txb += len(data) except Exception as e: RNS.log("Exception occurred while transmitting via "+str(self)+", tearing down interface", RNS.LOG_ERROR) @@ -670,8 +566,7 @@ class I2PInterfacePeer(Interface): def read_watchdog(self): - while self.wd_reset: - time.sleep(0.25) + while self.wd_reset: time.sleep(0.25) should_run = True try: @@ -679,17 +574,13 @@ class I2PInterfacePeer(Interface): time.sleep(1) if (time.time()-self.last_read > I2PInterfacePeer.I2P_PROBE_AFTER*2): - if self.i2p_tunnel_state != I2PInterfacePeer.TUNNEL_STATE_STALE: - RNS.log("I2P tunnel became unresponsive", RNS.LOG_DEBUG) - + if self.i2p_tunnel_state != I2PInterfacePeer.TUNNEL_STATE_STALE: RNS.log("I2P tunnel became unresponsive", RNS.LOG_DEBUG) self.i2p_tunnel_state = I2PInterfacePeer.TUNNEL_STATE_STALE - else: - self.i2p_tunnel_state = I2PInterfacePeer.TUNNEL_STATE_ACTIVE + else: self.i2p_tunnel_state = I2PInterfacePeer.TUNNEL_STATE_ACTIVE if (time.time()-self.last_write > I2PInterfacePeer.I2P_PROBE_AFTER*1): try: - if self.socket != None: - self.socket.sendall(bytes([HDLC.FLAG, HDLC.FLAG])) + if self.socket != None: self.socket.sendall(bytes([HDLC.FLAG, HDLC.FLAG])) except Exception as e: RNS.log("An error ocurred while sending I2P keepalive. The contained exception was: "+str(e), RNS.LOG_ERROR) self.shutdown_socket(self.socket) @@ -698,22 +589,17 @@ class I2PInterfacePeer(Interface): if (time.time()-self.last_read > I2PInterfacePeer.I2P_READ_TIMEOUT): RNS.log("I2P socket is unresponsive, restarting...", RNS.LOG_WARNING) if self.socket != None: - try: - self.socket.shutdown(socket.SHUT_RDWR) - except Exception as e: - RNS.log("Error while shutting down socket for "+str(self)+": "+str(e)) + try: self.socket.shutdown(socket.SHUT_RDWR) + except Exception as e: RNS.log("Error while shutting down socket for "+str(self)+": "+str(e)) - try: - self.socket.close() - except Exception as e: - RNS.log("Error while closing socket for "+str(self)+": "+str(e)) + try: self.socket.close() + except Exception as e: RNS.log("Error while closing socket for "+str(self)+": "+str(e)) should_run = False self.wd_reset = False - finally: - self.wd_reset = False + finally: self.wd_reset = False def read_loop(self): try: @@ -784,7 +670,6 @@ class I2PInterfacePeer(Interface): data_buffer = data_buffer+bytes([byte]) else: self.online = False - self.wd_reset = True time.sleep(2) self.wd_reset = False @@ -806,15 +691,12 @@ class I2PInterfacePeer(Interface): if self.initiator: RNS.log("Attempting to reconnect...", RNS.LOG_WARNING) self.reconnect() - else: - self.teardown() + else: self.teardown() def teardown(self): if self.initiator and not self.detached: RNS.log("The interface "+str(self)+" experienced an unrecoverable error and is being torn down. Restart Reticulum to attempt to open this interface again.", RNS.LOG_ERROR) - if RNS.Reticulum.panic_on_interface_error: - RNS.panic() - + if RNS.Reticulum.panic_on_interface_error: RNS.panic() else: RNS.log("The interface "+str(self)+" is being torn down.", RNS.LOG_VERBOSE) @@ -826,9 +708,7 @@ class I2PInterfacePeer(Interface): while self in self.parent_interface.spawned_interfaces: self.parent_interface.spawned_interfaces.remove(self) - if not self.initiator: - RNS.Transport.remove_interface(self) - + if not self.initiator: RNS.Transport.remove_interface(self) def __str__(self): return "I2PInterfacePeer["+str(self.name)+"]" @@ -839,8 +719,7 @@ class I2PInterface(Interface): DEFAULT_IFAC_SIZE = 16 @property - def clients(self): - return len(self.spawned_interfaces) + def clients(self): return len(self.spawned_interfaces) def __init__(self, owner, configuration): super().__init__() @@ -894,11 +773,9 @@ class I2PInterface(Interface): RNS.log("I2P controller did not become available in time, waiting for controller", RNS.LOG_VERBOSE) i2p_notready_warning = True - while not self.i2p.ready: - time.sleep(0.25) + while not self.i2p.ready: time.sleep(0.25) - if i2p_notready_warning == True: - RNS.log("I2P controller ready, continuing setup", RNS.LOG_VERBOSE) + if i2p_notready_warning == True: RNS.log("I2P controller ready, continuing setup", RNS.LOG_VERBOSE) def handlerFactory(callback): def createHandler(*args, **keys): @@ -971,18 +848,11 @@ class I2PInterface(Interface): spawned_interface.ifac_netkey = self.ifac_netkey if spawned_interface.ifac_netname != None or spawned_interface.ifac_netkey != None: ifac_origin = b"" - if spawned_interface.ifac_netname != None: - ifac_origin += RNS.Identity.full_hash(spawned_interface.ifac_netname.encode("utf-8")) - if spawned_interface.ifac_netkey != None: - ifac_origin += RNS.Identity.full_hash(spawned_interface.ifac_netkey.encode("utf-8")) + if spawned_interface.ifac_netname != None: ifac_origin += RNS.Identity.full_hash(spawned_interface.ifac_netname.encode("utf-8")) + if spawned_interface.ifac_netkey != None: ifac_origin += RNS.Identity.full_hash(spawned_interface.ifac_netkey.encode("utf-8")) ifac_origin_hash = RNS.Identity.full_hash(ifac_origin) - spawned_interface.ifac_key = RNS.Cryptography.hkdf( - length=64, - derive_from=ifac_origin_hash, - salt=RNS.Reticulum.IFAC_SALT, - context=None - ) + spawned_interface.ifac_key = RNS.Cryptography.hkdf(length=64, derive_from=ifac_origin_hash, salt=RNS.Reticulum.IFAC_SALT, context=None) spawned_interface.ifac_identity = RNS.Identity.from_bytes(spawned_interface.ifac_key) spawned_interface.ifac_signature = spawned_interface.ifac_identity.sign(RNS.Identity.full_hash(spawned_interface.ifac_key)) @@ -994,8 +864,7 @@ class I2PInterface(Interface): spawned_interface.HW_MTU = self.HW_MTU RNS.log("Spawned new I2PInterface Peer: "+str(spawned_interface), RNS.LOG_VERBOSE) RNS.Transport.add_interface(spawned_interface) - while spawned_interface in self.spawned_interfaces: - self.spawned_interfaces.remove(spawned_interface) + while spawned_interface in self.spawned_interfaces: self.spawned_interfaces.remove(spawned_interface) self.spawned_interfaces.append(spawned_interface) spawned_interface.read_loop() @@ -1018,13 +887,11 @@ class I2PInterface(Interface): RNS.log("Detaching "+str(self), RNS.LOG_DEBUG) self.i2p.stop() - def __str__(self): - return "I2PInterface["+self.name+"]" + def __str__(self): return "I2PInterface["+self.name+"]" class I2PInterfaceHandler(socketserver.BaseRequestHandler): def __init__(self, callback, *args, **keys): self.callback = callback socketserver.BaseRequestHandler.__init__(self, *args, **keys) - def handle(self): - self.callback(handler=self) + def handle(self): self.callback(handler=self) diff --git a/RNS/vendor/i2plib/aiosam.py b/RNS/vendor/i2plib/aiosam.py index bbb2c4aa..a12e0fc7 100644 --- a/RNS/vendor/i2plib/aiosam.py +++ b/RNS/vendor/i2plib/aiosam.py @@ -6,14 +6,12 @@ from . import utils from .log import logger def parse_reply(data): - if not data: - raise ConnectionAbortedError("Empty response: SAM API went offline") + if not data: raise ConnectionAbortedError("Empty response: SAM API went offline") try: msg = sam.Message(data.decode().strip()) logger.debug("SAM reply: "+str(msg)) - except: - raise ConnectionAbortedError("Invalid SAM response") + except: raise ConnectionAbortedError("Invalid SAM response") return msg @@ -28,8 +26,7 @@ async def get_sam_socket(sam_address=sam.DEFAULT_ADDRESS, loop=None): reader, writer = await asyncio.open_connection(*sam_address) writer.write(sam.hello("3.1", "3.1")) reply = parse_reply(await reader.readline()) - if reply.ok: - return (reader, writer) + if reply.ok: return (reader, writer) else: writer.close() raise exceptions.SAM_EXCEPTIONS[reply["RESULT"]]() @@ -49,10 +46,8 @@ async def dest_lookup(domain, sam_address=sam.DEFAULT_ADDRESS, writer.write(sam.naming_lookup(domain)) reply = parse_reply(await reader.readline()) writer.close() - if reply.ok: - return sam.Destination(reply["VALUE"]) - else: - raise exceptions.SAM_EXCEPTIONS[reply["RESULT"]]() + if reply.ok: return sam.Destination(reply["VALUE"]) + else: raise exceptions.SAM_EXCEPTIONS[reply["RESULT"]]() async def new_destination(sam_address=sam.DEFAULT_ADDRESS, loop=None, sig_type=sam.Destination.default_sig_type): @@ -91,15 +86,10 @@ async def create_session(session_name, sam_address=sam.DEFAULT_ADDRESS, """ logger.debug("Creating session {}".format(session_name)) if destination: - if type(destination) == sam.Destination: - destination = destination - else: - destination = sam.Destination( - destination, has_private_key=True) - + if type(destination) == sam.Destination: destination = destination + else: destination = sam.Destination(destination, has_private_key=True) dest_string = destination.private_key.base64 - else: - dest_string = sam.TRANSIENT_DESTINATION + else: dest_string = sam.TRANSIENT_DESTINATION options = " ".join(["{}={}".format(k, v) for k, v in options.items()]) @@ -109,9 +99,7 @@ async def create_session(session_name, sam_address=sam.DEFAULT_ADDRESS, reply = parse_reply(await reader.readline()) if reply.ok: - if not destination: - destination = sam.Destination( - reply["DESTINATION"], has_private_key=True) + if not destination: destination = sam.Destination(reply["DESTINATION"], has_private_key=True) logger.debug(destination.base32) logger.debug("Session created {}".format(session_name)) return (reader, writer) @@ -130,10 +118,8 @@ async def stream_connect(session_name, destination, :return: A (reader, writer) pair """ logger.debug("Connecting stream {}".format(session_name)) - if isinstance(destination, str) and not destination.endswith(".i2p"): - destination = sam.Destination(destination) - elif isinstance(destination, str): - destination = await dest_lookup(destination, sam_address, loop) + if isinstance(destination, str) and not destination.endswith(".i2p"): destination = sam.Destination(destination) + elif isinstance(destination, str): destination = await dest_lookup(destination, sam_address, loop) reader, writer = await get_sam_socket(sam_address, loop) writer.write(sam.stream_connect(session_name, destination.base64, @@ -158,8 +144,7 @@ async def stream_accept(session_name, sam_address=sam.DEFAULT_ADDRESS, reader, writer = await get_sam_socket(sam_address, loop) writer.write(sam.stream_accept(session_name, silent="false")) reply = parse_reply(await reader.readline()) - if reply.ok: - return (reader, writer) + if reply.ok: return (reader, writer) else: writer.close() raise exceptions.SAM_EXCEPTIONS[reply["RESULT"]]() @@ -183,9 +168,9 @@ class Session: :return: :class:`Session` object """ def __init__(self, session_name, sam_address=sam.DEFAULT_ADDRESS, - loop=None, style="STREAM", - signature_type=sam.Destination.default_sig_type, - destination=None, options={}): + loop=None, style="STREAM", signature_type=sam.Destination.default_sig_type, + destination=None, options={}): + self.session_name = session_name self.sam_address = sam_address self.loop = loop @@ -195,10 +180,9 @@ class Session: self.options = options async def __aenter__(self): - self.reader, self.writer = await create_session(self.session_name, - sam_address=self.sam_address, loop=self.loop, style=self.style, - signature_type=self.signature_type, - destination=self.destination, options=self.options) + self.reader, self.writer = await create_session(self.session_name, sam_address=self.sam_address, loop=self.loop, + style=self.style, signature_type=self.signature_type, + destination=self.destination, options=self.options) return self async def __aexit__(self, exc_type, exc, tb): @@ -214,16 +198,15 @@ class StreamConnection: :param loop: (optional) Event loop instance :return: :class:`StreamConnection` object """ - def __init__(self, session_name, destination, - sam_address=sam.DEFAULT_ADDRESS, loop=None): + def __init__(self, session_name, destination, sam_address=sam.DEFAULT_ADDRESS, loop=None): self.session_name = session_name self.sam_address = sam_address self.loop = loop self.destination = destination async def __aenter__(self): - self.reader, self.writer = await stream_connect(self.session_name, - self.destination, sam_address=self.sam_address, loop=self.loop) + self.reader, self.writer = await stream_connect(self.session_name, self.destination, + sam_address=self.sam_address, loop=self.loop) self.read = self.reader.read self.write = self.writer.write return self @@ -240,15 +223,13 @@ class StreamAcceptor: :param loop: (optional) Event loop instance :return: :class:`StreamAcceptor` object """ - def __init__(self, session_name, sam_address=sam.DEFAULT_ADDRESS, - loop=None): + def __init__(self, session_name, sam_address=sam.DEFAULT_ADDRESS, loop=None): self.session_name = session_name self.sam_address = sam_address self.loop = loop async def __aenter__(self): - self.reader, self.writer = await stream_accept(self.session_name, - sam_address=self.sam_address, loop=self.loop) + self.reader, self.writer = await stream_accept(self.session_name, sam_address=self.sam_address, loop=self.loop) self.read = self.reader.read self.write = self.writer.write return self diff --git a/RNS/vendor/i2plib/exceptions.py b/RNS/vendor/i2plib/exceptions.py index d693c841..477cfee9 100644 --- a/RNS/vendor/i2plib/exceptions.py +++ b/RNS/vendor/i2plib/exceptions.py @@ -30,15 +30,12 @@ class PeerNotFound(SAMException): class Timeout(SAMException): """The peer cannot be found on the network""" -SAM_EXCEPTIONS = { - "CANT_REACH_PEER": CantReachPeer, - "DUPLICATED_DEST": DuplicatedDest, - "DUPLICATED_ID": DuplicatedId, - "I2P_ERROR": I2PError, - "INVALID_ID": InvalidId, - "INVALID_KEY": InvalidKey, - "KEY_NOT_FOUND": KeyNotFound, - "PEER_NOT_FOUND": PeerNotFound, - "TIMEOUT": Timeout, -} - +SAM_EXCEPTIONS = { "CANT_REACH_PEER": CantReachPeer, + "DUPLICATED_DEST": DuplicatedDest, + "DUPLICATED_ID": DuplicatedId, + "I2P_ERROR": I2PError, + "INVALID_ID": InvalidId, + "INVALID_KEY": InvalidKey, + "KEY_NOT_FOUND": KeyNotFound, + "PEER_NOT_FOUND": PeerNotFound, + "TIMEOUT": Timeout } \ No newline at end of file diff --git a/RNS/vendor/i2plib/sam.py b/RNS/vendor/i2plib/sam.py index 1bad75d0..2176f0e3 100644 --- a/RNS/vendor/i2plib/sam.py +++ b/RNS/vendor/i2plib/sam.py @@ -3,7 +3,6 @@ from hashlib import sha256 import struct import re - I2P_B64_CHARS = "-~" def i2p_b64encode(x): @@ -41,11 +40,9 @@ class Message(object): return self.opts[key] @property - def ok(self): - return self["RESULT"] == "OK" + def ok(self): return self["RESULT"] == "OK" - def __repr__(self): - return self._reply_string + def __repr__(self): return self._reply_string # SAM request messages