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