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