1mod handoff;
4#[cfg(test)]
5mod tests;
6mod transfer;
7
8use anyhow::{Context, Result, bail, ensure};
9use mj_core::hex::lower_hex;
10use sha2::{Digest, Sha256};
11
12use super::lifecycle::SourceTargetDisposition;
13use super::{Controller, SessionResumeOptions, now};
14
15#[derive(Clone, Copy)]
17enum MoveOwnership {
18 Executing,
19 ExecutingQueue,
20 PendingQueue,
21}
22
23fn move_ownership() -> &'static std::sync::Mutex<std::collections::BTreeMap<String, MoveOwnership>>
24{
25 static OWNER: std::sync::OnceLock<
26 std::sync::Mutex<std::collections::BTreeMap<String, MoveOwnership>>,
27 > = std::sync::OnceLock::new();
28 OWNER.get_or_init(Default::default)
29}
30
31pub fn move_owns_session(session_id: &str) -> bool {
32 move_ownership()
33 .lock()
34 .unwrap_or_else(std::sync::PoisonError::into_inner)
35 .contains_key(session_id)
36}
37
38pub fn move_has_pending_queue(session_id: &str) -> bool {
39 matches!(
40 move_ownership()
41 .lock()
42 .unwrap_or_else(std::sync::PoisonError::into_inner)
43 .get(session_id),
44 Some(MoveOwnership::ExecutingQueue | MoveOwnership::PendingQueue)
45 )
46}
47
48pub fn release_move_queue_hold(session_id: &str) {
49 set_move_queue_hold(session_id, false);
50}
51
52fn set_move_queue_hold(session_id: &str, pending: bool) {
53 let mut owner = move_ownership()
54 .lock()
55 .unwrap_or_else(std::sync::PoisonError::into_inner);
56 let next = match (owner.get(session_id), pending) {
57 (Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue), true) => {
58 Some(MoveOwnership::ExecutingQueue)
59 }
60 (Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue), false) => {
61 Some(MoveOwnership::Executing)
62 }
63 (_, true) => Some(MoveOwnership::PendingQueue),
64 (_, false) => None,
65 };
66 if let Some(next) = next {
67 owner.insert(session_id.to_owned(), next);
68 } else {
69 owner.remove(session_id);
70 }
71}
72
73pub fn restore_move_queue_hold(operation: &MoveOperation) {
74 set_move_queue_hold(
75 &operation.selection.session_id,
76 operation.queue_admission_started && !operation.queue_admission_finished,
77 );
78}
79
80fn interrupted_source_stop_message(recovered: &str, in_place: bool) -> String {
84 if in_place {
85 format!("{recovered}; the in-place swap was interrupted; the environment was retained")
86 } else {
87 recovered.to_owned()
88 }
89}
90
91fn source_stopped_with_verified_checkpoint(record: &mj_core::state::SessionRecord) -> bool {
97 record.state == SessionState::Stopped && record.checkpoint.is_some()
98}
99
100fn stopped_source_recovery(
106 session_id: &str,
107 destination_profile: Option<&str>,
108 destination_target: Option<&str>,
109) -> String {
110 let flag = |name: &str, value: Option<&str>| {
111 value
112 .filter(|value| !value.is_empty())
113 .map(|value| format!(" --{name} {value}"))
114 .unwrap_or_default()
115 };
116 format!(
117 "Source is stopped with a verified checkpoint. Bring it back with \
118 `mj resume --session {session_id}{}{} --queue start`.",
119 flag("profile", destination_profile),
120 flag("target", destination_target),
121 )
122}
123
124fn failed_move_recovery(
133 operation: &MoveOperation,
134 record: Option<&mj_core::state::SessionRecord>,
135) -> String {
136 if operation.queue_admission_started {
137 return "Destination is live; retry queue admission on this same destination. Already accepted work may have effects.".to_owned();
138 }
139 let source_live = record.is_some_and(|record| {
140 matches!(
141 record.state,
142 SessionState::Running | SessionState::Disconnected
143 )
144 });
145 let environment_held = operation.holds_source_environment()
146 || (operation.in_place
147 && !source_live
148 && record.is_some_and(|record| record.target.is_some()));
149 if environment_held && !operation.checkpoint_retained() {
150 return missing_move_checkpoint_recovery(record);
151 }
152 if operation.in_place && record.is_some_and(|record| record.target.is_some()) {
153 return if source_live {
154 "Source retained and still running. Retry Move when ready.".to_owned()
155 } else {
156 "Environment and checkpoint retained. Retry Move on the same target; the checkout will not be recreated.".to_owned()
157 };
158 }
159 match record {
160 Some(record) if source_stopped_with_verified_checkpoint(record) => {
161 stopped_source_recovery(
162 &operation.selection.session_id,
163 operation.selection.profile_id.as_deref(),
164 operation.selection.target_template_id.as_deref(),
165 )
166 }
167 _ => "Source or partial destination is retained. Retry move after resolving the reported error.".to_owned(),
168 }
169}
170
171fn missing_move_checkpoint_recovery(record: Option<&mj_core::state::SessionRecord>) -> String {
176 let checkout = record
177 .and_then(|record| {
178 record
179 .managed_worktree
180 .as_ref()
181 .map(|checkout| checkout.worktree_root.clone())
182 .or_else(|| record.project_directory.clone())
183 })
184 .map(|path| format!(" in {}", path.display()))
185 .unwrap_or_default();
186 format!(
187 "The Move's checkpoint archive is missing, so the Move cannot be retried and Resume cannot restore the session. \
188 The environment and checkout{checkout} are retained; copy out anything you need, then destroy the session."
189 )
190}
191
192fn failed_move_message(
195 phase: Option<&str>,
196 cancelled: bool,
197 recovery: &str,
198 operation_id: &str,
199) -> String {
200 format!(
201 "{}{}{}. {recovery} The daemon log records the reason under reference {operation_id}",
202 mj_core::state::MOVE_FAILURE_PREFIX,
203 phase
204 .map(|phase| format!(" while {phase}"))
205 .unwrap_or_default(),
206 if cancelled { " (cancelled)" } else { "" },
207 )
208}
209
210pub(crate) fn forget_missing_move_archives(operation: &mut MoveOperation) -> bool {
218 let missing = |checkpoint: &Option<mj_core::state::CheckpointMetadata>| {
219 checkpoint.as_ref().is_some_and(|checkpoint| {
220 matches!(
221 std::fs::symlink_metadata(&checkpoint.archive_path),
222 Err(error) if error.kind() == std::io::ErrorKind::NotFound
223 )
224 })
225 };
226 let lost = if missing(&operation.handoff) {
227 [operation.handoff.take(), operation.checkpoint.take()]
228 } else if operation.handoff.is_none() && missing(&operation.checkpoint) {
229 [None, operation.checkpoint.take()]
230 } else {
231 return false;
232 };
233 for checkpoint in lost.iter().flatten() {
234 tracing::warn!(
235 session_id = %operation.selection.session_id,
236 reference = %operation.operation_id,
237 path = %checkpoint.archive_path.display(),
238 "a Move's checkpoint archive is missing; recording that the Move cannot restore it"
239 );
240 }
241 true
242}
243
244pub(crate) fn record_missing_move_archives(
248 state: &mj_core::state::State,
249 operations: Vec<MoveOperation>,
250) -> Result<()> {
251 for mut operation in operations {
252 if !operation.retains_checkpoint() || !forget_missing_move_archives(&mut operation) {
253 continue;
254 }
255 let record = state.sessions.get(&operation.selection.session_id);
256 let published = record.and_then(|record| record.last_error.as_deref());
257 if operation.is_active()
258 || !published
259 .is_some_and(|error| error.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
260 {
261 crate::database::save_move_operation(&operation)?;
262 continue;
263 }
264 let message = failed_move_message(
265 None,
266 operation.phase == MovePhase::Cancelled,
267 &failed_move_recovery(&operation, record),
268 &operation.operation_id,
269 );
270 operation.updated_at = now();
271 crate::database::save_move_outcome(&operation, Some(&message))?;
272 }
273 Ok(())
274}
275
276fn sealed_selection_difference(retained: &MoveSelection, requested: &MoveSelection) -> String {
279 let mut parts = Vec::new();
280 for (name, flag, retained, requested) in [
281 (
282 "profile",
283 "--profile",
284 &retained.profile_id,
285 &requested.profile_id,
286 ),
287 (
288 "target",
289 "--target",
290 &retained.target_template_id,
291 &requested.target_template_id,
292 ),
293 ] {
294 if retained != requested {
295 parts.push(format!(
296 "{name} (sealed with {0}; pass {flag} {0})",
297 retained.as_deref().unwrap_or_default()
298 ));
299 }
300 }
301 if retained.workspace.acknowledge_large_transfer
302 != requested.workspace.acknowledge_large_transfer
303 {
304 parts.push(if retained.workspace.acknowledge_large_transfer {
305 "large-transfer acknowledgement (the Move was sealed with it; pass --allow-large-transfer)".to_owned()
306 } else {
307 "large-transfer acknowledgement (the Move was sealed without it; omit --allow-large-transfer)".to_owned()
308 });
309 }
310 if retained.workspace.exclusions != requested.workspace.exclusions {
311 let excluded = retained
312 .workspace
313 .exclusions
314 .iter()
315 .map(|path| format!("{}:{}", path.repository, path.path.display()))
316 .collect::<Vec<_>>();
317 parts.push(if excluded.is_empty() {
318 "excluded files (the Move was sealed excluding none)".to_owned()
319 } else {
320 format!(
321 "excluded files (the Move was sealed excluding exactly {})",
322 excluded.join(", ")
323 )
324 });
325 }
326 if retained.additional_mounts != requested.additional_mounts {
327 parts.push("attached directories (send the ones the Move was sealed with)".to_owned());
328 }
329 if retained.resource_allocation != requested.resource_allocation
330 || retained.clear_resource_allocation != requested.clear_resource_allocation
331 {
332 parts.push("resource allocation (send the one the Move was sealed with)".to_owned());
333 }
334 if retained.subagents != requested.subagents {
335 parts.push("sub-agent policy (send the one the Move was sealed with)".to_owned());
336 }
337 if retained.session_id != requested.session_id || parts.is_empty() {
338 parts.push("the session".to_owned());
339 }
340 format!(
341 "a sealed Move must be retried with the selection it was sealed with; this request differs in its {}",
342 parts.join("; ")
343 )
344}
345
346pub struct MoveMutationGuard(String);
347
348impl MoveMutationGuard {
349 pub fn reserve(session_id: &str) -> Result<Self> {
350 let mut owner = move_ownership()
351 .lock()
352 .unwrap_or_else(std::sync::PoisonError::into_inner);
353 ensure!(
354 !matches!(
355 owner.get(session_id),
356 Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue)
357 ),
358 "session already has a move owner"
359 );
360 let next = if owner.contains_key(session_id) {
361 MoveOwnership::ExecutingQueue
362 } else {
363 MoveOwnership::Executing
364 };
365 owner.insert(session_id.to_owned(), next);
366 Ok(Self(session_id.to_owned()))
367 }
368}
369
370impl Drop for MoveMutationGuard {
371 fn drop(&mut self) {
372 let mut owner = move_ownership()
373 .lock()
374 .unwrap_or_else(std::sync::PoisonError::into_inner);
375 match owner.get(&self.0) {
376 Some(MoveOwnership::ExecutingQueue) => {
377 owner.insert(self.0.clone(), MoveOwnership::PendingQueue);
378 }
379 Some(MoveOwnership::Executing) => {
380 owner.remove(&self.0);
381 }
382 _ => {}
383 }
384 }
385}
386
387pub(crate) fn move_refuses_command(session_id: &str, command: &RelayCommand) -> bool {
388 move_owns_session(session_id)
389 && matches!(
390 command,
391 RelayCommand::Prompt { .. }
392 | RelayCommand::SetConfig { .. }
393 | RelayCommand::SetSessionMode { .. }
394 | RelayCommand::RestoreExecutionMode
395 | RelayCommand::RunUserShell { .. }
396 | RelayCommand::CancelUserShell { .. }
397 | RelayCommand::Cancel
398 | RelayCommand::RemoveQueuedPrompt { .. }
399 | RelayCommand::ClearQueuedPrompts
400 )
401}
402use crate::session_manager::{SessionManagerControl, StandaloneSession, new_command_id};
403use mj_checkpoint::archive::{CanonicalQueuedCommandKind, verify_archive_streaming};
404use mj_core::state::{MoveOperation, MovePhase, ResumeQueueDisposition, SessionState};
405
406pub use mj_core::state::{MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest};
407
408use crate::targets::{CommandExecutor, ProvisionStage, ProvisionStageGuard};
409use mj_core::relay::RelayCommand;
410
411pub async fn refresh_move_source(
413 manager: &SessionManagerControl,
414 id: &str,
415) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
416 Ok(MoveSourceRelay::lease(manager, id).await?.snapshot())
417}
418
419#[derive(Default)]
425pub(in crate::controller) struct MoveSourceRelay(Option<super::checkpoint::ControllerRelayLease>);
426
427impl MoveSourceRelay {
428 pub(in crate::controller) async fn lease(
433 manager: &SessionManagerControl,
434 id: &str,
435 ) -> Result<Self> {
436 let handle = manager
437 .wait_for_session(id, std::time::Duration::from_secs(5))
438 .await?;
439 match handle.lease_connection().await {
440 Ok(lease) => Ok(Self(Some(
441 super::checkpoint::ControllerRelayLease::Managed {
442 handle,
443 lease: Some(lease),
444 },
445 ))),
446 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
447 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
448 Ok(Self(None))
449 }
450 Err(error) => Err(error),
451 }
452 }
453
454 pub(in crate::controller) fn snapshot(
456 &mut self,
457 ) -> Option<mj_core::state::ManagedSessionSnapshot> {
458 self.0
459 .as_mut()
460 .map(|relay| relay.connection_mut().snapshot())
461 }
462
463 pub(in crate::controller) async fn sync(
466 &mut self,
467 id: &str,
468 ) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
469 let Some(relay) = self.0.as_mut() else {
470 return Ok(None);
471 };
472 match relay.connection_mut().sync().await {
473 Ok(snapshot) => Ok(Some(snapshot)),
474 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
475 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
476 self.0 = None;
477 Ok(None)
478 }
479 Err(error) => Err(error),
480 }
481 }
482
483 pub(in crate::controller) fn is_held(&self) -> bool {
484 self.0.is_some()
485 }
486
487 pub(in crate::controller) fn take(
489 &mut self,
490 ) -> Option<super::checkpoint::ControllerRelayLease> {
491 self.0.take()
492 }
493
494 pub(in crate::controller) fn replace_connection(
497 &mut self,
498 connection: crate::session_manager::StandaloneSession,
499 ) {
500 if let Some(relay) = self.0.as_mut() {
501 relay.replace_connection(connection);
502 }
503 }
504}
505
506impl Drop for MoveSourceRelay {
507 fn drop(&mut self) {
508 if let Some(relay) = self.0.take() {
509 relay.release();
510 }
511 }
512}
513
514fn digest(value: &impl serde::Serialize) -> Result<String> {
515 Ok(lower_hex(Sha256::digest(serde_json::to_vec(value)?)))
516}
517
518struct MovePhaseTimer<'a> {
519 session_id: &'a str,
520 phase: &'static str,
521 started: std::time::Instant,
522}
523
524impl<'a> MovePhaseTimer<'a> {
525 fn new(session_id: &'a str, phase: &'static str) -> Self {
526 Self {
527 session_id,
528 phase,
529 started: std::time::Instant::now(),
530 }
531 }
532}
533
534impl Drop for MovePhaseTimer<'_> {
535 fn drop(&mut self) {
536 tracing::info!(
537 session_id = self.session_id,
538 phase = self.phase,
539 elapsed_ms = self.started.elapsed().as_millis() as u64,
540 "move phase finished"
541 );
542 }
543}
544
545fn replace_queued_images_with_placeholders(
549 queued_commands: &mut [mj_core::state::MaterializedQueuedPrompt],
550) {
551 for block in queued_commands
552 .iter_mut()
553 .flat_map(|command| command.content.iter_mut())
554 {
555 if block.get("type").and_then(serde_json::Value::as_str) != Some("image") {
556 continue;
557 }
558 let mime = block
559 .get("mimeType")
560 .or_else(|| block.get("mime_type"))
561 .and_then(serde_json::Value::as_str)
562 .unwrap_or("image");
563 *block = serde_json::json!({"type": "text", "text": format!("[Image attachment: {mime}]")});
564 }
565}
566
567impl Controller {
568 fn validate_move_destination_paths(
572 &self,
573 source: &mj_core::state::SessionRecord,
574 target_id: &str,
575 executor: &(impl CommandExecutor + Sync),
576 ) -> Result<Option<super::worktree::RawToWorkspaceConversion>> {
577 use super::worktree::ResumePlan;
578 match super::worktree::resume_compatibility(source, &self.config, target_id)
579 .map_err(anyhow::Error::msg)?
580 {
581 ResumePlan::RawToWorkspace => {
582 return Ok(Some(super::worktree::plan_raw_to_workspace(
583 source,
584 &self.config,
585 executor,
586 )?));
587 }
588 ResumePlan::WorkspaceToRaw => {
589 self.plan_workspace_to_raw(source, target_id, executor)?;
590 }
591 ResumePlan::InPlace if source.managed_worktree.is_none() => {
592 if let Some(path) = &source.project_directory {
593 self.validate_project_directory(target_id, path, executor)?;
594 }
595 }
596 ResumePlan::InPlace => {}
597 }
598 Ok(None)
599 }
600 fn move_confirmation(
601 &self,
602 selection: &MoveSelection,
603 conversion: Option<&mj_core::state::RawConversionPreview>,
604 ) -> Result<(bool, Vec<mj_core::state::MaterializedQueuedPrompt>, String)> {
605 let source = self
606 .state
607 .sessions
608 .get(&selection.session_id)
609 .context("unknown move session")?;
610 let (mut active, mut queued) = crate::database::move_pending_work(&source.id)?;
611 if let Some(operation) = crate::database::load_move_operation(&source.id)?
612 && operation.queue_admission_started
613 && !operation.queue_admission_finished
614 {
615 ensure!(
616 operation.selection == *selection,
617 "queue admission is incomplete on the live destination; retry that move before selecting another destination"
618 );
619 let checkpoint = operation
620 .restore_artifact()
621 .context("retained queue checkpoint is missing")?;
622 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
623 ensure!(
624 verified.archive_sha256 == checkpoint.sha256
625 && verified.manifest.session.id == source.id,
626 "retained move checkpoint verification failed"
627 );
628 queued = verified
629 .canonical_session
630 .queued_prompts
631 .into_iter()
632 .map(|entry| mj_core::state::MaterializedQueuedPrompt {
633 accepted_ordinal: None,
634 command_id: entry.command_id,
635 kind: match entry.kind {
636 CanonicalQueuedCommandKind::Prompt => {
637 mj_core::state::QueuedCommandKind::Prompt
638 }
639 CanonicalQueuedCommandKind::SetConfig { key, value } => {
640 mj_core::state::QueuedCommandKind::SetConfig { key, value }
641 }
642 },
643 content: entry.content,
644 queued_at_ms: entry.queued_at_ms,
645 })
646 .collect();
647 active = false;
648 }
649 let fingerprint = digest(&(
650 &source.last_profile,
651 &source.target_template_id,
652 &source.target,
653 &source.target_runtime,
654 &source.native_session_id,
655 &source.resource_allocation,
656 &source.additional_mounts,
657 &source.container_cpus,
658 &source.container_memory,
659 selection,
660 self.move_configuration_fingerprint(selection)?,
661 &queued,
662 conversion.map(|preview| {
666 (
667 &preview.fetch_url,
668 &preview.push_urls,
669 &preview.branch,
670 &preview.destination,
671 )
672 }),
673 ))?;
674 Ok((active, queued, fingerprint))
675 }
676 pub(super) fn move_configuration_fingerprint(
677 &self,
678 selection: &MoveSelection,
679 ) -> Result<String> {
680 let profile = self
681 .config
682 .profiles
683 .get(
684 selection
685 .profile_id
686 .as_deref()
687 .context("move profile is unresolved")?,
688 )
689 .context("move profile no longer exists")?;
690 let target = self
691 .config
692 .targets
693 .get(
694 selection
695 .target_template_id
696 .as_deref()
697 .context("move target is unresolved")?,
698 )
699 .context("move target no longer exists")?;
700 let source = self
703 .state
704 .sessions
705 .get(&selection.session_id)
706 .context("unknown move session")?;
707 digest(&(
708 profile,
709 target,
710 &self.config.bundles,
711 self.config.targets.get(&source.target_template_id),
712 self.config.profiles.get(&source.last_profile),
713 ))
714 }
715
716 async fn validate_move_destination_configuration(
717 &self,
718 selection: &MoveSelection,
719 source_harness: mj_core::config::HarnessKind,
720 operational: &mj_core::relay::RelayOperationalState,
721 ) -> Result<()> {
722 let profile_id = selection
723 .profile_id
724 .as_deref()
725 .context("move profile is unresolved")?;
726 let profile = self
727 .config
728 .profiles
729 .get(profile_id)
730 .context("move profile no longer exists")?;
731 let source = self
732 .state
733 .sessions
734 .get(&selection.session_id)
735 .context("unknown move session")?;
736 if profile.kind != source_harness || profile_id == source.last_profile {
737 return Ok(());
738 }
739 let accepted = mj_core::acp::AcceptedSessionConfig::from_configuration(
740 &operational.config,
741 &operational.config_options,
742 );
743 if accepted.model.is_none() && accepted.effort.is_none() {
744 return Ok(());
745 }
746 let choices =
749 super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
750 .await
751 .with_context(|| {
752 format!("discover destination profile {profile_id:?} configuration")
753 })?;
754 validate_preserved_configuration(profile_id, &accepted, &choices)
755 }
756
757 pub async fn prepare_move_session_controlled(
758 &self,
759 mut selection: MoveSelection,
760 executor: &(impl CommandExecutor + Sync),
761 ) -> Result<MovePreparation> {
762 ensure!(
763 selection.profile_id.is_some() || selection.target_template_id.is_some(),
764 "move requires a target or profile selection"
765 );
766 let source = self
767 .state
768 .sessions
769 .get(&selection.session_id)
770 .context("unknown session")?;
771 ensure!(
772 !self.state.subagents.contains_key(&source.id),
773 "sub-agent sessions cannot move independently of their parent"
774 );
775 let previous = crate::database::load_move_operation(&source.id)?;
776 if previous
777 .as_ref()
778 .is_some_and(|op| op.holds_source_environment() && !op.checkpoint_retained())
779 {
780 bail!(
781 "this Move cannot be retried. {}",
782 missing_move_checkpoint_recovery(Some(source))
783 );
784 }
785 let retry = previous.as_ref().is_some_and(|op| {
786 !matches!(op.phase, MovePhase::Completed)
787 && (op.restore_artifact().is_some() || source.checkpoint.is_some())
788 });
789 ensure!(
790 matches!(
791 source.state,
792 SessionState::Running | SessionState::Disconnected
793 ) || retry,
794 "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
795 source.id
796 );
797 selection
798 .profile_id
799 .get_or_insert_with(|| source.last_profile.clone());
800 selection
801 .target_template_id
802 .get_or_insert_with(|| source.target_template_id.clone());
803 selection
804 .additional_mounts
805 .get_or_insert_with(|| source.additional_mounts.clone());
806 ensure!(
807 !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
808 "select resource allocation or explicitly clear it, not both"
809 );
810 if !selection.clear_resource_allocation {
811 selection.resource_allocation = selection
812 .resource_allocation
813 .or_else(|| source.resource_allocation.clone());
814 }
815 if selection.clear_resource_allocation
816 && source.resource_allocation.is_none()
817 && source.container_cpus.is_none()
818 && source.container_memory.is_none()
819 {
820 selection.clear_resource_allocation = false;
821 }
822 if let Some(retained) = previous.as_ref().filter(|op| op.holds_source_environment()) {
823 if retained.in_place {
827 selection.workspace = retained.selection.workspace.clone();
828 }
829 ensure!(
830 retained.selection == selection,
831 "{}",
832 sealed_selection_difference(&retained.selection, &selection)
833 );
834 ensure!(
835 !retained.in_place
836 || (retained.source_target.is_some()
837 && source.target == retained.source_target),
838 "retained Move target is missing or changed; refusing to recreate it"
839 );
840 }
841 let profile_id = selection.profile_id.as_deref().unwrap();
842 let target_id = selection.target_template_id.as_deref().unwrap();
843 let profile = self
844 .config
845 .profiles
846 .get(profile_id)
847 .context("unknown destination profile")?;
848 ensure!(
849 profile.enabled,
850 "destination profile {profile_id:?} is disabled"
851 );
852 if let Some(policy) = &selection.subagents {
853 super::profile_config::validate_session_subagent_policy(
854 &self.config,
855 profile_id,
856 policy,
857 )
858 .await?;
859 }
860 let target = self
861 .config
862 .targets
863 .get(target_id)
864 .context("unknown destination target")?;
865 self.validate_muse_resume_destination(source, profile.kind, target_id)?;
866 super::worktree::resume_compatibility(source, &self.config, target_id)
867 .map_err(anyhow::Error::msg)?;
868 super::backend::validate_resource_allocation(
869 target,
870 selection.resource_allocation.as_ref(),
871 )?;
872 let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
873 ensure!(
874 profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
875 "Muse Code ACP supports one workspace root; attached directories are unsupported"
876 );
877 ensure!(
878 mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
879 "attached resources are unsupported for this target; select compatible resources explicitly"
880 );
881 crate::targets::validate_additional_mounts(mounts)?;
882 for mount in mounts {
883 self.validate_mount_source(target_id, &mount.source, executor)?;
884 }
885 let planned_conversion =
886 self.validate_move_destination_paths(source, target_id, executor)?;
887 ensure!(
888 profile.home.is_dir(),
889 "destination profile home is unavailable; configure the profile before moving"
890 );
891 super::worker_binary::preflight_worker_binary(target, executor)?;
892 super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
893 let source_harness = previous
894 .as_ref()
895 .filter(|operation| {
896 !operation.queue_admission_started && operation.phase != MovePhase::Completed
897 })
898 .and_then(|operation| operation.recovery_session.as_ref())
899 .map_or(source.harness_kind, |record| record.harness_kind);
900 let cross_harness = profile.kind != source_harness;
901 if cross_harness {
902 let cancel = tokio_util::sync::CancellationToken::new();
903 let resolve =
904 crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
905 tokio::pin!(resolve);
906 loop {
907 tokio::select! {
908 result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
909 _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
910 if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
911 }
912 }
913 }
914 }
915 let conversion = planned_conversion
918 .map(|conversion| {
919 super::worktree::raw_conversion_preview(source, &conversion, executor)
920 .context("describe the move of this checkout into the target")
921 })
922 .transpose()?
923 .map(Box::new);
924 let (active, mut queued_commands, fingerprint) =
925 self.move_confirmation(&selection, conversion.as_deref())?;
926 replace_queued_images_with_placeholders(&mut queued_commands);
927 let operation_id = previous
928 .as_ref()
929 .filter(|operation| {
930 operation.selection == selection
931 && operation.phase != MovePhase::Completed
932 && operation.restore_artifact().is_some()
933 })
934 .map(|operation| operation.operation_id.clone())
935 .unwrap_or(new_command_id("move")?);
936 let in_place = previous
937 .as_ref()
938 .is_some_and(|op| retry && op.in_place && op.selection == selection)
939 || in_place_move_eligible(
940 source,
941 &selection,
942 &self.config.targets,
943 self.state.subagents.contains_key(&source.id),
944 retry,
945 );
946 let workspace = if in_place {
947 None
948 } else {
949 Some(self.assess_move_workspace(&selection, executor)?)
950 };
951 Ok(MovePreparation {
952 workspace,
953 source_unavailable: false,
954 in_place,
955 conversion,
956 selection,
957 source_profile_id: source.last_profile.clone(),
958 source_target_template_id: source.target_template_id.clone(),
959 cross_harness,
960 active,
961 queued_commands,
962 fingerprint,
963 operation_id,
964 })
965 }
966
967 pub async fn move_session_managed_controlled(
969 &mut self,
970 request: MoveSessionRequest,
971 executor: &(impl CommandExecutor + Sync),
972 manager: &SessionManagerControl,
973 ) -> Result<MoveOutcome> {
974 let prepared = &request.preparation;
975 let id = prepared.selection.session_id.clone();
976 let started = std::time::Instant::now();
977 executor.notify_notice("Checking destination");
978 let mut source_relay = MoveSourceRelay::default();
979 let checked = {
980 let _checking_destination =
981 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
982 let mut checked = {
983 let _timing = MovePhaseTimer::new(&id, "preflight destination checks");
984 self.prepare_move_session_controlled(prepared.selection.clone(), executor)
985 .await?
986 };
987 if matches!(
992 self.state.sessions[&id].state,
993 SessionState::Running | SessionState::Disconnected
994 ) {
995 let source_harness = self.state.sessions[&id].harness_kind;
996 source_relay = {
997 let _timing = MovePhaseTimer::new(&id, "preflight source lease");
998 MoveSourceRelay::lease(manager, &id).await?
999 };
1000 let snapshot = source_relay.snapshot();
1001 let (active, queue, fingerprint) =
1002 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
1003 checked.source_unavailable = snapshot
1004 .as_ref()
1005 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
1006 checked.active = active
1007 || checked.source_unavailable
1008 || snapshot.as_ref().is_some_and(|snapshot| {
1009 let mut operational = snapshot.operational.clone();
1010 operational.queued_prompts.clear();
1011 operational.checkpoint_barrier = None;
1012 !operational.safe_to_replace(source_harness)
1013 });
1014 if let Some(snapshot) = &snapshot {
1015 let _timing = MovePhaseTimer::new(&id, "preflight destination configuration");
1016 self.validate_move_destination_configuration(
1017 &checked.selection,
1018 source_harness,
1019 &snapshot.operational,
1020 )
1021 .await?;
1022 }
1023 checked.queued_commands = queue;
1024 checked.fingerprint = fingerprint;
1025 }
1026 checked
1027 };
1028 ensure!(
1029 checked.fingerprint == prepared.fingerprint,
1030 "session, pending work, or destination configuration changed; prepare and confirm Move again"
1031 );
1032 let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
1033 let source = self.state.sessions[&id].clone();
1034 if let Some(assessment) = &checked.workspace {
1035 checked.selection.workspace.validate(assessment)?;
1036 }
1037 let old_operation = crate::database::load_move_operation(&id)?;
1038 let retry = old_operation.filter(|op| {
1039 op.selection == prepared.selection
1040 && op.phase != MovePhase::Completed
1041 && op.restore_artifact().is_some()
1042 });
1043 if retry.is_none()
1044 && source.state == SessionState::Running
1045 && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
1046 && move_environment_change(&source, &checked.selection, false).is_none()
1047 {
1048 return Ok(outcome(
1049 &prepared.operation_id,
1050 &prepared.selection,
1051 "unchanged",
1052 None,
1053 None,
1054 ));
1055 }
1056 ensure!(
1057 !checked.active || request.acknowledge_interruption,
1058 "active work will be interrupted; confirm Move again with interruption acknowledgement"
1059 );
1060 ensure!(
1061 checked.queued_commands.is_empty() || request.queue.is_some(),
1062 "pending work requires an explicit queue choice: discard or start"
1063 );
1064 let timestamp = now();
1065 let mut operation = match retry {
1066 Some(mut op) => {
1067 ensure!(
1068 !op.queue_admission_started || op.queue == queue,
1069 "queued work may already have run; retry with the original queue choice on the same destination"
1070 );
1071 crate::database::clear_move_cancellation_for_retry(&id)?;
1072 op.configuration_fingerprint =
1073 self.move_configuration_fingerprint(&checked.selection)?;
1074 op.queue = queue;
1075 op.cancellation_requested = false;
1076 op.error = None;
1077 op
1078 }
1079 None => MoveOperation {
1080 workspace_transfer: if checked.in_place {
1081 None
1082 } else {
1083 Some(
1084 self.new_workspace_transfer(
1085 &id,
1086 &prepared.operation_id,
1087 checked
1088 .workspace
1089 .clone()
1090 .context("Move workspace assessment missing")?,
1091 executor,
1092 )?,
1093 )
1094 },
1095 handoff: None,
1096 source_checkpoint_only: false,
1097 in_place: in_place_move_eligible(
1098 &source,
1099 &checked.selection,
1100 &self.config.targets,
1101 self.state.subagents.contains_key(&id),
1102 false,
1103 ),
1104 operation_id: prepared.operation_id.clone(),
1105 selection: checked.selection.clone(),
1106 source_profile_id: source.last_profile.clone(),
1107 source_target_template_id: source.target_template_id.clone(),
1108 source_target: source.target.clone(),
1109 source_native_session_id: source.native_session_id.clone(),
1110 source_additional_mounts: source.additional_mounts.clone(),
1111 source_resource_allocation: source.resource_allocation.clone(),
1112 destination_target: None,
1113 destination_native_session_id: None,
1114 destination_store_id: None,
1115 configuration_fingerprint: self
1116 .move_configuration_fingerprint(&checked.selection)?,
1117 checkpoint: None,
1118 recovery_session: None,
1119 queue,
1120 phase: MovePhase::Preparing,
1121 queue_admission_started: false,
1122 queue_admission_finished: false,
1123 cancellation_requested: false,
1124 created_at: timestamp.clone(),
1125 updated_at: timestamp,
1126 error: None,
1127 },
1128 };
1129 crate::database::save_move_operation(&operation)?;
1130 tracing::info!(
1131 session_id = id,
1132 in_place = operation.in_place,
1133 reason = move_environment_change(
1134 &source,
1135 &checked.selection,
1136 bare_targets_share_environment(
1137 &self.config.targets,
1138 &source.target_template_id,
1139 checked
1140 .selection
1141 .target_template_id
1142 .as_deref()
1143 .unwrap_or_default(),
1144 ),
1145 )
1146 .unwrap_or(if operation.in_place {
1147 "environment unchanged"
1148 } else {
1149 "source environment unavailable or previously released"
1150 }),
1151 "move environment decision"
1152 );
1153 tracing::info!(
1154 session_id = id,
1155 phase = "preflight",
1156 elapsed_ms = started.elapsed().as_millis() as u64,
1157 "move phase completed"
1158 );
1159 let result = Box::pin(self.execute_move(
1162 &mut operation,
1163 Some(&checked),
1164 executor,
1165 manager,
1166 source_relay,
1167 ))
1168 .await;
1169 self.finish_move_result(&mut operation, result, executor)
1170 }
1171
1172 fn finish_move_result(
1173 &mut self,
1174 operation: &mut MoveOperation,
1175 result: Result<()>,
1176 executor: &impl CommandExecutor,
1177 ) -> Result<MoveOutcome> {
1178 if result.is_err()
1179 && operation.workspace_transfer.is_some()
1180 && !crate::upgrade::gate().is_open()
1181 {
1182 return Ok(outcome(
1185 &operation.operation_id,
1186 &operation.selection,
1187 "interrupted",
1188 None,
1189 Some("Move will continue after the daemon upgrade".into()),
1190 ));
1191 }
1192 let session_id = operation.selection.session_id.clone();
1193 let mut last_error = self
1194 .state
1195 .sessions
1196 .get(&session_id)
1197 .context("Move outcome session is missing")?
1198 .last_error
1199 .clone();
1200 let (status, error, recovery) = match result {
1201 Ok(()) => {
1202 operation.phase = MovePhase::Completed;
1203 operation.error = None;
1204 if last_error
1205 .as_deref()
1206 .is_some_and(|error| error.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
1207 {
1208 last_error = None;
1209 }
1210 ("completed", None, None)
1211 }
1212 Err(error) => {
1213 let phase = match operation.phase {
1214 MovePhase::Preparing => "preparing the destination",
1215 MovePhase::ClosingSource => "checkpointing and suspending the source",
1216 MovePhase::ResumingDestination => "resuming the destination",
1217 MovePhase::StartingQueue => "starting the destination queue",
1218 MovePhase::Completed | MovePhase::Failed | MovePhase::Cancelled => {
1219 "recovering the move"
1220 }
1221 };
1222 let cancelled =
1223 executor.cancellation_requested() || operation.cancellation_requested;
1224 operation.phase = if cancelled {
1225 MovePhase::Cancelled
1226 } else {
1227 MovePhase::Failed
1228 };
1229 operation.cancellation_requested = cancelled;
1230 let recovery = failed_move_recovery(
1231 operation,
1232 self.state.sessions.get(&operation.selection.session_id),
1233 );
1234 let error = format!("{error:#}");
1235 tracing::warn!(
1236 %session_id,
1237 reference = %operation.operation_id,
1238 phase,
1239 cancelled,
1240 %error,
1241 "session move did not finish"
1242 );
1243 last_error = Some(failed_move_message(
1244 Some(phase),
1245 cancelled,
1246 &recovery,
1247 &operation.operation_id,
1248 ));
1249 operation.error = Some(error.clone());
1250 (
1251 if cancelled { "cancelled" } else { "failed" },
1252 Some(error),
1253 Some(recovery),
1254 )
1255 }
1256 };
1257 operation.updated_at = now();
1258 crate::database::save_move_outcome(operation, last_error.as_deref())?;
1259 let record = self
1260 .state
1261 .sessions
1262 .get_mut(&session_id)
1263 .expect("Move outcome session");
1264 record.last_error = last_error;
1265 record.updated_at = operation.updated_at.clone();
1266 Ok(outcome(
1267 &operation.operation_id,
1268 &operation.selection,
1269 status,
1270 error,
1271 recovery,
1272 ))
1273 }
1274
1275 pub async fn recover_move_managed_controlled(
1276 &mut self,
1277 mut operation: MoveOperation,
1278 executor: &(impl CommandExecutor + Sync),
1279 manager: &SessionManagerControl,
1280 ) -> Result<MoveOutcome> {
1281 let id = operation.selection.session_id.clone();
1282 let result = async {
1283 let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1284 if operation.queue_admission_started {
1285 ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1286 self.finish_workspace_transfer(&mut operation, executor)?;
1287 return self.admit_move_queue(&mut operation, executor).await;
1288 }
1289 if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1290 && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1291 if session.state == SessionState::Provisioning || session.target != operation.source_target {
1292 self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1293 }
1294 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1295 }
1296 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1297 let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1298 let sealed = operation.in_place && operation.recovery_session.is_some() && session.state == SessionState::Closing;
1304 if !sealed {
1305 Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1306 }
1307 if !operation.retains_source_environment() {
1311 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1312 }
1313 if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1314 let previous = self.state.sessions[&id].clone();
1315 operation.recovery_session = Some(previous.clone());
1316 let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1317 return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1318 }
1319 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1320 self.cleanup_stopped_target(&id, &cleanup)?;
1321 }
1322 bail!("{}", interrupted_source_stop_message(
1323 "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1324 operation.in_place,
1325 ));
1326 }
1327 match operation.phase {
1328 MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1329 MovePhase::ResumingDestination if session.state == SessionState::Running => {
1330 ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1334 && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1335 "ready destination does not match the move intent");
1336 operation.destination_target = session.target.clone();
1337 operation.destination_native_session_id = session.native_session_id.clone();
1338 operation.queue_admission_started = true;
1339 operation.phase = MovePhase::StartingQueue;
1340 crate::database::save_move_operation(&operation)?;
1341 restore_move_queue_hold(&operation);
1342 ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1343 self.finish_workspace_transfer(&mut operation, executor)?;
1344 self.admit_move_queue(&mut operation, executor).await
1345 }
1346 MovePhase::ResumingDestination => {
1347 let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1348 let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1351 let error = if operation.in_place {
1352 self.retain_failed_in_place_move(&id, previous, cause)?
1353 } else if operation.workspace_transfer.is_some() {
1354 self.rollback_move_destination(&operation, cause, executor)?
1355 } else {
1356 self.rollback_failed_resume(&id, previous, false, cause, executor)?
1357 };
1358 Err(error)
1359 }
1360 MovePhase::ClosingSource => {
1361 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1362 Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1363 }
1364 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1365 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1366 self.cleanup_stopped_target(&id, executor)?;
1367 }
1368 bail!("{}", interrupted_source_stop_message(
1369 "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1370 operation.in_place,
1371 ))
1372 }
1373 _ => bail!("Move requires an explicit retry after the daemon restarted"),
1374 }
1375 }.await;
1376 self.finish_move_result(&mut operation, result, executor)
1377 }
1378
1379 async fn recover_move_source_stop(
1380 &mut self,
1381 operation: &mut MoveOperation,
1382 executor: &(impl CommandExecutor + Sync),
1383 manager: &SessionManagerControl,
1384 ) -> Result<()> {
1385 let id = operation.selection.session_id.clone();
1386 if self.state.sessions[&id].state == SessionState::Closing {
1387 self.prepare_move_source_checkpoint(
1388 &id,
1389 executor,
1390 manager,
1391 operation,
1392 &mut MoveSourceRelay::default(),
1393 )
1394 .await?;
1395 let handle = manager
1396 .wait_for_session(&id, std::time::Duration::from_secs(5))
1397 .await?;
1398 let mut lease = handle.lease_connection().await?;
1399 let execution = lease.connection_mut().sync().await?.operational.execution;
1400 if operation.retains_source_environment()
1401 && matches!(
1402 execution,
1403 mj_core::relay::RelayExecutionState::Closing
1404 | mj_core::relay::RelayExecutionState::Closed
1405 )
1406 {
1407 super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1408 lease.release();
1409 return Ok(());
1414 }
1415 lease.release();
1416 if matches!(
1417 execution,
1418 mj_core::relay::RelayExecutionState::Idle
1419 | mj_core::relay::RelayExecutionState::Running
1420 ) {
1421 if operation.cancellation_requested || operation.retains_source_environment() {
1422 let record = self.state.sessions.get_mut(&id).unwrap();
1423 record.state = SessionState::Running;
1424 record.updated_at = now();
1425 record.last_error = Some(
1426 "Move was interrupted before the source was sealed; source retained".into(),
1427 );
1428 crate::database::save_lifecycle_session(record)?;
1429 } else {
1430 Box::pin(self.suspend_session_for_move(
1432 &id,
1433 executor,
1434 manager,
1435 operation,
1436 None,
1437 SourceTargetDisposition::Destroy,
1438 MoveSourceRelay::default(),
1439 ))
1440 .await?;
1441 }
1442 return Ok(());
1443 }
1444 }
1445 ensure!(
1446 !operation.retains_source_environment(),
1447 "retained Move source cannot be proven; refusing target teardown"
1448 );
1449 self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1451 .await?;
1452 Ok(())
1453 }
1454
1455 async fn execute_move(
1458 &mut self,
1459 operation: &mut MoveOperation,
1460 preparation: Option<&MovePreparation>,
1461 executor: &(impl CommandExecutor + Sync),
1462 manager: &SessionManagerControl,
1463 mut source_relay: MoveSourceRelay,
1464 ) -> Result<()> {
1465 let id = operation.selection.session_id.clone();
1466 ensure!(
1467 !executor.cancellation_requested(),
1468 "move cancelled before source interruption"
1469 );
1470 if !operation.queue_admission_started {
1471 if forget_missing_move_archives(operation) {
1476 operation.updated_at = now();
1477 crate::database::save_move_operation(operation)?;
1478 bail!("the Move's checkpoint archive is missing; nothing is left to restore");
1479 }
1480 if operation.in_place {
1481 ensure!(
1482 operation.source_target.is_some()
1483 && self.state.sessions[&id].target == operation.source_target,
1484 "retained Move target is missing or changed; refusing to recreate the environment"
1485 );
1486 }
1487 if self.state.sessions[&id].state == SessionState::Error
1488 && let Some(previous) = operation.recovery_session.as_ref()
1489 {
1490 let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
1493 if operation.in_place {
1494 self.retain_failed_in_place_move(&id, previous, cause)?;
1500 } else if operation.workspace_transfer.is_some() {
1501 self.rollback_move_destination(operation, cause, executor)?;
1502 let record = self.state.sessions.get_mut(&id).unwrap();
1503 record.state = SessionState::Closing;
1504 crate::database::save_resumed_session(record, None)?;
1505 } else {
1506 let failure =
1507 self.rollback_failed_resume(&id, previous, false, cause, executor)?;
1508 ensure!(
1509 self.state.sessions[&id].state == SessionState::Stopped,
1510 "{failure:#}"
1511 );
1512 }
1513 }
1514 let state = self.state.sessions[&id].state;
1515 if matches!(state, SessionState::Closing | SessionState::Destroying)
1516 && !(operation.retains_source_environment() && operation.recovery_session.is_some())
1517 {
1518 Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
1519 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1520 }
1521 if matches!(
1522 self.state.sessions[&id].state,
1523 SessionState::Running | SessionState::Disconnected
1524 ) && operation.destination_target.is_none()
1525 {
1526 executor.notify_notice("Stopping source");
1527 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
1528 operation.phase = MovePhase::ClosingSource;
1529 operation.updated_at = now();
1530 crate::database::save_move_operation(operation)?;
1531 let disposition = if operation.in_place {
1535 SourceTargetDisposition::RetainForInPlaceSwap
1536 } else {
1537 SourceTargetDisposition::Destroy
1538 };
1539 Box::pin(self.suspend_session_for_move(
1540 &id,
1541 executor,
1542 manager,
1543 operation,
1544 preparation,
1545 disposition,
1546 std::mem::take(&mut source_relay),
1547 ))
1548 .await?;
1549 }
1550 drop(source_relay);
1552 if !operation.in_place
1555 && operation.workspace_transfer.is_none()
1556 && self.state.sessions[&id].state == SessionState::Stopped
1557 && self.state.sessions[&id].target.is_some()
1558 {
1559 executor.notify_notice("Cleaning up source");
1560 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
1561 self.cleanup_stopped_target(&id, executor)?;
1562 }
1563 operation.checkpoint = operation
1564 .checkpoint
1565 .clone()
1566 .or_else(|| self.state.sessions[&id].checkpoint.clone());
1567 ensure!(
1568 operation.restore_artifact().is_some(),
1569 "move has no verified checkpoint"
1570 );
1571 ensure!(
1572 !executor.cancellation_requested(),
1573 "move cancelled after source sealing; checkpoint and remaining environment retained"
1574 );
1575 operation.phase = MovePhase::ResumingDestination;
1576 operation.recovery_session = Some(self.state.sessions[&id].clone());
1577 operation.updated_at = now();
1578 crate::database::save_move_operation(operation)?;
1579 executor.reserve_move_destination();
1580 executor.notify_notice("Preparing destination");
1581 if let Some(policy) = &operation.selection.subagents {
1585 let session = self.state.sessions.get_mut(&id).unwrap();
1586 session.subagents = Some(policy.clone());
1587 crate::database::save_resumed_session(session, None)?;
1588 }
1589 if operation.in_place {
1590 Box::pin(self.restore_session_in_place(
1594 &id,
1595 operation.selection.profile_id.as_deref().unwrap(),
1596 operation.selection.target_template_id.as_deref().unwrap(),
1597 executor,
1598 ))
1599 .await?;
1600 } else {
1601 if operation.selection.clear_resource_allocation {
1602 let session = self.state.sessions.get_mut(&id).unwrap();
1603 session.resource_allocation = None;
1604 session.container_cpus = None;
1605 session.container_memory = None;
1606 crate::database::save_resumed_session(session, None)?;
1607 }
1608 if operation.workspace_transfer.is_some() {
1609 self.capture_move_workspace(operation, executor)?;
1610 Box::pin(self.resume_session_for_move(operation, executor)).await?;
1611 operation.workspace_transfer = crate::database::load_move_operation(&id)?
1612 .context("Move intent disappeared during restore")?
1613 .workspace_transfer;
1614 } else {
1615 Box::pin(self.resume_session_controlled(
1616 &id,
1617 operation.selection.profile_id.as_deref().unwrap(),
1618 operation.selection.target_template_id.as_deref().unwrap(),
1619 SessionResumeOptions {
1620 additional_mounts: operation.selection.additional_mounts.clone(),
1621 resource_allocation: operation.selection.resource_allocation.clone(),
1622 discard_queue: true,
1623 },
1624 executor,
1625 ))
1626 .await?;
1627 }
1628 }
1629 let destination = &self.state.sessions[&id];
1630 operation.destination_target = destination.target.clone();
1631 operation.destination_native_session_id = destination.native_session_id.clone();
1632 operation.phase = MovePhase::StartingQueue;
1634 operation.queue_admission_started = true;
1635 operation.updated_at = now();
1636 crate::database::save_move_operation(operation)?;
1637 }
1638 restore_move_queue_hold(operation);
1639 self.finish_workspace_transfer(operation, executor)?;
1640 self.admit_move_queue(operation, executor).await
1641 }
1642
1643 async fn admit_move_queue(
1644 &self,
1645 operation: &mut MoveOperation,
1646 executor: &(impl CommandExecutor + Sync),
1647 ) -> Result<()> {
1648 let timing_id = operation.selection.session_id.clone();
1649 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
1650 let id = &operation.selection.session_id;
1651 let mut relay = {
1652 let _checking_destination =
1653 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1654 let destination = &self.state.sessions[id];
1655 ensure!(
1656 destination.state == SessionState::Running
1657 && destination.target == operation.destination_target
1658 && destination.native_session_id == operation.destination_native_session_id,
1659 "cannot prove the same ready destination; refusing to replay potentially executed work"
1660 );
1661 let spec = self.reconnect_command(id)?;
1662 let relay = StandaloneSession::connect_command(&spec, id).await?;
1663 let store_id = relay.snapshot().operational.store_id.context("destination worker does not expose its durable store identity; upgrade the worker before admitting queued work")?;
1664 if let Some(expected) = &operation.destination_store_id {
1665 ensure!(
1666 *expected == store_id,
1667 "destination relay storage was replaced; refusing to replay potentially executed work"
1668 );
1669 } else {
1670 operation.destination_store_id = Some(store_id);
1672 crate::database::save_move_operation(operation)?;
1673 }
1674 ensure!(
1675 relay.snapshot().operational.native_session_id
1676 == operation.destination_native_session_id,
1677 "destination relay native identity changed; refusing queue replay"
1678 );
1679 relay
1680 };
1681 if operation.queue == ResumeQueueDisposition::Start {
1682 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
1683 executor.notify_notice("Starting queued work");
1684 let checkpoint = operation
1685 .restore_artifact()
1686 .context("move queue archive is missing")?;
1687 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1688 ensure!(
1689 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
1690 "move queue checkpoint verification failed"
1691 );
1692 for queued in verified.canonical_session.queued_prompts {
1693 if operation.queue_admission_finished {
1694 continue;
1695 }
1696 ensure!(
1697 !executor.cancellation_requested(),
1698 "move cancelled during queue admission; destination retained"
1699 );
1700 let command = match queued.kind {
1701 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
1702 prompt: queued
1703 .content
1704 .into_iter()
1705 .map(serde_json::from_value)
1706 .collect::<serde_json::Result<_>>()?,
1707 },
1708 CanonicalQueuedCommandKind::SetConfig { key, value } => {
1709 RelayCommand::SetConfig { key, value }
1710 }
1711 };
1712 relay.submit_accepted(queued.command_id, command).await?;
1713 }
1714 }
1715 operation.queue_admission_finished = true;
1716 crate::database::save_move_operation(operation)?;
1717 restore_move_queue_hold(operation);
1718 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
1719 "Queued work was discarded; ready and idle."
1720 } else {
1721 "Queued work was accepted."
1722 };
1723 let source_profile = &operation.source_profile_id;
1724 let source_target = &operation.source_target_template_id;
1725 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
1726 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
1727 let text = if operation.in_place {
1730 format!(
1731 "Switched from {source_profile} / {source_target} to {destination_profile} / {destination_target} in place; the workspace and environment were kept. {queue_sentence} The interrupted prompt was not replayed."
1732 )
1733 } else {
1734 format!(
1735 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
1736 )
1737 };
1738 relay
1739 .submit(
1740 format!("{}-notice", operation.operation_id),
1741 RelayCommand::RecordNotice { text },
1742 )
1743 .await?;
1744 Ok(())
1745 }
1746
1747 pub(super) fn validate_move_checkpoint(
1748 &self,
1749 operation: &MoveOperation,
1750 preparation: Option<&MovePreparation>,
1751 executor: &(impl CommandExecutor + Sync),
1752 ) -> Result<()> {
1753 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1754 let current = Controller {
1755 config: mj_core::config::Config::load()?,
1756 state: self.state.clone(),
1757 };
1758 ensure!(
1759 current.move_configuration_fingerprint(&operation.selection)?
1760 == operation.configuration_fingerprint,
1761 "destination configuration changed during move"
1762 );
1763 let id = &operation.selection.session_id;
1764 current.validate_move_destination_paths(
1765 &self.state.sessions[id],
1766 operation.selection.target_template_id.as_deref().unwrap(),
1767 executor,
1768 )?;
1769 if !operation.in_place
1770 && operation.workspace_transfer.is_none()
1771 && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
1772 .preflight_resume_repository_sources(
1773 id,
1774 operation.selection.target_template_id.as_deref().unwrap(),
1775 executor,
1776 )?
1777 {
1778 bail!(
1779 "destination repository source is missing checkpoint commit {}; source retained",
1780 mismatch.missing_commit
1781 );
1782 }
1783 if let Some(prepared) = preparation {
1784 let checkpoint = operation
1785 .restore_artifact()
1786 .or(self.state.sessions[id].checkpoint.as_ref())
1787 .context("no move checkpoint")?;
1788 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1789 let actual: Vec<_> = verified
1790 .canonical_session
1791 .queued_prompts
1792 .iter()
1793 .map(|p| p.command_id.as_str())
1794 .collect();
1795 let expected: Vec<_> = prepared
1796 .queued_commands
1797 .iter()
1798 .map(|p| p.command_id.as_str())
1799 .collect();
1800 ensure!(
1801 actual == expected,
1802 "pending queue changed before checkpoint capture; source retained, confirm Move again"
1803 );
1804 }
1805 Ok(())
1806 }
1807}
1808
1809fn validate_preserved_configuration(
1810 profile_id: &str,
1811 accepted: &mj_core::acp::AcceptedSessionConfig,
1812 choices: &mj_core::worker_launch::ProfileConfig,
1813) -> Result<()> {
1814 for (key, value, offered) in [
1815 ("model", accepted.model.as_deref(), &choices.models),
1816 ("effort", accepted.effort.as_deref(), &choices.efforts),
1817 ] {
1818 let Some(value) = value else { continue };
1819 ensure!(
1820 offered.iter().any(|choice| choice.value == value),
1821 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
1822 offered
1823 .iter()
1824 .map(|choice| choice.value.as_str())
1825 .collect::<Vec<_>>()
1826 .join(", ")
1827 );
1828 }
1829 Ok(())
1830}
1831
1832pub(super) fn in_place_move_eligible(
1839 source: &mj_core::state::SessionRecord,
1840 selection: &MoveSelection,
1841 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
1842 is_subagent: bool,
1843 retry: bool,
1844) -> bool {
1845 !retry
1846 && !is_subagent
1847 && source.target.is_some()
1848 && matches!(
1849 source.state,
1850 SessionState::Running | SessionState::Disconnected
1851 )
1852 && move_environment_change(
1853 source,
1854 selection,
1855 bare_targets_share_environment(
1856 targets,
1857 &source.target_template_id,
1858 selection.target_template_id.as_deref().unwrap_or_default(),
1859 ),
1860 )
1861 .is_none()
1862}
1863
1864fn bare_targets_share_environment(
1869 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
1870 source_id: &str,
1871 destination_id: &str,
1872) -> bool {
1873 use mj_core::config::TargetTemplate::{LocalBare, SshBare};
1874 match (targets.get(source_id), targets.get(destination_id)) {
1875 (Some(LocalBare), Some(LocalBare)) => true,
1876 (
1877 Some(SshBare { ssh: source, .. }),
1878 Some(SshBare {
1879 ssh: destination, ..
1880 }),
1881 ) => {
1882 crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
1883 }
1884 _ => false,
1885 }
1886}
1887fn outcome(
1888 operation_id: &str,
1889 selection: &MoveSelection,
1890 status: &str,
1891 error: Option<String>,
1892 recovery: Option<String>,
1893) -> MoveOutcome {
1894 MoveOutcome {
1895 operation_id: operation_id.into(),
1896 session_id: selection.session_id.clone(),
1897 profile_id: selection.profile_id.clone().unwrap_or_default(),
1898 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
1899 outcome: status.into(),
1900 error,
1901 recovery,
1902 }
1903}
1904
1905fn move_environment_change(
1909 source: &mj_core::state::SessionRecord,
1910 selection: &MoveSelection,
1911 same_environment: bool,
1912) -> Option<&'static str> {
1913 if Some(&source.target_template_id) != selection.target_template_id.as_ref()
1914 && !same_environment
1915 {
1916 Some("target changed")
1917 } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
1918 Some("attached mounts changed")
1919 } else if source.resource_allocation != selection.resource_allocation
1920 || (selection.clear_resource_allocation
1921 && (source.container_cpus.is_some() || source.container_memory.is_some()))
1922 {
1923 Some("resource allocation changed")
1924 } else {
1925 None
1926 }
1927}