Skip to main content

mj_controller/controller/
provisioning.rs

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