1#[cfg(test)]
4mod tests;
5
6use anyhow::{Context, Result, bail, ensure};
7use mj_core::hex::lower_hex;
8use sha2::{Digest, Sha256};
9
10use super::lifecycle::SourceTargetDisposition;
11use super::{Controller, SessionResumeOptions, now};
12
13fn mutation_holds() -> &'static std::sync::Mutex<std::collections::BTreeSet<String>> {
14 static HOLDS: std::sync::OnceLock<std::sync::Mutex<std::collections::BTreeSet<String>>> =
15 std::sync::OnceLock::new();
16 HOLDS.get_or_init(Default::default)
17}
18
19pub fn move_owns_session(session_id: &str) -> bool {
20 move_has_pending_queue(session_id)
21 || mutation_holds()
22 .lock()
23 .unwrap_or_else(std::sync::PoisonError::into_inner)
24 .contains(session_id)
25}
26
27fn queue_holds() -> &'static std::sync::Mutex<std::collections::BTreeSet<String>> {
28 static HOLDS: std::sync::OnceLock<std::sync::Mutex<std::collections::BTreeSet<String>>> =
29 std::sync::OnceLock::new();
30 HOLDS.get_or_init(Default::default)
31}
32
33pub fn move_has_pending_queue(session_id: &str) -> bool {
34 queue_holds()
35 .lock()
36 .unwrap_or_else(std::sync::PoisonError::into_inner)
37 .contains(session_id)
38}
39
40pub fn release_move_queue_hold(session_id: &str) {
41 queue_holds()
42 .lock()
43 .unwrap_or_else(std::sync::PoisonError::into_inner)
44 .remove(session_id);
45}
46
47pub fn restore_move_queue_hold(operation: &MoveOperation) {
48 let mut holds = queue_holds()
49 .lock()
50 .unwrap_or_else(std::sync::PoisonError::into_inner);
51 if operation.queue_admission_started && !operation.queue_admission_finished {
52 holds.insert(operation.selection.session_id.clone());
53 } else {
54 holds.remove(&operation.selection.session_id);
55 }
56}
57
58fn interrupted_source_stop_message(recovered: &str, in_place: bool) -> String {
64 if in_place {
65 format!("{recovered}; the in-place swap was interrupted; the environment was released")
66 } else {
67 recovered.to_owned()
68 }
69}
70
71fn source_stopped_with_verified_checkpoint(record: &mj_core::state::SessionRecord) -> bool {
77 record.state == SessionState::Stopped && record.checkpoint.is_some()
78}
79
80fn stopped_source_recovery(
86 session_id: &str,
87 destination_profile: Option<&str>,
88 destination_target: Option<&str>,
89) -> String {
90 let flag = |name: &str, value: Option<&str>| {
91 value
92 .filter(|value| !value.is_empty())
93 .map(|value| format!(" --{name} {value}"))
94 .unwrap_or_default()
95 };
96 format!(
97 "Source is stopped with a verified checkpoint. Bring it back with \
98 `mj resume --session {session_id}{}{} --queue start`.",
99 flag("profile", destination_profile),
100 flag("target", destination_target),
101 )
102}
103
104fn failed_move_recovery(
107 operation: &MoveOperation,
108 record: Option<&mj_core::state::SessionRecord>,
109) -> String {
110 if operation.queue_admission_started {
111 return "Destination is live; retry queue admission on this same destination. Already accepted work may have effects.".to_owned();
112 }
113 match record {
114 Some(record) if source_stopped_with_verified_checkpoint(record) => {
115 stopped_source_recovery(
116 &operation.selection.session_id,
117 operation.selection.profile_id.as_deref(),
118 operation.selection.target_template_id.as_deref(),
119 )
120 }
121 _ => "Source or partial destination is retained. Retry move after resolving the reported error.".to_owned(),
122 }
123}
124
125pub struct MoveMutationGuard(String);
126
127impl MoveMutationGuard {
128 pub fn reserve(session_id: &str) -> Result<Self> {
129 ensure!(
130 mutation_holds()
131 .lock()
132 .unwrap_or_else(std::sync::PoisonError::into_inner)
133 .insert(session_id.to_owned()),
134 "session already has a move owner"
135 );
136 Ok(Self(session_id.to_owned()))
137 }
138}
139
140impl Drop for MoveMutationGuard {
141 fn drop(&mut self) {
142 mutation_holds()
143 .lock()
144 .unwrap_or_else(std::sync::PoisonError::into_inner)
145 .remove(&self.0);
146 }
147}
148
149pub(crate) fn move_refuses_command(session_id: &str, command: &RelayCommand) -> bool {
150 move_owns_session(session_id)
151 && matches!(
152 command,
153 RelayCommand::Prompt { .. }
154 | RelayCommand::SetConfig { .. }
155 | RelayCommand::SetSessionMode { .. }
156 | RelayCommand::RunUserShell { .. }
157 | RelayCommand::CancelUserShell { .. }
158 | RelayCommand::Cancel
159 | RelayCommand::RemoveQueuedPrompt { .. }
160 | RelayCommand::ClearQueuedPrompts
161 )
162}
163use crate::session_manager::{SessionManagerControl, StandaloneSession, new_command_id};
164use mj_checkpoint::archive::{CanonicalQueuedCommandKind, verify_archive_streaming};
165use mj_core::state::{MoveOperation, MovePhase, ResumeQueueDisposition, SessionState};
166
167pub use mj_core::state::{MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest};
168
169use crate::targets::{CommandExecutor, ProvisionStage, ProvisionStageGuard};
170use mj_core::relay::RelayCommand;
171
172pub async fn refresh_move_source(
175 manager: &SessionManagerControl,
176 id: &str,
177) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
178 let handle = manager
179 .wait_for_session(id, std::time::Duration::from_secs(5))
180 .await?;
181 let result = async {
182 let mut lease = handle.lease_connection().await?;
183 let snapshot = lease.connection_mut().sync().await?;
184 lease.release();
185 Ok(snapshot)
186 }
187 .await;
188 match result {
189 Ok(snapshot) => Ok(Some(snapshot)),
190 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
191 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
192 Ok(None)
193 }
194 Err(error) => Err(error),
195 }
196}
197
198fn digest(value: &impl serde::Serialize) -> Result<String> {
199 Ok(lower_hex(Sha256::digest(serde_json::to_vec(value)?)))
200}
201
202struct MovePhaseTimer<'a> {
203 session_id: &'a str,
204 phase: &'static str,
205 started: std::time::Instant,
206}
207
208impl<'a> MovePhaseTimer<'a> {
209 fn new(session_id: &'a str, phase: &'static str) -> Self {
210 Self {
211 session_id,
212 phase,
213 started: std::time::Instant::now(),
214 }
215 }
216}
217
218impl Drop for MovePhaseTimer<'_> {
219 fn drop(&mut self) {
220 tracing::info!(
221 session_id = self.session_id,
222 phase = self.phase,
223 elapsed_ms = self.started.elapsed().as_millis() as u64,
224 "move phase finished"
225 );
226 }
227}
228
229fn replace_queued_images_with_placeholders(
233 queued_commands: &mut [mj_core::state::MaterializedQueuedPrompt],
234) {
235 for block in queued_commands
236 .iter_mut()
237 .flat_map(|command| command.content.iter_mut())
238 {
239 if block.get("type").and_then(serde_json::Value::as_str) != Some("image") {
240 continue;
241 }
242 let mime = block
243 .get("mimeType")
244 .or_else(|| block.get("mime_type"))
245 .and_then(serde_json::Value::as_str)
246 .unwrap_or("image");
247 *block = serde_json::json!({"type": "text", "text": format!("[Image attachment: {mime}]")});
248 }
249}
250
251impl Controller {
252 fn validate_move_destination_paths(
256 &self,
257 source: &mj_core::state::SessionRecord,
258 target_id: &str,
259 executor: &(impl CommandExecutor + Sync),
260 ) -> Result<Option<super::worktree::RawToWorkspaceConversion>> {
261 use super::worktree::ResumePlan;
262 match super::worktree::resume_compatibility(source, &self.config, target_id)
263 .map_err(anyhow::Error::msg)?
264 {
265 ResumePlan::RawToWorkspace => {
266 return Ok(Some(super::worktree::plan_raw_to_workspace(
267 source,
268 &self.config,
269 executor,
270 )?));
271 }
272 ResumePlan::WorkspaceToRaw => {
273 self.plan_workspace_to_raw(source, target_id, executor)?;
274 }
275 ResumePlan::InPlace if source.managed_worktree.is_none() => {
276 if let Some(path) = &source.project_directory {
277 self.validate_project_directory(target_id, path, executor)?;
278 }
279 }
280 ResumePlan::InPlace => {}
281 }
282 Ok(None)
283 }
284 fn move_confirmation(
285 &self,
286 selection: &MoveSelection,
287 conversion: Option<&mj_core::state::RawConversionPreview>,
288 ) -> Result<(bool, Vec<mj_core::state::MaterializedQueuedPrompt>, String)> {
289 let source = self
290 .state
291 .sessions
292 .get(&selection.session_id)
293 .context("unknown move session")?;
294 let (mut active, mut queued) = crate::database::move_pending_work(&source.id)?;
295 if let Some(operation) = crate::database::load_move_operation(&source.id)?
296 && operation.queue_admission_started
297 && !operation.queue_admission_finished
298 {
299 ensure!(
300 operation.selection == *selection,
301 "queue admission is incomplete on the live destination; retry that move before selecting another destination"
302 );
303 let checkpoint = operation
304 .checkpoint
305 .as_ref()
306 .context("retained queue checkpoint is missing")?;
307 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
308 ensure!(
309 verified.archive_sha256 == checkpoint.sha256
310 && verified.manifest.session.id == source.id,
311 "retained move checkpoint verification failed"
312 );
313 queued = verified
314 .canonical_session
315 .queued_prompts
316 .into_iter()
317 .map(|entry| mj_core::state::MaterializedQueuedPrompt {
318 accepted_ordinal: None,
319 command_id: entry.command_id,
320 kind: match entry.kind {
321 CanonicalQueuedCommandKind::Prompt => {
322 mj_core::state::QueuedCommandKind::Prompt
323 }
324 CanonicalQueuedCommandKind::SetConfig { key, value } => {
325 mj_core::state::QueuedCommandKind::SetConfig { key, value }
326 }
327 },
328 content: entry.content,
329 queued_at_ms: entry.queued_at_ms,
330 })
331 .collect();
332 active = false;
333 }
334 let fingerprint = digest(&(
335 &source.last_profile,
336 &source.target_template_id,
337 &source.target,
338 &source.target_runtime,
339 &source.native_session_id,
340 &source.resource_allocation,
341 &source.additional_mounts,
342 &source.container_cpus,
343 &source.container_memory,
344 selection,
345 self.move_configuration_fingerprint(selection)?,
346 &queued,
347 conversion.map(|preview| {
351 (
352 &preview.fetch_url,
353 &preview.push_urls,
354 &preview.branch,
355 &preview.destination,
356 )
357 }),
358 ))?;
359 Ok((active, queued, fingerprint))
360 }
361 pub(super) fn move_configuration_fingerprint(
362 &self,
363 selection: &MoveSelection,
364 ) -> Result<String> {
365 let profile = self
366 .config
367 .profiles
368 .get(
369 selection
370 .profile_id
371 .as_deref()
372 .context("move profile is unresolved")?,
373 )
374 .context("move profile no longer exists")?;
375 let target = self
376 .config
377 .targets
378 .get(
379 selection
380 .target_template_id
381 .as_deref()
382 .context("move target is unresolved")?,
383 )
384 .context("move target no longer exists")?;
385 let source = self
388 .state
389 .sessions
390 .get(&selection.session_id)
391 .context("unknown move session")?;
392 digest(&(
393 profile,
394 target,
395 &self.config.bundles,
396 self.config.targets.get(&source.target_template_id),
397 self.config.profiles.get(&source.last_profile),
398 ))
399 }
400
401 async fn validate_move_destination_configuration(
402 &self,
403 selection: &MoveSelection,
404 source_harness: mj_core::config::HarnessKind,
405 operational: &mj_core::relay::RelayOperationalState,
406 ) -> Result<()> {
407 let profile_id = selection
408 .profile_id
409 .as_deref()
410 .context("move profile is unresolved")?;
411 let profile = self
412 .config
413 .profiles
414 .get(profile_id)
415 .context("move profile no longer exists")?;
416 let source = self
417 .state
418 .sessions
419 .get(&selection.session_id)
420 .context("unknown move session")?;
421 if profile.kind != source_harness || profile_id == source.last_profile {
422 return Ok(());
423 }
424 let accepted = mj_core::acp::AcceptedSessionConfig::from_configuration(
425 &operational.config,
426 &operational.config_options,
427 );
428 if accepted.model.is_none() && accepted.effort.is_none() {
429 return Ok(());
430 }
431 let choices =
434 super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
435 .await
436 .with_context(|| {
437 format!("discover destination profile {profile_id:?} configuration")
438 })?;
439 validate_preserved_configuration(profile_id, &accepted, &choices)
440 }
441
442 pub async fn prepare_move_session_controlled(
443 &self,
444 mut selection: MoveSelection,
445 executor: &(impl CommandExecutor + Sync),
446 ) -> Result<MovePreparation> {
447 ensure!(
448 selection.profile_id.is_some() || selection.target_template_id.is_some(),
449 "move requires a target or profile selection"
450 );
451 let source = self
452 .state
453 .sessions
454 .get(&selection.session_id)
455 .context("unknown session")?;
456 ensure!(
457 !self.state.subagents.contains_key(&source.id),
458 "sub-agent sessions cannot move independently of their parent"
459 );
460 ensure!(
461 !self.state.subagents.values().any(|child| {
462 child.parent_session_id == source.id
463 && self
464 .state
465 .sessions
466 .get(&child.child_session_id)
467 .is_some_and(|session| session.state.is_active())
468 }),
469 "stop active sub-agents before moving their parent session"
470 );
471 let previous = crate::database::load_move_operation(&source.id)?;
472 let retry = previous.as_ref().is_some_and(|op| {
473 !matches!(op.phase, MovePhase::Completed) && source.checkpoint.is_some()
474 });
475 ensure!(
476 matches!(
477 source.state,
478 SessionState::Running | SessionState::Disconnected
479 ) || retry,
480 "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
481 source.id
482 );
483 selection
484 .profile_id
485 .get_or_insert_with(|| source.last_profile.clone());
486 selection
487 .target_template_id
488 .get_or_insert_with(|| source.target_template_id.clone());
489 selection
490 .additional_mounts
491 .get_or_insert_with(|| source.additional_mounts.clone());
492 ensure!(
493 !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
494 "select resource allocation or explicitly clear it, not both"
495 );
496 if !selection.clear_resource_allocation {
497 selection.resource_allocation = selection
498 .resource_allocation
499 .or_else(|| source.resource_allocation.clone());
500 }
501 let profile_id = selection.profile_id.as_deref().unwrap();
502 let target_id = selection.target_template_id.as_deref().unwrap();
503 let profile = self
504 .config
505 .profiles
506 .get(profile_id)
507 .context("unknown destination profile")?;
508 ensure!(
509 profile.enabled,
510 "destination profile {profile_id:?} is disabled"
511 );
512 let target = self
513 .config
514 .targets
515 .get(target_id)
516 .context("unknown destination target")?;
517 self.validate_muse_resume_destination(source, profile.kind, target_id)?;
518 super::worktree::resume_compatibility(source, &self.config, target_id)
519 .map_err(anyhow::Error::msg)?;
520 super::backend::validate_resource_allocation(
521 target,
522 selection.resource_allocation.as_ref(),
523 )?;
524 let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
525 ensure!(
526 profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
527 "Muse Code ACP supports one workspace root; attached directories are unsupported"
528 );
529 ensure!(
530 mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
531 "attached resources are unsupported for this target; select compatible resources explicitly"
532 );
533 crate::targets::validate_additional_mounts(mounts)?;
534 for mount in mounts {
535 self.validate_mount_source(target_id, &mount.source, executor)?;
536 }
537 let planned_conversion =
538 self.validate_move_destination_paths(source, target_id, executor)?;
539 ensure!(
540 profile.home.is_dir(),
541 "destination profile home is unavailable; configure the profile before moving"
542 );
543 super::worker_binary::preflight_worker_binary(target)?;
544 super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
545 let source_harness = previous
546 .as_ref()
547 .filter(|operation| {
548 !operation.queue_admission_started && operation.phase != MovePhase::Completed
549 })
550 .and_then(|operation| operation.recovery_session.as_ref())
551 .map_or(source.harness_kind, |record| record.harness_kind);
552 let cross_harness = profile.kind != source_harness;
553 if cross_harness {
554 let cancel = tokio_util::sync::CancellationToken::new();
555 let resolve =
556 crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
557 tokio::pin!(resolve);
558 loop {
559 tokio::select! {
560 result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
561 _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
562 if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
563 }
564 }
565 }
566 }
567 let conversion = planned_conversion
570 .map(|conversion| {
571 super::worktree::raw_conversion_preview(source, &conversion, executor)
572 .context("describe the move of this checkout into the target")
573 })
574 .transpose()?
575 .map(Box::new);
576 let (active, mut queued_commands, fingerprint) =
577 self.move_confirmation(&selection, conversion.as_deref())?;
578 replace_queued_images_with_placeholders(&mut queued_commands);
579 let operation_id = previous
580 .as_ref()
581 .filter(|operation| {
582 operation.selection == selection
583 && operation.phase != MovePhase::Completed
584 && operation.checkpoint.is_some()
585 })
586 .map(|operation| operation.operation_id.clone())
587 .unwrap_or(new_command_id("move")?);
588 Ok(MovePreparation {
589 source_unavailable: false,
590 in_place: in_place_move_eligible(
591 source,
592 &selection,
593 self.state.subagents.contains_key(&source.id),
594 retry,
595 ),
596 conversion,
597 selection,
598 source_profile_id: source.last_profile.clone(),
599 source_target_template_id: source.target_template_id.clone(),
600 cross_harness,
601 active,
602 queued_commands,
603 fingerprint,
604 operation_id,
605 })
606 }
607
608 pub async fn move_session_managed_controlled(
610 &mut self,
611 request: MoveSessionRequest,
612 executor: &(impl CommandExecutor + Sync),
613 manager: &SessionManagerControl,
614 ) -> Result<MoveOutcome> {
615 let prepared = &request.preparation;
616 let id = prepared.selection.session_id.clone();
617 let started = std::time::Instant::now();
618 executor.notify_notice("Checking destination");
619 let checked = {
620 let _checking_destination =
621 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
622 let mut checked = self
623 .prepare_move_session_controlled(prepared.selection.clone(), executor)
624 .await?;
625 if matches!(
628 self.state.sessions[&id].state,
629 SessionState::Running | SessionState::Disconnected
630 ) {
631 let source_harness = self.state.sessions[&id].harness_kind;
632 let snapshot = refresh_move_source(manager, &id).await?;
633 let (active, queue, fingerprint) =
634 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
635 checked.source_unavailable = snapshot
636 .as_ref()
637 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
638 checked.active = active
639 || checked.source_unavailable
640 || snapshot.as_ref().is_some_and(|snapshot| {
641 let mut operational = snapshot.operational.clone();
642 operational.queued_prompts.clear();
643 operational.checkpoint_barrier = None;
644 !operational.safe_to_replace(source_harness)
645 });
646 if let Some(snapshot) = &snapshot {
647 self.validate_move_destination_configuration(
648 &checked.selection,
649 source_harness,
650 &snapshot.operational,
651 )
652 .await?;
653 }
654 checked.queued_commands = queue;
655 checked.fingerprint = fingerprint;
656 }
657 checked
658 };
659 ensure!(
660 checked.fingerprint == prepared.fingerprint,
661 "session, pending work, or destination configuration changed; prepare and confirm Move again"
662 );
663 let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
664 let source = self.state.sessions[&id].clone();
665 let old_operation = crate::database::load_move_operation(&id)?;
666 let retry = old_operation.filter(|op| {
667 op.selection == prepared.selection
668 && op.phase != MovePhase::Completed
669 && op.checkpoint.is_some()
670 });
671 if retry.is_none()
672 && source.state == SessionState::Running
673 && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
674 && source.target_template_id == checked.selection.target_template_id.as_deref().unwrap()
675 && Some(&source.additional_mounts) == checked.selection.additional_mounts.as_ref()
676 && source.resource_allocation == checked.selection.resource_allocation
677 && (!checked.selection.clear_resource_allocation
678 || (source.container_cpus.is_none() && source.container_memory.is_none()))
679 {
680 return Ok(outcome(
681 &prepared.operation_id,
682 &prepared.selection,
683 "unchanged",
684 None,
685 None,
686 ));
687 }
688 ensure!(
689 !checked.active || request.acknowledge_interruption,
690 "active work will be interrupted; confirm Move again with interruption acknowledgement"
691 );
692 ensure!(
693 checked.queued_commands.is_empty() || request.queue.is_some(),
694 "pending work requires an explicit queue choice: discard or start"
695 );
696 let timestamp = now();
697 let mut operation = match retry {
698 Some(mut op) => {
699 ensure!(
700 !op.queue_admission_started || op.queue == queue,
701 "queued work may already have run; retry with the original queue choice on the same destination"
702 );
703 crate::database::clear_move_cancellation_for_retry(&id)?;
704 op.configuration_fingerprint =
705 self.move_configuration_fingerprint(&checked.selection)?;
706 op.queue = queue;
707 op.cancellation_requested = false;
708 op.error = None;
709 op
710 }
711 None => MoveOperation {
712 source_checkpoint_only: false,
713 in_place: in_place_move_eligible(
714 &source,
715 &checked.selection,
716 self.state.subagents.contains_key(&id),
717 false,
718 ),
719 operation_id: prepared.operation_id.clone(),
720 selection: checked.selection.clone(),
721 source_profile_id: source.last_profile.clone(),
722 source_target_template_id: source.target_template_id.clone(),
723 source_target: source.target.clone(),
724 source_native_session_id: source.native_session_id.clone(),
725 source_additional_mounts: source.additional_mounts.clone(),
726 source_resource_allocation: source.resource_allocation.clone(),
727 destination_target: None,
728 destination_native_session_id: None,
729 destination_store_id: None,
730 configuration_fingerprint: self
731 .move_configuration_fingerprint(&checked.selection)?,
732 checkpoint: None,
733 recovery_session: None,
734 queue,
735 phase: MovePhase::Preparing,
736 queue_admission_started: false,
737 queue_admission_finished: false,
738 cancellation_requested: false,
739 created_at: timestamp.clone(),
740 updated_at: timestamp,
741 error: None,
742 },
743 };
744 crate::database::save_move_operation(&operation)?;
745 tracing::info!(
746 session_id = id,
747 phase = "preflight",
748 elapsed_ms = started.elapsed().as_millis() as u64,
749 "move phase completed"
750 );
751 let result = self
752 .execute_move(&mut operation, Some(&checked), executor, manager)
753 .await;
754 self.finish_move_result(&mut operation, result, executor)
755 }
756
757 fn finish_move_result(
758 &self,
759 operation: &mut MoveOperation,
760 result: Result<()>,
761 executor: &impl CommandExecutor,
762 ) -> Result<MoveOutcome> {
763 let (status, error, recovery) = match result {
764 Ok(()) => {
765 operation.phase = MovePhase::Completed;
766 ("completed", None, None)
767 }
768 Err(error) => {
769 let cancelled =
770 executor.cancellation_requested() || operation.cancellation_requested;
771 operation.phase = if cancelled {
772 MovePhase::Cancelled
773 } else {
774 MovePhase::Failed
775 };
776 operation.cancellation_requested = cancelled;
777 let recovery = failed_move_recovery(
778 operation,
779 self.state.sessions.get(&operation.selection.session_id),
780 );
781 let error = format!("{error:#}");
782 operation.error = Some(error.clone());
783 (
784 if cancelled { "cancelled" } else { "failed" },
785 Some(error),
786 Some(recovery),
787 )
788 }
789 };
790 operation.updated_at = now();
791 crate::database::save_move_operation(operation)?;
792 Ok(outcome(
793 &operation.operation_id,
794 &operation.selection,
795 status,
796 error,
797 recovery,
798 ))
799 }
800
801 pub async fn recover_move_managed_controlled(
802 &mut self,
803 mut operation: MoveOperation,
804 executor: &(impl CommandExecutor + Sync),
805 manager: &SessionManagerControl,
806 ) -> Result<MoveOutcome> {
807 let id = operation.selection.session_id.clone();
808 let result = async {
809 let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
810 if operation.queue_admission_started {
811 ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
812 return self.admit_move_queue(&mut operation, executor).await;
813 }
814 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
815 let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
816 self.recover_move_source_stop(&mut operation, &cleanup, manager).await?;
817 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
818 if self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
819 self.cleanup_stopped_target(&id, &cleanup)?;
820 }
821 bail!("{}", interrupted_source_stop_message(
822 "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
823 operation.in_place,
824 ));
825 }
826 match operation.phase {
827 MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
828 MovePhase::ResumingDestination if session.state == SessionState::Running => {
829 ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
833 && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
834 "ready destination does not match the move intent");
835 operation.destination_target = session.target.clone();
836 operation.destination_native_session_id = session.native_session_id.clone();
837 operation.queue_admission_started = true;
838 operation.phase = MovePhase::StartingQueue;
839 crate::database::save_move_operation(&operation)?;
840 restore_move_queue_hold(&operation);
841 ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
842 self.admit_move_queue(&mut operation, executor).await
843 }
844 MovePhase::ResumingDestination => {
845 let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
846 let error = self.rollback_failed_resume(&id, previous, operation.in_place,
850 anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry"), executor)?;
851 Err(error)
852 }
853 MovePhase::ClosingSource => {
854 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
855 self.recover_move_source_stop(&mut operation, executor, manager).await?;
856 }
857 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
858 if self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
859 self.cleanup_stopped_target(&id, executor)?;
860 }
861 bail!("{}", interrupted_source_stop_message(
862 "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
863 operation.in_place,
864 ))
865 }
866 _ => bail!("Move requires an explicit retry after the daemon restarted"),
867 }
868 }.await;
869 self.finish_move_result(&mut operation, result, executor)
870 }
871
872 async fn recover_move_source_stop(
873 &mut self,
874 operation: &mut MoveOperation,
875 executor: &(impl CommandExecutor + Sync),
876 manager: &SessionManagerControl,
877 ) -> Result<()> {
878 let id = operation.selection.session_id.clone();
879 if self.state.sessions[&id].state == SessionState::Closing {
880 self.prepare_move_source_checkpoint(&id, executor, manager, operation)
881 .await?;
882 let handle = manager
883 .wait_for_session(&id, std::time::Duration::from_secs(5))
884 .await?;
885 let mut lease = handle.lease_connection().await?;
886 let execution = lease.connection_mut().sync().await?.operational.execution;
887 lease.release();
888 if matches!(
889 execution,
890 mj_core::relay::RelayExecutionState::Idle
891 | mj_core::relay::RelayExecutionState::Running
892 ) {
893 if operation.cancellation_requested {
894 let record = self.state.sessions.get_mut(&id).unwrap();
895 record.state = SessionState::Running;
896 record.updated_at = now();
897 record.last_error =
898 Some("Move was cancelled before the source was sealed".into());
899 crate::database::save_lifecycle_session(record)?;
900 } else {
901 self.suspend_session_for_move(
904 &id,
905 executor,
906 manager,
907 operation,
908 None,
909 SourceTargetDisposition::Destroy,
910 )
911 .await?;
912 }
913 return Ok(());
914 }
915 }
916 self.recover_interrupted_close_managed(&id, executor, manager, true, None)
918 .await?;
919 Ok(())
920 }
921
922 async fn execute_move(
923 &mut self,
924 operation: &mut MoveOperation,
925 preparation: Option<&MovePreparation>,
926 executor: &(impl CommandExecutor + Sync),
927 manager: &SessionManagerControl,
928 ) -> Result<()> {
929 let id = operation.selection.session_id.clone();
930 ensure!(
931 !executor.cancellation_requested(),
932 "move cancelled before source interruption"
933 );
934 if !operation.queue_admission_started {
935 if self.state.sessions[&id].state == SessionState::Error
936 && let Some(previous) = operation.recovery_session.as_ref()
937 {
938 let failure = self.rollback_failed_resume(
941 &id,
942 previous,
943 false,
944 anyhow::anyhow!("clean up the partial Move destination before retry"),
945 executor,
946 )?;
947 ensure!(
948 self.state.sessions[&id].state == SessionState::Stopped,
949 "{failure:#}"
950 );
951 }
952 let state = self.state.sessions[&id].state;
953 if matches!(state, SessionState::Closing | SessionState::Destroying) {
954 self.recover_move_source_stop(operation, executor, manager)
955 .await?;
956 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
957 } else if matches!(state, SessionState::Running | SessionState::Disconnected)
958 && operation.destination_target.is_none()
959 {
960 executor.notify_notice("Stopping source");
961 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
962 operation.phase = MovePhase::ClosingSource;
963 operation.updated_at = now();
964 crate::database::save_move_operation(operation)?;
965 let disposition = if operation.in_place {
969 SourceTargetDisposition::RetainForInPlaceSwap
970 } else {
971 SourceTargetDisposition::Destroy
972 };
973 self.suspend_session_for_move(
974 &id,
975 executor,
976 manager,
977 operation,
978 preparation,
979 disposition,
980 )
981 .await?;
982 }
983 if !operation.in_place
986 && self.state.sessions[&id].state == SessionState::Stopped
987 && self.state.sessions[&id].target.is_some()
988 {
989 executor.notify_notice("Cleaning up source");
990 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
991 self.cleanup_stopped_target(&id, executor)?;
992 }
993 operation.checkpoint = operation
994 .checkpoint
995 .clone()
996 .or_else(|| self.state.sessions[&id].checkpoint.clone());
997 ensure!(
998 operation.checkpoint.is_some(),
999 "move has no verified checkpoint"
1000 );
1001 ensure!(
1002 !executor.cancellation_requested(),
1003 "move cancelled after source teardown; session is stopped"
1004 );
1005 operation.phase = MovePhase::ResumingDestination;
1006 operation.recovery_session = Some(self.state.sessions[&id].clone());
1007 operation.updated_at = now();
1008 crate::database::save_move_operation(operation)?;
1009 executor.notify_notice("Preparing destination");
1010 if operation.in_place {
1011 self.restore_session_in_place(
1015 &id,
1016 operation.selection.profile_id.as_deref().unwrap(),
1017 executor,
1018 )
1019 .await?;
1020 } else {
1021 if operation.selection.clear_resource_allocation {
1022 let session = self.state.sessions.get_mut(&id).unwrap();
1023 session.resource_allocation = None;
1024 session.container_cpus = None;
1025 session.container_memory = None;
1026 crate::database::save_session(session)?;
1027 }
1028 self.resume_session_controlled(
1029 &id,
1030 operation.selection.profile_id.as_deref().unwrap(),
1031 operation.selection.target_template_id.as_deref().unwrap(),
1032 SessionResumeOptions {
1033 additional_mounts: operation.selection.additional_mounts.clone(),
1034 resource_allocation: operation.selection.resource_allocation.clone(),
1035 discard_queue: true,
1036 },
1037 executor,
1038 )
1039 .await?;
1040 }
1041 let destination = &self.state.sessions[&id];
1042 operation.destination_target = destination.target.clone();
1043 operation.destination_native_session_id = destination.native_session_id.clone();
1044 operation.phase = MovePhase::StartingQueue;
1046 operation.queue_admission_started = true;
1047 operation.updated_at = now();
1048 crate::database::save_move_operation(operation)?;
1049 }
1050 restore_move_queue_hold(operation);
1051 self.admit_move_queue(operation, executor).await
1052 }
1053
1054 async fn admit_move_queue(
1055 &self,
1056 operation: &mut MoveOperation,
1057 executor: &(impl CommandExecutor + Sync),
1058 ) -> Result<()> {
1059 let timing_id = operation.selection.session_id.clone();
1060 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
1061 let id = &operation.selection.session_id;
1062 let mut relay = {
1063 let _checking_destination =
1064 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1065 let destination = &self.state.sessions[id];
1066 ensure!(
1067 destination.state == SessionState::Running
1068 && destination.target == operation.destination_target
1069 && destination.native_session_id == operation.destination_native_session_id,
1070 "cannot prove the same ready destination; refusing to replay potentially executed work"
1071 );
1072 let spec = self.reconnect_command(id)?;
1073 let relay = StandaloneSession::connect_command(&spec, id).await?;
1074 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")?;
1075 if let Some(expected) = &operation.destination_store_id {
1076 ensure!(
1077 *expected == store_id,
1078 "destination relay storage was replaced; refusing to replay potentially executed work"
1079 );
1080 } else {
1081 operation.destination_store_id = Some(store_id);
1083 crate::database::save_move_operation(operation)?;
1084 }
1085 ensure!(
1086 relay.snapshot().operational.native_session_id
1087 == operation.destination_native_session_id,
1088 "destination relay native identity changed; refusing queue replay"
1089 );
1090 relay
1091 };
1092 if operation.queue == ResumeQueueDisposition::Start && !operation.queue_admission_finished {
1093 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
1094 executor.notify_notice("Starting queued work");
1095 let checkpoint = operation
1096 .checkpoint
1097 .as_ref()
1098 .context("move queue archive is missing")?;
1099 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1100 ensure!(
1101 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
1102 "move queue checkpoint verification failed"
1103 );
1104 for queued in verified.canonical_session.queued_prompts {
1105 ensure!(
1106 !executor.cancellation_requested(),
1107 "move cancelled during queue admission; destination retained"
1108 );
1109 let command = match queued.kind {
1110 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
1111 prompt: queued
1112 .content
1113 .into_iter()
1114 .map(serde_json::from_value)
1115 .collect::<serde_json::Result<_>>()?,
1116 },
1117 CanonicalQueuedCommandKind::SetConfig { key, value } => {
1118 RelayCommand::SetConfig { key, value }
1119 }
1120 };
1121 relay.submit(queued.command_id, command).await?;
1122 }
1123 }
1124 operation.queue_admission_finished = true;
1125 crate::database::save_move_operation(operation)?;
1126 restore_move_queue_hold(operation);
1127 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
1128 "Queued work was discarded; ready and idle."
1129 } else {
1130 "Queued work was accepted."
1131 };
1132 let source_profile = &operation.source_profile_id;
1133 let source_target = &operation.source_target_template_id;
1134 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
1135 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
1136 let text = if operation.in_place {
1139 format!(
1140 "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."
1141 )
1142 } else {
1143 format!(
1144 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
1145 )
1146 };
1147 relay
1148 .submit(
1149 format!("{}-notice", operation.operation_id),
1150 RelayCommand::RecordNotice { text },
1151 )
1152 .await?;
1153 Ok(())
1154 }
1155
1156 pub(super) fn validate_move_checkpoint(
1157 &self,
1158 operation: &MoveOperation,
1159 preparation: Option<&MovePreparation>,
1160 executor: &(impl CommandExecutor + Sync),
1161 ) -> Result<()> {
1162 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1163 let current = Controller {
1164 config: mj_core::config::Config::load()?,
1165 state: self.state.clone(),
1166 };
1167 ensure!(
1168 current.move_configuration_fingerprint(&operation.selection)?
1169 == operation.configuration_fingerprint,
1170 "destination configuration changed during move"
1171 );
1172 let id = &operation.selection.session_id;
1173 current.validate_move_destination_paths(
1174 &self.state.sessions[id],
1175 operation.selection.target_template_id.as_deref().unwrap(),
1176 executor,
1177 )?;
1178 if let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
1179 .preflight_resume_repository_sources(
1180 id,
1181 operation.selection.target_template_id.as_deref().unwrap(),
1182 executor,
1183 )?
1184 {
1185 bail!(
1186 "destination repository source is missing checkpoint commit {}; source retained",
1187 mismatch.missing_commit
1188 );
1189 }
1190 if let Some(prepared) = preparation {
1191 let checkpoint = self.state.sessions[id]
1192 .checkpoint
1193 .as_ref()
1194 .context("no move checkpoint")?;
1195 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1196 let actual: Vec<_> = verified
1197 .canonical_session
1198 .queued_prompts
1199 .iter()
1200 .map(|p| p.command_id.as_str())
1201 .collect();
1202 let expected: Vec<_> = prepared
1203 .queued_commands
1204 .iter()
1205 .map(|p| p.command_id.as_str())
1206 .collect();
1207 ensure!(
1208 actual == expected,
1209 "pending queue changed before checkpoint capture; source retained, confirm Move again"
1210 );
1211 }
1212 Ok(())
1213 }
1214}
1215
1216fn validate_preserved_configuration(
1217 profile_id: &str,
1218 accepted: &mj_core::acp::AcceptedSessionConfig,
1219 choices: &mj_core::worker_launch::ProfileConfig,
1220) -> Result<()> {
1221 for (key, value, offered) in [
1222 ("model", accepted.model.as_deref(), &choices.models),
1223 ("effort", accepted.effort.as_deref(), &choices.efforts),
1224 ] {
1225 let Some(value) = value else { continue };
1226 ensure!(
1227 offered.iter().any(|choice| choice.value == value),
1228 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
1229 offered
1230 .iter()
1231 .map(|choice| choice.value.as_str())
1232 .collect::<Vec<_>>()
1233 .join(", ")
1234 );
1235 }
1236 Ok(())
1237}
1238
1239pub(super) fn in_place_move_eligible(
1246 source: &mj_core::state::SessionRecord,
1247 selection: &MoveSelection,
1248 is_subagent: bool,
1249 retry: bool,
1250) -> bool {
1251 !retry
1252 && !is_subagent
1253 && source.target.is_some()
1254 && matches!(
1255 source.state,
1256 SessionState::Running | SessionState::Disconnected
1257 )
1258 && Some(&source.target_template_id) == selection.target_template_id.as_ref()
1259 && Some(&source.additional_mounts) == selection.additional_mounts.as_ref()
1260 && source.resource_allocation == selection.resource_allocation
1261 && !selection.clear_resource_allocation
1262}
1263
1264fn outcome(
1265 operation_id: &str,
1266 selection: &MoveSelection,
1267 status: &str,
1268 error: Option<String>,
1269 recovery: Option<String>,
1270) -> MoveOutcome {
1271 MoveOutcome {
1272 operation_id: operation_id.into(),
1273 session_id: selection.session_id.clone(),
1274 profile_id: selection.profile_id.clone().unwrap_or_default(),
1275 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
1276 outcome: status.into(),
1277 error,
1278 recovery,
1279 }
1280}