Skip to main content

mj_controller/controller/
resume.rs

1//! Resuming a stopped session onto a profile and target.
2
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5use std::time::{Duration, Instant};
6
7use agent_client_protocol::schema::v1::ContentBlock;
8use anyhow::{Context, Result, bail, ensure};
9use rayon::prelude::*;
10use tokio_util::sync::CancellationToken;
11
12use crate::checkpoint_transfer::restore_command;
13use crate::session_manager::new_command_id;
14use mj_checkpoint::archive::{
15    CanonicalQueuedCommandKind, CanonicalSessionSnapshot, CheckpointRepositoryBundle, SystemGit,
16    checkpoint_bundle_prerequisites, read_checkpoint_repository_bundles, verify_archive_streaming,
17};
18use mj_checkpoint::checkpoint::CheckpointRestoreSpec;
19use mj_core::config::{Config, HarnessKind, ProjectRepository, mount_history_host};
20use mj_core::state::{MaterializedSession, SessionRecord, SessionResourceAllocation, SessionState};
21use mj_transcript::projection::materialized_session_from_canonical;
22
23use crate::targets::{
24    self, AdditionalMount, CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec,
25    ProcessExecutor, ProvisionStage, ProvisionStageGuard,
26};
27use mj_core::relay::RelayCommand;
28
29use super::backend::{backend_locator, controller_github_token, validate_resource_allocation};
30use super::checkpoint::upload_checkpoint_spec;
31use super::provisioning::{
32    ProvisioningFailureDisposition, StagedExecutor, execute_concurrent_lanes,
33    install_attached_resources,
34};
35use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
36use super::worker_binary::{bridge_readiness_stage, start_worker_durably, worker_probe_diagnosis};
37use super::worktree::{
38    PrimaryCheckoutRequirement, ResumeConversion, ResumePlan, apply_raw_to_workspace,
39    apply_workspace_to_raw, cleanup_managed_worktree, create_managed_worktree,
40    managed_worktree_checkout_exists, managed_worktree_target, plan_raw_to_workspace,
41    preserve_retained_managed_worktree_branch, raw_checkout_divergence_notice,
42    raw_checkout_position, raw_checkout_snapshot, raw_conversion_preview, restore_managed_worktree,
43    resume_compatibility, retire_managed_worktree,
44};
45use super::{
46    Controller, SessionResumeOptions, execute_checked, now, selected_host_container_size,
47    target_profile_home,
48};
49
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct ResumeRepositorySourceMismatch {
52    pub session_id: String,
53    pub bundle_id: String,
54    pub repository_id: String,
55    pub missing_commit: String,
56    pub archived_origin: String,
57    pub configured_origin: String,
58}
59
60pub use mj_core::state::ResumeRepositorySourceReceipt;
61
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum ResumeRepositorySourcePreflight {
64    Ready(ResumeRepositorySourceReceipt),
65    RepositoryMoved(ResumeRepositorySourceMismatch),
66    /// The resume converts a local checkout into an isolated workspace. The
67    /// receipt is already valid; the preview is what a person has to confirm
68    /// before the checkout is snapshotted and the session moves off this
69    /// machine.
70    ConvertingRawCheckout {
71        receipt: ResumeRepositorySourceReceipt,
72        preview: Box<mj_core::state::RawConversionPreview>,
73    },
74}
75
76struct ResumeRepositoryBundles {
77    checkpoint_sha256: String,
78    repositories: Vec<CheckpointRepositoryBundle>,
79}
80
81/// A small timing scope for the expensive resume phases. Target commands
82/// already trace their own durations; this covers controller-side work and
83/// lets an operator see where a slow resume spent its wall-clock budget.
84struct ResumePhaseTimer<'a> {
85    session_id: &'a str,
86    phase: &'static str,
87    started: Instant,
88}
89
90impl<'a> ResumePhaseTimer<'a> {
91    fn new(session_id: &'a str, phase: &'static str) -> Self {
92        Self {
93            session_id,
94            phase,
95            started: Instant::now(),
96        }
97    }
98}
99
100impl Drop for ResumePhaseTimer<'_> {
101    fn drop(&mut self) {
102        tracing::debug!(
103            session_id = self.session_id,
104            phase = self.phase,
105            elapsed_ms = self.started.elapsed().as_millis(),
106            "resume phase completed"
107        );
108    }
109}
110
111impl Controller {
112    /// Muse has one workspace root; native restore relocates its metadata.
113    pub(super) fn validate_muse_resume_destination(
114        &self,
115        source: &SessionRecord,
116        destination_harness: HarnessKind,
117        _target_id: &str,
118    ) -> Result<()> {
119        if destination_harness != HarnessKind::Muse {
120            return Ok(());
121        }
122        ensure!(
123            source.project_directory.is_some()
124                || source
125                    .project_bundle(&self.config)
126                    .is_none_or(|bundle| bundle.repositories.len() == 1),
127            "Muse Code ACP supports one workspace root; use a single-repository bundle"
128        );
129        Ok(())
130    }
131    /// Prove that each configured repository source still supplies the commit
132    /// boundary its checkpoint bundle expects, before provisioning anything,
133    /// and describe a local checkout's conversion so a person can confirm it.
134    pub fn preflight_resume_repository_sources(
135        &self,
136        session_id: &str,
137        target_id: &str,
138        executor: &(impl CommandExecutor + Sync),
139    ) -> Result<ResumeRepositorySourcePreflight> {
140        self.preflight_repository_sources(session_id, target_id, true, executor)
141    }
142
143    /// `describe_conversion` buys the conversion preview with a read of the
144    /// checkout and a question to its remote. The resume itself only needs to
145    /// know whether a configured source moved, and the person has already
146    /// confirmed by then, so it asks for the cheap answer.
147    fn preflight_repository_sources(
148        &self,
149        session_id: &str,
150        target_id: &str,
151        describe_conversion: bool,
152        executor: &(impl CommandExecutor + Sync),
153    ) -> Result<ResumeRepositorySourcePreflight> {
154        let session = self
155            .state
156            .sessions
157            .get(session_id)
158            .with_context(|| format!("unknown session {session_id}"))?;
159        let checkpoint = session
160            .checkpoint
161            .as_ref()
162            .context("session has no checkpoint")?;
163        let plan = resume_compatibility(session, &self.config, target_id)
164            .map_err(|reason| anyhow::anyhow!(reason))?;
165        if session.project_directory.is_some() {
166            debug_assert!(matches!(
167                plan,
168                ResumePlan::InPlace | ResumePlan::RawToWorkspace
169            ));
170            // A raw session resumes from its live checkout. Its synthetic
171            // bundle is only a grouping identity and may no longer be in the
172            // config; neither an in-place resume nor a raw-to-workspace
173            // conversion restores repository contents from that bundle.
174            let receipt = ResumeRepositorySourceReceipt {
175                session_id: session_id.to_owned(),
176                bundle_id: session.bundle_id.clone(),
177                checkpoint_sha256: checkpoint.sha256.clone(),
178                repositories: Vec::new(),
179            };
180            // A conversion reads the checkout and its remote so the person
181            // sees what will travel. Planning failures (no network remote, an
182            // unreachable remote, a dirty submodule) are this preflight's
183            // error, which every surface already reports.
184            if describe_conversion && plan == ResumePlan::RawToWorkspace {
185                let preview = raw_conversion_preview_for(session, &self.config, executor)?;
186                return Ok(ResumeRepositorySourcePreflight::ConvertingRawCheckout {
187                    receipt,
188                    preview: Box::new(preview),
189                });
190            }
191            return Ok(ResumeRepositorySourcePreflight::Ready(receipt));
192        }
193        let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
194        self.preflight_verified_repository_sources(
195            session_id,
196            ResumeRepositoryBundles {
197                checkpoint_sha256: checkpoint.sha256.clone(),
198                repositories,
199            },
200            None,
201            plan != ResumePlan::WorkspaceToRaw,
202            executor,
203        )
204    }
205
206    fn preflight_verified_repository_sources(
207        &self,
208        session_id: &str,
209        verified: ResumeRepositoryBundles,
210        skip_repository_id: Option<&str>,
211        use_archived_network_sources: bool,
212        executor: &(impl CommandExecutor + Sync),
213    ) -> Result<ResumeRepositorySourcePreflight> {
214        let session = self
215            .state
216            .sessions
217            .get(session_id)
218            .with_context(|| format!("unknown session {session_id}"))?;
219        ensure!(
220            verified.repositories.iter().all(|repository| {
221                !repository.metadata.origin.starts_with("mj-local:")
222                    && !repository.metadata.origin.starts_with("ext::")
223            }),
224            "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
225        );
226        // The immutable archive supplies network provenance. Local source
227        // paths and later host configuration changes are irrelevant.
228        if use_archived_network_sources
229            && verified
230                .repositories
231                .iter()
232                .all(|repository| repository.metadata.remote_workspace)
233        {
234            return Ok(ResumeRepositorySourcePreflight::Ready(
235                ResumeRepositorySourceReceipt {
236                    session_id: session_id.to_owned(),
237                    bundle_id: session.bundle_id.clone(),
238                    checkpoint_sha256: verified.checkpoint_sha256,
239                    repositories: Vec::new(),
240                },
241            ));
242        }
243        if verified.repositories.is_empty() {
244            return Ok(ResumeRepositorySourcePreflight::Ready(
245                ResumeRepositorySourceReceipt {
246                    session_id: session_id.to_owned(),
247                    bundle_id: session.bundle_id.clone(),
248                    checkpoint_sha256: verified.checkpoint_sha256,
249                    repositories: Vec::new(),
250                },
251            ));
252        }
253        let bundle = session
254            .project_bundle(&self.config)
255            .with_context(|| format!("session bundle {:?} is missing", session.bundle_id))?;
256        let configured = verified
257            .repositories
258            .iter()
259            .map(|archived| {
260                bundle
261                    .repositories
262                    .iter()
263                    .find(|repository| repository.id == archived.metadata.id)
264                    .cloned()
265                    .with_context(|| {
266                        format!(
267                            "session bundle {:?} no longer contains repository {:?}",
268                            session.bundle_id, archived.metadata.id
269                        )
270                    })
271            })
272            .collect::<Result<Vec<_>>>()?;
273        let github_token = configured
274            .iter()
275            .any(|repository| repository.github.is_some())
276            .then(controller_github_token)
277            .flatten();
278        let outcomes = verified
279            .repositories
280            .par_iter()
281            .zip(configured.par_iter())
282            .map(|(archived, configured)| {
283                if skip_repository_id == Some(configured.id.as_str()) {
284                    return Ok(None);
285                }
286                checkpoint_source_missing_commit(
287                    configured,
288                    session
289                        .project
290                        .as_ref()
291                        .and_then(|project| project.network_sources.get(&configured.id)),
292                    archived,
293                    executor,
294                    github_token.as_deref(),
295                )
296                .map(|missing_commit| {
297                    missing_commit.map(|missing_commit| ResumeRepositorySourceMismatch {
298                        session_id: session_id.to_owned(),
299                        bundle_id: session.bundle_id.clone(),
300                        repository_id: configured.id.clone(),
301                        missing_commit,
302                        archived_origin: archived.metadata.origin.clone(),
303                        configured_origin: configured.source_label(),
304                    })
305                })
306            })
307            .collect::<Vec<Result<Option<ResumeRepositorySourceMismatch>>>>();
308        for outcome in outcomes {
309            if let Some(mismatch) = outcome? {
310                return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(mismatch));
311            }
312        }
313        Ok(ResumeRepositorySourcePreflight::Ready(
314            ResumeRepositorySourceReceipt {
315                session_id: session_id.to_owned(),
316                bundle_id: session.bundle_id.clone(),
317                checkpoint_sha256: verified.checkpoint_sha256,
318                repositories: configured,
319            },
320        ))
321    }
322
323    fn repository_source_receipt_is_current(
324        &self,
325        session_id: &str,
326        receipt: &ResumeRepositorySourceReceipt,
327    ) -> bool {
328        let Some(session) = self.state.sessions.get(session_id) else {
329            return false;
330        };
331        if receipt.session_id != session_id
332            || receipt.bundle_id != session.bundle_id
333            || session
334                .checkpoint
335                .as_ref()
336                .map(|checkpoint| &checkpoint.sha256)
337                != Some(&receipt.checkpoint_sha256)
338        {
339            return false;
340        }
341        if receipt.repositories.is_empty() {
342            return true;
343        }
344        let Some(bundle) = session.project_bundle(&self.config) else {
345            return false;
346        };
347        receipt.repositories.iter().all(|expected| {
348            bundle
349                .repositories
350                .iter()
351                .any(|configured| configured == expected)
352        })
353    }
354
355    /// Validate a replacement first, then atomically save it and check the
356    /// remaining sources so multi-repository bundles can report the next moved
357    /// repository without ever provisioning a partial target.
358    pub fn replace_resume_repository_origin(
359        &mut self,
360        session_id: &str,
361        repository_id: &str,
362        replacement: &str,
363        executor: &(impl CommandExecutor + Sync),
364    ) -> Result<ResumeRepositorySourcePreflight> {
365        let session = self
366            .state
367            .sessions
368            .get(session_id)
369            .with_context(|| format!("unknown session {session_id}"))?;
370        let bundle_id = session.bundle_id.clone();
371        let checkpoint = session
372            .checkpoint
373            .as_ref()
374            .context("session has no checkpoint")?;
375        let replacement = replacement_repository_source(repository_id, replacement)?;
376        let replacement_network = session
377            .project
378            .as_ref()
379            .map(|_| mj_core::remote_git::resolve_repository(&replacement, executor))
380            .transpose()?;
381        let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
382        let verified = ResumeRepositoryBundles {
383            checkpoint_sha256: checkpoint.sha256.clone(),
384            repositories,
385        };
386        let archived = verified
387            .repositories
388            .iter()
389            .find(|repository| repository.metadata.id == repository_id)
390            .with_context(|| format!("checkpoint does not contain repository {repository_id:?}"))?;
391        if let Some(missing_commit) = checkpoint_source_missing_commit(
392            &replacement,
393            replacement_network.as_ref(),
394            archived,
395            executor,
396            controller_github_token().as_deref(),
397        )? {
398            return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(
399                ResumeRepositorySourceMismatch {
400                    session_id: session_id.to_owned(),
401                    bundle_id,
402                    repository_id: repository_id.to_owned(),
403                    missing_commit,
404                    archived_origin: archived.metadata.origin.clone(),
405                    configured_origin: replacement.source_label(),
406                },
407            ));
408        }
409        let previous = session.clone();
410        let mut project = session.project.clone();
411        if let Some(project) = &mut project {
412            let repository = project
413                .bundle
414                .repositories
415                .iter_mut()
416                .find(|repository| repository.id == repository_id)
417                .context("accepted project no longer contains the repaired repository")?;
418            repository.github = replacement.github.clone();
419            repository.local = replacement.local.clone();
420            project.network_sources.insert(
421                repository_id.to_owned(),
422                replacement_network
423                    .clone()
424                    .expect("accepted project source resolved above"),
425            );
426        }
427        // A merged project can use different repository IDs. Its destination
428        // and identity identify the authoring entry; the session keeps its IDs.
429        let saved = crate::database::saved_project(&bundle_id)?;
430        let config_repository_id = saved
431            .as_ref()
432            .and_then(|(_, saved)| {
433                let accepted = previous.project.as_ref()?;
434                let original = accepted
435                    .bundle
436                    .repositories
437                    .iter()
438                    .find(|repository| repository.id == repository_id)?;
439                saved
440                    .bundle
441                    .repositories
442                    .iter()
443                    .find(|repository| {
444                        repository.destination == original.destination
445                            && saved.identities.get(&repository.id)
446                                == accepted.identities.get(repository_id)
447                    })
448                    .map(|repository| repository.id.clone())
449            })
450            .unwrap_or_else(|| repository_id.to_owned());
451        let config_bundle_id = saved
452            .as_ref()
453            .map_or(bundle_id.as_str(), |(id, _)| id.as_str());
454        let (config, ()) = Config::update(|config| {
455            if project.is_some() && !config.bundles.contains_key(config_bundle_id) {
456                return Ok(());
457            }
458            let bundle = config
459                .bundles
460                .get_mut(config_bundle_id)
461                .with_context(|| format!("session bundle {bundle_id:?} is missing"))?;
462            let repository = bundle
463                .repositories
464                .iter_mut()
465                .find(|repository| repository.id == config_repository_id)
466                .with_context(|| {
467                    format!(
468                        "session bundle {:?} no longer contains repository {repository_id:?}",
469                        bundle_id
470                    )
471                })?;
472            repository.github = replacement.github.clone();
473            repository.local = replacement.local.clone();
474            Ok(())
475        })?;
476        self.config = config;
477        if let Some(project) = project {
478            self.state
479                .sessions
480                .get_mut(session_id)
481                .expect("session checked above")
482                .project = Some(project);
483            self.persist_session_transition_or_restore(
484                session_id,
485                &previous,
486                "save repaired project source",
487            )?;
488        }
489        self.preflight_verified_repository_sources(
490            session_id,
491            verified,
492            Some(repository_id),
493            false,
494            executor,
495        )
496    }
497}
498
499/// Plan a local checkout's conversion into an isolated workspace and describe
500/// it, without changing anything.
501///
502/// The resume preflight and the browser's resume card both need this answer
503/// before a person confirms, and neither owns a [`Controller`] at that point,
504/// so it takes the record and the configuration directly.
505/// How a restore prepares the worker root it is about to install into.
506pub(super) enum WorkerRootReset {
507    /// Today's behaviour for a freshly provisioned target: clear leftover relay
508    /// state on a bare target, then make sure the worker root exists.
509    FreshTarget,
510    /// The environment is being kept and only the harness replaced. Stop the
511    /// live daemon, clear relay state, unlink the installed worker files, and
512    /// remove the previous per-session profile home. Runs on every locator.
513    InPlace {
514        /// The per-session profile root the replaced harness ran from.
515        previous_profile_root: String,
516    },
517}
518
519/// Everything the restore tail of a resume needs, once the destination target
520/// is ready and the session record already describes the destination.
521///
522/// The head of a resume differs a great deal between a fresh environment and an
523/// in-place harness swap; from here on the two are the same work.
524pub(super) struct RestoreIntoTarget<'a> {
525    /// The destination profile, which decides the harness home inside the target.
526    pub profile: &'a mj_core::config::HarnessProfile,
527    pub archive: &'a VerifiedResumeArchive,
528    /// The archive the target actually restores: the conversion's own archive
529    /// when a resume wrote one, otherwise the verified checkpoint.
530    pub restored_archive: &'a Path,
531    pub resumed_project_directory: Option<PathBuf>,
532    pub resumed_container_workspace: Option<PathBuf>,
533    pub restore_repositories: bool,
534    pub native_continuity: bool,
535    pub discard_queued_prompts: bool,
536    /// Whether the archived queue is resubmitted to a fresh native session.
537    pub replay_queue: bool,
538    /// The compacted handoff a cross-harness restore installs as its first
539    /// prompt context. Required whenever `native_continuity` is false.
540    pub utility_handoff: Option<String>,
541    pub projection_build: Option<tokio::task::JoinHandle<Result<MaterializedSession>>>,
542    /// Conversation lines to record once the destination answers.
543    pub resume_notices: Vec<String>,
544    pub install_attached_resources: bool,
545    pub worker_root_reset: WorkerRootReset,
546    /// A managed checkout to retire once, and only once, the restore succeeded.
547    pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
548}
549
550impl Controller {
551    /// Install the worker, restore the archive into the target, start the
552    /// harness, and report the projection the destination answers with.
553    pub(super) async fn restore_into_target(
554        &mut self,
555        session_id: &str,
556        restore: RestoreIntoTarget<'_>,
557        executor: &(impl CommandExecutor + Sync),
558    ) -> Result<MaterializedSession> {
559        crate::worker_lifecycle::run(session_id, "restore into target", executor, async {
560            let RestoreIntoTarget {
561                profile,
562                archive,
563                restored_archive,
564                resumed_project_directory,
565                resumed_container_workspace,
566                restore_repositories,
567                native_continuity,
568                discard_queued_prompts,
569                replay_queue,
570                utility_handoff,
571                projection_build,
572                mut resume_notices,
573                install_attached_resources: should_install_attached_resources,
574                worker_root_reset,
575                retire_after_ready,
576            } = restore;
577            let archive_manifest = &archive.manifest;
578            let canonical_session = &archive.canonical_session;
579            // Sub-agents the suspend stopped: the model hears about them on its
580            // first prompt, and the person in a conversation line. The list is
581            // kept until the relay has the note, so a resume that fails tells
582            // the next one.
583            let stopped_subagents = crate::database::load_stopped_subagents(session_id)
584                .unwrap_or_else(|error| {
585                    tracing::warn!(
586                        session_id,
587                        error = format!("{error:#}"),
588                        "could not read the sub-agents this session's suspend stopped"
589                    );
590                    Vec::new()
591                });
592            let stopped_subagents_context =
593                mj_core::subagent::stopped_subagents_prompt_context(&stopped_subagents);
594            resume_notices.extend(mj_core::subagent::stopped_subagents_notice(
595                &stopped_subagents,
596            ));
597            let (backend, worker_root) = self.worker_placement(session_id)?;
598            let harness_home = target_profile_home(&backend, session_id, profile);
599            let workspace_root = if let Some(project_directory) = &resumed_project_directory {
600                project_directory
601                    .parent()
602                    .context("bare project directory has no parent")?
603                    .to_string_lossy()
604                    .into_owned()
605            } else {
606                super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
607            };
608            let target_path = |path: &str| match &backend {
609                targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
610                    if !path.starts_with('/') =>
611                {
612                    PathBuf::from(format!("~/{path}"))
613                }
614                _ => PathBuf::from(path),
615            };
616            let remote_archive = format!("{worker_root}/restore.hel.zip");
617            let remote_spec = format!("{worker_root}/restore-spec.json");
618            use mj_checkpoint::checkpoint::QueueRestorePolicy;
619            let move_admission =
620                crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
621                    operation.phase == mj_core::state::MovePhase::ResumingDestination
622                        && !operation.queue_admission_started
623                        && operation.queue == mj_core::state::ResumeQueueDisposition::Start
624                });
625            let queue_policy = if move_admission || (!native_continuity && replay_queue) {
626                QueueRestorePolicy::Defer
627            } else if discard_queued_prompts {
628                QueueRestorePolicy::Discard
629            } else {
630                QueueRestorePolicy::Restore
631            };
632            let restore = CheckpointRestoreSpec {
633                archive_path: restore_archive_path(
634                    &backend,
635                    restored_archive,
636                    &target_path(&remote_archive),
637                ),
638                workspace_root: target_path(&workspace_root),
639                relay_root: target_path(&worker_root),
640                harness_home: target_path(&harness_home),
641                // A local checkout converting into a workspace arrives as a
642                // fresh clone of its own remote, and the conversion archive
643                // carries the commits, dirty files, and branch that go over it.
644                // An in-place managed checkout recreated from its retained
645                // branch still needs the archive's dirty state.
646                restore_repositories,
647                restore_native: native_continuity,
648                // A session with a project directory launches its harness there,
649                // spelled exactly as recorded (the worker launch configuration's
650                // `cwd`), so the restored harness session is keyed by that same
651                // text. Rebuilding it from the archive's repository layout loses a
652                // trailing separator, which Grok Build keys by, and cannot name
653                // the checkout a move put the session on. A session in a target
654                // workspace has no project directory, and the archive's layout
655                // names its directory under `/workspace`.
656                primary_repository_root: resumed_project_directory
657                    .as_ref()
658                    .map(|directory| target_path(&directory.to_string_lossy())),
659                queue_policy,
660            };
661            // Prepare the worker root before the worker binary is installed:
662            // a surviving daemon still holds the old binary open, and the
663            // install would land on a running executable.
664            {
665                let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
666                match &worker_root_reset {
667                    // A bare target keeps the closed session's worker root on
668                    // the host. Stop anything still writing there and clear the
669                    // leftover relay state, or the restore's seed loses to a
670                    // stale snapshot whose frontier no journal can support.
671                    WorkerRootReset::FreshTarget => {
672                        if let Some(command) =
673                            targets::clear_relay_state_plan(&backend, session_id)?
674                        {
675                            execute_checked(syncing, command)?;
676                        }
677                        // Both lanes below write into the worker root, so it
678                        // exists first.
679                        execute_checked(
680                            syncing,
681                            targets::command_on_locator(
682                                &backend,
683                                session_id,
684                                vec!["mkdir".into(), "-p".into(), worker_root.clone()],
685                                "create the session worker root",
686                            )?,
687                        )?;
688                    }
689                    // The environment survives this restore, so the old
690                    // harness has to be taken out of it: its daemon, its relay
691                    // state, the installed worker files, and its profile home.
692                    // The same command recreates the worker root.
693                    WorkerRootReset::InPlace {
694                        previous_profile_root,
695                    } => {
696                        execute_checked(
697                            syncing,
698                            targets::in_place_worker_reset_plan(
699                                &backend,
700                                session_id,
701                                previous_profile_root,
702                            )?,
703                        )?;
704                    }
705                }
706            }
707            let staging = tempfile::tempdir().context("create restore staging")?;
708            let local_spec = staging.path().join("restore-spec.json");
709            std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
710            // Two independent lanes into the target. The checkpoint transfer
711            // needs nothing from the worker install, and the worker install
712            // is independent of archive upload, so both run concurrently.
713            let controller = &*self;
714            let backend_ref = &backend;
715            let worker_root_ref = worker_root.as_str();
716            let local_spec_ref = local_spec.as_path();
717            execute_concurrent_lanes(
718                || {
719                    let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
720                    controller.prepare_worker_files(
721                        session_id,
722                        backend_ref,
723                        worker_root_ref,
724                        syncing,
725                    )?;
726                    super::provisioning::install_inherited_git_settings(
727                        syncing,
728                        backend_ref,
729                        session_id,
730                    )?;
731                    Ok(())
732                },
733                || {
734                    let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
735                    if should_upload_restore_archive(&backend) {
736                        upload_checkpoint_spec(
737                            restoring,
738                            backend_ref,
739                            session_id,
740                            restored_archive,
741                            &remote_archive,
742                        )?;
743                    }
744                    upload_checkpoint_spec(
745                        restoring,
746                        backend_ref,
747                        session_id,
748                        local_spec_ref,
749                        &remote_spec,
750                    )
751                },
752            )?;
753            {
754                let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
755                execute_checked(
756                    restoring,
757                    restore_command(&backend, session_id, &remote_spec)?,
758                )?;
759            }
760            if should_install_attached_resources {
761                let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
762                install_attached_resources(
763                    &self.state,
764                    session_id,
765                    &backend,
766                    &worker_root,
767                    syncing,
768                )?;
769            }
770            match projection_build {
771                Some(build) => {
772                    let mut restored_projection = build
773                        .await
774                        .context("rebuild the restored projection")?
775                        .context("rebuild the restored projection")?;
776                    if discard_queued_prompts {
777                        restored_projection.queued_prompts.clear();
778                    }
779                    crate::database::save_materialized_session(&restored_projection)?;
780                }
781                // The stored projection already is the archived one. Only the
782                // queue can still need changing.
783                None if discard_queued_prompts => {
784                    crate::database::replace_materialized_queued_prompts(session_id, &[])?;
785                }
786                None => {}
787            }
788            let readiness_stage = bridge_readiness_stage(profile);
789            let spec = self.reconnect_command(session_id)?;
790            let readiness = async {
791                let mut relay = {
792                    let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
793                    start_worker_durably(
794                        &crate::worker_lifecycle::require(session_id)?,
795                        self.state.sessions[session_id]
796                            .target
797                            .as_ref()
798                            .context("worker start has no durable target")?,
799                        executor,
800                        &backend,
801                        &worker_root,
802                    )?;
803                    connect_started_worker(&spec, session_id, executor, &backend, &worker_root)
804                        .await?
805                };
806                // Installed before the harness is ready, so a queued prompt the
807                // restored relay starts on its own cannot claim the hidden
808                // context first. The relay hands it to one prompt only.
809                if let Some(context) = &stopped_subagents_context {
810                    relay
811                        .install_prompt_context(context.clone())
812                        .await
813                        .context("tell the resumed session which sub-agents its suspend stopped")?;
814                }
815                let native_session_id =
816                    wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
817                let owner = crate::worker_lifecycle::require(session_id)?;
818                crate::database::finish_worker_restart(session_id, owner.operation_id())?;
819                Ok::<_, anyhow::Error>((relay, native_session_id))
820            }
821            .await;
822            let (mut relay, native_session_id) = readiness
823                .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
824            if native_continuity {
825                if !restored_native_session_accepted(
826                    &archive_manifest.session.native_session_id,
827                    &native_session_id,
828                    relay
829                        .operational()
830                        .replaced_unused_native_session_id
831                        .as_deref(),
832                ) {
833                    bail!(
834                        "ACP loaded native session {native_session_id}, expected {}",
835                        archive_manifest.session.native_session_id
836                    );
837                }
838            } else {
839                relay
840                    .install_prompt_context(
841                        utility_handoff
842                            .clone()
843                            .context("a resume into a fresh native session has no handoff")?,
844                    )
845                    .await?;
846                if replay_queue {
847                    for prompt in &canonical_session.queued_prompts {
848                        // A queued configuration change is replayed as itself;
849                        // rebuilding it as a prompt would send `/model x` to
850                        // the agent as text.
851                        let command = match &prompt.kind {
852                            CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
853                                prompt: prompt
854                                    .content
855                                    .iter()
856                                    .cloned()
857                                    .map(serde_json::from_value)
858                                    .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
859                            },
860                            CanonicalQueuedCommandKind::SetConfig { key, value } => {
861                                RelayCommand::SetConfig {
862                                    key: key.clone(),
863                                    value: value.clone(),
864                                }
865                            }
866                        };
867                        relay.submit(prompt.command_id.clone(), command).await?;
868                    }
869                }
870            }
871            // Last, and only once the resume has otherwise succeeded: a failure
872            // before this point rolls the record back to a session whose
873            // worktree still has to be there.
874            if let Some(worktree) = retire_after_ready
875                && let Err(error) = retire_managed_worktree(executor, worktree)
876            {
877                tracing::warn!(
878                    session_id,
879                    worktree = %worktree.worktree_root.display(),
880                    error = format!("{error:#}"),
881                    "could not retire the old managed worktree after resume"
882                );
883                resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
884            }
885            for notice in &resume_notices {
886                let submitted = async {
887                    let command_id = new_command_id("resume-notice")?;
888                    relay
889                        .submit(
890                            command_id,
891                            RelayCommand::RecordNotice {
892                                text: notice.clone(),
893                            },
894                        )
895                        .await
896                }
897                .await;
898                // The conversation line is a courtesy. A relay that refuses it
899                // has not damaged the resume, so report and carry on.
900                if let Err(error) = submitted {
901                    tracing::warn!(
902                        session_id,
903                        error = format!("{error:#}"),
904                        "could not record a resume notice in the conversation"
905                    );
906                }
907            }
908            self.mark_worker_connected(session_id, Some(native_session_id))?;
909            let materialized = relay.sync().await?.materialized;
910            if !stopped_subagents.is_empty() {
911                let delivered = stopped_subagents
912                    .iter()
913                    .map(|child| child.child_session_id.clone())
914                    .collect::<Vec<_>>();
915                // The relay owns the note now. Failing to forget the list only
916                // means a later resume tells the model again.
917                if let Err(error) = crate::database::clear_stopped_subagents(session_id, &delivered)
918                {
919                    tracing::warn!(
920                        session_id,
921                        error = format!("{error:#}"),
922                        "could not clear the stopped sub-agents after telling the resumed session"
923                    );
924                }
925            }
926            Ok(materialized)
927        })
928        .await
929    }
930}
931
932/// One checkpoint archive that has been read end to end and proved to be the
933/// one the session record names.
934///
935/// The canonical snapshot is shared behind an `Arc`: on a long session it is
936/// tens of megabytes, and a resume reads it from several places that would
937/// otherwise hold private copies.
938pub(super) struct VerifiedResumeArchive {
939    /// The archive's absolute, canonical path. A LocalBare worker shares the
940    /// controller's filesystem, so it can consume this file directly instead of
941    /// receiving a second large copy in its worker root.
942    pub archive_path: PathBuf,
943    pub manifest: mj_checkpoint::archive::ArchiveManifest,
944    pub canonical_session: Arc<CanonicalSessionSnapshot>,
945}
946
947/// Read and verify the archive a stopped session will be restored from.
948///
949/// Canonicalizing first keeps one exact absolute path for the restore. The
950/// digest and session id must both match the record, and a legacy host-bridge
951/// archive is refused here rather than part way through a restore.
952pub(super) fn verify_resume_checkpoint(
953    session_id: &str,
954    checkpoint: &mj_core::state::CheckpointMetadata,
955) -> Result<VerifiedResumeArchive> {
956    let archive_path = {
957        let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
958        checkpoint.archive_path.canonicalize().with_context(|| {
959            format!(
960                "resolve checkpoint archive {}",
961                checkpoint.archive_path.display()
962            )
963        })?
964    };
965    ensure!(
966        archive_path.is_absolute() && archive_path.is_file(),
967        "checkpoint archive path is not an absolute regular file: {}",
968        archive_path.display()
969    );
970    let mj_checkpoint::archive::VerifiedArchiveMetadata {
971        manifest,
972        canonical_session,
973        archive_sha256,
974    } = {
975        let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
976        verify_archive_streaming(&archive_path)?
977    };
978    if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
979        bail!("persisted checkpoint verification failed");
980    }
981    ensure!(
982        manifest.repositories.iter().all(|repository| {
983            !repository.metadata.origin.starts_with("mj-local:")
984                && !repository.metadata.origin.starts_with("ext::")
985        }),
986        "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
987    );
988    Ok(VerifiedResumeArchive {
989        archive_path,
990        manifest,
991        canonical_session: Arc::new(canonical_session),
992    })
993}
994
995pub fn raw_conversion_preview_for(
996    session: &SessionRecord,
997    config: &Config,
998    executor: &(impl CommandExecutor + Sync),
999) -> Result<mj_core::state::RawConversionPreview> {
1000    let conversion = plan_raw_to_workspace(session, config, executor)?;
1001    raw_conversion_preview(session, &conversion, executor)
1002}
1003
1004fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
1005    let replacement = replacement.trim();
1006    ensure!(!replacement.is_empty(), "enter the repository's new origin");
1007    let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
1008    let path = expanded.as_path();
1009    let (github, local) = if path.is_absolute() {
1010        ensure!(
1011            path.is_dir(),
1012            "local repository {replacement:?} is not a directory"
1013        );
1014        (None, Some(mj_core::local_git::canonical_repository(path)?))
1015    } else {
1016        let github = crate::setup::github_repository_from_origin(replacement)
1017            .context("origin must be a GitHub repository or an absolute local repository path")?;
1018        (
1019            Some(format!("{}/{}", github.owner, github.repository)),
1020            None,
1021        )
1022    };
1023    Ok(ProjectRepository {
1024        id: id.to_owned(),
1025        github,
1026        local,
1027        destination: PathBuf::from(id),
1028        git_ref: None,
1029    })
1030}
1031
1032fn checkpoint_source_missing_commit(
1033    configured: &ProjectRepository,
1034    accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1035    archived: &CheckpointRepositoryBundle,
1036    executor: &impl CommandExecutor,
1037    github_token: Option<&str>,
1038) -> Result<Option<String>> {
1039    let staging = tempfile::tempdir().context("create repository source preflight")?;
1040    let repository = staging.path().join("repository.git");
1041    checked_preflight_git(
1042        executor,
1043        CommandSpec::new(
1044            "git",
1045            [
1046                "init".to_owned(),
1047                "--bare".to_owned(),
1048                "--quiet".to_owned(),
1049                repository.to_string_lossy().into_owned(),
1050            ],
1051        )
1052        .purpose("initialize repository source preflight"),
1053    )?;
1054    let missing = checkpoint_bundle_prerequisites(archived)?;
1055    if missing.is_empty() {
1056        let bundle = staging.path().join("checkpoint.bundle");
1057        std::fs::write(&bundle, &archived.committed_bundle)
1058            .context("write self-contained checkpoint bundle for source preflight")?;
1059        checked_preflight_git(
1060            executor,
1061            checkpoint_bundle_import_command(&repository, &bundle),
1062        )?;
1063        return Ok(None);
1064    }
1065    // The restore clone will obtain the reachable ancestry. This probe only
1066    // needs to establish that the source still serves each boundary object, so
1067    // stop at that object instead of downloading and walking its whole graph.
1068    for commit in missing {
1069        let output = fetch_source_commit(
1070            executor,
1071            &repository,
1072            configured,
1073            accepted_source,
1074            &commit,
1075            github_token,
1076        )?;
1077        if output.status != 0 {
1078            let stderr = String::from_utf8_lossy(&output.stderr);
1079            if source_does_not_have_commit(&stderr) {
1080                return Ok(Some(commit));
1081            }
1082            bail!(
1083                "could not check configured source {:?}: {}",
1084                configured.source_label(),
1085                stderr.trim()
1086            );
1087        }
1088    }
1089    // Do not re-index the bundle against this deliberately shallow probe: the
1090    // shallow marker would make Git report artificial connectivity failures.
1091    // The real restore applies it to the full source clone.
1092    Ok(None)
1093}
1094
1095fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
1096    let mut command = CommandSpec::new(
1097        "git",
1098        [
1099            "-C".to_owned(),
1100            repository.to_string_lossy().into_owned(),
1101            "fetch".to_owned(),
1102            "--no-tags".to_owned(),
1103            bundle.to_string_lossy().into_owned(),
1104            "HEAD".to_owned(),
1105        ],
1106    )
1107    .purpose("validate self-contained checkpoint bundle");
1108    command
1109        .env
1110        .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1111    command
1112        .env
1113        .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1114    command
1115}
1116
1117fn fetch_source_commit(
1118    executor: &impl CommandExecutor,
1119    repository: &Path,
1120    configured: &ProjectRepository,
1121    accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1122    commit: &str,
1123    github_token: Option<&str>,
1124) -> Result<CommandOutput> {
1125    let mut arguments = Vec::new();
1126    let mut token_auth = false;
1127    let mut ssh_transport = false;
1128    let source = if let Some(source) = accepted_source {
1129        source.fetch_url.clone()
1130    } else if let Some(local) = &configured.local {
1131        local.to_string_lossy().into_owned()
1132    } else {
1133        let source = configured
1134            .github
1135            .as_deref()
1136            .context("repository source is missing")?;
1137        let github = crate::setup::github_repository_from_origin(source)
1138            .context("configured repository is not a GitHub source")?;
1139        if github_token.is_some() {
1140            token_auth = true;
1141            arguments.extend([
1142                "-c".to_owned(),
1143                "credential.helper=".to_owned(),
1144                "-c".to_owned(),
1145                "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1146            ]);
1147            format!(
1148                "https://github.com/{}/{}.git",
1149                github.owner, github.repository
1150            )
1151        } else {
1152            ssh_transport = true;
1153            format!("git@github.com:{}/{}.git", github.owner, github.repository)
1154        }
1155    };
1156    if accepted_source.is_some() {
1157        if source.starts_with("https://github.com/") && github_token.is_some() {
1158            token_auth = true;
1159            arguments.extend([
1160                "-c".to_owned(), "credential.helper=".to_owned(), "-c".to_owned(),
1161                "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1162            ]);
1163        }
1164        ssh_transport = source.starts_with("git@") || source.starts_with("ssh://");
1165    }
1166    arguments.extend([
1167        "-C".to_owned(),
1168        repository.to_string_lossy().into_owned(),
1169        "fetch".to_owned(),
1170        "--no-tags".to_owned(),
1171        "--depth=1".to_owned(),
1172        "--filter=blob:none".to_owned(),
1173        source,
1174        commit.to_owned(),
1175    ]);
1176    let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1177    command
1178        .env
1179        .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1180    command
1181        .env
1182        .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1183    if token_auth {
1184        let token = github_token.expect("token authentication requires a GitHub token");
1185        command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1186    }
1187    if ssh_transport {
1188        command.env.insert(
1189            "GIT_SSH_COMMAND".to_owned(),
1190            "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1191                .to_owned(),
1192        );
1193    }
1194    executor.execute(&command)
1195}
1196
1197fn source_does_not_have_commit(stderr: &str) -> bool {
1198    let stderr = stderr.to_ascii_lowercase();
1199    [
1200        "not our ref",
1201        "couldn't find remote ref",
1202        "not a valid object name",
1203        "no such ref was fetched",
1204    ]
1205    .iter()
1206    .any(|needle| stderr.contains(needle))
1207}
1208
1209fn checked_preflight_git(
1210    executor: &impl CommandExecutor,
1211    command: CommandSpec,
1212) -> Result<CommandOutput> {
1213    let output = executor.execute(&command)?;
1214    ensure!(
1215        output.status == 0,
1216        "{}: {}",
1217        command.purpose,
1218        String::from_utf8_lossy(&output.stderr).trim()
1219    );
1220    Ok(output)
1221}
1222
1223impl Controller {
1224    /// Resume a stopped logical session on any configured profile and
1225    /// target. Cross-harness resume restores Git and canonical history, starts
1226    /// a fresh native session, and supplies the prior transcript as its first
1227    /// context turn.
1228    pub async fn resume_session_with_options(
1229        &mut self,
1230        session_id: &str,
1231        profile_id: &str,
1232        target_id: &str,
1233        additional_mounts: Option<Vec<AdditionalMount>>,
1234        resource_allocation: Option<SessionResourceAllocation>,
1235    ) -> Result<MaterializedSession> {
1236        self.resume_session_with_options_and_queue_disposition(
1237            session_id,
1238            profile_id,
1239            target_id,
1240            additional_mounts,
1241            resource_allocation,
1242            false,
1243        )
1244        .await
1245    }
1246
1247    pub async fn resume_session_with_options_and_queue_disposition(
1248        &mut self,
1249        session_id: &str,
1250        profile_id: &str,
1251        target_id: &str,
1252        additional_mounts: Option<Vec<AdditionalMount>>,
1253        resource_allocation: Option<SessionResourceAllocation>,
1254        discard_queue: bool,
1255    ) -> Result<MaterializedSession> {
1256        self.resume_session_controlled(
1257            session_id,
1258            profile_id,
1259            target_id,
1260            SessionResumeOptions {
1261                additional_mounts,
1262                resource_allocation,
1263                discard_queue,
1264            },
1265            &ProcessExecutor,
1266        )
1267        .await
1268    }
1269
1270    pub async fn resume_session_controlled(
1271        &mut self,
1272        session_id: &str,
1273        profile_id: &str,
1274        target_id: &str,
1275        options: SessionResumeOptions,
1276        executor: &(impl CommandExecutor + Sync),
1277    ) -> Result<MaterializedSession> {
1278        self.resume_session_controlled_with_repository_preflight(
1279            session_id, profile_id, target_id, options, None, executor,
1280        )
1281        .await
1282    }
1283
1284    pub async fn resume_session_controlled_with_repository_preflight(
1285        &mut self,
1286        session_id: &str,
1287        profile_id: &str,
1288        target_id: &str,
1289        options: SessionResumeOptions,
1290        repository_preflight: Option<ResumeRepositorySourceReceipt>,
1291        executor: &(impl CommandExecutor + Sync),
1292    ) -> Result<MaterializedSession> {
1293        self.resume_session_with_origin(
1294            session_id,
1295            profile_id,
1296            target_id,
1297            options,
1298            repository_preflight,
1299            None,
1300            executor,
1301        )
1302        .await
1303    }
1304
1305    pub(in crate::controller) async fn resume_session_for_move(
1306        &mut self,
1307        operation: &mj_core::state::MoveOperation,
1308        executor: &(impl CommandExecutor + Sync),
1309    ) -> Result<MaterializedSession> {
1310        self.resume_session_with_origin(
1311            &operation.selection.session_id,
1312            operation
1313                .selection
1314                .profile_id
1315                .as_deref()
1316                .context("Move profile missing")?,
1317            operation
1318                .selection
1319                .target_template_id
1320                .as_deref()
1321                .context("Move target missing")?,
1322            SessionResumeOptions {
1323                additional_mounts: operation.selection.additional_mounts.clone(),
1324                resource_allocation: operation.selection.resource_allocation.clone(),
1325                discard_queue: true,
1326            },
1327            None,
1328            Some(operation),
1329            executor,
1330        )
1331        .await
1332    }
1333
1334    #[allow(clippy::too_many_arguments)]
1335    async fn resume_session_with_origin(
1336        &mut self,
1337        session_id: &str,
1338        profile_id: &str,
1339        target_id: &str,
1340        options: SessionResumeOptions,
1341        repository_preflight: Option<ResumeRepositorySourceReceipt>,
1342        move_operation: Option<&mj_core::state::MoveOperation>,
1343        executor: &(impl CommandExecutor + Sync),
1344    ) -> Result<MaterializedSession> {
1345        crate::worker_lifecycle::run(session_id, "resume session with origin", executor, async {
1346            crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1347        let transferring_workspace = move_operation.is_some();
1348        if let Some(operation) = crate::database::load_move_operation(session_id)? {
1349            ensure!(
1350                transferring_workspace || !operation.holds_source_environment(),
1351                "a profile switch retains this environment; retry Move instead of recreating it with Resume"
1352            );
1353        }
1354        let SessionResumeOptions {
1355            additional_mounts,
1356            resource_allocation,
1357            discard_queue,
1358        } = options;
1359        let mut previous = self
1360            .state
1361            .sessions
1362            .get(session_id)
1363            .with_context(|| format!("unknown session {session_id}"))?
1364            .clone();
1365        if !(matches!(
1366            previous.state,
1367            SessionState::Stopped | SessionState::Lost | SessionState::Error
1368        ) || (transferring_workspace && previous.state == SessionState::Closing))
1369        {
1370            bail!("session {session_id} is not stopped, lost, or retryable");
1371        }
1372        let checkpoint = move_operation
1373            .and_then(|op| op.handoff.as_ref())
1374            .or(previous.checkpoint.as_ref())
1375            .context("session has no checkpoint")?;
1376        // A receipt for an isolated destination says nothing about a host
1377        // checkout. Explicit moves to raw execution must check that source.
1378        let moving_to_raw = previous.project_directory.is_none()
1379            && self
1380                .config
1381                .targets
1382                .get(target_id)
1383                .is_some_and(mj_core::config::is_bare_project_target);
1384        if !transferring_workspace
1385            && (moving_to_raw
1386                || !repository_preflight.as_ref().is_some_and(|receipt| {
1387                    self.repository_source_receipt_is_current(session_id, receipt)
1388                }))
1389        {
1390            let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1391            if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1392                self.preflight_repository_sources(session_id, target_id, false, executor)?
1393            {
1394                bail!(
1395                    "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1396                    mismatch.missing_commit,
1397                    mismatch.configured_origin,
1398                    mismatch.repository_id,
1399                    mismatch.archived_origin,
1400                );
1401            }
1402        }
1403        let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1404        let archive_path = verified_archive.archive_path.clone();
1405        let archive_manifest = &verified_archive.manifest;
1406        let canonical_session = Arc::clone(&verified_archive.canonical_session);
1407        let profile = self
1408            .config
1409            .profiles
1410            .get(profile_id)
1411            .with_context(|| format!("unknown profile {profile_id:?}"))?
1412            .clone();
1413        ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1414        let target_template = self
1415            .config
1416            .targets
1417            .get(target_id)
1418            .with_context(|| format!("unknown target template {target_id:?}"))?
1419            .clone();
1420        // Decide the representation before the record changes, so an
1421        // incompatible target fails here instead of during provisioning.
1422        self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1423        ensure!(
1424            profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1425            "Muse Code ACP supports one workspace root; attached directories are unsupported"
1426        );
1427        let plan = resume_compatibility(&previous, &self.config, target_id)
1428            .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1429        // A worktree that stays put is reached with the machine's current ssh
1430        // options; only its location had to match (J-21).
1431        if plan == ResumePlan::InPlace
1432            && let Some(worktree) = previous.managed_worktree.as_mut()
1433            && let Ok(current) = managed_worktree_target(&target_template)
1434        {
1435            worktree.target = current;
1436        }
1437        // A converting resume writes its own archive below, from the host
1438        // checkout's own network remote, and provisioning reads that one. Every
1439        // other isolated resume clones what its stored archive already names.
1440        if !mj_core::config::is_bare_project_target(&target_template)
1441            && plan != ResumePlan::RawToWorkspace
1442            && !transferring_workspace
1443        {
1444            super::network_git::bundle_from_manifest(archive_manifest)?;
1445        }
1446        if plan == ResumePlan::InPlace
1447            && previous.managed_worktree.is_none()
1448            && let Some(project_directory) = &previous.project_directory
1449        {
1450            self.validate_project_directory(target_id, project_directory, executor)
1451                .context("raw project is unavailable for resume")?;
1452        }
1453        let conversion = match plan {
1454            ResumePlan::InPlace => None,
1455            ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(
1456                plan_raw_to_workspace(&previous, &self.config, executor)
1457                    .context("prepare the raw checkout for its new target")?,
1458            )),
1459            ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1460                self.plan_workspace_to_raw(&previous, target_id, executor)
1461                    .context("prepare a checkout for this session")?,
1462            )),
1463        };
1464        let resource_allocation =
1465            resource_allocation.or_else(|| previous.resource_allocation.clone());
1466        let additional_mounts =
1467            additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1468        validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1469        let selected_container_size =
1470            selected_host_container_size(&target_template, resource_allocation.as_ref());
1471        if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1472            bail!("attached resources are unsupported for this target");
1473        }
1474        targets::validate_additional_mounts(&additional_mounts)?;
1475        let history_host = mount_history_host(&target_template);
1476        let history_mounts = additional_mounts.clone();
1477        if !transferring_workspace
1478            && previous.state == SessionState::Error
1479            && let Some(locator) = &previous.target
1480        {
1481            let backend = backend_locator(locator, &previous, &self.config)?;
1482            targets::close_plan(&backend, session_id)?
1483                .execute(executor)
1484                .context("clean up target from failed resume")?;
1485        }
1486        let mut resume_notices = Vec::new();
1487        // The raw-to-workspace notice is written where the conversion snapshot
1488        // is taken, because it reports the branch the container arrives on.
1489        if let Some(conversion) = conversion
1490            .as_ref()
1491            .and_then(ResumeConversion::workspace_to_raw)
1492        {
1493            resume_notices.push(format!(
1494                "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1495                previous.target_template_id,
1496                conversion.worktree.worktree_root.display(),
1497                archive_manifest
1498                    .repositories
1499                    .first()
1500                    .and_then(|repository| repository.metadata.branch.as_deref())
1501                    .unwrap_or("a detached head"),
1502                conversion.worktree.branch,
1503            ));
1504        }
1505        let managed_checkout_present = previous
1506            .managed_worktree
1507            .as_ref()
1508            .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1509            .transpose()?
1510            .unwrap_or(true);
1511        // A checkout Mjolnir did not retire remains the truth for a raw session.
1512        // A retired checkout is recreated from the branch and archive below.
1513        if managed_checkout_present && let Some(project_directory) = &previous.project_directory {
1514            match raw_checkout_position(&previous, &self.config, project_directory, executor) {
1515                Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1516                    project_directory,
1517                    archive_manifest
1518                        .repositories
1519                        .first()
1520                        .map(|repository| &repository.metadata),
1521                    &live,
1522                )),
1523                // Informational only: a resume must not fail because Mjolnir could
1524                // not read where the checkout stands.
1525                Err(error) => tracing::warn!(
1526                    session_id,
1527                    error = format!("{error:#}"),
1528                    "could not read the raw checkout position for a resume notice"
1529                ),
1530            }
1531        }
1532        // Resolve the worker before compaction costs minutes and paid model
1533        // requests, including probing the platform of an existing SSH host. A resume
1534        // that could never install a worker fails here rather than after all
1535        // that work has been thrown away.
1536        super::worker_binary::preflight_worker_binary(&target_template, executor)?;
1537        let same_harness = profile.kind == archive_manifest.session.harness_kind;
1538        let native_continuity =
1539            native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1540        let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1541        // Cross-harness compaction is started alongside destination
1542        // provisioning below. Clone only the configuration it reads so the
1543        // controller can continue owning and mutating its session record.
1544        let utility_config = (!native_continuity).then(|| self.config.clone());
1545        let discard_queued_prompts = discard_queue || !same_harness;
1546        // When this controller archived the session, its durable projection is
1547        // already the archive's content. Reading one row decides that; a read
1548        // failure or any mismatch rebuilds as before.
1549        let stored_frontier = crate::database::materialized_event_frontier(session_id)
1550            .unwrap_or_else(|error| {
1551                tracing::warn!(
1552                    session_id,
1553                    error = format!("{error:#}"),
1554                    "could not read the stored projection frontier; rebuilding it from the archive"
1555                );
1556                None
1557            });
1558        let rebuild_projection = projection_rebuild_required(
1559            stored_frontier
1560                .as_ref()
1561                .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1562            canonical_session.event_frontier,
1563            &canonical_session.event_frontier_digest,
1564        );
1565        // Rebuilding the projection is a pure function of the archive and costs
1566        // seconds on a long session. Start it now so it runs while the target is
1567        // being provisioned; its result is awaited where it was consumed
1568        // before, and the writes it feeds have not moved.
1569        //
1570        // A resume that fails before the result is needed drops the handle.
1571        // `spawn_blocking` work cannot be cancelled, so the computation still
1572        // finishes on the blocking pool and its result is discarded; it owns
1573        // nothing but its own inputs, so nothing leaks beyond that CPU.
1574        let projection_build = rebuild_projection.then(|| {
1575            let canonical = Arc::clone(&canonical_session);
1576            let session_id = session_id.to_owned();
1577            tokio::task::spawn_blocking(move || {
1578                materialized_session_from_canonical(session_id, &canonical)
1579            })
1580        });
1581        let github_token = controller_github_token();
1582
1583        // The configuration gains the bundle before the record points at it, so
1584        // no persisted session ever names a bundle that is not there.
1585        if let Some(conversion) = conversion
1586            .as_ref()
1587            .and_then(ResumeConversion::raw_to_workspace)
1588            && let Some(bundle) = &conversion.new_bundle
1589        {
1590            let (config, ()) = Config::update(|config| {
1591                if let Some(existing) = config.bundles.get(&conversion.bundle_id) {
1592                    ensure!(
1593                        existing == bundle,
1594                        "bundle {:?} was configured concurrently with a different definition; retry the resume",
1595                        conversion.bundle_id
1596                    );
1597                } else {
1598                    config
1599                        .bundles
1600                        .insert(conversion.bundle_id.clone(), bundle.clone());
1601                }
1602                Ok(())
1603            })
1604            .context("save the bundle for a converted raw session")?;
1605            self.config = config;
1606        }
1607
1608        // A session that leaves a non-container target for a container one is
1609        // built a container it never had, so it starts using the per-session
1610        // workspace path here if it predates them. A session that already ran
1611        // in a container keeps its path: the harness's recorded working
1612        // directory has to survive the resume.
1613        let moving_into_first_container = self
1614            .config
1615            .targets
1616            .get(target_id)
1617            .is_some_and(mj_core::config::is_container_target)
1618            && !self
1619                .config
1620                .targets
1621                .get(&previous.target_template_id)
1622                .is_some_and(mj_core::config::is_container_target);
1623        let converted_project = conversion
1624            .as_ref()
1625            .and_then(ResumeConversion::raw_to_workspace)
1626            .map(|conversion| {
1627                crate::project_catalog::snapshot(
1628                    self.config
1629                        .bundles
1630                        .get(&conversion.bundle_id)
1631                        .context("converted project is missing")?,
1632                    executor,
1633                    true,
1634                )
1635            })
1636            .transpose()?;
1637        let record = self.state.sessions.get_mut(session_id).unwrap();
1638        if record.container_workspace.is_none() && moving_into_first_container {
1639            record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1640        }
1641        record.harness_kind = profile.kind;
1642        record.last_profile = profile_id.to_string();
1643        record.target_template_id = target_id.to_string();
1644        record.target_runtime = Some(
1645            self.config
1646                .targets
1647                .get(target_id)
1648                .context("resume target disappeared")?
1649                .into(),
1650        );
1651        record.resource_allocation = resource_allocation;
1652        record.additional_mounts = additional_mounts;
1653        record.target = None;
1654        record.native_session_id =
1655            native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1656        record.state = SessionState::Provisioning;
1657        record.updated_at = now();
1658        record.last_error = None;
1659        match &conversion {
1660            Some(ResumeConversion::RawToWorkspace(conversion)) => {
1661                apply_raw_to_workspace(record, conversion);
1662                record.project = converted_project;
1663            }
1664            Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1665                apply_workspace_to_raw(record, conversion);
1666            }
1667            None => {
1668                if let (Some(worktree), Some(refreshed)) = (
1669                    record.managed_worktree.as_mut(),
1670                    previous.managed_worktree.as_ref(),
1671                ) {
1672                    worktree.target = refreshed.target.clone();
1673                }
1674            }
1675        }
1676        let resumed_project_directory = record.project_directory.clone();
1677        let resumed_container_workspace = record.container_workspace.clone();
1678        if let Some(host) = history_host {
1679            self.state.remember_mount_sources(host, &history_mounts);
1680            crate::database::remember_mount_sources(host, &history_mounts)?;
1681        }
1682        // The session's prompt history is filed under its bundle, so a
1683        // conversion moves the history with it before the record is persisted.
1684        if let Some(conversion) = conversion
1685            .as_ref()
1686            .and_then(ResumeConversion::raw_to_workspace)
1687        {
1688            crate::database::rebind_session_bundle(session_id, &conversion.bundle_id)?;
1689        }
1690        // Resume changes target resources, while titles and drafts remain
1691        // owned by their independent writers throughout provisioning.
1692        crate::database::save_resumed_session(
1693            &self.state.sessions[session_id],
1694            selected_container_size
1695                .as_ref()
1696                .map(|(host, size)| (host.as_str(), *size)),
1697        )?;
1698        if let Some((host, size)) = selected_container_size.as_ref() {
1699            self.state.remember_container_size(host, *size);
1700        }
1701
1702        let mut recreated_managed_worktree = false;
1703        // Set once a conversion has written its archive, so the success path can
1704        // retire the archive it replaced and the failure path can remove it.
1705        let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1706        let result = async {
1707            if let Some(worktree) = previous.managed_worktree.as_ref() {
1708                recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1709                if !transferring_workspace && recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1710                    if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1711                        mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1712                            &archive_path, &worktree.worktree_root, &SystemGit,
1713                        )?;
1714                    } else {
1715                        mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1716                            &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1717                        )?;
1718                    }
1719                }
1720            }
1721            // The record already names the worktree, so a failure here rolls
1722            // back through the same path that cleans up a new session's.
1723            if let Some(conversion) = conversion
1724                .as_ref()
1725                .and_then(ResumeConversion::workspace_to_raw)
1726            {
1727                if conversion.reuse_existing_branch {
1728                    let recovery_ref =
1729                        preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1730                    restore_managed_worktree(executor, &conversion.worktree)?;
1731                    resume_notices.push(format!(
1732                        "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1733                    ));
1734                } else {
1735                    create_managed_worktree(
1736                        executor,
1737                        &conversion.worktree,
1738                        None,
1739                        PrimaryCheckoutRequirement::Any,
1740                    )?;
1741                }
1742                if transferring_workspace {
1743                    // Move installs its separately verified Git/file transfer below.
1744                } else if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1745                    mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1746                        &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1747                    )?;
1748                } else {
1749                    mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1750                        &archive_path,
1751                        &conversion.worktree.worktree_root,
1752                        &conversion.worktree.branch,
1753                        &SystemGit,
1754                    )?;
1755                }
1756            }
1757            // A local checkout becomes an isolated workspace by being
1758            // re-snapshotted into a new archive whose provenance is the
1759            // checkout's own network remote. Provisioning clones that remote,
1760            // and the restore below lays this snapshot over the fresh clone.
1761            if let Some(conversion) = conversion
1762                .as_ref()
1763                .and_then(ResumeConversion::raw_to_workspace)
1764                && !transferring_workspace
1765            {
1766                let destination = PathBuf::from(
1767                    previous
1768                        .project_directory
1769                        .as_deref()
1770                        .context("a raw session has no project directory")?
1771                        .file_name()
1772                        .context("a raw project directory cannot be the filesystem root")?,
1773                );
1774                let snapshot = raw_checkout_snapshot(
1775                    &conversion.checkout,
1776                    &conversion.source,
1777                    &destination,
1778                    &SystemGit,
1779                    conversion.retire.as_ref().is_some_and(|checkout| {
1780                        checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1781                    }),
1782                )
1783                .context("snapshot the host checkout for its new target")?;
1784                resume_notices.push(conversion_notice(
1785                    target_id,
1786                    previous
1787                        .project_directory
1788                        .as_deref()
1789                        .unwrap_or(&conversion.checkout),
1790                    snapshot.metadata.branch.as_deref(),
1791                    conversion.retire.as_ref(),
1792                ));
1793                let archives = mj_core::config::sessions_dir();
1794                std::fs::create_dir_all(&archives).with_context(|| {
1795                    format!("create the checkpoint directory {}", archives.display())
1796                })?;
1797                // Named like every other managed archive, so an interrupted
1798                // conversion's file is reconciled away, with `converted`
1799                // marking where it came from.
1800                let output = archives.join(format!(
1801                    "{session_id}-converted-{}-{}.hel.zip",
1802                    previous
1803                        .checkpoint
1804                        .as_ref()
1805                        .map_or(0, |checkpoint| checkpoint.event_frontier),
1806                    new_command_id("archive")?
1807                ));
1808                let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1809                conversion_checkpoint_written = Some(written.clone());
1810                let record = self.state.sessions.get_mut(session_id).unwrap();
1811                record.checkpoint = Some(written);
1812                record.updated_at = now();
1813                // Persist the converted archive before provisioning without
1814                // restoring client-owned fields from this earlier snapshot.
1815                crate::database::save_resumed_session(
1816                    &self.state.sessions[session_id],
1817                    selected_container_size.as_ref().map(|(host, size)| (host.as_str(), *size)),
1818                )?;
1819            }
1820            let utility_handoff = {
1821                let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1822                if let Some(config) = utility_config.as_ref() {
1823                    Some(
1824                        provision_with_cross_harness_handoff(
1825                            self,
1826                            session_id,
1827                            executor,
1828                            github_token.as_deref(),
1829                            config,
1830                            &canonical_session,
1831                            context_bytes,
1832                        )
1833                        .context("prepare the cross-harness destination")?,
1834                    )
1835                } else {
1836                    self.provision_session_with_failure_disposition(
1837                        session_id,
1838                        executor,
1839                        github_token.as_deref(),
1840                        ProvisioningFailureDisposition::Preserve,
1841                    )
1842                    .await?;
1843                    None
1844                }
1845            };
1846            if transferring_workspace {
1847                self.restore_move_workspace(session_id, executor)?;
1848            }
1849            let restore_repositories = !transferring_workspace && ((resumed_project_directory.is_none()
1850                && conversion.is_none())
1851                || plan == ResumePlan::RawToWorkspace
1852                || (recreated_managed_worktree && plan == ResumePlan::InPlace));
1853            // A conversion restores the archive it just wrote, not the raw one
1854            // the session was stopped with.
1855            let restored_archive = conversion_checkpoint_written
1856                .as_ref()
1857                .map_or(archive_path.as_path(), |checkpoint| {
1858                    checkpoint.archive_path.as_path()
1859                });
1860            self.restore_into_target(
1861                session_id,
1862                RestoreIntoTarget {
1863                    profile: &profile,
1864                    archive: &verified_archive,
1865                    restored_archive,
1866                    resumed_project_directory,
1867                    resumed_container_workspace,
1868                    restore_repositories,
1869                    native_continuity,
1870                    discard_queued_prompts,
1871                    replay_queue: !discard_queue,
1872                    utility_handoff,
1873                    projection_build,
1874                    resume_notices,
1875                    install_attached_resources: true,
1876                    worker_root_reset: WorkerRootReset::FreshTarget,
1877                    retire_after_ready: conversion
1878                        .as_ref()
1879                        .filter(|_| !transferring_workspace)
1880                        .and_then(ResumeConversion::raw_to_workspace)
1881                        .and_then(|plan| plan.retire.as_ref()),
1882                },
1883                executor,
1884            )
1885            .await
1886        }
1887        .await;
1888        match result {
1889            Ok(materialized) => {
1890                // The session now runs from the conversion archive, so the raw
1891                // one it replaced can go.
1892                if let Some(written) = &conversion_checkpoint_written {
1893                    super::checkpoint::prune_replaced_checkpoint(
1894                        previous.checkpoint.as_ref(),
1895                        written,
1896                    );
1897                }
1898                Ok(materialized)
1899            }
1900            Err(error) => {
1901                // The rollback puts back a record that names the previous
1902                // archive, so the conversion's archive is nothing but litter.
1903                if let Some(written) = &conversion_checkpoint_written
1904                    && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
1905                    && remove_error.kind() != std::io::ErrorKind::NotFound
1906                {
1907                    tracing::warn!(
1908                        session_id,
1909                        path = %written.archive_path.display(),
1910                        "could not remove the conversion checkpoint after resume failed: {remove_error}"
1911                    );
1912                }
1913                // Put back whatever this resume wrote to the durable
1914                // projection, including the failed worker's own lines.
1915                restore_projection_after_failed_resume(
1916                    session_id,
1917                    &canonical_session,
1918                    discard_queued_prompts,
1919                );
1920                if let Some(operation) = move_operation {
1921                    return Err(self.rollback_move_destination(operation, error, executor)?);
1922                }
1923                Err(self.rollback_failed_resume(
1924                    session_id,
1925                    &previous,
1926                    recreated_managed_worktree,
1927                    error,
1928                    executor,
1929                )?)
1930            }
1931        }
1932
1933        }).await
1934    }
1935
1936    pub(super) fn rollback_failed_resume(
1937        &mut self,
1938        session_id: &str,
1939        previous: &SessionRecord,
1940        recreated_managed_worktree: bool,
1941        error: anyhow::Error,
1942        _executor: &impl CommandExecutor,
1943    ) -> Result<anyhow::Error> {
1944        let current = self
1945            .state
1946            .sessions
1947            .get(session_id)
1948            .with_context(|| format!("unknown session {session_id}"))?
1949            .clone();
1950        let cleanup = match current.target.as_ref() {
1951            Some(locator) => (|| -> Result<()> {
1952                let backend = backend_locator(locator, &current, &self.config)?;
1953                targets::close_plan(&backend, session_id)?
1954                    // Use a fresh executor: cancellation applies to the
1955                    // requested operation, not to its compensating cleanup.
1956                    .execute(&CancellableProcessExecutor::with_timeout(
1957                        Duration::from_secs(15),
1958                    ))
1959                    .map(|_| ())
1960            })(),
1961            None => Ok(()),
1962        };
1963        // A failed target teardown may leave its harness writing. Keep its
1964        // checkout intact until a later retry proves the process is stopped.
1965        let worktree_cleanup = if cleanup.is_err() {
1966            Ok(())
1967        } else {
1968            match (
1969                current.managed_worktree.as_ref(),
1970                previous.managed_worktree.as_ref(),
1971            ) {
1972                (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
1973                    &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1974                    previous,
1975                ),
1976                (Some(current), Some(previous)) if current == previous => Ok(()),
1977                // The failed resume created this worktree and its branch, and
1978                // no harness ever ran in it, so the rollback removes both.
1979                (Some(worktree), _) => cleanup_managed_worktree(
1980                    &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1981                    worktree,
1982                    crate::controller::BranchDisposition::Delete,
1983                ),
1984                (None, _) => Ok(()),
1985            }
1986        };
1987        let cleanup_error = [cleanup, worktree_cleanup]
1988            .into_iter()
1989            .filter_map(Result::err)
1990            .map(|cleanup_error| format!("{cleanup_error:#}"))
1991            .collect::<Vec<_>>()
1992            .join("; ");
1993        if !cleanup_error.is_empty() {
1994            tracing::warn!(
1995                session_id,
1996                error = %cleanup_error,
1997                "resume rollback cleanup reported failures"
1998            );
1999        }
2000        // The session's error is one line; a worker's diagnostic dump goes to
2001        // the log (R8-2).
2002        let original = super::worker_binary::failure_line(&error);
2003        let detail = format!("{error:#}");
2004        if detail != original {
2005            tracing::warn!(session_id, error = %detail, "resume failed");
2006        }
2007        let record = self.state.sessions.get_mut(session_id).unwrap();
2008        let failure = apply_failed_resume_rollback(
2009            record,
2010            previous,
2011            &original,
2012            (!cleanup_error.is_empty()).then_some(cleanup_error),
2013        );
2014        // A conversion filed the session's prompt history under its new bundle.
2015        // The record went back, so the history goes back with it.
2016        if record.bundle_id != current.bundle_id {
2017            let bundle_id = record.bundle_id.clone();
2018            crate::database::rebind_session_bundle(session_id, &bundle_id)?;
2019        }
2020        // Rollback restores only resources owned by this attempt, preserving
2021        // client edits committed while provisioning was in flight.
2022        crate::database::save_resumed_session(&self.state.sessions[session_id], None)?;
2023        Ok(failure)
2024    }
2025}
2026
2027fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
2028    format!(
2029        "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
2030        worktree_root.display(),
2031        worktree_root.display()
2032    )
2033}
2034
2035pub(super) fn apply_failed_resume_rollback(
2036    current: &mut SessionRecord,
2037    previous: &SessionRecord,
2038    original_error: &str,
2039    cleanup_error: Option<String>,
2040) -> anyhow::Error {
2041    match cleanup_error {
2042        None => {
2043            *current = previous.clone();
2044            current.state = SessionState::Stopped;
2045            current.target = None;
2046            current.updated_at = now();
2047            current.last_error = Some(format!("resume failed: {original_error}"));
2048            anyhow::anyhow!(original_error.to_owned())
2049        }
2050        Some(cleanup_error) => {
2051            let failure = format!(
2052                "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
2053            );
2054            // Keep the exact partial target and checkout ownership until
2055            // cleanup succeeds; the harness may still be writing there.
2056            // A container conversion has no new managed host checkout, so
2057            // retain the original host checkout that it has not retired yet.
2058            if current.managed_worktree.is_none() {
2059                current
2060                    .project_directory
2061                    .clone_from(&previous.project_directory);
2062                current
2063                    .managed_worktree
2064                    .clone_from(&previous.managed_worktree);
2065                current.bundle_id.clone_from(&previous.bundle_id);
2066            }
2067            current.state = SessionState::Error;
2068            current.updated_at = now();
2069            current.last_error = Some(format!("resume failed: {failure}"));
2070            anyhow::anyhow!(failure)
2071        }
2072    }
2073}
2074
2075/// What the conversation is told when a local session moves into a target.
2076fn conversion_notice(
2077    target_id: &str,
2078    checkout: &Path,
2079    branch: Option<&str>,
2080    retire: Option<&mj_core::state::ManagedWorktree>,
2081) -> String {
2082    let branch = branch.unwrap_or("a detached head");
2083    match retire {
2084        Some(worktree) => format!(
2085            "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
2086            checkout.display(),
2087            worktree.branch,
2088            worktree.source_repository.display()
2089        ),
2090        None => format!(
2091            "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. The checkout on this machine stays where it is.",
2092            checkout.display()
2093        ),
2094    }
2095}
2096
2097/// The archive a converting resume provisions and restores from: the previous
2098/// archive's session, conversation, and native state, with the host checkout's
2099/// snapshot as its only repository.
2100///
2101/// The snapshot carries network provenance, so this archive is what lets the
2102/// destination clone a real remote and later checkpoint like any other
2103/// isolated session.
2104fn conversion_checkpoint(
2105    previous_archive: &Path,
2106    snapshot: mj_checkpoint::archive::RepositorySnapshot,
2107    output: &Path,
2108) -> Result<mj_core::state::CheckpointMetadata> {
2109    let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
2110        .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
2111    // Native harness state and relay attachments travel byte for byte: a
2112    // conversion replaces repository content and nothing else.
2113    let native_artifacts = previous
2114        .manifest
2115        .payloads
2116        .iter()
2117        .filter_map(|descriptor| match &descriptor.role {
2118            mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
2119                Some((relative_path, descriptor))
2120            }
2121            _ => None,
2122        })
2123        .map(|(relative_path, descriptor)| {
2124            Ok(mj_checkpoint::archive::NativeArtifact {
2125                relative_path: relative_path.clone(),
2126                data: previous.payload(descriptor)?.to_vec(),
2127                mode: descriptor.mode,
2128            })
2129        })
2130        .collect::<Result<Vec<_>>>()?;
2131    let canonical_session = previous.canonical_session()?;
2132    let event_frontier = canonical_session.event_frontier;
2133    let written = mj_checkpoint::archive::write_archive_atomic(
2134        output,
2135        &mj_checkpoint::archive::ArchiveInput {
2136            session: previous.manifest.session.clone(),
2137            // Provenance for a person reading the archive; a restore reads
2138            // nothing from it, so the target it was captured on stands.
2139            target: previous.manifest.target.clone(),
2140            bundle: mj_checkpoint::archive::BundleManifest {
2141                id: previous.manifest.bundle.id.clone(),
2142                // A restore finds the primary repository by this id, and with
2143                // it the working directory the native transcript is rewritten
2144                // to, so it has to name the snapshot.
2145                primary_repository: snapshot.metadata.id.clone(),
2146            },
2147            canonical_session,
2148            native_artifacts,
2149            repositories: vec![snapshot],
2150        },
2151    )
2152    .with_context(|| format!("write the conversion archive {}", output.display()))?;
2153    Ok(mj_core::state::CheckpointMetadata {
2154        archive_path: output.to_path_buf(),
2155        sha256: written.archive_sha256,
2156        created_at: now(),
2157        event_frontier,
2158    })
2159}
2160
2161/// Whether the native session a restored worker opened may stand in for the
2162/// archived one.
2163///
2164/// The archived session is the one the conversation lives in, so a different
2165/// one is refused: resuming in it would silently drop that history. The one
2166/// exception is a worker that says it replaced exactly the archived session
2167/// because the harness had no record of it and this session never used it.
2168/// Claude Code and Codex write nothing for a native session until its first
2169/// prompt, so a session suspended before one (or right after `/clear`) can
2170/// only ever come back this way (R7-5).
2171fn restored_native_session_accepted(
2172    archived: &str,
2173    opened: &str,
2174    replaced_unused: Option<&str>,
2175) -> bool {
2176    opened == archived || replaced_unused == Some(archived)
2177}
2178
2179/// Put the durable projection back to the archived one after a failed resume
2180/// or in-place swap.
2181///
2182/// The attempt can have changed it two ways: a rebuild from the archive, and
2183/// the events of the worker it started, which the controller projects while it
2184/// waits for the harness. Either leaves the stored frontier somewhere other
2185/// than the archive's, so the frontier decides, and the failed attempt leaves
2186/// no lines in the transcript for the next one to add to. A projection still
2187/// at the archived frontier only needs its queue back when the attempt
2188/// discarded it.
2189fn restore_projection_after_failed_resume(
2190    session_id: &str,
2191    canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
2192    discard_queued_prompts: bool,
2193) {
2194    let stored_frontier = crate::database::materialized_event_frontier(session_id)
2195        .unwrap_or_else(|error| {
2196            tracing::warn!(
2197                session_id,
2198                error = format!("{error:#}"),
2199                "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
2200            );
2201            None
2202        });
2203    if projection_rebuild_required(
2204        stored_frontier
2205            .as_ref()
2206            .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
2207        canonical_session.event_frontier,
2208        &canonical_session.event_frontier_digest,
2209    ) {
2210        match materialized_session_from_canonical(session_id, canonical_session) {
2211            Ok(previous_projection) => {
2212                if let Err(restore_error) =
2213                    crate::database::save_materialized_session(&previous_projection)
2214                {
2215                    tracing::error!(
2216                        session_id,
2217                        error = format!("{restore_error:#}"),
2218                        "could not restore the durable projection after a failed resume"
2219                    );
2220                }
2221            }
2222            Err(restore_error) => tracing::error!(
2223                session_id,
2224                error = format!("{restore_error:#}"),
2225                "could not rebuild the durable projection after a failed resume"
2226            ),
2227        }
2228    } else if discard_queued_prompts
2229        && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
2230            session_id,
2231            &mj_transcript::projection::materialized_queued_prompts_from_canonical(
2232                &canonical_session.queued_prompts,
2233            ),
2234        )
2235    {
2236        tracing::error!(
2237            session_id,
2238            error = format!("{restore_error:#}"),
2239            "could not restore queued prompts after a failed resume"
2240        );
2241    }
2242}
2243
2244/// Whether a resume has to rebuild the durable projection from its archive.
2245///
2246/// The projection is a deterministic fold of the relay event chain, so a stored
2247/// projection standing at the archive's frontier *and* carrying the archive's
2248/// frontier digest already holds the archived content: same chain, same
2249/// ordinal, same result. Anything else - no stored row, a different ordinal, a
2250/// different digest, or a frontier that could not be read - rebuilds.
2251fn projection_rebuild_required(
2252    stored: Option<(u64, &str)>,
2253    archive_frontier: u64,
2254    archive_frontier_digest: &str,
2255) -> bool {
2256    stored != Some((archive_frontier, archive_frontier_digest))
2257}
2258
2259fn restore_archive_path(
2260    backend: &targets::TargetLocator,
2261    verified_archive: &Path,
2262    remote_archive: &Path,
2263) -> PathBuf {
2264    if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2265        verified_archive.to_path_buf()
2266    } else {
2267        remote_archive.to_path_buf()
2268    }
2269}
2270
2271fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2272    !matches!(backend, targets::TargetLocator::LocalBare { .. })
2273}
2274
2275/// Provisioning currently performs its target plan synchronously inside an
2276/// async function. Run the network-bound cross-harness handoff on a joined
2277/// side runtime so it can make progress during that plan without borrowing
2278/// the mutable controller or leaving work behind on failure.
2279struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2280    inner: &'a E,
2281    cancellation: CancellationToken,
2282}
2283
2284impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2285    fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2286        if self.cancellation.is_cancelled() {
2287            bail!("operation cancelled while provisioning destination");
2288        }
2289        self.inner.execute(command)
2290    }
2291
2292    fn cancellation_requested(&self) -> bool {
2293        self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2294    }
2295
2296    fn stage_started(&self, stage: ProvisionStage) {
2297        self.inner.stage_started(stage);
2298    }
2299
2300    fn stage_finished(&self, stage: ProvisionStage) {
2301        self.inner.stage_finished(stage);
2302    }
2303
2304    fn notify_notice(&self, notice: &str) {
2305        self.inner.notify_notice(notice);
2306    }
2307
2308    fn execute_with_stdin(
2309        &self,
2310        command: &CommandSpec,
2311        input: &mut (dyn std::io::Read + Send),
2312    ) -> Result<CommandOutput> {
2313        if self.cancellation.is_cancelled() {
2314            bail!("operation cancelled while provisioning destination");
2315        }
2316        self.inner.execute_with_stdin(command, input)
2317    }
2318}
2319
2320fn provision_with_cross_harness_handoff(
2321    controller: &mut Controller,
2322    session_id: &str,
2323    executor: &(impl CommandExecutor + Sync),
2324    github_token: Option<&str>,
2325    config: &Config,
2326    snapshot: &CanonicalSessionSnapshot,
2327    context_bytes: usize,
2328) -> Result<String> {
2329    let (_provision, handoff) = execute_joined_cross_harness_work(
2330        "cross-harness provisioning",
2331        move |cancellation| {
2332            let provision_executor = CrossHarnessProvisionExecutor {
2333                inner: executor,
2334                cancellation,
2335            };
2336            futures::executor::block_on(controller.provision_session_with_failure_disposition(
2337                session_id,
2338                &provision_executor,
2339                github_token,
2340                ProvisioningFailureDisposition::Preserve,
2341            ))
2342        },
2343        "cross-harness handoff",
2344        move |cancellation| {
2345            let runtime = tokio::runtime::Builder::new_current_thread()
2346                .enable_all()
2347                .build()
2348                .context("create cross-harness handoff runtime")?;
2349            runtime.block_on(utility_handoff_while_cancellable(
2350                session_id,
2351                config,
2352                snapshot,
2353                context_bytes,
2354                executor,
2355                cancellation,
2356            ))
2357        },
2358    )?;
2359    ensure!(
2360        !executor.cancellation_requested(),
2361        "operation cancelled while provisioning destination"
2362    );
2363    Ok(handoff)
2364}
2365
2366/// Run the two independent cross-harness lanes together, cancelling and
2367/// joining the peer as soon as either lane fails. The lane order is stable so
2368/// diagnostics do not depend on which worker happened to finish first.
2369fn execute_joined_cross_harness_work<A: Send, B: Send>(
2370    first_name: &'static str,
2371    first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2372    second_name: &'static str,
2373    second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2374) -> Result<(A, B)> {
2375    let cancellation = CancellationToken::new();
2376    std::thread::scope(|scope| {
2377        let first_cancel = cancellation.clone();
2378        let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2379        let second_cancel = cancellation.clone();
2380        let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2381        let mut first_result = None;
2382        let mut second_result = None;
2383
2384        while first_result.is_none() || second_result.is_none() {
2385            if first_result.is_none()
2386                && first_handle
2387                    .as_ref()
2388                    .is_some_and(|handle| handle.is_finished())
2389            {
2390                let handle = first_handle.take().expect("first lane handle present");
2391                first_result = Some(match handle.join() {
2392                    Ok(result) => result,
2393                    Err(panic) => {
2394                        cancellation.cancel();
2395                        Err(anyhow::anyhow!(
2396                            "{first_name} thread panicked: {}",
2397                            targets::command_thread_panic_message(panic.as_ref())
2398                        ))
2399                    }
2400                });
2401                if first_result.as_ref().is_some_and(Result::is_err) {
2402                    cancellation.cancel();
2403                }
2404            }
2405            if second_result.is_none()
2406                && second_handle
2407                    .as_ref()
2408                    .is_some_and(|handle| handle.is_finished())
2409            {
2410                let handle = second_handle.take().expect("second lane handle present");
2411                second_result = Some(match handle.join() {
2412                    Ok(result) => result,
2413                    Err(panic) => {
2414                        cancellation.cancel();
2415                        Err(anyhow::anyhow!(
2416                            "{second_name} thread panicked: {}",
2417                            targets::command_thread_panic_message(panic.as_ref())
2418                        ))
2419                    }
2420                });
2421                if second_result.as_ref().is_some_and(Result::is_err) {
2422                    cancellation.cancel();
2423                }
2424            }
2425            if first_result.is_none() || second_result.is_none() {
2426                std::thread::sleep(Duration::from_millis(10));
2427            }
2428        }
2429
2430        match (
2431            first_result.expect("first lane result received after joined handle"),
2432            second_result.expect("second lane result received after joined handle"),
2433        ) {
2434            (Err(first), Err(second)) => {
2435                Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2436            }
2437            (Err(error), Ok(_)) => Err(error),
2438            (Ok(_), Err(error)) => Err(error),
2439            (Ok(first), Ok(second)) => Ok((first, second)),
2440        }
2441    })
2442}
2443
2444/// Whether the restored native session can carry the conversation into the
2445/// resumed session.
2446///
2447fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2448    profile_kind == archived_kind
2449}
2450
2451/// Discover a utility model and compact the cross-harness handoff while still
2452/// watching for cancellation. Discovery and compaction can both make several
2453/// network requests, so a cancelled resume must not wait them out.
2454async fn utility_handoff_while_cancellable(
2455    session_id: &str,
2456    config: &Config,
2457    snapshot: &CanonicalSessionSnapshot,
2458    context_bytes: usize,
2459    executor: &impl CommandExecutor,
2460    cancellation: CancellationToken,
2461) -> Result<String> {
2462    let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2463    if executor.cancellation_requested() {
2464        bail!("operation cancelled while compacting the cross-harness handoff");
2465    }
2466    let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2467    let cancel = cancellation.child_token();
2468    let operation =
2469        crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2470    tokio::pin!(operation);
2471    loop {
2472        tokio::select! {
2473            context = &mut operation => return context,
2474            _ = cancellation.cancelled() => {
2475                cancel.cancel();
2476                bail!("operation cancelled while compacting the cross-harness handoff");
2477            }
2478            _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2479                if executor.cancellation_requested() {
2480                    cancel.cancel();
2481                    bail!("operation cancelled while compacting the cross-harness handoff");
2482                }
2483            }
2484        }
2485    }
2486}
2487
2488mod in_place;
2489
2490#[cfg(test)]
2491mod tests;