1use 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 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
81struct 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 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 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 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 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 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 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 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
433pub(super) enum WorkerRootReset {
441 FreshTarget,
444 InPlace {
448 previous_profile_root: String,
450 },
451}
452
453pub(super) struct RestoreIntoTarget<'a> {
459 pub profile: &'a mj_core::config::HarnessProfile,
461 pub archive: &'a VerifiedResumeArchive,
462 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 pub primary_repository_root_from_conversion: bool,
472 pub native_continuity: bool,
473 pub discard_queued_prompts: bool,
474 pub replay_queue: bool,
476 pub utility_handoff: Option<String>,
479 pub projection_build: Option<tokio::task::JoinHandle<Result<MaterializedSession>>>,
480 pub resume_notices: Vec<String>,
482 pub install_attached_resources: bool,
483 pub worker_root_reset: WorkerRootReset,
484 pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
486}
487
488impl Controller {
489 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 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 restore_repositories,
571 restore_native: native_continuity,
572 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 {
588 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
589 match &worker_root_reset {
590 WorkerRootReset::FreshTarget => {
595 if let Some(command) = targets::clear_relay_state_plan(&backend, session_id)? {
596 execute_checked(syncing, command)?;
597 }
598 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 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 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 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 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 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 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 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 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
832pub(super) struct VerifiedResumeArchive {
839 pub archive_path: PathBuf,
843 pub manifest: mj_checkpoint::archive::ArchiveManifest,
844 pub canonical_session: Arc<CanonicalSessionSnapshot>,
845}
846
847pub(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 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 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 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 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 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 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 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 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 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 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 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 let utility_config = (!native_continuity).then(|| self.config.clone());
1359 let discard_queued_prompts = discard_queue || !same_harness;
1360 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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, ¤t, &self.config)?;
1749 targets::close_plan(&backend, session_id)?
1750 .execute(&CancellableProcessExecutor::with_timeout(
1753 Duration::from_secs(15),
1754 ))
1755 .map(|_| ())
1756 })(),
1757 None => Ok(()),
1758 };
1759 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 (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 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 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 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 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
1871fn 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
1893fn 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 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 target: previous.manifest.target.clone(),
1936 bundle: mj_checkpoint::archive::BundleManifest {
1937 id: previous.manifest.bundle.id.clone(),
1938 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
1957fn 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
1975fn 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
2040fn 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
2071struct 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
2162fn 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
2240fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2244 profile_kind == archived_kind
2245}
2246
2247async 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;