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