Skip to main content

mj_controller/controller/
provisioning.rs

1//! Session provisioning, rollback, and worker-side Git bootstrap.
2
3use std::collections::{BTreeMap, HashMap};
4use std::path::{Path, PathBuf};
5use std::process::{Command, Stdio};
6use std::sync::{Arc, Mutex, OnceLock};
7use std::time::Instant;
8
9use anyhow::{Context, Result, bail, ensure};
10
11use mj_core::config::{TargetTemplate, atomic_write, data_dir};
12use mj_core::state::{SessionState, State, TargetLocator};
13
14use crate::targets::{
15    self, CommandExecutor, CommandOutput, CommandSpec, ProvisionStage, ProvisionStageGuard,
16};
17
18use super::backend::{
19    ContainerOverrides, TargetCheck, backend_locator, backend_session_bundle, backend_target,
20    configure_github_token_environment, controller_github_token, preflight_target,
21    use_github_https_urls,
22};
23use super::git_cache;
24use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
25use super::worker_binary::{bridge_readiness_stage, start_worker, worker_probe_diagnosis};
26use super::{Controller, execute_checked, now};
27
28const INHERITED_GIT_SETTINGS: &[&str] = &[
29    "diff.algorithm",
30    "fetch.prune",
31    "fetch.prunetags",
32    "init.defaultbranch",
33    "merge.conflictstyle",
34    "pull.ff",
35    "pull.rebase",
36    "push.autosetupremote",
37    "push.default",
38    "rebase.autostash",
39    "rerere.autoupdate",
40    "rerere.enabled",
41    "user.email",
42    "user.name",
43];
44
45/// How many sub-agent children may be brought up inside one container at the
46/// same time.
47///
48/// Starting a child means starting a harness, and a harness start inside a
49/// container is expensive: the reviewer sidecar already caps its own
50/// specialist lanes at three for the same reason. Measured on a local Podman
51/// target, ten children started one after another each reached their harness
52/// in about seven seconds, while four started at once left two or three of
53/// them past the 300-second harness-startup wait. Admitting two at a time
54/// keeps a burst slower but finished, instead of fast and failed.
55const CONTAINER_START_ADMISSION: usize = 2;
56
57/// The admission gate for one container, created on first use.
58///
59/// Keyed by the container the children share. A bare target has no gate: a
60/// child there is an ordinary process on a whole machine, and twenty
61/// sequential and four concurrent starts measured 2.8 seconds each.
62pub(super) fn container_start_gate(
63    locator: &targets::TargetLocator,
64) -> Option<Arc<tokio::sync::Semaphore>> {
65    static GATES: OnceLock<Mutex<HashMap<String, Arc<tokio::sync::Semaphore>>>> = OnceLock::new();
66    let container = match locator {
67        targets::TargetLocator::LocalPodman { container_id, .. }
68        | targets::TargetLocator::LocalDocker { container_id, .. }
69        | targets::TargetLocator::AppleContainer { container_id, .. }
70        | targets::TargetLocator::SshPodman { container_id, .. }
71        | targets::TargetLocator::SshDocker { container_id, .. } => container_id.clone(),
72        targets::TargetLocator::LocalBare { .. }
73        | targets::TargetLocator::SshBare { .. }
74        | targets::TargetLocator::AwsEc2 { .. } => return None,
75    };
76    let gates = GATES.get_or_init(|| Mutex::new(HashMap::new()));
77    let mut gates = gates
78        .lock()
79        .unwrap_or_else(std::sync::PoisonError::into_inner);
80    Some(Arc::clone(gates.entry(container).or_insert_with(|| {
81        Arc::new(tokio::sync::Semaphore::new(CONTAINER_START_ADMISSION))
82    })))
83}
84
85/// Whether starting a child worker can be tried again.
86///
87/// Only a worker that provably never published its control socket qualifies:
88/// it owns no relay, no journal and no harness, so a second start cannot
89/// duplicate or corrupt work. A refusal is never retried, because it names a
90/// precondition that a second attempt would meet in exactly the same way, and
91/// a cancelled operation is not retried either.
92fn subagent_start_is_retryable(error: &anyhow::Error) -> bool {
93    if mj_core::refusal::Refusal::of(error).is_some() {
94        return false;
95    }
96    if format!("{error:#}").contains("operation cancelled") {
97        return false;
98    }
99    error
100        .downcast_ref::<super::readiness::WorkerStartupFailure>()
101        .is_some_and(|failure| !failure.reached_socket)
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq)]
105pub(super) enum ProvisioningFailureDisposition {
106    /// A freshly registered session has no durable history to retain.
107    Discard,
108    /// Resume owns rollback to the archived record and checkpoint lineage.
109    Preserve,
110}
111
112impl Controller {
113    pub async fn provision_session_controlled_with_commit(
114        &mut self,
115        session_id: &str,
116        executor: &(impl CommandExecutor + Sync),
117        grant_commit: impl FnOnce() -> Result<()>,
118    ) -> Result<()> {
119        let github_token = controller_github_token();
120        let repositories = self
121            .provision_session_target_with_failure_disposition(
122                session_id,
123                executor,
124                github_token.as_deref(),
125                ProvisioningFailureDisposition::Discard,
126            )
127            .await?;
128        let setup = execute_concurrent_lanes(
129            || execute_repository_setup(&repositories, executor),
130            || self.install_worker_payload(session_id, executor),
131        );
132        let result = match setup {
133            Ok(((), (backend, worker_root))) => {
134                self.connect_and_start_worker(session_id, executor, &backend, &worker_root, true)
135                    .await
136            }
137            Err(error) => Err(error),
138        };
139        match result {
140            Ok(native_session_id) => {
141                if let Err(error) = grant_commit() {
142                    return Err(self.rollback_failed_new_session(session_id, error)?);
143                }
144                self.mark_worker_connected(session_id, native_session_id)
145            }
146            Err(error) => Err(self.rollback_failed_new_session(session_id, error)?),
147        }
148    }
149
150    /// Start a child worker inside an already-provisioned parent target.
151    /// Repository, target, and mount setup belong exclusively to the parent.
152    pub async fn provision_subagent_session_controlled(
153        &mut self,
154        session_id: &str,
155        executor: &(impl CommandExecutor + Sync),
156    ) -> Result<()> {
157        // One retry, and only for a worker that provably never published a
158        // control socket. Such a worker has no relay, no durable journal and
159        // no harness, so starting another over the same root cannot duplicate
160        // or corrupt anything. A spawn is issued by a model that cannot see
161        // the target, so a transient start failure it could have retried by
162        // hand is better retried here.
163        let mut attempts: Vec<String> = Vec::new();
164        loop {
165            let (result, placement) = self.attempt_subagent_start(session_id, executor).await;
166            let error = match result {
167                Ok(native_session_id) => {
168                    return self.mark_worker_connected(session_id, native_session_id);
169                }
170                Err(error) => error,
171            };
172            let retry = attempts.is_empty() && subagent_start_is_retryable(&error);
173            let error = match &placement {
174                Some((backend, _)) => {
175                    super::subagent_park::explain_process_exhaustion(error, backend, session_id)
176                }
177                None => error,
178            };
179            attempts.push(format!("{error:#}"));
180            let error = if attempts.len() > 1 {
181                anyhow::anyhow!(
182                    "{}",
183                    attempts
184                        .iter()
185                        .enumerate()
186                        .map(|(index, error)| format!("attempt {}: {error}", index + 1))
187                        .collect::<Vec<_>>()
188                        .join("; ")
189                )
190            } else {
191                error
192            };
193            let diagnostic = note_new_session_launch_failure(session_id, &error);
194            let previous = self
195                .state
196                .sessions
197                .get(session_id)
198                .context("failed child disappeared")?
199                .clone();
200            let record = self.state.sessions.get_mut(session_id).unwrap();
201            record.state = SessionState::StartupCleanup;
202            record.updated_at = now();
203            record.last_error = Some(format!("sub-agent startup failed: {diagnostic}"));
204            self.persist_session_transition_or_restore(
205                session_id,
206                &previous,
207                "record failed startup before teardown",
208            )?;
209            let cleanup = super::failed_launch_cleanup_executor();
210            if let Err(cleanup_error) =
211                self.finish_failed_startup_controlled(session_id, &cleanup, retry)
212            {
213                return Err(error.context(format!(
214                    "startup cleanup remains pending: {cleanup_error:#}"
215                )));
216            }
217            if !retry {
218                return Err(error);
219            }
220            tracing::warn!(
221                session_id,
222                error = format!("{error:#}"),
223                "sub-agent worker never started and was stopped; retrying once"
224            );
225        }
226    }
227
228    /// Finish a failed startup without reopening its relay or replaying work.
229    pub(crate) fn cleanup_failed_startup_controlled(
230        &mut self,
231        session_id: &str,
232        executor: &impl CommandExecutor,
233    ) -> Result<()> {
234        self.finish_failed_startup_controlled(session_id, executor, false)
235    }
236
237    fn finish_failed_startup_controlled(
238        &mut self,
239        session_id: &str,
240        executor: &impl CommandExecutor,
241        retry_after_stop: bool,
242    ) -> Result<()> {
243        let previous = self
244            .state
245            .sessions
246            .get(session_id)
247            .context("cleanup session disappeared")?
248            .clone();
249        ensure!(
250            previous.state == SessionState::StartupCleanup,
251            "session is not awaiting startup cleanup"
252        );
253        let child = self.state.subagents.contains_key(session_id);
254        ensure!(
255            !retry_after_stop || child,
256            "only unaccepted child startup may retry"
257        );
258        let outcome = (|| -> Result<()> {
259            if let Some(locator) = &previous.target {
260                let backend = backend_locator(locator, &previous, &self.config)?;
261                if child {
262                    let root = targets::worker_root(&backend, session_id)?;
263                    super::worker_binary::stop_worker(executor, &backend, &root)?;
264                } else {
265                    targets::close_plan(&backend, session_id)?.execute(executor)?;
266                }
267            }
268            if !child {
269                self.cleanup_new_session_worktree_after_failure(session_id, executor)?;
270            }
271            Ok(())
272        })();
273        let record = self.state.sessions.get_mut(session_id).unwrap();
274        record.updated_at = now();
275        let cause = previous
276            .last_error
277            .as_deref()
278            .unwrap_or("startup failed")
279            .split("; startup cleanup failed:")
280            .next()
281            .unwrap()
282            .to_owned();
283        match &outcome {
284            Ok(()) => {
285                record.state = if retry_after_stop {
286                    SessionState::Provisioning
287                } else {
288                    SessionState::Error
289                };
290                record.last_error = (!retry_after_stop).then_some(cause);
291                if !child {
292                    record.target = None;
293                }
294            }
295            Err(error) => {
296                record.last_error = Some(format!("{cause}; startup cleanup failed: {error:#}"));
297            }
298        }
299        let record = self.state.sessions.get_mut(session_id).unwrap();
300        if outcome.is_ok() && child && !retry_after_stop {
301            record.archived = true;
302        }
303        if let Err(error) = crate::database::save_startup_cleanup_outcome(
304            record,
305            outcome.is_ok() && child && !retry_after_stop,
306        ) {
307            self.state.sessions.insert(session_id.to_owned(), previous);
308            return Err(error).context("record startup cleanup outcome");
309        }
310        outcome
311    }
312
313    /// One start of a child worker: place it, install its files, and wait for
314    /// its relay. The placement is returned even on failure, because the
315    /// caller needs it to stop a worker that may be half up.
316    async fn attempt_subagent_start(
317        &mut self,
318        session_id: &str,
319        executor: &(impl CommandExecutor + Sync),
320    ) -> (
321        Result<Option<String>>,
322        Option<(targets::TargetLocator, String)>,
323    ) {
324        // Placement failures must reach the same failure arm as startup
325        // failures; otherwise the child record stays `Provisioning` forever.
326        match self.worker_placement(session_id) {
327            Ok((backend, worker_root)) => {
328                let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
329                let prepared = self
330                    .prepare_worker_files(session_id, &backend, &worker_root, syncing)
331                    .and_then(|()| self.prepare_subagent_report_dir(session_id, &backend, syncing));
332                let result = match prepared {
333                    Ok(()) => {
334                        // Held across the harness startup wait, which is the
335                        // part that does not survive a crowd.
336                        let gate = container_start_gate(&backend);
337                        let _admitted = match &gate {
338                            Some(gate) => gate.acquire().await.ok(),
339                            None => None,
340                        };
341                        self.connect_and_start_worker(
342                            session_id,
343                            executor,
344                            &backend,
345                            &worker_root,
346                            false,
347                        )
348                        .await
349                    }
350                    Err(error) => Err(error),
351                };
352                (result, Some((backend, worker_root)))
353            }
354            Err(error) => (Err(error), None),
355        }
356    }
357
358    fn rollback_failed_new_session(
359        &mut self,
360        session_id: &str,
361        error: anyhow::Error,
362    ) -> Result<anyhow::Error> {
363        self.rollback_failed_new_session_with(
364            session_id,
365            error,
366            // Rollback must remain possible after the foreground operation's
367            // cancellation token has been set.
368            &super::failed_launch_cleanup_executor(),
369        )
370    }
371
372    fn rollback_failed_new_session_with(
373        &mut self,
374        session_id: &str,
375        error: anyhow::Error,
376        target_cleanup_executor: &impl CommandExecutor,
377    ) -> Result<anyhow::Error> {
378        let session = self
379            .state
380            .sessions
381            .get(session_id)
382            .with_context(|| format!("unknown session {session_id}"))?
383            .clone();
384        // The failure goes on record before anything the record points at is
385        // removed. Removing the worker and its checkout can outlast a daemon
386        // that is stopping, and a daemon that exits partway would otherwise
387        // leave a record naming a worker that no longer exists, which every
388        // later daemon keeps reconnecting to. A failed record that still
389        // names its target is one Destroy knows how to clean up.
390        let original = note_new_session_launch_failure(session_id, &error);
391        apply_failed_new_session_launch(&mut self.state, session_id, &original);
392        self.persist_session_transition_or_restore(
393            session_id,
394            &session,
395            "record the failed launch before removing its target",
396        )
397        .map_err(|persist_error| {
398            persist_error.context(format!(
399                "{original}; its target was left in place because the failure could not be recorded"
400            ))
401        })?;
402        let target_cleanup = match session.target.as_ref() {
403            Some(locator) => (|| -> Result<()> {
404                let backend = backend_locator(locator, &session, &self.config)?;
405                targets::close_plan(&backend, session_id)?
406                    .execute(target_cleanup_executor)
407                    .map(|_| ())
408            })(),
409            None => Ok(()),
410        };
411        // A failed stop leaves the harness able to write its checkout. Keep
412        // that checkout until a later cleanup confirms process termination.
413        let cleanup_error = target_cleanup
414            .and_then(|()| {
415                self.cleanup_new_session_worktree_after_failure(session_id, target_cleanup_executor)
416            })
417            .err()
418            .map(|error| format!("{error:#}"));
419        if let Some(cleanup_error) = &cleanup_error {
420            tracing::warn!(
421                session_id,
422                error = %cleanup_error,
423                "new-session rollback cleanup reported failures"
424            );
425        }
426        let failure = apply_failed_new_session_rollback(
427            &mut self.state,
428            session_id,
429            &original,
430            cleanup_error,
431        );
432        self.persist_session_state(session_id)?;
433        Ok(failure)
434    }
435
436    pub(super) async fn provision_session_with_failure_disposition(
437        &mut self,
438        session_id: &str,
439        executor: &(impl CommandExecutor + Sync),
440        github_token: Option<&str>,
441        failure_disposition: ProvisioningFailureDisposition,
442    ) -> Result<()> {
443        let repositories = self
444            .provision_session_target_with_failure_disposition(
445                session_id,
446                executor,
447                github_token,
448                failure_disposition,
449            )
450            .await?;
451        match execute_repository_setup(&repositories, executor) {
452            Ok(()) => Ok(()),
453            Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
454                Err(self.rollback_failed_new_session(session_id, error)?)
455            }
456            Err(error) => Err(error),
457        }
458    }
459
460    async fn provision_session_target_with_failure_disposition(
461        &mut self,
462        session_id: &str,
463        executor: &(impl CommandExecutor + Sync),
464        github_token: Option<&str>,
465        failure_disposition: ProvisioningFailureDisposition,
466    ) -> Result<targets::CommandPlan> {
467        let session = self
468            .state
469            .sessions
470            .get(session_id)
471            .with_context(|| format!("unknown session {session_id}"))?
472            .clone();
473        if session.state != SessionState::Provisioning {
474            bail!("session {session_id} is not provisioning");
475        }
476        let preparation = (|| {
477            let selected = self
478                .config
479                .targets
480                .get(&session.target_template_id)
481                .context("target template disappeared before provisioning")?;
482            let runtime = mj_core::state::TargetRuntimeSettings::from(selected);
483            if let Some(recorded) = &session.target_runtime {
484                ensure!(
485                    recorded == &runtime,
486                    "target access settings changed before provisioning; retry with the selected target"
487                );
488            } else {
489                self.state
490                    .sessions
491                    .get_mut(session_id)
492                    .unwrap()
493                    .target_runtime = Some(runtime);
494                self.persist_session_state(session_id)?;
495            }
496            let template = self
497                .config
498                .targets
499                .get(&session.target_template_id)
500                .context("target template disappeared during provisioning")?;
501            let profile = self
502                .config
503                .profiles
504                .get(&session.last_profile)
505                .context("harness profile disappeared during provisioning")?;
506            super::worker_binary::preflight_worker_binary(template, executor)?;
507            super::worker_binary::preflight_harness(template, profile, executor)?;
508            self.prepare_managed_raw_worktree(session_id, executor)
509        })();
510        let created_worktree = match preparation {
511            Ok(created) => created,
512            Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
513                return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
514            }
515            Err(error) => return Err(error),
516        };
517        let session = self
518            .state
519            .sessions
520            .get(session_id)
521            .expect("session retained after managed worktree preparation")
522            .clone();
523        // Keep planning, preflight, creation, and locator discovery in one
524        // result so the caller's failure disposition applies to every error.
525        let result = (|| {
526            let template = self
527                .config
528                .targets
529                .get(&session.target_template_id)
530                .context("target template disappeared during provisioning")?;
531            if matches!(template, TargetTemplate::AwsEc2 { .. }) {
532                for resource in &session.additional_mounts {
533                    ensure!(
534                        resource.source.is_dir(),
535                        "attached resource source is not a directory: {}",
536                        resource.source.display()
537                    );
538                }
539            }
540            let mut target = backend_target(
541                template,
542                session.resource_allocation.as_ref(),
543                ContainerOverrides::for_session(&session),
544            )?;
545            let mut runtime_mounts = if matches!(target, targets::TargetTemplate::AwsEc2(_)) {
546                Vec::new()
547            } else {
548                session.additional_mounts.clone()
549            };
550            // The mounts this container runs with, not the ones the session
551            // stores: a forced downgrade belongs to the host the container
552            // lands on, so it is decided here every time and never written
553            // over the user's choice.
554            for notice in enforce_overlay_capable_mounts(&target, &mut runtime_mounts, executor) {
555                executor.notify_notice(&notice);
556            }
557            // The image's user is a property of the host's copy of the image,
558            // so it is read here, once per image per daemon, and handed to the
559            // plan rather than stored on the session.
560            let image_user = podman_image_user(&target, executor);
561            let mut bundle = if session.project_directory.is_some() {
562                None
563            } else if let Some(bundle) = self.move_destination_bundle(session_id)? {
564                Some(bundle)
565            } else if failure_disposition == ProvisioningFailureDisposition::Preserve {
566                Some(super::network_git::checkpoint_bundle(&session)?)
567            } else {
568                Some(backend_session_bundle(&session, &self.config, executor)?)
569            };
570            let container_github_token =
571                github_token.filter(|_| configure_github_token_environment(&mut target));
572            if container_github_token.is_some()
573                && let Some(bundle) = bundle.as_mut()
574            {
575                use_github_https_urls(bundle);
576            }
577            preflight_target(template, executor, TargetCheck::Launch)?;
578            let resource_name = crate::database::load_move_operation(session_id)?
579                .filter(|op| {
580                    op.workspace_transfer.is_some()
581                        && op.phase == mj_core::state::MovePhase::ResumingDestination
582                })
583                .map(|op| targets::move_resource_name(session_id, &op.operation_id))
584                .unwrap_or_else(|| targets::resource_name(session_id))?;
585            let moving_workspace = resource_name != targets::resource_name(session_id)?;
586            let prepared_cache = bundle
587                .as_mut()
588                .filter(|_| !moving_workspace)
589                .and_then(|bundle| {
590                    git_cache::prepare(
591                        &target,
592                        session_id,
593                        bundle,
594                        &mut runtime_mounts,
595                        container_github_token,
596                        executor,
597                    )
598                });
599            // Mounts are fixed when the container is created, so the build
600            // cache is decided here, before the provisioning plan is built.
601            let build_cache = super::mbx::prepare(
602                &target,
603                &session,
604                bundle.as_ref(),
605                prepared_cache.as_ref(),
606                &mut runtime_mounts,
607                executor,
608            );
609            let provision = if let Some(project_directory) = &session.project_directory {
610                targets::provision_bare_project_plan(
611                    &target,
612                    session_id,
613                    &project_directory.to_string_lossy(),
614                )
615            } else {
616                bundle
617                    .as_ref()
618                    .context("project bundle disappeared during provisioning")
619                    .and_then(|bundle| {
620                        targets::provision_plan_named(
621                            &target,
622                            session_id,
623                            bundle,
624                            &runtime_mounts,
625                            image_user,
626                            session.container_workspace.as_deref(),
627                            &resource_name,
628                        )
629                    })
630            };
631            let mut provision = match provision {
632                Ok(provision) => provision,
633                Err(error) => {
634                    if let Some(cache) = &prepared_cache {
635                        let _ = cache.cleanup(executor);
636                    }
637                    return Err(error);
638                }
639            };
640            if let Some(token) = container_github_token
641                && let Err(error) =
642                    provision.provide_target_environment_secret(&target, "GH_TOKEN", token)
643            {
644                if let Some(cache) = &prepared_cache {
645                    let _ = cache.cleanup(executor);
646                }
647                return Err(error);
648            }
649
650            let started = Instant::now();
651            let result = provision_target_creation_named(
652                &provision,
653                &target,
654                session_id,
655                executor,
656                &resource_name,
657                |outputs| {
658                    super::backend::locator_after_provision_named(
659                        template,
660                        &target,
661                        session_id,
662                        outputs.first(),
663                        executor,
664                        &resource_name,
665                    )
666                },
667            )
668            .map(|(locator, remainder)| (locator, remainder, bundle, build_cache));
669            if result.is_err()
670                && let Some(cache) = &prepared_cache
671            {
672                if let Some(locator) =
673                    provisioned_locator_named(&target, session_id, None, &resource_name)
674                {
675                    let _ = targets::close_plan(&locator, session_id)
676                        .and_then(|plan| plan.execute(executor).map(|_| ()));
677                } else {
678                    let _ = cache.cleanup(executor);
679                }
680            }
681            tracing::debug!(
682                session_id,
683                elapsed_ms = started.elapsed().as_millis(),
684                "provisioning plan execution completed"
685            );
686            result
687        })();
688        let result = match result {
689            Err(error)
690                if created_worktree
691                    && failure_disposition == ProvisioningFailureDisposition::Discard =>
692            {
693                return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
694            }
695            Err(error) if failure_disposition == ProvisioningFailureDisposition::Preserve => {
696                Err(error)
697            }
698            Err(error) => {
699                // This arm (no managed worktree to unwind) is the one a
700                // provisioning failure such as a dropped target connection
701                // hits; record the diagnostic and the session-id log here too.
702                let detail = note_new_session_launch_failure(session_id, &error);
703                {
704                    let record = self.state.sessions.get_mut(session_id).unwrap();
705                    record.state = SessionState::Error;
706                    record.target = None;
707                    record.updated_at = super::now();
708                    record.last_error = Some(format!("session provisioning failed: {detail}"));
709                }
710                return match self.persist_session_state(session_id) {
711                    Ok(()) => Err(error),
712                    Err(persistence_error) => Err(error.context(format!(
713                        "persist removal of failed provisioning session {session_id}: {persistence_error:#}"
714                    ))),
715                };
716            }
717            Ok((locator, remainder, bundle, build_cache)) => {
718                apply_new_session_provisioning_result(&mut self.state, session_id, Ok(locator))?;
719                self.state
720                    .sessions
721                    .get_mut(session_id)
722                    .expect("session retained after provisioning")
723                    .build_cache = build_cache;
724                let session = &self.state.sessions[session_id];
725                let backend = backend_locator(
726                    session
727                        .target
728                        .as_ref()
729                        .context("provisioned target disappeared")?,
730                    session,
731                    &self.config,
732                )?;
733                if matches!(backend, targets::TargetLocator::AwsEc2 { .. }) {
734                    targets::provision_on_locator_plan(
735                        &backend,
736                        session_id,
737                        bundle
738                            .as_ref()
739                            .context("AWS provisioning requires a project bundle")?,
740                    )
741                } else {
742                    Ok(remainder)
743                }
744            }
745        };
746        let result = match result {
747            Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
748                return Err(self.rollback_failed_new_session(session_id, error)?);
749            }
750            result => result,
751        };
752        if result.is_ok()
753            && let Some(session) = self.state.sessions.get(session_id)
754            && let Some(directory) = session
755                .managed_worktree
756                .as_ref()
757                .map(|worktree| worktree.source_project_directory.clone())
758                .or_else(|| session.project_directory.clone())
759            && let Some(template) = self.config.targets.get(&session.target_template_id)
760        {
761            let host = match template {
762                TargetTemplate::LocalBare => Some("local"),
763                TargetTemplate::SshBare { ssh, .. } => Some(ssh.host.as_str()),
764                _ => None,
765            };
766            if let Some(host) = host {
767                self.state.remember_project_directory(host, &directory);
768                crate::database::remember_project_directory(host, &directory)?;
769            }
770        }
771        self.persist_session_state(session_id)?;
772        result
773    }
774
775    pub fn mark_worker_connected(
776        &mut self,
777        session_id: &str,
778        native_session_id: Option<String>,
779    ) -> Result<()> {
780        let session = self
781            .state
782            .sessions
783            .get(session_id)
784            .with_context(|| format!("unknown session {session_id}"))?;
785        if session.target.is_none() {
786            bail!("session {session_id} has no provisioned target");
787        }
788        let updated_at = now();
789        crate::database::mark_session_worker_connected(
790            session_id,
791            native_session_id.as_deref(),
792            &updated_at,
793        )?;
794        let session = self
795            .state
796            .sessions
797            .get_mut(session_id)
798            .expect("session disappeared after its worker connection was saved");
799        session.state = SessionState::Running;
800        if native_session_id.is_some() {
801            session.native_session_id = native_session_id;
802        }
803        session.updated_at = updated_at;
804        session.last_error = None;
805        Ok(())
806    }
807
808    fn install_worker_payload(
809        &self,
810        session_id: &str,
811        executor: &impl CommandExecutor,
812    ) -> Result<(targets::TargetLocator, String)> {
813        // Worker/profile installation is independent of repository cloning.
814        let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
815        let (backend, worker_root) = self.worker_placement(session_id)?;
816        self.prepare_worker_files(session_id, &backend, &worker_root, syncing)?;
817        install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
818        Ok((backend, worker_root))
819    }
820
821    async fn connect_and_start_worker(
822        &self,
823        session_id: &str,
824        executor: &impl CommandExecutor,
825        backend: &targets::TargetLocator,
826        worker_root: &str,
827        initialize_workspace: bool,
828    ) -> Result<Option<String>> {
829        let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
830        if initialize_workspace {
831            install_inherited_git_settings(executor, backend, session_id)?;
832            self.initialize_network_workspaces(session_id, backend, syncing)?;
833        }
834        let session = self
835            .state
836            .sessions
837            .get(session_id)
838            .with_context(|| format!("unknown session {session_id}"))?;
839        let profile = self
840            .config
841            .profiles
842            .get(&session.last_profile)
843            .with_context(|| format!("unknown profile {}", session.last_profile))?;
844        let readiness_stage = bridge_readiness_stage(profile);
845        let reconnect = &targets::reconnect_plan(backend, session_id)?.commands[0];
846        let readiness = async {
847            let mut relay = {
848                let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
849                start_worker(executor, backend, worker_root)?;
850                connect_started_worker(reconnect, session_id, executor, backend, worker_root)
851                    .await?
852            };
853            let native_session_id =
854                wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
855            Ok(Some(native_session_id))
856        }
857        .await;
858        match readiness {
859            Ok(native_session_id) => Ok(native_session_id),
860            Err(error) => {
861                // The diagnosis reads the worker's state, so it runs before the
862                // worker is stopped.
863                let error = worker_probe_diagnosis(executor, backend, worker_root, error);
864                Err(error)
865            }
866        }
867    }
868}
869
870const MAX_LAUNCH_DIAGNOSTIC_BYTES: usize = 64 * 1024;
871
872const RETAINED_LAUNCH_DIAGNOSTICS: usize = 20;
873
874/// Record a failed new-session launch consistently across every failure arm:
875/// log the failure with the session id (so `logs/mj-*.log` names the session)
876/// and save the local diagnostic file. Returns the underlying error chain
877/// annotated with the diagnostic path, which the caller stores in `last_error`
878/// so the reason travels to `mj sessions`, `mj wait`, and `mj events`.
879pub(super) fn note_new_session_launch_failure(session_id: &str, error: &anyhow::Error) -> String {
880    note_new_session_launch_failure_in(&data_dir().join("diagnostics"), session_id, error)
881}
882
883fn note_new_session_launch_failure_in(
884    directory: &Path,
885    session_id: &str,
886    error: &anyhow::Error,
887) -> String {
888    let original = format!("{error:#}");
889    tracing::warn!(session_id, error = %original, "session launch failed");
890    match persist_launch_failure_to(directory, session_id, &original) {
891        Ok(path) => format!("{original}; full diagnostic saved to {}", path.display()),
892        Err(save_error) => {
893            format!("{original}; saving the local diagnostic failed: {save_error:#}")
894        }
895    }
896}
897
898fn persist_launch_failure_to(directory: &Path, session_id: &str, detail: &str) -> Result<PathBuf> {
899    mj_core::config::validate_id("session", session_id)?;
900    std::fs::create_dir_all(directory).with_context(|| {
901        format!(
902            "create launch diagnostics directory {}",
903            directory.display()
904        )
905    })?;
906    #[cfg(unix)]
907    {
908        use std::os::unix::fs::PermissionsExt;
909        std::fs::set_permissions(directory, std::fs::Permissions::from_mode(0o700))?;
910    }
911    let path = directory.join(format!("{session_id}-launch-error.txt"));
912    let detail = bounded_launch_diagnostic(detail);
913    let body = format!(
914        "Mjolnir session launch failure\nsession: {session_id}\nat: {}\n\n{detail}\n",
915        now()
916    );
917    atomic_write(&path, body.as_bytes())?;
918    prune_launch_diagnostics(directory)?;
919    Ok(path)
920}
921
922fn bounded_launch_diagnostic(detail: &str) -> String {
923    if detail.len() <= MAX_LAUNCH_DIAGNOSTIC_BYTES {
924        return detail.to_owned();
925    }
926    let mut head_end = MAX_LAUNCH_DIAGNOSTIC_BYTES / 4;
927    while !detail.is_char_boundary(head_end) {
928        head_end -= 1;
929    }
930    let tail_bytes = MAX_LAUNCH_DIAGNOSTIC_BYTES - head_end;
931    let mut tail_start = detail.len() - tail_bytes;
932    while !detail.is_char_boundary(tail_start) {
933        tail_start += 1;
934    }
935    format!(
936        "{}\n\n[... launch diagnostic truncated ...]\n\n{}",
937        &detail[..head_end],
938        &detail[tail_start..]
939    )
940}
941
942fn prune_launch_diagnostics(directory: &Path) -> Result<()> {
943    let mut diagnostics = Vec::new();
944    for entry in std::fs::read_dir(directory)? {
945        let entry = entry?;
946        if !entry
947            .file_name()
948            .to_str()
949            .is_some_and(|name| name.ends_with("-launch-error.txt"))
950        {
951            continue;
952        }
953        diagnostics.push((entry.metadata()?.modified()?, entry.path()));
954    }
955    diagnostics.sort_by_key(|entry| std::cmp::Reverse(entry.0));
956    for (_, path) in diagnostics.into_iter().skip(RETAINED_LAUNCH_DIAGNOSTICS) {
957        std::fs::remove_file(&path)
958            .with_context(|| format!("prune old launch diagnostic {}", path.display()))?;
959    }
960    Ok(())
961}
962
963fn apply_new_session_provisioning_result(
964    state: &mut State,
965    session_id: &str,
966    result: Result<TargetLocator>,
967) -> Result<()> {
968    match result {
969        Ok(locator) => {
970            let record = state.sessions.get_mut(session_id).unwrap();
971            record.target = Some(locator);
972            // Provisioning has completed, but Running is reserved for a
973            // successful worker handshake.
974            record.state = SessionState::Disconnected;
975            record.updated_at = now();
976            record.last_error = None;
977            Ok(())
978        }
979        Err(error) => {
980            let record = state.sessions.get_mut(session_id).unwrap();
981            record.state = SessionState::Error;
982            record.target = None;
983            record.updated_at = now();
984            record.last_error = Some(format!("session provisioning failed: {error:#}"));
985            Err(error)
986        }
987    }
988}
989
990/// Mark a new session's launch failed while it still names its target, so the
991/// record is true while the rollback removes that target.
992fn apply_failed_new_session_launch(state: &mut State, session_id: &str, original_error: &str) {
993    let record = state.sessions.get_mut(session_id).unwrap();
994    record.state = SessionState::StartupCleanup;
995    record.updated_at = now();
996    record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
997}
998
999pub(super) fn apply_failed_new_session_rollback(
1000    state: &mut State,
1001    session_id: &str,
1002    original_error: &str,
1003    cleanup_error: Option<String>,
1004) -> anyhow::Error {
1005    match cleanup_error {
1006        None => {
1007            let record = state.sessions.get_mut(session_id).unwrap();
1008            record.state = SessionState::Error;
1009            record.target = None;
1010            record.updated_at = now();
1011            record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1012            anyhow::anyhow!("{original_error}; partial target removed and failed session retained")
1013        }
1014        Some(cleanup_error) => {
1015            let failure = format!(
1016                "{original_error}; cleanup of the failed session target failed: {cleanup_error}"
1017            );
1018            let record = state.sessions.get_mut(session_id).unwrap();
1019            record.state = SessionState::StartupCleanup;
1020            record.updated_at = now();
1021            record.last_error = Some(format!("worker bootstrap failed: {failure}"));
1022            anyhow::anyhow!(failure)
1023        }
1024    }
1025}
1026
1027pub(super) fn install_attached_resources(
1028    state: &State,
1029    session_id: &str,
1030    backend: &targets::TargetLocator,
1031    worker_root: &str,
1032    executor: &impl CommandExecutor,
1033) -> Result<()> {
1034    let targets::TargetLocator::AwsEc2 { .. } = backend else {
1035        return Ok(());
1036    };
1037    let session = state
1038        .sessions
1039        .get(session_id)
1040        .with_context(|| format!("unknown session {session_id}"))?;
1041    if session.additional_mounts.is_empty() {
1042        return Ok(());
1043    }
1044    for resource in &session.additional_mounts {
1045        let install = targets::command_on_locator(
1046            backend,
1047            session_id,
1048            vec![
1049                format!("{worker_root}/hel"),
1050                "worker".into(),
1051                "install-resource".into(),
1052                "--destination".into(),
1053                resource.destination.to_string_lossy().into_owned(),
1054            ],
1055            "stream attached resource",
1056        )?;
1057        mj_checkpoint::resources::stream_resource(&resource.source, |stream| {
1058            execute_checked_with_stdin(executor, &install, stream).map(|_| ())
1059        })
1060        .with_context(|| format!("stream attached resource {}", resource.source.display()))?;
1061    }
1062    Ok(())
1063}
1064
1065/// Run two independent target setup lanes at the same time and wait for both.
1066/// The first lane's failure wins deterministically when both fail, and neither
1067/// lane is abandoned while it may still own a transfer or subprocess.
1068pub(super) fn execute_concurrent_lanes<A: Send, B: Send>(
1069    first: impl FnOnce() -> Result<A> + Send,
1070    second: impl FnOnce() -> Result<B> + Send,
1071) -> Result<(A, B)> {
1072    std::thread::scope(|scope| {
1073        let second = scope.spawn(second);
1074        let first = first();
1075        let second = second.join().unwrap_or_else(|panic| {
1076            Err(anyhow::anyhow!(
1077                "concurrent target lane panicked: {}",
1078                targets::command_thread_panic_message(panic.as_ref())
1079            ))
1080        });
1081        match (first, second) {
1082            (Err(error), _) => Err(error),
1083            (Ok(_), Err(error)) => Err(error),
1084            (Ok(first), Ok(second)) => Ok((first, second)),
1085        }
1086    })
1087}
1088
1089fn execute_repository_setup(
1090    plan: &targets::CommandPlan,
1091    executor: &(impl CommandExecutor + Sync),
1092) -> Result<()> {
1093    if plan.commands.is_empty() {
1094        return Ok(());
1095    }
1096    let _cloning = ProvisionStageGuard::new(executor, ProvisionStage::Cloning);
1097    plan.execute_concurrent(executor).map(|_| ())
1098}
1099
1100/// Run a provisioning plan and discover the locator it produced, tearing the
1101/// target down again if anything after its creation fails.
1102///
1103/// Creation is the boundary that matters. A step that fails before the target
1104/// exists has left nothing behind; every failure after it — a later plan step
1105/// or locator discovery — owns a target no session record will point at.
1106#[cfg(test)]
1107fn provision_target(
1108    plan: &targets::CommandPlan,
1109    target: &targets::TargetTemplate,
1110    session_id: &str,
1111    executor: &(impl CommandExecutor + Sync),
1112    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1113) -> Result<TargetLocator> {
1114    let Some((creation, remainder)) = plan.split_at_target_creation() else {
1115        // Nothing this plan runs can leave a target behind.
1116        return discover(&plan.execute_concurrent(executor)?);
1117    };
1118    let mut outputs = creation.execute_concurrent(executor)?;
1119    let result = match remainder.execute_concurrent(executor) {
1120        Ok(rest) => {
1121            outputs.extend(rest);
1122            discover(&outputs)
1123        }
1124        Err(error) => Err(error),
1125    };
1126    result.map_err(|error| {
1127        match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1128            Some(note) => error.context(note),
1129            None => error,
1130        }
1131    })
1132}
1133
1134/// Bring the target into existence and return the commands that populate its
1135/// repositories. The caller may overlap that remainder with worker/profile
1136/// installation once it has persisted the discovered locator.
1137#[cfg(test)]
1138fn provision_target_creation(
1139    plan: &targets::CommandPlan,
1140    target: &targets::TargetTemplate,
1141    session_id: &str,
1142    executor: &(impl CommandExecutor + Sync),
1143    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1144) -> Result<(TargetLocator, targets::CommandPlan)> {
1145    provision_target_creation_named(
1146        plan,
1147        target,
1148        session_id,
1149        executor,
1150        &targets::resource_name(session_id)?,
1151        discover,
1152    )
1153}
1154
1155fn provision_target_creation_named(
1156    plan: &targets::CommandPlan,
1157    target: &targets::TargetTemplate,
1158    session_id: &str,
1159    executor: &(impl CommandExecutor + Sync),
1160    name: &str,
1161    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1162) -> Result<(TargetLocator, targets::CommandPlan)> {
1163    let Some((creation, remainder)) = plan.split_at_target_creation() else {
1164        // Nothing this plan runs can leave a target behind, so its commands
1165        // must still finish before the locator is usable.
1166        let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1167            plan.execute_concurrent(executor)
1168        })?;
1169        return discover(&outputs).map(|locator| {
1170            (
1171                locator,
1172                targets::CommandPlan {
1173                    description: plan.description.clone(),
1174                    commands: Vec::new(),
1175                },
1176            )
1177        });
1178    };
1179    // Creating the container is what downloads the image on Docker, and the
1180    // probe just before it is what downloads it on Podman. Either way, a
1181    // background download of the same image must finish first.
1182    let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1183        creation.execute_concurrent(executor)
1184    })?;
1185    discover(&outputs)
1186        .map(|locator| (locator, remainder))
1187        .map_err(|error| {
1188            match cleanup_failed_provision_named(
1189                target,
1190                session_id,
1191                outputs.first(),
1192                executor,
1193                name,
1194            ) {
1195                Some(note) => error.context(note),
1196                None => error,
1197            }
1198        })
1199}
1200
1201/// Best-effort teardown of a target whose creation succeeded but whose
1202/// provisioning failed before a locator was recorded. Returns a note
1203/// describing what happened for inclusion in the session error.
1204///
1205/// The teardown is the session's own close plan, so a failed launch and an
1206/// ordinary close can never disagree about what removing a target means.
1207#[cfg(test)]
1208fn cleanup_failed_provision(
1209    target: &targets::TargetTemplate,
1210    session_id: &str,
1211    create_output: Option<&CommandOutput>,
1212    executor: &impl CommandExecutor,
1213) -> Option<String> {
1214    cleanup_failed_provision_named(
1215        target,
1216        session_id,
1217        create_output,
1218        executor,
1219        &targets::resource_name(session_id).ok()?,
1220    )
1221}
1222
1223fn cleanup_failed_provision_named(
1224    target: &targets::TargetTemplate,
1225    session_id: &str,
1226    create_output: Option<&CommandOutput>,
1227    executor: &impl CommandExecutor,
1228    name: &str,
1229) -> Option<String> {
1230    let locator = provisioned_locator_named(target, session_id, create_output, name)?;
1231
1232    let leak = format!(
1233        "the resource may still exist; find it via its dev.mj.session={session_id} label/tag"
1234    );
1235    let cleanup = if name == targets::resource_name(session_id).ok()? {
1236        targets::close_plan(&locator, session_id)
1237    } else {
1238        targets::retire_move_target_plan(&locator, session_id)
1239    };
1240    let plan = match cleanup {
1241        Ok(plan) => plan,
1242        Err(error) => {
1243            tracing::warn!(
1244                session_id,
1245                error = format!("{error:#}"),
1246                "could not build provisioning cleanup plan"
1247            );
1248            return Some(format!("cleanup FAILED: {error:#}; {leak}"));
1249        }
1250    };
1251    let purpose = plan
1252        .commands
1253        .iter()
1254        .map(|command| command.purpose.clone())
1255        .collect::<Vec<_>>()
1256        .join("; ");
1257    let Err(error) = plan.execute(executor) else {
1258        return Some(format!("cleanup succeeded: {purpose}"));
1259    };
1260    match targets::cleanup_target_is_confirmed_absent(&locator, session_id, executor) {
1261        Ok(true) => Some(format!("cleanup succeeded: {purpose}")),
1262        Ok(false) => {
1263            tracing::warn!(
1264                session_id,
1265                error = format!("{error:#}"),
1266                "provisioning cleanup failed and the target may still exist"
1267            );
1268            Some(format!("cleanup FAILED ({purpose}): {error:#}; {leak}"))
1269        }
1270        Err(confirm_error) => {
1271            tracing::warn!(
1272                session_id,
1273                error = format!("{confirm_error:#}"),
1274                "could not confirm whether the failed provisioning target was removed"
1275            );
1276            Some(format!(
1277                "cleanup FAILED ({purpose}): {error:#}; checking whether it was removed also failed: {confirm_error:#}; {leak}"
1278            ))
1279        }
1280    }
1281}
1282
1283/// The locator a provisioning plan's creating command brought into existence.
1284///
1285/// Every target but AWS is named before its plan runs; an EC2 instance
1286/// reports its own ID in the launch response.
1287fn provisioned_locator_named(
1288    target: &targets::TargetTemplate,
1289    session_id: &str,
1290    create_output: Option<&CommandOutput>,
1291    name: &str,
1292) -> Option<targets::TargetLocator> {
1293    let container_id = || Some(name.to_owned());
1294
1295    Some(match target {
1296        // A bare project directory belongs to the user: provisioning creates
1297        // nothing that a failure could leak.
1298        targets::TargetTemplate::LocalBare => return None,
1299        targets::TargetTemplate::LocalPodman(container) => targets::TargetLocator::LocalPodman {
1300            borrowed_from: None,
1301            container_id: container_id()?,
1302            workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1303        },
1304        targets::TargetTemplate::LocalDocker(_) => targets::TargetLocator::LocalDocker {
1305            borrowed_from: None,
1306            container_id: container_id()?,
1307        },
1308        targets::TargetTemplate::AppleContainer(_) => targets::TargetLocator::AppleContainer {
1309            borrowed_from: None,
1310            container_id: container_id()?,
1311        },
1312        targets::TargetTemplate::SshPodman { ssh, container } => {
1313            targets::TargetLocator::SshPodman {
1314                borrowed_from: None,
1315                ssh: ssh.clone(),
1316                container_id: container_id()?,
1317                workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1318            }
1319        }
1320        targets::TargetTemplate::SshDocker { ssh, .. } => targets::TargetLocator::SshDocker {
1321            borrowed_from: None,
1322            ssh: ssh.clone(),
1323            container_id: container_id()?,
1324        },
1325        targets::TargetTemplate::SshBare { ssh, .. } => targets::TargetLocator::SshBare {
1326            ssh: ssh.clone(),
1327            workspace: targets::workspace_for(target, session_id).ok()?,
1328            worker_id: None,
1329        },
1330        targets::TargetTemplate::AwsEc2(aws) => targets::TargetLocator::AwsEc2 {
1331            profile: aws.profile.clone(),
1332            region: aws.region.clone(),
1333            instance_id: serde_json::from_slice::<serde_json::Value>(&create_output?.stdout)
1334                .ok()?
1335                .pointer("/Instances/0/InstanceId")?
1336                .as_str()?
1337                .to_owned(),
1338            ssh: aws.ssh.clone(),
1339            workspace: targets::workspace_for(target, session_id).ok()?,
1340        },
1341    })
1342}
1343
1344/// Attach read-only whatever the selected container overlay cannot hold, and
1345/// say so.
1346///
1347/// The filesystem is probed on the host that runs the container, because that
1348/// is where the overlay would be built. A probe that cannot answer leaves the
1349/// overlay alone: a failed probe is no evidence of an unsupported filesystem,
1350/// and refusing to provision over one would cost the user their session.
1351///
1352/// Apple's `container` engine already mounts every extra directory read-only,
1353/// and EC2 copies the directory instead of mounting it, so neither is probed.
1354pub(super) fn enforce_overlay_capable_mounts(
1355    target: &targets::TargetTemplate,
1356    mounts: &mut [targets::AdditionalMount],
1357    executor: &impl CommandExecutor,
1358) -> Vec<String> {
1359    let ssh = match target {
1360        targets::TargetTemplate::LocalPodman(_) | targets::TargetTemplate::LocalDocker(_) => None,
1361        targets::TargetTemplate::SshPodman { ssh, .. }
1362        | targets::TargetTemplate::SshDocker { ssh, .. } => Some(ssh),
1363        _ => return Vec::new(),
1364    };
1365    let overlaid = mounts
1366        .iter()
1367        .filter(|mount| mount.access == targets::MountAccess::Cow)
1368        .map(|mount| mount.source.clone())
1369        .collect::<Vec<_>>();
1370    if overlaid.is_empty() {
1371        return Vec::new();
1372    }
1373    // `stat` on this host cannot see what a VM-hosted Docker daemon mounts.
1374    if matches!(target, targets::TargetTemplate::LocalDocker(_)) {
1375        match targets::local_docker_vm_share(executor) {
1376            Ok(Some(reason)) => {
1377                return mounts
1378                    .iter_mut()
1379                    .filter(|mount| mount.access == targets::MountAccess::Cow)
1380                    .map(|mount| {
1381                        mount.access = mount.access.without_overlay();
1382                        format!(
1383                            "Mounted {} read-only: {reason}, which cannot back the \
1384                             copy-on-write overlay.",
1385                            mount.source.display()
1386                        )
1387                    })
1388                    .collect();
1389            }
1390            Ok(None) => {}
1391            Err(error) => tracing::warn!(
1392                error = format!("{error:#}"),
1393                "could not identify the Docker daemon platform; probing this host's filesystems"
1394            ),
1395        }
1396    }
1397    let filesystems = match targets::probe_filesystem_types(ssh, &overlaid, executor) {
1398        Ok(filesystems) => filesystems,
1399        Err(error) => {
1400            tracing::warn!(
1401                error = format!("{error:#}"),
1402                "could not probe attached-directory filesystems; preserving overlay mounts"
1403            );
1404            return vec![format!(
1405                "Could not read the filesystem under the attached directories, so they keep the \
1406                 copy-on-write overlay: {error:#}"
1407            )];
1408        }
1409    };
1410    let mut notices = Vec::new();
1411    for (mount, filesystem) in mounts
1412        .iter_mut()
1413        .filter(|mount| mount.access == targets::MountAccess::Cow)
1414        .zip(filesystems)
1415    {
1416        let Some(reason) = targets::overlay_unsupported_filesystem(&filesystem) else {
1417            continue;
1418        };
1419        mount.access = mount.access.without_overlay();
1420        notices.push(format!(
1421            "Mounted {} read-only: the overlay is unreliable on {filesystem} ({reason}).",
1422            mount.source.display()
1423        ));
1424    }
1425    notices
1426}
1427
1428/// The image users already probed, keyed by container host and image
1429/// reference. An image's
1430/// configured user does not change under a fixed reference, and reading it
1431/// costs a container start, so each daemon asks a host once.
1432static IMAGE_USERS: std::sync::LazyLock<std::sync::Mutex<BTreeMap<String, targets::ImageUser>>> =
1433    std::sync::LazyLock::new(std::sync::Mutex::default);
1434
1435/// The uid and gid a Podman session container maps onto the host user.
1436///
1437/// Only Podman is asked: Docker and Apple's `container` engine are left with
1438/// their own defaults. A probe that cannot answer is not a launch failure —
1439/// the container falls back to plain `--userns=keep-id`, which maps the
1440/// image's default user, and the user is told what happened.
1441pub(super) fn podman_image_user(
1442    target: &targets::TargetTemplate,
1443    executor: &impl CommandExecutor,
1444) -> Option<targets::ImageUser> {
1445    let (ssh, container) = match target {
1446        targets::TargetTemplate::LocalPodman(container) => (None, container),
1447        targets::TargetTemplate::SshPodman { ssh, container } => (Some(ssh), container),
1448        _ => return None,
1449    };
1450    let image = container.image.as_str();
1451    let key = format!(
1452        "{}|{image}",
1453        ssh.map_or("local", |ssh| ssh.destination.as_str())
1454    );
1455    if let Some(cached) = IMAGE_USERS.lock().expect("image user cache").get(&key) {
1456        return Some(*cached);
1457    }
1458    // The probe starts a container, so on Podman this is where a missing image
1459    // is actually downloaded. Wait for the daemon's own download instead of
1460    // starting a second one.
1461    match crate::image_pull_gate::with_image_ready(target, executor, || {
1462        targets::probe_image_user(ssh, container, executor)
1463    }) {
1464        Ok(user) => {
1465            IMAGE_USERS
1466                .lock()
1467                .expect("image user cache")
1468                .insert(key, user);
1469            Some(user)
1470        }
1471        Err(error) => {
1472            tracing::warn!(
1473                image,
1474                error = format!("{error:#}"),
1475                "could not read the container image user; keeping Podman's default user mapping"
1476            );
1477            executor.notify_notice(&format!(
1478                "Could not read the user of image {image}, so the container runs with Podman's \
1479                 default user mapping and may not be able to write to an attached directory: \
1480                 {error:#}"
1481            ));
1482            None
1483        }
1484    }
1485}
1486
1487/// Reports every command an installer issues as one launch stage, so progress
1488/// stays accurate without threading the stage through each `CommandSpec`.
1489/// A command that already names a stage keeps it.
1490pub(super) struct StagedExecutor<'a, E: CommandExecutor> {
1491    inner: &'a E,
1492    stage: ProvisionStage,
1493    _guard: ProvisionStageGuard<'a, E>,
1494}
1495
1496impl<'a, E: CommandExecutor> StagedExecutor<'a, E> {
1497    pub(crate) fn new(inner: &'a E, stage: ProvisionStage) -> Self {
1498        Self {
1499            inner,
1500            stage,
1501            _guard: ProvisionStageGuard::new(inner, stage),
1502        }
1503    }
1504
1505    fn staged(&self, command: &CommandSpec) -> CommandSpec {
1506        if command.stage.is_some() {
1507            return command.clone();
1508        }
1509        command.clone().stage(self.stage)
1510    }
1511}
1512
1513impl<E: CommandExecutor> CommandExecutor for StagedExecutor<'_, E> {
1514    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1515        self.inner.execute(&self.staged(command))
1516    }
1517
1518    fn cancellation_requested(&self) -> bool {
1519        self.inner.cancellation_requested()
1520    }
1521
1522    fn stage_started(&self, stage: ProvisionStage) {
1523        self.inner.stage_started(stage);
1524    }
1525
1526    fn stage_finished(&self, stage: ProvisionStage) {
1527        self.inner.stage_finished(stage);
1528    }
1529
1530    fn notify_notice(&self, notice: &str) {
1531        self.inner.notify_notice(notice);
1532    }
1533
1534    fn execute_with_stdin(
1535        &self,
1536        command: &CommandSpec,
1537        input: &mut (dyn std::io::Read + Send),
1538    ) -> Result<CommandOutput> {
1539        self.inner.execute_with_stdin(&self.staged(command), input)
1540    }
1541}
1542
1543fn execute_checked_with_stdin(
1544    executor: &impl CommandExecutor,
1545    command: &CommandSpec,
1546    input: &mut (dyn std::io::Read + Send),
1547) -> Result<CommandOutput> {
1548    let output = executor.execute_with_stdin(command, input)?;
1549    if output.status != 0 {
1550        bail!(
1551            "{} failed with status {}: {}",
1552            command.purpose,
1553            output.status,
1554            String::from_utf8_lossy(&output.stderr)
1555        );
1556    }
1557    Ok(output)
1558}
1559
1560pub(super) fn install_inherited_git_settings(
1561    executor: &impl CommandExecutor,
1562    locator: &targets::TargetLocator,
1563    session_id: &str,
1564) -> Result<()> {
1565    let settings = if inherits_controller_git_settings(locator) {
1566        controller_git_settings()?
1567    } else {
1568        BTreeMap::new()
1569    };
1570    for command in inherited_git_setting_commands(locator, session_id, settings)? {
1571        execute_checked(executor, command)?;
1572    }
1573    Ok(())
1574}
1575
1576fn inherits_controller_git_settings(locator: &targets::TargetLocator) -> bool {
1577    !matches!(
1578        locator,
1579        targets::TargetLocator::LocalBare { .. } | targets::TargetLocator::SshBare { .. }
1580    )
1581}
1582
1583fn inherited_git_setting_commands(
1584    locator: &targets::TargetLocator,
1585    session_id: &str,
1586    settings: BTreeMap<String, String>,
1587) -> Result<Vec<CommandSpec>> {
1588    if matches!(locator, targets::TargetLocator::SshBare { .. }) {
1589        return Ok(Vec::new());
1590    }
1591    settings
1592        .into_iter()
1593        .map(|(key, value)| {
1594            targets::command_on_locator(
1595                locator,
1596                session_id,
1597                vec![
1598                    "git".into(),
1599                    "config".into(),
1600                    "--global".into(),
1601                    "--replace-all".into(),
1602                    "--".into(),
1603                    key.clone(),
1604                    value,
1605                ],
1606                format!("inherit Git setting {key}"),
1607            )
1608        })
1609        .collect()
1610}
1611
1612fn controller_git_settings() -> Result<BTreeMap<String, String>> {
1613    let output = match Command::new("git")
1614        .args(["config", "--global", "--includes", "--null", "--list"])
1615        .stdin(Stdio::null())
1616        .output()
1617    {
1618        Ok(output) => output,
1619        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
1620        Err(error) => return Err(error).context("read controller Git configuration"),
1621    };
1622    if !output.status.success() {
1623        bail!(
1624            "read controller Git configuration failed with status {}: {}",
1625            output.status,
1626            String::from_utf8_lossy(&output.stderr).trim()
1627        );
1628    }
1629    parse_inherited_git_settings(&output.stdout)
1630}
1631
1632fn parse_inherited_git_settings(output: &[u8]) -> Result<BTreeMap<String, String>> {
1633    let mut settings = BTreeMap::new();
1634    for entry in output
1635        .split(|byte| *byte == 0)
1636        .filter(|entry| !entry.is_empty())
1637    {
1638        let entry = std::str::from_utf8(entry).context("decode controller Git configuration")?;
1639        let (key, value) = entry
1640            .split_once('\n')
1641            .with_context(|| format!("controller Git returned malformed entry {entry:?}"))?;
1642        let key = key.to_ascii_lowercase();
1643        if INHERITED_GIT_SETTINGS.contains(&key.as_str()) {
1644            settings.insert(key, value.to_owned());
1645        }
1646    }
1647    Ok(settings)
1648}
1649
1650#[cfg(test)]
1651mod tests;