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.
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)?;
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 failure_disposition == ProvisioningFailureDisposition::Preserve {
501                Some(super::network_git::checkpoint_bundle(&session)?)
502            } else {
503                Some(backend_bundle(
504                    self.config
505                        .bundles
506                        .get(&session.bundle_id)
507                        .context("session bundle is missing")?,
508                    executor,
509                )?)
510            };
511            let container_github_token =
512                github_token.filter(|_| configure_github_token_environment(&mut target));
513            if container_github_token.is_some()
514                && let Some(bundle) = bundle.as_mut()
515            {
516                use_github_https_urls(bundle);
517            }
518            preflight_target(template, executor, TargetCheck::Launch)?;
519            let prepared_cache = bundle.as_mut().and_then(|bundle| {
520                git_cache::prepare(
521                    &target,
522                    session_id,
523                    bundle,
524                    &mut runtime_mounts,
525                    container_github_token,
526                    executor,
527                )
528            });
529            // Mounts are fixed when the container is created, so the build
530            // cache is decided here, before the provisioning plan is built.
531            let build_cache = super::mbx::prepare(
532                &target,
533                &self.config.build_cache,
534                &session,
535                bundle.as_ref(),
536                prepared_cache.as_ref(),
537                &mut runtime_mounts,
538                executor,
539            );
540            let provision = if let Some(project_directory) = &session.project_directory {
541                targets::provision_bare_project_plan(
542                    &target,
543                    session_id,
544                    &project_directory.to_string_lossy(),
545                )
546            } else {
547                bundle
548                    .as_ref()
549                    .context("project bundle disappeared during provisioning")
550                    .and_then(|bundle| {
551                        targets::provision_plan(
552                            &target,
553                            session_id,
554                            bundle,
555                            &runtime_mounts,
556                            image_user,
557                            session.container_workspace.as_deref(),
558                        )
559                    })
560            };
561            let mut provision = match provision {
562                Ok(provision) => provision,
563                Err(error) => {
564                    if let Some(cache) = &prepared_cache {
565                        let _ = cache.cleanup(executor);
566                    }
567                    return Err(error);
568                }
569            };
570            if let Some(token) = container_github_token
571                && let Err(error) =
572                    provision.provide_target_environment_secret(&target, "GH_TOKEN", token)
573            {
574                if let Some(cache) = &prepared_cache {
575                    let _ = cache.cleanup(executor);
576                }
577                return Err(error);
578            }
579
580            let started = Instant::now();
581            let result =
582                provision_target_creation(&provision, &target, session_id, executor, |outputs| {
583                    locator_after_provision(
584                        template,
585                        &target,
586                        session_id,
587                        outputs.first(),
588                        executor,
589                    )
590                })
591                .map(|(locator, remainder)| (locator, remainder, bundle, build_cache));
592            if result.is_err()
593                && let Some(cache) = &prepared_cache
594            {
595                if let Some(locator) = provisioned_locator(&target, session_id, None) {
596                    let _ = targets::close_plan(&locator, session_id)
597                        .and_then(|plan| plan.execute(executor).map(|_| ()));
598                } else {
599                    let _ = cache.cleanup(executor);
600                }
601            }
602            tracing::debug!(
603                session_id,
604                elapsed_ms = started.elapsed().as_millis(),
605                "provisioning plan execution completed"
606            );
607            result
608        })();
609        let result = match result {
610            Err(error)
611                if created_worktree
612                    && failure_disposition == ProvisioningFailureDisposition::Discard =>
613            {
614                return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
615            }
616            Err(error) if failure_disposition == ProvisioningFailureDisposition::Preserve => {
617                Err(error)
618            }
619            Err(error) => {
620                // This arm (no managed worktree to unwind) is the one a
621                // provisioning failure such as a dropped target connection
622                // hits; record the diagnostic and the session-id log here too.
623                let detail = note_new_session_launch_failure(session_id, &error);
624                {
625                    let record = self.state.sessions.get_mut(session_id).unwrap();
626                    record.state = SessionState::Error;
627                    record.target = None;
628                    record.updated_at = super::now();
629                    record.last_error = Some(format!("session provisioning failed: {detail}"));
630                }
631                return match self.persist_session_state(session_id) {
632                    Ok(()) => Err(error),
633                    Err(persistence_error) => Err(error.context(format!(
634                        "persist removal of failed provisioning session {session_id}: {persistence_error:#}"
635                    ))),
636                };
637            }
638            Ok((locator, remainder, bundle, build_cache)) => {
639                apply_new_session_provisioning_result(&mut self.state, session_id, Ok(locator))?;
640                self.state
641                    .sessions
642                    .get_mut(session_id)
643                    .expect("session retained after provisioning")
644                    .build_cache = build_cache;
645                let session = &self.state.sessions[session_id];
646                let backend = backend_locator(
647                    session
648                        .target
649                        .as_ref()
650                        .context("provisioned target disappeared")?,
651                    session,
652                    &self.config,
653                )?;
654                if matches!(backend, targets::TargetLocator::AwsEc2 { .. }) {
655                    targets::provision_on_locator_plan(
656                        &backend,
657                        session_id,
658                        bundle
659                            .as_ref()
660                            .context("AWS provisioning requires a project bundle")?,
661                    )
662                } else {
663                    Ok(remainder)
664                }
665            }
666        };
667        let result = match result {
668            Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
669                return Err(self.rollback_failed_new_session(session_id, error, executor)?);
670            }
671            result => result,
672        };
673        if result.is_ok()
674            && let Some(session) = self.state.sessions.get(session_id)
675            && let Some(directory) = session
676                .managed_worktree
677                .as_ref()
678                .map(|worktree| worktree.source_project_directory.clone())
679                .or_else(|| session.project_directory.clone())
680            && let Some(template) = self.config.targets.get(&session.target_template_id)
681        {
682            let host = match template {
683                TargetTemplate::LocalBare => Some("local"),
684                TargetTemplate::SshBare { ssh, .. } => Some(ssh.host.as_str()),
685                _ => None,
686            };
687            if let Some(host) = host {
688                self.state.remember_project_directory(host, &directory);
689                crate::database::remember_project_directory(host, &directory)?;
690            }
691        }
692        self.persist_session_state(session_id)?;
693        result
694    }
695
696    pub fn mark_worker_connected(
697        &mut self,
698        session_id: &str,
699        native_session_id: Option<String>,
700    ) -> Result<()> {
701        let session = self
702            .state
703            .sessions
704            .get(session_id)
705            .with_context(|| format!("unknown session {session_id}"))?;
706        if session.target.is_none() {
707            bail!("session {session_id} has no provisioned target");
708        }
709        let updated_at = now();
710        crate::database::mark_session_worker_connected(
711            session_id,
712            native_session_id.as_deref(),
713            &updated_at,
714        )?;
715        let session = self
716            .state
717            .sessions
718            .get_mut(session_id)
719            .expect("session disappeared after its worker connection was saved");
720        session.state = SessionState::Running;
721        if native_session_id.is_some() {
722            session.native_session_id = native_session_id;
723        }
724        session.updated_at = updated_at;
725        session.last_error = None;
726        Ok(())
727    }
728
729    fn install_worker_payload(
730        &self,
731        session_id: &str,
732        executor: &impl CommandExecutor,
733    ) -> Result<(targets::TargetLocator, String)> {
734        // Worker/profile installation is independent of repository cloning.
735        let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
736        let (backend, worker_root) = self.worker_placement(session_id)?;
737        self.prepare_worker_files(session_id, &backend, &worker_root, syncing)?;
738        install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
739        Ok((backend, worker_root))
740    }
741
742    async fn connect_and_start_worker(
743        &self,
744        session_id: &str,
745        executor: &impl CommandExecutor,
746        backend: &targets::TargetLocator,
747        worker_root: &str,
748        initialize_workspace: bool,
749    ) -> Result<Option<String>> {
750        let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
751        if initialize_workspace {
752            install_inherited_git_settings(executor, backend, session_id)?;
753            self.initialize_network_workspaces(session_id, backend, syncing)?;
754        }
755        let session = self
756            .state
757            .sessions
758            .get(session_id)
759            .with_context(|| format!("unknown session {session_id}"))?;
760        let profile = self
761            .config
762            .profiles
763            .get(&session.last_profile)
764            .with_context(|| format!("unknown profile {}", session.last_profile))?;
765        let readiness_stage = bridge_readiness_stage(profile);
766        let reconnect = &targets::reconnect_plan(backend, session_id)?.commands[0];
767        let readiness = async {
768            let mut relay = {
769                let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
770                start_worker(executor, backend, worker_root)?;
771                connect_started_worker(reconnect, session_id, executor, backend, worker_root)
772                    .await?
773            };
774            let native_session_id =
775                wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
776            Ok(Some(native_session_id))
777        }
778        .await;
779        match readiness {
780            Ok(native_session_id) => Ok(native_session_id),
781            Err(error) => Err(worker_probe_diagnosis(
782                executor,
783                backend,
784                worker_root,
785                error,
786            )),
787        }
788    }
789}
790
791const MAX_LAUNCH_DIAGNOSTIC_BYTES: usize = 64 * 1024;
792
793const RETAINED_LAUNCH_DIAGNOSTICS: usize = 20;
794
795/// Record a failed new-session launch consistently across every failure arm:
796/// log the failure with the session id (so `logs/mj-*.log` names the session)
797/// and save the local diagnostic file. Returns the underlying error chain
798/// annotated with the diagnostic path, which the caller stores in `last_error`
799/// so the reason travels to `mj sessions`, `mj wait`, and `mj events`.
800pub(super) fn note_new_session_launch_failure(session_id: &str, error: &anyhow::Error) -> String {
801    note_new_session_launch_failure_in(&data_dir().join("diagnostics"), session_id, error)
802}
803
804fn note_new_session_launch_failure_in(
805    directory: &Path,
806    session_id: &str,
807    error: &anyhow::Error,
808) -> String {
809    let original = format!("{error:#}");
810    tracing::warn!(session_id, error = %original, "session launch failed");
811    match persist_launch_failure_to(directory, session_id, &original) {
812        Ok(path) => format!("{original}; full diagnostic saved to {}", path.display()),
813        Err(save_error) => {
814            format!("{original}; saving the local diagnostic failed: {save_error:#}")
815        }
816    }
817}
818
819fn persist_launch_failure_to(directory: &Path, session_id: &str, detail: &str) -> Result<PathBuf> {
820    mj_core::config::validate_id("session", session_id)?;
821    std::fs::create_dir_all(directory).with_context(|| {
822        format!(
823            "create launch diagnostics directory {}",
824            directory.display()
825        )
826    })?;
827    #[cfg(unix)]
828    {
829        use std::os::unix::fs::PermissionsExt;
830        std::fs::set_permissions(directory, std::fs::Permissions::from_mode(0o700))?;
831    }
832    let path = directory.join(format!("{session_id}-launch-error.txt"));
833    let detail = bounded_launch_diagnostic(detail);
834    let body = format!(
835        "Mjolnir session launch failure\nsession: {session_id}\nat: {}\n\n{detail}\n",
836        now()
837    );
838    atomic_write(&path, body.as_bytes())?;
839    prune_launch_diagnostics(directory)?;
840    Ok(path)
841}
842
843fn bounded_launch_diagnostic(detail: &str) -> String {
844    if detail.len() <= MAX_LAUNCH_DIAGNOSTIC_BYTES {
845        return detail.to_owned();
846    }
847    let mut head_end = MAX_LAUNCH_DIAGNOSTIC_BYTES / 4;
848    while !detail.is_char_boundary(head_end) {
849        head_end -= 1;
850    }
851    let tail_bytes = MAX_LAUNCH_DIAGNOSTIC_BYTES - head_end;
852    let mut tail_start = detail.len() - tail_bytes;
853    while !detail.is_char_boundary(tail_start) {
854        tail_start += 1;
855    }
856    format!(
857        "{}\n\n[... launch diagnostic truncated ...]\n\n{}",
858        &detail[..head_end],
859        &detail[tail_start..]
860    )
861}
862
863fn prune_launch_diagnostics(directory: &Path) -> Result<()> {
864    let mut diagnostics = Vec::new();
865    for entry in std::fs::read_dir(directory)? {
866        let entry = entry?;
867        if !entry
868            .file_name()
869            .to_str()
870            .is_some_and(|name| name.ends_with("-launch-error.txt"))
871        {
872            continue;
873        }
874        diagnostics.push((entry.metadata()?.modified()?, entry.path()));
875    }
876    diagnostics.sort_by_key(|entry| std::cmp::Reverse(entry.0));
877    for (_, path) in diagnostics.into_iter().skip(RETAINED_LAUNCH_DIAGNOSTICS) {
878        std::fs::remove_file(&path)
879            .with_context(|| format!("prune old launch diagnostic {}", path.display()))?;
880    }
881    Ok(())
882}
883
884fn apply_new_session_provisioning_result(
885    state: &mut State,
886    session_id: &str,
887    result: Result<TargetLocator>,
888) -> Result<()> {
889    match result {
890        Ok(locator) => {
891            let record = state.sessions.get_mut(session_id).unwrap();
892            record.target = Some(locator);
893            // Provisioning has completed, but Running is reserved for a
894            // successful worker handshake.
895            record.state = SessionState::Disconnected;
896            record.updated_at = now();
897            record.last_error = None;
898            Ok(())
899        }
900        Err(error) => {
901            let record = state.sessions.get_mut(session_id).unwrap();
902            record.state = SessionState::Error;
903            record.target = None;
904            record.updated_at = now();
905            record.last_error = Some(format!("session provisioning failed: {error:#}"));
906            Err(error)
907        }
908    }
909}
910
911/// Mark a new session's launch failed while it still names its target, so the
912/// record is true while the rollback removes that target.
913fn apply_failed_new_session_launch(state: &mut State, session_id: &str, original_error: &str) {
914    let record = state.sessions.get_mut(session_id).unwrap();
915    record.state = SessionState::Error;
916    record.updated_at = now();
917    record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
918}
919
920pub(super) fn apply_failed_new_session_rollback(
921    state: &mut State,
922    session_id: &str,
923    original_error: &str,
924    cleanup_error: Option<String>,
925) -> anyhow::Error {
926    match cleanup_error {
927        None => {
928            let record = state.sessions.get_mut(session_id).unwrap();
929            record.state = SessionState::Error;
930            record.target = None;
931            record.updated_at = now();
932            record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
933            anyhow::anyhow!("{original_error}; partial target removed and failed session retained")
934        }
935        Some(cleanup_error) => {
936            let failure = format!(
937                "{original_error}; cleanup of the failed session target failed: {cleanup_error}"
938            );
939            let record = state.sessions.get_mut(session_id).unwrap();
940            record.state = SessionState::Error;
941            record.updated_at = now();
942            record.last_error = Some(format!("worker bootstrap failed: {failure}"));
943            anyhow::anyhow!(failure)
944        }
945    }
946}
947
948pub(super) fn install_attached_resources(
949    state: &State,
950    session_id: &str,
951    backend: &targets::TargetLocator,
952    worker_root: &str,
953    executor: &impl CommandExecutor,
954) -> Result<()> {
955    let targets::TargetLocator::AwsEc2 { .. } = backend else {
956        return Ok(());
957    };
958    let session = state
959        .sessions
960        .get(session_id)
961        .with_context(|| format!("unknown session {session_id}"))?;
962    if session.additional_mounts.is_empty() {
963        return Ok(());
964    }
965    for resource in &session.additional_mounts {
966        let install = targets::command_on_locator(
967            backend,
968            session_id,
969            vec![
970                format!("{worker_root}/hel"),
971                "worker".into(),
972                "install-resource".into(),
973                "--destination".into(),
974                resource.destination.to_string_lossy().into_owned(),
975            ],
976            "stream attached resource",
977        )?;
978        mj_checkpoint::resources::stream_resource(&resource.source, |stream| {
979            execute_checked_with_stdin(executor, &install, stream).map(|_| ())
980        })
981        .with_context(|| format!("stream attached resource {}", resource.source.display()))?;
982    }
983    Ok(())
984}
985
986/// Run two independent target setup lanes at the same time and wait for both.
987/// The first lane's failure wins deterministically when both fail, and neither
988/// lane is abandoned while it may still own a transfer or subprocess.
989pub(super) fn execute_concurrent_lanes<A: Send, B: Send>(
990    first: impl FnOnce() -> Result<A> + Send,
991    second: impl FnOnce() -> Result<B> + Send,
992) -> Result<(A, B)> {
993    std::thread::scope(|scope| {
994        let second = scope.spawn(second);
995        let first = first();
996        let second = second.join().unwrap_or_else(|panic| {
997            Err(anyhow::anyhow!(
998                "concurrent target lane panicked: {}",
999                targets::command_thread_panic_message(panic.as_ref())
1000            ))
1001        });
1002        match (first, second) {
1003            (Err(error), _) => Err(error),
1004            (Ok(_), Err(error)) => Err(error),
1005            (Ok(first), Ok(second)) => Ok((first, second)),
1006        }
1007    })
1008}
1009
1010fn execute_repository_setup(
1011    plan: &targets::CommandPlan,
1012    executor: &(impl CommandExecutor + Sync),
1013) -> Result<()> {
1014    if plan.commands.is_empty() {
1015        return Ok(());
1016    }
1017    let _cloning = ProvisionStageGuard::new(executor, ProvisionStage::Cloning);
1018    plan.execute_concurrent(executor).map(|_| ())
1019}
1020
1021/// Run a provisioning plan and discover the locator it produced, tearing the
1022/// target down again if anything after its creation fails.
1023///
1024/// Creation is the boundary that matters. A step that fails before the target
1025/// exists has left nothing behind; every failure after it — a later plan step
1026/// or locator discovery — owns a target no session record will point at.
1027#[cfg(test)]
1028fn provision_target(
1029    plan: &targets::CommandPlan,
1030    target: &targets::TargetTemplate,
1031    session_id: &str,
1032    executor: &(impl CommandExecutor + Sync),
1033    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1034) -> Result<TargetLocator> {
1035    let Some((creation, remainder)) = plan.split_at_target_creation() else {
1036        // Nothing this plan runs can leave a target behind.
1037        return discover(&plan.execute_concurrent(executor)?);
1038    };
1039    let mut outputs = creation.execute_concurrent(executor)?;
1040    let result = match remainder.execute_concurrent(executor) {
1041        Ok(rest) => {
1042            outputs.extend(rest);
1043            discover(&outputs)
1044        }
1045        Err(error) => Err(error),
1046    };
1047    result.map_err(|error| {
1048        match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1049            Some(note) => error.context(note),
1050            None => error,
1051        }
1052    })
1053}
1054
1055/// Bring the target into existence and return the commands that populate its
1056/// repositories. The caller may overlap that remainder with worker/profile
1057/// installation once it has persisted the discovered locator.
1058fn provision_target_creation(
1059    plan: &targets::CommandPlan,
1060    target: &targets::TargetTemplate,
1061    session_id: &str,
1062    executor: &(impl CommandExecutor + Sync),
1063    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1064) -> Result<(TargetLocator, targets::CommandPlan)> {
1065    let Some((creation, remainder)) = plan.split_at_target_creation() else {
1066        // Nothing this plan runs can leave a target behind, so its commands
1067        // must still finish before the locator is usable.
1068        let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1069            plan.execute_concurrent(executor)
1070        })?;
1071        return discover(&outputs).map(|locator| {
1072            (
1073                locator,
1074                targets::CommandPlan {
1075                    description: plan.description.clone(),
1076                    commands: Vec::new(),
1077                },
1078            )
1079        });
1080    };
1081    // Creating the container is what downloads the image on Docker, and the
1082    // probe just before it is what downloads it on Podman. Either way, a
1083    // background download of the same image must finish first.
1084    let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1085        creation.execute_concurrent(executor)
1086    })?;
1087    discover(&outputs)
1088        .map(|locator| (locator, remainder))
1089        .map_err(|error| {
1090            match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1091                Some(note) => error.context(note),
1092                None => error,
1093            }
1094        })
1095}
1096
1097/// Best-effort teardown of a target whose creation succeeded but whose
1098/// provisioning failed before a locator was recorded. Returns a note
1099/// describing what happened for inclusion in the session error.
1100///
1101/// The teardown is the session's own close plan, so a failed launch and an
1102/// ordinary close can never disagree about what removing a target means.
1103fn cleanup_failed_provision(
1104    target: &targets::TargetTemplate,
1105    session_id: &str,
1106    create_output: Option<&CommandOutput>,
1107    executor: &impl CommandExecutor,
1108) -> Option<String> {
1109    let locator = provisioned_locator(target, session_id, create_output)?;
1110    let leak = format!(
1111        "the resource may still exist; find it via its dev.mj.session={session_id} label/tag"
1112    );
1113    let plan = match targets::close_plan(&locator, session_id) {
1114        Ok(plan) => plan,
1115        Err(error) => {
1116            tracing::warn!(
1117                session_id,
1118                error = format!("{error:#}"),
1119                "could not build provisioning cleanup plan"
1120            );
1121            return Some(format!("cleanup FAILED: {error:#}; {leak}"));
1122        }
1123    };
1124    let purpose = plan
1125        .commands
1126        .iter()
1127        .map(|command| command.purpose.clone())
1128        .collect::<Vec<_>>()
1129        .join("; ");
1130    let Err(error) = plan.execute(executor) else {
1131        return Some(format!("cleanup succeeded: {purpose}"));
1132    };
1133    match targets::cleanup_target_is_confirmed_absent(&locator, session_id, executor) {
1134        Ok(true) => Some(format!("cleanup succeeded: {purpose}")),
1135        Ok(false) => {
1136            tracing::warn!(
1137                session_id,
1138                error = format!("{error:#}"),
1139                "provisioning cleanup failed and the target may still exist"
1140            );
1141            Some(format!("cleanup FAILED ({purpose}): {error:#}; {leak}"))
1142        }
1143        Err(confirm_error) => {
1144            tracing::warn!(
1145                session_id,
1146                error = format!("{confirm_error:#}"),
1147                "could not confirm whether the failed provisioning target was removed"
1148            );
1149            Some(format!(
1150                "cleanup FAILED ({purpose}): {error:#}; checking whether it was removed also failed: {confirm_error:#}; {leak}"
1151            ))
1152        }
1153    }
1154}
1155
1156/// The locator a provisioning plan's creating command brought into existence.
1157///
1158/// Every target but AWS is named before its plan runs; an EC2 instance
1159/// reports its own ID in the launch response.
1160fn provisioned_locator(
1161    target: &targets::TargetTemplate,
1162    session_id: &str,
1163    create_output: Option<&CommandOutput>,
1164) -> Option<targets::TargetLocator> {
1165    let container_id = || targets::resource_name(session_id).ok();
1166    Some(match target {
1167        // A bare project directory belongs to the user: provisioning creates
1168        // nothing that a failure could leak.
1169        targets::TargetTemplate::LocalBare => return None,
1170        targets::TargetTemplate::LocalPodman(container) => targets::TargetLocator::LocalPodman {
1171            borrowed_from: None,
1172            container_id: container_id()?,
1173            workspace_storage: targets::podman_workspace_locator(container, session_id).ok()?,
1174        },
1175        targets::TargetTemplate::LocalDocker(_) => targets::TargetLocator::LocalDocker {
1176            borrowed_from: None,
1177            container_id: container_id()?,
1178        },
1179        targets::TargetTemplate::AppleContainer(_) => targets::TargetLocator::AppleContainer {
1180            borrowed_from: None,
1181            container_id: container_id()?,
1182        },
1183        targets::TargetTemplate::SshPodman { ssh, container } => {
1184            targets::TargetLocator::SshPodman {
1185                borrowed_from: None,
1186                ssh: ssh.clone(),
1187                container_id: container_id()?,
1188                workspace_storage: targets::podman_workspace_locator(container, session_id).ok()?,
1189            }
1190        }
1191        targets::TargetTemplate::SshDocker { ssh, .. } => targets::TargetLocator::SshDocker {
1192            borrowed_from: None,
1193            ssh: ssh.clone(),
1194            container_id: container_id()?,
1195        },
1196        targets::TargetTemplate::SshBare { ssh, .. } => targets::TargetLocator::SshBare {
1197            ssh: ssh.clone(),
1198            workspace: targets::workspace_for(target, session_id).ok()?,
1199            worker_id: None,
1200        },
1201        targets::TargetTemplate::AwsEc2(aws) => targets::TargetLocator::AwsEc2 {
1202            profile: aws.profile.clone(),
1203            region: aws.region.clone(),
1204            instance_id: serde_json::from_slice::<serde_json::Value>(&create_output?.stdout)
1205                .ok()?
1206                .pointer("/Instances/0/InstanceId")?
1207                .as_str()?
1208                .to_owned(),
1209            ssh: aws.ssh.clone(),
1210            workspace: targets::workspace_for(target, session_id).ok()?,
1211        },
1212    })
1213}
1214
1215/// Attach read-only whatever the selected container overlay cannot hold, and
1216/// say so.
1217///
1218/// The filesystem is probed on the host that runs the container, because that
1219/// is where the overlay would be built. A probe that cannot answer leaves the
1220/// overlay alone: a failed probe is no evidence of an unsupported filesystem,
1221/// and refusing to provision over one would cost the user their session.
1222///
1223/// Apple's `container` engine already mounts every extra directory read-only,
1224/// and EC2 copies the directory instead of mounting it, so neither is probed.
1225pub(super) fn enforce_overlay_capable_mounts(
1226    target: &targets::TargetTemplate,
1227    mounts: &mut [targets::AdditionalMount],
1228    executor: &impl CommandExecutor,
1229) -> Vec<String> {
1230    let ssh = match target {
1231        targets::TargetTemplate::LocalPodman(_) | targets::TargetTemplate::LocalDocker(_) => None,
1232        targets::TargetTemplate::SshPodman { ssh, .. }
1233        | targets::TargetTemplate::SshDocker { ssh, .. } => Some(ssh),
1234        _ => return Vec::new(),
1235    };
1236    let overlaid = mounts
1237        .iter()
1238        .filter(|mount| mount.access == targets::MountAccess::Cow)
1239        .map(|mount| mount.source.clone())
1240        .collect::<Vec<_>>();
1241    if overlaid.is_empty() {
1242        return Vec::new();
1243    }
1244    let filesystems = match targets::probe_filesystem_types(ssh, &overlaid, executor) {
1245        Ok(filesystems) => filesystems,
1246        Err(error) => {
1247            tracing::warn!(
1248                error = format!("{error:#}"),
1249                "could not probe attached-directory filesystems; preserving overlay mounts"
1250            );
1251            return vec![format!(
1252                "Could not read the filesystem under the attached directories, so they keep the \
1253                 copy-on-write overlay: {error:#}"
1254            )];
1255        }
1256    };
1257    let mut notices = Vec::new();
1258    for (mount, filesystem) in mounts
1259        .iter_mut()
1260        .filter(|mount| mount.access == targets::MountAccess::Cow)
1261        .zip(filesystems)
1262    {
1263        let Some(reason) = targets::overlay_unsupported_filesystem(&filesystem) else {
1264            continue;
1265        };
1266        mount.access = mount.access.without_overlay();
1267        notices.push(format!(
1268            "Mounted {} read-only: the overlay is unreliable on {filesystem} ({reason}).",
1269            mount.source.display()
1270        ));
1271    }
1272    notices
1273}
1274
1275/// The image users already probed, keyed by container host and image
1276/// reference. An image's
1277/// configured user does not change under a fixed reference, and reading it
1278/// costs a container start, so each daemon asks a host once.
1279static IMAGE_USERS: std::sync::LazyLock<std::sync::Mutex<BTreeMap<String, targets::ImageUser>>> =
1280    std::sync::LazyLock::new(std::sync::Mutex::default);
1281
1282/// The uid and gid a Podman session container maps onto the host user.
1283///
1284/// Only Podman is asked: Docker and Apple's `container` engine are left with
1285/// their own defaults. A probe that cannot answer is not a launch failure —
1286/// the container falls back to plain `--userns=keep-id`, which maps the
1287/// image's default user, and the user is told what happened.
1288pub(super) fn podman_image_user(
1289    target: &targets::TargetTemplate,
1290    executor: &impl CommandExecutor,
1291) -> Option<targets::ImageUser> {
1292    let (ssh, container) = match target {
1293        targets::TargetTemplate::LocalPodman(container) => (None, container),
1294        targets::TargetTemplate::SshPodman { ssh, container } => (Some(ssh), container),
1295        _ => return None,
1296    };
1297    let image = container.image.as_str();
1298    let key = format!(
1299        "{}|{image}",
1300        ssh.map_or("local", |ssh| ssh.destination.as_str())
1301    );
1302    if let Some(cached) = IMAGE_USERS.lock().expect("image user cache").get(&key) {
1303        return Some(*cached);
1304    }
1305    // The probe starts a container, so on Podman this is where a missing image
1306    // is actually downloaded. Wait for the daemon's own download instead of
1307    // starting a second one.
1308    match crate::image_pull_gate::with_image_ready(target, executor, || {
1309        targets::probe_image_user(ssh, container, executor)
1310    }) {
1311        Ok(user) => {
1312            IMAGE_USERS
1313                .lock()
1314                .expect("image user cache")
1315                .insert(key, user);
1316            Some(user)
1317        }
1318        Err(error) => {
1319            tracing::warn!(
1320                image,
1321                error = format!("{error:#}"),
1322                "could not read the container image user; keeping Podman's default user mapping"
1323            );
1324            executor.notify_notice(&format!(
1325                "Could not read the user of image {image}, so the container runs with Podman's \
1326                 default user mapping and may not be able to write to an attached directory: \
1327                 {error:#}"
1328            ));
1329            None
1330        }
1331    }
1332}
1333
1334/// Reports every command an installer issues as one launch stage, so progress
1335/// stays accurate without threading the stage through each `CommandSpec`.
1336/// A command that already names a stage keeps it.
1337pub(super) struct StagedExecutor<'a, E: CommandExecutor> {
1338    inner: &'a E,
1339    stage: ProvisionStage,
1340    _guard: ProvisionStageGuard<'a, E>,
1341}
1342
1343impl<'a, E: CommandExecutor> StagedExecutor<'a, E> {
1344    pub(crate) fn new(inner: &'a E, stage: ProvisionStage) -> Self {
1345        Self {
1346            inner,
1347            stage,
1348            _guard: ProvisionStageGuard::new(inner, stage),
1349        }
1350    }
1351
1352    fn staged(&self, command: &CommandSpec) -> CommandSpec {
1353        if command.stage.is_some() {
1354            return command.clone();
1355        }
1356        command.clone().stage(self.stage)
1357    }
1358}
1359
1360impl<E: CommandExecutor> CommandExecutor for StagedExecutor<'_, E> {
1361    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1362        self.inner.execute(&self.staged(command))
1363    }
1364
1365    fn cancellation_requested(&self) -> bool {
1366        self.inner.cancellation_requested()
1367    }
1368
1369    fn stage_started(&self, stage: ProvisionStage) {
1370        self.inner.stage_started(stage);
1371    }
1372
1373    fn stage_finished(&self, stage: ProvisionStage) {
1374        self.inner.stage_finished(stage);
1375    }
1376
1377    fn notify_notice(&self, notice: &str) {
1378        self.inner.notify_notice(notice);
1379    }
1380
1381    fn execute_with_stdin(
1382        &self,
1383        command: &CommandSpec,
1384        input: &mut (dyn std::io::Read + Send),
1385    ) -> Result<CommandOutput> {
1386        self.inner.execute_with_stdin(&self.staged(command), input)
1387    }
1388}
1389
1390fn execute_checked_with_stdin(
1391    executor: &impl CommandExecutor,
1392    command: &CommandSpec,
1393    input: &mut (dyn std::io::Read + Send),
1394) -> Result<CommandOutput> {
1395    let output = executor.execute_with_stdin(command, input)?;
1396    if output.status != 0 {
1397        bail!(
1398            "{} failed with status {}: {}",
1399            command.purpose,
1400            output.status,
1401            String::from_utf8_lossy(&output.stderr)
1402        );
1403    }
1404    Ok(output)
1405}
1406
1407pub(super) fn install_inherited_git_settings(
1408    executor: &impl CommandExecutor,
1409    locator: &targets::TargetLocator,
1410    session_id: &str,
1411) -> Result<()> {
1412    let settings = if inherits_controller_git_settings(locator) {
1413        controller_git_settings()?
1414    } else {
1415        BTreeMap::new()
1416    };
1417    for command in inherited_git_setting_commands(locator, session_id, settings)? {
1418        execute_checked(executor, command)?;
1419    }
1420    Ok(())
1421}
1422
1423fn inherits_controller_git_settings(locator: &targets::TargetLocator) -> bool {
1424    !matches!(
1425        locator,
1426        targets::TargetLocator::LocalBare { .. } | targets::TargetLocator::SshBare { .. }
1427    )
1428}
1429
1430fn inherited_git_setting_commands(
1431    locator: &targets::TargetLocator,
1432    session_id: &str,
1433    settings: BTreeMap<String, String>,
1434) -> Result<Vec<CommandSpec>> {
1435    if matches!(locator, targets::TargetLocator::SshBare { .. }) {
1436        return Ok(Vec::new());
1437    }
1438    settings
1439        .into_iter()
1440        .map(|(key, value)| {
1441            targets::command_on_locator(
1442                locator,
1443                session_id,
1444                vec![
1445                    "git".into(),
1446                    "config".into(),
1447                    "--global".into(),
1448                    "--replace-all".into(),
1449                    "--".into(),
1450                    key.clone(),
1451                    value,
1452                ],
1453                format!("inherit Git setting {key}"),
1454            )
1455        })
1456        .collect()
1457}
1458
1459fn controller_git_settings() -> Result<BTreeMap<String, String>> {
1460    let output = match Command::new("git")
1461        .args(["config", "--global", "--includes", "--null", "--list"])
1462        .stdin(Stdio::null())
1463        .output()
1464    {
1465        Ok(output) => output,
1466        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
1467        Err(error) => return Err(error).context("read controller Git configuration"),
1468    };
1469    if !output.status.success() {
1470        bail!(
1471            "read controller Git configuration failed with status {}: {}",
1472            output.status,
1473            String::from_utf8_lossy(&output.stderr).trim()
1474        );
1475    }
1476    parse_inherited_git_settings(&output.stdout)
1477}
1478
1479fn parse_inherited_git_settings(output: &[u8]) -> Result<BTreeMap<String, String>> {
1480    let mut settings = BTreeMap::new();
1481    for entry in output
1482        .split(|byte| *byte == 0)
1483        .filter(|entry| !entry.is_empty())
1484    {
1485        let entry = std::str::from_utf8(entry).context("decode controller Git configuration")?;
1486        let (key, value) = entry
1487            .split_once('\n')
1488            .with_context(|| format!("controller Git returned malformed entry {entry:?}"))?;
1489        let key = key.to_ascii_lowercase();
1490        if INHERITED_GIT_SETTINGS.contains(&key.as_str()) {
1491            settings.insert(key, value.to_owned());
1492        }
1493    }
1494    Ok(settings)
1495}
1496
1497#[cfg(test)]
1498mod tests;