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