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)
917 .await?;
918 Ok(())
919 }
920
921 async fn execute_move(
922 &mut self,
923 operation: &mut MoveOperation,
924 preparation: Option<&MovePreparation>,
925 executor: &(impl CommandExecutor + Sync),
926 manager: &SessionManagerControl,
927 ) -> Result<()> {
928 let id = operation.selection.session_id.clone();
929 ensure!(
930 !executor.cancellation_requested(),
931 "move cancelled before source interruption"
932 );
933 if !operation.queue_admission_started {
934 if self.state.sessions[&id].state == SessionState::Error
935 && let Some(previous) = operation.recovery_session.as_ref()
936 {
937 let failure = self.rollback_failed_resume(
940 &id,
941 previous,
942 false,
943 anyhow::anyhow!("clean up the partial Move destination before retry"),
944 executor,
945 )?;
946 ensure!(
947 self.state.sessions[&id].state == SessionState::Stopped,
948 "{failure:#}"
949 );
950 }
951 let state = self.state.sessions[&id].state;
952 if matches!(state, SessionState::Closing | SessionState::Destroying) {
953 self.recover_move_source_stop(operation, executor, manager)
954 .await?;
955 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
956 } else if matches!(state, SessionState::Running | SessionState::Disconnected)
957 && operation.destination_target.is_none()
958 {
959 executor.notify_notice("Stopping source");
960 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
961 operation.phase = MovePhase::ClosingSource;
962 operation.updated_at = now();
963 crate::database::save_move_operation(operation)?;
964 let disposition = if operation.in_place {
968 SourceTargetDisposition::RetainForInPlaceSwap
969 } else {
970 SourceTargetDisposition::Destroy
971 };
972 self.suspend_session_for_move(
973 &id,
974 executor,
975 manager,
976 operation,
977 preparation,
978 disposition,
979 )
980 .await?;
981 }
982 if !operation.in_place
985 && self.state.sessions[&id].state == SessionState::Stopped
986 && self.state.sessions[&id].target.is_some()
987 {
988 executor.notify_notice("Cleaning up source");
989 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
990 self.cleanup_stopped_target(&id, executor)?;
991 }
992 operation.checkpoint = operation
993 .checkpoint
994 .clone()
995 .or_else(|| self.state.sessions[&id].checkpoint.clone());
996 ensure!(
997 operation.checkpoint.is_some(),
998 "move has no verified checkpoint"
999 );
1000 ensure!(
1001 !executor.cancellation_requested(),
1002 "move cancelled after source teardown; session is stopped"
1003 );
1004 operation.phase = MovePhase::ResumingDestination;
1005 operation.recovery_session = Some(self.state.sessions[&id].clone());
1006 operation.updated_at = now();
1007 crate::database::save_move_operation(operation)?;
1008 executor.notify_notice("Preparing destination");
1009 if operation.in_place {
1010 self.restore_session_in_place(
1014 &id,
1015 operation.selection.profile_id.as_deref().unwrap(),
1016 executor,
1017 )
1018 .await?;
1019 } else {
1020 if operation.selection.clear_resource_allocation {
1021 let session = self.state.sessions.get_mut(&id).unwrap();
1022 session.resource_allocation = None;
1023 session.container_cpus = None;
1024 session.container_memory = None;
1025 crate::database::save_session(session)?;
1026 }
1027 self.resume_session_controlled(
1028 &id,
1029 operation.selection.profile_id.as_deref().unwrap(),
1030 operation.selection.target_template_id.as_deref().unwrap(),
1031 SessionResumeOptions {
1032 additional_mounts: operation.selection.additional_mounts.clone(),
1033 resource_allocation: operation.selection.resource_allocation.clone(),
1034 discard_queue: true,
1035 },
1036 executor,
1037 )
1038 .await?;
1039 }
1040 let destination = &self.state.sessions[&id];
1041 operation.destination_target = destination.target.clone();
1042 operation.destination_native_session_id = destination.native_session_id.clone();
1043 operation.phase = MovePhase::StartingQueue;
1045 operation.queue_admission_started = true;
1046 operation.updated_at = now();
1047 crate::database::save_move_operation(operation)?;
1048 }
1049 restore_move_queue_hold(operation);
1050 self.admit_move_queue(operation, executor).await
1051 }
1052
1053 async fn admit_move_queue(
1054 &self,
1055 operation: &mut MoveOperation,
1056 executor: &(impl CommandExecutor + Sync),
1057 ) -> Result<()> {
1058 let timing_id = operation.selection.session_id.clone();
1059 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
1060 let id = &operation.selection.session_id;
1061 let mut relay = {
1062 let _checking_destination =
1063 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1064 let destination = &self.state.sessions[id];
1065 ensure!(
1066 destination.state == SessionState::Running
1067 && destination.target == operation.destination_target
1068 && destination.native_session_id == operation.destination_native_session_id,
1069 "cannot prove the same ready destination; refusing to replay potentially executed work"
1070 );
1071 let spec = self.reconnect_command(id)?;
1072 let relay = StandaloneSession::connect_command(&spec, id).await?;
1073 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")?;
1074 if let Some(expected) = &operation.destination_store_id {
1075 ensure!(
1076 *expected == store_id,
1077 "destination relay storage was replaced; refusing to replay potentially executed work"
1078 );
1079 } else {
1080 operation.destination_store_id = Some(store_id);
1082 crate::database::save_move_operation(operation)?;
1083 }
1084 ensure!(
1085 relay.snapshot().operational.native_session_id
1086 == operation.destination_native_session_id,
1087 "destination relay native identity changed; refusing queue replay"
1088 );
1089 relay
1090 };
1091 if operation.queue == ResumeQueueDisposition::Start && !operation.queue_admission_finished {
1092 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
1093 executor.notify_notice("Starting queued work");
1094 let checkpoint = operation
1095 .checkpoint
1096 .as_ref()
1097 .context("move queue archive is missing")?;
1098 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1099 ensure!(
1100 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
1101 "move queue checkpoint verification failed"
1102 );
1103 for queued in verified.canonical_session.queued_prompts {
1104 ensure!(
1105 !executor.cancellation_requested(),
1106 "move cancelled during queue admission; destination retained"
1107 );
1108 let command = match queued.kind {
1109 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
1110 prompt: queued
1111 .content
1112 .into_iter()
1113 .map(serde_json::from_value)
1114 .collect::<serde_json::Result<_>>()?,
1115 },
1116 CanonicalQueuedCommandKind::SetConfig { key, value } => {
1117 RelayCommand::SetConfig { key, value }
1118 }
1119 };
1120 relay.submit(queued.command_id, command).await?;
1121 }
1122 }
1123 operation.queue_admission_finished = true;
1124 crate::database::save_move_operation(operation)?;
1125 restore_move_queue_hold(operation);
1126 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
1127 "Queued work was discarded; ready and idle."
1128 } else {
1129 "Queued work was accepted."
1130 };
1131 let source_profile = &operation.source_profile_id;
1132 let source_target = &operation.source_target_template_id;
1133 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
1134 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
1135 let text = if operation.in_place {
1138 format!(
1139 "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."
1140 )
1141 } else {
1142 format!(
1143 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
1144 )
1145 };
1146 relay
1147 .submit(
1148 format!("{}-notice", operation.operation_id),
1149 RelayCommand::RecordNotice { text },
1150 )
1151 .await?;
1152 Ok(())
1153 }
1154
1155 pub(super) fn validate_move_checkpoint(
1156 &self,
1157 operation: &MoveOperation,
1158 preparation: Option<&MovePreparation>,
1159 executor: &(impl CommandExecutor + Sync),
1160 ) -> Result<()> {
1161 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1162 let current = Controller {
1163 config: mj_core::config::Config::load()?,
1164 state: self.state.clone(),
1165 };
1166 ensure!(
1167 current.move_configuration_fingerprint(&operation.selection)?
1168 == operation.configuration_fingerprint,
1169 "destination configuration changed during move"
1170 );
1171 let id = &operation.selection.session_id;
1172 current.validate_move_destination_paths(
1173 &self.state.sessions[id],
1174 operation.selection.target_template_id.as_deref().unwrap(),
1175 executor,
1176 )?;
1177 if let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
1178 .preflight_resume_repository_sources(
1179 id,
1180 operation.selection.target_template_id.as_deref().unwrap(),
1181 executor,
1182 )?
1183 {
1184 bail!(
1185 "destination repository source is missing checkpoint commit {}; source retained",
1186 mismatch.missing_commit
1187 );
1188 }
1189 if let Some(prepared) = preparation {
1190 let checkpoint = self.state.sessions[id]
1191 .checkpoint
1192 .as_ref()
1193 .context("no move checkpoint")?;
1194 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1195 let actual: Vec<_> = verified
1196 .canonical_session
1197 .queued_prompts
1198 .iter()
1199 .map(|p| p.command_id.as_str())
1200 .collect();
1201 let expected: Vec<_> = prepared
1202 .queued_commands
1203 .iter()
1204 .map(|p| p.command_id.as_str())
1205 .collect();
1206 ensure!(
1207 actual == expected,
1208 "pending queue changed before checkpoint capture; source retained, confirm Move again"
1209 );
1210 }
1211 Ok(())
1212 }
1213}
1214
1215fn validate_preserved_configuration(
1216 profile_id: &str,
1217 accepted: &mj_core::acp::AcceptedSessionConfig,
1218 choices: &mj_core::worker_launch::ProfileConfig,
1219) -> Result<()> {
1220 for (key, value, offered) in [
1221 ("model", accepted.model.as_deref(), &choices.models),
1222 ("effort", accepted.effort.as_deref(), &choices.efforts),
1223 ] {
1224 let Some(value) = value else { continue };
1225 ensure!(
1226 offered.iter().any(|choice| choice.value == value),
1227 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
1228 offered
1229 .iter()
1230 .map(|choice| choice.value.as_str())
1231 .collect::<Vec<_>>()
1232 .join(", ")
1233 );
1234 }
1235 Ok(())
1236}
1237
1238pub(super) fn in_place_move_eligible(
1245 source: &mj_core::state::SessionRecord,
1246 selection: &MoveSelection,
1247 is_subagent: bool,
1248 retry: bool,
1249) -> bool {
1250 !retry
1251 && !is_subagent
1252 && source.target.is_some()
1253 && matches!(
1254 source.state,
1255 SessionState::Running | SessionState::Disconnected
1256 )
1257 && Some(&source.target_template_id) == selection.target_template_id.as_ref()
1258 && Some(&source.additional_mounts) == selection.additional_mounts.as_ref()
1259 && source.resource_allocation == selection.resource_allocation
1260 && !selection.clear_resource_allocation
1261}
1262
1263fn outcome(
1264 operation_id: &str,
1265 selection: &MoveSelection,
1266 status: &str,
1267 error: Option<String>,
1268 recovery: Option<String>,
1269) -> MoveOutcome {
1270 MoveOutcome {
1271 operation_id: operation_id.into(),
1272 session_id: selection.session_id.clone(),
1273 profile_id: selection.profile_id.clone().unwrap_or_default(),
1274 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
1275 outcome: status.into(),
1276 error,
1277 recovery,
1278 }
1279}