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