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 fn deployment_capacity_targets(&self) -> Vec<targets::DeploymentCapacityTarget> {
274 use targets::{DeploymentCapacityKind, DeploymentCapacityTarget};
275
276 let mut local_ids = Vec::new();
277 let mut ssh_hosts: BTreeMap<String, (Vec<String>, Vec<CommandSpec>)> = BTreeMap::new();
278 let mut targets = Vec::new();
279 for (target_id, template) in &self.config.targets {
280 match template {
281 TargetTemplate::LocalBare
282 | TargetTemplate::LocalPodman { .. }
283 | TargetTemplate::LocalDocker { .. }
284 | TargetTemplate::AppleContainer { .. } => {
285 local_ids.push(target_id.clone());
286 }
287 TargetTemplate::SshBare { ssh, .. }
288 | TargetTemplate::SshPodman { ssh, .. }
289 | TargetTemplate::SshDocker { ssh, .. } => {
290 let entry = ssh_hosts.entry(ssh.host.clone()).or_default();
291 entry.0.push(target_id.clone());
292 let command = targets::ssh_host_capacity_command(&SshTarget::from(ssh));
293 if !entry.1.contains(&command) {
294 entry.1.push(command);
295 }
296 }
297 TargetTemplate::AwsEc2 { .. } => {
298 let mut probes = Vec::new();
299 let mut probe_error = None;
300 for session in self.state.sessions.values().filter(|session| {
301 session.target_template_id == *target_id
302 && session.state.is_active()
303 && session.target.is_some()
304 }) {
305 let result = backend_locator(
306 session.target.as_ref().expect("filtered target"),
307 session,
308 &self.config,
309 )
310 .and_then(|locator| {
311 targets::aws_allocated_capacity_command(&locator, &session.id)
312 });
313 match result {
314 Ok(command) => probes.push(command),
315 Err(error) => probe_error = Some(format!("{error:#}")),
316 }
317 }
318 targets.push(DeploymentCapacityTarget {
319 id: format!("aws:{target_id}"),
320 host: target_id.clone(),
321 target_ids: vec![target_id.clone()],
322 kind: DeploymentCapacityKind::AwsFleet,
323 local: false,
324 probes,
325 probe_error,
326 });
327 }
328 }
329 }
330 if !local_ids.is_empty() {
331 targets.push(DeploymentCapacityTarget {
332 id: "local".into(),
333 host: "local".into(),
334 target_ids: local_ids,
335 kind: DeploymentCapacityKind::Host,
336 local: true,
337 probes: Vec::new(),
338 probe_error: None,
339 });
340 }
341 targets.extend(ssh_hosts.into_iter().map(|(host, (target_ids, probes))| {
342 DeploymentCapacityTarget {
343 id: format!("ssh:{host}"),
344 host,
345 target_ids,
346 kind: DeploymentCapacityKind::Host,
347 local: false,
348 probes,
349 probe_error: None,
350 }
351 }));
352 targets.sort_by(|left, right| left.id.cmp(&right.id));
353 targets
354 }
355
356 pub fn test_target(&self, target_id: &str, executor: &impl CommandExecutor) -> Result<()> {
358 verify_target(
359 self.configured_target(target_id)?,
360 executor,
361 TargetCheck::BeforeLaunch,
362 )
363 }
364
365 pub fn check_target_readiness(
368 &self,
369 target_id: &str,
370 executor: &impl CommandExecutor,
371 ) -> Result<()> {
372 preflight_target(
373 self.configured_target(target_id)?,
374 executor,
375 TargetCheck::BeforeLaunch,
376 )
377 }
378
379 fn configured_target(&self, target_id: &str) -> Result<&TargetTemplate> {
380 self.config
381 .targets
382 .get(target_id)
383 .with_context(|| format!("unknown target template {target_id:?}"))
384 }
385}
386
387#[derive(Debug, Clone, Copy, PartialEq, Eq)]
389pub(super) enum TargetCheck {
390 BeforeLaunch,
393 Launch,
396}
397
398impl TargetCheck {
399 const fn then_retry(self) -> &'static str {
402 match self {
403 Self::BeforeLaunch => "",
404 Self::Launch => ", then Retry launch",
405 }
406 }
407}
408
409pub(super) fn preflight_target(
416 template: &TargetTemplate,
417 executor: &impl CommandExecutor,
418 check: TargetCheck,
419) -> Result<()> {
420 match template {
421 TargetTemplate::SshBare { ssh, .. }
422 | TargetTemplate::SshPodman { ssh, .. }
423 | TargetTemplate::SshDocker { ssh, .. } => {
424 verify_ssh_connectivity(&SshTarget::from(ssh), executor)
425 }
426 _ => verify_target(template, executor, check),
427 }
428}
429
430#[derive(Debug, Clone, Copy, PartialEq, Eq)]
433pub(crate) enum LocalEngineReadiness {
434 Ready,
435 NotInstalled,
437 NotReady,
439}
440
441pub(crate) fn local_engine_readiness(
446 kind: &str,
447 executor: &impl CommandExecutor,
448) -> Option<LocalEngineReadiness> {
449 let result = match kind {
450 "local-podman" => targets::verify_local_podman(executor).map(|_| ()),
451 "local-docker" => targets::verify_local_docker(executor).map(|_| ()),
452 "apple-container" => executor
455 .execute(
456 &CommandSpec::new("container", ["system", "status"])
457 .purpose("preflight Apple container runtime")
458 .stage(ProvisionStage::Provisioning),
459 )
460 .and_then(|output| {
461 ensure!(output.status == 0, "container system status failed");
462 Ok(())
463 }),
464 _ => return None,
465 };
466 Some(match result {
467 Ok(()) => LocalEngineReadiness::Ready,
468 Err(error) if is_missing_command(&error) => LocalEngineReadiness::NotInstalled,
469 Err(_) => LocalEngineReadiness::NotReady,
470 })
471}
472
473fn is_missing_command(error: &anyhow::Error) -> bool {
475 error.chain().any(|cause| {
476 cause
477 .downcast_ref::<std::io::Error>()
478 .is_some_and(|io| io.kind() == std::io::ErrorKind::NotFound)
479 })
480}
481
482pub(super) fn verify_target(
484 template: &TargetTemplate,
485 executor: &impl CommandExecutor,
486 check: TargetCheck,
487) -> Result<()> {
488 let then_retry = check.then_retry();
489 match template {
490 TargetTemplate::LocalPodman { .. } => targets::verify_local_podman(executor)
491 .map(|_| ())
492 .map_err(|error| {
493 anyhow::anyhow!(
494 "local Podman is not ready. Fix the problem below{then_retry}: {error:#}"
495 )
496 }),
497 TargetTemplate::LocalDocker { .. } => targets::verify_local_docker(executor)
498 .map(|_| ())
499 .map_err(
500 |error| match error.downcast_ref::<targets::DockerUnavailable>() {
501 Some(problem) => anyhow::anyhow!(
505 "{problem} {}",
506 match check {
507 TargetCheck::BeforeLaunch => problem.remedy(),
508 TargetCheck::Launch => problem.launch_remedy(),
509 }
510 ),
511 None => anyhow::anyhow!(
512 "local Docker is not ready. Start Docker or fix the problem below{then_retry}: {error:#}"
513 ),
514 },
515 ),
516 TargetTemplate::SshPodman { ssh, .. } => {
517 let ssh = SshTarget::from(ssh);
518 targets::verify_ssh_podman(&ssh, executor)
519 .map(|preflight| {
520 for warning in preflight.warnings {
521 executor.notify_notice(&warning.notice());
522 }
523 })
524 .map_err(|error| {
525 anyhow::anyhow!(
526 "remote Podman is not ready on {}. Fix the problem below{then_retry}: {error:#}",
527 ssh.destination
528 )
529 })
530 }
531 TargetTemplate::SshDocker { ssh, .. } => {
532 let ssh = SshTarget::from(ssh);
533 targets::verify_ssh_docker(&ssh, executor)
534 .map(|_| ())
535 .map_err(|error| {
536 anyhow::anyhow!(
537 "remote Docker preflight failed for {}. Fix the problem below{then_retry}: {error:#}",
538 ssh.destination
539 )
540 })
541 }
542 TargetTemplate::AppleContainer { .. } => {
543 let command = CommandSpec::new("container", ["system", "status"])
544 .purpose("preflight Apple container runtime")
545 .stage(ProvisionStage::Provisioning);
546 let output = executor.execute(&command).map_err(|error| {
547 anyhow::anyhow!(
548 "Apple container is not ready. Fix the problem below{then_retry}: {error}"
549 )
550 })?;
551 if output.status != 0 {
552 bail!(
553 "Apple container is not ready. Start the runtime with `container system start`{then_retry}: container system status exited {}: {}",
554 output.status,
555 [
556 String::from_utf8_lossy(&output.stdout).trim(),
557 String::from_utf8_lossy(&output.stderr).trim(),
558 ]
559 .into_iter()
560 .filter(|message| !message.is_empty())
561 .collect::<Vec<_>>()
562 .join("\n")
563 );
564 }
565 Ok(())
566 }
567 TargetTemplate::SshBare { ssh, .. } => {
568 verify_ssh_connectivity(&SshTarget::from(ssh), executor)
569 }
570 TargetTemplate::AwsEc2 {
571 aws_profile,
572 region,
573 launch_template,
574 launch_template_version,
575 ..
576 } => {
577 let mut identity_args = vec!["sts".into(), "get-caller-identity".into()];
578 if let Some(profile) = aws_profile {
579 identity_args.extend(["--profile".into(), profile.clone()]);
580 }
581 let identity = CommandSpec::new("aws", identity_args)
582 .purpose("verify AWS credentials")
583 .stage(ProvisionStage::Provisioning);
584 let output = executor.execute(&identity)?;
585 ensure!(
586 output.status == 0,
587 "AWS credential test failed with status {}: {}",
588 output.status,
589 String::from_utf8_lossy(&output.stderr).trim()
590 );
591
592 let mut launch_args = vec![
593 "ec2".into(),
594 "describe-launch-template-versions".into(),
595 "--region".into(),
596 region.clone(),
597 if launch_template.starts_with("lt-") { "--launch-template-id" } else { "--launch-template-name" }.into(),
598 launch_template.clone(),
599 "--versions".into(),
600 launch_template_version
601 .clone()
602 .unwrap_or_else(|| "$Default".into()),
603 ];
604 if let Some(profile) = aws_profile {
605 launch_args.extend(["--profile".into(), profile.clone()]);
606 }
607 let launch = CommandSpec::new("aws", launch_args)
608 .purpose("verify AWS launch template")
609 .stage(ProvisionStage::Provisioning);
610 let output = executor.execute(&launch)?;
611 ensure!(
612 output.status == 0,
613 "AWS launch-template test failed with status {}: {}",
614 output.status,
615 String::from_utf8_lossy(&output.stderr).trim()
616 );
617 Ok(())
618 }
619 TargetTemplate::LocalBare => Ok(()),
620 }
621}
622
623fn verify_ssh_connectivity(ssh: &SshTarget, executor: &impl CommandExecutor) -> Result<()> {
624 let output = executor.execute(&targets::ssh_connectivity_probe(ssh))?;
625 ensure!(
626 output.status == 0,
627 "SSH connectivity test failed for {} with status {}: {}",
628 ssh.destination,
629 output.status,
630 String::from_utf8_lossy(&output.stderr).trim()
631 );
632 Ok(())
633}
634
635#[cfg(test)]
636pub(super) fn backend_bundle(
637 bundle: &ProjectBundle,
638 executor: &impl CommandExecutor,
639) -> Result<ProjectBundleSpec> {
640 backend_bundle_with_sources(bundle, None, executor)
641}
642
643pub(super) fn backend_session_bundle(
644 session: &SessionRecord,
645 config: &mj_core::config::Config,
646 executor: &impl CommandExecutor,
647) -> Result<ProjectBundleSpec> {
648 backend_bundle_with_sources(
649 session
650 .project_bundle(config)
651 .context("session bundle is missing")?,
652 session
653 .project
654 .as_ref()
655 .map(|project| &project.network_sources),
656 executor,
657 )
658}
659
660fn backend_bundle_with_sources(
661 bundle: &ProjectBundle,
662 sources: Option<&std::collections::BTreeMap<String, mj_core::remote_git::NetworkGitSource>>,
663 executor: &impl CommandExecutor,
664) -> Result<ProjectBundleSpec> {
665 let primary = bundle.primary().context("bundle primary is missing")?;
666 Ok(ProjectBundleSpec {
667 primary: primary.destination.to_string_lossy().into_owned(),
668 repositories: bundle
669 .repositories
670 .iter()
671 .map(|repository| {
672 let source = match sources.and_then(|sources| sources.get(&repository.id)) {
673 Some(source) => source.clone(),
674 None => mj_core::remote_git::resolve_repository(repository, executor)
675 .with_context(|| format!("repository {:?}", repository.id))?,
676 };
677 Ok(RepositorySpec {
678 url: Some(source.fetch_url),
679 push_urls: source.push_urls,
680 destination: repository.destination.to_string_lossy().into_owned(),
681 git_ref: None,
682 reference: None,
683 })
684 })
685 .collect::<Result<Vec<_>>>()?,
686 })
687}
688
689#[derive(Debug, Clone, Copy, Default)]
693pub(super) struct ContainerOverrides<'a> {
694 pub cpus: Option<&'a str>,
695 pub memory: Option<&'a str>,
696}
697
698impl<'a> ContainerOverrides<'a> {
699 pub(super) fn for_session(session: &'a SessionRecord) -> Self {
700 Self {
701 cpus: session.container_cpus.as_deref(),
702 memory: session.container_memory.as_deref(),
703 }
704 }
705}
706
707pub(super) fn backend_target(
708 template: &TargetTemplate,
709 allocation: Option<&SessionResourceAllocation>,
710 overrides: ContainerOverrides<'_>,
711) -> Result<targets::TargetTemplate> {
712 Ok(match template {
713 TargetTemplate::LocalBare => targets::TargetTemplate::LocalBare,
714 TargetTemplate::LocalPodman { container } => {
715 let mut backend = backend_container(container, allocation, overrides);
716 backend.workspace_storage = (&container.workspace_storage).into();
717 targets::TargetTemplate::LocalPodman(backend)
718 }
719 TargetTemplate::LocalDocker { container } => targets::TargetTemplate::LocalDocker(
720 backend_container(container, allocation, overrides),
721 ),
722 TargetTemplate::AppleContainer { container } => targets::TargetTemplate::AppleContainer(
723 backend_container(container, allocation, overrides),
724 ),
725 TargetTemplate::AwsEc2 {
726 aws_profile,
727 region,
728 launch_template,
729 launch_template_version,
730 ssh_user,
731 identity_file,
732 ssh_args,
733 ..
734 } => targets::TargetTemplate::AwsEc2(AwsTemplate {
735 profile: aws_profile.clone().unwrap_or_else(|| "default".into()),
736 region: region.clone(),
737 launch_template: launch_template.clone(),
738 launch_template_version: launch_template_version.clone(),
739 instance_type: match allocation {
740 Some(SessionResourceAllocation::AwsEc2 { instance_type, .. }) => {
741 Some(instance_type.clone())
742 }
743 _ => None,
744 },
745 ssh: SshTarget {
747 destination: format!("{ssh_user}@pending.invalid"),
748 ssh_args: targets::ssh_args_with_identity(ssh_args, identity_file.as_deref()),
749 },
750 }),
751 TargetTemplate::SshBare {
752 ssh,
753 workspace_prefix,
754 ..
755 } => targets::TargetTemplate::SshBare {
756 ssh: SshTarget::from(ssh),
757 workspace_prefix: workspace_prefix.to_string_lossy().into_owned(),
758 },
759 TargetTemplate::SshPodman { ssh, container, .. } => {
760 let mut backend = backend_container(container, allocation, overrides);
761 backend.workspace_storage = (&container.workspace_storage).into();
762 targets::TargetTemplate::SshPodman {
763 ssh: SshTarget::from(ssh),
764 container: backend,
765 }
766 }
767 TargetTemplate::SshDocker { ssh, container, .. } => targets::TargetTemplate::SshDocker {
768 ssh: SshTarget::from(ssh),
769 container: backend_container(container, allocation, overrides),
770 },
771 })
772}
773
774pub fn image_refresh_plan(config: &Config) -> Vec<ImageRefresh> {
786 let mut plan: Vec<ImageRefresh> = Vec::new();
787 for target in config.targets.values() {
788 let Some((host, container)) = target.image_host() else {
789 continue;
790 };
791 let Some(refresh) = targets::image_refresh(
792 host,
793 &container.image,
794 container.platform.as_deref(),
795 container.pull_policy,
796 ) else {
797 continue;
798 };
799 if let Some(existing) = plan.iter_mut().find(|entry| {
802 entry.host == refresh.host
803 && entry.image == refresh.image
804 && entry.platform == refresh.platform
805 }) {
806 existing.when = existing.when.max(refresh.when);
807 continue;
808 }
809 plan.push(refresh);
810 }
811 plan
812}
813
814pub(crate) fn controller_github_token() -> Option<String> {
815 for name in ["GH_TOKEN", "GITHUB_TOKEN"] {
816 if let Ok(token) = std::env::var(name)
817 && let Some(token) = usable_github_token(&token)
818 {
819 return Some(token.to_owned());
820 }
821 }
822 let output = match Command::new("gh")
823 .args(["auth", "token", "--hostname", "github.com"])
824 .stdin(Stdio::null())
825 .stderr(Stdio::null())
826 .output()
827 {
828 Ok(output) => output,
829 Err(error) => {
830 tracing::debug!(%error, "could not query the GitHub CLI for a token");
831 return None;
832 }
833 };
834 if !output.status.success() {
835 tracing::debug!(status = ?output.status, "GitHub CLI did not return an authenticated token");
836 return None;
837 }
838 let token = match std::str::from_utf8(&output.stdout) {
839 Ok(token) => token,
840 Err(error) => {
841 tracing::debug!(%error, "GitHub CLI returned a non-UTF-8 token");
842 return None;
843 }
844 };
845 let Some(token) = usable_github_token(token) else {
846 tracing::debug!("GitHub CLI returned an empty or invalid token");
847 return None;
848 };
849 Some(token.to_owned())
850}
851
852fn usable_github_token(token: &str) -> Option<&str> {
853 let token = token.trim();
854 (!token.is_empty() && !token.chars().any(char::is_whitespace)).then_some(token)
855}
856
857pub(super) fn configure_github_token_environment(target: &mut targets::TargetTemplate) -> bool {
858 let container = match target {
859 targets::TargetTemplate::LocalPodman(container)
860 | targets::TargetTemplate::LocalDocker(container)
861 | targets::TargetTemplate::AppleContainer(container)
862 | targets::TargetTemplate::SshPodman { container, .. }
863 | targets::TargetTemplate::SshDocker { container, .. } => container,
864 targets::TargetTemplate::LocalBare
865 | targets::TargetTemplate::AwsEc2(_)
866 | targets::TargetTemplate::SshBare { .. } => return false,
867 };
868 container
869 .extra_run_args
870 .extend(["--env".to_owned(), "GH_TOKEN".to_owned()]);
871 true
872}
873
874pub(super) fn use_github_https_urls(bundle: &mut targets::ProjectBundleSpec) {
875 for repository in &mut bundle.repositories {
876 for source in repository
877 .url
878 .iter_mut()
879 .chain(repository.push_urls.iter_mut())
880 {
881 if let Some(github) = crate::setup::github_repository_from_origin(source) {
882 *source = format!(
883 "https://github.com/{}/{}.git",
884 github.owner, github.repository
885 );
886 }
887 }
888 }
889}
890
891fn backend_container(
892 container: &mj_core::config::ContainerTemplate,
893 allocation: Option<&SessionResourceAllocation>,
894 overrides: ContainerOverrides<'_>,
895) -> ContainerTemplate {
896 let mut extra_run_args = Vec::new();
897 if let Some(platform) = &container.platform {
898 extra_run_args.push(format!("--platform={platform}"));
899 }
900 let (cpus, memory) = match allocation {
901 Some(SessionResourceAllocation::Container { cpus, memory_bytes }) => {
902 (Some(cpus.to_string()), Some(memory_bytes.to_string()))
903 }
904 _ => (container.cpus.clone(), container.memory.clone()),
905 };
906 let cpus = overrides.cpus.map(str::to_owned).or(cpus);
908 let memory = overrides.memory.map(str::to_owned).or(memory);
909 if let Some(cpus) = cpus {
910 extra_run_args.push(format!("--cpus={cpus}"));
911 }
912 if let Some(memory) = memory {
913 extra_run_args.push(format!("--memory={memory}"));
914 }
915 for (key, value) in &container.environment {
916 extra_run_args.extend(["--env".to_string(), format!("{key}={value}")]);
917 }
918 ContainerTemplate {
919 image: container.image.clone(),
920 pull_policy: container.pull_policy,
921 extra_run_args,
922 workspace_storage: targets::PodmanWorkspaceStorage::ContainerLayer,
923 build_cache: container.build_cache.clone(),
924 }
925}
926
927pub(super) fn validate_resource_allocation(
928 template: &TargetTemplate,
929 allocation: Option<&SessionResourceAllocation>,
930) -> Result<()> {
931 if let Some(allocation) = allocation {
932 allocation.validate()?;
933 }
934 match (template, allocation) {
935 (_, None)
936 | (
937 TargetTemplate::LocalPodman { .. }
938 | TargetTemplate::LocalDocker { .. }
939 | TargetTemplate::AppleContainer { .. }
940 | TargetTemplate::SshPodman { .. }
941 | TargetTemplate::SshDocker { .. },
942 Some(SessionResourceAllocation::Container { .. }),
943 )
944 | (TargetTemplate::AwsEc2 { .. }, Some(SessionResourceAllocation::AwsEc2 { .. })) => Ok(()),
945 (TargetTemplate::LocalBare | TargetTemplate::SshBare { .. }, Some(_)) => {
946 bail!(mj_core::state::BARE_TARGET_FIXED_RESOURCES)
947 }
948 _ => bail!("resource allocation does not match the selected target kind"),
949 }
950}
951
952const AWS_SSH_READY_TIMEOUT: Duration = Duration::from_secs(300);
954
955const AWS_SSH_READY_RETRY_DELAY: Duration = Duration::from_secs(3);
956
957fn wait_for_ssh_ready(
962 executor: &impl CommandExecutor,
963 probe: &CommandSpec,
964 timeout: Duration,
965 mut now: impl FnMut() -> Instant,
966 mut sleep: impl FnMut(Duration),
967) -> Result<()> {
968 let started = now();
969 loop {
970 if executor.cancellation_requested() {
971 bail!("cancelled while waiting for SSH on the new instance");
972 }
973 let failure = match executor.execute(probe) {
974 Ok(output) if output.status == 0 => return Ok(()),
975 Ok(output) => String::from_utf8_lossy(&output.stderr).trim().to_string(),
976 Err(error) => error.to_string(),
977 };
978 if now().duration_since(started) >= timeout {
979 bail!(
980 "{} timed out after {}s: {}",
981 probe.purpose,
982 timeout.as_secs(),
983 if failure.is_empty() {
984 "the SSH probe reported no error output"
985 } else {
986 failure.as_str()
987 }
988 );
989 }
990 sleep(AWS_SSH_READY_RETRY_DELAY);
991 }
992}
993
994pub(super) fn locator_after_provision_named(
995 canonical: &TargetTemplate,
996 backend: &targets::TargetTemplate,
997 session_id: &str,
998 first_output: Option<&CommandOutput>,
999 executor: &(impl CommandExecutor + Sync),
1000 name: &str,
1001) -> Result<TargetLocator> {
1002 let generated = name.to_owned();
1003
1004 Ok(match canonical {
1005 TargetTemplate::LocalBare => TargetLocator::LocalBare {
1006 worker_root: data_dir().join("workers").join(session_id),
1007 },
1008 TargetTemplate::LocalPodman { .. } => {
1009 let targets::TargetTemplate::LocalPodman(container) = backend else {
1010 bail!("session locator/template mismatch")
1011 };
1012 TargetLocator::LocalPodman {
1013 borrowed_from: None,
1014 container_id: generated,
1015 workspace_storage: PodmanWorkspaceLocator::from(
1016 targets::podman_workspace_locator_named(container, name)?,
1017 ),
1018 }
1019 }
1020 TargetTemplate::LocalDocker { .. } => TargetLocator::LocalDocker {
1021 borrowed_from: None,
1022 container_id: generated,
1023 },
1024 TargetTemplate::AppleContainer { .. } => TargetLocator::AppleContainer {
1025 borrowed_from: None,
1026 container_id: generated,
1027 },
1028 TargetTemplate::SshBare { ssh, .. } => TargetLocator::SshBare {
1029 host: ssh.host.clone(),
1030 workspace: PathBuf::from(targets::workspace_for(backend, session_id)?),
1031 worker_id: None,
1032 },
1033 TargetTemplate::SshPodman { ssh, .. } => {
1034 let targets::TargetTemplate::SshPodman { container, .. } = backend else {
1035 bail!("session locator/template mismatch")
1036 };
1037 TargetLocator::SshPodman {
1038 borrowed_from: None,
1039 host: ssh.host.clone(),
1040 container_id: generated,
1041 workspace_storage: PodmanWorkspaceLocator::from(
1042 targets::podman_workspace_locator_named(container, name)?,
1043 ),
1044 }
1045 }
1046 TargetTemplate::SshDocker { ssh, .. } => TargetLocator::SshDocker {
1047 borrowed_from: None,
1048 host: ssh.host.clone(),
1049 container_id: generated,
1050 },
1051 TargetTemplate::AwsEc2 { .. } => {
1052 let output = first_output.context("AWS launch produced no output")?;
1053 let json: serde_json::Value = serde_json::from_slice(&output.stdout)
1054 .context("parse aws ec2 run-instances response")?;
1055 let instance_id = json
1056 .pointer("/Instances/0/InstanceId")
1057 .and_then(serde_json::Value::as_str)
1058 .context("AWS response omitted instance ID")?
1059 .to_string();
1060 return ec2_locator_after_launch(canonical, instance_id, executor);
1061 }
1062 })
1063}
1064
1065pub(super) fn ec2_locator_after_launch(
1067 canonical: &TargetTemplate,
1068 instance_id: String,
1069 executor: &(impl CommandExecutor + Sync),
1070) -> Result<TargetLocator> {
1071 let TargetTemplate::AwsEc2 {
1072 aws_profile,
1073 region,
1074 ssh_user,
1075 address_source,
1076 identity_file,
1077 ssh_args,
1078 ..
1079 } = canonical
1080 else {
1081 bail!("EC2 locator requires an EC2 target");
1082 };
1083 let profile = aws_profile.clone().unwrap_or_else(|| "default".into());
1084 execute_checked(
1085 executor,
1086 CommandSpec::new(
1087 "aws",
1088 [
1089 "--profile".into(),
1090 profile.clone(),
1091 "--region".into(),
1092 region.clone(),
1093 "ec2".into(),
1094 "wait".into(),
1095 "instance-running".into(),
1096 "--instance-ids".into(),
1097 instance_id.clone(),
1098 ],
1099 )
1100 .purpose("wait for EC2 session instance to run")
1101 .stage(ProvisionStage::Booting),
1102 )?;
1103 let field = match address_source {
1104 AwsAddressSource::PublicDns => "PublicDnsName",
1105 AwsAddressSource::PublicIp => "PublicIpAddress",
1106 AwsAddressSource::PrivateDns => "PrivateDnsName",
1107 AwsAddressSource::PrivateIp => "PrivateIpAddress",
1108 };
1109 let address = execute_checked(
1110 executor,
1111 CommandSpec::new(
1112 "aws",
1113 [
1114 "--profile".into(),
1115 profile.clone(),
1116 "--region".into(),
1117 region.clone(),
1118 "ec2".into(),
1119 "describe-instances".into(),
1120 "--instance-ids".into(),
1121 instance_id.clone(),
1122 "--query".into(),
1123 format!("Reservations[0].Instances[0].{field}"),
1124 "--output".into(),
1125 "text".into(),
1126 ],
1127 )
1128 .purpose("resolve EC2 session address")
1129 .stage(ProvisionStage::Booting),
1130 )?;
1131 let address = String::from_utf8(address.stdout)
1132 .context("AWS address was not UTF-8")?
1133 .trim()
1134 .to_string();
1135 if address.is_empty() || address == "None" {
1136 bail!("AWS instance {instance_id} has no configured address");
1137 }
1138 let ssh = SshTarget {
1139 destination: format!("{ssh_user}@{address}"),
1140 ssh_args: targets::ssh_args_with_identity(ssh_args, identity_file.as_deref()),
1141 };
1142 wait_for_ssh_ready(
1143 executor,
1144 &crate::targets::ssh_command(&ssh, ["true"])
1145 .purpose("wait for EC2 SSH availability")
1146 .stage(ProvisionStage::Booting),
1147 AWS_SSH_READY_TIMEOUT,
1148 Instant::now,
1149 std::thread::sleep,
1150 )?;
1151 Ok(TargetLocator::AwsEc2 {
1152 instance_id,
1153 address: Some(address),
1154 })
1155}
1156
1157pub(super) fn backend_locator(
1162 locator: &TargetLocator,
1163 session: &SessionRecord,
1164 config: &Config,
1165) -> Result<targets::TargetLocator> {
1166 let runtime = if session.target_runtime.is_some() || targets::locator_needs_connection(locator)
1167 {
1168 Some(session.target_runtime_settings(config)?)
1169 } else {
1170 None
1171 };
1172 Ok(targets::TargetLocator::try_from(targets::RecordedTarget {
1173 locator,
1174 runtime: runtime.as_deref(),
1175 session_id: &session.id,
1176 })?)
1177}
1178
1179#[cfg(test)]
1180mod tests;