Compare commits

...
10 Commits
Author SHA1 Message Date
DuProcess 0cdcaaaec1 Release v1.0.0-rc45
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 20:13:08 -04:00
DuProcess 9245c666af Bound reliable stream bursts for terminal fairness
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:08:39 -04:00
DuProcess af11ab889e Keep large file regression test fast
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:03:05 -04:00
DuProcess 43a7a69b9b Keep reliable stream packets below path MTU
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:01:12 -04:00
DuProcess 86a1942aa0 Trace terminal loop failures
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:55:16 -04:00
DuProcess 7a9ad38657 Surface service forwarder failures
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:52:19 -04:00
DuProcess 9cec63aeb5 Trace forwarded stream lifecycles
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:50:49 -04:00
DuProcess 3364a7eb7b Release v1.0.0-rc44
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 19:36:21 -04:00
DuProcess 970d54b991 Preserve restart-orphaned sessions for full timeout
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:32:28 -04:00
DuProcess b8b56c0f32 Honor configured shell for remote exec
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:26:08 -04:00
6 changed files with 335 additions and 81 deletions
Generated
+1 -1
View File
@@ -436,7 +436,7 @@ dependencies = [
[[package]]
name = "dosh"
version = "1.0.0-rc43"
version = "1.0.0-rc45"
dependencies = [
"anyhow",
"base64",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "dosh"
version = "1.0.0-rc43"
version = "1.0.0-rc45"
edition = "2024"
license = "MIT"
+128 -66
View File
@@ -67,7 +67,7 @@ use tokio::net::windows::named_pipe::{ClientOptions, NamedPipeClient};
use tokio::net::{TcpListener, TcpStream, UdpSocket};
use tokio::sync::mpsc;
const STREAM_INITIAL_WINDOW: usize = 1024 * 1024;
const STREAM_INITIAL_WINDOW: usize = dosh::transport::DEFAULT_INITIAL_WINDOW;
const STREAM_RETIRED_TOMBSTONES: usize = 16 * 1024;
const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
const MAX_PENDING_USER_INPUT_BYTES: usize = 1024 * 1024;
@@ -486,22 +486,24 @@ async fn main() -> Result<()> {
detach_once(&socket, &cred, 2).await?;
return Ok(());
}
return run_terminal(
socket,
cred,
Some(target_udp_addr),
Some(frame),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
Vec::new(),
Vec::new(),
false,
None,
)
.await;
return trace_terminal_result(
run_terminal(
socket,
cred,
Some(target_udp_addr),
Some(frame),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
Vec::new(),
Vec::new(),
false,
None,
)
.await,
);
}
Err(err) => {
log_debug(
@@ -525,22 +527,24 @@ async fn main() -> Result<()> {
detach_once(&socket, &cred, 2).await?;
return Ok(());
}
return run_terminal(
socket,
cred,
Some(target_udp_addr),
Some(frame),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
Vec::new(),
Vec::new(),
false,
None,
)
.await;
return trace_terminal_result(
run_terminal(
socket,
cred,
Some(target_udp_addr),
Some(frame),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
Vec::new(),
Vec::new(),
false,
None,
)
.await,
);
}
Err(err) => {
log_debug(
@@ -601,22 +605,24 @@ async fn main() -> Result<()> {
detach_once(&socket, &cred, 2).await?;
return Ok(());
}
return run_terminal(
socket,
cred,
Some(addr),
Some(frame),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
local_forwards,
dynamic_forwards,
args.forward_only,
agent_sock,
)
.await;
return trace_terminal_result(
run_terminal(
socket,
cred,
Some(addr),
Some(frame),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
local_forwards,
dynamic_forwards,
args.forward_only,
agent_sock,
)
.await,
);
}
Err(err) if !forwarding_requested && auth_allows(&auth_preference, "ssh") => {
log_debug(
@@ -680,22 +686,24 @@ async fn main() -> Result<()> {
detach_once(&socket, &cred, 2).await?;
return Ok(());
}
run_terminal(
socket,
cred,
Some(target_udp_addr),
Some(first),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
Vec::new(),
Vec::new(),
false,
None,
trace_terminal_result(
run_terminal(
socket,
cred,
Some(target_udp_addr),
Some(first),
predict,
predict_mode,
escape_key,
startup_command,
config.reconnect_timeout_secs,
Vec::new(),
Vec::new(),
false,
None,
)
.await,
)
.await
}
fn run_trust_command(config: &dosh::config::ClientConfig, args: &Args) -> Result<()> {
@@ -2380,8 +2388,21 @@ struct FileForwarder {
impl Drop for FileForwarder {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
let exited = self.child.try_wait().ok().flatten();
if exited.is_none() {
let _ = self.child.kill();
let _ = self.child.wait();
return;
}
if exited.is_some_and(|status| !status.success()) {
let mut stderr = String::new();
if let Some(mut pipe) = self.child.stderr.take() {
let _ = pipe.read_to_string(&mut stderr);
}
if !stderr.trim().is_empty() {
eprintln!("Dosh service forwarder failed:\n{}", stderr.trim_end());
}
}
}
}
@@ -6482,6 +6503,13 @@ fn terminal_input_channel() -> (mpsc::Sender<Vec<u8>>, mpsc::Receiver<Vec<u8>>)
mpsc::channel(STDIN_QUEUE_CAPACITY)
}
fn trace_terminal_result(result: Result<()>) -> Result<()> {
if let Err(err) = &result {
dosh::trace::event("client.terminal_error", &[("error", format!("{err:#}"))]);
}
result
}
#[allow(clippy::too_many_arguments)]
async fn run_terminal(
socket: UdpSocket,
@@ -7500,6 +7528,10 @@ async fn run_terminal(
let Ok(ok) = protocol::from_body::<StreamOpenOk>(&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 +7581,13 @@ async fn run_terminal(
let Ok(reject) = protocol::from_body::<StreamOpenReject>(&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 +7697,10 @@ async fn run_terminal(
let Ok(eof) = protocol::from_body::<StreamEof>(&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 +7746,10 @@ async fn run_terminal(
let Ok(close) = protocol::from_body::<StreamClose>(&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 +7779,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 +7826,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 +7869,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,
+137 -11
View File
@@ -48,7 +48,7 @@ use tokio::net::{TcpListener, TcpStream, UdpSocket, UnixListener};
use tokio::process::Command as TokioCommand;
use tokio::sync::mpsc;
const STREAM_INITIAL_WINDOW: usize = 1024 * 1024;
const STREAM_INITIAL_WINDOW: usize = dosh::transport::DEFAULT_INITIAL_WINDOW;
const STREAM_RETIRED_TOMBSTONES: usize = 16 * 1024;
const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
@@ -398,6 +398,10 @@ struct Session {
holder_control: Option<StdUnixStream>,
/// Whether this session's shell lives in a holder process (persistent).
persistent: bool,
/// The holder survived a server restart but no client has reattached yet.
/// Such sessions need the full reconnect window because the client may be
/// asleep during an unattended server update.
restart_orphaned: bool,
/// Bytes of session output since the screen was last mirrored to disk, used
/// to throttle the (atomic) screen-persistence writes.
bytes_since_persist: usize,
@@ -569,6 +573,7 @@ impl ServerState {
empty_since: None,
holder_control: control,
persistent,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now() - SCREEN_PERSIST_MAX_AGE,
@@ -714,6 +719,7 @@ impl ServerState {
empty_since: Some(Instant::now()),
holder_control: Some(control),
persistent: true,
restart_orphaned: true,
bytes_since_persist: 0,
last_persisted_seq: output_seq,
last_screen_persist_at: Instant::now(),
@@ -732,6 +738,7 @@ impl ServerState {
if let Some(session) = self.sessions.get_mut(session_name) {
session.clients.insert(client_id, client);
session.empty_since = None;
session.restart_orphaned = false;
self.client_index
.insert(client_id, session_name.to_string());
}
@@ -2082,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::<Vec<u8>>(1024);
register_opened_stream(
state,
@@ -2106,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,
@@ -2457,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(());
@@ -2488,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(());
@@ -2880,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,
@@ -2893,6 +2926,10 @@ async fn run_file_stream_service(
}
}
}
dosh::trace::event(
"server.file_stream_input_closed",
&[("stream", stream_id.to_string())],
);
Ok(())
}
@@ -3271,15 +3308,25 @@ async fn run_exec_stream_service(
continue;
}
};
let output = TokioCommand::new("sh")
.arg("-lc")
let shell = {
state
.lock()
.expect("server state poisoned")
.config
.shell
.clone()
};
let output = TokioCommand::new(&shell)
.arg("-c")
.arg(&request.command)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()
.await
.with_context(|| format!("run command {:?}", request.command))?;
.with_context(|| {
format!("run command {:?} with shell {shell}", request.command)
})?;
for chunk in output.stdout.chunks(CHUNK_SIZE) {
send_exec_response_to_client(
&state,
@@ -3800,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) {
@@ -4447,7 +4501,8 @@ fn cleanup_disconnected_clients(state: &Arc<Mutex<ServerState>>) {
.filter(|(name, session)| {
!prewarm.contains(name.as_str())
&& session.empty_since.is_some_and(|since| {
now.duration_since(since) >= empty_session_timeout(name, timeout)
now.duration_since(since)
>= empty_session_timeout(name, timeout, session.restart_orphaned)
})
})
.map(|(name, _)| name.clone())
@@ -4467,8 +4522,12 @@ fn cleanup_disconnected_clients(state: &Arc<Mutex<ServerState>>) {
}
}
fn empty_session_timeout(name: &str, configured_timeout: Duration) -> Duration {
if protocol::is_implicit_session_name(name) {
fn empty_session_timeout(
name: &str,
configured_timeout: Duration,
restart_orphaned: bool,
) -> Duration {
if protocol::is_implicit_session_name(name) && !restart_orphaned {
configured_timeout.min(Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS))
} else {
configured_timeout
@@ -4613,6 +4672,7 @@ mod tests {
empty_since: Some(Instant::now()),
holder_control: None,
persistent: true,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 7,
last_screen_persist_at: Instant::now(),
@@ -4640,6 +4700,7 @@ mod tests {
empty_since: Some(Instant::now()),
holder_control: None,
persistent: true,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 7,
last_screen_persist_at: Instant::now(),
@@ -4687,6 +4748,7 @@ mod tests {
empty_since: Some(Instant::now()),
holder_control: None,
persistent: true,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 7,
last_screen_persist_at: Instant::now(),
@@ -5036,25 +5098,76 @@ mod tests {
}
#[test]
fn implicit_empty_session_timeout_is_bounded_for_update_reconnect() {
fn implicit_empty_session_timeout_preserves_restart_orphans() {
let configured = Duration::from_secs(2_592_000);
let implicit = protocol::generate_implicit_session_name();
assert_eq!(
empty_session_timeout(&implicit, configured),
empty_session_timeout(&implicit, configured, false),
Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS)
);
assert_eq!(
empty_session_timeout("work", configured),
empty_session_timeout(&implicit, configured, true),
configured,
"a sleeping client must retain its restart-surviving shell"
);
assert_eq!(
empty_session_timeout("work", configured, false),
configured,
"named sessions keep the normal long timeout"
);
assert_eq!(
empty_session_timeout(&implicit, Duration::from_secs(1)),
empty_session_timeout(&implicit, Duration::from_secs(1), false),
Duration::from_secs(1),
"tests/admins can still configure a shorter timeout"
);
}
#[test]
fn cleanup_keeps_restart_orphan_then_reaps_after_reattach_disconnect() {
let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(
ServerConfig {
client_timeout_secs: 2_592_000,
prewarm_sessions: Vec::new(),
..ServerConfig::default()
},
[0u8; 32],
pty_tx,
);
let implicit = protocol::generate_implicit_session_name();
state
.ensure_session(&implicit, 80, 24, "forward-only", &[])
.unwrap();
{
let session = state.sessions.get_mut(&implicit).unwrap();
session.empty_since = Some(
Instant::now()
- Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS + 10),
);
session.restart_orphaned = true;
}
let state = Arc::new(Mutex::new(state));
cleanup_disconnected_clients(&state);
assert!(
state.lock().unwrap().sessions.contains_key(&implicit),
"restart orphan was reaped before the reconnect timeout"
);
state
.lock()
.unwrap()
.sessions
.get_mut(&implicit)
.unwrap()
.restart_orphaned = false;
cleanup_disconnected_clients(&state);
assert!(
!state.lock().unwrap().sessions.contains_key(&implicit),
"ordinary abandoned implicit session was not reaped"
);
}
#[tokio::test]
async fn forged_plaintext_detach_does_not_remove_client() {
let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
@@ -5136,6 +5249,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5201,6 +5315,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5278,6 +5393,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5337,6 +5453,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5514,6 +5631,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5618,6 +5736,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5683,6 +5802,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5740,6 +5860,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5796,6 +5917,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5843,6 +5965,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5906,6 +6029,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -5982,6 +6106,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
@@ -6089,6 +6214,7 @@ mod tests {
empty_since: None,
holder_control: None,
persistent: false,
restart_orphaned: false,
bytes_since_persist: 0,
last_persisted_seq: 0,
last_screen_persist_at: Instant::now(),
+33 -2
View File
@@ -31,7 +31,10 @@ use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::net::UdpSocket;
pub const DEFAULT_INITIAL_WINDOW: usize = 1024 * 1024;
/// Per-stream bytes allowed in flight before window credit returns. This bounds
/// each sender burst to 64 MTU-safe packets so bulk streams cannot starve the
/// terminal while retaining useful bandwidth on high-latency links.
pub const DEFAULT_INITIAL_WINDOW: usize = 64 * 1024;
pub const DEFAULT_RETRANSMIT_AFTER: Duration = Duration::from_millis(200);
pub const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10);
pub const ADAPTIVE_RETRANSMIT_MIN: Duration = Duration::from_millis(10);
@@ -39,7 +42,10 @@ pub const DEFAULT_KEEPALIVE_AFTER: Duration = Duration::from_secs(2);
pub const DEFAULT_RETIRED_STREAM_TOMBSTONES: usize = 16 * 1024;
pub const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
pub const SERVICE_TARGET_PREFIX: &str = "@dosh-";
pub const MAX_STREAM_DATA_BYTES: usize = 60 * 1024;
/// Maximum application bytes in one encrypted UDP stream packet. Keeping this
/// aligned with terminal output framing avoids IP fragmentation and stays
/// below macOS route MTUs after Dosh, AEAD, UDP, and IP overhead.
pub const MAX_STREAM_DATA_BYTES: usize = 1024;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TransportConfig {
@@ -1675,6 +1681,31 @@ mod tests {
assert_eq!(third.bytes.len(), 13);
}
#[test]
fn maximum_stream_chunk_fits_a_safe_udp_datagram() {
let body = protocol::to_body(&StreamData {
stream_id: u64::MAX,
offset: u64::MAX,
bytes: vec![0xff; MAX_STREAM_DATA_BYTES],
})
.unwrap();
let packet = protocol::encode_encrypted(
PacketKind::StreamData,
[0xff; 16],
u64::MAX,
u64::MAX,
&[0xff; 32],
CLIENT_TO_SERVER,
&body,
)
.unwrap();
assert!(
packet.len() <= 1200,
"encrypted stream datagram is {} bytes",
packet.len()
);
}
#[test]
fn queued_large_write_flushes_in_window_sized_chunks_after_open() {
let mut mux = StreamMux::new(TransportConfig {
+35
View File
@@ -1027,6 +1027,8 @@ 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<u8> = (0..128 * 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 +1136,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")
@@ -1355,6 +1361,22 @@ fn native_exec_command_smoke() {
let dir = tempfile::tempdir().unwrap();
let port = free_udp_port();
let config = write_server_config(&dir, port);
let shell_log = dir.path().join("exec-shell.log");
let shell = dir.path().join("exec-shell");
fs::write(
&shell,
format!(
"#!/bin/sh\nprintf '%s\\n' \"$*\" >> '{}'\nexec /bin/sh \"$@\"\n",
shell_log.display()
),
)
.unwrap();
fs::set_permissions(&shell, fs::Permissions::from_mode(0o700)).unwrap();
let raw = fs::read_to_string(&config).unwrap().replace(
"shell = \"/bin/sh\"",
&format!("shell = {:?}", shell.display().to_string()),
);
fs::write(&config, raw).unwrap();
write_native_client_auth(&dir, &config);
let mut server = start_server(&dir, &config);
let client_bin = env!("CARGO_BIN_EXE_dosh-client");
@@ -1380,6 +1402,19 @@ fn native_exec_command_smoke() {
);
assert_eq!(String::from_utf8_lossy(&output.stdout), "out");
assert_eq!(String::from_utf8_lossy(&output.stderr), "err");
let shell_invocations = fs::read_to_string(shell_log).unwrap();
assert!(
shell_invocations
.lines()
.any(|line| line.starts_with("-c ")),
"configured shell was not used for exec: {shell_invocations:?}"
);
assert!(
!shell_invocations
.lines()
.any(|line| line.starts_with("-lc ")),
"exec must not start a login shell: {shell_invocations:?}"
);
}
#[test]