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, validate_resource_allocation};
30use super::checkpoint::upload_checkpoint_spec;
31use super::github_app::{GithubAppTokenProvider, github_repositories};
32use super::provisioning::{
33 ProvisioningFailureDisposition, StagedExecutor, execute_concurrent_lanes,
34 install_attached_resources,
35};
36use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
37use super::worker_binary::{bridge_readiness_stage, start_worker_durably, worker_probe_diagnosis};
38use super::worktree::{
39 PrimaryCheckoutRequirement, ResumeConversion, ResumePlan, apply_raw_to_workspace,
40 apply_workspace_to_raw, cleanup_managed_worktree, create_managed_worktree,
41 managed_worktree_checkout_exists, managed_worktree_target, plan_raw_to_workspace_for_session,
42 preserve_retained_managed_worktree_branch, raw_checkout_divergence_notice,
43 raw_checkout_position_with_checkout, raw_checkout_snapshot,
44 raw_conversion_preview_with_checkout, restore_managed_worktree,
45 resume_compatibility_with_checkout, retire_managed_worktree,
46};
47use super::{
48 Controller, SessionResumeOptions, execute_checked, now, selected_host_container_size,
49 target_profile_home,
50};
51
52#[derive(Debug, Clone, PartialEq, Eq)]
53pub struct ResumeRepositorySourceMismatch {
54 pub session_id: String,
55 pub bundle_id: String,
56 pub repository_id: String,
57 pub missing_commit: String,
58 pub archived_origin: String,
59 pub configured_origin: String,
60}
61
62pub use mj_core::state::ResumeRepositorySourceReceipt;
63
64#[derive(Debug, Clone, PartialEq, Eq)]
65pub enum ResumeRepositorySourcePreflight {
66 Ready(ResumeRepositorySourceReceipt),
67 RepositoryMoved(ResumeRepositorySourceMismatch),
68 ConvertingRawCheckout {
73 receipt: ResumeRepositorySourceReceipt,
74 preview: Box<mj_core::state::RawConversionPreview>,
75 },
76}
77
78struct ResumeRepositoryBundles {
79 checkpoint_sha256: String,
80 repositories: Vec<CheckpointRepositoryBundle>,
81}
82
83struct ResumePhaseTimer<'a> {
87 session_id: &'a str,
88 phase: &'static str,
89 started: Instant,
90}
91
92impl<'a> ResumePhaseTimer<'a> {
93 fn new(session_id: &'a str, phase: &'static str) -> Self {
94 Self {
95 session_id,
96 phase,
97 started: Instant::now(),
98 }
99 }
100}
101
102impl Drop for ResumePhaseTimer<'_> {
103 fn drop(&mut self) {
104 tracing::debug!(
105 session_id = self.session_id,
106 phase = self.phase,
107 elapsed_ms = self.started.elapsed().as_millis(),
108 "resume phase completed"
109 );
110 }
111}
112
113impl Controller {
114 pub(super) fn validate_muse_resume_destination(
116 &self,
117 source: &SessionRecord,
118 destination_harness: HarnessKind,
119 _target_id: &str,
120 ) -> Result<()> {
121 if destination_harness != HarnessKind::Muse {
122 return Ok(());
123 }
124 let checkout = self.state.checkout(&source.id)?;
125 ensure!(
126 checkout.project_directory().is_some()
127 || source
128 .project_bundle(&self.config)
129 .is_none_or(|bundle| bundle.repositories.len() == 1),
130 "Muse Code ACP supports one workspace root; use a single-repository bundle"
131 );
132 Ok(())
133 }
134 pub async fn preflight_resume_repository_sources(
138 &self,
139 session_id: &str,
140 target_id: &str,
141 executor: &(impl CommandExecutor + Sync),
142 ) -> Result<ResumeRepositorySourcePreflight> {
143 let github_token = self
144 .github_token_for_repository_preflight(session_id)
145 .await?;
146 self.preflight_repository_sources(
147 session_id,
148 target_id,
149 true,
150 github_token.as_deref(),
151 executor,
152 )
153 }
154
155 pub fn preflight_resume_repository_sources_with_token(
158 &self,
159 session_id: &str,
160 target_id: &str,
161 github_token: Option<&str>,
162 executor: &(impl CommandExecutor + Sync),
163 ) -> Result<ResumeRepositorySourcePreflight> {
164 self.preflight_repository_sources(session_id, target_id, true, github_token, executor)
165 }
166
167 fn preflight_repository_sources(
172 &self,
173 session_id: &str,
174 target_id: &str,
175 describe_conversion: bool,
176 github_token: Option<&str>,
177 executor: &(impl CommandExecutor + Sync),
178 ) -> Result<ResumeRepositorySourcePreflight> {
179 let session = self
180 .state
181 .sessions
182 .get(session_id)
183 .with_context(|| format!("unknown session {session_id}"))?;
184 let checkpoint = session
185 .checkpoint
186 .as_ref()
187 .context("session has no checkpoint")?;
188 let checkout = self.state.checkout(session_id)?;
189 let plan = resume_compatibility_with_checkout(session, &checkout, &self.config, target_id)
190 .map_err(|reason| anyhow::anyhow!(reason))?;
191 if checkout.project_directory().is_some() {
192 debug_assert!(matches!(
193 plan,
194 ResumePlan::InPlace | ResumePlan::RawToWorkspace
195 ));
196 let receipt = ResumeRepositorySourceReceipt {
200 session_id: session_id.to_owned(),
201 bundle_id: session.bundle_id.clone(),
202 checkpoint_sha256: checkpoint.sha256.clone(),
203 repositories: Vec::new(),
204 };
205 if describe_conversion && plan == ResumePlan::RawToWorkspace {
210 let preview =
211 raw_conversion_preview_for(session, &checkout, &self.config, executor)?;
212 return Ok(ResumeRepositorySourcePreflight::ConvertingRawCheckout {
213 receipt,
214 preview: Box::new(preview),
215 });
216 }
217 return Ok(ResumeRepositorySourcePreflight::Ready(receipt));
218 }
219 let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
220 self.preflight_verified_repository_sources(
221 session_id,
222 ResumeRepositoryBundles {
223 checkpoint_sha256: checkpoint.sha256.clone(),
224 repositories,
225 },
226 None,
227 plan != ResumePlan::WorkspaceToRaw,
228 github_token,
229 executor,
230 )
231 }
232
233 fn preflight_verified_repository_sources(
234 &self,
235 session_id: &str,
236 verified: ResumeRepositoryBundles,
237 skip_repository_id: Option<&str>,
238 use_archived_network_sources: bool,
239 github_token: Option<&str>,
240 executor: &(impl CommandExecutor + Sync),
241 ) -> Result<ResumeRepositorySourcePreflight> {
242 let session = self
243 .state
244 .sessions
245 .get(session_id)
246 .with_context(|| format!("unknown session {session_id}"))?;
247 ensure!(
248 verified.repositories.iter().all(|repository| {
249 !repository.metadata.origin.starts_with("mj-local:")
250 && !repository.metadata.origin.starts_with("ext::")
251 }),
252 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
253 );
254 if use_archived_network_sources
257 && verified
258 .repositories
259 .iter()
260 .all(|repository| repository.metadata.remote_workspace)
261 {
262 return Ok(ResumeRepositorySourcePreflight::Ready(
263 ResumeRepositorySourceReceipt {
264 session_id: session_id.to_owned(),
265 bundle_id: session.bundle_id.clone(),
266 checkpoint_sha256: verified.checkpoint_sha256,
267 repositories: Vec::new(),
268 },
269 ));
270 }
271 if verified.repositories.is_empty() {
272 return Ok(ResumeRepositorySourcePreflight::Ready(
273 ResumeRepositorySourceReceipt {
274 session_id: session_id.to_owned(),
275 bundle_id: session.bundle_id.clone(),
276 checkpoint_sha256: verified.checkpoint_sha256,
277 repositories: Vec::new(),
278 },
279 ));
280 }
281 let bundle = session
282 .project_bundle(&self.config)
283 .with_context(|| format!("session bundle {:?} is missing", session.bundle_id))?;
284 let configured = verified
285 .repositories
286 .iter()
287 .map(|archived| {
288 bundle
289 .repositories
290 .iter()
291 .find(|repository| repository.id == archived.metadata.id)
292 .cloned()
293 .with_context(|| {
294 format!(
295 "session bundle {:?} no longer contains repository {:?}",
296 session.bundle_id, archived.metadata.id
297 )
298 })
299 })
300 .collect::<Result<Vec<_>>>()?;
301 let outcomes = verified
302 .repositories
303 .par_iter()
304 .zip(configured.par_iter())
305 .map(move |(archived, configured)| {
306 if skip_repository_id == Some(configured.id.as_str()) {
307 return Ok(None);
308 }
309 checkpoint_source_missing_commit(
310 configured,
311 session
312 .project
313 .as_ref()
314 .and_then(|project| project.network_sources.get(&configured.id)),
315 archived,
316 executor,
317 github_token,
318 )
319 .map(|missing_commit| {
320 missing_commit.map(|missing_commit| ResumeRepositorySourceMismatch {
321 session_id: session_id.to_owned(),
322 bundle_id: session.bundle_id.clone(),
323 repository_id: configured.id.clone(),
324 missing_commit,
325 archived_origin: archived.metadata.origin.clone(),
326 configured_origin: configured.source_label(),
327 })
328 })
329 })
330 .collect::<Vec<Result<Option<ResumeRepositorySourceMismatch>>>>();
331 for outcome in outcomes {
332 if let Some(mismatch) = outcome? {
333 return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(mismatch));
334 }
335 }
336 Ok(ResumeRepositorySourcePreflight::Ready(
337 ResumeRepositorySourceReceipt {
338 session_id: session_id.to_owned(),
339 bundle_id: session.bundle_id.clone(),
340 checkpoint_sha256: verified.checkpoint_sha256,
341 repositories: configured,
342 },
343 ))
344 }
345
346 fn repository_source_receipt_is_current(
347 &self,
348 session_id: &str,
349 receipt: &ResumeRepositorySourceReceipt,
350 ) -> bool {
351 let Some(session) = self.state.sessions.get(session_id) else {
352 return false;
353 };
354 if receipt.session_id != session_id
355 || receipt.bundle_id != session.bundle_id
356 || session
357 .checkpoint
358 .as_ref()
359 .map(|checkpoint| &checkpoint.sha256)
360 != Some(&receipt.checkpoint_sha256)
361 {
362 return false;
363 }
364 if receipt.repositories.is_empty() {
365 return true;
366 }
367 let Some(bundle) = session.project_bundle(&self.config) else {
368 return false;
369 };
370 receipt.repositories.iter().all(|expected| {
371 bundle
372 .repositories
373 .iter()
374 .any(|configured| configured == expected)
375 })
376 }
377
378 pub async fn replace_resume_repository_origin(
382 &mut self,
383 session_id: &str,
384 repository_id: &str,
385 replacement: &str,
386 executor: &(impl CommandExecutor + Sync),
387 ) -> Result<ResumeRepositorySourcePreflight> {
388 self.replace_resume_repository_origin_with_token(
389 session_id,
390 repository_id,
391 replacement,
392 executor,
393 )
394 .await
395 }
396
397 async fn replace_resume_repository_origin_with_token(
398 &mut self,
399 session_id: &str,
400 repository_id: &str,
401 replacement: &str,
402 executor: &(impl CommandExecutor + Sync),
403 ) -> Result<ResumeRepositorySourcePreflight> {
404 let session = self
405 .state
406 .sessions
407 .get(session_id)
408 .with_context(|| format!("unknown session {session_id}"))?;
409 let bundle_id = session.bundle_id.clone();
410 let checkpoint = session
411 .checkpoint
412 .as_ref()
413 .context("session has no checkpoint")?;
414 let replacement = replacement_repository_source(repository_id, replacement)?;
415 let replacement_network = session
416 .project
417 .as_ref()
418 .map(|_| mj_core::remote_git::resolve_repository(&replacement, executor))
419 .transpose()?;
420 let previous = session.clone();
421 let mut project = session.project.clone();
422 if let Some(project) = &mut project {
423 let repository = project
424 .bundle
425 .repositories
426 .iter_mut()
427 .find(|repository| repository.id == repository_id)
428 .context("accepted project no longer contains the repaired repository")?;
429 repository.github = replacement.github.clone();
430 repository.local = replacement.local.clone();
431 project.network_sources.insert(
432 repository_id.to_owned(),
433 replacement_network
434 .clone()
435 .expect("accepted project source resolved above"),
436 );
437 }
438 let github_token = if let Some(app) = self.config.github.app.as_ref() {
439 if let Some(project) = &project {
440 let bundle = project.bundle.clone();
441 let network_sources = project.network_sources.clone();
442 let repositories = tokio::task::spawn_blocking(move || {
443 github_repositories(&bundle, Some(&network_sources), &ProcessExecutor)
444 })
445 .await
446 .context("replacement GitHub repository source task failed")??;
447 let provider = GithubAppTokenProvider::shared(app)?;
448 let scope = provider
449 .token_for_owner_repo_pairs(&bundle_id, &repositories)
450 .await
451 .map_err(super::GithubBundleSelectionError::into_anyhow)?;
452 match scope {
453 Some(scope) => Some(
454 provider
455 .token_for_installation(
456 scope.installation_id,
457 &scope.repositories,
458 app.session_permissions.as_ref(),
459 )
460 .await?,
461 ),
462 None => None,
463 }
464 } else {
465 None
466 }
467 } else {
468 self.github_token_for_session(session_id).await?
469 };
470 let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
471 let verified = ResumeRepositoryBundles {
472 checkpoint_sha256: checkpoint.sha256.clone(),
473 repositories,
474 };
475 let archived = verified
476 .repositories
477 .iter()
478 .find(|repository| repository.metadata.id == repository_id)
479 .with_context(|| format!("checkpoint does not contain repository {repository_id:?}"))?;
480 if let Some(missing_commit) = checkpoint_source_missing_commit(
481 &replacement,
482 replacement_network.as_ref(),
483 archived,
484 executor,
485 github_token.as_deref(),
486 )? {
487 return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(
488 ResumeRepositorySourceMismatch {
489 session_id: session_id.to_owned(),
490 bundle_id,
491 repository_id: repository_id.to_owned(),
492 missing_commit,
493 archived_origin: archived.metadata.origin.clone(),
494 configured_origin: replacement.source_label(),
495 },
496 ));
497 }
498 let saved = crate::database::saved_project(&bundle_id)?;
501 let config_repository_id = saved
502 .as_ref()
503 .and_then(|(_, saved)| {
504 let accepted = previous.project.as_ref()?;
505 let original = accepted
506 .bundle
507 .repositories
508 .iter()
509 .find(|repository| repository.id == repository_id)?;
510 saved
511 .bundle
512 .repositories
513 .iter()
514 .find(|repository| {
515 repository.destination == original.destination
516 && saved.identities.get(&repository.id)
517 == accepted.identities.get(repository_id)
518 })
519 .map(|repository| repository.id.clone())
520 })
521 .unwrap_or_else(|| repository_id.to_owned());
522 let config_bundle_id = saved
523 .as_ref()
524 .map_or(bundle_id.as_str(), |(id, _)| id.as_str());
525 let (config, ()) = Config::update(|config| {
526 if project.is_some() && !config.bundles.contains_key(config_bundle_id) {
527 return Ok(());
528 }
529 let bundle = config
530 .bundles
531 .get_mut(config_bundle_id)
532 .with_context(|| format!("session bundle {bundle_id:?} is missing"))?;
533 let repository = bundle
534 .repositories
535 .iter_mut()
536 .find(|repository| repository.id == config_repository_id)
537 .with_context(|| {
538 format!(
539 "session bundle {:?} no longer contains repository {repository_id:?}",
540 bundle_id
541 )
542 })?;
543 repository.github = replacement.github.clone();
544 repository.local = replacement.local.clone();
545 Ok(())
546 })?;
547 self.config = config;
548 if let Some(project) = project {
549 self.state
550 .sessions
551 .get_mut(session_id)
552 .expect("session checked above")
553 .project = Some(project);
554 self.persist_session_transition_or_restore(
555 session_id,
556 &previous,
557 "save repaired project source",
558 )?;
559 }
560 self.preflight_verified_repository_sources(
561 session_id,
562 verified,
563 Some(repository_id),
564 false,
565 github_token.as_deref(),
566 executor,
567 )
568 }
569}
570
571pub(super) enum WorkerRootReset {
579 FreshTarget,
582 InPlace {
586 previous_profile_root: String,
588 },
589}
590
591pub(super) struct RestoreIntoTarget<'a> {
597 pub profile: &'a mj_core::config::HarnessProfile,
599 pub archive: &'a VerifiedResumeArchive,
600 pub restored_archive: &'a Path,
603 pub resumed_project_directory: Option<PathBuf>,
604 pub resumed_container_workspace: Option<PathBuf>,
605 pub restore_repositories: bool,
606 pub native_continuity: bool,
607 pub discard_queued_prompts: bool,
608 pub replay_queue: bool,
610 pub utility_handoff: Option<String>,
613 pub projection_build: Option<tokio::task::JoinHandle<Result<MaterializedSession>>>,
614 pub resume_notices: Vec<String>,
616 pub include_in_place_subagents: bool,
619 pub install_attached_resources: bool,
620 pub worker_root_reset: WorkerRootReset,
621 pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
623}
624
625impl Controller {
626 pub(super) async fn restore_into_target(
629 &mut self,
630 session_id: &str,
631 restore: RestoreIntoTarget<'_>,
632 executor: &(impl CommandExecutor + Sync),
633 ) -> Result<MaterializedSession> {
634 crate::worker_lifecycle::run(session_id, "restore into target", executor, async {
635 let RestoreIntoTarget {
636 profile,
637 archive,
638 restored_archive,
639 resumed_project_directory,
640 resumed_container_workspace,
641 restore_repositories,
642 native_continuity,
643 discard_queued_prompts,
644 replay_queue,
645 utility_handoff,
646 projection_build,
647 mut resume_notices,
648 include_in_place_subagents,
649 install_attached_resources: should_install_attached_resources,
650 worker_root_reset,
651 retire_after_ready,
652 } = restore;
653 let archive_manifest = &archive.manifest;
654 let canonical_session = &archive.canonical_session;
655 let stopped_subagents = crate::database::load_stopped_subagents(session_id)
660 .unwrap_or_else(|error| {
661 tracing::warn!(
662 session_id,
663 error = format!("{error:#}"),
664 "could not read the sub-agents this session's suspend stopped"
665 );
666 Vec::new()
667 });
668 let stopped_subagents_context =
669 mj_core::subagent::stopped_subagents_prompt_context(&stopped_subagents);
670 let reopen_move_subagents =
671 matches!(&worker_root_reset, WorkerRootReset::InPlace { .. })
672 && self.state.sessions[session_id]
673 .subagents
674 .clone()
675 .unwrap_or_default()
676 .for_launch(profile.kind, false)
677 .parent_role()
678 .is_some();
679 resume_notices.extend(mj_core::subagent::stopped_subagents_notice(
680 &stopped_subagents,
681 ));
682 let mut prompt_contexts = Vec::new();
683 prompt_contexts.extend(stopped_subagents_context);
684 let (backend, worker_root) = self.worker_placement(session_id)?;
685 let harness_home = target_profile_home(&backend, session_id, profile);
686 let workspace_root = if let Some(project_directory) = &resumed_project_directory {
687 project_directory
688 .parent()
689 .context("bare project directory has no parent")?
690 .to_string_lossy()
691 .into_owned()
692 } else {
693 super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
694 };
695 let target_path = |path: &str| match &backend {
696 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
697 if !path.starts_with('/') =>
698 {
699 PathBuf::from(format!("~/{path}"))
700 }
701 _ => PathBuf::from(path),
702 };
703 let remote_archive = format!("{worker_root}/restore.hel.zip");
704 let remote_spec = format!("{worker_root}/restore-spec.json");
705 use mj_checkpoint::checkpoint::QueueRestorePolicy;
706 let move_admission =
707 crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
708 operation.phase == mj_core::state::MovePhase::ResumingDestination
709 && !operation.queue_admission_started
710 && operation.queue == mj_core::state::ResumeQueueDisposition::Start
711 });
712 let queue_policy = if move_admission || (!native_continuity && replay_queue) {
713 QueueRestorePolicy::Defer
714 } else if discard_queued_prompts {
715 QueueRestorePolicy::Discard
716 } else {
717 QueueRestorePolicy::Restore
718 };
719 let restore = CheckpointRestoreSpec {
720 archive_path: restore_archive_path(
721 &backend,
722 restored_archive,
723 &target_path(&remote_archive),
724 ),
725 workspace_root: target_path(&workspace_root),
726 relay_root: target_path(&worker_root),
727 harness_home: target_path(&harness_home),
728 restore_repositories,
734 restore_native: native_continuity,
735 primary_repository_root: resumed_project_directory
744 .as_ref()
745 .map(|directory| target_path(&directory.to_string_lossy())),
746 queue_policy,
747 };
748 {
752 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
753 match &worker_root_reset {
754 WorkerRootReset::FreshTarget => {
759 if let Some(command) =
760 targets::clear_relay_state_plan(&backend, session_id)?
761 {
762 execute_checked(syncing, command)?;
763 }
764 execute_checked(
767 syncing,
768 targets::command_on_locator(
769 &backend,
770 session_id,
771 vec!["mkdir".into(), "-p".into(), worker_root.clone()],
772 "create the session worker root",
773 )?,
774 )?;
775 }
776 WorkerRootReset::InPlace {
781 previous_profile_root,
782 } => {
783 execute_checked(
784 syncing,
785 targets::in_place_worker_reset_plan(
786 &backend,
787 session_id,
788 previous_profile_root,
789 )?,
790 )?;
791 }
792 }
793 }
794 let staging = tempfile::tempdir().context("create restore staging")?;
795 let local_spec = staging.path().join("restore-spec.json");
796 std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
797 let controller = &*self;
801 let backend_ref = &backend;
802 let worker_root_ref = worker_root.as_str();
803 let local_spec_ref = local_spec.as_path();
804 crate::target_storage::ensure_room_for(
807 &backend,
808 [
809 worker_root.as_str(),
810 workspace_root.as_str(),
811 harness_home.as_str(),
812 ],
813 || {
814 Ok(std::fs::metadata(restored_archive)
815 .with_context(|| format!("measure {}", restored_archive.display()))?
816 .len())
817 },
818 "restore the checkpoint",
819 )?;
820 execute_concurrent_lanes(
821 || {
822 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
823 controller.prepare_worker_files(
824 session_id,
825 backend_ref,
826 worker_root_ref,
827 syncing,
828 )?;
829 super::provisioning::install_inherited_git_settings(
830 syncing,
831 backend_ref,
832 session_id,
833 )?;
834 Ok(())
835 },
836 || {
837 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
838 if should_upload_restore_archive(&backend) {
839 upload_checkpoint_spec(
840 restoring,
841 backend_ref,
842 session_id,
843 restored_archive,
844 &remote_archive,
845 )?;
846 }
847 upload_checkpoint_spec(
848 restoring,
849 backend_ref,
850 session_id,
851 local_spec_ref,
852 &remote_spec,
853 )
854 },
855 )?;
856 {
857 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
858 execute_checked(
859 restoring,
860 restore_command(&backend, session_id, &remote_spec)?,
861 )?;
862 }
863 if should_install_attached_resources {
864 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
865 install_attached_resources(
866 &self.state,
867 session_id,
868 &backend,
869 &worker_root,
870 syncing,
871 )?;
872 }
873 match projection_build {
874 Some(build) => {
875 let mut restored_projection = build
876 .await
877 .context("rebuild the restored projection")?
878 .context("rebuild the restored projection")?;
879 if discard_queued_prompts {
880 restored_projection.queued_prompts.clear();
881 }
882 crate::database::save_materialized_session(&restored_projection)?;
883 }
884 None if discard_queued_prompts => {
887 crate::database::replace_materialized_queued_prompts(session_id, &[])?;
888 }
889 None => {}
890 }
891 let readiness_stage = bridge_readiness_stage(profile);
892 let spec = self.reconnect_command(session_id)?;
893 let readiness = async {
894 let mut relay = {
895 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
896 start_worker_durably(
897 &crate::worker_lifecycle::require(session_id)?,
898 self.state.sessions[session_id]
899 .target
900 .as_ref()
901 .context("worker start has no durable target")?,
902 executor,
903 &backend,
904 &worker_root,
905 )?;
906 connect_started_worker(&spec, session_id, executor, &backend, &worker_root)
907 .await?
908 };
909 if reopen_move_subagents {
910 relay
911 .set_subagent_admission(true)
912 .await
913 .context("reopen sub-agent requests on the replacement worker")?;
914 }
915 if include_in_place_subagents {
916 let session_id = session_id.to_owned();
917 let (context, _) = tokio::task::spawn_blocking(move || {
918 load_in_place_subagent_prompt_data(&session_id)
919 })
920 .await
921 .context("join the retained sub-agent roster read")??;
922 prompt_contexts.extend(context);
923 }
924 if !prompt_contexts.is_empty() {
928 relay
929 .install_prompt_context(prompt_contexts.join("\n\n"))
930 .await
931 .context("tell the resumed session about its sub-agents")?;
932 }
933 let native_session_id = wait_for_native_session_in_stage(
934 &mut relay,
935 executor,
936 readiness_stage,
937 profile.kind,
938 )
939 .await?;
940 let owner = crate::worker_lifecycle::require(session_id)?;
941 crate::database::finish_worker_restart(session_id, owner.operation_id())?;
942 Ok::<_, anyhow::Error>((relay, native_session_id))
943 }
944 .await;
945 let (mut relay, native_session_id) = readiness
946 .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
947 if native_continuity {
948 if !restored_native_session_accepted(
949 &archive_manifest.session.native_session_id,
950 &native_session_id,
951 relay
952 .operational()
953 .replaced_unused_native_session_id
954 .as_deref(),
955 ) {
956 bail!(
957 "ACP loaded native session {native_session_id}, expected {}",
958 archive_manifest.session.native_session_id
959 );
960 }
961 } else {
962 relay
963 .install_prompt_context(
964 utility_handoff
965 .clone()
966 .context("a resume into a fresh native session has no handoff")?,
967 )
968 .await?;
969 if replay_queue {
970 for prompt in &canonical_session.queued_prompts {
971 let command = match &prompt.kind {
975 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
976 prompt: prompt
977 .content
978 .iter()
979 .cloned()
980 .map(serde_json::from_value)
981 .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
982 },
983 CanonicalQueuedCommandKind::SetConfig { key, value } => {
984 RelayCommand::SetConfig {
985 key: key.clone(),
986 value: value.clone(),
987 }
988 }
989 };
990 relay.submit(prompt.command_id.clone(), command).await?;
991 }
992 }
993 }
994 if let Some(worktree) = retire_after_ready
998 && let Err(error) = retire_managed_worktree(executor, worktree)
999 {
1000 tracing::warn!(
1001 session_id,
1002 worktree = %worktree.worktree_root.display(),
1003 error = format!("{error:#}"),
1004 "could not retire the old managed worktree after resume"
1005 );
1006 resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
1007 }
1008 if include_in_place_subagents {
1009 let session_id = session_id.to_owned();
1010 let (_, notice) = tokio::task::spawn_blocking(move || {
1011 load_in_place_subagent_prompt_data(&session_id)
1012 })
1013 .await
1014 .context("join the retained sub-agent notice roster read")??;
1015 resume_notices.extend(notice);
1016 }
1017 for notice in &resume_notices {
1018 let submitted = async {
1019 let command_id = new_command_id("resume-notice")?;
1020 relay
1021 .submit(
1022 command_id,
1023 RelayCommand::RecordNotice {
1024 text: notice.clone(),
1025 },
1026 )
1027 .await
1028 }
1029 .await;
1030 if let Err(error) = submitted {
1033 tracing::warn!(
1034 session_id,
1035 error = format!("{error:#}"),
1036 "could not record a resume notice in the conversation"
1037 );
1038 }
1039 }
1040 self.mark_worker_connected(session_id, Some(native_session_id))?;
1041 let materialized = relay.sync().await?.materialized;
1042 if !stopped_subagents.is_empty() {
1043 let delivered = stopped_subagents
1044 .iter()
1045 .map(|child| child.child_session_id.clone())
1046 .collect::<Vec<_>>();
1047 if let Err(error) = crate::database::clear_stopped_subagents(session_id, &delivered)
1050 {
1051 tracing::warn!(
1052 session_id,
1053 error = format!("{error:#}"),
1054 "could not clear the stopped sub-agents after telling the resumed session"
1055 );
1056 }
1057 }
1058 Ok(materialized)
1059 })
1060 .await
1061 }
1062}
1063
1064pub(super) struct VerifiedResumeArchive {
1071 pub archive_path: PathBuf,
1075 pub manifest: mj_checkpoint::archive::ArchiveManifest,
1076 pub canonical_session: Arc<CanonicalSessionSnapshot>,
1077}
1078
1079fn load_in_place_subagent_prompt_data(
1080 session_id: &str,
1081) -> Result<(Option<String>, Option<String>)> {
1082 let state = crate::database::load_state()
1083 .context("load current retained sub-agents for the in-place resume")?;
1084 let children = super::move_session::move_children(&state, session_id);
1085 Ok((
1086 mj_core::subagent::in_place_subagents_prompt_context(&children),
1087 mj_core::subagent::in_place_subagents_notice(&children),
1088 ))
1089}
1090
1091pub(super) fn verify_resume_checkpoint(
1097 session_id: &str,
1098 checkpoint: &mj_core::state::CheckpointMetadata,
1099) -> Result<VerifiedResumeArchive> {
1100 let archive_path = {
1101 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
1102 checkpoint.archive_path.canonicalize().with_context(|| {
1103 format!(
1104 "resolve checkpoint archive {}",
1105 checkpoint.archive_path.display()
1106 )
1107 })?
1108 };
1109 ensure!(
1110 archive_path.is_absolute() && archive_path.is_file(),
1111 "checkpoint archive path is not an absolute regular file: {}",
1112 archive_path.display()
1113 );
1114 let mj_checkpoint::archive::VerifiedArchiveMetadata {
1115 manifest,
1116 canonical_session,
1117 archive_sha256,
1118 } = {
1119 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
1120 verify_archive_streaming(&archive_path)?
1121 };
1122 if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
1123 bail!("persisted checkpoint verification failed");
1124 }
1125 ensure!(
1126 manifest.repositories.iter().all(|repository| {
1127 !repository.metadata.origin.starts_with("mj-local:")
1128 && !repository.metadata.origin.starts_with("ext::")
1129 }),
1130 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
1131 );
1132 Ok(VerifiedResumeArchive {
1133 archive_path,
1134 manifest,
1135 canonical_session: Arc::new(canonical_session),
1136 })
1137}
1138
1139pub fn raw_conversion_preview_for(
1140 session: &SessionRecord,
1141 checkout: &mj_core::state::Checkout<'_>,
1142 config: &Config,
1143 executor: &(impl CommandExecutor + Sync),
1144) -> Result<mj_core::state::RawConversionPreview> {
1145 let conversion = plan_raw_to_workspace_for_session(session, checkout, config, executor)?;
1146 raw_conversion_preview_with_checkout(session, &conversion, executor)
1147}
1148
1149fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
1150 let replacement = replacement.trim();
1151 ensure!(!replacement.is_empty(), "enter the repository's new origin");
1152 let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
1153 let path = expanded.as_path();
1154 let (github, local) = if path.is_absolute() {
1155 ensure!(
1156 path.is_dir(),
1157 "local repository {replacement:?} is not a directory"
1158 );
1159 (None, Some(mj_core::local_git::canonical_repository(path)?))
1160 } else {
1161 let github = crate::setup::github_repository_from_origin(replacement)
1162 .context("origin must be a GitHub repository or an absolute local repository path")?;
1163 (
1164 Some(format!("{}/{}", github.owner, github.repository)),
1165 None,
1166 )
1167 };
1168 Ok(ProjectRepository {
1169 id: id.to_owned(),
1170 github,
1171 local,
1172 destination: PathBuf::from(id),
1173 git_ref: None,
1174 })
1175}
1176
1177fn checkpoint_source_missing_commit(
1178 configured: &ProjectRepository,
1179 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1180 archived: &CheckpointRepositoryBundle,
1181 executor: &impl CommandExecutor,
1182 github_token: Option<&str>,
1183) -> Result<Option<String>> {
1184 let staging = tempfile::tempdir().context("create repository source preflight")?;
1185 let repository = staging.path().join("repository.git");
1186 checked_preflight_git(
1187 executor,
1188 CommandSpec::new(
1189 "git",
1190 [
1191 "init".to_owned(),
1192 "--bare".to_owned(),
1193 "--quiet".to_owned(),
1194 repository.to_string_lossy().into_owned(),
1195 ],
1196 )
1197 .purpose("initialize repository source preflight"),
1198 )?;
1199 let missing = checkpoint_bundle_prerequisites(archived)?;
1200 if missing.is_empty() {
1201 let bundle = staging.path().join("checkpoint.bundle");
1202 std::fs::write(&bundle, &archived.committed_bundle)
1203 .context("write self-contained checkpoint bundle for source preflight")?;
1204 checked_preflight_git(
1205 executor,
1206 checkpoint_bundle_import_command(&repository, &bundle),
1207 )?;
1208 return Ok(None);
1209 }
1210 for commit in missing {
1214 let output = fetch_source_commit(
1215 executor,
1216 &repository,
1217 configured,
1218 accepted_source,
1219 &commit,
1220 github_token,
1221 )?;
1222 if output.status != 0 {
1223 let stderr = String::from_utf8_lossy(&output.stderr);
1224 if source_does_not_have_commit(&stderr) {
1225 return Ok(Some(commit));
1226 }
1227 bail!(
1228 "could not check configured source {:?}: {}",
1229 configured.source_label(),
1230 stderr.trim()
1231 );
1232 }
1233 }
1234 Ok(None)
1238}
1239
1240fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
1241 let mut command = CommandSpec::new(
1242 "git",
1243 [
1244 "-C".to_owned(),
1245 repository.to_string_lossy().into_owned(),
1246 "fetch".to_owned(),
1247 "--no-tags".to_owned(),
1248 bundle.to_string_lossy().into_owned(),
1249 "HEAD".to_owned(),
1250 ],
1251 )
1252 .purpose("validate self-contained checkpoint bundle");
1253 command
1254 .env
1255 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1256 command
1257 .env
1258 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1259 command
1260}
1261
1262fn fetch_source_commit(
1263 executor: &impl CommandExecutor,
1264 repository: &Path,
1265 configured: &ProjectRepository,
1266 accepted_source: Option<&mj_core::remote_git::NetworkGitSource>,
1267 commit: &str,
1268 github_token: Option<&str>,
1269) -> Result<CommandOutput> {
1270 let mut arguments = Vec::new();
1271 let mut token_auth = false;
1272 let mut ssh_transport = false;
1273 let source = if let Some(source) = accepted_source {
1274 source.fetch_url.clone()
1275 } else if let Some(local) = &configured.local {
1276 local.to_string_lossy().into_owned()
1277 } else {
1278 let source = configured
1279 .github
1280 .as_deref()
1281 .context("repository source is missing")?;
1282 let github = crate::setup::github_repository_from_origin(source)
1283 .context("configured repository is not a GitHub source")?;
1284 if github_token.is_some() {
1285 token_auth = true;
1286 arguments.extend([
1287 "-c".to_owned(),
1288 "credential.helper=".to_owned(),
1289 "-c".to_owned(),
1290 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1291 ]);
1292 format!(
1293 "https://github.com/{}/{}.git",
1294 github.owner, github.repository
1295 )
1296 } else {
1297 ssh_transport = true;
1298 format!("git@github.com:{}/{}.git", github.owner, github.repository)
1299 }
1300 };
1301 if accepted_source.is_some() {
1302 if source.starts_with("https://github.com/") && github_token.is_some() {
1303 token_auth = true;
1304 arguments.extend([
1305 "-c".to_owned(), "credential.helper=".to_owned(), "-c".to_owned(),
1306 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
1307 ]);
1308 }
1309 ssh_transport = source.starts_with("git@") || source.starts_with("ssh://");
1310 }
1311 arguments.extend([
1312 "-C".to_owned(),
1313 repository.to_string_lossy().into_owned(),
1314 "fetch".to_owned(),
1315 "--no-tags".to_owned(),
1316 "--depth=1".to_owned(),
1317 "--filter=blob:none".to_owned(),
1318 source,
1319 commit.to_owned(),
1320 ]);
1321 let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1322 command
1323 .env
1324 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1325 command
1326 .env
1327 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1328 if token_auth {
1329 let token = github_token.expect("token authentication requires a GitHub token");
1330 command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1331 }
1332 if ssh_transport {
1333 command.env.insert(
1334 "GIT_SSH_COMMAND".to_owned(),
1335 "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1336 .to_owned(),
1337 );
1338 }
1339 executor.execute(&command)
1340}
1341
1342fn source_does_not_have_commit(stderr: &str) -> bool {
1343 let stderr = stderr.to_ascii_lowercase();
1344 [
1345 "not our ref",
1346 "couldn't find remote ref",
1347 "not a valid object name",
1348 "no such ref was fetched",
1349 ]
1350 .iter()
1351 .any(|needle| stderr.contains(needle))
1352}
1353
1354fn checked_preflight_git(
1355 executor: &impl CommandExecutor,
1356 command: CommandSpec,
1357) -> Result<CommandOutput> {
1358 let output = executor.execute(&command)?;
1359 ensure!(
1360 output.status == 0,
1361 "{}: {}",
1362 command.purpose,
1363 String::from_utf8_lossy(&output.stderr).trim()
1364 );
1365 Ok(output)
1366}
1367
1368impl Controller {
1369 pub async fn resume_session_with_options(
1374 &mut self,
1375 session_id: &str,
1376 profile_id: &str,
1377 target_id: &str,
1378 additional_mounts: Option<Vec<AdditionalMount>>,
1379 resource_allocation: Option<SessionResourceAllocation>,
1380 ) -> Result<MaterializedSession> {
1381 self.resume_session_with_options_and_queue_disposition(
1382 session_id,
1383 profile_id,
1384 target_id,
1385 additional_mounts,
1386 resource_allocation,
1387 false,
1388 )
1389 .await
1390 }
1391
1392 pub async fn resume_session_with_options_and_queue_disposition(
1393 &mut self,
1394 session_id: &str,
1395 profile_id: &str,
1396 target_id: &str,
1397 additional_mounts: Option<Vec<AdditionalMount>>,
1398 resource_allocation: Option<SessionResourceAllocation>,
1399 discard_queue: bool,
1400 ) -> Result<MaterializedSession> {
1401 self.resume_session_controlled(
1402 session_id,
1403 profile_id,
1404 target_id,
1405 SessionResumeOptions {
1406 additional_mounts,
1407 resource_allocation,
1408 discard_queue,
1409 },
1410 &ProcessExecutor,
1411 )
1412 .await
1413 }
1414
1415 pub async fn resume_session_controlled(
1416 &mut self,
1417 session_id: &str,
1418 profile_id: &str,
1419 target_id: &str,
1420 options: SessionResumeOptions,
1421 executor: &(impl CommandExecutor + Sync),
1422 ) -> Result<MaterializedSession> {
1423 self.resume_session_controlled_with_repository_preflight(
1424 session_id, profile_id, target_id, options, None, executor,
1425 )
1426 .await
1427 }
1428
1429 pub async fn resume_session_controlled_with_repository_preflight(
1430 &mut self,
1431 session_id: &str,
1432 profile_id: &str,
1433 target_id: &str,
1434 options: SessionResumeOptions,
1435 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1436 executor: &(impl CommandExecutor + Sync),
1437 ) -> Result<MaterializedSession> {
1438 self.resume_session_with_origin(
1439 session_id,
1440 profile_id,
1441 target_id,
1442 options,
1443 repository_preflight,
1444 None,
1445 executor,
1446 )
1447 .await
1448 }
1449
1450 pub(in crate::controller) async fn resume_session_for_move(
1451 &mut self,
1452 operation: &mj_core::state::MoveOperation,
1453 executor: &(impl CommandExecutor + Sync),
1454 ) -> Result<MaterializedSession> {
1455 self.resume_session_with_origin(
1456 &operation.selection.session_id,
1457 operation
1458 .selection
1459 .profile_id
1460 .as_deref()
1461 .context("Move profile missing")?,
1462 operation
1463 .selection
1464 .target_template_id
1465 .as_deref()
1466 .context("Move target missing")?,
1467 SessionResumeOptions {
1468 additional_mounts: operation.selection.additional_mounts.clone(),
1469 resource_allocation: operation.selection.resource_allocation.clone(),
1470 discard_queue: true,
1471 },
1472 None,
1473 Some(operation),
1474 executor,
1475 )
1476 .await
1477 }
1478
1479 #[allow(clippy::too_many_arguments)]
1480 async fn resume_session_with_origin(
1481 &mut self,
1482 session_id: &str,
1483 profile_id: &str,
1484 target_id: &str,
1485 options: SessionResumeOptions,
1486 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1487 move_operation: Option<&mj_core::state::MoveOperation>,
1488 executor: &(impl CommandExecutor + Sync),
1489 ) -> Result<MaterializedSession> {
1490 crate::worker_lifecycle::run(session_id, "resume session with origin", executor, async {
1491 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1492 let transferring_workspace = move_operation.is_some();
1493 if let Some(operation) = crate::database::load_move_operation(session_id)? {
1494 ensure!(
1495 transferring_workspace || !operation.holds_source_environment(),
1496 "a profile switch retains this environment; retry Move instead of recreating it with Resume"
1497 );
1498 }
1499 let SessionResumeOptions {
1500 additional_mounts,
1501 resource_allocation,
1502 discard_queue,
1503 } = options;
1504 let mut previous = self
1505 .state
1506 .sessions
1507 .get(session_id)
1508 .with_context(|| format!("unknown session {session_id}"))?
1509 .clone();
1510 let previous_checkout = self.state.checkout(session_id)?;
1511 let previous_project_directory = previous_checkout
1512 .project_directory()
1513 .map(|path| path.to_path_buf());
1514 let mut previous_managed_worktree = match &previous_checkout {
1515 mj_core::state::Checkout::ManagedWorktree { worktree, .. } => Some((*worktree).clone()),
1516 _ => None,
1517 };
1518 drop(previous_checkout);
1519 if !(matches!(
1520 previous.state,
1521 SessionState::Stopped | SessionState::Lost | SessionState::Error
1522 ) || (transferring_workspace && previous.state == SessionState::Closing))
1523 {
1524 bail!("session {session_id} is not stopped, lost, or retryable");
1525 }
1526 let checkpoint = move_operation
1527 .and_then(|op| op.handoff.as_ref())
1528 .or(previous.checkpoint.as_ref())
1529 .context("session has no checkpoint")?;
1530 let moving_to_raw = previous_project_directory.is_none()
1533 && self
1534 .config
1535 .targets
1536 .get(target_id)
1537 .is_some_and(mj_core::config::is_bare_project_target);
1538 if !transferring_workspace
1539 && (moving_to_raw
1540 || !repository_preflight.as_ref().is_some_and(|receipt| {
1541 self.repository_source_receipt_is_current(session_id, receipt)
1542 }))
1543 {
1544 let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1545 if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1546 self.preflight_repository_sources(
1547 session_id,
1548 target_id,
1549 false,
1550 self.github_token_for_repository_preflight(session_id)
1551 .await?
1552 .as_deref(),
1553 executor,
1554 )?
1555 {
1556 bail!(
1557 "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1558 mismatch.missing_commit,
1559 mismatch.configured_origin,
1560 mismatch.repository_id,
1561 mismatch.archived_origin,
1562 );
1563 }
1564 }
1565 let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1566 let archive_path = verified_archive.archive_path.clone();
1567 let archive_manifest = &verified_archive.manifest;
1568 let canonical_session = Arc::clone(&verified_archive.canonical_session);
1569 let profile = self
1570 .config
1571 .profiles
1572 .get(profile_id)
1573 .with_context(|| format!("unknown profile {profile_id:?}"))?
1574 .clone();
1575 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1576 let target_template = self
1577 .config
1578 .targets
1579 .get(target_id)
1580 .with_context(|| format!("unknown target template {target_id:?}"))?
1581 .clone();
1582 self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1585 ensure!(
1586 profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1587 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1588 );
1589 let checkout = self.state.checkout(session_id)?;
1590 let plan = resume_compatibility_with_checkout(
1591 &previous,
1592 &checkout,
1593 &self.config,
1594 target_id,
1595 )
1596 .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1597 if plan == ResumePlan::InPlace
1600 && previous_managed_worktree.is_some()
1601 && let Some(worktree) = previous.managed_worktree.as_mut()
1602 && let Ok(current) = managed_worktree_target(&target_template)
1603 {
1604 worktree.target = current;
1605 let refreshed_checkout = self.state.checkout_for_record(session_id, &previous)?;
1606 previous_managed_worktree = refreshed_checkout.effective().managed_worktree().cloned();
1607 }
1608 if !mj_core::config::is_bare_project_target(&target_template)
1612 && plan != ResumePlan::RawToWorkspace
1613 && !transferring_workspace
1614 {
1615 super::network_git::bundle_from_manifest(archive_manifest)?;
1616 }
1617 if plan == ResumePlan::InPlace
1618 && matches!(checkout.effective(), mj_core::state::Checkout::Attached { .. })
1619 && let Some(project_directory) = checkout.project_directory()
1620 {
1621 self.validate_project_directory(target_id, project_directory, executor)
1622 .context("raw project is unavailable for resume")?;
1623 }
1624 let conversion = match plan {
1625 ResumePlan::InPlace => None,
1626 ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(Box::new(
1627 plan_raw_to_workspace_for_session(&previous, &checkout, &self.config, executor)
1628 .context("prepare the raw checkout for its new target")?,
1629 ))),
1630 ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1631 self.plan_workspace_to_raw(&previous, target_id, executor)
1632 .context("prepare a checkout for this session")?,
1633 )),
1634 };
1635 let resource_allocation =
1636 resource_allocation.or_else(|| previous.resource_allocation.clone());
1637 let additional_mounts =
1638 additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1639 validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1640 let selected_container_size =
1641 selected_host_container_size(&target_template, resource_allocation.as_ref());
1642 if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1643 bail!("attached resources are unsupported for this target");
1644 }
1645 targets::validate_additional_mounts(&additional_mounts)?;
1646 let history_host = mount_history_host(&target_template);
1647 let history_mounts = additional_mounts.clone();
1648 if !transferring_workspace
1649 && previous.state == SessionState::Error
1650 && let Some(locator) = &previous.target
1651 {
1652 let backend = backend_locator(locator, &previous, &self.config)?;
1653 targets::close_plan(&backend, session_id)?
1654 .execute(executor)
1655 .context("clean up target from failed resume")?;
1656 }
1657 let mut resume_notices = Vec::new();
1658 if let Some(conversion) = conversion
1661 .as_ref()
1662 .and_then(ResumeConversion::workspace_to_raw)
1663 {
1664 resume_notices.push(format!(
1665 "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1666 previous.target_template_id,
1667 conversion.worktree.worktree_root.display(),
1668 archive_manifest
1669 .repositories
1670 .first()
1671 .and_then(|repository| repository.metadata.branch.as_deref())
1672 .unwrap_or("a detached head"),
1673 conversion.worktree.branch,
1674 ));
1675 }
1676 let managed_checkout_present = previous_managed_worktree
1677 .as_ref()
1678 .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1679 .transpose()?
1680 .unwrap_or(true);
1681 if managed_checkout_present
1684 && let Some(project_directory) = previous_project_directory.as_deref()
1685 {
1686 let checkout = self.state.checkout(session_id)?;
1687 match raw_checkout_position_with_checkout(
1688 &previous,
1689 &checkout,
1690 &self.config,
1691 project_directory,
1692 executor,
1693 ) {
1694 Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1695 project_directory,
1696 archive_manifest
1697 .repositories
1698 .first()
1699 .map(|repository| &repository.metadata),
1700 &live,
1701 )),
1702 Err(error) => tracing::warn!(
1705 session_id,
1706 error = format!("{error:#}"),
1707 "could not read the raw checkout position for a resume notice"
1708 ),
1709 }
1710 }
1711 super::worker_binary::preflight_worker_binary(&target_template, executor)?;
1716 let same_harness = profile.kind == archive_manifest.session.harness_kind;
1717 let native_continuity =
1718 native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1719 let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1720 let utility_config = (!native_continuity).then(|| self.config.clone());
1724 let discard_queued_prompts = discard_queue || !same_harness;
1725 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1729 .unwrap_or_else(|error| {
1730 tracing::warn!(
1731 session_id,
1732 error = format!("{error:#}"),
1733 "could not read the stored projection frontier; rebuilding it from the archive"
1734 );
1735 None
1736 });
1737 let rebuild_projection = projection_rebuild_required(
1738 stored_frontier
1739 .as_ref()
1740 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1741 canonical_session.event_frontier,
1742 &canonical_session.event_frontier_digest,
1743 );
1744 let projection_build = rebuild_projection.then(|| {
1754 let canonical = Arc::clone(&canonical_session);
1755 let session_id = session_id.to_owned();
1756 tokio::task::spawn_blocking(move || {
1757 materialized_session_from_canonical(session_id, &canonical)
1758 })
1759 });
1760 let github_token = self.github_token_for_session(session_id).await?;
1761
1762 let moving_into_first_container = self
1768 .config
1769 .targets
1770 .get(target_id)
1771 .is_some_and(mj_core::config::is_container_target)
1772 && !self
1773 .config
1774 .targets
1775 .get(&previous.target_template_id)
1776 .is_some_and(mj_core::config::is_container_target);
1777 let record = self.state.sessions.get_mut(session_id).unwrap();
1778 if record.container_workspace.is_none() && moving_into_first_container {
1779 record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1780 }
1781 record.harness_kind = profile.kind;
1782 record.last_profile = profile_id.to_string();
1783 record.target_template_id = target_id.to_string();
1784 record.target_runtime = Some(
1785 self.config
1786 .targets
1787 .get(target_id)
1788 .context("resume target disappeared")?
1789 .into(),
1790 );
1791 record.resource_allocation = resource_allocation;
1792 record.additional_mounts = additional_mounts;
1793 record.target = None;
1794 record.native_session_id =
1795 native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1796 record.state = SessionState::Provisioning;
1797 record.updated_at = now();
1798 record.last_error = None;
1799 match &conversion {
1800 Some(ResumeConversion::RawToWorkspace(conversion)) => {
1801 apply_raw_to_workspace(record, conversion)?;
1802 }
1803 Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1804 apply_workspace_to_raw(record, conversion);
1805 }
1806 None => {
1807 if let (Some(worktree), Some(refreshed)) = (
1808 record.managed_worktree.as_mut(),
1809 previous_managed_worktree.as_ref(),
1810 ) {
1811 worktree.target = refreshed.target.clone();
1812 }
1813 }
1814 }
1815 let resumed_project_directory = record
1816 .checkout()
1817 .project_directory()
1818 .map(|path| path.to_path_buf());
1819 let resumed_container_workspace = record.container_workspace.clone();
1820 if let Some(host) = history_host {
1821 self.state.remember_mount_sources(host, &history_mounts);
1822 crate::database::remember_mount_sources(host, &history_mounts)?;
1823 }
1824 crate::database::save_resumed_session(
1827 &self.state.sessions[session_id],
1828 selected_container_size
1829 .as_ref()
1830 .map(|(host, size)| (host.as_str(), *size)),
1831 )?;
1832 if let Some((host, size)) = selected_container_size.as_ref() {
1833 self.state.remember_container_size(host, *size);
1834 }
1835
1836 let mut recreated_managed_worktree = false;
1837 let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1840 let result = async {
1841 if let Some(worktree) = previous_managed_worktree.as_ref() {
1842 recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1843 if !transferring_workspace && recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1844 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1845 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1846 &archive_path, &worktree.worktree_root, &SystemGit,
1847 )?;
1848 } else {
1849 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1850 &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1851 )?;
1852 }
1853 }
1854 }
1855 if let Some(conversion) = conversion
1858 .as_ref()
1859 .and_then(ResumeConversion::workspace_to_raw)
1860 {
1861 if conversion.reuse_existing_branch {
1862 let recovery_ref =
1863 preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1864 restore_managed_worktree(executor, &conversion.worktree)?;
1865 resume_notices.push(format!(
1866 "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1867 ));
1868 } else {
1869 create_managed_worktree(
1870 executor,
1871 &conversion.worktree,
1872 None,
1873 PrimaryCheckoutRequirement::Any,
1874 )?;
1875 }
1876 if transferring_workspace {
1877 } else if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1879 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1880 &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1881 )?;
1882 } else {
1883 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1884 &archive_path,
1885 &conversion.worktree.worktree_root,
1886 &conversion.worktree.branch,
1887 &SystemGit,
1888 )?;
1889 }
1890 }
1891 if let Some(conversion) = conversion
1896 .as_ref()
1897 .and_then(ResumeConversion::raw_to_workspace)
1898 && !transferring_workspace
1899 {
1900 let snapshot = raw_checkout_snapshot(
1901 &conversion.checkout,
1902 &conversion.repository_id,
1903 &conversion.source,
1904 &conversion.destination,
1905 &SystemGit,
1906 conversion.retire.as_ref().is_some_and(|checkout| {
1907 checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1908 }),
1909 )
1910 .context("snapshot the host checkout for its new target")?;
1911 resume_notices.push(conversion_notice(
1912 target_id,
1913 previous_project_directory
1914 .as_deref()
1915 .unwrap_or(&conversion.checkout),
1916 snapshot.metadata.branch.as_deref(),
1917 conversion.retire.as_ref(),
1918 ));
1919 let archives = mj_core::config::sessions_dir();
1920 std::fs::create_dir_all(&archives).with_context(|| {
1921 format!("create the checkpoint directory {}", archives.display())
1922 })?;
1923 let output = archives.join(format!(
1927 "{session_id}-converted-{}-{}.hel.zip",
1928 previous
1929 .checkpoint
1930 .as_ref()
1931 .map_or(0, |checkpoint| checkpoint.event_frontier),
1932 new_command_id("archive")?
1933 ));
1934 let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1935 conversion_checkpoint_written = Some(written.clone());
1936 let record = self.state.sessions.get_mut(session_id).unwrap();
1937 record.checkpoint = Some(written);
1938 record.updated_at = now();
1939 crate::database::save_resumed_session(
1942 &self.state.sessions[session_id],
1943 selected_container_size.as_ref().map(|(host, size)| (host.as_str(), *size)),
1944 )?;
1945 }
1946 let utility_handoff = {
1947 let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1948 if let Some(config) = utility_config.as_ref() {
1949 Some(
1950 provision_with_cross_harness_handoff(
1951 self,
1952 session_id,
1953 executor,
1954 github_token.as_deref(),
1955 config,
1956 &canonical_session,
1957 context_bytes,
1958 )
1959 .context("prepare the cross-harness destination")?,
1960 )
1961 } else {
1962 self.provision_session_with_failure_disposition(
1963 session_id,
1964 executor,
1965 github_token.as_deref(),
1966 ProvisioningFailureDisposition::Preserve,
1967 )
1968 .await?;
1969 None
1970 }
1971 };
1972 if transferring_workspace {
1973 self.restore_move_workspace(session_id, executor)?;
1974 }
1975 let restore_repositories = !transferring_workspace && ((resumed_project_directory.is_none()
1976 && conversion.is_none())
1977 || plan == ResumePlan::RawToWorkspace
1978 || (recreated_managed_worktree && plan == ResumePlan::InPlace));
1979 let restored_archive = conversion_checkpoint_written
1982 .as_ref()
1983 .map_or(archive_path.as_path(), |checkpoint| {
1984 checkpoint.archive_path.as_path()
1985 });
1986 self.restore_into_target(
1987 session_id,
1988 RestoreIntoTarget {
1989 profile: &profile,
1990 archive: &verified_archive,
1991 restored_archive,
1992 resumed_project_directory,
1993 resumed_container_workspace,
1994 restore_repositories,
1995 native_continuity,
1996 discard_queued_prompts,
1997 replay_queue: !discard_queue,
1998 utility_handoff,
1999 projection_build,
2000 resume_notices,
2001 include_in_place_subagents: false,
2002 install_attached_resources: true,
2003 worker_root_reset: WorkerRootReset::FreshTarget,
2004 retire_after_ready: conversion
2005 .as_ref()
2006 .filter(|_| !transferring_workspace)
2007 .and_then(ResumeConversion::raw_to_workspace)
2008 .and_then(|plan| plan.retire.as_ref()),
2009 },
2010 executor,
2011 )
2012 .await
2013 }
2014 .await;
2015 match result {
2016 Ok(materialized) => {
2017 if let Some(written) = &conversion_checkpoint_written {
2020 super::checkpoint::prune_replaced_checkpoint(
2021 previous.checkpoint.as_ref(),
2022 written,
2023 );
2024 }
2025 Ok(materialized)
2026 }
2027 Err(error) => {
2028 if let Some(written) = &conversion_checkpoint_written
2031 && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
2032 && remove_error.kind() != std::io::ErrorKind::NotFound
2033 {
2034 tracing::warn!(
2035 session_id,
2036 path = %written.archive_path.display(),
2037 "could not remove the conversion checkpoint after resume failed: {remove_error}"
2038 );
2039 }
2040 restore_projection_after_failed_resume(
2043 session_id,
2044 &canonical_session,
2045 discard_queued_prompts,
2046 );
2047 if let Some(operation) = move_operation {
2048 return Err(self.rollback_move_destination(operation, error, executor)?);
2049 }
2050 Err(self.rollback_failed_resume(
2051 session_id,
2052 &previous,
2053 recreated_managed_worktree,
2054 error,
2055 executor,
2056 )?)
2057 }
2058 }
2059
2060 }).await
2061 }
2062
2063 pub(super) fn rollback_failed_resume(
2064 &mut self,
2065 session_id: &str,
2066 previous: &SessionRecord,
2067 recreated_managed_worktree: bool,
2068 error: anyhow::Error,
2069 _executor: &impl CommandExecutor,
2070 ) -> Result<anyhow::Error> {
2071 let current = self
2072 .state
2073 .sessions
2074 .get(session_id)
2075 .with_context(|| format!("unknown session {session_id}"))?
2076 .clone();
2077 let cleanup = match current.target.as_ref() {
2078 Some(locator) => (|| -> Result<()> {
2079 let backend = backend_locator(locator, ¤t, &self.config)?;
2080 targets::close_plan(&backend, session_id)?
2081 .execute(&CancellableProcessExecutor::with_timeout(
2084 Duration::from_secs(15),
2085 ))
2086 .map(|_| ())
2087 })(),
2088 None => Ok(()),
2089 };
2090 let worktree_cleanup = if cleanup.is_err() {
2093 Ok(())
2094 } else {
2095 let current_worktree = self.state.checkout(session_id)?.managed_worktree().cloned();
2096 let previous_checkout = self.state.checkout_for_record(session_id, previous)?;
2097 let previous_worktree = previous_checkout.managed_worktree().cloned();
2098 match (current_worktree.as_ref(), previous_worktree.as_ref()) {
2099 (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
2100 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
2101 previous,
2102 ),
2103 (Some(current), Some(previous)) if current == previous => Ok(()),
2104 (Some(worktree), _) => cleanup_managed_worktree(
2107 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
2108 worktree,
2109 crate::controller::BranchDisposition::Delete,
2110 ),
2111 (None, _) => Ok(()),
2112 }
2113 };
2114 let cleanup_error = [cleanup, worktree_cleanup]
2115 .into_iter()
2116 .filter_map(Result::err)
2117 .map(|cleanup_error| format!("{cleanup_error:#}"))
2118 .collect::<Vec<_>>()
2119 .join("; ");
2120 if !cleanup_error.is_empty() {
2121 tracing::warn!(
2122 session_id,
2123 error = %cleanup_error,
2124 "resume rollback cleanup reported failures"
2125 );
2126 }
2127 let original = super::worker_binary::failure_line(&error);
2130 let detail = format!("{error:#}");
2131 if detail != original {
2132 tracing::warn!(session_id, error = %detail, "resume failed");
2133 }
2134 let current_owns_managed_checkout = self
2135 .state
2136 .checkout(session_id)?
2137 .managed_worktree()
2138 .is_some();
2139 let record = self.state.sessions.get_mut(session_id).unwrap();
2140 let failure = apply_failed_resume_rollback_with_ownership(
2141 record,
2142 previous,
2143 current_owns_managed_checkout,
2144 &original,
2145 (!cleanup_error.is_empty()).then_some(cleanup_error),
2146 );
2147 crate::database::save_resumed_session(&self.state.sessions[session_id], None)?;
2150 Ok(failure)
2151 }
2152}
2153
2154fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
2155 format!(
2156 "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
2157 worktree_root.display(),
2158 worktree_root.display()
2159 )
2160}
2161
2162#[cfg(test)]
2163pub(super) fn apply_failed_resume_rollback(
2164 current: &mut SessionRecord,
2165 previous: &SessionRecord,
2166 original_error: &str,
2167 cleanup_error: Option<String>,
2168) -> anyhow::Error {
2169 let current_owns_managed_checkout = current.checkout().managed_worktree().is_some();
2170 apply_failed_resume_rollback_with_ownership(
2171 current,
2172 previous,
2173 current_owns_managed_checkout,
2174 original_error,
2175 cleanup_error,
2176 )
2177}
2178
2179fn apply_failed_resume_rollback_with_ownership(
2180 current: &mut SessionRecord,
2181 previous: &SessionRecord,
2182 current_owns_managed_checkout: bool,
2183 original_error: &str,
2184 cleanup_error: Option<String>,
2185) -> anyhow::Error {
2186 match cleanup_error {
2187 None => {
2188 *current = previous.clone();
2189 current.state = SessionState::Stopped;
2190 current.target = None;
2191 current.updated_at = now();
2192 current.last_error = Some(format!("resume failed: {original_error}"));
2193 anyhow::anyhow!(original_error.to_owned())
2194 }
2195 Some(cleanup_error) => {
2196 let failure = format!(
2197 "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
2198 );
2199 if !current_owns_managed_checkout {
2204 current
2207 .project_directory
2208 .clone_from(&previous.project_directory);
2209 current
2210 .managed_worktree
2211 .clone_from(&previous.managed_worktree);
2212 current.bundle_id.clone_from(&previous.bundle_id);
2213 current.project.clone_from(&previous.project);
2214 }
2215 current.state = SessionState::Error;
2216 current.updated_at = now();
2217 current.last_error = Some(format!("resume failed: {failure}"));
2218 anyhow::anyhow!(failure)
2219 }
2220 }
2221}
2222
2223fn conversion_notice(
2225 target_id: &str,
2226 checkout: &Path,
2227 branch: Option<&str>,
2228 retire: Option<&mj_core::state::ManagedWorktree>,
2229) -> String {
2230 let branch = branch.unwrap_or("a detached head");
2231 match retire {
2232 Some(worktree) => format!(
2233 "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
2234 checkout.display(),
2235 worktree.branch,
2236 worktree.source_repository.display()
2237 ),
2238 None => format!(
2239 "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.",
2240 checkout.display()
2241 ),
2242 }
2243}
2244
2245fn conversion_checkpoint(
2253 previous_archive: &Path,
2254 snapshot: mj_checkpoint::archive::RepositorySnapshot,
2255 output: &Path,
2256) -> Result<mj_core::state::CheckpointMetadata> {
2257 let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
2258 .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
2259 let native_artifacts = previous
2262 .manifest
2263 .payloads
2264 .iter()
2265 .filter_map(|descriptor| match &descriptor.role {
2266 mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
2267 Some((relative_path, descriptor))
2268 }
2269 _ => None,
2270 })
2271 .map(|(relative_path, descriptor)| {
2272 Ok(mj_checkpoint::archive::NativeArtifact {
2273 relative_path: relative_path.clone(),
2274 data: previous.payload(descriptor)?.to_vec(),
2275 mode: descriptor.mode,
2276 })
2277 })
2278 .collect::<Result<Vec<_>>>()?;
2279 let canonical_session = previous.canonical_session()?;
2280 let event_frontier = canonical_session.event_frontier;
2281 let written = mj_checkpoint::archive::write_archive_atomic(
2282 output,
2283 &mj_checkpoint::archive::ArchiveInput {
2284 session: previous.manifest.session.clone(),
2285 target: previous.manifest.target.clone(),
2288 bundle: mj_checkpoint::archive::BundleManifest {
2289 id: previous.manifest.bundle.id.clone(),
2290 primary_repository: snapshot.metadata.id.clone(),
2294 },
2295 canonical_session,
2296 native_artifacts,
2297 repositories: vec![snapshot],
2298 },
2299 )
2300 .with_context(|| format!("write the conversion archive {}", output.display()))?;
2301 Ok(mj_core::state::CheckpointMetadata {
2302 archive_path: output.to_path_buf(),
2303 sha256: written.archive_sha256,
2304 created_at: now(),
2305 event_frontier,
2306 })
2307}
2308
2309fn restored_native_session_accepted(
2320 archived: &str,
2321 opened: &str,
2322 replaced_unused: Option<&str>,
2323) -> bool {
2324 opened == archived || replaced_unused == Some(archived)
2325}
2326
2327fn restore_projection_after_failed_resume(
2338 session_id: &str,
2339 canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
2340 discard_queued_prompts: bool,
2341) {
2342 let stored_frontier = crate::database::materialized_event_frontier(session_id)
2343 .unwrap_or_else(|error| {
2344 tracing::warn!(
2345 session_id,
2346 error = format!("{error:#}"),
2347 "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
2348 );
2349 None
2350 });
2351 if projection_rebuild_required(
2352 stored_frontier
2353 .as_ref()
2354 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
2355 canonical_session.event_frontier,
2356 &canonical_session.event_frontier_digest,
2357 ) {
2358 match materialized_session_from_canonical(session_id, canonical_session) {
2359 Ok(previous_projection) => {
2360 if let Err(restore_error) =
2361 crate::database::save_materialized_session(&previous_projection)
2362 {
2363 tracing::error!(
2364 session_id,
2365 error = format!("{restore_error:#}"),
2366 "could not restore the durable projection after a failed resume"
2367 );
2368 }
2369 }
2370 Err(restore_error) => tracing::error!(
2371 session_id,
2372 error = format!("{restore_error:#}"),
2373 "could not rebuild the durable projection after a failed resume"
2374 ),
2375 }
2376 } else if discard_queued_prompts
2377 && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
2378 session_id,
2379 &mj_transcript::projection::materialized_queued_prompts_from_canonical(
2380 &canonical_session.queued_prompts,
2381 ),
2382 )
2383 {
2384 tracing::error!(
2385 session_id,
2386 error = format!("{restore_error:#}"),
2387 "could not restore queued prompts after a failed resume"
2388 );
2389 }
2390}
2391
2392fn projection_rebuild_required(
2400 stored: Option<(u64, &str)>,
2401 archive_frontier: u64,
2402 archive_frontier_digest: &str,
2403) -> bool {
2404 stored != Some((archive_frontier, archive_frontier_digest))
2405}
2406
2407fn restore_archive_path(
2408 backend: &targets::TargetLocator,
2409 verified_archive: &Path,
2410 remote_archive: &Path,
2411) -> PathBuf {
2412 if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2413 verified_archive.to_path_buf()
2414 } else {
2415 remote_archive.to_path_buf()
2416 }
2417}
2418
2419fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2420 !matches!(backend, targets::TargetLocator::LocalBare { .. })
2421}
2422
2423struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2428 inner: &'a E,
2429 cancellation: CancellationToken,
2430}
2431
2432impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2433 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2434 if self.cancellation.is_cancelled() {
2435 bail!("operation cancelled while provisioning destination");
2436 }
2437 self.inner.execute(command)
2438 }
2439
2440 fn cancellation_requested(&self) -> bool {
2441 self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2442 }
2443
2444 fn stage_started(&self, stage: ProvisionStage) {
2445 self.inner.stage_started(stage);
2446 }
2447
2448 fn stage_finished(&self, stage: ProvisionStage) {
2449 self.inner.stage_finished(stage);
2450 }
2451
2452 fn notify_notice(&self, notice: &str) {
2453 self.inner.notify_notice(notice);
2454 }
2455
2456 fn execute_with_stdin(
2457 &self,
2458 command: &CommandSpec,
2459 input: &mut (dyn std::io::Read + Send),
2460 ) -> Result<CommandOutput> {
2461 if self.cancellation.is_cancelled() {
2462 bail!("operation cancelled while provisioning destination");
2463 }
2464 self.inner.execute_with_stdin(command, input)
2465 }
2466}
2467
2468fn provision_with_cross_harness_handoff(
2469 controller: &mut Controller,
2470 session_id: &str,
2471 executor: &(impl CommandExecutor + Sync),
2472 github_token: Option<&str>,
2473 config: &Config,
2474 snapshot: &CanonicalSessionSnapshot,
2475 context_bytes: usize,
2476) -> Result<String> {
2477 let (_provision, handoff) = execute_joined_cross_harness_work(
2478 "cross-harness provisioning",
2479 move |cancellation| {
2480 let provision_executor = CrossHarnessProvisionExecutor {
2481 inner: executor,
2482 cancellation,
2483 };
2484 futures::executor::block_on(controller.provision_session_with_failure_disposition(
2485 session_id,
2486 &provision_executor,
2487 github_token,
2488 ProvisioningFailureDisposition::Preserve,
2489 ))
2490 },
2491 "cross-harness handoff",
2492 move |cancellation| {
2493 let runtime = tokio::runtime::Builder::new_current_thread()
2494 .enable_all()
2495 .build()
2496 .context("create cross-harness handoff runtime")?;
2497 runtime.block_on(utility_handoff_while_cancellable(
2498 session_id,
2499 config,
2500 snapshot,
2501 context_bytes,
2502 executor,
2503 cancellation,
2504 ))
2505 },
2506 )?;
2507 ensure!(
2508 !executor.cancellation_requested(),
2509 "operation cancelled while provisioning destination"
2510 );
2511 Ok(handoff)
2512}
2513
2514fn execute_joined_cross_harness_work<A: Send, B: Send>(
2518 first_name: &'static str,
2519 first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2520 second_name: &'static str,
2521 second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2522) -> Result<(A, B)> {
2523 let cancellation = CancellationToken::new();
2524 std::thread::scope(|scope| {
2525 let first_cancel = cancellation.clone();
2526 let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2527 let second_cancel = cancellation.clone();
2528 let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2529 let mut first_result = None;
2530 let mut second_result = None;
2531
2532 while first_result.is_none() || second_result.is_none() {
2533 if first_result.is_none()
2534 && first_handle
2535 .as_ref()
2536 .is_some_and(|handle| handle.is_finished())
2537 {
2538 let handle = first_handle.take().expect("first lane handle present");
2539 first_result = Some(match handle.join() {
2540 Ok(result) => result,
2541 Err(panic) => {
2542 cancellation.cancel();
2543 Err(anyhow::anyhow!(
2544 "{first_name} thread panicked: {}",
2545 targets::command_thread_panic_message(panic.as_ref())
2546 ))
2547 }
2548 });
2549 if first_result.as_ref().is_some_and(Result::is_err) {
2550 cancellation.cancel();
2551 }
2552 }
2553 if second_result.is_none()
2554 && second_handle
2555 .as_ref()
2556 .is_some_and(|handle| handle.is_finished())
2557 {
2558 let handle = second_handle.take().expect("second lane handle present");
2559 second_result = Some(match handle.join() {
2560 Ok(result) => result,
2561 Err(panic) => {
2562 cancellation.cancel();
2563 Err(anyhow::anyhow!(
2564 "{second_name} thread panicked: {}",
2565 targets::command_thread_panic_message(panic.as_ref())
2566 ))
2567 }
2568 });
2569 if second_result.as_ref().is_some_and(Result::is_err) {
2570 cancellation.cancel();
2571 }
2572 }
2573 if first_result.is_none() || second_result.is_none() {
2574 std::thread::sleep(Duration::from_millis(10));
2575 }
2576 }
2577
2578 match (
2579 first_result.expect("first lane result received after joined handle"),
2580 second_result.expect("second lane result received after joined handle"),
2581 ) {
2582 (Err(first), Err(second)) => {
2583 Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2584 }
2585 (Err(error), Ok(_)) => Err(error),
2586 (Ok(_), Err(error)) => Err(error),
2587 (Ok(first), Ok(second)) => Ok((first, second)),
2588 }
2589 })
2590}
2591
2592fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2596 profile_kind == archived_kind
2597}
2598
2599async fn utility_handoff_while_cancellable(
2603 session_id: &str,
2604 config: &Config,
2605 snapshot: &CanonicalSessionSnapshot,
2606 context_bytes: usize,
2607 executor: &impl CommandExecutor,
2608 cancellation: CancellationToken,
2609) -> Result<String> {
2610 let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2611 if executor.cancellation_requested() {
2612 bail!("operation cancelled while compacting the cross-harness handoff");
2613 }
2614 let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2615 let cancel = cancellation.child_token();
2616 let operation =
2617 crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2618 tokio::pin!(operation);
2619 loop {
2620 tokio::select! {
2621 context = &mut operation => return context,
2622 _ = cancellation.cancelled() => {
2623 cancel.cancel();
2624 bail!("operation cancelled while compacting the cross-harness handoff");
2625 }
2626 _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2627 if executor.cancellation_requested() {
2628 cancel.cancel();
2629 bail!("operation cancelled while compacting the cross-harness handoff");
2630 }
2631 }
2632 }
2633 }
2634}
2635
2636mod in_place;
2637pub(crate) use in_place::InPlaceRestartError;
2638
2639#[cfg(test)]
2640mod tests;