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