Skip to main content

mj_controller/
controller.rs

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