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_durably, worker_probe_diagnosis};
37use super::worktree::{
38 PrimaryCheckoutRequirement, ResumeConversion, ResumePlan, apply_raw_to_workspace,
39 apply_workspace_to_raw, cleanup_managed_worktree, create_managed_worktree,
40 managed_worktree_checkout_exists, managed_worktree_target, plan_raw_to_workspace,
41 preserve_retained_managed_worktree_branch, raw_checkout_divergence_notice,
42 raw_checkout_position, raw_checkout_snapshot, raw_conversion_preview, restore_managed_worktree,
43 resume_compatibility, retire_managed_worktree,
44};
45use super::{
46 Controller, SessionResumeOptions, execute_checked, now, selected_host_container_size,
47 target_profile_home,
48};
49
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct ResumeRepositorySourceMismatch {
52 pub session_id: String,
53 pub bundle_id: String,
54 pub repository_id: String,
55 pub missing_commit: String,
56 pub archived_origin: String,
57 pub configured_origin: String,
58}
59
60pub use mj_core::state::ResumeRepositorySourceReceipt;
61
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum ResumeRepositorySourcePreflight {
64 Ready(ResumeRepositorySourceReceipt),
65 RepositoryMoved(ResumeRepositorySourceMismatch),
66 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 || source
125 .project_bundle(&self.config)
126 .is_none_or(|bundle| bundle.repositories.len() == 1),
127 "Muse Code ACP supports one workspace root; use a single-repository bundle"
128 );
129 Ok(())
130 }
131 pub fn preflight_resume_repository_sources(
135 &self,
136 session_id: &str,
137 target_id: &str,
138 executor: &(impl CommandExecutor + Sync),
139 ) -> Result<ResumeRepositorySourcePreflight> {
140 self.preflight_repository_sources(session_id, target_id, true, executor)
141 }
142
143 fn preflight_repository_sources(
148 &self,
149 session_id: &str,
150 target_id: &str,
151 describe_conversion: bool,
152 executor: &(impl CommandExecutor + Sync),
153 ) -> Result<ResumeRepositorySourcePreflight> {
154 let session = self
155 .state
156 .sessions
157 .get(session_id)
158 .with_context(|| format!("unknown session {session_id}"))?;
159 let checkpoint = session
160 .checkpoint
161 .as_ref()
162 .context("session has no checkpoint")?;
163 let plan = resume_compatibility(session, &self.config, target_id)
164 .map_err(|reason| anyhow::anyhow!(reason))?;
165 if session.project_directory.is_some() {
166 debug_assert!(matches!(
167 plan,
168 ResumePlan::InPlace | ResumePlan::RawToWorkspace
169 ));
170 let receipt = ResumeRepositorySourceReceipt {
175 session_id: session_id.to_owned(),
176 bundle_id: session.bundle_id.clone(),
177 checkpoint_sha256: checkpoint.sha256.clone(),
178 repositories: Vec::new(),
179 };
180 if describe_conversion && plan == ResumePlan::RawToWorkspace {
185 let preview = raw_conversion_preview_for(session, &self.config, executor)?;
186 return Ok(ResumeRepositorySourcePreflight::ConvertingRawCheckout {
187 receipt,
188 preview: Box::new(preview),
189 });
190 }
191 return Ok(ResumeRepositorySourcePreflight::Ready(receipt));
192 }
193 let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
194 self.preflight_verified_repository_sources(
195 session_id,
196 ResumeRepositoryBundles {
197 checkpoint_sha256: checkpoint.sha256.clone(),
198 repositories,
199 },
200 None,
201 plan != ResumePlan::WorkspaceToRaw,
202 executor,
203 )
204 }
205
206 fn preflight_verified_repository_sources(
207 &self,
208 session_id: &str,
209 verified: ResumeRepositoryBundles,
210 skip_repository_id: Option<&str>,
211 use_archived_network_sources: bool,
212 executor: &(impl CommandExecutor + Sync),
213 ) -> Result<ResumeRepositorySourcePreflight> {
214 let session = self
215 .state
216 .sessions
217 .get(session_id)
218 .with_context(|| format!("unknown session {session_id}"))?;
219 ensure!(
220 verified.repositories.iter().all(|repository| {
221 !repository.metadata.origin.starts_with("mj-local:")
222 && !repository.metadata.origin.starts_with("ext::")
223 }),
224 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
225 );
226 if use_archived_network_sources
229 && verified
230 .repositories
231 .iter()
232 .all(|repository| repository.metadata.remote_workspace)
233 {
234 return Ok(ResumeRepositorySourcePreflight::Ready(
235 ResumeRepositorySourceReceipt {
236 session_id: session_id.to_owned(),
237 bundle_id: session.bundle_id.clone(),
238 checkpoint_sha256: verified.checkpoint_sha256,
239 repositories: Vec::new(),
240 },
241 ));
242 }
243 if verified.repositories.is_empty() {
244 return Ok(ResumeRepositorySourcePreflight::Ready(
245 ResumeRepositorySourceReceipt {
246 session_id: session_id.to_owned(),
247 bundle_id: session.bundle_id.clone(),
248 checkpoint_sha256: verified.checkpoint_sha256,
249 repositories: Vec::new(),
250 },
251 ));
252 }
253 let bundle = session
254 .project_bundle(&self.config)
255 .with_context(|| format!("session bundle {:?} is missing", session.bundle_id))?;
256 let configured = verified
257 .repositories
258 .iter()
259 .map(|archived| {
260 bundle
261 .repositories
262 .iter()
263 .find(|repository| repository.id == archived.metadata.id)
264 .cloned()
265 .with_context(|| {
266 format!(
267 "session bundle {:?} no longer contains repository {:?}",
268 session.bundle_id, archived.metadata.id
269 )
270 })
271 })
272 .collect::<Result<Vec<_>>>()?;
273 let github_token = configured
274 .iter()
275 .any(|repository| repository.github.is_some())
276 .then(controller_github_token)
277 .flatten();
278 let outcomes = verified
279 .repositories
280 .par_iter()
281 .zip(configured.par_iter())
282 .map(|(archived, configured)| {
283 if skip_repository_id == Some(configured.id.as_str()) {
284 return Ok(None);
285 }
286 checkpoint_source_missing_commit(
287 configured,
288 session
289 .project
290 .as_ref()
291 .and_then(|project| project.network_sources.get(&configured.id)),
292 archived,
293 executor,
294 github_token.as_deref(),
295 )
296 .map(|missing_commit| {
297 missing_commit.map(|missing_commit| ResumeRepositorySourceMismatch {
298 session_id: session_id.to_owned(),
299 bundle_id: session.bundle_id.clone(),
300 repository_id: configured.id.clone(),
301 missing_commit,
302 archived_origin: archived.metadata.origin.clone(),
303 configured_origin: configured.source_label(),
304 })
305 })
306 })
307 .collect::<Vec<Result<Option<ResumeRepositorySourceMismatch>>>>();
308 for outcome in outcomes {
309 if let Some(mismatch) = outcome? {
310 return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(mismatch));
311 }
312 }
313 Ok(ResumeRepositorySourcePreflight::Ready(
314 ResumeRepositorySourceReceipt {
315 session_id: session_id.to_owned(),
316 bundle_id: session.bundle_id.clone(),
317 checkpoint_sha256: verified.checkpoint_sha256,
318 repositories: configured,
319 },
320 ))
321 }
322
323 fn repository_source_receipt_is_current(
324 &self,
325 session_id: &str,
326 receipt: &ResumeRepositorySourceReceipt,
327 ) -> bool {
328 let Some(session) = self.state.sessions.get(session_id) else {
329 return false;
330 };
331 if receipt.session_id != session_id
332 || receipt.bundle_id != session.bundle_id
333 || session
334 .checkpoint
335 .as_ref()
336 .map(|checkpoint| &checkpoint.sha256)
337 != Some(&receipt.checkpoint_sha256)
338 {
339 return false;
340 }
341 if receipt.repositories.is_empty() {
342 return true;
343 }
344 let Some(bundle) = session.project_bundle(&self.config) else {
345 return false;
346 };
347 receipt.repositories.iter().all(|expected| {
348 bundle
349 .repositories
350 .iter()
351 .any(|configured| configured == expected)
352 })
353 }
354
355 pub fn replace_resume_repository_origin(
359 &mut self,
360 session_id: &str,
361 repository_id: &str,
362 replacement: &str,
363 executor: &(impl CommandExecutor + Sync),
364 ) -> Result<ResumeRepositorySourcePreflight> {
365 let session = self
366 .state
367 .sessions
368 .get(session_id)
369 .with_context(|| format!("unknown session {session_id}"))?;
370 let bundle_id = session.bundle_id.clone();
371 let checkpoint = session
372 .checkpoint
373 .as_ref()
374 .context("session has no checkpoint")?;
375 let replacement = replacement_repository_source(repository_id, replacement)?;
376 let replacement_network = session
377 .project
378 .as_ref()
379 .map(|_| mj_core::remote_git::resolve_repository(&replacement, executor))
380 .transpose()?;
381 let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
382 let verified = ResumeRepositoryBundles {
383 checkpoint_sha256: checkpoint.sha256.clone(),
384 repositories,
385 };
386 let archived = verified
387 .repositories
388 .iter()
389 .find(|repository| repository.metadata.id == repository_id)
390 .with_context(|| format!("checkpoint does not contain repository {repository_id:?}"))?;
391 if let Some(missing_commit) = checkpoint_source_missing_commit(
392 &replacement,
393 replacement_network.as_ref(),
394 archived,
395 executor,
396 controller_github_token().as_deref(),
397 )? {
398 return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(
399 ResumeRepositorySourceMismatch {
400 session_id: session_id.to_owned(),
401 bundle_id,
402 repository_id: repository_id.to_owned(),
403 missing_commit,
404 archived_origin: archived.metadata.origin.clone(),
405 configured_origin: replacement.source_label(),
406 },
407 ));
408 }
409 let previous = session.clone();
410 let mut project = session.project.clone();
411 if let Some(project) = &mut project {
412 let repository = project
413 .bundle
414 .repositories
415 .iter_mut()
416 .find(|repository| repository.id == repository_id)
417 .context("accepted project no longer contains the repaired repository")?;
418 repository.github = replacement.github.clone();
419 repository.local = replacement.local.clone();
420 project.network_sources.insert(
421 repository_id.to_owned(),
422 replacement_network
423 .clone()
424 .expect("accepted project source resolved above"),
425 );
426 }
427 let saved = crate::database::saved_project(&bundle_id)?;
430 let config_repository_id = saved
431 .as_ref()
432 .and_then(|(_, saved)| {
433 let accepted = previous.project.as_ref()?;
434 let original = accepted
435 .bundle
436 .repositories
437 .iter()
438 .find(|repository| repository.id == repository_id)?;
439 saved
440 .bundle
441 .repositories
442 .iter()
443 .find(|repository| {
444 repository.destination == original.destination
445 && saved.identities.get(&repository.id)
446 == accepted.identities.get(repository_id)
447 })
448 .map(|repository| repository.id.clone())
449 })
450 .unwrap_or_else(|| repository_id.to_owned());
451 let config_bundle_id = saved
452 .as_ref()
453 .map_or(bundle_id.as_str(), |(id, _)| id.as_str());
454 let (config, ()) = Config::update(|config| {
455 if project.is_some() && !config.bundles.contains_key(config_bundle_id) {
456 return Ok(());
457 }
458 let bundle = config
459 .bundles
460 .get_mut(config_bundle_id)
461 .with_context(|| format!("session bundle {bundle_id:?} is missing"))?;
462 let repository = bundle
463 .repositories
464 .iter_mut()
465 .find(|repository| repository.id == config_repository_id)
466 .with_context(|| {
467 format!(
468 "session bundle {:?} no longer contains repository {repository_id:?}",
469 bundle_id
470 )
471 })?;
472 repository.github = replacement.github.clone();
473 repository.local = replacement.local.clone();
474 Ok(())
475 })?;
476 self.config = config;
477 if let Some(project) = project {
478 self.state
479 .sessions
480 .get_mut(session_id)
481 .expect("session checked above")
482 .project = Some(project);
483 self.persist_session_transition_or_restore(
484 session_id,
485 &previous,
486 "save repaired project source",
487 )?;
488 }
489 self.preflight_verified_repository_sources(
490 session_id,
491 verified,
492 Some(repository_id),
493 false,
494 executor,
495 )
496 }
497}
498
499pub(super) enum WorkerRootReset {
507 FreshTarget,
510 InPlace {
514 previous_profile_root: String,
516 },
517}
518
519pub(super) struct RestoreIntoTarget<'a> {
525 pub profile: &'a mj_core::config::HarnessProfile,
527 pub archive: &'a VerifiedResumeArchive,
528 pub restored_archive: &'a Path,
531 pub resumed_project_directory: Option<PathBuf>,
532 pub resumed_container_workspace: Option<PathBuf>,
533 pub restore_repositories: bool,
534 pub native_continuity: bool,
535 pub discard_queued_prompts: bool,
536 pub replay_queue: bool,
538 pub utility_handoff: Option<String>,
541 pub projection_build: Option<tokio::task::JoinHandle<Result<MaterializedSession>>>,
542 pub resume_notices: Vec<String>,
544 pub install_attached_resources: bool,
545 pub worker_root_reset: WorkerRootReset,
546 pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
548}
549
550impl Controller {
551 pub(super) async fn restore_into_target(
554 &mut self,
555 session_id: &str,
556 restore: RestoreIntoTarget<'_>,
557 executor: &(impl CommandExecutor + Sync),
558 ) -> Result<MaterializedSession> {
559 crate::worker_lifecycle::run(session_id, "restore into target", executor, async {
560 let RestoreIntoTarget {
561 profile,
562 archive,
563 restored_archive,
564 resumed_project_directory,
565 resumed_container_workspace,
566 restore_repositories,
567 native_continuity,
568 discard_queued_prompts,
569 replay_queue,
570 utility_handoff,
571 projection_build,
572 mut resume_notices,
573 install_attached_resources: should_install_attached_resources,
574 worker_root_reset,
575 retire_after_ready,
576 } = restore;
577 let archive_manifest = &archive.manifest;
578 let canonical_session = &archive.canonical_session;
579 let stopped_subagents = crate::database::load_stopped_subagents(session_id)
584 .unwrap_or_else(|error| {
585 tracing::warn!(
586 session_id,
587 error = format!("{error:#}"),
588 "could not read the sub-agents this session's suspend stopped"
589 );
590 Vec::new()
591 });
592 let stopped_subagents_context =
593 mj_core::subagent::stopped_subagents_prompt_context(&stopped_subagents);
594 resume_notices.extend(mj_core::subagent::stopped_subagents_notice(
595 &stopped_subagents,
596 ));
597 let (backend, worker_root) = self.worker_placement(session_id)?;
598 let harness_home = target_profile_home(&backend, session_id, profile);
599 let workspace_root = if let Some(project_directory) = &resumed_project_directory {
600 project_directory
601 .parent()
602 .context("bare project directory has no parent")?
603 .to_string_lossy()
604 .into_owned()
605 } else {
606 super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
607 };
608 let target_path = |path: &str| match &backend {
609 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
610 if !path.starts_with('/') =>
611 {
612 PathBuf::from(format!("~/{path}"))
613 }
614 _ => PathBuf::from(path),
615 };
616 let remote_archive = format!("{worker_root}/restore.hel.zip");
617 let remote_spec = format!("{worker_root}/restore-spec.json");
618 use mj_checkpoint::checkpoint::QueueRestorePolicy;
619 let move_admission =
620 crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
621 operation.phase == mj_core::state::MovePhase::ResumingDestination
622 && !operation.queue_admission_started
623 && operation.queue == mj_core::state::ResumeQueueDisposition::Start
624 });
625 let queue_policy = if move_admission || (!native_continuity && replay_queue) {
626 QueueRestorePolicy::Defer
627 } else if discard_queued_prompts {
628 QueueRestorePolicy::Discard
629 } else {
630 QueueRestorePolicy::Restore
631 };
632 let restore = CheckpointRestoreSpec {
633 archive_path: restore_archive_path(
634 &backend,
635 restored_archive,
636 &target_path(&remote_archive),
637 ),
638 workspace_root: target_path(&workspace_root),
639 relay_root: target_path(&worker_root),
640 harness_home: target_path(&harness_home),
641 restore_repositories,
647 restore_native: native_continuity,
648 primary_repository_root: resumed_project_directory
657 .as_ref()
658 .map(|directory| target_path(&directory.to_string_lossy())),
659 queue_policy,
660 };
661 {
665 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
666 match &worker_root_reset {
667 WorkerRootReset::FreshTarget => {
672 if let Some(command) =
673 targets::clear_relay_state_plan(&backend, session_id)?
674 {
675 execute_checked(syncing, command)?;
676 }
677 execute_checked(
680 syncing,
681 targets::command_on_locator(
682 &backend,
683 session_id,
684 vec!["mkdir".into(), "-p".into(), worker_root.clone()],
685 "create the session worker root",
686 )?,
687 )?;
688 }
689 WorkerRootReset::InPlace {
694 previous_profile_root,
695 } => {
696 execute_checked(
697 syncing,
698 targets::in_place_worker_reset_plan(
699 &backend,
700 session_id,
701 previous_profile_root,
702 )?,
703 )?;
704 }
705 }
706 }
707 let staging = tempfile::tempdir().context("create restore staging")?;
708 let local_spec = staging.path().join("restore-spec.json");
709 std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
710 let controller = &*self;
714 let backend_ref = &backend;
715 let worker_root_ref = worker_root.as_str();
716 let local_spec_ref = local_spec.as_path();
717 crate::target_storage::ensure_room_for(
720 &backend,
721 [
722 worker_root.as_str(),
723 workspace_root.as_str(),
724 harness_home.as_str(),
725 ],
726 || {
727 Ok(std::fs::metadata(restored_archive)
728 .with_context(|| format!("measure {}", restored_archive.display()))?
729 .len())
730 },
731 "restore the checkpoint",
732 )?;
733 execute_concurrent_lanes(
734 || {
735 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
736 controller.prepare_worker_files(
737 session_id,
738 backend_ref,
739 worker_root_ref,
740 syncing,
741 )?;
742 super::provisioning::install_inherited_git_settings(
743 syncing,
744 backend_ref,
745 session_id,
746 )?;
747 Ok(())
748 },
749 || {
750 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
751 if should_upload_restore_archive(&backend) {
752 upload_checkpoint_spec(
753 restoring,
754 backend_ref,
755 session_id,
756 restored_archive,
757 &remote_archive,
758 )?;
759 }
760 upload_checkpoint_spec(
761 restoring,
762 backend_ref,
763 session_id,
764 local_spec_ref,
765 &remote_spec,
766 )
767 },
768 )?;
769 {
770 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
771 execute_checked(
772 restoring,
773 restore_command(&backend, session_id, &remote_spec)?,
774 )?;
775 }
776 if should_install_attached_resources {
777 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
778 install_attached_resources(
779 &self.state,
780 session_id,
781 &backend,
782 &worker_root,
783 syncing,
784 )?;
785 }
786 match projection_build {
787 Some(build) => {
788 let mut restored_projection = build
789 .await
790 .context("rebuild the restored projection")?
791 .context("rebuild the restored projection")?;
792 if discard_queued_prompts {
793 restored_projection.queued_prompts.clear();
794 }
795 crate::database::save_materialized_session(&restored_projection)?;
796 }
797 None if discard_queued_prompts => {
800 crate::database::replace_materialized_queued_prompts(session_id, &[])?;
801 }
802 None => {}
803 }
804 let readiness_stage = bridge_readiness_stage(profile);
805 let spec = self.reconnect_command(session_id)?;
806 let readiness = async {
807 let mut relay = {
808 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
809 start_worker_durably(
810 &crate::worker_lifecycle::require(session_id)?,
811 self.state.sessions[session_id]
812 .target
813 .as_ref()
814 .context("worker start has no durable target")?,
815 executor,
816 &backend,
817 &worker_root,
818 )?;
819 connect_started_worker(&spec, session_id, executor, &backend, &worker_root)
820 .await?
821 };
822 if let Some(context) = &stopped_subagents_context {
826 relay
827 .install_prompt_context(context.clone())
828 .await
829 .context("tell the resumed session which sub-agents its suspend stopped")?;
830 }
831 let native_session_id =
832 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
833 let owner = crate::worker_lifecycle::require(session_id)?;
834 crate::database::finish_worker_restart(session_id, owner.operation_id())?;
835 Ok::<_, anyhow::Error>((relay, native_session_id))
836 }
837 .await;
838 let (mut relay, native_session_id) = readiness
839 .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
840 if native_continuity {
841 if !restored_native_session_accepted(
842 &archive_manifest.session.native_session_id,
843 &native_session_id,
844 relay
845 .operational()
846 .replaced_unused_native_session_id
847 .as_deref(),
848 ) {
849 bail!(
850 "ACP loaded native session {native_session_id}, expected {}",
851 archive_manifest.session.native_session_id
852 );
853 }
854 } else {
855 relay
856 .install_prompt_context(
857 utility_handoff
858 .clone()
859 .context("a resume into a fresh native session has no handoff")?,
860 )
861 .await?;
862 if replay_queue {
863 for prompt in &canonical_session.queued_prompts {
864 let command = match &prompt.kind {
868 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
869 prompt: prompt
870 .content
871 .iter()
872 .cloned()
873 .map(serde_json::from_value)
874 .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
875 },
876 CanonicalQueuedCommandKind::SetConfig { key, value } => {
877 RelayCommand::SetConfig {
878 key: key.clone(),
879 value: value.clone(),
880 }
881 }
882 };
883 relay.submit(prompt.command_id.clone(), command).await?;
884 }
885 }
886 }
887 if let Some(worktree) = retire_after_ready
891 && let Err(error) = retire_managed_worktree(executor, worktree)
892 {
893 tracing::warn!(
894 session_id,
895 worktree = %worktree.worktree_root.display(),
896 error = format!("{error:#}"),
897 "could not retire the old managed worktree after resume"
898 );
899 resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
900 }
901 for notice in &resume_notices {
902 let submitted = async {
903 let command_id = new_command_id("resume-notice")?;
904 relay
905 .submit(
906 command_id,
907 RelayCommand::RecordNotice {
908 text: notice.clone(),
909 },
910 )
911 .await
912 }
913 .await;
914 if let Err(error) = submitted {
917 tracing::warn!(
918 session_id,
919 error = format!("{error:#}"),
920 "could not record a resume notice in the conversation"
921 );
922 }
923 }
924 self.mark_worker_connected(session_id, Some(native_session_id))?;
925 let materialized = relay.sync().await?.materialized;
926 if !stopped_subagents.is_empty() {
927 let delivered = stopped_subagents
928 .iter()
929 .map(|child| child.child_session_id.clone())
930 .collect::<Vec<_>>();
931 if let Err(error) = crate::database::clear_stopped_subagents(session_id, &delivered)
934 {
935 tracing::warn!(
936 session_id,
937 error = format!("{error:#}"),
938 "could not clear the stopped sub-agents after telling the resumed session"
939 );
940 }
941 }
942 Ok(materialized)
943 })
944 .await
945 }
946}
947
948pub(super) struct VerifiedResumeArchive {
955 pub archive_path: PathBuf,
959 pub manifest: mj_checkpoint::archive::ArchiveManifest,
960 pub canonical_session: Arc<CanonicalSessionSnapshot>,
961}
962
963pub(super) fn verify_resume_checkpoint(
969 session_id: &str,
970 checkpoint: &mj_core::state::CheckpointMetadata,
971) -> Result<VerifiedResumeArchive> {
972 let archive_path = {
973 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
974 checkpoint.archive_path.canonicalize().with_context(|| {
975 format!(
976 "resolve checkpoint archive {}",
977 checkpoint.archive_path.display()
978 )
979 })?
980 };
981 ensure!(
982 archive_path.is_absolute() && archive_path.is_file(),
983 "checkpoint archive path is not an absolute regular file: {}",
984 archive_path.display()
985 );
986 let mj_checkpoint::archive::VerifiedArchiveMetadata {
987 manifest,
988 canonical_session,
989 archive_sha256,
990 } = {
991 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
992 verify_archive_streaming(&archive_path)?
993 };
994 if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
995 bail!("persisted checkpoint verification failed");
996 }
997 ensure!(
998 manifest.repositories.iter().all(|repository| {
999 !repository.metadata.origin.starts_with("mj-local:")
1000 && !repository.metadata.origin.starts_with("ext::")
1001 }),
1002 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
1003 );
1004 Ok(VerifiedResumeArchive {
1005 archive_path,
1006 manifest,
1007 canonical_session: Arc::new(canonical_session),
1008 })
1009}
1010
1011pub fn raw_conversion_preview_for(
1012 session: &SessionRecord,
1013 config: &Config,
1014 executor: &(impl CommandExecutor + Sync),
1015) -> Result<mj_core::state::RawConversionPreview> {
1016 let conversion = plan_raw_to_workspace(session, config, executor)?;
1017 raw_conversion_preview(session, &conversion, executor)
1018}
1019
1020fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
1021 let replacement = replacement.trim();
1022 ensure!(!replacement.is_empty(), "enter the repository's new origin");
1023 let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
1024 let path = expanded.as_path();
1025 let (github, local) = if path.is_absolute() {
1026 ensure!(
1027 path.is_dir(),
1028 "local repository {replacement:?} is not a directory"
1029 );
1030 (None, Some(mj_core::local_git::canonical_repository(path)?))
1031 } else {
1032 let github = crate::setup::github_repository_from_origin(replacement)
1033 .context("origin must be a GitHub repository or an absolute local repository path")?;
1034 (
1035 Some(format!("{}/{}", github.owner, github.repository)),
1036 None,
1037 )
1038 };
1039 Ok(ProjectRepository {
1040 id: id.to_owned(),
1041 github,
1042 local,
1043 destination: PathBuf::from(id),
1044 git_ref: None,
1045 })
1046}
1047
1048fn checkpoint_source_missing_commit(
1049 configured: &ProjectRepository,
1050 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1051 archived: &CheckpointRepositoryBundle,
1052 executor: &impl CommandExecutor,
1053 github_token: Option<&str>,
1054) -> Result<Option<String>> {
1055 let staging = tempfile::tempdir().context("create repository source preflight")?;
1056 let repository = staging.path().join("repository.git");
1057 checked_preflight_git(
1058 executor,
1059 CommandSpec::new(
1060 "git",
1061 [
1062 "init".to_owned(),
1063 "--bare".to_owned(),
1064 "--quiet".to_owned(),
1065 repository.to_string_lossy().into_owned(),
1066 ],
1067 )
1068 .purpose("initialize repository source preflight"),
1069 )?;
1070 let missing = checkpoint_bundle_prerequisites(archived)?;
1071 if missing.is_empty() {
1072 let bundle = staging.path().join("checkpoint.bundle");
1073 std::fs::write(&bundle, &archived.committed_bundle)
1074 .context("write self-contained checkpoint bundle for source preflight")?;
1075 checked_preflight_git(
1076 executor,
1077 checkpoint_bundle_import_command(&repository, &bundle),
1078 )?;
1079 return Ok(None);
1080 }
1081 for commit in missing {
1085 let output = fetch_source_commit(
1086 executor,
1087 &repository,
1088 configured,
1089 accepted_source,
1090 &commit,
1091 github_token,
1092 )?;
1093 if output.status != 0 {
1094 let stderr = String::from_utf8_lossy(&output.stderr);
1095 if source_does_not_have_commit(&stderr) {
1096 return Ok(Some(commit));
1097 }
1098 bail!(
1099 "could not check configured source {:?}: {}",
1100 configured.source_label(),
1101 stderr.trim()
1102 );
1103 }
1104 }
1105 Ok(None)
1109}
1110
1111fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
1112 let mut command = CommandSpec::new(
1113 "git",
1114 [
1115 "-C".to_owned(),
1116 repository.to_string_lossy().into_owned(),
1117 "fetch".to_owned(),
1118 "--no-tags".to_owned(),
1119 bundle.to_string_lossy().into_owned(),
1120 "HEAD".to_owned(),
1121 ],
1122 )
1123 .purpose("validate self-contained checkpoint bundle");
1124 command
1125 .env
1126 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1127 command
1128 .env
1129 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1130 command
1131}
1132
1133fn fetch_source_commit(
1134 executor: &impl CommandExecutor,
1135 repository: &Path,
1136 configured: &ProjectRepository,
1137 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1138 commit: &str,
1139 github_token: Option<&str>,
1140) -> Result<CommandOutput> {
1141 let mut arguments = Vec::new();
1142 let mut token_auth = false;
1143 let mut ssh_transport = false;
1144 let source = if let Some(source) = accepted_source {
1145 source.fetch_url.clone()
1146 } else if let Some(local) = &configured.local {
1147 local.to_string_lossy().into_owned()
1148 } else {
1149 let source = configured
1150 .github
1151 .as_deref()
1152 .context("repository source is missing")?;
1153 let github = crate::setup::github_repository_from_origin(source)
1154 .context("configured repository is not a GitHub source")?;
1155 if github_token.is_some() {
1156 token_auth = true;
1157 arguments.extend([
1158 "-c".to_owned(),
1159 "credential.helper=".to_owned(),
1160 "-c".to_owned(),
1161 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1162 ]);
1163 format!(
1164 "https://github.com/{}/{}.git",
1165 github.owner, github.repository
1166 )
1167 } else {
1168 ssh_transport = true;
1169 format!("git@github.com:{}/{}.git", github.owner, github.repository)
1170 }
1171 };
1172 if accepted_source.is_some() {
1173 if source.starts_with("https://github.com/") && github_token.is_some() {
1174 token_auth = true;
1175 arguments.extend([
1176 "-c".to_owned(), "credential.helper=".to_owned(), "-c".to_owned(),
1177 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1178 ]);
1179 }
1180 ssh_transport = source.starts_with("git@") || source.starts_with("ssh://");
1181 }
1182 arguments.extend([
1183 "-C".to_owned(),
1184 repository.to_string_lossy().into_owned(),
1185 "fetch".to_owned(),
1186 "--no-tags".to_owned(),
1187 "--depth=1".to_owned(),
1188 "--filter=blob:none".to_owned(),
1189 source,
1190 commit.to_owned(),
1191 ]);
1192 let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1193 command
1194 .env
1195 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1196 command
1197 .env
1198 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1199 if token_auth {
1200 let token = github_token.expect("token authentication requires a GitHub token");
1201 command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1202 }
1203 if ssh_transport {
1204 command.env.insert(
1205 "GIT_SSH_COMMAND".to_owned(),
1206 "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1207 .to_owned(),
1208 );
1209 }
1210 executor.execute(&command)
1211}
1212
1213fn source_does_not_have_commit(stderr: &str) -> bool {
1214 let stderr = stderr.to_ascii_lowercase();
1215 [
1216 "not our ref",
1217 "couldn't find remote ref",
1218 "not a valid object name",
1219 "no such ref was fetched",
1220 ]
1221 .iter()
1222 .any(|needle| stderr.contains(needle))
1223}
1224
1225fn checked_preflight_git(
1226 executor: &impl CommandExecutor,
1227 command: CommandSpec,
1228) -> Result<CommandOutput> {
1229 let output = executor.execute(&command)?;
1230 ensure!(
1231 output.status == 0,
1232 "{}: {}",
1233 command.purpose,
1234 String::from_utf8_lossy(&output.stderr).trim()
1235 );
1236 Ok(output)
1237}
1238
1239impl Controller {
1240 pub async fn resume_session_with_options(
1245 &mut self,
1246 session_id: &str,
1247 profile_id: &str,
1248 target_id: &str,
1249 additional_mounts: Option<Vec<AdditionalMount>>,
1250 resource_allocation: Option<SessionResourceAllocation>,
1251 ) -> Result<MaterializedSession> {
1252 self.resume_session_with_options_and_queue_disposition(
1253 session_id,
1254 profile_id,
1255 target_id,
1256 additional_mounts,
1257 resource_allocation,
1258 false,
1259 )
1260 .await
1261 }
1262
1263 pub async fn resume_session_with_options_and_queue_disposition(
1264 &mut self,
1265 session_id: &str,
1266 profile_id: &str,
1267 target_id: &str,
1268 additional_mounts: Option<Vec<AdditionalMount>>,
1269 resource_allocation: Option<SessionResourceAllocation>,
1270 discard_queue: bool,
1271 ) -> Result<MaterializedSession> {
1272 self.resume_session_controlled(
1273 session_id,
1274 profile_id,
1275 target_id,
1276 SessionResumeOptions {
1277 additional_mounts,
1278 resource_allocation,
1279 discard_queue,
1280 },
1281 &ProcessExecutor,
1282 )
1283 .await
1284 }
1285
1286 pub async fn resume_session_controlled(
1287 &mut self,
1288 session_id: &str,
1289 profile_id: &str,
1290 target_id: &str,
1291 options: SessionResumeOptions,
1292 executor: &(impl CommandExecutor + Sync),
1293 ) -> Result<MaterializedSession> {
1294 self.resume_session_controlled_with_repository_preflight(
1295 session_id, profile_id, target_id, options, None, executor,
1296 )
1297 .await
1298 }
1299
1300 pub async fn resume_session_controlled_with_repository_preflight(
1301 &mut self,
1302 session_id: &str,
1303 profile_id: &str,
1304 target_id: &str,
1305 options: SessionResumeOptions,
1306 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1307 executor: &(impl CommandExecutor + Sync),
1308 ) -> Result<MaterializedSession> {
1309 self.resume_session_with_origin(
1310 session_id,
1311 profile_id,
1312 target_id,
1313 options,
1314 repository_preflight,
1315 None,
1316 executor,
1317 )
1318 .await
1319 }
1320
1321 pub(in crate::controller) async fn resume_session_for_move(
1322 &mut self,
1323 operation: &mj_core::state::MoveOperation,
1324 executor: &(impl CommandExecutor + Sync),
1325 ) -> Result<MaterializedSession> {
1326 self.resume_session_with_origin(
1327 &operation.selection.session_id,
1328 operation
1329 .selection
1330 .profile_id
1331 .as_deref()
1332 .context("Move profile missing")?,
1333 operation
1334 .selection
1335 .target_template_id
1336 .as_deref()
1337 .context("Move target missing")?,
1338 SessionResumeOptions {
1339 additional_mounts: operation.selection.additional_mounts.clone(),
1340 resource_allocation: operation.selection.resource_allocation.clone(),
1341 discard_queue: true,
1342 },
1343 None,
1344 Some(operation),
1345 executor,
1346 )
1347 .await
1348 }
1349
1350 #[allow(clippy::too_many_arguments)]
1351 async fn resume_session_with_origin(
1352 &mut self,
1353 session_id: &str,
1354 profile_id: &str,
1355 target_id: &str,
1356 options: SessionResumeOptions,
1357 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1358 move_operation: Option<&mj_core::state::MoveOperation>,
1359 executor: &(impl CommandExecutor + Sync),
1360 ) -> Result<MaterializedSession> {
1361 crate::worker_lifecycle::run(session_id, "resume session with origin", executor, async {
1362 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1363 let transferring_workspace = move_operation.is_some();
1364 if let Some(operation) = crate::database::load_move_operation(session_id)? {
1365 ensure!(
1366 transferring_workspace || !operation.holds_source_environment(),
1367 "a profile switch retains this environment; retry Move instead of recreating it with Resume"
1368 );
1369 }
1370 let SessionResumeOptions {
1371 additional_mounts,
1372 resource_allocation,
1373 discard_queue,
1374 } = options;
1375 let mut previous = self
1376 .state
1377 .sessions
1378 .get(session_id)
1379 .with_context(|| format!("unknown session {session_id}"))?
1380 .clone();
1381 if !(matches!(
1382 previous.state,
1383 SessionState::Stopped | SessionState::Lost | SessionState::Error
1384 ) || (transferring_workspace && previous.state == SessionState::Closing))
1385 {
1386 bail!("session {session_id} is not stopped, lost, or retryable");
1387 }
1388 let checkpoint = move_operation
1389 .and_then(|op| op.handoff.as_ref())
1390 .or(previous.checkpoint.as_ref())
1391 .context("session has no checkpoint")?;
1392 let moving_to_raw = previous.project_directory.is_none()
1395 && self
1396 .config
1397 .targets
1398 .get(target_id)
1399 .is_some_and(mj_core::config::is_bare_project_target);
1400 if !transferring_workspace
1401 && (moving_to_raw
1402 || !repository_preflight.as_ref().is_some_and(|receipt| {
1403 self.repository_source_receipt_is_current(session_id, receipt)
1404 }))
1405 {
1406 let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1407 if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1408 self.preflight_repository_sources(session_id, target_id, false, executor)?
1409 {
1410 bail!(
1411 "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1412 mismatch.missing_commit,
1413 mismatch.configured_origin,
1414 mismatch.repository_id,
1415 mismatch.archived_origin,
1416 );
1417 }
1418 }
1419 let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1420 let archive_path = verified_archive.archive_path.clone();
1421 let archive_manifest = &verified_archive.manifest;
1422 let canonical_session = Arc::clone(&verified_archive.canonical_session);
1423 let profile = self
1424 .config
1425 .profiles
1426 .get(profile_id)
1427 .with_context(|| format!("unknown profile {profile_id:?}"))?
1428 .clone();
1429 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1430 let target_template = self
1431 .config
1432 .targets
1433 .get(target_id)
1434 .with_context(|| format!("unknown target template {target_id:?}"))?
1435 .clone();
1436 self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1439 ensure!(
1440 profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1441 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1442 );
1443 let plan = resume_compatibility(&previous, &self.config, target_id)
1444 .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1445 if plan == ResumePlan::InPlace
1448 && let Some(worktree) = previous.managed_worktree.as_mut()
1449 && let Ok(current) = managed_worktree_target(&target_template)
1450 {
1451 worktree.target = current;
1452 }
1453 if !mj_core::config::is_bare_project_target(&target_template)
1457 && plan != ResumePlan::RawToWorkspace
1458 && !transferring_workspace
1459 {
1460 super::network_git::bundle_from_manifest(archive_manifest)?;
1461 }
1462 if plan == ResumePlan::InPlace
1463 && previous.managed_worktree.is_none()
1464 && let Some(project_directory) = &previous.project_directory
1465 {
1466 self.validate_project_directory(target_id, project_directory, executor)
1467 .context("raw project is unavailable for resume")?;
1468 }
1469 let conversion = match plan {
1470 ResumePlan::InPlace => None,
1471 ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(
1472 plan_raw_to_workspace(&previous, &self.config, executor)
1473 .context("prepare the raw checkout for its new target")?,
1474 )),
1475 ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1476 self.plan_workspace_to_raw(&previous, target_id, executor)
1477 .context("prepare a checkout for this session")?,
1478 )),
1479 };
1480 let resource_allocation =
1481 resource_allocation.or_else(|| previous.resource_allocation.clone());
1482 let additional_mounts =
1483 additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1484 validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1485 let selected_container_size =
1486 selected_host_container_size(&target_template, resource_allocation.as_ref());
1487 if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1488 bail!("attached resources are unsupported for this target");
1489 }
1490 targets::validate_additional_mounts(&additional_mounts)?;
1491 let history_host = mount_history_host(&target_template);
1492 let history_mounts = additional_mounts.clone();
1493 if !transferring_workspace
1494 && previous.state == SessionState::Error
1495 && let Some(locator) = &previous.target
1496 {
1497 let backend = backend_locator(locator, &previous, &self.config)?;
1498 targets::close_plan(&backend, session_id)?
1499 .execute(executor)
1500 .context("clean up target from failed resume")?;
1501 }
1502 let mut resume_notices = Vec::new();
1503 if let Some(conversion) = conversion
1506 .as_ref()
1507 .and_then(ResumeConversion::workspace_to_raw)
1508 {
1509 resume_notices.push(format!(
1510 "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1511 previous.target_template_id,
1512 conversion.worktree.worktree_root.display(),
1513 archive_manifest
1514 .repositories
1515 .first()
1516 .and_then(|repository| repository.metadata.branch.as_deref())
1517 .unwrap_or("a detached head"),
1518 conversion.worktree.branch,
1519 ));
1520 }
1521 let managed_checkout_present = previous
1522 .managed_worktree
1523 .as_ref()
1524 .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1525 .transpose()?
1526 .unwrap_or(true);
1527 if managed_checkout_present && let Some(project_directory) = &previous.project_directory {
1530 match raw_checkout_position(&previous, &self.config, project_directory, executor) {
1531 Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1532 project_directory,
1533 archive_manifest
1534 .repositories
1535 .first()
1536 .map(|repository| &repository.metadata),
1537 &live,
1538 )),
1539 Err(error) => tracing::warn!(
1542 session_id,
1543 error = format!("{error:#}"),
1544 "could not read the raw checkout position for a resume notice"
1545 ),
1546 }
1547 }
1548 super::worker_binary::preflight_worker_binary(&target_template, executor)?;
1553 let same_harness = profile.kind == archive_manifest.session.harness_kind;
1554 let native_continuity =
1555 native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1556 let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1557 let utility_config = (!native_continuity).then(|| self.config.clone());
1561 let discard_queued_prompts = discard_queue || !same_harness;
1562 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1566 .unwrap_or_else(|error| {
1567 tracing::warn!(
1568 session_id,
1569 error = format!("{error:#}"),
1570 "could not read the stored projection frontier; rebuilding it from the archive"
1571 );
1572 None
1573 });
1574 let rebuild_projection = projection_rebuild_required(
1575 stored_frontier
1576 .as_ref()
1577 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1578 canonical_session.event_frontier,
1579 &canonical_session.event_frontier_digest,
1580 );
1581 let projection_build = rebuild_projection.then(|| {
1591 let canonical = Arc::clone(&canonical_session);
1592 let session_id = session_id.to_owned();
1593 tokio::task::spawn_blocking(move || {
1594 materialized_session_from_canonical(session_id, &canonical)
1595 })
1596 });
1597 let github_token = controller_github_token();
1598
1599 if let Some(conversion) = conversion
1602 .as_ref()
1603 .and_then(ResumeConversion::raw_to_workspace)
1604 && let Some(bundle) = &conversion.new_bundle
1605 {
1606 let (config, ()) = Config::update(|config| {
1607 if let Some(existing) = config.bundles.get(&conversion.bundle_id) {
1608 ensure!(
1609 existing == bundle,
1610 "bundle {:?} was configured concurrently with a different definition; retry the resume",
1611 conversion.bundle_id
1612 );
1613 } else {
1614 config
1615 .bundles
1616 .insert(conversion.bundle_id.clone(), bundle.clone());
1617 }
1618 Ok(())
1619 })
1620 .context("save the bundle for a converted raw session")?;
1621 self.config = config;
1622 }
1623
1624 let moving_into_first_container = self
1630 .config
1631 .targets
1632 .get(target_id)
1633 .is_some_and(mj_core::config::is_container_target)
1634 && !self
1635 .config
1636 .targets
1637 .get(&previous.target_template_id)
1638 .is_some_and(mj_core::config::is_container_target);
1639 let converted_project = conversion
1640 .as_ref()
1641 .and_then(ResumeConversion::raw_to_workspace)
1642 .map(|conversion| {
1643 crate::project_catalog::snapshot(
1644 self.config
1645 .bundles
1646 .get(&conversion.bundle_id)
1647 .context("converted project is missing")?,
1648 executor,
1649 true,
1650 )
1651 })
1652 .transpose()?;
1653 let record = self.state.sessions.get_mut(session_id).unwrap();
1654 if record.container_workspace.is_none() && moving_into_first_container {
1655 record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1656 }
1657 record.harness_kind = profile.kind;
1658 record.last_profile = profile_id.to_string();
1659 record.target_template_id = target_id.to_string();
1660 record.target_runtime = Some(
1661 self.config
1662 .targets
1663 .get(target_id)
1664 .context("resume target disappeared")?
1665 .into(),
1666 );
1667 record.resource_allocation = resource_allocation;
1668 record.additional_mounts = additional_mounts;
1669 record.target = None;
1670 record.native_session_id =
1671 native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1672 record.state = SessionState::Provisioning;
1673 record.updated_at = now();
1674 record.last_error = None;
1675 match &conversion {
1676 Some(ResumeConversion::RawToWorkspace(conversion)) => {
1677 apply_raw_to_workspace(record, conversion);
1678 record.project = converted_project;
1679 }
1680 Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1681 apply_workspace_to_raw(record, conversion);
1682 }
1683 None => {
1684 if let (Some(worktree), Some(refreshed)) = (
1685 record.managed_worktree.as_mut(),
1686 previous.managed_worktree.as_ref(),
1687 ) {
1688 worktree.target = refreshed.target.clone();
1689 }
1690 }
1691 }
1692 let resumed_project_directory = record.project_directory.clone();
1693 let resumed_container_workspace = record.container_workspace.clone();
1694 if let Some(host) = history_host {
1695 self.state.remember_mount_sources(host, &history_mounts);
1696 crate::database::remember_mount_sources(host, &history_mounts)?;
1697 }
1698 if let Some(conversion) = conversion
1701 .as_ref()
1702 .and_then(ResumeConversion::raw_to_workspace)
1703 {
1704 let bundle_id =
1709 crate::database::rebind_session_bundle(session_id, &conversion.bundle_id)?;
1710 self.state
1711 .sessions
1712 .get_mut(session_id)
1713 .expect("the resumed session is in state")
1714 .bundle_id = bundle_id;
1715 }
1716 crate::database::save_resumed_session(
1719 &self.state.sessions[session_id],
1720 selected_container_size
1721 .as_ref()
1722 .map(|(host, size)| (host.as_str(), *size)),
1723 )?;
1724 if let Some((host, size)) = selected_container_size.as_ref() {
1725 self.state.remember_container_size(host, *size);
1726 }
1727
1728 let mut recreated_managed_worktree = false;
1729 let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1732 let result = async {
1733 if let Some(worktree) = previous.managed_worktree.as_ref() {
1734 recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1735 if !transferring_workspace && recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1736 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1737 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1738 &archive_path, &worktree.worktree_root, &SystemGit,
1739 )?;
1740 } else {
1741 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1742 &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1743 )?;
1744 }
1745 }
1746 }
1747 if let Some(conversion) = conversion
1750 .as_ref()
1751 .and_then(ResumeConversion::workspace_to_raw)
1752 {
1753 if conversion.reuse_existing_branch {
1754 let recovery_ref =
1755 preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1756 restore_managed_worktree(executor, &conversion.worktree)?;
1757 resume_notices.push(format!(
1758 "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1759 ));
1760 } else {
1761 create_managed_worktree(
1762 executor,
1763 &conversion.worktree,
1764 None,
1765 PrimaryCheckoutRequirement::Any,
1766 )?;
1767 }
1768 if transferring_workspace {
1769 } else if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1771 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1772 &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1773 )?;
1774 } else {
1775 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1776 &archive_path,
1777 &conversion.worktree.worktree_root,
1778 &conversion.worktree.branch,
1779 &SystemGit,
1780 )?;
1781 }
1782 }
1783 if let Some(conversion) = conversion
1788 .as_ref()
1789 .and_then(ResumeConversion::raw_to_workspace)
1790 && !transferring_workspace
1791 {
1792 let destination = PathBuf::from(
1793 previous
1794 .project_directory
1795 .as_deref()
1796 .context("a raw session has no project directory")?
1797 .file_name()
1798 .context("a raw project directory cannot be the filesystem root")?,
1799 );
1800 let snapshot = raw_checkout_snapshot(
1801 &conversion.checkout,
1802 &conversion.source,
1803 &destination,
1804 &SystemGit,
1805 conversion.retire.as_ref().is_some_and(|checkout| {
1806 checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1807 }),
1808 )
1809 .context("snapshot the host checkout for its new target")?;
1810 resume_notices.push(conversion_notice(
1811 target_id,
1812 previous
1813 .project_directory
1814 .as_deref()
1815 .unwrap_or(&conversion.checkout),
1816 snapshot.metadata.branch.as_deref(),
1817 conversion.retire.as_ref(),
1818 ));
1819 let archives = mj_core::config::sessions_dir();
1820 std::fs::create_dir_all(&archives).with_context(|| {
1821 format!("create the checkpoint directory {}", archives.display())
1822 })?;
1823 let output = archives.join(format!(
1827 "{session_id}-converted-{}-{}.hel.zip",
1828 previous
1829 .checkpoint
1830 .as_ref()
1831 .map_or(0, |checkpoint| checkpoint.event_frontier),
1832 new_command_id("archive")?
1833 ));
1834 let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1835 conversion_checkpoint_written = Some(written.clone());
1836 let record = self.state.sessions.get_mut(session_id).unwrap();
1837 record.checkpoint = Some(written);
1838 record.updated_at = now();
1839 crate::database::save_resumed_session(
1842 &self.state.sessions[session_id],
1843 selected_container_size.as_ref().map(|(host, size)| (host.as_str(), *size)),
1844 )?;
1845 }
1846 let utility_handoff = {
1847 let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1848 if let Some(config) = utility_config.as_ref() {
1849 Some(
1850 provision_with_cross_harness_handoff(
1851 self,
1852 session_id,
1853 executor,
1854 github_token.as_deref(),
1855 config,
1856 &canonical_session,
1857 context_bytes,
1858 )
1859 .context("prepare the cross-harness destination")?,
1860 )
1861 } else {
1862 self.provision_session_with_failure_disposition(
1863 session_id,
1864 executor,
1865 github_token.as_deref(),
1866 ProvisioningFailureDisposition::Preserve,
1867 )
1868 .await?;
1869 None
1870 }
1871 };
1872 if transferring_workspace {
1873 self.restore_move_workspace(session_id, executor)?;
1874 }
1875 let restore_repositories = !transferring_workspace && ((resumed_project_directory.is_none()
1876 && conversion.is_none())
1877 || plan == ResumePlan::RawToWorkspace
1878 || (recreated_managed_worktree && plan == ResumePlan::InPlace));
1879 let restored_archive = conversion_checkpoint_written
1882 .as_ref()
1883 .map_or(archive_path.as_path(), |checkpoint| {
1884 checkpoint.archive_path.as_path()
1885 });
1886 self.restore_into_target(
1887 session_id,
1888 RestoreIntoTarget {
1889 profile: &profile,
1890 archive: &verified_archive,
1891 restored_archive,
1892 resumed_project_directory,
1893 resumed_container_workspace,
1894 restore_repositories,
1895 native_continuity,
1896 discard_queued_prompts,
1897 replay_queue: !discard_queue,
1898 utility_handoff,
1899 projection_build,
1900 resume_notices,
1901 install_attached_resources: true,
1902 worker_root_reset: WorkerRootReset::FreshTarget,
1903 retire_after_ready: conversion
1904 .as_ref()
1905 .filter(|_| !transferring_workspace)
1906 .and_then(ResumeConversion::raw_to_workspace)
1907 .and_then(|plan| plan.retire.as_ref()),
1908 },
1909 executor,
1910 )
1911 .await
1912 }
1913 .await;
1914 match result {
1915 Ok(materialized) => {
1916 if let Some(written) = &conversion_checkpoint_written {
1919 super::checkpoint::prune_replaced_checkpoint(
1920 previous.checkpoint.as_ref(),
1921 written,
1922 );
1923 }
1924 Ok(materialized)
1925 }
1926 Err(error) => {
1927 if let Some(written) = &conversion_checkpoint_written
1930 && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
1931 && remove_error.kind() != std::io::ErrorKind::NotFound
1932 {
1933 tracing::warn!(
1934 session_id,
1935 path = %written.archive_path.display(),
1936 "could not remove the conversion checkpoint after resume failed: {remove_error}"
1937 );
1938 }
1939 restore_projection_after_failed_resume(
1942 session_id,
1943 &canonical_session,
1944 discard_queued_prompts,
1945 );
1946 if let Some(operation) = move_operation {
1947 return Err(self.rollback_move_destination(operation, error, executor)?);
1948 }
1949 Err(self.rollback_failed_resume(
1950 session_id,
1951 &previous,
1952 recreated_managed_worktree,
1953 error,
1954 executor,
1955 )?)
1956 }
1957 }
1958
1959 }).await
1960 }
1961
1962 pub(super) fn rollback_failed_resume(
1963 &mut self,
1964 session_id: &str,
1965 previous: &SessionRecord,
1966 recreated_managed_worktree: bool,
1967 error: anyhow::Error,
1968 _executor: &impl CommandExecutor,
1969 ) -> Result<anyhow::Error> {
1970 let current = self
1971 .state
1972 .sessions
1973 .get(session_id)
1974 .with_context(|| format!("unknown session {session_id}"))?
1975 .clone();
1976 let cleanup = match current.target.as_ref() {
1977 Some(locator) => (|| -> Result<()> {
1978 let backend = backend_locator(locator, ¤t, &self.config)?;
1979 targets::close_plan(&backend, session_id)?
1980 .execute(&CancellableProcessExecutor::with_timeout(
1983 Duration::from_secs(15),
1984 ))
1985 .map(|_| ())
1986 })(),
1987 None => Ok(()),
1988 };
1989 let worktree_cleanup = if cleanup.is_err() {
1992 Ok(())
1993 } else {
1994 match (
1995 current.managed_worktree.as_ref(),
1996 previous.managed_worktree.as_ref(),
1997 ) {
1998 (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
1999 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
2000 previous,
2001 ),
2002 (Some(current), Some(previous)) if current == previous => Ok(()),
2003 (Some(worktree), _) => cleanup_managed_worktree(
2006 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
2007 worktree,
2008 crate::controller::BranchDisposition::Delete,
2009 ),
2010 (None, _) => Ok(()),
2011 }
2012 };
2013 let cleanup_error = [cleanup, worktree_cleanup]
2014 .into_iter()
2015 .filter_map(Result::err)
2016 .map(|cleanup_error| format!("{cleanup_error:#}"))
2017 .collect::<Vec<_>>()
2018 .join("; ");
2019 if !cleanup_error.is_empty() {
2020 tracing::warn!(
2021 session_id,
2022 error = %cleanup_error,
2023 "resume rollback cleanup reported failures"
2024 );
2025 }
2026 let original = super::worker_binary::failure_line(&error);
2029 let detail = format!("{error:#}");
2030 if detail != original {
2031 tracing::warn!(session_id, error = %detail, "resume failed");
2032 }
2033 let record = self.state.sessions.get_mut(session_id).unwrap();
2034 let failure = apply_failed_resume_rollback(
2035 record,
2036 previous,
2037 &original,
2038 (!cleanup_error.is_empty()).then_some(cleanup_error),
2039 );
2040 if record.bundle_id != current.bundle_id {
2043 let bundle_id = record.bundle_id.clone();
2044 record.bundle_id = crate::database::rebind_session_bundle(session_id, &bundle_id)?;
2045 }
2046 crate::database::save_resumed_session(&self.state.sessions[session_id], None)?;
2049 Ok(failure)
2050 }
2051}
2052
2053fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
2054 format!(
2055 "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
2056 worktree_root.display(),
2057 worktree_root.display()
2058 )
2059}
2060
2061pub(super) fn apply_failed_resume_rollback(
2062 current: &mut SessionRecord,
2063 previous: &SessionRecord,
2064 original_error: &str,
2065 cleanup_error: Option<String>,
2066) -> anyhow::Error {
2067 match cleanup_error {
2068 None => {
2069 *current = previous.clone();
2070 current.state = SessionState::Stopped;
2071 current.target = None;
2072 current.updated_at = now();
2073 current.last_error = Some(format!("resume failed: {original_error}"));
2074 anyhow::anyhow!(original_error.to_owned())
2075 }
2076 Some(cleanup_error) => {
2077 let failure = format!(
2078 "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
2079 );
2080 if current.managed_worktree.is_none() {
2085 current
2086 .project_directory
2087 .clone_from(&previous.project_directory);
2088 current
2089 .managed_worktree
2090 .clone_from(&previous.managed_worktree);
2091 current.bundle_id.clone_from(&previous.bundle_id);
2092 }
2093 current.state = SessionState::Error;
2094 current.updated_at = now();
2095 current.last_error = Some(format!("resume failed: {failure}"));
2096 anyhow::anyhow!(failure)
2097 }
2098 }
2099}
2100
2101fn conversion_notice(
2103 target_id: &str,
2104 checkout: &Path,
2105 branch: Option<&str>,
2106 retire: Option<&mj_core::state::ManagedWorktree>,
2107) -> String {
2108 let branch = branch.unwrap_or("a detached head");
2109 match retire {
2110 Some(worktree) => format!(
2111 "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
2112 checkout.display(),
2113 worktree.branch,
2114 worktree.source_repository.display()
2115 ),
2116 None => format!(
2117 "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.",
2118 checkout.display()
2119 ),
2120 }
2121}
2122
2123fn conversion_checkpoint(
2131 previous_archive: &Path,
2132 snapshot: mj_checkpoint::archive::RepositorySnapshot,
2133 output: &Path,
2134) -> Result<mj_core::state::CheckpointMetadata> {
2135 let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
2136 .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
2137 let native_artifacts = previous
2140 .manifest
2141 .payloads
2142 .iter()
2143 .filter_map(|descriptor| match &descriptor.role {
2144 mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
2145 Some((relative_path, descriptor))
2146 }
2147 _ => None,
2148 })
2149 .map(|(relative_path, descriptor)| {
2150 Ok(mj_checkpoint::archive::NativeArtifact {
2151 relative_path: relative_path.clone(),
2152 data: previous.payload(descriptor)?.to_vec(),
2153 mode: descriptor.mode,
2154 })
2155 })
2156 .collect::<Result<Vec<_>>>()?;
2157 let canonical_session = previous.canonical_session()?;
2158 let event_frontier = canonical_session.event_frontier;
2159 let written = mj_checkpoint::archive::write_archive_atomic(
2160 output,
2161 &mj_checkpoint::archive::ArchiveInput {
2162 session: previous.manifest.session.clone(),
2163 target: previous.manifest.target.clone(),
2166 bundle: mj_checkpoint::archive::BundleManifest {
2167 id: previous.manifest.bundle.id.clone(),
2168 primary_repository: snapshot.metadata.id.clone(),
2172 },
2173 canonical_session,
2174 native_artifacts,
2175 repositories: vec![snapshot],
2176 },
2177 )
2178 .with_context(|| format!("write the conversion archive {}", output.display()))?;
2179 Ok(mj_core::state::CheckpointMetadata {
2180 archive_path: output.to_path_buf(),
2181 sha256: written.archive_sha256,
2182 created_at: now(),
2183 event_frontier,
2184 })
2185}
2186
2187fn restored_native_session_accepted(
2198 archived: &str,
2199 opened: &str,
2200 replaced_unused: Option<&str>,
2201) -> bool {
2202 opened == archived || replaced_unused == Some(archived)
2203}
2204
2205fn restore_projection_after_failed_resume(
2216 session_id: &str,
2217 canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
2218 discard_queued_prompts: bool,
2219) {
2220 let stored_frontier = crate::database::materialized_event_frontier(session_id)
2221 .unwrap_or_else(|error| {
2222 tracing::warn!(
2223 session_id,
2224 error = format!("{error:#}"),
2225 "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
2226 );
2227 None
2228 });
2229 if projection_rebuild_required(
2230 stored_frontier
2231 .as_ref()
2232 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
2233 canonical_session.event_frontier,
2234 &canonical_session.event_frontier_digest,
2235 ) {
2236 match materialized_session_from_canonical(session_id, canonical_session) {
2237 Ok(previous_projection) => {
2238 if let Err(restore_error) =
2239 crate::database::save_materialized_session(&previous_projection)
2240 {
2241 tracing::error!(
2242 session_id,
2243 error = format!("{restore_error:#}"),
2244 "could not restore the durable projection after a failed resume"
2245 );
2246 }
2247 }
2248 Err(restore_error) => tracing::error!(
2249 session_id,
2250 error = format!("{restore_error:#}"),
2251 "could not rebuild the durable projection after a failed resume"
2252 ),
2253 }
2254 } else if discard_queued_prompts
2255 && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
2256 session_id,
2257 &mj_transcript::projection::materialized_queued_prompts_from_canonical(
2258 &canonical_session.queued_prompts,
2259 ),
2260 )
2261 {
2262 tracing::error!(
2263 session_id,
2264 error = format!("{restore_error:#}"),
2265 "could not restore queued prompts after a failed resume"
2266 );
2267 }
2268}
2269
2270fn projection_rebuild_required(
2278 stored: Option<(u64, &str)>,
2279 archive_frontier: u64,
2280 archive_frontier_digest: &str,
2281) -> bool {
2282 stored != Some((archive_frontier, archive_frontier_digest))
2283}
2284
2285fn restore_archive_path(
2286 backend: &targets::TargetLocator,
2287 verified_archive: &Path,
2288 remote_archive: &Path,
2289) -> PathBuf {
2290 if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2291 verified_archive.to_path_buf()
2292 } else {
2293 remote_archive.to_path_buf()
2294 }
2295}
2296
2297fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2298 !matches!(backend, targets::TargetLocator::LocalBare { .. })
2299}
2300
2301struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2306 inner: &'a E,
2307 cancellation: CancellationToken,
2308}
2309
2310impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2311 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2312 if self.cancellation.is_cancelled() {
2313 bail!("operation cancelled while provisioning destination");
2314 }
2315 self.inner.execute(command)
2316 }
2317
2318 fn cancellation_requested(&self) -> bool {
2319 self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2320 }
2321
2322 fn stage_started(&self, stage: ProvisionStage) {
2323 self.inner.stage_started(stage);
2324 }
2325
2326 fn stage_finished(&self, stage: ProvisionStage) {
2327 self.inner.stage_finished(stage);
2328 }
2329
2330 fn notify_notice(&self, notice: &str) {
2331 self.inner.notify_notice(notice);
2332 }
2333
2334 fn execute_with_stdin(
2335 &self,
2336 command: &CommandSpec,
2337 input: &mut (dyn std::io::Read + Send),
2338 ) -> Result<CommandOutput> {
2339 if self.cancellation.is_cancelled() {
2340 bail!("operation cancelled while provisioning destination");
2341 }
2342 self.inner.execute_with_stdin(command, input)
2343 }
2344}
2345
2346fn provision_with_cross_harness_handoff(
2347 controller: &mut Controller,
2348 session_id: &str,
2349 executor: &(impl CommandExecutor + Sync),
2350 github_token: Option<&str>,
2351 config: &Config,
2352 snapshot: &CanonicalSessionSnapshot,
2353 context_bytes: usize,
2354) -> Result<String> {
2355 let (_provision, handoff) = execute_joined_cross_harness_work(
2356 "cross-harness provisioning",
2357 move |cancellation| {
2358 let provision_executor = CrossHarnessProvisionExecutor {
2359 inner: executor,
2360 cancellation,
2361 };
2362 futures::executor::block_on(controller.provision_session_with_failure_disposition(
2363 session_id,
2364 &provision_executor,
2365 github_token,
2366 ProvisioningFailureDisposition::Preserve,
2367 ))
2368 },
2369 "cross-harness handoff",
2370 move |cancellation| {
2371 let runtime = tokio::runtime::Builder::new_current_thread()
2372 .enable_all()
2373 .build()
2374 .context("create cross-harness handoff runtime")?;
2375 runtime.block_on(utility_handoff_while_cancellable(
2376 session_id,
2377 config,
2378 snapshot,
2379 context_bytes,
2380 executor,
2381 cancellation,
2382 ))
2383 },
2384 )?;
2385 ensure!(
2386 !executor.cancellation_requested(),
2387 "operation cancelled while provisioning destination"
2388 );
2389 Ok(handoff)
2390}
2391
2392fn execute_joined_cross_harness_work<A: Send, B: Send>(
2396 first_name: &'static str,
2397 first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2398 second_name: &'static str,
2399 second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2400) -> Result<(A, B)> {
2401 let cancellation = CancellationToken::new();
2402 std::thread::scope(|scope| {
2403 let first_cancel = cancellation.clone();
2404 let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2405 let second_cancel = cancellation.clone();
2406 let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2407 let mut first_result = None;
2408 let mut second_result = None;
2409
2410 while first_result.is_none() || second_result.is_none() {
2411 if first_result.is_none()
2412 && first_handle
2413 .as_ref()
2414 .is_some_and(|handle| handle.is_finished())
2415 {
2416 let handle = first_handle.take().expect("first lane handle present");
2417 first_result = Some(match handle.join() {
2418 Ok(result) => result,
2419 Err(panic) => {
2420 cancellation.cancel();
2421 Err(anyhow::anyhow!(
2422 "{first_name} thread panicked: {}",
2423 targets::command_thread_panic_message(panic.as_ref())
2424 ))
2425 }
2426 });
2427 if first_result.as_ref().is_some_and(Result::is_err) {
2428 cancellation.cancel();
2429 }
2430 }
2431 if second_result.is_none()
2432 && second_handle
2433 .as_ref()
2434 .is_some_and(|handle| handle.is_finished())
2435 {
2436 let handle = second_handle.take().expect("second lane handle present");
2437 second_result = Some(match handle.join() {
2438 Ok(result) => result,
2439 Err(panic) => {
2440 cancellation.cancel();
2441 Err(anyhow::anyhow!(
2442 "{second_name} thread panicked: {}",
2443 targets::command_thread_panic_message(panic.as_ref())
2444 ))
2445 }
2446 });
2447 if second_result.as_ref().is_some_and(Result::is_err) {
2448 cancellation.cancel();
2449 }
2450 }
2451 if first_result.is_none() || second_result.is_none() {
2452 std::thread::sleep(Duration::from_millis(10));
2453 }
2454 }
2455
2456 match (
2457 first_result.expect("first lane result received after joined handle"),
2458 second_result.expect("second lane result received after joined handle"),
2459 ) {
2460 (Err(first), Err(second)) => {
2461 Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2462 }
2463 (Err(error), Ok(_)) => Err(error),
2464 (Ok(_), Err(error)) => Err(error),
2465 (Ok(first), Ok(second)) => Ok((first, second)),
2466 }
2467 })
2468}
2469
2470fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2474 profile_kind == archived_kind
2475}
2476
2477async fn utility_handoff_while_cancellable(
2481 session_id: &str,
2482 config: &Config,
2483 snapshot: &CanonicalSessionSnapshot,
2484 context_bytes: usize,
2485 executor: &impl CommandExecutor,
2486 cancellation: CancellationToken,
2487) -> Result<String> {
2488 let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2489 if executor.cancellation_requested() {
2490 bail!("operation cancelled while compacting the cross-harness handoff");
2491 }
2492 let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2493 let cancel = cancellation.child_token();
2494 let operation =
2495 crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2496 tokio::pin!(operation);
2497 loop {
2498 tokio::select! {
2499 context = &mut operation => return context,
2500 _ = cancellation.cancelled() => {
2501 cancel.cancel();
2502 bail!("operation cancelled while compacting the cross-harness handoff");
2503 }
2504 _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2505 if executor.cancellation_requested() {
2506 cancel.cancel();
2507 bail!("operation cancelled while compacting the cross-harness handoff");
2508 }
2509 }
2510 }
2511 }
2512}
2513
2514mod in_place;
2515pub(crate) use in_place::InPlaceRestartError;
2516
2517#[cfg(test)]
2518mod tests;