hopr_transport/constants.rs
1use std::time::Duration;
2
3/// The maximum waiting time for a message send to produce a half-key challenge reply
4pub const PACKET_QUEUE_TIMEOUT_MILLISECONDS: std::time::Duration = std::time::Duration::from_millis(15000);
5
6/// Maximum number of outgoing application-layer packets buffered before the writer
7/// observes backpressure (`Poll::Pending` on `poll_write`).
8///
9/// Caps the burst submitted to `buffered(8×N_cpus)` so tail packets never
10/// exceed `PACKET_ENCODING_TIMEOUT` (150 ms). With Sphinx encoding at ~21 ms/packet
11/// and 256 slots, the worst-case Rayon queue depth stays well under the timeout on
12/// machines with up to ~36 cores. Safe formula: `≤ 7 × available_parallelism()`.
13pub(crate) const MAXIMUM_MSG_OUTGOING_BUFFER_SIZE: usize = 256;
14
15/// Time within Start protocol must finish session initiation.
16/// This base value is always multiplied by the (max) number of hops, times 2 (for both-ways).
17pub(crate) const SESSION_INITIATION_TIMEOUT_BASE: Duration = Duration::from_secs(5);
18
19#[cfg(test)]
20mod tests {
21 use super::MAXIMUM_MSG_OUTGOING_BUFFER_SIZE;
22
23 /// Guard against accidentally inflating `MAXIMUM_MSG_OUTGOING_BUFFER_SIZE` back to a large
24 /// value that would overflow the Rayon encoding queue.
25 ///
26 /// Safe formula: ≤ 7 × available_parallelism. 256 covers machines with up to ~36 cores.
27 // Constant by construction — that is the whole point of the guard. Clippy suggests a `const`
28 // block instead, which would be stronger but cannot format the offending value into the
29 // message: `assert!` in const context takes a literal `&str` only. The value is worth more
30 // than the earlier failure here.
31 #[allow(clippy::assertions_on_constants)]
32 #[test]
33 fn outgoing_buffer_size_within_backpressure_limit() {
34 assert!(
35 MAXIMUM_MSG_OUTGOING_BUFFER_SIZE <= 256,
36 "MAXIMUM_MSG_OUTGOING_BUFFER_SIZE={MAXIMUM_MSG_OUTGOING_BUFFER_SIZE} exceeds the safe threshold; tail \
37 packets will exceed PACKET_ENCODING_TIMEOUT under burst writes"
38 );
39 }
40
41 /// Verify that `CrossfireSink` signals backpressure (`Poll::Pending`) once the channel
42 /// reaches `MAXIMUM_MSG_OUTGOING_BUFFER_SIZE`, preventing the Rayon encoding queue from
43 /// receiving an unbounded burst.
44 #[test]
45 fn outgoing_channel_signals_backpressure_when_full() {
46 use std::{
47 pin::Pin,
48 task::{Context, Poll},
49 };
50
51 use futures::Sink;
52 use hopr_utils::network_types::crossfire_sink::bounded_sink_channel;
53
54 let (mut sink, _rx) = bounded_sink_channel::<usize>(MAXIMUM_MSG_OUTGOING_BUFFER_SIZE);
55 let waker = futures::task::noop_waker_ref();
56 let mut cx = Context::from_waker(waker);
57
58 // Standard Sink protocol: each poll_ready → start_send pair sends one item.
59 // At i=0 the channel is empty; at i=N-1 the final item is buffered.
60 for i in 0..MAXIMUM_MSG_OUTGOING_BUFFER_SIZE {
61 assert!(
62 matches!(Pin::new(&mut sink).poll_ready(&mut cx), Poll::Ready(Ok(()))),
63 "poll_ready must be Ready before capacity is reached (item {i})"
64 );
65 Pin::new(&mut sink).start_send(i).unwrap();
66 }
67 // This poll_ready sends the last buffered item; channel is now exactly full.
68 assert!(
69 matches!(Pin::new(&mut sink).poll_ready(&mut cx), Poll::Ready(Ok(()))),
70 "poll_ready must be Ready when flushing the final item into a full-but-not-yet-full channel"
71 );
72
73 // One extra item: buffer it, then poll_ready must indicate the channel is saturated.
74 Pin::new(&mut sink)
75 .start_send(MAXIMUM_MSG_OUTGOING_BUFFER_SIZE)
76 .unwrap();
77 assert!(
78 matches!(Pin::new(&mut sink).poll_ready(&mut cx), Poll::Pending),
79 "CrossfireSink must return Poll::Pending when channel is at capacity ({})",
80 MAXIMUM_MSG_OUTGOING_BUFFER_SIZE
81 );
82 }
83}