diff --git a/RNS/Transport.py b/RNS/Transport.py index 4fe27aa6..589c7984 100755 --- a/RNS/Transport.py +++ b/RNS/Transport.py @@ -1710,830 +1710,828 @@ class Transport: elif Transport.interface_to_shared_instance(packet.receiving_interface): packet.hops -= 1 - # TODO: Clean up - if True: - # By default, remember packet hashes to avoid routing - # loops in the network, using the packet filter. - remember_packet_hash = True + # By default, remember packet hashes to avoid routing + # loops in the network, using the packet filter. + remember_packet_hash = True - # If this packet belongs to a link in our link table, - # we'll have to defer adding it to the filter list. - # In some cases, we might see a packet over a shared- - # medium interface, belonging to a link that transports - # or terminates with this instance, but before it would - # normally reach us. If the packet is appended to the - # filter list at this point, link transport will break. - if packet.destination_hash in Transport.link_table: remember_packet_hash = False + # If this packet belongs to a link in our link table, + # we'll have to defer adding it to the filter list. + # In some cases, we might see a packet over a shared- + # medium interface, belonging to a link that transports + # or terminates with this instance, but before it would + # normally reach us. If the packet is appended to the + # filter list at this point, link transport will break. + if packet.destination_hash in Transport.link_table: remember_packet_hash = False - # If this is a link request proof, don't add it until - # we are sure it's not actually somewhere else in the - # routing chain. - if packet.packet_type == RNS.Packet.PROOF and packet.context == RNS.Packet.LRPROOF: remember_packet_hash = False + # If this is a link request proof, don't add it until + # we are sure it's not actually somewhere else in the + # routing chain. + if packet.packet_type == RNS.Packet.PROOF and packet.context == RNS.Packet.LRPROOF: remember_packet_hash = False - if remember_packet_hash: - Transport.add_packet_hash(packet.packet_hash) - # TODO: Enable when caching has been redesigned - # Transport.cache(packet) - - # Check special conditions for local clients connected - # through a shared Reticulum instance - from_local_client = (packet.receiving_interface in Transport.local_client_interfaces) - for_local_client = (packet.packet_type != RNS.Packet.ANNOUNCE) and (packet.destination_hash in Transport.path_table and Transport.path_table[packet.destination_hash][IDX_PT_HOPS] == 0) - local_client_nh = (packet.packet_type != RNS.Packet.ANNOUNCE) and (packet.destination_hash in Transport.link_table and Transport.link_table[packet.destination_hash][IDX_LT_NH_IF] in Transport.local_client_interfaces) - local_client_rh = (packet.packet_type != RNS.Packet.ANNOUNCE) and (packet.destination_hash in Transport.link_table and Transport.link_table[packet.destination_hash][IDX_LT_RCVD_IF] in Transport.local_client_interfaces) - proof_for_local_client = (packet.destination_hash in Transport.reverse_table) and (Transport.reverse_table[packet.destination_hash][IDX_RT_RCVD_IF] in Transport.local_client_interfaces) - to_local_client = for_local_client or proof_for_local_client - instance_local_link = local_client_nh and local_client_rh - for_local_client_link = local_client_nh or local_client_rh - link_request_handled = False + if remember_packet_hash: + Transport.add_packet_hash(packet.packet_hash) + # TODO: Enable when caching has been redesigned + # Transport.cache(packet) + + # Check special conditions for local clients connected + # through a shared Reticulum instance + from_local_client = (packet.receiving_interface in Transport.local_client_interfaces) + for_local_client = (packet.packet_type != RNS.Packet.ANNOUNCE) and (packet.destination_hash in Transport.path_table and Transport.path_table[packet.destination_hash][IDX_PT_HOPS] == 0) + local_client_nh = (packet.packet_type != RNS.Packet.ANNOUNCE) and (packet.destination_hash in Transport.link_table and Transport.link_table[packet.destination_hash][IDX_LT_NH_IF] in Transport.local_client_interfaces) + local_client_rh = (packet.packet_type != RNS.Packet.ANNOUNCE) and (packet.destination_hash in Transport.link_table and Transport.link_table[packet.destination_hash][IDX_LT_RCVD_IF] in Transport.local_client_interfaces) + proof_for_local_client = (packet.destination_hash in Transport.reverse_table) and (Transport.reverse_table[packet.destination_hash][IDX_RT_RCVD_IF] in Transport.local_client_interfaces) + to_local_client = for_local_client or proof_for_local_client + instance_local_link = local_client_nh and local_client_rh + for_local_client_link = local_client_nh or local_client_rh + link_request_handled = False - # Plain broadcast packets from local clients are sent - # directly on all attached interfaces, since they are - # never injected into transport. - if not packet.destination_hash in Transport.control_hashes: - if packet.destination_type == RNS.Destination.PLAIN and packet.transport_type == Transport.BROADCAST: - # Send to all interfaces except the originator - if from_local_client: - for interface in Transport.interfaces: - if interface != packet.receiving_interface: - Transport.transmit(interface, packet.raw) - # If the packet was not from a local client, send - # it directly to all local clients - else: - for interface in Transport.local_client_interfaces: + # Plain broadcast packets from local clients are sent + # directly on all attached interfaces, since they are + # never injected into transport. + if not packet.destination_hash in Transport.control_hashes: + if packet.destination_type == RNS.Destination.PLAIN and packet.transport_type == Transport.BROADCAST: + # Send to all interfaces except the originator + if from_local_client: + for interface in Transport.interfaces: + if interface != packet.receiving_interface: Transport.transmit(interface, packet.raw) + # If the packet was not from a local client, send + # it directly to all local clients + else: + for interface in Transport.local_client_interfaces: + Transport.transmit(interface, packet.raw) - # General transport handling. Takes care of directing - # packets according to transport tables and recording - # entries in reverse and link tables. - if RNS.Reticulum.transport_enabled() or from_local_client or for_local_client or for_local_client_link: + # General transport handling. Takes care of directing + # packets according to transport tables and recording + # entries in reverse and link tables. + if RNS.Reticulum.transport_enabled() or from_local_client or for_local_client or for_local_client_link: - # If there is no transport id, but the packet is - # for a local client, we generate the transport - # id (it was stripped on the previous hop, since - # we "spoof" the hop count for clients behind a - # shared instance, so they look directly reach- - # able), and reinsert, so the normal transport - # implementation can handle the packet. - if packet.transport_id == None and for_local_client: - packet.transport_id = Transport.identity.hash + # If there is no transport id, but the packet is + # for a local client, we generate the transport + # id (it was stripped on the previous hop, since + # we "spoof" the hop count for clients behind a + # shared instance, so they look directly reach- + # able), and reinsert, so the normal transport + # implementation can handle the packet. + if packet.transport_id == None and for_local_client: + packet.transport_id = Transport.identity.hash - # If this is a cache request, and we can fullfill - # it, do so and stop processing. Otherwise resume - # normal processing. - if packet.context == RNS.Packet.CACHE_REQUEST: - if Transport.cache_request_packet(packet): return + # If this is a cache request, and we can fullfill + # it, do so and stop processing. Otherwise resume + # normal processing. + if packet.context == RNS.Packet.CACHE_REQUEST: + if Transport.cache_request_packet(packet): return - # If the packet is in transport, check whether we - # are the designated next hop, and process it - # accordingly if we are. - if packet.transport_id != None and packet.packet_type != RNS.Packet.ANNOUNCE: - if packet.transport_id == Transport.identity.hash: - if packet.destination_hash in Transport.path_table: - next_hop = Transport.path_table[packet.destination_hash][IDX_PT_NEXT_HOP] - remaining_hops = Transport.path_table[packet.destination_hash][IDX_PT_HOPS] - - if remaining_hops > 1: - # Just increase hop count and transmit - new_raw = packet.raw[0:1] - new_raw += struct.pack("!B", packet.hops) - new_raw += next_hop - new_raw += packet.raw[(RNS.Identity.TRUNCATED_HASHLENGTH//8)+2:] - elif remaining_hops == 1: - # Strip transport headers and transmit - new_flags = (RNS.Packet.HEADER_1) << 6 | (Transport.BROADCAST) << 4 | (packet.flags & 0b00001111) - new_raw = struct.pack("!B", new_flags) - new_raw += struct.pack("!B", packet.hops) - new_raw += packet.raw[(RNS.Identity.TRUNCATED_HASHLENGTH//8)+2:] - elif remaining_hops == 0: - if to_local_client and RNS.Transport.local_hops_delta != 0: - if packet.header_type == RNS.Packet.HEADER_2: - # Strip transport headers and transmit - new_flags = (RNS.Packet.HEADER_1) << 6 | (Transport.BROADCAST) << 4 | (packet.flags & 0b00001111) - new_raw = struct.pack("!B", new_flags) - new_raw += struct.pack("!B", packet.hops) - new_raw += packet.raw[(RNS.Identity.TRUNCATED_HASHLENGTH//8)+2:] - else: - # Just increase hop count and transmit - new_raw = packet.raw[0:1] - new_raw += struct.pack("!B", packet.hops) - new_raw += packet.raw[2:] + # If the packet is in transport, check whether we + # are the designated next hop, and process it + # accordingly if we are. + if packet.transport_id != None and packet.packet_type != RNS.Packet.ANNOUNCE: + if packet.transport_id == Transport.identity.hash: + if packet.destination_hash in Transport.path_table: + next_hop = Transport.path_table[packet.destination_hash][IDX_PT_NEXT_HOP] + remaining_hops = Transport.path_table[packet.destination_hash][IDX_PT_HOPS] + + if remaining_hops > 1: + # Just increase hop count and transmit + new_raw = packet.raw[0:1] + new_raw += struct.pack("!B", packet.hops) + new_raw += next_hop + new_raw += packet.raw[(RNS.Identity.TRUNCATED_HASHLENGTH//8)+2:] + elif remaining_hops == 1: + # Strip transport headers and transmit + new_flags = (RNS.Packet.HEADER_1) << 6 | (Transport.BROADCAST) << 4 | (packet.flags & 0b00001111) + new_raw = struct.pack("!B", new_flags) + new_raw += struct.pack("!B", packet.hops) + new_raw += packet.raw[(RNS.Identity.TRUNCATED_HASHLENGTH//8)+2:] + elif remaining_hops == 0: + if to_local_client and RNS.Transport.local_hops_delta != 0: + if packet.header_type == RNS.Packet.HEADER_2: + # Strip transport headers and transmit + new_flags = (RNS.Packet.HEADER_1) << 6 | (Transport.BROADCAST) << 4 | (packet.flags & 0b00001111) + new_raw = struct.pack("!B", new_flags) + new_raw += struct.pack("!B", packet.hops) + new_raw += packet.raw[(RNS.Identity.TRUNCATED_HASHLENGTH//8)+2:] else: # Just increase hop count and transmit new_raw = packet.raw[0:1] new_raw += struct.pack("!B", packet.hops) new_raw += packet.raw[2:] - - outbound_interface = Transport.path_table[packet.destination_hash][IDX_PT_RVCD_IF] - - if packet.packet_type == RNS.Packet.LINKREQUEST: - now = time.time() - proof_timeout = Transport.extra_link_proof_timeout(packet.receiving_interface) - proof_timeout += now + RNS.Link.ESTABLISHMENT_TIMEOUT_PER_HOP * max(1, remaining_hops) - - path_mtu = RNS.Link.mtu_from_lr_packet(packet) - mode = RNS.Link.mode_from_lr_packet(packet) - ph_mtu = packet.receiving_interface.HW_MTU if packet.receiving_interface else None - nh_mtu = outbound_interface.HW_MTU - if path_mtu: - if outbound_interface.HW_MTU == None: - RNS.log(f"No next-hop HW MTU, disabling link MTU upgrade", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - path_mtu = None - new_raw = new_raw[:-RNS.Link.LINK_MTU_SIZE] - elif not outbound_interface.AUTOCONFIGURE_MTU and not outbound_interface.FIXED_MTU: - RNS.log(f"Outbound interface doesn't support MTU autoconfiguration, disabling link MTU upgrade", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - path_mtu = None - new_raw = new_raw[:-RNS.Link.LINK_MTU_SIZE] - else: - if nh_mtu < path_mtu or (ph_mtu and ph_mtu < path_mtu): - try: - path_mtu = min(nh_mtu, ph_mtu) - clamped_mtu = RNS.Link.signalling_bytes(path_mtu, mode) - RNS.log(f"Clamping link MTU to {RNS.prettysize(path_mtu)}", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - new_raw = new_raw[:-RNS.Link.LINK_MTU_SIZE]+clamped_mtu - except Exception as e: - RNS.log(f"Dropping link request packet. The contained exception was: {e}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None - return - - # Entry format is - link_entry = [ now, # 0: Timestamp, - next_hop, # 1: Next-hop transport ID - outbound_interface, # 2: Next-hop interface - remaining_hops, # 3: Remaining hops - packet.receiving_interface, # 4: Received on interface - packet.hops, # 5: Taken hops - packet.destination_hash, # 6: Original destination hash - False, # 7: Validated - proof_timeout ] # 8: Proof timeout timestamp - - with Transport.link_table_lock: Transport.link_table[RNS.Link.link_id_from_lr_packet(packet)] = link_entry - link_request_handled = True - else: - # Entry format is - reverse_entry = [ packet.receiving_interface, # 0: Received on interface - outbound_interface, # 1: Outbound interface - time.time() ] # 2: Timestamp + # Just increase hop count and transmit + new_raw = packet.raw[0:1] + new_raw += struct.pack("!B", packet.hops) + new_raw += packet.raw[2:] - with Transport.reverse_table_lock: Transport.reverse_table[packet.getTruncatedHash()] = reverse_entry + outbound_interface = Transport.path_table[packet.destination_hash][IDX_PT_RVCD_IF] - if Transport.local_hops_delta != 0 and from_local_client and not to_local_client: new_raw = Transport.mangle_hops(new_raw, Transport.local_hops_delta) - Transport.transmit(outbound_interface, new_raw) - with Transport.path_table_lock: Transport.path_table[packet.destination_hash][IDX_PT_TIMESTAMP] = time.time() - - else: - # TODO: There should probably be some kind of REJECT - # mechanism here, to signal to the source that their - # expected path failed. - RNS.log("Got packet in transport, but no known path to final destination "+RNS.prettyhexrep(packet.destination_hash)+". Dropping packet.", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None - - # Link transport handling. Directs packets according - # to entries in the link tables - if packet.packet_type != RNS.Packet.ANNOUNCE and packet.packet_type != RNS.Packet.LINKREQUEST and packet.context != RNS.Packet.LRPROOF: - if packet.destination_hash in Transport.link_table: - link_entry = Transport.link_table[packet.destination_hash] - # If receiving and outbound interface is - # the same for this link, direction doesn't - # matter, and we simply repeat the packet. - outbound_interface = None - if link_entry[IDX_LT_NH_IF] == link_entry[IDX_LT_RCVD_IF]: - # But check that taken hops matches one - # of the expectede values. - if packet.hops == link_entry[IDX_LT_REM_HOPS] or packet.hops == link_entry[IDX_LT_HOPS]: - outbound_interface = link_entry[IDX_LT_NH_IF] - else: - # If interfaces differ, we transmit on - # the opposite interface of what the - # packet was received on. - if packet.receiving_interface == link_entry[IDX_LT_NH_IF]: - # Also check that expected hop count matches - if packet.hops == link_entry[IDX_LT_REM_HOPS]: - outbound_interface = link_entry[IDX_LT_RCVD_IF] - elif packet.receiving_interface == link_entry[IDX_LT_RCVD_IF]: - # Also check that expected hop count matches - if packet.hops == link_entry[IDX_LT_HOPS]: - outbound_interface = link_entry[IDX_LT_NH_IF] - - if outbound_interface != None: - # Add this packet to the filter hashlist if we - # have determined that it's actually our turn - # to process it. - Transport.add_packet_hash(packet.packet_hash) - - new_raw = packet.raw[0:1] - new_raw += struct.pack("!B", packet.hops if not from_local_client or instance_local_link or Transport.local_hops_delta == 0 else Transport.local_hops_delta) - new_raw += packet.raw[2:] - Transport.transmit(outbound_interface, new_raw) - Transport.link_table[packet.destination_hash][IDX_LT_TIMESTAMP] = time.time() - - # TODO: Can we return safely here? Test and possibly enable this at some point. - # return - - - # Announce handling. Handles logic related to incoming - # announces, queueing rebroadcasts of these, and removal - # of queued announce rebroadcasts once handed to the next node. - if packet.packet_type == RNS.Packet.ANNOUNCE: - local_destination = None - with Transport.destinations_map_lock: - if packet.destination_hash in Transport.destinations_map: - local_destination = Transport.destinations_map[packet.destination_hash] - - if local_destination == None and RNS.Identity.validate_announce(packet): - if packet.transport_id != None: - received_from = packet.transport_id - - # Check if this is a next retransmission from - # another node. If it is, we're removing the - # announce in question from our pending table - if RNS.Reticulum.transport_enabled() and packet.destination_hash in Transport.announce_table: - announce_entry = Transport.announce_table[packet.destination_hash] + if packet.packet_type == RNS.Packet.LINKREQUEST: + now = time.time() + proof_timeout = Transport.extra_link_proof_timeout(packet.receiving_interface) + proof_timeout += now + RNS.Link.ESTABLISHMENT_TIMEOUT_PER_HOP * max(1, remaining_hops) - if packet.hops-1 == announce_entry[IDX_AT_HOPS]: - RNS.log(f"Heard a rebroadcast of announce for {RNS.prettyhexrep(packet.destination_hash)} on {packet.receiving_interface}", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None - announce_entry[IDX_AT_LCL_RBRD] += 1 - if announce_entry[IDX_AT_RETRIES] > 0: - if announce_entry[IDX_AT_LCL_RBRD] >= Transport.LOCAL_REBROADCASTS_MAX: - RNS.log("Completed announce processing for "+RNS.prettyhexrep(packet.destination_hash)+", local rebroadcast limit reached", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None - with Transport.announce_table_lock: - if packet.destination_hash in Transport.announce_table: Transport.announce_table.pop(packet.destination_hash) + path_mtu = RNS.Link.mtu_from_lr_packet(packet) + mode = RNS.Link.mode_from_lr_packet(packet) + ph_mtu = packet.receiving_interface.HW_MTU if packet.receiving_interface else None + nh_mtu = outbound_interface.HW_MTU + if path_mtu: + if outbound_interface.HW_MTU == None: + RNS.log(f"No next-hop HW MTU, disabling link MTU upgrade", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + path_mtu = None + new_raw = new_raw[:-RNS.Link.LINK_MTU_SIZE] + elif not outbound_interface.AUTOCONFIGURE_MTU and not outbound_interface.FIXED_MTU: + RNS.log(f"Outbound interface doesn't support MTU autoconfiguration, disabling link MTU upgrade", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + path_mtu = None + new_raw = new_raw[:-RNS.Link.LINK_MTU_SIZE] + else: + if nh_mtu < path_mtu or (ph_mtu and ph_mtu < path_mtu): + try: + path_mtu = min(nh_mtu, ph_mtu) + clamped_mtu = RNS.Link.signalling_bytes(path_mtu, mode) + RNS.log(f"Clamping link MTU to {RNS.prettysize(path_mtu)}", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + new_raw = new_raw[:-RNS.Link.LINK_MTU_SIZE]+clamped_mtu + except Exception as e: + RNS.log(f"Dropping link request packet. The contained exception was: {e}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + return - if packet.hops-1 == announce_entry[IDX_AT_HOPS]+1 and announce_entry[IDX_AT_RETRIES] > 0: - now = time.time() - if now < announce_entry[IDX_AT_RTRNS_TMO]: - RNS.log("Rebroadcasted announce for "+RNS.prettyhexrep(packet.destination_hash)+" has been passed on to another node, no further tries needed", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None + # Entry format is + link_entry = [ now, # 0: Timestamp, + next_hop, # 1: Next-hop transport ID + outbound_interface, # 2: Next-hop interface + remaining_hops, # 3: Remaining hops + packet.receiving_interface, # 4: Received on interface + packet.hops, # 5: Taken hops + packet.destination_hash, # 6: Original destination hash + False, # 7: Validated + proof_timeout ] # 8: Proof timeout timestamp + + with Transport.link_table_lock: Transport.link_table[RNS.Link.link_id_from_lr_packet(packet)] = link_entry + link_request_handled = True + + else: + # Entry format is + reverse_entry = [ packet.receiving_interface, # 0: Received on interface + outbound_interface, # 1: Outbound interface + time.time() ] # 2: Timestamp + + with Transport.reverse_table_lock: Transport.reverse_table[packet.getTruncatedHash()] = reverse_entry + + if Transport.local_hops_delta != 0 and from_local_client and not to_local_client: new_raw = Transport.mangle_hops(new_raw, Transport.local_hops_delta) + Transport.transmit(outbound_interface, new_raw) + with Transport.path_table_lock: Transport.path_table[packet.destination_hash][IDX_PT_TIMESTAMP] = time.time() + + else: + # TODO: There should probably be some kind of REJECT + # mechanism here, to signal to the source that their + # expected path failed. + RNS.log("Got packet in transport, but no known path to final destination "+RNS.prettyhexrep(packet.destination_hash)+". Dropping packet.", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None + + # Link transport handling. Directs packets according + # to entries in the link tables + if packet.packet_type != RNS.Packet.ANNOUNCE and packet.packet_type != RNS.Packet.LINKREQUEST and packet.context != RNS.Packet.LRPROOF: + if packet.destination_hash in Transport.link_table: + link_entry = Transport.link_table[packet.destination_hash] + # If receiving and outbound interface is + # the same for this link, direction doesn't + # matter, and we simply repeat the packet. + outbound_interface = None + if link_entry[IDX_LT_NH_IF] == link_entry[IDX_LT_RCVD_IF]: + # But check that taken hops matches one + # of the expectede values. + if packet.hops == link_entry[IDX_LT_REM_HOPS] or packet.hops == link_entry[IDX_LT_HOPS]: + outbound_interface = link_entry[IDX_LT_NH_IF] + else: + # If interfaces differ, we transmit on + # the opposite interface of what the + # packet was received on. + if packet.receiving_interface == link_entry[IDX_LT_NH_IF]: + # Also check that expected hop count matches + if packet.hops == link_entry[IDX_LT_REM_HOPS]: + outbound_interface = link_entry[IDX_LT_RCVD_IF] + elif packet.receiving_interface == link_entry[IDX_LT_RCVD_IF]: + # Also check that expected hop count matches + if packet.hops == link_entry[IDX_LT_HOPS]: + outbound_interface = link_entry[IDX_LT_NH_IF] + + if outbound_interface != None: + # Add this packet to the filter hashlist if we + # have determined that it's actually our turn + # to process it. + Transport.add_packet_hash(packet.packet_hash) + + new_raw = packet.raw[0:1] + new_raw += struct.pack("!B", packet.hops if not from_local_client or instance_local_link or Transport.local_hops_delta == 0 else Transport.local_hops_delta) + new_raw += packet.raw[2:] + Transport.transmit(outbound_interface, new_raw) + Transport.link_table[packet.destination_hash][IDX_LT_TIMESTAMP] = time.time() + + # TODO: Can we return safely here? Test and possibly enable this at some point. + # return + + + # Announce handling. Handles logic related to incoming + # announces, queueing rebroadcasts of these, and removal + # of queued announce rebroadcasts once handed to the next node. + if packet.packet_type == RNS.Packet.ANNOUNCE: + local_destination = None + with Transport.destinations_map_lock: + if packet.destination_hash in Transport.destinations_map: + local_destination = Transport.destinations_map[packet.destination_hash] + + if local_destination == None and RNS.Identity.validate_announce(packet): + if packet.transport_id != None: + received_from = packet.transport_id + + # Check if this is a next retransmission from + # another node. If it is, we're removing the + # announce in question from our pending table + if RNS.Reticulum.transport_enabled() and packet.destination_hash in Transport.announce_table: + announce_entry = Transport.announce_table[packet.destination_hash] + + if packet.hops-1 == announce_entry[IDX_AT_HOPS]: + RNS.log(f"Heard a rebroadcast of announce for {RNS.prettyhexrep(packet.destination_hash)} on {packet.receiving_interface}", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None + announce_entry[IDX_AT_LCL_RBRD] += 1 + if announce_entry[IDX_AT_RETRIES] > 0: + if announce_entry[IDX_AT_LCL_RBRD] >= Transport.LOCAL_REBROADCASTS_MAX: + RNS.log("Completed announce processing for "+RNS.prettyhexrep(packet.destination_hash)+", local rebroadcast limit reached", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None with Transport.announce_table_lock: if packet.destination_hash in Transport.announce_table: Transport.announce_table.pop(packet.destination_hash) - else: received_from = packet.destination_hash + if packet.hops-1 == announce_entry[IDX_AT_HOPS]+1 and announce_entry[IDX_AT_RETRIES] > 0: + now = time.time() + if now < announce_entry[IDX_AT_RTRNS_TMO]: + RNS.log("Rebroadcasted announce for "+RNS.prettyhexrep(packet.destination_hash)+" has been passed on to another node, no further tries needed", RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None + with Transport.announce_table_lock: + if packet.destination_hash in Transport.announce_table: Transport.announce_table.pop(packet.destination_hash) - # Check if this announce should be inserted into - # announce and destination tables - should_add = False + else: received_from = packet.destination_hash - # First, check that the announce is not for a destination - # local to this system, and that hops are less than the max - with Transport.destinations_map_lock: - local_and_hops_condition = (packet.hops < Transport.PATHFINDER_M+1) and (not packet.destination_hash in Transport.destinations_map) + # Check if this announce should be inserted into + # announce and destination tables + should_add = False - if local_and_hops_condition: - announce_emitted = Transport.announce_emitted(packet) - - random_blob = packet.data[RNS.Identity.KEYSIZE//8+RNS.Identity.NAME_HASH_LENGTH//8:RNS.Identity.KEYSIZE//8+RNS.Identity.NAME_HASH_LENGTH//8+10] - random_blobs = [] - with Transport.inbound_announce_lock: - announced_destination_known = packet.destination_hash in Transport.path_table - if not announced_destination_known: - # If this destination is unknown in our table - # we should add it - should_add = True + # First, check that the announce is not for a destination + # local to this system, and that hops are less than the max + with Transport.destinations_map_lock: + local_and_hops_condition = (packet.hops < Transport.PATHFINDER_M+1) and (not packet.destination_hash in Transport.destinations_map) + if local_and_hops_condition: + announce_emitted = Transport.announce_emitted(packet) + + random_blob = packet.data[RNS.Identity.KEYSIZE//8+RNS.Identity.NAME_HASH_LENGTH//8:RNS.Identity.KEYSIZE//8+RNS.Identity.NAME_HASH_LENGTH//8+10] + random_blobs = [] + with Transport.inbound_announce_lock: + announced_destination_known = packet.destination_hash in Transport.path_table + if not announced_destination_known: + # If this destination is unknown in our table + # we should add it + should_add = True + + else: + random_blobs = Transport.path_table[packet.destination_hash][IDX_PT_RANDBLOBS] + current_gravity = Transport.path_table[packet.destination_hash][IDX_PT_RVCD_IF].gravity + announce_gravity = packet.receiving_interface.gravity if packet.receiving_interface != None else None + + # If we already have a path to the announced destination, + # but a more recently emitted announce arrives with a hop + # count equal to or less than the existing path, we will + # update our tables. + if packet.hops <= Transport.path_table[packet.destination_hash][IDX_PT_HOPS]: + path_timebase = Transport.timebase_from_random_blobs(random_blobs) + if not random_blob in random_blobs and announce_emitted > path_timebase: + Transport.mark_path_unknown_state(packet.destination_hash) + should_add = True + else: + # If the same announce is received later on an interface + # with higher gravity, allow updating the path table to + # use this interface instead. + if announce_emitted != path_timebase: should_add = False + elif announce_gravity == None or current_gravity == None: should_add = False + else: + if announce_gravity <= current_gravity: should_add = False + else: + RNS.log(f"Replacing path table entry for {RNS.prettyhexrep(packet.destination_hash)} with new announce due to higher gravity ({current_gravity}->{announce_gravity})", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + should_add = True else: - random_blobs = Transport.path_table[packet.destination_hash][IDX_PT_RANDBLOBS] - current_gravity = Transport.path_table[packet.destination_hash][IDX_PT_RVCD_IF].gravity - announce_gravity = packet.receiving_interface.gravity if packet.receiving_interface != None else None + # If an announce arrives with a larger hop + # count than we already have in the table, + # ignore it, unless the path is expired, or + # the emission timestamp is more recent. + now = time.time() + path_expires = Transport.path_table[packet.destination_hash][IDX_PT_EXPIRES] + + path_announce_emitted = 0 + for path_random_blob in random_blobs: + path_announce_emitted = max(path_announce_emitted, int.from_bytes(path_random_blob[5:10], "big")) + if path_announce_emitted >= announce_emitted: break - # If we already have a path to the announced destination, - # but a more recently emitted announce arrives with a hop - # count equal to or less than the existing path, we will - # update our tables. - if packet.hops <= Transport.path_table[packet.destination_hash][IDX_PT_HOPS]: - path_timebase = Transport.timebase_from_random_blobs(random_blobs) - if not random_blob in random_blobs and announce_emitted > path_timebase: + # If the path has expired, consider this + # announce for adding to the path table. + if (now >= path_expires): + # We check that the announce is + # different from ones we've already heard, + # to avoid loops in the network + if not random_blob in random_blobs: + RNS.log("Replacing path table entry for "+str(RNS.prettyhexrep(packet.destination_hash))+" with new announce due to expired path", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None Transport.mark_path_unknown_state(packet.destination_hash) should_add = True - else: - # If the same announce is received later on an interface - # with higher gravity, allow updating the path table to - # use this interface instead. - if announce_emitted != path_timebase: should_add = False - elif announce_gravity == None or current_gravity == None: should_add = False - else: - if announce_gravity <= current_gravity: should_add = False - else: - RNS.log(f"Replacing path table entry for {RNS.prettyhexrep(packet.destination_hash)} with new announce due to higher gravity ({current_gravity}->{announce_gravity})", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - should_add = True - else: - # If an announce arrives with a larger hop - # count than we already have in the table, - # ignore it, unless the path is expired, or - # the emission timestamp is more recent. - now = time.time() - path_expires = Transport.path_table[packet.destination_hash][IDX_PT_EXPIRES] - - path_announce_emitted = 0 - for path_random_blob in random_blobs: - path_announce_emitted = max(path_announce_emitted, int.from_bytes(path_random_blob[5:10], "big")) - if path_announce_emitted >= announce_emitted: break + else: should_add = False - # If the path has expired, consider this - # announce for adding to the path table. - if (now >= path_expires): - # We check that the announce is - # different from ones we've already heard, - # to avoid loops in the network + else: + # If the path is not expired, but the emission + # is more recent, and we haven't already heard + # this announce before, update the path table. + if (announce_emitted > path_announce_emitted): if not random_blob in random_blobs: - RNS.log("Replacing path table entry for "+str(RNS.prettyhexrep(packet.destination_hash))+" with new announce due to expired path", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + RNS.log("Replacing path table entry for "+str(RNS.prettyhexrep(packet.destination_hash))+" with new announce, since it was more recently emitted", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None Transport.mark_path_unknown_state(packet.destination_hash) should_add = True else: should_add = False + + # If we have already heard this announce before, + # but the path has been marked as unresponsive + # by a failed communications attempt or similar, + # allow updating the path table to this one. + elif announce_emitted == path_announce_emitted: + if Transport.path_is_unresponsive(packet.destination_hash): + RNS.log("Replacing path table entry for "+str(RNS.prettyhexrep(packet.destination_hash))+" with new announce, since previously tried path was unresponsive", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + should_add = True + else: should_add = False + + if should_add: + now = time.time() + is_from_local_client = Transport.from_local_client(packet) + + rate_blocked = False + if packet.context != RNS.Packet.PATH_RESPONSE and packet.receiving_interface.announce_rate_target != None: + with Transport.announce_rate_table_lock: + if not packet.destination_hash in Transport.announce_rate_table: + rate_entry = { "last": now, "rate_violations": 0, "blocked_until": 0, "timestamps": [now]} + Transport.announce_rate_table[packet.destination_hash] = rate_entry else: - # If the path is not expired, but the emission - # is more recent, and we haven't already heard - # this announce before, update the path table. - if (announce_emitted > path_announce_emitted): - if not random_blob in random_blobs: - RNS.log("Replacing path table entry for "+str(RNS.prettyhexrep(packet.destination_hash))+" with new announce, since it was more recently emitted", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - Transport.mark_path_unknown_state(packet.destination_hash) - should_add = True - else: should_add = False - - # If we have already heard this announce before, - # but the path has been marked as unresponsive - # by a failed communications attempt or similar, - # allow updating the path table to this one. - elif announce_emitted == path_announce_emitted: - if Transport.path_is_unresponsive(packet.destination_hash): - RNS.log("Replacing path table entry for "+str(RNS.prettyhexrep(packet.destination_hash))+" with new announce, since previously tried path was unresponsive", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - should_add = True - else: should_add = False + rate_entry = Transport.announce_rate_table[packet.destination_hash] + rate_entry["timestamps"].append(now) - if should_add: - now = time.time() - is_from_local_client = Transport.from_local_client(packet) + while len(rate_entry["timestamps"]) > Transport.MAX_RATE_TIMESTAMPS: + rate_entry["timestamps"].pop(0) - rate_blocked = False - if packet.context != RNS.Packet.PATH_RESPONSE and packet.receiving_interface.announce_rate_target != None: - with Transport.announce_rate_table_lock: - if not packet.destination_hash in Transport.announce_rate_table: - rate_entry = { "last": now, "rate_violations": 0, "blocked_until": 0, "timestamps": [now]} - Transport.announce_rate_table[packet.destination_hash] = rate_entry + current_rate = now - rate_entry["last"] - else: - rate_entry = Transport.announce_rate_table[packet.destination_hash] - rate_entry["timestamps"].append(now) + if now > rate_entry["blocked_until"]: + if current_rate < packet.receiving_interface.announce_rate_target: rate_entry["rate_violations"] += 1 + else: rate_entry["rate_violations"] = max(0, rate_entry["rate_violations"]-1) - while len(rate_entry["timestamps"]) > Transport.MAX_RATE_TIMESTAMPS: - rate_entry["timestamps"].pop(0) + if rate_entry["rate_violations"] > packet.receiving_interface.announce_rate_grace: + rate_target = packet.receiving_interface.announce_rate_target + rate_penalty = packet.receiving_interface.announce_rate_penalty + rate_entry["blocked_until"] = rate_entry["last"] + rate_target + rate_penalty + rate_blocked = True + else: + rate_entry["last"] = now - current_rate = now - rate_entry["last"] + else: rate_blocked = True - if now > rate_entry["blocked_until"]: - if current_rate < packet.receiving_interface.announce_rate_target: rate_entry["rate_violations"] += 1 - else: rate_entry["rate_violations"] = max(0, rate_entry["rate_violations"]-1) + retries = 0 + announce_hops = packet.hops + local_rebroadcasts = 0 + block_rebroadcasts = False + attached_interface = None + + retransmit_timeout = now + (RNS.rand() * Transport.PATHFINDER_RW) - if rate_entry["rate_violations"] > packet.receiving_interface.announce_rate_grace: - rate_target = packet.receiving_interface.announce_rate_target - rate_penalty = packet.receiving_interface.announce_rate_penalty - rate_entry["blocked_until"] = rate_entry["last"] + rate_target + rate_penalty - rate_blocked = True - else: - rate_entry["last"] = now + if hasattr(packet.receiving_interface, "mode") and packet.receiving_interface.mode == RNS.Interfaces.Interface.Interface.MODE_ACCESS_POINT: + expires = now + Transport.AP_PATH_TIME + elif hasattr(packet.receiving_interface, "mode") and packet.receiving_interface.mode == RNS.Interfaces.Interface.Interface.MODE_ROAMING: + expires = now + Transport.ROAMING_PATH_TIME + else: + expires = now + Transport.PATHFINDER_E + + if not random_blob in random_blobs: + random_blobs.append(random_blob) + random_blobs = random_blobs[-Transport.MAX_RANDOM_BLOBS:] - else: rate_blocked = True + if (RNS.Reticulum.transport_enabled() or is_from_local_client) and packet.context != RNS.Packet.PATH_RESPONSE: + # Insert announce into announce table for retransmission - retries = 0 - announce_hops = packet.hops - local_rebroadcasts = 0 - block_rebroadcasts = False - attached_interface = None - - retransmit_timeout = now + (RNS.rand() * Transport.PATHFINDER_RW) - - if hasattr(packet.receiving_interface, "mode") and packet.receiving_interface.mode == RNS.Interfaces.Interface.Interface.MODE_ACCESS_POINT: - expires = now + Transport.AP_PATH_TIME - elif hasattr(packet.receiving_interface, "mode") and packet.receiving_interface.mode == RNS.Interfaces.Interface.Interface.MODE_ROAMING: - expires = now + Transport.ROAMING_PATH_TIME + if rate_blocked: RNS.log("Blocking rebroadcast of announce from "+RNS.prettyhexrep(packet.destination_hash)+" due to excessive announce rate", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None else: - expires = now + Transport.PATHFINDER_E - - if not random_blob in random_blobs: - random_blobs.append(random_blob) - random_blobs = random_blobs[-Transport.MAX_RANDOM_BLOBS:] + if is_from_local_client: + # If the announce is from a local client, + # it is announced immediately, but only one time. + retransmit_timeout = now + retries = Transport.PATHFINDER_R - if (RNS.Reticulum.transport_enabled() or is_from_local_client) and packet.context != RNS.Packet.PATH_RESPONSE: - # Insert announce into announce table for retransmission + with Transport.announce_table_lock: + Transport.announce_table[packet.destination_hash] = [ + now, # 0: IDX_AT_TIMESTAMP + retransmit_timeout, # 1: IDX_AT_RTRNS_TMO + retries, # 2: IDX_AT_RETRIES + received_from, # 3: IDX_AT_RCVD_IF + announce_hops, # 4: IDX_AT_HOPS + packet, # 5: IDX_AT_PACKET + local_rebroadcasts, # 6: IDX_AT_LCL_RBRD + block_rebroadcasts, # 7: IDX_AT_BLCK_RBRD + attached_interface, # 8: IDX_AT_ATTCHD_IF + ] - if rate_blocked: RNS.log("Blocking rebroadcast of announce from "+RNS.prettyhexrep(packet.destination_hash)+" due to excessive announce rate", RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - else: - if is_from_local_client: - # If the announce is from a local client, - # it is announced immediately, but only one time. - retransmit_timeout = now - retries = Transport.PATHFINDER_R + elif is_from_local_client and packet.context == RNS.Packet.PATH_RESPONSE: + # If this is a path response from a local client, + # check if any external interfaces have pending + # path requests. + with Transport.pending_local_prs_lock: + if packet.destination_hash in Transport.pending_local_path_requests: + desiring_interface = Transport.pending_local_path_requests.pop(packet.destination_hash) + retransmit_timeout = now + retries = Transport.PATHFINDER_R with Transport.announce_table_lock: Transport.announce_table[packet.destination_hash] = [ - now, # 0: IDX_AT_TIMESTAMP - retransmit_timeout, # 1: IDX_AT_RTRNS_TMO - retries, # 2: IDX_AT_RETRIES - received_from, # 3: IDX_AT_RCVD_IF - announce_hops, # 4: IDX_AT_HOPS - packet, # 5: IDX_AT_PACKET - local_rebroadcasts, # 6: IDX_AT_LCL_RBRD - block_rebroadcasts, # 7: IDX_AT_BLCK_RBRD - attached_interface, # 8: IDX_AT_ATTCHD_IF + now, + retransmit_timeout, + retries, + received_from, + announce_hops, + packet, + local_rebroadcasts, + block_rebroadcasts, + attached_interface ] - elif is_from_local_client and packet.context == RNS.Packet.PATH_RESPONSE: - # If this is a path response from a local client, - # check if any external interfaces have pending - # path requests. - with Transport.pending_local_prs_lock: - if packet.destination_hash in Transport.pending_local_path_requests: - desiring_interface = Transport.pending_local_path_requests.pop(packet.destination_hash) - retransmit_timeout = now - retries = Transport.PATHFINDER_R + # If we have any local clients connected, we re- + # transmit the announce to them immediately + if (len(Transport.local_client_interfaces)): + announce_identity = RNS.Identity.recall(packet.destination_hash, _no_use=True) + announce_destination = RNS.Destination(announce_identity, RNS.Destination.OUT, RNS.Destination.SINGLE, "unknown", "unknown"); + announce_destination.hash = packet.destination_hash + announce_destination.hexhash = announce_destination.hash.hex() + announce_context = RNS.Packet.NONE + announce_data = packet.data - with Transport.announce_table_lock: - Transport.announce_table[packet.destination_hash] = [ - now, - retransmit_timeout, - retries, - received_from, - announce_hops, - packet, - local_rebroadcasts, - block_rebroadcasts, - attached_interface - ] + # TODO: Shouldn't the context be PATH_RESPONSE in the first case here? + if is_from_local_client and packet.context == RNS.Packet.PATH_RESPONSE: + for local_interface in Transport.local_client_interfaces: + if packet.receiving_interface != local_interface: + new_announce = RNS.Packet(announce_destination, announce_data, RNS.Packet.ANNOUNCE, # <-- This one? + context = announce_context, header_type = RNS.Packet.HEADER_2, + transport_type = Transport.TRANSPORT, transport_id = Transport.identity.hash, + attached_interface = local_interface, context_flag = packet.context_flag) + + new_announce.hops = packet.hops + new_announce.send() - # If we have any local clients connected, we re- - # transmit the announce to them immediately - if (len(Transport.local_client_interfaces)): - announce_identity = RNS.Identity.recall(packet.destination_hash, _no_use=True) - announce_destination = RNS.Destination(announce_identity, RNS.Destination.OUT, RNS.Destination.SINGLE, "unknown", "unknown"); - announce_destination.hash = packet.destination_hash - announce_destination.hexhash = announce_destination.hash.hex() - announce_context = RNS.Packet.NONE - announce_data = packet.data + else: + for local_interface in Transport.local_client_interfaces: + if packet.receiving_interface != local_interface: + new_announce = RNS.Packet(announce_destination, announce_data, RNS.Packet.ANNOUNCE, + context = announce_context, header_type = RNS.Packet.HEADER_2, + transport_type = Transport.TRANSPORT, transport_id = Transport.identity.hash, + attached_interface = local_interface, context_flag = packet.context_flag) - # TODO: Shouldn't the context be PATH_RESPONSE in the first case here? - if is_from_local_client and packet.context == RNS.Packet.PATH_RESPONSE: - for local_interface in Transport.local_client_interfaces: - if packet.receiving_interface != local_interface: - new_announce = RNS.Packet(announce_destination, announce_data, RNS.Packet.ANNOUNCE, # <-- This one? - context = announce_context, header_type = RNS.Packet.HEADER_2, - transport_type = Transport.TRANSPORT, transport_id = Transport.identity.hash, - attached_interface = local_interface, context_flag = packet.context_flag) - - new_announce.hops = packet.hops - new_announce.send() + new_announce.hops = packet.hops + new_announce.send() + # If we have any waiting discovery path requests + # for this destination, we retransmit to that + # interface immediately + if packet.destination_hash in Transport.discovery_path_requests: + pr_entry = Transport.discovery_path_requests[packet.destination_hash] + attached_interface = pr_entry["requesting_interface"] + + interface_str = " on "+str(attached_interface) + + RNS.log("Got matching announce, answering waiting discovery path request for "+RNS.prettyhexrep(packet.destination_hash)+interface_str, RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + announce_identity = RNS.Identity.recall(packet.destination_hash, _no_use=False) + announce_destination = RNS.Destination(announce_identity, RNS.Destination.OUT, RNS.Destination.SINGLE, "unknown", "unknown"); + announce_destination.hash = packet.destination_hash + announce_destination.hexhash = announce_destination.hash.hex() + announce_context = RNS.Packet.NONE + announce_data = packet.data + + new_announce = RNS.Packet(announce_destination, announce_data, RNS.Packet.ANNOUNCE, + context = RNS.Packet.PATH_RESPONSE, header_type = RNS.Packet.HEADER_2, + transport_type = Transport.TRANSPORT, transport_id = Transport.identity.hash, + attached_interface = attached_interface, context_flag = packet.context_flag) + + new_announce.hops = packet.hops + new_announce.send() + + if not Transport.owner.is_connected_to_shared_instance: Transport.cache(packet, force_cache=True, packet_type="announce") + path_table_entry = [now, received_from, announce_hops, expires, random_blobs, packet.receiving_interface, packet.packet_hash] + with Transport.path_table_lock: Transport.path_table[packet.destination_hash] = path_table_entry + Transport.mark_path_unknown_state(packet.destination_hash) + RNS.log("Destination "+RNS.prettyhexrep(packet.destination_hash)+" is now "+str(announce_hops)+" hops away via "+RNS.prettyhexrep(received_from)+" on "+str(packet.receiving_interface), RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None + if packet.destination_hash in Transport.path_requests: + RNS.Reticulum.get_instance()._used_destination_data(packet.destination_hash) + + # If the receiving interface is a tunnel, we add the + # announce to the tunnels table + if hasattr(packet.receiving_interface, "tunnel_id") and packet.receiving_interface.tunnel_id != None: + with Transport.tunnels_lock: + if not packet.receiving_interface.tunnel_id in Transport.tunnels: RNS.log(f"Tunnel ID for {packet.receiving_interface} was not found in tunnel table", RNS.LOG_WARNING) else: - for local_interface in Transport.local_client_interfaces: - if packet.receiving_interface != local_interface: - new_announce = RNS.Packet(announce_destination, announce_data, RNS.Packet.ANNOUNCE, - context = announce_context, header_type = RNS.Packet.HEADER_2, - transport_type = Transport.TRANSPORT, transport_id = Transport.identity.hash, - attached_interface = local_interface, context_flag = packet.context_flag) + tunnel_entry = Transport.tunnels[packet.receiving_interface.tunnel_id] + paths = tunnel_entry[IDX_TT_PATHS] + paths[packet.destination_hash] = [now, received_from, announce_hops, expires, random_blobs, None, packet.packet_hash] + expires = time.time() + Transport.TUNNEL_TIMEOUT + tunnel_entry[IDX_TT_EXPIRES] = expires + RNS.log("Path to "+RNS.prettyhexrep(packet.destination_hash)+" associated with tunnel "+RNS.prettyhexrep(packet.receiving_interface.tunnel_id), RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - new_announce.hops = packet.hops - new_announce.send() - - # If we have any waiting discovery path requests - # for this destination, we retransmit to that - # interface immediately - if packet.destination_hash in Transport.discovery_path_requests: - pr_entry = Transport.discovery_path_requests[packet.destination_hash] - attached_interface = pr_entry["requesting_interface"] - - interface_str = " on "+str(attached_interface) - - RNS.log("Got matching announce, answering waiting discovery path request for "+RNS.prettyhexrep(packet.destination_hash)+interface_str, RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - announce_identity = RNS.Identity.recall(packet.destination_hash, _no_use=False) - announce_destination = RNS.Destination(announce_identity, RNS.Destination.OUT, RNS.Destination.SINGLE, "unknown", "unknown"); - announce_destination.hash = packet.destination_hash - announce_destination.hexhash = announce_destination.hash.hex() - announce_context = RNS.Packet.NONE - announce_data = packet.data - - new_announce = RNS.Packet(announce_destination, announce_data, RNS.Packet.ANNOUNCE, - context = RNS.Packet.PATH_RESPONSE, header_type = RNS.Packet.HEADER_2, - transport_type = Transport.TRANSPORT, transport_id = Transport.identity.hash, - attached_interface = attached_interface, context_flag = packet.context_flag) - - new_announce.hops = packet.hops - new_announce.send() - - if not Transport.owner.is_connected_to_shared_instance: Transport.cache(packet, force_cache=True, packet_type="announce") - path_table_entry = [now, received_from, announce_hops, expires, random_blobs, packet.receiving_interface, packet.packet_hash] - with Transport.path_table_lock: Transport.path_table[packet.destination_hash] = path_table_entry - Transport.mark_path_unknown_state(packet.destination_hash) - RNS.log("Destination "+RNS.prettyhexrep(packet.destination_hash)+" is now "+str(announce_hops)+" hops away via "+RNS.prettyhexrep(received_from)+" on "+str(packet.receiving_interface), RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - if packet.destination_hash in Transport.path_requests: - RNS.Reticulum.get_instance()._used_destination_data(packet.destination_hash) - - # If the receiving interface is a tunnel, we add the - # announce to the tunnels table - if hasattr(packet.receiving_interface, "tunnel_id") and packet.receiving_interface.tunnel_id != None: - with Transport.tunnels_lock: - if not packet.receiving_interface.tunnel_id in Transport.tunnels: RNS.log(f"Tunnel ID for {packet.receiving_interface} was not found in tunnel table", RNS.LOG_WARNING) - else: - tunnel_entry = Transport.tunnels[packet.receiving_interface.tunnel_id] - paths = tunnel_entry[IDX_TT_PATHS] - paths[packet.destination_hash] = [now, received_from, announce_hops, expires, random_blobs, None, packet.packet_hash] - expires = time.time() + Transport.TUNNEL_TIMEOUT - tunnel_entry[IDX_TT_EXPIRES] = expires - RNS.log("Path to "+RNS.prettyhexrep(packet.destination_hash)+" associated with tunnel "+RNS.prettyhexrep(packet.receiving_interface.tunnel_id), RNS.LOG_PATHING) if RNS.sl(RNS.LOG_PATHING) else None - - # Call externally registered callbacks from apps - # wanting to know when an announce arrives - with Transport.announce_handler_lock: - for handler in Transport.announce_handlers: - try: - # Check that the announced destination matches - # the handlers aspect filter - execute_callback = False - announce_identity = RNS.Identity.recall(packet.destination_hash, _no_use=True) - - # If the handlers aspect filter is set to - # None, we execute the callback in all cases - if handler.aspect_filter == None: execute_callback = True - else: - handler_expected_hash = RNS.Destination.hash_from_name_and_identity(handler.aspect_filter, announce_identity) - if packet.destination_hash == handler_expected_hash: execute_callback = True - - # If this is a path response, check whether the - # handler wants to receive it. - if packet.context == RNS.Packet.PATH_RESPONSE: - if hasattr(handler, "receive_path_responses") and handler.receive_path_responses == True: pass - else: execute_callback = False - - if execute_callback: - if len(inspect.signature(handler.received_announce).parameters) == 3: - def job(handler=handler, packet=packet, announce_identity=announce_identity): - handler.received_announce(destination_hash=packet.destination_hash, - announced_identity=announce_identity, - app_data=RNS.Identity.recall_app_data(packet.destination_hash, _no_use=True)) - threading.Thread(target=job, daemon=True).start() - - elif len(inspect.signature(handler.received_announce).parameters) == 4: - def job(handler=handler, packet=packet, announce_identity=announce_identity): - handler.received_announce(destination_hash=packet.destination_hash, - announced_identity=announce_identity, - app_data=RNS.Identity.recall_app_data(packet.destination_hash, _no_use=True), - announce_packet_hash = packet.packet_hash) - threading.Thread(target=job, daemon=True).start() - - elif len(inspect.signature(handler.received_announce).parameters) == 5: - def job(handler=handler, packet=packet, announce_identity=announce_identity): - handler.received_announce(destination_hash=packet.destination_hash, - announced_identity=announce_identity, - app_data=RNS.Identity.recall_app_data(packet.destination_hash, _no_use=True), - announce_packet_hash = packet.packet_hash, - is_path_response = packet.context == RNS.Packet.PATH_RESPONSE) - threading.Thread(target=job, daemon=True).start() - - else: - raise TypeError("Invalid signature for announce handler callback") - - except Exception as e: - RNS.log("Error while processing external announce callback.", RNS.LOG_ERROR) - RNS.log("The contained exception was: "+str(e), RNS.LOG_ERROR) - RNS.trace_exception(e) - - # Handling for link requests to local destinations - elif packet.packet_type == RNS.Packet.LINKREQUEST and not link_request_handled: - if packet.transport_id == None or packet.transport_id == Transport.identity.hash: - destination = None - with Transport.destinations_map_lock: - if packet.destination_hash in Transport.destinations_map: - destination = Transport.destinations_map[packet.destination_hash] - - if destination and destination.type == packet.destination_type: - path_mtu = RNS.Link.mtu_from_lr_packet(packet) - mode = RNS.Link.mode_from_lr_packet(packet) - if packet.receiving_interface.AUTOCONFIGURE_MTU or packet.receiving_interface.FIXED_MTU: - nh_mtu = packet.receiving_interface.HW_MTU - else: - nh_mtu = RNS.Reticulum.MTU - - if path_mtu: - if packet.receiving_interface.HW_MTU == None: - path_mtu = None - packet.data = packet.data[:-RNS.Link.LINK_MTU_SIZE] - else: - if nh_mtu < path_mtu: + # Call externally registered callbacks from apps + # wanting to know when an announce arrives + with Transport.announce_handler_lock: + for handler in Transport.announce_handlers: try: - path_mtu = nh_mtu - clamped_mtu = RNS.Link.signalling_bytes(path_mtu, mode) - packet.data = packet.data[:-RNS.Link.LINK_MTU_SIZE]+clamped_mtu - except Exception as e: - RNS.log(f"Dropping link request packet to local destination. The contained exception was: {e}", RNS.LOG_WARNING) - return + # Check that the announced destination matches + # the handlers aspect filter + execute_callback = False + announce_identity = RNS.Identity.recall(packet.destination_hash, _no_use=True) - packet.destination = destination - destination.receive(packet) - - # Handling for local data packets - elif packet.packet_type == RNS.Packet.DATA: + # If the handlers aspect filter is set to + # None, we execute the callback in all cases + if handler.aspect_filter == None: execute_callback = True + else: + handler_expected_hash = RNS.Destination.hash_from_name_and_identity(handler.aspect_filter, announce_identity) + if packet.destination_hash == handler_expected_hash: execute_callback = True + + # If this is a path response, check whether the + # handler wants to receive it. + if packet.context == RNS.Packet.PATH_RESPONSE: + if hasattr(handler, "receive_path_responses") and handler.receive_path_responses == True: pass + else: execute_callback = False + + if execute_callback: + if len(inspect.signature(handler.received_announce).parameters) == 3: + def job(handler=handler, packet=packet, announce_identity=announce_identity): + handler.received_announce(destination_hash=packet.destination_hash, + announced_identity=announce_identity, + app_data=RNS.Identity.recall_app_data(packet.destination_hash, _no_use=True)) + threading.Thread(target=job, daemon=True).start() + + elif len(inspect.signature(handler.received_announce).parameters) == 4: + def job(handler=handler, packet=packet, announce_identity=announce_identity): + handler.received_announce(destination_hash=packet.destination_hash, + announced_identity=announce_identity, + app_data=RNS.Identity.recall_app_data(packet.destination_hash, _no_use=True), + announce_packet_hash = packet.packet_hash) + threading.Thread(target=job, daemon=True).start() + + elif len(inspect.signature(handler.received_announce).parameters) == 5: + def job(handler=handler, packet=packet, announce_identity=announce_identity): + handler.received_announce(destination_hash=packet.destination_hash, + announced_identity=announce_identity, + app_data=RNS.Identity.recall_app_data(packet.destination_hash, _no_use=True), + announce_packet_hash = packet.packet_hash, + is_path_response = packet.context == RNS.Packet.PATH_RESPONSE) + threading.Thread(target=job, daemon=True).start() + + else: + raise TypeError("Invalid signature for announce handler callback") + + except Exception as e: + RNS.log("Error while processing external announce callback.", RNS.LOG_ERROR) + RNS.log("The contained exception was: "+str(e), RNS.LOG_ERROR) + RNS.trace_exception(e) + + # Handling for link requests to local destinations + elif packet.packet_type == RNS.Packet.LINKREQUEST and not link_request_handled: + if packet.transport_id == None or packet.transport_id == Transport.identity.hash: + destination = None + with Transport.destinations_map_lock: + if packet.destination_hash in Transport.destinations_map: + destination = Transport.destinations_map[packet.destination_hash] + + if destination and destination.type == packet.destination_type: + path_mtu = RNS.Link.mtu_from_lr_packet(packet) + mode = RNS.Link.mode_from_lr_packet(packet) + if packet.receiving_interface.AUTOCONFIGURE_MTU or packet.receiving_interface.FIXED_MTU: + nh_mtu = packet.receiving_interface.HW_MTU + else: + nh_mtu = RNS.Reticulum.MTU + + if path_mtu: + if packet.receiving_interface.HW_MTU == None: + path_mtu = None + packet.data = packet.data[:-RNS.Link.LINK_MTU_SIZE] + else: + if nh_mtu < path_mtu: + try: + path_mtu = nh_mtu + clamped_mtu = RNS.Link.signalling_bytes(path_mtu, mode) + packet.data = packet.data[:-RNS.Link.LINK_MTU_SIZE]+clamped_mtu + except Exception as e: + RNS.log(f"Dropping link request packet to local destination. The contained exception was: {e}", RNS.LOG_WARNING) + return + + packet.destination = destination + destination.receive(packet) + + # Handling for local data packets + elif packet.packet_type == RNS.Packet.DATA: + if packet.destination_type == RNS.Destination.LINK: + with Transport.active_links_lock: + for link in Transport.active_links: + if link.link_id == packet.destination_hash: + if link.attached_interface == packet.receiving_interface: + packet.link = link + if packet.context == RNS.Packet.CACHE_REQUEST: + cached_packet = Transport.get_cached_packet(packet.data) + if cached_packet != None: + if not cached_packet.unpack(): return + RNS.Packet(destination=link, data=cached_packet.data, + packet_type=cached_packet.packet_type, context=cached_packet.context).send() + + else: link.receive(packet) + break + + else: + # In the strange and rare case that an interface + # is partly malfunctioning, and a link-associated + # packet is being received on an interface that + # has failed sending, and transport has failed over + # to another path, we remove this packet hash from + # the filter hashlist so the link can receive the + # packet when it finally arrives over another path. + while packet.packet_hash in Transport.packet_hashlist: + Transport.packet_hashlist.remove(packet.packet_hash) + else: + destination = None + with Transport.destinations_map_lock: + if packet.destination_hash in Transport.destinations_map: + destination = Transport.destinations_map[packet.destination_hash] + + if destination and destination.type == packet.destination_type: + packet.destination = destination + if destination.receive(packet): + if destination.proof_strategy == RNS.Destination.PROVE_ALL: packet.prove() + elif destination.proof_strategy == RNS.Destination.PROVE_APP: + if destination.callbacks.proof_requested: + try: + if destination.callbacks.proof_requested(packet): packet.prove() + except Exception as e: RNS.log("Error while executing proof request callback. The contained exception was: "+str(e), RNS.LOG_ERROR) + + # Handling for proofs and link-request proofs + elif packet.packet_type == RNS.Packet.PROOF: + if packet.context == RNS.Packet.LRPROOF: + # This is a link request proof, check if it needs to be transported + REBALANCE_LOGLEVEL = RNS.LOG_PATHING + if (RNS.Reticulum.transport_enabled() or for_local_client_link or from_local_client) and packet.destination_hash in Transport.link_table: + link_entry = Transport.link_table[packet.destination_hash] + if packet.hops != link_entry[IDX_LT_REM_HOPS] and Transport.ALLOW_LINK_PATH_REBALANCE: + if packet.receiving_interface == link_entry[IDX_LT_NH_IF]: + try: + if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2 or len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: + signalling_bytes = b"" + if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: + signalling_bytes = RNS.Link.signalling_bytes(RNS.Link.mtu_from_lp_packet(packet), RNS.Link.mode_from_lp_packet(packet)) + + peer_pub_bytes = packet.data[RNS.Identity.SIGLENGTH//8:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] + peer_identity = RNS.Identity.recall(link_entry[IDX_LT_DSTHASH], _no_use=True) + peer_sig_pub_bytes = peer_identity.get_public_key()[RNS.Link.ECPUBSIZE//2:RNS.Link.ECPUBSIZE] + + signed_data = packet.destination_hash+peer_pub_bytes+peer_sig_pub_bytes+signalling_bytes + signature = packet.data[:RNS.Identity.SIGLENGTH//8] + link_destination = link_entry[IDX_LT_DSTHASH] + + if peer_identity.validate(signature, signed_data) and not link_entry[IDX_LT_VALIDATED]: + RNS.log(f"Re-balancing path to {RNS.prettyhexrep(link_destination)} from link-request proof ({link_entry[IDX_LT_REM_HOPS]}->{packet.hops})", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + link_entry[IDX_LT_REM_HOPS] = packet.hops + with Transport.path_table_lock: + if link_destination in Transport.path_table: + path_entry = Transport.path_table[link_destination] + path_entry[IDX_PT_HOPS] = packet.hops + + elif not link_entry[IDX_LT_VALIDATED]: RNS.log(f"Aborting link request proof path re-balancing for {RNS.prettyhexrep(link_destination)} on link {RNS.prettyhexrep(packet.destination_hash)} due to invalid signature", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + else: pass + + except Exception as e: RNS.log(f"Error while re-balancing path from link request proof. The contained exception was: {e}", RNS.LOG_ERROR) if RNS.sl(RNS.LOG_ERROR) else None # TODO: Drop to DEBUG at some point + + if packet.hops == link_entry[IDX_LT_REM_HOPS]: + if packet.receiving_interface == link_entry[IDX_LT_NH_IF]: + try: + if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2 or len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: + signalling_bytes = b"" + if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: + signalling_bytes = RNS.Link.signalling_bytes(RNS.Link.mtu_from_lp_packet(packet), RNS.Link.mode_from_lp_packet(packet)) + + peer_pub_bytes = packet.data[RNS.Identity.SIGLENGTH//8:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] + peer_identity = RNS.Identity.recall(link_entry[IDX_LT_DSTHASH], _no_use=True) + peer_sig_pub_bytes = peer_identity.get_public_key()[RNS.Link.ECPUBSIZE//2:RNS.Link.ECPUBSIZE] + + signed_data = packet.destination_hash+peer_pub_bytes+peer_sig_pub_bytes+signalling_bytes + signature = packet.data[:RNS.Identity.SIGLENGTH//8] + + if peer_identity.validate(signature, signed_data): + RNS.log("Link request proof validated for transport via "+str(link_entry[IDX_LT_RCVD_IF]), RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None + new_raw = packet.raw[0:1] + new_raw += struct.pack("!B", packet.hops if not from_local_client or instance_local_link or Transport.local_hops_delta == 0 else Transport.local_hops_delta) + new_raw += packet.raw[2:] + Transport.link_table[packet.destination_hash][IDX_LT_VALIDATED] = True + Transport.transmit(link_entry[IDX_LT_RCVD_IF], new_raw) + if not Transport.owner.is_connected_to_shared_instance: + RNS.Identity._used_destination_data(link_entry[IDX_LT_DSTHASH]) + + else: RNS.log("Invalid link request proof in transport for link "+RNS.prettyhexrep(packet.destination_hash)+", dropping proof.", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + except Exception as e: RNS.log("Could not transport link request proof. The contained exception was: "+str(e), RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + else: RNS.log("Link request proof received on wrong interface, not transporting it.", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + else: RNS.log(f"Received link request proof with hop mismatch ({packet.hops}/{link_entry[IDX_LT_REM_HOPS]}:{link_entry[IDX_LT_NH_IF]}->{link_entry[IDX_LT_RCVD_IF]}), not transporting it", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + + else: + # Check if we can deliver it to a local pending link + pending_link = None + with Transport.pending_links_lock: + for link in Transport.pending_links: + if link.link_id == packet.destination_hash: + if packet.hops != link.expected_hops and link.status == RNS.Link.PENDING and Transport.ALLOW_LINK_PATH_REBALANCE: + RNS.log(f"Unbalanced link path ({packet.hops}/{link.expected_hops}) detected on link {link}, validating signature for re-balancing...", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + try: + if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2 or len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: + packet_data = packet.data + signalling_bytes = b"" + confirmed_mtu = None + mode = RNS.Link.mode_from_lp_packet(packet) + if mode != link.mode: raise TypeError(f"Invalid link mode {mode} in link request proof") + if len(packet_data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: + confirmed_mtu = RNS.Link.mtu_from_lp_packet(packet) + signalling_bytes = RNS.Link.signalling_bytes(confirmed_mtu, mode) + packet_data = packet_data[:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] + + peer_pub_bytes = packet_data[RNS.Identity.SIGLENGTH//8:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] + peer_sig_pub_bytes = link.destination.identity.get_public_key()[RNS.Link.ECPUBSIZE//2:RNS.Link.ECPUBSIZE] + + signed_data = link.link_id+peer_pub_bytes+peer_sig_pub_bytes+signalling_bytes + signature = packet_data[:RNS.Identity.SIGLENGTH//8] + + if link.destination.identity.validate(signature, signed_data): + with Transport.path_table_lock: + if not link.rebalanced: + RNS.log(f"Re-balancing path to {RNS.prettyhexrep(link.destination.hash)} at link terminus ({link.expected_hops}->{packet.hops})", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + link.rebalanced = time.time() + link.expected_hops = packet.hops + if link.destination.hash in Transport.path_table: + path_entry = Transport.path_table[link.destination.hash] + path_entry[IDX_PT_HOPS] = packet.hops + RNS.log(f"Path table re-balanced for {RNS.prettyhexrep(link.destination.hash)}", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + + else: RNS.log(f"Aborting path re-balancing at link terminus for {RNS.prettyhexrep(link.destination.hash)} on link {link} due to invalid signature", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + except Exception as e: RNS.log("Error while validating link request proof for path re-balancing at link terminus. The contained exception was: "+str(e), REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None + + if packet.hops == link.expected_hops: + # Add this packet to the filter hashlist if we + # have determined that it's actually destined + # for this system, and then validate the proof + Transport.add_packet_hash(packet.packet_hash) + pending_link = link + break + + if pending_link: pending_link.validate_proof(packet) + + elif packet.context == RNS.Packet.RESOURCE_PRF: + with Transport.active_links_lock: + for link in Transport.active_links: + if link.link_id == packet.destination_hash: + link.receive(packet) + break + else: if packet.destination_type == RNS.Destination.LINK: with Transport.active_links_lock: for link in Transport.active_links: if link.link_id == packet.destination_hash: - if link.attached_interface == packet.receiving_interface: - packet.link = link - if packet.context == RNS.Packet.CACHE_REQUEST: - cached_packet = Transport.get_cached_packet(packet.data) - if cached_packet != None: - if not cached_packet.unpack(): return - RNS.Packet(destination=link, data=cached_packet.data, - packet_type=cached_packet.packet_type, context=cached_packet.context).send() - - else: link.receive(packet) - break - - else: - # In the strange and rare case that an interface - # is partly malfunctioning, and a link-associated - # packet is being received on an interface that - # has failed sending, and transport has failed over - # to another path, we remove this packet hash from - # the filter hashlist so the link can receive the - # packet when it finally arrives over another path. - while packet.packet_hash in Transport.packet_hashlist: - Transport.packet_hashlist.remove(packet.packet_hash) - else: - destination = None - with Transport.destinations_map_lock: - if packet.destination_hash in Transport.destinations_map: - destination = Transport.destinations_map[packet.destination_hash] - - if destination and destination.type == packet.destination_type: - packet.destination = destination - if destination.receive(packet): - if destination.proof_strategy == RNS.Destination.PROVE_ALL: packet.prove() - elif destination.proof_strategy == RNS.Destination.PROVE_APP: - if destination.callbacks.proof_requested: - try: - if destination.callbacks.proof_requested(packet): packet.prove() - except Exception as e: RNS.log("Error while executing proof request callback. The contained exception was: "+str(e), RNS.LOG_ERROR) - - # Handling for proofs and link-request proofs - elif packet.packet_type == RNS.Packet.PROOF: - if packet.context == RNS.Packet.LRPROOF: - # This is a link request proof, check if it needs to be transported - REBALANCE_LOGLEVEL = RNS.LOG_PATHING - if (RNS.Reticulum.transport_enabled() or for_local_client_link or from_local_client) and packet.destination_hash in Transport.link_table: - link_entry = Transport.link_table[packet.destination_hash] - if packet.hops != link_entry[IDX_LT_REM_HOPS] and Transport.ALLOW_LINK_PATH_REBALANCE: - if packet.receiving_interface == link_entry[IDX_LT_NH_IF]: - try: - if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2 or len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: - signalling_bytes = b"" - if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: - signalling_bytes = RNS.Link.signalling_bytes(RNS.Link.mtu_from_lp_packet(packet), RNS.Link.mode_from_lp_packet(packet)) - - peer_pub_bytes = packet.data[RNS.Identity.SIGLENGTH//8:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] - peer_identity = RNS.Identity.recall(link_entry[IDX_LT_DSTHASH], _no_use=True) - peer_sig_pub_bytes = peer_identity.get_public_key()[RNS.Link.ECPUBSIZE//2:RNS.Link.ECPUBSIZE] - - signed_data = packet.destination_hash+peer_pub_bytes+peer_sig_pub_bytes+signalling_bytes - signature = packet.data[:RNS.Identity.SIGLENGTH//8] - link_destination = link_entry[IDX_LT_DSTHASH] - - if peer_identity.validate(signature, signed_data) and not link_entry[IDX_LT_VALIDATED]: - RNS.log(f"Re-balancing path to {RNS.prettyhexrep(link_destination)} from link-request proof ({link_entry[IDX_LT_REM_HOPS]}->{packet.hops})", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - link_entry[IDX_LT_REM_HOPS] = packet.hops - with Transport.path_table_lock: - if link_destination in Transport.path_table: - path_entry = Transport.path_table[link_destination] - path_entry[IDX_PT_HOPS] = packet.hops - - elif not link_entry[IDX_LT_VALIDATED]: RNS.log(f"Aborting link request proof path re-balancing for {RNS.prettyhexrep(link_destination)} on link {RNS.prettyhexrep(packet.destination_hash)} due to invalid signature", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - else: pass - - except Exception as e: RNS.log(f"Error while re-balancing path from link request proof. The contained exception was: {e}", RNS.LOG_ERROR) if RNS.sl(RNS.LOG_ERROR) else None # TODO: Drop to DEBUG at some point - - if packet.hops == link_entry[IDX_LT_REM_HOPS]: - if packet.receiving_interface == link_entry[IDX_LT_NH_IF]: - try: - if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2 or len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: - signalling_bytes = b"" - if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: - signalling_bytes = RNS.Link.signalling_bytes(RNS.Link.mtu_from_lp_packet(packet), RNS.Link.mode_from_lp_packet(packet)) - - peer_pub_bytes = packet.data[RNS.Identity.SIGLENGTH//8:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] - peer_identity = RNS.Identity.recall(link_entry[IDX_LT_DSTHASH], _no_use=True) - peer_sig_pub_bytes = peer_identity.get_public_key()[RNS.Link.ECPUBSIZE//2:RNS.Link.ECPUBSIZE] - - signed_data = packet.destination_hash+peer_pub_bytes+peer_sig_pub_bytes+signalling_bytes - signature = packet.data[:RNS.Identity.SIGLENGTH//8] - - if peer_identity.validate(signature, signed_data): - RNS.log("Link request proof validated for transport via "+str(link_entry[IDX_LT_RCVD_IF]), RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None - new_raw = packet.raw[0:1] - new_raw += struct.pack("!B", packet.hops if not from_local_client or instance_local_link or Transport.local_hops_delta == 0 else Transport.local_hops_delta) - new_raw += packet.raw[2:] - Transport.link_table[packet.destination_hash][IDX_LT_VALIDATED] = True - Transport.transmit(link_entry[IDX_LT_RCVD_IF], new_raw) - if not Transport.owner.is_connected_to_shared_instance: - RNS.Identity._used_destination_data(link_entry[IDX_LT_DSTHASH]) - - else: RNS.log("Invalid link request proof in transport for link "+RNS.prettyhexrep(packet.destination_hash)+", dropping proof.", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None - except Exception as e: RNS.log("Could not transport link request proof. The contained exception was: "+str(e), RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None - else: RNS.log("Link request proof received on wrong interface, not transporting it.", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None - else: RNS.log(f"Received link request proof with hop mismatch ({packet.hops}/{link_entry[IDX_LT_REM_HOPS]}:{link_entry[IDX_LT_NH_IF]}->{link_entry[IDX_LT_RCVD_IF]}), not transporting it", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - - else: - # Check if we can deliver it to a local pending link - pending_link = None - with Transport.pending_links_lock: - for link in Transport.pending_links: - if link.link_id == packet.destination_hash: - if packet.hops != link.expected_hops and link.status == RNS.Link.PENDING and Transport.ALLOW_LINK_PATH_REBALANCE: - RNS.log(f"Unbalanced link path ({packet.hops}/{link.expected_hops}) detected on link {link}, validating signature for re-balancing...", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - try: - if len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2 or len(packet.data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: - packet_data = packet.data - signalling_bytes = b"" - confirmed_mtu = None - mode = RNS.Link.mode_from_lp_packet(packet) - if mode != link.mode: raise TypeError(f"Invalid link mode {mode} in link request proof") - if len(packet_data) == RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2+RNS.Link.LINK_MTU_SIZE: - confirmed_mtu = RNS.Link.mtu_from_lp_packet(packet) - signalling_bytes = RNS.Link.signalling_bytes(confirmed_mtu, mode) - packet_data = packet_data[:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] - - peer_pub_bytes = packet_data[RNS.Identity.SIGLENGTH//8:RNS.Identity.SIGLENGTH//8+RNS.Link.ECPUBSIZE//2] - peer_sig_pub_bytes = link.destination.identity.get_public_key()[RNS.Link.ECPUBSIZE//2:RNS.Link.ECPUBSIZE] - - signed_data = link.link_id+peer_pub_bytes+peer_sig_pub_bytes+signalling_bytes - signature = packet_data[:RNS.Identity.SIGLENGTH//8] - - if link.destination.identity.validate(signature, signed_data): - with Transport.path_table_lock: - if not link.rebalanced: - RNS.log(f"Re-balancing path to {RNS.prettyhexrep(link.destination.hash)} at link terminus ({link.expected_hops}->{packet.hops})", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - link.rebalanced = time.time() - link.expected_hops = packet.hops - if link.destination.hash in Transport.path_table: - path_entry = Transport.path_table[link.destination.hash] - path_entry[IDX_PT_HOPS] = packet.hops - RNS.log(f"Path table re-balanced for {RNS.prettyhexrep(link.destination.hash)}", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - - else: RNS.log(f"Aborting path re-balancing at link terminus for {RNS.prettyhexrep(link.destination.hash)} on link {link} due to invalid signature", REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - except Exception as e: RNS.log("Error while validating link request proof for path re-balancing at link terminus. The contained exception was: "+str(e), REBALANCE_LOGLEVEL) if RNS.sl(REBALANCE_LOGLEVEL) else None - - if packet.hops == link.expected_hops: - # Add this packet to the filter hashlist if we - # have determined that it's actually destined - # for this system, and then validate the proof - Transport.add_packet_hash(packet.packet_hash) - pending_link = link - break - - if pending_link: pending_link.validate_proof(packet) - - elif packet.context == RNS.Packet.RESOURCE_PRF: - with Transport.active_links_lock: - for link in Transport.active_links: - if link.link_id == packet.destination_hash: - link.receive(packet) + packet.link = link break - else: - if packet.destination_type == RNS.Destination.LINK: - with Transport.active_links_lock: - for link in Transport.active_links: - if link.link_id == packet.destination_hash: - packet.link = link - break - - if len(packet.data) == RNS.PacketReceipt.EXPL_LENGTH: proof_hash = packet.data[:RNS.Identity.HASHLENGTH//8] - else: proof_hash = None + + if len(packet.data) == RNS.PacketReceipt.EXPL_LENGTH: proof_hash = packet.data[:RNS.Identity.HASHLENGTH//8] + else: proof_hash = None - # Check if this proof needs to be transported - if (RNS.Reticulum.transport_enabled() or from_local_client or proof_for_local_client) and packet.destination_hash in Transport.reverse_table: - with Transport.reverse_table_lock: reverse_entry = Transport.reverse_table.pop(packet.destination_hash) - if packet.receiving_interface == reverse_entry[IDX_RT_OUTB_IF]: - RNS.log("Proof received on correct interface, transporting it via "+str(reverse_entry[IDX_RT_RCVD_IF]), RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None - new_raw = packet.raw[0:1] - new_raw += struct.pack("!B", packet.hops if not from_local_client or proof_for_local_client or Transport.local_hops_delta == 0 else Transport.local_hops_delta) - new_raw += packet.raw[2:] - Transport.transmit(reverse_entry[IDX_RT_RCVD_IF], new_raw) - else: - RNS.log("Proof received on wrong interface, not transporting it.", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None + # Check if this proof needs to be transported + if (RNS.Reticulum.transport_enabled() or from_local_client or proof_for_local_client) and packet.destination_hash in Transport.reverse_table: + with Transport.reverse_table_lock: reverse_entry = Transport.reverse_table.pop(packet.destination_hash) + if packet.receiving_interface == reverse_entry[IDX_RT_OUTB_IF]: + RNS.log("Proof received on correct interface, transporting it via "+str(reverse_entry[IDX_RT_RCVD_IF]), RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None + new_raw = packet.raw[0:1] + new_raw += struct.pack("!B", packet.hops if not from_local_client or proof_for_local_client or Transport.local_hops_delta == 0 else Transport.local_hops_delta) + new_raw += packet.raw[2:] + Transport.transmit(reverse_entry[IDX_RT_RCVD_IF], new_raw) + else: + RNS.log("Proof received on wrong interface, not transporting it.", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None - # TODO: Lock contention could be improved here - # by a copy(), instead of running the list - # comprehension under lock. - with Transport.receipts_lock: - if proof_hash != None: - # Only test validation if hash matches - candidate_receipts = [r for r in Transport.receipts if r.hash == proof_hash] - else: - # In case of an implicit proof, we have - # to check every single outstanding receipt - candidate_receipts = Transport.receipts.copy() + # TODO: Lock contention could be improved here + # by a copy(), instead of running the list + # comprehension under lock. + with Transport.receipts_lock: + if proof_hash != None: + # Only test validation if hash matches + candidate_receipts = [r for r in Transport.receipts if r.hash == proof_hash] + else: + # In case of an implicit proof, we have + # to check every single outstanding receipt + candidate_receipts = Transport.receipts.copy() - for receipt in candidate_receipts: - if receipt.status != RNS.PacketReceipt.SENT: continue - if receipt.validate_proof_packet(packet): - with Transport.receipts_lock: - if receipt in Transport.receipts: Transport.receipts.remove(receipt) + for receipt in candidate_receipts: + if receipt.status != RNS.PacketReceipt.SENT: continue + if receipt.validate_proof_packet(packet): + with Transport.receipts_lock: + if receipt in Transport.receipts: Transport.receipts.remove(receipt) @staticmethod def synthesize_tunnel(interface):