diff --git a/easytier-core/src/peers/relay_peer_map.rs b/easytier-core/src/peers/relay_peer_map.rs index f3249cab..c6100376 100644 --- a/easytier-core/src/peers/relay_peer_map.rs +++ b/easytier-core/src/peers/relay_peer_map.rs @@ -229,6 +229,15 @@ impl RelayPeerMap { dst_peer_id: PeerId, policy: NextHopPolicy, ) -> Result<(), Error> { + // Forwarded packets belong to the original sender's end-to-end session. + // Do not establish another session or encrypt them at this hop. + if msg + .peer_manager_header() + .is_some_and(|hdr| hdr.from_peer_id.get() != self.my_peer_id) + { + return self.send_via_next_hop(msg, dst_peer_id, policy).await; + } + let now = Instant::now(); self.states.entry(dst_peer_id).or_default().last_active_at = now; @@ -865,6 +874,74 @@ mod tests { (Arc::new(context), public.as_bytes().to_vec()) } + #[derive(Default)] + struct ForwardedPacketTransport { + packets: StdMutex>, + } + + #[async_trait::async_trait] + impl RelayRouteTransport for ForwardedPacketTransport { + async fn get_route_peer_info(&self, _peer_id: PeerId) -> Option { + None + } + + async fn send_msg_to_next_hop( + &self, + msg: ZCPacket, + dst_peer_id: PeerId, + policy: NextHopPolicy, + ) -> Result<(), Error> { + self.packets + .lock() + .unwrap() + .push((msg, dst_peer_id, policy)); + Ok(()) + } + } + + #[tokio::test] + async fn secure_relay_forwards_remote_packets_without_changing_payload() { + let (context, _) = relay_test_context(2); + let transport = Arc::new(ForwardedPacketTransport::default()); + let relay = RelayPeerMap::new( + transport.clone(), + context, + 2, + Arc::new(PeerSessionStore::new()), + ); + + for packet_type in [ + PacketType::RelayHandshake, + PacketType::RelayHandshakeAck, + PacketType::Data, + ] { + let mut packet = ZCPacket::new_with_payload(b"end-to-end payload"); + packet.fill_peer_manager_hdr(1, 3, packet_type as u8); + if matches!(packet_type, PacketType::Data) { + packet + .mut_peer_manager_header() + .unwrap() + .set_encrypted(true); + } + let expected = packet.clone(); + relay + .send_msg(packet, 3, NextHopPolicy::LeastCost) + .await + .unwrap(); + + let mut packets = transport.packets.lock().unwrap(); + assert_eq!( + packets.len(), + 1, + "forwarded packet was held for a local handshake" + ); + let (forwarded, destination, policy) = packets.pop().unwrap(); + assert_eq!(forwarded.into_bytes(), expected.into_bytes()); + assert_eq!(destination, 3); + assert!(matches!(policy, NextHopPolicy::LeastCost)); + } + } + #[tokio::test] async fn latency_first_handshake_returns_ack_over_least_cost_route() { let (context_a, pubkey_a) = relay_test_context(1); diff --git a/easytier/src/tests/credential_tests.rs b/easytier/src/tests/credential_tests.rs index 867d0820..21c09a89 100644 --- a/easytier/src/tests/credential_tests.rs +++ b/easytier/src/tests/credential_tests.rs @@ -2527,9 +2527,12 @@ async fn credential_admin_shared_admin_credential_connectivity( prepare_credential_network(); let process_runtime = CoreProcessRuntime::new(); + // Keep traffic on the shared/admin relay path; direct P2P would bypass it. + // 10.1.1.1 let admin_a_config = create_admin_config("admin_a", Some("ns_adm"), "10.144.144.1", "fd00::1/64"); + disable_p2p(&admin_a_config); let mut admin_a_inst = Instance::new_with_process_runtime(admin_a_config, process_runtime.clone()); admin_a_inst.run().await.unwrap(); @@ -2537,6 +2540,7 @@ async fn credential_admin_shared_admin_credential_connectivity( // 10.1.1.2 let shared_b_config = create_shared_config("shared_b", Some("ns_c1"), "10.144.144.2", "fd00::2/64"); + disable_p2p(&shared_b_config); let mut shared_b_inst = Instance::new_with_process_runtime(shared_b_config, process_runtime.clone()); shared_b_inst.run().await.unwrap(); @@ -2544,6 +2548,7 @@ async fn credential_admin_shared_admin_credential_connectivity( // 10.1.1.4 let admin_c_config = create_admin_config("admin_c", Some("ns_c3"), "10.144.144.4", "fd00::4/64"); + disable_p2p(&admin_c_config); let mut admin_c_inst = Instance::new_with_process_runtime(admin_c_config, process_runtime.clone()); admin_c_inst.run().await.unwrap(); @@ -2585,6 +2590,7 @@ async fn credential_admin_shared_admin_credential_connectivity( .get_global_ctx() .issue_event(GlobalCtxEvent::CredentialChanged); + disable_p2p(&cred_d_config); let mut cred_d_inst = Instance::new_with_process_runtime(cred_d_config, process_runtime.clone()); cred_d_inst.run().await.unwrap();