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