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