Skip to content

Commit 6a7c9c6

Browse files
authored
Merge pull request #75 from block65/fix/relay-mode-and-peer-bugs
fix(relay): correct data plane wiring and peer management bugs
2 parents e11dece + 7a3b309 commit 6a7c9c6

8 files changed

Lines changed: 173 additions & 124 deletions

File tree

crates/core/src/client/client.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,8 @@ pub struct ConnectResult<T: wallhack_transport::Transport + ?Sized> {
4444
control_tx: mpsc::Sender<ControlMessage>,
4545
/// Receiver for the server's `Handshake` (delivered via the control loop).
4646
peer_handshake_rx: Option<oneshot::Receiver<Handshake>>,
47+
/// Pong-derived latency measurements from the control loop (milliseconds).
48+
latency_rx: Option<mpsc::Receiver<f64>>,
4749
}
4850

4951
impl<T: wallhack_transport::Transport + ?Sized> ConnectResult<T> {
@@ -55,6 +57,7 @@ impl<T: wallhack_transport::Transport + ?Sized> ConnectResult<T> {
5557
tasks: ConnectionTasks,
5658
control_tx: mpsc::Sender<ControlMessage>,
5759
peer_handshake_rx: Option<oneshot::Receiver<Handshake>>,
60+
latency_rx: Option<mpsc::Receiver<f64>>,
5861
) -> Self {
5962
Self {
6063
channels,
@@ -63,6 +66,7 @@ impl<T: wallhack_transport::Transport + ?Sized> ConnectResult<T> {
6366
transport,
6467
control_tx,
6568
peer_handshake_rx,
69+
latency_rx,
6670
}
6771
}
6872

@@ -106,6 +110,8 @@ pub struct ErasedConnectResult {
106110
pub control_tx: mpsc::Sender<ControlMessage>,
107111
pub peer_handshake_rx: Option<oneshot::Receiver<Handshake>>,
108112
pub peer_addr: String,
113+
/// Pong-derived latency measurements from the control loop (milliseconds).
114+
pub latency_rx: Option<mpsc::Receiver<f64>>,
109115
}
110116

111117
impl<T> ConnectResult<T>
@@ -129,6 +135,7 @@ where
129135
channels: self.channels,
130136
tasks: self.tasks,
131137
control_tx: self.control_tx,
138+
latency_rx: self.latency_rx,
132139
}
133140
}
134141
}

