Compare commits
23
Commits
v1.0.0-rc43
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
917d0b74b7 | ||
|
|
eb0c1db837 | ||
|
|
f8693f08b5 | ||
|
|
26532fc0e1 | ||
|
|
c5f699a6ef | ||
|
|
60403ba4c3 | ||
|
|
833ac1082f | ||
|
|
97cf165527 | ||
|
|
8f2d57d95e | ||
|
|
24180c5092 | ||
|
|
0fdfc0ee22 | ||
|
|
5dceb2792d | ||
|
|
58ac974fe2 | ||
|
|
0cdcaaaec1 | ||
|
|
9245c666af | ||
|
|
af11ab889e | ||
|
|
43a7a69b9b | ||
|
|
86a1942aa0 | ||
|
|
7a9ad38657 | ||
|
|
9cec63aeb5 | ||
|
|
3364a7eb7b | ||
|
|
970d54b991 | ||
|
|
b8b56c0f32 |
Generated
+1
-1
@@ -436,7 +436,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "dosh"
|
name = "dosh"
|
||||||
version = "1.0.0-rc43"
|
version = "1.0.0-rc49"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "dosh"
|
name = "dosh"
|
||||||
version = "1.0.0-rc43"
|
version = "1.0.0-rc49"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
license = "MIT"
|
license = "MIT"
|
||||||
|
|
||||||
|
|||||||
@@ -16,6 +16,9 @@ async fn main() -> Result<()> {
|
|||||||
client.user, client.session, client.conn_id
|
client.user, client.session, client.conn_id
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
DoshServerEvent::Disconnected(client) => {
|
||||||
|
eprintln!("disconnected conn={:?}", client.conn_id);
|
||||||
|
}
|
||||||
DoshServerEvent::Session {
|
DoshServerEvent::Session {
|
||||||
conn_id,
|
conn_id,
|
||||||
event: SessionEvent::Stream(TransportEvent::Open(open)),
|
event: SessionEvent::Stream(TransportEvent::Open(open)),
|
||||||
|
|||||||
+815
-381
File diff suppressed because it is too large
Load Diff
+137
-11
@@ -48,7 +48,7 @@ use tokio::net::{TcpListener, TcpStream, UdpSocket, UnixListener};
|
|||||||
use tokio::process::Command as TokioCommand;
|
use tokio::process::Command as TokioCommand;
|
||||||
use tokio::sync::mpsc;
|
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_RETIRED_TOMBSTONES: usize = 16 * 1024;
|
||||||
const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
|
const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
|
||||||
|
|
||||||
@@ -398,6 +398,10 @@ struct Session {
|
|||||||
holder_control: Option<StdUnixStream>,
|
holder_control: Option<StdUnixStream>,
|
||||||
/// Whether this session's shell lives in a holder process (persistent).
|
/// Whether this session's shell lives in a holder process (persistent).
|
||||||
persistent: bool,
|
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
|
/// Bytes of session output since the screen was last mirrored to disk, used
|
||||||
/// to throttle the (atomic) screen-persistence writes.
|
/// to throttle the (atomic) screen-persistence writes.
|
||||||
bytes_since_persist: usize,
|
bytes_since_persist: usize,
|
||||||
@@ -569,6 +573,7 @@ impl ServerState {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: control,
|
holder_control: control,
|
||||||
persistent,
|
persistent,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now() - SCREEN_PERSIST_MAX_AGE,
|
last_screen_persist_at: Instant::now() - SCREEN_PERSIST_MAX_AGE,
|
||||||
@@ -714,6 +719,7 @@ impl ServerState {
|
|||||||
empty_since: Some(Instant::now()),
|
empty_since: Some(Instant::now()),
|
||||||
holder_control: Some(control),
|
holder_control: Some(control),
|
||||||
persistent: true,
|
persistent: true,
|
||||||
|
restart_orphaned: true,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: output_seq,
|
last_persisted_seq: output_seq,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -732,6 +738,7 @@ impl ServerState {
|
|||||||
if let Some(session) = self.sessions.get_mut(session_name) {
|
if let Some(session) = self.sessions.get_mut(session_name) {
|
||||||
session.clients.insert(client_id, client);
|
session.clients.insert(client_id, client);
|
||||||
session.empty_since = None;
|
session.empty_since = None;
|
||||||
|
session.restart_orphaned = false;
|
||||||
self.client_index
|
self.client_index
|
||||||
.insert(client_id, session_name.to_string());
|
.insert(client_id, session_name.to_string());
|
||||||
}
|
}
|
||||||
@@ -2082,6 +2089,14 @@ async fn handle_stream_open(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if open.target_host == FILE_STREAM_SENTINEL {
|
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);
|
let (writer_tx, writer_rx) = mpsc::channel::<Vec<u8>>(1024);
|
||||||
register_opened_stream(
|
register_opened_stream(
|
||||||
state,
|
state,
|
||||||
@@ -2106,6 +2121,7 @@ async fn handle_stream_open(
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
|
eprintln!("file service stream {stream_id} failed: {err:#}");
|
||||||
let _ = send_file_response_to_client(
|
let _ = send_file_response_to_client(
|
||||||
&state,
|
&state,
|
||||||
&socket,
|
&socket,
|
||||||
@@ -2457,6 +2473,14 @@ async fn handle_stream_eof(
|
|||||||
let (key, session_name) = find_client_decrypt_key(state, &packet.header)?;
|
let (key, session_name) = find_client_decrypt_key(state, &packet.header)?;
|
||||||
let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?;
|
let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?;
|
||||||
let eof: StreamEof = protocol::from_body(&body)?;
|
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 mut locked = state.lock().expect("server state poisoned");
|
||||||
let Some(session) = locked.sessions.get_mut(&session_name) else {
|
let Some(session) = locked.sessions.get_mut(&session_name) else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
@@ -2488,6 +2512,14 @@ async fn handle_stream_close(
|
|||||||
let (key, session_name) = find_client_decrypt_key(state, &packet.header)?;
|
let (key, session_name) = find_client_decrypt_key(state, &packet.header)?;
|
||||||
let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?;
|
let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?;
|
||||||
let close: StreamClose = protocol::from_body(&body)?;
|
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 mut locked = state.lock().expect("server state poisoned");
|
||||||
let Some(session) = locked.sessions.get_mut(&session_name) else {
|
let Some(session) = locked.sessions.get_mut(&session_name) else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
@@ -2880,6 +2912,7 @@ async fn run_file_stream_service(
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
|
eprintln!("file service request on stream {stream_id} failed: {err:#}");
|
||||||
send_file_response_to_client(
|
send_file_response_to_client(
|
||||||
&state,
|
&state,
|
||||||
&socket,
|
&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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3271,15 +3308,25 @@ async fn run_exec_stream_service(
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
let output = TokioCommand::new("sh")
|
let shell = {
|
||||||
.arg("-lc")
|
state
|
||||||
|
.lock()
|
||||||
|
.expect("server state poisoned")
|
||||||
|
.config
|
||||||
|
.shell
|
||||||
|
.clone()
|
||||||
|
};
|
||||||
|
let output = TokioCommand::new(&shell)
|
||||||
|
.arg("-c")
|
||||||
.arg(&request.command)
|
.arg(&request.command)
|
||||||
.stdin(Stdio::null())
|
.stdin(Stdio::null())
|
||||||
.stdout(Stdio::piped())
|
.stdout(Stdio::piped())
|
||||||
.stderr(Stdio::piped())
|
.stderr(Stdio::piped())
|
||||||
.output()
|
.output()
|
||||||
.await
|
.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) {
|
for chunk in output.stdout.chunks(CHUNK_SIZE) {
|
||||||
send_exec_response_to_client(
|
send_exec_response_to_client(
|
||||||
&state,
|
&state,
|
||||||
@@ -3800,6 +3847,13 @@ async fn send_stream_close_to_client(
|
|||||||
client_id: [u8; 16],
|
client_id: [u8; 16],
|
||||||
stream_id: u64,
|
stream_id: u64,
|
||||||
) -> Result<()> {
|
) -> 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");
|
let mut locked = state.lock().expect("server state poisoned");
|
||||||
if let Some(client) = locked.client_mut(&client_id) {
|
if let Some(client) = locked.client_mut(&client_id) {
|
||||||
@@ -4447,7 +4501,8 @@ fn cleanup_disconnected_clients(state: &Arc<Mutex<ServerState>>) {
|
|||||||
.filter(|(name, session)| {
|
.filter(|(name, session)| {
|
||||||
!prewarm.contains(name.as_str())
|
!prewarm.contains(name.as_str())
|
||||||
&& session.empty_since.is_some_and(|since| {
|
&& 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())
|
.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 {
|
fn empty_session_timeout(
|
||||||
if protocol::is_implicit_session_name(name) {
|
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))
|
configured_timeout.min(Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS))
|
||||||
} else {
|
} else {
|
||||||
configured_timeout
|
configured_timeout
|
||||||
@@ -4613,6 +4672,7 @@ mod tests {
|
|||||||
empty_since: Some(Instant::now()),
|
empty_since: Some(Instant::now()),
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: true,
|
persistent: true,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 7,
|
last_persisted_seq: 7,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -4640,6 +4700,7 @@ mod tests {
|
|||||||
empty_since: Some(Instant::now()),
|
empty_since: Some(Instant::now()),
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: true,
|
persistent: true,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 7,
|
last_persisted_seq: 7,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -4687,6 +4748,7 @@ mod tests {
|
|||||||
empty_since: Some(Instant::now()),
|
empty_since: Some(Instant::now()),
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: true,
|
persistent: true,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 7,
|
last_persisted_seq: 7,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5036,25 +5098,76 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[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 configured = Duration::from_secs(2_592_000);
|
||||||
let implicit = protocol::generate_implicit_session_name();
|
let implicit = protocol::generate_implicit_session_name();
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
empty_session_timeout(&implicit, configured),
|
empty_session_timeout(&implicit, configured, false),
|
||||||
Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS)
|
Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS)
|
||||||
);
|
);
|
||||||
assert_eq!(
|
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,
|
configured,
|
||||||
"named sessions keep the normal long timeout"
|
"named sessions keep the normal long timeout"
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
empty_session_timeout(&implicit, Duration::from_secs(1)),
|
empty_session_timeout(&implicit, Duration::from_secs(1), false),
|
||||||
Duration::from_secs(1),
|
Duration::from_secs(1),
|
||||||
"tests/admins can still configure a shorter timeout"
|
"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]
|
#[tokio::test]
|
||||||
async fn forged_plaintext_detach_does_not_remove_client() {
|
async fn forged_plaintext_detach_does_not_remove_client() {
|
||||||
let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
|
let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
|
||||||
@@ -5136,6 +5249,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5201,6 +5315,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5278,6 +5393,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5337,6 +5453,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5514,6 +5631,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5618,6 +5736,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5683,6 +5802,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5740,6 +5860,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5796,6 +5917,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5843,6 +5965,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5906,6 +6029,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -5982,6 +6106,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
@@ -6089,6 +6214,7 @@ mod tests {
|
|||||||
empty_since: None,
|
empty_since: None,
|
||||||
holder_control: None,
|
holder_control: None,
|
||||||
persistent: false,
|
persistent: false,
|
||||||
|
restart_orphaned: false,
|
||||||
bytes_since_persist: 0,
|
bytes_since_persist: 0,
|
||||||
last_persisted_seq: 0,
|
last_persisted_seq: 0,
|
||||||
last_screen_persist_at: Instant::now(),
|
last_screen_persist_at: Instant::now(),
|
||||||
|
|||||||
+382
-21
@@ -16,8 +16,8 @@ use crate::transport::{
|
|||||||
use crate::udp::{is_transient_udp_error, is_transient_udp_send_error};
|
use crate::udp::{is_transient_udp_error, is_transient_udp_send_error};
|
||||||
use anyhow::{Context, Result, anyhow, bail};
|
use anyhow::{Context, Result, anyhow, bail};
|
||||||
use ed25519_dalek::SigningKey;
|
use ed25519_dalek::SigningKey;
|
||||||
use std::collections::{HashMap, HashSet};
|
use std::collections::{HashMap, HashSet, VecDeque};
|
||||||
use std::net::SocketAddr;
|
use std::net::{IpAddr, SocketAddr};
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
@@ -31,10 +31,13 @@ pub struct DoshServerConfig {
|
|||||||
pub transport: TransportConfig,
|
pub transport: TransportConfig,
|
||||||
pub require_current_user: bool,
|
pub require_current_user: bool,
|
||||||
pub auth_timeout: Duration,
|
pub auth_timeout: Duration,
|
||||||
|
pub max_pending_auth: usize,
|
||||||
|
pub connection_timeout: Duration,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DoshServerConfig {
|
impl DoshServerConfig {
|
||||||
pub fn new(server: ServerConfig) -> Self {
|
pub fn new(server: ServerConfig) -> Self {
|
||||||
|
let connection_timeout = Duration::from_secs(server.client_timeout_secs.max(1));
|
||||||
Self {
|
Self {
|
||||||
server,
|
server,
|
||||||
bind_addr: None,
|
bind_addr: None,
|
||||||
@@ -42,6 +45,8 @@ impl DoshServerConfig {
|
|||||||
transport: TransportConfig::default(),
|
transport: TransportConfig::default(),
|
||||||
require_current_user: true,
|
require_current_user: true,
|
||||||
auth_timeout: Duration::from_secs(30),
|
auth_timeout: Duration::from_secs(30),
|
||||||
|
max_pending_auth: 1024,
|
||||||
|
connection_timeout,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -71,6 +76,21 @@ impl DoshServerConfig {
|
|||||||
self.transport = transport;
|
self.transport = transport;
|
||||||
self
|
self
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn connection_timeout(mut self, timeout: Duration) -> Self {
|
||||||
|
self.connection_timeout = timeout.max(ADAPTIVE_RETRANSMIT_MIN);
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn auth_timeout(mut self, timeout: Duration) -> Self {
|
||||||
|
self.auth_timeout = timeout.max(ADAPTIVE_RETRANSMIT_MIN);
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn max_pending_auth(mut self, max_pending: usize) -> Self {
|
||||||
|
self.max_pending_auth = max_pending.max(1);
|
||||||
|
self
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for DoshServerConfig {
|
impl Default for DoshServerConfig {
|
||||||
@@ -91,6 +111,7 @@ pub struct DoshAccepted {
|
|||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub enum DoshServerEvent {
|
pub enum DoshServerEvent {
|
||||||
Accepted(DoshAccepted),
|
Accepted(DoshAccepted),
|
||||||
|
Disconnected(DoshAccepted),
|
||||||
Session {
|
Session {
|
||||||
conn_id: [u8; 16],
|
conn_id: [u8; 16],
|
||||||
event: SessionEvent,
|
event: SessionEvent,
|
||||||
@@ -107,13 +128,78 @@ struct PendingServerAuth {
|
|||||||
created_at: Instant,
|
created_at: Instant,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
struct AuthRateLimiter {
|
||||||
|
per_minute: u32,
|
||||||
|
max_sources: usize,
|
||||||
|
buckets: HashMap<IpAddr, AuthTokenBucket>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy)]
|
||||||
|
struct AuthTokenBucket {
|
||||||
|
tokens: f64,
|
||||||
|
last_refill: Instant,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AuthRateLimiter {
|
||||||
|
fn new(per_minute: u32, max_sources: usize) -> Self {
|
||||||
|
Self {
|
||||||
|
per_minute,
|
||||||
|
max_sources: max_sources.max(1),
|
||||||
|
buckets: HashMap::new(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn check(&mut self, ip: IpAddr, now: Instant) -> Result<u32, ()> {
|
||||||
|
if self.per_minute == 0 {
|
||||||
|
return Err(());
|
||||||
|
}
|
||||||
|
self.evict_full(now);
|
||||||
|
if !self.buckets.contains_key(&ip) && self.buckets.len() >= self.max_sources {
|
||||||
|
return Err(());
|
||||||
|
}
|
||||||
|
let capacity = self.per_minute as f64;
|
||||||
|
let refill_per_sec = capacity / 60.0;
|
||||||
|
let bucket = self.buckets.entry(ip).or_insert(AuthTokenBucket {
|
||||||
|
tokens: capacity,
|
||||||
|
last_refill: now,
|
||||||
|
});
|
||||||
|
let elapsed = now
|
||||||
|
.saturating_duration_since(bucket.last_refill)
|
||||||
|
.as_secs_f64();
|
||||||
|
bucket.tokens = (bucket.tokens + elapsed * refill_per_sec).min(capacity);
|
||||||
|
bucket.last_refill = now;
|
||||||
|
if bucket.tokens < 1.0 {
|
||||||
|
return Err(());
|
||||||
|
}
|
||||||
|
bucket.tokens -= 1.0;
|
||||||
|
Ok(bucket.tokens as u32)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn evict_full(&mut self, now: Instant) {
|
||||||
|
if self.per_minute == 0 {
|
||||||
|
self.buckets.clear();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let capacity = self.per_minute as f64;
|
||||||
|
let refill_per_sec = capacity / 60.0;
|
||||||
|
self.buckets.retain(|_, bucket| {
|
||||||
|
let elapsed = now
|
||||||
|
.saturating_duration_since(bucket.last_refill)
|
||||||
|
.as_secs_f64();
|
||||||
|
(bucket.tokens + elapsed * refill_per_sec) < capacity
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub struct DoshServer {
|
pub struct DoshServer {
|
||||||
socket: Arc<UdpSocket>,
|
socket: Arc<UdpSocket>,
|
||||||
config: DoshServerConfig,
|
config: DoshServerConfig,
|
||||||
host_signing: SigningKey,
|
host_signing: SigningKey,
|
||||||
|
auth_limiter: AuthRateLimiter,
|
||||||
pending: HashMap<[u8; 16], PendingServerAuth>,
|
pending: HashMap<[u8; 16], PendingServerAuth>,
|
||||||
transports: HashMap<[u8; 16], DoshTransport>,
|
transports: HashMap<[u8; 16], DoshTransport>,
|
||||||
accepted: HashMap<[u8; 16], DoshAccepted>,
|
accepted: HashMap<[u8; 16], DoshAccepted>,
|
||||||
|
disconnected: VecDeque<DoshAccepted>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DoshServer {
|
impl DoshServer {
|
||||||
@@ -139,13 +225,19 @@ impl DoshServer {
|
|||||||
};
|
};
|
||||||
let host_signing = load_or_create_host_key(&config.server)?;
|
let host_signing = load_or_create_host_key(&config.server)?;
|
||||||
let socket = Arc::new(UdpSocket::bind(bind_addr).await?);
|
let socket = Arc::new(UdpSocket::bind(bind_addr).await?);
|
||||||
|
let auth_limiter = AuthRateLimiter::new(
|
||||||
|
config.server.native_auth_rate_limit_per_minute,
|
||||||
|
config.max_pending_auth,
|
||||||
|
);
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
socket,
|
socket,
|
||||||
config,
|
config,
|
||||||
host_signing,
|
host_signing,
|
||||||
|
auth_limiter,
|
||||||
pending: HashMap::new(),
|
pending: HashMap::new(),
|
||||||
transports: HashMap::new(),
|
transports: HashMap::new(),
|
||||||
accepted: HashMap::new(),
|
accepted: HashMap::new(),
|
||||||
|
disconnected: VecDeque::new(),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -165,10 +257,19 @@ impl DoshServer {
|
|||||||
self.transports.get_mut(conn_id)
|
self.transports.get_mut(conn_id)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn remove_connection(&mut self, conn_id: &[u8; 16]) -> Option<DoshAccepted> {
|
||||||
|
self.transports.remove(conn_id);
|
||||||
|
self.accepted.remove(conn_id)
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn recv(&mut self) -> Result<DoshServerEvent> {
|
pub async fn recv(&mut self) -> Result<DoshServerEvent> {
|
||||||
let mut buf = vec![0u8; 65535];
|
let mut buf = vec![0u8; 65535];
|
||||||
loop {
|
loop {
|
||||||
self.expire_pending();
|
self.expire_pending();
|
||||||
|
self.expire_connections();
|
||||||
|
if let Some(connection) = self.disconnected.pop_front() {
|
||||||
|
return Ok(DoshServerEvent::Disconnected(connection));
|
||||||
|
}
|
||||||
let (n, peer) = match tokio::time::timeout(
|
let (n, peer) = match tokio::time::timeout(
|
||||||
ADAPTIVE_RETRANSMIT_MIN,
|
ADAPTIVE_RETRANSMIT_MIN,
|
||||||
self.socket.recv_from(&mut buf),
|
self.socket.recv_from(&mut buf),
|
||||||
@@ -207,6 +308,9 @@ impl DoshServer {
|
|||||||
}
|
}
|
||||||
if let Some(transport) = self.transports.get_mut(&packet.header.conn_id) {
|
if let Some(transport) = self.transports.get_mut(&packet.header.conn_id) {
|
||||||
let event = transport.handle_datagram(datagram, peer).await?;
|
let event = transport.handle_datagram(datagram, peer).await?;
|
||||||
|
if let Some(accepted) = self.accepted.get_mut(&packet.header.conn_id) {
|
||||||
|
accepted.peer_addr = transport.peer_addr();
|
||||||
|
}
|
||||||
return Ok(Some(DoshServerEvent::Session {
|
return Ok(Some(DoshServerEvent::Session {
|
||||||
conn_id: packet.header.conn_id,
|
conn_id: packet.header.conn_id,
|
||||||
event,
|
event,
|
||||||
@@ -270,7 +374,21 @@ impl DoshServer {
|
|||||||
self.send_reject(peer, [0u8; 16], &err.to_string()).await?;
|
self.send_reject(peer, [0u8; 16], &err.to_string()).await?;
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
let result = self.build_server_hello(req.hello, peer);
|
self.expire_pending();
|
||||||
|
if self.pending.len() >= self.config.max_pending_auth {
|
||||||
|
self.send_reject(peer, [0u8; 16], "native auth server busy")
|
||||||
|
.await?;
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
let rate_limit_remaining = match self.auth_limiter.check(peer.ip(), Instant::now()) {
|
||||||
|
Ok(remaining) => remaining,
|
||||||
|
Err(()) => {
|
||||||
|
self.send_reject(peer, [0u8; 16], "native auth rate limit exceeded")
|
||||||
|
.await?;
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let result = self.build_server_hello(req.hello, peer, Some(rate_limit_remaining));
|
||||||
let (pending_id, hello) = match result {
|
let (pending_id, hello) = match result {
|
||||||
Ok(value) => value,
|
Ok(value) => value,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -288,6 +406,7 @@ impl DoshServer {
|
|||||||
&mut self,
|
&mut self,
|
||||||
client: native::NativeClientHello,
|
client: native::NativeClientHello,
|
||||||
peer: SocketAddr,
|
peer: SocketAddr,
|
||||||
|
rate_limit_remaining: Option<u32>,
|
||||||
) -> Result<([u8; 16], NativeServerHello)> {
|
) -> Result<([u8; 16], NativeServerHello)> {
|
||||||
if !self.config.server.native_auth {
|
if !self.config.server.native_auth {
|
||||||
bail!("native auth disabled");
|
bail!("native auth disabled");
|
||||||
@@ -307,11 +426,14 @@ impl DoshServer {
|
|||||||
bail!("native auth requires a supported user key algorithm");
|
bail!("native auth requires a supported user key algorithm");
|
||||||
}
|
}
|
||||||
if self.config.require_current_user {
|
if self.config.require_current_user {
|
||||||
let current_user = std::env::var("USER").unwrap_or_else(|_| "unknown".to_string());
|
let current_user = local_username();
|
||||||
if client.requested_user != current_user {
|
if client.requested_user != current_user {
|
||||||
bail!("native auth user mismatch");
|
bail!("native auth user mismatch");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if self.pending.len() >= self.config.max_pending_auth {
|
||||||
|
bail!("native auth server busy");
|
||||||
|
}
|
||||||
|
|
||||||
let (server_secret, server_public) = generate_native_ephemeral();
|
let (server_secret, server_public) = generate_native_ephemeral();
|
||||||
let mut server = NativeServerHello {
|
let mut server = NativeServerHello {
|
||||||
@@ -322,7 +444,7 @@ impl DoshServer {
|
|||||||
chosen_aead: "chacha20poly1305".to_string(),
|
chosen_aead: "chacha20poly1305".to_string(),
|
||||||
server_key_epoch: 1,
|
server_key_epoch: 1,
|
||||||
auth_challenge: crypto::random_32(),
|
auth_challenge: crypto::random_32(),
|
||||||
rate_limit_remaining: None,
|
rate_limit_remaining,
|
||||||
host_signature: Vec::new(),
|
host_signature: Vec::new(),
|
||||||
};
|
};
|
||||||
sign_server_hello(&self.host_signing, &client, &mut server)?;
|
sign_server_hello(&self.host_signing, &client, &mut server)?;
|
||||||
@@ -477,7 +599,37 @@ impl DoshServer {
|
|||||||
let timeout = self.config.auth_timeout;
|
let timeout = self.config.auth_timeout;
|
||||||
self.pending
|
self.pending
|
||||||
.retain(|_, pending| pending.created_at.elapsed() <= timeout);
|
.retain(|_, pending| pending.created_at.elapsed() <= timeout);
|
||||||
|
self.auth_limiter.evict_full(Instant::now());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn expire_connections(&mut self) {
|
||||||
|
let timeout = self.config.connection_timeout;
|
||||||
|
let expired = self
|
||||||
|
.transports
|
||||||
|
.iter()
|
||||||
|
.filter_map(|(conn_id, transport)| {
|
||||||
|
(transport.stale_for() > timeout).then_some(*conn_id)
|
||||||
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
for conn_id in expired {
|
||||||
|
self.transports.remove(&conn_id);
|
||||||
|
if let Some(connection) = self.accepted.remove(&conn_id) {
|
||||||
|
self.disconnected.push_back(connection);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn local_username() -> String {
|
||||||
|
local_username_from_env(|name| std::env::var(name).ok())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn local_username_from_env(get: impl FnMut(&str) -> Option<String>) -> String {
|
||||||
|
["USER", "USERNAME"]
|
||||||
|
.into_iter()
|
||||||
|
.filter_map(get)
|
||||||
|
.find(|value| !value.is_empty())
|
||||||
|
.unwrap_or_else(|| "unknown".to_string())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn send_udp(socket: &UdpSocket, packet: &[u8], peer: SocketAddr) -> Result<bool> {
|
async fn send_udp(socket: &UdpSocket, packet: &[u8], peer: SocketAddr) -> Result<bool> {
|
||||||
@@ -530,6 +682,24 @@ mod tests {
|
|||||||
use crate::transport::TransportEvent;
|
use crate::transport::TransportEvent;
|
||||||
use ed25519_dalek::SigningKey;
|
use ed25519_dalek::SigningKey;
|
||||||
|
|
||||||
|
fn native_client_hello(public: [u8; 32]) -> native::NativeClientHello {
|
||||||
|
native::NativeClientHello {
|
||||||
|
protocol_version: native::NATIVE_PROTOCOL_VERSION,
|
||||||
|
client_random: crypto::random_32(),
|
||||||
|
client_ephemeral_public: public,
|
||||||
|
requested_host: "127.0.0.1".to_string(),
|
||||||
|
requested_user: "sdk-user".to_string(),
|
||||||
|
requested_session: "test".to_string(),
|
||||||
|
requested_mode: "forward-only".to_string(),
|
||||||
|
terminal_size: (80, 24),
|
||||||
|
supported_aead: vec!["chacha20poly1305".to_string()],
|
||||||
|
supported_user_key_algorithms: vec!["ssh-ed25519".to_string()],
|
||||||
|
cached_host_key_fingerprint: None,
|
||||||
|
attach_ticket_envelope: None,
|
||||||
|
requested_env: Vec::new(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn sdk_client_and_server_exchange_service_stream() {
|
async fn sdk_client_and_server_exchange_service_stream() {
|
||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
@@ -752,6 +922,209 @@ mod tests {
|
|||||||
assert_eq!(open.target_host, "@dosh-echo");
|
assert_eq!(open.target_host, "@dosh-echo");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn sdk_server_expires_disconnected_transport_and_reports_it() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let server_config = ServerConfig {
|
||||||
|
host_key: dir.path().join("host_key").to_string_lossy().to_string(),
|
||||||
|
..ServerConfig::default()
|
||||||
|
};
|
||||||
|
let server_config = DoshServerConfig::new(server_config)
|
||||||
|
.bind_addr("127.0.0.1:0".parse().unwrap())
|
||||||
|
.require_current_user(false)
|
||||||
|
.connection_timeout(Duration::from_millis(30));
|
||||||
|
let mut server = DoshServer::bind(server_config).await.unwrap();
|
||||||
|
let peer_socket = UdpSocket::bind("127.0.0.1:0").await.unwrap();
|
||||||
|
let peer_addr = peer_socket.local_addr().unwrap();
|
||||||
|
let conn_id = [81u8; 16];
|
||||||
|
let transport = DoshTransport::new(
|
||||||
|
Arc::clone(&server.socket),
|
||||||
|
SessionTransportConfig {
|
||||||
|
role: SessionRole::Server,
|
||||||
|
conn_id,
|
||||||
|
session_key: [82u8; 32],
|
||||||
|
peer_addr,
|
||||||
|
initial_send_seq: 1,
|
||||||
|
initial_ack: 0,
|
||||||
|
stream: server.config.transport.clone(),
|
||||||
|
},
|
||||||
|
);
|
||||||
|
let accepted = DoshAccepted {
|
||||||
|
conn_id,
|
||||||
|
user: "sdk-user".to_string(),
|
||||||
|
session: "mobile".to_string(),
|
||||||
|
services: vec!["echo".to_string()],
|
||||||
|
peer_addr,
|
||||||
|
};
|
||||||
|
server.transports.insert(conn_id, transport);
|
||||||
|
server.accepted.insert(conn_id, accepted.clone());
|
||||||
|
|
||||||
|
tokio::time::sleep(Duration::from_millis(40)).await;
|
||||||
|
let event = tokio::time::timeout(Duration::from_secs(1), server.recv())
|
||||||
|
.await
|
||||||
|
.expect("server did not report expired connection")
|
||||||
|
.unwrap();
|
||||||
|
match event {
|
||||||
|
DoshServerEvent::Disconnected(disconnected) => {
|
||||||
|
assert_eq!(disconnected.conn_id, conn_id);
|
||||||
|
assert_eq!(disconnected.session, "mobile");
|
||||||
|
}
|
||||||
|
other => panic!("unexpected expiry event {other:?}"),
|
||||||
|
}
|
||||||
|
assert!(server.connection(&conn_id).is_none());
|
||||||
|
assert!(server.transport(&conn_id).is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn sdk_server_updates_connection_metadata_after_authenticated_roam() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let server_config = ServerConfig {
|
||||||
|
host_key: dir.path().join("host_key").to_string_lossy().to_string(),
|
||||||
|
..ServerConfig::default()
|
||||||
|
};
|
||||||
|
let server_config = DoshServerConfig::new(server_config)
|
||||||
|
.bind_addr("127.0.0.1:0".parse().unwrap())
|
||||||
|
.require_current_user(false);
|
||||||
|
let mut server = DoshServer::bind(server_config).await.unwrap();
|
||||||
|
let original = UdpSocket::bind("127.0.0.1:0").await.unwrap();
|
||||||
|
let roaming = UdpSocket::bind("127.0.0.1:0").await.unwrap();
|
||||||
|
let original_addr = original.local_addr().unwrap();
|
||||||
|
let roaming_addr = roaming.local_addr().unwrap();
|
||||||
|
let conn_id = [83u8; 16];
|
||||||
|
let session_key = [84u8; 32];
|
||||||
|
let transport = DoshTransport::new(
|
||||||
|
Arc::clone(&server.socket),
|
||||||
|
SessionTransportConfig {
|
||||||
|
role: SessionRole::Server,
|
||||||
|
conn_id,
|
||||||
|
session_key,
|
||||||
|
peer_addr: original_addr,
|
||||||
|
initial_send_seq: 1,
|
||||||
|
initial_ack: 0,
|
||||||
|
stream: server.config.transport.clone(),
|
||||||
|
},
|
||||||
|
);
|
||||||
|
server.transports.insert(conn_id, transport);
|
||||||
|
server.accepted.insert(
|
||||||
|
conn_id,
|
||||||
|
DoshAccepted {
|
||||||
|
conn_id,
|
||||||
|
user: "sdk-user".to_string(),
|
||||||
|
session: "mobile".to_string(),
|
||||||
|
services: Vec::new(),
|
||||||
|
peer_addr: original_addr,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
let ping = protocol::encode_encrypted(
|
||||||
|
PacketKind::Ping,
|
||||||
|
conn_id,
|
||||||
|
1,
|
||||||
|
0,
|
||||||
|
&session_key,
|
||||||
|
CLIENT_TO_SERVER,
|
||||||
|
b"",
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
roaming
|
||||||
|
.send_to(&ping, server.local_addr().unwrap())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let event = tokio::time::timeout(Duration::from_secs(1), server.recv())
|
||||||
|
.await
|
||||||
|
.expect("server did not receive roaming ping")
|
||||||
|
.unwrap();
|
||||||
|
assert!(matches!(
|
||||||
|
event,
|
||||||
|
DoshServerEvent::Session {
|
||||||
|
event: SessionEvent::Ping,
|
||||||
|
..
|
||||||
|
}
|
||||||
|
));
|
||||||
|
assert_eq!(server.connection(&conn_id).unwrap().peer_addr, roaming_addr);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn sdk_server_bounds_and_rate_limits_pending_authentication() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let server_config = ServerConfig {
|
||||||
|
host_key: dir.path().join("host_key").to_string_lossy().to_string(),
|
||||||
|
native_auth_rate_limit_per_minute: 1,
|
||||||
|
..ServerConfig::default()
|
||||||
|
};
|
||||||
|
let server_config = DoshServerConfig::new(server_config)
|
||||||
|
.bind_addr("127.0.0.1:0".parse().unwrap())
|
||||||
|
.require_current_user(false)
|
||||||
|
.max_pending_auth(1);
|
||||||
|
let mut server = DoshServer::bind(server_config).await.unwrap();
|
||||||
|
let client = UdpSocket::bind("127.0.0.1:0").await.unwrap();
|
||||||
|
let peer = client.local_addr().unwrap();
|
||||||
|
let (_, public) = native::generate_native_ephemeral();
|
||||||
|
let body = protocol::to_body(&NativeClientHelloBody {
|
||||||
|
hello: native_client_hello(public),
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
server
|
||||||
|
.handle_client_hello(peer, body.clone())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let mut packet = [0u8; 65535];
|
||||||
|
let (n, _) = client.recv_from(&mut packet).await.unwrap();
|
||||||
|
let response = protocol::decode(&packet[..n]).unwrap();
|
||||||
|
assert_eq!(response.header.kind, PacketKind::NativeServerHello);
|
||||||
|
let hello: NativeServerHelloBody = protocol::from_body(&response.body).unwrap();
|
||||||
|
assert_eq!(hello.hello.rate_limit_remaining, Some(0));
|
||||||
|
assert_eq!(server.pending.len(), 1);
|
||||||
|
|
||||||
|
server
|
||||||
|
.handle_client_hello(peer, body.clone())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let (n, _) = client.recv_from(&mut packet).await.unwrap();
|
||||||
|
let response = protocol::decode(&packet[..n]).unwrap();
|
||||||
|
assert_eq!(response.header.kind, PacketKind::AttachReject);
|
||||||
|
let reject: AttachReject = protocol::from_body(&response.body).unwrap();
|
||||||
|
assert_eq!(reject.reason, "native auth server busy");
|
||||||
|
assert_eq!(server.pending.len(), 1);
|
||||||
|
|
||||||
|
server.pending.clear();
|
||||||
|
server.handle_client_hello(peer, body).await.unwrap();
|
||||||
|
let (n, _) = client.recv_from(&mut packet).await.unwrap();
|
||||||
|
let response = protocol::decode(&packet[..n]).unwrap();
|
||||||
|
assert_eq!(response.header.kind, PacketKind::AttachReject);
|
||||||
|
let reject: AttachReject = protocol::from_body(&response.body).unwrap();
|
||||||
|
assert_eq!(reject.reason, "native auth rate limit exceeded");
|
||||||
|
assert!(server.pending.is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn sdk_auth_rate_limiter_bounds_source_tracking_and_refills() {
|
||||||
|
let now = Instant::now();
|
||||||
|
let mut limiter = AuthRateLimiter::new(2, 1);
|
||||||
|
let first: IpAddr = "192.0.2.1".parse().unwrap();
|
||||||
|
let second: IpAddr = "192.0.2.2".parse().unwrap();
|
||||||
|
|
||||||
|
assert_eq!(limiter.check(first, now), Ok(1));
|
||||||
|
assert_eq!(limiter.check(first, now), Ok(0));
|
||||||
|
assert_eq!(limiter.check(first, now), Err(()));
|
||||||
|
assert_eq!(limiter.check(second, now), Err(()));
|
||||||
|
assert_eq!(limiter.check(second, now + Duration::from_secs(60)), Ok(1));
|
||||||
|
assert_eq!(limiter.buckets.len(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn sdk_server_username_supports_windows_environment() {
|
||||||
|
assert_eq!(
|
||||||
|
local_username_from_env(|name| match name {
|
||||||
|
"USER" => None,
|
||||||
|
"USERNAME" => Some("palav-win".to_string()),
|
||||||
|
_ => None,
|
||||||
|
}),
|
||||||
|
"palav-win"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn bad_native_auth_does_not_consume_pending_challenge() {
|
async fn bad_native_auth_does_not_consume_pending_challenge() {
|
||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
@@ -778,22 +1151,10 @@ mod tests {
|
|||||||
let mut server = DoshServer::bind(server_config).await.unwrap();
|
let mut server = DoshServer::bind(server_config).await.unwrap();
|
||||||
let peer: SocketAddr = "127.0.0.1:9".parse().unwrap();
|
let peer: SocketAddr = "127.0.0.1:9".parse().unwrap();
|
||||||
let (client_secret, client_public) = native::generate_native_ephemeral();
|
let (client_secret, client_public) = native::generate_native_ephemeral();
|
||||||
let hello = native::NativeClientHello {
|
let hello = native_client_hello(client_public);
|
||||||
protocol_version: native::NATIVE_PROTOCOL_VERSION,
|
let (pending_id, server_hello) = server
|
||||||
client_random: crypto::random_32(),
|
.build_server_hello(hello.clone(), peer, None)
|
||||||
client_ephemeral_public: client_public,
|
.unwrap();
|
||||||
requested_host: "127.0.0.1".to_string(),
|
|
||||||
requested_user: "sdk-user".to_string(),
|
|
||||||
requested_session: "test".to_string(),
|
|
||||||
requested_mode: "forward-only".to_string(),
|
|
||||||
terminal_size: (80, 24),
|
|
||||||
supported_aead: vec!["chacha20poly1305".to_string()],
|
|
||||||
supported_user_key_algorithms: vec!["ssh-ed25519".to_string()],
|
|
||||||
cached_host_key_fingerprint: None,
|
|
||||||
attach_ticket_envelope: None,
|
|
||||||
requested_env: Vec::new(),
|
|
||||||
};
|
|
||||||
let (pending_id, server_hello) = server.build_server_hello(hello.clone(), peer).unwrap();
|
|
||||||
let session_key = native::derive_native_session_key(
|
let session_key = native::derive_native_session_key(
|
||||||
&client_secret,
|
&client_secret,
|
||||||
server_hello.server_ephemeral_public,
|
server_hello.server_ephemeral_public,
|
||||||
|
|||||||
+33
-2
@@ -31,7 +31,10 @@ use std::sync::Arc;
|
|||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
use tokio::net::UdpSocket;
|
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 DEFAULT_RETRANSMIT_AFTER: Duration = Duration::from_millis(200);
|
||||||
pub const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10);
|
pub const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10);
|
||||||
pub const ADAPTIVE_RETRANSMIT_MIN: 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 DEFAULT_RETIRED_STREAM_TOMBSTONES: usize = 16 * 1024;
|
||||||
pub const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
|
pub const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
|
||||||
pub const SERVICE_TARGET_PREFIX: &str = "@dosh-";
|
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)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct TransportConfig {
|
pub struct TransportConfig {
|
||||||
@@ -1675,6 +1681,31 @@ mod tests {
|
|||||||
assert_eq!(third.bytes.len(), 13);
|
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]
|
#[test]
|
||||||
fn queued_large_write_flushes_in_window_sized_chunks_after_open() {
|
fn queued_large_write_flushes_in_window_sized_chunks_after_open() {
|
||||||
let mut mux = StreamMux::new(TransportConfig {
|
let mut mux = StreamMux::new(TransportConfig {
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
@@ -1027,6 +1027,8 @@ fn native_file_copy_recursive_round_trip() {
|
|||||||
fs::create_dir_all(src.join("nested")).unwrap();
|
fs::create_dir_all(src.join("nested")).unwrap();
|
||||||
fs::write(src.join("root.txt"), b"root file\n").unwrap();
|
fs::write(src.join("root.txt"), b"root file\n").unwrap();
|
||||||
fs::write(src.join("nested/child.txt"), b"child 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("root.txt", src.join("root-link")).unwrap();
|
||||||
std::os::unix::fs::symlink("nested/child.txt", src.join("child-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(),
|
fs::read_to_string(downloaded.join("nested/child.txt")).unwrap(),
|
||||||
"child file\n"
|
"child file\n"
|
||||||
);
|
);
|
||||||
|
assert_eq!(
|
||||||
|
fs::read(downloaded.join("large.bin")).unwrap(),
|
||||||
|
large_payload
|
||||||
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
fs::read_link(downloaded.join("root-link")).unwrap(),
|
fs::read_link(downloaded.join("root-link")).unwrap(),
|
||||||
PathBuf::from("root.txt")
|
PathBuf::from("root.txt")
|
||||||
@@ -1355,6 +1361,22 @@ fn native_exec_command_smoke() {
|
|||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
let port = free_udp_port();
|
let port = free_udp_port();
|
||||||
let config = write_server_config(&dir, 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);
|
write_native_client_auth(&dir, &config);
|
||||||
let mut server = start_server(&dir, &config);
|
let mut server = start_server(&dir, &config);
|
||||||
let client_bin = env!("CARGO_BIN_EXE_dosh-client");
|
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.stdout), "out");
|
||||||
assert_eq!(String::from_utf8_lossy(&output.stderr), "err");
|
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]
|
#[test]
|
||||||
|
|||||||
Reference in New Issue
Block a user