1mod backend;
4mod cache_host;
5pub(crate) mod checkpoint;
6mod git_cache;
7mod github_app;
8mod lifecycle;
9pub mod local_profile_homes;
10pub(crate) mod mbx;
11pub mod move_session;
12mod network_git;
13mod new_session_preflight;
14pub use new_session_preflight::{NewSessionPreflight, NewSessionRepository};
15mod path_completion;
16pub mod profile_config;
17mod provisioning;
18pub(crate) mod publication;
19mod readiness;
20pub(crate) use network_git::workspace_root;
21pub(crate) use readiness::HarnessPreparationFailure;
22pub(crate) use readiness::NATIVE_SESSION_STARTUP_TIMEOUT;
23mod recovery_scan;
24mod resume;
25mod reviewer;
26mod subagent_park;
27pub use subagent_park::ParkOutcome;
28pub(crate) use subagent_park::{FAILED_STARTUP_CLEANUP_TIMEOUT, failed_launch_cleanup_executor};
29mod subagents;
30#[cfg(test)]
31pub(crate) mod test_support;
32pub mod update;
33mod worker_binary;
34mod worker_restart;
35mod worktree;
36pub(crate) use worktree::RemoteGitExecutor;
37
38use std::collections::{BTreeMap, BTreeSet};
39use std::fs::{self, File, OpenOptions};
40use std::path::{Path, PathBuf};
41
42use anyhow::{Context, Result, bail, ensure};
43use chrono::Utc;
44
45use mj_core::config::{
46 Config, ProjectBundle, ProjectRepository, TargetTemplate, atomic_write, container_size_host,
47 data_dir, is_bare_project_target, mount_history_host,
48};
49
50use crate::import::{RepositoryIdentity, setup_style_id};
51use crate::setup::github_repository_from_origin;
52
53const CONFIG_RENAME_JOURNAL: &str = "config-rename.json";
54
55#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
56#[serde(rename_all = "snake_case")]
57enum ConfigRenameKind {
58 Profile,
59 Target,
60}
61
62#[derive(Debug, serde::Serialize, serde::Deserialize)]
63#[serde(deny_unknown_fields)]
64struct ConfigRenameJournal {
65 kind: ConfigRenameKind,
66 old_id: String,
67 new_id: String,
68}
69use mj_core::state::{
70 HostContainerSize, SessionRecord, SessionResourceAllocation, SessionState, State,
71 new_session_id, normalize_session_title,
72};
73
74use crate::targets::{
75 self, AdditionalMount, CommandExecutor, CommandOutput, CommandSpec, SshTarget,
76};
77
78pub(crate) use backend::backend_locator;
79pub(crate) use backend::controller_github_token;
80pub use backend::image_refresh_plan;
81pub(crate) use backend::validate_resource_allocation;
82pub(crate) use backend::{LocalEngineReadiness, local_engine_readiness};
83pub use github_app::GithubBundleSelectionError;
84pub(crate) use github_app::{
85 GithubAppTokenProvider, github_app_token_for_session, github_token_for_session,
86};
87pub use mbx::install::{
88 MbxInstallKind, MbxInstallResult, install_kind as mbx_install_kind, install_mbx,
89};
90pub use mbx::preview_build_cache;
91pub(crate) use mbx::release::{FAILURE_REPORT_SECS, ReleaseFailure, recent_release_failures};
92pub(crate) use mbx::{DoctorHostMbxStatus, MBX_VERSION, doctor_host_mbx};
93use provisioning::apply_failed_new_session_rollback;
94pub(crate) use worker_binary::{
95 prepare_recovery_worker_binary, recorded_exit_reason, refresh_target_worker_binary_if_stale,
96 worker_source_problem,
97};
98pub(crate) use worktree::path_exists_on_managed_target;
99
100pub use checkpoint::{
101 CheckpointArtifact, CheckpointDeferred, IdleWorkspaceLease, SessionExportLayout,
102 checkpoint_was_deferred, reconcile_managed_checkpoint_archives,
103 sweep_local_checkpoint_leftovers,
104};
105pub use lifecycle::{
106 BeforeClose, BranchDisposition, CheckoutDisposition, has_nothing_to_checkpoint,
107};
108pub use recovery_scan::{RecoveryCandidate, RecoveryScan};
109pub(crate) use resume::InPlaceRestartError;
110pub use resume::{
111 ResumeRepositorySourceMismatch, ResumeRepositorySourcePreflight, ResumeRepositorySourceReceipt,
112 raw_conversion_preview_for,
113};
114pub use reviewer::reviewer_stager;
115pub use subagents::{RegisterSubagentRequest, stopped_subagent, subagent_has_handed_back};
116pub use worker_binary::{
117 WorkerBinaryAvailability, native_worker_binary_prerequisite, pin_worker_binary_sources,
118 ssh_worker_binary_prerequisite, warm_worker_binary_sources,
119 worker_binary_prerequisite_for_arch,
120};
121pub use worker_restart::WorkerUpgradeOutcome;
122pub use worktree::{
123 ResumePlan, local_project_repository, resume_compatibility, resume_compatibility_with_checkout,
124};
125
126pub struct Controller {
127 pub config: Config,
128 pub state: State,
129}
130
131#[derive(Debug)]
135pub struct ControllerStoreGuard {
136 file: File,
137}
138
139impl ControllerStoreGuard {
140 pub fn acquire() -> Result<Self> {
141 let directory = data_dir();
142 Self::acquire_at(&directory)
143 }
144
145 fn acquire_at(directory: &Path) -> Result<Self> {
146 Self::try_acquire_at(directory)?.with_context(|| {
147 format!(
148 "another Mjolnir controller is already using {}; stop it before starting this command",
149 directory.display()
150 )
151 })
152 }
153
154 pub fn try_acquire() -> Result<Option<Self>> {
156 Self::try_acquire_at(&data_dir())
157 }
158
159 fn try_acquire_at(directory: &Path) -> Result<Option<Self>> {
160 mj_core::config::ensure_may_control_store(directory, "own this data store")?;
162 std::fs::create_dir_all(directory)
163 .with_context(|| format!("create controller data directory {}", directory.display()))?;
164 let path = directory.join("controller.lock");
165 let mut options = OpenOptions::new();
166 options.create(true).read(true).write(true);
167 #[cfg(unix)]
168 {
169 use std::os::unix::fs::OpenOptionsExt;
170 options.mode(0o600);
171 }
172 let file = options
173 .open(&path)
174 .with_context(|| format!("open controller lock {}", path.display()))?;
175 match file.try_lock() {
176 Ok(()) => {}
177 Err(std::fs::TryLockError::WouldBlock) => return Ok(None),
178 Err(std::fs::TryLockError::Error(error)) => {
179 return Err(error)
180 .with_context(|| format!("lock controller store {}", directory.display()));
181 }
182 }
183 Ok(Some(Self { file }))
184 }
185
186 pub fn start_database_writer(&self) -> Result<crate::database::DatabaseWriterOwner> {
189 crate::database::start_database_writer()
190 }
191}
192
193impl Drop for ControllerStoreGuard {
194 fn drop(&mut self) {
195 let _ = self.file.unlock();
198 }
199}
200
201#[derive(Debug)]
205pub struct QuickBundleCreation {
206 pub config: Config,
207 pub bundle_id: String,
208}
209
210#[derive(Debug)]
214pub enum QuickBundleFailure {
215 InvalidSource(anyhow::Error),
216 Persistence(anyhow::Error),
217}
218
219impl std::fmt::Display for QuickBundleFailure {
220 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
221 match self {
222 Self::InvalidSource(error) => write!(formatter, "invalid repository source: {error}"),
223 Self::Persistence(error) => write!(formatter, "persist quick bundle: {error}"),
224 }
225 }
226}
227
228impl std::error::Error for QuickBundleFailure {
229 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
230 match self {
231 Self::InvalidSource(error) | Self::Persistence(error) => Some(error.root_cause()),
232 }
233 }
234}
235
236pub fn create_quick_bundle(
242 source: &str,
243) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
244 create_bundle_from_sources(&[source.to_owned()])
245}
246
247pub fn create_quick_bundle_in_config(config: &mut Config, source: &str) -> Result<String> {
252 create_bundle_from_sources_in_config(config, &[source.to_owned()])
253}
254
255pub fn create_bundle_from_sources(
262 sources: &[String],
263) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
264 let (config, bundle_id) = Config::update(|config| {
265 let previous = config.bundles.clone();
266 let id = create_bundle_from_sources_in_config(config, sources)
267 .map_err(|error| anyhow::Error::new(QuickBundleFailure::InvalidSource(error)))?;
268 if let Some(old) = previous.get(&id)
269 && config.bundles.get(&id) != Some(old)
270 {
271 let accepted = crate::project_catalog::snapshot(
272 old,
273 &targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(
274 15,
275 )),
276 true,
277 )?;
278 crate::database::store_catalog_project(&id, &accepted, true)?;
279 }
280 let project = crate::project_catalog::snapshot(
281 config
282 .bundles
283 .get(&id)
284 .context("created project is missing")?,
285 &targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15)),
286 false,
287 )?;
288 let canonical = crate::database::store_catalog_project(&id, &project, true)?;
289 if canonical != id {
290 let mut saved = crate::database::read_project_catalog()?
291 .projects
292 .into_iter()
293 .find(|entry| entry.bundle_id == canonical)
294 .context("canonical project is missing")?
295 .project;
296 for repository in &mut saved.bundle.repositories {
297 let identity = saved
298 .identities
299 .get(&repository.id)
300 .context("canonical repository identity is missing")?;
301 if let Some(source) = project.bundle.repositories.iter().find(|source| {
302 project.identities.get(&source.id) == Some(identity)
303 && source.destination == repository.destination
304 }) {
305 repository.local = source.local.clone();
306 repository.github = source.github.clone();
307 }
308 }
309 config
310 .bundles
311 .insert(canonical.clone(), saved.bundle.clone());
312 config.bundles.remove(&id);
313 crate::database::store_catalog_project(&canonical, &saved, true)?;
314 }
315 Ok(canonical)
316 })
317 .map_err(|error| {
318 error
319 .downcast::<QuickBundleFailure>()
320 .unwrap_or_else(QuickBundleFailure::Persistence)
321 })?;
322 Ok(QuickBundleCreation { config, bundle_id })
323}
324
325pub fn remove_bundle(bundle_id: &str) -> Result<Config> {
330 let state = crate::database::load_state()?;
331 if let Some(refusal) = state.bundle_removal_refusal(bundle_id) {
332 bail!(refusal);
333 }
334 let (config, ()) = Config::update(|config| {
335 ensure!(
336 config.bundles.remove(bundle_id).is_some(),
337 "project {bundle_id:?} is no longer saved"
338 );
339 config.validate()
340 })?;
341 crate::database::hide_catalog_project(bundle_id)?;
342 Ok(config)
343}
344
345pub fn create_bundle_from_sources_in_config(
349 config: &mut Config,
350 sources: &[String],
351) -> Result<String> {
352 let sources = sources
353 .iter()
354 .map(|source| interpret_repository_source(source))
355 .collect::<Result<Vec<_>>>()?;
356 if sources.is_empty() {
357 bail!("at least one repository source is required");
358 }
359
360 let mut identities = BTreeSet::new();
361 for source in &sources {
362 if !identities.insert(source.identity()) {
363 bail!("duplicate repository source {:?}", source.display_name);
364 }
365 }
366
367 if let Some(existing) = exact_configured_bundle(config, &sources)? {
368 let bundle = config
369 .bundles
370 .get_mut(&existing)
371 .expect("matched project exists");
372 for repository in &mut bundle.repositories {
375 let Some(identity) =
376 crate::import::configured_repository_identity_for_lookup(&existing, repository)
377 else {
378 continue;
379 };
380 if let Some(source) = sources.iter().find(|source| {
381 source.identity == identity && matches!(source.kind, RepositorySourceKind::Local(_))
382 }) && let RepositorySourceKind::Local(root) = &source.kind
383 {
384 repository.local = Some(root.clone());
385 repository.github = None;
386 }
387 }
388 return Ok(existing);
389 }
390
391 let mut updated = config.clone();
394 let mut used_repository_ids = BTreeSet::new();
395 let mut repositories = Vec::with_capacity(sources.len());
396 for source in sources {
397 let base = setup_style_id(&source.name);
398 let repository_id = unique_id(&base, |candidate| used_repository_ids.contains(candidate));
399 used_repository_ids.insert(repository_id.clone());
400 repositories.push(source.into_project_repository(repository_id));
401 }
402 let primary_repo = repositories
403 .first()
404 .map(|repository| repository.id.clone())
405 .context("at least one repository source is required")?;
406 let bundle_id = unique_id(&primary_repo, |candidate| {
407 updated.bundles.contains_key(candidate)
408 });
409 updated.bundles.insert(
410 bundle_id.clone(),
411 ProjectBundle {
412 primary_repo,
413 repositories,
414 },
415 );
416 updated.validate()?;
417 *config = updated;
418 Ok(bundle_id)
419}
420
421#[derive(Debug, Clone)]
422enum RepositorySourceKind {
423 Github(crate::setup::GithubRepository),
424 Local(PathBuf),
425}
426
427#[derive(Debug, Clone)]
428struct InterpretedRepositorySource {
429 display_name: String,
430 name: String,
431 kind: RepositorySourceKind,
432 identity: RepositoryIdentity,
433}
434
435impl InterpretedRepositorySource {
436 fn identity(&self) -> RepositoryIdentity {
437 self.identity.clone()
438 }
439
440 fn into_project_repository(self, id: String) -> ProjectRepository {
441 let (github, local) = match self.kind {
442 RepositorySourceKind::Github(repository) => (
443 Some(format!("{}/{}", repository.owner, repository.repository)),
444 None,
445 ),
446 RepositorySourceKind::Local(root) => (None, Some(root)),
447 };
448 ProjectRepository {
449 id: id.clone(),
450 github,
451 local,
452 destination: PathBuf::from(id),
453 git_ref: None,
454 }
455 }
456}
457
458fn interpret_repository_source(source: &str) -> Result<InterpretedRepositorySource> {
461 let source = source.trim();
462 if source.is_empty() {
463 bail!("repository source cannot be empty");
464 }
465 let expanded = mj_core::path_input::expand_local(Path::new(source))?;
466 let candidate = expanded.as_path();
467 if candidate.exists() {
468 let resolved = mj_core::repository::resolve_directory(
469 candidate,
470 None,
471 &mj_core::targets::CancellableProcessExecutor::with_timeout(
472 std::time::Duration::from_secs(8),
473 ),
474 )?;
475 let identity = resolved.identity;
476 let name = identity.name();
477 return Ok(InterpretedRepositorySource {
478 display_name: source.to_owned(),
479 name,
480 identity,
481 kind: RepositorySourceKind::Local(resolved.checkout_root),
482 });
483 }
484 if candidate.is_absolute() || source.starts_with('.') || source.starts_with('~') {
485 bail!("local repository path {source:?} does not exist");
486 }
487 let repository = github_repository_from_origin(source).context(format!(
488 "{source:?} is not a GitHub owner/repository or URL"
489 ))?;
490 Ok(InterpretedRepositorySource {
491 display_name: source.to_owned(),
492 name: repository.repository.clone(),
493 identity: RepositoryIdentity::Github(
494 repository.owner.to_ascii_lowercase(),
495 repository.repository.to_ascii_lowercase(),
496 ),
497 kind: RepositorySourceKind::Github(repository),
498 })
499}
500
501fn exact_configured_bundle(
502 config: &Config,
503 requested: &[InterpretedRepositorySource],
504) -> Result<Option<String>> {
505 let primary = requested
506 .first()
507 .context("no repository sources")?
508 .identity();
509 let mut used = BTreeSet::new();
510 let requested_layout = requested
511 .iter()
512 .map(|source| {
513 let id = unique_id(&setup_style_id(&source.name), |id| used.contains(id));
514 used.insert(id.clone());
515 (source.identity(), PathBuf::from(id))
516 })
517 .collect::<BTreeSet<_>>();
518 for (id, bundle) in &config.bundles {
519 if bundle.repositories.len() != requested.len()
520 || bundle
521 .repositories
522 .iter()
523 .any(|repository| repository.git_ref.is_some())
524 {
525 continue;
526 }
527 let identity_results = bundle
528 .repositories
529 .iter()
530 .map(|repo| crate::import::configured_repository_identity_for_lookup(id, repo))
531 .collect::<Vec<_>>();
532 if identity_results.iter().any(Option::is_none) {
533 continue;
534 }
535 let identities = identity_results.into_iter().flatten().collect::<Vec<_>>();
536 let layout = bundle
537 .repositories
538 .iter()
539 .zip(&identities)
540 .map(|(repository, identity)| (identity.clone(), repository.destination.clone()))
541 .collect::<BTreeSet<_>>();
542 let Some(primary_index) = bundle
543 .repositories
544 .iter()
545 .position(|repository| repository.id == bundle.primary_repo)
546 else {
547 continue;
548 };
549 if layout == requested_layout && identities.get(primary_index) == Some(&primary) {
550 return Ok(Some(id.clone()));
551 }
552 }
553 Ok(None)
554}
555
556fn unique_id(base: &str, mut is_used: impl FnMut(&str) -> bool) -> String {
557 if !is_used(base) {
558 return base.to_owned();
559 }
560 for suffix in 2_u32.. {
561 let suffix = format!("-{suffix}");
562 let prefix_len = 64usize.saturating_sub(suffix.len());
563 let prefix = base.chars().take(prefix_len).collect::<String>();
564 let candidate = format!("{prefix}{suffix}");
565 if !is_used(&candidate) {
566 return candidate;
567 }
568 }
569 unreachable!("u32 repository/bundle id suffixes exhausted")
570}
571
572pub struct SessionLaunchOptions {
573 pub create_managed_worktree: Option<bool>,
574 pub at: Option<String>,
576 pub branch: Option<String>,
578 pub base: Option<String>,
580 pub subagents: Option<mj_core::subagent::SubagentPolicy>,
581 pub review: Option<mj_core::config::SessionReview>,
583 pub initial_prompt: Option<String>,
584 pub workspace_id: String,
585 pub additional_mounts: Vec<AdditionalMount>,
586 pub resource_allocation: Option<SessionResourceAllocation>,
587 pub project_directory: Option<PathBuf>,
588 pub session_title_override: Option<String>,
589}
590
591pub struct SessionResumeOptions {
592 pub additional_mounts: Option<Vec<AdditionalMount>>,
593 pub resource_allocation: Option<SessionResourceAllocation>,
594 pub discard_queue: bool,
595}
596
597fn selected_host_container_size(
598 template: &TargetTemplate,
599 allocation: Option<&SessionResourceAllocation>,
600) -> Option<(String, HostContainerSize)> {
601 let host = container_size_host(template)?;
602 let SessionResourceAllocation::Container { cpus, memory_bytes } = allocation? else {
603 return None;
604 };
605 Some((
606 host.to_owned(),
607 HostContainerSize {
608 cpus: *cpus,
609 memory_bytes: *memory_bytes,
610 },
611 ))
612}
613
614impl Controller {
615 pub fn load() -> Result<Self> {
616 let config = Config::load().map_err(|error| {
619 let refusal = mj_core::refusal::Refusal::precondition(format!("{error:#}"));
620 error.context(refusal)
621 })?;
622 let state = crate::database::load_state()?;
623 Ok(Self { config, state })
624 }
625
626 pub(crate) fn prepare_persisted_sessions(&mut self) -> Result<()> {
629 self.state.validate()?;
632 for session in self.state.sessions.values() {
633 if let Some(issue) = session.configuration_issue(&self.config) {
634 tracing::warn!(session_id = %session.id, "{issue}");
635 }
636 }
637 self.backfill_target_runtime()
638 }
639
640 pub fn reload(&mut self) -> Result<()> {
641 self.backfill_target_runtime()?;
644 *self = Self::load()?;
645 Ok(())
646 }
647
648 fn backfill_target_runtime(&mut self) -> Result<()> {
649 let missing: Vec<_> = self
650 .state
651 .sessions
652 .values()
653 .filter(|session| session.target_runtime.is_none())
654 .map(|session| session.id.clone())
655 .collect();
656 for id in missing {
657 let session = self.state.sessions.get_mut(&id).expect("collected session");
658 match session.target_runtime_settings(&self.config) {
659 Ok(runtime) => {
660 let runtime = runtime.into_owned();
661 session.target_runtime = crate::database::backfill_target_runtime(
662 &session.id,
663 &session.target_template_id,
664 &runtime,
665 )?;
666 }
667 Err(error) => {
668 tracing::warn!(session_id = %session.id, %error, "target access needs configuration repair")
669 }
670 }
671 }
672 Ok(())
673 }
674
675 fn persist_session_state(&self, session_id: &str) -> Result<()> {
676 match self.state.sessions.get(session_id) {
677 Some(session) => crate::database::save_lifecycle_session(session),
678 None => crate::database::delete_session(session_id),
679 }
680 }
681
682 fn persist_session_transition_or_restore(
683 &mut self,
684 session_id: &str,
685 previous: &SessionRecord,
686 context: &'static str,
687 ) -> Result<()> {
688 persist_session_record_transition_or_restore(
689 &mut self.state,
690 session_id,
691 previous,
692 context,
693 &crate::database::save_lifecycle_session,
694 )
695 }
696
697 fn restore_prior_session_after_persistence_failure(
698 &mut self,
699 session_id: &str,
700 previous: &SessionRecord,
701 primary: anyhow::Error,
702 ) -> anyhow::Error {
703 restore_session_after_persistence_failure(
704 &mut self.state,
705 session_id,
706 previous,
707 primary,
708 crate::database::save_lifecycle_session,
709 )
710 }
711
712 pub fn resolve_input_path(
714 &self,
715 target_id: &str,
716 path: &Path,
717 executor: &impl CommandExecutor,
718 ) -> Result<PathBuf> {
719 let target = self
720 .config
721 .targets
722 .get(target_id)
723 .context("Unknown path target")?;
724 resolve_target_input_path(target, path, executor)
725 }
726
727 pub fn validate_mount_source(
734 &self,
735 target_id: &str,
736 source: &Path,
737 executor: &impl CommandExecutor,
738 ) -> Result<Option<String>> {
739 let target = self
740 .config
741 .targets
742 .get(target_id)
743 .with_context(|| format!("unknown target template {target_id:?}"))?;
744 let exists = match target {
745 TargetTemplate::LocalPodman { .. }
746 | TargetTemplate::LocalDocker { .. }
747 | TargetTemplate::AppleContainer { .. }
748 | TargetTemplate::AwsEc2 { .. } => std::fs::metadata(source)
749 .map(|metadata| metadata.is_dir())
750 .or_else(|error| {
751 if error.kind() == std::io::ErrorKind::NotFound {
752 Ok(false)
753 } else {
754 Err(error)
755 }
756 })
757 .with_context(|| format!("inspect resource source {}", source.display()))?,
758 TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
759 targets::ssh_directory_exists(&SshTarget::from(ssh), source, executor)?
760 }
761 TargetTemplate::LocalBare | TargetTemplate::SshBare { .. } => {
762 bail!("resource attachments are unsupported for bare targets")
763 }
764 };
765 ensure!(
766 exists,
767 "source path {} does not exist or is not a directory",
768 source.display()
769 );
770 Ok(self.forced_read_only_reason(target, source, executor))
771 }
772
773 fn forced_read_only_reason(
775 &self,
776 target: &TargetTemplate,
777 source: &Path,
778 executor: &impl CommandExecutor,
779 ) -> Option<String> {
780 let ssh = match target {
781 TargetTemplate::LocalPodman { .. } | TargetTemplate::LocalDocker { .. } => None,
782 TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
783 Some(SshTarget::from(ssh))
784 }
785 _ => return None,
788 };
789 if matches!(target, TargetTemplate::LocalDocker { .. })
790 && let Ok(Some(_)) = targets::local_docker_vm_share(executor)
791 {
792 return Some("Docker VM file share (cannot back an overlay)".to_owned());
793 }
794 let filesystem = targets::probe_filesystem_types(
795 ssh.as_ref(),
796 std::slice::from_ref(&source.to_path_buf()),
797 executor,
798 )
799 .map_err(|error| {
800 tracing::debug!(
801 source = %source.display(),
802 error = format!("{error:#}"),
803 "could not probe the filesystem under a mount source"
804 );
805 })
806 .ok()?
807 .pop()?;
808 let reason = targets::overlay_unsupported_filesystem(&filesystem)?;
809 Some(format!("{filesystem} ({reason})"))
810 }
811
812 fn fail_new_session_with_cleanup(
813 &mut self,
814 session_id: &str,
815 error: anyhow::Error,
816 executor: &impl CommandExecutor,
817 ) -> Result<anyhow::Error> {
818 let original = provisioning::note_new_session_launch_failure(session_id, &error);
819 let cleanup_error = self
820 .cleanup_new_session_worktree_after_failure(session_id, executor)
821 .err()
822 .map(|cleanup_error| format!("{cleanup_error:#}"));
823 if let Some(cleanup_error) = &cleanup_error {
824 tracing::warn!(
825 session_id,
826 error = %cleanup_error,
827 "new-session worktree rollback reported a cleanup failure"
828 );
829 }
830 let failure = apply_failed_new_session_rollback(
831 &mut self.state,
832 session_id,
833 &original,
834 cleanup_error,
835 );
836 self.persist_session_state(session_id)?;
837 Ok(failure)
838 }
839
840 pub fn register_session_with_resources(
841 &mut self,
842 profile_id: &str,
843 bundle_id: &str,
844 target_id: &str,
845 title: impl Into<String>,
846 options: SessionLaunchOptions,
847 ) -> Result<String> {
848 self.register_session_with_executor(
849 profile_id,
850 bundle_id,
851 target_id,
852 title,
853 options,
854 &targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15)),
855 )
856 }
857
858 fn register_session_with_executor(
859 &mut self,
860 profile_id: &str,
861 bundle_id: &str,
862 target_id: &str,
863 title: impl Into<String>,
864 options: SessionLaunchOptions,
865 executor: &impl CommandExecutor,
866 ) -> Result<String> {
867 let SessionLaunchOptions {
868 create_managed_worktree,
869 at,
870 branch,
871 base,
872 subagents,
873 review,
874 initial_prompt,
875 workspace_id,
876 additional_mounts,
877 resource_allocation,
878 project_directory,
879 session_title_override,
880 } = options;
881 let launch_base = match base {
882 Some(base) => {
883 let base = base.trim();
884 if base.is_empty() {
885 bail!("`base` must not be empty");
886 }
887 Some(base.to_owned())
888 }
889 None => None,
890 };
891 let branch = match branch {
892 Some(branch) => {
893 let branch = branch.trim();
894 ensure!(!branch.is_empty(), "`branch` must not be empty");
895 Some(branch.to_owned())
896 }
897 None => None,
898 };
899 if launch_base.is_some() && create_managed_worktree == Some(false) {
900 bail!("`base` requires a managed worktree or a bundle session");
901 }
902 let session_title_override = match session_title_override {
903 Some(title) => {
904 Some(normalize_session_title(&title).context("session name cannot be empty")?)
905 }
906 None => None,
907 };
908 let profile = self
909 .config
910 .profiles
911 .get(profile_id)
912 .with_context(|| format!("unknown profile {profile_id:?}"))?;
913 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
914 let template = self
915 .config
916 .targets
917 .get(target_id)
918 .with_context(|| format!("unknown target template {target_id:?}"))?;
919 if create_managed_worktree == Some(true) && !is_bare_project_target(template) {
920 bail!("managed worktree creation requires a bare Git project");
921 }
922 if project_directory.is_some() != is_bare_project_target(template) {
923 bail!("raw project directories require a bare target, and bare targets require one");
924 }
925 if let Some(path) = &project_directory
926 && (!path.is_absolute()
927 || path
928 .components()
929 .any(|part| part == std::path::Component::ParentDir))
930 {
931 bail!("bare project directory must be an absolute safe path");
932 }
933 let bundle = project_directory
934 .is_none()
935 .then(|| self.config.bundles.get(bundle_id))
936 .flatten();
937 if project_directory.is_none() && bundle.is_none() {
938 bail!("unknown bundle {bundle_id:?}");
939 }
940 let (checkout, launch_branch) = match at {
944 Some(at) => {
945 let bundle = bundle.context("`at` requires a bundle-backed session")?;
946 ensure!(
947 create_managed_worktree != Some(false),
948 "`at` requires an isolated workspace"
949 );
950 let mut checkout = mj_core::remote_git::ExactCheckout {
951 repository_id: bundle.primary_repo.clone(),
952 commit: at.trim().to_owned(),
953 branch,
954 };
955 checkout.validate()?;
956 checkout.commit.make_ascii_lowercase();
957 (Some(checkout), None)
958 }
959 None => (None, branch),
960 };
961 if profile.kind == mj_core::config::HarnessKind::Muse
962 && (!additional_mounts.is_empty()
963 || bundle.is_some_and(|bundle| bundle.repositories.len() > 1))
964 {
965 bail!(
966 "{} ACP supports one workspace root; use a single-repository bundle without attached directories",
967 profile.kind.display_name()
968 );
969 }
970 let (bundle_id, project) = if let Some(bundle) = bundle {
971 let project = crate::project_catalog::snapshot(bundle, executor, true)?;
972 let canonical = crate::database::store_catalog_project(bundle_id, &project, true)?;
973 (canonical, Some(project))
974 } else if let Some(directory) = project_directory.as_deref() {
975 let git = match template {
976 TargetTemplate::LocalBare => {
977 executor
978 .execute(
979 &targets::CommandSpec::new(
980 "git",
981 [
982 "-C",
983 &directory.to_string_lossy(),
984 "rev-parse",
985 "--is-inside-work-tree",
986 ],
987 )
988 .purpose("inspect raw project repository"),
989 )?
990 .status
991 == 0
992 }
993 _ => true,
994 };
995 if git {
996 let (id, project) =
997 crate::project_catalog::accept_directory(template, directory, executor)?;
998 (id, Some(project))
999 } else {
1000 (bundle_id.to_owned(), None)
1001 }
1002 } else {
1003 (bundle_id.to_owned(), None)
1004 };
1005 validate_resource_allocation(template, resource_allocation.as_ref())?;
1006 let selected_container_size =
1007 selected_host_container_size(template, resource_allocation.as_ref());
1008 if !additional_mounts.is_empty() && mount_history_host(template).is_none() {
1009 bail!("attached resources are unsupported for this target");
1010 }
1011 targets::validate_additional_mounts(&additional_mounts)?;
1012 let id = new_session_id()?;
1013 let now = now();
1014 let record = SessionRecord {
1015 project,
1016 target_runtime: Some(template.into()),
1017 build_cache: None,
1018 create_managed_worktree,
1019 launch_base,
1020 launch_branch,
1021 checkout,
1022 publication: None,
1023 subagents: Some(subagents.unwrap_or_else(|| profile.subagents.clone())),
1024 archived: false,
1025 container_cpus: None,
1026 container_memory: None,
1027 container_workspace: Some(targets::new_container_workspace(&id)?),
1032 id: id.clone(),
1033 workspace_id,
1034 title: title.into(),
1035 harness_kind: profile.kind,
1036 last_profile: profile_id.to_string(),
1037 bundle_id,
1038 project_directory,
1039 managed_worktree: None,
1040 review,
1041 target_template_id: target_id.to_string(),
1042 resource_allocation,
1043 additional_mounts: additional_mounts.clone(),
1044 state: SessionState::Provisioning,
1045 target: None,
1046 native_session_id: None,
1047 acp_session_title: None,
1048 session_title_override,
1049 created_at: now.clone(),
1050 updated_at: now,
1051 viewed_through_event_ordinal: 0,
1052 draft_input: initial_prompt.unwrap_or_default(),
1053 last_error: None,
1054 last_checkpoint_error: None,
1055 checkpoint: None,
1056 };
1057 crate::database::save_new_session(&record, selected_container_size.clone())?;
1061 if record.harness_kind.supports_delegation_tools() {
1062 self.state.last_subagent_policy = record.subagents.clone().unwrap_or_default();
1063 }
1064 self.state.sessions.insert(id.clone(), record);
1065 if let Some((host, size)) = selected_container_size {
1066 self.state.remember_container_size(&host, size);
1067 }
1068 if let Some(host) = mount_history_host(template) {
1069 match crate::database::remember_mount_sources(host, &additional_mounts) {
1073 Ok(()) => self.state.remember_mount_sources(host, &additional_mounts),
1074 Err(error) => tracing::warn!(
1075 session_id = id,
1076 error = format!("{error:#}"),
1077 "could not remember the attached resource directories for later suggestions"
1078 ),
1079 }
1080 }
1081 Ok(id)
1082 }
1083
1084 pub fn rename_session(&mut self, session_id: &str, title: &str) -> Result<String> {
1085 let title = normalize_session_title(title).context("session name cannot be empty")?;
1086 ensure!(
1087 self.state.sessions.contains_key(session_id),
1088 "unknown session {session_id}"
1089 );
1090 let updated_at = now();
1091 crate::database::set_session_title_override(session_id, &title, &updated_at)?;
1092 let record = self
1093 .state
1094 .sessions
1095 .get_mut(session_id)
1096 .expect("session was checked before updating its title");
1097 record.session_title_override = Some(title.clone());
1098 record.updated_at = updated_at;
1099 Ok(title)
1100 }
1101
1102 pub fn set_session_workspace(&mut self, session_id: &str, workspace_id: &str) -> Result<usize> {
1105 ensure!(
1106 self.state.sessions.contains_key(session_id),
1107 "unknown session {session_id}"
1108 );
1109 let mut ids = vec![session_id.to_owned()];
1110 ids.extend(
1111 self.state
1112 .subagents
1113 .values()
1114 .filter(|child| child.parent_session_id == session_id)
1115 .map(|child| child.child_session_id.clone()),
1116 );
1117 crate::database::set_sessions_workspace(&ids, workspace_id)?;
1118 for id in &ids {
1119 if let Some(record) = self.state.sessions.get_mut(id) {
1120 record.workspace_id = workspace_id.to_owned();
1121 }
1122 }
1123 Ok(ids.len())
1124 }
1125
1126 pub fn rename_profile_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
1127 mj_core::config::validate_id("profile", new_id)?;
1128 if old_id == new_id {
1129 ensure!(
1130 self.config.profiles.contains_key(old_id),
1131 "unknown profile {old_id:?}"
1132 );
1133 return Ok(());
1134 }
1135 let journal = ConfigRenameJournal {
1136 kind: ConfigRenameKind::Profile,
1137 old_id: old_id.to_owned(),
1138 new_id: new_id.to_owned(),
1139 };
1140 write_config_rename_journal(&journal)?;
1141 let (config, ()) = match Config::update(|config| {
1142 ensure!(
1143 config.profiles.contains_key(old_id),
1144 "unknown profile {old_id:?}"
1145 );
1146 ensure!(
1147 !config.profiles.contains_key(new_id),
1148 "profile {new_id:?} already exists"
1149 );
1150 let profile = config
1151 .profiles
1152 .remove(old_id)
1153 .expect("profile was checked in the transaction");
1154 config.profiles.insert(new_id.to_owned(), profile);
1155 Ok(())
1156 }) {
1157 Ok(result) => result,
1158 Err(error) => {
1159 remove_config_rename_journal()
1160 .context("remove profile rename journal after config save failed")?;
1161 return Err(error).context("save renamed profile configuration");
1162 }
1163 };
1164 self.config = config;
1165 mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
1166 if let Err(error) = crate::database::rename_profile_references(old_id, new_id) {
1167 let restore = Config::update(|config| {
1168 let profile = config
1169 .profiles
1170 .remove(new_id)
1171 .with_context(|| format!("renamed profile {new_id:?} is missing"))?;
1172 ensure!(
1173 !config.profiles.contains_key(old_id),
1174 "cannot restore profile rename: both {old_id:?} and {new_id:?} exist"
1175 );
1176 config.profiles.insert(old_id.to_owned(), profile);
1177 Ok(())
1178 });
1179 let restored = match restore {
1180 Ok((config, ())) => config,
1181 Err(restore_error) => {
1182 return Err(error).context(format!(
1183 "rename profile references; additionally failed to restore config: {restore_error:#}"
1184 ));
1185 }
1186 };
1187 self.config = restored;
1188 if let Err(restore_error) = remove_config_rename_journal() {
1189 return Err(error).context(format!(
1190 "rename profile references; additionally failed to remove rename journal: {restore_error:#}"
1191 ));
1192 }
1193 return Err(error).context("rename profile references");
1194 }
1195 let affected: Vec<_> = self
1196 .state
1197 .sessions
1198 .values()
1199 .filter(|session| session.last_profile == old_id)
1200 .map(|session| session.id.clone())
1201 .collect();
1202 for id in affected {
1203 self.state
1204 .sessions
1205 .get_mut(&id)
1206 .expect("collected session")
1207 .last_profile = new_id.to_owned();
1208 }
1209 remove_config_rename_journal()?;
1210 Ok(())
1211 }
1212
1213 pub fn rename_target_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
1214 mj_core::config::validate_id("target template", new_id)?;
1215 if old_id == new_id {
1216 ensure!(
1217 self.config.targets.contains_key(old_id),
1218 "unknown target {old_id:?}"
1219 );
1220 return Ok(());
1221 }
1222 let journal = ConfigRenameJournal {
1223 kind: ConfigRenameKind::Target,
1224 old_id: old_id.to_owned(),
1225 new_id: new_id.to_owned(),
1226 };
1227 write_config_rename_journal(&journal)?;
1228 let (config, ()) = match Config::update(|config| {
1229 ensure!(
1230 config.targets.contains_key(old_id),
1231 "unknown target {old_id:?}"
1232 );
1233 ensure!(
1234 !config.targets.contains_key(new_id),
1235 "target {new_id:?} already exists"
1236 );
1237 let target = config
1238 .targets
1239 .remove(old_id)
1240 .expect("target was checked in the transaction");
1241 config.targets.insert(new_id.to_owned(), target);
1242 Ok(())
1243 }) {
1244 Ok(result) => result,
1245 Err(error) => {
1246 remove_config_rename_journal()
1247 .context("remove target rename journal after config save failed")?;
1248 return Err(error).context("save renamed target configuration");
1249 }
1250 };
1251 self.config = config;
1252 mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
1253 if let Err(error) = crate::database::rename_target_references(old_id, new_id) {
1254 let restore = Config::update(|config| {
1255 let target = config
1256 .targets
1257 .remove(new_id)
1258 .with_context(|| format!("renamed target {new_id:?} is missing"))?;
1259 ensure!(
1260 !config.targets.contains_key(old_id),
1261 "cannot restore target rename: both {old_id:?} and {new_id:?} exist"
1262 );
1263 config.targets.insert(old_id.to_owned(), target);
1264 Ok(())
1265 });
1266 let restored = match restore {
1267 Ok((config, ())) => config,
1268 Err(restore_error) => {
1269 return Err(error).context(format!(
1270 "rename target references; additionally failed to restore config: {restore_error:#}"
1271 ));
1272 }
1273 };
1274 self.config = restored;
1275 if let Err(restore_error) = remove_config_rename_journal() {
1276 return Err(error).context(format!(
1277 "rename target references; additionally failed to remove rename journal: {restore_error:#}"
1278 ));
1279 }
1280 return Err(error).context("rename target references");
1281 }
1282 let affected: Vec<_> = self
1283 .state
1284 .sessions
1285 .values()
1286 .filter(|session| session.target_template_id == old_id)
1287 .map(|session| session.id.clone())
1288 .collect();
1289 for id in affected {
1290 self.state
1291 .sessions
1292 .get_mut(&id)
1293 .expect("collected session")
1294 .target_template_id = new_id.to_owned();
1295 }
1296 remove_config_rename_journal()?;
1297 Ok(())
1298 }
1299
1300 pub fn recover_config_id_rename() -> Result<bool> {
1304 let path = config_rename_journal_path();
1305 let body = match fs::read(&path) {
1306 Ok(body) => body,
1307 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(false),
1308 Err(error) => return Err(error).context(format!("read {}", path.display())),
1309 };
1310 let journal: ConfigRenameJournal =
1311 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1312 match journal.kind {
1313 ConfigRenameKind::Profile => {
1314 Config::update(|config| {
1315 finish_config_map_rename(
1316 &mut config.profiles,
1317 &journal.old_id,
1318 &journal.new_id,
1319 "profile",
1320 )?;
1321 Ok(())
1322 })?;
1323 crate::database::rename_profile_references(&journal.old_id, &journal.new_id)?;
1324 }
1325 ConfigRenameKind::Target => {
1326 Config::update_to(&mj_core::config::config_path(), |config| {
1331 finish_config_map_rename(
1332 &mut config.targets,
1333 &journal.old_id,
1334 &journal.new_id,
1335 "target",
1336 )?;
1337 Ok(())
1338 })?;
1339 crate::database::rename_target_references(&journal.old_id, &journal.new_id)?;
1340 }
1341 }
1342 remove_config_rename_journal()?;
1343 Ok(true)
1344 }
1345
1346 pub fn update_session_container_settings(
1350 &mut self,
1351 session_id: &str,
1352 cpus: Option<String>,
1353 memory: Option<String>,
1354 additional_mounts: Vec<targets::AdditionalMount>,
1355 mount_history: Vec<std::path::PathBuf>,
1356 executor: &impl CommandExecutor,
1357 ) -> Result<()> {
1358 let session = self
1359 .state
1360 .sessions
1361 .get(session_id)
1362 .with_context(|| format!("unknown session {session_id}"))?;
1363 for mount in &additional_mounts {
1370 if session
1371 .additional_mounts
1372 .iter()
1373 .any(|attached| attached.source == mount.source)
1374 {
1375 continue;
1376 }
1377 self.validate_mount_source(&session.target_template_id, &mount.source, executor)?;
1378 }
1379 let cpus = cpus.filter(|value| !value.trim().is_empty());
1380 let memory = memory.filter(|value| !value.trim().is_empty());
1381 let updated_at = now();
1382 crate::database::set_session_container_settings(
1383 session_id,
1384 cpus.as_deref(),
1385 memory.as_deref(),
1386 &additional_mounts,
1387 &updated_at,
1388 )?;
1389 if let Some(host) = self
1390 .config
1391 .targets
1392 .get(
1393 &self.state.sessions[session_id]
1394 .target_template_id
1395 .to_owned(),
1396 )
1397 .and_then(mj_core::config::mount_history_host)
1398 {
1399 let host = host.to_owned();
1400 crate::database::replace_mount_history(&host, &mount_history)?;
1403 crate::database::remember_mount_sources(&host, &additional_mounts)?;
1404 self.state.mount_history.insert(host.clone(), mount_history);
1405 self.state.remember_mount_sources(&host, &additional_mounts);
1406 }
1407 let record = self
1408 .state
1409 .sessions
1410 .get_mut(session_id)
1411 .expect("session was checked before updating its container settings");
1412 record.container_cpus = cpus;
1413 record.container_memory = memory;
1414 record.additional_mounts = additional_mounts;
1415 record.updated_at = updated_at;
1416 Ok(())
1417 }
1418}
1419
1420fn config_rename_journal_path() -> PathBuf {
1421 data_dir().join(CONFIG_RENAME_JOURNAL)
1422}
1423
1424fn write_config_rename_journal(journal: &ConfigRenameJournal) -> Result<()> {
1425 let path = config_rename_journal_path();
1426 let body = serde_json::to_vec(journal).context("serialize config rename journal")?;
1427 atomic_write(&path, &body).with_context(|| format!("write {}", path.display()))
1428}
1429
1430fn remove_config_rename_journal() -> Result<()> {
1431 let path = config_rename_journal_path();
1432 match fs::remove_file(&path) {
1433 Ok(()) => Ok(()),
1434 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
1435 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
1436 }
1437}
1438
1439fn finish_config_map_rename<T>(
1440 entries: &mut BTreeMap<String, T>,
1441 old_id: &str,
1442 new_id: &str,
1443 kind: &str,
1444) -> Result<()> {
1445 if let Some(entry) = entries.remove(old_id) {
1446 ensure!(
1447 !entries.contains_key(new_id),
1448 "cannot recover {kind} rename: both {old_id:?} and {new_id:?} exist"
1449 );
1450 entries.insert(new_id.to_owned(), entry);
1451 } else {
1452 ensure!(
1453 entries.contains_key(new_id),
1454 "cannot recover {kind} rename: neither {old_id:?} nor {new_id:?} exists"
1455 );
1456 }
1457 Ok(())
1458}
1459
1460#[cfg(test)]
1468pub(crate) fn target_profile_home_for_test(
1469 locator: &targets::TargetLocator,
1470 session_id: &str,
1471 profile: &mj_core::config::HarnessProfile,
1472) -> String {
1473 target_profile_home(locator, session_id, profile)
1474}
1475
1476fn target_profile_home(
1477 locator: &targets::TargetLocator,
1478 session_id: &str,
1479 profile: &mj_core::config::HarnessProfile,
1480) -> String {
1481 let root = removable_profile_root(locator, session_id, profile);
1482 if profile.kind.nested_home() {
1483 PathBuf::from(root)
1484 .join(profile.kind.id())
1485 .to_string_lossy()
1486 .into_owned()
1487 } else {
1488 root
1489 }
1490}
1491
1492pub(super) fn removable_profile_root(
1500 locator: &targets::TargetLocator,
1501 session_id: &str,
1502 profile: &mj_core::config::HarnessProfile,
1503) -> String {
1504 match locator {
1505 targets::TargetLocator::LocalBare { .. } if profile.kind.nested_home() => {
1506 targets::local_nested_profile_root(session_id)
1507 .to_string_lossy()
1508 .into_owned()
1509 }
1510 targets::TargetLocator::LocalBare { worker_root } => Path::new(worker_root)
1511 .join("profile")
1512 .to_string_lossy()
1513 .into_owned(),
1514 targets::TargetLocator::LocalPodman { .. }
1515 | targets::TargetLocator::LocalDocker { .. }
1516 | targets::TargetLocator::AppleContainer { .. }
1517 | targets::TargetLocator::SshPodman { .. }
1518 | targets::TargetLocator::SshDocker { .. } => {
1519 format!("/var/lib/hel/profiles/{session_id}")
1520 }
1521 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. } => {
1522 format!(".local/share/hel/profiles/{session_id}")
1523 }
1524 }
1525}
1526
1527pub fn resolve_target_input_path(
1529 target: &TargetTemplate,
1530 path: &Path,
1531 executor: &impl CommandExecutor,
1532) -> Result<PathBuf> {
1533 if !mj_core::path_input::needs_home(path)? {
1534 return Ok(path.to_path_buf());
1535 }
1536 let host = cache_host::CacheHost::for_path_target(target)?;
1537 mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1538}
1539
1540pub fn resolve_machine_input_path(
1543 machine: &mj_core::config::Machine,
1544 path: &Path,
1545 executor: &impl CommandExecutor,
1546) -> Result<PathBuf> {
1547 if !mj_core::path_input::needs_home(path)? {
1548 return Ok(path.to_path_buf());
1549 }
1550 let host = cache_host::CacheHost::for_path_machine(machine)?;
1553 mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1554}
1555
1556fn execute_checked(executor: &impl CommandExecutor, command: CommandSpec) -> Result<CommandOutput> {
1557 let output = executor.execute(&command)?;
1558 if output.status != 0 {
1559 let detail = command_error_detail(&output.stderr);
1560 if detail.is_empty() {
1561 bail!("{} failed with status {}", command.purpose, output.status);
1562 }
1563 bail!("{detail}");
1564 }
1565 Ok(output)
1566}
1567
1568fn command_error_detail(stderr: &[u8]) -> String {
1569 let reported = String::from_utf8_lossy(stderr);
1570 let reported = reported.trim();
1571 let detail = reported
1572 .rsplit_once("\nCaused by:\n")
1573 .map_or(reported, |(_, causes)| causes);
1574 let detail = detail.strip_prefix("Error: ").unwrap_or(detail);
1575 detail
1576 .lines()
1577 .map(|line| line.strip_prefix(" ").unwrap_or(line))
1578 .collect::<Vec<_>>()
1579 .join("\n")
1580 .trim()
1581 .to_owned()
1582}
1583
1584fn now() -> String {
1585 Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1586}
1587
1588fn restore_session_after_persistence_failure(
1589 state: &mut State,
1590 session_id: &str,
1591 previous: &SessionRecord,
1592 primary: anyhow::Error,
1593 persist: impl FnOnce(&SessionRecord) -> Result<()>,
1594) -> anyhow::Error {
1595 state
1596 .sessions
1597 .insert(session_id.to_owned(), previous.clone());
1598 let restored = state
1599 .sessions
1600 .get(session_id)
1601 .expect("restored session record disappeared");
1602 match persist(restored) {
1603 Ok(()) => primary,
1604 Err(error) => primary.context(format!(
1605 "restored prior session state in memory, but failed to persist the rollback: {error:#}"
1606 )),
1607 }
1608}
1609
1610fn persist_session_record_transition_or_restore(
1611 state: &mut State,
1612 session_id: &str,
1613 previous: &SessionRecord,
1614 context: &'static str,
1615 persist: &impl Fn(&SessionRecord) -> Result<()>,
1616) -> Result<()> {
1617 let result = persist(
1618 state
1619 .sessions
1620 .get(session_id)
1621 .expect("checkpoint session disappeared before persistence"),
1622 );
1623 match result {
1624 Ok(()) => Ok(()),
1625 Err(error) => Err(restore_session_after_persistence_failure(
1626 state,
1627 session_id,
1628 previous,
1629 error.context(context),
1630 persist,
1631 )),
1632 }
1633}
1634
1635pub fn config_only_controller(config: Config) -> Controller {
1636 Controller {
1637 config,
1638 state: State::default(),
1639 }
1640}
1641
1642#[cfg(test)]
1643mod tests;