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