Skip to main content

mj_controller/controller/
provisioning.rs

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