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