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