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