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
24pub(in crate::controller) enum SubagentMutationDrain {
25 Drained,
26 UnsupportedWorkerProtocol(u32),
27}
28
29fn move_ownership() -> &'static std::sync::Mutex<std::collections::BTreeMap<String, MoveOwnership>>
30{
31 static OWNER: std::sync::OnceLock<
32 std::sync::Mutex<std::collections::BTreeMap<String, MoveOwnership>>,
33 > = std::sync::OnceLock::new();
34 OWNER.get_or_init(Default::default)
35}
36
37pub fn move_owns_session(session_id: &str) -> bool {
38 move_ownership()
39 .lock()
40 .unwrap_or_else(std::sync::PoisonError::into_inner)
41 .contains_key(session_id)
42}
43
44pub fn move_has_pending_queue(session_id: &str) -> bool {
45 matches!(
46 move_ownership()
47 .lock()
48 .unwrap_or_else(std::sync::PoisonError::into_inner)
49 .get(session_id),
50 Some(MoveOwnership::ExecutingQueue | MoveOwnership::PendingQueue)
51 )
52}
53
54pub fn release_move_queue_hold(session_id: &str) {
55 set_move_queue_hold(session_id, false);
56}
57
58fn set_move_queue_hold(session_id: &str, pending: bool) {
59 let mut owner = move_ownership()
60 .lock()
61 .unwrap_or_else(std::sync::PoisonError::into_inner);
62 let next = match (owner.get(session_id), pending) {
63 (Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue), true) => {
64 Some(MoveOwnership::ExecutingQueue)
65 }
66 (Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue), false) => {
67 Some(MoveOwnership::Executing)
68 }
69 (_, true) => Some(MoveOwnership::PendingQueue),
70 (_, false) => None,
71 };
72 if let Some(next) = next {
73 owner.insert(session_id.to_owned(), next);
74 } else {
75 owner.remove(session_id);
76 }
77}
78
79pub fn restore_move_queue_hold(operation: &MoveOperation) {
80 set_move_queue_hold(
81 &operation.selection.session_id,
82 operation.queue_admission_started && !operation.queue_admission_finished,
83 );
84}
85
86fn interrupted_source_stop_message(recovered: &str, in_place: bool) -> String {
90 if in_place {
91 format!("{recovered}; the in-place swap was interrupted; the environment was retained")
92 } else {
93 recovered.to_owned()
94 }
95}
96
97fn source_stopped_with_verified_checkpoint(record: &mj_core::state::SessionRecord) -> bool {
103 record.state == SessionState::Stopped && record.checkpoint.is_some()
104}
105
106fn stopped_source_recovery(
112 session_id: &str,
113 destination_profile: Option<&str>,
114 destination_target: Option<&str>,
115) -> String {
116 let flag = |name: &str, value: Option<&str>| {
117 value
118 .filter(|value| !value.is_empty())
119 .map(|value| format!(" --{name} {value}"))
120 .unwrap_or_default()
121 };
122 format!(
123 "Source is stopped with a verified checkpoint. Bring it back with \
124 `mj resume --session {session_id}{}{} --queue start`.",
125 flag("profile", destination_profile),
126 flag("target", destination_target),
127 )
128}
129
130fn failed_move_recovery(
139 operation: &MoveOperation,
140 record: Option<&mj_core::state::SessionRecord>,
141) -> String {
142 if let Some(destination) = &operation.prepared_destination
143 && matches!(
144 destination.state,
145 PreparedDestinationState::CleanupPending { .. }
146 )
147 {
148 return format!(
149 "Source and recovery data retained. Automatic EC2 cleanup is pending for {}; charges may continue until termination is confirmed.",
150 destination
151 .instance_id()
152 .unwrap_or("the recorded launch token")
153 );
154 }
155 if operation.prepared_destination.is_some()
156 && record.is_some_and(|r| {
157 matches!(r.state, SessionState::Running | SessionState::Disconnected)
158 && r.target == operation.source_target
159 })
160 {
161 return "Source retained and still running. EC2 destination cleaned up; prepare Move again when ready.".into();
162 }
163 if operation.queue_admission_started {
164 return "Destination is live; retry queue admission on this same destination. Already accepted work may have effects.".to_owned();
165 }
166 let source_live = record.is_some_and(|record| {
167 matches!(
168 record.state,
169 SessionState::Running | SessionState::Disconnected
170 )
171 });
172 let environment_held = operation.holds_source_environment()
173 || (operation.in_place
174 && !source_live
175 && record.is_some_and(|record| record.target.is_some()));
176 if environment_held && !operation.checkpoint_retained() {
177 return missing_move_checkpoint_recovery(record);
178 }
179 if operation.in_place && record.is_some_and(|record| record.target.is_some()) {
180 return if source_live {
181 "Source retained and still running. Retry Move when ready.".to_owned()
182 } else {
183 "Environment and checkpoint retained. Retry Move on the same target; the checkout will not be recreated.".to_owned()
184 };
185 }
186 match record {
187 Some(record) if source_stopped_with_verified_checkpoint(record) => {
188 stopped_source_recovery(
189 &operation.selection.session_id,
190 operation.selection.profile_id.as_deref(),
191 operation.selection.target_template_id.as_deref(),
192 )
193 }
194 _ => "Source or partial destination is retained. Retry move after resolving the reported error.".to_owned(),
195 }
196}
197
198fn missing_move_checkpoint_recovery(record: Option<&mj_core::state::SessionRecord>) -> String {
203 let checkout = record
204 .and_then(|record| {
205 record
206 .managed_worktree
207 .as_ref()
208 .map(|checkout| checkout.worktree_root.clone())
209 .or_else(|| record.project_directory.clone())
210 })
211 .map(|path| format!(" in {}", path.display()))
212 .unwrap_or_default();
213 format!(
214 "The Move's checkpoint archive is missing, so the Move cannot be retried and Resume cannot restore the session. \
215 The environment and checkout{checkout} are retained; copy out anything you need, then destroy the session."
216 )
217}
218
219fn failed_move_message(
222 phase: Option<&str>,
223 cancelled: bool,
224 recovery: &str,
225 operation_id: &str,
226) -> String {
227 format!(
228 "{}{}{}. {recovery} The daemon log records the reason under reference {operation_id}",
229 mj_core::state::MOVE_FAILURE_PREFIX,
230 phase
231 .map(|phase| format!(" while {phase}"))
232 .unwrap_or_default(),
233 if cancelled { " (cancelled)" } else { "" },
234 )
235}
236
237pub(crate) fn forget_missing_move_archives(operation: &mut MoveOperation) -> bool {
245 let missing = |checkpoint: &Option<mj_core::state::CheckpointMetadata>| {
246 checkpoint.as_ref().is_some_and(|checkpoint| {
247 matches!(
248 std::fs::symlink_metadata(&checkpoint.archive_path),
249 Err(error) if error.kind() == std::io::ErrorKind::NotFound
250 )
251 })
252 };
253 let lost = if missing(&operation.handoff) {
254 [operation.handoff.take(), operation.checkpoint.take()]
255 } else if operation.handoff.is_none() && missing(&operation.checkpoint) {
256 [None, operation.checkpoint.take()]
257 } else {
258 return false;
259 };
260 for checkpoint in lost.iter().flatten() {
261 tracing::warn!(
262 session_id = %operation.selection.session_id,
263 reference = %operation.operation_id,
264 path = %checkpoint.archive_path.display(),
265 "a Move's checkpoint archive is missing; recording that the Move cannot restore it"
266 );
267 }
268 true
269}
270
271pub(crate) fn record_missing_move_archives(
275 state: &mj_core::state::State,
276 operations: Vec<MoveOperation>,
277) -> Result<()> {
278 for mut operation in operations {
279 if !operation.retains_checkpoint() || !forget_missing_move_archives(&mut operation) {
280 continue;
281 }
282 let record = state.sessions.get(&operation.selection.session_id);
283 let published = record.and_then(|record| record.last_error.as_deref());
284 if operation.is_active()
285 || !published
286 .is_some_and(|error| error.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
287 {
288 crate::database::save_move_operation(&operation)?;
289 continue;
290 }
291 record_finished_move_recovery(state, &mut operation)?;
292 }
293 Ok(())
294}
295
296pub(crate) fn record_finished_move_recovery(
298 state: &mj_core::state::State,
299 operation: &mut MoveOperation,
300) -> Result<()> {
301 let record = state.sessions.get(&operation.selection.session_id);
302 operation.updated_at = now();
303 if !operation.is_active()
304 && record
305 .and_then(|record| record.last_error.as_deref())
306 .is_some_and(|message| message.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
307 {
308 let message = failed_move_message(
309 None,
310 operation.phase == MovePhase::Cancelled,
311 &failed_move_recovery(operation, record),
312 &operation.operation_id,
313 );
314 crate::database::save_move_outcome(operation, Some(&message))
315 } else {
316 crate::database::save_move_operation(operation)
317 }
318}
319
320fn sealed_selection_difference(retained: &MoveSelection, requested: &MoveSelection) -> String {
323 let mut parts = Vec::new();
324 for (name, flag, retained, requested) in [
325 (
326 "profile",
327 "--profile",
328 &retained.profile_id,
329 &requested.profile_id,
330 ),
331 (
332 "target",
333 "--target",
334 &retained.target_template_id,
335 &requested.target_template_id,
336 ),
337 ] {
338 if retained != requested {
339 parts.push(format!(
340 "{name} (sealed with {0}; pass {flag} {0})",
341 retained.as_deref().unwrap_or_default()
342 ));
343 }
344 }
345 if retained.workspace.acknowledge_large_transfer
346 != requested.workspace.acknowledge_large_transfer
347 {
348 parts.push(if retained.workspace.acknowledge_large_transfer {
349 "large-transfer acknowledgement (the Move was sealed with it; pass --allow-large-transfer)".to_owned()
350 } else {
351 "large-transfer acknowledgement (the Move was sealed without it; omit --allow-large-transfer)".to_owned()
352 });
353 }
354 if retained.workspace.exclusions != requested.workspace.exclusions {
355 let excluded = retained
356 .workspace
357 .exclusions
358 .iter()
359 .map(|path| format!("{}:{}", path.repository, path.path.display()))
360 .collect::<Vec<_>>();
361 parts.push(if excluded.is_empty() {
362 "excluded files (the Move was sealed excluding none)".to_owned()
363 } else {
364 format!(
365 "excluded files (the Move was sealed excluding exactly {})",
366 excluded.join(", ")
367 )
368 });
369 }
370 if retained.additional_mounts != requested.additional_mounts {
371 parts.push("attached directories (send the ones the Move was sealed with)".to_owned());
372 }
373 if retained.resource_allocation != requested.resource_allocation
374 || retained.clear_resource_allocation != requested.clear_resource_allocation
375 {
376 parts.push("resource allocation (send the one the Move was sealed with)".to_owned());
377 }
378 if retained.subagents != requested.subagents {
379 parts.push("sub-agent policy (send the one the Move was sealed with)".to_owned());
380 }
381 if retained.session_id != requested.session_id || parts.is_empty() {
382 parts.push("the session".to_owned());
383 }
384 format!(
385 "a sealed Move must be retried with the selection it was sealed with; this request differs in its {}",
386 parts.join("; ")
387 )
388}
389
390pub struct MoveMutationGuard(String);
391
392impl MoveMutationGuard {
393 pub fn reserve(session_id: &str) -> Result<Self> {
394 let mut owner = move_ownership()
395 .lock()
396 .unwrap_or_else(std::sync::PoisonError::into_inner);
397 ensure!(
398 !matches!(
399 owner.get(session_id),
400 Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue)
401 ),
402 "session already has a move owner"
403 );
404 let next = if owner.contains_key(session_id) {
405 MoveOwnership::ExecutingQueue
406 } else {
407 MoveOwnership::Executing
408 };
409 owner.insert(session_id.to_owned(), next);
410 Ok(Self(session_id.to_owned()))
411 }
412}
413
414impl Drop for MoveMutationGuard {
415 fn drop(&mut self) {
416 let mut owner = move_ownership()
417 .lock()
418 .unwrap_or_else(std::sync::PoisonError::into_inner);
419 match owner.get(&self.0) {
420 Some(MoveOwnership::ExecutingQueue) => {
421 owner.insert(self.0.clone(), MoveOwnership::PendingQueue);
422 }
423 Some(MoveOwnership::Executing) => {
424 owner.remove(&self.0);
425 }
426 _ => {}
427 }
428 }
429}
430
431pub(crate) fn move_refuses_command(session_id: &str, command: &RelayCommand) -> bool {
432 move_owns_session(session_id)
433 && matches!(
434 command,
435 RelayCommand::Prompt { .. }
436 | RelayCommand::SetConfig { .. }
437 | RelayCommand::SetSessionMode { .. }
438 | RelayCommand::RestoreExecutionMode
439 | RelayCommand::RunUserShell { .. }
440 | RelayCommand::CancelUserShell { .. }
441 | RelayCommand::Cancel
442 | RelayCommand::RemoveQueuedPrompt { .. }
443 | RelayCommand::ClearQueuedPrompts
444 )
445}
446use crate::session_manager::{SessionManagerControl, StandaloneSession, new_command_id};
447use mj_checkpoint::archive::{CanonicalQueuedCommandKind, verify_archive_streaming};
448use mj_core::state::{
449 DestinationChecks, MoveOperation, MovePhase, PreparedDestinationState, ResumeQueueDisposition,
450 SessionState,
451};
452
453pub use mj_core::state::{MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest};
454
455use crate::targets::{CommandExecutor, ProvisionStage, ProvisionStageGuard};
456use mj_core::relay::RelayCommand;
457
458pub async fn refresh_move_source(
460 manager: &SessionManagerControl,
461 id: &str,
462) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
463 Ok(MoveSourceRelay::lease(manager, id).await?.snapshot())
464}
465
466#[derive(Default)]
472pub(in crate::controller) struct MoveSourceRelay(
473 Option<super::checkpoint::ControllerRelayLease>,
474 Option<crate::worker_lifecycle::WorkerPermit>,
475);
476
477impl MoveSourceRelay {
478 pub(in crate::controller) fn owner(&self) -> Option<crate::worker_lifecycle::WorkerPermit> {
479 self.1.clone()
480 }
481
482 pub(in crate::controller) async fn lease(
487 manager: &SessionManagerControl,
488 id: &str,
489 ) -> Result<Self> {
490 let owner = crate::worker_lifecycle::WorkerPermit::acquire(
491 id,
492 "Move source relay",
493 &crate::targets::ProcessExecutor,
494 )
495 .await?;
496 let handle = manager
497 .wait_for_session(id, std::time::Duration::from_secs(5))
498 .await?;
499 match handle.lease_connection().await {
500 Ok(lease) => Ok(Self(
501 Some(super::checkpoint::ControllerRelayLease::Managed {
502 handle,
503 lease: Some(lease),
504 }),
505 Some(owner),
506 )),
507 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
508 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
509 Ok(Self(None, Some(owner)))
510 }
511 Err(error) => Err(error),
512 }
513 }
514
515 async fn set_subagent_admission_via_manager(
516 manager: &SessionManagerControl,
517 id: &str,
518 open: bool,
519 ) -> Result<bool> {
520 let handle = manager
521 .wait_for_session(id, std::time::Duration::from_secs(5))
522 .await?;
523 let mut lease = handle.lease_connection().await?;
524 if !(mj_core::relay::RelayRequest::SetSubagentAdmission { open })
525 .supported_at(lease.connection_mut().protocol_version())
526 {
527 lease.release();
528 return Ok(false);
529 }
530 let result = lease.connection_mut().set_subagent_admission(open).await;
531 lease.release();
532 result.map(|()| true)
533 }
534
535 pub(in crate::controller) fn snapshot(
537 &mut self,
538 ) -> Option<mj_core::state::ManagedSessionSnapshot> {
539 self.0
540 .as_mut()
541 .map(|relay| relay.connection_mut().snapshot())
542 }
543
544 pub(in crate::controller) async fn sync(
547 &mut self,
548 id: &str,
549 ) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
550 let Some(relay) = self.0.as_mut() else {
551 return Ok(None);
552 };
553 match relay.sync_snapshot().await {
554 Ok(snapshot) => Ok(Some(snapshot)),
555 Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
556 tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
557 self.0 = None;
558 Ok(None)
559 }
560 Err(error) => Err(error),
561 }
562 }
563
564 pub(in crate::controller) async fn drain_subagent_mutations(
568 &mut self,
569 id: &str,
570 executor: &(impl CommandExecutor + Sync),
571 ) -> Result<SubagentMutationDrain> {
572 ensure!(
573 self.0.is_some(),
574 "cannot safely drain sub-agent requests because the source worker is unavailable"
575 );
576 let protocol_version = self
577 .0
578 .as_mut()
579 .expect("held source relay checked")
580 .connection_mut()
581 .protocol_version();
582 if !(mj_core::relay::RelayRequest::SetSubagentAdmission { open: false })
583 .supported_at(protocol_version)
584 {
585 return Ok(SubagentMutationDrain::UnsupportedWorkerProtocol(
586 protocol_version,
587 ));
588 }
589 if let Err(error) = self
590 .0
591 .as_mut()
592 .expect("held source relay checked")
593 .connection_mut()
594 .set_subagent_admission(false)
595 .await
596 {
597 self.release_managed_connection();
601 if let Err(reopen) = self.reopen_subagent_mutations_with_retry(id).await {
602 return Err(error.context(format!(
603 "could not reopen sub-agent requests after an ambiguous admission close: {reopen:#}"
604 )));
605 }
606 return Err(error);
607 }
608 self.release_managed_connection();
609
610 let drain = async {
611 loop {
612 ensure!(
613 !executor.cancellation_requested(),
614 "Move cancelled while waiting for sub-agent requests"
615 );
616 let snapshot = self.sync(id).await?.context(
617 "source worker became unavailable while draining sub-agent requests",
618 )?;
619 let parent = id.to_owned();
620 let effects_pending = tokio::task::spawn_blocking(move || {
621 crate::database::has_pending_mutating_delegations(&parent)
622 })
623 .await??;
624 if !subagent_mutations_pending(&snapshot.subagent_requests, effects_pending) {
625 break;
626 }
627 executor.notify_notice("Waiting for subagent requests to finish");
628 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
629 }
630 Ok::<(), anyhow::Error>(())
631 }
632 .await;
633
634 if let Err(error) = drain {
635 if let Err(reopen) = self.reopen_subagent_mutations_with_retry(id).await {
636 return Err(error.context(format!(
637 "could not reopen sub-agent requests after Move stopped waiting: {reopen:#}"
638 )));
639 }
640 return Err(error);
641 }
642 self.reacquire_managed_connection().await?;
643 if executor.cancellation_requested() {
644 self.reopen_subagent_mutations_with_retry(id).await?;
645 bail!("Move cancelled before source interruption; sub-agent requests reopened");
646 }
647 Ok(SubagentMutationDrain::Drained)
648 }
649
650 pub(in crate::controller) async fn reopen_subagent_mutations(
651 &mut self,
652 id: &str,
653 ) -> Result<()> {
654 self.reopen_subagent_mutations_with_retry(id).await
655 }
656
657 async fn open_subagent_mutations(&mut self) -> Result<()> {
658 self.reacquire_managed_connection().await?;
659 let result = self
660 .0
661 .as_mut()
662 .context("source worker is unavailable")?
663 .connection_mut()
664 .set_subagent_admission(true)
665 .await;
666 self.release_managed_connection();
667 result
668 }
669
670 async fn reopen_subagent_mutations_with_retry(&mut self, id: &str) -> Result<()> {
671 let first_error = match self.open_subagent_mutations().await {
672 Ok(()) => return Ok(()),
673 Err(error) => error,
674 };
675 tracing::warn!(
676 session_id = id,
677 error = %first_error,
678 "sub-agent admission reopen failed; retrying on a fresh relay connection"
679 );
680 tokio::task::yield_now().await;
681 match self.open_subagent_mutations().await {
682 Ok(()) => {
683 tracing::info!(
684 session_id = id,
685 "reopened sub-agent requests on a replacement relay connection"
686 );
687 Ok(())
688 }
689 Err(retry_error) => {
690 tracing::warn!(
691 session_id = id,
692 error = %retry_error,
693 "could not reopen sub-agent requests on the replacement relay connection"
694 );
695 Err(first_error.context(format!("relay reconnect retry failed: {retry_error:#}")))
696 }
697 }
698 }
699
700 async fn reacquire_managed_connection(&mut self) -> Result<()> {
701 if let Some(super::checkpoint::ControllerRelayLease::Managed { handle, lease }) =
702 self.0.as_mut()
703 && lease.is_none()
704 {
705 *lease = Some(handle.lease_connection().await?);
706 }
707 Ok(())
708 }
709
710 fn release_managed_connection(&mut self) {
711 if let Some(super::checkpoint::ControllerRelayLease::Managed { lease, .. }) =
712 self.0.as_mut()
713 && let Some(lease) = lease.take()
714 {
715 lease.release();
716 }
717 }
718
719 pub(in crate::controller) fn is_held(&self) -> bool {
720 self.0.is_some()
721 }
722
723 pub(in crate::controller) fn take(
725 &mut self,
726 ) -> Option<super::checkpoint::ControllerRelayLease> {
727 self.0.take()
728 }
729
730 pub(in crate::controller) fn replace_connection(
733 &mut self,
734 connection: crate::session_manager::StandaloneSession,
735 ) {
736 if let Some(relay) = self.0.as_mut() {
737 relay.replace_connection(connection);
738 }
739 }
740}
741
742impl Drop for MoveSourceRelay {
743 fn drop(&mut self) {
744 if let Some(relay) = self.0.take() {
745 relay.release();
746 }
747 self.1.take();
748 }
749}
750
751fn digest(value: &impl serde::Serialize) -> Result<String> {
752 Ok(lower_hex(Sha256::digest(serde_json::to_vec(value)?)))
753}
754
755struct MovePhaseTimer<'a> {
756 session_id: &'a str,
757 phase: &'static str,
758 started: std::time::Instant,
759}
760
761impl<'a> MovePhaseTimer<'a> {
762 fn new(session_id: &'a str, phase: &'static str) -> Self {
763 Self {
764 session_id,
765 phase,
766 started: std::time::Instant::now(),
767 }
768 }
769}
770
771impl Drop for MovePhaseTimer<'_> {
772 fn drop(&mut self) {
773 tracing::info!(
774 session_id = self.session_id,
775 phase = self.phase,
776 elapsed_ms = self.started.elapsed().as_millis() as u64,
777 "move phase finished"
778 );
779 }
780}
781
782fn replace_queued_images_with_placeholders(
786 queued_commands: &mut [mj_core::state::MaterializedQueuedPrompt],
787) {
788 for block in queued_commands
789 .iter_mut()
790 .flat_map(|command| command.content.iter_mut())
791 {
792 if block.get("type").and_then(serde_json::Value::as_str) != Some("image") {
793 continue;
794 }
795 let mime = block
796 .get("mimeType")
797 .or_else(|| block.get("mime_type"))
798 .and_then(serde_json::Value::as_str)
799 .unwrap_or("image");
800 *block = serde_json::json!({"type": "text", "text": format!("[Image attachment: {mime}]")});
801 }
802}
803
804impl Controller {
805 fn validate_move_destination_paths(
809 &self,
810 source: &mj_core::state::SessionRecord,
811 target_id: &str,
812 executor: &(impl CommandExecutor + Sync),
813 ) -> Result<Option<super::worktree::RawToWorkspaceConversion>> {
814 use super::worktree::ResumePlan;
815 match super::worktree::resume_compatibility(source, &self.config, target_id)
816 .map_err(anyhow::Error::msg)?
817 {
818 ResumePlan::RawToWorkspace => {
819 return Ok(Some(super::worktree::plan_raw_to_workspace(
820 source,
821 &self.config,
822 executor,
823 )?));
824 }
825 ResumePlan::WorkspaceToRaw => {
826 self.plan_workspace_to_raw(source, target_id, executor)?;
827 }
828 ResumePlan::InPlace if source.managed_worktree.is_none() => {
829 if let Some(path) = &source.project_directory {
830 self.validate_project_directory(target_id, path, executor)?;
831 }
832 }
833 ResumePlan::InPlace => {}
834 }
835 Ok(None)
836 }
837 fn move_confirmation(
838 &self,
839 selection: &MoveSelection,
840 conversion: Option<&mj_core::state::RawConversionPreview>,
841 ) -> Result<(bool, Vec<mj_core::state::MaterializedQueuedPrompt>, String)> {
842 let source = self
843 .state
844 .sessions
845 .get(&selection.session_id)
846 .context("unknown move session")?;
847 let (mut active, mut queued) = crate::database::move_pending_work(&source.id)?;
848 if let Some(operation) = crate::database::load_move_operation(&source.id)?
849 && operation.queue_admission_started
850 && !operation.queue_admission_finished
851 {
852 ensure!(
853 operation.selection == *selection,
854 "queue admission is incomplete on the live destination; retry that move before selecting another destination"
855 );
856 let checkpoint = operation
857 .restore_artifact()
858 .context("retained queue checkpoint is missing")?;
859 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
860 ensure!(
861 verified.archive_sha256 == checkpoint.sha256
862 && verified.manifest.session.id == source.id,
863 "retained move checkpoint verification failed"
864 );
865 queued = verified
866 .canonical_session
867 .queued_prompts
868 .into_iter()
869 .map(|entry| mj_core::state::MaterializedQueuedPrompt {
870 accepted_ordinal: None,
871 command_id: entry.command_id,
872 kind: match entry.kind {
873 CanonicalQueuedCommandKind::Prompt => {
874 mj_core::state::QueuedCommandKind::Prompt
875 }
876 CanonicalQueuedCommandKind::SetConfig { key, value } => {
877 mj_core::state::QueuedCommandKind::SetConfig { key, value }
878 }
879 },
880 content: entry.content,
881 queued_at_ms: entry.queued_at_ms,
882 })
883 .collect();
884 active = false;
885 }
886 let fingerprint = digest(&(
887 &source.last_profile,
888 &source.target_template_id,
889 &source.target,
890 &source.target_runtime,
891 &source.native_session_id,
892 &source.resource_allocation,
893 &source.additional_mounts,
894 &source.container_cpus,
895 &source.container_memory,
896 selection,
897 self.move_configuration_fingerprint(selection)?,
898 &queued,
899 conversion.map(|preview| {
903 (
904 &preview.fetch_url,
905 &preview.push_urls,
906 &preview.branch,
907 &preview.destination,
908 )
909 }),
910 ))?;
911 Ok((active, queued, fingerprint))
912 }
913 fn verify_current_move_configuration(&self, operation: &MoveOperation) -> Result<()> {
914 let current = Controller {
915 config: mj_core::config::Config::load()?,
916 state: self.state.clone(),
917 };
918 ensure!(
919 current.move_configuration_fingerprint(&operation.selection)?
920 == operation.configuration_fingerprint,
921 "destination configuration changed during Move preparation; source retained, prepare and confirm Move again"
922 );
923 Ok(())
924 }
925
926 pub(super) fn move_configuration_fingerprint(
927 &self,
928 selection: &MoveSelection,
929 ) -> Result<String> {
930 let profile = self
931 .config
932 .profiles
933 .get(
934 selection
935 .profile_id
936 .as_deref()
937 .context("move profile is unresolved")?,
938 )
939 .context("move profile no longer exists")?;
940 let target = self
941 .config
942 .targets
943 .get(
944 selection
945 .target_template_id
946 .as_deref()
947 .context("move target is unresolved")?,
948 )
949 .context("move target no longer exists")?;
950 let source = self
953 .state
954 .sessions
955 .get(&selection.session_id)
956 .context("unknown move session")?;
957 digest(&(
958 profile,
959 target,
960 &self.config.bundles,
961 self.config.targets.get(&source.target_template_id),
962 self.config.profiles.get(&source.last_profile),
963 ))
964 }
965
966 async fn validate_move_destination_configuration(
967 &self,
968 selection: &MoveSelection,
969 source_harness: mj_core::config::HarnessKind,
970 operational: &mj_core::relay::RelayOperationalState,
971 ) -> Result<()> {
972 let profile_id = selection
973 .profile_id
974 .as_deref()
975 .context("move profile is unresolved")?;
976 let profile = self
977 .config
978 .profiles
979 .get(profile_id)
980 .context("move profile no longer exists")?;
981 let source = self
982 .state
983 .sessions
984 .get(&selection.session_id)
985 .context("unknown move session")?;
986 if profile.kind != source_harness || profile_id == source.last_profile {
987 return Ok(());
988 }
989 let accepted = mj_core::acp::AcceptedSessionConfig::from_configuration(
990 &operational.config,
991 &operational.config_options,
992 );
993 if accepted.model.is_none() && accepted.effort.is_none() {
994 return Ok(());
995 }
996 let choices =
999 super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
1000 .await
1001 .with_context(|| {
1002 format!("discover destination profile {profile_id:?} configuration")
1003 })?;
1004 validate_preserved_configuration(profile_id, &accepted, &choices)
1005 }
1006
1007 pub async fn prepare_move_session_controlled(
1008 &self,
1009 mut selection: MoveSelection,
1010 executor: &(impl CommandExecutor + Sync),
1011 ) -> Result<MovePreparation> {
1012 ensure!(
1013 selection.profile_id.is_some() || selection.target_template_id.is_some(),
1014 "move requires a target or profile selection"
1015 );
1016 let source = self
1017 .state
1018 .sessions
1019 .get(&selection.session_id)
1020 .context("unknown session")?;
1021 ensure!(
1022 !self.state.subagents.contains_key(&source.id),
1023 "sub-agent sessions cannot move independently of their parent"
1024 );
1025 let previous = crate::database::load_move_operation(&source.id)?;
1026 ensure!(
1027 previous
1028 .as_ref()
1029 .and_then(|op| op.prepared_destination.as_ref())
1030 .is_none_or(|d| !matches!(
1031 d.state,
1032 PreparedDestinationState::CleanupPending { .. }
1033 )),
1034 "EC2 destination cleanup is pending; source retained. Automatic cleanup will retry before another Move"
1035 );
1036
1037 if previous
1038 .as_ref()
1039 .is_some_and(|op| op.holds_source_environment() && !op.checkpoint_retained())
1040 {
1041 bail!(
1042 "this Move cannot be retried. {}",
1043 missing_move_checkpoint_recovery(Some(source))
1044 );
1045 }
1046 let retry = previous.as_ref().is_some_and(|op| {
1047 !matches!(op.phase, MovePhase::Completed)
1048 && (op.restore_artifact().is_some() || source.checkpoint.is_some())
1049 });
1050 ensure!(
1051 matches!(
1052 source.state,
1053 SessionState::Running | SessionState::Disconnected
1054 ) || retry,
1055 "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
1056 source.id
1057 );
1058 selection
1059 .profile_id
1060 .get_or_insert_with(|| source.last_profile.clone());
1061 selection
1062 .target_template_id
1063 .get_or_insert_with(|| source.target_template_id.clone());
1064 selection
1065 .additional_mounts
1066 .get_or_insert_with(|| source.additional_mounts.clone());
1067 ensure!(
1068 !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
1069 "select resource allocation or explicitly clear it, not both"
1070 );
1071 if selection.resource_allocation.is_none()
1072 && matches!(
1073 self.config
1074 .targets
1075 .get(selection.target_template_id.as_deref().unwrap()),
1076 Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1077 )
1078 && (matches!(
1079 source.resource_allocation,
1080 Some(mj_core::state::SessionResourceAllocation::Container { .. })
1081 ) || source.container_cpus.is_some()
1082 || source.container_memory.is_some())
1083 {
1084 selection.clear_resource_allocation = true;
1085 }
1086 if !selection.clear_resource_allocation {
1087 selection.resource_allocation = selection
1088 .resource_allocation
1089 .or_else(|| source.resource_allocation.clone());
1090 }
1091 if selection.clear_resource_allocation
1092 && source.resource_allocation.is_none()
1093 && source.container_cpus.is_none()
1094 && source.container_memory.is_none()
1095 {
1096 selection.clear_resource_allocation = false;
1097 }
1098 if let Some(retained) = previous.as_ref().filter(|op| op.holds_source_environment()) {
1099 if retained.in_place {
1103 selection.workspace = retained.selection.workspace.clone();
1104 }
1105 ensure!(
1106 retained.selection == selection,
1107 "{}",
1108 sealed_selection_difference(&retained.selection, &selection)
1109 );
1110 ensure!(
1111 !retained.in_place
1112 || (retained.source_target.is_some()
1113 && source.target == retained.source_target),
1114 "retained Move target is missing or changed; refusing to recreate it"
1115 );
1116 }
1117 let profile_id = selection.profile_id.as_deref().unwrap();
1118 let target_id = selection.target_template_id.as_deref().unwrap();
1119 let profile = self
1120 .config
1121 .profiles
1122 .get(profile_id)
1123 .context("unknown destination profile")?;
1124 ensure!(
1125 profile.enabled,
1126 "destination profile {profile_id:?} is disabled"
1127 );
1128 if let Some(policy) = &selection.subagents {
1129 super::profile_config::validate_session_subagent_policy(
1130 &self.config,
1131 profile_id,
1132 policy,
1133 )
1134 .await?;
1135 }
1136 let target = self
1137 .config
1138 .targets
1139 .get(target_id)
1140 .context("unknown destination target")?;
1141 self.validate_muse_resume_destination(source, profile.kind, target_id)?;
1142 super::worktree::resume_compatibility(source, &self.config, target_id)
1143 .map_err(anyhow::Error::msg)?;
1144 super::backend::validate_resource_allocation(
1145 target,
1146 selection.resource_allocation.as_ref(),
1147 )?;
1148 let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
1149 ensure!(
1150 profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
1151 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1152 );
1153 ensure!(
1154 mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
1155 "attached resources are unsupported for this target; select compatible resources explicitly"
1156 );
1157 crate::targets::validate_additional_mounts(mounts)?;
1158 for mount in mounts {
1159 self.validate_mount_source(target_id, &mount.source, executor)?;
1160 }
1161 let planned_conversion =
1162 self.validate_move_destination_paths(source, target_id, executor)?;
1163 ensure!(
1164 profile.home.is_dir(),
1165 "destination profile home is unavailable; configure the profile before moving"
1166 );
1167 super::worker_binary::preflight_worker_binary(target, executor)?;
1168 super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
1169 let source_harness = previous
1170 .as_ref()
1171 .filter(|operation| {
1172 !operation.queue_admission_started && operation.phase != MovePhase::Completed
1173 })
1174 .and_then(|operation| operation.recovery_session.as_ref())
1175 .map_or(source.harness_kind, |record| record.harness_kind);
1176 let cross_harness = profile.kind != source_harness;
1177 if cross_harness {
1178 let cancel = tokio_util::sync::CancellationToken::new();
1179 let resolve =
1180 crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
1181 tokio::pin!(resolve);
1182 loop {
1183 tokio::select! {
1184 result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
1185 _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
1186 if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
1187 }
1188 }
1189 }
1190 }
1191 let conversion = planned_conversion
1194 .map(|conversion| {
1195 super::worktree::raw_conversion_preview(source, &conversion, executor)
1196 .context("describe the move of this checkout into the target")
1197 })
1198 .transpose()?
1199 .map(Box::new);
1200 let (active, mut queued_commands, fingerprint) =
1201 self.move_confirmation(&selection, conversion.as_deref())?;
1202 replace_queued_images_with_placeholders(&mut queued_commands);
1203 let operation_id = previous
1204 .as_ref()
1205 .filter(|operation| {
1206 operation.selection == selection
1207 && operation.phase != MovePhase::Completed
1208 && operation.restore_artifact().is_some()
1209 })
1210 .map(|operation| operation.operation_id.clone())
1211 .unwrap_or(new_command_id("move")?);
1212 let in_place = previous
1213 .as_ref()
1214 .is_some_and(|op| retry && op.in_place && op.selection == selection)
1215 || in_place_move_eligible(
1216 source,
1217 &selection,
1218 &self.config.targets,
1219 self.state.subagents.contains_key(&source.id),
1220 retry,
1221 );
1222 let destination_has_parent_role =
1223 parent_tools_enabled(&move_subagent_policy(source, &selection), profile.kind);
1224 if in_place
1225 && !destination_has_parent_role
1226 && let Some(error) =
1227 roleless_move_children_error(&live_move_children(&self.state, &source.id))
1228 {
1229 bail!("{error}; close or finish those children before moving to this harness");
1230 }
1231 let workspace = if in_place {
1232 None
1233 } else {
1234 Some(self.assess_move_workspace(&selection, executor)?)
1235 };
1236 Ok(MovePreparation {
1237 destination_checks: if !in_place
1238 && matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. })
1239 {
1240 DestinationChecks::AfterProvisioning
1241 } else {
1242 DestinationChecks::Checked
1243 },
1244 workspace,
1245 source_unavailable: false,
1246 in_place,
1247 conversion,
1248 selection,
1249 source_profile_id: source.last_profile.clone(),
1250 source_target_template_id: source.target_template_id.clone(),
1251 cross_harness,
1252 active,
1253 queued_commands,
1254 fingerprint,
1255 operation_id,
1256 })
1257 }
1258
1259 pub async fn move_session_managed_controlled(
1261 &mut self,
1262 request: MoveSessionRequest,
1263 executor: &(impl CommandExecutor + Sync),
1264 manager: &SessionManagerControl,
1265 ) -> Result<MoveOutcome> {
1266 crate::worker_lifecycle::run(&request.preparation.selection.session_id.clone(), "move session managed controlled", executor, async {
1267 let prepared = &request.preparation;
1268 let id = prepared.selection.session_id.clone();
1269 let started = std::time::Instant::now();
1270 executor.notify_notice("Checking destination");
1271 let mut source_relay = MoveSourceRelay::default();
1272 let checked = {
1273 let _checking_destination =
1274 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1275 let mut checked = {
1276 let _timing = MovePhaseTimer::new(&id, "preflight destination checks");
1277 self.prepare_move_session_controlled(prepared.selection.clone(), executor)
1278 .await?
1279 };
1280 if matches!(
1285 self.state.sessions[&id].state,
1286 SessionState::Running | SessionState::Disconnected
1287 ) {
1288 let source_harness = self.state.sessions[&id].harness_kind;
1289 source_relay = {
1290 let _timing = MovePhaseTimer::new(&id, "preflight source lease");
1291 MoveSourceRelay::lease(manager, &id).await?
1292 };
1293 let snapshot = source_relay.snapshot();
1294 let (active, queue, fingerprint) =
1295 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
1296 checked.source_unavailable = snapshot
1297 .as_ref()
1298 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
1299 checked.active = active
1300 || checked.source_unavailable
1301 || snapshot.as_ref().is_some_and(|snapshot| {
1302 let mut operational = snapshot.operational.clone();
1303 operational.queued_prompts.clear();
1304 operational.checkpoint_barrier = None;
1305 !operational.safe_to_replace(source_harness)
1306 });
1307 if let Some(snapshot) = &snapshot {
1308 let _timing = MovePhaseTimer::new(&id, "preflight destination configuration");
1309 self.validate_move_destination_configuration(
1310 &checked.selection,
1311 source_harness,
1312 &snapshot.operational,
1313 )
1314 .await?;
1315 }
1316 checked.queued_commands = queue;
1317 checked.fingerprint = fingerprint;
1318 }
1319 checked
1320 };
1321 ensure!(
1322 checked.fingerprint == prepared.fingerprint,
1323 "session, pending work, or destination configuration changed; prepare and confirm Move again"
1324 );
1325 let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
1326 let source = self.state.sessions[&id].clone();
1327 if let Some(assessment) = &checked.workspace {
1328 checked.selection.workspace.validate(assessment)?;
1329 }
1330 let old_operation = crate::database::load_move_operation(&id)?;
1331 let retry = old_operation.filter(|op| {
1332 op.selection == prepared.selection
1333 && op.phase != MovePhase::Completed
1334 && op.restore_artifact().is_some()
1335 });
1336 if retry.is_none()
1337 && source.state == SessionState::Running
1338 && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
1339 && move_environment_change(&source, &checked.selection, false).is_none()
1340 {
1341 return Ok(outcome(
1342 &prepared.operation_id,
1343 &prepared.selection,
1344 "unchanged",
1345 None,
1346 None,
1347 ));
1348 }
1349 ensure!(
1350 !checked.active || request.acknowledge_interruption,
1351 "active work will be interrupted; confirm Move again with interruption acknowledgement"
1352 );
1353 ensure!(
1354 checked.queued_commands.is_empty() || request.queue.is_some(),
1355 "pending work requires an explicit queue choice: discard or start"
1356 );
1357 let timestamp = now();
1358 let mut operation = match retry {
1359 Some(mut op) => {
1360 ensure!(
1361 !op.queue_admission_started || op.queue == queue,
1362 "queued work may already have run; retry with the original queue choice on the same destination"
1363 );
1364 crate::database::clear_move_cancellation_for_retry(&id)?;
1365 op.configuration_fingerprint =
1366 self.move_configuration_fingerprint(&checked.selection)?;
1367 op.queue = queue;
1368 op.cancellation_requested = false;
1369 op.error = None;
1370 op
1371 }
1372 None => MoveOperation {
1373 prepared_destination: None,
1374 accepted_preparation: None,
1375 acknowledge_interruption: false,
1376 workspace_transfer: if checked.in_place {
1377 None
1378 } else {
1379 Some(
1380 self.new_workspace_transfer(
1381 &id,
1382 &prepared.operation_id,
1383 checked
1384 .workspace
1385 .clone()
1386 .context("Move workspace assessment missing")?,
1387 executor,
1388 )?,
1389 )
1390 },
1391 handoff: None,
1392 source_checkpoint_only: false,
1393 in_place: in_place_move_eligible(
1394 &source,
1395 &checked.selection,
1396 &self.config.targets,
1397 self.state.subagents.contains_key(&id),
1398 false,
1399 ),
1400 operation_id: prepared.operation_id.clone(),
1401 selection: checked.selection.clone(),
1402 source_profile_id: source.last_profile.clone(),
1403 source_target_template_id: source.target_template_id.clone(),
1404 source_target: source.target.clone(),
1405 source_native_session_id: source.native_session_id.clone(),
1406 source_additional_mounts: source.additional_mounts.clone(),
1407 source_resource_allocation: source.resource_allocation.clone(),
1408 destination_target: None,
1409 destination_native_session_id: None,
1410 destination_store_id: None,
1411 configuration_fingerprint: self
1412 .move_configuration_fingerprint(&checked.selection)?,
1413 checkpoint: None,
1414 recovery_session: None,
1415 queue,
1416 phase: MovePhase::Preparing,
1417 queue_admission_started: false,
1418 queue_admission_finished: false,
1419 cancellation_requested: false,
1420 created_at: timestamp.clone(),
1421 updated_at: timestamp,
1422 error: None,
1423 },
1424 };
1425 operation.accepted_preparation = Some(Box::new(checked.clone()));
1426 operation.acknowledge_interruption = request.acknowledge_interruption;
1427 crate::database::save_move_operation(&operation)?;
1428 tracing::info!(
1429 session_id = id,
1430 in_place = operation.in_place,
1431 reason = move_environment_change(
1432 &source,
1433 &checked.selection,
1434 bare_targets_share_environment(
1435 &self.config.targets,
1436 &source.target_template_id,
1437 checked
1438 .selection
1439 .target_template_id
1440 .as_deref()
1441 .unwrap_or_default(),
1442 ),
1443 )
1444 .unwrap_or(if operation.in_place {
1445 "environment unchanged"
1446 } else {
1447 "source environment unavailable or previously released"
1448 }),
1449 "move environment decision"
1450 );
1451 tracing::info!(
1452 session_id = id,
1453 phase = "preflight",
1454 elapsed_ms = started.elapsed().as_millis() as u64,
1455 "move phase completed"
1456 );
1457 let result = Box::pin(self.execute_move(
1460 &mut operation,
1461 Some(&checked),
1462 executor,
1463 manager,
1464 source_relay,
1465 ))
1466 .await;
1467 self.finish_move_result_with_subagent_recovery(
1468 &mut operation,
1469 result,
1470 executor,
1471 manager,
1472 )
1473 .await
1474
1475 }).await
1476 }
1477
1478 fn finish_move_result(
1479 &mut self,
1480 operation: &mut MoveOperation,
1481 result: Result<()>,
1482 executor: &impl CommandExecutor,
1483 ) -> Result<MoveOutcome> {
1484 if result.is_err()
1485 && operation.workspace_transfer.is_some()
1486 && !crate::upgrade::gate().is_open()
1487 {
1488 return Ok(outcome(
1491 &operation.operation_id,
1492 &operation.selection,
1493 "interrupted",
1494 None,
1495 Some("Move will continue after the daemon upgrade".into()),
1496 ));
1497 }
1498 if let Some(saved) = crate::database::load_move_operation(&operation.selection.session_id)?
1499 && saved.operation_id == operation.operation_id
1500 {
1501 operation.prepared_destination = saved.prepared_destination;
1502 }
1503 let result = match result {
1504 Err(error)
1505 if !operation.queue_admission_started
1506 && operation
1507 .prepared_destination
1508 .as_ref()
1509 .is_some_and(|d| d.owns_resource()) =>
1510 {
1511 executor.begin_resumable_move_work()?;
1512 let cleaned = self.cleanup_prepared_move_destination(operation, executor);
1513 executor.end_resumable_move_work()?;
1514 Err(match cleaned {
1515 Ok(()) => error,
1516 Err(cleanup_error) => error.context(format!("EC2 destination cleanup is pending and will retry automatically: {cleanup_error:#}")),
1517 })
1518 }
1519 result => result,
1520 };
1521 let session_id = operation.selection.session_id.clone();
1522 let mut last_error = self
1523 .state
1524 .sessions
1525 .get(&session_id)
1526 .context("Move outcome session is missing")?
1527 .last_error
1528 .clone();
1529 let (status, error, recovery) = match result {
1530 Ok(()) => {
1531 operation.phase = MovePhase::Completed;
1532 operation.error = None;
1533 if last_error
1534 .as_deref()
1535 .is_some_and(|error| error.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
1536 {
1537 last_error = None;
1538 }
1539 ("completed", None, None)
1540 }
1541 Err(error) => {
1542 let phase = match operation.phase {
1543 MovePhase::Preparing => "preparing the destination",
1544 MovePhase::ClosingSource => "checkpointing and suspending the source",
1545 MovePhase::ResumingDestination => "resuming the destination",
1546 MovePhase::StartingQueue => "starting the destination queue",
1547 MovePhase::Completed | MovePhase::Failed | MovePhase::Cancelled => {
1548 "recovering the move"
1549 }
1550 };
1551 let cancelled =
1552 executor.cancellation_requested() || operation.cancellation_requested;
1553 operation.phase = if cancelled {
1554 MovePhase::Cancelled
1555 } else {
1556 MovePhase::Failed
1557 };
1558 operation.cancellation_requested = cancelled;
1559 let recovery = failed_move_recovery(
1560 operation,
1561 self.state.sessions.get(&operation.selection.session_id),
1562 );
1563 let error = format!("{error:#}");
1564 tracing::warn!(
1565 %session_id,
1566 reference = %operation.operation_id,
1567 phase,
1568 cancelled,
1569 %error,
1570 "session move did not finish"
1571 );
1572 last_error = Some(failed_move_message(
1573 Some(phase),
1574 cancelled,
1575 &recovery,
1576 &operation.operation_id,
1577 ));
1578 operation.error = Some(error.clone());
1579 (
1580 if cancelled { "cancelled" } else { "failed" },
1581 Some(error),
1582 Some(recovery),
1583 )
1584 }
1585 };
1586 operation.updated_at = now();
1587 crate::database::save_move_outcome(operation, last_error.as_deref())?;
1588 let record = self
1589 .state
1590 .sessions
1591 .get_mut(&session_id)
1592 .expect("Move outcome session");
1593 record.last_error = last_error;
1594 record.updated_at = operation.updated_at.clone();
1595 Ok(outcome(
1596 &operation.operation_id,
1597 &operation.selection,
1598 status,
1599 error,
1600 recovery,
1601 ))
1602 }
1603
1604 async fn finish_move_result_with_subagent_recovery(
1605 &mut self,
1606 operation: &mut MoveOperation,
1607 result: Result<()>,
1608 executor: &(impl CommandExecutor + Sync),
1609 manager: &SessionManagerControl,
1610 ) -> Result<MoveOutcome> {
1611 let outcome = self.finish_move_result(operation, result, executor)?;
1612 if !operation.in_place
1613 || !matches!(operation.phase, MovePhase::Failed | MovePhase::Cancelled)
1614 {
1615 return Ok(outcome);
1616 }
1617 let session_id = &operation.selection.session_id;
1618 let Some(session) = self.state.sessions.get(session_id) else {
1619 return Ok(outcome);
1620 };
1621 if !matches!(
1622 session.state,
1623 SessionState::Running | SessionState::Disconnected
1624 ) || !parent_tools_enabled(
1625 &session.subagents.clone().unwrap_or_default(),
1626 session.harness_kind,
1627 ) {
1628 return Ok(outcome);
1629 }
1630
1631 for attempt in 1..=2 {
1635 match MoveSourceRelay::set_subagent_admission_via_manager(manager, session_id, true)
1636 .await
1637 {
1638 Ok(true) => {
1639 if attempt > 1 {
1640 tracing::info!(
1641 session_id,
1642 "reopened sub-agent requests on a replacement relay connection"
1643 );
1644 }
1645 break;
1646 }
1647 Ok(false) => break,
1648 Err(error) => {
1649 tracing::warn!(
1650 session_id,
1651 attempt,
1652 error = %error,
1653 "could not reopen sub-agent requests after the in-place Move left its source running"
1654 );
1655 if attempt == 1 {
1656 tokio::task::yield_now().await;
1657 }
1658 }
1659 }
1660 }
1661 Ok(outcome)
1662 }
1663
1664 pub async fn recover_move_managed_controlled(
1665 &mut self,
1666 mut operation: MoveOperation,
1667 executor: &(impl CommandExecutor + Sync),
1668 manager: &SessionManagerControl,
1669 ) -> Result<MoveOutcome> {
1670 crate::worker_lifecycle::run(&operation.selection.session_id.clone(), "recover move managed controlled", executor, async {
1671 let id = operation.selection.session_id.clone();
1672 let result = async {
1673 let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1674 if operation.phase == MovePhase::Preparing
1675 && operation.accepted_preparation.is_some()
1676 && matches!(self.config.targets.get(operation.selection.target_template_id.as_deref().unwrap_or_default()),
1677 Some(mj_core::config::TargetTemplate::AwsEc2 { .. }))
1678 && !operation.cancellation_requested {
1679 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1680 }
1681 if operation.queue_admission_started {
1682 ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1683 self.finish_workspace_transfer(&mut operation, executor)?;
1684 return self.admit_move_queue(&mut operation, executor).await;
1685 }
1686 if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1687 && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1688 if session.state == SessionState::Provisioning || session.target != operation.source_target {
1689 self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1690 }
1691 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1692 }
1693 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1694 let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1695 let sealed = operation.in_place && operation.recovery_session.is_some() && session.state == SessionState::Closing;
1701 if !sealed {
1702 Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1703 }
1704 if !operation.retains_source_environment() {
1708 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1709 }
1710 if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1711 let previous = self.state.sessions[&id].clone();
1712 operation.recovery_session = Some(previous.clone());
1713 let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1714 return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1715 }
1716 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1717 self.cleanup_stopped_target(&id, &cleanup)?;
1718 }
1719 bail!("{}", interrupted_source_stop_message(
1720 "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1721 operation.in_place,
1722 ));
1723 }
1724 match operation.phase {
1725 MovePhase::ClosingSource
1726 if operation.in_place
1727 && matches!(
1728 session.state,
1729 SessionState::Running | SessionState::Disconnected
1730 )
1731 && !operation.cancellation_requested =>
1732 {
1733 let relay = MoveSourceRelay::lease(manager, &id).await?;
1737 return Box::pin(self.execute_move(
1738 &mut operation,
1739 None,
1740 executor,
1741 manager,
1742 relay,
1743 ))
1744 .await;
1745 }
1746 MovePhase::ClosingSource
1747 if operation.in_place
1748 && matches!(
1749 session.state,
1750 SessionState::Running | SessionState::Disconnected
1751 )
1752 && operation.cancellation_requested =>
1753 {
1754 let source = self.state.sessions.get(&id).context("move source is missing")?;
1755 if parent_tools_enabled(
1756 &source.subagents.clone().unwrap_or_default(),
1757 source.harness_kind,
1758 ) {
1759 MoveSourceRelay::set_subagent_admission_via_manager(
1760 manager, &id, true,
1761 )
1762 .await?;
1763 }
1764 bail!("Move was cancelled before source interruption; source retained and sub-agent requests reopened")
1765 }
1766 MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1767 MovePhase::ResumingDestination if session.state == SessionState::Running => {
1768 ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1772 && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1773 "ready destination does not match the move intent");
1774 operation.destination_target = session.target.clone();
1775 operation.destination_native_session_id = session.native_session_id.clone();
1776 operation.queue_admission_started = true;
1777 operation.phase = MovePhase::StartingQueue;
1778 crate::database::save_move_operation(&operation)?;
1779 restore_move_queue_hold(&operation);
1780 ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1781 self.finish_workspace_transfer(&mut operation, executor)?;
1782 self.admit_move_queue(&mut operation, executor).await
1783 }
1784 MovePhase::ResumingDestination => {
1785 let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1786 let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1789 let error = if operation.in_place {
1790 self.retain_failed_in_place_move(&id, previous, cause)?
1791 } else if operation.workspace_transfer.is_some() {
1792 self.rollback_move_destination(&operation, cause, executor)?
1793 } else {
1794 self.rollback_failed_resume(&id, previous, false, cause, executor)?
1795 };
1796 Err(error)
1797 }
1798 MovePhase::ClosingSource => {
1799 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1800 Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1801 }
1802 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1803 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1804 self.cleanup_stopped_target(&id, executor)?;
1805 }
1806 bail!("{}", interrupted_source_stop_message(
1807 "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1808 operation.in_place,
1809 ))
1810 }
1811 _ => bail!("Move requires an explicit retry after the daemon restarted"),
1812 }
1813 }.await;
1814 self.finish_move_result_with_subagent_recovery(
1815 &mut operation,
1816 result,
1817 executor,
1818 manager,
1819 )
1820 .await
1821
1822 }).await
1823 }
1824
1825 async fn recover_move_source_stop(
1826 &mut self,
1827 operation: &mut MoveOperation,
1828 executor: &(impl CommandExecutor + Sync),
1829 manager: &SessionManagerControl,
1830 ) -> Result<()> {
1831 let id = operation.selection.session_id.clone();
1832 if self.state.sessions[&id].state == SessionState::Closing {
1833 self.prepare_move_source_checkpoint(
1834 &id,
1835 executor,
1836 manager,
1837 operation,
1838 &mut MoveSourceRelay::default(),
1839 )
1840 .await?;
1841 let handle = manager
1842 .wait_for_session(&id, std::time::Duration::from_secs(5))
1843 .await?;
1844 let mut lease = handle.lease_connection().await?;
1845 let execution = lease.connection_mut().sync().await?.operational.execution;
1846 if operation.retains_source_environment()
1847 && matches!(
1848 execution,
1849 mj_core::relay::RelayExecutionState::Closing
1850 | mj_core::relay::RelayExecutionState::Closed
1851 )
1852 {
1853 super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1854 lease.release();
1855 return Ok(());
1860 }
1861 lease.release();
1862 if matches!(
1863 execution,
1864 mj_core::relay::RelayExecutionState::Idle
1865 | mj_core::relay::RelayExecutionState::Running
1866 ) {
1867 if operation.cancellation_requested || operation.retains_source_environment() {
1868 let record = self.state.sessions.get_mut(&id).unwrap();
1869 record.state = SessionState::Running;
1870 record.updated_at = now();
1871 record.last_error = Some(
1872 "Move was interrupted before the source was sealed; source retained".into(),
1873 );
1874 crate::database::save_lifecycle_session(record)?;
1875 } else {
1876 Box::pin(self.suspend_session_for_move(
1878 &id,
1879 executor,
1880 manager,
1881 operation,
1882 None,
1883 SourceTargetDisposition::Destroy,
1884 MoveSourceRelay::default(),
1885 ))
1886 .await?;
1887 }
1888 return Ok(());
1889 }
1890 }
1891 ensure!(
1892 !operation.retains_source_environment(),
1893 "retained Move source cannot be proven; refusing target teardown"
1894 );
1895 self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1897 .await?;
1898 Ok(())
1899 }
1900
1901 async fn execute_move(
1904 &mut self,
1905 operation: &mut MoveOperation,
1906 preparation: Option<&MovePreparation>,
1907 executor: &(impl CommandExecutor + Sync),
1908 manager: &SessionManagerControl,
1909 mut source_relay: MoveSourceRelay,
1910 ) -> Result<()> {
1911 let retained = source_relay.owner();
1912 let admission_id = operation.selection.session_id.clone();
1913 crate::worker_lifecycle::run_with_owner(&admission_id, "execute move", executor, retained, async {
1914 let id = operation.selection.session_id.clone();
1915 let mut preparation = preparation.cloned();
1916 if let Some(saved) = crate::database::load_move_operation(&id)?
1917 && saved.operation_id == operation.operation_id
1918 {
1919 operation.prepared_destination = saved.prepared_destination;
1920 }
1921
1922 if !operation.in_place
1923 && !operation.queue_admission_started
1924 && matches!(
1925 self.config.targets.get(
1926 operation
1927 .selection
1928 .target_template_id
1929 .as_deref()
1930 .unwrap_or_default()
1931 ),
1932 Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1933 )
1934 {
1935 drop(source_relay);
1936 source_relay = MoveSourceRelay::default();
1937 self.verify_current_move_configuration(operation)?;
1938 if self.state.sessions[&id].state == SessionState::Error
1942 && operation.recovery_session.is_some()
1943 {
1944 ensure!(
1945 !forget_missing_move_archives(operation),
1946 "the Move's checkpoint archive is missing; nothing is left to restore"
1947 );
1948 self.rollback_move_destination(
1949 operation,
1950 anyhow::anyhow!("clean up the partial Move destination before retry"),
1951 executor,
1952 )?;
1953 let record = self.state.sessions.get_mut(&id).unwrap();
1954 record.state = SessionState::Closing;
1955 crate::database::save_resumed_session(record, None)?;
1956 operation.prepared_destination = crate::database::load_move_operation(&id)?
1957 .context("Move cleanup intent missing")?
1958 .prepared_destination;
1959 }
1960 self.prepare_ec2_move_destination(operation, executor)?;
1961 self.verify_current_move_configuration(operation)?;
1962 if matches!(
1963 self.state.sessions[&id].state,
1964 SessionState::Running | SessionState::Disconnected
1965 ) && operation.recovery_session.is_none()
1966 {
1967 let accepted = operation
1968 .accepted_preparation
1969 .as_ref()
1970 .context("EC2 Move confirmation missing")?;
1971 let mut checked = self
1972 .prepare_move_session_controlled(operation.selection.clone(), executor)
1973 .await?;
1974 let harness = self.state.sessions[&id].harness_kind;
1975 source_relay = MoveSourceRelay::lease(manager, &id).await?;
1976 let snapshot = source_relay.snapshot();
1977 let (active, queue, fingerprint) =
1978 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
1979 let unavailable = snapshot
1980 .as_ref()
1981 .is_none_or(|s| !s.operational.native_session_is_ready());
1982 let active = active
1983 || unavailable
1984 || snapshot.as_ref().is_some_and(|s| {
1985 let mut state = s.operational.clone();
1986 state.queued_prompts.clear();
1987 state.checkpoint_barrier = None;
1988 !state.safe_to_replace(harness)
1989 });
1990 ensure!(
1991 fingerprint == accepted.fingerprint,
1992 "session, pending work, or destination configuration changed; source retained, prepare and confirm Move again"
1993 );
1994 ensure!(
1995 !active || operation.acknowledge_interruption,
1996 "active work will be interrupted; source retained, confirm Move again with interruption acknowledgement"
1997 );
1998 if let Some(snapshot) = snapshot {
1999 self.validate_move_destination_configuration(
2000 &checked.selection,
2001 harness,
2002 &snapshot.operational,
2003 )
2004 .await?;
2005 }
2006 checked.active = active;
2007 checked.queued_commands = queue;
2008 checked.source_unavailable = unavailable;
2009 let assessment = checked
2010 .workspace
2011 .as_mut()
2012 .context("EC2 Move workspace missing")?;
2013 self.assess_prepared_destination(operation, assessment, executor)?;
2014 operation
2015 .workspace_transfer
2016 .as_mut()
2017 .context("EC2 Move transfer missing")?
2018 .assessment = assessment.clone();
2019 crate::database::save_move_operation(operation)?;
2020 preparation = Some(checked);
2021 } else {
2022 let mut assessment = operation
2023 .workspace_transfer
2024 .as_ref()
2025 .context("EC2 Move transfer missing")?
2026 .assessment
2027 .clone();
2028 assessment
2029 .storage
2030 .retain(|s| !s.allocations.contains("destination"));
2031 self.assess_prepared_destination(operation, &mut assessment, executor)?;
2032 }
2033 }
2034 ensure!(
2035 !executor.cancellation_requested(),
2036 "move cancelled before source interruption"
2037 );
2038 if !operation.queue_admission_started {
2039 if forget_missing_move_archives(operation) {
2044 operation.updated_at = now();
2045 crate::database::save_move_operation(operation)?;
2046 bail!("the Move's checkpoint archive is missing; nothing is left to restore");
2047 }
2048 if operation.in_place {
2049 ensure!(
2050 operation.source_target.is_some()
2051 && self.state.sessions[&id].target == operation.source_target,
2052 "retained Move target is missing or changed; refusing to recreate the environment"
2053 );
2054 }
2055 if self.state.sessions[&id].state == SessionState::Error
2056 && let Some(previous) = operation.recovery_session.as_ref()
2057 {
2058 let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
2061 if operation.in_place {
2062 self.retain_failed_in_place_move(&id, previous, cause)?;
2068 } else if operation.workspace_transfer.is_some() {
2069 self.rollback_move_destination(operation, cause, executor)?;
2070 let record = self.state.sessions.get_mut(&id).unwrap();
2071 record.state = SessionState::Closing;
2072 crate::database::save_resumed_session(record, None)?;
2073 } else {
2074 let failure =
2075 self.rollback_failed_resume(&id, previous, false, cause, executor)?;
2076 ensure!(
2077 self.state.sessions[&id].state == SessionState::Stopped,
2078 "{failure:#}"
2079 );
2080 }
2081 }
2082 let state = self.state.sessions[&id].state;
2083 if matches!(state, SessionState::Closing | SessionState::Destroying)
2084 && !(operation.retains_source_environment() && operation.recovery_session.is_some())
2085 {
2086 Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
2087 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
2088 }
2089 if matches!(
2090 self.state.sessions[&id].state,
2091 SessionState::Running | SessionState::Disconnected
2092 ) && operation.destination_target.is_none()
2093 {
2094 let source = self.state.sessions[&id].clone();
2095 let source_has_parent_role = parent_tools_enabled(
2096 &source.subagents.clone().unwrap_or_default(),
2097 source.harness_kind,
2098 );
2099 let destination_has_parent_role = parent_tools_enabled(
2100 &move_subagent_policy(&source, &operation.selection),
2101 operation
2102 .selection
2103 .profile_id
2104 .as_deref()
2105 .and_then(|profile| self.config.profiles.get(profile))
2106 .context("destination profile is missing")?
2107 .kind,
2108 );
2109 if operation.in_place {
2110 operation.phase = MovePhase::ClosingSource;
2114 operation.updated_at = now();
2115 crate::database::save_move_operation(operation)?;
2116 let mut stopped_subagents_for_legacy_worker = false;
2117 if source_has_parent_role {
2118 if !source_relay.is_held() {
2119 source_relay = MoveSourceRelay::lease(manager, &id).await?;
2120 }
2121 match source_relay.drain_subagent_mutations(&id, executor).await? {
2122 SubagentMutationDrain::Drained => {}
2123 SubagentMutationDrain::UnsupportedWorkerProtocol(version) => {
2124 executor.notify_notice(&format!(
2125 "The source worker uses relay protocol {version}; protocol 32 is required to safely drain sub-agent requests, so stopping sub-agents before this in-place Move"
2126 ));
2127 executor.before_move_source_stop().await?;
2128 stopped_subagents_for_legacy_worker = true;
2129 }
2130 }
2131 }
2132 if should_stop_move_subagents(true, destination_has_parent_role)
2133 && !stopped_subagents_for_legacy_worker
2134 {
2135 let current = crate::database::load_state()?;
2136 if let Some(error) =
2137 roleless_move_children_error(&live_move_children(¤t, &id))
2138 {
2139 if source_has_parent_role {
2140 source_relay
2141 .reopen_subagent_mutations(&id)
2142 .await
2143 .context("could not reopen sub-agent requests after refusing Move")?;
2144 }
2145 bail!("{error}; close or finish those children before moving to this harness");
2146 }
2147 if let Err(error) = executor.before_move_source_stop().await {
2148 if source_has_parent_role {
2149 source_relay
2150 .reopen_subagent_mutations(&id)
2151 .await
2152 .context("could not reopen sub-agent requests after stopping children failed")?;
2153 }
2154 return Err(error);
2155 }
2156 }
2157 } else {
2158 executor.before_move_source_stop().await?;
2161 }
2162 if executor.cancellation_requested() {
2163 if operation.in_place && source_has_parent_role {
2164 source_relay.reopen_subagent_mutations(&id).await?;
2165 }
2166 bail!("Move cancelled before source interruption");
2167 }
2168 executor.notify_notice("Stopping source");
2169 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
2170 if !operation.in_place {
2171 operation.phase = MovePhase::ClosingSource;
2172 operation.updated_at = now();
2173 crate::database::save_move_operation(operation)?;
2174 }
2175 let disposition = if operation.in_place {
2179 SourceTargetDisposition::RetainForInPlaceSwap
2180 } else {
2181 SourceTargetDisposition::Destroy
2182 };
2183 let suspended = Box::pin(self.suspend_session_for_move(
2184 &id,
2185 executor,
2186 manager,
2187 operation,
2188 preparation.as_ref(),
2189 disposition,
2190 std::mem::take(&mut source_relay),
2191 ))
2192 .await;
2193 if let Err(error) = suspended {
2194 if operation.in_place
2195 && source_has_parent_role
2196 && matches!(
2197 self.state.sessions[&id].state,
2198 SessionState::Running | SessionState::Disconnected
2199 )
2200 && let Err(reopen) = MoveSourceRelay::set_subagent_admission_via_manager(
2201 manager,
2202 &id,
2203 true,
2204 )
2205 .await
2206 {
2207 return Err(error.context(format!(
2208 "could not reopen sub-agent requests after source checkpoint failed: {reopen:#}"
2209 )));
2210 }
2211 return Err(error);
2212 }
2213 }
2214 drop(source_relay);
2216 if !operation.in_place
2219 && operation.workspace_transfer.is_none()
2220 && self.state.sessions[&id].state == SessionState::Stopped
2221 && self.state.sessions[&id].target.is_some()
2222 {
2223 executor.notify_notice("Cleaning up source");
2224 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
2225 self.cleanup_stopped_target(&id, executor)?;
2226 }
2227 operation.checkpoint = operation
2228 .checkpoint
2229 .clone()
2230 .or_else(|| self.state.sessions[&id].checkpoint.clone());
2231 ensure!(
2232 operation.restore_artifact().is_some(),
2233 "move has no verified checkpoint"
2234 );
2235 ensure!(
2236 !executor.cancellation_requested(),
2237 "move cancelled after source sealing; checkpoint and remaining environment retained"
2238 );
2239 operation.phase = MovePhase::ResumingDestination;
2240 operation.recovery_session = Some(self.state.sessions[&id].clone());
2241 operation.updated_at = now();
2242 crate::database::save_move_operation(operation)?;
2243 executor.reserve_move_destination();
2244 executor.notify_notice("Preparing destination");
2245 if let Some(policy) = &operation.selection.subagents {
2249 let session = self.state.sessions.get_mut(&id).unwrap();
2250 session.subagents = Some(policy.clone());
2251 crate::database::save_resumed_session(session, None)?;
2252 }
2253 if operation.in_place {
2254 Box::pin(self.restore_session_in_place(
2258 &id,
2259 operation.selection.profile_id.as_deref().unwrap(),
2260 operation.selection.target_template_id.as_deref().unwrap(),
2261 executor,
2262 ))
2263 .await?;
2264 } else {
2265 if operation.selection.clear_resource_allocation {
2266 let session = self.state.sessions.get_mut(&id).unwrap();
2267 session.resource_allocation = None;
2268 session.container_cpus = None;
2269 session.container_memory = None;
2270 crate::database::save_resumed_session(session, None)?;
2271 }
2272 if operation.workspace_transfer.is_some() {
2273 self.capture_move_workspace(operation, executor)?;
2274 Box::pin(self.resume_session_for_move(operation, executor)).await?;
2275 let saved = crate::database::load_move_operation(&id)?
2276 .context("Move intent disappeared during restore")?;
2277 operation.workspace_transfer = saved.workspace_transfer;
2278 operation.prepared_destination = saved.prepared_destination;
2279 } else {
2280 Box::pin(self.resume_session_controlled(
2281 &id,
2282 operation.selection.profile_id.as_deref().unwrap(),
2283 operation.selection.target_template_id.as_deref().unwrap(),
2284 SessionResumeOptions {
2285 additional_mounts: operation.selection.additional_mounts.clone(),
2286 resource_allocation: operation.selection.resource_allocation.clone(),
2287 discard_queue: true,
2288 },
2289 executor,
2290 ))
2291 .await?;
2292 }
2293 }
2294 let destination = &self.state.sessions[&id];
2295 operation.destination_target = destination.target.clone();
2296 operation.destination_native_session_id = destination.native_session_id.clone();
2297 operation.phase = MovePhase::StartingQueue;
2299 operation.queue_admission_started = true;
2300 operation.updated_at = now();
2301 crate::database::save_move_operation(operation)?;
2302 }
2303 restore_move_queue_hold(operation);
2304 self.finish_workspace_transfer(operation, executor)?;
2305 self.admit_move_queue(operation, executor).await
2306
2307 }).await
2308 }
2309
2310 async fn admit_move_queue(
2311 &self,
2312 operation: &mut MoveOperation,
2313 executor: &(impl CommandExecutor + Sync),
2314 ) -> Result<()> {
2315 let timing_id = operation.selection.session_id.clone();
2316 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
2317 let id = &operation.selection.session_id;
2318 let mut relay = {
2319 let _checking_destination =
2320 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2321 let destination = &self.state.sessions[id];
2322 ensure!(
2323 destination.state == SessionState::Running
2324 && destination.target == operation.destination_target
2325 && destination.native_session_id == operation.destination_native_session_id,
2326 "cannot prove the same ready destination; refusing to replay potentially executed work"
2327 );
2328 let spec = self.reconnect_command(id)?;
2329 let relay = StandaloneSession::connect_command(&spec, id).await?;
2330 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")?;
2331 if let Some(expected) = &operation.destination_store_id {
2332 ensure!(
2333 *expected == store_id,
2334 "destination relay storage was replaced; refusing to replay potentially executed work"
2335 );
2336 } else {
2337 operation.destination_store_id = Some(store_id);
2339 crate::database::save_move_operation(operation)?;
2340 }
2341 ensure!(
2342 relay.snapshot().operational.native_session_id
2343 == operation.destination_native_session_id,
2344 "destination relay native identity changed; refusing queue replay"
2345 );
2346 relay
2347 };
2348 if operation.queue == ResumeQueueDisposition::Start {
2349 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
2350 executor.notify_notice("Starting queued work");
2351 let checkpoint = operation
2352 .restore_artifact()
2353 .context("move queue archive is missing")?;
2354 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2355 ensure!(
2356 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
2357 "move queue checkpoint verification failed"
2358 );
2359 for queued in verified.canonical_session.queued_prompts {
2360 if operation.queue_admission_finished {
2361 continue;
2362 }
2363 ensure!(
2364 !executor.cancellation_requested(),
2365 "move cancelled during queue admission; destination retained"
2366 );
2367 let command = match queued.kind {
2368 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
2369 prompt: queued
2370 .content
2371 .into_iter()
2372 .map(serde_json::from_value)
2373 .collect::<serde_json::Result<_>>()?,
2374 },
2375 CanonicalQueuedCommandKind::SetConfig { key, value } => {
2376 RelayCommand::SetConfig { key, value }
2377 }
2378 };
2379 relay.submit_accepted(queued.command_id, command).await?;
2380 }
2381 }
2382 operation.queue_admission_finished = true;
2383 crate::database::save_move_operation(operation)?;
2384 restore_move_queue_hold(operation);
2385 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
2386 "Queued work was discarded; ready and idle."
2387 } else {
2388 "Queued work was accepted."
2389 };
2390 let source_profile = &operation.source_profile_id;
2391 let source_target = &operation.source_target_template_id;
2392 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
2393 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
2394 let text = if operation.in_place {
2397 format!(
2398 "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."
2399 )
2400 } else {
2401 format!(
2402 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
2403 )
2404 };
2405 relay
2406 .submit(
2407 format!("{}-notice", operation.operation_id),
2408 RelayCommand::RecordNotice { text },
2409 )
2410 .await?;
2411 Ok(())
2412 }
2413
2414 pub(super) fn validate_move_checkpoint(
2415 &self,
2416 operation: &MoveOperation,
2417 preparation: Option<&MovePreparation>,
2418 executor: &(impl CommandExecutor + Sync),
2419 ) -> Result<()> {
2420 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2421 let current = Controller {
2422 config: mj_core::config::Config::load()?,
2423 state: self.state.clone(),
2424 };
2425 ensure!(
2426 current.move_configuration_fingerprint(&operation.selection)?
2427 == operation.configuration_fingerprint,
2428 "destination configuration changed during move"
2429 );
2430 let id = &operation.selection.session_id;
2431 current.validate_move_destination_paths(
2432 &self.state.sessions[id],
2433 operation.selection.target_template_id.as_deref().unwrap(),
2434 executor,
2435 )?;
2436 if !operation.in_place
2437 && operation.workspace_transfer.is_none()
2438 && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
2439 .preflight_resume_repository_sources(
2440 id,
2441 operation.selection.target_template_id.as_deref().unwrap(),
2442 executor,
2443 )?
2444 {
2445 bail!(
2446 "destination repository source is missing checkpoint commit {}; source retained",
2447 mismatch.missing_commit
2448 );
2449 }
2450 if let Some(prepared) = preparation {
2451 let checkpoint = operation
2452 .restore_artifact()
2453 .or(self.state.sessions[id].checkpoint.as_ref())
2454 .context("no move checkpoint")?;
2455 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2456 let actual: Vec<_> = verified
2457 .canonical_session
2458 .queued_prompts
2459 .iter()
2460 .map(|p| p.command_id.as_str())
2461 .collect();
2462 let expected: Vec<_> = prepared
2463 .queued_commands
2464 .iter()
2465 .map(|p| p.command_id.as_str())
2466 .collect();
2467 ensure!(
2468 actual == expected,
2469 "pending queue changed before checkpoint capture; source retained, confirm Move again"
2470 );
2471 }
2472 Ok(())
2473 }
2474}
2475
2476fn validate_preserved_configuration(
2477 profile_id: &str,
2478 accepted: &mj_core::acp::AcceptedSessionConfig,
2479 choices: &mj_core::worker_launch::ProfileConfig,
2480) -> Result<()> {
2481 for (key, value, offered) in [
2482 ("model", accepted.model.as_deref(), &choices.models),
2483 ("effort", accepted.effort.as_deref(), &choices.efforts),
2484 ] {
2485 let Some(value) = value else { continue };
2486 ensure!(
2487 offered.iter().any(|choice| choice.value == value),
2488 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
2489 offered
2490 .iter()
2491 .map(|choice| choice.value.as_str())
2492 .collect::<Vec<_>>()
2493 .join(", ")
2494 );
2495 }
2496 Ok(())
2497}
2498
2499fn move_subagent_policy(
2500 source: &mj_core::state::SessionRecord,
2501 selection: &MoveSelection,
2502) -> mj_core::subagent::SubagentPolicy {
2503 selection
2504 .subagents
2505 .clone()
2506 .or_else(|| source.subagents.clone())
2507 .unwrap_or_default()
2508}
2509
2510pub(crate) fn parent_tools_enabled(
2511 policy: &mj_core::subagent::SubagentPolicy,
2512 harness: mj_core::config::HarnessKind,
2513) -> bool {
2514 policy.for_launch(harness, false).parent_role().is_some()
2515}
2516
2517pub(in crate::controller) fn move_children(
2518 state: &mj_core::state::State,
2519 parent_session_id: &str,
2520) -> Vec<mj_core::subagent::InPlaceSubagent> {
2521 state
2522 .subagents
2523 .values()
2524 .filter(|relation| relation.parent_session_id == parent_session_id)
2525 .filter_map(|relation| {
2526 let child = state.sessions.get(&relation.child_session_id)?;
2527 let state = if child.state.has_live_worker() {
2528 mj_core::subagent::InPlaceSubagentState::Running
2529 } else if child.state == SessionState::Parked {
2530 mj_core::subagent::InPlaceSubagentState::Parked
2531 } else {
2532 return None;
2533 };
2534 Some(mj_core::subagent::InPlaceSubagent {
2535 child_session_id: relation.child_session_id.clone(),
2536 task_name: relation.task_name.clone(),
2537 state,
2538 })
2539 })
2540 .collect()
2541}
2542
2543fn live_move_children(state: &mj_core::state::State, parent_session_id: &str) -> Vec<String> {
2544 move_children(state, parent_session_id)
2545 .into_iter()
2546 .filter(|child| child.state == mj_core::subagent::InPlaceSubagentState::Running)
2547 .map(|child| {
2548 format!(
2549 "{} (child_session_id {})",
2550 child.task_name, child.child_session_id
2551 )
2552 })
2553 .collect()
2554}
2555
2556fn roleless_move_children_error(children: &[String]) -> Option<String> {
2557 (!children.is_empty()).then(|| {
2558 format!(
2559 "cannot in-place Move to a harness without mj-agents parent tools while live sub-agents exist: {}",
2560 children.join(", ")
2561 )
2562 })
2563}
2564
2565fn should_stop_move_subagents(in_place: bool, destination_has_parent_role: bool) -> bool {
2566 !in_place || !destination_has_parent_role
2567}
2568
2569fn subagent_mutations_pending(
2570 requests: &[mj_core::subagent::SubagentToolRequest],
2571 durable_effect_pending: bool,
2572) -> bool {
2573 durable_effect_pending
2574 || requests
2575 .iter()
2576 .any(|request| request.action.mutates_child_state())
2577}
2578
2579pub(super) fn in_place_move_eligible(
2586 source: &mj_core::state::SessionRecord,
2587 selection: &MoveSelection,
2588 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2589 is_subagent: bool,
2590 retry: bool,
2591) -> bool {
2592 !retry
2593 && !is_subagent
2594 && source.target.is_some()
2595 && matches!(
2596 source.state,
2597 SessionState::Running | SessionState::Disconnected
2598 )
2599 && move_environment_change(
2600 source,
2601 selection,
2602 bare_targets_share_environment(
2603 targets,
2604 &source.target_template_id,
2605 selection.target_template_id.as_deref().unwrap_or_default(),
2606 ),
2607 )
2608 .is_none()
2609}
2610
2611fn bare_targets_share_environment(
2616 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2617 source_id: &str,
2618 destination_id: &str,
2619) -> bool {
2620 use mj_core::config::TargetTemplate::{LocalBare, SshBare};
2621 match (targets.get(source_id), targets.get(destination_id)) {
2622 (Some(LocalBare), Some(LocalBare)) => true,
2623 (
2624 Some(SshBare { ssh: source, .. }),
2625 Some(SshBare {
2626 ssh: destination, ..
2627 }),
2628 ) => {
2629 crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
2630 }
2631 _ => false,
2632 }
2633}
2634fn outcome(
2635 operation_id: &str,
2636 selection: &MoveSelection,
2637 status: &str,
2638 error: Option<String>,
2639 recovery: Option<String>,
2640) -> MoveOutcome {
2641 MoveOutcome {
2642 operation_id: operation_id.into(),
2643 session_id: selection.session_id.clone(),
2644 profile_id: selection.profile_id.clone().unwrap_or_default(),
2645 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
2646 outcome: status.into(),
2647 error,
2648 recovery,
2649 }
2650}
2651
2652fn move_environment_change(
2656 source: &mj_core::state::SessionRecord,
2657 selection: &MoveSelection,
2658 same_environment: bool,
2659) -> Option<&'static str> {
2660 if Some(&source.target_template_id) != selection.target_template_id.as_ref()
2661 && !same_environment
2662 {
2663 Some("target changed")
2664 } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
2665 Some("attached mounts changed")
2666 } else if source.resource_allocation != selection.resource_allocation
2667 || (selection.clear_resource_allocation
2668 && (source.container_cpus.is_some() || source.container_memory.is_some()))
2669 {
2670 Some("resource allocation changed")
2671 } else {
2672 None
2673 }
2674}