Skip to main content

mj_core/
targets.rs

1//! Declarative execution plans for Hel session targets.
2//!
3//! Plans deliberately contain argv vectors instead of local shell strings.  A
4//! shell is used only at the SSH boundary, where OpenSSH necessarily sends a
5//! command string; every remotely supplied argument is POSIX-quoted there.
6
7use std::collections::{BTreeMap, BTreeSet};
8use std::io::{Read, Write};
9use std::path::{Path, PathBuf};
10use std::process::{Command, Stdio};
11use std::sync::Arc;
12use std::sync::atomic::{AtomicBool, Ordering};
13use std::time::{Duration, Instant};
14
15use anyhow::{Context, Result, bail, ensure};
16use serde::{Deserialize, Serialize};
17use sha2::{Digest, Sha256};
18
19use crate::config::{HarnessHost, HarnessKind, ImagePullPolicy};
20
21pub const SESSION_LABEL: &str = "dev.mj.session";
22pub const MANAGED_LABEL: &str = "dev.mj.managed";
23pub const SESSION_TAG: &str = "dev.mj.session";
24pub const MANAGED_TAG: &str = "dev.mj.managed";
25/// Which Mjolnir instance (named `--instance` or data-directory fingerprint)
26/// created a worker; see `config::instance_identity`.
27pub const INSTANCE_LABEL: &str = "dev.mj.instance";
28pub const INSTANCE_TAG: &str = "dev.mj.instance";
29/// The shared in-container workspace every container session used before
30/// per-session workspaces existed. New sessions record a path under it; see
31/// `container_workspace_root`.
32pub const CONTAINER_WORKSPACE: &str = "/workspace";
33
34/// The launch phase a command belongs to, reported as launch progress.
35#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
36pub enum ProvisionStage {
37    /// Waiting for a background download of this target's image to finish,
38    /// rather than starting a second download of the same image.
39    PullingImage,
40    Provisioning,
41    Booting,
42    Cloning,
43    Syncing,
44    Restoring,
45    Starting,
46    Installing(HarnessKind),
47    Compacting,
48    RecoveryCopy,
49    Verifying,
50    Closing,
51    StoppingTarget,
52    RemovingContainer,
53    RemovingStorage,
54    CleaningCache,
55}
56
57impl ProvisionStage {
58    pub fn label(self) -> String {
59        match self {
60            Self::PullingImage => "Pull image".into(),
61            Self::Provisioning => "Provision".into(),
62            Self::Booting => "Boot".into(),
63            Self::Cloning => "Clone".into(),
64            Self::Syncing => "Sync".into(),
65            Self::Restoring => "Restore".into(),
66            Self::Starting => "Start".into(),
67            Self::Installing(harness) => format!("Installing {}", harness.display_name()),
68            Self::Compacting => "Compact".into(),
69            Self::RecoveryCopy => "Recovery copy".into(),
70            Self::Verifying => "Verify".into(),
71            Self::Closing => "Shut down agent".into(),
72            Self::StoppingTarget => "Stop target".into(),
73            Self::RemovingContainer => "Remove container".into(),
74            Self::RemovingStorage => "Remove container storage".into(),
75            Self::CleaningCache => "Clean cache".into(),
76        }
77    }
78}
79
80#[derive(Clone, PartialEq, Eq)]
81struct SensitiveCommandInput(Vec<u8>);
82
83impl std::fmt::Debug for SensitiveCommandInput {
84    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
85        formatter.write_str("<redacted>")
86    }
87}
88
89#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
90pub struct CommandSpec {
91    pub program: String,
92    pub args: Vec<String>,
93    #[serde(default)]
94    pub env: BTreeMap<String, String>,
95    /// Replace ambient process variables with the explicitly supplied environment.
96    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
97    pub clear_env: bool,
98    #[serde(default, skip_serializing_if = "Option::is_none")]
99    pub cwd: Option<std::path::PathBuf>,
100    pub purpose: String,
101    #[serde(default)]
102    pub stage: Option<ProvisionStage>,
103    /// Commands that share this marker and appear consecutively in a plan's
104    /// command list may run concurrently under
105    /// [`CommandPlan::execute_concurrent`]. Commands without a marker, or
106    /// whose neighbors do not share it, keep running strictly in plan order.
107    #[serde(default)]
108    pub parallel_group: Option<u32>,
109    /// Whether this command brings the session's target into existence. Every
110    /// command after it in a provisioning plan runs against a target that
111    /// already exists, so a later failure owes that target's teardown.
112    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
113    pub creates_target: bool,
114    /// The SSH destination this command opens a connection to, when it does.
115    /// Tagged commands pass through [`SshAdmission`] so the daemon never
116    /// exceeds the remote `sshd`'s `MaxStartups` budget, and a transport
117    /// rejection is retried rather than reported as a command failure.
118    #[serde(default, skip_serializing_if = "Option::is_none")]
119    pub ssh_destination: Option<String>,
120    /// The shared SSH connection this command runs on as one session, when it
121    /// does. Just before spawning, the executor leases a session on one of the
122    /// connection's masters with [`SshSessions::lease`] and puts the options
123    /// that bind the command to that master in front of its arguments (see
124    /// [`CommandSpec::open_ssh_session`]). The arguments stored here never
125    /// contain them.
126    #[serde(default, skip_serializing_if = "Option::is_none")]
127    pub ssh_session: Option<SshTarget>,
128    /// Input that must reach the child without becoming part of its arguments,
129    /// environment, serialized plan, or debug representation.
130    #[serde(skip)]
131    sensitive_stdin: Option<SensitiveCommandInput>,
132}
133
134impl CommandSpec {
135    pub fn new(
136        program: impl Into<String>,
137        args: impl IntoIterator<Item = impl Into<String>>,
138    ) -> Self {
139        Self {
140            program: program.into(),
141            args: args.into_iter().map(Into::into).collect(),
142            env: BTreeMap::new(),
143            clear_env: false,
144            cwd: None,
145            purpose: String::new(),
146            stage: None,
147            parallel_group: None,
148            creates_target: false,
149            ssh_destination: None,
150            ssh_session: None,
151            sensitive_stdin: None,
152        }
153    }
154
155    pub fn purpose(mut self, purpose: impl Into<String>) -> Self {
156        self.purpose = purpose.into();
157        self
158    }
159
160    pub fn stage(mut self, stage: ProvisionStage) -> Self {
161        self.stage = Some(stage);
162        self
163    }
164
165    /// Mark this command as eligible to run concurrently with its
166    /// plan-adjacent siblings that share the same group.
167    pub fn parallel_group(mut self, group: u32) -> Self {
168        self.parallel_group = Some(group);
169        self
170    }
171
172    /// Record that this command opens an `ssh` connection to `destination`,
173    /// which is the host (or `user@host`) argument, never an option value.
174    pub fn ssh_destination(mut self, destination: impl Into<String>) -> Self {
175        self.ssh_destination = Some(destination.into());
176        self
177    }
178
179    /// Record that this command is an `ssh` or `scp` invocation that runs as
180    /// one session on `ssh`'s shared connection. It also opens a connection
181    /// to that destination, for admission, when sharing is off.
182    pub fn ssh_session(mut self, ssh: &SshTarget) -> Self {
183        self.ssh_destination = Some(ssh.destination.clone());
184        self.ssh_session = Some(ssh.clone());
185        self
186    }
187
188    /// Lease the SSH session this command asks for and return the command as
189    /// it must be spawned, with the lease that must outlive the child.
190    ///
191    /// A command without a session request is returned unchanged. `executor`
192    /// runs the master check and, when needed, the opener.
193    pub fn open_ssh_session(&self, executor: &dyn CommandExecutor) -> Result<SessionCommand<'_>> {
194        let Some(ssh) = &self.ssh_session else {
195            return Ok(SessionCommand {
196                command: std::borrow::Cow::Borrowed(self),
197                lease: None,
198            });
199        };
200        let lease = SshSessions::lease(ssh, executor)?;
201        let mut command = self.clone();
202        command.ssh_session = None;
203        command.args = session_command_args(&self.program, &self.args, ssh, &lease);
204        Ok(SessionCommand {
205            command: std::borrow::Cow::Owned(command),
206            lease: Some(lease),
207        })
208    }
209
210    /// Mark this command as the one that creates the session's target.
211    pub fn creates_target(mut self) -> Self {
212        self.creates_target = true;
213        self
214    }
215
216    /// Feed private file content through the shared concurrent pipe handler.
217    /// The bytes stay out of argv, environments, serialization, and Debug.
218    pub fn with_sensitive_stdin(mut self, input: Vec<u8>) -> Self {
219        self.sensitive_stdin = Some(SensitiveCommandInput(input));
220        self
221    }
222}
223
224/// A command ready to spawn, together with the SSH session lease it runs on.
225/// Keep this value alive until the child has exited: dropping it frees the
226/// session slot.
227#[derive(Debug)]
228pub struct SessionCommand<'a> {
229    command: std::borrow::Cow<'a, CommandSpec>,
230    lease: Option<SshSessionLease>,
231}
232
233impl SessionCommand<'_> {
234    pub fn command(&self) -> &CommandSpec {
235        &self.command
236    }
237
238    pub fn lease(&self) -> Option<&SshSessionLease> {
239        self.lease.as_ref()
240    }
241
242    pub fn into_parts(self) -> (CommandSpec, Option<SshSessionLease>) {
243        (self.command.into_owned(), self.lease)
244    }
245}
246
247#[derive(Debug, Clone, PartialEq, Eq)]
248pub struct CommandOutput {
249    pub status: i32,
250    pub stdout: Vec<u8>,
251    pub stderr: Vec<u8>,
252}
253
254#[derive(Debug, Clone, PartialEq, Eq)]
255pub struct SessionResourceUsage {
256    pub cpu_percent: Option<u8>,
257    pub memory_current_bytes: u64,
258    pub memory_limit_bytes: Option<u64>,
259    pub swap_current_bytes: Option<u64>,
260    pub swap_limit_bytes: Option<u64>,
261    pub writable_disk_bytes: Option<u64>,
262}
263
264#[derive(Debug, Clone, PartialEq, Eq)]
265pub struct SessionResourceProbe {
266    pub memory: CommandSpec,
267    pub disk: Option<CommandSpec>,
268}
269
270#[derive(Debug, Clone, Copy, PartialEq, Eq)]
271pub enum DeploymentCapacityKind {
272    Host,
273    AwsFleet,
274}
275
276#[derive(Debug, Clone, PartialEq, Eq)]
277pub struct DeploymentCapacityTarget {
278    pub id: String,
279    pub host: String,
280    pub target_ids: Vec<String>,
281    pub kind: DeploymentCapacityKind,
282    pub local: bool,
283    /// Alternative commands for a host, or one command per live AWS instance.
284    pub probes: Vec<CommandSpec>,
285    /// Prevents a partial AWS fleet sample when one live instance cannot be probed yet.
286    pub probe_error: Option<String>,
287}
288
289#[derive(Debug, Clone, PartialEq, Eq)]
290pub struct DeploymentCapacityUsage {
291    pub cpu_percent: Option<u8>,
292    pub memory_used_bytes: u64,
293    pub memory_total_bytes: u64,
294    pub logical_cores: u64,
295    pub disk_total_bytes: Option<u64>,
296}
297
298/// An additional directory made available to one session.
299///
300/// Containers use isolated mounts. Remote targets may instead receive a
301/// controller-packed snapshot at the destination while retaining this shared
302/// persisted shape.
303#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
304#[serde(from = "AdditionalMountRepr", into = "AdditionalMountRepr")]
305pub struct AdditionalMount {
306    pub source: PathBuf,
307    pub destination: PathBuf,
308    pub access: MountAccess,
309}
310
311/// How a container sees an attached directory.
312#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
313#[serde(rename_all = "snake_case")]
314pub enum MountAccess {
315    /// The container can read the directory but not change it.
316    Ro,
317    /// The container writes to a private copy-on-write overlay; the host
318    /// directory never changes.
319    Cow,
320    /// The container writes straight through to the host directory.
321    Rw,
322}
323
324impl MountAccess {
325    pub const ALL: [Self; 3] = [Self::Ro, Self::Cow, Self::Rw];
326
327    pub fn label(self) -> &'static str {
328        match self {
329            Self::Ro => "ro",
330            Self::Cow => "cow",
331            Self::Rw => "rw",
332        }
333    }
334
335    /// The mode this attachment takes where the host filesystem cannot carry
336    /// a copy-on-write overlay. Read-only is the substitute, because the
337    /// alternative would let the container write through to the host
338    /// directory the user asked to keep unchanged.
339    ///
340    /// This is the one rule for an unavailable overlay: the wizard offers the
341    /// same modes [`Self::offered`] lists, and the runtime downgrades a stored
342    /// mount the same way.
343    pub fn without_overlay(self) -> Self {
344        match self {
345            Self::Cow => Self::Ro,
346            kept => kept,
347        }
348    }
349
350    /// The modes an attachment may be given, given whether the host filesystem
351    /// can carry the copy-on-write overlay.
352    pub fn offered(overlay_available: bool) -> Vec<Self> {
353        Self::ALL
354            .into_iter()
355            .filter(|access| overlay_available || access.without_overlay() == *access)
356            .collect()
357    }
358}
359
360/// The numeric identity of a container image's configured user, read from the
361/// image on the host that runs it.
362#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
363pub struct ImageUser {
364    pub uid: u32,
365    pub gid: u32,
366}
367
368/// Podman's user-namespace option for a session container, or `None` when the
369/// container keeps Podman's default mapping.
370///
371/// A rootless Podman container maps the image's user onto the host user
372/// running the daemon, so a file the container writes into an attached
373/// directory is owned by that host user instead of a subordinate id. The
374/// mapping has to name the image's own ids: plain `keep-id` maps the host user
375/// onto uid 1000 inside the container and demotes an image that runs as root,
376/// which is a change in how the container runs. Without the ids there is no
377/// safe option to pass, so the container runs the way it did before.
378pub fn podman_userns_option(image_user: Option<ImageUser>) -> Option<String> {
379    image_user.map(|ImageUser { uid, gid }| format!("--userns=keep-id:uid={uid},gid={gid}"))
380}
381
382/// The persisted shape. Read-only and copy-on-write mounts keep the original
383/// `read_only` boolean alone, so archives and API payloads that older builds
384/// read stay unchanged; only read-write adds `access`, which older builds
385/// cannot represent.
386#[derive(Serialize, Deserialize)]
387#[serde(deny_unknown_fields)]
388struct AdditionalMountRepr {
389    source: PathBuf,
390    destination: PathBuf,
391    #[serde(default)]
392    read_only: bool,
393    #[serde(default, skip_serializing_if = "Option::is_none")]
394    access: Option<MountAccess>,
395}
396
397impl From<AdditionalMountRepr> for AdditionalMount {
398    fn from(repr: AdditionalMountRepr) -> Self {
399        let access = repr.access.unwrap_or(if repr.read_only {
400            MountAccess::Ro
401        } else {
402            MountAccess::Cow
403        });
404        Self {
405            source: repr.source,
406            destination: repr.destination,
407            access,
408        }
409    }
410}
411
412impl From<AdditionalMount> for AdditionalMountRepr {
413    fn from(mount: AdditionalMount) -> Self {
414        Self {
415            source: mount.source,
416            destination: mount.destination,
417            read_only: mount.access == MountAccess::Ro,
418            access: (mount.access == MountAccess::Rw).then_some(MountAccess::Rw),
419        }
420    }
421}
422
423/// Why a filesystem cannot host a container target's copy-on-write overlay,
424/// or `None` when it can. Unknown types are allowed: the overlay is the better
425/// mount and only a filesystem known to break it is downgraded.
426///
427/// The names are those `stat -f -c %T` reports, matched case-insensitively.
428pub fn overlay_unsupported_filesystem(filesystem: &str) -> Option<&'static str> {
429    let name = filesystem.trim().to_ascii_lowercase();
430    // FUSE reports the backing driver as `fuse.sshfs`, `fuse.s3fs`, and so on.
431    if name == "fuse" || name == "fuseblk" || name.starts_with("fuse.") {
432        return Some("FUSE filesystem");
433    }
434    match name.as_str() {
435        "nfs" | "nfs4" | "cifs" | "smb2" | "smb3" | "9p" | "v9fs" | "virtiofs" | "ceph"
436        | "lustre" | "afs" | "glusterfs" | "ocfs2" | "gfs" | "gfs2" => Some("network filesystem"),
437        "msdos" | "vfat" | "fat" | "exfat" | "ntfs" | "ntfs3" => Some("no POSIX metadata"),
438        "overlayfs" => Some("overlay stacking limit"),
439        _ => None,
440    }
441}
442
443/// Container destinations cannot use the controller or login user's home.
444pub fn validate_mount_destination(path: &Path) -> Result<()> {
445    ensure!(
446        path.is_absolute()
447            && !path
448                .components()
449                .any(|part| part == std::path::Component::ParentDir),
450        "additional mount destination must be a safe absolute container path; ~ is not supported"
451    );
452    Ok(())
453}
454
455pub fn validate_additional_mounts(mounts: &[AdditionalMount]) -> Result<()> {
456    let mut destinations = BTreeSet::new();
457    for mount in mounts {
458        if !mount.source.is_absolute() || mount.source.as_os_str().is_empty() {
459            bail!("additional mount source must be an absolute directory path");
460        }
461        validate_mount_destination(&mount.destination)?;
462        if !destinations.insert(mount.destination.clone()) {
463            bail!(
464                "additional mount destination {:?} is configured more than once",
465                mount.destination
466            );
467        }
468    }
469    Ok(())
470}
471
472/// Choose the editable default destination for an additional host directory.
473pub fn default_mount_destination(source: &Path, existing: &[AdditionalMount]) -> PathBuf {
474    let basename = source
475        .file_name()
476        .filter(|name| !name.is_empty())
477        .unwrap_or_else(|| std::ffi::OsStr::new("mount"));
478    let base = PathBuf::from("/mnt").join(basename);
479    if !existing.iter().any(|mount| mount.destination == base) {
480        return base;
481    }
482    for number in 2.. {
483        let candidate =
484            PathBuf::from("/mnt").join(format!("{}-{number}", basename.to_string_lossy()));
485        if !existing.iter().any(|mount| mount.destination == candidate) {
486            return candidate;
487        }
488    }
489    unreachable!("a finite mount list always has an unused numbered destination")
490}
491
492pub trait CommandExecutor {
493    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput>;
494
495    /// Whether the operation supervising this executor has requested
496    /// cancellation. Test executors and ordinary process execution are not
497    /// cancellable unless they opt in.
498    fn cancellation_requested(&self) -> bool {
499        false
500    }
501
502    /// Report entry into a lifecycle stage. Callers that cover more than one
503    /// command should hold a [`ProvisionStageGuard`] for the whole operation
504    /// so concurrent stages remain visible between subprocesses.
505    fn stage_started(&self, _stage: ProvisionStage) {}
506
507    /// Report exit from a lifecycle stage previously passed to
508    /// [`Self::stage_started`].
509    fn stage_finished(&self, _stage: ProvisionStage) {}
510
511    /// Report a decision an operation made on the user's behalf. This is not a
512    /// failure: the work continues, and the user is told what changed.
513    fn notify_notice(&self, _notice: &str) {}
514
515    fn execute_with_stdin(
516        &self,
517        _command: &CommandSpec,
518        _input: &mut (dyn Read + Send),
519    ) -> Result<CommandOutput> {
520        bail!("this command executor does not support streamed stdin")
521    }
522}
523
524/// A scoped lifecycle-stage report for controller-side work or a sequence of
525/// commands. Dropping the guard reports completion even when the work returns
526/// early with an error.
527pub struct ProvisionStageGuard<'a, E: CommandExecutor + ?Sized> {
528    executor: &'a E,
529    stage: ProvisionStage,
530}
531
532impl<'a, E: CommandExecutor + ?Sized> ProvisionStageGuard<'a, E> {
533    pub fn new(executor: &'a E, stage: ProvisionStage) -> Self {
534        executor.stage_started(stage);
535        Self { executor, stage }
536    }
537}
538
539impl<E: CommandExecutor + ?Sized> Drop for ProvisionStageGuard<'_, E> {
540    fn drop(&mut self) {
541        self.executor.stage_finished(self.stage);
542    }
543}
544
545pub struct ProcessExecutor;
546
547/// Run one `ssh` invocation under process-wide admission control, retrying it
548/// when the server turned the connection away before authentication.
549///
550/// A permit is held only while a child is actually running and is released
551/// between attempts, because a waiting retry occupies no connection slot. A
552/// transport rejection means the remote command never started, so re-running
553/// the whole invocation cannot repeat a side effect.
554///
555/// A command that asks for a shared-connection session leases it on each
556/// attempt, before taking its permit: opening a master takes a permit of its
557/// own, and waiting for that while holding one could exhaust the gate. A
558/// session turned away by the transport usually means its master died, so
559/// the lease is invalidated and the retry checks and reopens the master.
560///
561/// Commands that are not tagged with a destination run untouched.
562fn with_ssh_admission(
563    command: &CommandSpec,
564    executor: &dyn CommandExecutor,
565    is_cancelled: &dyn Fn() -> bool,
566    mut run: impl FnMut(&CommandSpec) -> Result<CommandOutput>,
567) -> Result<CommandOutput> {
568    let Some(destination) = command.ssh_destination.as_deref() else {
569        return run(command);
570    };
571    for attempt in 1..=SSH_RETRY_ATTEMPTS {
572        let session = command.open_ssh_session(executor)?;
573        let output = {
574            let _permit = SshAdmission::acquire(destination);
575            run(session.command())?
576        };
577        let rejected =
578            is_transport_rejection(output.status, &String::from_utf8_lossy(&output.stderr));
579        if rejected && let Some(lease) = session.lease() {
580            lease.invalidate();
581        }
582        drop(session);
583        if attempt == SSH_RETRY_ATTEMPTS || !rejected {
584            return Ok(output);
585        }
586        let delay = ssh_retry_delay(attempt);
587        tracing::warn!(
588            destination,
589            purpose = command.purpose.as_str(),
590            attempt,
591            attempts = SSH_RETRY_ATTEMPTS,
592            delay_ms = delay.as_millis() as u64,
593            stderr = String::from_utf8_lossy(&output.stderr).trim(),
594            "ssh was refused by the server before authentication; retrying"
595        );
596        if !sleep_unless_cancelled(delay, is_cancelled) {
597            bail!("operation cancelled while {}", command.purpose);
598        }
599    }
600    unreachable!("the final attempt always returns");
601}
602
603/// Wait out `delay`, giving up early if the supervising operation is
604/// cancelled. Returns whether the wait completed.
605fn sleep_unless_cancelled(delay: Duration, is_cancelled: &dyn Fn() -> bool) -> bool {
606    let deadline = Instant::now() + delay;
607    loop {
608        if is_cancelled() {
609            return false;
610        }
611        let remaining = deadline.saturating_duration_since(Instant::now());
612        if remaining.is_zero() {
613            return true;
614        }
615        std::thread::sleep(remaining.min(Duration::from_millis(50)));
616    }
617}
618
619/// One debug line per finished target command, so a slow launch or resume
620/// phase can be attributed from logs instead of re-profiled by hand.
621pub fn trace_command_duration(command: &CommandSpec, started: Instant, status: i32) {
622    tracing::debug!(
623        purpose = command.purpose.as_str(),
624        program = command.program.as_str(),
625        status,
626        elapsed_ms = started.elapsed().as_millis() as u64,
627        "target command finished"
628    );
629}
630
631impl ProcessExecutor {
632    /// One attempt, with no admission or retry of its own.
633    fn run_once(&self, command: &CommandSpec) -> Result<CommandOutput> {
634        if let Some(input) = &command.sensitive_stdin {
635            let mut input = std::io::Cursor::new(input.0.as_slice());
636            // Owned bytes, so each attempt gets its own reader.
637            return stream_command_with_stdin(
638                configured_command(command),
639                command,
640                &mut input,
641                &|| false,
642            );
643        }
644        let started = Instant::now();
645        let output = configured_command(command)
646            .stdin(Stdio::null())
647            .output()
648            .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
649        let status = output.status.code().unwrap_or(-1);
650        trace_command_duration(command, started, status);
651        Ok(CommandOutput {
652            status,
653            stdout: output.stdout,
654            stderr: output.stderr,
655        })
656    }
657}
658
659impl CommandExecutor for ProcessExecutor {
660    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
661        with_ssh_admission(command, self, &|| false, |command| self.run_once(command))
662    }
663
664    fn execute_with_stdin(
665        &self,
666        command: &CommandSpec,
667        input: &mut (dyn Read + Send),
668    ) -> Result<CommandOutput> {
669        // A caller's stream cannot be replayed, so this path takes a session
670        // and a permit but never retries.
671        let session = command.open_ssh_session(self)?;
672        let _permit = command
673            .ssh_destination
674            .as_deref()
675            .map(SshAdmission::acquire);
676        let command = session.command();
677        let process = configured_command(command);
678        // Plain process execution is not cancellable, so the transfer only
679        // ends when the child does.
680        stream_command_with_stdin(process, command, input, &|| false)
681    }
682}
683
684/// Streams `input` into a freshly spawned child and collects its output.
685///
686/// Both executors share this one implementation because the pipe edge cases
687/// below are easy to get subtly wrong in a second copy.
688///
689/// `is_cancelled` reports whether the supervising operation wants the transfer
690/// abandoned; [`ProcessExecutor`] passes a check that is never true, which also
691/// makes the kill path below unreachable for it.
692fn stream_command_with_stdin(
693    mut process: Command,
694    command: &CommandSpec,
695    input: &mut (dyn Read + Send),
696    is_cancelled: &(dyn Fn() -> bool + Sync),
697) -> Result<CommandOutput> {
698    let started = Instant::now();
699    if is_cancelled() {
700        bail!("operation cancelled");
701    }
702    let mut child = process
703        .stdin(Stdio::piped())
704        .stdout(Stdio::piped())
705        .stderr(Stdio::piped())
706        .spawn()
707        .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
708    let stdin = child
709        .stdin
710        .take()
711        .context("streamed command stdin missing")?;
712    let mut stdout = child
713        .stdout
714        .take()
715        .context("streamed command stdout missing")?;
716    let mut stderr = child
717        .stderr
718        .take()
719        .context("streamed command stderr missing")?;
720    // Reader threads keep the child's output pipes drained; a child that fills
721    // one while nobody reads would block instead of exiting.
722    let stdout_reader = std::thread::spawn(move || {
723        let mut bytes = Vec::new();
724        std::io::copy(&mut stdout, &mut bytes).map(|_| bytes)
725    });
726    let stderr_reader = std::thread::spawn(move || {
727        let mut bytes = Vec::new();
728        std::io::copy(&mut stderr, &mut bytes).map(|_| bytes)
729    });
730    let process_result = std::thread::scope(|scope| -> Result<_> {
731        // Pipe writes can block forever when a remote helper stops reading.
732        // Keep the writer off the supervising thread so cancellation can kill
733        // the process group and thereby close the blocked pipe.
734        let input_writer = scope.spawn(move || -> Result<()> {
735            // Owning `stdin` here is what closes the pipe's write end once the
736            // transfer finishes. A child that reads to EOF, such as
737            // `mj worker export-checkpoint --spec -`, never exits while any
738            // copy of the write end is still open.
739            let mut stdin = stdin;
740            let mut buffer = [0_u8; 64 * 1024];
741            loop {
742                // Checking before each chunk makes large checkpoint copies
743                // cooperatively cancellable without changing the executor
744                // interface.
745                if is_cancelled() {
746                    bail!("operation cancelled");
747                }
748                let count = input.read(&mut buffer).context("read command input")?;
749                if count == 0 {
750                    break;
751                }
752                stdin
753                    .write_all(&buffer[..count])
754                    .context("stream command input")?;
755            }
756            stdin.flush().context("flush command input")
757        });
758        let status = loop {
759            if is_cancelled() {
760                terminate_cancellable_child(&mut child);
761                if let Err(error) = input_writer.join() {
762                    tracing::warn!(
763                        purpose = command.purpose.as_str(),
764                        "streamed command input writer panicked while cancelling: {error:?}"
765                    );
766                }
767                bail!("operation cancelled while {}", command.purpose);
768            }
769            match child.try_wait() {
770                Ok(Some(status)) => break status,
771                Ok(None) => std::thread::sleep(Duration::from_millis(25)),
772                Err(error) => {
773                    terminate_cancellable_child(&mut child);
774                    if let Err(join_error) = input_writer.join() {
775                        tracing::warn!(
776                            purpose = command.purpose.as_str(),
777                            "streamed command input writer panicked while waiting: {join_error:?}"
778                        );
779                    }
780                    return Err(error).with_context(|| format!("wait for {}", command.purpose));
781                }
782            }
783        };
784        let input_result = input_writer
785            .join()
786            .map_err(|_| anyhow::anyhow!("streamed command input writer panicked"))?;
787        Ok((status, input_result))
788    });
789    let stdout = stdout_reader
790        .join()
791        .map_err(|_| anyhow::anyhow!("streamed command stdout reader panicked"))??;
792    let stderr = stderr_reader
793        .join()
794        .map_err(|_| anyhow::anyhow!("streamed command stderr reader panicked"))??;
795    let (status, input_result) = process_result?;
796    if status.success() {
797        // A child that exited first explains the failure through its own
798        // status and stderr; the broken pipe that exit caused would only hide
799        // it. A successful child must not hide an input error.
800        input_result?;
801    }
802    let status = status.code().unwrap_or(-1);
803    trace_command_duration(command, started, status);
804    Ok(CommandOutput {
805        status,
806        stdout,
807        stderr,
808    })
809}
810
811#[derive(Clone)]
812pub struct CancellableProcessExecutor {
813    cancelled: Arc<AtomicBool>,
814    deadline: Option<Instant>,
815}
816
817impl CancellableProcessExecutor {
818    pub fn new(cancelled: Arc<AtomicBool>) -> Self {
819        Self {
820            cancelled,
821            deadline: None,
822        }
823    }
824
825    pub fn is_cancelled(&self) -> bool {
826        self.cancelled.load(Ordering::Acquire)
827            || self
828                .deadline
829                .is_some_and(|deadline| Instant::now() >= deadline)
830    }
831
832    pub fn with_timeout(timeout: Duration) -> Self {
833        Self {
834            cancelled: Arc::new(AtomicBool::new(false)),
835            deadline: Some(Instant::now() + timeout),
836        }
837    }
838
839    /// Bounds an existing flag-based executor with a deadline, so a wedged
840    /// child becomes a reported failure instead of running forever.
841    pub fn with_deadline(mut self, timeout: Duration) -> Self {
842        self.deadline = Some(Instant::now() + timeout);
843        self
844    }
845
846    fn check_cancelled(&self) -> Result<()> {
847        if self.is_cancelled() {
848            bail!("operation cancelled");
849        }
850        Ok(())
851    }
852}
853
854fn configured_command(command: &CommandSpec) -> Command {
855    let mut process = Command::new(&command.program);
856    if command.clear_env {
857        process.env_clear();
858    }
859    if let Some(cwd) = &command.cwd {
860        process.current_dir(cwd);
861    }
862    process.args(&command.args).envs(&command.env);
863    process
864}
865
866fn cancellable_command(command: &CommandSpec) -> Command {
867    #[cfg(unix)]
868    let mut process = configured_command(command);
869    #[cfg(not(unix))]
870    let process = configured_command(command);
871    #[cfg(unix)]
872    {
873        use std::os::unix::process::CommandExt as _;
874        process.process_group(0);
875    }
876    process
877}
878
879fn terminate_cancellable_child(child: &mut std::process::Child) {
880    #[cfg(unix)]
881    // The child owns a fresh process group, so descendants such as an SSH or
882    // shell helper cannot keep its output pipes open after cancellation. A
883    // group that is already gone is the wanted outcome, not a failure, so the
884    // shared helper decides what deserves a warning.
885    if let Err(error) = crate::subprocess::signal_process_group(child.id() as i32, libc::SIGKILL) {
886        tracing::warn!(pid = child.id(), %error, "could not terminate cancelled command process group");
887    }
888    #[cfg(not(unix))]
889    if let Err(error) = child.kill() {
890        tracing::warn!(pid = child.id(), %error, "could not terminate cancelled command");
891    }
892    if let Err(error) = child.wait() {
893        tracing::warn!(pid = child.id(), %error, "could not reap cancelled command");
894    }
895}
896
897impl CancellableProcessExecutor {
898    /// One attempt, with no admission or retry of its own.
899    fn run_once(&self, command: &CommandSpec) -> Result<CommandOutput> {
900        if let Some(input) = &command.sensitive_stdin {
901            let mut input = std::io::Cursor::new(input.0.as_slice());
902            // Owned bytes, so each attempt gets its own reader.
903            return stream_command_with_stdin(
904                cancellable_command(command),
905                command,
906                &mut input,
907                &|| self.is_cancelled(),
908            );
909        }
910        let started = Instant::now();
911        self.check_cancelled()?;
912        let mut child = cancellable_command(command)
913            .stdin(Stdio::null())
914            .stdout(Stdio::piped())
915            .stderr(Stdio::piped())
916            .spawn()
917            .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
918        let mut stdout = child.stdout.take().context("command stdout missing")?;
919        let mut stderr = child.stderr.take().context("command stderr missing")?;
920        let stdout_reader = std::thread::spawn(move || {
921            let mut bytes = Vec::new();
922            std::io::copy(&mut stdout, &mut bytes).map(|_| bytes)
923        });
924        let stderr_reader = std::thread::spawn(move || {
925            let mut bytes = Vec::new();
926            std::io::copy(&mut stderr, &mut bytes).map(|_| bytes)
927        });
928        let mut status = None;
929        let status = loop {
930            if self.is_cancelled() {
931                terminate_cancellable_child(&mut child);
932                for (stream, reader) in [("stdout", stdout_reader), ("stderr", stderr_reader)] {
933                    match reader.join() {
934                        Ok(Ok(_)) => {}
935                        Ok(Err(error)) => {
936                            tracing::warn!(stream, %error, "cancelled command reader failed")
937                        }
938                        Err(_) => tracing::warn!(stream, "cancelled command reader panicked"),
939                    }
940                }
941                bail!("operation cancelled while {}", command.purpose);
942            }
943            if status.is_none() {
944                status = child
945                    .try_wait()
946                    .with_context(|| format!("wait for {}", command.purpose))?;
947            }
948            // Descendants can retain these pipes after the shell exits. Keep
949            // enforcing the deadline until both readers have actually finished.
950            if let Some(status) = status
951                && stdout_reader.is_finished()
952                && stderr_reader.is_finished()
953            {
954                break status;
955            }
956            std::thread::sleep(Duration::from_millis(25));
957        };
958        let stdout = stdout_reader
959            .join()
960            .map_err(|_| anyhow::anyhow!("command stdout reader panicked"))??;
961        let stderr = stderr_reader
962            .join()
963            .map_err(|_| anyhow::anyhow!("command stderr reader panicked"))??;
964        let status = status.code().unwrap_or(-1);
965        trace_command_duration(command, started, status);
966        Ok(CommandOutput {
967            status,
968            stdout,
969            stderr,
970        })
971    }
972}
973
974impl CommandExecutor for CancellableProcessExecutor {
975    fn cancellation_requested(&self) -> bool {
976        self.is_cancelled()
977    }
978
979    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
980        with_ssh_admission(command, self, &|| self.is_cancelled(), |command| {
981            self.run_once(command)
982        })
983    }
984
985    fn execute_with_stdin(
986        &self,
987        command: &CommandSpec,
988        input: &mut (dyn Read + Send),
989    ) -> Result<CommandOutput> {
990        // A caller's stream cannot be replayed, so this path takes a session
991        // and a permit but never retries.
992        let session = command.open_ssh_session(self)?;
993        let _permit = command
994            .ssh_destination
995            .as_deref()
996            .map(SshAdmission::acquire);
997        let command = session.command();
998        // The child runs in its own process group so cancellation can kill the
999        // whole group, which is what releases a writer blocked on a full pipe.
1000        stream_command_with_stdin(cancellable_command(command), command, input, &|| {
1001            self.is_cancelled()
1002        })
1003    }
1004}
1005
1006/// A command supervised by [`BoundedProcessExecutor`] did not finish before
1007/// its deadline. Callers downcast to this to tell a hung probe apart from a
1008/// command that could not be started at all.
1009#[derive(Debug, Clone, PartialEq, Eq)]
1010pub struct CommandTimedOut {
1011    pub program: String,
1012    pub purpose: String,
1013    pub timeout: Duration,
1014}
1015
1016impl std::fmt::Display for CommandTimedOut {
1017    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1018        write!(
1019            formatter,
1020            "`{}` did not answer within {} seconds while trying to {}",
1021            self.program,
1022            self.timeout.as_secs(),
1023            self.purpose
1024        )
1025    }
1026}
1027
1028impl std::error::Error for CommandTimedOut {}
1029
1030/// Runs every command with its own deadline.
1031///
1032/// [`CancellableProcessExecutor::with_timeout`] bounds a whole operation from a
1033/// single shared deadline, which suits one provisioning run. Prerequisite
1034/// probes are different: each one is expected to answer quickly, and a wedged
1035/// socket or blackholed network must not stall the probes that follow it. A
1036/// timeout here names the probe that hung, so the caller can report it the same
1037/// way it reports any other probe failure.
1038#[derive(Debug, Clone, Copy)]
1039pub struct BoundedProcessExecutor {
1040    timeout: Duration,
1041}
1042
1043impl BoundedProcessExecutor {
1044    pub const fn new(timeout: Duration) -> Self {
1045        Self { timeout }
1046    }
1047}
1048
1049impl CommandExecutor for BoundedProcessExecutor {
1050    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1051        let executor = CancellableProcessExecutor::with_timeout(self.timeout);
1052        executor.execute(command).map_err(|error| {
1053            if executor.is_cancelled() {
1054                anyhow::Error::new(CommandTimedOut {
1055                    program: command.program.clone(),
1056                    purpose: command.purpose.clone(),
1057                    timeout: self.timeout,
1058                })
1059            } else {
1060                error
1061            }
1062        })
1063    }
1064
1065    fn execute_with_stdin(
1066        &self,
1067        command: &CommandSpec,
1068        input: &mut (dyn Read + Send),
1069    ) -> Result<CommandOutput> {
1070        CancellableProcessExecutor::with_timeout(self.timeout).execute_with_stdin(command, input)
1071    }
1072}
1073
1074#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1075pub struct CommandPlan {
1076    pub description: String,
1077    pub commands: Vec<CommandSpec>,
1078}
1079
1080impl CommandPlan {
1081    /// Supply one container environment value without placing it in the
1082    /// Podman/SSH argument vector. The target launcher reads the value from
1083    /// stdin, exports it, and asks the container engine to inherit it by name.
1084    pub fn provide_target_environment_secret(
1085        &mut self,
1086        target: &TargetTemplate,
1087        name: &str,
1088        value: &str,
1089    ) -> Result<()> {
1090        ensure!(
1091            !name.is_empty()
1092                && name.bytes().enumerate().all(|(index, byte)| byte == b'_'
1093                    || byte.is_ascii_alphabetic()
1094                    || (index > 0 && byte.is_ascii_digit())),
1095            "invalid secret environment variable name"
1096        );
1097        ensure!(
1098            !value.as_bytes().contains(&b'\n') && !value.as_bytes().contains(&b'\r'),
1099            "secret environment value cannot contain a newline"
1100        );
1101        let command = self
1102            .commands
1103            .iter_mut()
1104            .find(|command| command.creates_target)
1105            .context("provisioning plan has no target creation command")?;
1106        let read_and_export = format!("IFS= read -r {name} || exit 1; export {name};");
1107        match target {
1108            TargetTemplate::LocalPodman(_)
1109            | TargetTemplate::LocalDocker(_)
1110            | TargetTemplate::AppleContainer(_) => {
1111                let program = std::mem::replace(&mut command.program, "sh".to_owned());
1112                let args = std::mem::take(&mut command.args);
1113                command.args = vec![
1114                    "-c".to_owned(),
1115                    format!("{read_and_export} exec \"$@\""),
1116                    "mj-secret-env".to_owned(),
1117                    program,
1118                ];
1119                command.args.extend(args);
1120            }
1121            TargetTemplate::SshPodman { .. } | TargetTemplate::SshDocker { .. } => {
1122                let remote = command
1123                    .args
1124                    .last_mut()
1125                    .context("remote container command has no SSH command argument")?;
1126                *remote = format!("{read_and_export} exec {remote}");
1127            }
1128            TargetTemplate::LocalBare
1129            | TargetTemplate::AwsEc2(_)
1130            | TargetTemplate::SshBare { .. } => {
1131                bail!("target does not support inherited container environment")
1132            }
1133        }
1134        let mut input = value.as_bytes().to_vec();
1135        input.push(b'\n');
1136        command.sensitive_stdin = Some(SensitiveCommandInput(input));
1137        Ok(())
1138    }
1139
1140    pub fn execute(&self, executor: &impl CommandExecutor) -> Result<Vec<CommandOutput>> {
1141        let mut outputs = Vec::with_capacity(self.commands.len());
1142        for command in &self.commands {
1143            let output = executor.execute(command)?;
1144            if output.status != 0 {
1145                bail!(
1146                    "{} failed with status {}: {}",
1147                    command.purpose,
1148                    output.status,
1149                    String::from_utf8_lossy(&output.stderr)
1150                );
1151            }
1152            outputs.push(output);
1153        }
1154        Ok(outputs)
1155    }
1156
1157    /// Execute the plan the same way [`Self::execute`] does, except that
1158    /// commands sharing a [`CommandSpec::parallel_group`] marker and
1159    /// appearing consecutively in `commands` run concurrently as one batch.
1160    ///
1161    /// A batch starts only once every earlier command has succeeded, and a
1162    /// batch that fails reports the first failure in plan order regardless
1163    /// of which command finished first — the same fail-fast contract
1164    /// [`Self::execute`] provides between individual commands. This method
1165    /// requires a `Sync` executor because a batch shares it across threads;
1166    /// [`Self::execute`] keeps working with non-`Sync` executors such as
1167    /// test fakes built on `RefCell`.
1168    pub fn execute_concurrent(
1169        &self,
1170        executor: &(impl CommandExecutor + Sync),
1171    ) -> Result<Vec<CommandOutput>> {
1172        let mut outputs = Vec::with_capacity(self.commands.len());
1173        let mut index = 0;
1174        while index < self.commands.len() {
1175            let group = self.commands[index].parallel_group;
1176            let mut end = index + 1;
1177            if group.is_some() {
1178                while end < self.commands.len() && self.commands[end].parallel_group == group {
1179                    end += 1;
1180                }
1181            }
1182            let batch = &self.commands[index..end];
1183            if let [command] = batch {
1184                outputs.push(checked_command_output(command, executor.execute(command)?)?);
1185            } else {
1186                let results: Vec<Result<CommandOutput>> = std::thread::scope(|scope| {
1187                    let handles: Vec<_> = batch
1188                        .iter()
1189                        .map(|command| scope.spawn(|| executor.execute(command)))
1190                        .collect();
1191                    handles
1192                        .into_iter()
1193                        .map(|handle| match handle.join() {
1194                            Ok(result) => result,
1195                            Err(panic) => Err(anyhow::anyhow!(
1196                                "concurrent command thread panicked: {}",
1197                                command_thread_panic_message(panic.as_ref())
1198                            )),
1199                        })
1200                        .collect()
1201                });
1202                for (command, result) in batch.iter().zip(results) {
1203                    outputs.push(checked_command_output(command, result?)?);
1204                }
1205            }
1206            index = end;
1207        }
1208        Ok(outputs)
1209    }
1210
1211    /// Split the plan around the command that creates the session's target:
1212    /// the commands through that one, then the commands that run against a
1213    /// target which already exists.
1214    ///
1215    /// A plan that creates nothing — an existing project directory, say —
1216    /// splits into nothing, so a caller never arms a teardown for a target it
1217    /// did not bring into existence.
1218    pub fn split_at_target_creation(&self) -> Option<(Self, Self)> {
1219        let created = self
1220            .commands
1221            .iter()
1222            .position(|command| command.creates_target)?;
1223        let (creation, remainder) = self.commands.split_at(created + 1);
1224        Some((
1225            Self {
1226                description: self.description.clone(),
1227                commands: creation.to_vec(),
1228            },
1229            Self {
1230                description: self.description.clone(),
1231                commands: remainder.to_vec(),
1232            },
1233        ))
1234    }
1235}
1236
1237/// Fail the same way [`CommandPlan::execute`] does for a non-zero exit
1238/// status; kept as a shared helper so [`CommandPlan::execute_concurrent`]
1239/// reports identical error text.
1240pub fn checked_command_output(
1241    command: &CommandSpec,
1242    output: CommandOutput,
1243) -> Result<CommandOutput> {
1244    if output.status != 0 {
1245        bail!(
1246            "{} failed with status {}: {}",
1247            command.purpose,
1248            output.status,
1249            String::from_utf8_lossy(&output.stderr)
1250        );
1251    }
1252    Ok(output)
1253}
1254
1255/// Describe a spawned command thread's panic payload for error context.
1256pub fn command_thread_panic_message(payload: &(dyn std::any::Any + Send)) -> String {
1257    if let Some(message) = payload.downcast_ref::<&str>() {
1258        (*message).to_owned()
1259    } else if let Some(message) = payload.downcast_ref::<String>() {
1260        message.clone()
1261    } else {
1262        "non-string panic payload".to_owned()
1263    }
1264}
1265
1266#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1267pub struct RepositorySpec {
1268    /// Network clone URL. Managed workspaces require a configured remote.
1269    pub url: Option<String>,
1270    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1271    pub push_urls: Vec<String>,
1272    pub destination: String,
1273    pub git_ref: Option<String>,
1274    /// Read-only bare repository mounted into the target for Git object reuse.
1275    /// A missing or unusable reference is only an optimization miss.
1276    #[serde(default, skip_serializing_if = "Option::is_none")]
1277    pub reference: Option<String>,
1278}
1279
1280#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1281pub struct ProjectBundleSpec {
1282    pub primary: String,
1283    pub repositories: Vec<RepositorySpec>,
1284}
1285
1286impl ProjectBundleSpec {
1287    pub fn validate(&self) -> Result<()> {
1288        validate_relative_path(&self.primary)?;
1289        if self.repositories.is_empty() {
1290            bail!("a project bundle must contain at least one repository");
1291        }
1292        let mut destinations = std::collections::BTreeSet::new();
1293        for repository in &self.repositories {
1294            validate_relative_path(&repository.destination)?;
1295            ensure!(
1296                repository
1297                    .url
1298                    .as_deref()
1299                    .is_some_and(|url| !url.trim().is_empty() && !url.starts_with('-')),
1300                "isolated repositories require a network Git remote; configure a remote or use a raw local session"
1301            );
1302            crate::remote_git::validate_network_url(
1303                repository.url.as_deref().expect("checked above"),
1304            )?;
1305            for push_url in &repository.push_urls {
1306                crate::remote_git::validate_network_url(push_url)?;
1307            }
1308            ensure!(
1309                repository.git_ref.is_none(),
1310                "git_ref is no longer supported; remove it to start from the remote's default branch"
1311            );
1312            if !destinations.insert(&repository.destination) {
1313                bail!(
1314                    "duplicate repository destination {}",
1315                    repository.destination
1316                );
1317            }
1318        }
1319        if !destinations.contains(&self.primary) {
1320            bail!("primary repository is not present in the bundle");
1321        }
1322        Ok(())
1323    }
1324}
1325
1326#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1327#[serde(tag = "kind", rename_all = "snake_case")]
1328pub enum PodmanWorkspaceStorage {
1329    PodmanVolume,
1330    HostHelper {
1331        root: String,
1332        helper: Vec<String>,
1333    },
1334    #[default]
1335    ContainerLayer,
1336}
1337
1338#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1339pub struct ContainerTemplate {
1340    pub image: String,
1341    #[serde(default)]
1342    pub pull_policy: ImagePullPolicy,
1343    #[serde(default)]
1344    pub extra_run_args: Vec<String>,
1345    #[serde(default)]
1346    pub workspace_storage: PodmanWorkspaceStorage,
1347    /// Per-target mbx build cache overrides, carried from the user
1348    /// configuration so cache resolution can read them off a runtime target.
1349    #[serde(default)]
1350    pub build_cache: Option<crate::config::TargetBuildCache>,
1351}
1352
1353impl ImagePullPolicy {
1354    /// How fresh this target wants its image, with `Auto` read from the image
1355    /// reference. This is the freshness the background refresher acts on.
1356    pub fn resolve(self, image: &str) -> Self {
1357        if self != Self::Auto {
1358            return self;
1359        }
1360        if image_is_digest_pinned(image) {
1361            Self::Missing
1362        } else if image_is_remote(image) && image_uses_latest_tag(image) {
1363            Self::Newer
1364        } else {
1365            Self::Missing
1366        }
1367    }
1368
1369    /// How fresh a launch insists on being. `Auto` never pulls here: the daemon
1370    /// refreshes remote `:latest` images on its own schedule, so a session
1371    /// starts from the cached image instead of blocking a launch on a
1372    /// multi-gigabyte download. An explicit policy still means what it says.
1373    pub fn at_launch(self, image: &str) -> Self {
1374        if self == Self::Auto {
1375            Self::Missing
1376        } else {
1377            self.resolve(image)
1378        }
1379    }
1380
1381    /// What this policy does to `image`, in the words a settings screen can
1382    /// show. `Auto` is derived rather than an alias, so it names the launch
1383    /// behavior and the background refresh it implies, from the same image
1384    /// reading `resolve` uses.
1385    pub fn describe(self, image: &str) -> &'static str {
1386        match self {
1387            Self::Always => "Pull every launch",
1388            Self::Newer => "Pull when the registry is newer",
1389            Self::Missing => "Pull only if missing",
1390            Self::Never => "Never pull",
1391            Self::Auto => match self.resolve(image) {
1392                Self::Newer => "Pull if missing at launch; refresh :latest in background",
1393                _ => "Pull if missing",
1394            },
1395        }
1396    }
1397
1398    /// Podman's spelling of an already-resolved policy.
1399    pub fn podman_value(self) -> &'static str {
1400        match self {
1401            Self::Always => "always",
1402            Self::Newer => "newer",
1403            Self::Missing => "missing",
1404            Self::Never => "never",
1405            Self::Auto => unreachable!("auto pull policy must resolve"),
1406        }
1407    }
1408}
1409
1410/// Where a background image refresh runs.
1411///
1412/// The SSH form wraps commands the way provisioning does rather than the way
1413/// the preflight probes do: a pull runs for minutes, and the probes' two-second
1414/// keepalive would drop the connection underneath it.
1415#[derive(Debug, Clone, PartialEq, Eq)]
1416pub enum ImageHost {
1417    LocalPodman,
1418    LocalDocker,
1419    AppleContainer,
1420    SshPodman(SshTarget),
1421    SshDocker(SshTarget),
1422}
1423
1424impl ImageHost {
1425    pub const fn engine(&self) -> &'static str {
1426        match self {
1427            Self::LocalPodman | Self::SshPodman(_) => "podman",
1428            Self::LocalDocker | Self::SshDocker(_) => "docker",
1429            Self::AppleContainer => "container",
1430        }
1431    }
1432
1433    /// How this host is named in a log line.
1434    pub fn label(&self) -> String {
1435        match self {
1436            Self::LocalPodman => "local podman".to_owned(),
1437            Self::LocalDocker => "local docker".to_owned(),
1438            Self::AppleContainer => "apple container".to_owned(),
1439            Self::SshPodman(ssh) => format!("podman on {}", ssh.destination),
1440            Self::SshDocker(ssh) => format!("docker on {}", ssh.destination),
1441        }
1442    }
1443
1444    fn command(&self, args: Vec<String>, purpose: String) -> CommandSpec {
1445        match self {
1446            Self::LocalPodman | Self::LocalDocker | Self::AppleContainer => {
1447                CommandSpec::new(args[0].clone(), args[1..].iter().cloned())
1448            }
1449            Self::SshPodman(ssh) | Self::SshDocker(ssh) => ssh_command_owned(ssh, args),
1450        }
1451        .purpose(purpose)
1452    }
1453}
1454
1455/// When a background refresh downloads an image.
1456///
1457/// A host that already has the image is the common case, and most targets
1458/// only need the copy to exist. Ordering matters: merging two targets that
1459/// share an image takes the larger value, so one demanding target upgrades
1460/// the pair.
1461#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
1462pub enum RefreshWhen {
1463    /// Pull only when the host has no copy (missing, or auto on a fixed tag).
1464    WhenAbsent,
1465    /// Pull every tick (always, or newer / auto on a moving tag).
1466    Always,
1467}
1468
1469/// The commands that keep one host's copy of one image current, away from any
1470/// session launch. They run in this order, and only for that host.
1471#[derive(Debug, Clone, PartialEq, Eq)]
1472pub struct ImageRefresh {
1473    pub host: ImageHost,
1474    pub image: String,
1475    pub platform: Option<String>,
1476    /// Whether this image is downloaded only when the host lacks it, or on
1477    /// every refresh.
1478    pub when: RefreshWhen,
1479    /// Reads the cached image id, so a pull that changed nothing stays quiet.
1480    /// Run before and after the pull.
1481    pub image_id: CommandSpec,
1482    pub pull: CommandSpec,
1483    /// Dangling images only. Both engines keep an image a container still uses.
1484    /// `None` for Apple's `container` engine, which has no verified prune form.
1485    pub prune: Option<CommandSpec>,
1486}
1487
1488/// The background refresh for one configured container target, or `None` when
1489/// the target asked never to download the image in the background.
1490///
1491/// `never` is the only opt-out. Every other policy at least wants the host to
1492/// have a copy, so the daemon downloads a missing image once; `always` and
1493/// `newer` keep the hourly pull they always had.
1494pub fn image_refresh(
1495    host: ImageHost,
1496    image: &str,
1497    platform: Option<&str>,
1498    pull_policy: ImagePullPolicy,
1499) -> Option<ImageRefresh> {
1500    let when = match pull_policy.resolve(image) {
1501        ImagePullPolicy::Always | ImagePullPolicy::Newer => RefreshWhen::Always,
1502        ImagePullPolicy::Missing => RefreshWhen::WhenAbsent,
1503        ImagePullPolicy::Never => return None,
1504        ImagePullPolicy::Auto => unreachable!("auto pull policy must resolve"),
1505    };
1506    let engine = host.engine();
1507    // Apple's `container` CLI has no `--format` on `image inspect`: the exit
1508    // status is what says the image is present, and the printed description
1509    // stands in for the id when comparing before and after a pull.
1510    let apple = matches!(host, ImageHost::AppleContainer);
1511    let mut image_id_args = vec![engine.to_owned(), "image".to_owned(), "inspect".to_owned()];
1512    if !apple {
1513        image_id_args.push("--format".to_owned());
1514        image_id_args.push("{{.Id}}".to_owned());
1515    }
1516    image_id_args.push(image.to_owned());
1517    let image_id = host.command(
1518        image_id_args,
1519        format!("read the cached id of container image {image}"),
1520    );
1521    let mut pull_args = vec![engine.to_owned()];
1522    if apple {
1523        pull_args.push("image".to_owned());
1524    }
1525    pull_args.push("pull".to_owned());
1526    // Apple's engine runs only native images, so it takes no platform.
1527    if let Some(platform) = platform.filter(|_| !apple) {
1528        pull_args.push(format!("--platform={platform}"));
1529    }
1530    pull_args.push(image.to_owned());
1531    let pull = host.command(pull_args, format!("refresh container image {image}"));
1532    let prune = (!apple).then(|| {
1533        host.command(
1534            vec![
1535                engine.to_owned(),
1536                "image".to_owned(),
1537                "prune".to_owned(),
1538                "-f".to_owned(),
1539            ],
1540            "remove dangling container images".to_owned(),
1541        )
1542    });
1543    Some(ImageRefresh {
1544        host,
1545        image: image.to_owned(),
1546        platform: platform.map(str::to_owned),
1547        when,
1548        image_id,
1549        pull,
1550        prune,
1551    })
1552}
1553
1554fn image_is_digest_pinned(image: &str) -> bool {
1555    image
1556        .rsplit_once('@')
1557        .is_some_and(|(_, digest)| !digest.is_empty())
1558}
1559
1560fn image_is_remote(image: &str) -> bool {
1561    !image.starts_with("localhost/") && !image.starts_with("local/")
1562}
1563
1564fn image_uses_latest_tag(image: &str) -> bool {
1565    let name = image.split_once('@').map_or(image, |(name, _)| name);
1566    let final_component = name.rsplit('/').next().unwrap_or(name);
1567    !final_component.contains(':') || final_component.ends_with(":latest")
1568}
1569
1570#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1571pub struct SshTarget {
1572    pub destination: String,
1573    #[serde(default)]
1574    pub ssh_args: Vec<String>,
1575}
1576
1577#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1578pub struct AwsTemplate {
1579    pub profile: String,
1580    pub region: String,
1581    pub launch_template: String,
1582    pub launch_template_version: Option<String>,
1583    pub instance_type: Option<String>,
1584    pub ssh: SshTarget,
1585}
1586
1587#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1588#[serde(tag = "kind", rename_all = "snake_case")]
1589pub enum TargetTemplate {
1590    LocalBare,
1591    LocalPodman(ContainerTemplate),
1592    LocalDocker(ContainerTemplate),
1593    AppleContainer(ContainerTemplate),
1594    AwsEc2(AwsTemplate),
1595    SshBare {
1596        ssh: SshTarget,
1597        #[serde(default = "default_ssh_prefix")]
1598        workspace_prefix: String,
1599    },
1600    SshPodman {
1601        ssh: SshTarget,
1602        container: ContainerTemplate,
1603    },
1604    SshDocker {
1605        ssh: SshTarget,
1606        container: ContainerTemplate,
1607    },
1608}
1609
1610fn default_ssh_prefix() -> String {
1611    ".local/share/hel/workspaces".to_owned()
1612}
1613
1614#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1615#[serde(tag = "kind", rename_all = "snake_case")]
1616pub enum PodmanWorkspaceLocator {
1617    #[default]
1618    ContainerLayer,
1619    Volume {
1620        name: String,
1621    },
1622    HostPath {
1623        path: String,
1624        helper: Vec<String>,
1625        resource: String,
1626    },
1627}
1628
1629#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1630#[serde(tag = "kind", rename_all = "snake_case")]
1631pub enum TargetLocator {
1632    LocalBare {
1633        worker_root: String,
1634    },
1635    LocalPodman {
1636        container_id: String,
1637        #[serde(default)]
1638        workspace_storage: PodmanWorkspaceLocator,
1639        /// The session that owns the container when this locator is a
1640        /// sub-agent child borrowing its parent's container; `None` when the
1641        /// session owns the container itself.
1642        #[serde(default, skip_serializing_if = "Option::is_none")]
1643        borrowed_from: Option<String>,
1644    },
1645    LocalDocker {
1646        container_id: String,
1647        /// The session that owns the container when this locator is a
1648        /// sub-agent child borrowing its parent's container; `None` when the
1649        /// session owns the container itself.
1650        #[serde(default, skip_serializing_if = "Option::is_none")]
1651        borrowed_from: Option<String>,
1652    },
1653    AppleContainer {
1654        container_id: String,
1655        /// The session that owns the container when this locator is a
1656        /// sub-agent child borrowing its parent's container; `None` when the
1657        /// session owns the container itself.
1658        #[serde(default, skip_serializing_if = "Option::is_none")]
1659        borrowed_from: Option<String>,
1660    },
1661    AwsEc2 {
1662        profile: String,
1663        region: String,
1664        instance_id: String,
1665        ssh: SshTarget,
1666        workspace: String,
1667    },
1668    SshBare {
1669        ssh: SshTarget,
1670        workspace: String,
1671        /// A borrowed-target child has its own worker identity in the parent's workspace.
1672        #[serde(default, skip_serializing_if = "Option::is_none")]
1673        worker_id: Option<String>,
1674    },
1675    SshPodman {
1676        ssh: SshTarget,
1677        container_id: String,
1678        #[serde(default)]
1679        workspace_storage: PodmanWorkspaceLocator,
1680        /// The session that owns the container when this locator is a
1681        /// sub-agent child borrowing its parent's container; `None` when the
1682        /// session owns the container itself.
1683        #[serde(default, skip_serializing_if = "Option::is_none")]
1684        borrowed_from: Option<String>,
1685    },
1686    SshDocker {
1687        ssh: SshTarget,
1688        container_id: String,
1689        /// The session that owns the container when this locator is a
1690        /// sub-agent child borrowing its parent's container; `None` when the
1691        /// session owns the container itself.
1692        #[serde(default, skip_serializing_if = "Option::is_none")]
1693        borrowed_from: Option<String>,
1694    },
1695}
1696
1697impl TargetTemplate {
1698    pub const fn container_engine(&self) -> Option<&'static str> {
1699        match self {
1700            Self::LocalPodman(_) | Self::SshPodman { .. } => Some("podman"),
1701            Self::LocalDocker(_) | Self::SshDocker { .. } => Some("docker"),
1702            Self::AppleContainer(_) => Some("container"),
1703            _ => None,
1704        }
1705    }
1706
1707    /// The host that downloads this target's image, with the container
1708    /// settings that name the image. `None` for targets that run no image.
1709    pub fn image_host(&self) -> Option<(ImageHost, &ContainerTemplate)> {
1710        match self {
1711            Self::LocalPodman(container) => Some((ImageHost::LocalPodman, container)),
1712            Self::LocalDocker(container) => Some((ImageHost::LocalDocker, container)),
1713            Self::AppleContainer(container) => Some((ImageHost::AppleContainer, container)),
1714            Self::SshPodman { ssh, container } => {
1715                Some((ImageHost::SshPodman(ssh.clone()), container))
1716            }
1717            Self::SshDocker { ssh, container } => {
1718                Some((ImageHost::SshDocker(ssh.clone()), container))
1719            }
1720            Self::LocalBare | Self::AwsEc2(_) | Self::SshBare { .. } => None,
1721        }
1722    }
1723}
1724
1725impl TargetLocator {
1726    /// The operating system the harness will run on. Only a bare localhost
1727    /// target runs it on this machine; every other locator is a Linux
1728    /// container or a Linux host reached over SSH.
1729    pub const fn harness_host(&self) -> HarnessHost {
1730        match self {
1731            Self::LocalBare { .. } => HarnessHost::current(),
1732            _ => HarnessHost::Other,
1733        }
1734    }
1735
1736    /// The target kind spelling shared with [`crate::config::TargetTemplate`].
1737    pub const fn kind_name(&self) -> &'static str {
1738        match self {
1739            Self::LocalBare { .. } => "local-bare",
1740            Self::LocalPodman { .. } => "local-podman",
1741            Self::LocalDocker { .. } => "local-docker",
1742            Self::AppleContainer { .. } => "apple-container",
1743            Self::AwsEc2 { .. } => "aws-ec2",
1744            Self::SshBare { .. } => "ssh-bare",
1745            Self::SshPodman { .. } => "ssh-podman",
1746            Self::SshDocker { .. } => "ssh-docker",
1747        }
1748    }
1749
1750    pub const fn container_engine(&self) -> Option<&'static str> {
1751        match self {
1752            Self::LocalPodman { .. } | Self::SshPodman { .. } => Some("podman"),
1753            Self::LocalDocker { .. } | Self::SshDocker { .. } => Some("docker"),
1754            Self::AppleContainer { .. } => Some("container"),
1755            _ => None,
1756        }
1757    }
1758}
1759
1760/// Commands and identity needed to bring a stopped managed target back online.
1761/// Only runtimes whose stopped resources retain their durable files provide
1762/// one; callers leave every other target kind alone.
1763#[derive(Debug, Clone, PartialEq, Eq)]
1764pub struct TargetRecoveryPlan {
1765    pub exists: CommandSpec,
1766    pub inspect: CommandSpec,
1767    pub start: CommandSpec,
1768    pub session_id: String,
1769}
1770
1771#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1772pub enum TargetRecoveryOutcome {
1773    NotRequired,
1774    Missing,
1775    AlreadyRunning,
1776    Started,
1777}
1778
1779pub fn resource_name(session_id: &str) -> Result<String> {
1780    validate_session_id(session_id)?;
1781    let readable: String = session_id
1782        .chars()
1783        .filter(|character| character.is_ascii_alphanumeric())
1784        .take(12)
1785        .map(|character| character.to_ascii_lowercase())
1786        .collect();
1787    let digest = Sha256::digest(session_id.as_bytes());
1788    Ok(format!(
1789        "mj-{readable}-{:02x}{:02x}{:02x}",
1790        digest[0], digest[1], digest[2]
1791    ))
1792}
1793
1794pub fn podman_workspace_locator(
1795    template: &ContainerTemplate,
1796    session_id: &str,
1797) -> Result<PodmanWorkspaceLocator> {
1798    let resource = format!("{}-workspace", resource_name(session_id)?);
1799    match &template.workspace_storage {
1800        PodmanWorkspaceStorage::PodmanVolume => {
1801            Ok(PodmanWorkspaceLocator::Volume { name: resource })
1802        }
1803        PodmanWorkspaceStorage::HostHelper { root, helper } => {
1804            let root = Path::new(root);
1805            ensure!(
1806                root.is_absolute(),
1807                "Podman workspace storage root must be absolute"
1808            );
1809            ensure!(
1810                !helper.is_empty() && helper.iter().all(|argument| !argument.is_empty()),
1811                "Podman workspace storage helper must contain non-empty arguments"
1812            );
1813            Ok(PodmanWorkspaceLocator::HostPath {
1814                path: root.join(&resource).to_string_lossy().into_owned(),
1815                helper: helper.clone(),
1816                resource,
1817            })
1818        }
1819        PodmanWorkspaceStorage::ContainerLayer => Ok(PodmanWorkspaceLocator::ContainerLayer),
1820    }
1821}
1822
1823/// The in-container workspace root a session's repositories live under.
1824///
1825/// `recorded` is the session record's `container_workspace`. A session whose
1826/// container predates per-session workspaces has none and keeps the legacy
1827/// shared `/workspace`; every session created since records
1828/// `/workspace/<session id>`, so a build cache shared by every container on a
1829/// host never sees two checkouts of one project at the same absolute path.
1830pub fn container_workspace_root(recorded: Option<&Path>) -> String {
1831    recorded.map_or_else(
1832        || CONTAINER_WORKSPACE.to_owned(),
1833        |path| path.to_string_lossy().into_owned(),
1834    )
1835}
1836
1837/// The per-session container workspace recorded for a session created now.
1838pub fn new_container_workspace(session_id: &str) -> Result<PathBuf> {
1839    validate_session_id(session_id)?;
1840    Ok(Path::new(CONTAINER_WORKSPACE).join(session_id))
1841}
1842
1843/// The workspace directory an EC2 session owns in its login home.
1844///
1845/// EC2 instances are provisioned per session, so the path is decided by the
1846/// session id alone and is the same whether it is read from a target template
1847/// or rebuilt from a stored locator.
1848pub fn aws_workspace(session_id: &str) -> String {
1849    format!(".local/share/hel/workspaces/{session_id}")
1850}
1851
1852pub fn workspace_for(template: &TargetTemplate, session_id: &str) -> Result<String> {
1853    validate_session_id(session_id)?;
1854    match template {
1855        TargetTemplate::LocalBare => bail!("local bare projects use their selected directory"),
1856        // A container workspace is per session and recorded on the session
1857        // record, so it is read with `container_workspace_root` instead.
1858        TargetTemplate::LocalPodman(_)
1859        | TargetTemplate::LocalDocker(_)
1860        | TargetTemplate::AppleContainer(_)
1861        | TargetTemplate::SshPodman { .. }
1862        | TargetTemplate::SshDocker { .. } => {
1863            bail!("container targets use the session's recorded container workspace")
1864        }
1865        TargetTemplate::AwsEc2(_) => Ok(aws_workspace(session_id)),
1866        TargetTemplate::SshBare {
1867            workspace_prefix, ..
1868        } => {
1869            validate_workspace_prefix(workspace_prefix)?;
1870            // Interpret a leading "~/" as home-relative. Remote commands are
1871            // single-quoted, so a literal tilde would name a directory called
1872            // "~"; a relative path resolves against the login home for ssh
1873            // and scp alike.
1874            let prefix = workspace_prefix
1875                .strip_prefix("~/")
1876                .unwrap_or(workspace_prefix);
1877            Ok(format!("{}/{session_id}", prefix.trim_end_matches('/')))
1878        }
1879    }
1880}
1881
1882/// Wrap an argv vector for execution at a provisioned session target.
1883pub fn command_on_locator(
1884    locator: &TargetLocator,
1885    session_id: &str,
1886    args: Vec<String>,
1887    purpose: impl Into<String>,
1888) -> Result<CommandSpec> {
1889    verify_locator(locator, session_id)?;
1890    if args.is_empty() {
1891        bail!("target command must not be empty");
1892    }
1893    Ok(locator_command(locator, args).purpose(purpose))
1894}
1895
1896/// Wrap a non-empty argv vector for execution at a target, without checking
1897/// that the locator belongs to a particular session. Every per-target command
1898/// is built here, so each target kind is spelled out once.
1899pub fn locator_command(locator: &TargetLocator, args: Vec<String>) -> CommandSpec {
1900    match locator {
1901        TargetLocator::LocalBare { .. } => {
1902            let mut args = args.into_iter();
1903            let program = args.next().expect("target command must not be empty");
1904            CommandSpec::new(program, args)
1905        }
1906        TargetLocator::LocalPodman { container_id, .. }
1907        | TargetLocator::LocalDocker { container_id, .. }
1908        | TargetLocator::AppleContainer { container_id, .. } => container_exec(
1909            locator.container_engine().expect("local container"),
1910            container_id,
1911            args,
1912        ),
1913        TargetLocator::AwsEc2 { ssh, .. } | TargetLocator::SshBare { ssh, .. } => {
1914            ssh_command_owned(ssh, args)
1915        }
1916        TargetLocator::SshPodman {
1917            ssh, container_id, ..
1918        }
1919        | TargetLocator::SshDocker {
1920            ssh, container_id, ..
1921        } => {
1922            let mut remote = vec![
1923                locator
1924                    .container_engine()
1925                    .expect("remote container")
1926                    .to_owned(),
1927                "exec".to_owned(),
1928                "-i".to_owned(),
1929                container_id.to_owned(),
1930            ];
1931            remote.extend(args);
1932            ssh_command_owned(ssh, remote)
1933        }
1934    }
1935}
1936pub fn worker_root(locator: &TargetLocator, session_id: &str) -> Result<String> {
1937    verify_locator(locator, session_id)?;
1938    Ok(match locator {
1939        TargetLocator::LocalBare { worker_root } => worker_root.clone(),
1940        TargetLocator::LocalPodman { .. }
1941        | TargetLocator::LocalDocker { .. }
1942        | TargetLocator::AppleContainer { .. }
1943        | TargetLocator::SshPodman { .. }
1944        | TargetLocator::SshDocker { .. } => format!("/var/lib/hel/workers/{session_id}"),
1945        TargetLocator::AwsEc2 { .. } => format!(".local/share/hel/workers/{session_id}"),
1946        TargetLocator::SshBare { worker_id, .. } => format!(
1947            ".local/share/hel/workers/{}",
1948            worker_id.as_deref().unwrap_or(session_id)
1949        ),
1950    })
1951}
1952mod convert;
1953pub use convert::{StoredTarget, TargetConversionError, ssh_args_with_identity};
1954
1955mod ssh;
1956pub use ssh::*;
1957
1958pub fn container_exec(
1959    engine: &str,
1960    container_id: &str,
1961    args: impl IntoIterator<Item = impl Into<String>>,
1962) -> CommandSpec {
1963    let mut command_args = vec!["exec".to_owned(), "-i".to_owned(), container_id.to_owned()];
1964    command_args.extend(args.into_iter().map(Into::into));
1965    CommandSpec::new(engine, command_args)
1966}
1967
1968#[cfg(all(test, unix))]
1969mod executor_tests {
1970    use std::fs;
1971
1972    use super::*;
1973
1974    /// A stand-in for `ssh` that is refused by the server on its first call and
1975    /// connects on the next, the way a host at its `MaxStartups` ceiling
1976    /// behaves once the daemon's burst drains.
1977    fn flaky_ssh_script(directory: &Path) -> CommandSpec {
1978        let counter = directory.join("attempts");
1979        let script = format!(
1980            "count=$(cat {counter} 2>/dev/null || echo 0)\n\
1981             echo $((count + 1)) > {counter}\n\
1982             if [ \"$count\" -eq 0 ]; then\n\
1983             echo 'kex_exchange_identification: Connection closed by 10.0.0.1 port 22' >&2\n\
1984             exit 255\n\
1985             fi\n\
1986             echo connected\n",
1987            counter = counter.display()
1988        );
1989        CommandSpec::new("sh", ["-c".to_owned(), script])
1990            .ssh_destination("build@10.0.0.1")
1991            .purpose("run the flaky SSH fixture")
1992    }
1993
1994    fn attempts(directory: &Path) -> u32 {
1995        fs::read_to_string(directory.join("attempts"))
1996            .expect("the fixture records its attempts")
1997            .trim()
1998            .parse()
1999            .expect("attempt count is a number")
2000    }
2001
2002    /// An executor that models `ssh` masters and the sessions on them. The
2003    /// first session finds its master dead (it was killed after the ledger
2004    /// last checked it), which `ssh` reports the way a session guarded with
2005    /// `ProxyCommand=false` does.
2006    #[derive(Default)]
2007    struct MasterKilledOnce {
2008        running: std::cell::RefCell<BTreeSet<String>>,
2009        openers: std::cell::Cell<usize>,
2010        sessions: std::cell::RefCell<Vec<Vec<String>>>,
2011    }
2012
2013    impl CommandExecutor for MasterKilledOnce {
2014        fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2015            let reply = |status: i32, stderr: &str| CommandOutput {
2016                status,
2017                stdout: Vec::new(),
2018                stderr: stderr.as_bytes().to_vec(),
2019            };
2020            if command.ssh_session.is_some() {
2021                // A session. The first one arrives just after its master died.
2022                return with_ssh_admission(command, self, &|| false, |spawned| {
2023                    self.sessions.borrow_mut().push(spawned.args.clone());
2024                    if self.sessions.borrow().len() == 1 {
2025                        self.running.borrow_mut().clear();
2026                        return Ok(reply(255, "Connection closed by UNKNOWN port 65535"));
2027                    }
2028                    Ok(reply(0, ""))
2029                });
2030            }
2031            let socket = command
2032                .args
2033                .iter()
2034                .find_map(|arg| arg.strip_prefix("ControlPath="))
2035                .expect("a master command names its socket")
2036                .to_owned();
2037            if command.args.windows(2).any(|pair| pair == ["-O", "check"]) {
2038                return Ok(reply(
2039                    if self.running.borrow().contains(&socket) {
2040                        0
2041                    } else {
2042                        255
2043                    },
2044                    "",
2045                ));
2046            }
2047            assert!(command.args.contains(&"ControlMaster=yes".to_owned()));
2048            self.openers.set(self.openers.get() + 1);
2049            self.running.borrow_mut().insert(socket);
2050            Ok(reply(0, ""))
2051        }
2052    }
2053
2054    #[test]
2055    fn a_session_whose_master_died_is_retried_on_a_reopened_master() {
2056        let _guard = ssh::SHARING_TEST_LOCK
2057            .lock()
2058            .unwrap_or_else(std::sync::PoisonError::into_inner);
2059        set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2060        let socket_dir = tempfile::tempdir_in("/tmp").expect("short socket directory");
2061        set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2062            socket_dir.path().to_path_buf(),
2063        )));
2064        let ssh = SshTarget {
2065            destination: "master-killed-once-host".to_owned(),
2066            ssh_args: Vec::new(),
2067        };
2068        let executor = MasterKilledOnce::default();
2069        let output = executor.execute(&ssh_command(&ssh, ["true"]));
2070        set_ssh_connection_sharing_for_test(None);
2071        set_ssh_retry_backoff_for_test(None);
2072
2073        assert_eq!(output.expect("the retry succeeds").status, 0);
2074        let sessions = executor.sessions.borrow();
2075        assert_eq!(sessions.len(), 2);
2076        assert_eq!(
2077            executor.openers.get(),
2078            2,
2079            "the retry reopens the master instead of trusting the earlier check"
2080        );
2081        for args in sessions.iter() {
2082            assert_eq!(args[..6][5], "ProxyCommand=false", "{args:?}");
2083        }
2084    }
2085
2086    #[test]
2087    fn a_transport_rejected_ssh_command_is_retried_once_and_then_succeeds() {
2088        set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2089        let directory = tempfile::tempdir().expect("temp dir");
2090        let command = flaky_ssh_script(directory.path());
2091
2092        let output = ProcessExecutor
2093            .execute(&command)
2094            .expect("the retry must reach the successful attempt");
2095
2096        assert_eq!(output.status, 0);
2097        assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "connected");
2098        assert_eq!(attempts(directory.path()), 2);
2099        set_ssh_retry_backoff_for_test(None);
2100    }
2101
2102    #[test]
2103    fn an_untagged_command_is_not_retried_after_the_same_failure() {
2104        set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2105        let directory = tempfile::tempdir().expect("temp dir");
2106        let mut command = flaky_ssh_script(directory.path());
2107        command.ssh_destination = None;
2108
2109        let output = ProcessExecutor.execute(&command).expect("runs once");
2110
2111        assert_eq!(output.status, 255);
2112        assert_eq!(attempts(directory.path()), 1);
2113        set_ssh_retry_backoff_for_test(None);
2114    }
2115
2116    #[test]
2117    fn the_cancellable_executor_also_retries_a_transport_rejection() {
2118        set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2119        let directory = tempfile::tempdir().expect("temp dir");
2120        let command = flaky_ssh_script(directory.path());
2121
2122        let output = CancellableProcessExecutor::new(Arc::new(AtomicBool::new(false)))
2123            .execute(&command)
2124            .expect("the retry must reach the successful attempt");
2125
2126        assert_eq!(output.status, 0);
2127        assert_eq!(attempts(directory.path()), 2);
2128        set_ssh_retry_backoff_for_test(None);
2129    }
2130}