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