Skip to main content

mj_core/targets/
ssh.rs

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
9/// The connectivity probe `mj doctor` runs against an SSH target.
10///
11/// It reuses the provisioning argument order so the probe fails exactly where
12/// a real session would, with two deliberate overrides prepended. OpenSSH
13/// honours the first occurrence of an option, so these win over the
14/// provisioning defaults: `BatchMode=yes` never prompts for a password, and
15/// `StrictHostKeyChecking=yes` never accepts an unknown host key. Doctor
16/// diagnoses; the user decides whether to trust a key.
17///
18/// Those overrides are also why the probe joins a shared master but never
19/// opens one (see [`CommandSpec::ssh_probe_session`]): as the master it would hold
20/// them over every later session on that connection, and a plain `mj doctor`
21/// would leave an `ssh` process behind for the whole `ControlPersist` window
22/// even though the user asked only for a diagnosis.
23pub 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    // The socket is named after the target as configured, not after these
34    // probe-only overrides, so the probe finds the daemon's master.
35    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
52/// Build an `ssh` command that runs as one session on a shared connection.
53/// The executor leases the session and adds its options just before the
54/// command is spawned; see [`CommandSpec::ssh_session`].
55pub 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
62/// Home-relative directory on an SSH host where files bound for a remote
63/// container wait before the engine copies them in. Home-relative rather than
64/// `~/`, because `ssh_command` quotes every argument while `scp` expands `~`.
65pub const REMOTE_UPLOAD_STAGING: &str = ".cache/mjolnir/uploads";
66
67/// Upload a local file or directory to the SSH host.
68pub 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
78/// Download a remote file from the SSH host.
79pub 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
86/// The connection's `ssh` arguments rewritten for `scp`, which spells the port
87/// option `-P`; to `scp`, `-p` means "preserve file times".
88fn 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    // `scp` runs `ssh` underneath, so it takes a session on a shared
103    // connection and is admitted and retried the same way.
104    CommandSpec::new("scp", args).ssh_session(ssh)
105}
106
107/// How long a shared master connection stays alive after its last channel
108/// closes. The master is an `ssh` process that outlives the daemon by this
109/// long, so it is kept short enough to be unsurprising and long enough to
110/// cover a whole provision.
111#[cfg(unix)]
112const CONTROL_PERSIST: &str = "60";
113
114/// Environment override that turns connection sharing off. Any of `0`, `off`,
115/// `false`, or `no` disables it.
116pub const CONTROL_MASTER_ENV: &str = "MJ_SSH_CONTROL_MASTER";
117
118/// Longest `ControlPath` that still fits in a Unix socket address. `sun_path`
119/// holds 108 bytes on Linux and 104 on macOS, minus the terminating NUL.
120#[cfg(unix)]
121const MAX_CONTROL_PATH: usize = 103;
122
123/// Hex digits of the connection hash in a socket name. 64 bits keeps the
124/// names short while making a collision between two targets implausible.
125#[cfg(unix)]
126const CONNECTION_HASH_HEX: usize = 16;
127
128/// Bytes reserved after the directory for a socket's name: a separator, the
129/// connection hash, `-` and a shard index of up to four digits, and the
130/// `.` plus 16 random characters `ssh` appends to the path while it binds a
131/// new master.
132#[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/// How a test pins connection sharing instead of letting it resolve from the
147/// environment.
148#[doc(hidden)]
149#[derive(Debug, Clone)]
150pub enum SshSharingForTest {
151    /// Behave as though the escape hatch were set.
152    Disabled,
153    /// Keep the control sockets in this directory.
154    Directory(PathBuf),
155}
156
157static SHARING_OVERRIDE: Mutex<Option<SshSharingForTest>> = Mutex::new(None);
158
159/// The connection-sharing override is process-wide, so the tests that set it
160/// take turns.
161#[cfg(all(test, unix))]
162pub(super) static SHARING_TEST_LOCK: Mutex<()> = Mutex::new(());
163
164/// Pin connection sharing for a test, or restore the real resolution with
165/// `None`. Tests must not depend on the developer's `$XDG_RUNTIME_DIR` or home
166/// directory, so every test that inspects `ssh` arguments pins it. Not part of
167/// the daemon's behaviour.
168#[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/// The directory holding this instance's control sockets, or `None` when
184/// sharing is off.
185///
186/// `$XDG_RUNTIME_DIR/mjolnir/<identity>` is preferred because it is short,
187/// per-user, and on tmpfs; the instance's data directory is the fallback.
188/// Each instance gets its own directory because each daemon counts only its
189/// own sessions: two daemons sharing masters would together exceed the
190/// server's per-connection session limit. Neither location is
191/// world-writable, and the directory is created 0700 because `ssh` will not
192/// create it itself.
193#[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/// Where an instance keeps its sockets when no test pins the directory.
219#[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        // The data directory is already specific to the instance.
226        _ => crate::config::data_dir().join("ssh"),
227    }
228}
229
230/// The name of this instance's directory under `$XDG_RUNTIME_DIR/mjolnir`:
231/// the identity [`crate::config::instance_identity`] stamps on the
232/// instance's workers. A named instance uses its name, and a data-directory
233/// override (`MJ_DATA_DIR`, as every end-to-end lab sets) uses the
234/// fingerprint of that directory, so neither shares the default instance's
235/// masters or sweeps its lock files. The default instance keeps `default`,
236/// where its sockets have always been.
237#[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/// Create the socket directory 0700 and reject one whose sockets would not fit
247/// in a Unix socket address. Failure means no sharing, never a failed command.
248#[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/// The identity of one configured connection: its destination and the
278/// user's own `ssh` arguments, which can change the port, user, or route.
279/// Two targets that differ in either get separate masters.
280#[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/// The socket file name for one shard of a connection: `<hash>-<shard>`.
291///
292/// Mjolnir names sockets itself rather than using `ssh`'s `%C` so that the
293/// name follows exactly the key the daemon counts sessions under, and so the
294/// daemon knows the concrete path when it must remove a stale socket.
295#[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/// Whether the user's own `ssh` arguments already configure connection
308/// sharing. OpenSSH keeps the first value it sees, so Mjolnir cannot add its
309/// own sharing options without either overriding the user or being
310/// overridden; a user who configures sharing owns it, and Mjolnir adds none.
311#[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
327/// Append the options that let this command *reuse* a shared master without
328/// ever becoming one.
329///
330/// Use this for commands that carry deliberately impatient options, such as
331/// the short `ConnectTimeout` and one-miss `ServerAlive` keepalive of a
332/// validation probe or a Tab completion. Those settings belong to the one
333/// command that asked for them. If such a command opened the master, the
334/// master would enforce them for its whole lifetime and drop every later
335/// multiplexed session -- an upload, a `podman run`, the worker bootstrap --
336/// on a stall of a couple of seconds. With `ControlMaster=no` the command
337/// joins the connection's first master when one is up and otherwise opens its
338/// own direct connection, keeping its fail-fast options to itself.
339///
340/// These commands run in processes without the daemon's session ledger
341/// (`mj doctor`, completion) or only validate a target, so they are not
342/// counted and are not bound to a master with `ProxyCommand=false`: for a
343/// diagnosis, a direct connection is the stated behaviour.
344pub 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
368/// Check whether a directory exists on the configured SSH host.
369pub 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
401/// What the caller is told when `ssh` refused the host's key, so it does not
402/// get only a daemon log reference (launch finding R3-6).
403///
404/// OpenSSH's wording is the only signal: "Host key verification failed." ends
405/// both an unknown key under strict checking and a key that changed. The
406/// sentence quotes that line and names no host, so it may reach any client;
407/// the full ssh text stays on the error chain for the daemon log. It names
408/// the known_hosts file the machine's ssh options name (launch finding R6-5).
409fn 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/// The known_hosts file ssh checks a host key against, as far as the options
433/// on its command line tell.
434#[derive(Debug, Clone, PartialEq, Eq)]
435enum KnownHostsFile {
436    /// No option names one, so ssh uses its default.
437    Default,
438    /// `-o UserKnownHostsFile=` names this file, or these files separated
439    /// by spaces.
440    Named(String),
441    /// A config file given with `-F` may name any file, or the option names
442    /// no usable file.
443    Unknown,
444}
445
446impl KnownHostsFile {
447    /// Read `ssh`'s arguments the way OpenSSH does: the first
448    /// `UserKnownHostsFile` given with `-o` wins, the keyword is matched
449    /// without regard to case, and its value follows `=` or a space.
450    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    /// The file as a refusal names it.
480    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    /// The first file an option names, which is the one ssh adds keys to.
489    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
497/// The value of `keyword` in one ssh option (`Keyword=value`,
498/// `Keyword value`, or `Keyword = value`), or `None` for another keyword.
499fn 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
512/// Verify that a bare-SSH project path exists and has a committed Git HEAD.
513pub 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
590/// Wrap a value so a POSIX shell reads it as one literal argument. Used at the
591/// SSH boundary here and when Hel rebuilds an agent's terminal command line
592/// (`terminal::shell_line`).
593pub 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
687/// Whether this locator names a target another session owns: a sub-agent
688/// child either borrowing its parent's container or running as its own worker
689/// inside the parent's SSH workspace.
690pub 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
817/// `ssh` reserves exit status 255 for its own transport failures; a remote
818/// command never produces it, so the remote side provably never ran.
819pub const SSH_TRANSPORT_EXIT_STATUS: i32 = 255;
820
821/// Stderr fragments OpenSSH prints when the server hangs up before
822/// authentication. `sshd`'s `MaxStartups` produces exactly these when it drops
823/// an unauthenticated connection, and so does a server that is still starting.
824const 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
831/// Whether a finished `ssh` process was turned away by the transport rather
832/// than by the remote command.
833///
834/// The remote command never started in this case, so the caller may retry the
835/// whole invocation without worrying about repeating a side effect.
836pub fn is_transport_rejection(status: i32, stderr: &str) -> bool {
837    ssh_refusal(status, stderr).is_some()
838}
839
840/// What `ssh` prints when the server refuses a new session on an existing,
841/// authenticated shared connection. `sshd` does this once the connection
842/// carries `MaxSessions` sessions.
843const SESSION_REFUSAL_MARKER: &str = "Session open refused by peer";
844
845/// Why the SSH server turned an `ssh` invocation away before its remote
846/// command started. Either way the command never ran, so it may be retried.
847#[derive(Debug, Clone, Copy, PartialEq, Eq)]
848pub enum SshRefusal {
849    /// The server dropped a new connection before authentication, which is
850    /// what `MaxStartups` does, and a dead shared master looks the same to a
851    /// session bound to it.
852    BeforeAuthentication,
853    /// The shared connection is up and authenticated, but the server refused
854    /// one more session on it (`MaxSessions`).
855    SessionLimit,
856}
857
858impl SshRefusal {
859    /// A log message that names this refusal.
860    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    /// Log one retry of a command this refusal turned away.
872    ///
873    /// A refused session on a live shared connection is routine while many
874    /// sessions start at once, and the retry nearly always gets in (launch
875    /// finding R3-5 counted 98 in one dashboard log, all answered), so it is
876    /// logged at debug level. A connection closed before authentication can
877    /// mean a master died, so it stays a warning. A command still refused
878    /// after its last attempt is logged by [`Self::log_exhausted`].
879    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    /// Log a command the server still refused on its last attempt.
913    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
930/// Classify a finished `ssh` process that the server turned away. A refused
931/// session is checked first: a session bound to its master with
932/// `ProxyCommand=false` also reports a closed connection after the refusal.
933pub 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
946/// Default number of `ssh` processes this daemon will have in flight against
947/// one destination at a time.
948///
949/// `sshd` counts *unauthenticated* connections against `MaxStartups`, whose
950/// stock value is `10:30:100`: from the eleventh concurrent pre-auth connection
951/// it starts dropping them, and past a hundred it drops all of them. A daemon
952/// that spawns one fresh `ssh` per operation reaches that during startup, so it
953/// admits its own connections instead of letting the server refuse them.
954const DEFAULT_MAX_CONCURRENT_SSH: usize = 6;
955
956/// Environment override for [`DEFAULT_MAX_CONCURRENT_SSH`].
957pub 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
964/// A positive whole number from the environment variable `name`, or
965/// `default` when it is unset or invalid.
966fn 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
987/// A counting semaphore per SSH destination.
988///
989/// Deliberately built on `std::sync` rather than a runtime primitive: the
990/// blocking process executors are called from plain threads as well as from
991/// `spawn_blocking`, and both must share one gate.
992struct 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
1038/// One admitted `ssh` connection. The slot is returned on drop.
1039pub 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
1062/// Process-wide admission control for outbound `ssh` connections.
1063pub struct SshAdmission;
1064
1065impl SshAdmission {
1066    /// Block until this process may open another `ssh` connection to
1067    /// `destination`. The returned permit holds the slot until it is dropped.
1068    pub fn acquire(destination: &str) -> SshPermit {
1069        Self::gate(destination).acquire()
1070    }
1071
1072    /// Wait for admission without outliving the command's cancellation or deadline.
1073    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/// Default number of sessions the daemon places on one shared connection.
1092///
1093/// A stock `sshd` refuses the eleventh session on one connection
1094/// (`MaxSessions 10`). Two are left free for `ssh` commands from other
1095/// Mjolnir processes on this machine, such as `mj doctor` and Tab completion,
1096/// which join a master without being counted here.
1097#[cfg(unix)]
1098const DEFAULT_SESSIONS_PER_CONNECTION: usize = 8;
1099
1100/// Environment override for [`DEFAULT_SESSIONS_PER_CONNECTION`].
1101pub const SESSIONS_PER_CONNECTION_ENV: &str = "MJ_SSH_SESSIONS_PER_CONNECTION";
1102
1103/// How long a successful `ssh -O check` of a master is trusted before the
1104/// next lease on that shard checks again.
1105#[cfg(unix)]
1106const MASTER_CHECK_INTERVAL: Duration = Duration::from_secs(5);
1107
1108/// Per-command deadline for checking or opening a master from code that has
1109/// no executor of its own, such as the relay and the resource pollers.
1110pub 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/// One master connection and the sessions the daemon has placed on it.
1121#[cfg(unix)]
1122struct Shard {
1123    leased: usize,
1124    /// When `ssh -O check` last found this shard's master running.
1125    verified_at: Option<Instant>,
1126    /// Serializes checking and opening this shard's master, so concurrent
1127    /// leases never start two openers for one socket.
1128    opening: Arc<Mutex<()>>,
1129}
1130
1131/// The daemon's count of sessions per shard, keyed by connection.
1132#[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    /// Count one session on the lowest shard of `key` with room, adding a
1159    /// shard when all are full. Returns the shard and its opening lock.
1160    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    /// Lease a session on the lowest shard with room, opening that shard's
1182    /// master first when it is not known to be running.
1183    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        // From here on the slot is released on drop, including on error.
1192        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            // Another lease may have checked or opened the master while this
1213            // one waited for the lock.
1214            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    /// Count a fail-fast probe on the lowest shard with room, without
1226    /// checking or opening that shard's master.
1227    ///
1228    /// The probe joins the master when it is up and otherwise connects on
1229    /// its own, so its short timeouts never become a master's settings. It
1230    /// still uses one of the master's sessions while it runs, and a probe
1231    /// that was not counted could take the session a leased command was
1232    /// promised.
1233    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/// Make sure a master is listening on `socket`, opening one if needed.
1249///
1250/// The master is opened explicitly, with `ControlMaster=yes`, and then
1251/// checked again. Nothing else is attempted when that fails: the caller gets
1252/// an error naming the destination instead of a direct connection.
1253#[cfg(unix)]
1254fn ensure_master(ssh: &SshTarget, socket: &Path, executor: &dyn CommandExecutor) -> Result<()> {
1255    // Another process of this instance (the old daemon during a restart)
1256    // may be checking and opening the same socket. Without this lock both
1257    // find no master and both open one; the second finds the socket bound,
1258    // prints "already exists, disabling multiplexing", and keeps a plain
1259    // background connection that no ControlPersist ever closes (J-18). The
1260    // lock also keeps one process from removing, as stale, a socket the
1261    // other has just bound.
1262    let _opening = lock_master_opening_unless(socket, &|| executor.cancellation_requested())?;
1263    if master_running(ssh, socket, executor)? {
1264        return Ok(());
1265    }
1266    // A master that died without cleaning up leaves its socket behind, and
1267    // `ssh` will not bind over it: the opener would print "already exists,
1268    // disabling multiplexing" and hold a plain connection instead.
1269    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/// Take the file lock that serializes checking and opening the master on
1301/// `socket` across processes. It is held until the returned file is dropped.
1302/// The lock file sits beside the socket as `<socket>.lock`; `ssh` binds a
1303/// new master at `<socket>.<16 random characters>`, so the names never meet.
1304#[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        // The daemon removes stale locks when it starts. A lock taken on a
1337        // file removed meanwhile serializes nothing, so take the one at the
1338        // path now.
1339        if is_file_at(&file, &path) {
1340            return Ok(file);
1341        }
1342    }
1343}
1344
1345/// Remove the master lock files in `dir` whose master is gone.
1346///
1347/// A lock outlives the master it guarded: ssh removes its socket when the
1348/// master exits, but nothing removed `<socket>.lock` (launch finding R3-11).
1349/// A lock is removed only while its socket is absent and no other process
1350/// holds it, and only if the file locked is still the one at its path.
1351/// [`lock_master_opening`] checks the same after it locks, so an opener that
1352/// opened the file just before it was removed takes a fresh one instead.
1353#[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        // Held: another process is checking or opening this master now.
1379        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/// Whether `file` is the file now at `path`, rather than one removed from it.
1395#[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/// `ssh -O check` asks the master on `socket` whether it is alive. It opens
1410/// no network connection, so it is not admitted like one.
1411#[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/// Open a master on `socket` and return once it is authenticated.
1425///
1426/// `-f -N` backgrounds the master after authentication without keeping the
1427/// caller's output pipes open, and `ControlPersist` stops it on its own once
1428/// its last session has been gone that long. `BatchMode=yes` keeps the daemon
1429/// from ever waiting on a password prompt. This is a real connection, so it
1430/// is admitted and retried like one.
1431#[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/// The ledger entry a lease holds; dropping it frees the slot.
1455#[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
1501/// A leased session slot on one shard of a shared connection. Dropping it
1502/// frees the slot.
1503///
1504/// A lease without a socket stands for a command that runs on its own
1505/// connection: sharing is switched off with `MJ_SSH_CONTROL_MASTER`, the
1506/// user's `ssh_args` configure sharing themselves, or the platform has no
1507/// connection sharing.
1508pub struct SshSessionLease {
1509    #[cfg(unix)]
1510    slot: Option<LeasedSlot>,
1511    /// A counted fail-fast probe: it joins the master when one is up and
1512    /// otherwise connects directly, so it carries no `ProxyCommand=false`.
1513    #[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    /// The control socket this session must use, or `None` when the command
1537    /// runs on its own connection.
1538    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    /// Forget that this lease's master was verified, so the next lease on the
1546    /// shard checks it and reopens it if needed. Call this when a session
1547    /// failed in a way that suggests the master is gone.
1548    pub fn invalidate(&self) {
1549        #[cfg(unix)]
1550        if let Some(slot) = &self.slot {
1551            slot.set_verified(None);
1552        }
1553    }
1554}
1555
1556/// Process-wide placement of `ssh` sessions on shared connections.
1557///
1558/// A stock `sshd` allows ten sessions per connection. The daemon therefore
1559/// spreads its sessions for one connection across several masters (shards),
1560/// each opened explicitly, and every other command joins one of them with
1561/// options that make a direct connection impossible.
1562pub struct SshSessions;
1563
1564impl SshSessions {
1565    /// Reserve a session on a shard for `ssh`, opening that shard's master if
1566    /// it is not running. Blocks like [`SshAdmission::acquire`], so it is
1567    /// callable from plain threads and `spawn_blocking`. Returns an error only
1568    /// when a master could not be opened or verified.
1569    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    /// Remove this instance's master lock files whose master has exited.
1588    /// The daemon calls this when it starts; failures are only logged.
1589    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    /// Count a fail-fast probe, such as a target validation or the
1597    /// connectivity check, on a shard of `ssh`'s connection without opening
1598    /// a master. See [`CommandSpec::ssh_probe_session`].
1599    ///
1600    /// In a process other than the daemon (`mj doctor`) the ledger is empty,
1601    /// so the probe joins the first master; the two sessions per master that
1602    /// the daemon leaves free are for these.
1603    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
1622/// Options for a command that runs as one session on an already open
1623/// master. `ProxyCommand=false` makes it impossible for `ssh` to open a
1624/// direct connection: a multiplexed client never runs the proxy command, and
1625/// a client that fails to reach the master exits 255 instead of connecting on
1626/// its own. Appends nothing for a lease without a socket.
1627///
1628/// A probe lease gets only `ControlMaster=no` and the `ControlPath`: it joins
1629/// the master when one is up and otherwise connects on its own.
1630pub 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
1648/// The argument list for `program` (`ssh` or `scp`) running as one session
1649/// on `lease`'s master: the session options first, then the command's own
1650/// arguments.
1651///
1652/// OpenSSH keeps the first value it sees for an option, so leading with the
1653/// session options makes them win over the user's `ssh_args` and
1654/// `ssh_config`, including any `ProxyCommand` or `ProxyJump`. `ssh` refuses
1655/// the `-J` flag after a `ProxyCommand` outright, so a `-J` in the user's own
1656/// arguments is rewritten to the equivalent `-o ProxyJump=`, which the
1657/// session's `ProxyCommand=false` then overrides; a session never connects
1658/// by itself, and the master was opened with the user's jump host. `scp`
1659/// already passes its `-J` on as `-oProxyJump=`.
1660pub 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
1691/// How many times a transport-rejected `ssh` invocation is tried in total.
1692pub const SSH_RETRY_ATTEMPTS: usize = 3;
1693
1694/// Inclusive millisecond bounds the jittered delay is drawn from, indexed by
1695/// the number of attempts already made. `sshd` sheds load for as long as its
1696/// pre-auth queue stays full, so the second wait is a multiple of the first.
1697const SSH_RETRY_BACKOFF_MS: [(u64, u64); SSH_RETRY_ATTEMPTS - 1] = [(500, 2_000), (2_000, 4_000)];
1698
1699/// Test override collapsing every retry delay to this many milliseconds.
1700/// `u64::MAX` means "no override".
1701static SSH_RETRY_BACKOFF_OVERRIDE_MS: AtomicU64 = AtomicU64::new(u64::MAX);
1702
1703/// Shorten the retry backoff so tests can drive the retry path without
1704/// sleeping for seconds. Not part of the daemon's behaviour.
1705#[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
1713/// The jittered wait before retry number `attempts_made + 1`.
1714///
1715/// Jitter matters more than the mean here: every session's reconnect fails at
1716/// the same instant, so an unjittered schedule would simply re-send the whole
1717/// burst into the same full queue.
1718pub 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    // A failed draw only costs jitter, so fall back to the lower bound.
1729    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    /// Launch finding R3-11: `*.lock` files stayed in the instance's socket
1742    /// directory after their masters exited. The daemon clears, when it
1743    /// starts, each lock whose master's socket is gone, and leaves a lock that
1744    /// guards a live socket or that another process holds.
1745    #[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        // A process another test forks while the holder is open shares its
1764        // lock until that child execs, so the sweep may find the lock still
1765        // held for a few milliseconds; the daemon would try again at its next
1766        // start.
1767        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    /// Records the commands it is handed and reports an empty success.
1872    #[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        // macOS's default temporary path leaves too little room for SSH's hash.
1893        tempfile::tempdir_in("/tmp").expect("short control socket directory")
1894    }
1895
1896    /// Lease a session for `command` through the process-wide ledger with
1897    /// the sockets pinned to `dir`, and return the arguments it would be
1898    /// spawned with.
1899    #[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    /// A built command carries a session request instead of sharing
1915    /// options; at spawn time the session options go in front of everything
1916    /// the user configured, so OpenSSH honours them over any proxy setting.
1917    #[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    /// `ssh` refuses `-J` after a `ProxyCommand`, so a session rewrites the
1982    /// user's `-J` to the `ProxyJump` option the session's guard overrides.
1983    /// `scp` already turns its `-J` into that option.
1984    #[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    /// A user who configures sharing in `ssh_args` owns it: Mjolnir adds no
2030    /// sharing options of its own, in any spelling OpenSSH accepts, and opens
2031    /// no master.
2032    #[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    /// Each instance keeps its own sockets, because each daemon counts only
2070    /// its own sessions against a master.
2071    #[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    /// A daemon isolated only by `MJ_DATA_DIR` (every luna lab) is another
2086    /// instance, so it must not share the default instance's masters or
2087    /// sweep its lock files (launch finding R5-1). Its directory is named by
2088    /// the same fingerprint `instance_identity` stamps on its workers.
2089    #[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        // A named instance keeps its name whatever its data directory, as
2110        // `instance_identity` does.
2111        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    /// Sockets are named `<hash>-<shard>`, and the hash follows the whole
2118    /// configured connection, not only the destination.
2119    #[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    /// A command with a two-second keepalive must join a master, never open
2142    /// one: as the master it would impose that keepalive on every later
2143    /// session sharing the connection.
2144    #[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    /// `mj doctor` diagnoses and exits. Its connectivity probe must join a
2216    /// master when one is up and otherwise open a plain connection, so a
2217    /// doctor run never leaves a `ControlPersist` master behind, and the
2218    /// probe's own strict overrides never bind a shared connection.
2219    #[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    /// Against a real host: leasing a session opens a master that
2337    /// `ssh -O check` finds, and a command runs through it. Set
2338    /// `MJ_E2E_SSH_HOST` to a reachable destination to run it.
2339    #[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    /// Against a real host: with three sessions per connection, seven
2389    /// concurrent sessions open three masters, every session runs through
2390    /// its master, and a guarded session with no master fails instead of
2391    /// connecting on its own. Set `MJ_E2E_SSH_HOST` to run it.
2392    #[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        // Run all seven at once so they really share their masters.
2433        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    /// `scp` spells the port `-P`; passing an `ssh` `-p` through would ask it
2516    /// to preserve file times and read the port as a file name. Every `scp`
2517    /// also opens a connection, so it is admitted like `ssh`.
2518    #[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    /// A hand-written stand-in for `ssh` that models masters: `-O check`
2560    /// succeeds only for a socket whose master it opened, and an opener
2561    /// starts one unless told to refuse. It records every command.
2562    #[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        // Check, open, check again.
2660        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    /// Target validation and the connectivity probe run in the daemon while
2685    /// other sessions are being provisioned. Each one uses a session on the
2686    /// master it joins, so each is counted: a burst of probes cannot push a
2687    /// master past the server's `MaxSessions`. A probe never opens a master.
2688    #[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        // Without a failure report the recent check is still trusted.
2749        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        // The failed lease gave its slot back: with a cap of one, the next
2782        // lease still lands on the first shard.
2783        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    /// A stand-in for `ssh` shared by two threads that play two daemon
2789    /// processes (the old and new daemon during a restart). An opener takes a
2790    /// while to authenticate; if the socket is bound when it finishes, real
2791    /// `ssh` prints "already exists, disabling multiplexing" and keeps a plain
2792    /// background connection that nothing will ever close.
2793    #[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    /// Two daemon processes that lease on the same instance's sockets at
2829    /// once (J-18) must open one master between them, not one master and one
2830    /// orphaned plain connection.
2831    #[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                    // Each process has its own ledger.
2841                    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                        // Accept TCP but never send an SSH banner. ConnectTimeout
2931                        // must cover the handshake, not just the TCP connect.
2932                        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        // Exercise one open: the outer admission helper separately retries
2960        // pre-authentication timeouts with backoff.
2961        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    /// With sharing switched off, or configured by the user, a lease binds
2997    /// nothing and runs nothing.
2998    #[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    /// A session refused on a live shared connection is told apart from a
3031    /// connection dropped before authentication, even though a session bound
3032    /// with `ProxyCommand=false` prints a closed connection after the refusal.
3033    #[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        // The third acquire has nothing to take until a permit comes back.
3130        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}