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(
127 operation: &MoveOperation,
128 record: Option<&mj_core::state::SessionRecord>,
129) -> String {
130 if operation.queue_admission_started {
131 return "Destination is live; retry queue admission on this same destination. Already accepted work may have effects.".to_owned();
132 }
133 if operation.in_place && record.is_some_and(|record| record.target.is_some()) {
134 return "Environment and checkpoint retained. Retry Move on the same target; the checkout will not be recreated.".to_owned();
135 }
136 match record {
137 Some(record) if source_stopped_with_verified_checkpoint(record) => {
138 stopped_source_recovery(
139 &operation.selection.session_id,
140 operation.selection.profile_id.as_deref(),
141 operation.selection.target_template_id.as_deref(),
142 )
143 }
144 _ => "Source or partial destination is retained. Retry move after resolving the reported error.".to_owned(),
145 }
146}
147
148pub struct MoveMutationGuard(String);
149
150impl MoveMutationGuard {
151 pub fn reserve(session_id: &str) -> Result<Self> {
152 let mut owner = move_ownership()
153 .lock()
154 .unwrap_or_else(std::sync::PoisonError::into_inner);
155 ensure!(
156 !matches!(
157 owner.get(session_id),
158 Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue)
159 ),
160 "session already has a move owner"
161 );
162 let next = if owner.contains_key(session_id) {
163 MoveOwnership::ExecutingQueue
164 } else {
165 MoveOwnership::Executing
166 };
167 owner.insert(session_id.to_owned(), next);
168 Ok(Self(session_id.to_owned()))
169 }
170}
171
172impl Drop for MoveMutationGuard {
173 fn drop(&mut self) {
174 let mut owner = move_ownership()
175 .lock()
176 .unwrap_or_else(std::sync::PoisonError::into_inner);
177 match owner.get(&self.0) {
178 Some(MoveOwnership::ExecutingQueue) => {
179 owner.insert(self.0.clone(), MoveOwnership::PendingQueue);
180 }
181 Some(MoveOwnership::Executing) => {
182 owner.remove(&self.0);
183 }
184 _ => {}
185 }
186 }
187}
188
189pub(crate) fn move_refuses_command(session_id: &str, command: &RelayCommand) -> bool {
190 move_owns_session(session_id)
191 && matches!(
192 command,
193 RelayCommand::Prompt { .. }
194 | RelayCommand::SetConfig { .. }
195 | RelayCommand::SetSessionMode { .. }
196 | RelayCommand::RestoreExecutionMode
197 | RelayCommand::RunUserShell { .. }
198 | RelayCommand::CancelUserShell { .. }
199 | RelayCommand::Cancel
200 | RelayCommand::RemoveQueuedPrompt { .. }
201 | RelayCommand::ClearQueuedPrompts
202 )
203}
204use crate::session_manager::{SessionManagerControl, StandaloneSession, new_command_id};
205use mj_checkpoint::archive::{CanonicalQueuedCommandKind, verify_archive_streaming};
206use mj_core::state::{MoveOperation, MovePhase, ResumeQueueDisposition, SessionState};
207
208pub use mj_core::state::{MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest};
209
210use crate::targets::{CommandExecutor, ProvisionStage, ProvisionStageGuard};
211use mj_core::relay::RelayCommand;
212
213pub async fn refresh_move_source(
215 manager: &SessionManagerControl,
216 id: &str,
217) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
218 Ok(MoveSourceRelay::lease(manager, id).await?.snapshot())
219}
220
221#[derive(Default)]
227pub(in crate::controller) struct MoveSourceRelay(Option<super::checkpoint::ControllerRelayLease>);
228
229impl MoveSourceRelay {
230 pub(in crate::controller) async fn lease(
235 manager: &SessionManagerControl,
236 id: &str,
237 ) -> Result<Self> {
238 let handle = manager
239 .wait_for_session(id, std::time::Duration::from_secs(5))
240 .await?;
241 match handle.lease_connection().await {
242 Ok(lease) => Ok(Self(Some(
243 super::checkpoint::ControllerRelayLease::Managed {
244 handle,
245 lease: Some(lease),
246 },
247 ))),
248 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
249 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
250 Ok(Self(None))
251 }
252 Err(error) => Err(error),
253 }
254 }
255
256 pub(in crate::controller) fn snapshot(
258 &mut self,
259 ) -> Option<mj_core::state::ManagedSessionSnapshot> {
260 self.0
261 .as_mut()
262 .map(|relay| relay.connection_mut().snapshot())
263 }
264
265 pub(in crate::controller) async fn sync(
268 &mut self,
269 id: &str,
270 ) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
271 let Some(relay) = self.0.as_mut() else {
272 return Ok(None);
273 };
274 match relay.connection_mut().sync().await {
275 Ok(snapshot) => Ok(Some(snapshot)),
276 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
277 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
278 self.0 = None;
279 Ok(None)
280 }
281 Err(error) => Err(error),
282 }
283 }
284
285 pub(in crate::controller) fn is_held(&self) -> bool {
286 self.0.is_some()
287 }
288
289 pub(in crate::controller) fn take(
291 &mut self,
292 ) -> Option<super::checkpoint::ControllerRelayLease> {
293 self.0.take()
294 }
295
296 pub(in crate::controller) fn replace_connection(
299 &mut self,
300 connection: crate::session_manager::StandaloneSession,
301 ) {
302 if let Some(relay) = self.0.as_mut() {
303 relay.replace_connection(connection);
304 }
305 }
306}
307
308impl Drop for MoveSourceRelay {
309 fn drop(&mut self) {
310 if let Some(relay) = self.0.take() {
311 relay.release();
312 }
313 }
314}
315
316fn digest(value: &impl serde::Serialize) -> Result<String> {
317 Ok(lower_hex(Sha256::digest(serde_json::to_vec(value)?)))
318}
319
320struct MovePhaseTimer<'a> {
321 session_id: &'a str,
322 phase: &'static str,
323 started: std::time::Instant,
324}
325
326impl<'a> MovePhaseTimer<'a> {
327 fn new(session_id: &'a str, phase: &'static str) -> Self {
328 Self {
329 session_id,
330 phase,
331 started: std::time::Instant::now(),
332 }
333 }
334}
335
336impl Drop for MovePhaseTimer<'_> {
337 fn drop(&mut self) {
338 tracing::info!(
339 session_id = self.session_id,
340 phase = self.phase,
341 elapsed_ms = self.started.elapsed().as_millis() as u64,
342 "move phase finished"
343 );
344 }
345}
346
347fn replace_queued_images_with_placeholders(
351 queued_commands: &mut [mj_core::state::MaterializedQueuedPrompt],
352) {
353 for block in queued_commands
354 .iter_mut()
355 .flat_map(|command| command.content.iter_mut())
356 {
357 if block.get("type").and_then(serde_json::Value::as_str) != Some("image") {
358 continue;
359 }
360 let mime = block
361 .get("mimeType")
362 .or_else(|| block.get("mime_type"))
363 .and_then(serde_json::Value::as_str)
364 .unwrap_or("image");
365 *block = serde_json::json!({"type": "text", "text": format!("[Image attachment: {mime}]")});
366 }
367}
368
369impl Controller {
370 fn validate_move_destination_paths(
374 &self,
375 source: &mj_core::state::SessionRecord,
376 target_id: &str,
377 executor: &(impl CommandExecutor + Sync),
378 ) -> Result<Option<super::worktree::RawToWorkspaceConversion>> {
379 use super::worktree::ResumePlan;
380 match super::worktree::resume_compatibility(source, &self.config, target_id)
381 .map_err(anyhow::Error::msg)?
382 {
383 ResumePlan::RawToWorkspace => {
384 return Ok(Some(super::worktree::plan_raw_to_workspace(
385 source,
386 &self.config,
387 executor,
388 )?));
389 }
390 ResumePlan::WorkspaceToRaw => {
391 self.plan_workspace_to_raw(source, target_id, executor)?;
392 }
393 ResumePlan::InPlace if source.managed_worktree.is_none() => {
394 if let Some(path) = &source.project_directory {
395 self.validate_project_directory(target_id, path, executor)?;
396 }
397 }
398 ResumePlan::InPlace => {}
399 }
400 Ok(None)
401 }
402 fn move_confirmation(
403 &self,
404 selection: &MoveSelection,
405 conversion: Option<&mj_core::state::RawConversionPreview>,
406 ) -> Result<(bool, Vec<mj_core::state::MaterializedQueuedPrompt>, String)> {
407 let source = self
408 .state
409 .sessions
410 .get(&selection.session_id)
411 .context("unknown move session")?;
412 let (mut active, mut queued) = crate::database::move_pending_work(&source.id)?;
413 if let Some(operation) = crate::database::load_move_operation(&source.id)?
414 && operation.queue_admission_started
415 && !operation.queue_admission_finished
416 {
417 ensure!(
418 operation.selection == *selection,
419 "queue admission is incomplete on the live destination; retry that move before selecting another destination"
420 );
421 let checkpoint = operation
422 .restore_artifact()
423 .context("retained queue checkpoint is missing")?;
424 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
425 ensure!(
426 verified.archive_sha256 == checkpoint.sha256
427 && verified.manifest.session.id == source.id,
428 "retained move checkpoint verification failed"
429 );
430 queued = verified
431 .canonical_session
432 .queued_prompts
433 .into_iter()
434 .map(|entry| mj_core::state::MaterializedQueuedPrompt {
435 accepted_ordinal: None,
436 command_id: entry.command_id,
437 kind: match entry.kind {
438 CanonicalQueuedCommandKind::Prompt => {
439 mj_core::state::QueuedCommandKind::Prompt
440 }
441 CanonicalQueuedCommandKind::SetConfig { key, value } => {
442 mj_core::state::QueuedCommandKind::SetConfig { key, value }
443 }
444 },
445 content: entry.content,
446 queued_at_ms: entry.queued_at_ms,
447 })
448 .collect();
449 active = false;
450 }
451 let fingerprint = digest(&(
452 &source.last_profile,
453 &source.target_template_id,
454 &source.target,
455 &source.target_runtime,
456 &source.native_session_id,
457 &source.resource_allocation,
458 &source.additional_mounts,
459 &source.container_cpus,
460 &source.container_memory,
461 selection,
462 self.move_configuration_fingerprint(selection)?,
463 &queued,
464 conversion.map(|preview| {
468 (
469 &preview.fetch_url,
470 &preview.push_urls,
471 &preview.branch,
472 &preview.destination,
473 )
474 }),
475 ))?;
476 Ok((active, queued, fingerprint))
477 }
478 pub(super) fn move_configuration_fingerprint(
479 &self,
480 selection: &MoveSelection,
481 ) -> Result<String> {
482 let profile = self
483 .config
484 .profiles
485 .get(
486 selection
487 .profile_id
488 .as_deref()
489 .context("move profile is unresolved")?,
490 )
491 .context("move profile no longer exists")?;
492 let target = self
493 .config
494 .targets
495 .get(
496 selection
497 .target_template_id
498 .as_deref()
499 .context("move target is unresolved")?,
500 )
501 .context("move target no longer exists")?;
502 let source = self
505 .state
506 .sessions
507 .get(&selection.session_id)
508 .context("unknown move session")?;
509 digest(&(
510 profile,
511 target,
512 &self.config.bundles,
513 self.config.targets.get(&source.target_template_id),
514 self.config.profiles.get(&source.last_profile),
515 ))
516 }
517
518 async fn validate_move_destination_configuration(
519 &self,
520 selection: &MoveSelection,
521 source_harness: mj_core::config::HarnessKind,
522 operational: &mj_core::relay::RelayOperationalState,
523 ) -> Result<()> {
524 let profile_id = selection
525 .profile_id
526 .as_deref()
527 .context("move profile is unresolved")?;
528 let profile = self
529 .config
530 .profiles
531 .get(profile_id)
532 .context("move profile no longer exists")?;
533 let source = self
534 .state
535 .sessions
536 .get(&selection.session_id)
537 .context("unknown move session")?;
538 if profile.kind != source_harness || profile_id == source.last_profile {
539 return Ok(());
540 }
541 let accepted = mj_core::acp::AcceptedSessionConfig::from_configuration(
542 &operational.config,
543 &operational.config_options,
544 );
545 if accepted.model.is_none() && accepted.effort.is_none() {
546 return Ok(());
547 }
548 let choices =
551 super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
552 .await
553 .with_context(|| {
554 format!("discover destination profile {profile_id:?} configuration")
555 })?;
556 validate_preserved_configuration(profile_id, &accepted, &choices)
557 }
558
559 pub async fn prepare_move_session_controlled(
560 &self,
561 mut selection: MoveSelection,
562 executor: &(impl CommandExecutor + Sync),
563 ) -> Result<MovePreparation> {
564 ensure!(
565 selection.profile_id.is_some() || selection.target_template_id.is_some(),
566 "move requires a target or profile selection"
567 );
568 let source = self
569 .state
570 .sessions
571 .get(&selection.session_id)
572 .context("unknown session")?;
573 ensure!(
574 !self.state.subagents.contains_key(&source.id),
575 "sub-agent sessions cannot move independently of their parent"
576 );
577 let previous = crate::database::load_move_operation(&source.id)?;
578 let retry = previous.as_ref().is_some_and(|op| {
579 !matches!(op.phase, MovePhase::Completed)
580 && (op.restore_artifact().is_some() || source.checkpoint.is_some())
581 });
582 ensure!(
583 matches!(
584 source.state,
585 SessionState::Running | SessionState::Disconnected
586 ) || retry,
587 "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
588 source.id
589 );
590 selection
591 .profile_id
592 .get_or_insert_with(|| source.last_profile.clone());
593 selection
594 .target_template_id
595 .get_or_insert_with(|| source.target_template_id.clone());
596 selection
597 .additional_mounts
598 .get_or_insert_with(|| source.additional_mounts.clone());
599 ensure!(
600 !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
601 "select resource allocation or explicitly clear it, not both"
602 );
603 if !selection.clear_resource_allocation {
604 selection.resource_allocation = selection
605 .resource_allocation
606 .or_else(|| source.resource_allocation.clone());
607 }
608 if selection.clear_resource_allocation
609 && source.resource_allocation.is_none()
610 && source.container_cpus.is_none()
611 && source.container_memory.is_none()
612 {
613 selection.clear_resource_allocation = false;
614 }
615 if let Some(retained) = previous.as_ref().filter(|op| {
616 op.retains_source_environment()
617 && op.phase != MovePhase::Completed
618 && op.recovery_session.is_some()
619 }) {
620 ensure!(
621 retained.selection == selection,
622 "a sealed Move must be retried with its prepared destination and file selection"
623 );
624 ensure!(
625 !retained.in_place
626 || (retained.source_target.is_some()
627 && source.target == retained.source_target),
628 "retained Move target is missing or changed; refusing to recreate it"
629 );
630 }
631 let profile_id = selection.profile_id.as_deref().unwrap();
632 let target_id = selection.target_template_id.as_deref().unwrap();
633 let profile = self
634 .config
635 .profiles
636 .get(profile_id)
637 .context("unknown destination profile")?;
638 ensure!(
639 profile.enabled,
640 "destination profile {profile_id:?} is disabled"
641 );
642 if let Some(policy) = &selection.subagents {
643 super::profile_config::validate_session_subagent_policy(
644 &self.config,
645 profile_id,
646 policy,
647 )
648 .await?;
649 }
650 let target = self
651 .config
652 .targets
653 .get(target_id)
654 .context("unknown destination target")?;
655 self.validate_muse_resume_destination(source, profile.kind, target_id)?;
656 super::worktree::resume_compatibility(source, &self.config, target_id)
657 .map_err(anyhow::Error::msg)?;
658 super::backend::validate_resource_allocation(
659 target,
660 selection.resource_allocation.as_ref(),
661 )?;
662 let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
663 ensure!(
664 profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
665 "Muse Code ACP supports one workspace root; attached directories are unsupported"
666 );
667 ensure!(
668 mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
669 "attached resources are unsupported for this target; select compatible resources explicitly"
670 );
671 crate::targets::validate_additional_mounts(mounts)?;
672 for mount in mounts {
673 self.validate_mount_source(target_id, &mount.source, executor)?;
674 }
675 let planned_conversion =
676 self.validate_move_destination_paths(source, target_id, executor)?;
677 ensure!(
678 profile.home.is_dir(),
679 "destination profile home is unavailable; configure the profile before moving"
680 );
681 super::worker_binary::preflight_worker_binary(target, executor)?;
682 super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
683 let source_harness = previous
684 .as_ref()
685 .filter(|operation| {
686 !operation.queue_admission_started && operation.phase != MovePhase::Completed
687 })
688 .and_then(|operation| operation.recovery_session.as_ref())
689 .map_or(source.harness_kind, |record| record.harness_kind);
690 let cross_harness = profile.kind != source_harness;
691 if cross_harness {
692 let cancel = tokio_util::sync::CancellationToken::new();
693 let resolve =
694 crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
695 tokio::pin!(resolve);
696 loop {
697 tokio::select! {
698 result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
699 _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
700 if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
701 }
702 }
703 }
704 }
705 let conversion = planned_conversion
708 .map(|conversion| {
709 super::worktree::raw_conversion_preview(source, &conversion, executor)
710 .context("describe the move of this checkout into the target")
711 })
712 .transpose()?
713 .map(Box::new);
714 let (active, mut queued_commands, fingerprint) =
715 self.move_confirmation(&selection, conversion.as_deref())?;
716 replace_queued_images_with_placeholders(&mut queued_commands);
717 let operation_id = previous
718 .as_ref()
719 .filter(|operation| {
720 operation.selection == selection
721 && operation.phase != MovePhase::Completed
722 && operation.restore_artifact().is_some()
723 })
724 .map(|operation| operation.operation_id.clone())
725 .unwrap_or(new_command_id("move")?);
726 let in_place = previous
727 .as_ref()
728 .is_some_and(|op| retry && op.in_place && op.selection == selection)
729 || in_place_move_eligible(
730 source,
731 &selection,
732 &self.config.targets,
733 self.state.subagents.contains_key(&source.id),
734 retry,
735 );
736 let workspace = if in_place {
737 None
738 } else {
739 Some(self.assess_move_workspace(&selection, executor)?)
740 };
741 Ok(MovePreparation {
742 workspace,
743 source_unavailable: false,
744 in_place,
745 conversion,
746 selection,
747 source_profile_id: source.last_profile.clone(),
748 source_target_template_id: source.target_template_id.clone(),
749 cross_harness,
750 active,
751 queued_commands,
752 fingerprint,
753 operation_id,
754 })
755 }
756
757 pub async fn move_session_managed_controlled(
759 &mut self,
760 request: MoveSessionRequest,
761 executor: &(impl CommandExecutor + Sync),
762 manager: &SessionManagerControl,
763 ) -> Result<MoveOutcome> {
764 let prepared = &request.preparation;
765 let id = prepared.selection.session_id.clone();
766 let started = std::time::Instant::now();
767 executor.notify_notice("Checking destination");
768 let mut source_relay = MoveSourceRelay::default();
769 let checked = {
770 let _checking_destination =
771 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
772 let mut checked = self
773 .prepare_move_session_controlled(prepared.selection.clone(), executor)
774 .await?;
775 if matches!(
780 self.state.sessions[&id].state,
781 SessionState::Running | SessionState::Disconnected
782 ) {
783 let source_harness = self.state.sessions[&id].harness_kind;
784 source_relay = MoveSourceRelay::lease(manager, &id).await?;
785 let snapshot = source_relay.snapshot();
786 let (active, queue, fingerprint) =
787 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
788 checked.source_unavailable = snapshot
789 .as_ref()
790 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
791 checked.active = active
792 || checked.source_unavailable
793 || snapshot.as_ref().is_some_and(|snapshot| {
794 let mut operational = snapshot.operational.clone();
795 operational.queued_prompts.clear();
796 operational.checkpoint_barrier = None;
797 !operational.safe_to_replace(source_harness)
798 });
799 if let Some(snapshot) = &snapshot {
800 self.validate_move_destination_configuration(
801 &checked.selection,
802 source_harness,
803 &snapshot.operational,
804 )
805 .await?;
806 }
807 checked.queued_commands = queue;
808 checked.fingerprint = fingerprint;
809 }
810 checked
811 };
812 ensure!(
813 checked.fingerprint == prepared.fingerprint,
814 "session, pending work, or destination configuration changed; prepare and confirm Move again"
815 );
816 let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
817 let source = self.state.sessions[&id].clone();
818 if let Some(assessment) = &checked.workspace {
819 checked.selection.workspace.validate(assessment)?;
820 }
821 let old_operation = crate::database::load_move_operation(&id)?;
822 let retry = old_operation.filter(|op| {
823 op.selection == prepared.selection
824 && op.phase != MovePhase::Completed
825 && op.restore_artifact().is_some()
826 });
827 if retry.is_none()
828 && source.state == SessionState::Running
829 && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
830 && move_environment_change(&source, &checked.selection, false).is_none()
831 {
832 return Ok(outcome(
833 &prepared.operation_id,
834 &prepared.selection,
835 "unchanged",
836 None,
837 None,
838 ));
839 }
840 ensure!(
841 !checked.active || request.acknowledge_interruption,
842 "active work will be interrupted; confirm Move again with interruption acknowledgement"
843 );
844 ensure!(
845 checked.queued_commands.is_empty() || request.queue.is_some(),
846 "pending work requires an explicit queue choice: discard or start"
847 );
848 let timestamp = now();
849 let mut operation = match retry {
850 Some(mut op) => {
851 ensure!(
852 !op.queue_admission_started || op.queue == queue,
853 "queued work may already have run; retry with the original queue choice on the same destination"
854 );
855 crate::database::clear_move_cancellation_for_retry(&id)?;
856 op.configuration_fingerprint =
857 self.move_configuration_fingerprint(&checked.selection)?;
858 op.queue = queue;
859 op.cancellation_requested = false;
860 op.error = None;
861 op
862 }
863 None => MoveOperation {
864 workspace_transfer: if checked.in_place {
865 None
866 } else {
867 Some(
868 self.new_workspace_transfer(
869 &id,
870 &prepared.operation_id,
871 checked
872 .workspace
873 .clone()
874 .context("Move workspace assessment missing")?,
875 executor,
876 )?,
877 )
878 },
879 handoff: None,
880 source_checkpoint_only: false,
881 in_place: in_place_move_eligible(
882 &source,
883 &checked.selection,
884 &self.config.targets,
885 self.state.subagents.contains_key(&id),
886 false,
887 ),
888 operation_id: prepared.operation_id.clone(),
889 selection: checked.selection.clone(),
890 source_profile_id: source.last_profile.clone(),
891 source_target_template_id: source.target_template_id.clone(),
892 source_target: source.target.clone(),
893 source_native_session_id: source.native_session_id.clone(),
894 source_additional_mounts: source.additional_mounts.clone(),
895 source_resource_allocation: source.resource_allocation.clone(),
896 destination_target: None,
897 destination_native_session_id: None,
898 destination_store_id: None,
899 configuration_fingerprint: self
900 .move_configuration_fingerprint(&checked.selection)?,
901 checkpoint: None,
902 recovery_session: None,
903 queue,
904 phase: MovePhase::Preparing,
905 queue_admission_started: false,
906 queue_admission_finished: false,
907 cancellation_requested: false,
908 created_at: timestamp.clone(),
909 updated_at: timestamp,
910 error: None,
911 },
912 };
913 crate::database::save_move_operation(&operation)?;
914 tracing::info!(
915 session_id = id,
916 in_place = operation.in_place,
917 reason = move_environment_change(
918 &source,
919 &checked.selection,
920 bare_targets_share_environment(
921 &self.config.targets,
922 &source.target_template_id,
923 checked
924 .selection
925 .target_template_id
926 .as_deref()
927 .unwrap_or_default(),
928 ),
929 )
930 .unwrap_or(if operation.in_place {
931 "environment unchanged"
932 } else {
933 "source environment unavailable or previously released"
934 }),
935 "move environment decision"
936 );
937 tracing::info!(
938 session_id = id,
939 phase = "preflight",
940 elapsed_ms = started.elapsed().as_millis() as u64,
941 "move phase completed"
942 );
943 let result = Box::pin(self.execute_move(
946 &mut operation,
947 Some(&checked),
948 executor,
949 manager,
950 source_relay,
951 ))
952 .await;
953 self.finish_move_result(&mut operation, result, executor)
954 }
955
956 fn finish_move_result(
957 &self,
958 operation: &mut MoveOperation,
959 result: Result<()>,
960 executor: &impl CommandExecutor,
961 ) -> Result<MoveOutcome> {
962 if result.is_err()
963 && operation.workspace_transfer.is_some()
964 && !crate::upgrade::gate().is_open()
965 {
966 return Ok(outcome(
969 &operation.operation_id,
970 &operation.selection,
971 "interrupted",
972 None,
973 Some("Move will continue after the daemon upgrade".into()),
974 ));
975 }
976 let (status, error, recovery) = match result {
977 Ok(()) => {
978 operation.phase = MovePhase::Completed;
979 ("completed", None, None)
980 }
981 Err(error) => {
982 let cancelled =
983 executor.cancellation_requested() || operation.cancellation_requested;
984 operation.phase = if cancelled {
985 MovePhase::Cancelled
986 } else {
987 MovePhase::Failed
988 };
989 operation.cancellation_requested = cancelled;
990 let recovery = failed_move_recovery(
991 operation,
992 self.state.sessions.get(&operation.selection.session_id),
993 );
994 let error = format!("{error:#}");
995 operation.error = Some(error.clone());
996 (
997 if cancelled { "cancelled" } else { "failed" },
998 Some(error),
999 Some(recovery),
1000 )
1001 }
1002 };
1003 operation.updated_at = now();
1004 crate::database::save_move_operation(operation)?;
1005 Ok(outcome(
1006 &operation.operation_id,
1007 &operation.selection,
1008 status,
1009 error,
1010 recovery,
1011 ))
1012 }
1013
1014 pub async fn recover_move_managed_controlled(
1015 &mut self,
1016 mut operation: MoveOperation,
1017 executor: &(impl CommandExecutor + Sync),
1018 manager: &SessionManagerControl,
1019 ) -> Result<MoveOutcome> {
1020 let id = operation.selection.session_id.clone();
1021 let result = async {
1022 let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1023 if operation.queue_admission_started {
1024 ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1025 self.finish_workspace_transfer(&mut operation, executor)?;
1026 return self.admit_move_queue(&mut operation, executor).await;
1027 }
1028 if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1029 && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1030 if session.state == SessionState::Provisioning || session.target != operation.source_target {
1031 self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1032 }
1033 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1034 }
1035 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1036 let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1037 Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1038 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1039 if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1040 let previous = self.state.sessions[&id].clone();
1041 operation.recovery_session = Some(previous.clone());
1042 let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1043 return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1044 }
1045 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1046 self.cleanup_stopped_target(&id, &cleanup)?;
1047 }
1048 bail!("{}", interrupted_source_stop_message(
1049 "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1050 operation.in_place,
1051 ));
1052 }
1053 match operation.phase {
1054 MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1055 MovePhase::ResumingDestination if session.state == SessionState::Running => {
1056 ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1060 && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1061 "ready destination does not match the move intent");
1062 operation.destination_target = session.target.clone();
1063 operation.destination_native_session_id = session.native_session_id.clone();
1064 operation.queue_admission_started = true;
1065 operation.phase = MovePhase::StartingQueue;
1066 crate::database::save_move_operation(&operation)?;
1067 restore_move_queue_hold(&operation);
1068 ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1069 self.finish_workspace_transfer(&mut operation, executor)?;
1070 self.admit_move_queue(&mut operation, executor).await
1071 }
1072 MovePhase::ResumingDestination => {
1073 let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1074 let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1077 let error = if operation.in_place {
1078 self.retain_failed_in_place_move(&id, previous, cause)?
1079 } else if operation.workspace_transfer.is_some() {
1080 self.rollback_move_destination(&operation, cause, executor)?
1081 } else {
1082 self.rollback_failed_resume(&id, previous, false, cause, executor)?
1083 };
1084 Err(error)
1085 }
1086 MovePhase::ClosingSource => {
1087 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1088 Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1089 }
1090 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1091 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1092 self.cleanup_stopped_target(&id, executor)?;
1093 }
1094 bail!("{}", interrupted_source_stop_message(
1095 "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1096 operation.in_place,
1097 ))
1098 }
1099 _ => bail!("Move requires an explicit retry after the daemon restarted"),
1100 }
1101 }.await;
1102 self.finish_move_result(&mut operation, result, executor)
1103 }
1104
1105 async fn recover_move_source_stop(
1106 &mut self,
1107 operation: &mut MoveOperation,
1108 executor: &(impl CommandExecutor + Sync),
1109 manager: &SessionManagerControl,
1110 ) -> Result<()> {
1111 let id = operation.selection.session_id.clone();
1112 if self.state.sessions[&id].state == SessionState::Closing {
1113 self.prepare_move_source_checkpoint(
1114 &id,
1115 executor,
1116 manager,
1117 operation,
1118 &mut MoveSourceRelay::default(),
1119 )
1120 .await?;
1121 let handle = manager
1122 .wait_for_session(&id, std::time::Duration::from_secs(5))
1123 .await?;
1124 let mut lease = handle.lease_connection().await?;
1125 let execution = lease.connection_mut().sync().await?.operational.execution;
1126 if operation.retains_source_environment()
1127 && matches!(
1128 execution,
1129 mj_core::relay::RelayExecutionState::Closing
1130 | mj_core::relay::RelayExecutionState::Closed
1131 )
1132 {
1133 super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1134 lease.release();
1135 ensure!(
1136 operation.restore_artifact().is_some()
1137 || self.state.sessions[&id].checkpoint.is_some(),
1138 "sealed Move source has no checkpoint"
1139 );
1140 return Ok(());
1141 }
1142 lease.release();
1143 if matches!(
1144 execution,
1145 mj_core::relay::RelayExecutionState::Idle
1146 | mj_core::relay::RelayExecutionState::Running
1147 ) {
1148 if operation.cancellation_requested || operation.retains_source_environment() {
1149 let record = self.state.sessions.get_mut(&id).unwrap();
1150 record.state = SessionState::Running;
1151 record.updated_at = now();
1152 record.last_error = Some(
1153 "Move was interrupted before the source was sealed; source retained".into(),
1154 );
1155 crate::database::save_lifecycle_session(record)?;
1156 } else {
1157 Box::pin(self.suspend_session_for_move(
1159 &id,
1160 executor,
1161 manager,
1162 operation,
1163 None,
1164 SourceTargetDisposition::Destroy,
1165 MoveSourceRelay::default(),
1166 ))
1167 .await?;
1168 }
1169 return Ok(());
1170 }
1171 }
1172 ensure!(
1173 !operation.retains_source_environment(),
1174 "retained Move source cannot be proven; refusing target teardown"
1175 );
1176 self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1178 .await?;
1179 Ok(())
1180 }
1181
1182 async fn execute_move(
1185 &mut self,
1186 operation: &mut MoveOperation,
1187 preparation: Option<&MovePreparation>,
1188 executor: &(impl CommandExecutor + Sync),
1189 manager: &SessionManagerControl,
1190 mut source_relay: MoveSourceRelay,
1191 ) -> Result<()> {
1192 let id = operation.selection.session_id.clone();
1193 ensure!(
1194 !executor.cancellation_requested(),
1195 "move cancelled before source interruption"
1196 );
1197 if !operation.queue_admission_started {
1198 if operation.in_place {
1199 ensure!(
1200 operation.source_target.is_some()
1201 && self.state.sessions[&id].target == operation.source_target,
1202 "retained Move target is missing or changed; refusing to recreate the environment"
1203 );
1204 }
1205 if self.state.sessions[&id].state == SessionState::Error
1206 && let Some(previous) = operation.recovery_session.as_ref()
1207 {
1208 let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
1211 if operation.in_place {
1212 self.retain_failed_in_place_move(&id, previous, cause)?;
1213 let record = self.state.sessions.get_mut(&id).unwrap();
1214 record.state = SessionState::Closing;
1215 crate::database::save_resumed_session(record, None)?;
1216 } else if operation.workspace_transfer.is_some() {
1217 self.rollback_move_destination(operation, cause, executor)?;
1218 let record = self.state.sessions.get_mut(&id).unwrap();
1219 record.state = SessionState::Closing;
1220 crate::database::save_resumed_session(record, None)?;
1221 } else {
1222 let failure =
1223 self.rollback_failed_resume(&id, previous, false, cause, executor)?;
1224 ensure!(
1225 self.state.sessions[&id].state == SessionState::Stopped,
1226 "{failure:#}"
1227 );
1228 }
1229 }
1230 let state = self.state.sessions[&id].state;
1231 if matches!(state, SessionState::Closing | SessionState::Destroying)
1232 && !(operation.retains_source_environment() && operation.recovery_session.is_some())
1233 {
1234 Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
1235 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1236 }
1237 if matches!(
1238 self.state.sessions[&id].state,
1239 SessionState::Running | SessionState::Disconnected
1240 ) && operation.destination_target.is_none()
1241 {
1242 executor.notify_notice("Stopping source");
1243 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
1244 operation.phase = MovePhase::ClosingSource;
1245 operation.updated_at = now();
1246 crate::database::save_move_operation(operation)?;
1247 let disposition = if operation.in_place {
1251 SourceTargetDisposition::RetainForInPlaceSwap
1252 } else {
1253 SourceTargetDisposition::Destroy
1254 };
1255 Box::pin(self.suspend_session_for_move(
1256 &id,
1257 executor,
1258 manager,
1259 operation,
1260 preparation,
1261 disposition,
1262 std::mem::take(&mut source_relay),
1263 ))
1264 .await?;
1265 }
1266 drop(source_relay);
1268 if !operation.in_place
1271 && operation.workspace_transfer.is_none()
1272 && self.state.sessions[&id].state == SessionState::Stopped
1273 && self.state.sessions[&id].target.is_some()
1274 {
1275 executor.notify_notice("Cleaning up source");
1276 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
1277 self.cleanup_stopped_target(&id, executor)?;
1278 }
1279 operation.checkpoint = operation
1280 .checkpoint
1281 .clone()
1282 .or_else(|| self.state.sessions[&id].checkpoint.clone());
1283 ensure!(
1284 operation.restore_artifact().is_some(),
1285 "move has no verified checkpoint"
1286 );
1287 ensure!(
1288 !executor.cancellation_requested(),
1289 "move cancelled after source sealing; checkpoint and remaining environment retained"
1290 );
1291 operation.phase = MovePhase::ResumingDestination;
1292 operation.recovery_session = Some(self.state.sessions[&id].clone());
1293 operation.updated_at = now();
1294 crate::database::save_move_operation(operation)?;
1295 executor.reserve_move_destination();
1296 executor.notify_notice("Preparing destination");
1297 if let Some(policy) = &operation.selection.subagents {
1301 let session = self.state.sessions.get_mut(&id).unwrap();
1302 session.subagents = Some(policy.clone());
1303 crate::database::save_resumed_session(session, None)?;
1304 }
1305 if operation.in_place {
1306 Box::pin(self.restore_session_in_place(
1310 &id,
1311 operation.selection.profile_id.as_deref().unwrap(),
1312 operation.selection.target_template_id.as_deref().unwrap(),
1313 executor,
1314 ))
1315 .await?;
1316 } else {
1317 if operation.selection.clear_resource_allocation {
1318 let session = self.state.sessions.get_mut(&id).unwrap();
1319 session.resource_allocation = None;
1320 session.container_cpus = None;
1321 session.container_memory = None;
1322 crate::database::save_resumed_session(session, None)?;
1323 }
1324 if operation.workspace_transfer.is_some() {
1325 self.capture_move_workspace(operation, executor)?;
1326 Box::pin(self.resume_session_for_move(operation, executor)).await?;
1327 operation.workspace_transfer = crate::database::load_move_operation(&id)?
1328 .context("Move intent disappeared during restore")?
1329 .workspace_transfer;
1330 } else {
1331 Box::pin(self.resume_session_controlled(
1332 &id,
1333 operation.selection.profile_id.as_deref().unwrap(),
1334 operation.selection.target_template_id.as_deref().unwrap(),
1335 SessionResumeOptions {
1336 additional_mounts: operation.selection.additional_mounts.clone(),
1337 resource_allocation: operation.selection.resource_allocation.clone(),
1338 discard_queue: true,
1339 },
1340 executor,
1341 ))
1342 .await?;
1343 }
1344 }
1345 let destination = &self.state.sessions[&id];
1346 operation.destination_target = destination.target.clone();
1347 operation.destination_native_session_id = destination.native_session_id.clone();
1348 operation.phase = MovePhase::StartingQueue;
1350 operation.queue_admission_started = true;
1351 operation.updated_at = now();
1352 crate::database::save_move_operation(operation)?;
1353 }
1354 restore_move_queue_hold(operation);
1355 self.finish_workspace_transfer(operation, executor)?;
1356 self.admit_move_queue(operation, executor).await
1357 }
1358
1359 async fn admit_move_queue(
1360 &self,
1361 operation: &mut MoveOperation,
1362 executor: &(impl CommandExecutor + Sync),
1363 ) -> Result<()> {
1364 let timing_id = operation.selection.session_id.clone();
1365 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
1366 let id = &operation.selection.session_id;
1367 let mut relay = {
1368 let _checking_destination =
1369 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1370 let destination = &self.state.sessions[id];
1371 ensure!(
1372 destination.state == SessionState::Running
1373 && destination.target == operation.destination_target
1374 && destination.native_session_id == operation.destination_native_session_id,
1375 "cannot prove the same ready destination; refusing to replay potentially executed work"
1376 );
1377 let spec = self.reconnect_command(id)?;
1378 let relay = StandaloneSession::connect_command(&spec, id).await?;
1379 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")?;
1380 if let Some(expected) = &operation.destination_store_id {
1381 ensure!(
1382 *expected == store_id,
1383 "destination relay storage was replaced; refusing to replay potentially executed work"
1384 );
1385 } else {
1386 operation.destination_store_id = Some(store_id);
1388 crate::database::save_move_operation(operation)?;
1389 }
1390 ensure!(
1391 relay.snapshot().operational.native_session_id
1392 == operation.destination_native_session_id,
1393 "destination relay native identity changed; refusing queue replay"
1394 );
1395 relay
1396 };
1397 if operation.queue == ResumeQueueDisposition::Start {
1398 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
1399 executor.notify_notice("Starting queued work");
1400 let checkpoint = operation
1401 .restore_artifact()
1402 .context("move queue archive is missing")?;
1403 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1404 ensure!(
1405 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
1406 "move queue checkpoint verification failed"
1407 );
1408 for queued in verified.canonical_session.queued_prompts {
1409 if operation.queue_admission_finished {
1410 continue;
1411 }
1412 ensure!(
1413 !executor.cancellation_requested(),
1414 "move cancelled during queue admission; destination retained"
1415 );
1416 let command = match queued.kind {
1417 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
1418 prompt: queued
1419 .content
1420 .into_iter()
1421 .map(serde_json::from_value)
1422 .collect::<serde_json::Result<_>>()?,
1423 },
1424 CanonicalQueuedCommandKind::SetConfig { key, value } => {
1425 RelayCommand::SetConfig { key, value }
1426 }
1427 };
1428 relay.submit_accepted(queued.command_id, command).await?;
1429 }
1430 }
1431 operation.queue_admission_finished = true;
1432 crate::database::save_move_operation(operation)?;
1433 restore_move_queue_hold(operation);
1434 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
1435 "Queued work was discarded; ready and idle."
1436 } else {
1437 "Queued work was accepted."
1438 };
1439 let source_profile = &operation.source_profile_id;
1440 let source_target = &operation.source_target_template_id;
1441 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
1442 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
1443 let text = if operation.in_place {
1446 format!(
1447 "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."
1448 )
1449 } else {
1450 format!(
1451 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
1452 )
1453 };
1454 relay
1455 .submit(
1456 format!("{}-notice", operation.operation_id),
1457 RelayCommand::RecordNotice { text },
1458 )
1459 .await?;
1460 Ok(())
1461 }
1462
1463 pub(super) fn validate_move_checkpoint(
1464 &self,
1465 operation: &MoveOperation,
1466 preparation: Option<&MovePreparation>,
1467 executor: &(impl CommandExecutor + Sync),
1468 ) -> Result<()> {
1469 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1470 let current = Controller {
1471 config: mj_core::config::Config::load()?,
1472 state: self.state.clone(),
1473 };
1474 ensure!(
1475 current.move_configuration_fingerprint(&operation.selection)?
1476 == operation.configuration_fingerprint,
1477 "destination configuration changed during move"
1478 );
1479 let id = &operation.selection.session_id;
1480 current.validate_move_destination_paths(
1481 &self.state.sessions[id],
1482 operation.selection.target_template_id.as_deref().unwrap(),
1483 executor,
1484 )?;
1485 if !operation.in_place
1486 && operation.workspace_transfer.is_none()
1487 && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
1488 .preflight_resume_repository_sources(
1489 id,
1490 operation.selection.target_template_id.as_deref().unwrap(),
1491 executor,
1492 )?
1493 {
1494 bail!(
1495 "destination repository source is missing checkpoint commit {}; source retained",
1496 mismatch.missing_commit
1497 );
1498 }
1499 if let Some(prepared) = preparation {
1500 let checkpoint = operation
1501 .restore_artifact()
1502 .or(self.state.sessions[id].checkpoint.as_ref())
1503 .context("no move checkpoint")?;
1504 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1505 let actual: Vec<_> = verified
1506 .canonical_session
1507 .queued_prompts
1508 .iter()
1509 .map(|p| p.command_id.as_str())
1510 .collect();
1511 let expected: Vec<_> = prepared
1512 .queued_commands
1513 .iter()
1514 .map(|p| p.command_id.as_str())
1515 .collect();
1516 ensure!(
1517 actual == expected,
1518 "pending queue changed before checkpoint capture; source retained, confirm Move again"
1519 );
1520 }
1521 Ok(())
1522 }
1523}
1524
1525fn validate_preserved_configuration(
1526 profile_id: &str,
1527 accepted: &mj_core::acp::AcceptedSessionConfig,
1528 choices: &mj_core::worker_launch::ProfileConfig,
1529) -> Result<()> {
1530 for (key, value, offered) in [
1531 ("model", accepted.model.as_deref(), &choices.models),
1532 ("effort", accepted.effort.as_deref(), &choices.efforts),
1533 ] {
1534 let Some(value) = value else { continue };
1535 ensure!(
1536 offered.iter().any(|choice| choice.value == value),
1537 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
1538 offered
1539 .iter()
1540 .map(|choice| choice.value.as_str())
1541 .collect::<Vec<_>>()
1542 .join(", ")
1543 );
1544 }
1545 Ok(())
1546}
1547
1548pub(super) fn in_place_move_eligible(
1555 source: &mj_core::state::SessionRecord,
1556 selection: &MoveSelection,
1557 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
1558 is_subagent: bool,
1559 retry: bool,
1560) -> bool {
1561 !retry
1562 && !is_subagent
1563 && source.target.is_some()
1564 && matches!(
1565 source.state,
1566 SessionState::Running | SessionState::Disconnected
1567 )
1568 && move_environment_change(
1569 source,
1570 selection,
1571 bare_targets_share_environment(
1572 targets,
1573 &source.target_template_id,
1574 selection.target_template_id.as_deref().unwrap_or_default(),
1575 ),
1576 )
1577 .is_none()
1578}
1579
1580fn bare_targets_share_environment(
1585 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
1586 source_id: &str,
1587 destination_id: &str,
1588) -> bool {
1589 use mj_core::config::TargetTemplate::{LocalBare, SshBare};
1590 match (targets.get(source_id), targets.get(destination_id)) {
1591 (Some(LocalBare), Some(LocalBare)) => true,
1592 (
1593 Some(SshBare { ssh: source, .. }),
1594 Some(SshBare {
1595 ssh: destination, ..
1596 }),
1597 ) => {
1598 crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
1599 }
1600 _ => false,
1601 }
1602}
1603fn outcome(
1604 operation_id: &str,
1605 selection: &MoveSelection,
1606 status: &str,
1607 error: Option<String>,
1608 recovery: Option<String>,
1609) -> MoveOutcome {
1610 MoveOutcome {
1611 operation_id: operation_id.into(),
1612 session_id: selection.session_id.clone(),
1613 profile_id: selection.profile_id.clone().unwrap_or_default(),
1614 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
1615 outcome: status.into(),
1616 error,
1617 recovery,
1618 }
1619}
1620
1621fn move_environment_change(
1625 source: &mj_core::state::SessionRecord,
1626 selection: &MoveSelection,
1627 same_environment: bool,
1628) -> Option<&'static str> {
1629 if Some(&source.target_template_id) != selection.target_template_id.as_ref()
1630 && !same_environment
1631 {
1632 Some("target changed")
1633 } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
1634 Some("attached mounts changed")
1635 } else if source.resource_allocation != selection.resource_allocation
1636 || (selection.clear_resource_allocation
1637 && (source.container_cpus.is_some() || source.container_memory.is_some()))
1638 {
1639 Some("resource allocation changed")
1640 } else {
1641 None
1642 }
1643}