Skip to main content

mj_controller/
controller.rs

1//! Controller-side lifecycle transitions and canonical-to-backend conversion.
2
3mod backend;
4mod cache_host;
5pub(crate) mod checkpoint;
6mod git_cache;
7mod lifecycle;
8mod mbx;
9pub mod move_session;
10mod network_git;
11mod new_session_preflight;
12pub use new_session_preflight::{NewSessionPreflight, NewSessionRepository};
13mod path_completion;
14pub mod profile_config;
15mod provisioning;
16pub(crate) mod publication;
17mod readiness;
18pub(crate) use readiness::NATIVE_SESSION_STARTUP_TIMEOUT;
19mod recovery_scan;
20mod resume;
21mod reviewer;
22mod subagents;
23#[cfg(test)]
24pub(crate) mod test_support;
25pub mod update;
26mod worker_binary;
27mod worker_restart;
28mod worktree;
29
30use std::collections::{BTreeMap, BTreeSet};
31use std::fs::{self, File, OpenOptions};
32use std::path::{Path, PathBuf};
33
34use anyhow::{Context, Result, bail, ensure};
35use chrono::Utc;
36
37use mj_core::config::{
38    Config, ProjectBundle, ProjectRepository, TargetTemplate, atomic_write, container_size_host,
39    data_dir, is_bare_project_target, mount_history_host,
40};
41
42use crate::import::{
43    RepositoryIdentity, bundle_matches, configured_bundle_for_local, configured_bundle_for_origin,
44    setup_style_id,
45};
46use crate::setup::github_repository_from_origin;
47
48const CONFIG_RENAME_JOURNAL: &str = "config-rename.json";
49
50#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
51#[serde(rename_all = "snake_case")]
52enum ConfigRenameKind {
53    Profile,
54    Target,
55}
56
57#[derive(Debug, serde::Serialize, serde::Deserialize)]
58#[serde(deny_unknown_fields)]
59struct ConfigRenameJournal {
60    kind: ConfigRenameKind,
61    old_id: String,
62    new_id: String,
63}
64use mj_core::state::{
65    HostContainerSize, SessionRecord, SessionResourceAllocation, SessionState, State,
66    new_session_id, normalize_session_title,
67};
68
69use crate::targets::{
70    self, AdditionalMount, CommandExecutor, CommandOutput, CommandSpec, SshTarget,
71};
72
73pub(crate) use backend::controller_github_token;
74pub use backend::image_refresh_plan;
75use backend::validate_resource_allocation;
76pub(crate) use backend::{LocalEngineReadiness, local_engine_readiness};
77pub use mbx::preview_build_cache;
78pub(crate) use mbx::{DoctorHostMbxStatus, MBX_VERSION, doctor_host_mbx};
79use provisioning::apply_failed_new_session_rollback;
80pub(crate) use worker_binary::refresh_target_worker_binary_if_stale;
81pub(crate) use worktree::path_exists_on_managed_target;
82
83pub use checkpoint::{
84    CheckpointArtifact, CheckpointDeferred, IdleWorkspaceLease, SessionExportLayout,
85    checkpoint_was_deferred, reconcile_managed_checkpoint_archives,
86};
87pub use lifecycle::{BranchDisposition, CheckoutDisposition, has_nothing_to_checkpoint};
88pub use recovery_scan::{RecoveryCandidate, RecoveryScan};
89pub use resume::{
90    ResumeRepositorySourceMismatch, ResumeRepositorySourcePreflight, ResumeRepositorySourceReceipt,
91    raw_conversion_preview_for,
92};
93pub use reviewer::reviewer_stager;
94pub use subagents::RegisterSubagentRequest;
95pub use worker_binary::{
96    WorkerBinaryAvailability, native_worker_binary_prerequisite, pin_worker_binary_sources,
97    worker_binary_prerequisite_for_arch,
98};
99pub use worker_restart::WorkerUpgradeOutcome;
100pub use worktree::{ResumePlan, local_project_repository, resume_compatibility};
101
102pub struct Controller {
103    pub config: Config,
104    pub state: State,
105}
106
107/// Machine-wide advisory lock for one controller data store. This prevents a
108/// dashboard, server, or CLI lifecycle command from concurrently acting as a
109/// second controller against the same SQLite state and relay sessions.
110#[derive(Debug)]
111pub struct ControllerStoreGuard {
112    file: File,
113}
114
115impl ControllerStoreGuard {
116    pub fn acquire() -> Result<Self> {
117        let directory = data_dir();
118        Self::acquire_at(&directory)
119    }
120
121    fn acquire_at(directory: &Path) -> Result<Self> {
122        Self::try_acquire_at(directory)?.with_context(|| {
123            format!(
124                "another Mjolnir controller is already using {}; stop it before starting this command",
125                directory.display()
126            )
127        })
128    }
129
130    /// Probe exclusivity without treating an owner that is still exiting as an error.
131    pub fn try_acquire() -> Result<Option<Self>> {
132        Self::try_acquire_at(&data_dir())
133    }
134
135    fn try_acquire_at(directory: &Path) -> Result<Option<Self>> {
136        std::fs::create_dir_all(directory)
137            .with_context(|| format!("create controller data directory {}", directory.display()))?;
138        let path = directory.join("controller.lock");
139        let mut options = OpenOptions::new();
140        options.create(true).read(true).write(true);
141        #[cfg(unix)]
142        {
143            use std::os::unix::fs::OpenOptionsExt;
144            options.mode(0o600);
145        }
146        let file = options
147            .open(&path)
148            .with_context(|| format!("open controller lock {}", path.display()))?;
149        match file.try_lock() {
150            Ok(()) => {}
151            Err(std::fs::TryLockError::WouldBlock) => return Ok(None),
152            Err(std::fs::TryLockError::Error(error)) => {
153                return Err(error)
154                    .with_context(|| format!("lock controller store {}", directory.display()));
155            }
156        }
157        Ok(Some(Self { file }))
158    }
159
160    /// Start the sole production SQLite writer after controller exclusivity
161    /// has been established by this guard.
162    pub fn start_database_writer(&self) -> Result<crate::database::DatabaseWriterOwner> {
163        crate::database::start_database_writer()
164    }
165}
166
167impl Drop for ControllerStoreGuard {
168    fn drop(&mut self) {
169        // Make release explicit. `File` also unlocks on close, but an explicit
170        // unlock keeps same-process handoff deterministic across platforms.
171        let _ = self.file.unlock();
172    }
173}
174
175/// The durable result of creating a quick bundle. The returned config is the
176/// same fresh config that was written, allowing a serving projection to publish
177/// the new bundle before acknowledging the request that created it.
178#[derive(Debug)]
179pub struct QuickBundleCreation {
180    pub config: Config,
181    pub bundle_id: String,
182}
183
184/// Failure stages exposed to a viewer request without exposing the underlying
185/// filesystem/configuration error. The detailed error remains available to
186/// the caller for logs and terminal notices.
187#[derive(Debug)]
188pub enum QuickBundleFailure {
189    InvalidSource(anyhow::Error),
190    Persistence(anyhow::Error),
191}
192
193impl std::fmt::Display for QuickBundleFailure {
194    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
195        match self {
196            Self::InvalidSource(error) => write!(formatter, "invalid repository source: {error}"),
197            Self::Persistence(error) => write!(formatter, "persist quick bundle: {error}"),
198        }
199    }
200}
201
202impl std::error::Error for QuickBundleFailure {
203    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
204        match self {
205            Self::InvalidSource(error) | Self::Persistence(error) => Some(error.root_cause()),
206        }
207    }
208}
209
210/// Create a quick bundle from a local repository or GitHub source and persist
211/// it as one serialized fresh-config transaction. Identical sources reuse the
212/// existing configured bundle, matching the terminal's behavior. The returned
213/// config is the same fresh config that was written, allowing a serving
214/// projection to publish the new bundle before acknowledging the request.
215pub fn create_quick_bundle(
216    source: &str,
217) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
218    let (config, bundle_id) = Config::update(|config| {
219        create_quick_bundle_in_config(config, source)
220            .map_err(|error| anyhow::Error::new(QuickBundleFailure::InvalidSource(error)))
221    })
222    .map_err(|error| {
223        error
224            .downcast::<QuickBundleFailure>()
225            .unwrap_or_else(QuickBundleFailure::Persistence)
226    })?;
227    Ok(QuickBundleCreation { config, bundle_id })
228}
229
230/// Add a quick bundle to an already-loaded config. The helper still performs
231/// the local repository canonicalization/GitHub-source parsing, but callers
232/// that persist a config should use [`create_quick_bundle`] so concurrent saves
233/// cannot clobber one another.
234pub fn create_quick_bundle_in_config(config: &mut Config, source: &str) -> Result<String> {
235    let source = interpret_repository_source(source)?;
236    let existing = match &source.kind {
237        RepositorySourceKind::Local(root) => configured_bundle_for_local(config, root),
238        RepositorySourceKind::Github(repository) => {
239            configured_bundle_for_origin(config, repository)
240        }
241    };
242    if let Some(existing) = existing {
243        return Ok(existing);
244    }
245    let repository_id = setup_style_id(&source.name);
246    let mut bundle_id = repository_id.clone();
247    for suffix in 2_u32.. {
248        if !config.bundles.contains_key(&bundle_id) {
249            break;
250        }
251        bundle_id = format!("{repository_id}-{suffix}");
252    }
253    config.bundles.insert(
254        bundle_id.clone(),
255        ProjectBundle {
256            primary_repo: repository_id.clone(),
257            repositories: vec![source.into_project_repository(repository_id.clone())],
258        },
259    );
260    config.validate()?;
261    Ok(bundle_id)
262}
263
264/// Create one bundle from one or more local repositories or GitHub sources.
265///
266/// All sources are interpreted and checked before the config transaction can
267/// write anything. An existing bundle is reused only when its repository set
268/// and primary repository exactly match the request; this keeps selecting one
269/// repository from a larger bundle from silently changing the wizard's choice.
270pub fn create_bundle_from_sources(
271    sources: &[String],
272) -> std::result::Result<QuickBundleCreation, QuickBundleFailure> {
273    let (config, bundle_id) = Config::update(|config| {
274        create_bundle_from_sources_in_config(config, sources)
275            .map_err(|error| anyhow::Error::new(QuickBundleFailure::InvalidSource(error)))
276    })
277    .map_err(|error| {
278        error
279            .downcast::<QuickBundleFailure>()
280            .unwrap_or_else(QuickBundleFailure::Persistence)
281    })?;
282    Ok(QuickBundleCreation { config, bundle_id })
283}
284
285/// Add a bundle for all `sources` to an already-loaded config. The source
286/// interpretation is shared with the persisted [`create_bundle_from_sources`]
287/// entry point and the legacy quick-bundle helper.
288pub fn create_bundle_from_sources_in_config(
289    config: &mut Config,
290    sources: &[String],
291) -> Result<String> {
292    let sources = sources
293        .iter()
294        .map(|source| interpret_repository_source(source))
295        .collect::<Result<Vec<_>>>()?;
296    if sources.is_empty() {
297        bail!("at least one repository source is required");
298    }
299
300    let mut identities = BTreeSet::new();
301    for source in &sources {
302        if !identities.insert(source.identity()) {
303            bail!("duplicate repository source {:?}", source.display_name);
304        }
305    }
306
307    if let Some(existing) = exact_configured_bundle(config, &sources) {
308        return Ok(existing);
309    }
310
311    // Build and validate a candidate before replacing the caller's config, so
312    // a later validation error cannot leave an in-memory partial mutation.
313    let mut updated = config.clone();
314    let mut used_repository_ids = BTreeSet::new();
315    let mut repositories = Vec::with_capacity(sources.len());
316    for source in sources {
317        let base = setup_style_id(&source.name);
318        let repository_id = unique_id(&base, |candidate| used_repository_ids.contains(candidate));
319        used_repository_ids.insert(repository_id.clone());
320        repositories.push(source.into_project_repository(repository_id));
321    }
322    let primary_repo = repositories
323        .first()
324        .map(|repository| repository.id.clone())
325        .context("at least one repository source is required")?;
326    let bundle_id = unique_id(&primary_repo, |candidate| {
327        updated.bundles.contains_key(candidate)
328    });
329    updated.bundles.insert(
330        bundle_id.clone(),
331        ProjectBundle {
332            primary_repo,
333            repositories,
334        },
335    );
336    updated.validate()?;
337    *config = updated;
338    Ok(bundle_id)
339}
340
341#[derive(Debug, Clone)]
342enum RepositorySourceKind {
343    Github(crate::setup::GithubRepository),
344    Local(PathBuf),
345}
346
347#[derive(Debug, Clone)]
348struct InterpretedRepositorySource {
349    display_name: String,
350    name: String,
351    kind: RepositorySourceKind,
352}
353
354impl InterpretedRepositorySource {
355    fn identity(&self) -> RepositoryIdentity {
356        match &self.kind {
357            RepositorySourceKind::Github(repository) => RepositoryIdentity::Github(
358                repository.owner.to_ascii_lowercase(),
359                repository.repository.to_ascii_lowercase(),
360            ),
361            RepositorySourceKind::Local(root) => RepositoryIdentity::Local(root.clone()),
362        }
363    }
364
365    fn into_project_repository(self, id: String) -> ProjectRepository {
366        let (github, local) = match self.kind {
367            RepositorySourceKind::Github(repository) => (
368                Some(format!("{}/{}", repository.owner, repository.repository)),
369                None,
370            ),
371            RepositorySourceKind::Local(root) => (None, Some(root)),
372        };
373        ProjectRepository {
374            id: id.clone(),
375            github,
376            local,
377            destination: PathBuf::from(id),
378            git_ref: None,
379        }
380    }
381}
382
383/// Interpret a source once, including local Git canonicalization and GitHub
384/// parsing, so all creation paths use exactly the same source semantics.
385fn interpret_repository_source(source: &str) -> Result<InterpretedRepositorySource> {
386    let source = source.trim();
387    if source.is_empty() {
388        bail!("repository source cannot be empty");
389    }
390    let expanded = mj_core::path_input::expand_local(Path::new(source))?;
391    let candidate = expanded.as_path();
392    if candidate.exists() {
393        let root = mj_core::local_git::canonical_repository(candidate)?;
394        let name = root
395            .file_name()
396            .and_then(|name| name.to_str())
397            .context("local repository has no usable directory name")?
398            .to_owned();
399        return Ok(InterpretedRepositorySource {
400            display_name: source.to_owned(),
401            name,
402            kind: RepositorySourceKind::Local(root),
403        });
404    }
405    if candidate.is_absolute() || source.starts_with('.') || source.starts_with('~') {
406        bail!("local repository path {source:?} does not exist");
407    }
408    let repository = github_repository_from_origin(source).context(format!(
409        "{source:?} is not a GitHub owner/repository or URL"
410    ))?;
411    Ok(InterpretedRepositorySource {
412        display_name: source.to_owned(),
413        name: repository.repository.clone(),
414        kind: RepositorySourceKind::Github(repository),
415    })
416}
417
418fn exact_configured_bundle(
419    config: &Config,
420    requested: &[InterpretedRepositorySource],
421) -> Option<String> {
422    let requested_identities = requested
423        .iter()
424        .map(InterpretedRepositorySource::identity)
425        .collect::<BTreeSet<_>>();
426    let primary = requested.first()?.identity();
427    config.bundles.iter().find_map(|(id, bundle)| {
428        if bundle.repositories.len() != requested.len()
429            || bundle
430                .repositories
431                .iter()
432                .any(|repository| repository.git_ref.is_some())
433        {
434            return None;
435        }
436        bundle_matches(bundle, &requested_identities, &primary).then(|| id.clone())
437    })
438}
439
440fn unique_id(base: &str, mut is_used: impl FnMut(&str) -> bool) -> String {
441    if !is_used(base) {
442        return base.to_owned();
443    }
444    for suffix in 2_u32.. {
445        let suffix = format!("-{suffix}");
446        let prefix_len = 64usize.saturating_sub(suffix.len());
447        let prefix = base.chars().take(prefix_len).collect::<String>();
448        let candidate = format!("{prefix}{suffix}");
449        if !is_used(&candidate) {
450            return candidate;
451        }
452    }
453    unreachable!("u32 repository/bundle id suffixes exhausted")
454}
455
456pub struct SessionLaunchOptions {
457    pub create_managed_worktree: Option<bool>,
458    pub launch_base: Option<String>,
459    pub launch_branch: Option<String>,
460    pub mjolnir_subagents: Option<bool>,
461    pub initial_prompt: Option<String>,
462    pub workspace_id: String,
463    pub additional_mounts: Vec<AdditionalMount>,
464    pub resource_allocation: Option<SessionResourceAllocation>,
465    pub project_directory: Option<PathBuf>,
466    pub session_title_override: Option<String>,
467}
468
469pub struct SessionResumeOptions {
470    pub additional_mounts: Option<Vec<AdditionalMount>>,
471    pub resource_allocation: Option<SessionResourceAllocation>,
472    pub discard_queue: bool,
473}
474
475fn selected_host_container_size(
476    template: &TargetTemplate,
477    allocation: Option<&SessionResourceAllocation>,
478) -> Option<(String, HostContainerSize)> {
479    let host = container_size_host(template)?;
480    let SessionResourceAllocation::Container { cpus, memory_bytes } = allocation? else {
481        return None;
482    };
483    Some((
484        host.to_owned(),
485        HostContainerSize {
486            cpus: *cpus,
487            memory_bytes: *memory_bytes,
488        },
489    ))
490}
491
492impl Controller {
493    pub fn load() -> Result<Self> {
494        let config = Config::load()?;
495        let state = crate::database::load_state()?;
496        // Missing session dependencies must not lock users out of the tools
497        // needed to repair them. Operations validate the session they act on.
498        state.validate()?;
499        for session in state.sessions.values() {
500            if let Some(issue) = session.configuration_issue(&config) {
501                tracing::warn!(session_id = %session.id, "{issue}");
502            }
503        }
504        let mut controller = Self { config, state };
505        controller.backfill_target_runtime()?;
506        Ok(controller)
507    }
508
509    pub fn reload(&mut self) -> Result<()> {
510        // Capture legacy settings from the previous config before a manual
511        // edit removes them. Only fill missing metadata, never overwrite it.
512        self.backfill_target_runtime()?;
513        *self = Self::load()?;
514        Ok(())
515    }
516
517    fn backfill_target_runtime(&mut self) -> Result<()> {
518        for session in self.state.sessions.values_mut() {
519            if session.target_runtime.is_some() {
520                continue;
521            }
522            match session.target_runtime_settings(&self.config) {
523                Ok(runtime) => {
524                    let runtime = runtime.into_owned();
525                    session.target_runtime = crate::database::backfill_target_runtime(
526                        &session.id,
527                        &session.target_template_id,
528                        &runtime,
529                    )?;
530                }
531                Err(error) => {
532                    tracing::warn!(session_id = %session.id, %error, "target access needs configuration repair")
533                }
534            }
535        }
536        Ok(())
537    }
538
539    fn persist_session_state(&self, session_id: &str) -> Result<()> {
540        match self.state.sessions.get(session_id) {
541            Some(session) => crate::database::save_lifecycle_session(session),
542            None => crate::database::delete_session(session_id),
543        }
544    }
545
546    fn persist_session_transition_or_restore(
547        &mut self,
548        session_id: &str,
549        previous: &SessionRecord,
550        context: &'static str,
551    ) -> Result<()> {
552        persist_session_record_transition_or_restore(
553            &mut self.state,
554            session_id,
555            previous,
556            context,
557            &crate::database::save_lifecycle_session,
558        )
559    }
560
561    fn restore_prior_session_after_persistence_failure(
562        &mut self,
563        session_id: &str,
564        previous: &SessionRecord,
565        primary: anyhow::Error,
566    ) -> anyhow::Error {
567        restore_session_after_persistence_failure(
568            &mut self.state,
569            session_id,
570            previous,
571            primary,
572            crate::database::save_lifecycle_session,
573        )
574    }
575
576    /// Resolve an entered path on its owning host. Call only from background work.
577    pub fn resolve_input_path(
578        &self,
579        target_id: &str,
580        path: &Path,
581        executor: &impl CommandExecutor,
582    ) -> Result<PathBuf> {
583        let target = self
584            .config
585            .targets
586            .get(target_id)
587            .context("Unknown path target")?;
588        resolve_target_input_path(target, path, executor)
589    }
590
591    /// Verify a mount source on the host where Mjolnir will consume it, and report
592    /// the filesystem reason it must be attached read-only, if there is one.
593    ///
594    /// The probe runs in the same round trip as the existence check so the
595    /// editor learns both answers without a second wait. A probe that cannot
596    /// answer reports no reason: provisioning decides that authoritatively.
597    pub fn validate_mount_source(
598        &self,
599        target_id: &str,
600        source: &Path,
601        executor: &impl CommandExecutor,
602    ) -> Result<Option<String>> {
603        let target = self
604            .config
605            .targets
606            .get(target_id)
607            .with_context(|| format!("unknown target template {target_id:?}"))?;
608        let exists = match target {
609            TargetTemplate::LocalPodman { .. }
610            | TargetTemplate::LocalDocker { .. }
611            | TargetTemplate::AppleContainer { .. }
612            | TargetTemplate::AwsEc2 { .. } => std::fs::metadata(source)
613                .map(|metadata| metadata.is_dir())
614                .or_else(|error| {
615                    if error.kind() == std::io::ErrorKind::NotFound {
616                        Ok(false)
617                    } else {
618                        Err(error)
619                    }
620                })
621                .with_context(|| format!("inspect resource source {}", source.display()))?,
622            TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
623                targets::ssh_directory_exists(&SshTarget::from(ssh), source, executor)?
624            }
625            TargetTemplate::LocalBare | TargetTemplate::SshBare { .. } => {
626                bail!("resource attachments are unsupported for bare targets")
627            }
628        };
629        ensure!(
630            exists,
631            "source path {} does not exist or is not a directory",
632            source.display()
633        );
634        Ok(self.forced_read_only_reason(target, source, executor))
635    }
636
637    /// The `filesystem (reason)` label for a source the runtime cannot overlay.
638    fn forced_read_only_reason(
639        &self,
640        target: &TargetTemplate,
641        source: &Path,
642        executor: &impl CommandExecutor,
643    ) -> Option<String> {
644        let ssh = match target {
645            TargetTemplate::LocalPodman { .. } | TargetTemplate::LocalDocker { .. } => None,
646            TargetTemplate::SshPodman { ssh, .. } | TargetTemplate::SshDocker { ssh, .. } => {
647                Some(SshTarget::from(ssh))
648            }
649            // Apple Container already mounts read-only, and EC2 copies instead
650            // of mounting, so neither has an overlay to lose.
651            _ => return None,
652        };
653        let filesystem = targets::probe_filesystem_types(
654            ssh.as_ref(),
655            std::slice::from_ref(&source.to_path_buf()),
656            executor,
657        )
658        .map_err(|error| {
659            tracing::debug!(
660                source = %source.display(),
661                error = format!("{error:#}"),
662                "could not probe the filesystem under a mount source"
663            );
664        })
665        .ok()?
666        .pop()?;
667        let reason = targets::overlay_unsupported_filesystem(&filesystem)?;
668        Some(format!("{filesystem} ({reason})"))
669    }
670
671    fn fail_new_session_with_cleanup(
672        &mut self,
673        session_id: &str,
674        error: anyhow::Error,
675        executor: &impl CommandExecutor,
676    ) -> Result<anyhow::Error> {
677        let original = provisioning::note_new_session_launch_failure(session_id, &error);
678        let cleanup_error = self
679            .cleanup_new_session_worktree_after_failure(session_id, executor)
680            .err()
681            .map(|cleanup_error| format!("{cleanup_error:#}"));
682        if let Some(cleanup_error) = &cleanup_error {
683            tracing::warn!(
684                session_id,
685                error = %cleanup_error,
686                "new-session worktree rollback reported a cleanup failure"
687            );
688        }
689        let failure = apply_failed_new_session_rollback(
690            &mut self.state,
691            session_id,
692            &original,
693            cleanup_error,
694        );
695        self.persist_session_state(session_id)?;
696        Ok(failure)
697    }
698
699    pub fn register_session_with_resources(
700        &mut self,
701        profile_id: &str,
702        bundle_id: &str,
703        target_id: &str,
704        title: impl Into<String>,
705        options: SessionLaunchOptions,
706    ) -> Result<String> {
707        let SessionLaunchOptions {
708            create_managed_worktree,
709            launch_base,
710            launch_branch,
711            mjolnir_subagents,
712            initial_prompt,
713            workspace_id,
714            additional_mounts,
715            resource_allocation,
716            project_directory,
717            session_title_override,
718        } = options;
719        let launch_base = match launch_base {
720            Some(base) => {
721                let base = base.trim();
722                if base.is_empty() {
723                    bail!("launch base must not be empty");
724                }
725                Some(base.to_owned())
726            }
727            None => None,
728        };
729        let launch_branch = match launch_branch {
730            Some(branch) => {
731                let branch = branch.trim();
732                ensure!(!branch.is_empty(), "launch branch must not be empty");
733                Some(branch.to_owned())
734            }
735            None => None,
736        };
737        if launch_base.is_some() && create_managed_worktree == Some(false) {
738            bail!("a launch base requires a managed worktree or a bundle session");
739        }
740        let session_title_override = match session_title_override {
741            Some(title) => {
742                Some(normalize_session_title(&title).context("session name cannot be empty")?)
743            }
744            None => None,
745        };
746        let profile = self
747            .config
748            .profiles
749            .get(profile_id)
750            .with_context(|| format!("unknown profile {profile_id:?}"))?;
751        ensure!(profile.enabled, "profile {profile_id:?} is disabled");
752        let template = self
753            .config
754            .targets
755            .get(target_id)
756            .with_context(|| format!("unknown target template {target_id:?}"))?;
757        if create_managed_worktree == Some(true) && !is_bare_project_target(template) {
758            bail!("managed worktree creation requires a bare Git project");
759        }
760        if project_directory.is_some() != is_bare_project_target(template) {
761            bail!("raw project directories require a bare target, and bare targets require one");
762        }
763        if let Some(path) = &project_directory
764            && (!path.is_absolute()
765                || path
766                    .components()
767                    .any(|part| part == std::path::Component::ParentDir))
768        {
769            bail!("bare project directory must be an absolute safe path");
770        }
771        let bundle = project_directory
772            .is_none()
773            .then(|| self.config.bundles.get(bundle_id))
774            .flatten();
775        if project_directory.is_none() && bundle.is_none() {
776            bail!("unknown bundle {bundle_id:?}");
777        }
778        if profile.kind == mj_core::config::HarnessKind::Muse
779            && (!additional_mounts.is_empty()
780                || bundle.is_some_and(|bundle| bundle.repositories.len() > 1))
781        {
782            bail!(
783                "{} ACP supports one workspace root; use a single-repository bundle without attached directories",
784                profile.kind.display_name()
785            );
786        }
787        if let Some(bundle) = bundle {
788            for repository in &bundle.repositories {
789                mj_core::remote_git::resolve_repository(
790                    repository,
791                    &targets::CancellableProcessExecutor::with_timeout(
792                        std::time::Duration::from_secs(15),
793                    ),
794                )
795                .with_context(|| format!("repository {:?}", repository.id))?;
796            }
797        }
798        validate_resource_allocation(template, resource_allocation.as_ref())?;
799        let selected_container_size =
800            selected_host_container_size(template, resource_allocation.as_ref());
801        if !additional_mounts.is_empty() && mount_history_host(template).is_none() {
802            bail!("attached resources are unsupported for this target");
803        }
804        targets::validate_additional_mounts(&additional_mounts)?;
805        let id = new_session_id()?;
806        let now = now();
807        let record = SessionRecord {
808            target_runtime: Some(template.into()),
809            build_cache: None,
810            create_managed_worktree,
811            launch_base,
812            launch_branch,
813            publication: None,
814            mjolnir_subagents,
815            archived: false,
816            container_cpus: None,
817            container_memory: None,
818            // Recorded for every new session, container-backed or not, so a
819            // later move into a container already knows the path its checkout
820            // will occupy. Only sessions that predate per-session container
821            // workspaces leave it unset.
822            container_workspace: Some(targets::new_container_workspace(&id)?),
823            id: id.clone(),
824            workspace_id,
825            title: title.into(),
826            harness_kind: profile.kind,
827            last_profile: profile_id.to_string(),
828            bundle_id: bundle_id.to_string(),
829            project_directory,
830            managed_worktree: None,
831            target_template_id: target_id.to_string(),
832            resource_allocation,
833            additional_mounts: additional_mounts.clone(),
834            state: SessionState::Provisioning,
835            target: None,
836            native_session_id: None,
837            acp_session_title: None,
838            session_title_override,
839            created_at: now.clone(),
840            updated_at: now,
841            viewed_through_event_ordinal: 0,
842            draft_input: initial_prompt.unwrap_or_default(),
843            last_error: None,
844            last_checkpoint_error: None,
845            checkpoint: None,
846        };
847        // Creation authors the whole record, so it writes the whole row. The
848        // record reaches memory only once it is durable: a session this process
849        // alone knows about is one the database can never resume or clean up.
850        if let Some((host, size)) = selected_container_size.as_ref() {
851            crate::database::save_session_with_container_size(&record, host, *size)?;
852        } else {
853            crate::database::save_session(&record)?;
854        }
855        self.state.sessions.insert(id.clone(), record);
856        if let Some((host, size)) = selected_container_size {
857            self.state.remember_container_size(&host, size);
858        }
859        if let Some(host) = mount_history_host(template) {
860            // Mount history only seeds the attach dialog's suggestions. The
861            // session row is already committed, so a failed suggestion write is
862            // reported rather than turned into a failed registration.
863            match crate::database::remember_mount_sources(host, &additional_mounts) {
864                Ok(()) => self.state.remember_mount_sources(host, &additional_mounts),
865                Err(error) => tracing::warn!(
866                    session_id = id,
867                    error = format!("{error:#}"),
868                    "could not remember the attached resource directories for later suggestions"
869                ),
870            }
871        }
872        Ok(id)
873    }
874
875    pub fn rename_session(&mut self, session_id: &str, title: &str) -> Result<String> {
876        let title = normalize_session_title(title).context("session name cannot be empty")?;
877        ensure!(
878            self.state.sessions.contains_key(session_id),
879            "unknown session {session_id}"
880        );
881        let updated_at = now();
882        crate::database::set_session_title_override(session_id, &title, &updated_at)?;
883        let record = self
884            .state
885            .sessions
886            .get_mut(session_id)
887            .expect("session was checked before updating its title");
888        record.session_title_override = Some(title.clone());
889        record.updated_at = updated_at;
890        Ok(title)
891    }
892
893    pub fn rename_profile_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
894        mj_core::config::validate_id("profile", new_id)?;
895        if old_id == new_id {
896            ensure!(
897                self.config.profiles.contains_key(old_id),
898                "unknown profile {old_id:?}"
899            );
900            return Ok(());
901        }
902        let journal = ConfigRenameJournal {
903            kind: ConfigRenameKind::Profile,
904            old_id: old_id.to_owned(),
905            new_id: new_id.to_owned(),
906        };
907        write_config_rename_journal(&journal)?;
908        let (config, ()) = match Config::update(|config| {
909            ensure!(
910                config.profiles.contains_key(old_id),
911                "unknown profile {old_id:?}"
912            );
913            ensure!(
914                !config.profiles.contains_key(new_id),
915                "profile {new_id:?} already exists"
916            );
917            let profile = config
918                .profiles
919                .remove(old_id)
920                .expect("profile was checked in the transaction");
921            config.profiles.insert(new_id.to_owned(), profile);
922            Ok(())
923        }) {
924            Ok(result) => result,
925            Err(error) => {
926                remove_config_rename_journal()
927                    .context("remove profile rename journal after config save failed")?;
928                return Err(error).context("save renamed profile configuration");
929            }
930        };
931        self.config = config;
932        mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
933        if let Err(error) = crate::database::rename_profile_references(old_id, new_id) {
934            let restore = Config::update(|config| {
935                let profile = config
936                    .profiles
937                    .remove(new_id)
938                    .with_context(|| format!("renamed profile {new_id:?} is missing"))?;
939                ensure!(
940                    !config.profiles.contains_key(old_id),
941                    "cannot restore profile rename: both {old_id:?} and {new_id:?} exist"
942                );
943                config.profiles.insert(old_id.to_owned(), profile);
944                Ok(())
945            });
946            let restored = match restore {
947                Ok((config, ())) => config,
948                Err(restore_error) => {
949                    return Err(error).context(format!(
950                        "rename profile references; additionally failed to restore config: {restore_error:#}"
951                    ));
952                }
953            };
954            self.config = restored;
955            if let Err(restore_error) = remove_config_rename_journal() {
956                return Err(error).context(format!(
957                    "rename profile references; additionally failed to remove rename journal: {restore_error:#}"
958                ));
959            }
960            return Err(error).context("rename profile references");
961        }
962        for session in self.state.sessions.values_mut() {
963            if session.last_profile == old_id {
964                session.last_profile = new_id.to_owned();
965            }
966        }
967        remove_config_rename_journal()?;
968        Ok(())
969    }
970
971    pub fn rename_target_id(&mut self, old_id: &str, new_id: &str) -> Result<()> {
972        mj_core::config::validate_id("target template", new_id)?;
973        if old_id == new_id {
974            ensure!(
975                self.config.targets.contains_key(old_id),
976                "unknown target {old_id:?}"
977            );
978            return Ok(());
979        }
980        let journal = ConfigRenameJournal {
981            kind: ConfigRenameKind::Target,
982            old_id: old_id.to_owned(),
983            new_id: new_id.to_owned(),
984        };
985        write_config_rename_journal(&journal)?;
986        let (config, ()) = match Config::update(|config| {
987            ensure!(
988                config.targets.contains_key(old_id),
989                "unknown target {old_id:?}"
990            );
991            ensure!(
992                !config.targets.contains_key(new_id),
993                "target {new_id:?} already exists"
994            );
995            let target = config
996                .targets
997                .remove(old_id)
998                .expect("target was checked in the transaction");
999            config.targets.insert(new_id.to_owned(), target);
1000            Ok(())
1001        }) {
1002            Ok(result) => result,
1003            Err(error) => {
1004                remove_config_rename_journal()
1005                    .context("remove target rename journal after config save failed")?;
1006                return Err(error).context("save renamed target configuration");
1007            }
1008        };
1009        self.config = config;
1010        mj_core::test_hooks::reach_test_hook("config_replacement_before_reference_migration")?;
1011        if let Err(error) = crate::database::rename_target_references(old_id, new_id) {
1012            let restore = Config::update(|config| {
1013                let target = config
1014                    .targets
1015                    .remove(new_id)
1016                    .with_context(|| format!("renamed target {new_id:?} is missing"))?;
1017                ensure!(
1018                    !config.targets.contains_key(old_id),
1019                    "cannot restore target rename: both {old_id:?} and {new_id:?} exist"
1020                );
1021                config.targets.insert(old_id.to_owned(), target);
1022                Ok(())
1023            });
1024            let restored = match restore {
1025                Ok((config, ())) => config,
1026                Err(restore_error) => {
1027                    return Err(error).context(format!(
1028                        "rename target references; additionally failed to restore config: {restore_error:#}"
1029                    ));
1030                }
1031            };
1032            self.config = restored;
1033            if let Err(restore_error) = remove_config_rename_journal() {
1034                return Err(error).context(format!(
1035                    "rename target references; additionally failed to remove rename journal: {restore_error:#}"
1036                ));
1037            }
1038            return Err(error).context("rename target references");
1039        }
1040        for session in self.state.sessions.values_mut() {
1041            if session.target_template_id == old_id {
1042                session.target_template_id = new_id.to_owned();
1043            }
1044        }
1045        remove_config_rename_journal()?;
1046        Ok(())
1047    }
1048
1049    /// Finish a profile/target id rename interrupted between the atomic config
1050    /// replacement and SQLite transaction. Each step is idempotent, so a
1051    /// second crash leaves the same intent available for the next startup.
1052    pub fn recover_config_id_rename() -> Result<bool> {
1053        let path = config_rename_journal_path();
1054        let body = match fs::read(&path) {
1055            Ok(body) => body,
1056            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(false),
1057            Err(error) => return Err(error).context(format!("read {}", path.display())),
1058        };
1059        let journal: ConfigRenameJournal =
1060            serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1061        match journal.kind {
1062            ConfigRenameKind::Profile => {
1063                Config::update(|config| {
1064                    finish_config_map_rename(
1065                        &mut config.profiles,
1066                        &journal.old_id,
1067                        &journal.new_id,
1068                        "profile",
1069                    )?;
1070                    Ok(())
1071                })?;
1072                crate::database::rename_profile_references(&journal.old_id, &journal.new_id)?;
1073            }
1074            ConfigRenameKind::Target => {
1075                // The file's own entries, not `Config::update`'s view: that one
1076                // adds the standard local targets, so after a built-in id such
1077                // as `localhost` was renamed away, its default would reappear
1078                // beside the new id and read as a rename that cannot finish.
1079                Config::update_to(&mj_core::config::config_path(), |config| {
1080                    finish_config_map_rename(
1081                        &mut config.targets,
1082                        &journal.old_id,
1083                        &journal.new_id,
1084                        "target",
1085                    )?;
1086                    Ok(())
1087                })?;
1088                crate::database::rename_target_references(&journal.old_id, &journal.new_id)?;
1089            }
1090        }
1091        remove_config_rename_journal()?;
1092        Ok(true)
1093    }
1094
1095    /// Record the per-session container size overrides and attached
1096    /// directories. Nothing is applied to a running container: the values are
1097    /// read the next time the session's container is created.
1098    pub fn update_session_container_settings(
1099        &mut self,
1100        session_id: &str,
1101        cpus: Option<String>,
1102        memory: Option<String>,
1103        additional_mounts: Vec<targets::AdditionalMount>,
1104        mount_history: Vec<std::path::PathBuf>,
1105        executor: &impl CommandExecutor,
1106    ) -> Result<()> {
1107        let session = self
1108            .state
1109            .sessions
1110            .get(session_id)
1111            .with_context(|| format!("unknown session {session_id}"))?;
1112        // A directory the runtime cannot find is created by Docker, owned by
1113        // root, when the container is recreated (launch finding J-24). Check
1114        // each newly attached source where the container will read it: on
1115        // this host for a local runtime, over the machine's connection for an
1116        // SSH one. A directory already attached was checked when it was
1117        // added, and a size change must not fail because it has since gone.
1118        for mount in &additional_mounts {
1119            if session
1120                .additional_mounts
1121                .iter()
1122                .any(|attached| attached.source == mount.source)
1123            {
1124                continue;
1125            }
1126            self.validate_mount_source(&session.target_template_id, &mount.source, executor)?;
1127        }
1128        let cpus = cpus.filter(|value| !value.trim().is_empty());
1129        let memory = memory.filter(|value| !value.trim().is_empty());
1130        let updated_at = now();
1131        crate::database::set_session_container_settings(
1132            session_id,
1133            cpus.as_deref(),
1134            memory.as_deref(),
1135            &additional_mounts,
1136            &updated_at,
1137        )?;
1138        if let Some(host) = self
1139            .config
1140            .targets
1141            .get(
1142                &self.state.sessions[session_id]
1143                    .target_template_id
1144                    .to_owned(),
1145            )
1146            .and_then(mj_core::config::mount_history_host)
1147        {
1148            let host = host.to_owned();
1149            // The dialog owns the suggestion list, so forgetting a directory
1150            // there has to survive the mounts being remembered right after.
1151            crate::database::replace_mount_history(&host, &mount_history)?;
1152            crate::database::remember_mount_sources(&host, &additional_mounts)?;
1153            self.state.mount_history.insert(host.clone(), mount_history);
1154            self.state.remember_mount_sources(&host, &additional_mounts);
1155        }
1156        let record = self
1157            .state
1158            .sessions
1159            .get_mut(session_id)
1160            .expect("session was checked before updating its container settings");
1161        record.container_cpus = cpus;
1162        record.container_memory = memory;
1163        record.additional_mounts = additional_mounts;
1164        record.updated_at = updated_at;
1165        Ok(())
1166    }
1167}
1168
1169fn config_rename_journal_path() -> PathBuf {
1170    data_dir().join(CONFIG_RENAME_JOURNAL)
1171}
1172
1173fn write_config_rename_journal(journal: &ConfigRenameJournal) -> Result<()> {
1174    let path = config_rename_journal_path();
1175    let body = serde_json::to_vec(journal).context("serialize config rename journal")?;
1176    atomic_write(&path, &body).with_context(|| format!("write {}", path.display()))
1177}
1178
1179fn remove_config_rename_journal() -> Result<()> {
1180    let path = config_rename_journal_path();
1181    match fs::remove_file(&path) {
1182        Ok(()) => Ok(()),
1183        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
1184        Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
1185    }
1186}
1187
1188fn finish_config_map_rename<T>(
1189    entries: &mut BTreeMap<String, T>,
1190    old_id: &str,
1191    new_id: &str,
1192    kind: &str,
1193) -> Result<()> {
1194    if let Some(entry) = entries.remove(old_id) {
1195        ensure!(
1196            !entries.contains_key(new_id),
1197            "cannot recover {kind} rename: both {old_id:?} and {new_id:?} exist"
1198        );
1199        entries.insert(new_id.to_owned(), entry);
1200    } else {
1201        ensure!(
1202            entries.contains_key(new_id),
1203            "cannot recover {kind} rename: neither {old_id:?} nor {new_id:?} exists"
1204        );
1205    }
1206    Ok(())
1207}
1208
1209/// Whether this profile must run from a private staged copy of its home even on
1210/// a local bare target, where a session would otherwise use the profile home
1211/// directly.
1212///
1213/// A Codex profile with a custom model provider qualifies: Mjolnir generates the
1214/// provider's model catalog for each launch and points the staged `config.toml`
1215/// at it, and it must never write either into the user's own profile home.
1216pub(crate) fn requires_private_profile_home(profile: &mj_core::config::HarnessProfile) -> bool {
1217    profile.codex_provider().ok().flatten().is_some()
1218}
1219
1220/// Whether this session's harness home belongs to the session rather than to
1221/// the profile.
1222///
1223/// Only such a session has anywhere to stage files into. Staging when this is
1224/// false would copy the daemon's staged profile straight over the user's own
1225/// harness configuration, and nothing would ever remove it again: teardown
1226/// removes what [`removable_profile_root`] names, which is nothing here.
1227pub(crate) fn session_owns_profile_home(
1228    locator: &targets::TargetLocator,
1229    session_id: &str,
1230    profile: &mj_core::config::HarnessProfile,
1231) -> bool {
1232    removable_profile_root(locator, session_id, profile).is_some()
1233}
1234
1235/// Whether a Claude session on a local bare target runs from a private staged
1236/// home. It can only do so where `CLAUDE_CONFIG_DIR` is what points Claude at
1237/// that home; on macOS the variable scopes nothing, so a private copy would be
1238/// a home Claude never reads. See
1239/// [`HarnessKind::scopes_home_with_environment`](mj_core::config::HarnessKind::scopes_home_with_environment).
1240fn claude_takes_a_private_home(
1241    profile: &mj_core::config::HarnessProfile,
1242    locator: &targets::TargetLocator,
1243) -> bool {
1244    profile.kind == mj_core::config::HarnessKind::Claude
1245        && profile
1246            .kind
1247            .scopes_home_with_environment(locator.harness_host())
1248}
1249
1250/// Where this session's harness reads and writes its profile inside the target.
1251///
1252/// Every case but one is the per-session root `removable_profile_root` names; a
1253/// profile that runs straight out of the user's own home has no per-session root
1254/// and uses that home. Muse keeps its state in a `muse` subdirectory of the
1255/// root, because its ACP adapter owns the directory it is given.
1256#[cfg(test)]
1257pub(crate) fn target_profile_home_for_test(
1258    locator: &targets::TargetLocator,
1259    session_id: &str,
1260    profile: &mj_core::config::HarnessProfile,
1261) -> String {
1262    target_profile_home(locator, session_id, profile)
1263}
1264
1265fn target_profile_home(
1266    locator: &targets::TargetLocator,
1267    session_id: &str,
1268    profile: &mj_core::config::HarnessProfile,
1269) -> String {
1270    let root = removable_profile_root(locator, session_id, profile)
1271        .unwrap_or_else(|| profile.home.to_string_lossy().into_owned());
1272    if profile.kind == mj_core::config::HarnessKind::Muse {
1273        PathBuf::from(root)
1274            .join("muse")
1275            .to_string_lossy()
1276            .into_owned()
1277    } else {
1278        root
1279    }
1280}
1281
1282/// The per-session profile directory an in-place harness replacement may delete,
1283/// or `None` when the session runs straight out of the user's own profile home.
1284///
1285/// This is the root that [`target_profile_home`] derives its answer from, not
1286/// that answer itself: a Muse session's home is a `muse` subdirectory of a
1287/// per-session root, and the whole root is what belongs to the session.
1288pub(super) fn removable_profile_root(
1289    locator: &targets::TargetLocator,
1290    session_id: &str,
1291    profile: &mj_core::config::HarnessProfile,
1292) -> Option<String> {
1293    match locator {
1294        targets::TargetLocator::LocalBare { worker_root } => {
1295            if profile.kind == mj_core::config::HarnessKind::Muse {
1296                Some(
1297                    mj_core::config::data_dir()
1298                        .join("profiles")
1299                        .join(session_id)
1300                        .to_string_lossy()
1301                        .into_owned(),
1302                )
1303            } else if claude_takes_a_private_home(profile, locator)
1304                || requires_private_profile_home(profile)
1305            {
1306                Some(
1307                    Path::new(worker_root)
1308                        .join("profile")
1309                        .to_string_lossy()
1310                        .into_owned(),
1311                )
1312            } else {
1313                // The session reads and writes the user's own profile home.
1314                // Nothing here belongs to the session, so nothing is removed.
1315                None
1316            }
1317        }
1318        targets::TargetLocator::LocalPodman { .. }
1319        | targets::TargetLocator::LocalDocker { .. }
1320        | targets::TargetLocator::AppleContainer { .. }
1321        | targets::TargetLocator::SshPodman { .. }
1322        | targets::TargetLocator::SshDocker { .. } => {
1323            Some(format!("/var/lib/hel/profiles/{session_id}"))
1324        }
1325        targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. } => {
1326            Some(format!(".local/share/hel/profiles/{session_id}"))
1327        }
1328    }
1329}
1330
1331/// Resolve the login home on the machine that owns an editable path.
1332pub fn resolve_target_input_path(
1333    target: &TargetTemplate,
1334    path: &Path,
1335    executor: &impl CommandExecutor,
1336) -> Result<PathBuf> {
1337    if !mj_core::path_input::needs_home(path)? {
1338        return Ok(path.to_path_buf());
1339    }
1340    let host = cache_host::CacheHost::for_path_target(target)?;
1341    mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1342}
1343
1344/// Resolve the login home on a configured machine, for the Settings screen's
1345/// path fields. A machine, not a runtime, is what owns a home directory.
1346pub fn resolve_machine_input_path(
1347    machine: &mj_core::config::Machine,
1348    path: &Path,
1349    executor: &impl CommandExecutor,
1350) -> Result<PathBuf> {
1351    if !mj_core::path_input::needs_home(path)? {
1352        return Ok(path.to_path_buf());
1353    }
1354    // An EC2 instance does not exist until a session starts, so the only home
1355    // this screen can resolve is this machine's.
1356    let host = cache_host::CacheHost::for_path_machine(machine)?;
1357    mj_core::path_input::expand_home(path, Some(&host.home(executor)?))
1358}
1359
1360fn execute_checked(executor: &impl CommandExecutor, command: CommandSpec) -> Result<CommandOutput> {
1361    let output = executor.execute(&command)?;
1362    if output.status != 0 {
1363        let detail = command_error_detail(&output.stderr);
1364        if detail.is_empty() {
1365            bail!("{} failed with status {}", command.purpose, output.status);
1366        }
1367        bail!("{detail}");
1368    }
1369    Ok(output)
1370}
1371
1372fn command_error_detail(stderr: &[u8]) -> String {
1373    let reported = String::from_utf8_lossy(stderr);
1374    let reported = reported.trim();
1375    let detail = reported
1376        .rsplit_once("\nCaused by:\n")
1377        .map_or(reported, |(_, causes)| causes);
1378    let detail = detail.strip_prefix("Error: ").unwrap_or(detail);
1379    detail
1380        .lines()
1381        .map(|line| line.strip_prefix("    ").unwrap_or(line))
1382        .collect::<Vec<_>>()
1383        .join("\n")
1384        .trim()
1385        .to_owned()
1386}
1387
1388fn now() -> String {
1389    Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1390}
1391
1392fn restore_session_after_persistence_failure(
1393    state: &mut State,
1394    session_id: &str,
1395    previous: &SessionRecord,
1396    primary: anyhow::Error,
1397    persist: impl FnOnce(&SessionRecord) -> Result<()>,
1398) -> anyhow::Error {
1399    state
1400        .sessions
1401        .insert(session_id.to_owned(), previous.clone());
1402    let restored = state
1403        .sessions
1404        .get(session_id)
1405        .expect("restored session record disappeared");
1406    match persist(restored) {
1407        Ok(()) => primary,
1408        Err(error) => primary.context(format!(
1409            "restored prior session state in memory, but failed to persist the rollback: {error:#}"
1410        )),
1411    }
1412}
1413
1414fn persist_session_record_transition_or_restore(
1415    state: &mut State,
1416    session_id: &str,
1417    previous: &SessionRecord,
1418    context: &'static str,
1419    persist: &impl Fn(&SessionRecord) -> Result<()>,
1420) -> Result<()> {
1421    let result = persist(
1422        state
1423            .sessions
1424            .get(session_id)
1425            .expect("checkpoint session disappeared before persistence"),
1426    );
1427    match result {
1428        Ok(()) => Ok(()),
1429        Err(error) => Err(restore_session_after_persistence_failure(
1430            state,
1431            session_id,
1432            previous,
1433            error.context(context),
1434            persist,
1435        )),
1436    }
1437}
1438
1439pub fn config_only_controller(config: Config) -> Controller {
1440    Controller {
1441        config,
1442        state: State::default(),
1443    }
1444}
1445
1446#[cfg(test)]
1447mod tests;