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(default, skip_serializing_if = "std::ops::Not::not")]
131 pub ssh_session_probe: bool,
132 #[serde(skip)]
135 sensitive_stdin: Option<SensitiveCommandInput>,
136}
137
138impl CommandSpec {
139 pub fn new(
140 program: impl Into<String>,
141 args: impl IntoIterator<Item = impl Into<String>>,
142 ) -> Self {
143 Self {
144 program: program.into(),
145 args: args.into_iter().map(Into::into).collect(),
146 env: BTreeMap::new(),
147 clear_env: false,
148 cwd: None,
149 purpose: String::new(),
150 stage: None,
151 parallel_group: None,
152 creates_target: false,
153 ssh_destination: None,
154 ssh_session: None,
155 ssh_session_probe: false,
156 sensitive_stdin: None,
157 }
158 }
159
160 pub fn purpose(mut self, purpose: impl Into<String>) -> Self {
161 self.purpose = purpose.into();
162 self
163 }
164
165 pub fn stage(mut self, stage: ProvisionStage) -> Self {
166 self.stage = Some(stage);
167 self
168 }
169
170 pub fn parallel_group(mut self, group: u32) -> Self {
173 self.parallel_group = Some(group);
174 self
175 }
176
177 pub fn ssh_destination(mut self, destination: impl Into<String>) -> Self {
180 self.ssh_destination = Some(destination.into());
181 self
182 }
183
184 pub fn ssh_session(mut self, ssh: &SshTarget) -> Self {
188 self.ssh_destination = Some(ssh.destination.clone());
189 self.ssh_session = Some(ssh.clone());
190 self
191 }
192
193 pub fn ssh_probe_session(mut self, ssh: &SshTarget) -> Self {
200 self = self.ssh_session(ssh);
201 self.ssh_session_probe = true;
202 self
203 }
204
205 pub fn open_ssh_session(&self, executor: &dyn CommandExecutor) -> Result<SessionCommand<'_>> {
211 let Some(ssh) = &self.ssh_session else {
212 return Ok(SessionCommand {
213 command: std::borrow::Cow::Borrowed(self),
214 lease: None,
215 });
216 };
217 let lease = if self.ssh_session_probe {
218 SshSessions::lease_probe(ssh)
219 } else {
220 SshSessions::lease(ssh, executor)?
221 };
222 let mut command = self.clone();
223 command.ssh_session = None;
224 command.ssh_session_probe = false;
225 command.args = session_command_args(&self.program, &self.args, ssh, &lease);
226 Ok(SessionCommand {
227 command: std::borrow::Cow::Owned(command),
228 lease: Some(lease),
229 })
230 }
231
232 pub fn creates_target(mut self) -> Self {
234 self.creates_target = true;
235 self
236 }
237
238 pub fn with_sensitive_stdin(mut self, input: Vec<u8>) -> Self {
241 self.sensitive_stdin = Some(SensitiveCommandInput(input));
242 self
243 }
244}
245
246#[derive(Debug)]
250pub struct SessionCommand<'a> {
251 command: std::borrow::Cow<'a, CommandSpec>,
252 lease: Option<SshSessionLease>,
253}
254
255impl SessionCommand<'_> {
256 pub fn command(&self) -> &CommandSpec {
257 &self.command
258 }
259
260 pub fn lease(&self) -> Option<&SshSessionLease> {
261 self.lease.as_ref()
262 }
263
264 pub fn into_parts(self) -> (CommandSpec, Option<SshSessionLease>) {
265 (self.command.into_owned(), self.lease)
266 }
267}
268
269#[derive(Debug, Clone, PartialEq, Eq)]
270pub struct CommandOutput {
271 pub status: i32,
272 pub stdout: Vec<u8>,
273 pub stderr: Vec<u8>,
274}
275
276#[derive(Debug, Clone, PartialEq, Eq)]
277pub struct SessionResourceUsage {
278 pub cpu_percent: Option<u8>,
279 pub memory_current_bytes: u64,
280 pub memory_limit_bytes: Option<u64>,
281 pub swap_current_bytes: Option<u64>,
282 pub swap_limit_bytes: Option<u64>,
283 pub writable_disk_bytes: Option<u64>,
284}
285
286#[derive(Debug, Clone, PartialEq, Eq)]
287pub struct SessionResourceProbe {
288 pub memory: CommandSpec,
289 pub disk: Option<CommandSpec>,
290}
291
292#[derive(Debug, Clone, Copy, PartialEq, Eq)]
293pub enum DeploymentCapacityKind {
294 Host,
295 AwsFleet,
296}
297
298#[derive(Debug, Clone, PartialEq, Eq)]
299pub struct DeploymentCapacityTarget {
300 pub id: String,
301 pub host: String,
302 pub target_ids: Vec<String>,
303 pub kind: DeploymentCapacityKind,
304 pub local: bool,
305 pub probes: Vec<CommandSpec>,
307 pub probe_error: Option<String>,
309}
310
311#[derive(Debug, Clone, PartialEq, Eq)]
312pub struct DeploymentCapacityUsage {
313 pub cpu_percent: Option<u8>,
314 pub memory_used_bytes: u64,
315 pub memory_total_bytes: u64,
316 pub logical_cores: u64,
317 pub disk_total_bytes: Option<u64>,
318}
319
320#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
326#[serde(from = "AdditionalMountRepr", into = "AdditionalMountRepr")]
327pub struct AdditionalMount {
328 pub source: PathBuf,
329 pub destination: PathBuf,
330 pub access: MountAccess,
331}
332
333#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
335#[serde(rename_all = "snake_case")]
336pub enum MountAccess {
337 Ro,
339 Cow,
342 Rw,
344}
345
346impl MountAccess {
347 pub const ALL: [Self; 3] = [Self::Ro, Self::Cow, Self::Rw];
348
349 pub fn label(self) -> &'static str {
350 match self {
351 Self::Ro => "ro",
352 Self::Cow => "cow",
353 Self::Rw => "rw",
354 }
355 }
356
357 pub fn without_overlay(self) -> Self {
366 match self {
367 Self::Cow => Self::Ro,
368 kept => kept,
369 }
370 }
371
372 pub fn offered(overlay_available: bool) -> Vec<Self> {
375 Self::ALL
376 .into_iter()
377 .filter(|access| overlay_available || access.without_overlay() == *access)
378 .collect()
379 }
380}
381
382#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
385pub struct ImageUser {
386 pub uid: u32,
387 pub gid: u32,
388}
389
390pub fn podman_userns_option(image_user: Option<ImageUser>) -> Option<String> {
401 image_user.map(|ImageUser { uid, gid }| format!("--userns=keep-id:uid={uid},gid={gid}"))
402}
403
404#[derive(Serialize, Deserialize)]
409#[serde(deny_unknown_fields)]
410struct AdditionalMountRepr {
411 source: PathBuf,
412 destination: PathBuf,
413 #[serde(default)]
414 read_only: bool,
415 #[serde(default, skip_serializing_if = "Option::is_none")]
416 access: Option<MountAccess>,
417}
418
419impl From<AdditionalMountRepr> for AdditionalMount {
420 fn from(repr: AdditionalMountRepr) -> Self {
421 let access = repr.access.unwrap_or(if repr.read_only {
422 MountAccess::Ro
423 } else {
424 MountAccess::Cow
425 });
426 Self {
427 source: repr.source,
428 destination: repr.destination,
429 access,
430 }
431 }
432}
433
434impl From<AdditionalMount> for AdditionalMountRepr {
435 fn from(mount: AdditionalMount) -> Self {
436 Self {
437 source: mount.source,
438 destination: mount.destination,
439 read_only: mount.access == MountAccess::Ro,
440 access: (mount.access == MountAccess::Rw).then_some(MountAccess::Rw),
441 }
442 }
443}
444
445pub fn overlay_unsupported_filesystem(filesystem: &str) -> Option<&'static str> {
451 let name = filesystem.trim().to_ascii_lowercase();
452 if name == "fuse" || name == "fuseblk" || name.starts_with("fuse.") {
454 return Some("FUSE filesystem");
455 }
456 match name.as_str() {
457 "nfs" | "nfs4" | "cifs" | "smb2" | "smb3" | "9p" | "v9fs" | "virtiofs" | "ceph"
458 | "lustre" | "afs" | "glusterfs" | "ocfs2" | "gfs" | "gfs2" => Some("network filesystem"),
459 "msdos" | "vfat" | "fat" | "exfat" | "ntfs" | "ntfs3" => Some("no POSIX metadata"),
460 "overlayfs" => Some("overlay stacking limit"),
461 _ => None,
462 }
463}
464
465pub fn validate_mount_destination(path: &Path) -> Result<()> {
467 ensure!(
468 path.is_absolute()
469 && !path
470 .components()
471 .any(|part| part == std::path::Component::ParentDir),
472 "additional mount destination must be a safe absolute container path; ~ is not supported"
473 );
474 Ok(())
475}
476
477pub fn validate_additional_mounts(mounts: &[AdditionalMount]) -> Result<()> {
478 let mut destinations = BTreeSet::new();
479 for mount in mounts {
480 if !mount.source.is_absolute() || mount.source.as_os_str().is_empty() {
481 bail!("additional mount source must be an absolute directory path");
482 }
483 validate_mount_destination(&mount.destination)?;
484 if !destinations.insert(mount.destination.clone()) {
485 bail!(
486 "additional mount destination {:?} is configured more than once",
487 mount.destination
488 );
489 }
490 }
491 Ok(())
492}
493
494pub fn default_mount_destination(source: &Path, existing: &[AdditionalMount]) -> PathBuf {
496 let basename = source
497 .file_name()
498 .filter(|name| !name.is_empty())
499 .unwrap_or_else(|| std::ffi::OsStr::new("mount"));
500 let base = PathBuf::from("/mnt").join(basename);
501 if !existing.iter().any(|mount| mount.destination == base) {
502 return base;
503 }
504 for number in 2.. {
505 let candidate =
506 PathBuf::from("/mnt").join(format!("{}-{number}", basename.to_string_lossy()));
507 if !existing.iter().any(|mount| mount.destination == candidate) {
508 return candidate;
509 }
510 }
511 unreachable!("a finite mount list always has an unused numbered destination")
512}
513
514pub trait CommandExecutor {
515 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput>;
516
517 fn cancellation_requested(&self) -> bool {
521 false
522 }
523
524 fn stage_started(&self, _stage: ProvisionStage) {}
528
529 fn stage_finished(&self, _stage: ProvisionStage) {}
532
533 fn notify_notice(&self, _notice: &str) {}
536
537 fn execute_with_stdin(
538 &self,
539 _command: &CommandSpec,
540 _input: &mut (dyn Read + Send),
541 ) -> Result<CommandOutput> {
542 bail!("this command executor does not support streamed stdin")
543 }
544}
545
546pub struct ProvisionStageGuard<'a, E: CommandExecutor + ?Sized> {
550 executor: &'a E,
551 stage: ProvisionStage,
552}
553
554impl<'a, E: CommandExecutor + ?Sized> ProvisionStageGuard<'a, E> {
555 pub fn new(executor: &'a E, stage: ProvisionStage) -> Self {
556 executor.stage_started(stage);
557 Self { executor, stage }
558 }
559}
560
561impl<E: CommandExecutor + ?Sized> Drop for ProvisionStageGuard<'_, E> {
562 fn drop(&mut self) {
563 self.executor.stage_finished(self.stage);
564 }
565}
566
567pub struct ProcessExecutor;
568
569fn with_ssh_admission(
585 command: &CommandSpec,
586 executor: &dyn CommandExecutor,
587 is_cancelled: &dyn Fn() -> bool,
588 mut run: impl FnMut(&CommandSpec) -> Result<CommandOutput>,
589) -> Result<CommandOutput> {
590 let Some(destination) = command.ssh_destination.as_deref() else {
591 return run(command);
592 };
593 for attempt in 1..=SSH_RETRY_ATTEMPTS {
594 let session = command.open_ssh_session(executor)?;
595 let output = {
596 let _permit = SshAdmission::acquire(destination);
597 run(session.command())?
598 };
599 let refusal = ssh_refusal(output.status, &String::from_utf8_lossy(&output.stderr));
600 if refusal == Some(SshRefusal::BeforeAuthentication)
603 && let Some(lease) = session.lease()
604 {
605 lease.invalidate();
606 }
607 drop(session);
608 let Some(refusal) = refusal else {
609 return Ok(output);
610 };
611 let stderr = String::from_utf8_lossy(&output.stderr);
612 if attempt == SSH_RETRY_ATTEMPTS {
613 refusal.log_exhausted(destination, &command.purpose, stderr.trim());
614 return Ok(output);
615 }
616 let delay = ssh_retry_delay(attempt);
617 refusal.log_retry(destination, &command.purpose, attempt, delay, stderr.trim());
618 if !sleep_unless_cancelled(delay, is_cancelled) {
619 bail!("operation cancelled while {}", command.purpose);
620 }
621 }
622 unreachable!("the final attempt always returns");
623}
624
625fn sleep_unless_cancelled(delay: Duration, is_cancelled: &dyn Fn() -> bool) -> bool {
628 let deadline = Instant::now() + delay;
629 loop {
630 if is_cancelled() {
631 return false;
632 }
633 let remaining = deadline.saturating_duration_since(Instant::now());
634 if remaining.is_zero() {
635 return true;
636 }
637 std::thread::sleep(remaining.min(Duration::from_millis(50)));
638 }
639}
640
641pub fn trace_command_duration(command: &CommandSpec, started: Instant, status: i32) {
644 tracing::debug!(
645 purpose = command.purpose.as_str(),
646 program = command.program.as_str(),
647 status,
648 elapsed_ms = started.elapsed().as_millis() as u64,
649 "target command finished"
650 );
651}
652
653impl ProcessExecutor {
654 fn run_once(&self, command: &CommandSpec) -> Result<CommandOutput> {
656 if let Some(input) = &command.sensitive_stdin {
657 let mut input = std::io::Cursor::new(input.0.as_slice());
658 return stream_command_with_stdin(
660 configured_command(command),
661 command,
662 &mut input,
663 &|| false,
664 );
665 }
666 let started = Instant::now();
667 let output = configured_command(command)
668 .stdin(Stdio::null())
669 .output()
670 .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
671 let status = output.status.code().unwrap_or(-1);
672 trace_command_duration(command, started, status);
673 Ok(CommandOutput {
674 status,
675 stdout: output.stdout,
676 stderr: output.stderr,
677 })
678 }
679}
680
681impl CommandExecutor for ProcessExecutor {
682 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
683 with_ssh_admission(command, self, &|| false, |command| self.run_once(command))
684 }
685
686 fn execute_with_stdin(
687 &self,
688 command: &CommandSpec,
689 input: &mut (dyn Read + Send),
690 ) -> Result<CommandOutput> {
691 let session = command.open_ssh_session(self)?;
694 let _permit = command
695 .ssh_destination
696 .as_deref()
697 .map(SshAdmission::acquire);
698 let command = session.command();
699 let process = configured_command(command);
700 stream_command_with_stdin(process, command, input, &|| false)
703 }
704}
705
706fn stream_command_with_stdin(
715 mut process: Command,
716 command: &CommandSpec,
717 input: &mut (dyn Read + Send),
718 is_cancelled: &(dyn Fn() -> bool + Sync),
719) -> Result<CommandOutput> {
720 let started = Instant::now();
721 if is_cancelled() {
722 bail!("operation cancelled");
723 }
724 let mut child = process
725 .stdin(Stdio::piped())
726 .stdout(Stdio::piped())
727 .stderr(Stdio::piped())
728 .spawn()
729 .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
730 let stdin = child
731 .stdin
732 .take()
733 .context("streamed command stdin missing")?;
734 let mut stdout = child
735 .stdout
736 .take()
737 .context("streamed command stdout missing")?;
738 let mut stderr = child
739 .stderr
740 .take()
741 .context("streamed command stderr missing")?;
742 let stdout_reader = std::thread::spawn(move || {
745 let mut bytes = Vec::new();
746 std::io::copy(&mut stdout, &mut bytes).map(|_| bytes)
747 });
748 let stderr_reader = std::thread::spawn(move || {
749 let mut bytes = Vec::new();
750 std::io::copy(&mut stderr, &mut bytes).map(|_| bytes)
751 });
752 let process_result = std::thread::scope(|scope| -> Result<_> {
753 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 status = loop {
781 if is_cancelled() {
782 terminate_cancellable_child(&mut child);
783 if let Err(error) = input_writer.join() {
784 tracing::warn!(
785 purpose = command.purpose.as_str(),
786 "streamed command input writer panicked while cancelling: {error:?}"
787 );
788 }
789 bail!("operation cancelled while {}", command.purpose);
790 }
791 match child.try_wait() {
792 Ok(Some(status)) => break status,
793 Ok(None) => std::thread::sleep(Duration::from_millis(25)),
794 Err(error) => {
795 terminate_cancellable_child(&mut child);
796 if let Err(join_error) = input_writer.join() {
797 tracing::warn!(
798 purpose = command.purpose.as_str(),
799 "streamed command input writer panicked while waiting: {join_error:?}"
800 );
801 }
802 return Err(error).with_context(|| format!("wait for {}", command.purpose));
803 }
804 }
805 };
806 let input_result = input_writer
807 .join()
808 .map_err(|_| anyhow::anyhow!("streamed command input writer panicked"))?;
809 Ok((status, input_result))
810 });
811 let stdout = stdout_reader
812 .join()
813 .map_err(|_| anyhow::anyhow!("streamed command stdout reader panicked"))??;
814 let stderr = stderr_reader
815 .join()
816 .map_err(|_| anyhow::anyhow!("streamed command stderr reader panicked"))??;
817 let (status, input_result) = process_result?;
818 if status.success() {
819 input_result?;
823 }
824 let status = status.code().unwrap_or(-1);
825 trace_command_duration(command, started, status);
826 Ok(CommandOutput {
827 status,
828 stdout,
829 stderr,
830 })
831}
832
833#[derive(Clone)]
834pub struct CancellableProcessExecutor {
835 cancelled: Arc<AtomicBool>,
836 deadline: Option<Instant>,
837}
838
839impl CancellableProcessExecutor {
840 pub fn new(cancelled: Arc<AtomicBool>) -> Self {
841 Self {
842 cancelled,
843 deadline: None,
844 }
845 }
846
847 pub fn is_cancelled(&self) -> bool {
848 self.cancelled.load(Ordering::Acquire)
849 || self
850 .deadline
851 .is_some_and(|deadline| Instant::now() >= deadline)
852 }
853
854 pub fn with_timeout(timeout: Duration) -> Self {
855 Self {
856 cancelled: Arc::new(AtomicBool::new(false)),
857 deadline: Some(Instant::now() + timeout),
858 }
859 }
860
861 pub fn with_deadline(mut self, timeout: Duration) -> Self {
864 self.deadline = Some(Instant::now() + timeout);
865 self
866 }
867
868 fn check_cancelled(&self) -> Result<()> {
869 if self.is_cancelled() {
870 bail!("operation cancelled");
871 }
872 Ok(())
873 }
874}
875
876fn configured_command(command: &CommandSpec) -> Command {
877 let mut process = Command::new(&command.program);
878 if command.clear_env {
879 process.env_clear();
880 }
881 if let Some(cwd) = &command.cwd {
882 process.current_dir(cwd);
883 }
884 process.args(&command.args).envs(&command.env);
885 process
886}
887
888fn cancellable_command(command: &CommandSpec) -> Command {
889 #[cfg(unix)]
890 let mut process = configured_command(command);
891 #[cfg(not(unix))]
892 let process = configured_command(command);
893 #[cfg(unix)]
894 {
895 use std::os::unix::process::CommandExt as _;
896 process.process_group(0);
897 }
898 process
899}
900
901fn terminate_cancellable_child(child: &mut std::process::Child) {
902 #[cfg(unix)]
903 if let Err(error) = crate::subprocess::signal_process_group(child.id() as i32, libc::SIGKILL) {
908 tracing::warn!(pid = child.id(), %error, "could not terminate cancelled command process group");
909 }
910 #[cfg(not(unix))]
911 if let Err(error) = child.kill() {
912 tracing::warn!(pid = child.id(), %error, "could not terminate cancelled command");
913 }
914 if let Err(error) = child.wait() {
915 tracing::warn!(pid = child.id(), %error, "could not reap cancelled command");
916 }
917}
918
919impl CancellableProcessExecutor {
920 fn run_once(&self, command: &CommandSpec) -> Result<CommandOutput> {
922 if let Some(input) = &command.sensitive_stdin {
923 let mut input = std::io::Cursor::new(input.0.as_slice());
924 return stream_command_with_stdin(
926 cancellable_command(command),
927 command,
928 &mut input,
929 &|| self.is_cancelled(),
930 );
931 }
932 let started = Instant::now();
933 self.check_cancelled()?;
934 let mut child = cancellable_command(command)
935 .stdin(Stdio::null())
936 .stdout(Stdio::piped())
937 .stderr(Stdio::piped())
938 .spawn()
939 .with_context(|| format!("run {} for {}", command.program, command.purpose))?;
940 let mut stdout = child.stdout.take().context("command stdout missing")?;
941 let mut stderr = child.stderr.take().context("command stderr missing")?;
942 let stdout_reader = std::thread::spawn(move || {
943 let mut bytes = Vec::new();
944 std::io::copy(&mut stdout, &mut bytes).map(|_| bytes)
945 });
946 let stderr_reader = std::thread::spawn(move || {
947 let mut bytes = Vec::new();
948 std::io::copy(&mut stderr, &mut bytes).map(|_| bytes)
949 });
950 let mut status = None;
951 let status = loop {
952 if self.is_cancelled() {
953 terminate_cancellable_child(&mut child);
954 for (stream, reader) in [("stdout", stdout_reader), ("stderr", stderr_reader)] {
955 match reader.join() {
956 Ok(Ok(_)) => {}
957 Ok(Err(error)) => {
958 tracing::warn!(stream, %error, "cancelled command reader failed")
959 }
960 Err(_) => tracing::warn!(stream, "cancelled command reader panicked"),
961 }
962 }
963 bail!("operation cancelled while {}", command.purpose);
964 }
965 if status.is_none() {
966 status = child
967 .try_wait()
968 .with_context(|| format!("wait for {}", command.purpose))?;
969 }
970 if let Some(status) = status
973 && stdout_reader.is_finished()
974 && stderr_reader.is_finished()
975 {
976 break status;
977 }
978 std::thread::sleep(Duration::from_millis(25));
979 };
980 let stdout = stdout_reader
981 .join()
982 .map_err(|_| anyhow::anyhow!("command stdout reader panicked"))??;
983 let stderr = stderr_reader
984 .join()
985 .map_err(|_| anyhow::anyhow!("command stderr reader panicked"))??;
986 let status = status.code().unwrap_or(-1);
987 trace_command_duration(command, started, status);
988 Ok(CommandOutput {
989 status,
990 stdout,
991 stderr,
992 })
993 }
994}
995
996impl CommandExecutor for CancellableProcessExecutor {
997 fn cancellation_requested(&self) -> bool {
998 self.is_cancelled()
999 }
1000
1001 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1002 with_ssh_admission(command, self, &|| self.is_cancelled(), |command| {
1003 self.run_once(command)
1004 })
1005 }
1006
1007 fn execute_with_stdin(
1008 &self,
1009 command: &CommandSpec,
1010 input: &mut (dyn Read + Send),
1011 ) -> Result<CommandOutput> {
1012 let session = command.open_ssh_session(self)?;
1015 let _permit = command
1016 .ssh_destination
1017 .as_deref()
1018 .map(SshAdmission::acquire);
1019 let command = session.command();
1020 stream_command_with_stdin(cancellable_command(command), command, input, &|| {
1023 self.is_cancelled()
1024 })
1025 }
1026}
1027
1028#[derive(Debug, Clone, PartialEq, Eq)]
1032pub struct CommandTimedOut {
1033 pub program: String,
1034 pub purpose: String,
1035 pub timeout: Duration,
1036}
1037
1038impl std::fmt::Display for CommandTimedOut {
1039 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1040 write!(
1041 formatter,
1042 "`{}` did not answer within {} seconds while trying to {}",
1043 self.program,
1044 self.timeout.as_secs(),
1045 self.purpose
1046 )
1047 }
1048}
1049
1050impl std::error::Error for CommandTimedOut {}
1051
1052#[derive(Debug, Clone, Copy)]
1061pub struct BoundedProcessExecutor {
1062 timeout: Duration,
1063}
1064
1065impl BoundedProcessExecutor {
1066 pub const fn new(timeout: Duration) -> Self {
1067 Self { timeout }
1068 }
1069}
1070
1071impl CommandExecutor for BoundedProcessExecutor {
1072 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1073 let executor = CancellableProcessExecutor::with_timeout(self.timeout);
1074 executor.execute(command).map_err(|error| {
1075 if executor.is_cancelled() {
1076 anyhow::Error::new(CommandTimedOut {
1077 program: command.program.clone(),
1078 purpose: command.purpose.clone(),
1079 timeout: self.timeout,
1080 })
1081 } else {
1082 error
1083 }
1084 })
1085 }
1086
1087 fn execute_with_stdin(
1088 &self,
1089 command: &CommandSpec,
1090 input: &mut (dyn Read + Send),
1091 ) -> Result<CommandOutput> {
1092 CancellableProcessExecutor::with_timeout(self.timeout).execute_with_stdin(command, input)
1093 }
1094}
1095
1096#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1097pub struct CommandPlan {
1098 pub description: String,
1099 pub commands: Vec<CommandSpec>,
1100}
1101
1102impl CommandPlan {
1103 pub fn provide_target_environment_secret(
1107 &mut self,
1108 target: &TargetTemplate,
1109 name: &str,
1110 value: &str,
1111 ) -> Result<()> {
1112 ensure!(
1113 !name.is_empty()
1114 && name.bytes().enumerate().all(|(index, byte)| byte == b'_'
1115 || byte.is_ascii_alphabetic()
1116 || (index > 0 && byte.is_ascii_digit())),
1117 "invalid secret environment variable name"
1118 );
1119 ensure!(
1120 !value.as_bytes().contains(&b'\n') && !value.as_bytes().contains(&b'\r'),
1121 "secret environment value cannot contain a newline"
1122 );
1123 let command = self
1124 .commands
1125 .iter_mut()
1126 .find(|command| command.creates_target)
1127 .context("provisioning plan has no target creation command")?;
1128 let read_and_export = format!("IFS= read -r {name} || exit 1; export {name};");
1129 match target {
1130 TargetTemplate::LocalPodman(_)
1131 | TargetTemplate::LocalDocker(_)
1132 | TargetTemplate::AppleContainer(_) => {
1133 let program = std::mem::replace(&mut command.program, "sh".to_owned());
1134 let args = std::mem::take(&mut command.args);
1135 command.args = vec![
1136 "-c".to_owned(),
1137 format!("{read_and_export} exec \"$@\""),
1138 "mj-secret-env".to_owned(),
1139 program,
1140 ];
1141 command.args.extend(args);
1142 }
1143 TargetTemplate::SshPodman { .. } | TargetTemplate::SshDocker { .. } => {
1144 let remote = command
1145 .args
1146 .last_mut()
1147 .context("remote container command has no SSH command argument")?;
1148 *remote = format!("{read_and_export} exec {remote}");
1149 }
1150 TargetTemplate::LocalBare
1151 | TargetTemplate::AwsEc2(_)
1152 | TargetTemplate::SshBare { .. } => {
1153 bail!("target does not support inherited container environment")
1154 }
1155 }
1156 let mut input = value.as_bytes().to_vec();
1157 input.push(b'\n');
1158 command.sensitive_stdin = Some(SensitiveCommandInput(input));
1159 Ok(())
1160 }
1161
1162 pub fn execute(&self, executor: &impl CommandExecutor) -> Result<Vec<CommandOutput>> {
1163 let mut outputs = Vec::with_capacity(self.commands.len());
1164 for command in &self.commands {
1165 let output = executor.execute(command)?;
1166 if output.status != 0 {
1167 bail!(
1168 "{} failed with status {}: {}",
1169 command.purpose,
1170 output.status,
1171 String::from_utf8_lossy(&output.stderr)
1172 );
1173 }
1174 outputs.push(output);
1175 }
1176 Ok(outputs)
1177 }
1178
1179 pub fn execute_concurrent(
1191 &self,
1192 executor: &(impl CommandExecutor + Sync),
1193 ) -> Result<Vec<CommandOutput>> {
1194 let mut outputs = Vec::with_capacity(self.commands.len());
1195 let mut index = 0;
1196 while index < self.commands.len() {
1197 let group = self.commands[index].parallel_group;
1198 let mut end = index + 1;
1199 if group.is_some() {
1200 while end < self.commands.len() && self.commands[end].parallel_group == group {
1201 end += 1;
1202 }
1203 }
1204 let batch = &self.commands[index..end];
1205 if let [command] = batch {
1206 outputs.push(checked_command_output(command, executor.execute(command)?)?);
1207 } else {
1208 let results: Vec<Result<CommandOutput>> = std::thread::scope(|scope| {
1209 let handles: Vec<_> = batch
1210 .iter()
1211 .map(|command| scope.spawn(|| executor.execute(command)))
1212 .collect();
1213 handles
1214 .into_iter()
1215 .map(|handle| match handle.join() {
1216 Ok(result) => result,
1217 Err(panic) => Err(anyhow::anyhow!(
1218 "concurrent command thread panicked: {}",
1219 command_thread_panic_message(panic.as_ref())
1220 )),
1221 })
1222 .collect()
1223 });
1224 for (command, result) in batch.iter().zip(results) {
1225 outputs.push(checked_command_output(command, result?)?);
1226 }
1227 }
1228 index = end;
1229 }
1230 Ok(outputs)
1231 }
1232
1233 pub fn split_at_target_creation(&self) -> Option<(Self, Self)> {
1241 let created = self
1242 .commands
1243 .iter()
1244 .position(|command| command.creates_target)?;
1245 let (creation, remainder) = self.commands.split_at(created + 1);
1246 Some((
1247 Self {
1248 description: self.description.clone(),
1249 commands: creation.to_vec(),
1250 },
1251 Self {
1252 description: self.description.clone(),
1253 commands: remainder.to_vec(),
1254 },
1255 ))
1256 }
1257}
1258
1259pub fn checked_command_output(
1263 command: &CommandSpec,
1264 output: CommandOutput,
1265) -> Result<CommandOutput> {
1266 if output.status != 0 {
1267 bail!(
1268 "{} failed with status {}: {}",
1269 command.purpose,
1270 output.status,
1271 String::from_utf8_lossy(&output.stderr)
1272 );
1273 }
1274 Ok(output)
1275}
1276
1277pub fn command_thread_panic_message(payload: &(dyn std::any::Any + Send)) -> String {
1279 if let Some(message) = payload.downcast_ref::<&str>() {
1280 (*message).to_owned()
1281 } else if let Some(message) = payload.downcast_ref::<String>() {
1282 message.clone()
1283 } else {
1284 "non-string panic payload".to_owned()
1285 }
1286}
1287
1288#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1289pub struct RepositorySpec {
1290 pub url: Option<String>,
1292 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1293 pub push_urls: Vec<String>,
1294 pub destination: String,
1295 pub git_ref: Option<String>,
1296 #[serde(default, skip_serializing_if = "Option::is_none")]
1299 pub reference: Option<String>,
1300}
1301
1302#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1303pub struct ProjectBundleSpec {
1304 pub primary: String,
1305 pub repositories: Vec<RepositorySpec>,
1306}
1307
1308impl ProjectBundleSpec {
1309 pub fn validate(&self) -> Result<()> {
1310 validate_relative_path(&self.primary)?;
1311 if self.repositories.is_empty() {
1312 bail!("a project bundle must contain at least one repository");
1313 }
1314 let mut destinations = std::collections::BTreeSet::new();
1315 for repository in &self.repositories {
1316 validate_relative_path(&repository.destination)?;
1317 ensure!(
1318 repository
1319 .url
1320 .as_deref()
1321 .is_some_and(|url| !url.trim().is_empty() && !url.starts_with('-')),
1322 "isolated repositories require a network Git remote; configure a remote or use a raw local session"
1323 );
1324 crate::remote_git::validate_network_url(
1325 repository.url.as_deref().expect("checked above"),
1326 )?;
1327 for push_url in &repository.push_urls {
1328 crate::remote_git::validate_network_url(push_url)?;
1329 }
1330 ensure!(
1331 repository.git_ref.is_none(),
1332 "git_ref is no longer supported; remove it to start from the remote's default branch"
1333 );
1334 if !destinations.insert(&repository.destination) {
1335 bail!(
1336 "duplicate repository destination {}",
1337 repository.destination
1338 );
1339 }
1340 }
1341 if !destinations.contains(&self.primary) {
1342 bail!("primary repository is not present in the bundle");
1343 }
1344 Ok(())
1345 }
1346}
1347
1348#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1349#[serde(tag = "kind", rename_all = "snake_case")]
1350pub enum PodmanWorkspaceStorage {
1351 PodmanVolume,
1352 HostHelper {
1353 root: String,
1354 helper: Vec<String>,
1355 },
1356 #[default]
1357 ContainerLayer,
1358}
1359
1360#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1361pub struct ContainerTemplate {
1362 pub image: String,
1363 #[serde(default)]
1364 pub pull_policy: ImagePullPolicy,
1365 #[serde(default)]
1366 pub extra_run_args: Vec<String>,
1367 #[serde(default)]
1368 pub workspace_storage: PodmanWorkspaceStorage,
1369 #[serde(default)]
1372 pub build_cache: Option<crate::config::TargetBuildCache>,
1373}
1374
1375impl ImagePullPolicy {
1376 pub fn resolve(self, image: &str) -> Self {
1379 if self != Self::Auto {
1380 return self;
1381 }
1382 if image_is_digest_pinned(image) {
1383 Self::Missing
1384 } else if image_is_remote(image) && image_uses_latest_tag(image) {
1385 Self::Newer
1386 } else {
1387 Self::Missing
1388 }
1389 }
1390
1391 pub fn at_launch(self, image: &str) -> Self {
1396 if self == Self::Auto {
1397 Self::Missing
1398 } else {
1399 self.resolve(image)
1400 }
1401 }
1402
1403 pub fn describe(self, image: &str) -> &'static str {
1408 match self {
1409 Self::Always => "Pull every launch",
1410 Self::Newer => "Pull when the registry is newer",
1411 Self::Missing => "Pull only if missing",
1412 Self::Never => "Never pull",
1413 Self::Auto => match self.resolve(image) {
1414 Self::Newer => "Pull if missing at launch; refresh :latest in background",
1415 _ => "Pull if missing",
1416 },
1417 }
1418 }
1419
1420 pub fn podman_value(self) -> &'static str {
1422 match self {
1423 Self::Always => "always",
1424 Self::Newer => "newer",
1425 Self::Missing => "missing",
1426 Self::Never => "never",
1427 Self::Auto => unreachable!("auto pull policy must resolve"),
1428 }
1429 }
1430}
1431
1432#[derive(Debug, Clone, PartialEq, Eq)]
1438pub enum ImageHost {
1439 LocalPodman,
1440 LocalDocker,
1441 AppleContainer,
1442 SshPodman(SshTarget),
1443 SshDocker(SshTarget),
1444}
1445
1446impl ImageHost {
1447 pub const fn engine(&self) -> &'static str {
1448 match self {
1449 Self::LocalPodman | Self::SshPodman(_) => "podman",
1450 Self::LocalDocker | Self::SshDocker(_) => "docker",
1451 Self::AppleContainer => "container",
1452 }
1453 }
1454
1455 pub fn label(&self) -> String {
1457 match self {
1458 Self::LocalPodman => "local podman".to_owned(),
1459 Self::LocalDocker => "local docker".to_owned(),
1460 Self::AppleContainer => "apple container".to_owned(),
1461 Self::SshPodman(ssh) => format!("podman on {}", ssh.destination),
1462 Self::SshDocker(ssh) => format!("docker on {}", ssh.destination),
1463 }
1464 }
1465
1466 fn command(&self, args: Vec<String>, purpose: String) -> CommandSpec {
1467 match self {
1468 Self::LocalPodman | Self::LocalDocker | Self::AppleContainer => {
1469 CommandSpec::new(args[0].clone(), args[1..].iter().cloned())
1470 }
1471 Self::SshPodman(ssh) | Self::SshDocker(ssh) => ssh_command_owned(ssh, args),
1472 }
1473 .purpose(purpose)
1474 }
1475}
1476
1477#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
1484pub enum RefreshWhen {
1485 WhenAbsent,
1487 Always,
1489}
1490
1491#[derive(Debug, Clone, PartialEq, Eq)]
1494pub struct ImageRefresh {
1495 pub host: ImageHost,
1496 pub image: String,
1497 pub platform: Option<String>,
1498 pub when: RefreshWhen,
1501 pub image_id: CommandSpec,
1504 pub pull: CommandSpec,
1505 pub prune: Option<CommandSpec>,
1508}
1509
1510pub fn image_refresh(
1517 host: ImageHost,
1518 image: &str,
1519 platform: Option<&str>,
1520 pull_policy: ImagePullPolicy,
1521) -> Option<ImageRefresh> {
1522 let when = match pull_policy.resolve(image) {
1523 ImagePullPolicy::Always | ImagePullPolicy::Newer => RefreshWhen::Always,
1524 ImagePullPolicy::Missing => RefreshWhen::WhenAbsent,
1525 ImagePullPolicy::Never => return None,
1526 ImagePullPolicy::Auto => unreachable!("auto pull policy must resolve"),
1527 };
1528 let engine = host.engine();
1529 let apple = matches!(host, ImageHost::AppleContainer);
1533 let mut image_id_args = vec![engine.to_owned(), "image".to_owned(), "inspect".to_owned()];
1534 if !apple {
1535 image_id_args.push("--format".to_owned());
1536 image_id_args.push("{{.Id}}".to_owned());
1537 }
1538 image_id_args.push(image.to_owned());
1539 let image_id = host.command(
1540 image_id_args,
1541 format!("read the cached id of container image {image}"),
1542 );
1543 let mut pull_args = vec![engine.to_owned()];
1544 if apple {
1545 pull_args.push("image".to_owned());
1546 }
1547 pull_args.push("pull".to_owned());
1548 if let Some(platform) = platform.filter(|_| !apple) {
1550 pull_args.push(format!("--platform={platform}"));
1551 }
1552 pull_args.push(image.to_owned());
1553 let pull = host.command(pull_args, format!("refresh container image {image}"));
1554 let prune = (!apple).then(|| {
1555 host.command(
1556 vec![
1557 engine.to_owned(),
1558 "image".to_owned(),
1559 "prune".to_owned(),
1560 "-f".to_owned(),
1561 ],
1562 "remove dangling container images".to_owned(),
1563 )
1564 });
1565 Some(ImageRefresh {
1566 host,
1567 image: image.to_owned(),
1568 platform: platform.map(str::to_owned),
1569 when,
1570 image_id,
1571 pull,
1572 prune,
1573 })
1574}
1575
1576fn image_is_digest_pinned(image: &str) -> bool {
1577 image
1578 .rsplit_once('@')
1579 .is_some_and(|(_, digest)| !digest.is_empty())
1580}
1581
1582fn image_is_remote(image: &str) -> bool {
1583 !image.starts_with("localhost/") && !image.starts_with("local/")
1584}
1585
1586fn image_uses_latest_tag(image: &str) -> bool {
1587 let name = image.split_once('@').map_or(image, |(name, _)| name);
1588 let final_component = name.rsplit('/').next().unwrap_or(name);
1589 !final_component.contains(':') || final_component.ends_with(":latest")
1590}
1591
1592#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1593pub struct SshTarget {
1594 pub destination: String,
1595 #[serde(default)]
1596 pub ssh_args: Vec<String>,
1597}
1598
1599#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1600pub struct AwsTemplate {
1601 pub profile: String,
1602 pub region: String,
1603 pub launch_template: String,
1604 pub launch_template_version: Option<String>,
1605 pub instance_type: Option<String>,
1606 pub ssh: SshTarget,
1607}
1608
1609#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1610#[serde(tag = "kind", rename_all = "snake_case")]
1611pub enum TargetTemplate {
1612 LocalBare,
1613 LocalPodman(ContainerTemplate),
1614 LocalDocker(ContainerTemplate),
1615 AppleContainer(ContainerTemplate),
1616 AwsEc2(AwsTemplate),
1617 SshBare {
1618 ssh: SshTarget,
1619 #[serde(default = "default_ssh_prefix")]
1620 workspace_prefix: String,
1621 },
1622 SshPodman {
1623 ssh: SshTarget,
1624 container: ContainerTemplate,
1625 },
1626 SshDocker {
1627 ssh: SshTarget,
1628 container: ContainerTemplate,
1629 },
1630}
1631
1632fn default_ssh_prefix() -> String {
1633 ".local/share/hel/workspaces".to_owned()
1634}
1635
1636#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1637#[serde(tag = "kind", rename_all = "snake_case")]
1638pub enum PodmanWorkspaceLocator {
1639 #[default]
1640 ContainerLayer,
1641 Volume {
1642 name: String,
1643 },
1644 HostPath {
1645 path: String,
1646 helper: Vec<String>,
1647 resource: String,
1648 },
1649}
1650
1651#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1652#[serde(tag = "kind", rename_all = "snake_case")]
1653pub enum TargetLocator {
1654 LocalBare {
1655 worker_root: String,
1656 },
1657 LocalPodman {
1658 container_id: String,
1659 #[serde(default)]
1660 workspace_storage: PodmanWorkspaceLocator,
1661 #[serde(default, skip_serializing_if = "Option::is_none")]
1665 borrowed_from: Option<String>,
1666 },
1667 LocalDocker {
1668 container_id: String,
1669 #[serde(default, skip_serializing_if = "Option::is_none")]
1673 borrowed_from: Option<String>,
1674 },
1675 AppleContainer {
1676 container_id: String,
1677 #[serde(default, skip_serializing_if = "Option::is_none")]
1681 borrowed_from: Option<String>,
1682 },
1683 AwsEc2 {
1684 profile: String,
1685 region: String,
1686 instance_id: String,
1687 ssh: SshTarget,
1688 workspace: String,
1689 },
1690 SshBare {
1691 ssh: SshTarget,
1692 workspace: String,
1693 #[serde(default, skip_serializing_if = "Option::is_none")]
1695 worker_id: Option<String>,
1696 },
1697 SshPodman {
1698 ssh: SshTarget,
1699 container_id: String,
1700 #[serde(default)]
1701 workspace_storage: PodmanWorkspaceLocator,
1702 #[serde(default, skip_serializing_if = "Option::is_none")]
1706 borrowed_from: Option<String>,
1707 },
1708 SshDocker {
1709 ssh: SshTarget,
1710 container_id: String,
1711 #[serde(default, skip_serializing_if = "Option::is_none")]
1715 borrowed_from: Option<String>,
1716 },
1717}
1718
1719impl TargetTemplate {
1720 pub const fn container_engine(&self) -> Option<&'static str> {
1721 match self {
1722 Self::LocalPodman(_) | Self::SshPodman { .. } => Some("podman"),
1723 Self::LocalDocker(_) | Self::SshDocker { .. } => Some("docker"),
1724 Self::AppleContainer(_) => Some("container"),
1725 _ => None,
1726 }
1727 }
1728
1729 pub fn image_host(&self) -> Option<(ImageHost, &ContainerTemplate)> {
1732 match self {
1733 Self::LocalPodman(container) => Some((ImageHost::LocalPodman, container)),
1734 Self::LocalDocker(container) => Some((ImageHost::LocalDocker, container)),
1735 Self::AppleContainer(container) => Some((ImageHost::AppleContainer, container)),
1736 Self::SshPodman { ssh, container } => {
1737 Some((ImageHost::SshPodman(ssh.clone()), container))
1738 }
1739 Self::SshDocker { ssh, container } => {
1740 Some((ImageHost::SshDocker(ssh.clone()), container))
1741 }
1742 Self::LocalBare | Self::AwsEc2(_) | Self::SshBare { .. } => None,
1743 }
1744 }
1745}
1746
1747impl TargetLocator {
1748 pub const fn harness_host(&self) -> HarnessHost {
1752 match self {
1753 Self::LocalBare { .. } => HarnessHost::current(),
1754 _ => HarnessHost::Other,
1755 }
1756 }
1757
1758 pub const fn kind_name(&self) -> &'static str {
1760 match self {
1761 Self::LocalBare { .. } => "local-bare",
1762 Self::LocalPodman { .. } => "local-podman",
1763 Self::LocalDocker { .. } => "local-docker",
1764 Self::AppleContainer { .. } => "apple-container",
1765 Self::AwsEc2 { .. } => "aws-ec2",
1766 Self::SshBare { .. } => "ssh-bare",
1767 Self::SshPodman { .. } => "ssh-podman",
1768 Self::SshDocker { .. } => "ssh-docker",
1769 }
1770 }
1771
1772 pub const fn container_engine(&self) -> Option<&'static str> {
1773 match self {
1774 Self::LocalPodman { .. } | Self::SshPodman { .. } => Some("podman"),
1775 Self::LocalDocker { .. } | Self::SshDocker { .. } => Some("docker"),
1776 Self::AppleContainer { .. } => Some("container"),
1777 _ => None,
1778 }
1779 }
1780}
1781
1782#[derive(Debug, Clone, PartialEq, Eq)]
1786pub struct TargetRecoveryPlan {
1787 pub exists: CommandSpec,
1788 pub inspect: CommandSpec,
1789 pub start: CommandSpec,
1790 pub session_id: String,
1791}
1792
1793#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1794pub enum TargetRecoveryOutcome {
1795 NotRequired,
1796 Missing,
1797 AlreadyRunning,
1798 Started,
1799}
1800
1801pub fn resource_name(session_id: &str) -> Result<String> {
1802 validate_session_id(session_id)?;
1803 let readable: String = session_id
1804 .chars()
1805 .filter(|character| character.is_ascii_alphanumeric())
1806 .take(12)
1807 .map(|character| character.to_ascii_lowercase())
1808 .collect();
1809 let digest = Sha256::digest(session_id.as_bytes());
1810 Ok(format!(
1811 "mj-{readable}-{:02x}{:02x}{:02x}",
1812 digest[0], digest[1], digest[2]
1813 ))
1814}
1815
1816pub fn podman_workspace_locator(
1817 template: &ContainerTemplate,
1818 session_id: &str,
1819) -> Result<PodmanWorkspaceLocator> {
1820 let resource = format!("{}-workspace", resource_name(session_id)?);
1821 match &template.workspace_storage {
1822 PodmanWorkspaceStorage::PodmanVolume => {
1823 Ok(PodmanWorkspaceLocator::Volume { name: resource })
1824 }
1825 PodmanWorkspaceStorage::HostHelper { root, helper } => {
1826 let root = Path::new(root);
1827 ensure!(
1828 root.is_absolute(),
1829 "Podman workspace storage root must be absolute"
1830 );
1831 ensure!(
1832 !helper.is_empty() && helper.iter().all(|argument| !argument.is_empty()),
1833 "Podman workspace storage helper must contain non-empty arguments"
1834 );
1835 Ok(PodmanWorkspaceLocator::HostPath {
1836 path: root.join(&resource).to_string_lossy().into_owned(),
1837 helper: helper.clone(),
1838 resource,
1839 })
1840 }
1841 PodmanWorkspaceStorage::ContainerLayer => Ok(PodmanWorkspaceLocator::ContainerLayer),
1842 }
1843}
1844
1845pub fn container_workspace_root(recorded: Option<&Path>) -> String {
1853 recorded.map_or_else(
1854 || CONTAINER_WORKSPACE.to_owned(),
1855 |path| path.to_string_lossy().into_owned(),
1856 )
1857}
1858
1859pub fn new_container_workspace(session_id: &str) -> Result<PathBuf> {
1861 validate_session_id(session_id)?;
1862 Ok(Path::new(CONTAINER_WORKSPACE).join(session_id))
1863}
1864
1865pub fn aws_workspace(session_id: &str) -> String {
1871 format!(".local/share/hel/workspaces/{session_id}")
1872}
1873
1874pub fn workspace_for(template: &TargetTemplate, session_id: &str) -> Result<String> {
1875 validate_session_id(session_id)?;
1876 match template {
1877 TargetTemplate::LocalBare => bail!("local bare projects use their selected directory"),
1878 TargetTemplate::LocalPodman(_)
1881 | TargetTemplate::LocalDocker(_)
1882 | TargetTemplate::AppleContainer(_)
1883 | TargetTemplate::SshPodman { .. }
1884 | TargetTemplate::SshDocker { .. } => {
1885 bail!("container targets use the session's recorded container workspace")
1886 }
1887 TargetTemplate::AwsEc2(_) => Ok(aws_workspace(session_id)),
1888 TargetTemplate::SshBare {
1889 workspace_prefix, ..
1890 } => {
1891 validate_workspace_prefix(workspace_prefix)?;
1892 let prefix = workspace_prefix
1897 .strip_prefix("~/")
1898 .unwrap_or(workspace_prefix);
1899 Ok(format!("{}/{session_id}", prefix.trim_end_matches('/')))
1900 }
1901 }
1902}
1903
1904pub fn command_on_locator(
1906 locator: &TargetLocator,
1907 session_id: &str,
1908 args: Vec<String>,
1909 purpose: impl Into<String>,
1910) -> Result<CommandSpec> {
1911 verify_locator(locator, session_id)?;
1912 if args.is_empty() {
1913 bail!("target command must not be empty");
1914 }
1915 Ok(locator_command(locator, args).purpose(purpose))
1916}
1917
1918pub fn locator_command(locator: &TargetLocator, args: Vec<String>) -> CommandSpec {
1922 match locator {
1923 TargetLocator::LocalBare { .. } => {
1924 let mut args = args.into_iter();
1925 let program = args.next().expect("target command must not be empty");
1926 CommandSpec::new(program, args)
1927 }
1928 TargetLocator::LocalPodman { container_id, .. }
1929 | TargetLocator::LocalDocker { container_id, .. }
1930 | TargetLocator::AppleContainer { container_id, .. } => container_exec(
1931 locator.container_engine().expect("local container"),
1932 container_id,
1933 args,
1934 ),
1935 TargetLocator::AwsEc2 { ssh, .. } | TargetLocator::SshBare { ssh, .. } => {
1936 ssh_command_owned(ssh, args)
1937 }
1938 TargetLocator::SshPodman {
1939 ssh, container_id, ..
1940 }
1941 | TargetLocator::SshDocker {
1942 ssh, container_id, ..
1943 } => {
1944 let mut remote = vec![
1945 locator
1946 .container_engine()
1947 .expect("remote container")
1948 .to_owned(),
1949 "exec".to_owned(),
1950 "-i".to_owned(),
1951 container_id.to_owned(),
1952 ];
1953 remote.extend(args);
1954 ssh_command_owned(ssh, remote)
1955 }
1956 }
1957}
1958pub fn worker_root(locator: &TargetLocator, session_id: &str) -> Result<String> {
1959 verify_locator(locator, session_id)?;
1960 Ok(match locator {
1961 TargetLocator::LocalBare { worker_root } => worker_root.clone(),
1962 TargetLocator::LocalPodman { .. }
1963 | TargetLocator::LocalDocker { .. }
1964 | TargetLocator::AppleContainer { .. }
1965 | TargetLocator::SshPodman { .. }
1966 | TargetLocator::SshDocker { .. } => format!("/var/lib/hel/workers/{session_id}"),
1967 TargetLocator::AwsEc2 { .. } => format!(".local/share/hel/workers/{session_id}"),
1968 TargetLocator::SshBare { worker_id, .. } => format!(
1969 ".local/share/hel/workers/{}",
1970 worker_id.as_deref().unwrap_or(session_id)
1971 ),
1972 })
1973}
1974mod convert;
1975pub use convert::{
1976 RecordedTarget, StoredTarget, TargetConversionError, locator_needs_connection,
1977 ssh_args_with_identity,
1978};
1979
1980mod ssh;
1981pub use ssh::*;
1982
1983pub fn container_exec(
1984 engine: &str,
1985 container_id: &str,
1986 args: impl IntoIterator<Item = impl Into<String>>,
1987) -> CommandSpec {
1988 let mut command_args = vec!["exec".to_owned(), "-i".to_owned(), container_id.to_owned()];
1989 command_args.extend(args.into_iter().map(Into::into));
1990 CommandSpec::new(engine, command_args)
1991}
1992
1993#[cfg(all(test, unix))]
1994mod executor_tests {
1995 use std::fs;
1996
1997 use super::*;
1998
1999 fn flaky_ssh_script(directory: &Path) -> CommandSpec {
2003 let counter = directory.join("attempts");
2004 let script = format!(
2005 "count=$(cat {counter} 2>/dev/null || echo 0)\n\
2006 echo $((count + 1)) > {counter}\n\
2007 if [ \"$count\" -eq 0 ]; then\n\
2008 echo 'kex_exchange_identification: Connection closed by 10.0.0.1 port 22' >&2\n\
2009 exit 255\n\
2010 fi\n\
2011 echo connected\n",
2012 counter = counter.display()
2013 );
2014 CommandSpec::new("sh", ["-c".to_owned(), script])
2015 .ssh_destination("build@10.0.0.1")
2016 .purpose("run the flaky SSH fixture")
2017 }
2018
2019 fn attempts(directory: &Path) -> u32 {
2020 fs::read_to_string(directory.join("attempts"))
2021 .expect("the fixture records its attempts")
2022 .trim()
2023 .parse()
2024 .expect("attempt count is a number")
2025 }
2026
2027 #[derive(Default)]
2032 struct MasterKilledOnce {
2033 running: std::cell::RefCell<BTreeSet<String>>,
2034 openers: std::cell::Cell<usize>,
2035 sessions: std::cell::RefCell<Vec<Vec<String>>>,
2036 }
2037
2038 impl CommandExecutor for MasterKilledOnce {
2039 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2040 let reply = |status: i32, stderr: &str| CommandOutput {
2041 status,
2042 stdout: Vec::new(),
2043 stderr: stderr.as_bytes().to_vec(),
2044 };
2045 if command.ssh_session.is_some() {
2046 return with_ssh_admission(command, self, &|| false, |spawned| {
2048 self.sessions.borrow_mut().push(spawned.args.clone());
2049 if self.sessions.borrow().len() == 1 {
2050 self.running.borrow_mut().clear();
2051 return Ok(reply(255, "Connection closed by UNKNOWN port 65535"));
2052 }
2053 Ok(reply(0, ""))
2054 });
2055 }
2056 let socket = command
2057 .args
2058 .iter()
2059 .find_map(|arg| arg.strip_prefix("ControlPath="))
2060 .expect("a master command names its socket")
2061 .to_owned();
2062 if command.args.windows(2).any(|pair| pair == ["-O", "check"]) {
2063 return Ok(reply(
2064 if self.running.borrow().contains(&socket) {
2065 0
2066 } else {
2067 255
2068 },
2069 "",
2070 ));
2071 }
2072 assert!(command.args.contains(&"ControlMaster=yes".to_owned()));
2073 self.openers.set(self.openers.get() + 1);
2074 self.running.borrow_mut().insert(socket);
2075 Ok(reply(0, ""))
2076 }
2077 }
2078
2079 #[test]
2080 fn a_session_whose_master_died_is_retried_on_a_reopened_master() {
2081 let _guard = ssh::SHARING_TEST_LOCK
2082 .lock()
2083 .unwrap_or_else(std::sync::PoisonError::into_inner);
2084 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2085 let socket_dir = tempfile::tempdir_in("/tmp").expect("short socket directory");
2086 set_ssh_connection_sharing_for_test(Some(SshSharingForTest::Directory(
2087 socket_dir.path().to_path_buf(),
2088 )));
2089 let ssh = SshTarget {
2090 destination: "master-killed-once-host".to_owned(),
2091 ssh_args: Vec::new(),
2092 };
2093 let executor = MasterKilledOnce::default();
2094 let output = executor.execute(&ssh_command(&ssh, ["true"]));
2095 set_ssh_connection_sharing_for_test(None);
2096 set_ssh_retry_backoff_for_test(None);
2097
2098 assert_eq!(output.expect("the retry succeeds").status, 0);
2099 let sessions = executor.sessions.borrow();
2100 assert_eq!(sessions.len(), 2);
2101 assert_eq!(
2102 executor.openers.get(),
2103 2,
2104 "the retry reopens the master instead of trusting the earlier check"
2105 );
2106 for args in sessions.iter() {
2107 assert_eq!(args[..6][5], "ProxyCommand=false", "{args:?}");
2108 }
2109 }
2110
2111 #[test]
2112 fn a_transport_rejected_ssh_command_is_retried_once_and_then_succeeds() {
2113 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2114 let directory = tempfile::tempdir().expect("temp dir");
2115 let command = flaky_ssh_script(directory.path());
2116
2117 let output = ProcessExecutor
2118 .execute(&command)
2119 .expect("the retry must reach the successful attempt");
2120
2121 assert_eq!(output.status, 0);
2122 assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "connected");
2123 assert_eq!(attempts(directory.path()), 2);
2124 set_ssh_retry_backoff_for_test(None);
2125 }
2126
2127 #[test]
2128 fn an_untagged_command_is_not_retried_after_the_same_failure() {
2129 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2130 let directory = tempfile::tempdir().expect("temp dir");
2131 let mut command = flaky_ssh_script(directory.path());
2132 command.ssh_destination = None;
2133
2134 let output = ProcessExecutor.execute(&command).expect("runs once");
2135
2136 assert_eq!(output.status, 255);
2137 assert_eq!(attempts(directory.path()), 1);
2138 set_ssh_retry_backoff_for_test(None);
2139 }
2140
2141 #[test]
2142 fn the_cancellable_executor_also_retries_a_transport_rejection() {
2143 set_ssh_retry_backoff_for_test(Some(Duration::from_millis(5)));
2144 let directory = tempfile::tempdir().expect("temp dir");
2145 let command = flaky_ssh_script(directory.path());
2146
2147 let output = CancellableProcessExecutor::new(Arc::new(AtomicBool::new(false)))
2148 .execute(&command)
2149 .expect("the retry must reach the successful attempt");
2150
2151 assert_eq!(output.status, 0);
2152 assert_eq!(attempts(directory.path()), 2);
2153 set_ssh_retry_backoff_for_test(None);
2154 }
2155}