From 9cec63aeb57c80d45246e3648360953cd0334ff5 Mon Sep 17 00:00:00 2001 From: DuProcess <273172371+DuProcess@users.noreply.github.com> Date: Fri, 17 Jul 2026 19:50:49 -0400 Subject: [PATCH] Trace forwarded stream lifecycles --- src/bin/dosh-client.rs | 34 ++++++++++++++++++++++++++++++++++ src/bin/dosh-server/unix.rs | 37 +++++++++++++++++++++++++++++++++++++ tests/integration_smoke.rs | 8 ++++++++ 3 files changed, 79 insertions(+) diff --git a/src/bin/dosh-client.rs b/src/bin/dosh-client.rs index 458bc87..1b9ed67 100644 --- a/src/bin/dosh-client.rs +++ b/src/bin/dosh-client.rs @@ -7500,6 +7500,10 @@ async fn run_terminal( let Ok(ok) = protocol::from_body::(&plain) else { continue; }; + dosh::trace::event( + "client.stream_open_ok", + &[("stream", ok.stream_id.to_string())], + ); last_packet_at = Instant::now(); let Some(pending_open) = stream_pending_opens.remove(&ok.stream_id) else { continue; @@ -7549,6 +7553,13 @@ async fn run_terminal( let Ok(reject) = protocol::from_body::(&plain) else { continue; }; + dosh::trace::event( + "client.stream_open_reject", + &[ + ("stream", reject.stream_id.to_string()), + ("reason", reject.reason.clone()), + ], + ); last_packet_at = Instant::now(); stream_pending_opens.remove(&reject.stream_id); if pending_socks_replies.remove(&reject.stream_id) @@ -7658,6 +7669,10 @@ async fn run_terminal( let Ok(eof) = protocol::from_body::(&plain) else { continue; }; + dosh::trace::event( + "client.stream_eof_received", + &[("stream", eof.stream_id.to_string())], + ); last_packet_at = Instant::now(); if !opened_streams.contains(&eof.stream_id) { continue; @@ -7703,6 +7718,10 @@ async fn run_terminal( let Ok(close) = protocol::from_body::(&plain) else { continue; }; + dosh::trace::event( + "client.stream_close_received", + &[("stream", close.stream_id.to_string())], + ); last_packet_at = Instant::now(); pending_socks_replies.remove(&close.stream_id); stream_close_retransmit.remove(&close.stream_id); @@ -7732,6 +7751,13 @@ async fn run_terminal( forward_event = forward_rx.recv() => { match forward_event { Some(ForwardEvent::Open { stream_id, target_host, target_port, writer, socks5_reply }) => { + dosh::trace::event( + "client.forward_open", + &[ + ("stream", stream_id.to_string()), + ("target", format!("{target_host}:{target_port}")), + ], + ); stream_writers.insert(stream_id, writer); stream_send_credit.insert(stream_id, STREAM_INITIAL_WINDOW); stream_next_send_offset.entry(stream_id).or_insert(0); @@ -7772,6 +7798,10 @@ async fn run_terminal( ).await?; } Some(ForwardEvent::Eof { stream_id }) => { + dosh::trace::event( + "client.forward_eof", + &[("stream", stream_id.to_string())], + ); stream_local_eof_sent.insert(stream_id); if opened_streams.contains(&stream_id) { send_stream_eof(&socket, addr, &cred, &mut send_seq, stream_id).await?; @@ -7811,6 +7841,10 @@ async fn run_terminal( } } Some(ForwardEvent::Close { stream_id }) => { + dosh::trace::event( + "client.forward_close", + &[("stream", stream_id.to_string())], + ); pending_socks_replies.remove(&stream_id); retire_stream_state( stream_id, diff --git a/src/bin/dosh-server/unix.rs b/src/bin/dosh-server/unix.rs index a33455a..bd470a0 100644 --- a/src/bin/dosh-server/unix.rs +++ b/src/bin/dosh-server/unix.rs @@ -2089,6 +2089,14 @@ async fn handle_stream_open( } if open.target_host == FILE_STREAM_SENTINEL { + dosh::trace::event( + "server.file_stream_open", + &[ + ("session", session_name.clone()), + ("stream", open.stream_id.to_string()), + ("peer", peer.to_string()), + ], + ); let (writer_tx, writer_rx) = mpsc::channel::>(1024); register_opened_stream( state, @@ -2113,6 +2121,7 @@ async fn handle_stream_open( ) .await { + eprintln!("file service stream {stream_id} failed: {err:#}"); let _ = send_file_response_to_client( &state, &socket, @@ -2464,6 +2473,14 @@ async fn handle_stream_eof( let (key, session_name) = find_client_decrypt_key(state, &packet.header)?; let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?; let eof: StreamEof = protocol::from_body(&body)?; + dosh::trace::event( + "server.stream_eof_received", + &[ + ("session", session_name.clone()), + ("stream", eof.stream_id.to_string()), + ("peer", peer.to_string()), + ], + ); let mut locked = state.lock().expect("server state poisoned"); let Some(session) = locked.sessions.get_mut(&session_name) else { return Ok(()); @@ -2495,6 +2512,14 @@ async fn handle_stream_close( let (key, session_name) = find_client_decrypt_key(state, &packet.header)?; let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?; let close: StreamClose = protocol::from_body(&body)?; + dosh::trace::event( + "server.stream_close_received", + &[ + ("session", session_name.clone()), + ("stream", close.stream_id.to_string()), + ("peer", peer.to_string()), + ], + ); let mut locked = state.lock().expect("server state poisoned"); let Some(session) = locked.sessions.get_mut(&session_name) else { return Ok(()); @@ -2887,6 +2912,7 @@ async fn run_file_stream_service( ) .await { + eprintln!("file service request on stream {stream_id} failed: {err:#}"); send_file_response_to_client( &state, &socket, @@ -2900,6 +2926,10 @@ async fn run_file_stream_service( } } } + dosh::trace::event( + "server.file_stream_input_closed", + &[("stream", stream_id.to_string())], + ); Ok(()) } @@ -3817,6 +3847,13 @@ async fn send_stream_close_to_client( client_id: [u8; 16], stream_id: u64, ) -> Result<()> { + dosh::trace::event( + "server.stream_close_sent", + &[ + ("stream", stream_id.to_string()), + ("client", hex_id(client_id)), + ], + ); { let mut locked = state.lock().expect("server state poisoned"); if let Some(client) = locked.client_mut(&client_id) { diff --git a/tests/integration_smoke.rs b/tests/integration_smoke.rs index dfc9d67..0d53fba 100644 --- a/tests/integration_smoke.rs +++ b/tests/integration_smoke.rs @@ -1027,6 +1027,10 @@ fn native_file_copy_recursive_round_trip() { fs::create_dir_all(src.join("nested")).unwrap(); fs::write(src.join("root.txt"), b"root file\n").unwrap(); fs::write(src.join("nested/child.txt"), b"child file\n").unwrap(); + let large_payload: Vec = (0..2 * 1024 * 1024) + .map(|index| (index % 251) as u8) + .collect(); + fs::write(src.join("large.bin"), &large_payload).unwrap(); std::os::unix::fs::symlink("root.txt", src.join("root-link")).unwrap(); std::os::unix::fs::symlink("nested/child.txt", src.join("child-link")).unwrap(); @@ -1134,6 +1138,10 @@ fn native_file_copy_recursive_round_trip() { fs::read_to_string(downloaded.join("nested/child.txt")).unwrap(), "child file\n" ); + assert_eq!( + fs::read(downloaded.join("large.bin")).unwrap(), + large_payload + ); assert_eq!( fs::read_link(downloaded.join("root-link")).unwrap(), PathBuf::from("root.txt")