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