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::{HarnessHost, HarnessKind, ImagePullPolicy};
20
21pub const SESSION_LABEL: &str = "dev.mj.session";
22pub const MANAGED_LABEL: &str = "dev.mj.managed";
23pub const SESSION_TAG: &str = "dev.mj.session";
24pub const MANAGED_TAG: &str = "dev.mj.managed";
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 = "Option::is_none")]
119 pub ssh_destination: Option<String>,
120 #[serde(default, skip_serializing_if = "Option::is_none")]
127 pub ssh_session: Option<SshTarget>,
128 #[serde(skip)]
131 sensitive_stdin: Option<SensitiveCommandInput>,
132}
133
134impl CommandSpec {
135 pub fn new(
136 program: impl Into<String>,
137 args: impl IntoIterator<Item = impl Into<String>>,
138 ) -> Self {
139 Self {
140 program: program.into(),
141 args: args.into_iter().map(Into::into).collect(),
142 env: BTreeMap::new(),
143 clear_env: false,
144 cwd: None,
145 purpose: String::new(),
146 stage: None,
147 parallel_group: None,
148 creates_target: false,
149 ssh_destination: None,
150 ssh_session: None,
151 sensitive_stdin: None,
152 }
153 }
154
155 pub fn purpose(mut self, purpose: impl Into<String>) -> Self {
156 self.purpose = purpose.into();
157 self
158 }
159
160 pub fn stage(mut self, stage: ProvisionStage) -> Self {
161 self.stage = Some(stage);
162 self
163 }
164
165 pub fn parallel_group(mut self, group: u32) -> Self {
168 self.parallel_group = Some(group);
169 self
170 }
171
172 pub fn ssh_destination(mut self, destination: impl Into<String>) -> Self {
175 self.ssh_destination = Some(destination.into());
176 self
177 }
178
179 pub fn ssh_session(mut self, ssh: &SshTarget) -> Self {
183 self.ssh_destination = Some(ssh.destination.clone());
184 self.ssh_session = Some(ssh.clone());
185 self
186 }
187
188 pub fn open_ssh_session(&self, executor: &dyn CommandExecutor) -> Result<SessionCommand<'_>> {
194 let Some(ssh) = &self.ssh_session else {
195 return Ok(SessionCommand {
196 command: std::borrow::Cow::Borrowed(self),
197 lease: None,
198 });
199 };
200 let lease = SshSessions::lease(ssh, executor)?;
201 let mut command = self.clone();
202 command.ssh_session = None;
203 command.args = session_command_args(&self.program, &self.args, ssh, &lease);
204 Ok(SessionCommand {
205 command: std::borrow::Cow::Owned(command),
206 lease: Some(lease),
207 })
208 }
209
210 pub fn creates_target(mut self) -> Self {
212 self.creates_target = true;
213 self
214 }
215
216 pub fn with_sensitive_stdin(mut self, input: Vec<u8>) -> Self {
219 self.sensitive_stdin = Some(SensitiveCommandInput(input));
220 self
221 }
222}
223
224#[derive(Debug)]
228pub struct SessionCommand<'a> {
229 command: std::borrow::Cow<'a, CommandSpec>,
230 lease: Option<SshSessionLease>,
231}
232
233impl SessionCommand<'_> {
234 pub fn command(&self) -> &CommandSpec {
235 &self.command
236 }
237
238 pub fn lease(&self) -> Option<&SshSessionLease> {
239 self.lease.as_ref()
240 }
241
242 pub fn into_parts(self) -> (CommandSpec, Option<SshSessionLease>) {
243 (self.command.into_owned(), self.lease)
244 }
245}
246
247#[derive(Debug, Clone, PartialEq, Eq)]
248pub struct CommandOutput {
249 pub status: i32,
250 pub stdout: Vec<u8>,
251 pub stderr: Vec<u8>,
252}
253
254#[derive(Debug, Clone, PartialEq, Eq)]
255pub struct SessionResourceUsage {
256 pub cpu_percent: Option<u8>,
257 pub memory_current_bytes: u64,
258 pub memory_limit_bytes: Option<u64>,
259 pub swap_current_bytes: Option<u64>,
260 pub swap_limit_bytes: Option<u64>,
261 pub writable_disk_bytes: Option<u64>,
262}
263
264#[derive(Debug, Clone, PartialEq, Eq)]
265pub struct SessionResourceProbe {
266 pub memory: CommandSpec,
267 pub disk: Option<CommandSpec>,
268}
269
270#[derive(Debug, Clone, Copy, PartialEq, Eq)]
271pub enum DeploymentCapacityKind {
272 Host,
273 AwsFleet,
274}
275
276#[derive(Debug, Clone, PartialEq, Eq)]
277pub struct DeploymentCapacityTarget {
278 pub id: String,
279 pub host: String,
280 pub target_ids: Vec<String>,
281 pub kind: DeploymentCapacityKind,
282 pub local: bool,
283 pub probes: Vec<CommandSpec>,
285 pub probe_error: Option<String>,
287}
288
289#[derive(Debug, Clone, PartialEq, Eq)]
290pub struct DeploymentCapacityUsage {
291 pub cpu_percent: Option<u8>,
292 pub memory_used_bytes: u64,
293 pub memory_total_bytes: u64,
294 pub logical_cores: u64,
295 pub disk_total_bytes: Option<u64>,
296}
297
298#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
304#[serde(from = "AdditionalMountRepr", into = "AdditionalMountRepr")]
305pub struct AdditionalMount {
306 pub source: PathBuf,
307 pub destination: PathBuf,
308 pub access: MountAccess,
309}
310
311#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
313#[serde(rename_all = "snake_case")]
314pub enum MountAccess {
315 Ro,
317 Cow,
320 Rw,
322}
323
324impl MountAccess {
325 pub const ALL: [Self; 3] = [Self::Ro, Self::Cow, Self::Rw];
326
327 pub fn label(self) -> &'static str {
328 match self {
329 Self::Ro => "ro",
330 Self::Cow => "cow",
331 Self::Rw => "rw",
332 }
333 }
334
335 pub fn without_overlay(self) -> Self {
344 match self {
345 Self::Cow => Self::Ro,
346 kept => kept,
347 }
348 }
349
350 pub fn offered(overlay_available: bool) -> Vec<Self> {
353 Self::ALL
354 .into_iter()
355 .filter(|access| overlay_available || access.without_overlay() == *access)
356 .collect()
357 }
358}
359
360#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
363pub struct ImageUser {
364 pub uid: u32,
365 pub gid: u32,
366}
367
368pub fn podman_userns_option(image_user: Option<ImageUser>) -> Option<String> {
379 image_user.map(|ImageUser { uid, gid }| format!("--userns=keep-id:uid={uid},gid={gid}"))
380}
381
382#[derive(Serialize, Deserialize)]
387#[serde(deny_unknown_fields)]
388struct AdditionalMountRepr {
389 source: PathBuf,
390 destination: PathBuf,
391 #[serde(default)]
392 read_only: bool,
393 #[serde(default, skip_serializing_if = "Option::is_none")]
394 access: Option<MountAccess>,
395}
396
397impl From<AdditionalMountRepr> for AdditionalMount {
398 fn from(repr: AdditionalMountRepr) -> Self {
399 let access = repr.access.unwrap_or(if repr.read_only {
400 MountAccess::Ro
401 } else {
402 MountAccess::Cow
403 });
404 Self {
405 source: repr.source,
406 destination: repr.destination,
407 access,
408 }
409 }
410}
411
412impl From<AdditionalMount> for AdditionalMountRepr {
413 fn from(mount: AdditionalMount) -> Self {
414 Self {
415 source: mount.source,
416 destination: mount.destination,
417 read_only: mount.access == MountAccess::Ro,
418 access: (mount.access == MountAccess::Rw).then_some(MountAccess::Rw),
419 }
420 }
421}
422
423pub fn overlay_unsupported_filesystem(filesystem: &str) -> Option<&'static str> {
429 let name = filesystem.trim().to_ascii_lowercase();
430 if name == "fuse" || name == "fuseblk" || name.starts_with("fuse.") {
432 return Some("FUSE filesystem");
433 }
434 match name.as_str() {
435 "nfs" | "nfs4" | "cifs" | "smb2" | "smb3" | "9p" | "v9fs" | "virtiofs" | "ceph"
436 | "lustre" | "afs" | "glusterfs" | "ocfs2" | "gfs" | "gfs2" => Some("network filesystem"),
437 "msdos" | "vfat" | "fat" | "exfat" | "ntfs" | "ntfs3" => Some("no POSIX metadata"),
438 "overlayfs" => Some("overlay stacking limit"),
439 _ => None,
440 }
441}
442
443pub fn validate_mount_destination(path: &Path) -> Result<()> {
445 ensure!(
446 path.is_absolute()
447 && !path
448 .components()
449 .any(|part| part == std::path::Component::ParentDir),
450 "additional mount destination must be a safe absolute container path; ~ is not supported"
451 );
452 Ok(())
453}
454
455pub fn validate_additional_mounts(mounts: &[AdditionalMount]) -> Result<()> {
456 let mut destinations = BTreeSet::new();
457 for mount in mounts {
458 if !mount.source.is_absolute() || mount.source.as_os_str().is_empty() {
459 bail!("additional mount source must be an absolute directory path");
460 }
461 validate_mount_destination(&mount.destination)?;
462 if !destinations.insert(mount.destination.clone()) {
463 bail!(
464 "additional mount destination {:?} is configured more than once",
465 mount.destination
466 );
467 }
468 }
469 Ok(())
470}
471
472pub fn default_mount_destination(source: &Path, existing: &[AdditionalMount]) -> PathBuf {
474 let basename = source
475 .file_name()
476 .filter(|name| !name.is_empty())
477 .unwrap_or_else(|| std::ffi::OsStr::new("mount"));
478 let base = PathBuf::from("/mnt").join(basename);
479 if !existing.iter().any(|mount| mount.destination == base) {
480 return base;
481 }
482 for number in 2.. {
483 let candidate =
484 PathBuf::from("/mnt").join(format!("{}-{number}", basename.to_string_lossy()));
485 if !existing.iter().any(|mount| mount.destination == candidate) {
486 return candidate;
487 }
488 }
489 unreachable!("a finite mount list always has an unused numbered destination")
490}
491
492pub trait CommandExecutor {
493 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput>;
494
495 fn cancellation_requested(&self) -> bool {
499 false
500 }
501
502 fn stage_started(&self, _stage: ProvisionStage) {}
506
507 fn stage_finished(&self, _stage: ProvisionStage) {}
510
511 fn notify_notice(&self, _notice: &str) {}
514
515 fn execute_with_stdin(
516 &self,
517 _command: &CommandSpec,
518 _input: &mut (dyn Read + Send),
519 ) -> Result<CommandOutput> {
520 bail!("this command executor does not support streamed stdin")
521 }
522}
523
524pub struct ProvisionStageGuard<'a, E: CommandExecutor + ?Sized> {
528 executor: &'a E,
529 stage: ProvisionStage,
530}
531
532impl<'a, E: CommandExecutor + ?Sized> ProvisionStageGuard<'a, E> {
533 pub fn new(executor: &'a E, stage: ProvisionStage) -> Self {
534 executor.stage_started(stage);
535 Self { executor, stage }
536 }
537}
538
539impl<E: CommandExecutor + ?Sized> Drop for ProvisionStageGuard<'_, E> {
540 fn drop(&mut self) {
541 self.executor.stage_finished(self.stage);
542 }
543}
544
545pub struct ProcessExecutor;
546
547fn with_ssh_admission(
563 command: &CommandSpec,
564 executor: &dyn CommandExecutor,
565 is_cancelled: &dyn Fn() -> bool,
566 mut run: impl FnMut(&CommandSpec) -> Result<CommandOutput>,
567) -> Result<CommandOutput> {
568 let Some(destination) = command.ssh_destination.as_deref() else {
569 return run(command);
570 };
571 for attempt in 1..=SSH_RETRY_ATTEMPTS {
572 let session = command.open_ssh_session(executor)?;
573 let output = {
574 let _permit = SshAdmission::acquire(destination);
575 run(session.command())?
576 };
577 let rejected =
578 is_transport_rejection(output.status, &String::from_utf8_lossy(&output.stderr));
579 if rejected && let Some(lease) = session.lease() {
580 lease.invalidate();
581 }
582 drop(session);
583 if attempt == SSH_RETRY_ATTEMPTS || !rejected {
584 return Ok(output);
585 }
586 let delay = ssh_retry_delay(attempt);
587 tracing::warn!(
588 destination,
589 purpose = command.purpose.as_str(),
590 attempt,
591 attempts = SSH_RETRY_ATTEMPTS,
592 delay_ms = delay.as_millis() as u64,
593 stderr = String::from_utf8_lossy(&output.stderr).trim(),
594 "ssh was refused by the server before authentication; retrying"
595 );
596 if !sleep_unless_cancelled(delay, is_cancelled) {
597 bail!("operation cancelled while {}", command.purpose);
598 }
599 }
600 unreachable!("the final attempt always returns");
601}
602
603fn sleep_unless_cancelled(delay: Duration, is_cancelled: &dyn Fn() -> bool) -> bool {
606 let deadline = Instant::now() + delay;
607 loop {
608 if is_cancelled() {
609 return false;
610 }
611 let remaining = deadline.saturating_duration_since(Instant::now());
612 if remaining.is_zero() {
613 return true;
614 }
615 std::thread::sleep(remaining.min(Duration::from_millis(50)));
616 }
617}
618
619pub fn trace_command_duration(command: &CommandSpec, started: Instant, status: i32) {
622 tracing::debug!(
623 purpose = command.purpose.as_str(),
624 program = command.program.as_str(),
625 status,
626 elapsed_ms = started.elapsed().as_millis() as u64,
627 "target command finished"
628 );
629}
630
631impl ProcessExecutor {
632 fn run_once(&self, command: &CommandSpec) -> Result<CommandOutput> {
634 if let Some(input) = &command.sensitive_stdin {
635 let mut input = std::io::Cursor::new(input.0.as_slice());
636 return stream_command_with_stdin(
638 configured_command(command),
639 command,
640 &mut input,
641 &|| false,
642 );
643 }
644 let started = Instant::now();
645 let output = configured_command(command)
646 .stdin(Stdio::null())
647 .output()
648 .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
649 let status = output.status.code().unwrap_or(-1);
650 trace_command_duration(command, started, status);
651 Ok(CommandOutput {
652 status,
653 stdout: output.stdout,
654 stderr: output.stderr,
655 })
656 }
657}
658
659impl CommandExecutor for ProcessExecutor {
660 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
661 with_ssh_admission(command, self, &|| false, |command| self.run_once(command))
662 }
663
664 fn execute_with_stdin(
665 &self,
666 command: &CommandSpec,
667 input: &mut (dyn Read + Send),
668 ) -> Result<CommandOutput> {
669 let session = command.open_ssh_session(self)?;
672 let _permit = command
673 .ssh_destination
674 .as_deref()
675 .map(SshAdmission::acquire);
676 let command = session.command();
677 let process = configured_command(command);
678 stream_command_with_stdin(process, command, input, &|| false)
681 }
682}
683
684fn stream_command_with_stdin(
693 mut process: Command,
694 command: &CommandSpec,
695 input: &mut (dyn Read + Send),
696 is_cancelled: &(dyn Fn() -> bool + Sync),
697) -> Result<CommandOutput> {
698 let started = Instant::now();
699 if is_cancelled() {
700 bail!("operation cancelled");
701 }
702 let mut child = process
703 .stdin(Stdio::piped())
704 .stdout(Stdio::piped())
705 .stderr(Stdio::piped())
706 .spawn()
707 .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
708 let stdin = child
709 .stdin
710 .take()
711 .context("streamed command stdin missing")?;
712 let mut stdout = child
713 .stdout
714 .take()
715 .context("streamed command stdout missing")?;
716 let mut stderr = child
717 .stderr
718 .take()
719 .context("streamed command stderr missing")?;
720 let stdout_reader = std::thread::spawn(move || {
723 let mut bytes = Vec::new();
724 std::io::copy(&mut stdout, &mut bytes).map(|_| bytes)
725 });
726 let stderr_reader = std::thread::spawn(move || {
727 let mut bytes = Vec::new();
728 std::io::copy(&mut stderr, &mut bytes).map(|_| bytes)
729 });
730 let process_result = std::thread::scope(|scope| -> Result<_> {
731 let input_writer = scope.spawn(move || -> Result<()> {
735 let mut stdin = stdin;
740 let mut buffer = [0_u8; 64 * 1024];
741 loop {
742 if is_cancelled() {
746 bail!("operation cancelled");
747 }
748 let count = input.read(&mut buffer).context("read command input")?;
749 if count == 0 {
750 break;
751 }
752 stdin
753 .write_all(&buffer[..count])
754 .context("stream command input")?;
755 }
756 stdin.flush().context("flush command input")
757 });
758 let status = loop {
759 if is_cancelled() {
760 terminate_cancellable_child(&mut child);
761 if let Err(error) = input_writer.join() {
762 tracing::warn!(
763 purpose = command.purpose.as_str(),
764 "streamed command input writer panicked while cancelling: {error:?}"
765 );
766 }
767 bail!("operation cancelled while {}", command.purpose);
768 }
769 match child.try_wait() {
770 Ok(Some(status)) => break status,
771 Ok(None) => std::thread::sleep(Duration::from_millis(25)),
772 Err(error) => {
773 terminate_cancellable_child(&mut child);
774 if let Err(join_error) = input_writer.join() {
775 tracing::warn!(
776 purpose = command.purpose.as_str(),
777 "streamed command input writer panicked while waiting: {join_error:?}"
778 );
779 }
780 return Err(error).with_context(|| format!("wait for {}", command.purpose));
781 }
782 }
783 };
784 let input_result = input_writer
785 .join()
786 .map_err(|_| anyhow::anyhow!("streamed command input writer panicked"))?;
787 Ok((status, input_result))
788 });
789 let stdout = stdout_reader
790 .join()
791 .map_err(|_| anyhow::anyhow!("streamed command stdout reader panicked"))??;
792 let stderr = stderr_reader
793 .join()
794 .map_err(|_| anyhow::anyhow!("streamed command stderr reader panicked"))??;
795 let (status, input_result) = process_result?;
796 if status.success() {
797 input_result?;
801 }
802 let status = status.code().unwrap_or(-1);
803 trace_command_duration(command, started, status);
804 Ok(CommandOutput {
805 status,
806 stdout,
807 stderr,
808 })
809}
810
811#[derive(Clone)]
812pub struct CancellableProcessExecutor {
813 cancelled: Arc<AtomicBool>,
814 deadline: Option<Instant>,
815}
816
817impl CancellableProcessExecutor {
818 pub fn new(cancelled: Arc<AtomicBool>) -> Self {
819 Self {
820 cancelled,
821 deadline: None,
822 }
823 }
824
825 pub fn is_cancelled(&self) -> bool {
826 self.cancelled.load(Ordering::Acquire)
827 || self
828 .deadline
829 .is_some_and(|deadline| Instant::now() >= deadline)
830 }
831
832 pub fn with_timeout(timeout: Duration) -> Self {
833 Self {
834 cancelled: Arc::new(AtomicBool::new(false)),
835 deadline: Some(Instant::now() + timeout),
836 }
837 }
838
839 pub fn with_deadline(mut self, timeout: Duration) -> Self {
842 self.deadline = Some(Instant::now() + timeout);
843 self
844 }
845
846 fn check_cancelled(&self) -> Result<()> {
847 if self.is_cancelled() {
848 bail!("operation cancelled");
849 }
850 Ok(())
851 }
852}
853
854fn configured_command(command: &CommandSpec) -> Command {
855 let mut process = Command::new(&command.program);
856 if command.clear_env {
857 process.env_clear();
858 }
859 if let Some(cwd) = &command.cwd {
860 process.current_dir(cwd);
861 }
862 process.args(&command.args).envs(&command.env);
863 process
864}
865
866fn cancellable_command(command: &CommandSpec) -> Command {
867 #[cfg(unix)]
868 let mut process = configured_command(command);
869 #[cfg(not(unix))]
870 let process = configured_command(command);
871 #[cfg(unix)]
872 {
873 use std::os::unix::process::CommandExt as _;
874 process.process_group(0);
875 }
876 process
877}
878
879fn terminate_cancellable_child(child: &mut std::process::Child) {
880 #[cfg(unix)]
881 if let Err(error) = crate::subprocess::signal_process_group(child.id() as i32, libc::SIGKILL) {
886 tracing::warn!(pid = child.id(), %error, "could not terminate cancelled command process group");
887 }
888 #[cfg(not(unix))]
889 if let Err(error) = child.kill() {
890 tracing::warn!(pid = child.id(), %error, "could not terminate cancelled command");
891 }
892 if let Err(error) = child.wait() {
893 tracing::warn!(pid = child.id(), %error, "could not reap cancelled command");
894 }
895}
896
897impl CancellableProcessExecutor {
898 fn run_once(&self, command: &CommandSpec) -> Result<CommandOutput> {
900 if let Some(input) = &command.sensitive_stdin {
901 let mut input = std::io::Cursor::new(input.0.as_slice());
902 return stream_command_with_stdin(
904 cancellable_command(command),
905 command,
906 &mut input,
907 &|| self.is_cancelled(),
908 );
909 }
910 let started = Instant::now();
911 self.check_cancelled()?;
912 let mut child = cancellable_command(command)
913 .stdin(Stdio::null())
914 .stdout(Stdio::piped())
915 .stderr(Stdio::piped())
916 .spawn()
917 .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
918 let mut stdout = child.stdout.take().context("command stdout missing")?;
919 let mut stderr = child.stderr.take().context("command stderr missing")?;
920 let stdout_reader = std::thread::spawn(move || {
921 let mut bytes = Vec::new();
922 std::io::copy(&mut stdout, &mut bytes).map(|_| bytes)
923 });
924 let stderr_reader = std::thread::spawn(move || {
925 let mut bytes = Vec::new();
926 std::io::copy(&mut stderr, &mut bytes).map(|_| bytes)
927 });
928 let mut status = None;
929 let status = loop {
930 if self.is_cancelled() {
931 terminate_cancellable_child(&mut child);
932 for (stream, reader) in [("stdout", stdout_reader), ("stderr", stderr_reader)] {
933 match reader.join() {
934 Ok(Ok(_)) => {}
935 Ok(Err(error)) => {
936 tracing::warn!(stream, %error, "cancelled command reader failed")
937 }
938 Err(_) => tracing::warn!(stream, "cancelled command reader panicked"),
939 }
940 }
941 bail!("operation cancelled while {}", command.purpose);
942 }
943 if status.is_none() {
944 status = child
945 .try_wait()
946 .with_context(|| format!("wait for {}", command.purpose))?;
947 }
948 if let Some(status) = status
951 && stdout_reader.is_finished()
952 && stderr_reader.is_finished()
953 {
954 break status;
955 }
956 std::thread::sleep(Duration::from_millis(25));
957 };
958 let stdout = stdout_reader
959 .join()
960 .map_err(|_| anyhow::anyhow!("command stdout reader panicked"))??;
961 let stderr = stderr_reader
962 .join()
963 .map_err(|_| anyhow::anyhow!("command stderr reader panicked"))??;
964 let status = status.code().unwrap_or(-1);
965 trace_command_duration(command, started, status);
966 Ok(CommandOutput {
967 status,
968 stdout,
969 stderr,
970 })
971 }
972}
973
974impl CommandExecutor for CancellableProcessExecutor {
975 fn cancellation_requested(&self) -> bool {
976 self.is_cancelled()
977 }
978
979 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
980 with_ssh_admission(command, self, &|| self.is_cancelled(), |command| {
981 self.run_once(command)
982 })
983 }
984
985 fn execute_with_stdin(
986 &self,
987 command: &CommandSpec,
988 input: &mut (dyn Read + Send),
989 ) -> Result<CommandOutput> {
990 let session = command.open_ssh_session(self)?;
993 let _permit = command
994 .ssh_destination
995 .as_deref()
996 .map(SshAdmission::acquire);
997 let command = session.command();
998 stream_command_with_stdin(cancellable_command(command), command, input, &|| {
1001 self.is_cancelled()
1002 })
1003 }
1004}
1005
1006#[derive(Debug, Clone, PartialEq, Eq)]
1010pub struct CommandTimedOut {
1011 pub program: String,
1012 pub purpose: String,
1013 pub timeout: Duration,
1014}
1015
1016impl std::fmt::Display for CommandTimedOut {
1017 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1018 write!(
1019 formatter,
1020 "`{}` did not answer within {} seconds while trying to {}",
1021 self.program,
1022 self.timeout.as_secs(),
1023 self.purpose
1024 )
1025 }
1026}
1027
1028impl std::error::Error for CommandTimedOut {}
1029
1030#[derive(Debug, Clone, Copy)]
1039pub struct BoundedProcessExecutor {
1040 timeout: Duration,
1041}
1042
1043impl BoundedProcessExecutor {
1044 pub const fn new(timeout: Duration) -> Self {
1045 Self { timeout }
1046 }
1047}
1048
1049impl CommandExecutor for BoundedProcessExecutor {
1050 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1051 let executor = CancellableProcessExecutor::with_timeout(self.timeout);
1052 executor.execute(command).map_err(|error| {
1053 if executor.is_cancelled() {
1054 anyhow::Error::new(CommandTimedOut {
1055 program: command.program.clone(),
1056 purpose: command.purpose.clone(),
1057 timeout: self.timeout,
1058 })
1059 } else {
1060 error
1061 }
1062 })
1063 }
1064
1065 fn execute_with_stdin(
1066 &self,
1067 command: &CommandSpec,
1068 input: &mut (dyn Read + Send),
1069 ) -> Result<CommandOutput> {
1070 CancellableProcessExecutor::with_timeout(self.timeout).execute_with_stdin(command, input)
1071 }
1072}
1073
1074#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1075pub struct CommandPlan {
1076 pub description: String,
1077 pub commands: Vec<CommandSpec>,
1078}
1079
1080impl CommandPlan {
1081 pub fn provide_target_environment_secret(
1085 &mut self,
1086 target: &TargetTemplate,
1087 name: &str,
1088 value: &str,
1089 ) -> Result<()> {
1090 ensure!(
1091 !name.is_empty()
1092 && name.bytes().enumerate().all(|(index, byte)| byte == b'_'
1093 || byte.is_ascii_alphabetic()
1094 || (index > 0 && byte.is_ascii_digit())),
1095 "invalid secret environment variable name"
1096 );
1097 ensure!(
1098 !value.as_bytes().contains(&b'\n') && !value.as_bytes().contains(&b'\r'),
1099 "secret environment value cannot contain a newline"
1100 );
1101 let command = self
1102 .commands
1103 .iter_mut()
1104 .find(|command| command.creates_target)
1105 .context("provisioning plan has no target creation command")?;
1106 let read_and_export = format!("IFS= read -r {name} || exit 1; export {name};");
1107 match target {
1108 TargetTemplate::LocalPodman(_)
1109 | TargetTemplate::LocalDocker(_)
1110 | TargetTemplate::AppleContainer(_) => {
1111 let program = std::mem::replace(&mut command.program, "sh".to_owned());
1112 let args = std::mem::take(&mut command.args);
1113 command.args = vec![
1114 "-c".to_owned(),
1115 format!("{read_and_export} exec \"$@\""),
1116 "mj-secret-env".to_owned(),
1117 program,
1118 ];
1119 command.args.extend(args);
1120 }
1121 TargetTemplate::SshPodman { .. } | TargetTemplate::SshDocker { .. } => {
1122 let remote = command
1123 .args
1124 .last_mut()
1125 .context("remote container command has no SSH command argument")?;
1126 *remote = format!("{read_and_export} exec {remote}");
1127 }
1128 TargetTemplate::LocalBare
1129 | TargetTemplate::AwsEc2(_)
1130 | TargetTemplate::SshBare { .. } => {
1131 bail!("target does not support inherited container environment")
1132 }
1133 }
1134 let mut input = value.as_bytes().to_vec();
1135 input.push(b'\n');
1136 command.sensitive_stdin = Some(SensitiveCommandInput(input));
1137 Ok(())
1138 }
1139
1140 pub fn execute(&self, executor: &impl CommandExecutor) -> Result<Vec<CommandOutput>> {
1141 let mut outputs = Vec::with_capacity(self.commands.len());
1142 for command in &self.commands {
1143 let output = executor.execute(command)?;
1144 if output.status != 0 {
1145 bail!(
1146 "{} failed with status {}: {}",
1147 command.purpose,
1148 output.status,
1149 String::from_utf8_lossy(&output.stderr)
1150 );
1151 }
1152 outputs.push(output);
1153 }
1154 Ok(outputs)
1155 }
1156
1157 pub fn execute_concurrent(
1169 &self,
1170 executor: &(impl CommandExecutor + Sync),
1171 ) -> Result<Vec<CommandOutput>> {
1172 let mut outputs = Vec::with_capacity(self.commands.len());
1173 let mut index = 0;
1174 while index < self.commands.len() {
1175 let group = self.commands[index].parallel_group;
1176 let mut end = index + 1;
1177 if group.is_some() {
1178 while end < self.commands.len() && self.commands[end].parallel_group == group {
1179 end += 1;
1180 }
1181 }
1182 let batch = &self.commands[index..end];
1183 if let [command] = batch {
1184 outputs.push(checked_command_output(command, executor.execute(command)?)?);
1185 } else {
1186 let results: Vec<Result<CommandOutput>> = std::thread::scope(|scope| {
1187 let handles: Vec<_> = batch
1188 .iter()
1189 .map(|command| scope.spawn(|| executor.execute(command)))
1190 .collect();
1191 handles
1192 .into_iter()
1193 .map(|handle| match handle.join() {
1194 Ok(result) => result,
1195 Err(panic) => Err(anyhow::anyhow!(
1196 "concurrent command thread panicked: {}",
1197 command_thread_panic_message(panic.as_ref())
1198 )),
1199 })
1200 .collect()
1201 });
1202 for (command, result) in batch.iter().zip(results) {
1203 outputs.push(checked_command_output(command, result?)?);
1204 }
1205 }
1206 index = end;
1207 }
1208 Ok(outputs)
1209 }
1210
1211 pub fn split_at_target_creation(&self) -> Option<(Self, Self)> {
1219 let created = self
1220 .commands
1221 .iter()
1222 .position(|command| command.creates_target)?;
1223 let (creation, remainder) = self.commands.split_at(created + 1);
1224 Some((
1225 Self {
1226 description: self.description.clone(),
1227 commands: creation.to_vec(),
1228 },
1229 Self {
1230 description: self.description.clone(),
1231 commands: remainder.to_vec(),
1232 },
1233 ))
1234 }
1235}
1236
1237pub fn checked_command_output(
1241 command: &CommandSpec,
1242 output: CommandOutput,
1243) -> Result<CommandOutput> {
1244 if output.status != 0 {
1245 bail!(
1246 "{} failed with status {}: {}",
1247 command.purpose,
1248 output.status,
1249 String::from_utf8_lossy(&output.stderr)
1250 );
1251 }
1252 Ok(output)
1253}
1254
1255pub fn command_thread_panic_message(payload: &(dyn std::any::Any + Send)) -> String {
1257 if let Some(message) = payload.downcast_ref::<&str>() {
1258 (*message).to_owned()
1259 } else if let Some(message) = payload.downcast_ref::<String>() {
1260 message.clone()
1261 } else {
1262 "non-string panic payload".to_owned()
1263 }
1264}
1265
1266#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1267pub struct RepositorySpec {
1268 pub url: Option<String>,
1270 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1271 pub push_urls: Vec<String>,
1272 pub destination: String,
1273 pub git_ref: Option<String>,
1274 #[serde(default, skip_serializing_if = "Option::is_none")]
1277 pub reference: Option<String>,
1278}
1279
1280#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1281pub struct ProjectBundleSpec {
1282 pub primary: String,
1283 pub repositories: Vec<RepositorySpec>,
1284}
1285
1286impl ProjectBundleSpec {
1287 pub fn validate(&self) -> Result<()> {
1288 validate_relative_path(&self.primary)?;
1289 if self.repositories.is_empty() {
1290 bail!("a project bundle must contain at least one repository");
1291 }
1292 let mut destinations = std::collections::BTreeSet::new();
1293 for repository in &self.repositories {
1294 validate_relative_path(&repository.destination)?;
1295 ensure!(
1296 repository
1297 .url
1298 .as_deref()
1299 .is_some_and(|url| !url.trim().is_empty() && !url.starts_with('-')),
1300 "isolated repositories require a network Git remote; configure a remote or use a raw local session"
1301 );
1302 crate::remote_git::validate_network_url(
1303 repository.url.as_deref().expect("checked above"),
1304 )?;
1305 for push_url in &repository.push_urls {
1306 crate::remote_git::validate_network_url(push_url)?;
1307 }
1308 ensure!(
1309 repository.git_ref.is_none(),
1310 "git_ref is no longer supported; remove it to start from the remote's default branch"
1311 );
1312 if !destinations.insert(&repository.destination) {
1313 bail!(
1314 "duplicate repository destination {}",
1315 repository.destination
1316 );
1317 }
1318 }
1319 if !destinations.contains(&self.primary) {
1320 bail!("primary repository is not present in the bundle");
1321 }
1322 Ok(())
1323 }
1324}
1325
1326#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1327#[serde(tag = "kind", rename_all = "snake_case")]
1328pub enum PodmanWorkspaceStorage {
1329 PodmanVolume,
1330 HostHelper {
1331 root: String,
1332 helper: Vec<String>,
1333 },
1334 #[default]
1335 ContainerLayer,
1336}
1337
1338#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1339pub struct ContainerTemplate {
1340 pub image: String,
1341 #[serde(default)]
1342 pub pull_policy: ImagePullPolicy,
1343 #[serde(default)]
1344 pub extra_run_args: Vec<String>,
1345 #[serde(default)]
1346 pub workspace_storage: PodmanWorkspaceStorage,
1347 #[serde(default)]
1350 pub build_cache: Option<crate::config::TargetBuildCache>,
1351}
1352
1353impl ImagePullPolicy {
1354 pub fn resolve(self, image: &str) -> Self {
1357 if self != Self::Auto {
1358 return self;
1359 }
1360 if image_is_digest_pinned(image) {
1361 Self::Missing
1362 } else if image_is_remote(image) && image_uses_latest_tag(image) {
1363 Self::Newer
1364 } else {
1365 Self::Missing
1366 }
1367 }
1368
1369 pub fn at_launch(self, image: &str) -> Self {
1374 if self == Self::Auto {
1375 Self::Missing
1376 } else {
1377 self.resolve(image)
1378 }
1379 }
1380
1381 pub fn describe(self, image: &str) -> &'static str {
1386 match self {
1387 Self::Always => "Pull every launch",
1388 Self::Newer => "Pull when the registry is newer",
1389 Self::Missing => "Pull only if missing",
1390 Self::Never => "Never pull",
1391 Self::Auto => match self.resolve(image) {
1392 Self::Newer => "Pull if missing at launch; refresh :latest in background",
1393 _ => "Pull if missing",
1394 },
1395 }
1396 }
1397
1398 pub fn podman_value(self) -> &'static str {
1400 match self {
1401 Self::Always => "always",
1402 Self::Newer => "newer",
1403 Self::Missing => "missing",
1404 Self::Never => "never",
1405 Self::Auto => unreachable!("auto pull policy must resolve"),
1406 }
1407 }
1408}
1409
1410#[derive(Debug, Clone, PartialEq, Eq)]
1416pub enum ImageHost {
1417 LocalPodman,
1418 LocalDocker,
1419 AppleContainer,
1420 SshPodman(SshTarget),
1421 SshDocker(SshTarget),
1422}
1423
1424impl ImageHost {
1425 pub const fn engine(&self) -> &'static str {
1426 match self {
1427 Self::LocalPodman | Self::SshPodman(_) => "podman",
1428 Self::LocalDocker | Self::SshDocker(_) => "docker",
1429 Self::AppleContainer => "container",
1430 }
1431 }
1432
1433 pub fn label(&self) -> String {
1435 match self {
1436 Self::LocalPodman => "local podman".to_owned(),
1437 Self::LocalDocker => "local docker".to_owned(),
1438 Self::AppleContainer => "apple container".to_owned(),
1439 Self::SshPodman(ssh) => format!("podman on {}", ssh.destination),
1440 Self::SshDocker(ssh) => format!("docker on {}", ssh.destination),
1441 }
1442 }
1443
1444 fn command(&self, args: Vec<String>, purpose: String) -> CommandSpec {
1445 match self {
1446 Self::LocalPodman | Self::LocalDocker | Self::AppleContainer => {
1447 CommandSpec::new(args[0].clone(), args[1..].iter().cloned())
1448 }
1449 Self::SshPodman(ssh) | Self::SshDocker(ssh) => ssh_command_owned(ssh, args),
1450 }
1451 .purpose(purpose)
1452 }
1453}
1454
1455#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
1462pub enum RefreshWhen {
1463 WhenAbsent,
1465 Always,
1467}
1468
1469#[derive(Debug, Clone, PartialEq, Eq)]
1472pub struct ImageRefresh {
1473 pub host: ImageHost,
1474 pub image: String,
1475 pub platform: Option<String>,
1476 pub when: RefreshWhen,
1479 pub image_id: CommandSpec,
1482 pub pull: CommandSpec,
1483 pub prune: Option<CommandSpec>,
1486}
1487
1488pub fn image_refresh(
1495 host: ImageHost,
1496 image: &str,
1497 platform: Option<&str>,
1498 pull_policy: ImagePullPolicy,
1499) -> Option<ImageRefresh> {
1500 let when = match pull_policy.resolve(image) {
1501 ImagePullPolicy::Always | ImagePullPolicy::Newer => RefreshWhen::Always,
1502 ImagePullPolicy::Missing => RefreshWhen::WhenAbsent,
1503 ImagePullPolicy::Never => return None,
1504 ImagePullPolicy::Auto => unreachable!("auto pull policy must resolve"),
1505 };
1506 let engine = host.engine();
1507 let apple = matches!(host, ImageHost::AppleContainer);
1511 let mut image_id_args = vec![engine.to_owned(), "image".to_owned(), "inspect".to_owned()];
1512 if !apple {
1513 image_id_args.push("--format".to_owned());
1514 image_id_args.push("{{.Id}}".to_owned());
1515 }
1516 image_id_args.push(image.to_owned());
1517 let image_id = host.command(
1518 image_id_args,
1519 format!("read the cached id of container image {image}"),
1520 );
1521 let mut pull_args = vec![engine.to_owned()];
1522 if apple {
1523 pull_args.push("image".to_owned());
1524 }
1525 pull_args.push("pull".to_owned());
1526 if let Some(platform) = platform.filter(|_| !apple) {
1528 pull_args.push(format!("--platform={platform}"));
1529 }
1530 pull_args.push(image.to_owned());
1531 let pull = host.command(pull_args, format!("refresh container image {image}"));
1532 let prune = (!apple).then(|| {
1533 host.command(
1534 vec![
1535 engine.to_owned(),
1536 "image".to_owned(),
1537 "prune".to_owned(),
1538 "-f".to_owned(),
1539 ],
1540 "remove dangling container images".to_owned(),
1541 )
1542 });
1543 Some(ImageRefresh {
1544 host,
1545 image: image.to_owned(),
1546 platform: platform.map(str::to_owned),
1547 when,
1548 image_id,
1549 pull,
1550 prune,
1551 })
1552}
1553
1554fn image_is_digest_pinned(image: &str) -> bool {
1555 image
1556 .rsplit_once('@')
1557 .is_some_and(|(_, digest)| !digest.is_empty())
1558}
1559
1560fn image_is_remote(image: &str) -> bool {
1561 !image.starts_with("localhost/") && !image.starts_with("local/")
1562}
1563
1564fn image_uses_latest_tag(image: &str) -> bool {
1565 let name = image.split_once('@').map_or(image, |(name, _)| name);
1566 let final_component = name.rsplit('/').next().unwrap_or(name);
1567 !final_component.contains(':') || final_component.ends_with(":latest")
1568}
1569
1570#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1571pub struct SshTarget {
1572 pub destination: String,
1573 #[serde(default)]
1574 pub ssh_args: Vec<String>,
1575}
1576
1577#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1578pub struct AwsTemplate {
1579 pub profile: String,
1580 pub region: String,
1581 pub launch_template: String,
1582 pub launch_template_version: Option<String>,
1583 pub instance_type: Option<String>,
1584 pub ssh: SshTarget,
1585}
1586
1587#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1588#[serde(tag = "kind", rename_all = "snake_case")]
1589pub enum TargetTemplate {
1590 LocalBare,
1591 LocalPodman(ContainerTemplate),
1592 LocalDocker(ContainerTemplate),
1593 AppleContainer(ContainerTemplate),
1594 AwsEc2(AwsTemplate),
1595 SshBare {
1596 ssh: SshTarget,
1597 #[serde(default = "default_ssh_prefix")]
1598 workspace_prefix: String,
1599 },
1600 SshPodman {
1601 ssh: SshTarget,
1602 container: ContainerTemplate,
1603 },
1604 SshDocker {
1605 ssh: SshTarget,
1606 container: ContainerTemplate,
1607 },
1608}
1609
1610fn default_ssh_prefix() -> String {
1611 ".local/share/hel/workspaces".to_owned()
1612}
1613
1614#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1615#[serde(tag = "kind", rename_all = "snake_case")]
1616pub enum PodmanWorkspaceLocator {
1617 #[default]
1618 ContainerLayer,
1619 Volume {
1620 name: String,
1621 },
1622 HostPath {
1623 path: String,
1624 helper: Vec<String>,
1625 resource: String,
1626 },
1627}
1628
1629#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1630#[serde(tag = "kind", rename_all = "snake_case")]
1631pub enum TargetLocator {
1632 LocalBare {
1633 worker_root: String,
1634 },
1635 LocalPodman {
1636 container_id: String,
1637 #[serde(default)]
1638 workspace_storage: PodmanWorkspaceLocator,
1639 #[serde(default, skip_serializing_if = "Option::is_none")]
1643 borrowed_from: Option<String>,
1644 },
1645 LocalDocker {
1646 container_id: String,
1647 #[serde(default, skip_serializing_if = "Option::is_none")]
1651 borrowed_from: Option<String>,
1652 },
1653 AppleContainer {
1654 container_id: String,
1655 #[serde(default, skip_serializing_if = "Option::is_none")]
1659 borrowed_from: Option<String>,
1660 },
1661 AwsEc2 {
1662 profile: String,
1663 region: String,
1664 instance_id: String,
1665 ssh: SshTarget,
1666 workspace: String,
1667 },
1668 SshBare {
1669 ssh: SshTarget,
1670 workspace: String,
1671 #[serde(default, skip_serializing_if = "Option::is_none")]
1673 worker_id: Option<String>,
1674 },
1675 SshPodman {
1676 ssh: SshTarget,
1677 container_id: String,
1678 #[serde(default)]
1679 workspace_storage: PodmanWorkspaceLocator,
1680 #[serde(default, skip_serializing_if = "Option::is_none")]
1684 borrowed_from: Option<String>,
1685 },
1686 SshDocker {
1687 ssh: SshTarget,
1688 container_id: String,
1689 #[serde(default, skip_serializing_if = "Option::is_none")]
1693 borrowed_from: Option<String>,
1694 },
1695}
1696
1697impl TargetTemplate {
1698 pub const fn container_engine(&self) -> Option<&'static str> {
1699 match self {
1700 Self::LocalPodman(_) | Self::SshPodman { .. } => Some("podman"),
1701 Self::LocalDocker(_) | Self::SshDocker { .. } => Some("docker"),
1702 Self::AppleContainer(_) => Some("container"),
1703 _ => None,
1704 }
1705 }
1706
1707 pub fn image_host(&self) -> Option<(ImageHost, &ContainerTemplate)> {
1710 match self {
1711 Self::LocalPodman(container) => Some((ImageHost::LocalPodman, container)),
1712 Self::LocalDocker(container) => Some((ImageHost::LocalDocker, container)),
1713 Self::AppleContainer(container) => Some((ImageHost::AppleContainer, container)),
1714 Self::SshPodman { ssh, container } => {
1715 Some((ImageHost::SshPodman(ssh.clone()), container))
1716 }
1717 Self::SshDocker { ssh, container } => {
1718 Some((ImageHost::SshDocker(ssh.clone()), container))
1719 }
1720 Self::LocalBare | Self::AwsEc2(_) | Self::SshBare { .. } => None,
1721 }
1722 }
1723}
1724
1725impl TargetLocator {
1726 pub const fn harness_host(&self) -> HarnessHost {
1730 match self {
1731 Self::LocalBare { .. } => HarnessHost::current(),
1732 _ => HarnessHost::Other,
1733 }
1734 }
1735
1736 pub const fn kind_name(&self) -> &'static str {
1738 match self {
1739 Self::LocalBare { .. } => "local-bare",
1740 Self::LocalPodman { .. } => "local-podman",
1741 Self::LocalDocker { .. } => "local-docker",
1742 Self::AppleContainer { .. } => "apple-container",
1743 Self::AwsEc2 { .. } => "aws-ec2",
1744 Self::SshBare { .. } => "ssh-bare",
1745 Self::SshPodman { .. } => "ssh-podman",
1746 Self::SshDocker { .. } => "ssh-docker",
1747 }
1748 }
1749
1750 pub const fn container_engine(&self) -> Option<&'static str> {
1751 match self {
1752 Self::LocalPodman { .. } | Self::SshPodman { .. } => Some("podman"),
1753 Self::LocalDocker { .. } | Self::SshDocker { .. } => Some("docker"),
1754 Self::AppleContainer { .. } => Some("container"),
1755 _ => None,
1756 }
1757 }
1758}
1759
1760#[derive(Debug, Clone, PartialEq, Eq)]
1764pub struct TargetRecoveryPlan {
1765 pub exists: CommandSpec,
1766 pub inspect: CommandSpec,
1767 pub start: CommandSpec,
1768 pub session_id: String,
1769}
1770
1771#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1772pub enum TargetRecoveryOutcome {
1773 NotRequired,
1774 Missing,
1775 AlreadyRunning,
1776 Started,
1777}
1778
1779pub fn resource_name(session_id: &str) -> Result<String> {
1780 validate_session_id(session_id)?;
1781 let readable: String = session_id
1782 .chars()
1783 .filter(|character| character.is_ascii_alphanumeric())
1784 .take(12)
1785 .map(|character| character.to_ascii_lowercase())
1786 .collect();
1787 let digest = Sha256::digest(session_id.as_bytes());
1788 Ok(format!(
1789 "mj-{readable}-{:02x}{:02x}{:02x}",
1790 digest[0], digest[1], digest[2]
1791 ))
1792}
1793
1794pub fn podman_workspace_locator(
1795 template: &ContainerTemplate,
1796 session_id: &str,
1797) -> Result<PodmanWorkspaceLocator> {
1798 let resource = format!("{}-workspace", resource_name(session_id)?);
1799 match &template.workspace_storage {
1800 PodmanWorkspaceStorage::PodmanVolume => {
1801 Ok(PodmanWorkspaceLocator::Volume { name: resource })
1802 }
1803 PodmanWorkspaceStorage::HostHelper { root, helper } => {
1804 let root = Path::new(root);
1805 ensure!(
1806 root.is_absolute(),
1807 "Podman workspace storage root must be absolute"
1808 );
1809 ensure!(
1810 !helper.is_empty() && helper.iter().all(|argument| !argument.is_empty()),
1811 "Podman workspace storage helper must contain non-empty arguments"
1812 );
1813 Ok(PodmanWorkspaceLocator::HostPath {
1814 path: root.join(&resource).to_string_lossy().into_owned(),
1815 helper: helper.clone(),
1816 resource,
1817 })
1818 }
1819 PodmanWorkspaceStorage::ContainerLayer => Ok(PodmanWorkspaceLocator::ContainerLayer),
1820 }
1821}
1822
1823pub fn container_workspace_root(recorded: Option<&Path>) -> String {
1831 recorded.map_or_else(
1832 || CONTAINER_WORKSPACE.to_owned(),
1833 |path| path.to_string_lossy().into_owned(),
1834 )
1835}
1836
1837pub fn new_container_workspace(session_id: &str) -> Result<PathBuf> {
1839 validate_session_id(session_id)?;
1840 Ok(Path::new(CONTAINER_WORKSPACE).join(session_id))
1841}
1842
1843pub fn aws_workspace(session_id: &str) -> String {
1849 format!(".local/share/hel/workspaces/{session_id}")
1850}
1851
1852pub fn workspace_for(template: &TargetTemplate, session_id: &str) -> Result<String> {
1853 validate_session_id(session_id)?;
1854 match template {
1855 TargetTemplate::LocalBare => bail!("local bare projects use their selected directory"),
1856 TargetTemplate::LocalPodman(_)
1859 | TargetTemplate::LocalDocker(_)
1860 | TargetTemplate::AppleContainer(_)
1861 | TargetTemplate::SshPodman { .. }
1862 | TargetTemplate::SshDocker { .. } => {
1863 bail!("container targets use the session's recorded container workspace")
1864 }
1865 TargetTemplate::AwsEc2(_) => Ok(aws_workspace(session_id)),
1866 TargetTemplate::SshBare {
1867 workspace_prefix, ..
1868 } => {
1869 validate_workspace_prefix(workspace_prefix)?;
1870 let prefix = workspace_prefix
1875 .strip_prefix("~/")
1876 .unwrap_or(workspace_prefix);
1877 Ok(format!("{}/{session_id}", prefix.trim_end_matches('/')))
1878 }
1879 }
1880}
1881
1882pub fn command_on_locator(
1884 locator: &TargetLocator,
1885 session_id: &str,
1886 args: Vec<String>,
1887 purpose: impl Into<String>,
1888) -> Result<CommandSpec> {
1889 verify_locator(locator, session_id)?;
1890 if args.is_empty() {
1891 bail!("target command must not be empty");
1892 }
1893 Ok(locator_command(locator, args).purpose(purpose))
1894}
1895
1896pub fn locator_command(locator: &TargetLocator, args: Vec<String>) -> CommandSpec {
1900 match locator {
1901 TargetLocator::LocalBare { .. } => {
1902 let mut args = args.into_iter();
1903 let program = args.next().expect("target command must not be empty");
1904 CommandSpec::new(program, args)
1905 }
1906 TargetLocator::LocalPodman { container_id, .. }
1907 | TargetLocator::LocalDocker { container_id, .. }
1908 | TargetLocator::AppleContainer { container_id, .. } => container_exec(
1909 locator.container_engine().expect("local container"),
1910 container_id,
1911 args,
1912 ),
1913 TargetLocator::AwsEc2 { ssh, .. } | TargetLocator::SshBare { ssh, .. } => {
1914 ssh_command_owned(ssh, args)
1915 }
1916 TargetLocator::SshPodman {
1917 ssh, container_id, ..
1918 }
1919 | TargetLocator::SshDocker {
1920 ssh, container_id, ..
1921 } => {
1922 let mut remote = vec![
1923 locator
1924 .container_engine()
1925 .expect("remote container")
1926 .to_owned(),
1927 "exec".to_owned(),
1928 "-i".to_owned(),
1929 container_id.to_owned(),
1930 ];
1931 remote.extend(args);
1932 ssh_command_owned(ssh, remote)
1933 }
1934 }
1935}
1936pub fn worker_root(locator: &TargetLocator, session_id: &str) -> Result<String> {
1937 verify_locator(locator, session_id)?;
1938 Ok(match locator {
1939 TargetLocator::LocalBare { worker_root } => worker_root.clone(),
1940 TargetLocator::LocalPodman { .. }
1941 | TargetLocator::LocalDocker { .. }
1942 | TargetLocator::AppleContainer { .. }
1943 | TargetLocator::SshPodman { .. }
1944 | TargetLocator::SshDocker { .. } => format!("/var/lib/hel/workers/{session_id}"),
1945 TargetLocator::AwsEc2 { .. } => format!(".local/share/hel/workers/{session_id}"),
1946 TargetLocator::SshBare { worker_id, .. } => format!(
1947 ".local/share/hel/workers/{}",
1948 worker_id.as_deref().unwrap_or(session_id)
1949 ),
1950 })
1951}
1952mod convert;
1953pub use convert::{StoredTarget, TargetConversionError, ssh_args_with_identity};
1954
1955mod ssh;
1956pub use ssh::*;
1957
1958pub fn container_exec(
1959 engine: &str,
1960 container_id: &str,
1961 args: impl IntoIterator<Item = impl Into<String>>,
1962) -> CommandSpec {
1963 let mut command_args = vec!["exec".to_owned(), "-i".to_owned(), container_id.to_owned()];
1964 command_args.extend(args.into_iter().map(Into::into));
1965 CommandSpec::new(engine, command_args)
1966}
1967
1968#[cfg(all(test, unix))]
1969mod executor_tests {
1970 use std::fs;
1971
1972 use super::*;
1973
1974 fn flaky_ssh_script(directory: &Path) -> CommandSpec {
1978 let counter = directory.join("attempts");
1979 let script = format!(
1980 "count=$(cat {counter} 2>/dev/null || echo 0)\n\
1981 echo $((count + 1)) > {counter}\n\
1982 if [ \"$count\" -eq 0 ]; then\n\
1983 echo 'kex_exchange_identification: Connection closed by 10.0.0.1 port 22' >&2\n\
1984 exit 255\n\
1985 fi\n\
1986 echo connected\n",
1987 counter = counter.display()
1988 );
1989 CommandSpec::new("sh", ["-c".to_owned(), script])
1990 .ssh_destination("build@10.0.0.1")
1991 .purpose("run the flaky SSH fixture")
1992 }
1993
1994 fn attempts(directory: &Path) -> u32 {
1995 fs::read_to_string(directory.join("attempts"))
1996 .expect("the fixture records its attempts")
1997 .trim()
1998 .parse()
1999 .expect("attempt count is a number")
2000 }
2001
2002 #[derive(Default)]
2007 struct MasterKilledOnce {
2008 running: std::cell::RefCell<BTreeSet<String>>,
2009 openers: std::cell::Cell<usize>,
2010 sessions: std::cell::RefCell<Vec<Vec<String>>>,
2011 }
2012
2013 impl CommandExecutor for MasterKilledOnce {
2014 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2015 let reply = |status: i32, stderr: &str| CommandOutput {
2016 status,
2017 stdout: Vec::new(),
2018 stderr: stderr.as_bytes().to_vec(),
2019 };
2020 if command.ssh_session.is_some() {
2021 return with_ssh_admission(command, self, &|| false, |spawned| {
2023 self.sessions.borrow_mut().push(spawned.args.clone());
2024 if self.sessions.borrow().len() == 1 {
2025 self.running.borrow_mut().clear();
2026 return Ok(reply(255, "Connection closed by UNKNOWN port 65535"));
2027 }
2028 Ok(reply(0, ""))
2029 });
2030 }
2031 let socket = command
2032 .args
2033 .iter()
2034 .find_map(|arg| arg.strip_prefix("ControlPath="))
2035 .expect("a master command names its socket")
2036 .to_owned();
2037 if command.args.windows(2).any(|pair| pair == ["-O", "check"]) {
2038 return Ok(reply(
2039 if self.running.borrow().contains(&socket) {
2040 0
2041 } else {
2042 255
2043 },
2044 "",
2045 ));
2046 }
2047 assert!(command.args.contains(&"ControlMaster=yes".to_owned()));
2048 self.openers.set(self.openers.get() + 1);
2049 self.running.borrow_mut().insert(socket);
2050 Ok(reply(0, ""))
2051 }
2052 }
2053
2054 #[test]
2055 fn a_session_whose_master_died_is_retried_on_a_reopened_master() {
2056 let _guard = ssh::SHARING_TEST_LOCK
2057 .lock()
2058 .unwrap_or_else(std::sync::PoisonError::into_inner);
2059 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2060 let socket_dir = tempfile::tempdir_in("/tmp").expect("short socket directory");
2061 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2062 socket_dir.path().to_path_buf(),
2063 )));
2064 let ssh = SshTarget {
2065 destination: "master-killed-once-host".to_owned(),
2066 ssh_args: Vec::new(),
2067 };
2068 let executor = MasterKilledOnce::default();
2069 let output = executor.execute(&ssh_command(&ssh, ["true"]));
2070 set_ssh_connection_sharing_for_test(None);
2071 set_ssh_retry_backoff_for_test(None);
2072
2073 assert_eq!(output.expect("the retry succeeds").status, 0);
2074 let sessions = executor.sessions.borrow();
2075 assert_eq!(sessions.len(), 2);
2076 assert_eq!(
2077 executor.openers.get(),
2078 2,
2079 "the retry reopens the master instead of trusting the earlier check"
2080 );
2081 for args in sessions.iter() {
2082 assert_eq!(args[..6][5], "ProxyCommand=false", "{args:?}");
2083 }
2084 }
2085
2086 #[test]
2087 fn a_transport_rejected_ssh_command_is_retried_once_and_then_succeeds() {
2088 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2089 let directory = tempfile::tempdir().expect("temp dir");
2090 let command = flaky_ssh_script(directory.path());
2091
2092 let output = ProcessExecutor
2093 .execute(&command)
2094 .expect("the retry must reach the successful attempt");
2095
2096 assert_eq!(output.status, 0);
2097 assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "connected");
2098 assert_eq!(attempts(directory.path()), 2);
2099 set_ssh_retry_backoff_for_test(None);
2100 }
2101
2102 #[test]
2103 fn an_untagged_command_is_not_retried_after_the_same_failure() {
2104 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2105 let directory = tempfile::tempdir().expect("temp dir");
2106 let mut command = flaky_ssh_script(directory.path());
2107 command.ssh_destination = None;
2108
2109 let output = ProcessExecutor.execute(&command).expect("runs once");
2110
2111 assert_eq!(output.status, 255);
2112 assert_eq!(attempts(directory.path()), 1);
2113 set_ssh_retry_backoff_for_test(None);
2114 }
2115
2116 #[test]
2117 fn the_cancellable_executor_also_retries_a_transport_rejection() {
2118 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2119 let directory = tempfile::tempdir().expect("temp dir");
2120 let command = flaky_ssh_script(directory.path());
2121
2122 let output = CancellableProcessExecutor::new(Arc::new(AtomicBool::new(false)))
2123 .execute(&command)
2124 .expect("the retry must reach the successful attempt");
2125
2126 assert_eq!(output.status, 0);
2127 assert_eq!(attempts(directory.path()), 2);
2128 set_ssh_retry_backoff_for_test(None);
2129 }
2130}