1use std::path::{Path, PathBuf};
4use std::sync::Arc;
5use std::time::{Duration, Instant};
6
7use agent_client_protocol::schema::v1::ContentBlock;
8use anyhow::{Context, Result, bail, ensure};
9use rayon::prelude::*;
10use tokio_util::sync::CancellationToken;
11
12use crate::checkpoint_transfer::restore_command;
13use crate::session_manager::new_command_id;
14use mj_checkpoint::archive::{
15 CanonicalQueuedCommandKind, CanonicalSessionSnapshot, CheckpointRepositoryBundle, SystemGit,
16 checkpoint_bundle_prerequisites, read_checkpoint_repository_bundles, verify_archive_streaming,
17};
18use mj_checkpoint::checkpoint::CheckpointRestoreSpec;
19use mj_core::config::{Config, HarnessKind, ProjectRepository, mount_history_host};
20use mj_core::state::{MaterializedSession, SessionRecord, SessionResourceAllocation, SessionState};
21use mj_transcript::projection::materialized_session_from_canonical;
22
23use crate::targets::{
24 self, AdditionalMount, CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec,
25 ProcessExecutor, ProvisionStage, ProvisionStageGuard,
26};
27use mj_core::relay::RelayCommand;
28
29use super::backend::{backend_locator, controller_github_token, validate_resource_allocation};
30use super::checkpoint::upload_checkpoint_spec;
31use super::provisioning::{
32 ProvisioningFailureDisposition, StagedExecutor, execute_concurrent_lanes,
33 install_attached_resources,
34};
35use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
36use super::worker_binary::{bridge_readiness_stage, start_worker, worker_probe_diagnosis};
37use super::worktree::{
38 PrimaryCheckoutRequirement, ResumeConversion, ResumePlan, apply_raw_to_workspace,
39 apply_workspace_to_raw, cleanup_managed_worktree, create_managed_worktree,
40 managed_worktree_checkout_exists, managed_worktree_target, plan_raw_to_workspace,
41 preserve_retained_managed_worktree_branch, raw_checkout_divergence_notice,
42 raw_checkout_position, raw_checkout_snapshot, raw_conversion_preview, restore_managed_worktree,
43 resume_compatibility, retire_managed_worktree,
44};
45use super::{
46 Controller, SessionResumeOptions, execute_checked, now, selected_host_container_size,
47 target_profile_home,
48};
49
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct ResumeRepositorySourceMismatch {
52 pub session_id: String,
53 pub bundle_id: String,
54 pub repository_id: String,
55 pub missing_commit: String,
56 pub archived_origin: String,
57 pub configured_origin: String,
58}
59
60pub use mj_core::state::ResumeRepositorySourceReceipt;
61
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum ResumeRepositorySourcePreflight {
64 Ready(ResumeRepositorySourceReceipt),
65 RepositoryMoved(ResumeRepositorySourceMismatch),
66 ConvertingRawCheckout {
71 receipt: ResumeRepositorySourceReceipt,
72 preview: Box<mj_core::state::RawConversionPreview>,
73 },
74}
75
76struct ResumeRepositoryBundles {
77 checkpoint_sha256: String,
78 repositories: Vec<CheckpointRepositoryBundle>,
79}
80
81struct ResumePhaseTimer<'a> {
85 session_id: &'a str,
86 phase: &'static str,
87 started: Instant,
88}
89
90impl<'a> ResumePhaseTimer<'a> {
91 fn new(session_id: &'a str, phase: &'static str) -> Self {
92 Self {
93 session_id,
94 phase,
95 started: Instant::now(),
96 }
97 }
98}
99
100impl Drop for ResumePhaseTimer<'_> {
101 fn drop(&mut self) {
102 tracing::debug!(
103 session_id = self.session_id,
104 phase = self.phase,
105 elapsed_ms = self.started.elapsed().as_millis(),
106 "resume phase completed"
107 );
108 }
109}
110
111impl Controller {
112 pub(super) fn validate_muse_resume_destination(
114 &self,
115 source: &SessionRecord,
116 destination_harness: HarnessKind,
117 _target_id: &str,
118 ) -> Result<()> {
119 if destination_harness != HarnessKind::Muse {
120 return Ok(());
121 }
122 ensure!(
123 source.project_directory.is_some()
124 || self
125 .config
126 .bundles
127 .get(&source.bundle_id)
128 .is_none_or(|bundle| bundle.repositories.len() == 1),
129 "Muse Code ACP supports one workspace root; use a single-repository bundle"
130 );
131 Ok(())
132 }
133 pub fn preflight_resume_repository_sources(
137 &self,
138 session_id: &str,
139 target_id: &str,
140 executor: &(impl CommandExecutor + Sync),
141 ) -> Result<ResumeRepositorySourcePreflight> {
142 self.preflight_repository_sources(session_id, target_id, true, executor)
143 }
144
145 fn preflight_repository_sources(
150 &self,
151 session_id: &str,
152 target_id: &str,
153 describe_conversion: bool,
154 executor: &(impl CommandExecutor + Sync),
155 ) -> Result<ResumeRepositorySourcePreflight> {
156 let session = self
157 .state
158 .sessions
159 .get(session_id)
160 .with_context(|| format!("unknown session {session_id}"))?;
161 let checkpoint = session
162 .checkpoint
163 .as_ref()
164 .context("session has no checkpoint")?;
165 let plan = resume_compatibility(session, &self.config, target_id)
166 .map_err(|reason| anyhow::anyhow!(reason))?;
167 if session.project_directory.is_some() {
168 debug_assert!(matches!(
169 plan,
170 ResumePlan::InPlace | ResumePlan::RawToWorkspace
171 ));
172 let receipt = ResumeRepositorySourceReceipt {
177 session_id: session_id.to_owned(),
178 bundle_id: session.bundle_id.clone(),
179 checkpoint_sha256: checkpoint.sha256.clone(),
180 repositories: Vec::new(),
181 };
182 if describe_conversion && plan == ResumePlan::RawToWorkspace {
187 let preview = raw_conversion_preview_for(session, &self.config, executor)?;
188 return Ok(ResumeRepositorySourcePreflight::ConvertingRawCheckout {
189 receipt,
190 preview: Box::new(preview),
191 });
192 }
193 return Ok(ResumeRepositorySourcePreflight::Ready(receipt));
194 }
195 let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
196 self.preflight_verified_repository_sources(
197 session_id,
198 ResumeRepositoryBundles {
199 checkpoint_sha256: checkpoint.sha256.clone(),
200 repositories,
201 },
202 None,
203 plan != ResumePlan::WorkspaceToRaw,
204 executor,
205 )
206 }
207
208 fn preflight_verified_repository_sources(
209 &self,
210 session_id: &str,
211 verified: ResumeRepositoryBundles,
212 skip_repository_id: Option<&str>,
213 use_archived_network_sources: bool,
214 executor: &(impl CommandExecutor + Sync),
215 ) -> Result<ResumeRepositorySourcePreflight> {
216 let session = self
217 .state
218 .sessions
219 .get(session_id)
220 .with_context(|| format!("unknown session {session_id}"))?;
221 ensure!(
222 verified.repositories.iter().all(|repository| {
223 !repository.metadata.origin.starts_with("mj-local:")
224 && !repository.metadata.origin.starts_with("ext::")
225 }),
226 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
227 );
228 if use_archived_network_sources
231 && verified
232 .repositories
233 .iter()
234 .all(|repository| repository.metadata.remote_workspace)
235 {
236 return Ok(ResumeRepositorySourcePreflight::Ready(
237 ResumeRepositorySourceReceipt {
238 session_id: session_id.to_owned(),
239 bundle_id: session.bundle_id.clone(),
240 checkpoint_sha256: verified.checkpoint_sha256,
241 repositories: Vec::new(),
242 },
243 ));
244 }
245 if verified.repositories.is_empty() {
246 return Ok(ResumeRepositorySourcePreflight::Ready(
247 ResumeRepositorySourceReceipt {
248 session_id: session_id.to_owned(),
249 bundle_id: session.bundle_id.clone(),
250 checkpoint_sha256: verified.checkpoint_sha256,
251 repositories: Vec::new(),
252 },
253 ));
254 }
255 let bundle = self
256 .config
257 .bundles
258 .get(&session.bundle_id)
259 .with_context(|| format!("session bundle {:?} is missing", session.bundle_id))?;
260 let configured = verified
261 .repositories
262 .iter()
263 .map(|archived| {
264 bundle
265 .repositories
266 .iter()
267 .find(|repository| repository.id == archived.metadata.id)
268 .cloned()
269 .with_context(|| {
270 format!(
271 "session bundle {:?} no longer contains repository {:?}",
272 session.bundle_id, archived.metadata.id
273 )
274 })
275 })
276 .collect::<Result<Vec<_>>>()?;
277 let github_token = configured
278 .iter()
279 .any(|repository| repository.github.is_some())
280 .then(controller_github_token)
281 .flatten();
282 let outcomes = verified
283 .repositories
284 .par_iter()
285 .zip(configured.par_iter())
286 .map(|(archived, configured)| {
287 if skip_repository_id == Some(configured.id.as_str()) {
288 return Ok(None);
289 }
290 checkpoint_source_missing_commit(
291 configured,
292 archived,
293 executor,
294 github_token.as_deref(),
295 )
296 .map(|missing_commit| {
297 missing_commit.map(|missing_commit| ResumeRepositorySourceMismatch {
298 session_id: session_id.to_owned(),
299 bundle_id: session.bundle_id.clone(),
300 repository_id: configured.id.clone(),
301 missing_commit,
302 archived_origin: archived.metadata.origin.clone(),
303 configured_origin: configured.source_label(),
304 })
305 })
306 })
307 .collect::<Vec<Result<Option<ResumeRepositorySourceMismatch>>>>();
308 for outcome in outcomes {
309 if let Some(mismatch) = outcome? {
310 return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(mismatch));
311 }
312 }
313 Ok(ResumeRepositorySourcePreflight::Ready(
314 ResumeRepositorySourceReceipt {
315 session_id: session_id.to_owned(),
316 bundle_id: session.bundle_id.clone(),
317 checkpoint_sha256: verified.checkpoint_sha256,
318 repositories: configured,
319 },
320 ))
321 }
322
323 fn repository_source_receipt_is_current(
324 &self,
325 session_id: &str,
326 receipt: &ResumeRepositorySourceReceipt,
327 ) -> bool {
328 let Some(session) = self.state.sessions.get(session_id) else {
329 return false;
330 };
331 if receipt.session_id != session_id
332 || receipt.bundle_id != session.bundle_id
333 || session
334 .checkpoint
335 .as_ref()
336 .map(|checkpoint| &checkpoint.sha256)
337 != Some(&receipt.checkpoint_sha256)
338 {
339 return false;
340 }
341 if receipt.repositories.is_empty() {
342 return true;
343 }
344 let Some(bundle) = self.config.bundles.get(&session.bundle_id) else {
345 return false;
346 };
347 receipt.repositories.iter().all(|expected| {
348 bundle
349 .repositories
350 .iter()
351 .any(|configured| configured == expected)
352 })
353 }
354
355 pub fn replace_resume_repository_origin(
359 &mut self,
360 session_id: &str,
361 repository_id: &str,
362 replacement: &str,
363 executor: &(impl CommandExecutor + Sync),
364 ) -> Result<ResumeRepositorySourcePreflight> {
365 let session = self
366 .state
367 .sessions
368 .get(session_id)
369 .with_context(|| format!("unknown session {session_id}"))?;
370 let bundle_id = session.bundle_id.clone();
371 let checkpoint = session
372 .checkpoint
373 .as_ref()
374 .context("session has no checkpoint")?;
375 let replacement = replacement_repository_source(repository_id, replacement)?;
376 let repositories = read_checkpoint_repository_bundles(&checkpoint.archive_path)?;
377 let verified = ResumeRepositoryBundles {
378 checkpoint_sha256: checkpoint.sha256.clone(),
379 repositories,
380 };
381 let archived = verified
382 .repositories
383 .iter()
384 .find(|repository| repository.metadata.id == repository_id)
385 .with_context(|| format!("checkpoint does not contain repository {repository_id:?}"))?;
386 if let Some(missing_commit) = checkpoint_source_missing_commit(
387 &replacement,
388 archived,
389 executor,
390 controller_github_token().as_deref(),
391 )? {
392 return Ok(ResumeRepositorySourcePreflight::RepositoryMoved(
393 ResumeRepositorySourceMismatch {
394 session_id: session_id.to_owned(),
395 bundle_id,
396 repository_id: repository_id.to_owned(),
397 missing_commit,
398 archived_origin: archived.metadata.origin.clone(),
399 configured_origin: replacement.source_label(),
400 },
401 ));
402 }
403 let (config, ()) = Config::update(|config| {
404 let bundle = config
405 .bundles
406 .get_mut(&bundle_id)
407 .with_context(|| format!("session bundle {bundle_id:?} is missing"))?;
408 let repository = bundle
409 .repositories
410 .iter_mut()
411 .find(|repository| repository.id == repository_id)
412 .with_context(|| {
413 format!(
414 "session bundle {:?} no longer contains repository {repository_id:?}",
415 bundle_id
416 )
417 })?;
418 repository.github = replacement.github.clone();
419 repository.local = replacement.local.clone();
420 Ok(())
421 })?;
422 self.config = config;
423 self.preflight_verified_repository_sources(
424 session_id,
425 verified,
426 Some(repository_id),
427 false,
428 executor,
429 )
430 }
431}
432
433pub(super) enum WorkerRootReset {
441 FreshTarget,
444 InPlace {
448 previous_profile_root: Option<String>,
451 },
452}
453
454pub(super) struct RestoreIntoTarget<'a> {
460 pub profile: &'a mj_core::config::HarnessProfile,
462 pub archive: &'a VerifiedResumeArchive,
463 pub restored_archive: &'a Path,
466 pub resumed_project_directory: Option<PathBuf>,
467 pub resumed_container_workspace: Option<PathBuf>,
468 pub restore_repositories: bool,
469 pub primary_repository_root_from_conversion: bool,
473 pub native_continuity: bool,
474 pub discard_queued_prompts: bool,
475 pub replay_queue: bool,
477 pub utility_handoff: Option<String>,
480 pub projection_build: Option<tokio::task::JoinHandle<Result<MaterializedSession>>>,
481 pub resume_notices: Vec<String>,
483 pub install_attached_resources: bool,
484 pub worker_root_reset: WorkerRootReset,
485 pub retire_after_ready: Option<&'a mj_core::state::ManagedWorktree>,
487}
488
489impl Controller {
490 pub(super) async fn restore_into_target(
493 &mut self,
494 session_id: &str,
495 restore: RestoreIntoTarget<'_>,
496 executor: &(impl CommandExecutor + Sync),
497 ) -> Result<MaterializedSession> {
498 let RestoreIntoTarget {
499 profile,
500 archive,
501 restored_archive,
502 resumed_project_directory,
503 resumed_container_workspace,
504 restore_repositories,
505 primary_repository_root_from_conversion,
506 native_continuity,
507 discard_queued_prompts,
508 replay_queue,
509 utility_handoff,
510 projection_build,
511 mut resume_notices,
512 install_attached_resources: should_install_attached_resources,
513 worker_root_reset,
514 retire_after_ready,
515 } = restore;
516 let archive_manifest = &archive.manifest;
517 let canonical_session = &archive.canonical_session;
518 let (backend, worker_root) = self.worker_placement(session_id)?;
519 let harness_home = target_profile_home(&backend, session_id, profile);
520 let workspace_root = if let Some(project_directory) = &resumed_project_directory {
521 project_directory
522 .parent()
523 .context("bare project directory has no parent")?
524 .to_string_lossy()
525 .into_owned()
526 } else {
527 super::network_git::workspace_root(&backend, resumed_container_workspace.as_deref())
528 };
529 let target_path = |path: &str| match &backend {
530 targets::TargetLocator::AwsEc2 { .. } | targets::TargetLocator::SshBare { .. }
531 if !path.starts_with('/') =>
532 {
533 PathBuf::from(format!("~/{path}"))
534 }
535 _ => PathBuf::from(path),
536 };
537 let remote_archive = format!("{worker_root}/restore.hel.zip");
538 let remote_spec = format!("{worker_root}/restore-spec.json");
539 let restore = CheckpointRestoreSpec {
540 archive_path: restore_archive_path(
541 &backend,
542 restored_archive,
543 &target_path(&remote_archive),
544 ),
545 workspace_root: target_path(&workspace_root),
546 relay_root: target_path(&worker_root),
547 harness_home: target_path(&harness_home),
548 restore_repositories,
554 restore_native: native_continuity,
555 primary_repository_root: primary_repository_root_from_conversion
562 .then(|| resumed_project_directory.clone())
563 .flatten()
564 .map(|directory| target_path(&directory.to_string_lossy())),
565 discard_queued_prompts,
566 };
567 {
571 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
572 match &worker_root_reset {
573 WorkerRootReset::FreshTarget => {
578 if let Some(command) = targets::clear_relay_state_plan(&backend, session_id)? {
579 execute_checked(syncing, command)?;
580 }
581 execute_checked(
584 syncing,
585 targets::command_on_locator(
586 &backend,
587 session_id,
588 vec!["mkdir".into(), "-p".into(), worker_root.clone()],
589 "create the session worker root",
590 )?,
591 )?;
592 }
593 WorkerRootReset::InPlace {
598 previous_profile_root,
599 } => {
600 execute_checked(
601 syncing,
602 targets::in_place_worker_reset_plan(
603 &backend,
604 session_id,
605 previous_profile_root.as_deref(),
606 )?,
607 )?;
608 }
609 }
610 }
611 let staging = tempfile::tempdir().context("create restore staging")?;
612 let local_spec = staging.path().join("restore-spec.json");
613 std::fs::write(&local_spec, serde_json::to_vec_pretty(&restore)?)?;
614 let controller = &*self;
618 let backend_ref = &backend;
619 let worker_root_ref = worker_root.as_str();
620 let local_spec_ref = local_spec.as_path();
621 execute_concurrent_lanes(
622 || {
623 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
624 controller.prepare_worker_files(
625 session_id,
626 backend_ref,
627 worker_root_ref,
628 syncing,
629 )?;
630 super::provisioning::install_inherited_git_settings(
631 syncing,
632 backend_ref,
633 session_id,
634 )?;
635 Ok(())
636 },
637 || {
638 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
639 if should_upload_restore_archive(&backend) {
640 upload_checkpoint_spec(
641 restoring,
642 backend_ref,
643 session_id,
644 restored_archive,
645 &remote_archive,
646 )?;
647 }
648 upload_checkpoint_spec(
649 restoring,
650 backend_ref,
651 session_id,
652 local_spec_ref,
653 &remote_spec,
654 )
655 },
656 )?;
657 {
658 let restoring = &StagedExecutor::new(executor, ProvisionStage::Restoring);
659 execute_checked(
660 restoring,
661 restore_command(&backend, session_id, &remote_spec)?,
662 )?;
663 }
664 if should_install_attached_resources {
665 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
666 install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
667 }
668 match projection_build {
669 Some(build) => {
670 let mut restored_projection = build
671 .await
672 .context("rebuild the restored projection")?
673 .context("rebuild the restored projection")?;
674 if discard_queued_prompts {
675 restored_projection.queued_prompts.clear();
676 }
677 crate::database::save_materialized_session(&restored_projection)?;
678 }
679 None if discard_queued_prompts => {
682 crate::database::replace_materialized_queued_prompts(session_id, &[])?;
683 }
684 None => {}
685 }
686 let readiness_stage = bridge_readiness_stage(profile);
687 let spec = self.reconnect_command(session_id)?;
688 let readiness = async {
689 let mut relay = {
690 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
691 start_worker(executor, &backend, &worker_root)?;
692 connect_started_worker(&spec, session_id, executor, &backend, &worker_root).await?
693 };
694 let native_session_id =
695 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
696 Ok::<_, anyhow::Error>((relay, native_session_id))
697 }
698 .await;
699 let (mut relay, native_session_id) = readiness
700 .map_err(|error| worker_probe_diagnosis(executor, &backend, &worker_root, error))?;
701 if native_continuity {
702 if !restored_native_session_accepted(
703 &archive_manifest.session.native_session_id,
704 &native_session_id,
705 relay
706 .operational()
707 .replaced_unused_native_session_id
708 .as_deref(),
709 ) {
710 bail!(
711 "ACP loaded native session {native_session_id}, expected {}",
712 archive_manifest.session.native_session_id
713 );
714 }
715 } else {
716 relay
717 .install_prompt_context(
718 utility_handoff
719 .clone()
720 .context("a resume into a fresh native session has no handoff")?,
721 )
722 .await?;
723 if replay_queue {
724 for prompt in &canonical_session.queued_prompts {
725 let command = match &prompt.kind {
729 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
730 prompt: prompt
731 .content
732 .iter()
733 .cloned()
734 .map(serde_json::from_value)
735 .collect::<serde_json::Result<Vec<ContentBlock>>>()?,
736 },
737 CanonicalQueuedCommandKind::SetConfig { key, value } => {
738 RelayCommand::SetConfig {
739 key: key.clone(),
740 value: value.clone(),
741 }
742 }
743 };
744 relay.submit(prompt.command_id.clone(), command).await?;
745 }
746 }
747 }
748 if let Some(worktree) = retire_after_ready
752 && let Err(error) = retire_managed_worktree(executor, worktree)
753 {
754 tracing::warn!(
755 session_id,
756 worktree = %worktree.worktree_root.display(),
757 error = format!("{error:#}"),
758 "could not retire the old managed worktree after resume"
759 );
760 resume_notices.push(worktree_cleanup_notice(&worktree.worktree_root, &error));
761 }
762 for notice in &resume_notices {
763 let submitted = async {
764 let command_id = new_command_id("resume-notice")?;
765 relay
766 .submit(
767 command_id,
768 RelayCommand::RecordNotice {
769 text: notice.clone(),
770 },
771 )
772 .await
773 }
774 .await;
775 if let Err(error) = submitted {
778 tracing::warn!(
779 session_id,
780 error = format!("{error:#}"),
781 "could not record a resume notice in the conversation"
782 );
783 }
784 }
785 self.mark_worker_connected(session_id, Some(native_session_id))?;
786 Ok(relay.sync().await?.materialized)
787 }
788}
789
790pub(super) struct VerifiedResumeArchive {
797 pub archive_path: PathBuf,
801 pub manifest: mj_checkpoint::archive::ArchiveManifest,
802 pub canonical_session: Arc<CanonicalSessionSnapshot>,
803}
804
805pub(super) fn verify_resume_checkpoint(
811 session_id: &str,
812 checkpoint: &mj_core::state::CheckpointMetadata,
813) -> Result<VerifiedResumeArchive> {
814 let archive_path = {
815 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive");
816 checkpoint.archive_path.canonicalize().with_context(|| {
817 format!(
818 "resolve checkpoint archive {}",
819 checkpoint.archive_path.display()
820 )
821 })?
822 };
823 ensure!(
824 archive_path.is_absolute() && archive_path.is_file(),
825 "checkpoint archive path is not an absolute regular file: {}",
826 archive_path.display()
827 );
828 let mj_checkpoint::archive::VerifiedArchiveMetadata {
829 manifest,
830 canonical_session,
831 archive_sha256,
832 } = {
833 let _phase = ResumePhaseTimer::new(session_id, "verify checkpoint archive contents");
834 verify_archive_streaming(&archive_path)?
835 };
836 if archive_sha256 != checkpoint.sha256 || manifest.session.id != session_id {
837 bail!("persisted checkpoint verification failed");
838 }
839 ensure!(
840 manifest.repositories.iter().all(|repository| {
841 !repository.metadata.origin.starts_with("mj-local:")
842 && !repository.metadata.origin.starts_with("ext::")
843 }),
844 "resuming legacy host-bridge sessions is not supported; start a new network-backed session"
845 );
846 Ok(VerifiedResumeArchive {
847 archive_path,
848 manifest,
849 canonical_session: Arc::new(canonical_session),
850 })
851}
852
853pub fn raw_conversion_preview_for(
854 session: &SessionRecord,
855 config: &Config,
856 executor: &(impl CommandExecutor + Sync),
857) -> Result<mj_core::state::RawConversionPreview> {
858 let conversion = plan_raw_to_workspace(session, config, executor)?;
859 raw_conversion_preview(session, &conversion, executor)
860}
861
862fn replacement_repository_source(id: &str, replacement: &str) -> Result<ProjectRepository> {
863 let replacement = replacement.trim();
864 ensure!(!replacement.is_empty(), "enter the repository's new origin");
865 let expanded = mj_core::path_input::expand_local(Path::new(replacement))?;
866 let path = expanded.as_path();
867 let (github, local) = if path.is_absolute() {
868 ensure!(
869 path.is_dir(),
870 "local repository {replacement:?} is not a directory"
871 );
872 (None, Some(mj_core::local_git::canonical_repository(path)?))
873 } else {
874 let github = crate::setup::github_repository_from_origin(replacement)
875 .context("origin must be a GitHub repository or an absolute local repository path")?;
876 (
877 Some(format!("{}/{}", github.owner, github.repository)),
878 None,
879 )
880 };
881 Ok(ProjectRepository {
882 id: id.to_owned(),
883 github,
884 local,
885 destination: PathBuf::from(id),
886 git_ref: None,
887 })
888}
889
890fn checkpoint_source_missing_commit(
891 configured: &ProjectRepository,
892 archived: &CheckpointRepositoryBundle,
893 executor: &impl CommandExecutor,
894 github_token: Option<&str>,
895) -> Result<Option<String>> {
896 let staging = tempfile::tempdir().context("create repository source preflight")?;
897 let repository = staging.path().join("repository.git");
898 checked_preflight_git(
899 executor,
900 CommandSpec::new(
901 "git",
902 [
903 "init".to_owned(),
904 "--bare".to_owned(),
905 "--quiet".to_owned(),
906 repository.to_string_lossy().into_owned(),
907 ],
908 )
909 .purpose("initialize repository source preflight"),
910 )?;
911 let missing = checkpoint_bundle_prerequisites(archived)?;
912 if missing.is_empty() {
913 let bundle = staging.path().join("checkpoint.bundle");
914 std::fs::write(&bundle, &archived.committed_bundle)
915 .context("write self-contained checkpoint bundle for source preflight")?;
916 checked_preflight_git(
917 executor,
918 checkpoint_bundle_import_command(&repository, &bundle),
919 )?;
920 return Ok(None);
921 }
922 for commit in missing {
926 let output = fetch_source_commit(executor, &repository, configured, &commit, github_token)?;
927 if output.status != 0 {
928 let stderr = String::from_utf8_lossy(&output.stderr);
929 if source_does_not_have_commit(&stderr) {
930 return Ok(Some(commit));
931 }
932 bail!(
933 "could not check configured source {:?}: {}",
934 configured.source_label(),
935 stderr.trim()
936 );
937 }
938 }
939 Ok(None)
943}
944
945fn checkpoint_bundle_import_command(repository: &Path, bundle: &Path) -> CommandSpec {
946 let mut command = CommandSpec::new(
947 "git",
948 [
949 "-C".to_owned(),
950 repository.to_string_lossy().into_owned(),
951 "fetch".to_owned(),
952 "--no-tags".to_owned(),
953 bundle.to_string_lossy().into_owned(),
954 "HEAD".to_owned(),
955 ],
956 )
957 .purpose("validate self-contained checkpoint bundle");
958 command
959 .env
960 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
961 command
962 .env
963 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
964 command
965}
966
967fn fetch_source_commit(
968 executor: &impl CommandExecutor,
969 repository: &Path,
970 configured: &ProjectRepository,
971 commit: &str,
972 github_token: Option<&str>,
973) -> Result<CommandOutput> {
974 let mut arguments = Vec::new();
975 let mut token_auth = false;
976 let mut ssh_transport = false;
977 let source = if let Some(local) = &configured.local {
978 local.to_string_lossy().into_owned()
979 } else {
980 let source = configured
981 .github
982 .as_deref()
983 .context("repository source is missing")?;
984 let github = crate::setup::github_repository_from_origin(source)
985 .context("configured repository is not a GitHub source")?;
986 if github_token.is_some() {
987 token_auth = true;
988 arguments.extend([
989 "-c".to_owned(),
990 "credential.helper=".to_owned(),
991 "-c".to_owned(),
992 "credential.helper=!f() { if [ \"$1\" = get ]; then echo username=x-access-token; echo \"password=$GH_TOKEN\"; fi; }; f".to_owned(),
993 ]);
994 format!(
995 "https://github.com/{}/{}.git",
996 github.owner, github.repository
997 )
998 } else {
999 ssh_transport = true;
1000 format!("git@github.com:{}/{}.git", github.owner, github.repository)
1001 }
1002 };
1003 arguments.extend([
1004 "-C".to_owned(),
1005 repository.to_string_lossy().into_owned(),
1006 "fetch".to_owned(),
1007 "--no-tags".to_owned(),
1008 "--depth=1".to_owned(),
1009 "--filter=blob:none".to_owned(),
1010 source,
1011 commit.to_owned(),
1012 ]);
1013 let mut command = CommandSpec::new("git", arguments).purpose("check checkpoint base commit");
1014 command
1015 .env
1016 .insert("GIT_NO_LAZY_FETCH".to_owned(), "1".to_owned());
1017 command
1018 .env
1019 .insert("GIT_TERMINAL_PROMPT".to_owned(), "0".to_owned());
1020 if token_auth {
1021 let token = github_token.expect("token authentication requires a GitHub token");
1022 command.env.insert("GH_TOKEN".to_owned(), token.to_owned());
1023 }
1024 if ssh_transport {
1025 command.env.insert(
1026 "GIT_SSH_COMMAND".to_owned(),
1027 "ssh -o BatchMode=yes -o StrictHostKeyChecking=accept-new -o ConnectTimeout=15"
1028 .to_owned(),
1029 );
1030 }
1031 executor.execute(&command)
1032}
1033
1034fn source_does_not_have_commit(stderr: &str) -> bool {
1035 let stderr = stderr.to_ascii_lowercase();
1036 [
1037 "not our ref",
1038 "couldn't find remote ref",
1039 "not a valid object name",
1040 "no such ref was fetched",
1041 ]
1042 .iter()
1043 .any(|needle| stderr.contains(needle))
1044}
1045
1046fn checked_preflight_git(
1047 executor: &impl CommandExecutor,
1048 command: CommandSpec,
1049) -> Result<CommandOutput> {
1050 let output = executor.execute(&command)?;
1051 ensure!(
1052 output.status == 0,
1053 "{}: {}",
1054 command.purpose,
1055 String::from_utf8_lossy(&output.stderr).trim()
1056 );
1057 Ok(output)
1058}
1059
1060impl Controller {
1061 pub async fn resume_session_with_options(
1066 &mut self,
1067 session_id: &str,
1068 profile_id: &str,
1069 target_id: &str,
1070 additional_mounts: Option<Vec<AdditionalMount>>,
1071 resource_allocation: Option<SessionResourceAllocation>,
1072 ) -> Result<MaterializedSession> {
1073 self.resume_session_with_options_and_queue_disposition(
1074 session_id,
1075 profile_id,
1076 target_id,
1077 additional_mounts,
1078 resource_allocation,
1079 false,
1080 )
1081 .await
1082 }
1083
1084 pub async fn resume_session_with_options_and_queue_disposition(
1085 &mut self,
1086 session_id: &str,
1087 profile_id: &str,
1088 target_id: &str,
1089 additional_mounts: Option<Vec<AdditionalMount>>,
1090 resource_allocation: Option<SessionResourceAllocation>,
1091 discard_queue: bool,
1092 ) -> Result<MaterializedSession> {
1093 self.resume_session_controlled(
1094 session_id,
1095 profile_id,
1096 target_id,
1097 SessionResumeOptions {
1098 additional_mounts,
1099 resource_allocation,
1100 discard_queue,
1101 },
1102 &ProcessExecutor,
1103 )
1104 .await
1105 }
1106
1107 pub async fn resume_session_controlled(
1108 &mut self,
1109 session_id: &str,
1110 profile_id: &str,
1111 target_id: &str,
1112 options: SessionResumeOptions,
1113 executor: &(impl CommandExecutor + Sync),
1114 ) -> Result<MaterializedSession> {
1115 self.resume_session_controlled_with_repository_preflight(
1116 session_id, profile_id, target_id, options, None, executor,
1117 )
1118 .await
1119 }
1120
1121 pub async fn resume_session_controlled_with_repository_preflight(
1122 &mut self,
1123 session_id: &str,
1124 profile_id: &str,
1125 target_id: &str,
1126 options: SessionResumeOptions,
1127 repository_preflight: Option<ResumeRepositorySourceReceipt>,
1128 executor: &(impl CommandExecutor + Sync),
1129 ) -> Result<MaterializedSession> {
1130 let SessionResumeOptions {
1131 additional_mounts,
1132 resource_allocation,
1133 discard_queue,
1134 } = options;
1135 let mut previous = self
1136 .state
1137 .sessions
1138 .get(session_id)
1139 .with_context(|| format!("unknown session {session_id}"))?
1140 .clone();
1141 if !matches!(
1142 previous.state,
1143 SessionState::Stopped | SessionState::Lost | SessionState::Error
1144 ) {
1145 bail!("session {session_id} is not stopped, lost, or retryable");
1146 }
1147 let checkpoint = previous
1148 .checkpoint
1149 .as_ref()
1150 .context("session has no checkpoint")?;
1151 let moving_to_raw = previous.project_directory.is_none()
1154 && self
1155 .config
1156 .targets
1157 .get(target_id)
1158 .is_some_and(mj_core::config::is_bare_project_target);
1159 if moving_to_raw
1160 || !repository_preflight.as_ref().is_some_and(|receipt| {
1161 self.repository_source_receipt_is_current(session_id, receipt)
1162 })
1163 {
1164 let _phase = ResumePhaseTimer::new(session_id, "preflight repository sources");
1165 if let ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) =
1166 self.preflight_repository_sources(session_id, target_id, false, executor)?
1167 {
1168 bail!(
1169 "checkpoint base commit {} is missing from configured source {:?} for repository {:?}; the repository may have moved (archived origin: {:?})",
1170 mismatch.missing_commit,
1171 mismatch.configured_origin,
1172 mismatch.repository_id,
1173 mismatch.archived_origin,
1174 );
1175 }
1176 }
1177 let verified_archive = verify_resume_checkpoint(session_id, checkpoint)?;
1178 let archive_path = verified_archive.archive_path.clone();
1179 let archive_manifest = &verified_archive.manifest;
1180 let canonical_session = Arc::clone(&verified_archive.canonical_session);
1181 let profile = self
1182 .config
1183 .profiles
1184 .get(profile_id)
1185 .with_context(|| format!("unknown profile {profile_id:?}"))?
1186 .clone();
1187 ensure!(profile.enabled, "profile {profile_id:?} is disabled");
1188 let target_template = self
1189 .config
1190 .targets
1191 .get(target_id)
1192 .with_context(|| format!("unknown target template {target_id:?}"))?
1193 .clone();
1194 self.validate_muse_resume_destination(&previous, profile.kind, target_id)?;
1197 ensure!(
1198 profile.kind != HarnessKind::Muse || previous.additional_mounts.is_empty(),
1199 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1200 );
1201 let plan = resume_compatibility(&previous, &self.config, target_id)
1202 .map_err(|reason| anyhow::anyhow!("{reason}"))?;
1203 if plan == ResumePlan::InPlace
1206 && let Some(worktree) = previous.managed_worktree.as_mut()
1207 && let Ok(current) = managed_worktree_target(&target_template)
1208 {
1209 worktree.target = current;
1210 }
1211 if !mj_core::config::is_bare_project_target(&target_template)
1215 && plan != ResumePlan::RawToWorkspace
1216 {
1217 super::network_git::bundle_from_manifest(archive_manifest)?;
1218 }
1219 if plan == ResumePlan::InPlace
1220 && previous.managed_worktree.is_none()
1221 && let Some(project_directory) = &previous.project_directory
1222 {
1223 self.validate_project_directory(target_id, project_directory, executor)
1224 .context("raw project is unavailable for resume")?;
1225 }
1226 let conversion = match plan {
1227 ResumePlan::InPlace => None,
1228 ResumePlan::RawToWorkspace => Some(ResumeConversion::RawToWorkspace(
1229 plan_raw_to_workspace(&previous, &self.config, executor)
1230 .context("prepare the raw checkout for its new target")?,
1231 )),
1232 ResumePlan::WorkspaceToRaw => Some(ResumeConversion::WorkspaceToRaw(
1233 self.plan_workspace_to_raw(&previous, target_id, executor)
1234 .context("prepare a checkout for this session")?,
1235 )),
1236 };
1237 let resource_allocation =
1238 resource_allocation.or_else(|| previous.resource_allocation.clone());
1239 let additional_mounts =
1240 additional_mounts.unwrap_or_else(|| previous.additional_mounts.clone());
1241 validate_resource_allocation(&target_template, resource_allocation.as_ref())?;
1242 let selected_container_size =
1243 selected_host_container_size(&target_template, resource_allocation.as_ref());
1244 if !additional_mounts.is_empty() && mount_history_host(&target_template).is_none() {
1245 bail!("attached resources are unsupported for this target");
1246 }
1247 targets::validate_additional_mounts(&additional_mounts)?;
1248 let history_host = mount_history_host(&target_template);
1249 let history_mounts = additional_mounts.clone();
1250 if previous.state == SessionState::Error
1251 && let Some(locator) = &previous.target
1252 {
1253 let backend = backend_locator(locator, &previous, &self.config)?;
1254 targets::close_plan(&backend, session_id)?
1255 .execute(executor)
1256 .context("clean up target from failed resume")?;
1257 }
1258 let mut resume_notices = Vec::new();
1259 if let Some(conversion) = conversion
1262 .as_ref()
1263 .and_then(ResumeConversion::workspace_to_raw)
1264 {
1265 resume_notices.push(format!(
1266 "This session moved out of its {} target and into {}. Its branch {} is now {}.",
1267 previous.target_template_id,
1268 conversion.worktree.worktree_root.display(),
1269 archive_manifest
1270 .repositories
1271 .first()
1272 .and_then(|repository| repository.metadata.branch.as_deref())
1273 .unwrap_or("a detached head"),
1274 conversion.worktree.branch,
1275 ));
1276 }
1277 let managed_checkout_present = previous
1278 .managed_worktree
1279 .as_ref()
1280 .map(|worktree| managed_worktree_checkout_exists(executor, worktree))
1281 .transpose()?
1282 .unwrap_or(true);
1283 if managed_checkout_present && let Some(project_directory) = &previous.project_directory {
1286 match raw_checkout_position(&previous, &self.config, project_directory, executor) {
1287 Ok(live) => resume_notices.extend(raw_checkout_divergence_notice(
1288 project_directory,
1289 archive_manifest
1290 .repositories
1291 .first()
1292 .map(|repository| &repository.metadata),
1293 &live,
1294 )),
1295 Err(error) => tracing::warn!(
1298 session_id,
1299 error = format!("{error:#}"),
1300 "could not read the raw checkout position for a resume notice"
1301 ),
1302 }
1303 }
1304 super::worker_binary::preflight_worker_binary(&target_template)?;
1309 let same_harness = profile.kind == archive_manifest.session.harness_kind;
1310 let native_continuity =
1311 native_continuity_preserved(profile.kind, archive_manifest.session.harness_kind);
1312 let context_bytes = crate::handoff::profile_handoff_bytes(Some(&profile));
1313 let utility_config = (!native_continuity).then(|| self.config.clone());
1317 let discard_queued_prompts = discard_queue || !same_harness;
1318 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1322 .unwrap_or_else(|error| {
1323 tracing::warn!(
1324 session_id,
1325 error = format!("{error:#}"),
1326 "could not read the stored projection frontier; rebuilding it from the archive"
1327 );
1328 None
1329 });
1330 let rebuild_projection = projection_rebuild_required(
1331 stored_frontier
1332 .as_ref()
1333 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1334 canonical_session.event_frontier,
1335 &canonical_session.event_frontier_digest,
1336 );
1337 let projection_build = rebuild_projection.then(|| {
1347 let canonical = Arc::clone(&canonical_session);
1348 let session_id = session_id.to_owned();
1349 tokio::task::spawn_blocking(move || {
1350 materialized_session_from_canonical(session_id, &canonical)
1351 })
1352 });
1353 let github_token = controller_github_token();
1354
1355 if let Some(conversion) = conversion
1358 .as_ref()
1359 .and_then(ResumeConversion::raw_to_workspace)
1360 && let Some(bundle) = &conversion.new_bundle
1361 {
1362 let (config, ()) = Config::update(|config| {
1363 if let Some(existing) = config.bundles.get(&conversion.bundle_id) {
1364 ensure!(
1365 existing == bundle,
1366 "bundle {:?} was configured concurrently with a different definition; retry the resume",
1367 conversion.bundle_id
1368 );
1369 } else {
1370 config
1371 .bundles
1372 .insert(conversion.bundle_id.clone(), bundle.clone());
1373 }
1374 Ok(())
1375 })
1376 .context("save the bundle for a converted raw session")?;
1377 self.config = config;
1378 }
1379
1380 let moving_into_first_container = self
1386 .config
1387 .targets
1388 .get(target_id)
1389 .is_some_and(mj_core::config::is_container_target)
1390 && !self
1391 .config
1392 .targets
1393 .get(&previous.target_template_id)
1394 .is_some_and(mj_core::config::is_container_target);
1395 let record = self.state.sessions.get_mut(session_id).unwrap();
1396 if record.container_workspace.is_none() && moving_into_first_container {
1397 record.container_workspace = Some(targets::new_container_workspace(session_id)?);
1398 }
1399 record.harness_kind = profile.kind;
1400 record.last_profile = profile_id.to_string();
1401 record.target_template_id = target_id.to_string();
1402 record.target_runtime = Some(
1403 self.config
1404 .targets
1405 .get(target_id)
1406 .context("resume target disappeared")?
1407 .into(),
1408 );
1409 record.resource_allocation = resource_allocation;
1410 record.additional_mounts = additional_mounts;
1411 record.target = None;
1412 record.native_session_id =
1413 native_continuity.then(|| archive_manifest.session.native_session_id.clone());
1414 record.state = SessionState::Provisioning;
1415 record.updated_at = now();
1416 record.last_error = None;
1417 match &conversion {
1418 Some(ResumeConversion::RawToWorkspace(conversion)) => {
1419 apply_raw_to_workspace(record, conversion);
1420 }
1421 Some(ResumeConversion::WorkspaceToRaw(conversion)) => {
1422 apply_workspace_to_raw(record, conversion);
1423 }
1424 None => {
1425 if let (Some(worktree), Some(refreshed)) = (
1426 record.managed_worktree.as_mut(),
1427 previous.managed_worktree.as_ref(),
1428 ) {
1429 worktree.target = refreshed.target.clone();
1430 }
1431 }
1432 }
1433 let resumed_project_directory = record.project_directory.clone();
1434 let resumed_container_workspace = record.container_workspace.clone();
1435 if let Some(host) = history_host {
1436 self.state.remember_mount_sources(host, &history_mounts);
1437 crate::database::remember_mount_sources(host, &history_mounts)?;
1438 }
1439 if let Some(conversion) = conversion
1442 .as_ref()
1443 .and_then(ResumeConversion::raw_to_workspace)
1444 {
1445 crate::database::rebind_session_bundle(session_id, &conversion.bundle_id)?;
1446 }
1447 if let Some((host, size)) = selected_container_size.as_ref() {
1450 crate::database::save_session_with_container_size(
1451 &self.state.sessions[session_id],
1452 host,
1453 *size,
1454 )?;
1455 } else {
1456 crate::database::save_session(&self.state.sessions[session_id])?;
1457 }
1458 if let Some((host, size)) = selected_container_size.as_ref() {
1459 self.state.remember_container_size(host, *size);
1460 }
1461
1462 let mut recreated_managed_worktree = false;
1463 let mut conversion_checkpoint_written: Option<mj_core::state::CheckpointMetadata> = None;
1466 let result = async {
1467 if let Some(worktree) = previous.managed_worktree.as_ref() {
1468 recreated_managed_worktree = restore_managed_worktree(executor, worktree)?;
1469 if recreated_managed_worktree && plan == ResumePlan::RawToWorkspace {
1470 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1471 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1472 &archive_path, &worktree.worktree_root, &SystemGit,
1473 )?;
1474 } else {
1475 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1476 &archive_path, &worktree.worktree_root, &worktree.branch, &SystemGit,
1477 )?;
1478 }
1479 }
1480 }
1481 if let Some(conversion) = conversion
1484 .as_ref()
1485 .and_then(ResumeConversion::workspace_to_raw)
1486 {
1487 if conversion.reuse_existing_branch {
1488 let recovery_ref =
1489 preserve_retained_managed_worktree_branch(executor, &conversion.worktree)?;
1490 restore_managed_worktree(executor, &conversion.worktree)?;
1491 resume_notices.push(format!(
1492 "Before restoring this session's retained branch, Mjolnir preserved its tip at {recovery_ref}."
1493 ));
1494 } else {
1495 create_managed_worktree(
1496 executor,
1497 &conversion.worktree,
1498 None,
1499 PrimaryCheckoutRequirement::Any,
1500 )?;
1501 }
1502 if conversion.worktree.kind == mj_core::state::ManagedCheckoutKind::Clone {
1503 mj_checkpoint::checkpoint::restore_single_repository_into_checkout(
1504 &archive_path, &conversion.worktree.worktree_root, &SystemGit,
1505 )?;
1506 } else {
1507 mj_checkpoint::checkpoint::restore_single_repository_onto_branch(
1508 &archive_path,
1509 &conversion.worktree.worktree_root,
1510 &conversion.worktree.branch,
1511 &SystemGit,
1512 )?;
1513 }
1514 }
1515 if let Some(conversion) = conversion
1520 .as_ref()
1521 .and_then(ResumeConversion::raw_to_workspace)
1522 {
1523 let destination = PathBuf::from(
1524 previous
1525 .project_directory
1526 .as_deref()
1527 .context("a raw session has no project directory")?
1528 .file_name()
1529 .context("a raw project directory cannot be the filesystem root")?,
1530 );
1531 let snapshot = raw_checkout_snapshot(
1532 &conversion.checkout,
1533 &conversion.source,
1534 &destination,
1535 &SystemGit,
1536 conversion.retire.as_ref().is_some_and(|checkout| {
1537 checkout.kind == mj_core::state::ManagedCheckoutKind::Clone
1538 }),
1539 )
1540 .context("snapshot the host checkout for its new target")?;
1541 resume_notices.push(conversion_notice(
1542 target_id,
1543 previous
1544 .project_directory
1545 .as_deref()
1546 .unwrap_or(&conversion.checkout),
1547 snapshot.metadata.branch.as_deref(),
1548 conversion.retire.as_ref(),
1549 ));
1550 let archives = mj_core::config::sessions_dir();
1551 std::fs::create_dir_all(&archives).with_context(|| {
1552 format!("create the checkpoint directory {}", archives.display())
1553 })?;
1554 let output = archives.join(format!(
1558 "{session_id}-converted-{}-{}.hel.zip",
1559 previous
1560 .checkpoint
1561 .as_ref()
1562 .map_or(0, |checkpoint| checkpoint.event_frontier),
1563 new_command_id("archive")?
1564 ));
1565 let written = conversion_checkpoint(&archive_path, snapshot, &output)?;
1566 conversion_checkpoint_written = Some(written.clone());
1567 let record = self.state.sessions.get_mut(session_id).unwrap();
1568 record.checkpoint = Some(written);
1569 record.updated_at = now();
1570 if let Some((host, size)) = selected_container_size.as_ref() {
1573 crate::database::save_session_with_container_size(
1574 &self.state.sessions[session_id],
1575 host,
1576 *size,
1577 )?;
1578 } else {
1579 crate::database::save_session(&self.state.sessions[session_id])?;
1580 }
1581 }
1582 let utility_handoff = {
1583 let _provisioning = ResumePhaseTimer::new(session_id, "provision destination");
1584 if let Some(config) = utility_config.as_ref() {
1585 Some(
1586 provision_with_cross_harness_handoff(
1587 self,
1588 session_id,
1589 executor,
1590 github_token.as_deref(),
1591 config,
1592 &canonical_session,
1593 context_bytes,
1594 )
1595 .context("prepare the cross-harness destination")?,
1596 )
1597 } else {
1598 self.provision_session_with_failure_disposition(
1599 session_id,
1600 executor,
1601 github_token.as_deref(),
1602 ProvisioningFailureDisposition::Preserve,
1603 )
1604 .await?;
1605 None
1606 }
1607 };
1608 let restore_repositories = (resumed_project_directory.is_none()
1609 && conversion.is_none())
1610 || plan == ResumePlan::RawToWorkspace
1611 || (recreated_managed_worktree && plan == ResumePlan::InPlace);
1612 let restored_archive = conversion_checkpoint_written
1615 .as_ref()
1616 .map_or(archive_path.as_path(), |checkpoint| {
1617 checkpoint.archive_path.as_path()
1618 });
1619 self.restore_into_target(
1620 session_id,
1621 RestoreIntoTarget {
1622 profile: &profile,
1623 archive: &verified_archive,
1624 restored_archive,
1625 resumed_project_directory,
1626 resumed_container_workspace,
1627 restore_repositories,
1628 primary_repository_root_from_conversion: conversion.is_some(),
1629 native_continuity,
1630 discard_queued_prompts,
1631 replay_queue: !discard_queue,
1632 utility_handoff,
1633 projection_build,
1634 resume_notices,
1635 install_attached_resources: true,
1636 worker_root_reset: WorkerRootReset::FreshTarget,
1637 retire_after_ready: conversion
1638 .as_ref()
1639 .and_then(ResumeConversion::raw_to_workspace)
1640 .and_then(|plan| plan.retire.as_ref()),
1641 },
1642 executor,
1643 )
1644 .await
1645 }
1646 .await;
1647 match result {
1648 Ok(materialized) => {
1649 if let Some(written) = &conversion_checkpoint_written {
1652 super::checkpoint::prune_replaced_checkpoint(
1653 previous.checkpoint.as_ref(),
1654 written,
1655 );
1656 }
1657 Ok(materialized)
1658 }
1659 Err(error) => {
1660 if let Some(written) = &conversion_checkpoint_written
1663 && let Err(remove_error) = std::fs::remove_file(&written.archive_path)
1664 && remove_error.kind() != std::io::ErrorKind::NotFound
1665 {
1666 tracing::warn!(
1667 session_id,
1668 path = %written.archive_path.display(),
1669 "could not remove the conversion checkpoint after resume failed: {remove_error}"
1670 );
1671 }
1672 restore_projection_after_failed_resume(
1675 session_id,
1676 &canonical_session,
1677 discard_queued_prompts,
1678 );
1679 Err(self.rollback_failed_resume(
1680 session_id,
1681 &previous,
1682 recreated_managed_worktree,
1683 error,
1684 executor,
1685 )?)
1686 }
1687 }
1688 }
1689
1690 pub(super) fn rollback_failed_resume(
1691 &mut self,
1692 session_id: &str,
1693 previous: &SessionRecord,
1694 recreated_managed_worktree: bool,
1695 error: anyhow::Error,
1696 _executor: &impl CommandExecutor,
1697 ) -> Result<anyhow::Error> {
1698 let current = self
1699 .state
1700 .sessions
1701 .get(session_id)
1702 .with_context(|| format!("unknown session {session_id}"))?
1703 .clone();
1704 let cleanup = match current.target.as_ref() {
1705 Some(locator) => (|| -> Result<()> {
1706 let backend = backend_locator(locator, ¤t, &self.config)?;
1707 targets::close_plan(&backend, session_id)?
1708 .execute(&CancellableProcessExecutor::with_timeout(
1711 Duration::from_secs(15),
1712 ))
1713 .map(|_| ())
1714 })(),
1715 None => Ok(()),
1716 };
1717 let worktree_cleanup = if cleanup.is_err() {
1720 Ok(())
1721 } else {
1722 match (
1723 current.managed_worktree.as_ref(),
1724 previous.managed_worktree.as_ref(),
1725 ) {
1726 (_, Some(previous)) if recreated_managed_worktree => retire_managed_worktree(
1727 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1728 previous,
1729 ),
1730 (Some(current), Some(previous)) if current == previous => Ok(()),
1731 (Some(worktree), _) => cleanup_managed_worktree(
1734 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
1735 worktree,
1736 crate::controller::BranchDisposition::Delete,
1737 ),
1738 (None, _) => Ok(()),
1739 }
1740 };
1741 let cleanup_error = [cleanup, worktree_cleanup]
1742 .into_iter()
1743 .filter_map(Result::err)
1744 .map(|cleanup_error| format!("{cleanup_error:#}"))
1745 .collect::<Vec<_>>()
1746 .join("; ");
1747 if !cleanup_error.is_empty() {
1748 tracing::warn!(
1749 session_id,
1750 error = %cleanup_error,
1751 "resume rollback cleanup reported failures"
1752 );
1753 }
1754 let original = super::worker_binary::failure_line(&error);
1757 let detail = format!("{error:#}");
1758 if detail != original {
1759 tracing::warn!(session_id, error = %detail, "resume failed");
1760 }
1761 let record = self.state.sessions.get_mut(session_id).unwrap();
1762 let failure = apply_failed_resume_rollback(
1763 record,
1764 previous,
1765 &original,
1766 (!cleanup_error.is_empty()).then_some(cleanup_error),
1767 );
1768 if record.bundle_id != current.bundle_id {
1771 let bundle_id = record.bundle_id.clone();
1772 crate::database::rebind_session_bundle(session_id, &bundle_id)?;
1773 }
1774 crate::database::save_session(&self.state.sessions[session_id])?;
1777 Ok(failure)
1778 }
1779}
1780
1781fn worktree_cleanup_notice(worktree_root: &Path, error: &anyhow::Error) -> String {
1782 format!(
1783 "Mjolnir could not remove the worktree at {}: {error:#}. Remove it with `git worktree remove --force {}`.",
1784 worktree_root.display(),
1785 worktree_root.display()
1786 )
1787}
1788
1789pub(super) fn apply_failed_resume_rollback(
1790 current: &mut SessionRecord,
1791 previous: &SessionRecord,
1792 original_error: &str,
1793 cleanup_error: Option<String>,
1794) -> anyhow::Error {
1795 match cleanup_error {
1796 None => {
1797 *current = previous.clone();
1798 current.state = SessionState::Stopped;
1799 current.target = None;
1800 current.updated_at = now();
1801 current.last_error = Some(format!("resume failed: {original_error}"));
1802 anyhow::anyhow!(original_error.to_owned())
1803 }
1804 Some(cleanup_error) => {
1805 let failure = format!(
1806 "{original_error}; cleanup of the partial resume target failed: {cleanup_error}"
1807 );
1808 if current.managed_worktree.is_none() {
1813 current
1814 .project_directory
1815 .clone_from(&previous.project_directory);
1816 current
1817 .managed_worktree
1818 .clone_from(&previous.managed_worktree);
1819 current.bundle_id.clone_from(&previous.bundle_id);
1820 }
1821 current.state = SessionState::Error;
1822 current.updated_at = now();
1823 current.last_error = Some(format!("resume failed: {failure}"));
1824 anyhow::anyhow!(failure)
1825 }
1826 }
1827}
1828
1829fn conversion_notice(
1831 target_id: &str,
1832 checkout: &Path,
1833 branch: Option<&str>,
1834 retire: Option<&mj_core::state::ManagedWorktree>,
1835) -> String {
1836 let branch = branch.unwrap_or("a detached head");
1837 match retire {
1838 Some(worktree) => format!(
1839 "This session moved out of {} and into the {target_id} target, where its checkout is on {branch}. Its branch {} stays in {}.",
1840 checkout.display(),
1841 worktree.branch,
1842 worktree.source_repository.display()
1843 ),
1844 None => format!(
1845 "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.",
1846 checkout.display()
1847 ),
1848 }
1849}
1850
1851fn conversion_checkpoint(
1859 previous_archive: &Path,
1860 snapshot: mj_checkpoint::archive::RepositorySnapshot,
1861 output: &Path,
1862) -> Result<mj_core::state::CheckpointMetadata> {
1863 let previous = mj_checkpoint::archive::read_archive_verified(previous_archive)
1864 .with_context(|| format!("read checkpoint archive {}", previous_archive.display()))?;
1865 let native_artifacts = previous
1868 .manifest
1869 .payloads
1870 .iter()
1871 .filter_map(|descriptor| match &descriptor.role {
1872 mj_checkpoint::archive::PayloadRole::NativeArtifact { relative_path } => {
1873 Some((relative_path, descriptor))
1874 }
1875 _ => None,
1876 })
1877 .map(|(relative_path, descriptor)| {
1878 Ok(mj_checkpoint::archive::NativeArtifact {
1879 relative_path: relative_path.clone(),
1880 data: previous.payload(descriptor)?.to_vec(),
1881 mode: descriptor.mode,
1882 })
1883 })
1884 .collect::<Result<Vec<_>>>()?;
1885 let canonical_session = previous.canonical_session()?;
1886 let event_frontier = canonical_session.event_frontier;
1887 let written = mj_checkpoint::archive::write_archive_atomic(
1888 output,
1889 &mj_checkpoint::archive::ArchiveInput {
1890 session: previous.manifest.session.clone(),
1891 target: previous.manifest.target.clone(),
1894 bundle: mj_checkpoint::archive::BundleManifest {
1895 id: previous.manifest.bundle.id.clone(),
1896 primary_repository: snapshot.metadata.id.clone(),
1900 },
1901 canonical_session,
1902 native_artifacts,
1903 repositories: vec![snapshot],
1904 },
1905 )
1906 .with_context(|| format!("write the conversion archive {}", output.display()))?;
1907 Ok(mj_core::state::CheckpointMetadata {
1908 archive_path: output.to_path_buf(),
1909 sha256: written.archive_sha256,
1910 created_at: now(),
1911 event_frontier,
1912 })
1913}
1914
1915fn restored_native_session_accepted(
1926 archived: &str,
1927 opened: &str,
1928 replaced_unused: Option<&str>,
1929) -> bool {
1930 opened == archived || replaced_unused == Some(archived)
1931}
1932
1933fn restore_projection_after_failed_resume(
1944 session_id: &str,
1945 canonical_session: &mj_checkpoint::archive::CanonicalSessionSnapshot,
1946 discard_queued_prompts: bool,
1947) {
1948 let stored_frontier = crate::database::materialized_event_frontier(session_id)
1949 .unwrap_or_else(|error| {
1950 tracing::warn!(
1951 session_id,
1952 error = format!("{error:#}"),
1953 "could not read the stored projection frontier after a failed resume; rebuilding it from the archive"
1954 );
1955 None
1956 });
1957 if projection_rebuild_required(
1958 stored_frontier
1959 .as_ref()
1960 .map(|(ordinal, digest)| (*ordinal, digest.as_str())),
1961 canonical_session.event_frontier,
1962 &canonical_session.event_frontier_digest,
1963 ) {
1964 match materialized_session_from_canonical(session_id, canonical_session) {
1965 Ok(previous_projection) => {
1966 if let Err(restore_error) =
1967 crate::database::save_materialized_session(&previous_projection)
1968 {
1969 tracing::error!(
1970 session_id,
1971 error = format!("{restore_error:#}"),
1972 "could not restore the durable projection after a failed resume"
1973 );
1974 }
1975 }
1976 Err(restore_error) => tracing::error!(
1977 session_id,
1978 error = format!("{restore_error:#}"),
1979 "could not rebuild the durable projection after a failed resume"
1980 ),
1981 }
1982 } else if discard_queued_prompts
1983 && let Err(restore_error) = crate::database::replace_materialized_queued_prompts(
1984 session_id,
1985 &mj_transcript::projection::materialized_queued_prompts_from_canonical(
1986 &canonical_session.queued_prompts,
1987 ),
1988 )
1989 {
1990 tracing::error!(
1991 session_id,
1992 error = format!("{restore_error:#}"),
1993 "could not restore queued prompts after a failed resume"
1994 );
1995 }
1996}
1997
1998fn projection_rebuild_required(
2006 stored: Option<(u64, &str)>,
2007 archive_frontier: u64,
2008 archive_frontier_digest: &str,
2009) -> bool {
2010 stored != Some((archive_frontier, archive_frontier_digest))
2011}
2012
2013fn restore_archive_path(
2014 backend: &targets::TargetLocator,
2015 verified_archive: &Path,
2016 remote_archive: &Path,
2017) -> PathBuf {
2018 if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
2019 verified_archive.to_path_buf()
2020 } else {
2021 remote_archive.to_path_buf()
2022 }
2023}
2024
2025fn should_upload_restore_archive(backend: &targets::TargetLocator) -> bool {
2026 !matches!(backend, targets::TargetLocator::LocalBare { .. })
2027}
2028
2029struct CrossHarnessProvisionExecutor<'a, E: CommandExecutor + ?Sized> {
2034 inner: &'a E,
2035 cancellation: CancellationToken,
2036}
2037
2038impl<E: CommandExecutor + ?Sized> CommandExecutor for CrossHarnessProvisionExecutor<'_, E> {
2039 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2040 if self.cancellation.is_cancelled() {
2041 bail!("operation cancelled while provisioning destination");
2042 }
2043 self.inner.execute(command)
2044 }
2045
2046 fn cancellation_requested(&self) -> bool {
2047 self.cancellation.is_cancelled() || self.inner.cancellation_requested()
2048 }
2049
2050 fn stage_started(&self, stage: ProvisionStage) {
2051 self.inner.stage_started(stage);
2052 }
2053
2054 fn stage_finished(&self, stage: ProvisionStage) {
2055 self.inner.stage_finished(stage);
2056 }
2057
2058 fn notify_notice(&self, notice: &str) {
2059 self.inner.notify_notice(notice);
2060 }
2061
2062 fn execute_with_stdin(
2063 &self,
2064 command: &CommandSpec,
2065 input: &mut (dyn std::io::Read + Send),
2066 ) -> Result<CommandOutput> {
2067 if self.cancellation.is_cancelled() {
2068 bail!("operation cancelled while provisioning destination");
2069 }
2070 self.inner.execute_with_stdin(command, input)
2071 }
2072}
2073
2074fn provision_with_cross_harness_handoff(
2075 controller: &mut Controller,
2076 session_id: &str,
2077 executor: &(impl CommandExecutor + Sync),
2078 github_token: Option<&str>,
2079 config: &Config,
2080 snapshot: &CanonicalSessionSnapshot,
2081 context_bytes: usize,
2082) -> Result<String> {
2083 let (_provision, handoff) = execute_joined_cross_harness_work(
2084 "cross-harness provisioning",
2085 move |cancellation| {
2086 let provision_executor = CrossHarnessProvisionExecutor {
2087 inner: executor,
2088 cancellation,
2089 };
2090 futures::executor::block_on(controller.provision_session_with_failure_disposition(
2091 session_id,
2092 &provision_executor,
2093 github_token,
2094 ProvisioningFailureDisposition::Preserve,
2095 ))
2096 },
2097 "cross-harness handoff",
2098 move |cancellation| {
2099 let runtime = tokio::runtime::Builder::new_current_thread()
2100 .enable_all()
2101 .build()
2102 .context("create cross-harness handoff runtime")?;
2103 runtime.block_on(utility_handoff_while_cancellable(
2104 session_id,
2105 config,
2106 snapshot,
2107 context_bytes,
2108 executor,
2109 cancellation,
2110 ))
2111 },
2112 )?;
2113 ensure!(
2114 !executor.cancellation_requested(),
2115 "operation cancelled while provisioning destination"
2116 );
2117 Ok(handoff)
2118}
2119
2120fn execute_joined_cross_harness_work<A: Send, B: Send>(
2124 first_name: &'static str,
2125 first: impl FnOnce(CancellationToken) -> Result<A> + Send,
2126 second_name: &'static str,
2127 second: impl FnOnce(CancellationToken) -> Result<B> + Send,
2128) -> Result<(A, B)> {
2129 let cancellation = CancellationToken::new();
2130 std::thread::scope(|scope| {
2131 let first_cancel = cancellation.clone();
2132 let mut first_handle = Some(scope.spawn(move || first(first_cancel)));
2133 let second_cancel = cancellation.clone();
2134 let mut second_handle = Some(scope.spawn(move || second(second_cancel)));
2135 let mut first_result = None;
2136 let mut second_result = None;
2137
2138 while first_result.is_none() || second_result.is_none() {
2139 if first_result.is_none()
2140 && first_handle
2141 .as_ref()
2142 .is_some_and(|handle| handle.is_finished())
2143 {
2144 let handle = first_handle.take().expect("first lane handle present");
2145 first_result = Some(match handle.join() {
2146 Ok(result) => result,
2147 Err(panic) => {
2148 cancellation.cancel();
2149 Err(anyhow::anyhow!(
2150 "{first_name} thread panicked: {}",
2151 targets::command_thread_panic_message(panic.as_ref())
2152 ))
2153 }
2154 });
2155 if first_result.as_ref().is_some_and(Result::is_err) {
2156 cancellation.cancel();
2157 }
2158 }
2159 if second_result.is_none()
2160 && second_handle
2161 .as_ref()
2162 .is_some_and(|handle| handle.is_finished())
2163 {
2164 let handle = second_handle.take().expect("second lane handle present");
2165 second_result = Some(match handle.join() {
2166 Ok(result) => result,
2167 Err(panic) => {
2168 cancellation.cancel();
2169 Err(anyhow::anyhow!(
2170 "{second_name} thread panicked: {}",
2171 targets::command_thread_panic_message(panic.as_ref())
2172 ))
2173 }
2174 });
2175 if second_result.as_ref().is_some_and(Result::is_err) {
2176 cancellation.cancel();
2177 }
2178 }
2179 if first_result.is_none() || second_result.is_none() {
2180 std::thread::sleep(Duration::from_millis(10));
2181 }
2182 }
2183
2184 match (
2185 first_result.expect("first lane result received after joined handle"),
2186 second_result.expect("second lane result received after joined handle"),
2187 ) {
2188 (Err(first), Err(second)) => {
2189 Err(first.context(format!("{second_name} lane also failed: {second:#}")))
2190 }
2191 (Err(error), Ok(_)) => Err(error),
2192 (Ok(_), Err(error)) => Err(error),
2193 (Ok(first), Ok(second)) => Ok((first, second)),
2194 }
2195 })
2196}
2197
2198fn native_continuity_preserved(profile_kind: HarnessKind, archived_kind: HarnessKind) -> bool {
2202 profile_kind == archived_kind
2203}
2204
2205async fn utility_handoff_while_cancellable(
2209 session_id: &str,
2210 config: &Config,
2211 snapshot: &CanonicalSessionSnapshot,
2212 context_bytes: usize,
2213 executor: &impl CommandExecutor,
2214 cancellation: CancellationToken,
2215) -> Result<String> {
2216 let _phase = ResumePhaseTimer::new(session_id, "cross-harness handoff");
2217 if executor.cancellation_requested() {
2218 bail!("operation cancelled while compacting the cross-harness handoff");
2219 }
2220 let _compacting = ProvisionStageGuard::new(executor, ProvisionStage::Compacting);
2221 let cancel = cancellation.child_token();
2222 let operation =
2223 crate::handoff::build_handoff_context(session_id, config, snapshot, context_bytes, &cancel);
2224 tokio::pin!(operation);
2225 loop {
2226 tokio::select! {
2227 context = &mut operation => return context,
2228 _ = cancellation.cancelled() => {
2229 cancel.cancel();
2230 bail!("operation cancelled while compacting the cross-harness handoff");
2231 }
2232 _ = tokio::time::sleep(super::readiness::CANCELLATION_POLL_INTERVAL) => {
2233 if executor.cancellation_requested() {
2234 cancel.cancel();
2235 bail!("operation cancelled while compacting the cross-harness handoff");
2236 }
2237 }
2238 }
2239 }
2240}
2241
2242mod in_place;
2243
2244#[cfg(test)]
2245mod tests;