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_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#[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 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 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 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 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 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 _ => 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 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 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 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 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 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 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
1144pub(crate) fn requires_private_profile_home(profile: &mj_core::config::HarnessProfile) -> bool {
1152 profile.codex_provider().ok().flatten().is_some()
1153}
1154
1155pub(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
1170fn 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#[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
1217pub(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 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
1266pub 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
1279pub 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 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;