mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-08 18:37:46 +08:00
fix(tunnel): bound upstream clients and heartbeat deltas
This commit is contained in:
@@ -239,7 +239,10 @@ where
|
||||
|
||||
// Create body channel and spawn handler
|
||||
let (body_tx, body_rx) = mpsc::channel::<Frame>(64);
|
||||
streams.insert(frame.stream_id, body_tx);
|
||||
let request_headers_end_stream = frame.is_end_stream();
|
||||
if !request_headers_end_stream {
|
||||
streams.insert(frame.stream_id, body_tx);
|
||||
}
|
||||
|
||||
let state_clone = Arc::clone(&state);
|
||||
let server_clone = Arc::clone(&server);
|
||||
@@ -326,8 +329,16 @@ where
|
||||
// Trigger every 64 frames OR when the count exceeds max_streams.
|
||||
frames_since_cleanup += 1;
|
||||
if frames_since_cleanup >= 64 || handler_handles.len() > max_streams {
|
||||
let closed_streams = prune_closed_stream_senders(&mut streams);
|
||||
if closed_streams > 0 {
|
||||
debug!(closed_streams, "removed closed request body stream senders");
|
||||
}
|
||||
handler_handles.retain(|h| !h.is_finished());
|
||||
frames_since_cleanup = 0;
|
||||
if draining && streams.is_empty() {
|
||||
info!("tunnel drained after cleanup");
|
||||
break None;
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
@@ -397,6 +408,12 @@ fn try_send_stream_error(frame_tx: &FrameSender, stream_id: u32, message: &'stat
|
||||
}
|
||||
}
|
||||
|
||||
fn prune_closed_stream_senders(streams: &mut HashMap<u32, mpsc::Sender<Frame>>) -> usize {
|
||||
let before = streams.len();
|
||||
streams.retain(|_, tx| !tx.is_closed());
|
||||
before.saturating_sub(streams.len())
|
||||
}
|
||||
|
||||
/// Wait for all active stream handlers to finish (with a timeout).
|
||||
async fn drain_handlers(handles: Vec<JoinHandle<()>>) {
|
||||
if handles.is_empty() {
|
||||
@@ -470,4 +487,18 @@ mod tests {
|
||||
Bytes::from_static(b"tunnel request body dispatch stalled")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn prune_closed_stream_senders_drops_streams_with_closed_receivers() {
|
||||
let (closed_tx, closed_rx) = mpsc::channel::<Frame>(1);
|
||||
let (open_tx, _open_rx) = mpsc::channel::<Frame>(1);
|
||||
drop(closed_rx);
|
||||
let mut streams = HashMap::from([(7, closed_tx), (9, open_tx)]);
|
||||
|
||||
let removed = prune_closed_stream_senders(&mut streams);
|
||||
|
||||
assert_eq!(removed, 1);
|
||||
assert!(!streams.contains_key(&7));
|
||||
assert!(streams.contains_key(&9));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -532,6 +532,7 @@ mod tests {
|
||||
upstream_connect_timeout_secs: 30,
|
||||
upstream_pool_max_idle_per_host: 4,
|
||||
upstream_pool_idle_timeout_secs: 60,
|
||||
upstream_client_pool_capacity: crate::config::DEFAULT_UPSTREAM_CLIENT_POOL_CAPACITY,
|
||||
upstream_tcp_keepalive_secs: 60,
|
||||
upstream_tcp_nodelay: true,
|
||||
upstream_proxy_url: None,
|
||||
|
||||
@@ -521,6 +521,21 @@ fn prepare_request_body(
|
||||
}
|
||||
}
|
||||
|
||||
fn prepare_bodyless_request_body(
|
||||
body_rx: mpsc::Receiver<TunnelFrame>,
|
||||
follow_redirects: bool,
|
||||
) -> PreparedRequestBody {
|
||||
drop(body_rx);
|
||||
PreparedRequestBody {
|
||||
first_request_body: Some(empty_request_body()),
|
||||
replay_body: if follow_redirects {
|
||||
ReplayableRequestBody::None
|
||||
} else {
|
||||
ReplayableRequestBody::NonReplayable
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
async fn collect_request_body_for_replay(
|
||||
mut body_rx: mpsc::Receiver<TunnelFrame>,
|
||||
body_size: Arc<AtomicUsize>,
|
||||
@@ -1379,17 +1394,7 @@ async fn handle_stream_inner(
|
||||
0,
|
||||
)
|
||||
} else {
|
||||
PreparedRequestBody {
|
||||
first_request_body: Some(build_streaming_request_body(
|
||||
body_rx,
|
||||
Arc::clone(&request_body_size),
|
||||
)),
|
||||
replay_body: if follow_redirects {
|
||||
ReplayableRequestBody::None
|
||||
} else {
|
||||
ReplayableRequestBody::NonReplayable
|
||||
},
|
||||
}
|
||||
prepare_bodyless_request_body(body_rx, follow_redirects)
|
||||
};
|
||||
|
||||
let mut total_dns_ms = 0u64;
|
||||
@@ -1794,6 +1799,22 @@ mod tests {
|
||||
assert_eq!(body_size.load(Ordering::Relaxed), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn bodyless_request_body_completes_without_waiting_for_tunnel_sender() {
|
||||
let (_tx, rx) = mpsc::channel(4);
|
||||
let mut prepared = prepare_bodyless_request_body(rx, true);
|
||||
let mut body = prepared
|
||||
.first_request_body
|
||||
.take()
|
||||
.expect("bodyless request should have an initial body");
|
||||
|
||||
let frame = tokio::time::timeout(Duration::from_millis(25), body.frame())
|
||||
.await
|
||||
.expect("bodyless request body should not wait for tunnel body frames");
|
||||
assert!(frame.is_none());
|
||||
assert!(matches!(prepared.replay_body, ReplayableRequestBody::None));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn prepare_request_body_streams_immediately_and_replays_after_completion() {
|
||||
let (tx, rx) = mpsc::channel(4);
|
||||
@@ -2690,6 +2711,7 @@ mod tests {
|
||||
upstream_connect_timeout_secs: 30,
|
||||
upstream_pool_max_idle_per_host: 4,
|
||||
upstream_pool_idle_timeout_secs: 60,
|
||||
upstream_client_pool_capacity: crate::config::DEFAULT_UPSTREAM_CLIENT_POOL_CAPACITY,
|
||||
upstream_tcp_keepalive_secs: 60,
|
||||
upstream_tcp_nodelay: true,
|
||||
upstream_proxy_url: None,
|
||||
|
||||
Reference in New Issue
Block a user