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 native_continuity: bool,
469 pub discard_queued_prompts: bool,
470 pub replay_queue: bool,
472 pub utility_handoff: Option<String>,
475 pub projection_build: Option<tokio::task::JoinHandle<Result<MaterializedSession>>>,
476 pub resume_notices: Vec<String>,
478 pub install_attached_resources: bool,
479 pub worker_root_reset: WorkerRootReset,
480 pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
482}
483
484impl Controller {
485 pub(super) async fn restore_into_target(
488 &mut self,
489 session_id: &str,
490 restore: RestoreIntoTarget<'_>,
491 executor: &(impl CommandExecutor + Sync),
492 ) -> Result<MaterializedSession> {
493 let RestoreIntoTarget {
494 profile,
495 archive,
496 restored_archive,
497 resumed_project_directory,
498 resumed_container_workspace,
499 restore_repositories,
500 native_continuity,
501 discard_queued_prompts,
502 replay_queue,
503 utility_handoff,
504 projection_build,
505 mut resume_notices,
506 install_attached_resources: should_install_attached_resources,
507 worker_root_reset,
508 retire_after_ready,
509 } = restore;
510 let archive_manifest = &archive.manifest;
511 let canonical_session = &archive.canonical_session;
512 let stopped_subagents =
517 crate::database::load_stopped_subagents(session_id).unwrap_or_else(|error| {
518 tracing::warn!(
519 session_id,
520 error = format!("{error:#}"),
521 "could not read the sub-agents this session's suspend stopped"
522 );
523 Vec::new()
524 });
525 let stopped_subagents_context =
526 mj_core::subagent::stopped_subagents_prompt_context(&stopped_subagents);
527 resume_notices.extend(mj_core::subagent::stopped_subagents_notice(
528 &stopped_subagents,
529 ));
530 let (backend, worker_root) = self.worker_placement(session_id)?;
531 let harness_home = target_profile_home(&backend, session_id, profile);
532 let workspace_root = if let Some(project_directory) = &resumed_project_directory {
533 project_directory
534 .parent()
535 .context("bare project directory has no parent")?
536 .to_string_lossy()
537 .into_owned()
538 } else {
539 super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
540 };
541 let target_path = |path: &str| match &backend {
542 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
543 if !path.starts_with('/') =>
544 {
545 PathBuf::from(format!("~/{path}"))
546 }
547 _ => PathBuf::from(path),
548 };
549 let remote_archive = format!("{worker_root}/restore.hel.zip");
550 let remote_spec = format!("{worker_root}/restore-spec.json");
551 use mj_checkpoint::checkpoint::QueueRestorePolicy;
552 let move_admission =
553 crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
554 operation.phase == mj_core::state::MovePhase::ResumingDestination
555 && !operation.queue_admission_started
556 && operation.queue == mj_core::state::ResumeQueueDisposition::Start
557 });
558 let queue_policy = if move_admission || (!native_continuity && replay_queue) {
559 QueueRestorePolicy::Defer
560 } else if discard_queued_prompts {
561 QueueRestorePolicy::Discard
562 } else {
563 QueueRestorePolicy::Restore
564 };
565 let restore = CheckpointRestoreSpec {
566 archive_path: restore_archive_path(
567 &backend,
568 restored_archive,
569 &target_path(&remote_archive),
570 ),
571 workspace_root: target_path(&workspace_root),
572 relay_root: target_path(&worker_root),
573 harness_home: target_path(&harness_home),
574 restore_repositories,
580 restore_native: native_continuity,
581 primary_repository_root: resumed_project_directory
590 .as_ref()
591 .map(|directory| target_path(&directory.to_string_lossy())),
592 queue_policy,
593 };
594 {
598 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
599 match &worker_root_reset {
600 WorkerRootReset::FreshTarget => {
605 if let Some(command) = targets::clear_relay_state_plan(&backend, session_id)? {
606 execute_checked(syncing, command)?;
607 }
608 execute_checked(
611 syncing,
612 targets::command_on_locator(
613 &backend,
614 session_id,
615 vec!["mkdir".into(), "-p".into(), worker_root.clone()],
616 "create the session worker root",
617 )?,
618 )?;
619 }
620 WorkerRootReset::InPlace {
625 previous_profile_root,
626 } => {
627 execute_checked(
628 syncing,
629 targets::in_place_worker_reset_plan(
630 &backend,
631 session_id,
632 previous_profile_root,
633 )?,
634 )?;
635 }
636 }
637 }
638 let staging = tempfile::tempdir().context("create restore staging")?;
639 let local_spec = staging.path().join("restore-spec.json");
640 std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
641 let controller = &*self;
645 let backend_ref = &backend;
646 let worker_root_ref = worker_root.as_str();
647 let local_spec_ref = local_spec.as_path();
648 execute_concurrent_lanes(
649 || {
650 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
651 controller.prepare_worker_files(
652 session_id,
653 backend_ref,
654 worker_root_ref,
655 syncing,
656 )?;
657 super::provisioning::install_inherited_git_settings(
658 syncing,
659 backend_ref,
660 session_id,
661 )?;
662 Ok(())
663 },
664 || {
665 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
666 if should_upload_restore_archive(&backend) {
667 upload_checkpoint_spec(
668 restoring,
669 backend_ref,
670 session_id,
671 restored_archive,
672 &remote_archive,
673 )?;
674 }
675 upload_checkpoint_spec(
676 restoring,
677 backend_ref,
678 session_id,
679 local_spec_ref,
680 &remote_spec,
681 )
682 },
683 )?;
684 {
685 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
686 execute_checked(
687 restoring,
688 restore_command(&backend, session_id, &remote_spec)?,
689 )?;
690 }
691 if should_install_attached_resources {
692 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
693 install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
694 }
695 match projection_build {
696 Some(build) => {
697 let mut restored_projection = build
698 .await
699 .context("rebuild the restored projection")?
700 .context("rebuild the restored projection")?;
701 if discard_queued_prompts {
702 restored_projection.queued_prompts.clear();
703 }
704 crate::database::save_materialized_session(&restored_projection)?;
705 }
706 None if discard_queued_prompts => {
709 crate::database::replace_materialized_queued_prompts(session_id, &[])?;
710 }
711 None => {}
712 }
713 let readiness_stage = bridge_readiness_stage(profile);
714 let spec = self.reconnect_command(session_id)?;
715 let readiness = async {
716 let mut relay = {
717 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
718 start_worker(executor, &backend, &worker_root)?;
719 connect_started_worker(&spec, session_id, executor, &backend, &worker_root).await?
720 };
721 if let Some(context) = &stopped_subagents_context {
725 relay
726 .install_prompt_context(context.clone())
727 .await
728 .context("tell the resumed session which sub-agents its suspend stopped")?;
729 }
730 let native_session_id =
731 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
732 Ok::<_, anyhow::Error>((relay, native_session_id))
733 }
734 .await;
735 let (mut relay, native_session_id) = readiness
736 .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
737 if native_continuity {
738 if !restored_native_session_accepted(
739 &archive_manifest.session.native_session_id,
740 &native_session_id,
741 relay
742 .operational()
743 .replaced_unused_native_session_id
744 .as_deref(),
745 ) {
746 bail!(
747 "ACP loaded native session {native_session_id}, expected {}",
748 archive_manifest.session.native_session_id
749 );
750 }
751 } else {
752 relay
753 .install_prompt_context(
754 utility_handoff
755 .clone()
756 .context("a resume into a fresh native session has no handoff")?,
757 )
758 .await?;
759 if replay_queue {
760 for prompt in &canonical_session.queued_prompts {
761 let command = match &prompt.kind {
765 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
766 prompt: prompt
767 .content
768 .iter()
769 .cloned()
770 .map(serde_json::from_value)
771 .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
772 },
773 CanonicalQueuedCommandKind::SetConfig { key, value } => {
774 RelayCommand::SetConfig {
775 key: key.clone(),
776 value: value.clone(),
777 }
778 }
779 };
780 relay.submit(prompt.command_id.clone(), command).await?;
781 }
782 }
783 }
784 if let Some(worktree) = retire_after_ready
788 && let Err(error) = retire_managed_worktree(executor, worktree)
789 {
790 tracing::warn!(
791 session_id,
792 worktree = %worktree.worktree_root.display(),
793 error = format!("{error:#}"),
794 "could not retire the old managed worktree after resume"
795 );
796 resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
797 }
798 for notice in &resume_notices {
799 let submitted = async {
800 let command_id = new_command_id("resume-notice")?;
801 relay
802 .submit(
803 command_id,
804 RelayCommand::RecordNotice {
805 text: notice.clone(),
806 },
807 )
808 .await
809 }
810 .await;
811 if let Err(error) = submitted {
814 tracing::warn!(
815 session_id,
816 error = format!("{error:#}"),
817 "could not record a resume notice in the conversation"
818 );
819 }
820 }
821 self.mark_worker_connected(session_id, Some(native_session_id))?;
822 let materialized = relay.sync().await?.materialized;
823 if !stopped_subagents.is_empty() {
824 let delivered = stopped_subagents
825 .iter()
826 .map(|child| child.child_session_id.clone())
827 .collect::<Vec<_>>();
828 if let Err(error) = crate::database::clear_stopped_subagents(session_id, &delivered) {
831 tracing::warn!(
832 session_id,
833 error = format!("{error:#}"),
834 "could not clear the stopped sub-agents after telling the resumed session"
835 );
836 }
837 }
838 Ok(materialized)
839 }
840}
841
842pub(super) struct VerifiedResumeArchive {
849 pub archive_path: PathBuf,
853 pub manifest: mj_checkpoint::archive::ArchiveManifest,
854 pub canonical_session: Arc<CanonicalSessionSnapshot>,
855}
856
857pub(super) fn verify_resume_checkpoint(
863 session_id: &str,
864 checkpoint: &mj_core::state::CheckpointMetadata,
865) -> Result<VerifiedResumeArchive> {
866 let archive_path = {
867 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
868 checkpoint.archive_path.canonicalize().with_context(|| {
869 format!(
870 "resolve checkpoint archive {}",
871 checkpoint.archive_path.display()
872 )
873 })?
874 };
875 ensure!(
876 archive_path.is_absolute() && archive_path.is_file(),
877 "checkpoint archive path is not an absolute regular file: {}",
878 archive_path.display()
879 );
880 let mj_checkpoint::archive::VerifiedArchiveMetadata {
881 manifest,
882 canonical_session,
883 archive_sha256,
884 } = {
885 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
886 verify_archive_streaming(&archive_path)?
887 };
888 if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
889 bail!("persisted checkpoint verification failed");
890 }
891 ensure!(
892 manifest.repositories.iter().all(|repository| {
893 !repository.metadata.origin.starts_with("mj-local:")
894 && !repository.metadata.origin.starts_with("ext::")
895 }),
896 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
897 );
898 Ok(VerifiedResumeArchive {
899 archive_path,
900 manifest,
901 canonical_session: Arc::new(canonical_session),
902 })
903}
904
905pub fn raw_conversion_preview_for(
906 session: &SessionRecord,
907 config: &Config,
908 executor: &(impl CommandExecutor + Sync),
909) -> Result<mj_core::state::RawConversionPreview> {
910 let conversion = plan_raw_to_workspace(session, config, executor)?;
911 raw_conversion_preview(session, &conversion, executor)
912}
913
914fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
915 let replacement = replacement.trim();
916 ensure!(!replacement.is_empty(), "enter the repository's new origin");
917 let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
918 let path = expanded.as_path();
919 let (github, local) = if path.is_absolute() {
920 ensure!(
921 path.is_dir(),
922 "local repository {replacement:?} is not a directory"
923 );
924 (None, Some(mj_core::local_git::canonical_repository(path)?))
925 } else {
926 let github = crate::setup::github_repository_from_origin(replacement)
927 .context("origin must be a GitHub repository or an absolute local repository path")?;
928 (
929 Some(format!("{}/{}", github.owner, github.repository)),
930 None,
931 )
932 };
933 Ok(ProjectRepository {
934 id: id.to_owned(),
935 github,
936 local,
937 destination: PathBuf::from(id),
938 git_ref: None,
939 })
940}
941
942fn checkpoint_source_missing_commit(
943 configured: &ProjectRepository,
944 archived: &CheckpointRepositoryBundle,
945 executor: &impl CommandExecutor,
946 github_token: Option<&str>,
947) -> Result<Option<String>> {
948 let staging = tempfile::tempdir().context("create repository source preflight")?;
949 let repository = staging.path().join("repository.git");
950 checked_preflight_git(
951 executor,
952 CommandSpec::new(
953 "git",
954 [
955 "init".to_owned(),
956 "--bare".to_owned(),
957 "--quiet".to_owned(),
958 repository.to_string_lossy().into_owned(),
959 ],
960 )
961 .purpose("initialize repository source preflight"),
962 )?;
963 let missing = checkpoint_bundle_prerequisites(archived)?;
964 if missing.is_empty() {
965 let bundle = staging.path().join("checkpoint.bundle");
966 std::fs::write(&bundle, &archived.committed_bundle)
967 .context("write self-contained checkpoint bundle for source preflight")?;
968 checked_preflight_git(
969 executor,
970 checkpoint_bundle_import_command(&repository, &bundle),
971 )?;
972 return Ok(None);
973 }
974 for commit in missing {
978 let output = fetch_source_commit(executor, &repository, configured, &commit, github_token)?;
979 if output.status != 0 {
980 let stderr = String::from_utf8_lossy(&output.stderr);
981 if source_does_not_have_commit(&stderr) {
982 return Ok(Some(commit));
983 }
984 bail!(
985 "could not check configured source {:?}: {}",
986 configured.source_label(),
987 stderr.trim()
988 );
989 }
990 }
991 Ok(None)
995}
996
997fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
998 let mut command = CommandSpec::new(
999 "git",
1000 [
1001 "-C".to_owned(),
1002 repository.to_string_lossy().into_owned(),
1003 "fetch".to_owned(),
1004 "--no-tags".to_owned(),
1005 bundle.to_string_lossy().into_owned(),
1006 "HEAD".to_owned(),
1007 ],
1008 )
1009 .purpose("validate self-contained checkpoint bundle");
1010 command
1011 .env
1012 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1013 command
1014 .env
1015 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1016 command
1017}
1018
1019fn fetch_source_commit(
1020 executor: &impl CommandExecutor,
1021 repository: &Path,
1022 configured: &ProjectRepository,
1023 commit: &str,
1024 github_token: Option<&str>,
1025) -> Result<CommandOutput> {
1026 let mut arguments = Vec::new();
1027 let mut token_auth = false;
1028 let mut ssh_transport = false;
1029 let source = if let Some(local) = &configured.local {
1030 local.to_string_lossy().into_owned()
1031 } else {
1032 let source = configured
1033 .github
1034 .as_deref()
1035 .context("repository source is missing")?;
1036 let github = crate::setup::github_repository_from_origin(source)
1037 .context("configured repository is not a GitHub source")?;
1038 if github_token.is_some() {
1039 token_auth = true;
1040 arguments.extend([
1041 "-c".to_owned(),
1042 "credential.helper=".to_owned(),
1043 "-c".to_owned(),
1044 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1045 ]);
1046 format!(
1047 "https://github.com/{}/{}.git",
1048 github.owner, github.repository
1049 )
1050 } else {
1051 ssh_transport = true;
1052 format!("git@github.com:{}/{}.git", github.owner, github.repository)
1053 }
1054 };
1055 arguments.extend([
1056 "-C".to_owned(),
1057 repository.to_string_lossy().into_owned(),
1058 "fetch".to_owned(),
1059 "--no-tags".to_owned(),
1060 "--depth=1".to_owned(),
1061 "--filter=blob:none".to_owned(),
1062 source,
1063 commit.to_owned(),
1064 ]);
1065 let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1066 command
1067 .env
1068 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1069 command
1070 .env
1071 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1072 if token_auth {
1073 let token = github_token.expect("token authentication requires a GitHub token");
1074 command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1075 }
1076 if ssh_transport {
1077 command.env.insert(
1078 "GIT_SSH_COMMAND".to_owned(),
1079 "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1080 .to_owned(),
1081 );
1082 }
1083 executor.execute(&command)
1084}
1085
1086fn source_does_not_have_commit(stderr: &str) -> bool {
1087 let stderr = stderr.to_ascii_lowercase();
1088 [
1089 "not our ref",
1090 "couldn't find remote ref",
1091 "not a valid object name",
1092 "no such ref was fetched",
1093 ]
1094 .iter()
1095 .any(|needle| stderr.contains(needle))
1096}
1097
1098fn checked_preflight_git(
1099 executor: &impl CommandExecutor,
1100 command: CommandSpec,
1101) -> Result<CommandOutput> {
1102 let output = executor.execute(&command)?;
1103 ensure!(
1104 output.status == 0,
1105 "{}: {}",
1106 command.purpose,
1107 String::from_utf8_lossy(&output.stderr).trim()
1108 );
1109 Ok(output)
1110}
1111
1112impl Controller {
1113 pub async fn resume_session_with_options(
1118 &mut self,
1119 session_id: &str,
1120 profile_id: &str,
1121 target_id: &str,
1122 additional_mounts: Option<Vec<AdditionalMount>>,
1123 resource_allocation: Option<SessionResourceAllocation>,
1124 ) -> Result<MaterializedSession> {
1125 self.resume_session_with_options_and_queue_disposition(
1126 session_id,
1127 profile_id,
1128 target_id,
1129 additional_mounts,
1130 resource_allocation,
1131 false,
1132 )
1133 .await
1134 }
1135
1136 pub async fn resume_session_with_options_and_queue_disposition(
1137 &mut self,
1138 session_id: &str,
1139 profile_id: &str,
1140 target_id: &str,
1141 additional_mounts: Option<Vec<AdditionalMount>>,
1142 resource_allocation: Option<SessionResourceAllocation>,
1143 discard_queue: bool,
1144 ) -> Result<MaterializedSession> {
1145 self.resume_session_controlled(
1146 session_id,
1147 profile_id,
1148 target_id,
1149 SessionResumeOptions {
1150 additional_mounts,
1151 resource_allocation,
1152 discard_queue,
1153 },
1154 &ProcessExecutor,
1155 )
1156 .await
1157 }
1158
1159 pub async fn resume_session_controlled(
1160 &mut self,
1161 session_id: &str,
1162 profile_id: &str,
1163 target_id: &str,
1164 options: SessionResumeOptions,
1165 executor: &(impl CommandExecutor + Sync),
1166 ) -> Result<MaterializedSession> {
1167 self.resume_session_controlled_with_repository_preflight(
1168 session_id, profile_id, target_id, options, None, executor,
1169 )
1170 .await
1171 }
1172
1173 pub async fn resume_session_controlled_with_repository_preflight(
1174 &mut self,
1175 session_id: &str,
1176 profile_id: &str,
1177 target_id: &str,
1178 options: SessionResumeOptions,
1179 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1180 executor: &(impl CommandExecutor + Sync),
1181 ) -> Result<MaterializedSession> {
1182 self.resume_session_with_origin(
1183 session_id,
1184 profile_id,
1185 target_id,
1186 options,
1187 repository_preflight,
1188 None,
1189 executor,
1190 )
1191 .await
1192 }
1193
1194 pub(in crate::controller) async fn resume_session_for_move(
1195 &mut self,
1196 operation: &mj_core::state::MoveOperation,
1197 executor: &(impl CommandExecutor + Sync),
1198 ) -> Result<MaterializedSession> {
1199 self.resume_session_with_origin(
1200 &operation.selection.session_id,
1201 operation
1202 .selection
1203 .profile_id
1204 .as_deref()
1205 .context("Move profile missing")?,
1206 operation
1207 .selection
1208 .target_template_id
1209 .as_deref()
1210 .context("Move target missing")?,
1211 SessionResumeOptions {
1212 additional_mounts: operation.selection.additional_mounts.clone(),
1213 resource_allocation: operation.selection.resource_allocation.clone(),
1214 discard_queue: true,
1215 },
1216 None,
1217 Some(operation),
1218 executor,
1219 )
1220 .await
1221 }
1222
1223 #[allow(clippy::too_many_arguments)]
1224 async fn resume_session_with_origin(
1225 &mut self,
1226 session_id: &str,
1227 profile_id: &str,
1228 target_id: &str,
1229 options: SessionResumeOptions,
1230 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1231 move_operation: Option<&mj_core::state::MoveOperation>,
1232 executor: &(impl CommandExecutor + Sync),
1233 ) -> Result<MaterializedSession> {
1234 let transferring_workspace = move_operation.is_some();
1235 if let Some(operation) = crate::database::load_move_operation(session_id)? {
1236 ensure!(
1237 transferring_workspace
1238 || !((operation.in_place || operation.workspace_transfer.is_some())
1239 && operation.phase != mj_core::state::MovePhase::Completed
1240 && operation.recovery_session.is_some()),
1241 "a profile switch retains this environment; retry Move instead of recreating it with Resume"
1242 );
1243 }
1244 let SessionResumeOptions {
1245 additional_mounts,
1246 resource_allocation,
1247 discard_queue,
1248 } = options;
1249 let mut previous = self
1250 .state
1251 .sessions
1252 .get(session_id)
1253 .with_context(|| format!("unknown session {session_id}"))?
1254 .clone();
1255 if !(matches!(
1256 previous.state,
1257 SessionState::Stopped | SessionState::Lost | SessionState::Error
1258 ) || (transferring_workspace && previous.state == SessionState::Closing))
1259 {
1260 bail!("session {session_id} is not stopped, lost, or retryable");
1261 }
1262 let checkpoint = move_operation
1263 .and_then(|op| op.handoff.as_ref())
1264 .or(previous.checkpoint.as_ref())
1265 .context("session has no checkpoint")?;
1266 let moving_to_raw = previous.project_directory.is_none()
1269 && self
1270 .config
1271 .targets
1272 .get(target_id)
1273 .is_some_and(mj_core::config::is_bare_project_target);
1274 if !transferring_workspace
1275 && (moving_to_raw
1276 || !repository_preflight.as_ref().is_some_and(|receipt| {
1277 self.repository_source_receipt_is_current(session_id, receipt)
1278 }))
1279 {
1280 let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1281 if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1282 self.preflight_repository_sources(session_id, target_id, false, executor)?
1283 {
1284 bail!(
1285 "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1286 mismatch.missing_commit,
1287 mismatch.configured_origin,
1288 mismatch.repository_id,
1289 mismatch.archived_origin,
1290 );
1291 }
1292 }
1293 let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1294 let archive_path = verified_archive.archive_path.clone();
1295 let archive_manifest = &verified_archive.manifest;
1296 let canonical_session = Arc::clone(&verified_archive.canonical_session);
1297 let profile = self
1298 .config
1299 .profiles
1300 .get(profile_id)
1301 .with_context(|| format!("unknown profile {profile_id:?}"))?
1302 .clone();
1303 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1304 let target_template = self
1305 .config
1306 .targets
1307 .get(target_id)
1308 .with_context(|| format!("unknown target template {target_id:?}"))?
1309 .clone();
1310 self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1313 ensure!(
1314 profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1315 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1316 );
1317 let plan = resume_compatibility(&previous, &self.config, target_id)
1318 .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1319 if plan == ResumePlan::InPlace
1322 && let Some(worktree) = previous.managed_worktree.as_mut()
1323 && let Ok(current) = managed_worktree_target(&target_template)
1324 {
1325 worktree.target = current;
1326 }
1327 if !mj_core::config::is_bare_project_target(&target_template)
1331 && plan != ResumePlan::RawToWorkspace
1332 && !transferring_workspace
1333 {
1334 super::network_git::bundle_from_manifest(archive_manifest)?;
1335 }
1336 if plan == ResumePlan::InPlace
1337 && previous.managed_worktree.is_none()
1338 && let Some(project_directory) = &previous.project_directory
1339 {
1340 self.validate_project_directory(target_id, project_directory, executor)
1341 .context("raw project is unavailable for resume")?;
1342 }
1343 let conversion = match plan {
1344 ResumePlan::InPlace => None,
1345 ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(
1346 plan_raw_to_workspace(&previous, &self.config, executor)
1347 .context("prepare the raw checkout for its new target")?,
1348 )),
1349 ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1350 self.plan_workspace_to_raw(&previous, target_id, executor)
1351 .context("prepare a checkout for this session")?,
1352 )),
1353 };
1354 let resource_allocation =
1355 resource_allocation.or_else(|| previous.resource_allocation.clone());
1356 let additional_mounts =
1357 additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1358 validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1359 let selected_container_size =
1360 selected_host_container_size(&target_template, resource_allocation.as_ref());
1361 if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1362 bail!("attached resources are unsupported for this target");
1363 }
1364 targets::validate_additional_mounts(&additional_mounts)?;
1365 let history_host = mount_history_host(&target_template);
1366 let history_mounts = additional_mounts.clone();
1367 if !transferring_workspace
1368 && previous.state == SessionState::Error
1369 && let Some(locator) = &previous.target
1370 {
1371 let backend = backend_locator(locator, &previous, &self.config)?;
1372 targets::close_plan(&backend, session_id)?
1373 .execute(executor)
1374 .context("clean up target from failed resume")?;
1375 }
1376 let mut resume_notices = Vec::new();
1377 if let Some(conversion) = conversion
1380 .as_ref()
1381 .and_then(ResumeConversion::workspace_to_raw)
1382 {
1383 resume_notices.push(format!(
1384 "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1385 previous.target_template_id,
1386 conversion.worktree.worktree_root.display(),
1387 archive_manifest
1388 .repositories
1389 .first()
1390 .and_then(|repository| repository.metadata.branch.as_deref())
1391 .unwrap_or("a detached head"),
1392 conversion.worktree.branch,
1393 ));
1394 }
1395 let managed_checkout_present = previous
1396 .managed_worktree
1397 .as_ref()
1398 .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1399 .transpose()?
1400 .unwrap_or(true);
1401 if managed_checkout_present && let Some(project_directory) = &previous.project_directory {
1404 match raw_checkout_position(&previous, &self.config, project_directory, executor) {
1405 Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1406 project_directory,
1407 archive_manifest
1408 .repositories
1409 .first()
1410 .map(|repository| &repository.metadata),
1411 &live,
1412 )),
1413 Err(error) => tracing::warn!(
1416 session_id,
1417 error = format!("{error:#}"),
1418 "could not read the raw checkout position for a resume notice"
1419 ),
1420 }
1421 }
1422 super::worker_binary::preflight_worker_binary(&target_template, executor)?;
1427 let same_harness = profile.kind == archive_manifest.session.harness_kind;
1428 let native_continuity =
1429 native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1430 let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1431 let utility_config = (!native_continuity).then(|| self.config.clone());
1435 let discard_queued_prompts = discard_queue || !same_harness;
1436 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1440 .unwrap_or_else(|error| {
1441 tracing::warn!(
1442 session_id,
1443 error = format!("{error:#}"),
1444 "could not read the stored projection frontier; rebuilding it from the archive"
1445 );
1446 None
1447 });
1448 let rebuild_projection = projection_rebuild_required(
1449 stored_frontier
1450 .as_ref()
1451 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1452 canonical_session.event_frontier,
1453 &canonical_session.event_frontier_digest,
1454 );
1455 let projection_build = rebuild_projection.then(|| {
1465 let canonical = Arc::clone(&canonical_session);
1466 let session_id = session_id.to_owned();
1467 tokio::task::spawn_blocking(move || {
1468 materialized_session_from_canonical(session_id, &canonical)
1469 })
1470 });
1471 let github_token = controller_github_token();
1472
1473 if let Some(conversion) = conversion
1476 .as_ref()
1477 .and_then(ResumeConversion::raw_to_workspace)
1478 && let Some(bundle) = &conversion.new_bundle
1479 {
1480 let (config, ()) = Config::update(|config| {
1481 if let Some(existing) = config.bundles.get(&conversion.bundle_id) {
1482 ensure!(
1483 existing == bundle,
1484 "bundle {:?} was configured concurrently with a different definition; retry the resume",
1485 conversion.bundle_id
1486 );
1487 } else {
1488 config
1489 .bundles
1490 .insert(conversion.bundle_id.clone(), bundle.clone());
1491 }
1492 Ok(())
1493 })
1494 .context("save the bundle for a converted raw session")?;
1495 self.config = config;
1496 }
1497
1498 let moving_into_first_container = self
1504 .config
1505 .targets
1506 .get(target_id)
1507 .is_some_and(mj_core::config::is_container_target)
1508 && !self
1509 .config
1510 .targets
1511 .get(&previous.target_template_id)
1512 .is_some_and(mj_core::config::is_container_target);
1513 let record = self.state.sessions.get_mut(session_id).unwrap();
1514 if record.container_workspace.is_none() && moving_into_first_container {
1515 record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1516 }
1517 record.harness_kind = profile.kind;
1518 record.last_profile = profile_id.to_string();
1519 record.target_template_id = target_id.to_string();
1520 record.target_runtime = Some(
1521 self.config
1522 .targets
1523 .get(target_id)
1524 .context("resume target disappeared")?
1525 .into(),
1526 );
1527 record.resource_allocation = resource_allocation;
1528 record.additional_mounts = additional_mounts;
1529 record.target = None;
1530 record.native_session_id =
1531 native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1532 record.state = SessionState::Provisioning;
1533 record.updated_at = now();
1534 record.last_error = None;
1535 match &conversion {
1536 Some(ResumeConversion::RawToWorkspace(conversion)) => {
1537 apply_raw_to_workspace(record, conversion);
1538 }
1539 Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1540 apply_workspace_to_raw(record, conversion);
1541 }
1542 None => {
1543 if let (Some(worktree), Some(refreshed)) = (
1544 record.managed_worktree.as_mut(),
1545 previous.managed_worktree.as_ref(),
1546 ) {
1547 worktree.target = refreshed.target.clone();
1548 }
1549 }
1550 }
1551 let resumed_project_directory = record.project_directory.clone();
1552 let resumed_container_workspace = record.container_workspace.clone();
1553 if let Some(host) = history_host {
1554 self.state.remember_mount_sources(host, &history_mounts);
1555 crate::database::remember_mount_sources(host, &history_mounts)?;
1556 }
1557 if let Some(conversion) = conversion
1560 .as_ref()
1561 .and_then(ResumeConversion::raw_to_workspace)
1562 {
1563 crate::database::rebind_session_bundle(session_id, &conversion.bundle_id)?;
1564 }
1565 crate::database::save_resumed_session(
1568 &self.state.sessions[session_id],
1569 selected_container_size
1570 .as_ref()
1571 .map(|(host, size)| (host.as_str(), *size)),
1572 )?;
1573 if let Some((host, size)) = selected_container_size.as_ref() {
1574 self.state.remember_container_size(host, *size);
1575 }
1576
1577 let mut recreated_managed_worktree = false;
1578 let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1581 let result = async {
1582 if let Some(worktree) = previous.managed_worktree.as_ref() {
1583 recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1584 if !transferring_workspace && recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1585 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1586 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1587 &archive_path, &worktree.worktree_root, &SystemGit,
1588 )?;
1589 } else {
1590 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1591 &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1592 )?;
1593 }
1594 }
1595 }
1596 if let Some(conversion) = conversion
1599 .as_ref()
1600 .and_then(ResumeConversion::workspace_to_raw)
1601 {
1602 if conversion.reuse_existing_branch {
1603 let recovery_ref =
1604 preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1605 restore_managed_worktree(executor, &conversion.worktree)?;
1606 resume_notices.push(format!(
1607 "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1608 ));
1609 } else {
1610 create_managed_worktree(
1611 executor,
1612 &conversion.worktree,
1613 None,
1614 PrimaryCheckoutRequirement::Any,
1615 )?;
1616 }
1617 if transferring_workspace {
1618 } else if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1620 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1621 &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1622 )?;
1623 } else {
1624 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1625 &archive_path,
1626 &conversion.worktree.worktree_root,
1627 &conversion.worktree.branch,
1628 &SystemGit,
1629 )?;
1630 }
1631 }
1632 if let Some(conversion) = conversion
1637 .as_ref()
1638 .and_then(ResumeConversion::raw_to_workspace)
1639 && !transferring_workspace
1640 {
1641 let destination = PathBuf::from(
1642 previous
1643 .project_directory
1644 .as_deref()
1645 .context("a raw session has no project directory")?
1646 .file_name()
1647 .context("a raw project directory cannot be the filesystem root")?,
1648 );
1649 let snapshot = raw_checkout_snapshot(
1650 &conversion.checkout,
1651 &conversion.source,
1652 &destination,
1653 &SystemGit,
1654 conversion.retire.as_ref().is_some_and(|checkout| {
1655 checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1656 }),
1657 )
1658 .context("snapshot the host checkout for its new target")?;
1659 resume_notices.push(conversion_notice(
1660 target_id,
1661 previous
1662 .project_directory
1663 .as_deref()
1664 .unwrap_or(&conversion.checkout),
1665 snapshot.metadata.branch.as_deref(),
1666 conversion.retire.as_ref(),
1667 ));
1668 let archives = mj_core::config::sessions_dir();
1669 std::fs::create_dir_all(&archives).with_context(|| {
1670 format!("create the checkpoint directory {}", archives.display())
1671 })?;
1672 let output = archives.join(format!(
1676 "{session_id}-converted-{}-{}.hel.zip",
1677 previous
1678 .checkpoint
1679 .as_ref()
1680 .map_or(0, |checkpoint| checkpoint.event_frontier),
1681 new_command_id("archive")?
1682 ));
1683 let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1684 conversion_checkpoint_written = Some(written.clone());
1685 let record = self.state.sessions.get_mut(session_id).unwrap();
1686 record.checkpoint = Some(written);
1687 record.updated_at = now();
1688 crate::database::save_resumed_session(
1691 &self.state.sessions[session_id],
1692 selected_container_size.as_ref().map(|(host, size)| (host.as_str(), *size)),
1693 )?;
1694 }
1695 let utility_handoff = {
1696 let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1697 if let Some(config) = utility_config.as_ref() {
1698 Some(
1699 provision_with_cross_harness_handoff(
1700 self,
1701 session_id,
1702 executor,
1703 github_token.as_deref(),
1704 config,
1705 &canonical_session,
1706 context_bytes,
1707 )
1708 .context("prepare the cross-harness destination")?,
1709 )
1710 } else {
1711 self.provision_session_with_failure_disposition(
1712 session_id,
1713 executor,
1714 github_token.as_deref(),
1715 ProvisioningFailureDisposition::Preserve,
1716 )
1717 .await?;
1718 None
1719 }
1720 };
1721 if transferring_workspace {
1722 self.restore_move_workspace(session_id, executor)?;
1723 }
1724 let restore_repositories = !transferring_workspace && ((resumed_project_directory.is_none()
1725 && conversion.is_none())
1726 || plan == ResumePlan::RawToWorkspace
1727 || (recreated_managed_worktree && plan == ResumePlan::InPlace));
1728 let restored_archive = conversion_checkpoint_written
1731 .as_ref()
1732 .map_or(archive_path.as_path(), |checkpoint| {
1733 checkpoint.archive_path.as_path()
1734 });
1735 self.restore_into_target(
1736 session_id,
1737 RestoreIntoTarget {
1738 profile: &profile,
1739 archive: &verified_archive,
1740 restored_archive,
1741 resumed_project_directory,
1742 resumed_container_workspace,
1743 restore_repositories,
1744 native_continuity,
1745 discard_queued_prompts,
1746 replay_queue: !discard_queue,
1747 utility_handoff,
1748 projection_build,
1749 resume_notices,
1750 install_attached_resources: true,
1751 worker_root_reset: WorkerRootReset::FreshTarget,
1752 retire_after_ready: conversion
1753 .as_ref()
1754 .filter(|_| !transferring_workspace)
1755 .and_then(ResumeConversion::raw_to_workspace)
1756 .and_then(|plan| plan.retire.as_ref()),
1757 },
1758 executor,
1759 )
1760 .await
1761 }
1762 .await;
1763 match result {
1764 Ok(materialized) => {
1765 if let Some(written) = &conversion_checkpoint_written {
1768 super::checkpoint::prune_replaced_checkpoint(
1769 previous.checkpoint.as_ref(),
1770 written,
1771 );
1772 }
1773 Ok(materialized)
1774 }
1775 Err(error) => {
1776 if let Some(written) = &conversion_checkpoint_written
1779 && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
1780 && remove_error.kind() != std::io::ErrorKind::NotFound
1781 {
1782 tracing::warn!(
1783 session_id,
1784 path = %written.archive_path.display(),
1785 "could not remove the conversion checkpoint after resume failed: {remove_error}"
1786 );
1787 }
1788 restore_projection_after_failed_resume(
1791 session_id,
1792 &canonical_session,
1793 discard_queued_prompts,
1794 );
1795 if let Some(operation) = move_operation {
1796 return Err(self.rollback_move_destination(operation, error, executor)?);
1797 }
1798 Err(self.rollback_failed_resume(
1799 session_id,
1800 &previous,
1801 recreated_managed_worktree,
1802 error,
1803 executor,
1804 )?)
1805 }
1806 }
1807 }
1808
1809 pub(super) fn rollback_failed_resume(
1810 &mut self,
1811 session_id: &str,
1812 previous: &SessionRecord,
1813 recreated_managed_worktree: bool,
1814 error: anyhow::Error,
1815 _executor: &impl CommandExecutor,
1816 ) -> Result<anyhow::Error> {
1817 let current = self
1818 .state
1819 .sessions
1820 .get(session_id)
1821 .with_context(|| format!("unknown session {session_id}"))?
1822 .clone();
1823 let cleanup = match current.target.as_ref() {
1824 Some(locator) => (|| -> Result<()> {
1825 let backend = backend_locator(locator, ¤t, &self.config)?;
1826 targets::close_plan(&backend, session_id)?
1827 .execute(&CancellableProcessExecutor::with_timeout(
1830 Duration::from_secs(15),
1831 ))
1832 .map(|_| ())
1833 })(),
1834 None => Ok(()),
1835 };
1836 let worktree_cleanup = if cleanup.is_err() {
1839 Ok(())
1840 } else {
1841 match (
1842 current.managed_worktree.as_ref(),
1843 previous.managed_worktree.as_ref(),
1844 ) {
1845 (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
1846 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1847 previous,
1848 ),
1849 (Some(current), Some(previous)) if current == previous => Ok(()),
1850 (Some(worktree), _) => cleanup_managed_worktree(
1853 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1854 worktree,
1855 crate::controller::BranchDisposition::Delete,
1856 ),
1857 (None, _) => Ok(()),
1858 }
1859 };
1860 let cleanup_error = [cleanup, worktree_cleanup]
1861 .into_iter()
1862 .filter_map(Result::err)
1863 .map(|cleanup_error| format!("{cleanup_error:#}"))
1864 .collect::<Vec<_>>()
1865 .join("; ");
1866 if !cleanup_error.is_empty() {
1867 tracing::warn!(
1868 session_id,
1869 error = %cleanup_error,
1870 "resume rollback cleanup reported failures"
1871 );
1872 }
1873 let original = super::worker_binary::failure_line(&error);
1876 let detail = format!("{error:#}");
1877 if detail != original {
1878 tracing::warn!(session_id, error = %detail, "resume failed");
1879 }
1880 let record = self.state.sessions.get_mut(session_id).unwrap();
1881 let failure = apply_failed_resume_rollback(
1882 record,
1883 previous,
1884 &original,
1885 (!cleanup_error.is_empty()).then_some(cleanup_error),
1886 );
1887 if record.bundle_id != current.bundle_id {
1890 let bundle_id = record.bundle_id.clone();
1891 crate::database::rebind_session_bundle(session_id, &bundle_id)?;
1892 }
1893 crate::database::save_resumed_session(&self.state.sessions[session_id], None)?;
1896 Ok(failure)
1897 }
1898}
1899
1900fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
1901 format!(
1902 "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
1903 worktree_root.display(),
1904 worktree_root.display()
1905 )
1906}
1907
1908pub(super) fn apply_failed_resume_rollback(
1909 current: &mut SessionRecord,
1910 previous: &SessionRecord,
1911 original_error: &str,
1912 cleanup_error: Option<String>,
1913) -> anyhow::Error {
1914 match cleanup_error {
1915 None => {
1916 *current = previous.clone();
1917 current.state = SessionState::Stopped;
1918 current.target = None;
1919 current.updated_at = now();
1920 current.last_error = Some(format!("resume failed: {original_error}"));
1921 anyhow::anyhow!(original_error.to_owned())
1922 }
1923 Some(cleanup_error) => {
1924 let failure = format!(
1925 "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
1926 );
1927 if current.managed_worktree.is_none() {
1932 current
1933 .project_directory
1934 .clone_from(&previous.project_directory);
1935 current
1936 .managed_worktree
1937 .clone_from(&previous.managed_worktree);
1938 current.bundle_id.clone_from(&previous.bundle_id);
1939 }
1940 current.state = SessionState::Error;
1941 current.updated_at = now();
1942 current.last_error = Some(format!("resume failed: {failure}"));
1943 anyhow::anyhow!(failure)
1944 }
1945 }
1946}
1947
1948fn conversion_notice(
1950 target_id: &str,
1951 checkout: &Path,
1952 branch: Option<&str>,
1953 retire: Option<&mj_core::state::ManagedWorktree>,
1954) -> String {
1955 let branch = branch.unwrap_or("a detached head");
1956 match retire {
1957 Some(worktree) => format!(
1958 "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
1959 checkout.display(),
1960 worktree.branch,
1961 worktree.source_repository.display()
1962 ),
1963 None => format!(
1964 "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.",
1965 checkout.display()
1966 ),
1967 }
1968}
1969
1970fn conversion_checkpoint(
1978 previous_archive: &Path,
1979 snapshot: mj_checkpoint::archive::RepositorySnapshot,
1980 output: &Path,
1981) -> Result<mj_core::state::CheckpointMetadata> {
1982 let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
1983 .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
1984 let native_artifacts = previous
1987 .manifest
1988 .payloads
1989 .iter()
1990 .filter_map(|descriptor| match &descriptor.role {
1991 mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
1992 Some((relative_path, descriptor))
1993 }
1994 _ => None,
1995 })
1996 .map(|(relative_path, descriptor)| {
1997 Ok(mj_checkpoint::archive::NativeArtifact {
1998 relative_path: relative_path.clone(),
1999 data: previous.payload(descriptor)?.to_vec(),
2000 mode: descriptor.mode,
2001 })
2002 })
2003 .collect::<Result<Vec<_>>>()?;
2004 let canonical_session = previous.canonical_session()?;
2005 let event_frontier = canonical_session.event_frontier;
2006 let written = mj_checkpoint::archive::write_archive_atomic(
2007 output,
2008 &mj_checkpoint::archive::ArchiveInput {
2009 session: previous.manifest.session.clone(),
2010 target: previous.manifest.target.clone(),
2013 bundle: mj_checkpoint::archive::BundleManifest {
2014 id: previous.manifest.bundle.id.clone(),
2015 primary_repository: snapshot.metadata.id.clone(),
2019 },
2020 canonical_session,
2021 native_artifacts,
2022 repositories: vec![snapshot],
2023 },
2024 )
2025 .with_context(|| format!("write the conversion archive {}", output.display()))?;
2026 Ok(mj_core::state::CheckpointMetadata {
2027 archive_path: output.to_path_buf(),
2028 sha256: written.archive_sha256,
2029 created_at: now(),
2030 event_frontier,
2031 })
2032}
2033
2034fn restored_native_session_accepted(
2045 archived: &str,
2046 opened: &str,
2047 replaced_unused: Option<&str>,
2048) -> bool {
2049 opened == archived || replaced_unused == Some(archived)
2050}
2051
2052fn restore_projection_after_failed_resume(
2063 session_id: &str,
2064 canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
2065 discard_queued_prompts: bool,
2066) {
2067 let stored_frontier = crate::database::materialized_event_frontier(session_id)
2068 .unwrap_or_else(|error| {
2069 tracing::warn!(
2070 session_id,
2071 error = format!("{error:#}"),
2072 "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
2073 );
2074 None
2075 });
2076 if projection_rebuild_required(
2077 stored_frontier
2078 .as_ref()
2079 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
2080 canonical_session.event_frontier,
2081 &canonical_session.event_frontier_digest,
2082 ) {
2083 match materialized_session_from_canonical(session_id, canonical_session) {
2084 Ok(previous_projection) => {
2085 if let Err(restore_error) =
2086 crate::database::save_materialized_session(&previous_projection)
2087 {
2088 tracing::error!(
2089 session_id,
2090 error = format!("{restore_error:#}"),
2091 "could not restore the durable projection after a failed resume"
2092 );
2093 }
2094 }
2095 Err(restore_error) => tracing::error!(
2096 session_id,
2097 error = format!("{restore_error:#}"),
2098 "could not rebuild the durable projection after a failed resume"
2099 ),
2100 }
2101 } else if discard_queued_prompts
2102 && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
2103 session_id,
2104 &mj_transcript::projection::materialized_queued_prompts_from_canonical(
2105 &canonical_session.queued_prompts,
2106 ),
2107 )
2108 {
2109 tracing::error!(
2110 session_id,
2111 error = format!("{restore_error:#}"),
2112 "could not restore queued prompts after a failed resume"
2113 );
2114 }
2115}
2116
2117fn projection_rebuild_required(
2125 stored: Option<(u64, &str)>,
2126 archive_frontier: u64,
2127 archive_frontier_digest: &str,
2128) -> bool {
2129 stored != Some((archive_frontier, archive_frontier_digest))
2130}
2131
2132fn restore_archive_path(
2133 backend: &targets::TargetLocator,
2134 verified_archive: &Path,
2135 remote_archive: &Path,
2136) -> PathBuf {
2137 if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2138 verified_archive.to_path_buf()
2139 } else {
2140 remote_archive.to_path_buf()
2141 }
2142}
2143
2144fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2145 !matches!(backend, targets::TargetLocator::LocalBare { .. })
2146}
2147
2148struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2153 inner: &'a E,
2154 cancellation: CancellationToken,
2155}
2156
2157impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2158 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2159 if self.cancellation.is_cancelled() {
2160 bail!("operation cancelled while provisioning destination");
2161 }
2162 self.inner.execute(command)
2163 }
2164
2165 fn cancellation_requested(&self) -> bool {
2166 self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2167 }
2168
2169 fn stage_started(&self, stage: ProvisionStage) {
2170 self.inner.stage_started(stage);
2171 }
2172
2173 fn stage_finished(&self, stage: ProvisionStage) {
2174 self.inner.stage_finished(stage);
2175 }
2176
2177 fn notify_notice(&self, notice: &str) {
2178 self.inner.notify_notice(notice);
2179 }
2180
2181 fn execute_with_stdin(
2182 &self,
2183 command: &CommandSpec,
2184 input: &mut (dyn std::io::Read + Send),
2185 ) -> Result<CommandOutput> {
2186 if self.cancellation.is_cancelled() {
2187 bail!("operation cancelled while provisioning destination");
2188 }
2189 self.inner.execute_with_stdin(command, input)
2190 }
2191}
2192
2193fn provision_with_cross_harness_handoff(
2194 controller: &mut Controller,
2195 session_id: &str,
2196 executor: &(impl CommandExecutor + Sync),
2197 github_token: Option<&str>,
2198 config: &Config,
2199 snapshot: &CanonicalSessionSnapshot,
2200 context_bytes: usize,
2201) -> Result<String> {
2202 let (_provision, handoff) = execute_joined_cross_harness_work(
2203 "cross-harness provisioning",
2204 move |cancellation| {
2205 let provision_executor = CrossHarnessProvisionExecutor {
2206 inner: executor,
2207 cancellation,
2208 };
2209 futures::executor::block_on(controller.provision_session_with_failure_disposition(
2210 session_id,
2211 &provision_executor,
2212 github_token,
2213 ProvisioningFailureDisposition::Preserve,
2214 ))
2215 },
2216 "cross-harness handoff",
2217 move |cancellation| {
2218 let runtime = tokio::runtime::Builder::new_current_thread()
2219 .enable_all()
2220 .build()
2221 .context("create cross-harness handoff runtime")?;
2222 runtime.block_on(utility_handoff_while_cancellable(
2223 session_id,
2224 config,
2225 snapshot,
2226 context_bytes,
2227 executor,
2228 cancellation,
2229 ))
2230 },
2231 )?;
2232 ensure!(
2233 !executor.cancellation_requested(),
2234 "operation cancelled while provisioning destination"
2235 );
2236 Ok(handoff)
2237}
2238
2239fn execute_joined_cross_harness_work<A: Send, B: Send>(
2243 first_name: &'static str,
2244 first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2245 second_name: &'static str,
2246 second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2247) -> Result<(A, B)> {
2248 let cancellation = CancellationToken::new();
2249 std::thread::scope(|scope| {
2250 let first_cancel = cancellation.clone();
2251 let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2252 let second_cancel = cancellation.clone();
2253 let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2254 let mut first_result = None;
2255 let mut second_result = None;
2256
2257 while first_result.is_none() || second_result.is_none() {
2258 if first_result.is_none()
2259 && first_handle
2260 .as_ref()
2261 .is_some_and(|handle| handle.is_finished())
2262 {
2263 let handle = first_handle.take().expect("first lane handle present");
2264 first_result = Some(match handle.join() {
2265 Ok(result) => result,
2266 Err(panic) => {
2267 cancellation.cancel();
2268 Err(anyhow::anyhow!(
2269 "{first_name} thread panicked: {}",
2270 targets::command_thread_panic_message(panic.as_ref())
2271 ))
2272 }
2273 });
2274 if first_result.as_ref().is_some_and(Result::is_err) {
2275 cancellation.cancel();
2276 }
2277 }
2278 if second_result.is_none()
2279 && second_handle
2280 .as_ref()
2281 .is_some_and(|handle| handle.is_finished())
2282 {
2283 let handle = second_handle.take().expect("second lane handle present");
2284 second_result = Some(match handle.join() {
2285 Ok(result) => result,
2286 Err(panic) => {
2287 cancellation.cancel();
2288 Err(anyhow::anyhow!(
2289 "{second_name} thread panicked: {}",
2290 targets::command_thread_panic_message(panic.as_ref())
2291 ))
2292 }
2293 });
2294 if second_result.as_ref().is_some_and(Result::is_err) {
2295 cancellation.cancel();
2296 }
2297 }
2298 if first_result.is_none() || second_result.is_none() {
2299 std::thread::sleep(Duration::from_millis(10));
2300 }
2301 }
2302
2303 match (
2304 first_result.expect("first lane result received after joined handle"),
2305 second_result.expect("second lane result received after joined handle"),
2306 ) {
2307 (Err(first), Err(second)) => {
2308 Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2309 }
2310 (Err(error), Ok(_)) => Err(error),
2311 (Ok(_), Err(error)) => Err(error),
2312 (Ok(first), Ok(second)) => Ok((first, second)),
2313 }
2314 })
2315}
2316
2317fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2321 profile_kind == archived_kind
2322}
2323
2324async fn utility_handoff_while_cancellable(
2328 session_id: &str,
2329 config: &Config,
2330 snapshot: &CanonicalSessionSnapshot,
2331 context_bytes: usize,
2332 executor: &impl CommandExecutor,
2333 cancellation: CancellationToken,
2334) -> Result<String> {
2335 let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2336 if executor.cancellation_requested() {
2337 bail!("operation cancelled while compacting the cross-harness handoff");
2338 }
2339 let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2340 let cancel = cancellation.child_token();
2341 let operation =
2342 crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2343 tokio::pin!(operation);
2344 loop {
2345 tokio::select! {
2346 context = &mut operation => return context,
2347 _ = cancellation.cancelled() => {
2348 cancel.cancel();
2349 bail!("operation cancelled while compacting the cross-harness handoff");
2350 }
2351 _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2352 if executor.cancellation_requested() {
2353 cancel.cancel();
2354 bail!("operation cancelled while compacting the cross-harness handoff");
2355 }
2356 }
2357 }
2358 }
2359}
2360
2361mod in_place;
2362
2363#[cfg(test)]
2364mod tests;