1use 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";
25pub const INSTANCE_LABEL: &str = "dev.mj.instance";
28pub const INSTANCE_TAG: &str = "dev.mj.instance";
29pub const CONTAINER_WORKSPACE: &str = "/workspace";
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
36pub enum ProvisionStage {
37 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 #[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 #[serde(default)]
108 pub parallel_group: Option<u32>,
109 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
113 pub creates_target: bool,
114 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
120 pub detaches: bool,
121 #[serde(default, skip_serializing_if = "Option::is_none")]
126 pub ssh_destination: Option<String>,
127 #[serde(default, skip_serializing_if = "Option::is_none")]
134 pub ssh_session: Option<SshTarget>,
135 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
138 pub ssh_session_probe: bool,
139 #[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 pub fn parallel_group(mut self, group: u32) -> Self {
181 self.parallel_group = Some(group);
182 self
183 }
184
185 pub fn ssh_destination(mut self, destination: impl Into<String>) -> Self {
188 self.ssh_destination = Some(destination.into());
189 self
190 }
191
192 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 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 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 pub fn creates_target(mut self) -> Self {
242 self.creates_target = true;
243 self
244 }
245
246 pub fn with_sensitive_stdin(mut self, input: Vec<u8>) -> Self {
249 self.sensitive_stdin = Some(SensitiveCommandInput(input));
250 self
251 }
252}
253
254#[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 pub probes: Vec<CommandSpec>,
299 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#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
327#[serde(rename_all = "snake_case")]
328pub enum MountAccess {
329 Ro,
331 Cow,
334 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 pub fn without_overlay(self) -> Self {
358 match self {
359 Self::Cow => Self::Ro,
360 kept => kept,
361 }
362 }
363
364 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
377pub struct ImageUser {
378 pub uid: u32,
379 pub gid: u32,
380}
381
382pub 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#[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
437pub fn overlay_unsupported_filesystem(filesystem: &str) -> Option<&'static str> {
443 let name = filesystem.trim().to_ascii_lowercase();
444 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
457pub 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
486pub 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 fn cancellation_requested(&self) -> bool {
513 false
514 }
515
516 fn stage_started(&self, _stage: ProvisionStage) {}
520
521 fn stage_finished(&self, _stage: ProvisionStage) {}
524
525 fn notify_notice(&self, _notice: &str) {}
528
529 fn reserve_move_destination(&self) {}
532
533 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
550pub 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
573fn 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 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
629fn 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
645pub 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 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 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 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 stream_command_with_stdin(process, command, input, &|| false)
707 }
708}
709
710fn 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 let stdout_reader = PipeCollector::spawn(stdout);
751 let stderr_reader = PipeCollector::spawn(stderr);
752 let process_result = std::thread::scope(|scope| -> Result<_> {
753 let input_writer = scope.spawn(move || -> Result<()> {
757 let mut stdin = stdin;
762 let mut buffer = [0_u8; 64 * 1024];
763 loop {
764 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 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 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
862pub 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 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
937const IO_DRAIN_TIMEOUT: Duration = Duration::from_secs(2);
943
944trait 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 unsafe { libc::poll(&raw mut poll, 1, millis) > 0 }
964 }
965
966 #[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
978struct 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 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 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 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 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 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 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 stream_command_with_stdin(cancellable_command(command), command, input, &|| {
1167 self.is_cancelled()
1168 })
1169 }
1170}
1171
1172#[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#[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 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 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 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
1403pub 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
1421pub 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 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 #[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 #[serde(default)]
1516 pub build_cache: Option<crate::config::TargetBuildCache>,
1517}
1518
1519impl ImagePullPolicy {
1520 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 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 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 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#[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 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#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
1628pub enum RefreshWhen {
1629 WhenAbsent,
1631 Always,
1633}
1634
1635#[derive(Debug, Clone, PartialEq, Eq)]
1638pub struct ImageRefresh {
1639 pub host: ImageHost,
1640 pub image: String,
1641 pub platform: Option<String>,
1642 pub when: RefreshWhen,
1645 pub image_id: CommandSpec,
1648 pub pull: CommandSpec,
1649 pub prune: Option<CommandSpec>,
1652}
1653
1654pub 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 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 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 #[serde(default, skip_serializing_if = "Option::is_none")]
1809 borrowed_from: Option<String>,
1810 },
1811 LocalDocker {
1812 container_id: String,
1813 #[serde(default, skip_serializing_if = "Option::is_none")]
1817 borrowed_from: Option<String>,
1818 },
1819 AppleContainer {
1820 container_id: String,
1821 #[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 #[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 #[serde(default, skip_serializing_if = "Option::is_none")]
1850 borrowed_from: Option<String>,
1851 },
1852 SshDocker {
1853 ssh: SshTarget,
1854 container_id: String,
1855 #[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 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 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#[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
1950pub 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
2006pub 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
2020pub 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
2026pub 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 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 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
2065pub 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
2079pub 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 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 #[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 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 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 #[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 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}