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 include_in_place_subagents: bool,
547 pub install_attached_resources: bool,
548 pub worker_root_reset: WorkerRootReset,
549 pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
551}
552
553impl Controller {
554 pub(super) async fn restore_into_target(
557 &mut self,
558 session_id: &str,
559 restore: RestoreIntoTarget<'_>,
560 executor: &(impl CommandExecutor + Sync),
561 ) -> Result<MaterializedSession> {
562 crate::worker_lifecycle::run(session_id, "restore into target", executor, async {
563 let RestoreIntoTarget {
564 profile,
565 archive,
566 restored_archive,
567 resumed_project_directory,
568 resumed_container_workspace,
569 restore_repositories,
570 native_continuity,
571 discard_queued_prompts,
572 replay_queue,
573 utility_handoff,
574 projection_build,
575 mut resume_notices,
576 include_in_place_subagents,
577 install_attached_resources: should_install_attached_resources,
578 worker_root_reset,
579 retire_after_ready,
580 } = restore;
581 let archive_manifest = &archive.manifest;
582 let canonical_session = &archive.canonical_session;
583 let stopped_subagents = crate::database::load_stopped_subagents(session_id)
588 .unwrap_or_else(|error| {
589 tracing::warn!(
590 session_id,
591 error = format!("{error:#}"),
592 "could not read the sub-agents this session's suspend stopped"
593 );
594 Vec::new()
595 });
596 let stopped_subagents_context =
597 mj_core::subagent::stopped_subagents_prompt_context(&stopped_subagents);
598 let reopen_move_subagents =
599 matches!(&worker_root_reset, WorkerRootReset::InPlace { .. })
600 && self.state.sessions[session_id]
601 .subagents
602 .clone()
603 .unwrap_or_default()
604 .for_launch(profile.kind, false)
605 .parent_role()
606 .is_some();
607 resume_notices.extend(mj_core::subagent::stopped_subagents_notice(
608 &stopped_subagents,
609 ));
610 let mut prompt_contexts = Vec::new();
611 prompt_contexts.extend(stopped_subagents_context);
612 let (backend, worker_root) = self.worker_placement(session_id)?;
613 let harness_home = target_profile_home(&backend, session_id, profile);
614 let workspace_root = if let Some(project_directory) = &resumed_project_directory {
615 project_directory
616 .parent()
617 .context("bare project directory has no parent")?
618 .to_string_lossy()
619 .into_owned()
620 } else {
621 super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
622 };
623 let target_path = |path: &str| match &backend {
624 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
625 if !path.starts_with('/') =>
626 {
627 PathBuf::from(format!("~/{path}"))
628 }
629 _ => PathBuf::from(path),
630 };
631 let remote_archive = format!("{worker_root}/restore.hel.zip");
632 let remote_spec = format!("{worker_root}/restore-spec.json");
633 use mj_checkpoint::checkpoint::QueueRestorePolicy;
634 let move_admission =
635 crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
636 operation.phase == mj_core::state::MovePhase::ResumingDestination
637 && !operation.queue_admission_started
638 && operation.queue == mj_core::state::ResumeQueueDisposition::Start
639 });
640 let queue_policy = if move_admission || (!native_continuity && replay_queue) {
641 QueueRestorePolicy::Defer
642 } else if discard_queued_prompts {
643 QueueRestorePolicy::Discard
644 } else {
645 QueueRestorePolicy::Restore
646 };
647 let restore = CheckpointRestoreSpec {
648 archive_path: restore_archive_path(
649 &backend,
650 restored_archive,
651 &target_path(&remote_archive),
652 ),
653 workspace_root: target_path(&workspace_root),
654 relay_root: target_path(&worker_root),
655 harness_home: target_path(&harness_home),
656 restore_repositories,
662 restore_native: native_continuity,
663 primary_repository_root: resumed_project_directory
672 .as_ref()
673 .map(|directory| target_path(&directory.to_string_lossy())),
674 queue_policy,
675 };
676 {
680 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
681 match &worker_root_reset {
682 WorkerRootReset::FreshTarget => {
687 if let Some(command) =
688 targets::clear_relay_state_plan(&backend, session_id)?
689 {
690 execute_checked(syncing, command)?;
691 }
692 execute_checked(
695 syncing,
696 targets::command_on_locator(
697 &backend,
698 session_id,
699 vec!["mkdir".into(), "-p".into(), worker_root.clone()],
700 "create the session worker root",
701 )?,
702 )?;
703 }
704 WorkerRootReset::InPlace {
709 previous_profile_root,
710 } => {
711 execute_checked(
712 syncing,
713 targets::in_place_worker_reset_plan(
714 &backend,
715 session_id,
716 previous_profile_root,
717 )?,
718 )?;
719 }
720 }
721 }
722 let staging = tempfile::tempdir().context("create restore staging")?;
723 let local_spec = staging.path().join("restore-spec.json");
724 std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
725 let controller = &*self;
729 let backend_ref = &backend;
730 let worker_root_ref = worker_root.as_str();
731 let local_spec_ref = local_spec.as_path();
732 crate::target_storage::ensure_room_for(
735 &backend,
736 [
737 worker_root.as_str(),
738 workspace_root.as_str(),
739 harness_home.as_str(),
740 ],
741 || {
742 Ok(std::fs::metadata(restored_archive)
743 .with_context(|| format!("measure {}", restored_archive.display()))?
744 .len())
745 },
746 "restore the checkpoint",
747 )?;
748 execute_concurrent_lanes(
749 || {
750 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
751 controller.prepare_worker_files(
752 session_id,
753 backend_ref,
754 worker_root_ref,
755 syncing,
756 )?;
757 super::provisioning::install_inherited_git_settings(
758 syncing,
759 backend_ref,
760 session_id,
761 )?;
762 Ok(())
763 },
764 || {
765 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
766 if should_upload_restore_archive(&backend) {
767 upload_checkpoint_spec(
768 restoring,
769 backend_ref,
770 session_id,
771 restored_archive,
772 &remote_archive,
773 )?;
774 }
775 upload_checkpoint_spec(
776 restoring,
777 backend_ref,
778 session_id,
779 local_spec_ref,
780 &remote_spec,
781 )
782 },
783 )?;
784 {
785 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
786 execute_checked(
787 restoring,
788 restore_command(&backend, session_id, &remote_spec)?,
789 )?;
790 }
791 if should_install_attached_resources {
792 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
793 install_attached_resources(
794 &self.state,
795 session_id,
796 &backend,
797 &worker_root,
798 syncing,
799 )?;
800 }
801 match projection_build {
802 Some(build) => {
803 let mut restored_projection = build
804 .await
805 .context("rebuild the restored projection")?
806 .context("rebuild the restored projection")?;
807 if discard_queued_prompts {
808 restored_projection.queued_prompts.clear();
809 }
810 crate::database::save_materialized_session(&restored_projection)?;
811 }
812 None if discard_queued_prompts => {
815 crate::database::replace_materialized_queued_prompts(session_id, &[])?;
816 }
817 None => {}
818 }
819 let readiness_stage = bridge_readiness_stage(profile);
820 let spec = self.reconnect_command(session_id)?;
821 let readiness = async {
822 let mut relay = {
823 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
824 start_worker_durably(
825 &crate::worker_lifecycle::require(session_id)?,
826 self.state.sessions[session_id]
827 .target
828 .as_ref()
829 .context("worker start has no durable target")?,
830 executor,
831 &backend,
832 &worker_root,
833 )?;
834 connect_started_worker(&spec, session_id, executor, &backend, &worker_root)
835 .await?
836 };
837 if reopen_move_subagents {
838 relay
839 .set_subagent_admission(true)
840 .await
841 .context("reopen sub-agent requests on the replacement worker")?;
842 }
843 if include_in_place_subagents {
844 let session_id = session_id.to_owned();
845 let (context, _) = tokio::task::spawn_blocking(move || {
846 load_in_place_subagent_prompt_data(&session_id)
847 })
848 .await
849 .context("join the retained sub-agent roster read")??;
850 prompt_contexts.extend(context);
851 }
852 if !prompt_contexts.is_empty() {
856 relay
857 .install_prompt_context(prompt_contexts.join("\n\n"))
858 .await
859 .context("tell the resumed session about its sub-agents")?;
860 }
861 let native_session_id = wait_for_native_session_in_stage(
862 &mut relay,
863 executor,
864 readiness_stage,
865 profile.kind,
866 )
867 .await?;
868 let owner = crate::worker_lifecycle::require(session_id)?;
869 crate::database::finish_worker_restart(session_id, owner.operation_id())?;
870 Ok::<_, anyhow::Error>((relay, native_session_id))
871 }
872 .await;
873 let (mut relay, native_session_id) = readiness
874 .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
875 if native_continuity {
876 if !restored_native_session_accepted(
877 &archive_manifest.session.native_session_id,
878 &native_session_id,
879 relay
880 .operational()
881 .replaced_unused_native_session_id
882 .as_deref(),
883 ) {
884 bail!(
885 "ACP loaded native session {native_session_id}, expected {}",
886 archive_manifest.session.native_session_id
887 );
888 }
889 } else {
890 relay
891 .install_prompt_context(
892 utility_handoff
893 .clone()
894 .context("a resume into a fresh native session has no handoff")?,
895 )
896 .await?;
897 if replay_queue {
898 for prompt in &canonical_session.queued_prompts {
899 let command = match &prompt.kind {
903 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
904 prompt: prompt
905 .content
906 .iter()
907 .cloned()
908 .map(serde_json::from_value)
909 .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
910 },
911 CanonicalQueuedCommandKind::SetConfig { key, value } => {
912 RelayCommand::SetConfig {
913 key: key.clone(),
914 value: value.clone(),
915 }
916 }
917 };
918 relay.submit(prompt.command_id.clone(), command).await?;
919 }
920 }
921 }
922 if let Some(worktree) = retire_after_ready
926 && let Err(error) = retire_managed_worktree(executor, worktree)
927 {
928 tracing::warn!(
929 session_id,
930 worktree = %worktree.worktree_root.display(),
931 error = format!("{error:#}"),
932 "could not retire the old managed worktree after resume"
933 );
934 resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
935 }
936 if include_in_place_subagents {
937 let session_id = session_id.to_owned();
938 let (_, notice) = tokio::task::spawn_blocking(move || {
939 load_in_place_subagent_prompt_data(&session_id)
940 })
941 .await
942 .context("join the retained sub-agent notice roster read")??;
943 resume_notices.extend(notice);
944 }
945 for notice in &resume_notices {
946 let submitted = async {
947 let command_id = new_command_id("resume-notice")?;
948 relay
949 .submit(
950 command_id,
951 RelayCommand::RecordNotice {
952 text: notice.clone(),
953 },
954 )
955 .await
956 }
957 .await;
958 if let Err(error) = submitted {
961 tracing::warn!(
962 session_id,
963 error = format!("{error:#}"),
964 "could not record a resume notice in the conversation"
965 );
966 }
967 }
968 self.mark_worker_connected(session_id, Some(native_session_id))?;
969 let materialized = relay.sync().await?.materialized;
970 if !stopped_subagents.is_empty() {
971 let delivered = stopped_subagents
972 .iter()
973 .map(|child| child.child_session_id.clone())
974 .collect::<Vec<_>>();
975 if let Err(error) = crate::database::clear_stopped_subagents(session_id, &delivered)
978 {
979 tracing::warn!(
980 session_id,
981 error = format!("{error:#}"),
982 "could not clear the stopped sub-agents after telling the resumed session"
983 );
984 }
985 }
986 Ok(materialized)
987 })
988 .await
989 }
990}
991
992pub(super) struct VerifiedResumeArchive {
999 pub archive_path: PathBuf,
1003 pub manifest: mj_checkpoint::archive::ArchiveManifest,
1004 pub canonical_session: Arc<CanonicalSessionSnapshot>,
1005}
1006
1007fn load_in_place_subagent_prompt_data(
1008 session_id: &str,
1009) -> Result<(Option<String>, Option<String>)> {
1010 let state = crate::database::load_state()
1011 .context("load current retained sub-agents for the in-place resume")?;
1012 let children = super::move_session::move_children(&state, session_id);
1013 Ok((
1014 mj_core::subagent::in_place_subagents_prompt_context(&children),
1015 mj_core::subagent::in_place_subagents_notice(&children),
1016 ))
1017}
1018
1019pub(super) fn verify_resume_checkpoint(
1025 session_id: &str,
1026 checkpoint: &mj_core::state::CheckpointMetadata,
1027) -> Result<VerifiedResumeArchive> {
1028 let archive_path = {
1029 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
1030 checkpoint.archive_path.canonicalize().with_context(|| {
1031 format!(
1032 "resolve checkpoint archive {}",
1033 checkpoint.archive_path.display()
1034 )
1035 })?
1036 };
1037 ensure!(
1038 archive_path.is_absolute() && archive_path.is_file(),
1039 "checkpoint archive path is not an absolute regular file: {}",
1040 archive_path.display()
1041 );
1042 let mj_checkpoint::archive::VerifiedArchiveMetadata {
1043 manifest,
1044 canonical_session,
1045 archive_sha256,
1046 } = {
1047 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
1048 verify_archive_streaming(&archive_path)?
1049 };
1050 if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
1051 bail!("persisted checkpoint verification failed");
1052 }
1053 ensure!(
1054 manifest.repositories.iter().all(|repository| {
1055 !repository.metadata.origin.starts_with("mj-local:")
1056 && !repository.metadata.origin.starts_with("ext::")
1057 }),
1058 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
1059 );
1060 Ok(VerifiedResumeArchive {
1061 archive_path,
1062 manifest,
1063 canonical_session: Arc::new(canonical_session),
1064 })
1065}
1066
1067pub fn raw_conversion_preview_for(
1068 session: &SessionRecord,
1069 config: &Config,
1070 executor: &(impl CommandExecutor + Sync),
1071) -> Result<mj_core::state::RawConversionPreview> {
1072 let conversion = plan_raw_to_workspace(session, config, executor)?;
1073 raw_conversion_preview(session, &conversion, executor)
1074}
1075
1076fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
1077 let replacement = replacement.trim();
1078 ensure!(!replacement.is_empty(), "enter the repository's new origin");
1079 let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
1080 let path = expanded.as_path();
1081 let (github, local) = if path.is_absolute() {
1082 ensure!(
1083 path.is_dir(),
1084 "local repository {replacement:?} is not a directory"
1085 );
1086 (None, Some(mj_core::local_git::canonical_repository(path)?))
1087 } else {
1088 let github = crate::setup::github_repository_from_origin(replacement)
1089 .context("origin must be a GitHub repository or an absolute local repository path")?;
1090 (
1091 Some(format!("{}/{}", github.owner, github.repository)),
1092 None,
1093 )
1094 };
1095 Ok(ProjectRepository {
1096 id: id.to_owned(),
1097 github,
1098 local,
1099 destination: PathBuf::from(id),
1100 git_ref: None,
1101 })
1102}
1103
1104fn checkpoint_source_missing_commit(
1105 configured: &ProjectRepository,
1106 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1107 archived: &CheckpointRepositoryBundle,
1108 executor: &impl CommandExecutor,
1109 github_token: Option<&str>,
1110) -> Result<Option<String>> {
1111 let staging = tempfile::tempdir().context("create repository source preflight")?;
1112 let repository = staging.path().join("repository.git");
1113 checked_preflight_git(
1114 executor,
1115 CommandSpec::new(
1116 "git",
1117 [
1118 "init".to_owned(),
1119 "--bare".to_owned(),
1120 "--quiet".to_owned(),
1121 repository.to_string_lossy().into_owned(),
1122 ],
1123 )
1124 .purpose("initialize repository source preflight"),
1125 )?;
1126 let missing = checkpoint_bundle_prerequisites(archived)?;
1127 if missing.is_empty() {
1128 let bundle = staging.path().join("checkpoint.bundle");
1129 std::fs::write(&bundle, &archived.committed_bundle)
1130 .context("write self-contained checkpoint bundle for source preflight")?;
1131 checked_preflight_git(
1132 executor,
1133 checkpoint_bundle_import_command(&repository, &bundle),
1134 )?;
1135 return Ok(None);
1136 }
1137 for commit in missing {
1141 let output = fetch_source_commit(
1142 executor,
1143 &repository,
1144 configured,
1145 accepted_source,
1146 &commit,
1147 github_token,
1148 )?;
1149 if output.status != 0 {
1150 let stderr = String::from_utf8_lossy(&output.stderr);
1151 if source_does_not_have_commit(&stderr) {
1152 return Ok(Some(commit));
1153 }
1154 bail!(
1155 "could not check configured source {:?}: {}",
1156 configured.source_label(),
1157 stderr.trim()
1158 );
1159 }
1160 }
1161 Ok(None)
1165}
1166
1167fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
1168 let mut command = CommandSpec::new(
1169 "git",
1170 [
1171 "-C".to_owned(),
1172 repository.to_string_lossy().into_owned(),
1173 "fetch".to_owned(),
1174 "--no-tags".to_owned(),
1175 bundle.to_string_lossy().into_owned(),
1176 "HEAD".to_owned(),
1177 ],
1178 )
1179 .purpose("validate self-contained checkpoint bundle");
1180 command
1181 .env
1182 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1183 command
1184 .env
1185 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1186 command
1187}
1188
1189fn fetch_source_commit(
1190 executor: &impl CommandExecutor,
1191 repository: &Path,
1192 configured: &ProjectRepository,
1193 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1194 commit: &str,
1195 github_token: Option<&str>,
1196) -> Result<CommandOutput> {
1197 let mut arguments = Vec::new();
1198 let mut token_auth = false;
1199 let mut ssh_transport = false;
1200 let source = if let Some(source) = accepted_source {
1201 source.fetch_url.clone()
1202 } else if let Some(local) = &configured.local {
1203 local.to_string_lossy().into_owned()
1204 } else {
1205 let source = configured
1206 .github
1207 .as_deref()
1208 .context("repository source is missing")?;
1209 let github = crate::setup::github_repository_from_origin(source)
1210 .context("configured repository is not a GitHub source")?;
1211 if github_token.is_some() {
1212 token_auth = true;
1213 arguments.extend([
1214 "-c".to_owned(),
1215 "credential.helper=".to_owned(),
1216 "-c".to_owned(),
1217 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1218 ]);
1219 format!(
1220 "https://github.com/{}/{}.git",
1221 github.owner, github.repository
1222 )
1223 } else {
1224 ssh_transport = true;
1225 format!("git@github.com:{}/{}.git", github.owner, github.repository)
1226 }
1227 };
1228 if accepted_source.is_some() {
1229 if source.starts_with("https://github.com/") && github_token.is_some() {
1230 token_auth = true;
1231 arguments.extend([
1232 "-c".to_owned(), "credential.helper=".to_owned(), "-c".to_owned(),
1233 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1234 ]);
1235 }
1236 ssh_transport = source.starts_with("git@") || source.starts_with("ssh://");
1237 }
1238 arguments.extend([
1239 "-C".to_owned(),
1240 repository.to_string_lossy().into_owned(),
1241 "fetch".to_owned(),
1242 "--no-tags".to_owned(),
1243 "--depth=1".to_owned(),
1244 "--filter=blob:none".to_owned(),
1245 source,
1246 commit.to_owned(),
1247 ]);
1248 let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1249 command
1250 .env
1251 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1252 command
1253 .env
1254 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1255 if token_auth {
1256 let token = github_token.expect("token authentication requires a GitHub token");
1257 command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1258 }
1259 if ssh_transport {
1260 command.env.insert(
1261 "GIT_SSH_COMMAND".to_owned(),
1262 "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1263 .to_owned(),
1264 );
1265 }
1266 executor.execute(&command)
1267}
1268
1269fn source_does_not_have_commit(stderr: &str) -> bool {
1270 let stderr = stderr.to_ascii_lowercase();
1271 [
1272 "not our ref",
1273 "couldn't find remote ref",
1274 "not a valid object name",
1275 "no such ref was fetched",
1276 ]
1277 .iter()
1278 .any(|needle| stderr.contains(needle))
1279}
1280
1281fn checked_preflight_git(
1282 executor: &impl CommandExecutor,
1283 command: CommandSpec,
1284) -> Result<CommandOutput> {
1285 let output = executor.execute(&command)?;
1286 ensure!(
1287 output.status == 0,
1288 "{}: {}",
1289 command.purpose,
1290 String::from_utf8_lossy(&output.stderr).trim()
1291 );
1292 Ok(output)
1293}
1294
1295impl Controller {
1296 pub async fn resume_session_with_options(
1301 &mut self,
1302 session_id: &str,
1303 profile_id: &str,
1304 target_id: &str,
1305 additional_mounts: Option<Vec<AdditionalMount>>,
1306 resource_allocation: Option<SessionResourceAllocation>,
1307 ) -> Result<MaterializedSession> {
1308 self.resume_session_with_options_and_queue_disposition(
1309 session_id,
1310 profile_id,
1311 target_id,
1312 additional_mounts,
1313 resource_allocation,
1314 false,
1315 )
1316 .await
1317 }
1318
1319 pub async fn resume_session_with_options_and_queue_disposition(
1320 &mut self,
1321 session_id: &str,
1322 profile_id: &str,
1323 target_id: &str,
1324 additional_mounts: Option<Vec<AdditionalMount>>,
1325 resource_allocation: Option<SessionResourceAllocation>,
1326 discard_queue: bool,
1327 ) -> Result<MaterializedSession> {
1328 self.resume_session_controlled(
1329 session_id,
1330 profile_id,
1331 target_id,
1332 SessionResumeOptions {
1333 additional_mounts,
1334 resource_allocation,
1335 discard_queue,
1336 },
1337 &ProcessExecutor,
1338 )
1339 .await
1340 }
1341
1342 pub async fn resume_session_controlled(
1343 &mut self,
1344 session_id: &str,
1345 profile_id: &str,
1346 target_id: &str,
1347 options: SessionResumeOptions,
1348 executor: &(impl CommandExecutor + Sync),
1349 ) -> Result<MaterializedSession> {
1350 self.resume_session_controlled_with_repository_preflight(
1351 session_id, profile_id, target_id, options, None, executor,
1352 )
1353 .await
1354 }
1355
1356 pub async fn resume_session_controlled_with_repository_preflight(
1357 &mut self,
1358 session_id: &str,
1359 profile_id: &str,
1360 target_id: &str,
1361 options: SessionResumeOptions,
1362 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1363 executor: &(impl CommandExecutor + Sync),
1364 ) -> Result<MaterializedSession> {
1365 self.resume_session_with_origin(
1366 session_id,
1367 profile_id,
1368 target_id,
1369 options,
1370 repository_preflight,
1371 None,
1372 executor,
1373 )
1374 .await
1375 }
1376
1377 pub(in crate::controller) async fn resume_session_for_move(
1378 &mut self,
1379 operation: &mj_core::state::MoveOperation,
1380 executor: &(impl CommandExecutor + Sync),
1381 ) -> Result<MaterializedSession> {
1382 self.resume_session_with_origin(
1383 &operation.selection.session_id,
1384 operation
1385 .selection
1386 .profile_id
1387 .as_deref()
1388 .context("Move profile missing")?,
1389 operation
1390 .selection
1391 .target_template_id
1392 .as_deref()
1393 .context("Move target missing")?,
1394 SessionResumeOptions {
1395 additional_mounts: operation.selection.additional_mounts.clone(),
1396 resource_allocation: operation.selection.resource_allocation.clone(),
1397 discard_queue: true,
1398 },
1399 None,
1400 Some(operation),
1401 executor,
1402 )
1403 .await
1404 }
1405
1406 #[allow(clippy::too_many_arguments)]
1407 async fn resume_session_with_origin(
1408 &mut self,
1409 session_id: &str,
1410 profile_id: &str,
1411 target_id: &str,
1412 options: SessionResumeOptions,
1413 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1414 move_operation: Option<&mj_core::state::MoveOperation>,
1415 executor: &(impl CommandExecutor + Sync),
1416 ) -> Result<MaterializedSession> {
1417 crate::worker_lifecycle::run(session_id, "resume session with origin", executor, async {
1418 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1419 let transferring_workspace = move_operation.is_some();
1420 if let Some(operation) = crate::database::load_move_operation(session_id)? {
1421 ensure!(
1422 transferring_workspace || !operation.holds_source_environment(),
1423 "a profile switch retains this environment; retry Move instead of recreating it with Resume"
1424 );
1425 }
1426 let SessionResumeOptions {
1427 additional_mounts,
1428 resource_allocation,
1429 discard_queue,
1430 } = options;
1431 let mut previous = self
1432 .state
1433 .sessions
1434 .get(session_id)
1435 .with_context(|| format!("unknown session {session_id}"))?
1436 .clone();
1437 if !(matches!(
1438 previous.state,
1439 SessionState::Stopped | SessionState::Lost | SessionState::Error
1440 ) || (transferring_workspace && previous.state == SessionState::Closing))
1441 {
1442 bail!("session {session_id} is not stopped, lost, or retryable");
1443 }
1444 let checkpoint = move_operation
1445 .and_then(|op| op.handoff.as_ref())
1446 .or(previous.checkpoint.as_ref())
1447 .context("session has no checkpoint")?;
1448 let moving_to_raw = previous.project_directory.is_none()
1451 && self
1452 .config
1453 .targets
1454 .get(target_id)
1455 .is_some_and(mj_core::config::is_bare_project_target);
1456 if !transferring_workspace
1457 && (moving_to_raw
1458 || !repository_preflight.as_ref().is_some_and(|receipt| {
1459 self.repository_source_receipt_is_current(session_id, receipt)
1460 }))
1461 {
1462 let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1463 if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1464 self.preflight_repository_sources(session_id, target_id, false, executor)?
1465 {
1466 bail!(
1467 "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1468 mismatch.missing_commit,
1469 mismatch.configured_origin,
1470 mismatch.repository_id,
1471 mismatch.archived_origin,
1472 );
1473 }
1474 }
1475 let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1476 let archive_path = verified_archive.archive_path.clone();
1477 let archive_manifest = &verified_archive.manifest;
1478 let canonical_session = Arc::clone(&verified_archive.canonical_session);
1479 let profile = self
1480 .config
1481 .profiles
1482 .get(profile_id)
1483 .with_context(|| format!("unknown profile {profile_id:?}"))?
1484 .clone();
1485 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1486 let target_template = self
1487 .config
1488 .targets
1489 .get(target_id)
1490 .with_context(|| format!("unknown target template {target_id:?}"))?
1491 .clone();
1492 self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1495 ensure!(
1496 profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1497 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1498 );
1499 let plan = resume_compatibility(&previous, &self.config, target_id)
1500 .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1501 if plan == ResumePlan::InPlace
1504 && let Some(worktree) = previous.managed_worktree.as_mut()
1505 && let Ok(current) = managed_worktree_target(&target_template)
1506 {
1507 worktree.target = current;
1508 }
1509 if !mj_core::config::is_bare_project_target(&target_template)
1513 && plan != ResumePlan::RawToWorkspace
1514 && !transferring_workspace
1515 {
1516 super::network_git::bundle_from_manifest(archive_manifest)?;
1517 }
1518 if plan == ResumePlan::InPlace
1519 && previous.managed_worktree.is_none()
1520 && let Some(project_directory) = &previous.project_directory
1521 {
1522 self.validate_project_directory(target_id, project_directory, executor)
1523 .context("raw project is unavailable for resume")?;
1524 }
1525 let conversion = match plan {
1526 ResumePlan::InPlace => None,
1527 ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(
1528 plan_raw_to_workspace(&previous, &self.config, executor)
1529 .context("prepare the raw checkout for its new target")?,
1530 )),
1531 ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1532 self.plan_workspace_to_raw(&previous, target_id, executor)
1533 .context("prepare a checkout for this session")?,
1534 )),
1535 };
1536 let resource_allocation =
1537 resource_allocation.or_else(|| previous.resource_allocation.clone());
1538 let additional_mounts =
1539 additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1540 validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1541 let selected_container_size =
1542 selected_host_container_size(&target_template, resource_allocation.as_ref());
1543 if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1544 bail!("attached resources are unsupported for this target");
1545 }
1546 targets::validate_additional_mounts(&additional_mounts)?;
1547 let history_host = mount_history_host(&target_template);
1548 let history_mounts = additional_mounts.clone();
1549 if !transferring_workspace
1550 && previous.state == SessionState::Error
1551 && let Some(locator) = &previous.target
1552 {
1553 let backend = backend_locator(locator, &previous, &self.config)?;
1554 targets::close_plan(&backend, session_id)?
1555 .execute(executor)
1556 .context("clean up target from failed resume")?;
1557 }
1558 let mut resume_notices = Vec::new();
1559 if let Some(conversion) = conversion
1562 .as_ref()
1563 .and_then(ResumeConversion::workspace_to_raw)
1564 {
1565 resume_notices.push(format!(
1566 "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1567 previous.target_template_id,
1568 conversion.worktree.worktree_root.display(),
1569 archive_manifest
1570 .repositories
1571 .first()
1572 .and_then(|repository| repository.metadata.branch.as_deref())
1573 .unwrap_or("a detached head"),
1574 conversion.worktree.branch,
1575 ));
1576 }
1577 let managed_checkout_present = previous
1578 .managed_worktree
1579 .as_ref()
1580 .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1581 .transpose()?
1582 .unwrap_or(true);
1583 if managed_checkout_present && let Some(project_directory) = &previous.project_directory {
1586 match raw_checkout_position(&previous, &self.config, project_directory, executor) {
1587 Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1588 project_directory,
1589 archive_manifest
1590 .repositories
1591 .first()
1592 .map(|repository| &repository.metadata),
1593 &live,
1594 )),
1595 Err(error) => tracing::warn!(
1598 session_id,
1599 error = format!("{error:#}"),
1600 "could not read the raw checkout position for a resume notice"
1601 ),
1602 }
1603 }
1604 super::worker_binary::preflight_worker_binary(&target_template, executor)?;
1609 let same_harness = profile.kind == archive_manifest.session.harness_kind;
1610 let native_continuity =
1611 native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1612 let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1613 let utility_config = (!native_continuity).then(|| self.config.clone());
1617 let discard_queued_prompts = discard_queue || !same_harness;
1618 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1622 .unwrap_or_else(|error| {
1623 tracing::warn!(
1624 session_id,
1625 error = format!("{error:#}"),
1626 "could not read the stored projection frontier; rebuilding it from the archive"
1627 );
1628 None
1629 });
1630 let rebuild_projection = projection_rebuild_required(
1631 stored_frontier
1632 .as_ref()
1633 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1634 canonical_session.event_frontier,
1635 &canonical_session.event_frontier_digest,
1636 );
1637 let projection_build = rebuild_projection.then(|| {
1647 let canonical = Arc::clone(&canonical_session);
1648 let session_id = session_id.to_owned();
1649 tokio::task::spawn_blocking(move || {
1650 materialized_session_from_canonical(session_id, &canonical)
1651 })
1652 });
1653 let github_token = controller_github_token();
1654
1655 if let Some(conversion) = conversion
1658 .as_ref()
1659 .and_then(ResumeConversion::raw_to_workspace)
1660 && let Some(bundle) = &conversion.new_bundle
1661 {
1662 let (config, ()) = Config::update(|config| {
1663 if let Some(existing) = config.bundles.get(&conversion.bundle_id) {
1664 ensure!(
1665 existing == bundle,
1666 "bundle {:?} was configured concurrently with a different definition; retry the resume",
1667 conversion.bundle_id
1668 );
1669 } else {
1670 config
1671 .bundles
1672 .insert(conversion.bundle_id.clone(), bundle.clone());
1673 }
1674 Ok(())
1675 })
1676 .context("save the bundle for a converted raw session")?;
1677 self.config = config;
1678 }
1679
1680 let moving_into_first_container = self
1686 .config
1687 .targets
1688 .get(target_id)
1689 .is_some_and(mj_core::config::is_container_target)
1690 && !self
1691 .config
1692 .targets
1693 .get(&previous.target_template_id)
1694 .is_some_and(mj_core::config::is_container_target);
1695 let converted_project = conversion
1696 .as_ref()
1697 .and_then(ResumeConversion::raw_to_workspace)
1698 .map(|conversion| {
1699 crate::project_catalog::snapshot(
1700 self.config
1701 .bundles
1702 .get(&conversion.bundle_id)
1703 .context("converted project is missing")?,
1704 executor,
1705 true,
1706 )
1707 })
1708 .transpose()?;
1709 let record = self.state.sessions.get_mut(session_id).unwrap();
1710 if record.container_workspace.is_none() && moving_into_first_container {
1711 record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1712 }
1713 record.harness_kind = profile.kind;
1714 record.last_profile = profile_id.to_string();
1715 record.target_template_id = target_id.to_string();
1716 record.target_runtime = Some(
1717 self.config
1718 .targets
1719 .get(target_id)
1720 .context("resume target disappeared")?
1721 .into(),
1722 );
1723 record.resource_allocation = resource_allocation;
1724 record.additional_mounts = additional_mounts;
1725 record.target = None;
1726 record.native_session_id =
1727 native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1728 record.state = SessionState::Provisioning;
1729 record.updated_at = now();
1730 record.last_error = None;
1731 match &conversion {
1732 Some(ResumeConversion::RawToWorkspace(conversion)) => {
1733 apply_raw_to_workspace(record, conversion);
1734 record.project = converted_project;
1735 }
1736 Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1737 apply_workspace_to_raw(record, conversion);
1738 }
1739 None => {
1740 if let (Some(worktree), Some(refreshed)) = (
1741 record.managed_worktree.as_mut(),
1742 previous.managed_worktree.as_ref(),
1743 ) {
1744 worktree.target = refreshed.target.clone();
1745 }
1746 }
1747 }
1748 let resumed_project_directory = record.project_directory.clone();
1749 let resumed_container_workspace = record.container_workspace.clone();
1750 if let Some(host) = history_host {
1751 self.state.remember_mount_sources(host, &history_mounts);
1752 crate::database::remember_mount_sources(host, &history_mounts)?;
1753 }
1754 if let Some(conversion) = conversion
1757 .as_ref()
1758 .and_then(ResumeConversion::raw_to_workspace)
1759 {
1760 let bundle_id =
1765 crate::database::rebind_session_bundle(session_id, &conversion.bundle_id)?;
1766 self.state
1767 .sessions
1768 .get_mut(session_id)
1769 .expect("the resumed session is in state")
1770 .bundle_id = bundle_id;
1771 }
1772 crate::database::save_resumed_session(
1775 &self.state.sessions[session_id],
1776 selected_container_size
1777 .as_ref()
1778 .map(|(host, size)| (host.as_str(), *size)),
1779 )?;
1780 if let Some((host, size)) = selected_container_size.as_ref() {
1781 self.state.remember_container_size(host, *size);
1782 }
1783
1784 let mut recreated_managed_worktree = false;
1785 let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1788 let result = async {
1789 if let Some(worktree) = previous.managed_worktree.as_ref() {
1790 recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1791 if !transferring_workspace && recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1792 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1793 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1794 &archive_path, &worktree.worktree_root, &SystemGit,
1795 )?;
1796 } else {
1797 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1798 &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1799 )?;
1800 }
1801 }
1802 }
1803 if let Some(conversion) = conversion
1806 .as_ref()
1807 .and_then(ResumeConversion::workspace_to_raw)
1808 {
1809 if conversion.reuse_existing_branch {
1810 let recovery_ref =
1811 preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1812 restore_managed_worktree(executor, &conversion.worktree)?;
1813 resume_notices.push(format!(
1814 "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1815 ));
1816 } else {
1817 create_managed_worktree(
1818 executor,
1819 &conversion.worktree,
1820 None,
1821 PrimaryCheckoutRequirement::Any,
1822 )?;
1823 }
1824 if transferring_workspace {
1825 } else if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1827 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1828 &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1829 )?;
1830 } else {
1831 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1832 &archive_path,
1833 &conversion.worktree.worktree_root,
1834 &conversion.worktree.branch,
1835 &SystemGit,
1836 )?;
1837 }
1838 }
1839 if let Some(conversion) = conversion
1844 .as_ref()
1845 .and_then(ResumeConversion::raw_to_workspace)
1846 && !transferring_workspace
1847 {
1848 let destination = PathBuf::from(
1849 previous
1850 .project_directory
1851 .as_deref()
1852 .context("a raw session has no project directory")?
1853 .file_name()
1854 .context("a raw project directory cannot be the filesystem root")?,
1855 );
1856 let snapshot = raw_checkout_snapshot(
1857 &conversion.checkout,
1858 &conversion.source,
1859 &destination,
1860 &SystemGit,
1861 conversion.retire.as_ref().is_some_and(|checkout| {
1862 checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1863 }),
1864 )
1865 .context("snapshot the host checkout for its new target")?;
1866 resume_notices.push(conversion_notice(
1867 target_id,
1868 previous
1869 .project_directory
1870 .as_deref()
1871 .unwrap_or(&conversion.checkout),
1872 snapshot.metadata.branch.as_deref(),
1873 conversion.retire.as_ref(),
1874 ));
1875 let archives = mj_core::config::sessions_dir();
1876 std::fs::create_dir_all(&archives).with_context(|| {
1877 format!("create the checkpoint directory {}", archives.display())
1878 })?;
1879 let output = archives.join(format!(
1883 "{session_id}-converted-{}-{}.hel.zip",
1884 previous
1885 .checkpoint
1886 .as_ref()
1887 .map_or(0, |checkpoint| checkpoint.event_frontier),
1888 new_command_id("archive")?
1889 ));
1890 let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1891 conversion_checkpoint_written = Some(written.clone());
1892 let record = self.state.sessions.get_mut(session_id).unwrap();
1893 record.checkpoint = Some(written);
1894 record.updated_at = now();
1895 crate::database::save_resumed_session(
1898 &self.state.sessions[session_id],
1899 selected_container_size.as_ref().map(|(host, size)| (host.as_str(), *size)),
1900 )?;
1901 }
1902 let utility_handoff = {
1903 let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1904 if let Some(config) = utility_config.as_ref() {
1905 Some(
1906 provision_with_cross_harness_handoff(
1907 self,
1908 session_id,
1909 executor,
1910 github_token.as_deref(),
1911 config,
1912 &canonical_session,
1913 context_bytes,
1914 )
1915 .context("prepare the cross-harness destination")?,
1916 )
1917 } else {
1918 self.provision_session_with_failure_disposition(
1919 session_id,
1920 executor,
1921 github_token.as_deref(),
1922 ProvisioningFailureDisposition::Preserve,
1923 )
1924 .await?;
1925 None
1926 }
1927 };
1928 if transferring_workspace {
1929 self.restore_move_workspace(session_id, executor)?;
1930 }
1931 let restore_repositories = !transferring_workspace && ((resumed_project_directory.is_none()
1932 && conversion.is_none())
1933 || plan == ResumePlan::RawToWorkspace
1934 || (recreated_managed_worktree && plan == ResumePlan::InPlace));
1935 let restored_archive = conversion_checkpoint_written
1938 .as_ref()
1939 .map_or(archive_path.as_path(), |checkpoint| {
1940 checkpoint.archive_path.as_path()
1941 });
1942 self.restore_into_target(
1943 session_id,
1944 RestoreIntoTarget {
1945 profile: &profile,
1946 archive: &verified_archive,
1947 restored_archive,
1948 resumed_project_directory,
1949 resumed_container_workspace,
1950 restore_repositories,
1951 native_continuity,
1952 discard_queued_prompts,
1953 replay_queue: !discard_queue,
1954 utility_handoff,
1955 projection_build,
1956 resume_notices,
1957 include_in_place_subagents: false,
1958 install_attached_resources: true,
1959 worker_root_reset: WorkerRootReset::FreshTarget,
1960 retire_after_ready: conversion
1961 .as_ref()
1962 .filter(|_| !transferring_workspace)
1963 .and_then(ResumeConversion::raw_to_workspace)
1964 .and_then(|plan| plan.retire.as_ref()),
1965 },
1966 executor,
1967 )
1968 .await
1969 }
1970 .await;
1971 match result {
1972 Ok(materialized) => {
1973 if let Some(written) = &conversion_checkpoint_written {
1976 super::checkpoint::prune_replaced_checkpoint(
1977 previous.checkpoint.as_ref(),
1978 written,
1979 );
1980 }
1981 Ok(materialized)
1982 }
1983 Err(error) => {
1984 if let Some(written) = &conversion_checkpoint_written
1987 && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
1988 && remove_error.kind() != std::io::ErrorKind::NotFound
1989 {
1990 tracing::warn!(
1991 session_id,
1992 path = %written.archive_path.display(),
1993 "could not remove the conversion checkpoint after resume failed: {remove_error}"
1994 );
1995 }
1996 restore_projection_after_failed_resume(
1999 session_id,
2000 &canonical_session,
2001 discard_queued_prompts,
2002 );
2003 if let Some(operation) = move_operation {
2004 return Err(self.rollback_move_destination(operation, error, executor)?);
2005 }
2006 Err(self.rollback_failed_resume(
2007 session_id,
2008 &previous,
2009 recreated_managed_worktree,
2010 error,
2011 executor,
2012 )?)
2013 }
2014 }
2015
2016 }).await
2017 }
2018
2019 pub(super) fn rollback_failed_resume(
2020 &mut self,
2021 session_id: &str,
2022 previous: &SessionRecord,
2023 recreated_managed_worktree: bool,
2024 error: anyhow::Error,
2025 _executor: &impl CommandExecutor,
2026 ) -> Result<anyhow::Error> {
2027 let current = self
2028 .state
2029 .sessions
2030 .get(session_id)
2031 .with_context(|| format!("unknown session {session_id}"))?
2032 .clone();
2033 let cleanup = match current.target.as_ref() {
2034 Some(locator) => (|| -> Result<()> {
2035 let backend = backend_locator(locator, ¤t, &self.config)?;
2036 targets::close_plan(&backend, session_id)?
2037 .execute(&CancellableProcessExecutor::with_timeout(
2040 Duration::from_secs(15),
2041 ))
2042 .map(|_| ())
2043 })(),
2044 None => Ok(()),
2045 };
2046 let worktree_cleanup = if cleanup.is_err() {
2049 Ok(())
2050 } else {
2051 match (
2052 current.managed_worktree.as_ref(),
2053 previous.managed_worktree.as_ref(),
2054 ) {
2055 (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
2056 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
2057 previous,
2058 ),
2059 (Some(current), Some(previous)) if current == previous => Ok(()),
2060 (Some(worktree), _) => cleanup_managed_worktree(
2063 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
2064 worktree,
2065 crate::controller::BranchDisposition::Delete,
2066 ),
2067 (None, _) => Ok(()),
2068 }
2069 };
2070 let cleanup_error = [cleanup, worktree_cleanup]
2071 .into_iter()
2072 .filter_map(Result::err)
2073 .map(|cleanup_error| format!("{cleanup_error:#}"))
2074 .collect::<Vec<_>>()
2075 .join("; ");
2076 if !cleanup_error.is_empty() {
2077 tracing::warn!(
2078 session_id,
2079 error = %cleanup_error,
2080 "resume rollback cleanup reported failures"
2081 );
2082 }
2083 let original = super::worker_binary::failure_line(&error);
2086 let detail = format!("{error:#}");
2087 if detail != original {
2088 tracing::warn!(session_id, error = %detail, "resume failed");
2089 }
2090 let record = self.state.sessions.get_mut(session_id).unwrap();
2091 let failure = apply_failed_resume_rollback(
2092 record,
2093 previous,
2094 &original,
2095 (!cleanup_error.is_empty()).then_some(cleanup_error),
2096 );
2097 if record.bundle_id != current.bundle_id {
2100 let bundle_id = record.bundle_id.clone();
2101 record.bundle_id = crate::database::rebind_session_bundle(session_id, &bundle_id)?;
2102 }
2103 crate::database::save_resumed_session(&self.state.sessions[session_id], None)?;
2106 Ok(failure)
2107 }
2108}
2109
2110fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
2111 format!(
2112 "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
2113 worktree_root.display(),
2114 worktree_root.display()
2115 )
2116}
2117
2118pub(super) fn apply_failed_resume_rollback(
2119 current: &mut SessionRecord,
2120 previous: &SessionRecord,
2121 original_error: &str,
2122 cleanup_error: Option<String>,
2123) -> anyhow::Error {
2124 match cleanup_error {
2125 None => {
2126 *current = previous.clone();
2127 current.state = SessionState::Stopped;
2128 current.target = None;
2129 current.updated_at = now();
2130 current.last_error = Some(format!("resume failed: {original_error}"));
2131 anyhow::anyhow!(original_error.to_owned())
2132 }
2133 Some(cleanup_error) => {
2134 let failure = format!(
2135 "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
2136 );
2137 if current.managed_worktree.is_none() {
2142 current
2143 .project_directory
2144 .clone_from(&previous.project_directory);
2145 current
2146 .managed_worktree
2147 .clone_from(&previous.managed_worktree);
2148 current.bundle_id.clone_from(&previous.bundle_id);
2149 }
2150 current.state = SessionState::Error;
2151 current.updated_at = now();
2152 current.last_error = Some(format!("resume failed: {failure}"));
2153 anyhow::anyhow!(failure)
2154 }
2155 }
2156}
2157
2158fn conversion_notice(
2160 target_id: &str,
2161 checkout: &Path,
2162 branch: Option<&str>,
2163 retire: Option<&mj_core::state::ManagedWorktree>,
2164) -> String {
2165 let branch = branch.unwrap_or("a detached head");
2166 match retire {
2167 Some(worktree) => format!(
2168 "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
2169 checkout.display(),
2170 worktree.branch,
2171 worktree.source_repository.display()
2172 ),
2173 None => format!(
2174 "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.",
2175 checkout.display()
2176 ),
2177 }
2178}
2179
2180fn conversion_checkpoint(
2188 previous_archive: &Path,
2189 snapshot: mj_checkpoint::archive::RepositorySnapshot,
2190 output: &Path,
2191) -> Result<mj_core::state::CheckpointMetadata> {
2192 let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
2193 .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
2194 let native_artifacts = previous
2197 .manifest
2198 .payloads
2199 .iter()
2200 .filter_map(|descriptor| match &descriptor.role {
2201 mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
2202 Some((relative_path, descriptor))
2203 }
2204 _ => None,
2205 })
2206 .map(|(relative_path, descriptor)| {
2207 Ok(mj_checkpoint::archive::NativeArtifact {
2208 relative_path: relative_path.clone(),
2209 data: previous.payload(descriptor)?.to_vec(),
2210 mode: descriptor.mode,
2211 })
2212 })
2213 .collect::<Result<Vec<_>>>()?;
2214 let canonical_session = previous.canonical_session()?;
2215 let event_frontier = canonical_session.event_frontier;
2216 let written = mj_checkpoint::archive::write_archive_atomic(
2217 output,
2218 &mj_checkpoint::archive::ArchiveInput {
2219 session: previous.manifest.session.clone(),
2220 target: previous.manifest.target.clone(),
2223 bundle: mj_checkpoint::archive::BundleManifest {
2224 id: previous.manifest.bundle.id.clone(),
2225 primary_repository: snapshot.metadata.id.clone(),
2229 },
2230 canonical_session,
2231 native_artifacts,
2232 repositories: vec![snapshot],
2233 },
2234 )
2235 .with_context(|| format!("write the conversion archive {}", output.display()))?;
2236 Ok(mj_core::state::CheckpointMetadata {
2237 archive_path: output.to_path_buf(),
2238 sha256: written.archive_sha256,
2239 created_at: now(),
2240 event_frontier,
2241 })
2242}
2243
2244fn restored_native_session_accepted(
2255 archived: &str,
2256 opened: &str,
2257 replaced_unused: Option<&str>,
2258) -> bool {
2259 opened == archived || replaced_unused == Some(archived)
2260}
2261
2262fn restore_projection_after_failed_resume(
2273 session_id: &str,
2274 canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
2275 discard_queued_prompts: bool,
2276) {
2277 let stored_frontier = crate::database::materialized_event_frontier(session_id)
2278 .unwrap_or_else(|error| {
2279 tracing::warn!(
2280 session_id,
2281 error = format!("{error:#}"),
2282 "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
2283 );
2284 None
2285 });
2286 if projection_rebuild_required(
2287 stored_frontier
2288 .as_ref()
2289 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
2290 canonical_session.event_frontier,
2291 &canonical_session.event_frontier_digest,
2292 ) {
2293 match materialized_session_from_canonical(session_id, canonical_session) {
2294 Ok(previous_projection) => {
2295 if let Err(restore_error) =
2296 crate::database::save_materialized_session(&previous_projection)
2297 {
2298 tracing::error!(
2299 session_id,
2300 error = format!("{restore_error:#}"),
2301 "could not restore the durable projection after a failed resume"
2302 );
2303 }
2304 }
2305 Err(restore_error) => tracing::error!(
2306 session_id,
2307 error = format!("{restore_error:#}"),
2308 "could not rebuild the durable projection after a failed resume"
2309 ),
2310 }
2311 } else if discard_queued_prompts
2312 && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
2313 session_id,
2314 &mj_transcript::projection::materialized_queued_prompts_from_canonical(
2315 &canonical_session.queued_prompts,
2316 ),
2317 )
2318 {
2319 tracing::error!(
2320 session_id,
2321 error = format!("{restore_error:#}"),
2322 "could not restore queued prompts after a failed resume"
2323 );
2324 }
2325}
2326
2327fn projection_rebuild_required(
2335 stored: Option<(u64, &str)>,
2336 archive_frontier: u64,
2337 archive_frontier_digest: &str,
2338) -> bool {
2339 stored != Some((archive_frontier, archive_frontier_digest))
2340}
2341
2342fn restore_archive_path(
2343 backend: &targets::TargetLocator,
2344 verified_archive: &Path,
2345 remote_archive: &Path,
2346) -> PathBuf {
2347 if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2348 verified_archive.to_path_buf()
2349 } else {
2350 remote_archive.to_path_buf()
2351 }
2352}
2353
2354fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2355 !matches!(backend, targets::TargetLocator::LocalBare { .. })
2356}
2357
2358struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2363 inner: &'a E,
2364 cancellation: CancellationToken,
2365}
2366
2367impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2368 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2369 if self.cancellation.is_cancelled() {
2370 bail!("operation cancelled while provisioning destination");
2371 }
2372 self.inner.execute(command)
2373 }
2374
2375 fn cancellation_requested(&self) -> bool {
2376 self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2377 }
2378
2379 fn stage_started(&self, stage: ProvisionStage) {
2380 self.inner.stage_started(stage);
2381 }
2382
2383 fn stage_finished(&self, stage: ProvisionStage) {
2384 self.inner.stage_finished(stage);
2385 }
2386
2387 fn notify_notice(&self, notice: &str) {
2388 self.inner.notify_notice(notice);
2389 }
2390
2391 fn execute_with_stdin(
2392 &self,
2393 command: &CommandSpec,
2394 input: &mut (dyn std::io::Read + Send),
2395 ) -> Result<CommandOutput> {
2396 if self.cancellation.is_cancelled() {
2397 bail!("operation cancelled while provisioning destination");
2398 }
2399 self.inner.execute_with_stdin(command, input)
2400 }
2401}
2402
2403fn provision_with_cross_harness_handoff(
2404 controller: &mut Controller,
2405 session_id: &str,
2406 executor: &(impl CommandExecutor + Sync),
2407 github_token: Option<&str>,
2408 config: &Config,
2409 snapshot: &CanonicalSessionSnapshot,
2410 context_bytes: usize,
2411) -> Result<String> {
2412 let (_provision, handoff) = execute_joined_cross_harness_work(
2413 "cross-harness provisioning",
2414 move |cancellation| {
2415 let provision_executor = CrossHarnessProvisionExecutor {
2416 inner: executor,
2417 cancellation,
2418 };
2419 futures::executor::block_on(controller.provision_session_with_failure_disposition(
2420 session_id,
2421 &provision_executor,
2422 github_token,
2423 ProvisioningFailureDisposition::Preserve,
2424 ))
2425 },
2426 "cross-harness handoff",
2427 move |cancellation| {
2428 let runtime = tokio::runtime::Builder::new_current_thread()
2429 .enable_all()
2430 .build()
2431 .context("create cross-harness handoff runtime")?;
2432 runtime.block_on(utility_handoff_while_cancellable(
2433 session_id,
2434 config,
2435 snapshot,
2436 context_bytes,
2437 executor,
2438 cancellation,
2439 ))
2440 },
2441 )?;
2442 ensure!(
2443 !executor.cancellation_requested(),
2444 "operation cancelled while provisioning destination"
2445 );
2446 Ok(handoff)
2447}
2448
2449fn execute_joined_cross_harness_work<A: Send, B: Send>(
2453 first_name: &'static str,
2454 first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2455 second_name: &'static str,
2456 second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2457) -> Result<(A, B)> {
2458 let cancellation = CancellationToken::new();
2459 std::thread::scope(|scope| {
2460 let first_cancel = cancellation.clone();
2461 let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2462 let second_cancel = cancellation.clone();
2463 let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2464 let mut first_result = None;
2465 let mut second_result = None;
2466
2467 while first_result.is_none() || second_result.is_none() {
2468 if first_result.is_none()
2469 && first_handle
2470 .as_ref()
2471 .is_some_and(|handle| handle.is_finished())
2472 {
2473 let handle = first_handle.take().expect("first lane handle present");
2474 first_result = Some(match handle.join() {
2475 Ok(result) => result,
2476 Err(panic) => {
2477 cancellation.cancel();
2478 Err(anyhow::anyhow!(
2479 "{first_name} thread panicked: {}",
2480 targets::command_thread_panic_message(panic.as_ref())
2481 ))
2482 }
2483 });
2484 if first_result.as_ref().is_some_and(Result::is_err) {
2485 cancellation.cancel();
2486 }
2487 }
2488 if second_result.is_none()
2489 && second_handle
2490 .as_ref()
2491 .is_some_and(|handle| handle.is_finished())
2492 {
2493 let handle = second_handle.take().expect("second lane handle present");
2494 second_result = Some(match handle.join() {
2495 Ok(result) => result,
2496 Err(panic) => {
2497 cancellation.cancel();
2498 Err(anyhow::anyhow!(
2499 "{second_name} thread panicked: {}",
2500 targets::command_thread_panic_message(panic.as_ref())
2501 ))
2502 }
2503 });
2504 if second_result.as_ref().is_some_and(Result::is_err) {
2505 cancellation.cancel();
2506 }
2507 }
2508 if first_result.is_none() || second_result.is_none() {
2509 std::thread::sleep(Duration::from_millis(10));
2510 }
2511 }
2512
2513 match (
2514 first_result.expect("first lane result received after joined handle"),
2515 second_result.expect("second lane result received after joined handle"),
2516 ) {
2517 (Err(first), Err(second)) => {
2518 Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2519 }
2520 (Err(error), Ok(_)) => Err(error),
2521 (Ok(_), Err(error)) => Err(error),
2522 (Ok(first), Ok(second)) => Ok((first, second)),
2523 }
2524 })
2525}
2526
2527fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2531 profile_kind == archived_kind
2532}
2533
2534async fn utility_handoff_while_cancellable(
2538 session_id: &str,
2539 config: &Config,
2540 snapshot: &CanonicalSessionSnapshot,
2541 context_bytes: usize,
2542 executor: &impl CommandExecutor,
2543 cancellation: CancellationToken,
2544) -> Result<String> {
2545 let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2546 if executor.cancellation_requested() {
2547 bail!("operation cancelled while compacting the cross-harness handoff");
2548 }
2549 let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2550 let cancel = cancellation.child_token();
2551 let operation =
2552 crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2553 tokio::pin!(operation);
2554 loop {
2555 tokio::select! {
2556 context = &mut operation => return context,
2557 _ = cancellation.cancelled() => {
2558 cancel.cancel();
2559 bail!("operation cancelled while compacting the cross-harness handoff");
2560 }
2561 _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2562 if executor.cancellation_requested() {
2563 cancel.cancel();
2564 bail!("operation cancelled while compacting the cross-harness handoff");
2565 }
2566 }
2567 }
2568 }
2569}
2570
2571mod in_place;
2572pub(crate) use in_place::InPlaceRestartError;
2573
2574#[cfg(test)]
2575mod tests;