1use std::path::{Path, PathBuf};
4use std::sync::Arc;
5use std::time::{Duration, Instant};
6
7use agent_client_protocol::schema::v1::ContentBlock;
8use anyhow::{Context, Result, bail, ensure};
9use rayon::prelude::*;
10use tokio_util::sync::CancellationToken;
11
12use crate::checkpoint_transfer::restore_command;
13use crate::session_manager::new_command_id;
14use mj_checkpoint::archive::{
15 CanonicalQueuedCommandKind, CanonicalSessionSnapshot, CheckpointRepositoryBundle, SystemGit,
16 checkpoint_bundle_prerequisites, read_checkpoint_repository_bundles, verify_archive_streaming,
17};
18use mj_checkpoint::checkpoint::CheckpointRestoreSpec;
19use mj_core::config::{Config, HarnessKind, ProjectRepository, mount_history_host};
20use mj_core::state::{MaterializedSession, SessionRecord, SessionResourceAllocation, SessionState};
21use mj_transcript::projection::materialized_session_from_canonical;
22
23use crate::targets::{
24 self, AdditionalMount, CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec,
25 ProcessExecutor, ProvisionStage, ProvisionStageGuard,
26};
27use mj_core::relay::RelayCommand;
28
29use super::backend::{backend_locator, controller_github_token, validate_resource_allocation};
30use super::checkpoint::upload_checkpoint_spec;
31use super::provisioning::{
32 ProvisioningFailureDisposition, StagedExecutor, execute_concurrent_lanes,
33 install_attached_resources,
34};
35use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
36use super::worker_binary::{bridge_readiness_stage, start_worker, worker_probe_diagnosis};
37use super::worktree::{
38 PrimaryCheckoutRequirement, ResumeConversion, ResumePlan, apply_raw_to_workspace,
39 apply_workspace_to_raw, cleanup_managed_worktree, create_managed_worktree,
40 managed_worktree_checkout_exists, managed_worktree_target, plan_raw_to_workspace,
41 preserve_retained_managed_worktree_branch, raw_checkout_divergence_notice,
42 raw_checkout_position, raw_checkout_snapshot, raw_conversion_preview, restore_managed_worktree,
43 resume_compatibility, retire_managed_worktree,
44};
45use super::{
46 Controller, SessionResumeOptions, execute_checked, now, selected_host_container_size,
47 target_profile_home,
48};
49
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct ResumeRepositorySourceMismatch {
52 pub session_id: String,
53 pub bundle_id: String,
54 pub repository_id: String,
55 pub missing_commit: String,
56 pub archived_origin: String,
57 pub configured_origin: String,
58}
59
60pub use mj_core::state::ResumeRepositorySourceReceipt;
61
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum ResumeRepositorySourcePreflight {
64 Ready(ResumeRepositorySourceReceipt),
65 RepositoryMoved(ResumeRepositorySourceMismatch),
66 ConvertingRawCheckout {
71 receipt: ResumeRepositorySourceReceipt,
72 preview: Box<mj_core::state::RawConversionPreview>,
73 },
74}
75
76struct ResumeRepositoryBundles {
77 checkpoint_sha256: String,
78 repositories: Vec<CheckpointRepositoryBundle>,
79}
80
81struct ResumePhaseTimer<'a> {
85 session_id: &'a str,
86 phase: &'static str,
87 started: Instant,
88}
89
90impl<'a> ResumePhaseTimer<'a> {
91 fn new(session_id: &'a str, phase: &'static str) -> Self {
92 Self {
93 session_id,
94 phase,
95 started: Instant::now(),
96 }
97 }
98}
99
100impl Drop for ResumePhaseTimer<'_> {
101 fn drop(&mut self) {
102 tracing::debug!(
103 session_id = self.session_id,
104 phase = self.phase,
105 elapsed_ms = self.started.elapsed().as_millis(),
106 "resume phase completed"
107 );
108 }
109}
110
111impl Controller {
112 pub(super) fn validate_muse_resume_destination(
114 &self,
115 source: &SessionRecord,
116 destination_harness: HarnessKind,
117 _target_id: &str,
118 ) -> Result<()> {
119 if destination_harness != HarnessKind::Muse {
120 return Ok(());
121 }
122 ensure!(
123 source.project_directory.is_some()
124 || 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 let RestoreIntoTarget {
560 profile,
561 archive,
562 restored_archive,
563 resumed_project_directory,
564 resumed_container_workspace,
565 restore_repositories,
566 native_continuity,
567 discard_queued_prompts,
568 replay_queue,
569 utility_handoff,
570 projection_build,
571 mut resume_notices,
572 install_attached_resources: should_install_attached_resources,
573 worker_root_reset,
574 retire_after_ready,
575 } = restore;
576 let archive_manifest = &archive.manifest;
577 let canonical_session = &archive.canonical_session;
578 let stopped_subagents =
583 crate::database::load_stopped_subagents(session_id).unwrap_or_else(|error| {
584 tracing::warn!(
585 session_id,
586 error = format!("{error:#}"),
587 "could not read the sub-agents this session's suspend stopped"
588 );
589 Vec::new()
590 });
591 let stopped_subagents_context =
592 mj_core::subagent::stopped_subagents_prompt_context(&stopped_subagents);
593 resume_notices.extend(mj_core::subagent::stopped_subagents_notice(
594 &stopped_subagents,
595 ));
596 let (backend, worker_root) = self.worker_placement(session_id)?;
597 let harness_home = target_profile_home(&backend, session_id, profile);
598 let workspace_root = if let Some(project_directory) = &resumed_project_directory {
599 project_directory
600 .parent()
601 .context("bare project directory has no parent")?
602 .to_string_lossy()
603 .into_owned()
604 } else {
605 super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
606 };
607 let target_path = |path: &str| match &backend {
608 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
609 if !path.starts_with('/') =>
610 {
611 PathBuf::from(format!("~/{path}"))
612 }
613 _ => PathBuf::from(path),
614 };
615 let remote_archive = format!("{worker_root}/restore.hel.zip");
616 let remote_spec = format!("{worker_root}/restore-spec.json");
617 use mj_checkpoint::checkpoint::QueueRestorePolicy;
618 let move_admission =
619 crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
620 operation.phase == mj_core::state::MovePhase::ResumingDestination
621 && !operation.queue_admission_started
622 && operation.queue == mj_core::state::ResumeQueueDisposition::Start
623 });
624 let queue_policy = if move_admission || (!native_continuity && replay_queue) {
625 QueueRestorePolicy::Defer
626 } else if discard_queued_prompts {
627 QueueRestorePolicy::Discard
628 } else {
629 QueueRestorePolicy::Restore
630 };
631 let restore = CheckpointRestoreSpec {
632 archive_path: restore_archive_path(
633 &backend,
634 restored_archive,
635 &target_path(&remote_archive),
636 ),
637 workspace_root: target_path(&workspace_root),
638 relay_root: target_path(&worker_root),
639 harness_home: target_path(&harness_home),
640 restore_repositories,
646 restore_native: native_continuity,
647 primary_repository_root: resumed_project_directory
656 .as_ref()
657 .map(|directory| target_path(&directory.to_string_lossy())),
658 queue_policy,
659 };
660 {
664 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
665 match &worker_root_reset {
666 WorkerRootReset::FreshTarget => {
671 if let Some(command) = targets::clear_relay_state_plan(&backend, session_id)? {
672 execute_checked(syncing, command)?;
673 }
674 execute_checked(
677 syncing,
678 targets::command_on_locator(
679 &backend,
680 session_id,
681 vec!["mkdir".into(), "-p".into(), worker_root.clone()],
682 "create the session worker root",
683 )?,
684 )?;
685 }
686 WorkerRootReset::InPlace {
691 previous_profile_root,
692 } => {
693 execute_checked(
694 syncing,
695 targets::in_place_worker_reset_plan(
696 &backend,
697 session_id,
698 previous_profile_root,
699 )?,
700 )?;
701 }
702 }
703 }
704 let staging = tempfile::tempdir().context("create restore staging")?;
705 let local_spec = staging.path().join("restore-spec.json");
706 std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
707 let controller = &*self;
711 let backend_ref = &backend;
712 let worker_root_ref = worker_root.as_str();
713 let local_spec_ref = local_spec.as_path();
714 execute_concurrent_lanes(
715 || {
716 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
717 controller.prepare_worker_files(
718 session_id,
719 backend_ref,
720 worker_root_ref,
721 syncing,
722 )?;
723 super::provisioning::install_inherited_git_settings(
724 syncing,
725 backend_ref,
726 session_id,
727 )?;
728 Ok(())
729 },
730 || {
731 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
732 if should_upload_restore_archive(&backend) {
733 upload_checkpoint_spec(
734 restoring,
735 backend_ref,
736 session_id,
737 restored_archive,
738 &remote_archive,
739 )?;
740 }
741 upload_checkpoint_spec(
742 restoring,
743 backend_ref,
744 session_id,
745 local_spec_ref,
746 &remote_spec,
747 )
748 },
749 )?;
750 {
751 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
752 execute_checked(
753 restoring,
754 restore_command(&backend, session_id, &remote_spec)?,
755 )?;
756 }
757 if should_install_attached_resources {
758 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
759 install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
760 }
761 match projection_build {
762 Some(build) => {
763 let mut restored_projection = build
764 .await
765 .context("rebuild the restored projection")?
766 .context("rebuild the restored projection")?;
767 if discard_queued_prompts {
768 restored_projection.queued_prompts.clear();
769 }
770 crate::database::save_materialized_session(&restored_projection)?;
771 }
772 None if discard_queued_prompts => {
775 crate::database::replace_materialized_queued_prompts(session_id, &[])?;
776 }
777 None => {}
778 }
779 let readiness_stage = bridge_readiness_stage(profile);
780 let spec = self.reconnect_command(session_id)?;
781 let readiness = async {
782 let mut relay = {
783 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
784 start_worker(executor, &backend, &worker_root)?;
785 connect_started_worker(&spec, session_id, executor, &backend, &worker_root).await?
786 };
787 if let Some(context) = &stopped_subagents_context {
791 relay
792 .install_prompt_context(context.clone())
793 .await
794 .context("tell the resumed session which sub-agents its suspend stopped")?;
795 }
796 let native_session_id =
797 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
798 Ok::<_, anyhow::Error>((relay, native_session_id))
799 }
800 .await;
801 let (mut relay, native_session_id) = readiness
802 .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
803 if native_continuity {
804 if !restored_native_session_accepted(
805 &archive_manifest.session.native_session_id,
806 &native_session_id,
807 relay
808 .operational()
809 .replaced_unused_native_session_id
810 .as_deref(),
811 ) {
812 bail!(
813 "ACP loaded native session {native_session_id}, expected {}",
814 archive_manifest.session.native_session_id
815 );
816 }
817 } else {
818 relay
819 .install_prompt_context(
820 utility_handoff
821 .clone()
822 .context("a resume into a fresh native session has no handoff")?,
823 )
824 .await?;
825 if replay_queue {
826 for prompt in &canonical_session.queued_prompts {
827 let command = match &prompt.kind {
831 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
832 prompt: prompt
833 .content
834 .iter()
835 .cloned()
836 .map(serde_json::from_value)
837 .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
838 },
839 CanonicalQueuedCommandKind::SetConfig { key, value } => {
840 RelayCommand::SetConfig {
841 key: key.clone(),
842 value: value.clone(),
843 }
844 }
845 };
846 relay.submit(prompt.command_id.clone(), command).await?;
847 }
848 }
849 }
850 if let Some(worktree) = retire_after_ready
854 && let Err(error) = retire_managed_worktree(executor, worktree)
855 {
856 tracing::warn!(
857 session_id,
858 worktree = %worktree.worktree_root.display(),
859 error = format!("{error:#}"),
860 "could not retire the old managed worktree after resume"
861 );
862 resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
863 }
864 for notice in &resume_notices {
865 let submitted = async {
866 let command_id = new_command_id("resume-notice")?;
867 relay
868 .submit(
869 command_id,
870 RelayCommand::RecordNotice {
871 text: notice.clone(),
872 },
873 )
874 .await
875 }
876 .await;
877 if let Err(error) = submitted {
880 tracing::warn!(
881 session_id,
882 error = format!("{error:#}"),
883 "could not record a resume notice in the conversation"
884 );
885 }
886 }
887 self.mark_worker_connected(session_id, Some(native_session_id))?;
888 let materialized = relay.sync().await?.materialized;
889 if !stopped_subagents.is_empty() {
890 let delivered = stopped_subagents
891 .iter()
892 .map(|child| child.child_session_id.clone())
893 .collect::<Vec<_>>();
894 if let Err(error) = crate::database::clear_stopped_subagents(session_id, &delivered) {
897 tracing::warn!(
898 session_id,
899 error = format!("{error:#}"),
900 "could not clear the stopped sub-agents after telling the resumed session"
901 );
902 }
903 }
904 Ok(materialized)
905 }
906}
907
908pub(super) struct VerifiedResumeArchive {
915 pub archive_path: PathBuf,
919 pub manifest: mj_checkpoint::archive::ArchiveManifest,
920 pub canonical_session: Arc<CanonicalSessionSnapshot>,
921}
922
923pub(super) fn verify_resume_checkpoint(
929 session_id: &str,
930 checkpoint: &mj_core::state::CheckpointMetadata,
931) -> Result<VerifiedResumeArchive> {
932 let archive_path = {
933 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
934 checkpoint.archive_path.canonicalize().with_context(|| {
935 format!(
936 "resolve checkpoint archive {}",
937 checkpoint.archive_path.display()
938 )
939 })?
940 };
941 ensure!(
942 archive_path.is_absolute() && archive_path.is_file(),
943 "checkpoint archive path is not an absolute regular file: {}",
944 archive_path.display()
945 );
946 let mj_checkpoint::archive::VerifiedArchiveMetadata {
947 manifest,
948 canonical_session,
949 archive_sha256,
950 } = {
951 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
952 verify_archive_streaming(&archive_path)?
953 };
954 if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
955 bail!("persisted checkpoint verification failed");
956 }
957 ensure!(
958 manifest.repositories.iter().all(|repository| {
959 !repository.metadata.origin.starts_with("mj-local:")
960 && !repository.metadata.origin.starts_with("ext::")
961 }),
962 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
963 );
964 Ok(VerifiedResumeArchive {
965 archive_path,
966 manifest,
967 canonical_session: Arc::new(canonical_session),
968 })
969}
970
971pub fn raw_conversion_preview_for(
972 session: &SessionRecord,
973 config: &Config,
974 executor: &(impl CommandExecutor + Sync),
975) -> Result<mj_core::state::RawConversionPreview> {
976 let conversion = plan_raw_to_workspace(session, config, executor)?;
977 raw_conversion_preview(session, &conversion, executor)
978}
979
980fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
981 let replacement = replacement.trim();
982 ensure!(!replacement.is_empty(), "enter the repository's new origin");
983 let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
984 let path = expanded.as_path();
985 let (github, local) = if path.is_absolute() {
986 ensure!(
987 path.is_dir(),
988 "local repository {replacement:?} is not a directory"
989 );
990 (None, Some(mj_core::local_git::canonical_repository(path)?))
991 } else {
992 let github = crate::setup::github_repository_from_origin(replacement)
993 .context("origin must be a GitHub repository or an absolute local repository path")?;
994 (
995 Some(format!("{}/{}", github.owner, github.repository)),
996 None,
997 )
998 };
999 Ok(ProjectRepository {
1000 id: id.to_owned(),
1001 github,
1002 local,
1003 destination: PathBuf::from(id),
1004 git_ref: None,
1005 })
1006}
1007
1008fn checkpoint_source_missing_commit(
1009 configured: &ProjectRepository,
1010 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1011 archived: &CheckpointRepositoryBundle,
1012 executor: &impl CommandExecutor,
1013 github_token: Option<&str>,
1014) -> Result<Option<String>> {
1015 let staging = tempfile::tempdir().context("create repository source preflight")?;
1016 let repository = staging.path().join("repository.git");
1017 checked_preflight_git(
1018 executor,
1019 CommandSpec::new(
1020 "git",
1021 [
1022 "init".to_owned(),
1023 "--bare".to_owned(),
1024 "--quiet".to_owned(),
1025 repository.to_string_lossy().into_owned(),
1026 ],
1027 )
1028 .purpose("initialize repository source preflight"),
1029 )?;
1030 let missing = checkpoint_bundle_prerequisites(archived)?;
1031 if missing.is_empty() {
1032 let bundle = staging.path().join("checkpoint.bundle");
1033 std::fs::write(&bundle, &archived.committed_bundle)
1034 .context("write self-contained checkpoint bundle for source preflight")?;
1035 checked_preflight_git(
1036 executor,
1037 checkpoint_bundle_import_command(&repository, &bundle),
1038 )?;
1039 return Ok(None);
1040 }
1041 for commit in missing {
1045 let output = fetch_source_commit(
1046 executor,
1047 &repository,
1048 configured,
1049 accepted_source,
1050 &commit,
1051 github_token,
1052 )?;
1053 if output.status != 0 {
1054 let stderr = String::from_utf8_lossy(&output.stderr);
1055 if source_does_not_have_commit(&stderr) {
1056 return Ok(Some(commit));
1057 }
1058 bail!(
1059 "could not check configured source {:?}: {}",
1060 configured.source_label(),
1061 stderr.trim()
1062 );
1063 }
1064 }
1065 Ok(None)
1069}
1070
1071fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
1072 let mut command = CommandSpec::new(
1073 "git",
1074 [
1075 "-C".to_owned(),
1076 repository.to_string_lossy().into_owned(),
1077 "fetch".to_owned(),
1078 "--no-tags".to_owned(),
1079 bundle.to_string_lossy().into_owned(),
1080 "HEAD".to_owned(),
1081 ],
1082 )
1083 .purpose("validate self-contained checkpoint bundle");
1084 command
1085 .env
1086 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1087 command
1088 .env
1089 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1090 command
1091}
1092
1093fn fetch_source_commit(
1094 executor: &impl CommandExecutor,
1095 repository: &Path,
1096 configured: &ProjectRepository,
1097 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1098 commit: &str,
1099 github_token: Option<&str>,
1100) -> Result<CommandOutput> {
1101 let mut arguments = Vec::new();
1102 let mut token_auth = false;
1103 let mut ssh_transport = false;
1104 let source = if let Some(source) = accepted_source {
1105 source.fetch_url.clone()
1106 } else if let Some(local) = &configured.local {
1107 local.to_string_lossy().into_owned()
1108 } else {
1109 let source = configured
1110 .github
1111 .as_deref()
1112 .context("repository source is missing")?;
1113 let github = crate::setup::github_repository_from_origin(source)
1114 .context("configured repository is not a GitHub source")?;
1115 if github_token.is_some() {
1116 token_auth = true;
1117 arguments.extend([
1118 "-c".to_owned(),
1119 "credential.helper=".to_owned(),
1120 "-c".to_owned(),
1121 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1122 ]);
1123 format!(
1124 "https://github.com/{}/{}.git",
1125 github.owner, github.repository
1126 )
1127 } else {
1128 ssh_transport = true;
1129 format!("git@github.com:{}/{}.git", github.owner, github.repository)
1130 }
1131 };
1132 if accepted_source.is_some() {
1133 if source.starts_with("https://github.com/") && github_token.is_some() {
1134 token_auth = true;
1135 arguments.extend([
1136 "-c".to_owned(), "credential.helper=".to_owned(), "-c".to_owned(),
1137 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1138 ]);
1139 }
1140 ssh_transport = source.starts_with("git@") || source.starts_with("ssh://");
1141 }
1142 arguments.extend([
1143 "-C".to_owned(),
1144 repository.to_string_lossy().into_owned(),
1145 "fetch".to_owned(),
1146 "--no-tags".to_owned(),
1147 "--depth=1".to_owned(),
1148 "--filter=blob:none".to_owned(),
1149 source,
1150 commit.to_owned(),
1151 ]);
1152 let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1153 command
1154 .env
1155 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1156 command
1157 .env
1158 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1159 if token_auth {
1160 let token = github_token.expect("token authentication requires a GitHub token");
1161 command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1162 }
1163 if ssh_transport {
1164 command.env.insert(
1165 "GIT_SSH_COMMAND".to_owned(),
1166 "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1167 .to_owned(),
1168 );
1169 }
1170 executor.execute(&command)
1171}
1172
1173fn source_does_not_have_commit(stderr: &str) -> bool {
1174 let stderr = stderr.to_ascii_lowercase();
1175 [
1176 "not our ref",
1177 "couldn't find remote ref",
1178 "not a valid object name",
1179 "no such ref was fetched",
1180 ]
1181 .iter()
1182 .any(|needle| stderr.contains(needle))
1183}
1184
1185fn checked_preflight_git(
1186 executor: &impl CommandExecutor,
1187 command: CommandSpec,
1188) -> Result<CommandOutput> {
1189 let output = executor.execute(&command)?;
1190 ensure!(
1191 output.status == 0,
1192 "{}: {}",
1193 command.purpose,
1194 String::from_utf8_lossy(&output.stderr).trim()
1195 );
1196 Ok(output)
1197}
1198
1199impl Controller {
1200 pub async fn resume_session_with_options(
1205 &mut self,
1206 session_id: &str,
1207 profile_id: &str,
1208 target_id: &str,
1209 additional_mounts: Option<Vec<AdditionalMount>>,
1210 resource_allocation: Option<SessionResourceAllocation>,
1211 ) -> Result<MaterializedSession> {
1212 self.resume_session_with_options_and_queue_disposition(
1213 session_id,
1214 profile_id,
1215 target_id,
1216 additional_mounts,
1217 resource_allocation,
1218 false,
1219 )
1220 .await
1221 }
1222
1223 pub async fn resume_session_with_options_and_queue_disposition(
1224 &mut self,
1225 session_id: &str,
1226 profile_id: &str,
1227 target_id: &str,
1228 additional_mounts: Option<Vec<AdditionalMount>>,
1229 resource_allocation: Option<SessionResourceAllocation>,
1230 discard_queue: bool,
1231 ) -> Result<MaterializedSession> {
1232 self.resume_session_controlled(
1233 session_id,
1234 profile_id,
1235 target_id,
1236 SessionResumeOptions {
1237 additional_mounts,
1238 resource_allocation,
1239 discard_queue,
1240 },
1241 &ProcessExecutor,
1242 )
1243 .await
1244 }
1245
1246 pub async fn resume_session_controlled(
1247 &mut self,
1248 session_id: &str,
1249 profile_id: &str,
1250 target_id: &str,
1251 options: SessionResumeOptions,
1252 executor: &(impl CommandExecutor + Sync),
1253 ) -> Result<MaterializedSession> {
1254 self.resume_session_controlled_with_repository_preflight(
1255 session_id, profile_id, target_id, options, None, executor,
1256 )
1257 .await
1258 }
1259
1260 pub async fn resume_session_controlled_with_repository_preflight(
1261 &mut self,
1262 session_id: &str,
1263 profile_id: &str,
1264 target_id: &str,
1265 options: SessionResumeOptions,
1266 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1267 executor: &(impl CommandExecutor + Sync),
1268 ) -> Result<MaterializedSession> {
1269 self.resume_session_with_origin(
1270 session_id,
1271 profile_id,
1272 target_id,
1273 options,
1274 repository_preflight,
1275 None,
1276 executor,
1277 )
1278 .await
1279 }
1280
1281 pub(in crate::controller) async fn resume_session_for_move(
1282 &mut self,
1283 operation: &mj_core::state::MoveOperation,
1284 executor: &(impl CommandExecutor + Sync),
1285 ) -> Result<MaterializedSession> {
1286 self.resume_session_with_origin(
1287 &operation.selection.session_id,
1288 operation
1289 .selection
1290 .profile_id
1291 .as_deref()
1292 .context("Move profile missing")?,
1293 operation
1294 .selection
1295 .target_template_id
1296 .as_deref()
1297 .context("Move target missing")?,
1298 SessionResumeOptions {
1299 additional_mounts: operation.selection.additional_mounts.clone(),
1300 resource_allocation: operation.selection.resource_allocation.clone(),
1301 discard_queue: true,
1302 },
1303 None,
1304 Some(operation),
1305 executor,
1306 )
1307 .await
1308 }
1309
1310 #[allow(clippy::too_many_arguments)]
1311 async fn resume_session_with_origin(
1312 &mut self,
1313 session_id: &str,
1314 profile_id: &str,
1315 target_id: &str,
1316 options: SessionResumeOptions,
1317 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1318 move_operation: Option<&mj_core::state::MoveOperation>,
1319 executor: &(impl CommandExecutor + Sync),
1320 ) -> Result<MaterializedSession> {
1321 let transferring_workspace = move_operation.is_some();
1322 if let Some(operation) = crate::database::load_move_operation(session_id)? {
1323 ensure!(
1324 transferring_workspace || !operation.holds_source_environment(),
1325 "a profile switch retains this environment; retry Move instead of recreating it with Resume"
1326 );
1327 }
1328 let SessionResumeOptions {
1329 additional_mounts,
1330 resource_allocation,
1331 discard_queue,
1332 } = options;
1333 let mut previous = self
1334 .state
1335 .sessions
1336 .get(session_id)
1337 .with_context(|| format!("unknown session {session_id}"))?
1338 .clone();
1339 if !(matches!(
1340 previous.state,
1341 SessionState::Stopped | SessionState::Lost | SessionState::Error
1342 ) || (transferring_workspace && previous.state == SessionState::Closing))
1343 {
1344 bail!("session {session_id} is not stopped, lost, or retryable");
1345 }
1346 let checkpoint = move_operation
1347 .and_then(|op| op.handoff.as_ref())
1348 .or(previous.checkpoint.as_ref())
1349 .context("session has no checkpoint")?;
1350 let moving_to_raw = previous.project_directory.is_none()
1353 && self
1354 .config
1355 .targets
1356 .get(target_id)
1357 .is_some_and(mj_core::config::is_bare_project_target);
1358 if !transferring_workspace
1359 && (moving_to_raw
1360 || !repository_preflight.as_ref().is_some_and(|receipt| {
1361 self.repository_source_receipt_is_current(session_id, receipt)
1362 }))
1363 {
1364 let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1365 if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1366 self.preflight_repository_sources(session_id, target_id, false, executor)?
1367 {
1368 bail!(
1369 "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1370 mismatch.missing_commit,
1371 mismatch.configured_origin,
1372 mismatch.repository_id,
1373 mismatch.archived_origin,
1374 );
1375 }
1376 }
1377 let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1378 let archive_path = verified_archive.archive_path.clone();
1379 let archive_manifest = &verified_archive.manifest;
1380 let canonical_session = Arc::clone(&verified_archive.canonical_session);
1381 let profile = self
1382 .config
1383 .profiles
1384 .get(profile_id)
1385 .with_context(|| format!("unknown profile {profile_id:?}"))?
1386 .clone();
1387 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1388 let target_template = self
1389 .config
1390 .targets
1391 .get(target_id)
1392 .with_context(|| format!("unknown target template {target_id:?}"))?
1393 .clone();
1394 self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1397 ensure!(
1398 profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1399 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1400 );
1401 let plan = resume_compatibility(&previous, &self.config, target_id)
1402 .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1403 if plan == ResumePlan::InPlace
1406 && let Some(worktree) = previous.managed_worktree.as_mut()
1407 && let Ok(current) = managed_worktree_target(&target_template)
1408 {
1409 worktree.target = current;
1410 }
1411 if !mj_core::config::is_bare_project_target(&target_template)
1415 && plan != ResumePlan::RawToWorkspace
1416 && !transferring_workspace
1417 {
1418 super::network_git::bundle_from_manifest(archive_manifest)?;
1419 }
1420 if plan == ResumePlan::InPlace
1421 && previous.managed_worktree.is_none()
1422 && let Some(project_directory) = &previous.project_directory
1423 {
1424 self.validate_project_directory(target_id, project_directory, executor)
1425 .context("raw project is unavailable for resume")?;
1426 }
1427 let conversion = match plan {
1428 ResumePlan::InPlace => None,
1429 ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(
1430 plan_raw_to_workspace(&previous, &self.config, executor)
1431 .context("prepare the raw checkout for its new target")?,
1432 )),
1433 ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1434 self.plan_workspace_to_raw(&previous, target_id, executor)
1435 .context("prepare a checkout for this session")?,
1436 )),
1437 };
1438 let resource_allocation =
1439 resource_allocation.or_else(|| previous.resource_allocation.clone());
1440 let additional_mounts =
1441 additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1442 validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1443 let selected_container_size =
1444 selected_host_container_size(&target_template, resource_allocation.as_ref());
1445 if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1446 bail!("attached resources are unsupported for this target");
1447 }
1448 targets::validate_additional_mounts(&additional_mounts)?;
1449 let history_host = mount_history_host(&target_template);
1450 let history_mounts = additional_mounts.clone();
1451 if !transferring_workspace
1452 && previous.state == SessionState::Error
1453 && let Some(locator) = &previous.target
1454 {
1455 let backend = backend_locator(locator, &previous, &self.config)?;
1456 targets::close_plan(&backend, session_id)?
1457 .execute(executor)
1458 .context("clean up target from failed resume")?;
1459 }
1460 let mut resume_notices = Vec::new();
1461 if let Some(conversion) = conversion
1464 .as_ref()
1465 .and_then(ResumeConversion::workspace_to_raw)
1466 {
1467 resume_notices.push(format!(
1468 "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1469 previous.target_template_id,
1470 conversion.worktree.worktree_root.display(),
1471 archive_manifest
1472 .repositories
1473 .first()
1474 .and_then(|repository| repository.metadata.branch.as_deref())
1475 .unwrap_or("a detached head"),
1476 conversion.worktree.branch,
1477 ));
1478 }
1479 let managed_checkout_present = previous
1480 .managed_worktree
1481 .as_ref()
1482 .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1483 .transpose()?
1484 .unwrap_or(true);
1485 if managed_checkout_present && let Some(project_directory) = &previous.project_directory {
1488 match raw_checkout_position(&previous, &self.config, project_directory, executor) {
1489 Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1490 project_directory,
1491 archive_manifest
1492 .repositories
1493 .first()
1494 .map(|repository| &repository.metadata),
1495 &live,
1496 )),
1497 Err(error) => tracing::warn!(
1500 session_id,
1501 error = format!("{error:#}"),
1502 "could not read the raw checkout position for a resume notice"
1503 ),
1504 }
1505 }
1506 super::worker_binary::preflight_worker_binary(&target_template, executor)?;
1511 let same_harness = profile.kind == archive_manifest.session.harness_kind;
1512 let native_continuity =
1513 native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1514 let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1515 let utility_config = (!native_continuity).then(|| self.config.clone());
1519 let discard_queued_prompts = discard_queue || !same_harness;
1520 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1524 .unwrap_or_else(|error| {
1525 tracing::warn!(
1526 session_id,
1527 error = format!("{error:#}"),
1528 "could not read the stored projection frontier; rebuilding it from the archive"
1529 );
1530 None
1531 });
1532 let rebuild_projection = projection_rebuild_required(
1533 stored_frontier
1534 .as_ref()
1535 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1536 canonical_session.event_frontier,
1537 &canonical_session.event_frontier_digest,
1538 );
1539 let projection_build = rebuild_projection.then(|| {
1549 let canonical = Arc::clone(&canonical_session);
1550 let session_id = session_id.to_owned();
1551 tokio::task::spawn_blocking(move || {
1552 materialized_session_from_canonical(session_id, &canonical)
1553 })
1554 });
1555 let github_token = controller_github_token();
1556
1557 if let Some(conversion) = conversion
1560 .as_ref()
1561 .and_then(ResumeConversion::raw_to_workspace)
1562 && let Some(bundle) = &conversion.new_bundle
1563 {
1564 let (config, ()) = Config::update(|config| {
1565 if let Some(existing) = config.bundles.get(&conversion.bundle_id) {
1566 ensure!(
1567 existing == bundle,
1568 "bundle {:?} was configured concurrently with a different definition; retry the resume",
1569 conversion.bundle_id
1570 );
1571 } else {
1572 config
1573 .bundles
1574 .insert(conversion.bundle_id.clone(), bundle.clone());
1575 }
1576 Ok(())
1577 })
1578 .context("save the bundle for a converted raw session")?;
1579 self.config = config;
1580 }
1581
1582 let moving_into_first_container = self
1588 .config
1589 .targets
1590 .get(target_id)
1591 .is_some_and(mj_core::config::is_container_target)
1592 && !self
1593 .config
1594 .targets
1595 .get(&previous.target_template_id)
1596 .is_some_and(mj_core::config::is_container_target);
1597 let converted_project = conversion
1598 .as_ref()
1599 .and_then(ResumeConversion::raw_to_workspace)
1600 .map(|conversion| {
1601 crate::project_catalog::snapshot(
1602 self.config
1603 .bundles
1604 .get(&conversion.bundle_id)
1605 .context("converted project is missing")?,
1606 executor,
1607 true,
1608 )
1609 })
1610 .transpose()?;
1611 let record = self.state.sessions.get_mut(session_id).unwrap();
1612 if record.container_workspace.is_none() && moving_into_first_container {
1613 record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1614 }
1615 record.harness_kind = profile.kind;
1616 record.last_profile = profile_id.to_string();
1617 record.target_template_id = target_id.to_string();
1618 record.target_runtime = Some(
1619 self.config
1620 .targets
1621 .get(target_id)
1622 .context("resume target disappeared")?
1623 .into(),
1624 );
1625 record.resource_allocation = resource_allocation;
1626 record.additional_mounts = additional_mounts;
1627 record.target = None;
1628 record.native_session_id =
1629 native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1630 record.state = SessionState::Provisioning;
1631 record.updated_at = now();
1632 record.last_error = None;
1633 match &conversion {
1634 Some(ResumeConversion::RawToWorkspace(conversion)) => {
1635 apply_raw_to_workspace(record, conversion);
1636 record.project = converted_project;
1637 }
1638 Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1639 apply_workspace_to_raw(record, conversion);
1640 }
1641 None => {
1642 if let (Some(worktree), Some(refreshed)) = (
1643 record.managed_worktree.as_mut(),
1644 previous.managed_worktree.as_ref(),
1645 ) {
1646 worktree.target = refreshed.target.clone();
1647 }
1648 }
1649 }
1650 let resumed_project_directory = record.project_directory.clone();
1651 let resumed_container_workspace = record.container_workspace.clone();
1652 if let Some(host) = history_host {
1653 self.state.remember_mount_sources(host, &history_mounts);
1654 crate::database::remember_mount_sources(host, &history_mounts)?;
1655 }
1656 if let Some(conversion) = conversion
1659 .as_ref()
1660 .and_then(ResumeConversion::raw_to_workspace)
1661 {
1662 crate::database::rebind_session_bundle(session_id, &conversion.bundle_id)?;
1663 }
1664 crate::database::save_resumed_session(
1667 &self.state.sessions[session_id],
1668 selected_container_size
1669 .as_ref()
1670 .map(|(host, size)| (host.as_str(), *size)),
1671 )?;
1672 if let Some((host, size)) = selected_container_size.as_ref() {
1673 self.state.remember_container_size(host, *size);
1674 }
1675
1676 let mut recreated_managed_worktree = false;
1677 let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1680 let result = async {
1681 if let Some(worktree) = previous.managed_worktree.as_ref() {
1682 recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1683 if !transferring_workspace && recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1684 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1685 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1686 &archive_path, &worktree.worktree_root, &SystemGit,
1687 )?;
1688 } else {
1689 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1690 &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1691 )?;
1692 }
1693 }
1694 }
1695 if let Some(conversion) = conversion
1698 .as_ref()
1699 .and_then(ResumeConversion::workspace_to_raw)
1700 {
1701 if conversion.reuse_existing_branch {
1702 let recovery_ref =
1703 preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1704 restore_managed_worktree(executor, &conversion.worktree)?;
1705 resume_notices.push(format!(
1706 "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1707 ));
1708 } else {
1709 create_managed_worktree(
1710 executor,
1711 &conversion.worktree,
1712 None,
1713 PrimaryCheckoutRequirement::Any,
1714 )?;
1715 }
1716 if transferring_workspace {
1717 } else if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1719 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1720 &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1721 )?;
1722 } else {
1723 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1724 &archive_path,
1725 &conversion.worktree.worktree_root,
1726 &conversion.worktree.branch,
1727 &SystemGit,
1728 )?;
1729 }
1730 }
1731 if let Some(conversion) = conversion
1736 .as_ref()
1737 .and_then(ResumeConversion::raw_to_workspace)
1738 && !transferring_workspace
1739 {
1740 let destination = PathBuf::from(
1741 previous
1742 .project_directory
1743 .as_deref()
1744 .context("a raw session has no project directory")?
1745 .file_name()
1746 .context("a raw project directory cannot be the filesystem root")?,
1747 );
1748 let snapshot = raw_checkout_snapshot(
1749 &conversion.checkout,
1750 &conversion.source,
1751 &destination,
1752 &SystemGit,
1753 conversion.retire.as_ref().is_some_and(|checkout| {
1754 checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1755 }),
1756 )
1757 .context("snapshot the host checkout for its new target")?;
1758 resume_notices.push(conversion_notice(
1759 target_id,
1760 previous
1761 .project_directory
1762 .as_deref()
1763 .unwrap_or(&conversion.checkout),
1764 snapshot.metadata.branch.as_deref(),
1765 conversion.retire.as_ref(),
1766 ));
1767 let archives = mj_core::config::sessions_dir();
1768 std::fs::create_dir_all(&archives).with_context(|| {
1769 format!("create the checkpoint directory {}", archives.display())
1770 })?;
1771 let output = archives.join(format!(
1775 "{session_id}-converted-{}-{}.hel.zip",
1776 previous
1777 .checkpoint
1778 .as_ref()
1779 .map_or(0, |checkpoint| checkpoint.event_frontier),
1780 new_command_id("archive")?
1781 ));
1782 let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1783 conversion_checkpoint_written = Some(written.clone());
1784 let record = self.state.sessions.get_mut(session_id).unwrap();
1785 record.checkpoint = Some(written);
1786 record.updated_at = now();
1787 crate::database::save_resumed_session(
1790 &self.state.sessions[session_id],
1791 selected_container_size.as_ref().map(|(host, size)| (host.as_str(), *size)),
1792 )?;
1793 }
1794 let utility_handoff = {
1795 let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1796 if let Some(config) = utility_config.as_ref() {
1797 Some(
1798 provision_with_cross_harness_handoff(
1799 self,
1800 session_id,
1801 executor,
1802 github_token.as_deref(),
1803 config,
1804 &canonical_session,
1805 context_bytes,
1806 )
1807 .context("prepare the cross-harness destination")?,
1808 )
1809 } else {
1810 self.provision_session_with_failure_disposition(
1811 session_id,
1812 executor,
1813 github_token.as_deref(),
1814 ProvisioningFailureDisposition::Preserve,
1815 )
1816 .await?;
1817 None
1818 }
1819 };
1820 if transferring_workspace {
1821 self.restore_move_workspace(session_id, executor)?;
1822 }
1823 let restore_repositories = !transferring_workspace && ((resumed_project_directory.is_none()
1824 && conversion.is_none())
1825 || plan == ResumePlan::RawToWorkspace
1826 || (recreated_managed_worktree && plan == ResumePlan::InPlace));
1827 let restored_archive = conversion_checkpoint_written
1830 .as_ref()
1831 .map_or(archive_path.as_path(), |checkpoint| {
1832 checkpoint.archive_path.as_path()
1833 });
1834 self.restore_into_target(
1835 session_id,
1836 RestoreIntoTarget {
1837 profile: &profile,
1838 archive: &verified_archive,
1839 restored_archive,
1840 resumed_project_directory,
1841 resumed_container_workspace,
1842 restore_repositories,
1843 native_continuity,
1844 discard_queued_prompts,
1845 replay_queue: !discard_queue,
1846 utility_handoff,
1847 projection_build,
1848 resume_notices,
1849 install_attached_resources: true,
1850 worker_root_reset: WorkerRootReset::FreshTarget,
1851 retire_after_ready: conversion
1852 .as_ref()
1853 .filter(|_| !transferring_workspace)
1854 .and_then(ResumeConversion::raw_to_workspace)
1855 .and_then(|plan| plan.retire.as_ref()),
1856 },
1857 executor,
1858 )
1859 .await
1860 }
1861 .await;
1862 match result {
1863 Ok(materialized) => {
1864 if let Some(written) = &conversion_checkpoint_written {
1867 super::checkpoint::prune_replaced_checkpoint(
1868 previous.checkpoint.as_ref(),
1869 written,
1870 );
1871 }
1872 Ok(materialized)
1873 }
1874 Err(error) => {
1875 if let Some(written) = &conversion_checkpoint_written
1878 && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
1879 && remove_error.kind() != std::io::ErrorKind::NotFound
1880 {
1881 tracing::warn!(
1882 session_id,
1883 path = %written.archive_path.display(),
1884 "could not remove the conversion checkpoint after resume failed: {remove_error}"
1885 );
1886 }
1887 restore_projection_after_failed_resume(
1890 session_id,
1891 &canonical_session,
1892 discard_queued_prompts,
1893 );
1894 if let Some(operation) = move_operation {
1895 return Err(self.rollback_move_destination(operation, error, executor)?);
1896 }
1897 Err(self.rollback_failed_resume(
1898 session_id,
1899 &previous,
1900 recreated_managed_worktree,
1901 error,
1902 executor,
1903 )?)
1904 }
1905 }
1906 }
1907
1908 pub(super) fn rollback_failed_resume(
1909 &mut self,
1910 session_id: &str,
1911 previous: &SessionRecord,
1912 recreated_managed_worktree: bool,
1913 error: anyhow::Error,
1914 _executor: &impl CommandExecutor,
1915 ) -> Result<anyhow::Error> {
1916 let current = self
1917 .state
1918 .sessions
1919 .get(session_id)
1920 .with_context(|| format!("unknown session {session_id}"))?
1921 .clone();
1922 let cleanup = match current.target.as_ref() {
1923 Some(locator) => (|| -> Result<()> {
1924 let backend = backend_locator(locator, ¤t, &self.config)?;
1925 targets::close_plan(&backend, session_id)?
1926 .execute(&CancellableProcessExecutor::with_timeout(
1929 Duration::from_secs(15),
1930 ))
1931 .map(|_| ())
1932 })(),
1933 None => Ok(()),
1934 };
1935 let worktree_cleanup = if cleanup.is_err() {
1938 Ok(())
1939 } else {
1940 match (
1941 current.managed_worktree.as_ref(),
1942 previous.managed_worktree.as_ref(),
1943 ) {
1944 (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
1945 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1946 previous,
1947 ),
1948 (Some(current), Some(previous)) if current == previous => Ok(()),
1949 (Some(worktree), _) => cleanup_managed_worktree(
1952 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1953 worktree,
1954 crate::controller::BranchDisposition::Delete,
1955 ),
1956 (None, _) => Ok(()),
1957 }
1958 };
1959 let cleanup_error = [cleanup, worktree_cleanup]
1960 .into_iter()
1961 .filter_map(Result::err)
1962 .map(|cleanup_error| format!("{cleanup_error:#}"))
1963 .collect::<Vec<_>>()
1964 .join("; ");
1965 if !cleanup_error.is_empty() {
1966 tracing::warn!(
1967 session_id,
1968 error = %cleanup_error,
1969 "resume rollback cleanup reported failures"
1970 );
1971 }
1972 let original = super::worker_binary::failure_line(&error);
1975 let detail = format!("{error:#}");
1976 if detail != original {
1977 tracing::warn!(session_id, error = %detail, "resume failed");
1978 }
1979 let record = self.state.sessions.get_mut(session_id).unwrap();
1980 let failure = apply_failed_resume_rollback(
1981 record,
1982 previous,
1983 &original,
1984 (!cleanup_error.is_empty()).then_some(cleanup_error),
1985 );
1986 if record.bundle_id != current.bundle_id {
1989 let bundle_id = record.bundle_id.clone();
1990 crate::database::rebind_session_bundle(session_id, &bundle_id)?;
1991 }
1992 crate::database::save_resumed_session(&self.state.sessions[session_id], None)?;
1995 Ok(failure)
1996 }
1997}
1998
1999fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
2000 format!(
2001 "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
2002 worktree_root.display(),
2003 worktree_root.display()
2004 )
2005}
2006
2007pub(super) fn apply_failed_resume_rollback(
2008 current: &mut SessionRecord,
2009 previous: &SessionRecord,
2010 original_error: &str,
2011 cleanup_error: Option<String>,
2012) -> anyhow::Error {
2013 match cleanup_error {
2014 None => {
2015 *current = previous.clone();
2016 current.state = SessionState::Stopped;
2017 current.target = None;
2018 current.updated_at = now();
2019 current.last_error = Some(format!("resume failed: {original_error}"));
2020 anyhow::anyhow!(original_error.to_owned())
2021 }
2022 Some(cleanup_error) => {
2023 let failure = format!(
2024 "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
2025 );
2026 if current.managed_worktree.is_none() {
2031 current
2032 .project_directory
2033 .clone_from(&previous.project_directory);
2034 current
2035 .managed_worktree
2036 .clone_from(&previous.managed_worktree);
2037 current.bundle_id.clone_from(&previous.bundle_id);
2038 }
2039 current.state = SessionState::Error;
2040 current.updated_at = now();
2041 current.last_error = Some(format!("resume failed: {failure}"));
2042 anyhow::anyhow!(failure)
2043 }
2044 }
2045}
2046
2047fn conversion_notice(
2049 target_id: &str,
2050 checkout: &Path,
2051 branch: Option<&str>,
2052 retire: Option<&mj_core::state::ManagedWorktree>,
2053) -> String {
2054 let branch = branch.unwrap_or("a detached head");
2055 match retire {
2056 Some(worktree) => format!(
2057 "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
2058 checkout.display(),
2059 worktree.branch,
2060 worktree.source_repository.display()
2061 ),
2062 None => format!(
2063 "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.",
2064 checkout.display()
2065 ),
2066 }
2067}
2068
2069fn conversion_checkpoint(
2077 previous_archive: &Path,
2078 snapshot: mj_checkpoint::archive::RepositorySnapshot,
2079 output: &Path,
2080) -> Result<mj_core::state::CheckpointMetadata> {
2081 let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
2082 .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
2083 let native_artifacts = previous
2086 .manifest
2087 .payloads
2088 .iter()
2089 .filter_map(|descriptor| match &descriptor.role {
2090 mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
2091 Some((relative_path, descriptor))
2092 }
2093 _ => None,
2094 })
2095 .map(|(relative_path, descriptor)| {
2096 Ok(mj_checkpoint::archive::NativeArtifact {
2097 relative_path: relative_path.clone(),
2098 data: previous.payload(descriptor)?.to_vec(),
2099 mode: descriptor.mode,
2100 })
2101 })
2102 .collect::<Result<Vec<_>>>()?;
2103 let canonical_session = previous.canonical_session()?;
2104 let event_frontier = canonical_session.event_frontier;
2105 let written = mj_checkpoint::archive::write_archive_atomic(
2106 output,
2107 &mj_checkpoint::archive::ArchiveInput {
2108 session: previous.manifest.session.clone(),
2109 target: previous.manifest.target.clone(),
2112 bundle: mj_checkpoint::archive::BundleManifest {
2113 id: previous.manifest.bundle.id.clone(),
2114 primary_repository: snapshot.metadata.id.clone(),
2118 },
2119 canonical_session,
2120 native_artifacts,
2121 repositories: vec![snapshot],
2122 },
2123 )
2124 .with_context(|| format!("write the conversion archive {}", output.display()))?;
2125 Ok(mj_core::state::CheckpointMetadata {
2126 archive_path: output.to_path_buf(),
2127 sha256: written.archive_sha256,
2128 created_at: now(),
2129 event_frontier,
2130 })
2131}
2132
2133fn restored_native_session_accepted(
2144 archived: &str,
2145 opened: &str,
2146 replaced_unused: Option<&str>,
2147) -> bool {
2148 opened == archived || replaced_unused == Some(archived)
2149}
2150
2151fn restore_projection_after_failed_resume(
2162 session_id: &str,
2163 canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
2164 discard_queued_prompts: bool,
2165) {
2166 let stored_frontier = crate::database::materialized_event_frontier(session_id)
2167 .unwrap_or_else(|error| {
2168 tracing::warn!(
2169 session_id,
2170 error = format!("{error:#}"),
2171 "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
2172 );
2173 None
2174 });
2175 if projection_rebuild_required(
2176 stored_frontier
2177 .as_ref()
2178 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
2179 canonical_session.event_frontier,
2180 &canonical_session.event_frontier_digest,
2181 ) {
2182 match materialized_session_from_canonical(session_id, canonical_session) {
2183 Ok(previous_projection) => {
2184 if let Err(restore_error) =
2185 crate::database::save_materialized_session(&previous_projection)
2186 {
2187 tracing::error!(
2188 session_id,
2189 error = format!("{restore_error:#}"),
2190 "could not restore the durable projection after a failed resume"
2191 );
2192 }
2193 }
2194 Err(restore_error) => tracing::error!(
2195 session_id,
2196 error = format!("{restore_error:#}"),
2197 "could not rebuild the durable projection after a failed resume"
2198 ),
2199 }
2200 } else if discard_queued_prompts
2201 && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
2202 session_id,
2203 &mj_transcript::projection::materialized_queued_prompts_from_canonical(
2204 &canonical_session.queued_prompts,
2205 ),
2206 )
2207 {
2208 tracing::error!(
2209 session_id,
2210 error = format!("{restore_error:#}"),
2211 "could not restore queued prompts after a failed resume"
2212 );
2213 }
2214}
2215
2216fn projection_rebuild_required(
2224 stored: Option<(u64, &str)>,
2225 archive_frontier: u64,
2226 archive_frontier_digest: &str,
2227) -> bool {
2228 stored != Some((archive_frontier, archive_frontier_digest))
2229}
2230
2231fn restore_archive_path(
2232 backend: &targets::TargetLocator,
2233 verified_archive: &Path,
2234 remote_archive: &Path,
2235) -> PathBuf {
2236 if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2237 verified_archive.to_path_buf()
2238 } else {
2239 remote_archive.to_path_buf()
2240 }
2241}
2242
2243fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2244 !matches!(backend, targets::TargetLocator::LocalBare { .. })
2245}
2246
2247struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2252 inner: &'a E,
2253 cancellation: CancellationToken,
2254}
2255
2256impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2257 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2258 if self.cancellation.is_cancelled() {
2259 bail!("operation cancelled while provisioning destination");
2260 }
2261 self.inner.execute(command)
2262 }
2263
2264 fn cancellation_requested(&self) -> bool {
2265 self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2266 }
2267
2268 fn stage_started(&self, stage: ProvisionStage) {
2269 self.inner.stage_started(stage);
2270 }
2271
2272 fn stage_finished(&self, stage: ProvisionStage) {
2273 self.inner.stage_finished(stage);
2274 }
2275
2276 fn notify_notice(&self, notice: &str) {
2277 self.inner.notify_notice(notice);
2278 }
2279
2280 fn execute_with_stdin(
2281 &self,
2282 command: &CommandSpec,
2283 input: &mut (dyn std::io::Read + Send),
2284 ) -> Result<CommandOutput> {
2285 if self.cancellation.is_cancelled() {
2286 bail!("operation cancelled while provisioning destination");
2287 }
2288 self.inner.execute_with_stdin(command, input)
2289 }
2290}
2291
2292fn provision_with_cross_harness_handoff(
2293 controller: &mut Controller,
2294 session_id: &str,
2295 executor: &(impl CommandExecutor + Sync),
2296 github_token: Option<&str>,
2297 config: &Config,
2298 snapshot: &CanonicalSessionSnapshot,
2299 context_bytes: usize,
2300) -> Result<String> {
2301 let (_provision, handoff) = execute_joined_cross_harness_work(
2302 "cross-harness provisioning",
2303 move |cancellation| {
2304 let provision_executor = CrossHarnessProvisionExecutor {
2305 inner: executor,
2306 cancellation,
2307 };
2308 futures::executor::block_on(controller.provision_session_with_failure_disposition(
2309 session_id,
2310 &provision_executor,
2311 github_token,
2312 ProvisioningFailureDisposition::Preserve,
2313 ))
2314 },
2315 "cross-harness handoff",
2316 move |cancellation| {
2317 let runtime = tokio::runtime::Builder::new_current_thread()
2318 .enable_all()
2319 .build()
2320 .context("create cross-harness handoff runtime")?;
2321 runtime.block_on(utility_handoff_while_cancellable(
2322 session_id,
2323 config,
2324 snapshot,
2325 context_bytes,
2326 executor,
2327 cancellation,
2328 ))
2329 },
2330 )?;
2331 ensure!(
2332 !executor.cancellation_requested(),
2333 "operation cancelled while provisioning destination"
2334 );
2335 Ok(handoff)
2336}
2337
2338fn execute_joined_cross_harness_work<A: Send, B: Send>(
2342 first_name: &'static str,
2343 first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2344 second_name: &'static str,
2345 second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2346) -> Result<(A, B)> {
2347 let cancellation = CancellationToken::new();
2348 std::thread::scope(|scope| {
2349 let first_cancel = cancellation.clone();
2350 let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2351 let second_cancel = cancellation.clone();
2352 let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2353 let mut first_result = None;
2354 let mut second_result = None;
2355
2356 while first_result.is_none() || second_result.is_none() {
2357 if first_result.is_none()
2358 && first_handle
2359 .as_ref()
2360 .is_some_and(|handle| handle.is_finished())
2361 {
2362 let handle = first_handle.take().expect("first lane handle present");
2363 first_result = Some(match handle.join() {
2364 Ok(result) => result,
2365 Err(panic) => {
2366 cancellation.cancel();
2367 Err(anyhow::anyhow!(
2368 "{first_name} thread panicked: {}",
2369 targets::command_thread_panic_message(panic.as_ref())
2370 ))
2371 }
2372 });
2373 if first_result.as_ref().is_some_and(Result::is_err) {
2374 cancellation.cancel();
2375 }
2376 }
2377 if second_result.is_none()
2378 && second_handle
2379 .as_ref()
2380 .is_some_and(|handle| handle.is_finished())
2381 {
2382 let handle = second_handle.take().expect("second lane handle present");
2383 second_result = Some(match handle.join() {
2384 Ok(result) => result,
2385 Err(panic) => {
2386 cancellation.cancel();
2387 Err(anyhow::anyhow!(
2388 "{second_name} thread panicked: {}",
2389 targets::command_thread_panic_message(panic.as_ref())
2390 ))
2391 }
2392 });
2393 if second_result.as_ref().is_some_and(Result::is_err) {
2394 cancellation.cancel();
2395 }
2396 }
2397 if first_result.is_none() || second_result.is_none() {
2398 std::thread::sleep(Duration::from_millis(10));
2399 }
2400 }
2401
2402 match (
2403 first_result.expect("first lane result received after joined handle"),
2404 second_result.expect("second lane result received after joined handle"),
2405 ) {
2406 (Err(first), Err(second)) => {
2407 Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2408 }
2409 (Err(error), Ok(_)) => Err(error),
2410 (Ok(_), Err(error)) => Err(error),
2411 (Ok(first), Ok(second)) => Ok((first, second)),
2412 }
2413 })
2414}
2415
2416fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2420 profile_kind == archived_kind
2421}
2422
2423async fn utility_handoff_while_cancellable(
2427 session_id: &str,
2428 config: &Config,
2429 snapshot: &CanonicalSessionSnapshot,
2430 context_bytes: usize,
2431 executor: &impl CommandExecutor,
2432 cancellation: CancellationToken,
2433) -> Result<String> {
2434 let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2435 if executor.cancellation_requested() {
2436 bail!("operation cancelled while compacting the cross-harness handoff");
2437 }
2438 let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2439 let cancel = cancellation.child_token();
2440 let operation =
2441 crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2442 tokio::pin!(operation);
2443 loop {
2444 tokio::select! {
2445 context = &mut operation => return context,
2446 _ = cancellation.cancelled() => {
2447 cancel.cancel();
2448 bail!("operation cancelled while compacting the cross-harness handoff");
2449 }
2450 _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2451 if executor.cancellation_requested() {
2452 cancel.cancel();
2453 bail!("operation cancelled while compacting the cross-harness handoff");
2454 }
2455 }
2456 }
2457 }
2458}
2459
2460mod in_place;
2461
2462#[cfg(test)]
2463mod tests;