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