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