Skip to main content

mj_controller/
controller.rs

1//! Controller-side lifecycle transitions and canonical-to-backend conversion.
2
3mod backend;
4mod cache_host;
5pub(crate) mod checkpoint;
6mod git_cache;
7mod lifecycle;
8pub mod local_profile_homes;
9pub(crate) mod mbx;
10pub mod move_session;
11mod network_git;
12mod new_session_preflight;
13pub use new_session_preflight::{NewSessionPreflight, NewSessionRepository};
14mod path_completion;
15pub mod profile_config;
16mod provisioning;
17pub(crate) mod publication;
18mod readiness;
19pub(crate) use readiness::NATIVE_SESSION_STARTUP_TIMEOUT;
20mod recovery_scan;
21mod resume;
22mod reviewer;
23mod subagent_park;
24pub use subagent_park::ParkOutcome;
25pub(crate) use subagent_park::{FAILED_STARTUP_CLEANUP_TIMEOUT, failed_launch_cleanup_executor};
26mod subagents;
27#[cfg(test)]
28pub(crate) mod test_support;
29pub mod update;
30mod worker_binary;
31mod worker_restart;
32mod worktree;
33pub(crate) use worktree::RemoteGitExecutor;
34
35use std::collections::{BTreeMap, BTreeSet};
36use std::fs::{self, File, OpenOptions};
37use std::path::{Path, PathBuf};
38
39use anyhow::{Context, Result, bail, ensure};
40use chrono::Utc;
41
42use mj_core::config::{
43    Config, ProjectBundle, ProjectRepository, TargetTemplate, atomic_write, container_size_host,
44    data_dir, is_bare_project_target, mount_history_host,
45};
46
47use crate::import::{RepositoryIdentity, bundle_matches, setup_style_id};
48use crate::setup::github_repository_from_origin;
49
50const CONFIG_RENAME_JOURNAL: &str = "config-rename.json";
51
52#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
53#[serde(rename_all = "snake_case")]
54enum ConfigRenameKind {
55    Profile,
56    Target,
57}
58
59#[derive(Debug, serde::Serialize, serde::Deserialize)]
60#[serde(deny_unknown_fields)]
61struct ConfigRenameJournal {
62    kind: ConfigRenameKind,
63    old_id: String,
64    new_id: String,
65}
66use mj_core::state::{
67    HostContainerSize, SessionRecord, SessionResourceAllocation, SessionState, State,
68    new_session_id, normalize_session_title,
69};
70
71use crate::targets::{
72    self, AdditionalMount, CommandExecutor, CommandOutput, CommandSpec, SshTarget,
73};
74
75pub(crate) use backend::controller_github_token;
76pub use backend::image_refresh_plan;
77use backend::validate_resource_allocation;
78pub(crate) use backend::{LocalEngineReadiness, local_engine_readiness};
79pub use mbx::preview_build_cache;
80pub(crate) use mbx::{DoctorHostMbxStatus, MBX_VERSION, doctor_host_mbx};
81use provisioning::apply_failed_new_session_rollback;
82pub(crate) use worker_binary::{
83    prepare_recovery_worker_binary, refresh_target_worker_binary_if_stale, worker_source_problem,
84};
85pub(crate) use worktree::path_exists_on_managed_target;
86
87pub use checkpoint::{
88    CheckpointArtifact, CheckpointDeferred, IdleWorkspaceLease, SessionExportLayout,
89    checkpoint_was_deferred, reconcile_managed_checkpoint_archives,
90    sweep_local_checkpoint_leftovers,
91};
92pub use lifecycle::{
93    BeforeClose, BranchDisposition, CheckoutDisposition, has_nothing_to_checkpoint,
94};
95pub use recovery_scan::{RecoveryCandidate, RecoveryScan};
96pub use resume::{
97    ResumeRepositorySourceMismatch, ResumeRepositorySourcePreflight, ResumeRepositorySourceReceipt,
98    raw_conversion_preview_for,
99};
100pub use reviewer::reviewer_stager;
101pub use subagents::{RegisterSubagentRequest, stopped_subagent, subagent_has_handed_back};
102pub use worker_binary::{
103    WorkerBinaryAvailability, native_worker_binary_prerequisite, pin_worker_binary_sources,
104    ssh_worker_binary_prerequisite, worker_binary_prerequisite_for_arch,
105};
106pub use worker_restart::WorkerUpgradeOutcome;
107pub use worktree::{ResumePlan, local_project_repository, resume_compatibility};
108
109pub struct Controller {
110    pub config: Config,
111    pub state: State,
112}
113
114/// Machine-wide advisory lock for one controller data store. This prevents a
115/// dashboard, server, or CLI lifecycle command from concurrently acting as a
116/// second controller against the same SQLite state and relay sessions.
117#[derive(Debug)]
118pub struct ControllerStoreGuard {
119    file: File,
120}
121
122impl ControllerStoreGuard {
123    pub fn acquire() -> Result<Self> {
124        let directory = data_dir();
125        Self::acquire_at(&directory)
126    }
127
128    fn acquire_at(directory: &Path) -> Result<Self> {
129        Self::try_acquire_at(directory)?.with_context(|| {
130            format!(
131                "another Mjolnir controller is already using {}; stop it before starting this command",
132                directory.display()
133            )
134        })
135    }
136
137    /// Probe exclusivity without treating an owner that is still exiting as an error.
138    pub fn try_acquire() -> Result<Option<Self>> {
139        Self::try_acquire_at(&data_dir())
140    }
141
142    fn try_acquire_at(directory: &Path) -> Result<Option<Self>> {
143        // The lock holder starts the database writer, which migrates the store.
144        mj_core::config::ensure_may_control_store(directory, "own this data store")?;
145        std::fs::create_dir_all(directory)
146            .with_context(|| format!("create controller data directory {}", directory.display()))?;
147        let path = directory.join("controller.lock");
148        let mut options = OpenOptions::new();
149        options.create(true).read(true).write(true);
150        #[cfg(unix)]
151        {
152            use std::os::unix::fs::OpenOptionsExt;
153            options.mode(0o600);
154        }
155        let file = options
156            .open(&path)
157            .with_context(|| format!("open controller lock {}", path.display()))?;
158        match file.try_lock() {
159            Ok(()) => {}
160            Err(std::fs::TryLockError::WouldBlock) => return Ok(None),
161            Err(std::fs::TryLockError::Error(error)) => {
162                return Err(error)
163                    .with_context(|| format!("lock controller store {}", directory.display()));
164            }
165        }
166        Ok(Some(Self { file }))
167    }
168
169    /// Start the sole production SQLite writer after controller exclusivity
170    /// has been established by this guard.
171    pub fn start_database_writer(&self) -> Result<crate::database::DatabaseWriterOwner> {
172        crate::database::start_database_writer()
173    }
174}
175
176impl Drop for ControllerStoreGuard {
177    fn drop(&mut self) {
178        // Make release explicit. `File` also unlocks on close, but an explicit
179        // unlock keeps same-process handoff deterministic across platforms.
180        let _ = self.file.unlock();
181    }
182}
183
184/// The durable result of creating a quick bundle. The returned config is the
185/// same fresh config that was written, allowing a serving projection to publish
186/// the new bundle before acknowledging the request that created it.
187#[derive(Debug)]
188pub struct QuickBundleCreation {
189    pub config: Config,
190    pub bundle_id: String,
191}
192
193/// Failure stages exposed to a viewer request without exposing the underlying
194/// filesystem/configuration error. The detailed error remains available to
195/// the caller for logs and terminal notices.
196#[derive(Debug)]
197pub enum QuickBundleFailure {
198    InvalidSource(anyhow::Error),
199    Persistence(anyhow::Error),
200}
201
202impl std::fmt::Display for QuickBundleFailure {
203    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
204        match self {
205            Self::InvalidSource(error) => write!(formatter, "invalid repository source: {error}"),
206            Self::Persistence(error) => write!(formatter, "persist quick bundle: {error}"),
207        }
208    }
209}
210
211impl std::error::Error for QuickBundleFailure {
212    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
213        match self {
214            Self::InvalidSource(error) | Self::Persistence(error) => Some(error.root_cause()),
215        }
216    }
217}
218
219/// Create a quick bundle from a local repository or GitHub source and persist
220/// it as one serialized fresh-config transaction. Identical sources reuse the
221/// existing configured bundle, matching the terminal's behavior. The returned
222/// config is the same fresh config that was written, allowing a serving
223/// projection to publish the new bundle before acknowledging the request.
224pub fn create_quick_bundle(
225    source: &str,
226) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
227    create_bundle_from_sources(&[source.to_owned()])
228}
229
230/// Add a quick bundle to an already-loaded config. The helper still performs
231/// the local repository canonicalization/GitHub-source parsing, but callers
232/// that persist a config should use [`create_quick_bundle`] so concurrent saves
233/// cannot clobber one another.
234pub fn create_quick_bundle_in_config(config: &mut Config, source: &str) -> Result<String> {
235    create_bundle_from_sources_in_config(config, &[source.to_owned()])
236}
237
238/// Create one bundle from one or more local repositories or GitHub sources.
239///
240/// All sources are interpreted and checked before the config transaction can
241/// write anything. An existing bundle is reused only when its repository set
242/// and primary repository exactly match the request; this keeps selecting one
243/// repository from a larger bundle from silently changing the wizard's choice.
244pub fn create_bundle_from_sources(
245    sources: &[String],
246) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
247    let (config, bundle_id) = Config::update(|config| {
248        let previous = config.bundles.clone();
249        let id = create_bundle_from_sources_in_config(config, sources)
250            .map_err(|error| anyhow::Error::new(QuickBundleFailure::InvalidSource(error)))?;
251        if let Some(old) = previous.get(&id)
252            && config.bundles.get(&id) != Some(old)
253        {
254            let accepted = crate::project_catalog::snapshot(
255                old,
256                &targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(
257                    15,
258                )),
259                true,
260            )?;
261            crate::database::store_catalog_project(&id, &accepted, true)?;
262        }
263        let project = crate::project_catalog::snapshot(
264            config
265                .bundles
266                .get(&id)
267                .context("created project is missing")?,
268            &targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15)),
269            false,
270        )?;
271        let canonical = crate::database::store_catalog_project(&id, &project, true)?;
272        if canonical != id {
273            let mut saved = crate::database::read_project_catalog()?
274                .projects
275                .into_iter()
276                .find(|entry| entry.bundle_id == canonical)
277                .context("canonical project is missing")?
278                .project;
279            for repository in &mut saved.bundle.repositories {
280                let identity = saved
281                    .identities
282                    .get(&repository.id)
283                    .context("canonical repository identity is missing")?;
284                if let Some(source) = project.bundle.repositories.iter().find(|source| {
285                    project.identities.get(&source.id) == Some(identity)
286                        && source.destination == repository.destination
287                }) {
288                    repository.local = source.local.clone();
289                    repository.github = source.github.clone();
290                }
291            }
292            config
293                .bundles
294                .insert(canonical.clone(), saved.bundle.clone());
295            config.bundles.remove(&id);
296            crate::database::store_catalog_project(&canonical, &saved, true)?;
297        }
298        Ok(canonical)
299    })
300    .map_err(|error| {
301        error
302            .downcast::<QuickBundleFailure>()
303            .unwrap_or_else(QuickBundleFailure::Persistence)
304    })?;
305    Ok(QuickBundleCreation { config, bundle_id })
306}
307
308/// Remove a saved project from the config. Only the config entry goes: no
309/// checkout, clone or file is touched. Refused while any session still uses
310/// the project, checked against fresh state so a dashboard's stale view
311/// cannot remove a project a new session just started with.
312pub fn remove_bundle(bundle_id: &str) -> Result<Config> {
313    let state = crate::database::load_state()?;
314    if let Some(refusal) = state.bundle_removal_refusal(bundle_id) {
315        bail!(refusal);
316    }
317    let (config, ()) = Config::update(|config| {
318        ensure!(
319            config.bundles.remove(bundle_id).is_some(),
320            "project {bundle_id:?} is no longer saved"
321        );
322        config.validate()
323    })?;
324    crate::database::hide_catalog_project(bundle_id)?;
325    Ok(config)
326}
327
328/// Add a bundle for all `sources` to an already-loaded config. The source
329/// interpretation is shared with the persisted [`create_bundle_from_sources`]
330/// entry point and the legacy quick-bundle helper.
331pub fn create_bundle_from_sources_in_config(
332    config: &mut Config,
333    sources: &[String],
334) -> Result<String> {
335    let sources = sources
336        .iter()
337        .map(|source| interpret_repository_source(source))
338        .collect::<Result<Vec<_>>>()?;
339    if sources.is_empty() {
340        bail!("at least one repository source is required");
341    }
342
343    let mut identities = BTreeSet::new();
344    for source in &sources {
345        if !identities.insert(source.identity()) {
346            bail!("duplicate repository source {:?}", source.display_name);
347        }
348    }
349
350    if let Some(existing) = exact_configured_bundle(config, &sources)? {
351        let bundle = config
352            .bundles
353            .get_mut(&existing)
354            .expect("matched project exists");
355        // Local selection carries its own fetch/push settings. Keep the saved
356        // repository IDs and layout; earlier sessions hold accepted snapshots.
357        for repository in &mut bundle.repositories {
358            let identity = crate::import::configured_repository_identity(repository)?
359                .context("configured identity is missing")?;
360            if let Some(source) = sources.iter().find(|source| {
361                source.identity == identity && matches!(source.kind, RepositorySourceKind::Local(_))
362            }) && let RepositorySourceKind::Local(root) = &source.kind
363            {
364                repository.local = Some(root.clone());
365                repository.github = None;
366            }
367        }
368        return Ok(existing);
369    }
370
371    // Build and validate a candidate before replacing the caller's config, so
372    // a later validation error cannot leave an in-memory partial mutation.
373    let mut updated = config.clone();
374    let mut used_repository_ids = BTreeSet::new();
375    let mut repositories = Vec::with_capacity(sources.len());
376    for source in sources {
377        let base = setup_style_id(&source.name);
378        let repository_id = unique_id(&base, |candidate| used_repository_ids.contains(candidate));
379        used_repository_ids.insert(repository_id.clone());
380        repositories.push(source.into_project_repository(repository_id));
381    }
382    let primary_repo = repositories
383        .first()
384        .map(|repository| repository.id.clone())
385        .context("at least one repository source is required")?;
386    let bundle_id = unique_id(&primary_repo, |candidate| {
387        updated.bundles.contains_key(candidate)
388    });
389    updated.bundles.insert(
390        bundle_id.clone(),
391        ProjectBundle {
392            primary_repo,
393            repositories,
394        },
395    );
396    updated.validate()?;
397    *config = updated;
398    Ok(bundle_id)
399}
400
401#[derive(Debug, Clone)]
402enum RepositorySourceKind {
403    Github(crate::setup::GithubRepository),
404    Local(PathBuf),
405}
406
407#[derive(Debug, Clone)]
408struct InterpretedRepositorySource {
409    display_name: String,
410    name: String,
411    kind: RepositorySourceKind,
412    identity: RepositoryIdentity,
413}
414
415impl InterpretedRepositorySource {
416    fn identity(&self) -> RepositoryIdentity {
417        self.identity.clone()
418    }
419
420    fn into_project_repository(self, id: String) -> ProjectRepository {
421        let (github, local) = match self.kind {
422            RepositorySourceKind::Github(repository) => (
423                Some(format!("{}/{}", repository.owner, repository.repository)),
424                None,
425            ),
426            RepositorySourceKind::Local(root) => (None, Some(root)),
427        };
428        ProjectRepository {
429            id: id.clone(),
430            github,
431            local,
432            destination: PathBuf::from(id),
433            git_ref: None,
434        }
435    }
436}
437
438/// Interpret a source once, including local Git canonicalization and GitHub
439/// parsing, so all creation paths use exactly the same source semantics.
440fn interpret_repository_source(source: &str) -> Result<InterpretedRepositorySource> {
441    let source = source.trim();
442    if source.is_empty() {
443        bail!("repository source cannot be empty");
444    }
445    let expanded = mj_core::path_input::expand_local(Path::new(source))?;
446    let candidate = expanded.as_path();
447    if candidate.exists() {
448        let resolved = mj_core::repository::resolve_directory(
449            candidate,
450            None,
451            &mj_core::targets::CancellableProcessExecutor::with_timeout(
452                std::time::Duration::from_secs(8),
453            ),
454        )?;
455        let identity = resolved.identity;
456        let name = identity.name();
457        return Ok(InterpretedRepositorySource {
458            display_name: source.to_owned(),
459            name,
460            identity,
461            kind: RepositorySourceKind::Local(resolved.checkout_root),
462        });
463    }
464    if candidate.is_absolute() || source.starts_with('.') || source.starts_with('~') {
465        bail!("local repository path {source:?} does not exist");
466    }
467    let repository = github_repository_from_origin(source).context(format!(
468        "{source:?} is not a GitHub owner/repository or URL"
469    ))?;
470    Ok(InterpretedRepositorySource {
471        display_name: source.to_owned(),
472        name: repository.repository.clone(),
473        identity: RepositoryIdentity::Github(
474            repository.owner.to_ascii_lowercase(),
475            repository.repository.to_ascii_lowercase(),
476        ),
477        kind: RepositorySourceKind::Github(repository),
478    })
479}
480
481fn exact_configured_bundle(
482    config: &Config,
483    requested: &[InterpretedRepositorySource],
484) -> Result<Option<String>> {
485    let requested_identities = requested
486        .iter()
487        .map(InterpretedRepositorySource::identity)
488        .collect::<BTreeSet<_>>();
489    let primary = requested
490        .first()
491        .context("no repository sources")?
492        .identity();
493    let mut used = BTreeSet::new();
494    let requested_layout = requested
495        .iter()
496        .map(|source| {
497            let id = unique_id(&setup_style_id(&source.name), |id| used.contains(id));
498            used.insert(id.clone());
499            (source.identity(), PathBuf::from(id))
500        })
501        .collect::<BTreeSet<_>>();
502    for (id, bundle) in &config.bundles {
503        if bundle.repositories.len() != requested.len()
504            || bundle
505                .repositories
506                .iter()
507                .any(|repository| repository.git_ref.is_some())
508        {
509            continue;
510        }
511        let layout = bundle
512            .repositories
513            .iter()
514            .map(|repo| {
515                Ok((
516                    crate::import::configured_repository_identity(repo)?
517                        .context("configured identity is missing")?,
518                    repo.destination.clone(),
519                ))
520            })
521            .collect::<Result<BTreeSet<_>>>()?;
522        if layout == requested_layout && bundle_matches(bundle, &requested_identities, &primary)? {
523            return Ok(Some(id.clone()));
524        }
525    }
526    Ok(None)
527}
528
529fn unique_id(base: &str, mut is_used: impl FnMut(&str) -> bool) -> String {
530    if !is_used(base) {
531        return base.to_owned();
532    }
533    for suffix in 2_u32.. {
534        let suffix = format!("-{suffix}");
535        let prefix_len = 64usize.saturating_sub(suffix.len());
536        let prefix = base.chars().take(prefix_len).collect::<String>();
537        let candidate = format!("{prefix}{suffix}");
538        if !is_used(&candidate) {
539            return candidate;
540        }
541    }
542    unreachable!("u32 repository/bundle id suffixes exhausted")
543}
544
545pub struct SessionLaunchOptions {
546    pub create_managed_worktree: Option<bool>,
547    /// Full commit ID to start the bundle's primary repository at.
548    pub at: Option<String>,
549    /// With `at`, the new branch created there; otherwise an existing branch.
550    pub branch: Option<String>,
551    /// Diff base; defaults to `at`.
552    pub base: Option<String>,
553    pub subagents: Option<mj_core::subagent::SubagentPolicy>,
554    pub initial_prompt: Option<String>,
555    pub workspace_id: String,
556    pub additional_mounts: Vec<AdditionalMount>,
557    pub resource_allocation: Option<SessionResourceAllocation>,
558    pub project_directory: Option<PathBuf>,
559    pub session_title_override: Option<String>,
560}
561
562pub struct SessionResumeOptions {
563    pub additional_mounts: Option<Vec<AdditionalMount>>,
564    pub resource_allocation: Option<SessionResourceAllocation>,
565    pub discard_queue: bool,
566}
567
568fn selected_host_container_size(
569    template: &TargetTemplate,
570    allocation: Option<&SessionResourceAllocation>,
571) -> Option<(String, HostContainerSize)> {
572    let host = container_size_host(template)?;
573    let SessionResourceAllocation::Container { cpus, memory_bytes } = allocation? else {
574        return None;
575    };
576    Some((
577        host.to_owned(),
578        HostContainerSize {
579            cpus: *cpus,
580            memory_bytes: *memory_bytes,
581        },
582    ))
583}
584
585impl Controller {
586    pub fn load() -> Result<Self> {
587        // A configuration the person has to fix, such as an entry whose secret
588        // is missing, is the reason a request fails, not an internal error.
589        let config = Config::load().map_err(|error| {
590            let refusal = mj_core::refusal::Refusal::precondition(format!("{error:#}"));
591            error.context(refusal)
592        })?;
593        let state = crate::database::load_state()?;
594        Ok(Self { config, state })
595    }
596
597    /// Startup-only traversal for legacy metadata and configuration diagnostics.
598    /// Ordinary readers consume the writer's committed snapshot directly.
599    pub(crate) fn prepare_persisted_sessions(&mut self) -> Result<()> {
600        // Missing session dependencies must not lock users out of the tools
601        // needed to repair them. Operations validate the session they act on.
602        self.state.validate()?;
603        for session in self.state.sessions.values() {
604            if let Some(issue) = session.configuration_issue(&self.config) {
605                tracing::warn!(session_id = %session.id, "{issue}");
606            }
607        }
608        self.backfill_target_runtime()
609    }
610
611    pub fn reload(&mut self) -> Result<()> {
612        // Capture legacy settings from the previous config before a manual
613        // edit removes them. Only fill missing metadata, never overwrite it.
614        self.backfill_target_runtime()?;
615        *self = Self::load()?;
616        Ok(())
617    }
618
619    fn backfill_target_runtime(&mut self) -> Result<()> {
620        let missing: Vec<_> = self
621            .state
622            .sessions
623            .values()
624            .filter(|session| session.target_runtime.is_none())
625            .map(|session| session.id.clone())
626            .collect();
627        for id in missing {
628            let session = self.state.sessions.get_mut(&id).expect("collected session");
629            match session.target_runtime_settings(&self.config) {
630                Ok(runtime) => {
631                    let runtime = runtime.into_owned();
632                    session.target_runtime = crate::database::backfill_target_runtime(
633                        &session.id,
634                        &session.target_template_id,
635                        &runtime,
636                    )?;
637                }
638                Err(error) => {
639                    tracing::warn!(session_id = %session.id, %error, "target access needs configuration repair")
640                }
641            }
642        }
643        Ok(())
644    }
645
646    fn persist_session_state(&self, session_id: &str) -> Result<()> {
647        match self.state.sessions.get(session_id) {
648            Some(session) => crate::database::save_lifecycle_session(session),
649            None => crate::database::delete_session(session_id),
650        }
651    }
652
653    fn persist_session_transition_or_restore(
654        &mut self,
655        session_id: &str,
656        previous: &SessionRecord,
657        context: &'static str,
658    ) -> Result<()> {
659        persist_session_record_transition_or_restore(
660            &mut self.state,
661            session_id,
662            previous,
663            context,
664            &crate::database::save_lifecycle_session,
665        )
666    }
667
668    fn restore_prior_session_after_persistence_failure(
669        &mut self,
670        session_id: &str,
671        previous: &SessionRecord,
672        primary: anyhow::Error,
673    ) -> anyhow::Error {
674        restore_session_after_persistence_failure(
675            &mut self.state,
676            session_id,
677            previous,
678            primary,
679            crate::database::save_lifecycle_session,
680        )
681    }
682
683    /// Resolve an entered path on its owning host. Call only from background work.
684    pub fn resolve_input_path(
685        &self,
686        target_id: &str,
687        path: &Path,
688        executor: &impl CommandExecutor,
689    ) -> Result<PathBuf> {
690        let target = self
691            .config
692            .targets
693            .get(target_id)
694            .context("Unknown path target")?;
695        resolve_target_input_path(target, path, executor)
696    }
697
698    /// Verify a mount source on the host where Mjolnir will consume it, and report
699    /// the filesystem reason it must be attached read-only, if there is one.
700    ///
701    /// The probe runs in the same round trip as the existence check so the
702    /// editor learns both answers without a second wait. A probe that cannot
703    /// answer reports no reason: provisioning decides that authoritatively.
704    pub fn validate_mount_source(
705        &self,
706        target_id: &str,
707        source: &Path,
708        executor: &impl CommandExecutor,
709    ) -> Result<Option<String>> {
710        let target = self
711            .config
712            .targets
713            .get(target_id)
714            .with_context(|| format!("unknown target template {target_id:?}"))?;
715        let exists = match target {
716            TargetTemplate::LocalPodman { .. }
717            | TargetTemplate::LocalDocker { .. }
718            | TargetTemplate::AppleContainer { .. }
719            | TargetTemplate::AwsEc2 { .. } => std::fs::metadata(source)
720                .map(|metadata| metadata.is_dir())
721                .or_else(|error| {
722                    if error.kind() == std::io::ErrorKind::NotFound {
723                        Ok(false)
724                    } else {
725                        Err(error)
726                    }
727                })
728                .with_context(|| format!("inspect resource source {}", source.display()))?,
729            TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
730                targets::ssh_directory_exists(&SshTarget::from(ssh), source, executor)?
731            }
732            TargetTemplate::LocalBare | TargetTemplate::SshBare { .. } => {
733                bail!("resource attachments are unsupported for bare targets")
734            }
735        };
736        ensure!(
737            exists,
738            "source path {} does not exist or is not a directory",
739            source.display()
740        );
741        Ok(self.forced_read_only_reason(target, source, executor))
742    }
743
744    /// The `filesystem (reason)` label for a source the runtime cannot overlay.
745    fn forced_read_only_reason(
746        &self,
747        target: &TargetTemplate,
748        source: &Path,
749        executor: &impl CommandExecutor,
750    ) -> Option<String> {
751        let ssh = match target {
752            TargetTemplate::LocalPodman { .. } | TargetTemplate::LocalDocker { .. } => None,
753            TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
754                Some(SshTarget::from(ssh))
755            }
756            // Apple Container already mounts read-only, and EC2 copies instead
757            // of mounting, so neither has an overlay to lose.
758            _ => return None,
759        };
760        if matches!(target, TargetTemplate::LocalDocker { .. })
761            && let Ok(Some(_)) = targets::local_docker_vm_share(executor)
762        {
763            return Some("Docker VM file share (cannot back an overlay)".to_owned());
764        }
765        let filesystem = targets::probe_filesystem_types(
766            ssh.as_ref(),
767            std::slice::from_ref(&source.to_path_buf()),
768            executor,
769        )
770        .map_err(|error| {
771            tracing::debug!(
772                source = %source.display(),
773                error = format!("{error:#}"),
774                "could not probe the filesystem under a mount source"
775            );
776        })
777        .ok()?
778        .pop()?;
779        let reason = targets::overlay_unsupported_filesystem(&filesystem)?;
780        Some(format!("{filesystem} ({reason})"))
781    }
782
783    fn fail_new_session_with_cleanup(
784        &mut self,
785        session_id: &str,
786        error: anyhow::Error,
787        executor: &impl CommandExecutor,
788    ) -> Result<anyhow::Error> {
789        let original = provisioning::note_new_session_launch_failure(session_id, &error);
790        let cleanup_error = self
791            .cleanup_new_session_worktree_after_failure(session_id, executor)
792            .err()
793            .map(|cleanup_error| format!("{cleanup_error:#}"));
794        if let Some(cleanup_error) = &cleanup_error {
795            tracing::warn!(
796                session_id,
797                error = %cleanup_error,
798                "new-session worktree rollback reported a cleanup failure"
799            );
800        }
801        let failure = apply_failed_new_session_rollback(
802            &mut self.state,
803            session_id,
804            &original,
805            cleanup_error,
806        );
807        self.persist_session_state(session_id)?;
808        Ok(failure)
809    }
810
811    pub fn register_session_with_resources(
812        &mut self,
813        profile_id: &str,
814        bundle_id: &str,
815        target_id: &str,
816        title: impl Into<String>,
817        options: SessionLaunchOptions,
818    ) -> Result<String> {
819        self.register_session_with_executor(
820            profile_id,
821            bundle_id,
822            target_id,
823            title,
824            options,
825            &targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15)),
826        )
827    }
828
829    fn register_session_with_executor(
830        &mut self,
831        profile_id: &str,
832        bundle_id: &str,
833        target_id: &str,
834        title: impl Into<String>,
835        options: SessionLaunchOptions,
836        executor: &impl CommandExecutor,
837    ) -> Result<String> {
838        let SessionLaunchOptions {
839            create_managed_worktree,
840            at,
841            branch,
842            base,
843            subagents,
844            initial_prompt,
845            workspace_id,
846            additional_mounts,
847            resource_allocation,
848            project_directory,
849            session_title_override,
850        } = options;
851        let launch_base = match base {
852            Some(base) => {
853                let base = base.trim();
854                if base.is_empty() {
855                    bail!("`base` must not be empty");
856                }
857                Some(base.to_owned())
858            }
859            None => None,
860        };
861        let branch = match branch {
862            Some(branch) => {
863                let branch = branch.trim();
864                ensure!(!branch.is_empty(), "`branch` must not be empty");
865                Some(branch.to_owned())
866            }
867            None => None,
868        };
869        if launch_base.is_some() && create_managed_worktree == Some(false) {
870            bail!("`base` requires a managed worktree or a bundle session");
871        }
872        let session_title_override = match session_title_override {
873            Some(title) => {
874                Some(normalize_session_title(&title).context("session name cannot be empty")?)
875            }
876            None => None,
877        };
878        let profile = self
879            .config
880            .profiles
881            .get(profile_id)
882            .with_context(|| format!("unknown profile {profile_id:?}"))?;
883        ensure!(profile.enabled, "profile {profile_id:?} is disabled");
884        let template = self
885            .config
886            .targets
887            .get(target_id)
888            .with_context(|| format!("unknown target template {target_id:?}"))?;
889        if create_managed_worktree == Some(true) && !is_bare_project_target(template) {
890            bail!("managed worktree creation requires a bare Git project");
891        }
892        if project_directory.is_some() != is_bare_project_target(template) {
893            bail!("raw project directories require a bare target, and bare targets require one");
894        }
895        if let Some(path) = &project_directory
896            && (!path.is_absolute()
897                || path
898                    .components()
899                    .any(|part| part == std::path::Component::ParentDir))
900        {
901            bail!("bare project directory must be an absolute safe path");
902        }
903        let bundle = project_directory
904            .is_none()
905            .then(|| self.config.bundles.get(bundle_id))
906            .flatten();
907        if project_directory.is_none() && bundle.is_none() {
908            bail!("unknown bundle {bundle_id:?}");
909        }
910        // `at` becomes an exact checkout of the bundle's primary repository,
911        // which then owns the branch. Without it, the branch is an existing
912        // one to check out.
913        let (checkout, launch_branch) = match at {
914            Some(at) => {
915                let bundle = bundle.context("`at` requires a bundle-backed session")?;
916                ensure!(
917                    create_managed_worktree != Some(false),
918                    "`at` requires an isolated workspace"
919                );
920                let mut checkout = mj_core::remote_git::ExactCheckout {
921                    repository_id: bundle.primary_repo.clone(),
922                    commit: at.trim().to_owned(),
923                    branch,
924                };
925                checkout.validate()?;
926                checkout.commit.make_ascii_lowercase();
927                (Some(checkout), None)
928            }
929            None => (None, branch),
930        };
931        if profile.kind == mj_core::config::HarnessKind::Muse
932            && (!additional_mounts.is_empty()
933                || bundle.is_some_and(|bundle| bundle.repositories.len() > 1))
934        {
935            bail!(
936                "{} ACP supports one workspace root; use a single-repository bundle without attached directories",
937                profile.kind.display_name()
938            );
939        }
940        let (bundle_id, project) = if let Some(bundle) = bundle {
941            let project = crate::project_catalog::snapshot(bundle, executor, true)?;
942            let canonical = crate::database::store_catalog_project(bundle_id, &project, true)?;
943            (canonical, Some(project))
944        } else if let Some(directory) = project_directory.as_deref() {
945            let git = match template {
946                TargetTemplate::LocalBare => {
947                    executor
948                        .execute(
949                            &targets::CommandSpec::new(
950                                "git",
951                                [
952                                    "-C",
953                                    &directory.to_string_lossy(),
954                                    "rev-parse",
955                                    "--is-inside-work-tree",
956                                ],
957                            )
958                            .purpose("inspect raw project repository"),
959                        )?
960                        .status
961                        == 0
962                }
963                _ => true,
964            };
965            if git {
966                let (id, project) =
967                    crate::project_catalog::accept_directory(template, directory, executor)?;
968                (id, Some(project))
969            } else {
970                (bundle_id.to_owned(), None)
971            }
972        } else {
973            (bundle_id.to_owned(), None)
974        };
975        validate_resource_allocation(template, resource_allocation.as_ref())?;
976        let selected_container_size =
977            selected_host_container_size(template, resource_allocation.as_ref());
978        if !additional_mounts.is_empty() && mount_history_host(template).is_none() {
979            bail!("attached resources are unsupported for this target");
980        }
981        targets::validate_additional_mounts(&additional_mounts)?;
982        let id = new_session_id()?;
983        let now = now();
984        let record = SessionRecord {
985            project,
986            target_runtime: Some(template.into()),
987            build_cache: None,
988            create_managed_worktree,
989            launch_base,
990            launch_branch,
991            checkout,
992            publication: None,
993            subagents: Some(subagents.unwrap_or_else(|| profile.subagents.clone())),
994            archived: false,
995            container_cpus: None,
996            container_memory: None,
997            // Recorded for every new session, container-backed or not, so a
998            // later move into a container already knows the path its checkout
999            // will occupy. Only sessions that predate per-session container
1000            // workspaces leave it unset.
1001            container_workspace: Some(targets::new_container_workspace(&id)?),
1002            id: id.clone(),
1003            workspace_id,
1004            title: title.into(),
1005            harness_kind: profile.kind,
1006            last_profile: profile_id.to_string(),
1007            bundle_id,
1008            project_directory,
1009            managed_worktree: None,
1010            target_template_id: target_id.to_string(),
1011            resource_allocation,
1012            additional_mounts: additional_mounts.clone(),
1013            state: SessionState::Provisioning,
1014            target: None,
1015            native_session_id: None,
1016            acp_session_title: None,
1017            session_title_override,
1018            created_at: now.clone(),
1019            updated_at: now,
1020            viewed_through_event_ordinal: 0,
1021            draft_input: initial_prompt.unwrap_or_default(),
1022            last_error: None,
1023            last_checkpoint_error: None,
1024            checkpoint: None,
1025        };
1026        // Creation authors the whole record, so it writes the whole row. The
1027        // record reaches memory only once it is durable: a session this process
1028        // alone knows about is one the database can never resume or clean up.
1029        crate::database::save_new_session(&record, selected_container_size.clone())?;
1030        if record.harness_kind.supports_delegation_tools() {
1031            self.state.last_subagent_policy = record.subagents.clone().unwrap_or_default();
1032        }
1033        self.state.sessions.insert(id.clone(), record);
1034        if let Some((host, size)) = selected_container_size {
1035            self.state.remember_container_size(&host, size);
1036        }
1037        if let Some(host) = mount_history_host(template) {
1038            // Mount history only seeds the attach dialog's suggestions. The
1039            // session row is already committed, so a failed suggestion write is
1040            // reported rather than turned into a failed registration.
1041            match crate::database::remember_mount_sources(host, &additional_mounts) {
1042                Ok(()) => self.state.remember_mount_sources(host, &additional_mounts),
1043                Err(error) => tracing::warn!(
1044                    session_id = id,
1045                    error = format!("{error:#}"),
1046                    "could not remember the attached resource directories for later suggestions"
1047                ),
1048            }
1049        }
1050        Ok(id)
1051    }
1052
1053    pub fn rename_session(&mut self, session_id: &str, title: &str) -> Result<String> {
1054        let title = normalize_session_title(title).context("session name cannot be empty")?;
1055        ensure!(
1056            self.state.sessions.contains_key(session_id),
1057            "unknown session {session_id}"
1058        );
1059        let updated_at = now();
1060        crate::database::set_session_title_override(session_id, &title, &updated_at)?;
1061        let record = self
1062            .state
1063            .sessions
1064            .get_mut(session_id)
1065            .expect("session was checked before updating its title");
1066        record.session_title_override = Some(title.clone());
1067        record.updated_at = updated_at;
1068        Ok(title)
1069    }
1070
1071    /// Move a session, and the sub-agents under it, to another workspace.
1072    /// Returns how many sessions moved.
1073    pub fn set_session_workspace(&mut self, session_id: &str, workspace_id: &str) -> Result<usize> {
1074        ensure!(
1075            self.state.sessions.contains_key(session_id),
1076            "unknown session {session_id}"
1077        );
1078        let mut ids = vec![session_id.to_owned()];
1079        ids.extend(
1080            self.state
1081                .subagents
1082                .values()
1083                .filter(|child| child.parent_session_id == session_id)
1084                .map(|child| child.child_session_id.clone()),
1085        );
1086        crate::database::set_sessions_workspace(&ids, workspace_id)?;
1087        for id in &ids {
1088            if let Some(record) = self.state.sessions.get_mut(id) {
1089                record.workspace_id = workspace_id.to_owned();
1090            }
1091        }
1092        Ok(ids.len())
1093    }
1094
1095    pub fn rename_profile_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
1096        mj_core::config::validate_id("profile", new_id)?;
1097        if old_id == new_id {
1098            ensure!(
1099                self.config.profiles.contains_key(old_id),
1100                "unknown profile {old_id:?}"
1101            );
1102            return Ok(());
1103        }
1104        let journal = ConfigRenameJournal {
1105            kind: ConfigRenameKind::Profile,
1106            old_id: old_id.to_owned(),
1107            new_id: new_id.to_owned(),
1108        };
1109        write_config_rename_journal(&journal)?;
1110        let (config, ()) = match Config::update(|config| {
1111            ensure!(
1112                config.profiles.contains_key(old_id),
1113                "unknown profile {old_id:?}"
1114            );
1115            ensure!(
1116                !config.profiles.contains_key(new_id),
1117                "profile {new_id:?} already exists"
1118            );
1119            let profile = config
1120                .profiles
1121                .remove(old_id)
1122                .expect("profile was checked in the transaction");
1123            config.profiles.insert(new_id.to_owned(), profile);
1124            Ok(())
1125        }) {
1126            Ok(result) => result,
1127            Err(error) => {
1128                remove_config_rename_journal()
1129                    .context("remove profile rename journal after config save failed")?;
1130                return Err(error).context("save renamed profile configuration");
1131            }
1132        };
1133        self.config = config;
1134        mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
1135        if let Err(error) = crate::database::rename_profile_references(old_id, new_id) {
1136            let restore = Config::update(|config| {
1137                let profile = config
1138                    .profiles
1139                    .remove(new_id)
1140                    .with_context(|| format!("renamed profile {new_id:?} is missing"))?;
1141                ensure!(
1142                    !config.profiles.contains_key(old_id),
1143                    "cannot restore profile rename: both {old_id:?} and {new_id:?} exist"
1144                );
1145                config.profiles.insert(old_id.to_owned(), profile);
1146                Ok(())
1147            });
1148            let restored = match restore {
1149                Ok((config, ())) => config,
1150                Err(restore_error) => {
1151                    return Err(error).context(format!(
1152                        "rename profile references; additionally failed to restore config: {restore_error:#}"
1153                    ));
1154                }
1155            };
1156            self.config = restored;
1157            if let Err(restore_error) = remove_config_rename_journal() {
1158                return Err(error).context(format!(
1159                    "rename profile references; additionally failed to remove rename journal: {restore_error:#}"
1160                ));
1161            }
1162            return Err(error).context("rename profile references");
1163        }
1164        let affected: Vec<_> = self
1165            .state
1166            .sessions
1167            .values()
1168            .filter(|session| session.last_profile == old_id)
1169            .map(|session| session.id.clone())
1170            .collect();
1171        for id in affected {
1172            self.state
1173                .sessions
1174                .get_mut(&id)
1175                .expect("collected session")
1176                .last_profile = new_id.to_owned();
1177        }
1178        remove_config_rename_journal()?;
1179        Ok(())
1180    }
1181
1182    pub fn rename_target_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
1183        mj_core::config::validate_id("target template", new_id)?;
1184        if old_id == new_id {
1185            ensure!(
1186                self.config.targets.contains_key(old_id),
1187                "unknown target {old_id:?}"
1188            );
1189            return Ok(());
1190        }
1191        let journal = ConfigRenameJournal {
1192            kind: ConfigRenameKind::Target,
1193            old_id: old_id.to_owned(),
1194            new_id: new_id.to_owned(),
1195        };
1196        write_config_rename_journal(&journal)?;
1197        let (config, ()) = match Config::update(|config| {
1198            ensure!(
1199                config.targets.contains_key(old_id),
1200                "unknown target {old_id:?}"
1201            );
1202            ensure!(
1203                !config.targets.contains_key(new_id),
1204                "target {new_id:?} already exists"
1205            );
1206            let target = config
1207                .targets
1208                .remove(old_id)
1209                .expect("target was checked in the transaction");
1210            config.targets.insert(new_id.to_owned(), target);
1211            Ok(())
1212        }) {
1213            Ok(result) => result,
1214            Err(error) => {
1215                remove_config_rename_journal()
1216                    .context("remove target rename journal after config save failed")?;
1217                return Err(error).context("save renamed target configuration");
1218            }
1219        };
1220        self.config = config;
1221        mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
1222        if let Err(error) = crate::database::rename_target_references(old_id, new_id) {
1223            let restore = Config::update(|config| {
1224                let target = config
1225                    .targets
1226                    .remove(new_id)
1227                    .with_context(|| format!("renamed target {new_id:?} is missing"))?;
1228                ensure!(
1229                    !config.targets.contains_key(old_id),
1230                    "cannot restore target rename: both {old_id:?} and {new_id:?} exist"
1231                );
1232                config.targets.insert(old_id.to_owned(), target);
1233                Ok(())
1234            });
1235            let restored = match restore {
1236                Ok((config, ())) => config,
1237                Err(restore_error) => {
1238                    return Err(error).context(format!(
1239                        "rename target references; additionally failed to restore config: {restore_error:#}"
1240                    ));
1241                }
1242            };
1243            self.config = restored;
1244            if let Err(restore_error) = remove_config_rename_journal() {
1245                return Err(error).context(format!(
1246                    "rename target references; additionally failed to remove rename journal: {restore_error:#}"
1247                ));
1248            }
1249            return Err(error).context("rename target references");
1250        }
1251        let affected: Vec<_> = self
1252            .state
1253            .sessions
1254            .values()
1255            .filter(|session| session.target_template_id == old_id)
1256            .map(|session| session.id.clone())
1257            .collect();
1258        for id in affected {
1259            self.state
1260                .sessions
1261                .get_mut(&id)
1262                .expect("collected session")
1263                .target_template_id = new_id.to_owned();
1264        }
1265        remove_config_rename_journal()?;
1266        Ok(())
1267    }
1268
1269    /// Finish a profile/target id rename interrupted between the atomic config
1270    /// replacement and SQLite transaction. Each step is idempotent, so a
1271    /// second crash leaves the same intent available for the next startup.
1272    pub fn recover_config_id_rename() -> Result<bool> {
1273        let path = config_rename_journal_path();
1274        let body = match fs::read(&path) {
1275            Ok(body) => body,
1276            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(false),
1277            Err(error) => return Err(error).context(format!("read {}", path.display())),
1278        };
1279        let journal: ConfigRenameJournal =
1280            serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1281        match journal.kind {
1282            ConfigRenameKind::Profile => {
1283                Config::update(|config| {
1284                    finish_config_map_rename(
1285                        &mut config.profiles,
1286                        &journal.old_id,
1287                        &journal.new_id,
1288                        "profile",
1289                    )?;
1290                    Ok(())
1291                })?;
1292                crate::database::rename_profile_references(&journal.old_id, &journal.new_id)?;
1293            }
1294            ConfigRenameKind::Target => {
1295                // The file's own entries, not `Config::update`'s view: that one
1296                // adds the standard local targets, so after a built-in id such
1297                // as `localhost` was renamed away, its default would reappear
1298                // beside the new id and read as a rename that cannot finish.
1299                Config::update_to(&mj_core::config::config_path(), |config| {
1300                    finish_config_map_rename(
1301                        &mut config.targets,
1302                        &journal.old_id,
1303                        &journal.new_id,
1304                        "target",
1305                    )?;
1306                    Ok(())
1307                })?;
1308                crate::database::rename_target_references(&journal.old_id, &journal.new_id)?;
1309            }
1310        }
1311        remove_config_rename_journal()?;
1312        Ok(true)
1313    }
1314
1315    /// Record the per-session container size overrides and attached
1316    /// directories. Nothing is applied to a running container: the values are
1317    /// read the next time the session's container is created.
1318    pub fn update_session_container_settings(
1319        &mut self,
1320        session_id: &str,
1321        cpus: Option<String>,
1322        memory: Option<String>,
1323        additional_mounts: Vec<targets::AdditionalMount>,
1324        mount_history: Vec<std::path::PathBuf>,
1325        executor: &impl CommandExecutor,
1326    ) -> Result<()> {
1327        let session = self
1328            .state
1329            .sessions
1330            .get(session_id)
1331            .with_context(|| format!("unknown session {session_id}"))?;
1332        // A directory the runtime cannot find is created by Docker, owned by
1333        // root, when the container is recreated (launch finding J-24). Check
1334        // each newly attached source where the container will read it: on
1335        // this host for a local runtime, over the machine's connection for an
1336        // SSH one. A directory already attached was checked when it was
1337        // added, and a size change must not fail because it has since gone.
1338        for mount in &additional_mounts {
1339            if session
1340                .additional_mounts
1341                .iter()
1342                .any(|attached| attached.source == mount.source)
1343            {
1344                continue;
1345            }
1346            self.validate_mount_source(&session.target_template_id, &mount.source, executor)?;
1347        }
1348        let cpus = cpus.filter(|value| !value.trim().is_empty());
1349        let memory = memory.filter(|value| !value.trim().is_empty());
1350        let updated_at = now();
1351        crate::database::set_session_container_settings(
1352            session_id,
1353            cpus.as_deref(),
1354            memory.as_deref(),
1355            &additional_mounts,
1356            &updated_at,
1357        )?;
1358        if let Some(host) = self
1359            .config
1360            .targets
1361            .get(
1362                &self.state.sessions[session_id]
1363                    .target_template_id
1364                    .to_owned(),
1365            )
1366            .and_then(mj_core::config::mount_history_host)
1367        {
1368            let host = host.to_owned();
1369            // The dialog owns the suggestion list, so forgetting a directory
1370            // there has to survive the mounts being remembered right after.
1371            crate::database::replace_mount_history(&host, &mount_history)?;
1372            crate::database::remember_mount_sources(&host, &additional_mounts)?;
1373            self.state.mount_history.insert(host.clone(), mount_history);
1374            self.state.remember_mount_sources(&host, &additional_mounts);
1375        }
1376        let record = self
1377            .state
1378            .sessions
1379            .get_mut(session_id)
1380            .expect("session was checked before updating its container settings");
1381        record.container_cpus = cpus;
1382        record.container_memory = memory;
1383        record.additional_mounts = additional_mounts;
1384        record.updated_at = updated_at;
1385        Ok(())
1386    }
1387}
1388
1389fn config_rename_journal_path() -> PathBuf {
1390    data_dir().join(CONFIG_RENAME_JOURNAL)
1391}
1392
1393fn write_config_rename_journal(journal: &ConfigRenameJournal) -> Result<()> {
1394    let path = config_rename_journal_path();
1395    let body = serde_json::to_vec(journal).context("serialize config rename journal")?;
1396    atomic_write(&path, &body).with_context(|| format!("write {}", path.display()))
1397}
1398
1399fn remove_config_rename_journal() -> Result<()> {
1400    let path = config_rename_journal_path();
1401    match fs::remove_file(&path) {
1402        Ok(()) => Ok(()),
1403        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
1404        Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
1405    }
1406}
1407
1408fn finish_config_map_rename<T>(
1409    entries: &mut BTreeMap<String, T>,
1410    old_id: &str,
1411    new_id: &str,
1412    kind: &str,
1413) -> Result<()> {
1414    if let Some(entry) = entries.remove(old_id) {
1415        ensure!(
1416            !entries.contains_key(new_id),
1417            "cannot recover {kind} rename: both {old_id:?} and {new_id:?} exist"
1418        );
1419        entries.insert(new_id.to_owned(), entry);
1420    } else {
1421        ensure!(
1422            entries.contains_key(new_id),
1423            "cannot recover {kind} rename: neither {old_id:?} nor {new_id:?} exists"
1424        );
1425    }
1426    Ok(())
1427}
1428
1429/// Where this session's harness reads and writes its profile inside the target.
1430///
1431/// This is the per-session root [`removable_profile_root`] names, on every
1432/// target: a session always runs from a staged copy of its profile, never from
1433/// the profile home itself. Muse keeps its state in a `muse` subdirectory of
1434/// the root, because its ACP adapter owns the directory it is given.
1435#[cfg(test)]
1436pub(crate) fn target_profile_home_for_test(
1437    locator: &targets::TargetLocator,
1438    session_id: &str,
1439    profile: &mj_core::config::HarnessProfile,
1440) -> String {
1441    target_profile_home(locator, session_id, profile)
1442}
1443
1444fn target_profile_home(
1445    locator: &targets::TargetLocator,
1446    session_id: &str,
1447    profile: &mj_core::config::HarnessProfile,
1448) -> String {
1449    let root = removable_profile_root(locator, session_id, profile);
1450    if profile.kind == mj_core::config::HarnessKind::Muse {
1451        PathBuf::from(root)
1452            .join("muse")
1453            .to_string_lossy()
1454            .into_owned()
1455    } else {
1456        root
1457    }
1458}
1459
1460/// The per-session profile directory a session's harness home lives in, which
1461/// the session owns and an in-place harness replacement or teardown may delete.
1462///
1463/// This is the root that [`target_profile_home`] derives its answer from, not
1464/// that answer itself: a Muse session's home is a `muse` subdirectory of a
1465/// per-session root, and the whole root is what belongs to the session.
1466pub(super) fn removable_profile_root(
1467    locator: &targets::TargetLocator,
1468    session_id: &str,
1469    profile: &mj_core::config::HarnessProfile,
1470) -> String {
1471    match locator {
1472        targets::TargetLocator::LocalBare { .. }
1473            if profile.kind == mj_core::config::HarnessKind::Muse =>
1474        {
1475            targets::local_muse_profile_root(session_id)
1476                .to_string_lossy()
1477                .into_owned()
1478        }
1479        targets::TargetLocator::LocalBare { worker_root } => Path::new(worker_root)
1480            .join("profile")
1481            .to_string_lossy()
1482            .into_owned(),
1483        targets::TargetLocator::LocalPodman { .. }
1484        | targets::TargetLocator::LocalDocker { .. }
1485        | targets::TargetLocator::AppleContainer { .. }
1486        | targets::TargetLocator::SshPodman { .. }
1487        | targets::TargetLocator::SshDocker { .. } => {
1488            format!("/var/lib/hel/profiles/{session_id}")
1489        }
1490        targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. } => {
1491            format!(".local/share/hel/profiles/{session_id}")
1492        }
1493    }
1494}
1495
1496/// Resolve the login home on the machine that owns an editable path.
1497pub fn resolve_target_input_path(
1498    target: &TargetTemplate,
1499    path: &Path,
1500    executor: &impl CommandExecutor,
1501) -> Result<PathBuf> {
1502    if !mj_core::path_input::needs_home(path)? {
1503        return Ok(path.to_path_buf());
1504    }
1505    let host = cache_host::CacheHost::for_path_target(target)?;
1506    mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1507}
1508
1509/// Resolve the login home on a configured machine, for the Settings screen's
1510/// path fields. A machine, not a runtime, is what owns a home directory.
1511pub fn resolve_machine_input_path(
1512    machine: &mj_core::config::Machine,
1513    path: &Path,
1514    executor: &impl CommandExecutor,
1515) -> Result<PathBuf> {
1516    if !mj_core::path_input::needs_home(path)? {
1517        return Ok(path.to_path_buf());
1518    }
1519    // An EC2 instance does not exist until a session starts, so the only home
1520    // this screen can resolve is this machine's.
1521    let host = cache_host::CacheHost::for_path_machine(machine)?;
1522    mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1523}
1524
1525fn execute_checked(executor: &impl CommandExecutor, command: CommandSpec) -> Result<CommandOutput> {
1526    let output = executor.execute(&command)?;
1527    if output.status != 0 {
1528        let detail = command_error_detail(&output.stderr);
1529        if detail.is_empty() {
1530            bail!("{} failed with status {}", command.purpose, output.status);
1531        }
1532        bail!("{detail}");
1533    }
1534    Ok(output)
1535}
1536
1537fn command_error_detail(stderr: &[u8]) -> String {
1538    let reported = String::from_utf8_lossy(stderr);
1539    let reported = reported.trim();
1540    let detail = reported
1541        .rsplit_once("\nCaused by:\n")
1542        .map_or(reported, |(_, causes)| causes);
1543    let detail = detail.strip_prefix("Error: ").unwrap_or(detail);
1544    detail
1545        .lines()
1546        .map(|line| line.strip_prefix("    ").unwrap_or(line))
1547        .collect::<Vec<_>>()
1548        .join("\n")
1549        .trim()
1550        .to_owned()
1551}
1552
1553fn now() -> String {
1554    Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1555}
1556
1557fn restore_session_after_persistence_failure(
1558    state: &mut State,
1559    session_id: &str,
1560    previous: &SessionRecord,
1561    primary: anyhow::Error,
1562    persist: impl FnOnce(&SessionRecord) -> Result<()>,
1563) -> anyhow::Error {
1564    state
1565        .sessions
1566        .insert(session_id.to_owned(), previous.clone());
1567    let restored = state
1568        .sessions
1569        .get(session_id)
1570        .expect("restored session record disappeared");
1571    match persist(restored) {
1572        Ok(()) => primary,
1573        Err(error) => primary.context(format!(
1574            "restored prior session state in memory, but failed to persist the rollback: {error:#}"
1575        )),
1576    }
1577}
1578
1579fn persist_session_record_transition_or_restore(
1580    state: &mut State,
1581    session_id: &str,
1582    previous: &SessionRecord,
1583    context: &'static str,
1584    persist: &impl Fn(&SessionRecord) -> Result<()>,
1585) -> Result<()> {
1586    let result = persist(
1587        state
1588            .sessions
1589            .get(session_id)
1590            .expect("checkpoint session disappeared before persistence"),
1591    );
1592    match result {
1593        Ok(()) => Ok(()),
1594        Err(error) => Err(restore_session_after_persistence_failure(
1595            state,
1596            session_id,
1597            previous,
1598            error.context(context),
1599            persist,
1600        )),
1601    }
1602}
1603
1604pub fn config_only_controller(config: Config) -> Controller {
1605    Controller {
1606        config,
1607        state: State::default(),
1608    }
1609}
1610
1611#[cfg(test)]
1612mod tests;