Skip to main content

mj_controller/controller/
backend.rs

1//! Backend target, locator, and capacity conversion for provisioned sessions.
2
3use 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    /// Inspect the actual execution checkout in a background worker.
25    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                // A plain folder opened through `mj go` has no branch. Say so
65                // briefly rather than echoing multi-line git stderr, which wraps
66                // the banner.
67                "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    /// The session checkout's branch, distance from upstream, and changed
79    /// files, read with the target's own `git`. Builds on
80    /// [`Self::session_working_context`], so it works wherever that does: a
81    /// local checkout, a container, or an SSH host.
82    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        // No upstream is an ordinary state, so a failing count is `None`
113        // rather than an error.
114        let ahead_behind = run(
115            &["rev-list", "--left-right", "--count", "@{upstream}...HEAD"],
116            "count commits against upstream",
117        )?;
118        // A repository with no commit yet has no HEAD to diff against.
119        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    /// The full check behind the Targets pane's Test action.
357    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    /// The check a session wizard runs before it offers a target: the same
366    /// one a launch runs, worded for a launch that has not happened yet.
367    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/// When a target is checked, which decides how its refusal ends.
388#[derive(Debug, Clone, Copy, PartialEq, Eq)]
389pub(super) enum TargetCheck {
390    /// Before anything is launched: the session wizard's target row, the
391    /// Targets pane.
392    BeforeLaunch,
393    /// Inside a launch, resume, or move, whose failure the launch-failure
394    /// dialog reports with its Retry launch button.
395    Launch,
396}
397
398impl TargetCheck {
399    /// The clause that points at the failure dialog's Retry launch. Before a
400    /// launch there is no such dialog to point at (launch finding R6-3).
401    const fn then_retry(self) -> &'static str {
402        match self {
403            Self::BeforeLaunch => "",
404            Self::Launch => ", then Retry launch",
405        }
406    }
407}
408
409/// Check a target before a launch, resume, or move uses it.
410///
411/// An SSH host's container runtime is verified by `mj setup`, `mj doctor`,
412/// and the Targets pane's Test action; here it is assumed to still be as
413/// they left it, and only the host's reachability is checked. That probe
414/// joins an open shared connection, so it usually costs one SSH channel.
415pub(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/// Whether a local container engine can run sessions, as the launch
431/// preflight and the dashboard's Targets pane judge it.
432#[derive(Debug, Clone, Copy, PartialEq, Eq)]
433pub(crate) enum LocalEngineReadiness {
434    Ready,
435    /// The engine's command is not on this host.
436    NotInstalled,
437    /// The command is there, but the engine did not answer its check.
438    NotReady,
439}
440
441/// Run the launch preflight's engine check for a local container target
442/// kind (`local-podman`, `local-docker`, or `apple-container`).
443///
444/// Answers `None` for any other kind: their readiness is not a local engine's.
445pub(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        // The same command `verify_target` runs, without its sentence, which
453        // flattens the cause this needs to tell a missing command apart.
454        "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
473/// Whether a command failed because its program is not installed.
474fn 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
482/// Check that a target's host and runtime can run sessions.
483pub(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                    // Leads with the sentence the launch options use, so a
502                    // missing Docker reads "not installed" in the wizard's
503                    // row too (launch finding R5-3).
504                    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/// Per-session container size overrides. They win over both the target
690/// template's values and any recorded resource allocation, and they are read
691/// only while a container is being created.
692#[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            // The address is filled after describe-instances.
746            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
774/// Every container image the daemon downloads in the background, once per
775/// (host, image, platform).
776///
777/// Every container target is covered, including Apple's `container` engine.
778/// A `never` policy is the one opt-out. The rest differ only in when they
779/// download: `always` and `newer` pull on every refresh, while the others
780/// pull only when the host has no copy of the image.
781///
782/// Several targets often share one image on one host, and that needs one
783/// download. When two such targets disagree about when to pull, the merged
784/// entry takes the more eager of the two.
785pub 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        // Commands are decided by the host, the image, and the platform alone,
800        // so those three identify the duplicates worth collapsing.
801        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    // The session's own overrides are the last word on size.
907    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
952/// How long a freshly launched EC2 instance may take to accept SSH.
953const AWS_SSH_READY_TIMEOUT: Duration = Duration::from_secs(300);
954
955const AWS_SSH_READY_RETRY_DELAY: Duration = Duration::from_secs(3);
956
957/// Poll a remote host until it accepts SSH, or until the deadline passes.
958///
959/// `now` and `sleep` are injected so tests can drive the deadline without
960/// waiting in real time.
961fn 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
1065/// Resolve a created instance. Callers persist its ID before these long waits.
1066pub(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
1157/// The execution-plan locator for a session's stored target.
1158///
1159/// The mapping itself lives in `mj_core::targets`; this only pairs the stored
1160/// locator with the target template the session was created against.
1161pub(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;