1mod 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_remote_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#[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 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 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 let _ = self.file.unlock();
169 }
170}
171
172#[derive(Debug)]
176pub struct QuickBundleCreation {
177 pub config: Config,
178 pub bundle_id: String,
179}
180
181#[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
207pub 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
227pub 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
261pub 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
282pub 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 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
380fn 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 mjolnir_subagents: Option<bool>,
456 pub initial_prompt: Option<String>,
457 pub workspace_id: String,
458 pub additional_mounts: Vec<AdditionalMount>,
459 pub resource_allocation: Option<SessionResourceAllocation>,
460 pub project_directory: Option<PathBuf>,
461 pub session_title_override: Option<String>,
462}
463
464pub struct SessionResumeOptions {
465 pub additional_mounts: Option<Vec<AdditionalMount>>,
466 pub resource_allocation: Option<SessionResourceAllocation>,
467 pub discard_queue: bool,
468}
469
470fn selected_host_container_size(
471 template: &TargetTemplate,
472 allocation: Option<&SessionResourceAllocation>,
473) -> Option<(String, HostContainerSize)> {
474 let host = container_size_host(template)?;
475 let SessionResourceAllocation::Container { cpus, memory_bytes } = allocation? else {
476 return None;
477 };
478 Some((
479 host.to_owned(),
480 HostContainerSize {
481 cpus: *cpus,
482 memory_bytes: *memory_bytes,
483 },
484 ))
485}
486
487impl Controller {
488 pub fn load() -> Result<Self> {
489 let config = Config::load()?;
490 let state = crate::database::load_state()?;
491 state.validate()?;
494 for session in state.sessions.values() {
495 if let Some(issue) = session.configuration_issue(&config) {
496 tracing::warn!(session_id = %session.id, "{issue}");
497 }
498 }
499 Ok(Self { config, state })
500 }
501
502 pub fn reload(&mut self) -> Result<()> {
503 *self = Self::load()?;
504 Ok(())
505 }
506
507 fn persist_session_state(&self, session_id: &str) -> Result<()> {
508 match self.state.sessions.get(session_id) {
509 Some(session) => crate::database::save_lifecycle_session(session),
510 None => crate::database::delete_session(session_id),
511 }
512 }
513
514 fn persist_session_transition_or_restore(
515 &mut self,
516 session_id: &str,
517 previous: &SessionRecord,
518 context: &'static str,
519 ) -> Result<()> {
520 persist_session_record_transition_or_restore(
521 &mut self.state,
522 session_id,
523 previous,
524 context,
525 &crate::database::save_lifecycle_session,
526 )
527 }
528
529 fn restore_prior_session_after_persistence_failure(
530 &mut self,
531 session_id: &str,
532 previous: &SessionRecord,
533 primary: anyhow::Error,
534 ) -> anyhow::Error {
535 restore_session_after_persistence_failure(
536 &mut self.state,
537 session_id,
538 previous,
539 primary,
540 crate::database::save_lifecycle_session,
541 )
542 }
543
544 pub fn resolve_input_path(
546 &self,
547 target_id: &str,
548 path: &Path,
549 executor: &impl CommandExecutor,
550 ) -> Result<PathBuf> {
551 let target = self
552 .config
553 .targets
554 .get(target_id)
555 .context("Unknown path target")?;
556 resolve_target_input_path(target, path, executor)
557 }
558
559 pub fn validate_mount_source(
566 &self,
567 target_id: &str,
568 source: &Path,
569 executor: &impl CommandExecutor,
570 ) -> Result<Option<String>> {
571 let target = self
572 .config
573 .targets
574 .get(target_id)
575 .with_context(|| format!("unknown target template {target_id:?}"))?;
576 let exists = match target {
577 TargetTemplate::LocalPodman { .. }
578 | TargetTemplate::LocalDocker { .. }
579 | TargetTemplate::AppleContainer { .. }
580 | TargetTemplate::AwsEc2 { .. } => std::fs::metadata(source)
581 .map(|metadata| metadata.is_dir())
582 .or_else(|error| {
583 if error.kind() == std::io::ErrorKind::NotFound {
584 Ok(false)
585 } else {
586 Err(error)
587 }
588 })
589 .with_context(|| format!("inspect resource source {}", source.display()))?,
590 TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
591 targets::ssh_directory_exists(&SshTarget::from(ssh), source, executor)?
592 }
593 TargetTemplate::LocalBare | TargetTemplate::SshBare { .. } => {
594 bail!("resource attachments are unsupported for bare targets")
595 }
596 };
597 ensure!(
598 exists,
599 "source path {} does not exist or is not a directory",
600 source.display()
601 );
602 Ok(self.forced_read_only_reason(target, source, executor))
603 }
604
605 fn forced_read_only_reason(
607 &self,
608 target: &TargetTemplate,
609 source: &Path,
610 executor: &impl CommandExecutor,
611 ) -> Option<String> {
612 let ssh = match target {
613 TargetTemplate::LocalPodman { .. } | TargetTemplate::LocalDocker { .. } => None,
614 TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
615 Some(SshTarget::from(ssh))
616 }
617 _ => return None,
620 };
621 let filesystem = targets::probe_filesystem_types(
622 ssh.as_ref(),
623 std::slice::from_ref(&source.to_path_buf()),
624 executor,
625 )
626 .map_err(|error| {
627 tracing::debug!(
628 source = %source.display(),
629 error = format!("{error:#}"),
630 "could not probe the filesystem under a mount source"
631 );
632 })
633 .ok()?
634 .pop()?;
635 let reason = targets::overlay_unsupported_filesystem(&filesystem)?;
636 Some(format!("{filesystem} ({reason})"))
637 }
638
639 fn fail_new_session_with_cleanup(
640 &mut self,
641 session_id: &str,
642 error: anyhow::Error,
643 executor: &impl CommandExecutor,
644 ) -> Result<anyhow::Error> {
645 let original = provisioning::note_new_session_launch_failure(session_id, &error);
646 let cleanup_error = self
647 .cleanup_new_session_worktree_after_failure(session_id, executor)
648 .err()
649 .map(|cleanup_error| format!("{cleanup_error:#}"));
650 if let Some(cleanup_error) = &cleanup_error {
651 tracing::warn!(
652 session_id,
653 error = %cleanup_error,
654 "new-session worktree rollback reported a cleanup failure"
655 );
656 }
657 let failure = apply_failed_new_session_rollback(
658 &mut self.state,
659 session_id,
660 &original,
661 cleanup_error,
662 );
663 self.persist_session_state(session_id)?;
664 Ok(failure)
665 }
666
667 pub fn register_session_with_resources(
668 &mut self,
669 profile_id: &str,
670 bundle_id: &str,
671 target_id: &str,
672 title: impl Into<String>,
673 options: SessionLaunchOptions,
674 ) -> Result<String> {
675 let SessionLaunchOptions {
676 create_managed_worktree,
677 mjolnir_subagents,
678 initial_prompt,
679 workspace_id,
680 additional_mounts,
681 resource_allocation,
682 project_directory,
683 session_title_override,
684 } = options;
685 let session_title_override = match session_title_override {
686 Some(title) => {
687 Some(normalize_session_title(&title).context("session name cannot be empty")?)
688 }
689 None => None,
690 };
691 let profile = self
692 .config
693 .profiles
694 .get(profile_id)
695 .with_context(|| format!("unknown profile {profile_id:?}"))?;
696 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
697 let template = self
698 .config
699 .targets
700 .get(target_id)
701 .with_context(|| format!("unknown target template {target_id:?}"))?;
702 if create_managed_worktree == Some(true) && !is_bare_project_target(template) {
703 bail!("managed worktree creation requires a bare Git project");
704 }
705 if project_directory.is_some() != is_bare_project_target(template) {
706 bail!("raw project directories require a bare target, and bare targets require one");
707 }
708 if let Some(path) = &project_directory
709 && (!path.is_absolute()
710 || path
711 .components()
712 .any(|part| part == std::path::Component::ParentDir))
713 {
714 bail!("bare project directory must be an absolute safe path");
715 }
716 let bundle = project_directory
717 .is_none()
718 .then(|| self.config.bundles.get(bundle_id))
719 .flatten();
720 if project_directory.is_none() && bundle.is_none() {
721 bail!("unknown bundle {bundle_id:?}");
722 }
723 if profile.kind == mj_core::config::HarnessKind::Muse
724 && (!additional_mounts.is_empty()
725 || bundle.is_some_and(|bundle| bundle.repositories.len() > 1))
726 {
727 bail!(
728 "{} ACP supports one workspace root; use a single-repository bundle without attached directories",
729 profile.kind.display_name()
730 );
731 }
732 if let Some(bundle) = bundle {
733 for repository in &bundle.repositories {
734 mj_core::remote_git::resolve_repository(
735 repository,
736 &targets::CancellableProcessExecutor::with_timeout(
737 std::time::Duration::from_secs(15),
738 ),
739 )
740 .with_context(|| format!("repository {:?}", repository.id))?;
741 }
742 }
743 validate_resource_allocation(template, resource_allocation.as_ref())?;
744 let selected_container_size =
745 selected_host_container_size(template, resource_allocation.as_ref());
746 if !additional_mounts.is_empty() && mount_history_host(template).is_none() {
747 bail!("attached resources are unsupported for this target");
748 }
749 targets::validate_additional_mounts(&additional_mounts)?;
750 let id = new_session_id()?;
751 let now = now();
752 let record = SessionRecord {
753 build_cache: None,
754 create_managed_worktree,
755 mjolnir_subagents,
756 archived: false,
757 container_cpus: None,
758 container_memory: None,
759 container_workspace: Some(targets::new_container_workspace(&id)?),
764 id: id.clone(),
765 workspace_id,
766 title: title.into(),
767 harness_kind: profile.kind,
768 last_profile: profile_id.to_string(),
769 bundle_id: bundle_id.to_string(),
770 project_directory,
771 managed_worktree: None,
772 target_template_id: target_id.to_string(),
773 resource_allocation,
774 additional_mounts: additional_mounts.clone(),
775 state: SessionState::Provisioning,
776 target: None,
777 native_session_id: None,
778 acp_session_title: None,
779 session_title_override,
780 created_at: now.clone(),
781 updated_at: now,
782 viewed_through_event_ordinal: 0,
783 draft_input: initial_prompt.unwrap_or_default(),
784 last_error: None,
785 last_checkpoint_error: None,
786 checkpoint: None,
787 };
788 if let Some((host, size)) = selected_container_size.as_ref() {
792 crate::database::save_session_with_container_size(&record, host, *size)?;
793 } else {
794 crate::database::save_session(&record)?;
795 }
796 self.state.sessions.insert(id.clone(), record);
797 if let Some((host, size)) = selected_container_size {
798 self.state.remember_container_size(&host, size);
799 }
800 if let Some(host) = mount_history_host(template) {
801 match crate::database::remember_mount_sources(host, &additional_mounts) {
805 Ok(()) => self.state.remember_mount_sources(host, &additional_mounts),
806 Err(error) => tracing::warn!(
807 session_id = id,
808 error = format!("{error:#}"),
809 "could not remember the attached resource directories for later suggestions"
810 ),
811 }
812 }
813 Ok(id)
814 }
815
816 pub fn rename_session(&mut self, session_id: &str, title: &str) -> Result<String> {
817 let title = normalize_session_title(title).context("session name cannot be empty")?;
818 ensure!(
819 self.state.sessions.contains_key(session_id),
820 "unknown session {session_id}"
821 );
822 let updated_at = now();
823 crate::database::set_session_title_override(session_id, &title, &updated_at)?;
824 let record = self
825 .state
826 .sessions
827 .get_mut(session_id)
828 .expect("session was checked before updating its title");
829 record.session_title_override = Some(title.clone());
830 record.updated_at = updated_at;
831 Ok(title)
832 }
833
834 pub fn rename_profile_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
835 mj_core::config::validate_id("profile", new_id)?;
836 if old_id == new_id {
837 ensure!(
838 self.config.profiles.contains_key(old_id),
839 "unknown profile {old_id:?}"
840 );
841 return Ok(());
842 }
843 let journal = ConfigRenameJournal {
844 kind: ConfigRenameKind::Profile,
845 old_id: old_id.to_owned(),
846 new_id: new_id.to_owned(),
847 };
848 write_config_rename_journal(&journal)?;
849 let (config, ()) = match Config::update(|config| {
850 ensure!(
851 config.profiles.contains_key(old_id),
852 "unknown profile {old_id:?}"
853 );
854 ensure!(
855 !config.profiles.contains_key(new_id),
856 "profile {new_id:?} already exists"
857 );
858 let profile = config
859 .profiles
860 .remove(old_id)
861 .expect("profile was checked in the transaction");
862 config.profiles.insert(new_id.to_owned(), profile);
863 Ok(())
864 }) {
865 Ok(result) => result,
866 Err(error) => {
867 remove_config_rename_journal()
868 .context("remove profile rename journal after config save failed")?;
869 return Err(error).context("save renamed profile configuration");
870 }
871 };
872 self.config = config;
873 mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
874 if let Err(error) = crate::database::rename_profile_references(old_id, new_id) {
875 let restore = Config::update(|config| {
876 let profile = config
877 .profiles
878 .remove(new_id)
879 .with_context(|| format!("renamed profile {new_id:?} is missing"))?;
880 ensure!(
881 !config.profiles.contains_key(old_id),
882 "cannot restore profile rename: both {old_id:?} and {new_id:?} exist"
883 );
884 config.profiles.insert(old_id.to_owned(), profile);
885 Ok(())
886 });
887 let restored = match restore {
888 Ok((config, ())) => config,
889 Err(restore_error) => {
890 return Err(error).context(format!(
891 "rename profile references; additionally failed to restore config: {restore_error:#}"
892 ));
893 }
894 };
895 self.config = restored;
896 if let Err(restore_error) = remove_config_rename_journal() {
897 return Err(error).context(format!(
898 "rename profile references; additionally failed to remove rename journal: {restore_error:#}"
899 ));
900 }
901 return Err(error).context("rename profile references");
902 }
903 for session in self.state.sessions.values_mut() {
904 if session.last_profile == old_id {
905 session.last_profile = new_id.to_owned();
906 }
907 }
908 remove_config_rename_journal()?;
909 Ok(())
910 }
911
912 pub fn rename_target_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
913 mj_core::config::validate_id("target template", new_id)?;
914 if old_id == new_id {
915 ensure!(
916 self.config.targets.contains_key(old_id),
917 "unknown target {old_id:?}"
918 );
919 return Ok(());
920 }
921 let journal = ConfigRenameJournal {
922 kind: ConfigRenameKind::Target,
923 old_id: old_id.to_owned(),
924 new_id: new_id.to_owned(),
925 };
926 write_config_rename_journal(&journal)?;
927 let (config, ()) = match Config::update(|config| {
928 ensure!(
929 config.targets.contains_key(old_id),
930 "unknown target {old_id:?}"
931 );
932 ensure!(
933 !config.targets.contains_key(new_id),
934 "target {new_id:?} already exists"
935 );
936 let target = config
937 .targets
938 .remove(old_id)
939 .expect("target was checked in the transaction");
940 config.targets.insert(new_id.to_owned(), target);
941 Ok(())
942 }) {
943 Ok(result) => result,
944 Err(error) => {
945 remove_config_rename_journal()
946 .context("remove target rename journal after config save failed")?;
947 return Err(error).context("save renamed target configuration");
948 }
949 };
950 self.config = config;
951 mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
952 if let Err(error) = crate::database::rename_target_references(old_id, new_id) {
953 let restore = Config::update(|config| {
954 let target = config
955 .targets
956 .remove(new_id)
957 .with_context(|| format!("renamed target {new_id:?} is missing"))?;
958 ensure!(
959 !config.targets.contains_key(old_id),
960 "cannot restore target rename: both {old_id:?} and {new_id:?} exist"
961 );
962 config.targets.insert(old_id.to_owned(), target);
963 Ok(())
964 });
965 let restored = match restore {
966 Ok((config, ())) => config,
967 Err(restore_error) => {
968 return Err(error).context(format!(
969 "rename target references; additionally failed to restore config: {restore_error:#}"
970 ));
971 }
972 };
973 self.config = restored;
974 if let Err(restore_error) = remove_config_rename_journal() {
975 return Err(error).context(format!(
976 "rename target references; additionally failed to remove rename journal: {restore_error:#}"
977 ));
978 }
979 return Err(error).context("rename target references");
980 }
981 for session in self.state.sessions.values_mut() {
982 if session.target_template_id == old_id {
983 session.target_template_id = new_id.to_owned();
984 }
985 }
986 remove_config_rename_journal()?;
987 Ok(())
988 }
989
990 pub fn recover_config_id_rename() -> Result<bool> {
994 let path = config_rename_journal_path();
995 let body = match fs::read(&path) {
996 Ok(body) => body,
997 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(false),
998 Err(error) => return Err(error).context(format!("read {}", path.display())),
999 };
1000 let journal: ConfigRenameJournal =
1001 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1002 match journal.kind {
1003 ConfigRenameKind::Profile => {
1004 Config::update(|config| {
1005 finish_config_map_rename(
1006 &mut config.profiles,
1007 &journal.old_id,
1008 &journal.new_id,
1009 "profile",
1010 )?;
1011 Ok(())
1012 })?;
1013 crate::database::rename_profile_references(&journal.old_id, &journal.new_id)?;
1014 }
1015 ConfigRenameKind::Target => {
1016 Config::update(|config| {
1017 finish_config_map_rename(
1018 &mut config.targets,
1019 &journal.old_id,
1020 &journal.new_id,
1021 "target",
1022 )?;
1023 Ok(())
1024 })?;
1025 crate::database::rename_target_references(&journal.old_id, &journal.new_id)?;
1026 }
1027 }
1028 remove_config_rename_journal()?;
1029 Ok(true)
1030 }
1031
1032 pub fn update_session_container_settings(
1036 &mut self,
1037 session_id: &str,
1038 cpus: Option<String>,
1039 memory: Option<String>,
1040 additional_mounts: Vec<targets::AdditionalMount>,
1041 mount_history: Vec<std::path::PathBuf>,
1042 ) -> Result<()> {
1043 ensure!(
1044 self.state.sessions.contains_key(session_id),
1045 "unknown session {session_id}"
1046 );
1047 let cpus = cpus.filter(|value| !value.trim().is_empty());
1048 let memory = memory.filter(|value| !value.trim().is_empty());
1049 let updated_at = now();
1050 crate::database::set_session_container_settings(
1051 session_id,
1052 cpus.as_deref(),
1053 memory.as_deref(),
1054 &additional_mounts,
1055 &updated_at,
1056 )?;
1057 if let Some(host) = self
1058 .config
1059 .targets
1060 .get(
1061 &self.state.sessions[session_id]
1062 .target_template_id
1063 .to_owned(),
1064 )
1065 .and_then(mj_core::config::mount_history_host)
1066 {
1067 let host = host.to_owned();
1068 crate::database::replace_mount_history(&host, &mount_history)?;
1071 crate::database::remember_mount_sources(&host, &additional_mounts)?;
1072 self.state.mount_history.insert(host.clone(), mount_history);
1073 self.state.remember_mount_sources(&host, &additional_mounts);
1074 }
1075 let record = self
1076 .state
1077 .sessions
1078 .get_mut(session_id)
1079 .expect("session was checked before updating its container settings");
1080 record.container_cpus = cpus;
1081 record.container_memory = memory;
1082 record.additional_mounts = additional_mounts;
1083 record.updated_at = updated_at;
1084 Ok(())
1085 }
1086}
1087
1088fn config_rename_journal_path() -> PathBuf {
1089 data_dir().join(CONFIG_RENAME_JOURNAL)
1090}
1091
1092fn write_config_rename_journal(journal: &ConfigRenameJournal) -> Result<()> {
1093 let path = config_rename_journal_path();
1094 let body = serde_json::to_vec(journal).context("serialize config rename journal")?;
1095 atomic_write(&path, &body).with_context(|| format!("write {}", path.display()))
1096}
1097
1098fn remove_config_rename_journal() -> Result<()> {
1099 let path = config_rename_journal_path();
1100 match fs::remove_file(&path) {
1101 Ok(()) => Ok(()),
1102 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
1103 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
1104 }
1105}
1106
1107fn finish_config_map_rename<T>(
1108 entries: &mut BTreeMap<String, T>,
1109 old_id: &str,
1110 new_id: &str,
1111 kind: &str,
1112) -> Result<()> {
1113 if let Some(entry) = entries.remove(old_id) {
1114 ensure!(
1115 !entries.contains_key(new_id),
1116 "cannot recover {kind} rename: both {old_id:?} and {new_id:?} exist"
1117 );
1118 entries.insert(new_id.to_owned(), entry);
1119 } else {
1120 ensure!(
1121 entries.contains_key(new_id),
1122 "cannot recover {kind} rename: neither {old_id:?} nor {new_id:?} exists"
1123 );
1124 }
1125 Ok(())
1126}
1127
1128pub(crate) fn requires_private_profile_home(profile: &mj_core::config::HarnessProfile) -> bool {
1136 profile.codex_provider().ok().flatten().is_some()
1137}
1138
1139pub(crate) fn session_owns_profile_home(
1147 locator: &targets::TargetLocator,
1148 session_id: &str,
1149 profile: &mj_core::config::HarnessProfile,
1150) -> bool {
1151 removable_profile_root(locator, session_id, profile).is_some()
1152}
1153
1154fn claude_takes_a_private_home(
1160 profile: &mj_core::config::HarnessProfile,
1161 locator: &targets::TargetLocator,
1162) -> bool {
1163 profile.kind == mj_core::config::HarnessKind::Claude
1164 && profile
1165 .kind
1166 .scopes_home_with_environment(locator.harness_host())
1167}
1168
1169#[cfg(test)]
1176pub(crate) fn target_profile_home_for_test(
1177 locator: &targets::TargetLocator,
1178 session_id: &str,
1179 profile: &mj_core::config::HarnessProfile,
1180) -> String {
1181 target_profile_home(locator, session_id, profile)
1182}
1183
1184fn target_profile_home(
1185 locator: &targets::TargetLocator,
1186 session_id: &str,
1187 profile: &mj_core::config::HarnessProfile,
1188) -> String {
1189 let root = removable_profile_root(locator, session_id, profile)
1190 .unwrap_or_else(|| profile.home.to_string_lossy().into_owned());
1191 if profile.kind == mj_core::config::HarnessKind::Muse {
1192 PathBuf::from(root)
1193 .join("muse")
1194 .to_string_lossy()
1195 .into_owned()
1196 } else {
1197 root
1198 }
1199}
1200
1201pub(super) fn removable_profile_root(
1208 locator: &targets::TargetLocator,
1209 session_id: &str,
1210 profile: &mj_core::config::HarnessProfile,
1211) -> Option<String> {
1212 match locator {
1213 targets::TargetLocator::LocalBare { worker_root } => {
1214 if profile.kind == mj_core::config::HarnessKind::Muse {
1215 Some(
1216 mj_core::config::data_dir()
1217 .join("profiles")
1218 .join(session_id)
1219 .to_string_lossy()
1220 .into_owned(),
1221 )
1222 } else if claude_takes_a_private_home(profile, locator)
1223 || requires_private_profile_home(profile)
1224 {
1225 Some(
1226 Path::new(worker_root)
1227 .join("profile")
1228 .to_string_lossy()
1229 .into_owned(),
1230 )
1231 } else {
1232 None
1235 }
1236 }
1237 targets::TargetLocator::LocalPodman { .. }
1238 | targets::TargetLocator::LocalDocker { .. }
1239 | targets::TargetLocator::AppleContainer { .. }
1240 | targets::TargetLocator::SshPodman { .. }
1241 | targets::TargetLocator::SshDocker { .. } => {
1242 Some(format!("/var/lib/hel/profiles/{session_id}"))
1243 }
1244 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. } => {
1245 Some(format!(".local/share/hel/profiles/{session_id}"))
1246 }
1247 }
1248}
1249
1250pub fn resolve_target_input_path(
1252 target: &TargetTemplate,
1253 path: &Path,
1254 executor: &impl CommandExecutor,
1255) -> Result<PathBuf> {
1256 if !mj_core::path_input::needs_home(path)? {
1257 return Ok(path.to_path_buf());
1258 }
1259 let host = cache_host::CacheHost::for_path_target(target)?;
1260 mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1261}
1262
1263pub fn resolve_machine_input_path(
1266 machine: &mj_core::config::Machine,
1267 path: &Path,
1268 executor: &impl CommandExecutor,
1269) -> Result<PathBuf> {
1270 if !mj_core::path_input::needs_home(path)? {
1271 return Ok(path.to_path_buf());
1272 }
1273 let host = cache_host::CacheHost::for_path_machine(machine)?;
1276 mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1277}
1278
1279fn execute_checked(executor: &impl CommandExecutor, command: CommandSpec) -> Result<CommandOutput> {
1280 let output = executor.execute(&command)?;
1281 if output.status != 0 {
1282 let detail = command_error_detail(&output.stderr);
1283 if detail.is_empty() {
1284 bail!("{} failed with status {}", command.purpose, output.status);
1285 }
1286 bail!("{detail}");
1287 }
1288 Ok(output)
1289}
1290
1291fn command_error_detail(stderr: &[u8]) -> String {
1292 let reported = String::from_utf8_lossy(stderr);
1293 let reported = reported.trim();
1294 let detail = reported
1295 .rsplit_once("\nCaused by:\n")
1296 .map_or(reported, |(_, causes)| causes);
1297 let detail = detail.strip_prefix("Error: ").unwrap_or(detail);
1298 detail
1299 .lines()
1300 .map(|line| line.strip_prefix(" ").unwrap_or(line))
1301 .collect::<Vec<_>>()
1302 .join("\n")
1303 .trim()
1304 .to_owned()
1305}
1306
1307fn now() -> String {
1308 Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1309}
1310
1311fn restore_session_after_persistence_failure(
1312 state: &mut State,
1313 session_id: &str,
1314 previous: &SessionRecord,
1315 primary: anyhow::Error,
1316 persist: impl FnOnce(&SessionRecord) -> Result<()>,
1317) -> anyhow::Error {
1318 state
1319 .sessions
1320 .insert(session_id.to_owned(), previous.clone());
1321 let restored = state
1322 .sessions
1323 .get(session_id)
1324 .expect("restored session record disappeared");
1325 match persist(restored) {
1326 Ok(()) => primary,
1327 Err(error) => primary.context(format!(
1328 "restored prior session state in memory, but failed to persist the rollback: {error:#}"
1329 )),
1330 }
1331}
1332
1333fn persist_session_record_transition_or_restore(
1334 state: &mut State,
1335 session_id: &str,
1336 previous: &SessionRecord,
1337 context: &'static str,
1338 persist: &impl Fn(&SessionRecord) -> Result<()>,
1339) -> Result<()> {
1340 let result = persist(
1341 state
1342 .sessions
1343 .get(session_id)
1344 .expect("checkpoint session disappeared before persistence"),
1345 );
1346 match result {
1347 Ok(()) => Ok(()),
1348 Err(error) => Err(restore_session_after_persistence_failure(
1349 state,
1350 session_id,
1351 previous,
1352 error.context(context),
1353 persist,
1354 )),
1355 }
1356}
1357
1358pub fn config_only_controller(config: Config) -> Controller {
1359 Controller {
1360 config,
1361 state: State::default(),
1362 }
1363}
1364
1365#[cfg(test)]
1366mod tests;