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