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