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 let source = self
1013 .state
1014 .sessions
1015 .get(&selection.session_id)
1016 .context("unknown session")?;
1017 ensure!(
1018 !self.state.subagents.contains_key(&source.id),
1019 "sub-agent sessions cannot move independently of their parent"
1020 );
1021 let previous = crate::database::load_move_operation(&source.id)?;
1022 if let Some(retained) = previous.as_ref().filter(|op| op.holds_source_environment()) {
1025 selection.profile_id = selection
1026 .profile_id
1027 .or(retained.selection.profile_id.clone());
1028 selection.target_template_id = selection
1029 .target_template_id
1030 .or(retained.selection.target_template_id.clone());
1031 selection.subagents = selection.subagents.or(retained.selection.subagents.clone());
1032 if retained.selection.subagents.is_none()
1035 && selection.subagents.as_ref()
1036 == Some(&source.subagents.clone().unwrap_or_default())
1037 {
1038 selection.subagents = None;
1039 }
1040 }
1041 ensure!(
1042 selection.profile_id.is_some() || selection.target_template_id.is_some(),
1043 "move requires a target or profile selection"
1044 );
1045 ensure!(
1046 previous
1047 .as_ref()
1048 .and_then(|op| op.prepared_destination.as_ref())
1049 .is_none_or(|d| !matches!(
1050 d.state,
1051 PreparedDestinationState::CleanupPending { .. }
1052 )),
1053 "EC2 destination cleanup is pending; source retained. Automatic cleanup will retry before another Move"
1054 );
1055
1056 if previous
1057 .as_ref()
1058 .is_some_and(|op| op.holds_source_environment() && !op.checkpoint_retained())
1059 {
1060 bail!(
1061 "this Move cannot be retried. {}",
1062 missing_move_checkpoint_recovery(Some(source))
1063 );
1064 }
1065 let retry = previous.as_ref().is_some_and(|op| {
1066 !matches!(op.phase, MovePhase::Completed)
1067 && (op.restore_artifact().is_some() || source.checkpoint.is_some())
1068 });
1069 ensure!(
1070 matches!(
1071 source.state,
1072 SessionState::Running | SessionState::Disconnected
1073 ) || retry,
1074 "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
1075 source.id
1076 );
1077 selection
1078 .profile_id
1079 .get_or_insert_with(|| source.last_profile.clone());
1080 selection
1081 .target_template_id
1082 .get_or_insert_with(|| source.target_template_id.clone());
1083 selection
1084 .additional_mounts
1085 .get_or_insert_with(|| source.additional_mounts.clone());
1086 ensure!(
1087 !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
1088 "select resource allocation or explicitly clear it, not both"
1089 );
1090 if selection.resource_allocation.is_none()
1091 && matches!(
1092 self.config
1093 .targets
1094 .get(selection.target_template_id.as_deref().unwrap()),
1095 Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1096 )
1097 && (matches!(
1098 source.resource_allocation,
1099 Some(mj_core::state::SessionResourceAllocation::Container { .. })
1100 ) || source.container_cpus.is_some()
1101 || source.container_memory.is_some())
1102 {
1103 selection.clear_resource_allocation = true;
1104 }
1105 if !selection.clear_resource_allocation {
1106 selection.resource_allocation = selection
1107 .resource_allocation
1108 .or_else(|| source.resource_allocation.clone());
1109 }
1110 if selection.clear_resource_allocation
1111 && source.resource_allocation.is_none()
1112 && source.container_cpus.is_none()
1113 && source.container_memory.is_none()
1114 {
1115 selection.clear_resource_allocation = false;
1116 }
1117 if let Some(retained) = previous.as_ref().filter(|op| op.holds_source_environment()) {
1118 if retained.in_place {
1122 selection.workspace = retained.selection.workspace.clone();
1123 }
1124 ensure!(
1125 retained.selection == selection,
1126 "{}",
1127 sealed_selection_difference(&retained.selection, &selection)
1128 );
1129 ensure!(
1130 !retained.in_place
1131 || (retained.source_target.is_some()
1132 && source.target == retained.source_target),
1133 "retained Move target is missing or changed; refusing to recreate it"
1134 );
1135 }
1136 let profile_id = selection.profile_id.as_deref().unwrap();
1137 let target_id = selection.target_template_id.as_deref().unwrap();
1138 let profile = self
1139 .config
1140 .profiles
1141 .get(profile_id)
1142 .context("unknown destination profile")?;
1143 ensure!(
1144 profile.enabled,
1145 "destination profile {profile_id:?} is disabled"
1146 );
1147 if let Some(policy) = &selection.subagents
1148 && policy != &source.subagents.clone().unwrap_or_default()
1149 {
1150 super::profile_config::validate_session_subagent_policy(
1151 &self.config,
1152 profile_id,
1153 policy,
1154 )
1155 .await?;
1156 }
1157 let target = self
1158 .config
1159 .targets
1160 .get(target_id)
1161 .context("unknown destination target")?;
1162 self.validate_muse_resume_destination(source, profile.kind, target_id)?;
1163 super::worktree::resume_compatibility(source, &self.config, target_id)
1164 .map_err(anyhow::Error::msg)?;
1165 super::backend::validate_resource_allocation(
1166 target,
1167 selection.resource_allocation.as_ref(),
1168 )?;
1169 let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
1170 ensure!(
1171 profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
1172 "Muse Code ACP supports one workspace root; attached directories are unsupported"
1173 );
1174 ensure!(
1175 mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
1176 "attached resources are unsupported for this target; select compatible resources explicitly"
1177 );
1178 crate::targets::validate_additional_mounts(mounts)?;
1179 for mount in mounts {
1180 self.validate_mount_source(target_id, &mount.source, executor)?;
1181 }
1182 let planned_conversion =
1183 self.validate_move_destination_paths(source, target_id, executor)?;
1184 ensure!(
1185 profile.home.is_dir(),
1186 "destination profile home is unavailable; configure the profile before moving"
1187 );
1188 super::worker_binary::preflight_worker_binary(target, executor)?;
1189 super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
1190 let source_harness = previous
1191 .as_ref()
1192 .filter(|operation| {
1193 !operation.queue_admission_started && operation.phase != MovePhase::Completed
1194 })
1195 .and_then(|operation| operation.recovery_session.as_ref())
1196 .map_or(source.harness_kind, |record| record.harness_kind);
1197 let cross_harness = profile.kind != source_harness;
1198 if cross_harness {
1199 let cancel = tokio_util::sync::CancellationToken::new();
1200 let resolve =
1201 crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
1202 tokio::pin!(resolve);
1203 loop {
1204 tokio::select! {
1205 result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
1206 _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
1207 if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
1208 }
1209 }
1210 }
1211 }
1212 let conversion = planned_conversion
1215 .map(|conversion| {
1216 super::worktree::raw_conversion_preview(source, &conversion, executor)
1217 .context("describe the move of this checkout into the target")
1218 })
1219 .transpose()?
1220 .map(Box::new);
1221 let (active, mut queued_commands, fingerprint) =
1222 self.move_confirmation(&selection, conversion.as_deref())?;
1223 replace_queued_images_with_placeholders(&mut queued_commands);
1224 let operation_id = previous
1225 .as_ref()
1226 .filter(|operation| {
1227 operation.selection == selection
1228 && operation.phase != MovePhase::Completed
1229 && operation.restore_artifact().is_some()
1230 })
1231 .map(|operation| operation.operation_id.clone())
1232 .unwrap_or(new_command_id("move")?);
1233 let in_place = previous
1234 .as_ref()
1235 .is_some_and(|op| retry && op.in_place && op.selection == selection)
1236 || in_place_move_eligible(
1237 source,
1238 &selection,
1239 &self.config.targets,
1240 self.state.subagents.contains_key(&source.id),
1241 retry,
1242 );
1243 let destination_has_parent_role =
1244 parent_tools_enabled(&move_subagent_policy(source, &selection), profile.kind);
1245 if in_place
1246 && !destination_has_parent_role
1247 && let Some(error) =
1248 roleless_move_children_error(&live_move_children(&self.state, &source.id))
1249 {
1250 bail!("{error}; close or finish those children before moving to this harness");
1251 }
1252 let workspace = if in_place {
1253 None
1254 } else {
1255 Some(self.assess_move_workspace(&selection, executor)?)
1256 };
1257 Ok(MovePreparation {
1258 destination_checks: if !in_place
1259 && matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. })
1260 {
1261 DestinationChecks::AfterProvisioning
1262 } else {
1263 DestinationChecks::Checked
1264 },
1265 workspace,
1266 source_unavailable: false,
1267 in_place,
1268 conversion,
1269 selection,
1270 source_profile_id: source.last_profile.clone(),
1271 source_target_template_id: source.target_template_id.clone(),
1272 cross_harness,
1273 active,
1274 queued_commands,
1275 fingerprint,
1276 operation_id,
1277 })
1278 }
1279
1280 pub async fn move_session_managed_controlled(
1282 &mut self,
1283 request: MoveSessionRequest,
1284 executor: &(impl CommandExecutor + Sync),
1285 manager: &SessionManagerControl,
1286 ) -> Result<MoveOutcome> {
1287 crate::worker_lifecycle::run(&request.preparation.selection.session_id.clone(), "move session managed controlled", executor, async {
1288 let prepared = &request.preparation;
1289 let id = prepared.selection.session_id.clone();
1290 let started = std::time::Instant::now();
1291 executor.notify_notice("Checking destination");
1292 let mut source_relay = MoveSourceRelay::default();
1293 let checked = {
1294 let _checking_destination =
1295 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1296 let mut checked = {
1297 let _timing = MovePhaseTimer::new(&id, "preflight destination checks");
1298 self.prepare_move_session_controlled(prepared.selection.clone(), executor)
1299 .await?
1300 };
1301 if matches!(
1306 self.state.sessions[&id].state,
1307 SessionState::Running | SessionState::Disconnected
1308 ) {
1309 let source_harness = self.state.sessions[&id].harness_kind;
1310 source_relay = {
1311 let _timing = MovePhaseTimer::new(&id, "preflight source lease");
1312 MoveSourceRelay::lease(manager, &id).await?
1313 };
1314 let snapshot = source_relay.snapshot();
1315 let (active, queue, fingerprint) =
1316 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
1317 checked.source_unavailable = snapshot
1318 .as_ref()
1319 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
1320 checked.active = active
1321 || checked.source_unavailable
1322 || snapshot.as_ref().is_some_and(|snapshot| {
1323 let mut operational = snapshot.operational.clone();
1324 operational.queued_prompts.clear();
1325 operational.checkpoint_barrier = None;
1326 !operational.safe_to_replace(source_harness)
1327 });
1328 if let Some(snapshot) = &snapshot {
1329 let _timing = MovePhaseTimer::new(&id, "preflight destination configuration");
1330 self.validate_move_destination_configuration(
1331 &checked.selection,
1332 source_harness,
1333 &snapshot.operational,
1334 )
1335 .await?;
1336 }
1337 checked.queued_commands = queue;
1338 checked.fingerprint = fingerprint;
1339 }
1340 checked
1341 };
1342 ensure!(
1343 checked.fingerprint == prepared.fingerprint,
1344 "session, pending work, or destination configuration changed; prepare and confirm Move again"
1345 );
1346 let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
1347 let source = self.state.sessions[&id].clone();
1348 if let Some(assessment) = &checked.workspace {
1349 checked.selection.workspace.validate(assessment)?;
1350 }
1351 let old_operation = crate::database::load_move_operation(&id)?;
1352 let retry = old_operation.filter(|op| {
1353 op.selection == prepared.selection
1354 && op.phase != MovePhase::Completed
1355 && op.restore_artifact().is_some()
1356 });
1357 if retry.is_none()
1358 && source.state == SessionState::Running
1359 && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
1360 && move_environment_change(&source, &checked.selection, false).is_none()
1361 {
1362 return Ok(outcome(
1363 &prepared.operation_id,
1364 &prepared.selection,
1365 "unchanged",
1366 None,
1367 None,
1368 ));
1369 }
1370 ensure!(
1371 !checked.active || request.acknowledge_interruption,
1372 "active work will be interrupted; confirm Move again with interruption acknowledgement"
1373 );
1374 ensure!(
1375 checked.queued_commands.is_empty() || request.queue.is_some(),
1376 "pending work requires an explicit queue choice: discard or start"
1377 );
1378 let timestamp = now();
1379 let mut operation = match retry {
1380 Some(mut op) => {
1381 ensure!(
1382 !op.queue_admission_started || op.queue == queue,
1383 "queued work may already have run; retry with the original queue choice on the same destination"
1384 );
1385 crate::database::clear_move_cancellation_for_retry(&id)?;
1386 op.configuration_fingerprint =
1387 self.move_configuration_fingerprint(&checked.selection)?;
1388 op.queue = queue;
1389 op.cancellation_requested = false;
1390 op.error = None;
1391 op
1392 }
1393 None => MoveOperation {
1394 prepared_destination: None,
1395 accepted_preparation: None,
1396 acknowledge_interruption: false,
1397 workspace_transfer: if checked.in_place {
1398 None
1399 } else {
1400 Some(
1401 self.new_workspace_transfer(
1402 &id,
1403 &prepared.operation_id,
1404 checked
1405 .workspace
1406 .clone()
1407 .context("Move workspace assessment missing")?,
1408 executor,
1409 )?,
1410 )
1411 },
1412 handoff: None,
1413 source_checkpoint_only: false,
1414 in_place: in_place_move_eligible(
1415 &source,
1416 &checked.selection,
1417 &self.config.targets,
1418 self.state.subagents.contains_key(&id),
1419 false,
1420 ),
1421 operation_id: prepared.operation_id.clone(),
1422 selection: checked.selection.clone(),
1423 source_profile_id: source.last_profile.clone(),
1424 source_target_template_id: source.target_template_id.clone(),
1425 source_target: source.target.clone(),
1426 source_native_session_id: source.native_session_id.clone(),
1427 source_additional_mounts: source.additional_mounts.clone(),
1428 source_resource_allocation: source.resource_allocation.clone(),
1429 destination_target: None,
1430 destination_native_session_id: None,
1431 destination_store_id: None,
1432 configuration_fingerprint: self
1433 .move_configuration_fingerprint(&checked.selection)?,
1434 checkpoint: None,
1435 recovery_session: None,
1436 queue,
1437 phase: MovePhase::Preparing,
1438 queue_admission_started: false,
1439 queue_admission_finished: false,
1440 cancellation_requested: false,
1441 created_at: timestamp.clone(),
1442 updated_at: timestamp,
1443 error: None,
1444 },
1445 };
1446 operation.accepted_preparation = Some(Box::new(checked.clone()));
1447 operation.acknowledge_interruption = request.acknowledge_interruption;
1448 crate::database::save_move_operation(&operation)?;
1449 tracing::info!(
1450 session_id = id,
1451 in_place = operation.in_place,
1452 reason = move_environment_change(
1453 &source,
1454 &checked.selection,
1455 bare_targets_share_environment(
1456 &self.config.targets,
1457 &source.target_template_id,
1458 checked
1459 .selection
1460 .target_template_id
1461 .as_deref()
1462 .unwrap_or_default(),
1463 ),
1464 )
1465 .unwrap_or(if operation.in_place {
1466 "environment unchanged"
1467 } else {
1468 "source environment unavailable or previously released"
1469 }),
1470 "move environment decision"
1471 );
1472 tracing::info!(
1473 session_id = id,
1474 phase = "preflight",
1475 elapsed_ms = started.elapsed().as_millis() as u64,
1476 "move phase completed"
1477 );
1478 let result = Box::pin(self.execute_move(
1481 &mut operation,
1482 Some(&checked),
1483 executor,
1484 manager,
1485 source_relay,
1486 ))
1487 .await;
1488 self.finish_move_result_with_subagent_recovery(
1489 &mut operation,
1490 result,
1491 executor,
1492 manager,
1493 )
1494 .await
1495
1496 }).await
1497 }
1498
1499 fn finish_move_result(
1500 &mut self,
1501 operation: &mut MoveOperation,
1502 result: Result<()>,
1503 executor: &impl CommandExecutor,
1504 ) -> Result<MoveOutcome> {
1505 if result.is_err()
1506 && operation.workspace_transfer.is_some()
1507 && !crate::upgrade::gate().is_open()
1508 {
1509 return Ok(outcome(
1512 &operation.operation_id,
1513 &operation.selection,
1514 "interrupted",
1515 None,
1516 Some("Move will continue after the daemon upgrade".into()),
1517 ));
1518 }
1519 if let Some(saved) = crate::database::load_move_operation(&operation.selection.session_id)?
1520 && saved.operation_id == operation.operation_id
1521 {
1522 operation.prepared_destination = saved.prepared_destination;
1523 }
1524 let result = match result {
1525 Err(error)
1526 if !operation.queue_admission_started
1527 && operation
1528 .prepared_destination
1529 .as_ref()
1530 .is_some_and(|d| d.owns_resource()) =>
1531 {
1532 executor.begin_resumable_move_work()?;
1533 let cleaned = self.cleanup_prepared_move_destination(operation, executor);
1534 executor.end_resumable_move_work()?;
1535 Err(match cleaned {
1536 Ok(()) => error,
1537 Err(cleanup_error) => error.context(format!("EC2 destination cleanup is pending and will retry automatically: {cleanup_error:#}")),
1538 })
1539 }
1540 result => result,
1541 };
1542 let session_id = operation.selection.session_id.clone();
1543 let mut last_error = self
1544 .state
1545 .sessions
1546 .get(&session_id)
1547 .context("Move outcome session is missing")?
1548 .last_error
1549 .clone();
1550 let (status, error, recovery) = match result {
1551 Ok(()) => {
1552 operation.phase = MovePhase::Completed;
1553 operation.error = None;
1554 if last_error
1555 .as_deref()
1556 .is_some_and(|error| error.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
1557 {
1558 last_error = None;
1559 }
1560 ("completed", None, None)
1561 }
1562 Err(error) => {
1563 let phase = match operation.phase {
1564 MovePhase::Preparing => "preparing the destination",
1565 MovePhase::ClosingSource => "checkpointing and suspending the source",
1566 MovePhase::ResumingDestination => "resuming the destination",
1567 MovePhase::StartingQueue => "starting the destination queue",
1568 MovePhase::Completed | MovePhase::Failed | MovePhase::Cancelled => {
1569 "recovering the move"
1570 }
1571 };
1572 let cancelled =
1573 executor.cancellation_requested() || operation.cancellation_requested;
1574 operation.phase = if cancelled {
1575 MovePhase::Cancelled
1576 } else {
1577 MovePhase::Failed
1578 };
1579 operation.cancellation_requested = cancelled;
1580 let recovery = failed_move_recovery(
1581 operation,
1582 self.state.sessions.get(&operation.selection.session_id),
1583 );
1584 let error = format!("{error:#}");
1585 tracing::warn!(
1586 %session_id,
1587 reference = %operation.operation_id,
1588 phase,
1589 cancelled,
1590 %error,
1591 "session move did not finish"
1592 );
1593 last_error = Some(failed_move_message(
1594 Some(phase),
1595 cancelled,
1596 &recovery,
1597 &operation.operation_id,
1598 ));
1599 operation.error = Some(error.clone());
1600 (
1601 if cancelled { "cancelled" } else { "failed" },
1602 Some(error),
1603 Some(recovery),
1604 )
1605 }
1606 };
1607 operation.updated_at = now();
1608 crate::database::save_move_outcome(operation, last_error.as_deref())?;
1609 let record = self
1610 .state
1611 .sessions
1612 .get_mut(&session_id)
1613 .expect("Move outcome session");
1614 record.last_error = last_error;
1615 record.updated_at = operation.updated_at.clone();
1616 Ok(outcome(
1617 &operation.operation_id,
1618 &operation.selection,
1619 status,
1620 error,
1621 recovery,
1622 ))
1623 }
1624
1625 async fn finish_move_result_with_subagent_recovery(
1626 &mut self,
1627 operation: &mut MoveOperation,
1628 result: Result<()>,
1629 executor: &(impl CommandExecutor + Sync),
1630 manager: &SessionManagerControl,
1631 ) -> Result<MoveOutcome> {
1632 let outcome = self.finish_move_result(operation, result, executor)?;
1633 if !operation.in_place
1634 || !matches!(operation.phase, MovePhase::Failed | MovePhase::Cancelled)
1635 {
1636 return Ok(outcome);
1637 }
1638 let session_id = &operation.selection.session_id;
1639 let Some(session) = self.state.sessions.get(session_id) else {
1640 return Ok(outcome);
1641 };
1642 if !matches!(
1643 session.state,
1644 SessionState::Running | SessionState::Disconnected
1645 ) || !parent_tools_enabled(
1646 &session.subagents.clone().unwrap_or_default(),
1647 session.harness_kind,
1648 ) {
1649 return Ok(outcome);
1650 }
1651
1652 for attempt in 1..=2 {
1656 match MoveSourceRelay::set_subagent_admission_via_manager(manager, session_id, true)
1657 .await
1658 {
1659 Ok(true) => {
1660 if attempt > 1 {
1661 tracing::info!(
1662 session_id,
1663 "reopened sub-agent requests on a replacement relay connection"
1664 );
1665 }
1666 break;
1667 }
1668 Ok(false) => break,
1669 Err(error) => {
1670 tracing::warn!(
1671 session_id,
1672 attempt,
1673 error = %error,
1674 "could not reopen sub-agent requests after the in-place Move left its source running"
1675 );
1676 if attempt == 1 {
1677 tokio::task::yield_now().await;
1678 }
1679 }
1680 }
1681 }
1682 Ok(outcome)
1683 }
1684
1685 pub async fn recover_move_managed_controlled(
1686 &mut self,
1687 mut operation: MoveOperation,
1688 executor: &(impl CommandExecutor + Sync),
1689 manager: &SessionManagerControl,
1690 ) -> Result<MoveOutcome> {
1691 crate::worker_lifecycle::run(&operation.selection.session_id.clone(), "recover move managed controlled", executor, async {
1692 let id = operation.selection.session_id.clone();
1693 let result = async {
1694 let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1695 if operation.phase == MovePhase::Preparing
1696 && operation.accepted_preparation.is_some()
1697 && matches!(self.config.targets.get(operation.selection.target_template_id.as_deref().unwrap_or_default()),
1698 Some(mj_core::config::TargetTemplate::AwsEc2 { .. }))
1699 && !operation.cancellation_requested {
1700 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1701 }
1702 if operation.queue_admission_started {
1703 ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1704 self.finish_workspace_transfer(&mut operation, executor)?;
1705 return self.admit_move_queue(&mut operation, executor).await;
1706 }
1707 if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1708 && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1709 if session.state == SessionState::Provisioning || session.target != operation.source_target {
1710 self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1711 }
1712 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1713 }
1714 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1715 let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1716 let sealed = operation.in_place && operation.recovery_session.is_some() && session.state == SessionState::Closing;
1722 if !sealed {
1723 Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1724 }
1725 if !operation.retains_source_environment() {
1729 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1730 }
1731 if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1732 let previous = self.state.sessions[&id].clone();
1733 operation.recovery_session = Some(previous.clone());
1734 let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1735 return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1736 }
1737 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1738 self.cleanup_stopped_target(&id, &cleanup)?;
1739 }
1740 bail!("{}", interrupted_source_stop_message(
1741 "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1742 operation.in_place,
1743 ));
1744 }
1745 match operation.phase {
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 relay = MoveSourceRelay::lease(manager, &id).await?;
1758 return Box::pin(self.execute_move(
1759 &mut operation,
1760 None,
1761 executor,
1762 manager,
1763 relay,
1764 ))
1765 .await;
1766 }
1767 MovePhase::ClosingSource
1768 if operation.in_place
1769 && matches!(
1770 session.state,
1771 SessionState::Running | SessionState::Disconnected
1772 )
1773 && operation.cancellation_requested =>
1774 {
1775 let source = self.state.sessions.get(&id).context("move source is missing")?;
1776 if parent_tools_enabled(
1777 &source.subagents.clone().unwrap_or_default(),
1778 source.harness_kind,
1779 ) {
1780 MoveSourceRelay::set_subagent_admission_via_manager(
1781 manager, &id, true,
1782 )
1783 .await?;
1784 }
1785 bail!("Move was cancelled before source interruption; source retained and sub-agent requests reopened")
1786 }
1787 MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1788 MovePhase::ResumingDestination if session.state == SessionState::Running => {
1789 ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1793 && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1794 "ready destination does not match the move intent");
1795 operation.destination_target = session.target.clone();
1796 operation.destination_native_session_id = session.native_session_id.clone();
1797 operation.queue_admission_started = true;
1798 operation.phase = MovePhase::StartingQueue;
1799 crate::database::save_move_operation(&operation)?;
1800 restore_move_queue_hold(&operation);
1801 ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1802 self.finish_workspace_transfer(&mut operation, executor)?;
1803 self.admit_move_queue(&mut operation, executor).await
1804 }
1805 MovePhase::ResumingDestination => {
1806 let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1807 let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1810 let error = if operation.in_place {
1811 self.retain_failed_in_place_move(&id, previous, cause)?
1812 } else if operation.workspace_transfer.is_some() {
1813 self.rollback_move_destination(&operation, cause, executor)?
1814 } else {
1815 self.rollback_failed_resume(&id, previous, false, cause, executor)?
1816 };
1817 Err(error)
1818 }
1819 MovePhase::ClosingSource => {
1820 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1821 Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1822 }
1823 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1824 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1825 self.cleanup_stopped_target(&id, executor)?;
1826 }
1827 bail!("{}", interrupted_source_stop_message(
1828 "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1829 operation.in_place,
1830 ))
1831 }
1832 _ => bail!("Move requires an explicit retry after the daemon restarted"),
1833 }
1834 }.await;
1835 self.finish_move_result_with_subagent_recovery(
1836 &mut operation,
1837 result,
1838 executor,
1839 manager,
1840 )
1841 .await
1842
1843 }).await
1844 }
1845
1846 async fn recover_move_source_stop(
1847 &mut self,
1848 operation: &mut MoveOperation,
1849 executor: &(impl CommandExecutor + Sync),
1850 manager: &SessionManagerControl,
1851 ) -> Result<()> {
1852 let id = operation.selection.session_id.clone();
1853 if self.state.sessions[&id].state == SessionState::Closing {
1854 self.prepare_move_source_checkpoint(
1855 &id,
1856 executor,
1857 manager,
1858 operation,
1859 &mut MoveSourceRelay::default(),
1860 )
1861 .await?;
1862 let handle = manager
1863 .wait_for_session(&id, std::time::Duration::from_secs(5))
1864 .await?;
1865 let mut lease = handle.lease_connection().await?;
1866 let execution = lease.connection_mut().sync().await?.operational.execution;
1867 if operation.retains_source_environment()
1868 && matches!(
1869 execution,
1870 mj_core::relay::RelayExecutionState::Closing
1871 | mj_core::relay::RelayExecutionState::Closed
1872 )
1873 {
1874 super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1875 lease.release();
1876 return Ok(());
1881 }
1882 lease.release();
1883 if matches!(
1884 execution,
1885 mj_core::relay::RelayExecutionState::Idle
1886 | mj_core::relay::RelayExecutionState::Running
1887 ) {
1888 if operation.cancellation_requested || operation.retains_source_environment() {
1889 let record = self.state.sessions.get_mut(&id).unwrap();
1890 record.state = SessionState::Running;
1891 record.updated_at = now();
1892 record.last_error = Some(
1893 "Move was interrupted before the source was sealed; source retained".into(),
1894 );
1895 crate::database::save_lifecycle_session(record)?;
1896 } else {
1897 Box::pin(self.suspend_session_for_move(
1899 &id,
1900 executor,
1901 manager,
1902 operation,
1903 None,
1904 SourceTargetDisposition::Destroy,
1905 MoveSourceRelay::default(),
1906 ))
1907 .await?;
1908 }
1909 return Ok(());
1910 }
1911 }
1912 ensure!(
1913 !operation.retains_source_environment(),
1914 "retained Move source cannot be proven; refusing target teardown"
1915 );
1916 self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1918 .await?;
1919 Ok(())
1920 }
1921
1922 async fn execute_move(
1925 &mut self,
1926 operation: &mut MoveOperation,
1927 preparation: Option<&MovePreparation>,
1928 executor: &(impl CommandExecutor + Sync),
1929 manager: &SessionManagerControl,
1930 mut source_relay: MoveSourceRelay,
1931 ) -> Result<()> {
1932 let retained = source_relay.owner();
1933 let admission_id = operation.selection.session_id.clone();
1934 crate::worker_lifecycle::run_with_owner(&admission_id, "execute move", executor, retained, async {
1935 let id = operation.selection.session_id.clone();
1936 let mut preparation = preparation.cloned();
1937 if let Some(saved) = crate::database::load_move_operation(&id)?
1938 && saved.operation_id == operation.operation_id
1939 {
1940 operation.prepared_destination = saved.prepared_destination;
1941 }
1942
1943 if !operation.in_place
1944 && !operation.queue_admission_started
1945 && matches!(
1946 self.config.targets.get(
1947 operation
1948 .selection
1949 .target_template_id
1950 .as_deref()
1951 .unwrap_or_default()
1952 ),
1953 Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1954 )
1955 {
1956 drop(source_relay);
1957 source_relay = MoveSourceRelay::default();
1958 self.verify_current_move_configuration(operation)?;
1959 if self.state.sessions[&id].state == SessionState::Error
1963 && operation.recovery_session.is_some()
1964 {
1965 ensure!(
1966 !forget_missing_move_archives(operation),
1967 "the Move's checkpoint archive is missing; nothing is left to restore"
1968 );
1969 self.rollback_move_destination(
1970 operation,
1971 anyhow::anyhow!("clean up the partial Move destination before retry"),
1972 executor,
1973 )?;
1974 let record = self.state.sessions.get_mut(&id).unwrap();
1975 record.state = SessionState::Closing;
1976 crate::database::save_resumed_session(record, None)?;
1977 operation.prepared_destination = crate::database::load_move_operation(&id)?
1978 .context("Move cleanup intent missing")?
1979 .prepared_destination;
1980 }
1981 self.prepare_ec2_move_destination(operation, executor)?;
1982 self.verify_current_move_configuration(operation)?;
1983 if matches!(
1984 self.state.sessions[&id].state,
1985 SessionState::Running | SessionState::Disconnected
1986 ) && operation.recovery_session.is_none()
1987 {
1988 let accepted = operation
1989 .accepted_preparation
1990 .as_ref()
1991 .context("EC2 Move confirmation missing")?;
1992 let mut checked = self
1993 .prepare_move_session_controlled(operation.selection.clone(), executor)
1994 .await?;
1995 let harness = self.state.sessions[&id].harness_kind;
1996 source_relay = MoveSourceRelay::lease(manager, &id).await?;
1997 let snapshot = source_relay.snapshot();
1998 let (active, queue, fingerprint) =
1999 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
2000 let unavailable = snapshot
2001 .as_ref()
2002 .is_none_or(|s| !s.operational.native_session_is_ready());
2003 let active = active
2004 || unavailable
2005 || snapshot.as_ref().is_some_and(|s| {
2006 let mut state = s.operational.clone();
2007 state.queued_prompts.clear();
2008 state.checkpoint_barrier = None;
2009 !state.safe_to_replace(harness)
2010 });
2011 ensure!(
2012 fingerprint == accepted.fingerprint,
2013 "session, pending work, or destination configuration changed; source retained, prepare and confirm Move again"
2014 );
2015 ensure!(
2016 !active || operation.acknowledge_interruption,
2017 "active work will be interrupted; source retained, confirm Move again with interruption acknowledgement"
2018 );
2019 if let Some(snapshot) = snapshot {
2020 self.validate_move_destination_configuration(
2021 &checked.selection,
2022 harness,
2023 &snapshot.operational,
2024 )
2025 .await?;
2026 }
2027 checked.active = active;
2028 checked.queued_commands = queue;
2029 checked.source_unavailable = unavailable;
2030 let assessment = checked
2031 .workspace
2032 .as_mut()
2033 .context("EC2 Move workspace missing")?;
2034 self.assess_prepared_destination(operation, assessment, executor)?;
2035 operation
2036 .workspace_transfer
2037 .as_mut()
2038 .context("EC2 Move transfer missing")?
2039 .assessment = assessment.clone();
2040 crate::database::save_move_operation(operation)?;
2041 preparation = Some(checked);
2042 } else {
2043 let mut assessment = operation
2044 .workspace_transfer
2045 .as_ref()
2046 .context("EC2 Move transfer missing")?
2047 .assessment
2048 .clone();
2049 assessment
2050 .storage
2051 .retain(|s| !s.allocations.contains("destination"));
2052 self.assess_prepared_destination(operation, &mut assessment, executor)?;
2053 }
2054 }
2055 ensure!(
2056 !executor.cancellation_requested(),
2057 "move cancelled before source interruption"
2058 );
2059 if !operation.queue_admission_started {
2060 if forget_missing_move_archives(operation) {
2065 operation.updated_at = now();
2066 crate::database::save_move_operation(operation)?;
2067 bail!("the Move's checkpoint archive is missing; nothing is left to restore");
2068 }
2069 if operation.in_place {
2070 ensure!(
2071 operation.source_target.is_some()
2072 && self.state.sessions[&id].target == operation.source_target,
2073 "retained Move target is missing or changed; refusing to recreate the environment"
2074 );
2075 }
2076 if self.state.sessions[&id].state == SessionState::Error
2077 && let Some(previous) = operation.recovery_session.as_ref()
2078 {
2079 let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
2082 if operation.in_place {
2083 self.retain_failed_in_place_move(&id, previous, cause)?;
2089 } else if operation.workspace_transfer.is_some() {
2090 self.rollback_move_destination(operation, cause, executor)?;
2091 let record = self.state.sessions.get_mut(&id).unwrap();
2092 record.state = SessionState::Closing;
2093 crate::database::save_resumed_session(record, None)?;
2094 } else {
2095 let failure =
2096 self.rollback_failed_resume(&id, previous, false, cause, executor)?;
2097 ensure!(
2098 self.state.sessions[&id].state == SessionState::Stopped,
2099 "{failure:#}"
2100 );
2101 }
2102 }
2103 let state = self.state.sessions[&id].state;
2104 if matches!(state, SessionState::Closing | SessionState::Destroying)
2105 && !(operation.retains_source_environment() && operation.recovery_session.is_some())
2106 {
2107 Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
2108 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
2109 }
2110 if matches!(
2111 self.state.sessions[&id].state,
2112 SessionState::Running | SessionState::Disconnected
2113 ) && operation.destination_target.is_none()
2114 {
2115 let source = self.state.sessions[&id].clone();
2116 let source_has_parent_role = parent_tools_enabled(
2117 &source.subagents.clone().unwrap_or_default(),
2118 source.harness_kind,
2119 );
2120 let destination_has_parent_role = parent_tools_enabled(
2121 &move_subagent_policy(&source, &operation.selection),
2122 operation
2123 .selection
2124 .profile_id
2125 .as_deref()
2126 .and_then(|profile| self.config.profiles.get(profile))
2127 .context("destination profile is missing")?
2128 .kind,
2129 );
2130 if operation.in_place {
2131 operation.phase = MovePhase::ClosingSource;
2135 operation.updated_at = now();
2136 crate::database::save_move_operation(operation)?;
2137 let mut stopped_subagents_for_legacy_worker = false;
2138 if source_has_parent_role {
2139 if !source_relay.is_held() {
2140 source_relay = MoveSourceRelay::lease(manager, &id).await?;
2141 }
2142 match source_relay.drain_subagent_mutations(&id, executor).await? {
2143 SubagentMutationDrain::Drained => {}
2144 SubagentMutationDrain::UnsupportedWorkerProtocol(version) => {
2145 executor.notify_notice(&format!(
2146 "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"
2147 ));
2148 executor.before_move_source_stop().await?;
2149 stopped_subagents_for_legacy_worker = true;
2150 }
2151 }
2152 }
2153 if should_stop_move_subagents(true, destination_has_parent_role)
2154 && !stopped_subagents_for_legacy_worker
2155 {
2156 let current = crate::database::load_state()?;
2157 if let Some(error) =
2158 roleless_move_children_error(&live_move_children(¤t, &id))
2159 {
2160 if source_has_parent_role {
2161 source_relay
2162 .reopen_subagent_mutations(&id)
2163 .await
2164 .context("could not reopen sub-agent requests after refusing Move")?;
2165 }
2166 bail!("{error}; close or finish those children before moving to this harness");
2167 }
2168 if let Err(error) = executor.before_move_source_stop().await {
2169 if source_has_parent_role {
2170 source_relay
2171 .reopen_subagent_mutations(&id)
2172 .await
2173 .context("could not reopen sub-agent requests after stopping children failed")?;
2174 }
2175 return Err(error);
2176 }
2177 }
2178 } else {
2179 executor.before_move_source_stop().await?;
2182 }
2183 if executor.cancellation_requested() {
2184 if operation.in_place && source_has_parent_role {
2185 source_relay.reopen_subagent_mutations(&id).await?;
2186 }
2187 bail!("Move cancelled before source interruption");
2188 }
2189 executor.notify_notice("Stopping source");
2190 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
2191 if !operation.in_place {
2192 operation.phase = MovePhase::ClosingSource;
2193 operation.updated_at = now();
2194 crate::database::save_move_operation(operation)?;
2195 }
2196 let disposition = if operation.in_place {
2200 SourceTargetDisposition::RetainForInPlaceSwap
2201 } else {
2202 SourceTargetDisposition::Destroy
2203 };
2204 let suspended = Box::pin(self.suspend_session_for_move(
2205 &id,
2206 executor,
2207 manager,
2208 operation,
2209 preparation.as_ref(),
2210 disposition,
2211 std::mem::take(&mut source_relay),
2212 ))
2213 .await;
2214 if let Err(error) = suspended {
2215 if operation.in_place
2216 && source_has_parent_role
2217 && matches!(
2218 self.state.sessions[&id].state,
2219 SessionState::Running | SessionState::Disconnected
2220 )
2221 && let Err(reopen) = MoveSourceRelay::set_subagent_admission_via_manager(
2222 manager,
2223 &id,
2224 true,
2225 )
2226 .await
2227 {
2228 return Err(error.context(format!(
2229 "could not reopen sub-agent requests after source checkpoint failed: {reopen:#}"
2230 )));
2231 }
2232 return Err(error);
2233 }
2234 }
2235 drop(source_relay);
2237 if !operation.in_place
2240 && operation.workspace_transfer.is_none()
2241 && self.state.sessions[&id].state == SessionState::Stopped
2242 && self.state.sessions[&id].target.is_some()
2243 {
2244 executor.notify_notice("Cleaning up source");
2245 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
2246 self.cleanup_stopped_target(&id, executor)?;
2247 }
2248 operation.checkpoint = operation
2249 .checkpoint
2250 .clone()
2251 .or_else(|| self.state.sessions[&id].checkpoint.clone());
2252 ensure!(
2253 operation.restore_artifact().is_some(),
2254 "move has no verified checkpoint"
2255 );
2256 ensure!(
2257 !executor.cancellation_requested(),
2258 "move cancelled after source sealing; checkpoint and remaining environment retained"
2259 );
2260 operation.phase = MovePhase::ResumingDestination;
2261 operation.recovery_session = Some(self.state.sessions[&id].clone());
2262 operation.updated_at = now();
2263 crate::database::save_move_operation(operation)?;
2264 executor.reserve_move_destination();
2265 executor.notify_notice("Preparing destination");
2266 if let Some(policy) = &operation.selection.subagents {
2270 let session = self.state.sessions.get_mut(&id).unwrap();
2271 session.subagents = Some(policy.clone());
2272 crate::database::save_resumed_session(session, None)?;
2273 }
2274 if operation.in_place {
2275 Box::pin(self.restore_session_in_place(
2279 &id,
2280 operation.selection.profile_id.as_deref().unwrap(),
2281 operation.selection.target_template_id.as_deref().unwrap(),
2282 executor,
2283 ))
2284 .await?;
2285 } else {
2286 if operation.selection.clear_resource_allocation {
2287 let session = self.state.sessions.get_mut(&id).unwrap();
2288 session.resource_allocation = None;
2289 session.container_cpus = None;
2290 session.container_memory = None;
2291 crate::database::save_resumed_session(session, None)?;
2292 }
2293 if operation.workspace_transfer.is_some() {
2294 self.capture_move_workspace(operation, executor)?;
2295 Box::pin(self.resume_session_for_move(operation, executor)).await?;
2296 let saved = crate::database::load_move_operation(&id)?
2297 .context("Move intent disappeared during restore")?;
2298 operation.workspace_transfer = saved.workspace_transfer;
2299 operation.prepared_destination = saved.prepared_destination;
2300 } else {
2301 Box::pin(self.resume_session_controlled(
2302 &id,
2303 operation.selection.profile_id.as_deref().unwrap(),
2304 operation.selection.target_template_id.as_deref().unwrap(),
2305 SessionResumeOptions {
2306 additional_mounts: operation.selection.additional_mounts.clone(),
2307 resource_allocation: operation.selection.resource_allocation.clone(),
2308 discard_queue: true,
2309 },
2310 executor,
2311 ))
2312 .await?;
2313 }
2314 }
2315 let destination = &self.state.sessions[&id];
2316 operation.destination_target = destination.target.clone();
2317 operation.destination_native_session_id = destination.native_session_id.clone();
2318 operation.phase = MovePhase::StartingQueue;
2320 operation.queue_admission_started = true;
2321 operation.updated_at = now();
2322 crate::database::save_move_operation(operation)?;
2323 }
2324 restore_move_queue_hold(operation);
2325 self.finish_workspace_transfer(operation, executor)?;
2326 self.admit_move_queue(operation, executor).await
2327
2328 }).await
2329 }
2330
2331 async fn admit_move_queue(
2332 &self,
2333 operation: &mut MoveOperation,
2334 executor: &(impl CommandExecutor + Sync),
2335 ) -> Result<()> {
2336 let timing_id = operation.selection.session_id.clone();
2337 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
2338 let id = &operation.selection.session_id;
2339 let mut relay = {
2340 let _checking_destination =
2341 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2342 let destination = &self.state.sessions[id];
2343 ensure!(
2344 destination.state == SessionState::Running
2345 && destination.target == operation.destination_target
2346 && destination.native_session_id == operation.destination_native_session_id,
2347 "cannot prove the same ready destination; refusing to replay potentially executed work"
2348 );
2349 let spec = self.reconnect_command(id)?;
2350 let relay = StandaloneSession::connect_command(&spec, id).await?;
2351 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")?;
2352 if let Some(expected) = &operation.destination_store_id {
2353 ensure!(
2354 *expected == store_id,
2355 "destination relay storage was replaced; refusing to replay potentially executed work"
2356 );
2357 } else {
2358 operation.destination_store_id = Some(store_id);
2360 crate::database::save_move_operation(operation)?;
2361 }
2362 ensure!(
2363 relay.snapshot().operational.native_session_id
2364 == operation.destination_native_session_id,
2365 "destination relay native identity changed; refusing queue replay"
2366 );
2367 relay
2368 };
2369 if operation.queue == ResumeQueueDisposition::Start {
2370 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
2371 executor.notify_notice("Starting queued work");
2372 let checkpoint = operation
2373 .restore_artifact()
2374 .context("move queue archive is missing")?;
2375 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2376 ensure!(
2377 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
2378 "move queue checkpoint verification failed"
2379 );
2380 for queued in verified.canonical_session.queued_prompts {
2381 if operation.queue_admission_finished {
2382 continue;
2383 }
2384 ensure!(
2385 !executor.cancellation_requested(),
2386 "move cancelled during queue admission; destination retained"
2387 );
2388 let command = match queued.kind {
2389 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
2390 prompt: queued
2391 .content
2392 .into_iter()
2393 .map(serde_json::from_value)
2394 .collect::<serde_json::Result<_>>()?,
2395 },
2396 CanonicalQueuedCommandKind::SetConfig { key, value } => {
2397 RelayCommand::SetConfig { key, value }
2398 }
2399 };
2400 relay.submit_accepted(queued.command_id, command).await?;
2401 }
2402 }
2403 operation.queue_admission_finished = true;
2404 crate::database::save_move_operation(operation)?;
2405 restore_move_queue_hold(operation);
2406 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
2407 "Queued work was discarded; ready and idle."
2408 } else {
2409 "Queued work was accepted."
2410 };
2411 let source_profile = &operation.source_profile_id;
2412 let source_target = &operation.source_target_template_id;
2413 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
2414 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
2415 let text = if operation.in_place {
2418 format!(
2419 "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."
2420 )
2421 } else {
2422 format!(
2423 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
2424 )
2425 };
2426 relay
2427 .submit(
2428 format!("{}-notice", operation.operation_id),
2429 RelayCommand::RecordNotice { text },
2430 )
2431 .await?;
2432 Ok(())
2433 }
2434
2435 pub(super) fn validate_move_checkpoint(
2436 &self,
2437 operation: &MoveOperation,
2438 preparation: Option<&MovePreparation>,
2439 github_token: Option<&str>,
2440 executor: &(impl CommandExecutor + Sync),
2441 ) -> Result<()> {
2442 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2443 let current = Controller {
2444 config: mj_core::config::Config::load()?,
2445 state: self.state.clone(),
2446 };
2447 ensure!(
2448 current.move_configuration_fingerprint(&operation.selection)?
2449 == operation.configuration_fingerprint,
2450 "destination configuration changed during move"
2451 );
2452 let id = &operation.selection.session_id;
2453 current.validate_move_destination_paths(
2454 &self.state.sessions[id],
2455 operation.selection.target_template_id.as_deref().unwrap(),
2456 executor,
2457 )?;
2458 if !operation.in_place
2459 && operation.workspace_transfer.is_none()
2460 && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
2461 .preflight_resume_repository_sources_with_token(
2462 id,
2463 operation.selection.target_template_id.as_deref().unwrap(),
2464 github_token,
2465 executor,
2466 )?
2467 {
2468 bail!(
2469 "destination repository source is missing checkpoint commit {}; source retained",
2470 mismatch.missing_commit
2471 );
2472 }
2473 if let Some(prepared) = preparation {
2474 let checkpoint = operation
2475 .restore_artifact()
2476 .or(self.state.sessions[id].checkpoint.as_ref())
2477 .context("no move checkpoint")?;
2478 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2479 let actual: Vec<_> = verified
2480 .canonical_session
2481 .queued_prompts
2482 .iter()
2483 .map(|p| p.command_id.as_str())
2484 .collect();
2485 let expected: Vec<_> = prepared
2486 .queued_commands
2487 .iter()
2488 .map(|p| p.command_id.as_str())
2489 .collect();
2490 ensure!(
2491 actual == expected,
2492 "pending queue changed before checkpoint capture; source retained, confirm Move again"
2493 );
2494 }
2495 Ok(())
2496 }
2497}
2498
2499fn validate_preserved_configuration(
2500 profile_id: &str,
2501 accepted: &mj_core::acp::AcceptedSessionConfig,
2502 choices: &mj_core::worker_launch::ProfileConfig,
2503) -> Result<()> {
2504 for (key, value, offered) in [
2505 ("model", accepted.model.as_deref(), &choices.models),
2506 ("effort", accepted.effort.as_deref(), &choices.efforts),
2507 ] {
2508 let Some(value) = value else { continue };
2509 ensure!(
2510 offered.iter().any(|choice| choice.value == value),
2511 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
2512 offered
2513 .iter()
2514 .map(|choice| choice.value.as_str())
2515 .collect::<Vec<_>>()
2516 .join(", ")
2517 );
2518 }
2519 Ok(())
2520}
2521
2522fn move_subagent_policy(
2523 source: &mj_core::state::SessionRecord,
2524 selection: &MoveSelection,
2525) -> mj_core::subagent::SubagentPolicy {
2526 selection
2527 .subagents
2528 .clone()
2529 .or_else(|| source.subagents.clone())
2530 .unwrap_or_default()
2531}
2532
2533pub(crate) fn parent_tools_enabled(
2534 policy: &mj_core::subagent::SubagentPolicy,
2535 harness: mj_core::config::HarnessKind,
2536) -> bool {
2537 policy.for_launch(harness, false).parent_role().is_some()
2538}
2539
2540pub(in crate::controller) fn move_children(
2541 state: &mj_core::state::State,
2542 parent_session_id: &str,
2543) -> Vec<mj_core::subagent::InPlaceSubagent> {
2544 state
2545 .subagents
2546 .values()
2547 .filter(|relation| relation.parent_session_id == parent_session_id)
2548 .filter_map(|relation| {
2549 let child = state.sessions.get(&relation.child_session_id)?;
2550 let state = if child.state.has_live_worker() {
2551 mj_core::subagent::InPlaceSubagentState::Running
2552 } else if child.state == SessionState::Parked {
2553 mj_core::subagent::InPlaceSubagentState::Parked
2554 } else {
2555 return None;
2556 };
2557 Some(mj_core::subagent::InPlaceSubagent {
2558 child_session_id: relation.child_session_id.clone(),
2559 task_name: relation.task_name.clone(),
2560 state,
2561 })
2562 })
2563 .collect()
2564}
2565
2566fn live_move_children(state: &mj_core::state::State, parent_session_id: &str) -> Vec<String> {
2567 move_children(state, parent_session_id)
2568 .into_iter()
2569 .filter(|child| child.state == mj_core::subagent::InPlaceSubagentState::Running)
2570 .map(|child| {
2571 format!(
2572 "{} (child_session_id {})",
2573 child.task_name, child.child_session_id
2574 )
2575 })
2576 .collect()
2577}
2578
2579fn roleless_move_children_error(children: &[String]) -> Option<String> {
2580 (!children.is_empty()).then(|| {
2581 format!(
2582 "cannot in-place Move to a harness without mj-agents parent tools while live sub-agents exist: {}",
2583 children.join(", ")
2584 )
2585 })
2586}
2587
2588fn should_stop_move_subagents(in_place: bool, destination_has_parent_role: bool) -> bool {
2589 !in_place || !destination_has_parent_role
2590}
2591
2592fn subagent_mutations_pending(
2593 requests: &[mj_core::subagent::SubagentToolRequest],
2594 durable_effect_pending: bool,
2595) -> bool {
2596 durable_effect_pending
2597 || requests
2598 .iter()
2599 .any(|request| request.action.mutates_child_state())
2600}
2601
2602pub(super) fn in_place_move_eligible(
2609 source: &mj_core::state::SessionRecord,
2610 selection: &MoveSelection,
2611 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2612 is_subagent: bool,
2613 retry: bool,
2614) -> bool {
2615 !retry
2616 && !is_subagent
2617 && source.target.is_some()
2618 && matches!(
2619 source.state,
2620 SessionState::Running | SessionState::Disconnected
2621 )
2622 && move_environment_change(
2623 source,
2624 selection,
2625 bare_targets_share_environment(
2626 targets,
2627 &source.target_template_id,
2628 selection.target_template_id.as_deref().unwrap_or_default(),
2629 ),
2630 )
2631 .is_none()
2632}
2633
2634fn bare_targets_share_environment(
2639 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2640 source_id: &str,
2641 destination_id: &str,
2642) -> bool {
2643 use mj_core::config::TargetTemplate::{LocalBare, SshBare};
2644 match (targets.get(source_id), targets.get(destination_id)) {
2645 (Some(LocalBare), Some(LocalBare)) => true,
2646 (
2647 Some(SshBare { ssh: source, .. }),
2648 Some(SshBare {
2649 ssh: destination, ..
2650 }),
2651 ) => {
2652 crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
2653 }
2654 _ => false,
2655 }
2656}
2657fn outcome(
2658 operation_id: &str,
2659 selection: &MoveSelection,
2660 status: &str,
2661 error: Option<String>,
2662 recovery: Option<String>,
2663) -> MoveOutcome {
2664 MoveOutcome {
2665 operation_id: operation_id.into(),
2666 session_id: selection.session_id.clone(),
2667 profile_id: selection.profile_id.clone().unwrap_or_default(),
2668 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
2669 outcome: status.into(),
2670 error,
2671 recovery,
2672 }
2673}
2674
2675fn move_environment_change(
2679 source: &mj_core::state::SessionRecord,
2680 selection: &MoveSelection,
2681 same_environment: bool,
2682) -> Option<&'static str> {
2683 if Some(&source.target_template_id) != selection.target_template_id.as_ref()
2684 && !same_environment
2685 {
2686 Some("target changed")
2687 } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
2688 Some("attached mounts changed")
2689 } else if source.resource_allocation != selection.resource_allocation
2690 || (selection.clear_resource_allocation
2691 && (source.container_cpus.is_some() || source.container_memory.is_some()))
2692 {
2693 Some("resource allocation changed")
2694 } else {
2695 None
2696 }
2697}