1use super::*;
2
3#[cfg(unix)]
4use std::fs;
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::sync::{Condvar, Mutex, OnceLock};
7use std::time::Duration;
8
9pub fn ssh_connectivity_probe(ssh: &SshTarget) -> CommandSpec {
24 let mut args = vec![
25 "-o".to_owned(),
26 "BatchMode=yes".to_owned(),
27 "-o".to_owned(),
28 "StrictHostKeyChecking=yes".to_owned(),
29 ];
30 args.extend(ssh.ssh_args.iter().cloned());
31 args.push(ssh.destination.clone());
32 args.push(join_remote_command(&["true".to_owned()]));
33 CommandSpec::new("ssh", args)
36 .ssh_probe_session(ssh)
37 .purpose("verify SSH connectivity")
38}
39
40pub fn ssh_command(
41 ssh: &SshTarget,
42 args: impl IntoIterator<Item = impl AsRef<str>>,
43) -> CommandSpec {
44 ssh_command_owned(
45 ssh,
46 args.into_iter()
47 .map(|arg| arg.as_ref().to_owned())
48 .collect(),
49 )
50}
51
52pub fn ssh_command_owned(ssh: &SshTarget, remote_args: Vec<String>) -> CommandSpec {
56 let mut args = ssh.ssh_args.clone();
57 args.push(ssh.destination.clone());
58 args.push(join_remote_command(&remote_args));
59 CommandSpec::new("ssh", args).ssh_session(ssh)
60}
61
62pub const REMOTE_UPLOAD_STAGING: &str = ".cache/mjolnir/uploads";
66
67pub fn scp_upload(ssh: &SshTarget, source: &Path, remote: &str, recursive: bool) -> CommandSpec {
69 let mut args = scp_args(ssh);
70 if recursive {
71 args.push("-r".into());
72 }
73 args.push(source.to_string_lossy().into_owned());
74 args.push(format!("{}:{remote}", ssh.destination));
75 scp_command(ssh, args)
76}
77
78pub fn scp_download(ssh: &SshTarget, remote: &str, local: &str) -> CommandSpec {
80 let mut args = scp_args(ssh);
81 args.push(format!("{}:{remote}", ssh.destination));
82 args.push(local.into());
83 scp_command(ssh, args)
84}
85
86fn scp_args(ssh: &SshTarget) -> Vec<String> {
89 ssh.ssh_args
90 .iter()
91 .map(|argument| {
92 if argument == "-p" {
93 "-P".to_owned()
94 } else {
95 argument.clone()
96 }
97 })
98 .collect()
99}
100
101fn scp_command(ssh: &SshTarget, args: Vec<String>) -> CommandSpec {
102 CommandSpec::new("scp", args).ssh_session(ssh)
105}
106
107#[cfg(unix)]
112const CONTROL_PERSIST: &str = "60";
113
114pub const CONTROL_MASTER_ENV: &str = "MJ_SSH_CONTROL_MASTER";
117
118#[cfg(unix)]
121const MAX_CONTROL_PATH: usize = 103;
122
123#[cfg(unix)]
126const CONNECTION_HASH_HEX: usize = 16;
127
128#[cfg(unix)]
133const CONTROL_SOCKET_NAME_RESERVE: usize = 1 + CONNECTION_HASH_HEX + 1 + 4 + 17;
134
135#[cfg(unix)]
136fn sharing_disabled(value: Option<&std::ffi::OsStr>) -> bool {
137 let Some(value) = value else {
138 return false;
139 };
140 matches!(
141 value.to_string_lossy().trim().to_ascii_lowercase().as_str(),
142 "0" | "off" | "false" | "no"
143 )
144}
145
146#[doc(hidden)]
149#[derive(Debug, Clone)]
150pub enum SshSharingForTest {
151 Disabled,
153 Directory(PathBuf),
155}
156
157static SHARING_OVERRIDE: Mutex<Option<SshSharingForTest>> = Mutex::new(None);
158
159#[cfg(all(test, unix))]
162pub(super) static SHARING_TEST_LOCK: Mutex<()> = Mutex::new(());
163
164#[doc(hidden)]
169pub fn set_ssh_connection_sharing_for_test(setting: Option<SshSharingForTest>) {
170 *SHARING_OVERRIDE
171 .lock()
172 .unwrap_or_else(std::sync::PoisonError::into_inner) = setting;
173}
174
175#[cfg(unix)]
176fn sharing_override() -> Option<SshSharingForTest> {
177 SHARING_OVERRIDE
178 .lock()
179 .unwrap_or_else(std::sync::PoisonError::into_inner)
180 .clone()
181}
182
183#[cfg(unix)]
194fn control_socket_dir() -> Option<PathBuf> {
195 match sharing_override() {
196 Some(SshSharingForTest::Disabled) => return None,
197 Some(SshSharingForTest::Directory(dir)) => return prepare_control_dir(dir),
198 None => {}
199 }
200 static DIR: OnceLock<Option<PathBuf>> = OnceLock::new();
201 DIR.get_or_init(|| {
202 if sharing_disabled(std::env::var_os(CONTROL_MASTER_ENV).as_deref()) {
203 return None;
204 }
205 let data_dir_override = crate::config::env_override_os("DATA_DIR").map(PathBuf::from);
206 let identity = control_dir_identity(
207 crate::config::instance_name().as_deref(),
208 data_dir_override.as_deref(),
209 );
210 prepare_control_dir(default_control_dir(
211 std::env::var_os("XDG_RUNTIME_DIR"),
212 &identity,
213 ))
214 })
215 .clone()
216}
217
218#[cfg(unix)]
220fn default_control_dir(runtime: Option<std::ffi::OsString>, identity: &str) -> PathBuf {
221 match runtime {
222 Some(runtime) if !runtime.is_empty() => {
223 PathBuf::from(runtime).join("mjolnir").join(identity)
224 }
225 _ => crate::config::data_dir().join("ssh"),
227 }
228}
229
230#[cfg(unix)]
238fn control_dir_identity(instance: Option<&str>, data_dir_override: Option<&Path>) -> String {
239 if instance.is_none() && data_dir_override.is_none() {
240 return "default".to_owned();
241 }
242 let data_dir = data_dir_override.map_or_else(crate::config::data_dir, Path::to_path_buf);
243 crate::config::instance_identity_for(instance, &data_dir)
244}
245
246#[cfg(unix)]
249fn prepare_control_dir(dir: PathBuf) -> Option<PathBuf> {
250 if dir.as_os_str().len() + CONTROL_SOCKET_NAME_RESERVE > MAX_CONTROL_PATH {
251 tracing::debug!(
252 directory = %dir.display(),
253 "skipping SSH connection sharing: control socket path would be too long"
254 );
255 return None;
256 }
257 if let Err(error) = fs::create_dir_all(&dir) {
258 tracing::debug!(
259 directory = %dir.display(),
260 %error,
261 "skipping SSH connection sharing: control directory is unavailable"
262 );
263 return None;
264 }
265 use std::os::unix::fs::PermissionsExt;
266 if let Err(error) = fs::set_permissions(&dir, fs::Permissions::from_mode(0o700)) {
267 tracing::debug!(
268 directory = %dir.display(),
269 %error,
270 "skipping SSH connection sharing: cannot restrict control directory"
271 );
272 return None;
273 }
274 Some(dir)
275}
276
277#[cfg(unix)]
281fn connection_key(ssh: &SshTarget) -> String {
282 let mut key = ssh.destination.clone();
283 for argument in &ssh.ssh_args {
284 key.push('\0');
285 key.push_str(argument);
286 }
287 key
288}
289
290#[cfg(unix)]
296fn control_socket_name(ssh: &SshTarget, shard: usize) -> String {
297 use sha2::{Digest, Sha256};
298 let digest = Sha256::digest(connection_key(ssh).as_bytes());
299 let mut name = String::with_capacity(CONNECTION_HASH_HEX + 5);
300 for byte in digest.iter().take(CONNECTION_HASH_HEX / 2) {
301 name.push_str(&format!("{byte:02x}"));
302 }
303 name.push_str(&format!("-{shard}"));
304 name
305}
306
307#[cfg(unix)]
312fn user_configures_sharing(ssh_args: &[String]) -> bool {
313 ssh_args.iter().any(|argument| {
314 if argument.starts_with("-S") {
315 return true;
316 }
317 let option = argument.strip_prefix("-o").unwrap_or(argument).trim_start();
318 let option = option.to_ascii_lowercase();
319 ["controlmaster", "controlpath"].iter().any(|name| {
320 option
321 .strip_prefix(name)
322 .is_some_and(|rest| rest.starts_with(['=', ' ', '\t']))
323 })
324 })
325}
326
327pub fn push_connection_reuse_args(args: &mut Vec<String>, ssh: &SshTarget) {
345 #[cfg(unix)]
346 if !user_configures_sharing(&ssh.ssh_args)
347 && let Some(dir) = control_socket_dir()
348 {
349 let socket = dir.join(control_socket_name(ssh, 0));
350 args.extend([
351 "-o".to_owned(),
352 "ControlMaster=no".to_owned(),
353 "-o".to_owned(),
354 format!("ControlPath={}", socket.display()),
355 ]);
356 }
357 #[cfg(not(unix))]
358 let _ = (args, ssh);
359}
360
361pub fn join_remote_command(args: &[String]) -> String {
362 args.iter()
363 .map(|arg| posix_quote(arg))
364 .collect::<Vec<_>>()
365 .join(" ")
366}
367
368pub fn ssh_directory_exists(
370 ssh: &SshTarget,
371 path: &Path,
372 executor: &impl CommandExecutor,
373) -> Result<bool> {
374 let command = ssh_validation_command(
375 ssh,
376 vec![
377 "test".into(),
378 "-d".into(),
379 path.to_string_lossy().into_owned(),
380 ],
381 "validate remote directory",
382 );
383 let output = executor.execute(&command)?;
384 match output.status {
385 0 => Ok(true),
386 1 => Ok(false),
387 status => {
388 let stderr = String::from_utf8_lossy(&output.stderr);
389 let error = anyhow::anyhow!(
390 "remote directory check failed with status {status}: {}",
391 stderr.trim()
392 );
393 Err(match host_key_refusal(&stderr, &ssh.ssh_args) {
394 Some(refusal) => error.context(refusal),
395 None => error,
396 })
397 }
398 }
399}
400
401fn host_key_refusal(stderr: &str, ssh_args: &[String]) -> Option<crate::refusal::Refusal> {
410 if !stderr.contains("Host key verification failed") {
411 return None;
412 }
413 let known_hosts = KnownHostsFile::from_ssh_args(ssh_args);
414 let file = known_hosts.phrase();
415 Some(crate::refusal::Refusal::precondition(
416 if stderr.contains("REMOTE HOST IDENTIFICATION HAS CHANGED") {
417 let keygen = match known_hosts.first_named() {
418 Some(path) => format!("`ssh-keygen -f {path} -R`"),
419 None => "`ssh-keygen -R`".to_owned(),
420 };
421 format!(
422 "ssh reported \"Host key verification failed\": the machine's host key is not the one saved in {file}. If you expected the change, remove the old entry with {keygen} and the host name, add the new key, and try again."
423 )
424 } else {
425 format!(
426 "ssh reported \"Host key verification failed\": the machine's host key is not in {file}, and its ssh options require a known key. Add the host key (for example with `ssh-keyscan`, after checking the fingerprint), or put `-o StrictHostKeyChecking=accept-new` in the machine's extra_args, and try again."
427 )
428 },
429 ))
430}
431
432#[derive(Debug, Clone, PartialEq, Eq)]
435enum KnownHostsFile {
436 Default,
438 Named(String),
441 Unknown,
444}
445
446impl KnownHostsFile {
447 fn from_ssh_args(args: &[String]) -> Self {
451 let mut config_file = false;
452 let mut args = args.iter();
453 while let Some(arg) = args.next() {
454 let option = match arg.as_str() {
455 "-o" => args.next().map(String::as_str),
456 other => other.strip_prefix("-o"),
457 };
458 if arg.starts_with("-F") {
459 config_file = true;
460 }
461 let Some(value) = option.and_then(|option| option_value(option, "UserKnownHostsFile"))
462 else {
463 continue;
464 };
465 let value = value.trim_matches('"').trim();
466 return if value.is_empty() || value.eq_ignore_ascii_case("none") {
467 Self::Unknown
468 } else {
469 Self::Named(value.to_owned())
470 };
471 }
472 if config_file {
473 Self::Unknown
474 } else {
475 Self::Default
476 }
477 }
478
479 fn phrase(&self) -> String {
481 match self {
482 Self::Default => "~/.ssh/known_hosts".to_owned(),
483 Self::Named(paths) => paths.split_whitespace().collect::<Vec<_>>().join(" or "),
484 Self::Unknown => "the known_hosts file ssh uses".to_owned(),
485 }
486 }
487
488 fn first_named(&self) -> Option<&str> {
490 match self {
491 Self::Named(paths) => paths.split_whitespace().next(),
492 Self::Default | Self::Unknown => None,
493 }
494 }
495}
496
497fn option_value<'a>(option: &'a str, keyword: &str) -> Option<&'a str> {
500 let option = option.trim_start();
501 let end = option
502 .find(|character: char| character == '=' || character.is_whitespace())
503 .unwrap_or(option.len());
504 let (name, rest) = option.split_at(end);
505 if !name.eq_ignore_ascii_case(keyword) {
506 return None;
507 }
508 let rest = rest.trim_start();
509 Some(rest.strip_prefix('=').unwrap_or(rest).trim())
510}
511
512pub fn validate_bare_project_directory(
514 ssh: &SshTarget,
515 path: &Path,
516 executor: &impl CommandExecutor,
517) -> Result<()> {
518 validate_bare_project_path(path)?;
519 if !ssh_directory_exists(ssh, path, executor)? {
520 return Err(anyhow::Error::new(crate::refusal::Refusal::unusable(
521 format!(
522 "remote project directory {} does not exist or is not a directory on {}",
523 path.display(),
524 ssh.destination
525 ),
526 )));
527 }
528 let output = executor.execute(&ssh_validation_command(
529 ssh,
530 vec![
531 "git".into(),
532 "-C".into(),
533 path.to_string_lossy().into_owned(),
534 "rev-parse".into(),
535 "--verify".into(),
536 "HEAD".into(),
537 ],
538 "validate bare SSH Git project",
539 ))?;
540 if output.status != 0 {
541 let detail = String::from_utf8_lossy(&output.stderr);
542 let detail = detail.trim();
543 if detail.is_empty() {
544 bail!(
545 "remote project directory {} has no valid Git HEAD",
546 path.display()
547 );
548 }
549 bail!(
550 "remote project directory {} has no valid Git HEAD: {detail}",
551 path.display()
552 );
553 }
554 Ok(())
555}
556
557pub fn validate_bare_project_path(path: &Path) -> Result<()> {
558 if !path.is_absolute()
559 || path
560 .components()
561 .any(|part| part == std::path::Component::ParentDir)
562 {
563 bail!("bare project directory must be an absolute safe path");
564 }
565 Ok(())
566}
567
568pub fn ssh_validation_command(
569 ssh: &SshTarget,
570 remote_args: Vec<String>,
571 purpose: &'static str,
572) -> CommandSpec {
573 let mut args = ssh.ssh_args.clone();
574 args.extend([
575 "-o".into(),
576 "BatchMode=yes".into(),
577 "-o".into(),
578 "ConnectTimeout=3".into(),
579 "-o".into(),
580 "ServerAliveInterval=2".into(),
581 "-o".into(),
582 "ServerAliveCountMax=1".into(),
583 ]);
584 args.extend([ssh.destination.clone(), join_remote_command(&remote_args)]);
585 CommandSpec::new("ssh", args)
586 .ssh_probe_session(ssh)
587 .purpose(purpose)
588}
589
590pub fn posix_quote(value: &str) -> String {
594 format!("'{}'", value.replace('\'', "'\\''"))
595}
596
597pub fn verify_locator(locator: &TargetLocator, session_id: &str) -> Result<()> {
598 validate_session_id(session_id)?;
599 match locator {
600 TargetLocator::LocalBare { worker_root } => {
601 let path = Path::new(worker_root);
602 if !path.is_absolute()
603 || path
604 .components()
605 .any(|part| part == std::path::Component::ParentDir)
606 || !path.ends_with(session_id)
607 {
608 bail!("refusing cleanup: invalid local bare worker root");
609 }
610 }
611 TargetLocator::LocalPodman {
612 container_id,
613 borrowed_from,
614 ..
615 }
616 | TargetLocator::LocalDocker {
617 container_id,
618 borrowed_from,
619 }
620 | TargetLocator::AppleContainer {
621 container_id,
622 borrowed_from,
623 }
624 | TargetLocator::SshPodman {
625 container_id,
626 borrowed_from,
627 ..
628 }
629 | TargetLocator::SshDocker {
630 container_id,
631 borrowed_from,
632 ..
633 } => match borrowed_from {
634 Some(owner) => {
635 validate_session_id(owner)?;
636 if owner == session_id {
637 bail!(
638 "refusing cleanup: a borrowed container cannot be owned by the borrowing session"
639 );
640 }
641 if !resource_name_belongs_to(container_id, owner)?
642 && !is_runtime_container_id(container_id)
643 {
644 bail!(
645 "refusing cleanup: borrowed container locator is neither the owning session's generated name nor an immutable runtime ID"
646 );
647 }
648 }
649 None => {
650 if !resource_name_belongs_to(container_id, session_id)?
651 && !is_runtime_container_id(container_id)
652 {
653 bail!(
654 "refusing cleanup: container locator is neither the generated name nor an immutable runtime ID"
655 );
656 }
657 }
658 },
659 TargetLocator::AwsEc2 {
660 instance_id,
661 workspace,
662 ..
663 } => {
664 if !valid_ec2_instance_id(instance_id) {
665 bail!("refusing cleanup: invalid EC2 instance ID");
666 }
667 verify_session_workspace(workspace, session_id)?;
668 }
669 TargetLocator::SshBare {
670 workspace,
671 worker_id,
672 ..
673 } => match worker_id {
674 Some(worker_id) => {
675 validate_session_id(worker_id)?;
676 if worker_id != session_id {
677 bail!("refusing cleanup: SSH worker identity does not match session ID");
678 }
679 validate_workspace_prefix(workspace)?;
680 }
681 None => verify_session_workspace(workspace, session_id)?,
682 },
683 }
684 Ok(())
685}
686
687pub fn is_borrowed(locator: &TargetLocator) -> bool {
691 match locator {
692 TargetLocator::LocalPodman { borrowed_from, .. }
693 | TargetLocator::LocalDocker { borrowed_from, .. }
694 | TargetLocator::AppleContainer { borrowed_from, .. }
695 | TargetLocator::SshPodman { borrowed_from, .. }
696 | TargetLocator::SshDocker { borrowed_from, .. } => borrowed_from.is_some(),
697 TargetLocator::SshBare { worker_id, .. } => worker_id.is_some(),
698 TargetLocator::LocalBare { .. } | TargetLocator::AwsEc2 { .. } => false,
699 }
700}
701
702pub fn verify_session_workspace(workspace: &str, session_id: &str) -> Result<()> {
703 validate_workspace_prefix(workspace)?;
704 let final_component = workspace.trim_end_matches('/').rsplit('/').next();
705 if final_component != Some(session_id) {
706 bail!("refusing cleanup: workspace does not end in the exact session ID");
707 }
708 Ok(())
709}
710
711pub fn validate_session_id(value: &str) -> Result<()> {
712 if value.len() < 8
713 || value.len() > 128
714 || !value
715 .chars()
716 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_'))
717 {
718 bail!("session ID must be 8-128 ASCII letters, digits, '-' or '_'");
719 }
720 Ok(())
721}
722
723pub fn validate_relative_path(value: &str) -> Result<()> {
724 let path = std::path::Path::new(value);
725 if value.is_empty()
726 || path.is_absolute()
727 || path
728 .components()
729 .any(|part| !matches!(part, std::path::Component::Normal(_)))
730 {
731 bail!("unsafe relative bundle path {value:?}");
732 }
733 Ok(())
734}
735
736pub fn validate_workspace_prefix(value: &str) -> Result<()> {
737 if value.is_empty()
738 || value == "/"
739 || value == "~"
740 || value == "~/"
741 || value.contains('\0')
742 || value.split('/').any(|part| part == "..")
743 {
744 bail!("unsafe workspace path");
745 }
746 Ok(())
747}
748
749pub fn validate_container_template(template: &ContainerTemplate) -> Result<()> {
750 if template.image.trim().is_empty() || template.image.starts_with('-') {
751 bail!("invalid container image");
752 }
753 if template
754 .extra_run_args
755 .iter()
756 .any(|arg| arg == "--name" || arg.starts_with("--name="))
757 {
758 bail!("container template may not override the generated name");
759 }
760 if template.extra_run_args.iter().any(|arg| {
761 arg == "--label"
762 || [SESSION_LABEL, MANAGED_LABEL, INSTANCE_LABEL]
763 .iter()
764 .any(|label| arg.starts_with(&format!("--label={label}=")))
765 }) {
766 bail!("container template may not override Mjolnir ownership labels");
767 }
768 Ok(())
769}
770
771pub fn validate_ssh(ssh: &SshTarget) -> Result<()> {
772 if ssh.destination.trim().is_empty()
773 || ssh.destination.starts_with('-')
774 || ssh.destination.chars().any(char::is_whitespace)
775 {
776 bail!("invalid SSH destination");
777 }
778 Ok(())
779}
780
781pub fn validate_aws(aws: &AwsTemplate) -> Result<()> {
782 validate_ssh(&aws.ssh)?;
783 for (name, value) in [
784 ("AWS profile", &aws.profile),
785 ("AWS region", &aws.region),
786 ("launch template", &aws.launch_template),
787 ] {
788 if value.is_empty()
789 || value.starts_with('-')
790 || !value
791 .chars()
792 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.' | '/'))
793 {
794 bail!("invalid {name}");
795 }
796 }
797 Ok(())
798}
799
800pub fn validate_executable(value: &str) -> Result<()> {
801 if value.is_empty() || value.starts_with('-') || value.chars().any(char::is_whitespace) {
802 bail!("invalid executable name");
803 }
804 Ok(())
805}
806
807pub fn valid_ec2_instance_id(value: &str) -> bool {
808 value
809 .strip_prefix("i-")
810 .is_some_and(|rest| rest.len() >= 8 && rest.chars().all(|c| c.is_ascii_hexdigit()))
811}
812
813pub fn is_runtime_container_id(value: &str) -> bool {
814 value.len() >= 12 && value.len() <= 128 && value.chars().all(|c| c.is_ascii_hexdigit())
815}
816
817pub const SSH_TRANSPORT_EXIT_STATUS: i32 = 255;
820
821const TRANSPORT_REJECTION_MARKERS: [&str; 4] = [
825 "Connection closed by",
826 "Connection reset by",
827 "kex_exchange_identification",
828 "Connection timed out during banner exchange",
829];
830
831pub fn is_transport_rejection(status: i32, stderr: &str) -> bool {
837 ssh_refusal(status, stderr).is_some()
838}
839
840const SESSION_REFUSAL_MARKER: &str = "Session open refused by peer";
844
845#[derive(Debug, Clone, Copy, PartialEq, Eq)]
848pub enum SshRefusal {
849 BeforeAuthentication,
853 SessionLimit,
856}
857
858impl SshRefusal {
859 pub fn retry_message(self) -> &'static str {
861 match self {
862 Self::BeforeAuthentication => {
863 "the SSH server closed the connection before authentication; retrying"
864 }
865 Self::SessionLimit => {
866 "the SSH server refused another session on a shared connection (MaxSessions); retrying"
867 }
868 }
869 }
870
871 pub fn log_retry(
880 self,
881 destination: &str,
882 purpose: &str,
883 attempt: usize,
884 delay: Duration,
885 stderr: &str,
886 ) {
887 let delay_ms = delay.as_millis() as u64;
888 match self {
889 Self::SessionLimit => tracing::debug!(
890 destination,
891 purpose,
892 attempt,
893 attempts = SSH_RETRY_ATTEMPTS,
894 delay_ms,
895 stderr,
896 "{}",
897 self.retry_message()
898 ),
899 Self::BeforeAuthentication => tracing::warn!(
900 destination,
901 purpose,
902 attempt,
903 attempts = SSH_RETRY_ATTEMPTS,
904 delay_ms,
905 stderr,
906 "{}",
907 self.retry_message()
908 ),
909 }
910 }
911
912 pub fn log_exhausted(self, destination: &str, purpose: &str, stderr: &str) {
914 tracing::warn!(
915 destination,
916 purpose,
917 attempts = SSH_RETRY_ATTEMPTS,
918 stderr,
919 "{}",
920 match self {
921 Self::BeforeAuthentication =>
922 "the SSH server closed the connection before authentication on every attempt",
923 Self::SessionLimit =>
924 "the SSH server refused another session on a shared connection (MaxSessions) on every attempt",
925 }
926 );
927 }
928}
929
930pub fn ssh_refusal(status: i32, stderr: &str) -> Option<SshRefusal> {
934 if status != SSH_TRANSPORT_EXIT_STATUS {
935 return None;
936 }
937 if stderr.contains(SESSION_REFUSAL_MARKER) {
938 return Some(SshRefusal::SessionLimit);
939 }
940 TRANSPORT_REJECTION_MARKERS
941 .iter()
942 .any(|marker| stderr.contains(marker))
943 .then_some(SshRefusal::BeforeAuthentication)
944}
945
946const DEFAULT_MAX_CONCURRENT_SSH: usize = 6;
955
956pub const MAX_CONCURRENT_SSH_ENV: &str = "MJ_SSH_MAX_CONCURRENT";
958
959fn max_concurrent_ssh() -> usize {
960 static LIMIT: OnceLock<usize> = OnceLock::new();
961 *LIMIT.get_or_init(|| positive_env_limit(MAX_CONCURRENT_SSH_ENV, DEFAULT_MAX_CONCURRENT_SSH))
962}
963
964fn positive_env_limit(name: &str, default: usize) -> usize {
967 let Some(raw) = std::env::var_os(name) else {
968 return default;
969 };
970 match raw
971 .to_str()
972 .and_then(|value| value.trim().parse::<usize>().ok())
973 {
974 Some(limit) if limit > 0 => limit,
975 _ => {
976 tracing::warn!(
977 variable = name,
978 value = %raw.to_string_lossy(),
979 default,
980 "ignoring invalid SSH limit"
981 );
982 default
983 }
984 }
985}
986
987struct DestinationGate {
993 limit: usize,
994 in_flight: Mutex<usize>,
995 released: Condvar,
996}
997
998impl DestinationGate {
999 fn new(limit: usize) -> Arc<Self> {
1000 Arc::new(Self {
1001 limit,
1002 in_flight: Mutex::new(0),
1003 released: Condvar::new(),
1004 })
1005 }
1006
1007 fn acquire(self: &Arc<Self>) -> SshPermit {
1008 self.acquire_unless(&|| false)
1009 .expect("unconditional SSH admission cannot be cancelled")
1010 }
1011
1012 fn acquire_unless(self: &Arc<Self>, cancelled: &dyn Fn() -> bool) -> Result<SshPermit> {
1013 let mut in_flight = self
1014 .in_flight
1015 .lock()
1016 .unwrap_or_else(std::sync::PoisonError::into_inner);
1017 loop {
1018 ensure!(
1019 !cancelled(),
1020 "operation cancelled while waiting for SSH admission"
1021 );
1022 if *in_flight < self.limit {
1023 break;
1024 }
1025 (in_flight, _) = self
1026 .released
1027 .wait_timeout(in_flight, Duration::from_millis(25))
1028 .unwrap_or_else(std::sync::PoisonError::into_inner);
1029 }
1030 *in_flight += 1;
1031 drop(in_flight);
1032 Ok(SshPermit {
1033 gate: Arc::clone(self),
1034 })
1035 }
1036}
1037
1038pub struct SshPermit {
1040 gate: Arc<DestinationGate>,
1041}
1042
1043impl std::fmt::Debug for SshPermit {
1044 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1045 formatter.write_str("SshPermit")
1046 }
1047}
1048
1049impl Drop for SshPermit {
1050 fn drop(&mut self) {
1051 let mut in_flight = self
1052 .gate
1053 .in_flight
1054 .lock()
1055 .unwrap_or_else(std::sync::PoisonError::into_inner);
1056 *in_flight = in_flight.saturating_sub(1);
1057 drop(in_flight);
1058 self.gate.released.notify_one();
1059 }
1060}
1061
1062pub struct SshAdmission;
1064
1065impl SshAdmission {
1066 pub fn acquire(destination: &str) -> SshPermit {
1069 Self::gate(destination).acquire()
1070 }
1071
1072 pub fn acquire_unless(destination: &str, cancelled: &dyn Fn() -> bool) -> Result<SshPermit> {
1074 Self::gate(destination).acquire_unless(cancelled)
1075 }
1076
1077 fn gate(destination: &str) -> Arc<DestinationGate> {
1078 static GATES: OnceLock<Mutex<BTreeMap<String, Arc<DestinationGate>>>> = OnceLock::new();
1079 let mut gates = GATES
1080 .get_or_init(|| Mutex::new(BTreeMap::new()))
1081 .lock()
1082 .unwrap_or_else(std::sync::PoisonError::into_inner);
1083 Arc::clone(
1084 gates
1085 .entry(destination.to_owned())
1086 .or_insert_with(|| DestinationGate::new(max_concurrent_ssh())),
1087 )
1088 }
1089}
1090
1091#[cfg(unix)]
1098const DEFAULT_SESSIONS_PER_CONNECTION: usize = 8;
1099
1100pub const SESSIONS_PER_CONNECTION_ENV: &str = "MJ_SSH_SESSIONS_PER_CONNECTION";
1102
1103#[cfg(unix)]
1106const MASTER_CHECK_INTERVAL: Duration = Duration::from_secs(5);
1107
1108pub const SSH_MASTER_OPEN_TIMEOUT: Duration = Duration::from_secs(60);
1111
1112#[cfg(unix)]
1113fn sessions_per_connection() -> usize {
1114 static LIMIT: OnceLock<usize> = OnceLock::new();
1115 *LIMIT.get_or_init(|| {
1116 positive_env_limit(SESSIONS_PER_CONNECTION_ENV, DEFAULT_SESSIONS_PER_CONNECTION)
1117 })
1118}
1119
1120#[cfg(unix)]
1122struct Shard {
1123 leased: usize,
1124 verified_at: Option<Instant>,
1126 opening: Arc<Mutex<()>>,
1129}
1130
1131#[cfg(unix)]
1133struct SessionLedger {
1134 per_connection: usize,
1135 connections: Mutex<BTreeMap<String, Vec<Shard>>>,
1136}
1137
1138#[cfg(unix)]
1139impl SessionLedger {
1140 fn new(per_connection: usize) -> Arc<Self> {
1141 Arc::new(Self {
1142 per_connection: per_connection.max(1),
1143 connections: Mutex::new(BTreeMap::new()),
1144 })
1145 }
1146
1147 fn global() -> Arc<Self> {
1148 static LEDGER: OnceLock<Arc<SessionLedger>> = OnceLock::new();
1149 Arc::clone(LEDGER.get_or_init(|| Self::new(sessions_per_connection())))
1150 }
1151
1152 fn connections(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Vec<Shard>>> {
1153 self.connections
1154 .lock()
1155 .unwrap_or_else(std::sync::PoisonError::into_inner)
1156 }
1157
1158 fn reserve(&self, key: &str) -> (usize, Arc<Mutex<()>>) {
1161 let mut connections = self.connections();
1162 let shards = connections.entry(key.to_owned()).or_default();
1163 let index = match shards
1164 .iter()
1165 .position(|shard| shard.leased < self.per_connection)
1166 {
1167 Some(index) => index,
1168 None => {
1169 shards.push(Shard {
1170 leased: 0,
1171 verified_at: None,
1172 opening: Arc::new(Mutex::new(())),
1173 });
1174 shards.len() - 1
1175 }
1176 };
1177 shards[index].leased += 1;
1178 (index, Arc::clone(&shards[index].opening))
1179 }
1180
1181 fn lease(
1184 self: &Arc<Self>,
1185 ssh: &SshTarget,
1186 dir: &Path,
1187 executor: &dyn CommandExecutor,
1188 ) -> Result<SshSessionLease> {
1189 let key = connection_key(ssh);
1190 let (shard, opening) = self.reserve(&key);
1191 let slot = LeasedSlot {
1193 ledger: Arc::clone(self),
1194 key,
1195 shard,
1196 socket: dir.join(control_socket_name(ssh, shard)),
1197 };
1198 if slot.needs_check() {
1199 let _opening = loop {
1200 ensure!(
1201 !executor.cancellation_requested(),
1202 "operation cancelled while waiting for SSH master"
1203 );
1204 match opening.try_lock() {
1205 Ok(guard) => break guard,
1206 Err(std::sync::TryLockError::Poisoned(error)) => break error.into_inner(),
1207 Err(std::sync::TryLockError::WouldBlock) => {
1208 std::thread::sleep(Duration::from_millis(25));
1209 }
1210 }
1211 };
1212 if slot.needs_check() {
1215 ensure_master(ssh, &slot.socket, executor)?;
1216 slot.set_verified(Some(Instant::now()));
1217 }
1218 }
1219 Ok(SshSessionLease {
1220 slot: Some(slot),
1221 probe: false,
1222 })
1223 }
1224
1225 fn lease_probe(self: &Arc<Self>, ssh: &SshTarget, dir: &Path) -> SshSessionLease {
1234 let key = connection_key(ssh);
1235 let (shard, _) = self.reserve(&key);
1236 SshSessionLease {
1237 slot: Some(LeasedSlot {
1238 ledger: Arc::clone(self),
1239 key,
1240 shard,
1241 socket: dir.join(control_socket_name(ssh, shard)),
1242 }),
1243 probe: true,
1244 }
1245 }
1246}
1247
1248#[cfg(unix)]
1254fn ensure_master(ssh: &SshTarget, socket: &Path, executor: &dyn CommandExecutor) -> Result<()> {
1255 let _opening = lock_master_opening_unless(socket, &|| executor.cancellation_requested())?;
1263 if master_running(ssh, socket, executor)? {
1264 return Ok(());
1265 }
1266 match fs::remove_file(socket) {
1270 Ok(()) => tracing::debug!(
1271 socket = %socket.display(),
1272 "removed a stale SSH control socket"
1273 ),
1274 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
1275 Err(error) => {
1276 return Err(error)
1277 .with_context(|| format!("remove stale SSH control socket {}", socket.display()));
1278 }
1279 }
1280 let opened = executor.execute(&master_open_command(ssh, socket))?;
1281 if master_running(ssh, socket, executor)? {
1282 tracing::info!(
1283 destination = ssh.destination.as_str(),
1284 socket = %socket.display(),
1285 "opened a shared SSH connection"
1286 );
1287 return Ok(());
1288 }
1289 let stderr = String::from_utf8_lossy(&opened.stderr);
1290 let detail = match stderr.trim() {
1291 "" => format!("ssh exited with status {}", opened.status),
1292 stderr => stderr.to_owned(),
1293 };
1294 bail!(
1295 "could not open a shared SSH connection to {}: {detail}",
1296 ssh.destination
1297 )
1298}
1299
1300#[cfg(all(unix, test))]
1305fn lock_master_opening(socket: &Path) -> Result<fs::File> {
1306 lock_master_opening_unless(socket, &|| false)
1307}
1308
1309#[cfg(unix)]
1310fn lock_master_opening_unless(socket: &Path, cancelled: &dyn Fn() -> bool) -> Result<fs::File> {
1311 let mut path = socket.as_os_str().to_owned();
1312 path.push(".lock");
1313 let path = PathBuf::from(path);
1314 loop {
1315 ensure!(
1316 !cancelled(),
1317 "operation cancelled while waiting for SSH master lock"
1318 );
1319 let file = fs::OpenOptions::new()
1320 .create(true)
1321 .truncate(false)
1322 .write(true)
1323 .open(&path)
1324 .with_context(|| format!("open SSH master lock {}", path.display()))?;
1325 match file.try_lock() {
1326 Ok(()) => {}
1327 Err(std::fs::TryLockError::WouldBlock) => {
1328 std::thread::sleep(Duration::from_millis(25));
1329 continue;
1330 }
1331 Err(std::fs::TryLockError::Error(error)) => {
1332 return Err(error)
1333 .with_context(|| format!("lock SSH master lock {}", path.display()));
1334 }
1335 }
1336 if is_file_at(&file, &path) {
1340 return Ok(file);
1341 }
1342 }
1343}
1344
1345#[cfg(unix)]
1354fn remove_stale_master_locks_in(dir: &Path) {
1355 let entries = match fs::read_dir(dir) {
1356 Ok(entries) => entries,
1357 Err(error) => {
1358 tracing::debug!(directory = %dir.display(), %error, "cannot list SSH master locks");
1359 return;
1360 }
1361 };
1362 for entry in entries.flatten() {
1363 let lock = entry.path();
1364 let Some(socket) = lock
1365 .file_name()
1366 .and_then(std::ffi::OsStr::to_str)
1367 .and_then(|name| name.strip_suffix(".lock"))
1368 .map(|name| dir.join(name))
1369 else {
1370 continue;
1371 };
1372 if fs::symlink_metadata(&socket).is_ok() {
1373 continue;
1374 }
1375 let Ok(file) = fs::OpenOptions::new().write(true).open(&lock) else {
1376 continue;
1377 };
1378 if file.try_lock().is_err() {
1380 continue;
1381 }
1382 if fs::symlink_metadata(&socket).is_ok() || !is_file_at(&file, &lock) {
1383 continue;
1384 }
1385 match fs::remove_file(&lock) {
1386 Ok(()) => tracing::debug!(lock = %lock.display(), "removed a stale SSH master lock"),
1387 Err(error) => {
1388 tracing::debug!(lock = %lock.display(), %error, "cannot remove a stale SSH master lock")
1389 }
1390 }
1391 }
1392}
1393
1394#[cfg(unix)]
1396fn is_file_at(file: &fs::File, path: &Path) -> bool {
1397 use std::os::unix::fs::MetadataExt;
1398 match (file.metadata(), fs::metadata(path)) {
1399 (Ok(open), Ok(named)) => open.dev() == named.dev() && open.ino() == named.ino(),
1400 _ => false,
1401 }
1402}
1403
1404#[cfg(unix)]
1405fn master_running(ssh: &SshTarget, socket: &Path, executor: &dyn CommandExecutor) -> Result<bool> {
1406 Ok(executor.execute(&master_check_command(ssh, socket))?.status == 0)
1407}
1408
1409#[cfg(unix)]
1412fn master_check_command(ssh: &SshTarget, socket: &Path) -> CommandSpec {
1413 let mut args = ssh.ssh_args.clone();
1414 args.extend([
1415 "-o".to_owned(),
1416 format!("ControlPath={}", socket.display()),
1417 "-O".to_owned(),
1418 "check".to_owned(),
1419 ssh.destination.clone(),
1420 ]);
1421 CommandSpec::new("ssh", args).purpose("check a shared SSH connection")
1422}
1423
1424#[cfg(unix)]
1432fn master_open_command(ssh: &SshTarget, socket: &Path) -> CommandSpec {
1433 let mut args = ssh.ssh_args.clone();
1434 args.extend([
1435 "-o".to_owned(),
1436 "BatchMode=yes".to_owned(),
1437 "-o".to_owned(),
1438 "ConnectTimeout=10".to_owned(),
1439 "-f".to_owned(),
1440 "-N".to_owned(),
1441 "-o".to_owned(),
1442 "ControlMaster=yes".to_owned(),
1443 "-o".to_owned(),
1444 format!("ControlPath={}", socket.display()),
1445 "-o".to_owned(),
1446 format!("ControlPersist={CONTROL_PERSIST}"),
1447 ssh.destination.clone(),
1448 ]);
1449 CommandSpec::new("ssh", args)
1450 .ssh_destination(ssh.destination.clone())
1451 .purpose("open a shared SSH connection")
1452}
1453
1454#[cfg(unix)]
1456struct LeasedSlot {
1457 ledger: Arc<SessionLedger>,
1458 key: String,
1459 shard: usize,
1460 socket: PathBuf,
1461}
1462
1463#[cfg(unix)]
1464impl LeasedSlot {
1465 fn needs_check(&self) -> bool {
1466 let connections = self.ledger.connections();
1467 connections
1468 .get(&self.key)
1469 .and_then(|shards| shards.get(self.shard))
1470 .is_none_or(|shard| {
1471 shard
1472 .verified_at
1473 .is_none_or(|verified| verified.elapsed() >= MASTER_CHECK_INTERVAL)
1474 })
1475 }
1476
1477 fn set_verified(&self, verified_at: Option<Instant>) {
1478 let mut connections = self.ledger.connections();
1479 if let Some(shard) = connections
1480 .get_mut(&self.key)
1481 .and_then(|shards| shards.get_mut(self.shard))
1482 {
1483 shard.verified_at = verified_at;
1484 }
1485 }
1486}
1487
1488#[cfg(unix)]
1489impl Drop for LeasedSlot {
1490 fn drop(&mut self) {
1491 let mut connections = self.ledger.connections();
1492 if let Some(shard) = connections
1493 .get_mut(&self.key)
1494 .and_then(|shards| shards.get_mut(self.shard))
1495 {
1496 shard.leased = shard.leased.saturating_sub(1);
1497 }
1498 }
1499}
1500
1501pub struct SshSessionLease {
1509 #[cfg(unix)]
1510 slot: Option<LeasedSlot>,
1511 #[cfg(unix)]
1514 probe: bool,
1515}
1516
1517impl std::fmt::Debug for SshSessionLease {
1518 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1519 formatter
1520 .debug_struct("SshSessionLease")
1521 .field("control_path", &self.control_path())
1522 .finish()
1523 }
1524}
1525
1526impl SshSessionLease {
1527 fn unshared() -> Self {
1528 Self {
1529 #[cfg(unix)]
1530 slot: None,
1531 #[cfg(unix)]
1532 probe: false,
1533 }
1534 }
1535
1536 pub fn control_path(&self) -> Option<&Path> {
1539 #[cfg(unix)]
1540 return self.slot.as_ref().map(|slot| slot.socket.as_path());
1541 #[cfg(not(unix))]
1542 None
1543 }
1544
1545 pub fn invalidate(&self) {
1549 #[cfg(unix)]
1550 if let Some(slot) = &self.slot {
1551 slot.set_verified(None);
1552 }
1553 }
1554}
1555
1556pub struct SshSessions;
1563
1564impl SshSessions {
1565 pub fn lease(ssh: &SshTarget, executor: &dyn CommandExecutor) -> Result<SshSessionLease> {
1570 #[cfg(unix)]
1571 {
1572 if user_configures_sharing(&ssh.ssh_args) {
1573 return Ok(SshSessionLease::unshared());
1574 }
1575 let Some(dir) = control_socket_dir() else {
1576 return Ok(SshSessionLease::unshared());
1577 };
1578 SessionLedger::global().lease(ssh, &dir, executor)
1579 }
1580 #[cfg(not(unix))]
1581 {
1582 let _ = (ssh, executor);
1583 Ok(SshSessionLease::unshared())
1584 }
1585 }
1586
1587 pub fn remove_stale_master_locks() {
1590 #[cfg(unix)]
1591 if let Some(dir) = control_socket_dir() {
1592 remove_stale_master_locks_in(&dir);
1593 }
1594 }
1595
1596 pub fn lease_probe(ssh: &SshTarget) -> SshSessionLease {
1604 #[cfg(unix)]
1605 {
1606 if user_configures_sharing(&ssh.ssh_args) {
1607 return SshSessionLease::unshared();
1608 }
1609 let Some(dir) = control_socket_dir() else {
1610 return SshSessionLease::unshared();
1611 };
1612 SessionLedger::global().lease_probe(ssh, &dir)
1613 }
1614 #[cfg(not(unix))]
1615 {
1616 let _ = ssh;
1617 SshSessionLease::unshared()
1618 }
1619 }
1620}
1621
1622pub fn push_session_args(args: &mut Vec<String>, lease: &SshSessionLease) {
1631 if let Some(socket) = lease.control_path() {
1632 args.extend([
1633 "-o".to_owned(),
1634 "ControlMaster=no".to_owned(),
1635 "-o".to_owned(),
1636 format!("ControlPath={}", socket.display()),
1637 ]);
1638 #[cfg(unix)]
1639 let probe = lease.probe;
1640 #[cfg(not(unix))]
1641 let probe = false;
1642 if !probe {
1643 args.extend(["-o".to_owned(), "ProxyCommand=false".to_owned()]);
1644 }
1645 }
1646}
1647
1648pub fn session_command_args(
1661 program: &str,
1662 args: &[String],
1663 ssh: &SshTarget,
1664 lease: &SshSessionLease,
1665) -> Vec<String> {
1666 let mut session = Vec::with_capacity(args.len() + 6);
1667 push_session_args(&mut session, lease);
1668 if session.is_empty() {
1669 return args.to_vec();
1670 }
1671 if program == "ssh" && args.starts_with(&ssh.ssh_args) {
1672 let (user, rest) = args.split_at(ssh.ssh_args.len());
1673 let mut user = user.iter();
1674 while let Some(argument) = user.next() {
1675 match argument.strip_prefix("-J") {
1676 Some("") => match user.next() {
1677 Some(jump) => session.extend(["-o".to_owned(), format!("ProxyJump={jump}")]),
1678 None => session.push(argument.clone()),
1679 },
1680 Some(jump) => session.extend(["-o".to_owned(), format!("ProxyJump={jump}")]),
1681 None => session.push(argument.clone()),
1682 }
1683 }
1684 session.extend(rest.iter().cloned());
1685 } else {
1686 session.extend(args.iter().cloned());
1687 }
1688 session
1689}
1690
1691pub const SSH_RETRY_ATTEMPTS: usize = 3;
1693
1694const SSH_RETRY_BACKOFF_MS: [(u64, u64); SSH_RETRY_ATTEMPTS - 1] = [(500, 2_000), (2_000, 4_000)];
1698
1699static SSH_RETRY_BACKOFF_OVERRIDE_MS: AtomicU64 = AtomicU64::new(u64::MAX);
1702
1703#[doc(hidden)]
1706pub fn set_ssh_retry_backoff_for_test(delay: Option<Duration>) {
1707 SSH_RETRY_BACKOFF_OVERRIDE_MS.store(
1708 delay.map_or(u64::MAX, |delay| delay.as_millis() as u64),
1709 Ordering::Relaxed,
1710 );
1711}
1712
1713pub fn ssh_retry_delay(attempts_made: usize) -> Duration {
1719 let override_ms = SSH_RETRY_BACKOFF_OVERRIDE_MS.load(Ordering::Relaxed);
1720 if override_ms != u64::MAX {
1721 return Duration::from_millis(override_ms);
1722 }
1723 let (low, high) = SSH_RETRY_BACKOFF_MS
1724 .get(attempts_made.saturating_sub(1))
1725 .copied()
1726 .unwrap_or(*SSH_RETRY_BACKOFF_MS.last().expect("non-empty schedule"));
1727 let mut bytes = [0_u8; 8];
1728 let spread = if getrandom::fill(&mut bytes).is_ok() {
1730 u64::from_le_bytes(bytes) % (high - low + 1)
1731 } else {
1732 0
1733 };
1734 Duration::from_millis(low + spread)
1735}
1736
1737#[cfg(test)]
1738mod tests {
1739 use super::*;
1740
1741 #[cfg(unix)]
1746 #[test]
1747 fn stale_master_locks_are_removed_but_live_or_held_ones_stay() {
1748 let dir = tempfile::tempdir().unwrap();
1749 let stale = dir.path().join("aaaaaaaaaaaaaaaa-0.lock");
1750 fs::write(&stale, b"").unwrap();
1751 fs::write(dir.path().join("bbbbbbbbbbbbbbbb-0"), b"").unwrap();
1752 let live = dir.path().join("bbbbbbbbbbbbbbbb-0.lock");
1753 fs::write(&live, b"").unwrap();
1754 let held = dir.path().join("cccccccccccccccc-0.lock");
1755 let holder = lock_master_opening(&dir.path().join("cccccccccccccccc-0")).unwrap();
1756
1757 remove_stale_master_locks_in(dir.path());
1758 assert!(!stale.exists(), "a lock whose master is gone is removed");
1759 assert!(live.exists(), "a lock beside a live socket stays");
1760 assert!(held.exists(), "a lock another opener holds stays");
1761
1762 drop(holder);
1763 for _ in 0..200 {
1768 remove_stale_master_locks_in(dir.path());
1769 if !held.exists() {
1770 break;
1771 }
1772 std::thread::sleep(Duration::from_millis(10));
1773 }
1774 assert!(!held.exists(), "a lock nobody holds any more is removed");
1775 }
1776
1777 const BORROW_PARENT: &str = "0123456789abcdef0123456789abcdef";
1778 const BORROW_CHILD: &str = "fedcba9876543210fedcba9876543210";
1779
1780 fn borrowed_podman(owner: &str) -> TargetLocator {
1781 TargetLocator::LocalPodman {
1782 container_id: crate::targets::resource_name(owner).unwrap(),
1783 workspace_storage: PodmanWorkspaceLocator::default(),
1784 borrowed_from: Some(owner.to_owned()),
1785 }
1786 }
1787
1788 #[test]
1789 fn verify_locator_accepts_a_container_borrowed_from_its_owner() {
1790 verify_locator(&borrowed_podman(BORROW_PARENT), BORROW_CHILD)
1791 .expect("a child may borrow its parent's container");
1792 }
1793
1794 #[test]
1795 fn verify_locator_rejects_a_container_borrowed_from_the_checking_session() {
1796 let error = verify_locator(&borrowed_podman(BORROW_PARENT), BORROW_PARENT)
1797 .expect_err("a session cannot borrow from itself");
1798 assert!(
1799 format!("{error:#}").contains("cannot be owned by the borrowing session"),
1800 "unexpected error: {error:#}"
1801 );
1802 }
1803
1804 #[test]
1805 fn verify_locator_rejects_a_borrowed_container_naming_another_session() {
1806 let locator = TargetLocator::LocalPodman {
1807 container_id: crate::targets::resource_name(BORROW_CHILD).unwrap(),
1808 workspace_storage: PodmanWorkspaceLocator::default(),
1809 borrowed_from: Some(BORROW_PARENT.to_owned()),
1810 };
1811 let error = verify_locator(&locator, BORROW_CHILD)
1812 .expect_err("the container must belong to the recorded owner");
1813 assert!(
1814 format!("{error:#}").contains("borrowed container locator"),
1815 "unexpected error: {error:#}"
1816 );
1817 }
1818
1819 #[test]
1820 fn worker_root_of_a_borrowed_container_is_the_childs_own_directory() {
1821 assert_eq!(
1822 crate::targets::worker_root(&borrowed_podman(BORROW_PARENT), BORROW_CHILD).unwrap(),
1823 format!("/var/lib/hel/workers/{BORROW_CHILD}")
1824 );
1825 }
1826
1827 #[test]
1828 fn is_borrowed_distinguishes_borrowed_targets_from_owned_ones() {
1829 assert!(is_borrowed(&borrowed_podman(BORROW_PARENT)));
1830 assert!(is_borrowed(&TargetLocator::SshBare {
1831 ssh: SshTarget {
1832 destination: "host".to_owned(),
1833 ssh_args: Vec::new(),
1834 },
1835 workspace: format!(".local/share/hel/workspaces/{BORROW_PARENT}"),
1836 worker_id: Some(BORROW_CHILD.to_owned()),
1837 }));
1838 assert!(!is_borrowed(&TargetLocator::LocalPodman {
1839 container_id: crate::targets::resource_name(BORROW_CHILD).unwrap(),
1840 workspace_storage: PodmanWorkspaceLocator::default(),
1841 borrowed_from: None,
1842 }));
1843 }
1844
1845 #[test]
1846 fn an_owned_container_locator_serializes_without_a_borrowed_from_key() {
1847 let owned = TargetLocator::LocalDocker {
1848 container_id: crate::targets::resource_name(BORROW_CHILD).unwrap(),
1849 borrowed_from: None,
1850 };
1851 let serialized = serde_json::to_string(&owned).unwrap();
1852 assert!(
1853 !serialized.contains("borrowed_from"),
1854 "owned locators must stay byte-identical for older readers: {serialized}"
1855 );
1856 assert_eq!(
1857 serde_json::from_str::<TargetLocator>(&serialized).unwrap(),
1858 owned
1859 );
1860
1861 let borrowed = borrowed_podman(BORROW_PARENT);
1862 let serialized = serde_json::to_string(&borrowed).unwrap();
1863 assert!(serialized.contains("borrowed_from"));
1864 assert_eq!(
1865 serde_json::from_str::<TargetLocator>(&serialized).unwrap(),
1866 borrowed
1867 );
1868 }
1869 use std::sync::atomic::{AtomicUsize, Ordering};
1870
1871 #[cfg(unix)]
1873 #[derive(Default)]
1874 struct RecordingExecutor {
1875 seen: std::cell::RefCell<Vec<CommandSpec>>,
1876 }
1877
1878 #[cfg(unix)]
1879 impl CommandExecutor for RecordingExecutor {
1880 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1881 self.seen.borrow_mut().push(command.clone());
1882 Ok(CommandOutput {
1883 status: 0,
1884 stdout: Vec::new(),
1885 stderr: Vec::new(),
1886 })
1887 }
1888 }
1889
1890 #[cfg(unix)]
1891 fn sharing_socket_dir() -> tempfile::TempDir {
1892 tempfile::tempdir_in("/tmp").expect("short control socket directory")
1894 }
1895
1896 #[cfg(unix)]
1900 fn spawned_args(
1901 command: &CommandSpec,
1902 dir: Option<&Path>,
1903 masters: &FakeMasters,
1904 ) -> Vec<String> {
1905 set_ssh_connection_sharing_for_test(Some(match dir {
1906 Some(dir) => SshSharingForTest::Directory(dir.to_path_buf()),
1907 None => SshSharingForTest::Disabled,
1908 }));
1909 let session = command.open_ssh_session(masters);
1910 set_ssh_connection_sharing_for_test(None);
1911 session.expect("session").command().args.clone()
1912 }
1913
1914 #[test]
1918 #[cfg(unix)]
1919 fn session_options_lead_the_spawned_command() {
1920 let _guard = SHARING_TEST_LOCK
1921 .lock()
1922 .unwrap_or_else(std::sync::PoisonError::into_inner);
1923 let socket_dir = sharing_socket_dir();
1924 let ssh = SshTarget {
1925 destination: "session-options-host".to_owned(),
1926 ssh_args: vec![
1927 "-p".to_owned(),
1928 "2222".to_owned(),
1929 "-o".to_owned(),
1930 "ProxyCommand=nc %h %p".to_owned(),
1931 ],
1932 };
1933 let command = ssh_command(&ssh, ["true"]);
1934 assert_eq!(
1935 command.args,
1936 [
1937 "-p",
1938 "2222",
1939 "-o",
1940 "ProxyCommand=nc %h %p",
1941 "session-options-host",
1942 "'true'"
1943 ],
1944 "stored arguments never contain sharing options"
1945 );
1946 assert_eq!(command.ssh_session.as_ref(), Some(&ssh));
1947 assert_eq!(
1948 command.ssh_destination.as_deref(),
1949 Some("session-options-host")
1950 );
1951
1952 let masters = FakeMasters::default();
1953 let args = spawned_args(&command, Some(socket_dir.path()), &masters);
1954 let socket = socket_dir.path().join(control_socket_name(&ssh, 0));
1955 assert_eq!(
1956 args,
1957 [
1958 "-o".to_owned(),
1959 "ControlMaster=no".to_owned(),
1960 "-o".to_owned(),
1961 format!("ControlPath={}", socket.display()),
1962 "-o".to_owned(),
1963 "ProxyCommand=false".to_owned(),
1964 "-p".to_owned(),
1965 "2222".to_owned(),
1966 "-o".to_owned(),
1967 "ProxyCommand=nc %h %p".to_owned(),
1968 "session-options-host".to_owned(),
1969 "'true'".to_owned(),
1970 ]
1971 );
1972 assert_eq!(masters.openers(), 1);
1973 assert_eq!(
1974 std::os::unix::fs::MetadataExt::mode(
1975 &fs::metadata(socket_dir.path()).expect("socket directory")
1976 ) & 0o777,
1977 0o700
1978 );
1979 }
1980
1981 #[test]
1985 #[cfg(unix)]
1986 fn a_jump_host_flag_becomes_an_option_the_session_guard_overrides() {
1987 let _guard = SHARING_TEST_LOCK
1988 .lock()
1989 .unwrap_or_else(std::sync::PoisonError::into_inner);
1990 let socket_dir = sharing_socket_dir();
1991 let ssh = SshTarget {
1992 destination: "jump-rewrite-host".to_owned(),
1993 ssh_args: vec!["-J".to_owned(), "bastion".to_owned(), "-Jother".to_owned()],
1994 };
1995 let masters = FakeMasters::default();
1996 let args = spawned_args(
1997 &ssh_command(&ssh, ["-J"]),
1998 Some(socket_dir.path()),
1999 &masters,
2000 );
2001 assert_eq!(
2002 args[6..],
2003 [
2004 "-o",
2005 "ProxyJump=bastion",
2006 "-o",
2007 "ProxyJump=other",
2008 "jump-rewrite-host",
2009 "'-J'",
2010 ]
2011 );
2012 let upload = spawned_args(
2013 &scp_upload(&ssh, Path::new("/tmp/file"), "file", false),
2014 Some(socket_dir.path()),
2015 &masters,
2016 );
2017 assert_eq!(
2018 upload[6..],
2019 [
2020 "-J",
2021 "bastion",
2022 "-Jother",
2023 "/tmp/file",
2024 "jump-rewrite-host:file"
2025 ]
2026 );
2027 }
2028
2029 #[test]
2033 #[cfg(unix)]
2034 fn user_configured_sharing_suppresses_mjolnir_sharing() {
2035 let _guard = SHARING_TEST_LOCK
2036 .lock()
2037 .unwrap_or_else(std::sync::PoisonError::into_inner);
2038 let socket_dir = sharing_socket_dir();
2039 let spellings: [&[&str]; 5] = [
2040 &["-o", "ControlMaster=no"],
2041 &["-o", "controlpath /tmp/mine"],
2042 &["-oControlPath=/tmp/mine"],
2043 &["-S", "/tmp/mine"],
2044 &["-S/tmp/mine"],
2045 ];
2046 let masters = FakeMasters::default();
2047 for user in spellings {
2048 let ssh = SshTarget {
2049 destination: "user-sharing-host".to_owned(),
2050 ssh_args: user.iter().map(|arg| (*arg).to_owned()).collect(),
2051 };
2052 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2053 socket_dir.path().to_path_buf(),
2054 )));
2055 let validation = ssh_validation_command(&ssh, vec!["true".to_owned()], "test");
2056 let validation = spawned_args(&validation, Some(socket_dir.path()), &masters);
2057 let command = ssh_command(&ssh, ["true"]);
2058 let args = spawned_args(&command, Some(socket_dir.path()), &masters);
2059 assert_eq!(args, command.args, "user args {user:?}");
2060 let socket_dir_text = socket_dir.path().display().to_string();
2061 assert!(
2062 !validation.iter().any(|arg| arg.contains(&socket_dir_text)),
2063 "user args {user:?}: {validation:?}"
2064 );
2065 }
2066 assert_eq!(masters.commands(), 0);
2067 }
2068
2069 #[test]
2072 #[cfg(unix)]
2073 fn control_sockets_live_in_a_directory_per_instance() {
2074 let runtime = Some(std::ffi::OsString::from("/run/user/1000"));
2075 assert_eq!(
2076 default_control_dir(runtime.clone(), &control_dir_identity(Some("hel2"), None)),
2077 PathBuf::from("/run/user/1000/mjolnir/hel2")
2078 );
2079 assert_eq!(
2080 default_control_dir(runtime, &control_dir_identity(None, None)),
2081 PathBuf::from("/run/user/1000/mjolnir/default")
2082 );
2083 }
2084
2085 #[test]
2090 #[cfg(unix)]
2091 fn a_data_directory_override_gets_its_own_socket_directory() {
2092 let runtime = Some(std::ffi::OsString::from("/run/user/1000"));
2093 let lab = Path::new("/tmp/lab-a/data");
2094 let other_lab = Path::new("/tmp/lab-b/data");
2095 let lab_dir = default_control_dir(runtime.clone(), &control_dir_identity(None, Some(lab)));
2096 assert_ne!(lab_dir, PathBuf::from("/run/user/1000/mjolnir/default"));
2097 assert_eq!(
2098 lab_dir,
2099 PathBuf::from("/run/user/1000/mjolnir")
2100 .join(crate::config::instance_identity_for(None, lab))
2101 );
2102 assert_ne!(
2103 lab_dir,
2104 default_control_dir(
2105 runtime.clone(),
2106 &control_dir_identity(None, Some(other_lab))
2107 )
2108 );
2109 assert_eq!(
2112 default_control_dir(runtime, &control_dir_identity(Some("hel2"), Some(lab))),
2113 PathBuf::from("/run/user/1000/mjolnir/hel2")
2114 );
2115 }
2116
2117 #[test]
2120 #[cfg(unix)]
2121 fn socket_names_identify_the_connection_and_the_shard() {
2122 let plain = SshTarget {
2123 destination: "host".to_owned(),
2124 ssh_args: Vec::new(),
2125 };
2126 let other_port = SshTarget {
2127 destination: "host".to_owned(),
2128 ssh_args: vec!["-p".to_owned(), "2222".to_owned()],
2129 };
2130 let first = control_socket_name(&plain, 0);
2131 let second = control_socket_name(&plain, 1);
2132 assert_eq!(first.len(), CONNECTION_HASH_HEX + 2, "{first}");
2133 assert!(first.ends_with("-0") && second.ends_with("-1"));
2134 assert_eq!(first[..CONNECTION_HASH_HEX], second[..CONNECTION_HASH_HEX]);
2135 assert_ne!(
2136 first[..CONNECTION_HASH_HEX],
2137 control_socket_name(&other_port, 0)[..CONNECTION_HASH_HEX]
2138 );
2139 }
2140
2141 #[test]
2145 #[cfg(unix)]
2146 fn fail_fast_commands_reuse_a_master_without_becoming_one() {
2147 let _guard = SHARING_TEST_LOCK
2148 .lock()
2149 .unwrap_or_else(std::sync::PoisonError::into_inner);
2150 let socket_dir = sharing_socket_dir();
2151 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2152 socket_dir.path().to_path_buf(),
2153 )));
2154 let ssh = SshTarget {
2155 destination: "host".to_owned(),
2156 ssh_args: Vec::new(),
2157 };
2158 let masters = FakeMasters::default();
2159 let validation = spawned_args(
2160 &ssh_validation_command(&ssh, vec!["true".to_owned()], "test"),
2161 Some(socket_dir.path()),
2162 &masters,
2163 );
2164 assert_eq!(
2165 masters.commands(),
2166 0,
2167 "a probe never checks or opens a master"
2168 );
2169 assert!(!validation.contains(&"ProxyCommand=false".to_owned()));
2170 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2171 socket_dir.path().to_path_buf(),
2172 )));
2173 let executor = RecordingExecutor::default();
2174 crate::path_completion::ssh_completions(
2175 &ssh,
2176 "/srv/pr",
2177 crate::path_completion::CompletionKind::Directories,
2178 &executor,
2179 )
2180 .expect("completion runs");
2181 let completion = executor.seen.borrow()[0].args.clone();
2182 set_ssh_connection_sharing_for_test(None);
2183
2184 let control_path = format!(
2185 "ControlPath={}/{}",
2186 socket_dir.path().display(),
2187 control_socket_name(&ssh, 0)
2188 );
2189 for args in [&validation, &completion] {
2190 assert!(args.contains(&"ControlMaster=no".to_owned()), "{args:?}");
2191 assert!(args.contains(&control_path), "{args:?}");
2192 assert!(
2193 !args.iter().any(|arg| arg.starts_with("ControlPersist")),
2194 "a fail-fast command must not set how long a master lingers: {args:?}"
2195 );
2196 let master = args
2197 .iter()
2198 .position(|arg| arg == "ControlMaster=no")
2199 .expect("sharing options");
2200 assert!(
2201 args.contains(&"ServerAliveCountMax=1".to_owned()),
2202 "its own keepalive: {args:?}"
2203 );
2204 assert!(
2205 master
2206 < args
2207 .iter()
2208 .position(|arg| arg == "host")
2209 .expect("destination"),
2210 "{args:?}"
2211 );
2212 }
2213 }
2214
2215 #[test]
2220 #[cfg(unix)]
2221 fn connectivity_probe_joins_a_master_without_becoming_one() {
2222 let _guard = SHARING_TEST_LOCK
2223 .lock()
2224 .unwrap_or_else(std::sync::PoisonError::into_inner);
2225 let socket_dir = sharing_socket_dir();
2226 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2227 socket_dir.path().to_path_buf(),
2228 )));
2229 let ssh = SshTarget {
2230 destination: "host".to_owned(),
2231 ssh_args: Vec::new(),
2232 };
2233 let masters = FakeMasters::default();
2234 let args = spawned_args(
2235 &ssh_connectivity_probe(&ssh),
2236 Some(socket_dir.path()),
2237 &masters,
2238 );
2239
2240 assert_eq!(masters.commands(), 0, "a probe never opens a master");
2241 assert!(args.contains(&"ControlMaster=no".to_owned()), "{args:?}");
2242 assert!(
2243 args.contains(&format!(
2244 "ControlPath={}/{}",
2245 socket_dir.path().display(),
2246 control_socket_name(&ssh, 0)
2247 )),
2248 "the probe must still join an existing master: {args:?}"
2249 );
2250 assert!(
2251 !args.iter().any(|arg| arg.starts_with("ControlPersist")),
2252 "a doctor probe must not set how long a master lingers: {args:?}"
2253 );
2254 let master = args
2255 .iter()
2256 .position(|arg| arg == "ControlMaster=no")
2257 .expect("sharing options");
2258 let strict = args
2259 .iter()
2260 .position(|arg| arg == "StrictHostKeyChecking=yes")
2261 .expect("its own host key policy");
2262 assert!(master < strict, "{args:?}");
2263 }
2264
2265 #[test]
2266 #[cfg(unix)]
2267 fn connection_sharing_is_absent_when_turned_off() {
2268 let _guard = SHARING_TEST_LOCK
2269 .lock()
2270 .unwrap_or_else(std::sync::PoisonError::into_inner);
2271 let ssh = SshTarget {
2272 destination: "sharing-off-host".to_owned(),
2273 ssh_args: Vec::new(),
2274 };
2275 let masters = FakeMasters::default();
2276 let args = spawned_args(&ssh_command(&ssh, ["true"]), None, &masters);
2277 let validation = spawned_args(
2278 &ssh_validation_command(&ssh, vec!["true".to_owned()], "test"),
2279 None,
2280 &masters,
2281 );
2282 assert_eq!(args, ["sharing-off-host", "'true'"]);
2283 assert!(!validation.iter().any(|arg| arg.starts_with("Control")));
2284 assert_eq!(masters.commands(), 0);
2285 }
2286
2287 #[test]
2288 #[cfg(unix)]
2289 fn a_control_path_that_cannot_fit_a_socket_address_is_skipped() {
2290 let _guard = SHARING_TEST_LOCK
2291 .lock()
2292 .unwrap_or_else(std::sync::PoisonError::into_inner);
2293 let root = tempfile::tempdir().expect("temp dir");
2294 let long = root.path().join("a".repeat(MAX_CONTROL_PATH));
2295 let ssh = SshTarget {
2296 destination: "long-path-host".to_owned(),
2297 ssh_args: Vec::new(),
2298 };
2299 let masters = FakeMasters::default();
2300 let args = spawned_args(&ssh_command(&ssh, ["true"]), Some(&long), &masters);
2301 assert_eq!(args, ["long-path-host", "'true'"]);
2302 assert_eq!(masters.commands(), 0);
2303 assert!(!long.exists(), "an unusable directory must not be created");
2304 }
2305
2306 #[test]
2307 #[cfg(not(unix))]
2308 fn connection_sharing_is_unix_only() {
2309 let mut args = vec!["-o".to_owned(), "BatchMode=yes".to_owned()];
2310 let ssh = SshTarget {
2311 destination: "host".to_owned(),
2312 ssh_args: Vec::new(),
2313 };
2314 push_connection_reuse_args(&mut args, &ssh);
2315 assert_eq!(args, vec!["-o".to_owned(), "BatchMode=yes".to_owned()]);
2316 }
2317
2318 #[test]
2319 #[cfg(unix)]
2320 fn the_escape_hatch_accepts_the_usual_off_spellings() {
2321 for value in ["0", "off", "FALSE", " no "] {
2322 assert!(
2323 sharing_disabled(Some(std::ffi::OsStr::new(value))),
2324 "{value:?} must disable connection sharing"
2325 );
2326 }
2327 for value in ["1", "auto", "", "yes"] {
2328 assert!(
2329 !sharing_disabled(Some(std::ffi::OsStr::new(value))),
2330 "{value:?} must leave connection sharing on"
2331 );
2332 }
2333 assert!(!sharing_disabled(None));
2334 }
2335
2336 #[test]
2340 #[cfg(unix)]
2341 fn a_leased_session_runs_through_an_opened_master_on_a_real_host() {
2342 let _guard = SHARING_TEST_LOCK
2343 .lock()
2344 .unwrap_or_else(std::sync::PoisonError::into_inner);
2345 let Some(host) = std::env::var_os("MJ_E2E_SSH_HOST") else {
2346 return;
2347 };
2348 let host = host.to_string_lossy().into_owned();
2349 let socket_dir = sharing_socket_dir();
2350 let ssh = SshTarget {
2351 destination: host.clone(),
2352 ssh_args: vec!["-o".to_owned(), "BatchMode=yes".to_owned()],
2353 };
2354 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2355 socket_dir.path().to_path_buf(),
2356 )));
2357 let output = ProcessExecutor.execute(&ssh_command(&ssh, ["true"]));
2358 set_ssh_connection_sharing_for_test(None);
2359 let socket = socket_dir.path().join(control_socket_name(&ssh, 0));
2360 let check = ProcessExecutor
2361 .execute(&master_check_command(&ssh, &socket))
2362 .expect("ssh -O check must run");
2363 let exit = std::process::Command::new("ssh")
2364 .args([
2365 "-O",
2366 "exit",
2367 "-o",
2368 &format!("ControlPath={}", socket.display()),
2369 &host,
2370 ])
2371 .output();
2372 let output = output.expect("ssh must run");
2373 assert_eq!(
2374 output.status,
2375 0,
2376 "ssh {host} true failed: {}",
2377 String::from_utf8_lossy(&output.stderr)
2378 );
2379 assert_eq!(
2380 check.status,
2381 0,
2382 "no master is running: {}",
2383 String::from_utf8_lossy(&check.stderr)
2384 );
2385 drop(exit);
2386 }
2387
2388 #[test]
2393 #[cfg(unix)]
2394 fn sessions_shard_across_masters_on_a_real_host() {
2395 let _guard = SHARING_TEST_LOCK
2396 .lock()
2397 .unwrap_or_else(std::sync::PoisonError::into_inner);
2398 let Some(host) = std::env::var_os("MJ_E2E_SSH_HOST") else {
2399 return;
2400 };
2401 let host = host.to_string_lossy().into_owned();
2402 let socket_dir = sharing_socket_dir();
2403 let ssh = SshTarget {
2404 destination: host.clone(),
2405 ssh_args: vec!["-o".to_owned(), "BatchMode=yes".to_owned()],
2406 };
2407 let ledger = SessionLedger::new(3);
2408 let leases: Vec<SshSessionLease> = (0..7)
2409 .map(|_| {
2410 ledger
2411 .lease(&ssh, socket_dir.path(), &ProcessExecutor)
2412 .expect("lease a session on a real host")
2413 })
2414 .collect();
2415 let sockets: Vec<PathBuf> = (0..4)
2416 .map(|shard| socket_dir.path().join(control_socket_name(&ssh, shard)))
2417 .collect();
2418 let exit_all = || {
2419 for socket in &sockets {
2420 let _ = std::process::Command::new("ssh")
2421 .args([
2422 "-O",
2423 "exit",
2424 "-o",
2425 &format!("ControlPath={}", socket.display()),
2426 &host,
2427 ])
2428 .output();
2429 }
2430 };
2431
2432 let base = ssh_command(&ssh, ["sleep", "2"]);
2434 let children: Vec<std::io::Result<std::process::Output>> = std::thread::scope(|scope| {
2435 let handles: Vec<_> = leases
2436 .iter()
2437 .map(|lease| {
2438 let args = session_command_args(&base.program, &base.args, &ssh, lease);
2439 scope.spawn(move || {
2440 std::process::Command::new("ssh")
2441 .args(args)
2442 .stdin(std::process::Stdio::null())
2443 .output()
2444 })
2445 })
2446 .collect();
2447 handles
2448 .into_iter()
2449 .map(|handle| handle.join().expect("session thread"))
2450 .collect()
2451 });
2452 let running: Vec<bool> = sockets
2453 .iter()
2454 .map(|socket| {
2455 ProcessExecutor
2456 .execute(&master_check_command(&ssh, socket))
2457 .map(|output| output.status == 0)
2458 .unwrap_or(false)
2459 })
2460 .collect();
2461 let orphan = std::process::Command::new("ssh")
2462 .args(session_command_args(
2463 &base.program,
2464 &base.args,
2465 &ssh,
2466 &SshSessionLease {
2467 slot: Some(LeasedSlot {
2468 ledger: Arc::clone(&ledger),
2469 key: connection_key(&ssh),
2470 shard: 9,
2471 socket: socket_dir.path().join(control_socket_name(&ssh, 9)),
2472 }),
2473 probe: false,
2474 },
2475 ))
2476 .stdin(std::process::Stdio::null())
2477 .output();
2478 drop(leases);
2479 exit_all();
2480
2481 let shards: Vec<usize> = leases_per_shard(&ledger, &ssh);
2482 assert_eq!(shards, [0, 0, 0], "every slot is freed on drop");
2483 for (index, output) in children.iter().enumerate() {
2484 let output = output.as_ref().expect("ssh must run");
2485 assert_eq!(
2486 output.status.code(),
2487 Some(0),
2488 "session {index} failed: {}",
2489 String::from_utf8_lossy(&output.stderr)
2490 );
2491 }
2492 assert_eq!(
2493 running,
2494 [true, true, true, false],
2495 "seven sessions at three per master"
2496 );
2497 let orphan = orphan.expect("ssh must run");
2498 assert_eq!(
2499 orphan.status.code(),
2500 Some(255),
2501 "a guarded session with no master must not connect: {}",
2502 String::from_utf8_lossy(&orphan.stderr)
2503 );
2504 }
2505
2506 #[cfg(unix)]
2507 fn leases_per_shard(ledger: &SessionLedger, ssh: &SshTarget) -> Vec<usize> {
2508 ledger
2509 .connections()
2510 .get(&connection_key(ssh))
2511 .map(|shards| shards.iter().map(|shard| shard.leased).collect())
2512 .unwrap_or_default()
2513 }
2514
2515 #[test]
2519 #[cfg(unix)]
2520 fn scp_translates_the_ssh_port_option_and_is_tagged_with_its_destination() {
2521 let _guard = SHARING_TEST_LOCK
2522 .lock()
2523 .unwrap_or_else(std::sync::PoisonError::into_inner);
2524 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Disabled));
2525 let ssh = SshTarget {
2526 destination: "build@10.0.0.1".into(),
2527 ssh_args: vec!["-p".into(), "2222".into()],
2528 };
2529
2530 let upload = scp_upload(&ssh, Path::new("/tmp/local"), "remote/path", true);
2531 let download = scp_download(&ssh, "remote/archive.zip", "/tmp/local.zip");
2532 set_ssh_connection_sharing_for_test(None);
2533
2534 assert_eq!(
2535 upload.args,
2536 [
2537 "-P",
2538 "2222",
2539 "-r",
2540 "/tmp/local",
2541 "build@10.0.0.1:remote/path"
2542 ]
2543 );
2544 assert_eq!(
2545 download.args,
2546 [
2547 "-P",
2548 "2222",
2549 "build@10.0.0.1:remote/archive.zip",
2550 "/tmp/local.zip"
2551 ]
2552 );
2553 for command in [upload, download] {
2554 assert_eq!(command.program, "scp");
2555 assert_eq!(command.ssh_destination.as_deref(), Some("build@10.0.0.1"));
2556 }
2557 }
2558
2559 #[cfg(unix)]
2563 #[derive(Default)]
2564 struct FakeMasters {
2565 running: std::cell::RefCell<BTreeSet<String>>,
2566 refuse_open: std::cell::Cell<Option<&'static str>>,
2567 socket_existed_at_open: std::cell::RefCell<Vec<bool>>,
2568 seen: std::cell::RefCell<Vec<CommandSpec>>,
2569 }
2570
2571 #[cfg(unix)]
2572 impl FakeMasters {
2573 fn socket(command: &CommandSpec) -> String {
2574 command
2575 .args
2576 .iter()
2577 .find_map(|arg| arg.strip_prefix("ControlPath="))
2578 .expect("every master command names its socket")
2579 .to_owned()
2580 }
2581
2582 fn kill(&self, socket: &Path) {
2583 self.running
2584 .borrow_mut()
2585 .remove(&socket.display().to_string());
2586 }
2587
2588 fn commands(&self) -> usize {
2589 self.seen.borrow().len()
2590 }
2591
2592 fn openers(&self) -> usize {
2593 self.seen
2594 .borrow()
2595 .iter()
2596 .filter(|command| command.args.contains(&"ControlMaster=yes".to_owned()))
2597 .count()
2598 }
2599 }
2600
2601 #[cfg(unix)]
2602 impl CommandExecutor for FakeMasters {
2603 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2604 self.seen.borrow_mut().push(command.clone());
2605 assert_eq!(command.program, "ssh");
2606 let socket = Self::socket(command);
2607 let (status, stderr) = if command.args.windows(2).any(|pair| pair == ["-O", "check"]) {
2608 if self.running.borrow().contains(&socket) {
2609 (0, "")
2610 } else {
2611 (255, "Control socket connect: No such file or directory")
2612 }
2613 } else if command.args.contains(&"ControlMaster=yes".to_owned()) {
2614 self.socket_existed_at_open
2615 .borrow_mut()
2616 .push(Path::new(&socket).exists());
2617 match self.refuse_open.get() {
2618 Some(stderr) => (255, stderr),
2619 None => {
2620 self.running.borrow_mut().insert(socket);
2621 (0, "")
2622 }
2623 }
2624 } else {
2625 panic!("the ledger ran an unexpected command: {command:?}");
2626 };
2627 Ok(CommandOutput {
2628 status,
2629 stdout: Vec::new(),
2630 stderr: stderr.as_bytes().to_vec(),
2631 })
2632 }
2633 }
2634
2635 #[cfg(unix)]
2636 fn shard_of(lease: &SshSessionLease) -> String {
2637 let path = lease.control_path().expect("a shared lease has a socket");
2638 let name = path.file_name().unwrap().to_string_lossy().into_owned();
2639 name.rsplit('-').next().unwrap().to_owned()
2640 }
2641
2642 #[cfg(unix)]
2643 fn plain_target(destination: &str) -> SshTarget {
2644 SshTarget {
2645 destination: destination.to_owned(),
2646 ssh_args: Vec::new(),
2647 }
2648 }
2649
2650 #[test]
2651 #[cfg(unix)]
2652 fn leases_fill_the_lowest_shard_and_open_another_at_the_cap() {
2653 let dir = sharing_socket_dir();
2654 let ledger = SessionLedger::new(2);
2655 let ssh = plain_target("host");
2656 let masters = FakeMasters::default();
2657
2658 let first = ledger.lease(&ssh, dir.path(), &masters).expect("first");
2659 assert_eq!(masters.commands(), 3);
2661 let second = ledger.lease(&ssh, dir.path(), &masters).expect("second");
2662 assert_eq!(
2663 masters.commands(),
2664 3,
2665 "a master verified moments ago is not checked again"
2666 );
2667 let third = ledger.lease(&ssh, dir.path(), &masters).expect("third");
2668 assert_eq!(
2669 [&first, &second, &third].map(shard_of),
2670 ["0", "0", "1"].map(str::to_owned)
2671 );
2672 assert_eq!(masters.openers(), 2, "one master per shard");
2673 assert_eq!(
2674 third.control_path().unwrap(),
2675 dir.path().join(control_socket_name(&ssh, 1))
2676 );
2677
2678 drop(first);
2679 let fourth = ledger.lease(&ssh, dir.path(), &masters).expect("fourth");
2680 assert_eq!(shard_of(&fourth), "0", "a freed slot is reused first");
2681 assert_eq!(masters.openers(), 2);
2682 }
2683
2684 #[test]
2689 #[cfg(unix)]
2690 fn probes_are_counted_on_the_shard_they_join_without_opening_it() {
2691 let dir = sharing_socket_dir();
2692 let ledger = SessionLedger::new(2);
2693 let ssh = plain_target("probe-host");
2694 let masters = FakeMasters::default();
2695
2696 let session = ledger.lease(&ssh, dir.path(), &masters).expect("session");
2697 let probe = ledger.lease_probe(&ssh, dir.path());
2698 assert_eq!(leases_per_shard(&ledger, &ssh), [2]);
2699 let second_probe = ledger.lease_probe(&ssh, dir.path());
2700 assert_eq!(
2701 [&session, &probe, &second_probe].map(shard_of),
2702 ["0", "0", "1"].map(str::to_owned)
2703 );
2704 assert_eq!(masters.openers(), 1, "a probe never opens a master");
2705
2706 let mut args = Vec::new();
2707 push_session_args(&mut args, &probe);
2708 assert!(
2709 !args.contains(&"ProxyCommand=false".to_owned()),
2710 "a probe may connect directly when its master is down: {args:?}"
2711 );
2712 drop(probe);
2713 drop(second_probe);
2714 assert_eq!(leases_per_shard(&ledger, &ssh), [1, 0]);
2715 let next = ledger.lease(&ssh, dir.path(), &masters).expect("next");
2716 assert_eq!(shard_of(&next), "0");
2717 }
2718
2719 #[test]
2720 #[cfg(unix)]
2721 fn separate_connections_are_counted_separately() {
2722 let dir = sharing_socket_dir();
2723 let ledger = SessionLedger::new(1);
2724 let masters = FakeMasters::default();
2725 let first = ledger
2726 .lease(&plain_target("one"), dir.path(), &masters)
2727 .expect("one");
2728 let second = ledger
2729 .lease(&plain_target("two"), dir.path(), &masters)
2730 .expect("two");
2731 assert_eq!(
2732 [&first, &second].map(shard_of),
2733 ["0", "0"].map(str::to_owned)
2734 );
2735 assert_ne!(first.control_path(), second.control_path());
2736 }
2737
2738 #[test]
2739 #[cfg(unix)]
2740 fn an_invalidated_lease_makes_the_next_lease_reopen_a_dead_master() {
2741 let dir = sharing_socket_dir();
2742 let ledger = SessionLedger::new(8);
2743 let ssh = plain_target("host");
2744 let masters = FakeMasters::default();
2745 let first = ledger.lease(&ssh, dir.path(), &masters).expect("first");
2746 masters.kill(first.control_path().unwrap());
2747
2748 drop(ledger.lease(&ssh, dir.path(), &masters).expect("trusted"));
2750 assert_eq!(masters.openers(), 1);
2751
2752 first.invalidate();
2753 let second = ledger.lease(&ssh, dir.path(), &masters).expect("reopened");
2754 assert_eq!(masters.openers(), 2);
2755 assert_eq!(first.control_path(), second.control_path());
2756 }
2757
2758 #[test]
2759 #[cfg(unix)]
2760 fn a_master_that_cannot_be_opened_is_an_error_naming_the_destination() {
2761 let dir = sharing_socket_dir();
2762 let ledger = SessionLedger::new(1);
2763 let ssh = plain_target("build@10.0.0.1");
2764 let masters = FakeMasters::default();
2765 masters
2766 .refuse_open
2767 .set(Some("Permission denied (publickey)."));
2768
2769 let error = ledger
2770 .lease(&ssh, dir.path(), &masters)
2771 .expect_err("no master means no session");
2772 let message = format!("{error:#}");
2773 assert!(message.contains("build@10.0.0.1"), "{message}");
2774 assert!(message.contains("Permission denied"), "{message}");
2775 assert_eq!(
2776 masters.openers(),
2777 1,
2778 "the opener is not retried by the ledger"
2779 );
2780
2781 masters.refuse_open.set(None);
2784 let lease = ledger.lease(&ssh, dir.path(), &masters).expect("opens");
2785 assert_eq!(shard_of(&lease), "0");
2786 }
2787
2788 #[cfg(unix)]
2794 #[derive(Default)]
2795 struct RacingMasters {
2796 bound: Mutex<BTreeSet<String>>,
2797 masters: AtomicUsize,
2798 orphans: AtomicUsize,
2799 }
2800
2801 #[cfg(unix)]
2802 impl CommandExecutor for RacingMasters {
2803 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2804 let socket = FakeMasters::socket(command);
2805 let status = if command.args.windows(2).any(|pair| pair == ["-O", "check"]) {
2806 if self.bound.lock().unwrap().contains(&socket) {
2807 0
2808 } else {
2809 255
2810 }
2811 } else {
2812 std::thread::sleep(Duration::from_millis(100));
2813 if self.bound.lock().unwrap().insert(socket) {
2814 self.masters.fetch_add(1, Ordering::SeqCst);
2815 } else {
2816 self.orphans.fetch_add(1, Ordering::SeqCst);
2817 }
2818 0
2819 };
2820 Ok(CommandOutput {
2821 status,
2822 stdout: Vec::new(),
2823 stderr: Vec::new(),
2824 })
2825 }
2826 }
2827
2828 #[test]
2832 #[cfg(unix)]
2833 fn two_processes_opening_one_socket_open_one_master() {
2834 let dir = sharing_socket_dir();
2835 let ssh = plain_target("racing-host");
2836 let fake = RacingMasters::default();
2837 std::thread::scope(|scope| {
2838 for _ in 0..2 {
2839 scope.spawn(|| {
2840 let ledger = SessionLedger::new(8);
2842 ledger.lease(&ssh, dir.path(), &fake).expect("lease");
2843 });
2844 }
2845 });
2846 assert_eq!(fake.masters.load(Ordering::SeqCst), 1);
2847 assert_eq!(fake.orphans.load(Ordering::SeqCst), 0);
2848 }
2849
2850 #[test]
2851 #[cfg(unix)]
2852 fn a_stale_socket_is_removed_before_the_master_is_opened() {
2853 let dir = sharing_socket_dir();
2854 let ledger = SessionLedger::new(8);
2855 let ssh = plain_target("host");
2856 let socket = dir.path().join(control_socket_name(&ssh, 0));
2857 fs::write(&socket, b"").expect("stale socket stand-in");
2858 let masters = FakeMasters::default();
2859
2860 ledger.lease(&ssh, dir.path(), &masters).expect("opens");
2861
2862 assert_eq!(*masters.socket_existed_at_open.borrow(), [false]);
2863 }
2864
2865 #[test]
2866 #[cfg(unix)]
2867 fn the_opener_is_an_admitted_batch_master_and_the_check_is_local() {
2868 let ssh = SshTarget {
2869 destination: "host".to_owned(),
2870 ssh_args: vec!["-J".to_owned(), "jump".to_owned()],
2871 };
2872 let socket = Path::new("/run/mj/abc-0");
2873 let open = master_open_command(&ssh, socket);
2874 assert_eq!(
2875 open.args,
2876 [
2877 "-J",
2878 "jump",
2879 "-o",
2880 "BatchMode=yes",
2881 "-o",
2882 "ConnectTimeout=10",
2883 "-f",
2884 "-N",
2885 "-o",
2886 "ControlMaster=yes",
2887 "-o",
2888 "ControlPath=/run/mj/abc-0",
2889 "-o",
2890 &format!("ControlPersist={CONTROL_PERSIST}"),
2891 "host",
2892 ]
2893 );
2894 assert_eq!(open.ssh_destination.as_deref(), Some("host"));
2895
2896 let check = master_check_command(&ssh, socket);
2897 assert_eq!(
2898 check.args,
2899 [
2900 "-J",
2901 "jump",
2902 "-o",
2903 "ControlPath=/run/mj/abc-0",
2904 "-O",
2905 "check",
2906 "host"
2907 ]
2908 );
2909 assert_eq!(
2910 check.ssh_destination, None,
2911 "a check opens no connection and takes no admission permit"
2912 );
2913 }
2914
2915 #[test]
2916 #[cfg(unix)]
2917 fn a_master_open_times_out_during_handshake_and_honors_the_users_shorter_budget() {
2918 use std::net::TcpListener;
2919 use std::sync::mpsc;
2920
2921 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
2922 let port = listener.local_addr().unwrap().port();
2923 listener.set_nonblocking(true).unwrap();
2924 let (finish, stopping) = mpsc::channel();
2925 let server = std::thread::spawn(move || {
2926 let deadline = Instant::now() + Duration::from_secs(5);
2927 loop {
2928 match listener.accept() {
2929 Ok((connection, _)) => {
2930 let _ = stopping.recv_timeout(Duration::from_secs(5));
2933 drop(connection);
2934 return;
2935 }
2936 Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
2937 if stopping.try_recv().is_ok() || Instant::now() >= deadline {
2938 return;
2939 }
2940 std::thread::sleep(Duration::from_millis(5));
2941 }
2942 Err(error) => panic!("accept stalled SSH handshake: {error}"),
2943 }
2944 }
2945 });
2946 let directory = tempfile::tempdir().unwrap();
2947 let ssh = SshTarget {
2948 destination: "127.0.0.1".into(),
2949 ssh_args: vec![
2950 "-F".into(),
2951 "/dev/null".into(),
2952 "-p".into(),
2953 port.to_string(),
2954 "-o".into(),
2955 "ConnectTimeout=1".into(),
2956 ],
2957 };
2958 let started = Instant::now();
2959 let result = CancellableProcessExecutor::with_timeout(Duration::from_secs(3))
2962 .run_once(&master_open_command(&ssh, &directory.path().join("master")));
2963 let _ = finish.send(());
2964 server.join().unwrap();
2965 let output = result.expect("SSH's handshake timeout must beat the executor deadline");
2966 assert_eq!(output.status, 255);
2967 let stderr = String::from_utf8_lossy(&output.stderr);
2968 assert!(stderr.contains("timed out"), "{stderr}");
2969 assert!(started.elapsed() < Duration::from_secs(3));
2970 }
2971
2972 #[test]
2973 #[cfg(unix)]
2974 fn session_args_forbid_a_direct_connection() {
2975 let dir = sharing_socket_dir();
2976 let ledger = SessionLedger::new(8);
2977 let masters = FakeMasters::default();
2978 let lease = ledger
2979 .lease(&plain_target("host"), dir.path(), &masters)
2980 .expect("lease");
2981 let mut args = Vec::new();
2982 push_session_args(&mut args, &lease);
2983 assert_eq!(
2984 args,
2985 [
2986 "-o".to_owned(),
2987 "ControlMaster=no".to_owned(),
2988 "-o".to_owned(),
2989 format!("ControlPath={}", lease.control_path().unwrap().display()),
2990 "-o".to_owned(),
2991 "ProxyCommand=false".to_owned(),
2992 ]
2993 );
2994 }
2995
2996 #[test]
2999 #[cfg(unix)]
3000 fn unshared_connections_lease_without_a_socket() {
3001 let _guard = SHARING_TEST_LOCK
3002 .lock()
3003 .unwrap_or_else(std::sync::PoisonError::into_inner);
3004 let masters = FakeMasters::default();
3005 let dir = sharing_socket_dir();
3006 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
3007 dir.path().to_path_buf(),
3008 )));
3009 let user_owned = SshSessions::lease(
3010 &SshTarget {
3011 destination: "unshared-user-host".to_owned(),
3012 ssh_args: vec!["-S".to_owned(), "/tmp/mine".to_owned()],
3013 },
3014 &masters,
3015 );
3016 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Disabled));
3017 let disabled = SshSessions::lease(&plain_target("unshared-disabled-host"), &masters);
3018 set_ssh_connection_sharing_for_test(None);
3019
3020 for lease in [user_owned, disabled] {
3021 let lease = lease.expect("an unshared lease never fails");
3022 assert_eq!(lease.control_path(), None);
3023 let mut args = Vec::new();
3024 push_session_args(&mut args, &lease);
3025 assert!(args.is_empty());
3026 }
3027 assert_eq!(masters.commands(), 0);
3028 }
3029
3030 #[test]
3034 fn a_refused_session_is_named_apart_from_a_pre_authentication_hangup() {
3035 assert_eq!(
3036 ssh_refusal(
3037 255,
3038 "mux_client_request_session: session request failed: Session open refused by peer\n\
3039 kex_exchange_identification: Connection closed by remote host\n\
3040 Connection closed by UNKNOWN port 65535"
3041 ),
3042 Some(SshRefusal::SessionLimit)
3043 );
3044 assert_eq!(
3045 ssh_refusal(255, "Connection closed by 192.168.1.77 port 22"),
3046 Some(SshRefusal::BeforeAuthentication)
3047 );
3048 assert_eq!(ssh_refusal(1, "Session open refused by peer"), None);
3049 assert!(
3050 !SshRefusal::SessionLimit
3051 .retry_message()
3052 .contains("before authentication")
3053 );
3054 }
3055
3056 #[test]
3057 fn transport_rejection_matches_only_sshd_hangups() {
3058 let cases: [(i32, &str, bool); 7] = [
3059 (255, "Connection closed by 192.168.1.77 port 22", true),
3060 (
3061 255,
3062 "kex_exchange_identification: read: Connection reset by peer",
3063 true,
3064 ),
3065 (255, "ssh: Connection reset by 10.0.0.1 port 22", true),
3066 (255, "Connection timed out during banner exchange", true),
3067 (255, "Permission denied (publickey).", false),
3068 (
3069 255,
3070 "ssh: connect to host h port 22: Connection refused",
3071 false,
3072 ),
3073 (1, "Connection closed by 192.168.1.77 port 22", false),
3074 ];
3075 for (status, stderr, expected) in cases {
3076 assert_eq!(
3077 is_transport_rejection(status, stderr),
3078 expected,
3079 "status {status} stderr {stderr:?}"
3080 );
3081 }
3082 }
3083
3084 #[test]
3085 fn admission_never_admits_more_than_the_limit() {
3086 let gate = DestinationGate::new(2);
3087 let in_flight = Arc::new(AtomicUsize::new(0));
3088 let peak = Arc::new(AtomicUsize::new(0));
3089 let threads: Vec<_> = (0..12)
3090 .map(|_| {
3091 let gate = Arc::clone(&gate);
3092 let in_flight = Arc::clone(&in_flight);
3093 let peak = Arc::clone(&peak);
3094 std::thread::spawn(move || {
3095 for _ in 0..25 {
3096 let permit = gate.acquire();
3097 let now = in_flight.fetch_add(1, Ordering::SeqCst) + 1;
3098 peak.fetch_max(now, Ordering::SeqCst);
3099 std::thread::yield_now();
3100 in_flight.fetch_sub(1, Ordering::SeqCst);
3101 drop(permit);
3102 }
3103 })
3104 })
3105 .collect();
3106 for thread in threads {
3107 thread.join().expect("admission worker must not panic");
3108 }
3109 assert!(
3110 peak.load(Ordering::SeqCst) <= 2,
3111 "admission let {} connections run against a 2-permit gate",
3112 peak.load(Ordering::SeqCst)
3113 );
3114 assert_eq!(in_flight.load(Ordering::SeqCst), 0);
3115 }
3116
3117 #[test]
3118 fn admission_blocks_once_every_permit_is_held() {
3119 let gate = DestinationGate::new(2);
3120 let first = gate.acquire();
3121 let second = gate.acquire();
3122 let waiter = {
3123 let gate = Arc::clone(&gate);
3124 std::thread::spawn(move || {
3125 let permit = gate.acquire();
3126 drop(permit);
3127 })
3128 };
3129 std::thread::sleep(std::time::Duration::from_millis(50));
3131 assert!(!waiter.is_finished());
3132 drop(first);
3133 waiter
3134 .join()
3135 .expect("waiter must be admitted once a permit frees");
3136 drop(second);
3137 }
3138
3139 #[test]
3140 fn cancelled_admission_does_not_wait_for_the_holder_or_consume_a_slot() {
3141 let gate = DestinationGate::new(1);
3142 let held = gate.acquire();
3143 let executor = CancellableProcessExecutor::with_timeout(Duration::from_millis(50));
3144 let error = gate
3145 .acquire_unless(&|| executor.is_cancelled())
3146 .unwrap_err();
3147 assert!(error.to_string().contains("cancelled"));
3148 assert_eq!(*gate.in_flight.lock().unwrap(), 1);
3149 drop(held);
3150 assert!(gate.acquire_unless(&|| false).is_ok());
3151 }
3152
3153 #[cfg(unix)]
3154 #[test]
3155 fn master_file_and_thread_admission_honor_the_executor_deadline() {
3156 let dir = sharing_socket_dir();
3157 let socket = dir.path().join("held-master");
3158 let _held = lock_master_opening(&socket).unwrap();
3159 let executor = CancellableProcessExecutor::with_timeout(Duration::from_millis(50));
3160 assert!(lock_master_opening_unless(&socket, &|| executor.is_cancelled()).is_err());
3161
3162 let ledger = SessionLedger::new(2);
3163 let ssh = plain_target("cancelled-master-host");
3164 let (_, opening) = ledger.reserve(&connection_key(&ssh));
3165 let _opening = opening.lock().unwrap();
3166 let executor = CancellableProcessExecutor::with_timeout(Duration::from_millis(50));
3167 assert!(ledger.lease(&ssh, dir.path(), &executor).is_err());
3168 assert_eq!(ledger.connections()[&connection_key(&ssh)][0].leased, 1);
3169 }
3170}