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