From a55efe858cec2c4aaee85e4f97d0c6e7484ce928 Mon Sep 17 00:00:00 2001 From: Jeongkyu Shin Date: Fri, 28 Aug 2026 13:26:54 +0900 Subject: [PATCH 1/5] feat(ssh): enforce RekeyLimit transport policy (#301) - Replace raw RekeyLimit strings with a checked byte/time policy and preserve OpenSSH's `default none` semantics. - Apply the effective policy to russh read, write, and time limits while capping data at the backend's nonce-safety ceiling. - Forward piped stdin for a single raw SSH destination so bounded transfer regressions exercise real transport rekeys without racing multi-host readers. - `cargo test rekey --lib` - focused ssh_config parser, resolver, and runtime registry tests - `cargo check --lib --bins --tests` - `cargo clippy --lib --bins --tests -- -D warnings` - `cargo fmt --all -- --check` - `git diff --check` Refs #301 --- src/executor/parallel.rs | 26 +- src/ssh/client/command.rs | 73 ++++- src/ssh/ssh_config/mod.rs | 36 ++- .../ssh_config/parser/options/connection.rs | 162 +--------- src/ssh/ssh_config/parser/options/support.rs | 5 +- src/ssh/ssh_config/rekey.rs | 303 ++++++++++++++++++ src/ssh/ssh_config/resolver.rs | 2 +- src/ssh/ssh_config/types.rs | 4 +- src/ssh/tokio_client/channel_manager.rs | 125 +++++--- src/ssh/tokio_client/connection.rs | 18 +- src/ssh/tokio_client/connection_tests.rs | 37 ++- src/ssh/tokio_client/session.rs | 60 +++- 12 files changed, 633 insertions(+), 218 deletions(-) create mode 100644 src/ssh/ssh_config/rekey.rs diff --git a/src/executor/parallel.rs b/src/executor/parallel.rs index 948a23d6..13c0edaf 100644 --- a/src/executor/parallel.rs +++ b/src/executor/parallel.rs @@ -1238,6 +1238,10 @@ impl ParallelExecutor { let semaphore = Arc::new(Semaphore::new(self.max_parallel)); let raw_stream = output_mode.is_raw_stream(); + // A single OpenSSH-compatible destination owns piped stdin. Never let + // parallel fan-out race multiple readers over the process input, and + // leave terminal input to the interactive/PTY path. + let forward_stdin = raw_stream && self.nodes.len() == 1 && !std::io::stdin().is_terminal(); let mut manager = MultiNodeStreamManager::new(); let mut handles = Vec::new(); let stdin_is_terminal = std::io::stdin().is_terminal(); @@ -1364,10 +1368,24 @@ impl ParallelExecutor { } } } else { - match client - .connect_and_execute_with_output_streaming(&command, &config, tx.clone()) - .await - { + let execution = if forward_stdin { + client + .connect_and_execute_with_output_streaming_and_stdin( + &command, + &config, + tx.clone(), + ) + .await + } else { + client + .connect_and_execute_with_output_streaming( + &command, + &config, + tx.clone(), + ) + .await + }; + match execution { Ok(exit_status) => { tracing::debug!( "Command completed for {}: exit code {}", diff --git a/src/ssh/client/command.rs b/src/ssh/client/command.rs index 80c7e7ee..ea715362 100644 --- a/src/ssh/client/command.rs +++ b/src/ssh/client/command.rs @@ -254,6 +254,28 @@ impl SshClient { command: &str, config: &ConnectionConfig<'_>, output_sender: Sender, + ) -> Result { + self.connect_and_execute_with_output_streaming_inner(command, config, output_sender, false) + .await + } + + /// Execute a command with byte-transparent stdout/stderr and piped stdin. + pub async fn connect_and_execute_with_output_streaming_and_stdin( + &mut self, + command: &str, + config: &ConnectionConfig<'_>, + output_sender: Sender, + ) -> Result { + self.connect_and_execute_with_output_streaming_inner(command, config, output_sender, true) + .await + } + + async fn connect_and_execute_with_output_streaming_inner( + &mut self, + command: &str, + config: &ConnectionConfig<'_>, + output_sender: Sender, + forward_stdin: bool, ) -> Result { tracing::debug!("Connecting to {}:{}", self.host, self.port); @@ -305,6 +327,7 @@ impl SshClient { config.timeout_seconds, output_sender, config.session_policy, + forward_stdin, ) .await } @@ -321,14 +344,22 @@ impl SshClient { command: &str, output_sender: Sender, session_policy: Option<&crate::ssh::SessionPolicy>, + forward_stdin: bool, ) -> Result { match session_policy { + Some(policy) if forward_stdin => { + client + .execute_session_streaming_with_stdin(policy, output_sender) + .await + } Some(policy) => { client .execute_session_streaming(policy, output_sender) .await } - None => client.execute_streaming(command, output_sender).await, + None => { + Self::execute_streaming_once(client, command, output_sender, forward_stdin).await + } } } @@ -340,12 +371,19 @@ impl SshClient { timeout_seconds: Option, output_sender: Sender, session_policy: Option<&crate::ssh::SessionPolicy>, + forward_stdin: bool, ) -> Result { if let Some(timeout_secs) = timeout_seconds { if timeout_secs == 0 { // No timeout (unlimited) tracing::debug!("Executing command with streaming, no timeout (unlimited)"); - Self::execute_resolved_session_streaming(client, command, output_sender, session_policy) + Self::execute_resolved_session_streaming( + client, + command, + output_sender, + session_policy, + forward_stdin, + ) .await .with_context(|| format!("Failed to execute command '{}' on {}:{}. The SSH connection was successful but the command could not be executed.", command, self.host, self.port)) } else { @@ -357,7 +395,13 @@ impl SshClient { ); tokio::time::timeout( command_timeout, - Self::execute_resolved_session_streaming(client, command, output_sender, session_policy) + Self::execute_resolved_session_streaming( + client, + command, + output_sender, + session_policy, + forward_stdin, + ) ) .await .with_context(|| format!("Command execution timeout: The command '{}' did not complete within {} seconds on {}:{}", command, timeout_secs, self.host, self.port))? @@ -369,7 +413,13 @@ impl SshClient { tracing::debug!("Executing command with streaming, default timeout of 300 seconds"); tokio::time::timeout( command_timeout, - Self::execute_resolved_session_streaming(client, command, output_sender, session_policy) + Self::execute_resolved_session_streaming( + client, + command, + output_sender, + session_policy, + forward_stdin, + ) ) .await .with_context(|| format!("Command execution timeout: The command '{}' did not complete within 5 minutes on {}:{}", command, self.host, self.port))? @@ -377,6 +427,21 @@ impl SshClient { } } + async fn execute_streaming_once( + client: &crate::ssh::tokio_client::Client, + command: &str, + output_sender: Sender, + forward_stdin: bool, + ) -> Result { + if forward_stdin { + client + .execute_streaming_with_stdin(command, output_sender) + .await + } else { + client.execute_streaming(command, output_sender).await + } + } + /// Execute a command with sudo password support and streaming output. /// /// This method handles automatic sudo password injection when sudo prompts are detected diff --git a/src/ssh/ssh_config/mod.rs b/src/ssh/ssh_config/mod.rs index eb8ecece..1d619c56 100644 --- a/src/ssh/ssh_config/mod.rs +++ b/src/ssh/ssh_config/mod.rs @@ -30,6 +30,7 @@ mod match_directive; mod parser; mod path; mod pattern; +mod rekey; mod resolver; #[cfg(test)] mod resolver_tests; @@ -40,6 +41,9 @@ mod types; // Re-export public types pub use ip_qos::{IpQosParseError, IpQosPolicy, IpQosValue}; +pub use rekey::{ + RUSSH_REKEY_BYTE_CEILING, RekeyDataLimit, RekeyLimit, RekeyLimitParseError, RekeyTimeLimit, +}; pub use types::SshHostConfig; /// SSH configuration parser and resolver @@ -635,7 +639,13 @@ Host backup-server bulk: IpQosValue::None, }) ); - assert_eq!(host1.rekey_limit, Some("1G 1h".to_string())); + assert_eq!( + host1.rekey_limit, + Some(RekeyLimit { + data: RekeyDataLimit::Bytes(1 << 30), + time: RekeyTimeLimit::Seconds(3_600), + }) + ); // Verify backup-server config let host2 = &config.hosts[1]; @@ -647,7 +657,13 @@ Host backup-server bulk: IpQosValue::Class(0x48), }) ); - assert_eq!(host2.rekey_limit, Some("default none".to_string())); + assert_eq!( + host2.rekey_limit, + Some(RekeyLimit { + data: RekeyDataLimit::Default, + time: RekeyTimeLimit::None, + }) + ); } #[test] @@ -740,7 +756,13 @@ Host web1.example.com ); // RekeyLimit should be from web1.example.com (most specific) - assert_eq!(host_config.rekey_limit, Some("1G 2h".to_string())); + assert_eq!( + host_config.rekey_limit, + Some(RekeyLimit { + data: RekeyDataLimit::Bytes(1 << 30), + time: RekeyTimeLimit::Seconds(7_200), + }) + ); // ForwardX11Timeout should be from *.example.com assert_eq!(host_config.forward_x11_timeout, Some("30m".to_string())); @@ -917,7 +939,13 @@ Host test RekeyLimit 1G 1h "#; let config = SshConfig::parse(config_content).unwrap(); - assert_eq!(config.hosts[0].rekey_limit, Some("1G 1h".to_string())); + assert_eq!( + config.hosts[0].rekey_limit, + Some(RekeyLimit { + data: RekeyDataLimit::Bytes(1 << 30), + time: RekeyTimeLimit::Seconds(3_600), + }) + ); // Test ForwardX11Timeout - should reject invalid format let config_content = r#" diff --git a/src/ssh/ssh_config/parser/options/connection.rs b/src/ssh/ssh_config/parser/options/connection.rs index 72c70572..5e839142 100644 --- a/src/ssh/ssh_config/parser/options/connection.rs +++ b/src/ssh/ssh_config/parser/options/connection.rs @@ -17,7 +17,7 @@ //! Handles connection-related configuration options including keepalive //! settings, timeouts, compression, and network settings. -use crate::ssh::ssh_config::IpQosPolicy; +use crate::ssh::ssh_config::{IpQosPolicy, RekeyLimit}; use crate::ssh::ssh_config::parser::helpers::parse_yes_no; use crate::ssh::ssh_config::types::SshHostConfig; use anyhow::{Context, Result}; @@ -177,162 +177,20 @@ pub(super) fn parse_connection_option( if args.is_empty() { anyhow::bail!("RekeyLimit requires a value at line {line_number}"); } - // RekeyLimit can have one or two values (data limit and time limit) - // Format: [