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;
16pub(crate) mod publication;
17mod readiness;
18pub(crate) use readiness::NATIVE_SESSION_STARTUP_TIMEOUT;
19mod recovery_scan;
20mod resume;
21mod reviewer;
22mod subagents;
23#[cfg(test)]
24pub(crate) mod test_support;
25pub mod update;
26mod worker_binary;
27mod worker_restart;
28mod worktree;
29
30use std::collections::{BTreeMap, BTreeSet};
31use std::fs::{self, File, OpenOptions};
32use std::path::{Path, PathBuf};
33
34use anyhow::{Context, Result, bail, ensure};
35use chrono::Utc;
36
37use mj_core::config::{
38 Config, ProjectBundle, ProjectRepository, TargetTemplate, atomic_write, container_size_host,
39 data_dir, is_bare_project_target, mount_history_host,
40};
41
42use crate::import::{
43 RepositoryIdentity, bundle_matches, configured_bundle_for_local, configured_bundle_for_origin,
44 setup_style_id,
45};
46use crate::setup::github_repository_from_origin;
47
48const CONFIG_RENAME_JOURNAL: &str = "config-rename.json";
49
50#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
51#[serde(rename_all = "snake_case")]
52enum ConfigRenameKind {
53 Profile,
54 Target,
55}
56
57#[derive(Debug, serde::Serialize, serde::Deserialize)]
58#[serde(deny_unknown_fields)]
59struct ConfigRenameJournal {
60 kind: ConfigRenameKind,
61 old_id: String,
62 new_id: String,
63}
64use mj_core::state::{
65 HostContainerSize, SessionRecord, SessionResourceAllocation, SessionState, State,
66 new_session_id, normalize_session_title,
67};
68
69use crate::targets::{
70 self, AdditionalMount, CommandExecutor, CommandOutput, CommandSpec, SshTarget,
71};
72
73pub(crate) use backend::controller_github_token;
74pub use backend::image_refresh_plan;
75use backend::validate_resource_allocation;
76pub(crate) use backend::{LocalEngineReadiness, local_engine_readiness};
77pub use mbx::preview_build_cache;
78pub(crate) use mbx::{DoctorHostMbxStatus, MBX_VERSION, doctor_host_mbx};
79use provisioning::apply_failed_new_session_rollback;
80pub(crate) use worker_binary::refresh_target_worker_binary_if_stale;
81pub(crate) use worktree::path_exists_on_managed_target;
82
83pub use checkpoint::{
84 CheckpointArtifact, CheckpointDeferred, IdleWorkspaceLease, SessionExportLayout,
85 checkpoint_was_deferred, reconcile_managed_checkpoint_archives,
86};
87pub use lifecycle::{BranchDisposition, CheckoutDisposition, has_nothing_to_checkpoint};
88pub use recovery_scan::{RecoveryCandidate, RecoveryScan};
89pub use resume::{
90 ResumeRepositorySourceMismatch, ResumeRepositorySourcePreflight, ResumeRepositorySourceReceipt,
91 raw_conversion_preview_for,
92};
93pub use reviewer::reviewer_stager;
94pub use subagents::RegisterSubagentRequest;
95pub use worker_binary::{
96 WorkerBinaryAvailability, native_worker_binary_prerequisite, pin_worker_binary_sources,
97 worker_binary_prerequisite_for_arch,
98};
99pub use worker_restart::WorkerUpgradeOutcome;
100pub use worktree::{ResumePlan, local_project_repository, resume_compatibility};
101
102pub struct Controller {
103 pub config: Config,
104 pub state: State,
105}
106
107#[derive(Debug)]
111pub struct ControllerStoreGuard {
112 file: File,
113}
114
115impl ControllerStoreGuard {
116 pub fn acquire() -> Result<Self> {
117 let directory = data_dir();
118 Self::acquire_at(&directory)
119 }
120
121 fn acquire_at(directory: &Path) -> Result<Self> {
122 Self::try_acquire_at(directory)?.with_context(|| {
123 format!(
124 "another Mjolnir controller is already using {}; stop it before starting this command",
125 directory.display()
126 )
127 })
128 }
129
130 pub fn try_acquire() -> Result<Option<Self>> {
132 Self::try_acquire_at(&data_dir())
133 }
134
135 fn try_acquire_at(directory: &Path) -> Result<Option<Self>> {
136 std::fs::create_dir_all(directory)
137 .with_context(|| format!("create controller data directory {}", directory.display()))?;
138 let path = directory.join("controller.lock");
139 let mut options = OpenOptions::new();
140 options.create(true).read(true).write(true);
141 #[cfg(unix)]
142 {
143 use std::os::unix::fs::OpenOptionsExt;
144 options.mode(0o600);
145 }
146 let file = options
147 .open(&path)
148 .with_context(|| format!("open controller lock {}", path.display()))?;
149 match file.try_lock() {
150 Ok(()) => {}
151 Err(std::fs::TryLockError::WouldBlock) => return Ok(None),
152 Err(std::fs::TryLockError::Error(error)) => {
153 return Err(error)
154 .with_context(|| format!("lock controller store {}", directory.display()));
155 }
156 }
157 Ok(Some(Self { file }))
158 }
159
160 pub fn start_database_writer(&self) -> Result<crate::database::DatabaseWriterOwner> {
163 crate::database::start_database_writer()
164 }
165}
166
167impl Drop for ControllerStoreGuard {
168 fn drop(&mut self) {
169 let _ = self.file.unlock();
172 }
173}
174
175#[derive(Debug)]
179pub struct QuickBundleCreation {
180 pub config: Config,
181 pub bundle_id: String,
182}
183
184#[derive(Debug)]
188pub enum QuickBundleFailure {
189 InvalidSource(anyhow::Error),
190 Persistence(anyhow::Error),
191}
192
193impl std::fmt::Display for QuickBundleFailure {
194 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
195 match self {
196 Self::InvalidSource(error) => write!(formatter, "invalid repository source: {error}"),
197 Self::Persistence(error) => write!(formatter, "persist quick bundle: {error}"),
198 }
199 }
200}
201
202impl std::error::Error for QuickBundleFailure {
203 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
204 match self {
205 Self::InvalidSource(error) | Self::Persistence(error) => Some(error.root_cause()),
206 }
207 }
208}
209
210pub fn create_quick_bundle(
216 source: &str,
217) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
218 let (config, bundle_id) = Config::update(|config| {
219 create_quick_bundle_in_config(config, source)
220 .map_err(|error| anyhow::Error::new(QuickBundleFailure::InvalidSource(error)))
221 })
222 .map_err(|error| {
223 error
224 .downcast::<QuickBundleFailure>()
225 .unwrap_or_else(QuickBundleFailure::Persistence)
226 })?;
227 Ok(QuickBundleCreation { config, bundle_id })
228}
229
230pub fn create_quick_bundle_in_config(config: &mut Config, source: &str) -> Result<String> {
235 let source = interpret_repository_source(source)?;
236 let existing = match &source.kind {
237 RepositorySourceKind::Local(root) => configured_bundle_for_local(config, root),
238 RepositorySourceKind::Github(repository) => {
239 configured_bundle_for_origin(config, repository)
240 }
241 };
242 if let Some(existing) = existing {
243 return Ok(existing);
244 }
245 let repository_id = setup_style_id(&source.name);
246 let mut bundle_id = repository_id.clone();
247 for suffix in 2_u32.. {
248 if !config.bundles.contains_key(&bundle_id) {
249 break;
250 }
251 bundle_id = format!("{repository_id}-{suffix}");
252 }
253 config.bundles.insert(
254 bundle_id.clone(),
255 ProjectBundle {
256 primary_repo: repository_id.clone(),
257 repositories: vec![source.into_project_repository(repository_id.clone())],
258 },
259 );
260 config.validate()?;
261 Ok(bundle_id)
262}
263
264pub fn create_bundle_from_sources(
271 sources: &[String],
272) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
273 let (config, bundle_id) = Config::update(|config| {
274 create_bundle_from_sources_in_config(config, sources)
275 .map_err(|error| anyhow::Error::new(QuickBundleFailure::InvalidSource(error)))
276 })
277 .map_err(|error| {
278 error
279 .downcast::<QuickBundleFailure>()
280 .unwrap_or_else(QuickBundleFailure::Persistence)
281 })?;
282 Ok(QuickBundleCreation { config, bundle_id })
283}
284
285pub fn create_bundle_from_sources_in_config(
289 config: &mut Config,
290 sources: &[String],
291) -> Result<String> {
292 let sources = sources
293 .iter()
294 .map(|source| interpret_repository_source(source))
295 .collect::<Result<Vec<_>>>()?;
296 if sources.is_empty() {
297 bail!("at least one repository source is required");
298 }
299
300 let mut identities = BTreeSet::new();
301 for source in &sources {
302 if !identities.insert(source.identity()) {
303 bail!("duplicate repository source {:?}", source.display_name);
304 }
305 }
306
307 if let Some(existing) = exact_configured_bundle(config, &sources) {
308 return Ok(existing);
309 }
310
311 let mut updated = config.clone();
314 let mut used_repository_ids = BTreeSet::new();
315 let mut repositories = Vec::with_capacity(sources.len());
316 for source in sources {
317 let base = setup_style_id(&source.name);
318 let repository_id = unique_id(&base, |candidate| used_repository_ids.contains(candidate));
319 used_repository_ids.insert(repository_id.clone());
320 repositories.push(source.into_project_repository(repository_id));
321 }
322 let primary_repo = repositories
323 .first()
324 .map(|repository| repository.id.clone())
325 .context("at least one repository source is required")?;
326 let bundle_id = unique_id(&primary_repo, |candidate| {
327 updated.bundles.contains_key(candidate)
328 });
329 updated.bundles.insert(
330 bundle_id.clone(),
331 ProjectBundle {
332 primary_repo,
333 repositories,
334 },
335 );
336 updated.validate()?;
337 *config = updated;
338 Ok(bundle_id)
339}
340
341#[derive(Debug, Clone)]
342enum RepositorySourceKind {
343 Github(crate::setup::GithubRepository),
344 Local(PathBuf),
345}
346
347#[derive(Debug, Clone)]
348struct InterpretedRepositorySource {
349 display_name: String,
350 name: String,
351 kind: RepositorySourceKind,
352}
353
354impl InterpretedRepositorySource {
355 fn identity(&self) -> RepositoryIdentity {
356 match &self.kind {
357 RepositorySourceKind::Github(repository) => RepositoryIdentity::Github(
358 repository.owner.to_ascii_lowercase(),
359 repository.repository.to_ascii_lowercase(),
360 ),
361 RepositorySourceKind::Local(root) => RepositoryIdentity::Local(root.clone()),
362 }
363 }
364
365 fn into_project_repository(self, id: String) -> ProjectRepository {
366 let (github, local) = match self.kind {
367 RepositorySourceKind::Github(repository) => (
368 Some(format!("{}/{}", repository.owner, repository.repository)),
369 None,
370 ),
371 RepositorySourceKind::Local(root) => (None, Some(root)),
372 };
373 ProjectRepository {
374 id: id.clone(),
375 github,
376 local,
377 destination: PathBuf::from(id),
378 git_ref: None,
379 }
380 }
381}
382
383fn interpret_repository_source(source: &str) -> Result<InterpretedRepositorySource> {
386 let source = source.trim();
387 if source.is_empty() {
388 bail!("repository source cannot be empty");
389 }
390 let expanded = mj_core::path_input::expand_local(Path::new(source))?;
391 let candidate = expanded.as_path();
392 if candidate.exists() {
393 let root = mj_core::local_git::canonical_repository(candidate)?;
394 let name = root
395 .file_name()
396 .and_then(|name| name.to_str())
397 .context("local repository has no usable directory name")?
398 .to_owned();
399 return Ok(InterpretedRepositorySource {
400 display_name: source.to_owned(),
401 name,
402 kind: RepositorySourceKind::Local(root),
403 });
404 }
405 if candidate.is_absolute() || source.starts_with('.') || source.starts_with('~') {
406 bail!("local repository path {source:?} does not exist");
407 }
408 let repository = github_repository_from_origin(source).context(format!(
409 "{source:?} is not a GitHub owner/repository or URL"
410 ))?;
411 Ok(InterpretedRepositorySource {
412 display_name: source.to_owned(),
413 name: repository.repository.clone(),
414 kind: RepositorySourceKind::Github(repository),
415 })
416}
417
418fn exact_configured_bundle(
419 config: &Config,
420 requested: &[InterpretedRepositorySource],
421) -> Option<String> {
422 let requested_identities = requested
423 .iter()
424 .map(InterpretedRepositorySource::identity)
425 .collect::<BTreeSet<_>>();
426 let primary = requested.first()?.identity();
427 config.bundles.iter().find_map(|(id, bundle)| {
428 if bundle.repositories.len() != requested.len()
429 || bundle
430 .repositories
431 .iter()
432 .any(|repository| repository.git_ref.is_some())
433 {
434 return None;
435 }
436 bundle_matches(bundle, &requested_identities, &primary).then(|| id.clone())
437 })
438}
439
440fn unique_id(base: &str, mut is_used: impl FnMut(&str) -> bool) -> String {
441 if !is_used(base) {
442 return base.to_owned();
443 }
444 for suffix in 2_u32.. {
445 let suffix = format!("-{suffix}");
446 let prefix_len = 64usize.saturating_sub(suffix.len());
447 let prefix = base.chars().take(prefix_len).collect::<String>();
448 let candidate = format!("{prefix}{suffix}");
449 if !is_used(&candidate) {
450 return candidate;
451 }
452 }
453 unreachable!("u32 repository/bundle id suffixes exhausted")
454}
455
456pub struct SessionLaunchOptions {
457 pub create_managed_worktree: Option<bool>,
458 pub launch_base: Option<String>,
459 pub launch_branch: Option<String>,
460 pub mjolnir_subagents: Option<bool>,
461 pub initial_prompt: Option<String>,
462 pub workspace_id: String,
463 pub additional_mounts: Vec<AdditionalMount>,
464 pub resource_allocation: Option<SessionResourceAllocation>,
465 pub project_directory: Option<PathBuf>,
466 pub session_title_override: Option<String>,
467}
468
469pub struct SessionResumeOptions {
470 pub additional_mounts: Option<Vec<AdditionalMount>>,
471 pub resource_allocation: Option<SessionResourceAllocation>,
472 pub discard_queue: bool,
473}
474
475fn selected_host_container_size(
476 template: &TargetTemplate,
477 allocation: Option<&SessionResourceAllocation>,
478) -> Option<(String, HostContainerSize)> {
479 let host = container_size_host(template)?;
480 let SessionResourceAllocation::Container { cpus, memory_bytes } = allocation? else {
481 return None;
482 };
483 Some((
484 host.to_owned(),
485 HostContainerSize {
486 cpus: *cpus,
487 memory_bytes: *memory_bytes,
488 },
489 ))
490}
491
492impl Controller {
493 pub fn load() -> Result<Self> {
494 let config = Config::load()?;
495 let state = crate::database::load_state()?;
496 state.validate()?;
499 for session in state.sessions.values() {
500 if let Some(issue) = session.configuration_issue(&config) {
501 tracing::warn!(session_id = %session.id, "{issue}");
502 }
503 }
504 let mut controller = Self { config, state };
505 controller.backfill_target_runtime()?;
506 Ok(controller)
507 }
508
509 pub fn reload(&mut self) -> Result<()> {
510 self.backfill_target_runtime()?;
513 *self = Self::load()?;
514 Ok(())
515 }
516
517 fn backfill_target_runtime(&mut self) -> Result<()> {
518 for session in self.state.sessions.values_mut() {
519 if session.target_runtime.is_some() {
520 continue;
521 }
522 match session.target_runtime_settings(&self.config) {
523 Ok(runtime) => {
524 let runtime = runtime.into_owned();
525 session.target_runtime = crate::database::backfill_target_runtime(
526 &session.id,
527 &session.target_template_id,
528 &runtime,
529 )?;
530 }
531 Err(error) => {
532 tracing::warn!(session_id = %session.id, %error, "target access needs configuration repair")
533 }
534 }
535 }
536 Ok(())
537 }
538
539 fn persist_session_state(&self, session_id: &str) -> Result<()> {
540 match self.state.sessions.get(session_id) {
541 Some(session) => crate::database::save_lifecycle_session(session),
542 None => crate::database::delete_session(session_id),
543 }
544 }
545
546 fn persist_session_transition_or_restore(
547 &mut self,
548 session_id: &str,
549 previous: &SessionRecord,
550 context: &'static str,
551 ) -> Result<()> {
552 persist_session_record_transition_or_restore(
553 &mut self.state,
554 session_id,
555 previous,
556 context,
557 &crate::database::save_lifecycle_session,
558 )
559 }
560
561 fn restore_prior_session_after_persistence_failure(
562 &mut self,
563 session_id: &str,
564 previous: &SessionRecord,
565 primary: anyhow::Error,
566 ) -> anyhow::Error {
567 restore_session_after_persistence_failure(
568 &mut self.state,
569 session_id,
570 previous,
571 primary,
572 crate::database::save_lifecycle_session,
573 )
574 }
575
576 pub fn resolve_input_path(
578 &self,
579 target_id: &str,
580 path: &Path,
581 executor: &impl CommandExecutor,
582 ) -> Result<PathBuf> {
583 let target = self
584 .config
585 .targets
586 .get(target_id)
587 .context("Unknown path target")?;
588 resolve_target_input_path(target, path, executor)
589 }
590
591 pub fn validate_mount_source(
598 &self,
599 target_id: &str,
600 source: &Path,
601 executor: &impl CommandExecutor,
602 ) -> Result<Option<String>> {
603 let target = self
604 .config
605 .targets
606 .get(target_id)
607 .with_context(|| format!("unknown target template {target_id:?}"))?;
608 let exists = match target {
609 TargetTemplate::LocalPodman { .. }
610 | TargetTemplate::LocalDocker { .. }
611 | TargetTemplate::AppleContainer { .. }
612 | TargetTemplate::AwsEc2 { .. } => std::fs::metadata(source)
613 .map(|metadata| metadata.is_dir())
614 .or_else(|error| {
615 if error.kind() == std::io::ErrorKind::NotFound {
616 Ok(false)
617 } else {
618 Err(error)
619 }
620 })
621 .with_context(|| format!("inspect resource source {}", source.display()))?,
622 TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
623 targets::ssh_directory_exists(&SshTarget::from(ssh), source, executor)?
624 }
625 TargetTemplate::LocalBare | TargetTemplate::SshBare { .. } => {
626 bail!("resource attachments are unsupported for bare targets")
627 }
628 };
629 ensure!(
630 exists,
631 "source path {} does not exist or is not a directory",
632 source.display()
633 );
634 Ok(self.forced_read_only_reason(target, source, executor))
635 }
636
637 fn forced_read_only_reason(
639 &self,
640 target: &TargetTemplate,
641 source: &Path,
642 executor: &impl CommandExecutor,
643 ) -> Option<String> {
644 let ssh = match target {
645 TargetTemplate::LocalPodman { .. } | TargetTemplate::LocalDocker { .. } => None,
646 TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
647 Some(SshTarget::from(ssh))
648 }
649 _ => return None,
652 };
653 let filesystem = targets::probe_filesystem_types(
654 ssh.as_ref(),
655 std::slice::from_ref(&source.to_path_buf()),
656 executor,
657 )
658 .map_err(|error| {
659 tracing::debug!(
660 source = %source.display(),
661 error = format!("{error:#}"),
662 "could not probe the filesystem under a mount source"
663 );
664 })
665 .ok()?
666 .pop()?;
667 let reason = targets::overlay_unsupported_filesystem(&filesystem)?;
668 Some(format!("{filesystem} ({reason})"))
669 }
670
671 fn fail_new_session_with_cleanup(
672 &mut self,
673 session_id: &str,
674 error: anyhow::Error,
675 executor: &impl CommandExecutor,
676 ) -> Result<anyhow::Error> {
677 let original = provisioning::note_new_session_launch_failure(session_id, &error);
678 let cleanup_error = self
679 .cleanup_new_session_worktree_after_failure(session_id, executor)
680 .err()
681 .map(|cleanup_error| format!("{cleanup_error:#}"));
682 if let Some(cleanup_error) = &cleanup_error {
683 tracing::warn!(
684 session_id,
685 error = %cleanup_error,
686 "new-session worktree rollback reported a cleanup failure"
687 );
688 }
689 let failure = apply_failed_new_session_rollback(
690 &mut self.state,
691 session_id,
692 &original,
693 cleanup_error,
694 );
695 self.persist_session_state(session_id)?;
696 Ok(failure)
697 }
698
699 pub fn register_session_with_resources(
700 &mut self,
701 profile_id: &str,
702 bundle_id: &str,
703 target_id: &str,
704 title: impl Into<String>,
705 options: SessionLaunchOptions,
706 ) -> Result<String> {
707 let SessionLaunchOptions {
708 create_managed_worktree,
709 launch_base,
710 launch_branch,
711 mjolnir_subagents,
712 initial_prompt,
713 workspace_id,
714 additional_mounts,
715 resource_allocation,
716 project_directory,
717 session_title_override,
718 } = options;
719 let launch_base = match launch_base {
720 Some(base) => {
721 let base = base.trim();
722 if base.is_empty() {
723 bail!("launch base must not be empty");
724 }
725 Some(base.to_owned())
726 }
727 None => None,
728 };
729 let launch_branch = match launch_branch {
730 Some(branch) => {
731 let branch = branch.trim();
732 ensure!(!branch.is_empty(), "launch branch must not be empty");
733 Some(branch.to_owned())
734 }
735 None => None,
736 };
737 if launch_base.is_some() && create_managed_worktree == Some(false) {
738 bail!("a launch base requires a managed worktree or a bundle session");
739 }
740 let session_title_override = match session_title_override {
741 Some(title) => {
742 Some(normalize_session_title(&title).context("session name cannot be empty")?)
743 }
744 None => None,
745 };
746 let profile = self
747 .config
748 .profiles
749 .get(profile_id)
750 .with_context(|| format!("unknown profile {profile_id:?}"))?;
751 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
752 let template = self
753 .config
754 .targets
755 .get(target_id)
756 .with_context(|| format!("unknown target template {target_id:?}"))?;
757 if create_managed_worktree == Some(true) && !is_bare_project_target(template) {
758 bail!("managed worktree creation requires a bare Git project");
759 }
760 if project_directory.is_some() != is_bare_project_target(template) {
761 bail!("raw project directories require a bare target, and bare targets require one");
762 }
763 if let Some(path) = &project_directory
764 && (!path.is_absolute()
765 || path
766 .components()
767 .any(|part| part == std::path::Component::ParentDir))
768 {
769 bail!("bare project directory must be an absolute safe path");
770 }
771 let bundle = project_directory
772 .is_none()
773 .then(|| self.config.bundles.get(bundle_id))
774 .flatten();
775 if project_directory.is_none() && bundle.is_none() {
776 bail!("unknown bundle {bundle_id:?}");
777 }
778 if profile.kind == mj_core::config::HarnessKind::Muse
779 && (!additional_mounts.is_empty()
780 || bundle.is_some_and(|bundle| bundle.repositories.len() > 1))
781 {
782 bail!(
783 "{} ACP supports one workspace root; use a single-repository bundle without attached directories",
784 profile.kind.display_name()
785 );
786 }
787 if let Some(bundle) = bundle {
788 for repository in &bundle.repositories {
789 mj_core::remote_git::resolve_repository(
790 repository,
791 &targets::CancellableProcessExecutor::with_timeout(
792 std::time::Duration::from_secs(15),
793 ),
794 )
795 .with_context(|| format!("repository {:?}", repository.id))?;
796 }
797 }
798 validate_resource_allocation(template, resource_allocation.as_ref())?;
799 let selected_container_size =
800 selected_host_container_size(template, resource_allocation.as_ref());
801 if !additional_mounts.is_empty() && mount_history_host(template).is_none() {
802 bail!("attached resources are unsupported for this target");
803 }
804 targets::validate_additional_mounts(&additional_mounts)?;
805 let id = new_session_id()?;
806 let now = now();
807 let record = SessionRecord {
808 target_runtime: Some(template.into()),
809 build_cache: None,
810 create_managed_worktree,
811 launch_base,
812 launch_branch,
813 publication: None,
814 mjolnir_subagents,
815 archived: false,
816 container_cpus: None,
817 container_memory: None,
818 container_workspace: Some(targets::new_container_workspace(&id)?),
823 id: id.clone(),
824 workspace_id,
825 title: title.into(),
826 harness_kind: profile.kind,
827 last_profile: profile_id.to_string(),
828 bundle_id: bundle_id.to_string(),
829 project_directory,
830 managed_worktree: None,
831 target_template_id: target_id.to_string(),
832 resource_allocation,
833 additional_mounts: additional_mounts.clone(),
834 state: SessionState::Provisioning,
835 target: None,
836 native_session_id: None,
837 acp_session_title: None,
838 session_title_override,
839 created_at: now.clone(),
840 updated_at: now,
841 viewed_through_event_ordinal: 0,
842 draft_input: initial_prompt.unwrap_or_default(),
843 last_error: None,
844 last_checkpoint_error: None,
845 checkpoint: None,
846 };
847 if let Some((host, size)) = selected_container_size.as_ref() {
851 crate::database::save_session_with_container_size(&record, host, *size)?;
852 } else {
853 crate::database::save_session(&record)?;
854 }
855 self.state.sessions.insert(id.clone(), record);
856 if let Some((host, size)) = selected_container_size {
857 self.state.remember_container_size(&host, size);
858 }
859 if let Some(host) = mount_history_host(template) {
860 match crate::database::remember_mount_sources(host, &additional_mounts) {
864 Ok(()) => self.state.remember_mount_sources(host, &additional_mounts),
865 Err(error) => tracing::warn!(
866 session_id = id,
867 error = format!("{error:#}"),
868 "could not remember the attached resource directories for later suggestions"
869 ),
870 }
871 }
872 Ok(id)
873 }
874
875 pub fn rename_session(&mut self, session_id: &str, title: &str) -> Result<String> {
876 let title = normalize_session_title(title).context("session name cannot be empty")?;
877 ensure!(
878 self.state.sessions.contains_key(session_id),
879 "unknown session {session_id}"
880 );
881 let updated_at = now();
882 crate::database::set_session_title_override(session_id, &title, &updated_at)?;
883 let record = self
884 .state
885 .sessions
886 .get_mut(session_id)
887 .expect("session was checked before updating its title");
888 record.session_title_override = Some(title.clone());
889 record.updated_at = updated_at;
890 Ok(title)
891 }
892
893 pub fn rename_profile_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
894 mj_core::config::validate_id("profile", new_id)?;
895 if old_id == new_id {
896 ensure!(
897 self.config.profiles.contains_key(old_id),
898 "unknown profile {old_id:?}"
899 );
900 return Ok(());
901 }
902 let journal = ConfigRenameJournal {
903 kind: ConfigRenameKind::Profile,
904 old_id: old_id.to_owned(),
905 new_id: new_id.to_owned(),
906 };
907 write_config_rename_journal(&journal)?;
908 let (config, ()) = match Config::update(|config| {
909 ensure!(
910 config.profiles.contains_key(old_id),
911 "unknown profile {old_id:?}"
912 );
913 ensure!(
914 !config.profiles.contains_key(new_id),
915 "profile {new_id:?} already exists"
916 );
917 let profile = config
918 .profiles
919 .remove(old_id)
920 .expect("profile was checked in the transaction");
921 config.profiles.insert(new_id.to_owned(), profile);
922 Ok(())
923 }) {
924 Ok(result) => result,
925 Err(error) => {
926 remove_config_rename_journal()
927 .context("remove profile rename journal after config save failed")?;
928 return Err(error).context("save renamed profile configuration");
929 }
930 };
931 self.config = config;
932 mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
933 if let Err(error) = crate::database::rename_profile_references(old_id, new_id) {
934 let restore = Config::update(|config| {
935 let profile = config
936 .profiles
937 .remove(new_id)
938 .with_context(|| format!("renamed profile {new_id:?} is missing"))?;
939 ensure!(
940 !config.profiles.contains_key(old_id),
941 "cannot restore profile rename: both {old_id:?} and {new_id:?} exist"
942 );
943 config.profiles.insert(old_id.to_owned(), profile);
944 Ok(())
945 });
946 let restored = match restore {
947 Ok((config, ())) => config,
948 Err(restore_error) => {
949 return Err(error).context(format!(
950 "rename profile references; additionally failed to restore config: {restore_error:#}"
951 ));
952 }
953 };
954 self.config = restored;
955 if let Err(restore_error) = remove_config_rename_journal() {
956 return Err(error).context(format!(
957 "rename profile references; additionally failed to remove rename journal: {restore_error:#}"
958 ));
959 }
960 return Err(error).context("rename profile references");
961 }
962 for session in self.state.sessions.values_mut() {
963 if session.last_profile == old_id {
964 session.last_profile = new_id.to_owned();
965 }
966 }
967 remove_config_rename_journal()?;
968 Ok(())
969 }
970
971 pub fn rename_target_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
972 mj_core::config::validate_id("target template", new_id)?;
973 if old_id == new_id {
974 ensure!(
975 self.config.targets.contains_key(old_id),
976 "unknown target {old_id:?}"
977 );
978 return Ok(());
979 }
980 let journal = ConfigRenameJournal {
981 kind: ConfigRenameKind::Target,
982 old_id: old_id.to_owned(),
983 new_id: new_id.to_owned(),
984 };
985 write_config_rename_journal(&journal)?;
986 let (config, ()) = match Config::update(|config| {
987 ensure!(
988 config.targets.contains_key(old_id),
989 "unknown target {old_id:?}"
990 );
991 ensure!(
992 !config.targets.contains_key(new_id),
993 "target {new_id:?} already exists"
994 );
995 let target = config
996 .targets
997 .remove(old_id)
998 .expect("target was checked in the transaction");
999 config.targets.insert(new_id.to_owned(), target);
1000 Ok(())
1001 }) {
1002 Ok(result) => result,
1003 Err(error) => {
1004 remove_config_rename_journal()
1005 .context("remove target rename journal after config save failed")?;
1006 return Err(error).context("save renamed target configuration");
1007 }
1008 };
1009 self.config = config;
1010 mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
1011 if let Err(error) = crate::database::rename_target_references(old_id, new_id) {
1012 let restore = Config::update(|config| {
1013 let target = config
1014 .targets
1015 .remove(new_id)
1016 .with_context(|| format!("renamed target {new_id:?} is missing"))?;
1017 ensure!(
1018 !config.targets.contains_key(old_id),
1019 "cannot restore target rename: both {old_id:?} and {new_id:?} exist"
1020 );
1021 config.targets.insert(old_id.to_owned(), target);
1022 Ok(())
1023 });
1024 let restored = match restore {
1025 Ok((config, ())) => config,
1026 Err(restore_error) => {
1027 return Err(error).context(format!(
1028 "rename target references; additionally failed to restore config: {restore_error:#}"
1029 ));
1030 }
1031 };
1032 self.config = restored;
1033 if let Err(restore_error) = remove_config_rename_journal() {
1034 return Err(error).context(format!(
1035 "rename target references; additionally failed to remove rename journal: {restore_error:#}"
1036 ));
1037 }
1038 return Err(error).context("rename target references");
1039 }
1040 for session in self.state.sessions.values_mut() {
1041 if session.target_template_id == old_id {
1042 session.target_template_id = new_id.to_owned();
1043 }
1044 }
1045 remove_config_rename_journal()?;
1046 Ok(())
1047 }
1048
1049 pub fn recover_config_id_rename() -> Result<bool> {
1053 let path = config_rename_journal_path();
1054 let body = match fs::read(&path) {
1055 Ok(body) => body,
1056 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(false),
1057 Err(error) => return Err(error).context(format!("read {}", path.display())),
1058 };
1059 let journal: ConfigRenameJournal =
1060 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1061 match journal.kind {
1062 ConfigRenameKind::Profile => {
1063 Config::update(|config| {
1064 finish_config_map_rename(
1065 &mut config.profiles,
1066 &journal.old_id,
1067 &journal.new_id,
1068 "profile",
1069 )?;
1070 Ok(())
1071 })?;
1072 crate::database::rename_profile_references(&journal.old_id, &journal.new_id)?;
1073 }
1074 ConfigRenameKind::Target => {
1075 Config::update_to(&mj_core::config::config_path(), |config| {
1080 finish_config_map_rename(
1081 &mut config.targets,
1082 &journal.old_id,
1083 &journal.new_id,
1084 "target",
1085 )?;
1086 Ok(())
1087 })?;
1088 crate::database::rename_target_references(&journal.old_id, &journal.new_id)?;
1089 }
1090 }
1091 remove_config_rename_journal()?;
1092 Ok(true)
1093 }
1094
1095 pub fn update_session_container_settings(
1099 &mut self,
1100 session_id: &str,
1101 cpus: Option<String>,
1102 memory: Option<String>,
1103 additional_mounts: Vec<targets::AdditionalMount>,
1104 mount_history: Vec<std::path::PathBuf>,
1105 executor: &impl CommandExecutor,
1106 ) -> Result<()> {
1107 let session = self
1108 .state
1109 .sessions
1110 .get(session_id)
1111 .with_context(|| format!("unknown session {session_id}"))?;
1112 for mount in &additional_mounts {
1119 if session
1120 .additional_mounts
1121 .iter()
1122 .any(|attached| attached.source == mount.source)
1123 {
1124 continue;
1125 }
1126 self.validate_mount_source(&session.target_template_id, &mount.source, executor)?;
1127 }
1128 let cpus = cpus.filter(|value| !value.trim().is_empty());
1129 let memory = memory.filter(|value| !value.trim().is_empty());
1130 let updated_at = now();
1131 crate::database::set_session_container_settings(
1132 session_id,
1133 cpus.as_deref(),
1134 memory.as_deref(),
1135 &additional_mounts,
1136 &updated_at,
1137 )?;
1138 if let Some(host) = self
1139 .config
1140 .targets
1141 .get(
1142 &self.state.sessions[session_id]
1143 .target_template_id
1144 .to_owned(),
1145 )
1146 .and_then(mj_core::config::mount_history_host)
1147 {
1148 let host = host.to_owned();
1149 crate::database::replace_mount_history(&host, &mount_history)?;
1152 crate::database::remember_mount_sources(&host, &additional_mounts)?;
1153 self.state.mount_history.insert(host.clone(), mount_history);
1154 self.state.remember_mount_sources(&host, &additional_mounts);
1155 }
1156 let record = self
1157 .state
1158 .sessions
1159 .get_mut(session_id)
1160 .expect("session was checked before updating its container settings");
1161 record.container_cpus = cpus;
1162 record.container_memory = memory;
1163 record.additional_mounts = additional_mounts;
1164 record.updated_at = updated_at;
1165 Ok(())
1166 }
1167}
1168
1169fn config_rename_journal_path() -> PathBuf {
1170 data_dir().join(CONFIG_RENAME_JOURNAL)
1171}
1172
1173fn write_config_rename_journal(journal: &ConfigRenameJournal) -> Result<()> {
1174 let path = config_rename_journal_path();
1175 let body = serde_json::to_vec(journal).context("serialize config rename journal")?;
1176 atomic_write(&path, &body).with_context(|| format!("write {}", path.display()))
1177}
1178
1179fn remove_config_rename_journal() -> Result<()> {
1180 let path = config_rename_journal_path();
1181 match fs::remove_file(&path) {
1182 Ok(()) => Ok(()),
1183 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
1184 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
1185 }
1186}
1187
1188fn finish_config_map_rename<T>(
1189 entries: &mut BTreeMap<String, T>,
1190 old_id: &str,
1191 new_id: &str,
1192 kind: &str,
1193) -> Result<()> {
1194 if let Some(entry) = entries.remove(old_id) {
1195 ensure!(
1196 !entries.contains_key(new_id),
1197 "cannot recover {kind} rename: both {old_id:?} and {new_id:?} exist"
1198 );
1199 entries.insert(new_id.to_owned(), entry);
1200 } else {
1201 ensure!(
1202 entries.contains_key(new_id),
1203 "cannot recover {kind} rename: neither {old_id:?} nor {new_id:?} exists"
1204 );
1205 }
1206 Ok(())
1207}
1208
1209pub(crate) fn requires_private_profile_home(profile: &mj_core::config::HarnessProfile) -> bool {
1217 profile.codex_provider().ok().flatten().is_some()
1218}
1219
1220pub(crate) fn session_owns_profile_home(
1228 locator: &targets::TargetLocator,
1229 session_id: &str,
1230 profile: &mj_core::config::HarnessProfile,
1231) -> bool {
1232 removable_profile_root(locator, session_id, profile).is_some()
1233}
1234
1235fn claude_takes_a_private_home(
1241 profile: &mj_core::config::HarnessProfile,
1242 locator: &targets::TargetLocator,
1243) -> bool {
1244 profile.kind == mj_core::config::HarnessKind::Claude
1245 && profile
1246 .kind
1247 .scopes_home_with_environment(locator.harness_host())
1248}
1249
1250#[cfg(test)]
1257pub(crate) fn target_profile_home_for_test(
1258 locator: &targets::TargetLocator,
1259 session_id: &str,
1260 profile: &mj_core::config::HarnessProfile,
1261) -> String {
1262 target_profile_home(locator, session_id, profile)
1263}
1264
1265fn target_profile_home(
1266 locator: &targets::TargetLocator,
1267 session_id: &str,
1268 profile: &mj_core::config::HarnessProfile,
1269) -> String {
1270 let root = removable_profile_root(locator, session_id, profile)
1271 .unwrap_or_else(|| profile.home.to_string_lossy().into_owned());
1272 if profile.kind == mj_core::config::HarnessKind::Muse {
1273 PathBuf::from(root)
1274 .join("muse")
1275 .to_string_lossy()
1276 .into_owned()
1277 } else {
1278 root
1279 }
1280}
1281
1282pub(super) fn removable_profile_root(
1289 locator: &targets::TargetLocator,
1290 session_id: &str,
1291 profile: &mj_core::config::HarnessProfile,
1292) -> Option<String> {
1293 match locator {
1294 targets::TargetLocator::LocalBare { worker_root } => {
1295 if profile.kind == mj_core::config::HarnessKind::Muse {
1296 Some(
1297 mj_core::config::data_dir()
1298 .join("profiles")
1299 .join(session_id)
1300 .to_string_lossy()
1301 .into_owned(),
1302 )
1303 } else if claude_takes_a_private_home(profile, locator)
1304 || requires_private_profile_home(profile)
1305 {
1306 Some(
1307 Path::new(worker_root)
1308 .join("profile")
1309 .to_string_lossy()
1310 .into_owned(),
1311 )
1312 } else {
1313 None
1316 }
1317 }
1318 targets::TargetLocator::LocalPodman { .. }
1319 | targets::TargetLocator::LocalDocker { .. }
1320 | targets::TargetLocator::AppleContainer { .. }
1321 | targets::TargetLocator::SshPodman { .. }
1322 | targets::TargetLocator::SshDocker { .. } => {
1323 Some(format!("/var/lib/hel/profiles/{session_id}"))
1324 }
1325 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. } => {
1326 Some(format!(".local/share/hel/profiles/{session_id}"))
1327 }
1328 }
1329}
1330
1331pub fn resolve_target_input_path(
1333 target: &TargetTemplate,
1334 path: &Path,
1335 executor: &impl CommandExecutor,
1336) -> Result<PathBuf> {
1337 if !mj_core::path_input::needs_home(path)? {
1338 return Ok(path.to_path_buf());
1339 }
1340 let host = cache_host::CacheHost::for_path_target(target)?;
1341 mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1342}
1343
1344pub fn resolve_machine_input_path(
1347 machine: &mj_core::config::Machine,
1348 path: &Path,
1349 executor: &impl CommandExecutor,
1350) -> Result<PathBuf> {
1351 if !mj_core::path_input::needs_home(path)? {
1352 return Ok(path.to_path_buf());
1353 }
1354 let host = cache_host::CacheHost::for_path_machine(machine)?;
1357 mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1358}
1359
1360fn execute_checked(executor: &impl CommandExecutor, command: CommandSpec) -> Result<CommandOutput> {
1361 let output = executor.execute(&command)?;
1362 if output.status != 0 {
1363 let detail = command_error_detail(&output.stderr);
1364 if detail.is_empty() {
1365 bail!("{} failed with status {}", command.purpose, output.status);
1366 }
1367 bail!("{detail}");
1368 }
1369 Ok(output)
1370}
1371
1372fn command_error_detail(stderr: &[u8]) -> String {
1373 let reported = String::from_utf8_lossy(stderr);
1374 let reported = reported.trim();
1375 let detail = reported
1376 .rsplit_once("\nCaused by:\n")
1377 .map_or(reported, |(_, causes)| causes);
1378 let detail = detail.strip_prefix("Error: ").unwrap_or(detail);
1379 detail
1380 .lines()
1381 .map(|line| line.strip_prefix(" ").unwrap_or(line))
1382 .collect::<Vec<_>>()
1383 .join("\n")
1384 .trim()
1385 .to_owned()
1386}
1387
1388fn now() -> String {
1389 Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1390}
1391
1392fn restore_session_after_persistence_failure(
1393 state: &mut State,
1394 session_id: &str,
1395 previous: &SessionRecord,
1396 primary: anyhow::Error,
1397 persist: impl FnOnce(&SessionRecord) -> Result<()>,
1398) -> anyhow::Error {
1399 state
1400 .sessions
1401 .insert(session_id.to_owned(), previous.clone());
1402 let restored = state
1403 .sessions
1404 .get(session_id)
1405 .expect("restored session record disappeared");
1406 match persist(restored) {
1407 Ok(()) => primary,
1408 Err(error) => primary.context(format!(
1409 "restored prior session state in memory, but failed to persist the rollback: {error:#}"
1410 )),
1411 }
1412}
1413
1414fn persist_session_record_transition_or_restore(
1415 state: &mut State,
1416 session_id: &str,
1417 previous: &SessionRecord,
1418 context: &'static str,
1419 persist: &impl Fn(&SessionRecord) -> Result<()>,
1420) -> Result<()> {
1421 let result = persist(
1422 state
1423 .sessions
1424 .get(session_id)
1425 .expect("checkpoint session disappeared before persistence"),
1426 );
1427 match result {
1428 Ok(()) => Ok(()),
1429 Err(error) => Err(restore_session_after_persistence_failure(
1430 state,
1431 session_id,
1432 previous,
1433 error.context(context),
1434 persist,
1435 )),
1436 }
1437}
1438
1439pub fn config_only_controller(config: Config) -> Controller {
1440 Controller {
1441 config,
1442 state: State::default(),
1443 }
1444}
1445
1446#[cfg(test)]
1447mod tests;