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