crates/core/src/client/quic/mod.rs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -169,15 +169,16 @@ impl Client for QuicClient {
169169

170170
// Create oneshot for receiving server's Handshake via the control loop.
171171
let (handshake_tx, handshake_rx) = tokio::sync::oneshot::channel::<Handshake>();
172+
let (latency_tx, latency_rx) = tokio::sync::mpsc::channel::<f64>(4);
172173

173174
// Spawn control stream task
174175
let transport_ctrl = Arc::clone(&transport);
175176
let control_handle = tokio::spawn(async move {
176177
let mut channels = protocol::ControlChannels {
177178
outgoing_rx: control_rx,
178179
handshake_tx: Some(handshake_tx), // receive server's Handshake
179-
latency_tx: None, // pong handled inline
180-
control_response_tx: None, // no ControlResponse channel needed now
180+
latency_tx: Some(latency_tx),
181+
control_response_tx: None,
181182
role_transition_tx: None,
182183
};
183184
match protocol::run_control_stream_initiator(
@@ -230,6 +231,7 @@ impl Client for QuicClient {
230231
tasks,
231232
control_tx,
232233
Some(handshake_rx),
234+
Some(latency_rx),
233235
))
234236
}
235237

crates/core/src/client/ws/mod.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -386,14 +386,15 @@ impl WsClient {
386386

387387
// Create oneshot for receiving server's Handshake via the control loop.
388388
let (handshake_tx, handshake_rx) = tokio::sync::oneshot::channel::<Handshake>();
389+
let (latency_tx, latency_rx) = tokio::sync::mpsc::channel::<f64>(4);
389390

390391
// Spawn control stream task
391392
let transport_ctrl = Arc::clone(&transport);
392393
let control_handle = tokio::spawn(async move {
393394
let mut channels = protocol::ControlChannels {
394395
outgoing_rx: control_rx,
395396
handshake_tx: Some(handshake_tx), // receive server's Handshake
396-
latency_tx: None,
397+
latency_tx: Some(latency_tx),
397398
control_response_tx: None,
398399
role_transition_tx: None,
399400
};
@@ -447,6 +448,7 @@ impl WsClient {
447448
tasks,
448449
control_tx,
449450
Some(handshake_rx),
451+
Some(latency_rx),
450452
))
451453
}
452454
}

crates/daemon/src/mode/auto.rs

Lines changed: 15 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -259,17 +259,8 @@ async fn run_auto_connector(
259259
let r = Arc::clone(&routes);
260260
let ru = ru.resubscribe();
261261
async move {
262-
let peer_handshake_rx = e.peer_handshake_rx;
263-
let transport = e.transport;
264-
let channels = e.channels;
265-
let tasks = e.tasks;
266-
let control_tx = e.control_tx;
267262
run_auto_connect_session_dispatch(
268-
peer_handshake_rx,
269-
transport,
270-
channels,
271-
tasks,
272-
control_tx,
263+
e,
273264
&lhs,
274265
&pa,
275266
metrics,
@@ -319,17 +310,8 @@ async fn run_auto_connector(
319310
let r = Arc::clone(&routes);
320311
let ru = ru.resubscribe();
321312
async move {
322-
let peer_handshake_rx = e.peer_handshake_rx;
323-
let transport = e.transport;
324-
let channels = e.channels;
325-
let tasks = e.tasks;
326-
let control_tx = e.control_tx;
327313
run_auto_connect_session_dispatch(
328-
peer_handshake_rx,
329-
transport,
330-
channels,
331-
tasks,
332-
control_tx,
314+
e,
333315
&lhs,
334316
&pa,
335317
metrics,
@@ -354,11 +336,7 @@ async fn run_auto_connector(
354336
/// Non-generic auto-connector dispatch: negotiates role and runs the session.
355337
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
356338
async fn run_auto_connect_session_dispatch(
357-
peer_handshake_rx: Option<tokio::sync::oneshot::Receiver<Handshake>>,
358-
transport: Arc<dyn ErasedTransport>,
359-
channels: DataChannels,
360-
tasks: wallhack_core::client::client::ConnectionTasks,
361-
control_tx: tokio::sync::mpsc::Sender<wallhack_wire::control::ControlMessage>,
339+
connect_result: wallhack_core::client::client::ErasedConnectResult,
362340
local_hs: &Handshake,
363341
peer_addr: &str,
364342
metrics: Arc<Metrics>,
@@ -369,6 +347,16 @@ async fn run_auto_connect_session_dispatch(
369347
tokio::sync::broadcast::Receiver<wallhack_core::control::routes::RouteUpdate>,
370348
>,
371349
) -> Result<(), NodeError> {
350+
let wallhack_core::client::client::ErasedConnectResult {
351+
peer_handshake_rx,
352+
transport,
353+
channels,
354+
tasks,
355+
control_tx,
356+
peer_addr: _,
357+
latency_rx,
358+
} = connect_result;
359+
372360
let DataChannels {
373361
instructions_tx,
374362
instructions_rx,
@@ -435,6 +423,8 @@ async fn run_auto_connect_session_dispatch(
435423
&metrics,
436424
peer_addr,
437425
Some(peer_name),
426+
Some(Arc::clone(&peers)),
427+
latency_rx,
438428
routes,
439429
route_updates,
440430
)

crates/daemon/src/mode/entry.rs

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -316,6 +316,7 @@ pub(crate) async fn run_entry_connect(
316316
}
317317

318318
let peer_addr = endpoint.to_string();
319+
let peers = Arc::clone(&res.peers);
319320
let security = SecurityParams {
320321
psk: global.psk.clone(),
321322
accept_fingerprint: None,
@@ -352,11 +353,12 @@ pub(crate) async fn run_entry_connect(
352353
move |connect_result| {
353354
let e = connect_result.erase();
354355
let m = Arc::clone(&metrics);
356+
let p = Arc::clone(&peers);
355357
let r = Arc::clone(&routes);
356358
let ru = r_updates.resubscribe();
357359
let pa = peer_addr.clone();
358360
async move {
359-
run_entry_connected_erased(e, &m, &pa, Some(r), Some(ru)).await
361+
run_entry_connected_erased(e, &m, &pa, Some(p), Some(r), Some(ru)).await
360362
}
361363
},
362364
RECONNECT_DELAY,
@@ -388,11 +390,12 @@ pub(crate) async fn run_entry_connect(
388390
move |connect_result| {
389391
let e = connect_result.erase();
390392
let m = Arc::clone(&metrics);
393+
let p = Arc::clone(&peers);
391394
let r = Arc::clone(&routes);
392395
let ru = r_updates.resubscribe();
393396
let pa = peer_addr.clone();
394397
async move {
395-
run_entry_connected_erased(e, &m, &pa, Some(r), Some(ru)).await
398+
run_entry_connected_erased(e, &m, &pa, Some(p), Some(r), Some(ru)).await
396399
}
397400
},
398401
RECONNECT_DELAY,
@@ -429,6 +432,7 @@ pub(crate) async fn run_entry_connected_erased(
429432
connect_result: wallhack_core::client::client::ErasedConnectResult,
430433
metrics: &Arc<Metrics>,
431434
peer_addr: &str,
435+
peers: Option<Arc<Registry>>,
432436
routes: Option<SharedRouteTable>,
433437
route_updates: Option<
434438
tokio::sync::broadcast::Receiver<wallhack_core::control::routes::RouteUpdate>,
@@ -443,6 +447,7 @@ pub(crate) async fn run_entry_connected_erased(
443447
tasks: _tasks,
444448
control_tx,
445449
peer_addr: _,
450+
latency_rx,
446451
} = connect_result;
447452

448453
// Wait for the server's handshake to get the peer name
@@ -478,6 +483,8 @@ pub(crate) async fn run_entry_connected_erased(
478483
metrics,
479484
peer_addr,
480485
peer_name.as_deref(),
486+
peers,
487+
latency_rx,
481488
routes,
482489
route_updates,
483490
)
@@ -495,6 +502,8 @@ pub(crate) async fn run_entry_connected_inner(
495502
metrics: &Arc<Metrics>,
496503
peer_addr: &str,
497504
peer_name: Option<&str>,
505+
peers: Option<Arc<Registry>>,
506+
latency_rx: Option<tokio::sync::mpsc::Receiver<f64>>,
498507
routes: Option<SharedRouteTable>,
499508
route_updates: Option<
500509
tokio::sync::broadcast::Receiver<wallhack_core::control::routes::RouteUpdate>,
@@ -574,6 +583,19 @@ pub(crate) async fn run_entry_connected_inner(
574583
tracing::debug!("Initial ping failed: {e}");
575584
}
576585

586+
// Consume latency measurements from the control loop and update the registry.
587+
if let (Some(mut lat_rx), Some(reg)) = (latency_rx, &peers) {
588+
let peer_id = peer_name.map(std::string::ToString::to_string);
589+
let reg = Arc::clone(reg);
590+
tokio::spawn(async move {
591+
while let Some(ms) = lat_rx.recv().await {
592+
if let Some(ref id) = peer_id {
593+
reg.update_latency(id, ms);
594+
}
595+
}
596+
});
597+
}
598+
577599
let manager_handle = tokio::spawn(async move { manager.run().await });
578600

579601
match manager_handle.await {

crates/daemon/src/mode/exit.rs

Lines changed: 23 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -416,9 +416,13 @@ async fn run_exit_loop_inner(
416416

417417
tracing::info!("Connected to {peer_addr}");
418418

419-
// Resolve the peer's handshake name (delivered asynchronously via the
419+
// Resolve the peer's handshake (delivered asynchronously via the
420420
// control loop). Fall back to the address if unavailable.
421-
let peer_name = resolve_peer_name(peer_handshake_rx, peer_addr).await;
421+
let peer_handshake = resolve_peer_handshake(peer_handshake_rx).await;
422+
let peer_name = peer_handshake
423+
.as_ref()
424+
.filter(|h| !h.name.is_empty())
425+
.map_or_else(|| peer_addr.to_string(), |h| h.name.clone());
422426

423427
// We connected to the entry peer, so side=Connect.
424428
ctx.peers.register(
@@ -428,6 +432,14 @@ async fn run_exit_loop_inner(
428432
ConnectionSide::Connect,
429433
);
430434

435+
// Apply the peer's advertised capabilities from the handshake.
436+
let peer_capabilities = peer_handshake
437+
.as_ref()
438+
.and_then(|h| h.capabilities)
439+
.unwrap_or_default();
440+
ctx.peers
441+
.update_capabilities(&peer_name, &peer_capabilities);
442+
431443
// Outgoing: open uni stream to entry peer, send responses.
432444
let transport_out = Arc::clone(&transport);
433445
tokio::spawn(async move {
@@ -475,20 +487,18 @@ async fn run_exit_loop_inner(
475487
Ok(())
476488
}
477489

478-
/// Await the peer's handshake name from a oneshot receiver (with timeout).
490+
/// Await the peer's handshake from a oneshot receiver (with timeout).
479491
///
480-
/// Returns the handshake name if non-empty, otherwise falls back to the
481-
/// peer address.
482-
async fn resolve_peer_name(
492+
/// Returns the received `Handshake` if available within the timeout, or
493+
/// `None` if the receiver is absent, the sender dropped, or the timeout
494+
/// expires.
495+
async fn resolve_peer_handshake(
483496
rx: Option<tokio::sync::oneshot::Receiver<wallhack_wire::data::Handshake>>,
484-
peer_addr: &str,
485-
) -> String {
486-
let Some(rx) = rx else {
487-
return peer_addr.to_string();
488-
};
497+
) -> Option<wallhack_wire::data::Handshake> {
498+
let rx = rx?;
489499
match tokio::time::timeout(std::time::Duration::from_secs(10), rx).await {
490-
Ok(Ok(hs)) if !hs.name.is_empty() => hs.name,
491-
_ => peer_addr.to_string(),
500+
Ok(Ok(hs)) => Some(hs),
501+
_ => None,
492502
}
493503
}
494504

0 commit comments

Comments
 (0)