Skip to content

Commit 7d02af1

Browse files
authored
Merge pull request #79 from saorsa-labs/fix/stability-improvements
fix(p2p): improve send, dial, and relay stability
2 parents 5d2408e + 55f423c commit 7d02af1

4 files changed

Lines changed: 1078 additions & 280 deletions

File tree

src/masque/relay_server.rs

Lines changed: 92 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,8 @@ use crate::upnp::{UpnpConfig, UpnpMappingService};
6060
/// `nf_conntrack_udp_timeout_stream` is 120 s on Linux) and prevents
6161
/// the QUIC idle timeout from firing on the underlying connection.
6262
const RELAY_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(15);
63+
const RELAY_STREAM_BATCH_MAX_FRAMES: usize = 64;
64+
const RELAY_STREAM_BATCH_MAX_BYTES: usize = 64 * 1024;
6365

6466
/// Capacity of the bounded channel between the UDP reader task and the
6567
/// stream writer task in Direction 1 (external → NATted node). Sized
@@ -204,6 +206,19 @@ enum WriterItem {
204206
Control(Bytes),
205207
}
206208

209+
fn append_relay_frame(out: &mut Vec<u8>, encoded: &Bytes) {
210+
let frame_len = encoded.len() as u32;
211+
out.extend_from_slice(&frame_len.to_be_bytes());
212+
out.extend_from_slice(encoded);
213+
}
214+
215+
fn append_control_frame(out: &mut Vec<u8>, body: &Bytes) {
216+
let body_len = body.len() as u32;
217+
out.extend_from_slice(&CONTROL_FRAME_MARKER.to_be_bytes());
218+
out.extend_from_slice(&body_len.to_be_bytes());
219+
out.extend_from_slice(body);
220+
}
221+
207222
/// Configuration for the MASQUE relay server
208223
#[derive(Debug, Clone)]
209224
pub struct MasqueRelayConfig {
@@ -1043,7 +1058,7 @@ impl MasqueRelayServer {
10431058
match socket.recv_from(&mut buf).await {
10441059
Ok((len, source)) => {
10451060
let payload = Bytes::copy_from_slice(&buf[..len]);
1046-
tracing::debug!(
1061+
tracing::trace!(
10471062
session_id,
10481063
source = %source,
10491064
len,
@@ -1116,7 +1131,7 @@ impl MasqueRelayServer {
11161131
};
11171132
match resolved {
11181133
Some((target, payload)) => {
1119-
tracing::debug!(
1134+
tracing::trace!(
11201135
session_id,
11211136
target = %target,
11221137
len = payload.len(),
@@ -1126,7 +1141,7 @@ impl MasqueRelayServer {
11261141
server2.stats.record_datagram();
11271142
match socket2.send_to(&payload, target).await {
11281143
Ok(n) => {
1129-
tracing::debug!(
1144+
tracing::trace!(
11301145
session_id,
11311146
target = %target,
11321147
len = payload.len(),
@@ -1244,7 +1259,7 @@ impl MasqueRelayServer {
12441259
match socket.recv_from(&mut buf).await {
12451260
Ok((len, source)) => {
12461261
let payload = Bytes::copy_from_slice(&buf[..len]);
1247-
tracing::debug!(
1262+
tracing::trace!(
12481263
session_id, source = %source, len,
12491264
"RELAY_TUNNEL[srv]: stream-loop dir1 recv UDP → forwarding to relay-client"
12501265
);
@@ -1254,12 +1269,12 @@ impl MasqueRelayServer {
12541269
stats.record_bytes(encoded.len() as u64);
12551270
stats.record_datagram();
12561271
if fwd_tx.send(WriterItem::Data(encoded)).await.is_err() {
1257-
break; // writer closed
1272+
return "writer_channel_closed";
12581273
}
12591274
}
12601275
Err(e) => {
12611276
tracing::debug!(session_id, error = %e, "UDP recv error");
1262-
break;
1277+
return "udp_recv_error";
12631278
}
12641279
}
12651280
}
@@ -1273,62 +1288,92 @@ impl MasqueRelayServer {
12731288
loop {
12741289
tokio::select! {
12751290
item = fwd_rx.recv() => {
1276-
let Some(item) = item else { break };
1291+
let Some(item) = item else {
1292+
return "forward_channel_closed";
1293+
};
12771294
match item {
12781295
WriterItem::Data(encoded) => {
1279-
let frame_len = encoded.len() as u32;
1280-
if let Err(e) = send_stream.write_all(&frame_len.to_be_bytes()).await {
1281-
tracing::debug!(session_id, error = %e, "Stream write error (length)");
1282-
break;
1296+
let mut batch = Vec::with_capacity(
1297+
encoded
1298+
.len()
1299+
.saturating_add(std::mem::size_of::<u32>()),
1300+
);
1301+
append_relay_frame(&mut batch, &encoded);
1302+
1303+
let mut frames = 1usize;
1304+
while frames < RELAY_STREAM_BATCH_MAX_FRAMES
1305+
&& batch.len() < RELAY_STREAM_BATCH_MAX_BYTES
1306+
{
1307+
match fwd_rx.try_recv() {
1308+
Ok(WriterItem::Data(next)) => {
1309+
append_relay_frame(&mut batch, &next);
1310+
frames += 1;
1311+
}
1312+
Ok(WriterItem::Control(body)) => {
1313+
append_control_frame(&mut batch, &body);
1314+
break;
1315+
}
1316+
Err(tokio::sync::mpsc::error::TryRecvError::Empty) => break,
1317+
Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => break,
1318+
}
12831319
}
1284-
if let Err(e) = send_stream.write_all(&encoded).await {
1285-
tracing::debug!(session_id, error = %e, "Stream write error (data)");
1286-
break;
1320+
1321+
if let Err(e) = send_stream.write_all(&batch).await {
1322+
tracing::debug!(session_id, error = %e, frames, bytes = batch.len(), "Stream batch write error");
1323+
return "stream_batch_write_error";
12871324
}
12881325
}
12891326
WriterItem::Control(body) => {
1290-
// Wire format: [4-byte BE marker][4-byte BE body_len][body]
1291-
let body_len = body.len() as u32;
1292-
if let Err(e) = send_stream
1293-
.write_all(&CONTROL_FRAME_MARKER.to_be_bytes())
1294-
.await
1295-
{
1296-
tracing::debug!(session_id, error = %e, "Stream write error (control marker)");
1297-
break;
1298-
}
1299-
if let Err(e) = send_stream.write_all(&body_len.to_be_bytes()).await {
1300-
tracing::debug!(session_id, error = %e, "Stream write error (control len)");
1301-
break;
1302-
}
1303-
if let Err(e) = send_stream.write_all(&body).await {
1304-
tracing::debug!(session_id, error = %e, "Stream write error (control body)");
1305-
break;
1327+
let mut frame = Vec::with_capacity(
1328+
body.len()
1329+
.saturating_add(std::mem::size_of::<u64>()),
1330+
);
1331+
append_control_frame(&mut frame, &body);
1332+
if let Err(e) = send_stream.write_all(&frame).await {
1333+
tracing::debug!(session_id, error = %e, "Stream write error (control frame)");
1334+
return "stream_control_write_error";
13061335
}
13071336
}
13081337
}
13091338
}
13101339
_ = keepalive.tick() => {
13111340
if let Err(e) = send_stream.write_all(&keepalive_bytes).await {
13121341
tracing::debug!(session_id, error = %e, "Keepalive write error");
1313-
break;
1342+
return "keepalive_write_error";
13141343
}
13151344
}
13161345
}
13171346
}
13181347
});
13191348

1320-
tokio::select! {
1321-
_ = reader_handle => {},
1322-
_ = writer_handle => {},
1349+
let close_reason = tokio::select! {
1350+
result = reader_handle => {
1351+
match result {
1352+
Ok(reason) => reason,
1353+
Err(e) => {
1354+
tracing::debug!(session_id, error = %e, "Relay UDP reader task join error");
1355+
"reader_task_join_error"
1356+
}
1357+
}
1358+
},
1359+
result = writer_handle => {
1360+
match result {
1361+
Ok(reason) => reason,
1362+
Err(e) => {
1363+
tracing::debug!(session_id, error = %e, "Relay stream writer task join error");
1364+
"writer_task_join_error"
1365+
}
1366+
}
1367+
},
13231368

13241369
// Direction 2: Stream → UDP (client → relay → target)
1325-
_ = async {
1370+
reason = async {
13261371
loop {
13271372
// Read 4-byte length prefix
13281373
let mut len_buf = [0u8; 4];
13291374
if let Err(e) = recv_stream.read_exact(&mut len_buf).await {
13301375
tracing::debug!(session_id, error = %e, "Stream read error (length)");
1331-
break;
1376+
return "stream_read_length_error";
13321377
}
13331378
let frame_len = u32::from_be_bytes(len_buf) as usize;
13341379

@@ -1348,21 +1393,21 @@ impl MasqueRelayServer {
13481393
frame_len,
13491394
"Corrupt stream frame length, closing session"
13501395
);
1351-
break;
1396+
return "corrupt_frame_length";
13521397
}
13531398

13541399
// Read frame data
13551400
let mut frame_buf = vec![0u8; frame_len];
13561401
if let Err(e) = recv_stream.read_exact(&mut frame_buf).await {
13571402
tracing::debug!(session_id, error = %e, "Stream read error (data)");
1358-
break;
1403+
return "stream_read_data_error";
13591404
}
13601405

13611406
// Decode and forward
13621407
let mut cursor = Bytes::from(frame_buf);
13631408
match UncompressedDatagram::decode(&mut cursor) {
13641409
Ok(datagram) => {
1365-
tracing::debug!(
1410+
tracing::trace!(
13661411
session_id, target = %datagram.target,
13671412
len = datagram.payload.len(),
13681413
"RELAY_TUNNEL[srv]: stream-loop dir2 recv from relay-client → sendto target"
@@ -1373,7 +1418,7 @@ impl MasqueRelayServer {
13731418
let payload_len = datagram.payload.len();
13741419
match socket2.send_to(&datagram.payload, target).await {
13751420
Ok(n) => {
1376-
tracing::debug!(
1421+
tracing::trace!(
13771422
session_id,
13781423
target = %target,
13791424
len = payload_len,
@@ -1425,10 +1470,14 @@ impl MasqueRelayServer {
14251470
}
14261471
}
14271472
}
1428-
} => {},
1429-
}
1473+
} => reason,
1474+
};
14301475

1431-
tracing::info!(session_id, "Stream-based relay forwarding loop ended");
1476+
tracing::info!(
1477+
session_id,
1478+
close_reason,
1479+
"Stream-based relay forwarding loop ended"
1480+
);
14321481
if let Err(e) = self.close_session(session_id).await {
14331482
tracing::debug!(session_id, error = %e, "Error closing session");
14341483
}

src/masque/relay_socket.rs

Lines changed: 30 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,8 @@ const RELAY_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(15);
7070
/// docs), which matches UDP's lossy semantics. 8192 × ~1200 B ≈ 10 MB
7171
/// of worst-case buffering before drops begin.
7272
const SEND_QUEUE_CAPACITY: usize = 8192;
73+
const RELAY_STREAM_BATCH_MAX_FRAMES: usize = 64;
74+
const RELAY_STREAM_BATCH_MAX_BYTES: usize = 64 * 1024;
7375

7476
/// Upper bound on decoded inbound packets queued for `poll_recv`.
7577
/// The reader task awaits on `Sender::send`, so when this fills up the
@@ -81,6 +83,12 @@ const RECV_QUEUE_CAPACITY: usize = 8192;
8183
/// framing error or corruption and closes the session.
8284
const MAX_RELAY_FRAME: usize = 512 * 1024;
8385

86+
fn append_relay_frame(out: &mut Vec<u8>, encoded: &Bytes) {
87+
let frame_len = encoded.len() as u32;
88+
out.extend_from_slice(&frame_len.to_be_bytes());
89+
out.extend_from_slice(encoded);
90+
}
91+
8492
/// Raw QUIC streams from a relay session, before socket construction.
8593
///
8694
/// Returned by `establish_relay_session` so the caller can construct a
@@ -274,7 +282,7 @@ impl MasqueRelaySocket {
274282
// the original frame buffer — no clone needed.
275283
let inbound_source = datagram.target;
276284
let inbound_len = datagram.payload.len();
277-
tracing::debug!(
285+
tracing::trace!(
278286
relay = %relay_public_addr,
279287
source = %inbound_source,
280288
len = inbound_len,
@@ -315,21 +323,30 @@ impl MasqueRelaySocket {
315323
// slow) stream write — so the Quinn endpoint can start
316324
// assembling the next packet concurrently with this
317325
// frame going out on the wire.
326+
let mut batch =
327+
Vec::with_capacity(encoded.len().saturating_add(std::mem::size_of::<u32>()));
328+
append_relay_frame(&mut batch, &encoded);
318329
writer_capacity.notify_one();
319330

320-
let frame_len = encoded.len() as u32;
321-
if let Err(e) = send_stream.write_all(&frame_len.to_be_bytes()).await {
322-
tracing::debug!(error = %e, "MasqueRelaySocket: stream write error (length)");
323-
break;
324-
}
325-
// Zero-length frames (keepalives) need only the length
326-
// prefix — skip the empty write_all.
327-
if !encoded.is_empty() {
328-
if let Err(e) = send_stream.write_all(&encoded).await {
329-
tracing::debug!(error = %e, "MasqueRelaySocket: stream write error (data)");
330-
break;
331+
let mut frames = 1usize;
332+
while frames < RELAY_STREAM_BATCH_MAX_FRAMES
333+
&& batch.len() < RELAY_STREAM_BATCH_MAX_BYTES
334+
{
335+
match send_rx.try_recv() {
336+
Ok(next) => {
337+
append_relay_frame(&mut batch, &next);
338+
writer_capacity.notify_one();
339+
frames += 1;
340+
}
341+
Err(mpsc::error::TryRecvError::Empty) => break,
342+
Err(mpsc::error::TryRecvError::Disconnected) => break,
331343
}
332344
}
345+
346+
if let Err(e) = send_stream.write_all(&batch).await {
347+
tracing::debug!(error = %e, frames, bytes = batch.len(), "MasqueRelaySocket: stream batch write error");
348+
break;
349+
}
333350
}
334351
// Writer exited (stream error or receiver dropped). Dropping
335352
// `send_rx` closes the channel so subsequent `try_send`
@@ -410,7 +427,7 @@ impl AsyncUdpSocket for MasqueRelaySocket {
410427
// of `segment_size` bytes. Each segment must be sent as its
411428
// own tunnel frame — the relay server has a per-frame size
412429
// limit and cannot handle the entire batch as one.
413-
tracing::debug!(
430+
tracing::trace!(
414431
relay = %self.relay_public_addr,
415432
destination = %transmit.destination,
416433
len = transmit.contents.len(),

0 commit comments

Comments
 (0)