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