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