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            // Before the runtime is recorded, so a container never starts
525            // without an entry whose secret is missing.
526            selected.ensure_ready(&session.target_template_id)?;
527            self.config
528                .profiles
529                .get(&session.last_profile)
530                .context("harness profile disappeared before provisioning")?
531                .ensure_ready(&session.last_profile)?;
532            let runtime = mj_core::state::TargetRuntimeSettings::from(selected);
533            if let Some(recorded) = &session.target_runtime {
534                ensure!(
535                    recorded == &runtime,
536                    "target access settings changed before provisioning; retry with the selected target"
537                );
538            } else {
539                self.state
540                    .sessions
541                    .get_mut(session_id)
542                    .unwrap()
543                    .target_runtime = Some(runtime);
544                self.persist_session_state(session_id)?;
545            }
546            let template = self
547                .config
548                .targets
549                .get(&session.target_template_id)
550                .context("target template disappeared during provisioning")?;
551            let profile = self
552                .config
553                .profiles
554                .get(&session.last_profile)
555                .context("harness profile disappeared during provisioning")?;
556            super::worker_binary::preflight_worker_binary(template, executor)?;
557            super::worker_binary::preflight_harness(template, profile, executor)?;
558            self.prepare_managed_raw_worktree(session_id, executor)
559        })();
560        let created_worktree = match preparation {
561            Ok(created) => created,
562            Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
563                return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
564            }
565            Err(error) => return Err(error),
566        };
567        let session = self
568            .state
569            .sessions
570            .get(session_id)
571            .expect("session retained after managed worktree preparation")
572            .clone();
573        // Keep planning, preflight, creation, and locator discovery in one
574        // result so the caller's failure disposition applies to every error.
575        let result = (|| {
576            let template = self
577                .config
578                .targets
579                .get(&session.target_template_id)
580                .context("target template disappeared during provisioning")?;
581            if matches!(template, TargetTemplate::AwsEc2 { .. }) {
582                for resource in &session.additional_mounts {
583                    ensure!(
584                        resource.source.is_dir(),
585                        "attached resource source is not a directory: {}",
586                        resource.source.display()
587                    );
588                }
589            }
590            let mut target = backend_target(
591                template,
592                session.resource_allocation.as_ref(),
593                ContainerOverrides::for_session(&session),
594            )?;
595            let mut runtime_mounts = if matches!(target, targets::TargetTemplate::AwsEc2(_)) {
596                Vec::new()
597            } else {
598                session.additional_mounts.clone()
599            };
600            // The mounts this container runs with, not the ones the session
601            // stores: a forced downgrade belongs to the host the container
602            // lands on, so it is decided here every time and never written
603            // over the user's choice.
604            for notice in enforce_overlay_capable_mounts(&target, &mut runtime_mounts, executor) {
605                executor.notify_notice(&notice);
606            }
607            // The image's user is a property of the host's copy of the image,
608            // so it is read here, once per image per daemon, and handed to the
609            // plan rather than stored on the session.
610            let image_user = podman_image_user(&target, executor);
611            let mut bundle = if session.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) = &session.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                            image_user,
676                            session.container_workspace.as_deref(),
677                            &resource_name,
678                        )
679                    })
680            };
681            let mut provision = match provision {
682                Ok(provision) => provision,
683                Err(error) => {
684                    if let Some(cache) = &prepared_cache {
685                        let _ = cache.cleanup(executor);
686                    }
687                    return Err(error);
688                }
689            };
690            if let Some(token) = container_github_token
691                && let Err(error) =
692                    provision.provide_target_environment_secret(&target, "GH_TOKEN", token)
693            {
694                if let Some(cache) = &prepared_cache {
695                    let _ = cache.cleanup(executor);
696                }
697                return Err(error);
698            }
699
700            let started = Instant::now();
701            let result = provision_target_creation_named(
702                &provision,
703                &target,
704                session_id,
705                executor,
706                &resource_name,
707                |outputs| {
708                    super::backend::locator_after_provision_named(
709                        template,
710                        &target,
711                        session_id,
712                        outputs.first(),
713                        executor,
714                        &resource_name,
715                    )
716                },
717            )
718            .map(|(locator, remainder)| (locator, remainder, bundle, build_cache));
719            if result.is_err()
720                && let Some(cache) = &prepared_cache
721            {
722                if let Some(locator) =
723                    provisioned_locator_named(&target, session_id, None, &resource_name)
724                {
725                    let _ = targets::close_plan(&locator, session_id)
726                        .and_then(|plan| plan.execute(executor).map(|_| ()));
727                } else {
728                    let _ = cache.cleanup(executor);
729                }
730            }
731            tracing::debug!(
732                session_id,
733                elapsed_ms = started.elapsed().as_millis(),
734                "provisioning plan execution completed"
735            );
736            result
737        })();
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                    )
791                } else {
792                    Ok(remainder)
793                }
794            }
795        };
796        let result = match result {
797            Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
798                return Err(self.rollback_failed_new_session(session_id, error)?);
799            }
800            result => result,
801        };
802        if result.is_ok()
803            && let Some(session) = self.state.sessions.get(session_id)
804            && let Some(directory) = session
805                .managed_worktree
806                .as_ref()
807                .map(|worktree| worktree.source_project_directory.clone())
808                .or_else(|| session.project_directory.clone())
809            && let Some(template) = self.config.targets.get(&session.target_template_id)
810        {
811            let host = match template {
812                TargetTemplate::LocalBare => Some("local"),
813                TargetTemplate::SshBare { ssh, .. } => Some(ssh.host.as_str()),
814                _ => None,
815            };
816            if let Some(host) = host {
817                self.state.remember_project_directory(host, &directory);
818                crate::database::remember_project_directory(host, &directory)?;
819            }
820        }
821        self.persist_session_state(session_id)?;
822        result
823
824        }).await
825    }
826
827    pub fn mark_worker_connected(
828        &mut self,
829        session_id: &str,
830        native_session_id: Option<String>,
831    ) -> Result<()> {
832        let session = self
833            .state
834            .sessions
835            .get(session_id)
836            .with_context(|| format!("unknown session {session_id}"))?;
837        if session.target.is_none() {
838            bail!("session {session_id} has no provisioned target");
839        }
840        let updated_at = now();
841        crate::database::mark_session_worker_connected(
842            session_id,
843            native_session_id.as_deref(),
844            &updated_at,
845        )?;
846        let session = self
847            .state
848            .sessions
849            .get_mut(session_id)
850            .expect("session disappeared after its worker connection was saved");
851        session.state = SessionState::Running;
852        if native_session_id.is_some() {
853            session.native_session_id = native_session_id;
854        }
855        session.updated_at = updated_at;
856        session.last_error = None;
857        Ok(())
858    }
859
860    fn install_worker_payload(
861        &self,
862        session_id: &str,
863        executor: &impl CommandExecutor,
864    ) -> Result<(targets::TargetLocator, String)> {
865        // Worker/profile installation is independent of repository cloning.
866        let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
867        let (backend, worker_root) = self.worker_placement(session_id)?;
868        self.prepare_worker_files(session_id, &backend, &worker_root, syncing)?;
869        install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
870        Ok((backend, worker_root))
871    }
872
873    async fn connect_and_start_worker(
874        &self,
875        session_id: &str,
876        executor: &impl CommandExecutor,
877        backend: &targets::TargetLocator,
878        worker_root: &str,
879        initialize_workspace: bool,
880    ) -> Result<Option<String>> {
881        crate::worker_lifecycle::run(session_id, "connect and start worker", executor, async {
882            let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
883            if initialize_workspace {
884                install_inherited_git_settings(executor, backend, session_id)?;
885                self.initialize_network_workspaces(session_id, backend, syncing)?;
886            }
887            let session = self
888                .state
889                .sessions
890                .get(session_id)
891                .with_context(|| format!("unknown session {session_id}"))?;
892            let profile = self
893                .config
894                .profiles
895                .get(&session.last_profile)
896                .with_context(|| format!("unknown profile {}", session.last_profile))?;
897            let readiness_stage = bridge_readiness_stage(profile);
898            let reconnect = &targets::reconnect_plan(backend, session_id)?.commands[0];
899            let readiness = async {
900                let mut relay = {
901                    let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
902                    start_worker_durably(
903                        &crate::worker_lifecycle::require(session_id)?,
904                        self.state.sessions[session_id]
905                            .target
906                            .as_ref()
907                            .context("worker start has no durable target")?,
908                        executor,
909                        backend,
910                        worker_root,
911                    )?;
912                    connect_started_worker(reconnect, session_id, executor, backend, worker_root)
913                        .await?
914                };
915                let native_session_id =
916                    wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
917                let owner = crate::worker_lifecycle::require(session_id)?;
918                crate::database::finish_worker_restart(session_id, owner.operation_id())?;
919                Ok(Some(native_session_id))
920            }
921            .await;
922            match readiness {
923                Ok(native_session_id) => Ok(native_session_id),
924                Err(error) => {
925                    // The diagnosis reads the worker's state, so it runs before the
926                    // worker is stopped.
927                    let error = worker_probe_diagnosis(executor, backend, worker_root, error);
928                    Err(error)
929                }
930            }
931        })
932        .await
933    }
934}
935
936const MAX_LAUNCH_DIAGNOSTIC_BYTES: usize = 64 * 1024;
937
938const RETAINED_LAUNCH_DIAGNOSTICS: usize = 20;
939
940/// Record a failed new-session launch consistently across every failure arm:
941/// log the failure with the session id (so `logs/mj-*.log` names the session)
942/// and save the local diagnostic file. Returns the underlying error chain
943/// annotated with the diagnostic path, which the caller stores in `last_error`
944/// so the reason travels to `mj sessions`, `mj wait`, and `mj events`.
945pub(super) fn note_new_session_launch_failure(session_id: &str, error: &anyhow::Error) -> String {
946    note_new_session_launch_failure_in(&data_dir().join("diagnostics"), session_id, error)
947}
948
949fn note_new_session_launch_failure_in(
950    directory: &Path,
951    session_id: &str,
952    error: &anyhow::Error,
953) -> String {
954    let original = format!("{error:#}");
955    tracing::warn!(session_id, error = %original, "session launch failed");
956    match persist_launch_failure_to(directory, session_id, &original) {
957        Ok(path) => format!("{original}; full diagnostic saved to {}", path.display()),
958        Err(save_error) => {
959            format!("{original}; saving the local diagnostic failed: {save_error:#}")
960        }
961    }
962}
963
964fn persist_launch_failure_to(directory: &Path, session_id: &str, detail: &str) -> Result<PathBuf> {
965    mj_core::config::validate_id("session", session_id)?;
966    std::fs::create_dir_all(directory).with_context(|| {
967        format!(
968            "create launch diagnostics directory {}",
969            directory.display()
970        )
971    })?;
972    #[cfg(unix)]
973    {
974        use std::os::unix::fs::PermissionsExt;
975        std::fs::set_permissions(directory, std::fs::Permissions::from_mode(0o700))?;
976    }
977    let path = directory.join(format!("{session_id}-launch-error.txt"));
978    let detail = bounded_launch_diagnostic(detail);
979    let body = format!(
980        "Mjolnir session launch failure\nsession: {session_id}\nat: {}\n\n{detail}\n",
981        now()
982    );
983    atomic_write(&path, body.as_bytes())?;
984    prune_launch_diagnostics(directory)?;
985    Ok(path)
986}
987
988fn bounded_launch_diagnostic(detail: &str) -> String {
989    if detail.len() <= MAX_LAUNCH_DIAGNOSTIC_BYTES {
990        return detail.to_owned();
991    }
992    let mut head_end = MAX_LAUNCH_DIAGNOSTIC_BYTES / 4;
993    while !detail.is_char_boundary(head_end) {
994        head_end -= 1;
995    }
996    let tail_bytes = MAX_LAUNCH_DIAGNOSTIC_BYTES - head_end;
997    let mut tail_start = detail.len() - tail_bytes;
998    while !detail.is_char_boundary(tail_start) {
999        tail_start += 1;
1000    }
1001    format!(
1002        "{}\n\n[... launch diagnostic truncated ...]\n\n{}",
1003        &detail[..head_end],
1004        &detail[tail_start..]
1005    )
1006}
1007
1008fn prune_launch_diagnostics(directory: &Path) -> Result<()> {
1009    let mut diagnostics = Vec::new();
1010    for entry in std::fs::read_dir(directory)? {
1011        let entry = entry?;
1012        if !entry
1013            .file_name()
1014            .to_str()
1015            .is_some_and(|name| name.ends_with("-launch-error.txt"))
1016        {
1017            continue;
1018        }
1019        diagnostics.push((entry.metadata()?.modified()?, entry.path()));
1020    }
1021    diagnostics.sort_by_key(|entry| std::cmp::Reverse(entry.0));
1022    for (_, path) in diagnostics.into_iter().skip(RETAINED_LAUNCH_DIAGNOSTICS) {
1023        std::fs::remove_file(&path)
1024            .with_context(|| format!("prune old launch diagnostic {}", path.display()))?;
1025    }
1026    Ok(())
1027}
1028
1029fn apply_new_session_provisioning_result(
1030    state: &mut State,
1031    session_id: &str,
1032    result: Result<TargetLocator>,
1033) -> Result<()> {
1034    match result {
1035        Ok(locator) => {
1036            let record = state.sessions.get_mut(session_id).unwrap();
1037            record.target = Some(locator);
1038            // Provisioning has completed, but Running is reserved for a
1039            // successful worker handshake.
1040            record.state = SessionState::Disconnected;
1041            record.updated_at = now();
1042            record.last_error = None;
1043            Ok(())
1044        }
1045        Err(error) => {
1046            let record = state.sessions.get_mut(session_id).unwrap();
1047            record.state = SessionState::Error;
1048            record.target = None;
1049            record.updated_at = now();
1050            record.last_error = Some(format!("session provisioning failed: {error:#}"));
1051            Err(error)
1052        }
1053    }
1054}
1055
1056/// Mark a new session's launch failed while it still names its target, so the
1057/// record is true while the rollback removes that target.
1058fn apply_failed_new_session_launch(state: &mut State, session_id: &str, original_error: &str) {
1059    let record = state.sessions.get_mut(session_id).unwrap();
1060    record.state = SessionState::StartupCleanup;
1061    record.updated_at = now();
1062    record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1063}
1064
1065pub(super) fn apply_failed_new_session_rollback(
1066    state: &mut State,
1067    session_id: &str,
1068    original_error: &str,
1069    cleanup_error: Option<String>,
1070) -> anyhow::Error {
1071    match cleanup_error {
1072        None => {
1073            let record = state.sessions.get_mut(session_id).unwrap();
1074            record.state = SessionState::Error;
1075            record.target = None;
1076            record.updated_at = now();
1077            record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1078            anyhow::anyhow!("{original_error}; partial target removed and failed session retained")
1079        }
1080        Some(cleanup_error) => {
1081            let failure = format!(
1082                "{original_error}; cleanup of the failed session target failed: {cleanup_error}"
1083            );
1084            let record = state.sessions.get_mut(session_id).unwrap();
1085            record.state = SessionState::StartupCleanup;
1086            record.updated_at = now();
1087            record.last_error = Some(format!("worker bootstrap failed: {failure}"));
1088            anyhow::anyhow!(failure)
1089        }
1090    }
1091}
1092
1093pub(super) fn install_attached_resources(
1094    state: &State,
1095    session_id: &str,
1096    backend: &targets::TargetLocator,
1097    worker_root: &str,
1098    executor: &impl CommandExecutor,
1099) -> Result<()> {
1100    let targets::TargetLocator::AwsEc2 { .. } = backend else {
1101        return Ok(());
1102    };
1103    let session = state
1104        .sessions
1105        .get(session_id)
1106        .with_context(|| format!("unknown session {session_id}"))?;
1107    if session.additional_mounts.is_empty() {
1108        return Ok(());
1109    }
1110    for resource in &session.additional_mounts {
1111        let install = targets::command_on_locator(
1112            backend,
1113            session_id,
1114            vec![
1115                format!("{worker_root}/hel"),
1116                "worker".into(),
1117                "install-resource".into(),
1118                "--destination".into(),
1119                resource.destination.to_string_lossy().into_owned(),
1120            ],
1121            "stream attached resource",
1122        )?;
1123        mj_checkpoint::resources::stream_resource(&resource.source, |stream| {
1124            execute_checked_with_stdin(executor, &install, stream).map(|_| ())
1125        })
1126        .with_context(|| format!("stream attached resource {}", resource.source.display()))?;
1127    }
1128    Ok(())
1129}
1130
1131/// Run two independent target setup lanes at the same time and wait for both.
1132/// The first lane's failure wins deterministically when both fail, and neither
1133/// lane is abandoned while it may still own a transfer or subprocess.
1134pub(super) fn execute_concurrent_lanes<A: Send, B: Send>(
1135    first: impl FnOnce() -> Result<A> + Send,
1136    second: impl FnOnce() -> Result<B> + Send,
1137) -> Result<(A, B)> {
1138    std::thread::scope(|scope| {
1139        let owner = crate::worker_lifecycle::capture();
1140        let second = scope.spawn(move || match owner {
1141            Some(owner) => owner.scope_blocking(second),
1142            None => second(),
1143        });
1144        let first = first();
1145        let second = second.join().unwrap_or_else(|panic| {
1146            Err(anyhow::anyhow!(
1147                "concurrent target lane panicked: {}",
1148                targets::command_thread_panic_message(panic.as_ref())
1149            ))
1150        });
1151        match (first, second) {
1152            (Err(error), _) => Err(error),
1153            (Ok(_), Err(error)) => Err(error),
1154            (Ok(first), Ok(second)) => Ok((first, second)),
1155        }
1156    })
1157}
1158
1159fn execute_repository_setup(
1160    plan: &targets::CommandPlan,
1161    executor: &(impl CommandExecutor + Sync),
1162) -> Result<()> {
1163    if plan.commands.is_empty() {
1164        return Ok(());
1165    }
1166    let _cloning = ProvisionStageGuard::new(executor, ProvisionStage::Cloning);
1167    plan.execute_concurrent(executor).map(|_| ())
1168}
1169
1170/// Run a provisioning plan and discover the locator it produced, tearing the
1171/// target down again if anything after its creation fails.
1172///
1173/// Creation is the boundary that matters. A step that fails before the target
1174/// exists has left nothing behind; every failure after it — a later plan step
1175/// or locator discovery — owns a target no session record will point at.
1176#[cfg(test)]
1177fn provision_target(
1178    plan: &targets::CommandPlan,
1179    target: &targets::TargetTemplate,
1180    session_id: &str,
1181    executor: &(impl CommandExecutor + Sync),
1182    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1183) -> Result<TargetLocator> {
1184    let Some((creation, remainder)) = plan.split_at_target_creation() else {
1185        // Nothing this plan runs can leave a target behind.
1186        return discover(&plan.execute_concurrent(executor)?);
1187    };
1188    let mut outputs = creation.execute_concurrent(executor)?;
1189    let result = match remainder.execute_concurrent(executor) {
1190        Ok(rest) => {
1191            outputs.extend(rest);
1192            discover(&outputs)
1193        }
1194        Err(error) => Err(error),
1195    };
1196    result.map_err(|error| {
1197        match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1198            Some(note) => error.context(note),
1199            None => error,
1200        }
1201    })
1202}
1203
1204/// Bring the target into existence and return the commands that populate its
1205/// repositories. The caller may overlap that remainder with worker/profile
1206/// installation once it has persisted the discovered locator.
1207#[cfg(test)]
1208fn provision_target_creation(
1209    plan: &targets::CommandPlan,
1210    target: &targets::TargetTemplate,
1211    session_id: &str,
1212    executor: &(impl CommandExecutor + Sync),
1213    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1214) -> Result<(TargetLocator, targets::CommandPlan)> {
1215    provision_target_creation_named(
1216        plan,
1217        target,
1218        session_id,
1219        executor,
1220        &targets::resource_name(session_id)?,
1221        discover,
1222    )
1223}
1224
1225fn provision_target_creation_named(
1226    plan: &targets::CommandPlan,
1227    target: &targets::TargetTemplate,
1228    session_id: &str,
1229    executor: &(impl CommandExecutor + Sync),
1230    name: &str,
1231    discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1232) -> Result<(TargetLocator, targets::CommandPlan)> {
1233    let Some((creation, remainder)) = plan.split_at_target_creation() else {
1234        // Nothing this plan runs can leave a target behind, so its commands
1235        // must still finish before the locator is usable.
1236        let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1237            plan.execute_concurrent(executor)
1238        })?;
1239        return discover(&outputs).map(|locator| {
1240            (
1241                locator,
1242                targets::CommandPlan {
1243                    description: plan.description.clone(),
1244                    commands: Vec::new(),
1245                },
1246            )
1247        });
1248    };
1249    // Creating the container is what downloads the image on Docker, and the
1250    // probe just before it is what downloads it on Podman. Either way, a
1251    // background download of the same image must finish first.
1252    let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1253        creation.execute_concurrent(executor)
1254    })?;
1255    discover(&outputs)
1256        .map(|locator| (locator, remainder))
1257        .map_err(|error| {
1258            match cleanup_failed_provision_named(
1259                target,
1260                session_id,
1261                outputs.first(),
1262                executor,
1263                name,
1264            ) {
1265                Some(note) => error.context(note),
1266                None => error,
1267            }
1268        })
1269}
1270
1271/// Best-effort teardown of a target whose creation succeeded but whose
1272/// provisioning failed before a locator was recorded. Returns a note
1273/// describing what happened for inclusion in the session error.
1274///
1275/// The teardown is the session's own close plan, so a failed launch and an
1276/// ordinary close can never disagree about what removing a target means.
1277#[cfg(test)]
1278fn cleanup_failed_provision(
1279    target: &targets::TargetTemplate,
1280    session_id: &str,
1281    create_output: Option<&CommandOutput>,
1282    executor: &impl CommandExecutor,
1283) -> Option<String> {
1284    cleanup_failed_provision_named(
1285        target,
1286        session_id,
1287        create_output,
1288        executor,
1289        &targets::resource_name(session_id).ok()?,
1290    )
1291}
1292
1293fn cleanup_failed_provision_named(
1294    target: &targets::TargetTemplate,
1295    session_id: &str,
1296    create_output: Option<&CommandOutput>,
1297    executor: &impl CommandExecutor,
1298    name: &str,
1299) -> Option<String> {
1300    let locator = provisioned_locator_named(target, session_id, create_output, name)?;
1301
1302    let leak = format!(
1303        "the resource may still exist; find it via its dev.mj.session={session_id} label/tag"
1304    );
1305    let cleanup = if name == targets::resource_name(session_id).ok()? {
1306        targets::close_plan(&locator, session_id)
1307    } else {
1308        targets::retire_move_target_plan(&locator, session_id)
1309    };
1310    let plan = match cleanup {
1311        Ok(plan) => plan,
1312        Err(error) => {
1313            tracing::warn!(
1314                session_id,
1315                error = format!("{error:#}"),
1316                "could not build provisioning cleanup plan"
1317            );
1318            return Some(format!("cleanup FAILED: {error:#}; {leak}"));
1319        }
1320    };
1321    let purpose = plan
1322        .commands
1323        .iter()
1324        .map(|command| command.purpose.clone())
1325        .collect::<Vec<_>>()
1326        .join("; ");
1327    let Err(error) = plan.execute(executor) else {
1328        return Some(format!("cleanup succeeded: {purpose}"));
1329    };
1330    match targets::cleanup_target_is_confirmed_absent(&locator, session_id, executor) {
1331        Ok(true) => Some(format!("cleanup succeeded: {purpose}")),
1332        Ok(false) => {
1333            tracing::warn!(
1334                session_id,
1335                error = format!("{error:#}"),
1336                "provisioning cleanup failed and the target may still exist"
1337            );
1338            Some(format!("cleanup FAILED ({purpose}): {error:#}; {leak}"))
1339        }
1340        Err(confirm_error) => {
1341            tracing::warn!(
1342                session_id,
1343                error = format!("{confirm_error:#}"),
1344                "could not confirm whether the failed provisioning target was removed"
1345            );
1346            Some(format!(
1347                "cleanup FAILED ({purpose}): {error:#}; checking whether it was removed also failed: {confirm_error:#}; {leak}"
1348            ))
1349        }
1350    }
1351}
1352
1353/// The locator a provisioning plan's creating command brought into existence.
1354///
1355/// Every target but AWS is named before its plan runs; an EC2 instance
1356/// reports its own ID in the launch response.
1357fn provisioned_locator_named(
1358    target: &targets::TargetTemplate,
1359    session_id: &str,
1360    create_output: Option<&CommandOutput>,
1361    name: &str,
1362) -> Option<targets::TargetLocator> {
1363    let container_id = || Some(name.to_owned());
1364
1365    Some(match target {
1366        // A bare project directory belongs to the user: provisioning creates
1367        // nothing that a failure could leak.
1368        targets::TargetTemplate::LocalBare => return None,
1369        targets::TargetTemplate::LocalPodman(container) => targets::TargetLocator::LocalPodman {
1370            borrowed_from: None,
1371            container_id: container_id()?,
1372            workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1373        },
1374        targets::TargetTemplate::LocalDocker(_) => targets::TargetLocator::LocalDocker {
1375            borrowed_from: None,
1376            container_id: container_id()?,
1377        },
1378        targets::TargetTemplate::AppleContainer(_) => targets::TargetLocator::AppleContainer {
1379            borrowed_from: None,
1380            container_id: container_id()?,
1381        },
1382        targets::TargetTemplate::SshPodman { ssh, container } => {
1383            targets::TargetLocator::SshPodman {
1384                borrowed_from: None,
1385                ssh: ssh.clone(),
1386                container_id: container_id()?,
1387                workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1388            }
1389        }
1390        targets::TargetTemplate::SshDocker { ssh, .. } => targets::TargetLocator::SshDocker {
1391            borrowed_from: None,
1392            ssh: ssh.clone(),
1393            container_id: container_id()?,
1394        },
1395        targets::TargetTemplate::SshBare { ssh, .. } => targets::TargetLocator::SshBare {
1396            ssh: ssh.clone(),
1397            workspace: targets::workspace_for(target, session_id).ok()?,
1398            worker_id: None,
1399        },
1400        targets::TargetTemplate::AwsEc2(aws) => targets::TargetLocator::AwsEc2 {
1401            profile: aws.profile.clone(),
1402            region: aws.region.clone(),
1403            instance_id: serde_json::from_slice::<serde_json::Value>(&create_output?.stdout)
1404                .ok()?
1405                .pointer("/Instances/0/InstanceId")?
1406                .as_str()?
1407                .to_owned(),
1408            ssh: aws.ssh.clone(),
1409            workspace: targets::workspace_for(target, session_id).ok()?,
1410        },
1411    })
1412}
1413
1414/// Attach read-only whatever the selected container overlay cannot hold, and
1415/// say so.
1416///
1417/// The filesystem is probed on the host that runs the container, because that
1418/// is where the overlay would be built. A probe that cannot answer leaves the
1419/// overlay alone: a failed probe is no evidence of an unsupported filesystem,
1420/// and refusing to provision over one would cost the user their session.
1421///
1422/// Apple's `container` engine already mounts every extra directory read-only,
1423/// and EC2 copies the directory instead of mounting it, so neither is probed.
1424pub(super) fn enforce_overlay_capable_mounts(
1425    target: &targets::TargetTemplate,
1426    mounts: &mut [targets::AdditionalMount],
1427    executor: &impl CommandExecutor,
1428) -> Vec<String> {
1429    let ssh = match target {
1430        targets::TargetTemplate::LocalPodman(_) | targets::TargetTemplate::LocalDocker(_) => None,
1431        targets::TargetTemplate::SshPodman { ssh, .. }
1432        | targets::TargetTemplate::SshDocker { ssh, .. } => Some(ssh),
1433        _ => return Vec::new(),
1434    };
1435    let overlaid = mounts
1436        .iter()
1437        .filter(|mount| mount.access == targets::MountAccess::Cow)
1438        .map(|mount| mount.source.clone())
1439        .collect::<Vec<_>>();
1440    if overlaid.is_empty() {
1441        return Vec::new();
1442    }
1443    // `stat` on this host cannot see what a VM-hosted Docker daemon mounts.
1444    if matches!(target, targets::TargetTemplate::LocalDocker(_)) {
1445        match targets::local_docker_vm_share(executor) {
1446            Ok(Some(reason)) => {
1447                return mounts
1448                    .iter_mut()
1449                    .filter(|mount| mount.access == targets::MountAccess::Cow)
1450                    .map(|mount| {
1451                        mount.access = mount.access.without_overlay();
1452                        format!(
1453                            "Mounted {} read-only: {reason}, which cannot back the \
1454                             copy-on-write overlay.",
1455                            mount.source.display()
1456                        )
1457                    })
1458                    .collect();
1459            }
1460            Ok(None) => {}
1461            Err(error) => tracing::warn!(
1462                error = format!("{error:#}"),
1463                "could not identify the Docker daemon platform; probing this host's filesystems"
1464            ),
1465        }
1466    }
1467    let filesystems = match targets::probe_filesystem_types(ssh, &overlaid, executor) {
1468        Ok(filesystems) => filesystems,
1469        Err(error) => {
1470            tracing::warn!(
1471                error = format!("{error:#}"),
1472                "could not probe attached-directory filesystems; preserving overlay mounts"
1473            );
1474            return vec![format!(
1475                "Could not read the filesystem under the attached directories, so they keep the \
1476                 copy-on-write overlay: {error:#}"
1477            )];
1478        }
1479    };
1480    let mut notices = Vec::new();
1481    for (mount, filesystem) in mounts
1482        .iter_mut()
1483        .filter(|mount| mount.access == targets::MountAccess::Cow)
1484        .zip(filesystems)
1485    {
1486        let Some(reason) = targets::overlay_unsupported_filesystem(&filesystem) else {
1487            continue;
1488        };
1489        mount.access = mount.access.without_overlay();
1490        notices.push(format!(
1491            "Mounted {} read-only: the overlay is unreliable on {filesystem} ({reason}).",
1492            mount.source.display()
1493        ));
1494    }
1495    notices
1496}
1497
1498/// The image users already probed, keyed by container host and image
1499/// reference. An image's
1500/// configured user does not change under a fixed reference, and reading it
1501/// costs a container start, so each daemon asks a host once.
1502static IMAGE_USERS: std::sync::LazyLock<std::sync::Mutex<BTreeMap<String, targets::ImageUser>>> =
1503    std::sync::LazyLock::new(std::sync::Mutex::default);
1504
1505/// The uid and gid a Podman session container maps onto the host user.
1506///
1507/// Only Podman is asked: Docker and Apple's `container` engine are left with
1508/// their own defaults. A probe that cannot answer is not a launch failure —
1509/// the container falls back to plain `--userns=keep-id`, which maps the
1510/// image's default user, and the user is told what happened.
1511pub(super) fn podman_image_user(
1512    target: &targets::TargetTemplate,
1513    executor: &impl CommandExecutor,
1514) -> Option<targets::ImageUser> {
1515    let (ssh, container) = match target {
1516        targets::TargetTemplate::LocalPodman(container) => (None, container),
1517        targets::TargetTemplate::SshPodman { ssh, container } => (Some(ssh), container),
1518        _ => return None,
1519    };
1520    let image = container.image.as_str();
1521    let key = format!(
1522        "{}|{image}",
1523        ssh.map_or("local", |ssh| ssh.destination.as_str())
1524    );
1525    if let Some(cached) = IMAGE_USERS.lock().expect("image user cache").get(&key) {
1526        return Some(*cached);
1527    }
1528    // The probe starts a container, so on Podman this is where a missing image
1529    // is actually downloaded. Wait for the daemon's own download instead of
1530    // starting a second one.
1531    match crate::image_pull_gate::with_image_ready(target, executor, || {
1532        targets::probe_image_user(ssh, container, executor)
1533    }) {
1534        Ok(user) => {
1535            IMAGE_USERS
1536                .lock()
1537                .expect("image user cache")
1538                .insert(key, user);
1539            Some(user)
1540        }
1541        Err(error) => {
1542            tracing::warn!(
1543                image,
1544                error = format!("{error:#}"),
1545                "could not read the container image user; keeping Podman's default user mapping"
1546            );
1547            executor.notify_notice(&format!(
1548                "Could not read the user of image {image}, so the container runs with Podman's \
1549                 default user mapping and may not be able to write to an attached directory: \
1550                 {error:#}"
1551            ));
1552            None
1553        }
1554    }
1555}
1556
1557/// Reports every command an installer issues as one launch stage, so progress
1558/// stays accurate without threading the stage through each `CommandSpec`.
1559/// A command that already names a stage keeps it.
1560pub(super) struct StagedExecutor<'a, E: CommandExecutor> {
1561    inner: &'a E,
1562    stage: ProvisionStage,
1563    _guard: ProvisionStageGuard<'a, E>,
1564}
1565
1566impl<'a, E: CommandExecutor> StagedExecutor<'a, E> {
1567    pub(crate) fn new(inner: &'a E, stage: ProvisionStage) -> Self {
1568        Self {
1569            inner,
1570            stage,
1571            _guard: ProvisionStageGuard::new(inner, stage),
1572        }
1573    }
1574
1575    fn staged(&self, command: &CommandSpec) -> CommandSpec {
1576        if command.stage.is_some() {
1577            return command.clone();
1578        }
1579        command.clone().stage(self.stage)
1580    }
1581}
1582
1583impl<E: CommandExecutor> CommandExecutor for StagedExecutor<'_, E> {
1584    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1585        self.inner.execute(&self.staged(command))
1586    }
1587
1588    fn cancellation_requested(&self) -> bool {
1589        self.inner.cancellation_requested()
1590    }
1591
1592    fn stage_started(&self, stage: ProvisionStage) {
1593        self.inner.stage_started(stage);
1594    }
1595
1596    fn stage_finished(&self, stage: ProvisionStage) {
1597        self.inner.stage_finished(stage);
1598    }
1599
1600    fn notify_notice(&self, notice: &str) {
1601        self.inner.notify_notice(notice);
1602    }
1603
1604    fn execute_with_stdin(
1605        &self,
1606        command: &CommandSpec,
1607        input: &mut (dyn std::io::Read + Send),
1608    ) -> Result<CommandOutput> {
1609        self.inner.execute_with_stdin(&self.staged(command), input)
1610    }
1611}
1612
1613fn execute_checked_with_stdin(
1614    executor: &impl CommandExecutor,
1615    command: &CommandSpec,
1616    input: &mut (dyn std::io::Read + Send),
1617) -> Result<CommandOutput> {
1618    let output = executor.execute_with_stdin(command, input)?;
1619    if output.status != 0 {
1620        bail!(
1621            "{} failed with status {}: {}",
1622            command.purpose,
1623            output.status,
1624            String::from_utf8_lossy(&output.stderr)
1625        );
1626    }
1627    Ok(output)
1628}
1629
1630pub(super) fn install_inherited_git_settings(
1631    executor: &impl CommandExecutor,
1632    locator: &targets::TargetLocator,
1633    session_id: &str,
1634) -> Result<()> {
1635    let settings = if inherits_controller_git_settings(locator) {
1636        controller_git_settings()?
1637    } else {
1638        BTreeMap::new()
1639    };
1640    for command in inherited_git_setting_commands(locator, session_id, settings)? {
1641        execute_checked(executor, command)?;
1642    }
1643    Ok(())
1644}
1645
1646fn inherits_controller_git_settings(locator: &targets::TargetLocator) -> bool {
1647    !matches!(
1648        locator,
1649        targets::TargetLocator::LocalBare { .. } | targets::TargetLocator::SshBare { .. }
1650    )
1651}
1652
1653fn inherited_git_setting_commands(
1654    locator: &targets::TargetLocator,
1655    session_id: &str,
1656    settings: BTreeMap<String, String>,
1657) -> Result<Vec<CommandSpec>> {
1658    if matches!(locator, targets::TargetLocator::SshBare { .. }) {
1659        return Ok(Vec::new());
1660    }
1661    settings
1662        .into_iter()
1663        .map(|(key, value)| {
1664            targets::command_on_locator(
1665                locator,
1666                session_id,
1667                vec![
1668                    "git".into(),
1669                    "config".into(),
1670                    "--global".into(),
1671                    "--replace-all".into(),
1672                    "--".into(),
1673                    key.clone(),
1674                    value,
1675                ],
1676                format!("inherit Git setting {key}"),
1677            )
1678        })
1679        .collect()
1680}
1681
1682fn controller_git_settings() -> Result<BTreeMap<String, String>> {
1683    let output = match Command::new("git")
1684        .args(["config", "--global", "--includes", "--null", "--list"])
1685        .stdin(Stdio::null())
1686        .output()
1687    {
1688        Ok(output) => output,
1689        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
1690        Err(error) => return Err(error).context("read controller Git configuration"),
1691    };
1692    if !output.status.success() {
1693        bail!(
1694            "read controller Git configuration failed with status {}: {}",
1695            output.status,
1696            String::from_utf8_lossy(&output.stderr).trim()
1697        );
1698    }
1699    parse_inherited_git_settings(&output.stdout)
1700}
1701
1702fn parse_inherited_git_settings(output: &[u8]) -> Result<BTreeMap<String, String>> {
1703    let mut settings = BTreeMap::new();
1704    for entry in output
1705        .split(|byte| *byte == 0)
1706        .filter(|entry| !entry.is_empty())
1707    {
1708        let entry = std::str::from_utf8(entry).context("decode controller Git configuration")?;
1709        let (key, value) = entry
1710            .split_once('\n')
1711            .with_context(|| format!("controller Git returned malformed entry {entry:?}"))?;
1712        let key = key.to_ascii_lowercase();
1713        if INHERITED_GIT_SETTINGS.contains(&key.as_str()) {
1714            settings.insert(key, value.to_owned());
1715        }
1716    }
1717    Ok(settings)
1718}
1719
1720#[cfg(test)]
1721mod tests;