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