1use std::collections::BTreeMap;
4use std::path::PathBuf;
5use std::process::{Command, Stdio};
6use std::time::{Duration, Instant};
7
8use anyhow::{Context, Result, bail, ensure};
9
10use mj_core::config::{AwsAddressSource, Config, ProjectBundle, TargetTemplate, data_dir};
11use mj_core::state::{
12 PodmanWorkspaceLocator, SessionRecord, SessionResourceAllocation, TargetLocator,
13 allocation_cpus,
14};
15
16use crate::targets::{
17 self, AwsTemplate, CommandExecutor, CommandOutput, CommandSpec, ContainerTemplate,
18 ImageRefresh, ProjectBundleSpec, ProvisionStage, RepositorySpec, SshTarget,
19};
20
21use super::{Controller, execute_checked};
22
23impl Controller {
24 pub fn session_working_context(
26 &self,
27 session_id: &str,
28 executor: &impl CommandExecutor,
29 ) -> Result<(PathBuf, String)> {
30 let session = self
31 .state
32 .sessions
33 .get(session_id)
34 .context("session is missing")?;
35 let locator = session
36 .target
37 .as_ref()
38 .context("target is still starting")?;
39 let backend = backend_locator(locator, session, &self.config)?;
40 let launch = self.current_worker_launch_config(session_id, &backend)?;
41 let output = executor.execute(&targets::command_on_locator(
42 &backend,
43 session_id,
44 vec![
45 "git".into(),
46 "-C".into(),
47 launch.cwd.to_string_lossy().into_owned(),
48 "rev-parse".into(),
49 "--abbrev-ref".into(),
50 "HEAD".into(),
51 ],
52 "read current session branch",
53 )?)?;
54 let branch = if output.status == 0 {
55 let branch = String::from_utf8(output.stdout).context("decode session branch")?;
56 if branch.trim() == "HEAD" {
57 "detached HEAD".to_owned()
58 } else {
59 branch.trim().to_owned()
60 }
61 } else {
62 let stderr = String::from_utf8_lossy(&output.stderr);
63 if stderr.contains("not a git repository") {
64 "not a git checkout".to_owned()
68 } else {
69 format!(
70 "unavailable: {}",
71 stderr.lines().next().unwrap_or("").trim()
72 )
73 }
74 };
75 Ok((launch.cwd, branch))
76 }
77
78 pub fn session_git_status(
83 &self,
84 session_id: &str,
85 executor: &impl CommandExecutor,
86 ) -> Result<mj_core::local_git::SessionGitStatus> {
87 let (cwd, branch) = self.session_working_context(session_id, executor)?;
88 if branch.starts_with("not a git") || branch.starts_with("unavailable") {
89 return Ok(mj_core::local_git::parse_git_status(
90 cwd, &branch, None, "", "",
91 ));
92 }
93 let session = self
94 .state
95 .sessions
96 .get(session_id)
97 .context("session is missing")?;
98 let locator = session
99 .target
100 .as_ref()
101 .context("target is still starting")?;
102 let backend = backend_locator(locator, session, &self.config)?;
103 let cwd_text = cwd.to_string_lossy().into_owned();
104 let run = |args: &[&str], purpose: &str| -> Result<Option<String>> {
105 let mut command = vec!["git".to_owned(), "-C".to_owned(), cwd_text.clone()];
106 command.extend(args.iter().map(|arg| (*arg).to_owned()));
107 let output = executor.execute(&targets::command_on_locator(
108 &backend, session_id, command, purpose,
109 )?)?;
110 Ok((output.status == 0).then(|| String::from_utf8_lossy(&output.stdout).into_owned()))
111 };
112 let ahead_behind = run(
115 &["rev-list", "--left-right", "--count", "@{upstream}...HEAD"],
116 "count commits against upstream",
117 )?;
118 let numstat = run(
120 &["--no-optional-locks", "diff", "--numstat", "HEAD"],
121 "count changed lines",
122 )?
123 .unwrap_or_default();
124 let porcelain = run(
125 &[
126 "--no-optional-locks",
127 "status",
128 "--porcelain",
129 "--untracked-files=normal",
130 ],
131 "list changed files",
132 )?
133 .unwrap_or_default();
134 Ok(mj_core::local_git::parse_git_status(
135 cwd,
136 &branch,
137 ahead_behind.as_deref(),
138 &numstat,
139 &porcelain,
140 ))
141 }
142
143 pub fn resolve_aws_resource_options(
144 &self,
145 target_id: &str,
146 executor: &impl CommandExecutor,
147 ) -> Result<Vec<SessionResourceAllocation>> {
148 let TargetTemplate::AwsEc2 {
149 aws_profile,
150 region,
151 launch_template,
152 launch_template_version,
153 ..
154 } = self
155 .config
156 .targets
157 .get(target_id)
158 .with_context(|| format!("unknown target template {target_id:?}"))?
159 else {
160 bail!("target {target_id:?} is not an AWS EC2 target");
161 };
162 let profile = aws_profile.as_deref().unwrap_or("default");
163 let launch_key = if launch_template.starts_with("lt-") {
164 "--launch-template-id"
165 } else {
166 "--launch-template-name"
167 };
168 let version = launch_template_version.as_deref().unwrap_or("$Default");
169 let describe_template = CommandSpec::new(
170 "aws",
171 [
172 "--profile",
173 profile,
174 "--region",
175 region,
176 "ec2",
177 "describe-launch-template-versions",
178 launch_key,
179 launch_template,
180 "--versions",
181 version,
182 "--output",
183 "json",
184 ],
185 )
186 .purpose("resolve EC2 launch template instance family");
187 let output = executor.execute(&describe_template)?;
188 if output.status != 0 {
189 bail!(
190 "{} failed with status {}: {}",
191 describe_template.purpose,
192 output.status,
193 String::from_utf8_lossy(&output.stderr).trim()
194 );
195 }
196 let response: serde_json::Value =
197 serde_json::from_slice(&output.stdout).context("parse EC2 launch template response")?;
198 let instance_type = response
199 .pointer("/LaunchTemplateVersions/0/LaunchTemplateData/InstanceType")
200 .and_then(serde_json::Value::as_str)
201 .context("launch template does not specify a concrete instance type")?;
202 let family = instance_type
203 .rsplit_once('.')
204 .map(|(family, _)| family)
205 .context("launch template instance type has no size suffix")?;
206 let filter = format!("Name=instance-type,Values={family}.*");
207 let describe_types = CommandSpec::new(
208 "aws",
209 [
210 "--profile",
211 profile,
212 "--region",
213 region,
214 "ec2",
215 "describe-instance-types",
216 "--filters",
217 &filter,
218 "--output",
219 "json",
220 ],
221 )
222 .purpose("discover EC2 instance sizes");
223 let output = executor.execute(&describe_types)?;
224 if output.status != 0 {
225 bail!(
226 "{} failed with status {}: {}",
227 describe_types.purpose,
228 output.status,
229 String::from_utf8_lossy(&output.stderr).trim()
230 );
231 }
232 let response: serde_json::Value =
233 serde_json::from_slice(&output.stdout).context("parse EC2 instance type response")?;
234 let mut options = response
235 .get("InstanceTypes")
236 .and_then(serde_json::Value::as_array)
237 .context("EC2 instance type response omitted InstanceTypes")?
238 .iter()
239 .filter_map(|entry| {
240 Some(SessionResourceAllocation::AwsEc2 {
241 instance_type: entry.get("InstanceType")?.as_str()?.to_owned(),
242 vcpus: entry.pointer("/VCpuInfo/DefaultVCpus")?.as_u64()?,
243 memory_bytes: entry
244 .pointer("/MemoryInfo/SizeInMiB")?
245 .as_u64()?
246 .checked_mul(1024 * 1024)?,
247 })
248 })
249 .collect::<Vec<_>>();
250 options.sort_by_key(allocation_cpus);
251 if !options.iter().any(|option| allocation_cpus(option) == 8) {
252 bail!("EC2 family {family:?} has no exact 8-vCPU baseline size");
253 }
254 Ok(options)
255 }
256
257 pub fn reconnect_command(&self, session_id: &str) -> Result<CommandSpec> {
258 let session = self
259 .state
260 .sessions
261 .get(session_id)
262 .with_context(|| format!("unknown session {session_id}"))?;
263 session.validate_configuration(&self.config)?;
264 let locator = session.target.as_ref().context("session has no target")?;
265 let backend = backend_locator(locator, session, &self.config)?;
266 targets::reconnect_plan(&backend, session_id)?
267 .commands
268 .into_iter()
269 .next()
270 .context("reconnect plan is empty")
271 }
272
273 pub(crate) fn session_owns_profile_home(&self, session_id: &str) -> Result<bool> {
280 let session = self
281 .state
282 .sessions
283 .get(session_id)
284 .with_context(|| format!("unknown session {session_id}"))?;
285 let profile = self
286 .config
287 .profiles
288 .get(&session.last_profile)
289 .with_context(|| format!("unknown profile {}", session.last_profile))?;
290 let locator = session.target.as_ref().context("session has no target")?;
291 let backend = backend_locator(locator, session, &self.config)?;
292 Ok(crate::controller::session_owns_profile_home(
293 &backend, session_id, profile,
294 ))
295 }
296
297 pub fn resource_probe(&self, session_id: &str) -> Result<targets::SessionResourceProbe> {
298 let session = self
299 .state
300 .sessions
301 .get(session_id)
302 .with_context(|| format!("unknown session {session_id}"))?;
303 let locator = session.target.as_ref().context("session has no target")?;
304 let backend = backend_locator(locator, session, &self.config)?;
305 targets::resource_probe(&backend, session_id)
306 }
307
308 pub fn deployment_capacity_targets(&self) -> Vec<targets::DeploymentCapacityTarget> {
309 use targets::{DeploymentCapacityKind, DeploymentCapacityTarget};
310
311 let mut local_ids = Vec::new();
312 let mut ssh_hosts: BTreeMap<String, (Vec<String>, Vec<CommandSpec>)> = BTreeMap::new();
313 let mut targets = Vec::new();
314 for (target_id, template) in &self.config.targets {
315 match template {
316 TargetTemplate::LocalBare
317 | TargetTemplate::LocalPodman { .. }
318 | TargetTemplate::LocalDocker { .. }
319 | TargetTemplate::AppleContainer { .. } => {
320 local_ids.push(target_id.clone());
321 }
322 TargetTemplate::SshBare { ssh, .. }
323 | TargetTemplate::SshPodman { ssh, .. }
324 | TargetTemplate::SshDocker { ssh, .. } => {
325 let entry = ssh_hosts.entry(ssh.host.clone()).or_default();
326 entry.0.push(target_id.clone());
327 let command = targets::ssh_host_capacity_command(&SshTarget::from(ssh));
328 if !entry.1.contains(&command) {
329 entry.1.push(command);
330 }
331 }
332 TargetTemplate::AwsEc2 { .. } => {
333 let mut probes = Vec::new();
334 let mut probe_error = None;
335 for session in self.state.sessions.values().filter(|session| {
336 session.target_template_id == *target_id
337 && session.state.is_active()
338 && session.target.is_some()
339 }) {
340 let result = backend_locator(
341 session.target.as_ref().expect("filtered target"),
342 session,
343 &self.config,
344 )
345 .and_then(|locator| {
346 targets::aws_allocated_capacity_command(&locator, &session.id)
347 });
348 match result {
349 Ok(command) => probes.push(command),
350 Err(error) => probe_error = Some(format!("{error:#}")),
351 }
352 }
353 targets.push(DeploymentCapacityTarget {
354 id: format!("aws:{target_id}"),
355 host: target_id.clone(),
356 target_ids: vec![target_id.clone()],
357 kind: DeploymentCapacityKind::AwsFleet,
358 local: false,
359 probes,
360 probe_error,
361 });
362 }
363 }
364 }
365 if !local_ids.is_empty() {
366 targets.push(DeploymentCapacityTarget {
367 id: "local".into(),
368 host: "local".into(),
369 target_ids: local_ids,
370 kind: DeploymentCapacityKind::Host,
371 local: true,
372 probes: Vec::new(),
373 probe_error: None,
374 });
375 }
376 targets.extend(ssh_hosts.into_iter().map(|(host, (target_ids, probes))| {
377 DeploymentCapacityTarget {
378 id: format!("ssh:{host}"),
379 host,
380 target_ids,
381 kind: DeploymentCapacityKind::Host,
382 local: false,
383 probes,
384 probe_error: None,
385 }
386 }));
387 targets.sort_by(|left, right| left.id.cmp(&right.id));
388 targets
389 }
390
391 pub fn test_target(&self, target_id: &str, executor: &impl CommandExecutor) -> Result<()> {
393 verify_target(
394 self.configured_target(target_id)?,
395 executor,
396 TargetCheck::BeforeLaunch,
397 )
398 }
399
400 pub fn check_target_readiness(
403 &self,
404 target_id: &str,
405 executor: &impl CommandExecutor,
406 ) -> Result<()> {
407 preflight_target(
408 self.configured_target(target_id)?,
409 executor,
410 TargetCheck::BeforeLaunch,
411 )
412 }
413
414 fn configured_target(&self, target_id: &str) -> Result<&TargetTemplate> {
415 self.config
416 .targets
417 .get(target_id)
418 .with_context(|| format!("unknown target template {target_id:?}"))
419 }
420}
421
422#[derive(Debug, Clone, Copy, PartialEq, Eq)]
424pub(super) enum TargetCheck {
425 BeforeLaunch,
428 Launch,
431}
432
433impl TargetCheck {
434 const fn then_retry(self) -> &'static str {
437 match self {
438 Self::BeforeLaunch => "",
439 Self::Launch => ", then Retry launch",
440 }
441 }
442}
443
444pub(super) fn preflight_target(
451 template: &TargetTemplate,
452 executor: &impl CommandExecutor,
453 check: TargetCheck,
454) -> Result<()> {
455 match template {
456 TargetTemplate::SshBare { ssh, .. }
457 | TargetTemplate::SshPodman { ssh, .. }
458 | TargetTemplate::SshDocker { ssh, .. } => {
459 verify_ssh_connectivity(&SshTarget::from(ssh), executor)
460 }
461 _ => verify_target(template, executor, check),
462 }
463}
464
465#[derive(Debug, Clone, Copy, PartialEq, Eq)]
468pub(crate) enum LocalEngineReadiness {
469 Ready,
470 NotInstalled,
472 NotReady,
474}
475
476pub(crate) fn local_engine_readiness(
481 kind: &str,
482 executor: &impl CommandExecutor,
483) -> Option<LocalEngineReadiness> {
484 let result = match kind {
485 "local-podman" => targets::verify_local_podman(executor).map(|_| ()),
486 "local-docker" => targets::verify_local_docker(executor).map(|_| ()),
487 "apple-container" => executor
490 .execute(
491 &CommandSpec::new("container", ["system", "status"])
492 .purpose("preflight Apple container runtime")
493 .stage(ProvisionStage::Provisioning),
494 )
495 .and_then(|output| {
496 ensure!(output.status == 0, "container system status failed");
497 Ok(())
498 }),
499 _ => return None,
500 };
501 Some(match result {
502 Ok(()) => LocalEngineReadiness::Ready,
503 Err(error) if is_missing_command(&error) => LocalEngineReadiness::NotInstalled,
504 Err(_) => LocalEngineReadiness::NotReady,
505 })
506}
507
508fn is_missing_command(error: &anyhow::Error) -> bool {
510 error.chain().any(|cause| {
511 cause
512 .downcast_ref::<std::io::Error>()
513 .is_some_and(|io| io.kind() == std::io::ErrorKind::NotFound)
514 })
515}
516
517pub(super) fn verify_target(
519 template: &TargetTemplate,
520 executor: &impl CommandExecutor,
521 check: TargetCheck,
522) -> Result<()> {
523 let then_retry = check.then_retry();
524 match template {
525 TargetTemplate::LocalPodman { .. } => targets::verify_local_podman(executor)
526 .map(|_| ())
527 .map_err(|error| {
528 anyhow::anyhow!(
529 "local Podman is not ready. Fix the problem below{then_retry}: {error:#}"
530 )
531 }),
532 TargetTemplate::LocalDocker { .. } => targets::verify_local_docker(executor)
533 .map(|_| ())
534 .map_err(
535 |error| match error.downcast_ref::<targets::DockerUnavailable>() {
536 Some(problem) => anyhow::anyhow!(
540 "{problem} {}",
541 match check {
542 TargetCheck::BeforeLaunch => problem.remedy(),
543 TargetCheck::Launch => problem.launch_remedy(),
544 }
545 ),
546 None => anyhow::anyhow!(
547 "local Docker is not ready. Start Docker or fix the problem below{then_retry}: {error:#}"
548 ),
549 },
550 ),
551 TargetTemplate::SshPodman { ssh, .. } => {
552 let ssh = SshTarget::from(ssh);
553 targets::verify_ssh_podman(&ssh, executor)
554 .map(|preflight| {
555 for warning in preflight.warnings {
556 executor.notify_notice(&warning.notice());
557 }
558 })
559 .map_err(|error| {
560 anyhow::anyhow!(
561 "remote Podman is not ready on {}. Fix the problem below{then_retry}: {error:#}",
562 ssh.destination
563 )
564 })
565 }
566 TargetTemplate::SshDocker { ssh, .. } => {
567 let ssh = SshTarget::from(ssh);
568 targets::verify_ssh_docker(&ssh, executor)
569 .map(|_| ())
570 .map_err(|error| {
571 anyhow::anyhow!(
572 "remote Docker preflight failed for {}. Fix the problem below{then_retry}: {error:#}",
573 ssh.destination
574 )
575 })
576 }
577 TargetTemplate::AppleContainer { .. } => {
578 let command = CommandSpec::new("container", ["system", "status"])
579 .purpose("preflight Apple container runtime")
580 .stage(ProvisionStage::Provisioning);
581 let output = executor.execute(&command).map_err(|error| {
582 anyhow::anyhow!(
583 "Apple container is not ready. Fix the problem below{then_retry}: {error}"
584 )
585 })?;
586 if output.status != 0 {
587 bail!(
588 "Apple container is not ready. Start the runtime with `container system start`{then_retry}: container system status exited {}: {}",
589 output.status,
590 [
591 String::from_utf8_lossy(&output.stdout).trim(),
592 String::from_utf8_lossy(&output.stderr).trim(),
593 ]
594 .into_iter()
595 .filter(|message| !message.is_empty())
596 .collect::<Vec<_>>()
597 .join("\n")
598 );
599 }
600 Ok(())
601 }
602 TargetTemplate::SshBare { ssh, .. } => {
603 verify_ssh_connectivity(&SshTarget::from(ssh), executor)
604 }
605 TargetTemplate::AwsEc2 {
606 aws_profile,
607 region,
608 launch_template,
609 launch_template_version,
610 ..
611 } => {
612 let mut identity_args = vec!["sts".into(), "get-caller-identity".into()];
613 if let Some(profile) = aws_profile {
614 identity_args.extend(["--profile".into(), profile.clone()]);
615 }
616 let identity = CommandSpec::new("aws", identity_args)
617 .purpose("verify AWS credentials")
618 .stage(ProvisionStage::Provisioning);
619 let output = executor.execute(&identity)?;
620 ensure!(
621 output.status == 0,
622 "AWS credential test failed with status {}: {}",
623 output.status,
624 String::from_utf8_lossy(&output.stderr).trim()
625 );
626
627 let mut launch_args = vec![
628 "ec2".into(),
629 "describe-launch-template-versions".into(),
630 "--region".into(),
631 region.clone(),
632 "--launch-template-name".into(),
633 launch_template.clone(),
634 "--versions".into(),
635 launch_template_version
636 .clone()
637 .unwrap_or_else(|| "$Default".into()),
638 ];
639 if let Some(profile) = aws_profile {
640 launch_args.extend(["--profile".into(), profile.clone()]);
641 }
642 let launch = CommandSpec::new("aws", launch_args)
643 .purpose("verify AWS launch template")
644 .stage(ProvisionStage::Provisioning);
645 let output = executor.execute(&launch)?;
646 ensure!(
647 output.status == 0,
648 "AWS launch-template test failed with status {}: {}",
649 output.status,
650 String::from_utf8_lossy(&output.stderr).trim()
651 );
652 Ok(())
653 }
654 TargetTemplate::LocalBare => Ok(()),
655 }
656}
657
658fn verify_ssh_connectivity(ssh: &SshTarget, executor: &impl CommandExecutor) -> Result<()> {
659 let output = executor.execute(&targets::ssh_connectivity_probe(ssh))?;
660 ensure!(
661 output.status == 0,
662 "SSH connectivity test failed for {} with status {}: {}",
663 ssh.destination,
664 output.status,
665 String::from_utf8_lossy(&output.stderr).trim()
666 );
667 Ok(())
668}
669
670pub(super) fn backend_bundle(
671 bundle: &ProjectBundle,
672 executor: &impl CommandExecutor,
673) -> Result<ProjectBundleSpec> {
674 let primary = bundle.primary().context("bundle primary is missing")?;
675 Ok(ProjectBundleSpec {
676 primary: primary.destination.to_string_lossy().into_owned(),
677 repositories: bundle
678 .repositories
679 .iter()
680 .map(|repository| {
681 let source = mj_core::remote_git::resolve_repository(repository, executor)
682 .with_context(|| format!("repository {:?}", repository.id))?;
683 Ok(RepositorySpec {
684 url: Some(source.fetch_url),
685 push_urls: source.push_urls,
686 destination: repository.destination.to_string_lossy().into_owned(),
687 git_ref: None,
688 reference: None,
689 })
690 })
691 .collect::<Result<Vec<_>>>()?,
692 })
693}
694
695#[derive(Debug, Clone, Copy, Default)]
699pub(super) struct ContainerOverrides<'a> {
700 pub cpus: Option<&'a str>,
701 pub memory: Option<&'a str>,
702}
703
704impl<'a> ContainerOverrides<'a> {
705 pub(super) fn for_session(session: &'a SessionRecord) -> Self {
706 Self {
707 cpus: session.container_cpus.as_deref(),
708 memory: session.container_memory.as_deref(),
709 }
710 }
711}
712
713pub(super) fn backend_target(
714 template: &TargetTemplate,
715 allocation: Option<&SessionResourceAllocation>,
716 overrides: ContainerOverrides<'_>,
717) -> Result<targets::TargetTemplate> {
718 Ok(match template {
719 TargetTemplate::LocalBare => targets::TargetTemplate::LocalBare,
720 TargetTemplate::LocalPodman { container } => {
721 let mut backend = backend_container(container, allocation, overrides);
722 backend.workspace_storage = (&container.workspace_storage).into();
723 targets::TargetTemplate::LocalPodman(backend)
724 }
725 TargetTemplate::LocalDocker { container } => targets::TargetTemplate::LocalDocker(
726 backend_container(container, allocation, overrides),
727 ),
728 TargetTemplate::AppleContainer { container } => targets::TargetTemplate::AppleContainer(
729 backend_container(container, allocation, overrides),
730 ),
731 TargetTemplate::AwsEc2 {
732 aws_profile,
733 region,
734 launch_template,
735 launch_template_version,
736 ssh_user,
737 identity_file,
738 ssh_args,
739 ..
740 } => targets::TargetTemplate::AwsEc2(AwsTemplate {
741 profile: aws_profile.clone().unwrap_or_else(|| "default".into()),
742 region: region.clone(),
743 launch_template: launch_template.clone(),
744 launch_template_version: launch_template_version.clone(),
745 instance_type: match allocation {
746 Some(SessionResourceAllocation::AwsEc2 { instance_type, .. }) => {
747 Some(instance_type.clone())
748 }
749 _ => None,
750 },
751 ssh: SshTarget {
753 destination: format!("{ssh_user}@pending.invalid"),
754 ssh_args: targets::ssh_args_with_identity(ssh_args, identity_file.as_deref()),
755 },
756 }),
757 TargetTemplate::SshBare {
758 ssh,
759 workspace_prefix,
760 ..
761 } => targets::TargetTemplate::SshBare {
762 ssh: SshTarget::from(ssh),
763 workspace_prefix: workspace_prefix.to_string_lossy().into_owned(),
764 },
765 TargetTemplate::SshPodman { ssh, container, .. } => {
766 let mut backend = backend_container(container, allocation, overrides);
767 backend.workspace_storage = (&container.workspace_storage).into();
768 targets::TargetTemplate::SshPodman {
769 ssh: SshTarget::from(ssh),
770 container: backend,
771 }
772 }
773 TargetTemplate::SshDocker { ssh, container, .. } => targets::TargetTemplate::SshDocker {
774 ssh: SshTarget::from(ssh),
775 container: backend_container(container, allocation, overrides),
776 },
777 })
778}
779
780pub fn image_refresh_plan(config: &Config) -> Vec<ImageRefresh> {
792 let mut plan: Vec<ImageRefresh> = Vec::new();
793 for target in config.targets.values() {
794 let Some((host, container)) = target.image_host() else {
795 continue;
796 };
797 let Some(refresh) = targets::image_refresh(
798 host,
799 &container.image,
800 container.platform.as_deref(),
801 container.pull_policy,
802 ) else {
803 continue;
804 };
805 if let Some(existing) = plan.iter_mut().find(|entry| {
808 entry.host == refresh.host
809 && entry.image == refresh.image
810 && entry.platform == refresh.platform
811 }) {
812 existing.when = existing.when.max(refresh.when);
813 continue;
814 }
815 plan.push(refresh);
816 }
817 plan
818}
819
820pub(crate) fn controller_github_token() -> Option<String> {
821 for name in ["GH_TOKEN", "GITHUB_TOKEN"] {
822 if let Ok(token) = std::env::var(name)
823 && let Some(token) = usable_github_token(&token)
824 {
825 return Some(token.to_owned());
826 }
827 }
828 let output = match Command::new("gh")
829 .args(["auth", "token", "--hostname", "github.com"])
830 .stdin(Stdio::null())
831 .stderr(Stdio::null())
832 .output()
833 {
834 Ok(output) => output,
835 Err(error) => {
836 tracing::debug!(%error, "could not query the GitHub CLI for a token");
837 return None;
838 }
839 };
840 if !output.status.success() {
841 tracing::debug!(status = ?output.status, "GitHub CLI did not return an authenticated token");
842 return None;
843 }
844 let token = match std::str::from_utf8(&output.stdout) {
845 Ok(token) => token,
846 Err(error) => {
847 tracing::debug!(%error, "GitHub CLI returned a non-UTF-8 token");
848 return None;
849 }
850 };
851 let Some(token) = usable_github_token(token) else {
852 tracing::debug!("GitHub CLI returned an empty or invalid token");
853 return None;
854 };
855 Some(token.to_owned())
856}
857
858fn usable_github_token(token: &str) -> Option<&str> {
859 let token = token.trim();
860 (!token.is_empty() && !token.chars().any(char::is_whitespace)).then_some(token)
861}
862
863pub(super) fn configure_github_token_environment(target: &mut targets::TargetTemplate) -> bool {
864 let container = match target {
865 targets::TargetTemplate::LocalPodman(container)
866 | targets::TargetTemplate::LocalDocker(container)
867 | targets::TargetTemplate::AppleContainer(container)
868 | targets::TargetTemplate::SshPodman { container, .. }
869 | targets::TargetTemplate::SshDocker { container, .. } => container,
870 targets::TargetTemplate::LocalBare
871 | targets::TargetTemplate::AwsEc2(_)
872 | targets::TargetTemplate::SshBare { .. } => return false,
873 };
874 container
875 .extra_run_args
876 .extend(["--env".to_owned(), "GH_TOKEN".to_owned()]);
877 true
878}
879
880pub(super) fn use_github_https_urls(bundle: &mut targets::ProjectBundleSpec) {
881 for repository in &mut bundle.repositories {
882 for source in repository
883 .url
884 .iter_mut()
885 .chain(repository.push_urls.iter_mut())
886 {
887 if let Some(github) = crate::setup::github_repository_from_origin(source) {
888 *source = format!(
889 "https://github.com/{}/{}.git",
890 github.owner, github.repository
891 );
892 }
893 }
894 }
895}
896
897fn backend_container(
898 container: &mj_core::config::ContainerTemplate,
899 allocation: Option<&SessionResourceAllocation>,
900 overrides: ContainerOverrides<'_>,
901) -> ContainerTemplate {
902 let mut extra_run_args = Vec::new();
903 if let Some(platform) = &container.platform {
904 extra_run_args.push(format!("--platform={platform}"));
905 }
906 let (cpus, memory) = match allocation {
907 Some(SessionResourceAllocation::Container { cpus, memory_bytes }) => {
908 (Some(cpus.to_string()), Some(memory_bytes.to_string()))
909 }
910 _ => (container.cpus.clone(), container.memory.clone()),
911 };
912 let cpus = overrides.cpus.map(str::to_owned).or(cpus);
914 let memory = overrides.memory.map(str::to_owned).or(memory);
915 if let Some(cpus) = cpus {
916 extra_run_args.push(format!("--cpus={cpus}"));
917 }
918 if let Some(memory) = memory {
919 extra_run_args.push(format!("--memory={memory}"));
920 }
921 for (key, value) in &container.environment {
922 extra_run_args.extend(["--env".to_string(), format!("{key}={value}")]);
923 }
924 ContainerTemplate {
925 image: container.image.clone(),
926 pull_policy: container.pull_policy,
927 extra_run_args,
928 workspace_storage: targets::PodmanWorkspaceStorage::ContainerLayer,
929 build_cache: container.build_cache.clone(),
930 }
931}
932
933pub(super) fn validate_resource_allocation(
934 template: &TargetTemplate,
935 allocation: Option<&SessionResourceAllocation>,
936) -> Result<()> {
937 if let Some(allocation) = allocation {
938 allocation.validate()?;
939 }
940 match (template, allocation) {
941 (_, None)
942 | (
943 TargetTemplate::LocalPodman { .. }
944 | TargetTemplate::LocalDocker { .. }
945 | TargetTemplate::AppleContainer { .. }
946 | TargetTemplate::SshPodman { .. }
947 | TargetTemplate::SshDocker { .. },
948 Some(SessionResourceAllocation::Container { .. }),
949 )
950 | (TargetTemplate::AwsEc2 { .. }, Some(SessionResourceAllocation::AwsEc2 { .. })) => Ok(()),
951 (TargetTemplate::LocalBare | TargetTemplate::SshBare { .. }, Some(_)) => {
952 bail!("bare targets have fixed host resources")
953 }
954 _ => bail!("resource allocation does not match the selected target kind"),
955 }
956}
957
958const AWS_SSH_READY_TIMEOUT: Duration = Duration::from_secs(300);
960
961const AWS_SSH_READY_RETRY_DELAY: Duration = Duration::from_secs(3);
962
963fn wait_for_ssh_ready(
968 executor: &impl CommandExecutor,
969 probe: &CommandSpec,
970 timeout: Duration,
971 mut now: impl FnMut() -> Instant,
972 mut sleep: impl FnMut(Duration),
973) -> Result<()> {
974 let started = now();
975 loop {
976 if executor.cancellation_requested() {
977 bail!("cancelled while waiting for SSH on the new instance");
978 }
979 let failure = match executor.execute(probe) {
980 Ok(output) if output.status == 0 => return Ok(()),
981 Ok(output) => String::from_utf8_lossy(&output.stderr).trim().to_string(),
982 Err(error) => error.to_string(),
983 };
984 if now().duration_since(started) >= timeout {
985 bail!(
986 "{} timed out after {}s: {}",
987 probe.purpose,
988 timeout.as_secs(),
989 if failure.is_empty() {
990 "the SSH probe reported no error output"
991 } else {
992 failure.as_str()
993 }
994 );
995 }
996 sleep(AWS_SSH_READY_RETRY_DELAY);
997 }
998}
999
1000pub(super) fn locator_after_provision(
1001 canonical: &TargetTemplate,
1002 backend: &targets::TargetTemplate,
1003 session_id: &str,
1004 first_output: Option<&CommandOutput>,
1005 executor: &(impl CommandExecutor + Sync),
1006) -> Result<TargetLocator> {
1007 let generated = targets::resource_name(session_id)?;
1008 Ok(match canonical {
1009 TargetTemplate::LocalBare => TargetLocator::LocalBare {
1010 worker_root: data_dir().join("workers").join(session_id),
1011 },
1012 TargetTemplate::LocalPodman { .. } => {
1013 let targets::TargetTemplate::LocalPodman(container) = backend else {
1014 bail!("session locator/template mismatch")
1015 };
1016 TargetLocator::LocalPodman {
1017 borrowed_from: None,
1018 container_id: generated,
1019 workspace_storage: PodmanWorkspaceLocator::from(targets::podman_workspace_locator(
1020 container, session_id,
1021 )?),
1022 }
1023 }
1024 TargetTemplate::LocalDocker { .. } => TargetLocator::LocalDocker {
1025 borrowed_from: None,
1026 container_id: generated,
1027 },
1028 TargetTemplate::AppleContainer { .. } => TargetLocator::AppleContainer {
1029 borrowed_from: None,
1030 container_id: generated,
1031 },
1032 TargetTemplate::SshBare { ssh, .. } => TargetLocator::SshBare {
1033 host: ssh.host.clone(),
1034 workspace: PathBuf::from(targets::workspace_for(backend, session_id)?),
1035 worker_id: None,
1036 },
1037 TargetTemplate::SshPodman { ssh, .. } => {
1038 let targets::TargetTemplate::SshPodman { container, .. } = backend else {
1039 bail!("session locator/template mismatch")
1040 };
1041 TargetLocator::SshPodman {
1042 borrowed_from: None,
1043 host: ssh.host.clone(),
1044 container_id: generated,
1045 workspace_storage: PodmanWorkspaceLocator::from(targets::podman_workspace_locator(
1046 container, session_id,
1047 )?),
1048 }
1049 }
1050 TargetTemplate::SshDocker { ssh, .. } => TargetLocator::SshDocker {
1051 borrowed_from: None,
1052 host: ssh.host.clone(),
1053 container_id: generated,
1054 },
1055 TargetTemplate::AwsEc2 {
1056 aws_profile,
1057 region,
1058 ssh_user,
1059 address_source,
1060 identity_file,
1061 ssh_args,
1062 ..
1063 } => {
1064 let output = first_output.context("AWS launch produced no output")?;
1065 let json: serde_json::Value = serde_json::from_slice(&output.stdout)
1066 .context("parse aws ec2 run-instances response")?;
1067 let instance_id = json
1068 .pointer("/Instances/0/InstanceId")
1069 .and_then(serde_json::Value::as_str)
1070 .context("AWS response omitted instance ID")?
1071 .to_string();
1072 let profile = aws_profile.clone().unwrap_or_else(|| "default".into());
1073 execute_checked(
1074 executor,
1075 CommandSpec::new(
1076 "aws",
1077 [
1078 "--profile".into(),
1079 profile.clone(),
1080 "--region".into(),
1081 region.clone(),
1082 "ec2".into(),
1083 "wait".into(),
1084 "instance-running".into(),
1085 "--instance-ids".into(),
1086 instance_id.clone(),
1087 ],
1088 )
1089 .purpose("wait for EC2 session instance to run")
1090 .stage(ProvisionStage::Booting),
1091 )?;
1092 let field = match address_source {
1093 AwsAddressSource::PublicDns => "PublicDnsName",
1094 AwsAddressSource::PublicIp => "PublicIpAddress",
1095 AwsAddressSource::PrivateDns => "PrivateDnsName",
1096 AwsAddressSource::PrivateIp => "PrivateIpAddress",
1097 };
1098 let address = execute_checked(
1099 executor,
1100 CommandSpec::new(
1101 "aws",
1102 [
1103 "--profile".into(),
1104 profile.clone(),
1105 "--region".into(),
1106 region.clone(),
1107 "ec2".into(),
1108 "describe-instances".into(),
1109 "--instance-ids".into(),
1110 instance_id.clone(),
1111 "--query".into(),
1112 format!("Reservations[0].Instances[0].{field}"),
1113 "--output".into(),
1114 "text".into(),
1115 ],
1116 )
1117 .purpose("resolve EC2 session address")
1118 .stage(ProvisionStage::Booting),
1119 )?;
1120 let address = String::from_utf8(address.stdout)
1121 .context("AWS address was not UTF-8")?
1122 .trim()
1123 .to_string();
1124 if address.is_empty() || address == "None" {
1125 bail!("AWS instance {instance_id} has no configured address");
1126 }
1127 let ssh = SshTarget {
1128 destination: format!("{ssh_user}@{address}"),
1129 ssh_args: targets::ssh_args_with_identity(ssh_args, identity_file.as_deref()),
1130 };
1131 wait_for_ssh_ready(
1132 executor,
1133 &crate::targets::ssh_command(&ssh, ["true"])
1134 .purpose("wait for EC2 SSH availability")
1135 .stage(ProvisionStage::Booting),
1136 AWS_SSH_READY_TIMEOUT,
1137 Instant::now,
1138 std::thread::sleep,
1139 )?;
1140 TargetLocator::AwsEc2 {
1141 instance_id,
1142 address: Some(address),
1143 }
1144 }
1145 })
1146}
1147
1148pub(super) fn backend_locator(
1153 locator: &TargetLocator,
1154 session: &SessionRecord,
1155 config: &Config,
1156) -> Result<targets::TargetLocator> {
1157 let runtime = if session.target_runtime.is_some() || targets::locator_needs_connection(locator)
1158 {
1159 Some(session.target_runtime_settings(config)?)
1160 } else {
1161 None
1162 };
1163 Ok(targets::TargetLocator::try_from(targets::RecordedTarget {
1164 locator,
1165 runtime: runtime.as_deref(),
1166 session_id: &session.id,
1167 })?)
1168}
1169
1170#[cfg(test)]
1171mod tests;