Added hash map lookups for pending and active links
parent
7311dc8544
commit
629e4fde2d
|
|
@ -164,6 +164,8 @@ class Transport:
|
|||
destinations_map = {} # Destination hash map of active destinations
|
||||
pending_links = [] # Links that are being established
|
||||
active_links = [] # Links that are active
|
||||
pending_links_map = {} # Link ID hash map of pending links
|
||||
active_links_map = {} # Link ID hash map of active links
|
||||
packet_hashlist = set() # A list of packet hashes for duplicate detection
|
||||
packet_hashlist_prev = set()
|
||||
receipts = [] # Receipts of all outgoing packets for proof processing
|
||||
|
|
@ -714,14 +716,18 @@ class Transport:
|
|||
|
||||
closed_pending_links.append(link)
|
||||
|
||||
for closed_link in closed_pending_links: Transport.pending_links.remove(closed_link)
|
||||
for closed_link in closed_pending_links:
|
||||
Transport.pending_links.remove(closed_link)
|
||||
Transport.pending_links_map.pop(closed_link.link_id, None)
|
||||
|
||||
with Transport.active_links_lock:
|
||||
closed_links = []
|
||||
for link in Transport.active_links:
|
||||
if link.status == RNS.Link.CLOSED: closed_links.append(link)
|
||||
|
||||
for closed_link in closed_links: Transport.active_links.remove(closed_link)
|
||||
for closed_link in closed_links:
|
||||
Transport.active_links.remove(closed_link)
|
||||
Transport.active_links_map.pop(closed_link.link_id, None)
|
||||
|
||||
Transport.links_last_checked = time.time()
|
||||
|
||||
|
|
@ -2572,32 +2578,29 @@ class Transport:
|
|||
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)
|
||||
while packet.packet_hash in Transport.packet_hashlist_prev:
|
||||
Transport.packet_hashlist_prev.remove(packet.packet_hash)
|
||||
link = Transport.active_links_map.get(packet.destination_hash)
|
||||
if link != None:
|
||||
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)
|
||||
|
||||
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)
|
||||
while packet.packet_hash in Transport.packet_hashlist_prev: Transport.packet_hashlist_prev.remove(packet.packet_hash)
|
||||
else:
|
||||
destination = None
|
||||
with Transport.destinations_map_lock:
|
||||
|
|
@ -2685,8 +2688,9 @@ class Transport:
|
|||
# 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:
|
||||
link = Transport.pending_links_map.get(packet.destination_hash)
|
||||
if link != None:
|
||||
# TODO: Cleanup indentation
|
||||
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:
|
||||
|
|
@ -2726,23 +2730,16 @@ class Transport:
|
|||
# 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
|
||||
link = Transport.active_links_map.get(packet.destination_hash)
|
||||
if link != None: link.receive(packet)
|
||||
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
|
||||
link = Transport.active_links_map.get(packet.destination_hash)
|
||||
if link != None: packet.link = link
|
||||
|
||||
if len(packet.data) == RNS.PacketReceipt.EXPL_LENGTH: proof_hash = packet.data[:RNS.Identity.HASHLENGTH//8]
|
||||
else: proof_hash = None
|
||||
|
|
@ -2945,9 +2942,13 @@ class Transport:
|
|||
def register_link(link):
|
||||
RNS.log("Registering link "+str(link), RNS.LOG_EXTREME) if RNS.sl(RNS.LOG_EXTREME) else None
|
||||
if link.initiator:
|
||||
with Transport.pending_links_lock: Transport.pending_links.append(link)
|
||||
with Transport.pending_links_lock:
|
||||
Transport.pending_links.append(link)
|
||||
Transport.pending_links_map[link.link_id] = link
|
||||
else:
|
||||
with Transport.active_links_lock: Transport.active_links.append(link)
|
||||
with Transport.active_links_lock:
|
||||
Transport.active_links.append(link)
|
||||
Transport.active_links_map[link.link_id] = link
|
||||
|
||||
@staticmethod
|
||||
def activate_link(link):
|
||||
|
|
@ -2956,10 +2957,14 @@ class Transport:
|
|||
if link in Transport.pending_links:
|
||||
if link.status != RNS.Link.ACTIVE: raise IOError("Invalid link state for link activation: "+str(link.status))
|
||||
Transport.pending_links.remove(link)
|
||||
with Transport.active_links_lock: Transport.active_links.append(link)
|
||||
Transport.pending_links_map.pop(link.link_id, None)
|
||||
with Transport.active_links_lock:
|
||||
Transport.active_links.append(link)
|
||||
Transport.active_links_map[link.link_id] = link
|
||||
|
||||
link.status = RNS.Link.ACTIVE
|
||||
else:
|
||||
RNS.log("Attempted to activate a link that was not in the pending table", RNS.LOG_ERROR)
|
||||
|
||||
else: RNS.log("Attempted to activate a link that was not in the pending table", RNS.LOG_ERROR)
|
||||
|
||||
@staticmethod
|
||||
def register_announce_handler(handler):
|
||||
|
|
|
|||
Loading…
Reference in New Issue