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 ) || !worker_subagent_queue_enabled(
1658 &session.subagents.clone().unwrap_or_default(),
1659 session.harness_kind,
1660 self.state.subagents.get(session_id),
1661 self.config.agent_mailboxes_enabled(),
1662 ) {
1663 return Ok(outcome);
1664 }
1665
1666 for attempt in 1..=2 {
1670 match MoveSourceRelay::set_subagent_admission_via_manager(manager, session_id, true)
1671 .await
1672 {
1673 Ok(true) => {
1674 if attempt > 1 {
1675 tracing::info!(
1676 session_id,
1677 "reopened sub-agent requests on a replacement relay connection"
1678 );
1679 }
1680 break;
1681 }
1682 Ok(false) => break,
1683 Err(error) => {
1684 tracing::warn!(
1685 session_id,
1686 attempt,
1687 error = %error,
1688 "could not reopen sub-agent requests after the in-place Move left its source running"
1689 );
1690 if attempt == 1 {
1691 tokio::task::yield_now().await;
1692 }
1693 }
1694 }
1695 }
1696 Ok(outcome)
1697 }
1698
1699 pub async fn recover_move_managed_controlled(
1700 &mut self,
1701 mut operation: MoveOperation,
1702 executor: &(impl CommandExecutor + Sync),
1703 manager: &SessionManagerControl,
1704 ) -> Result<MoveOutcome> {
1705 crate::worker_lifecycle::run(&operation.selection.session_id.clone(), "recover move managed controlled", executor, async {
1706 let id = operation.selection.session_id.clone();
1707 let result = async {
1708 let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1709 if operation.phase == MovePhase::Preparing
1710 && operation.accepted_preparation.is_some()
1711 && matches!(self.config.targets.get(operation.selection.target_template_id.as_deref().unwrap_or_default()),
1712 Some(mj_core::config::TargetTemplate::AwsEc2 { .. }))
1713 && !operation.cancellation_requested {
1714 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1715 }
1716 if operation.queue_admission_started {
1717 ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1718 self.finish_workspace_transfer(&mut operation, executor)?;
1719 return self.admit_move_queue(&mut operation, executor).await;
1720 }
1721 if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1722 && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1723 if session.state == SessionState::Provisioning || session.target != operation.source_target {
1724 self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1725 }
1726 return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1727 }
1728 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1729 let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1730 let sealed = operation.in_place && operation.recovery_session.is_some() && session.state == SessionState::Closing;
1736 if !sealed {
1737 Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1738 }
1739 if !operation.retains_source_environment() {
1743 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1744 }
1745 if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1746 let previous = self.state.sessions[&id].clone();
1747 operation.recovery_session = Some(previous.clone());
1748 let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1749 return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1750 }
1751 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1752 self.cleanup_stopped_target(&id, &cleanup)?;
1753 }
1754 bail!("{}", interrupted_source_stop_message(
1755 "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1756 operation.in_place,
1757 ));
1758 }
1759 match operation.phase {
1760 MovePhase::ClosingSource
1761 if operation.in_place
1762 && matches!(
1763 session.state,
1764 SessionState::Running | SessionState::Disconnected
1765 )
1766 && !operation.cancellation_requested =>
1767 {
1768 let relay = MoveSourceRelay::lease(manager, &id).await?;
1772 return Box::pin(self.execute_move(
1773 &mut operation,
1774 None,
1775 executor,
1776 manager,
1777 relay,
1778 ))
1779 .await;
1780 }
1781 MovePhase::ClosingSource
1782 if operation.in_place
1783 && matches!(
1784 session.state,
1785 SessionState::Running | SessionState::Disconnected
1786 )
1787 && operation.cancellation_requested =>
1788 {
1789 let source = self.state.sessions.get(&id).context("move source is missing")?;
1790 if parent_tools_enabled(
1791 &source.subagents.clone().unwrap_or_default(),
1792 source.harness_kind,
1793 ) {
1794 MoveSourceRelay::set_subagent_admission_via_manager(
1795 manager, &id, true,
1796 )
1797 .await?;
1798 }
1799 bail!("Move was cancelled before source interruption; source retained and sub-agent requests reopened")
1800 }
1801 MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1802 MovePhase::ResumingDestination if session.state == SessionState::Running => {
1803 ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1807 && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1808 "ready destination does not match the move intent");
1809 operation.destination_target = session.target.clone();
1810 operation.destination_native_session_id = session.native_session_id.clone();
1811 operation.queue_admission_started = true;
1812 operation.phase = MovePhase::StartingQueue;
1813 crate::database::save_move_operation(&operation)?;
1814 restore_move_queue_hold(&operation);
1815 ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1816 self.finish_workspace_transfer(&mut operation, executor)?;
1817 self.admit_move_queue(&mut operation, executor).await
1818 }
1819 MovePhase::ResumingDestination => {
1820 let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1821 let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1824 let error = if operation.in_place {
1825 self.retain_failed_in_place_move(&id, previous, cause)?
1826 } else if operation.workspace_transfer.is_some() {
1827 self.rollback_move_destination(&operation, cause, executor)?
1828 } else {
1829 self.rollback_failed_resume(&id, previous, false, cause, executor)?
1830 };
1831 Err(error)
1832 }
1833 MovePhase::ClosingSource => {
1834 if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1835 Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1836 }
1837 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1838 if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1839 self.cleanup_stopped_target(&id, executor)?;
1840 }
1841 bail!("{}", interrupted_source_stop_message(
1842 "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1843 operation.in_place,
1844 ))
1845 }
1846 _ => bail!("Move requires an explicit retry after the daemon restarted"),
1847 }
1848 }.await;
1849 self.finish_move_result_with_subagent_recovery(
1850 &mut operation,
1851 result,
1852 executor,
1853 manager,
1854 )
1855 .await
1856
1857 }).await
1858 }
1859
1860 async fn recover_move_source_stop(
1861 &mut self,
1862 operation: &mut MoveOperation,
1863 executor: &(impl CommandExecutor + Sync),
1864 manager: &SessionManagerControl,
1865 ) -> Result<()> {
1866 let id = operation.selection.session_id.clone();
1867 if self.state.sessions[&id].state == SessionState::Closing {
1868 self.prepare_move_source_checkpoint(
1869 &id,
1870 executor,
1871 manager,
1872 operation,
1873 &mut MoveSourceRelay::default(),
1874 )
1875 .await?;
1876 let handle = manager
1877 .wait_for_session(&id, std::time::Duration::from_secs(5))
1878 .await?;
1879 let mut lease = handle.lease_connection().await?;
1880 let execution = lease.connection_mut().sync().await?.operational.execution;
1881 if operation.retains_source_environment()
1882 && matches!(
1883 execution,
1884 mj_core::relay::RelayExecutionState::Closing
1885 | mj_core::relay::RelayExecutionState::Closed
1886 )
1887 {
1888 super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1889 lease.release();
1890 return Ok(());
1895 }
1896 lease.release();
1897 if matches!(
1898 execution,
1899 mj_core::relay::RelayExecutionState::Idle
1900 | mj_core::relay::RelayExecutionState::Running
1901 ) {
1902 if operation.cancellation_requested || operation.retains_source_environment() {
1903 let record = self.state.sessions.get_mut(&id).unwrap();
1904 record.state = SessionState::Running;
1905 record.updated_at = now();
1906 record.last_error = Some(
1907 "Move was interrupted before the source was sealed; source retained".into(),
1908 );
1909 crate::database::save_lifecycle_session(record)?;
1910 } else {
1911 Box::pin(self.suspend_session_for_move(
1913 &id,
1914 executor,
1915 manager,
1916 operation,
1917 None,
1918 SourceTargetDisposition::Destroy,
1919 MoveSourceRelay::default(),
1920 ))
1921 .await?;
1922 }
1923 return Ok(());
1924 }
1925 }
1926 ensure!(
1927 !operation.retains_source_environment(),
1928 "retained Move source cannot be proven; refusing target teardown"
1929 );
1930 self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1932 .await?;
1933 Ok(())
1934 }
1935
1936 async fn execute_move(
1939 &mut self,
1940 operation: &mut MoveOperation,
1941 preparation: Option<&MovePreparation>,
1942 executor: &(impl CommandExecutor + Sync),
1943 manager: &SessionManagerControl,
1944 mut source_relay: MoveSourceRelay,
1945 ) -> Result<()> {
1946 let retained = source_relay.owner();
1947 let admission_id = operation.selection.session_id.clone();
1948 crate::worker_lifecycle::run_with_owner(&admission_id, "execute move", executor, retained, async {
1949 let id = operation.selection.session_id.clone();
1950 let mut preparation = preparation.cloned();
1951 if let Some(saved) = crate::database::load_move_operation(&id)?
1952 && saved.operation_id == operation.operation_id
1953 {
1954 operation.prepared_destination = saved.prepared_destination;
1955 }
1956
1957 if !operation.in_place
1958 && !operation.queue_admission_started
1959 && matches!(
1960 self.config.targets.get(
1961 operation
1962 .selection
1963 .target_template_id
1964 .as_deref()
1965 .unwrap_or_default()
1966 ),
1967 Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1968 )
1969 {
1970 drop(source_relay);
1971 source_relay = MoveSourceRelay::default();
1972 self.verify_current_move_configuration(operation)?;
1973 if self.state.sessions[&id].state == SessionState::Error
1977 && operation.recovery_session.is_some()
1978 {
1979 ensure!(
1980 !forget_missing_move_archives(operation),
1981 "the Move's checkpoint archive is missing; nothing is left to restore"
1982 );
1983 self.rollback_move_destination(
1984 operation,
1985 anyhow::anyhow!("clean up the partial Move destination before retry"),
1986 executor,
1987 )?;
1988 let record = self.state.sessions.get_mut(&id).unwrap();
1989 record.state = SessionState::Closing;
1990 crate::database::save_resumed_session(record, None)?;
1991 operation.prepared_destination = crate::database::load_move_operation(&id)?
1992 .context("Move cleanup intent missing")?
1993 .prepared_destination;
1994 }
1995 self.prepare_ec2_move_destination(operation, executor)?;
1996 self.verify_current_move_configuration(operation)?;
1997 if matches!(
1998 self.state.sessions[&id].state,
1999 SessionState::Running | SessionState::Disconnected
2000 ) && operation.recovery_session.is_none()
2001 {
2002 let accepted = operation
2003 .accepted_preparation
2004 .as_ref()
2005 .context("EC2 Move confirmation missing")?;
2006 let mut checked = self
2007 .prepare_move_session_controlled(operation.selection.clone(), executor)
2008 .await?;
2009 let harness = self.state.sessions[&id].harness_kind;
2010 source_relay = MoveSourceRelay::lease(manager, &id).await?;
2011 let snapshot = source_relay.snapshot();
2012 let (active, queue, fingerprint) =
2013 self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
2014 let unavailable = snapshot
2015 .as_ref()
2016 .is_none_or(|s| !s.operational.native_session_is_ready());
2017 let active = active
2018 || unavailable
2019 || snapshot.as_ref().is_some_and(|s| {
2020 let mut state = s.operational.clone();
2021 state.queued_prompts.clear();
2022 state.checkpoint_barrier = None;
2023 !state.safe_to_replace(harness)
2024 });
2025 ensure!(
2026 fingerprint == accepted.fingerprint,
2027 "session, pending work, or destination configuration changed; source retained, prepare and confirm Move again"
2028 );
2029 ensure!(
2030 !active || operation.acknowledge_interruption,
2031 "active work will be interrupted; source retained, confirm Move again with interruption acknowledgement"
2032 );
2033 if let Some(snapshot) = snapshot {
2034 self.validate_move_destination_configuration(
2035 &checked.selection,
2036 harness,
2037 &snapshot.operational,
2038 )
2039 .await?;
2040 }
2041 checked.active = active;
2042 checked.queued_commands = queue;
2043 checked.source_unavailable = unavailable;
2044 let assessment = checked
2045 .workspace
2046 .as_mut()
2047 .context("EC2 Move workspace missing")?;
2048 self.assess_prepared_destination(operation, assessment, executor)?;
2049 operation
2050 .workspace_transfer
2051 .as_mut()
2052 .context("EC2 Move transfer missing")?
2053 .assessment = assessment.clone();
2054 crate::database::save_move_operation(operation)?;
2055 preparation = Some(checked);
2056 } else {
2057 let mut assessment = operation
2058 .workspace_transfer
2059 .as_ref()
2060 .context("EC2 Move transfer missing")?
2061 .assessment
2062 .clone();
2063 assessment
2064 .storage
2065 .retain(|s| !s.allocations.contains("destination"));
2066 self.assess_prepared_destination(operation, &mut assessment, executor)?;
2067 }
2068 }
2069 ensure!(
2070 !executor.cancellation_requested(),
2071 "move cancelled before source interruption"
2072 );
2073 if !operation.queue_admission_started {
2074 if forget_missing_move_archives(operation) {
2079 operation.updated_at = now();
2080 crate::database::save_move_operation(operation)?;
2081 bail!("the Move's checkpoint archive is missing; nothing is left to restore");
2082 }
2083 if operation.in_place {
2084 ensure!(
2085 operation.source_target.is_some()
2086 && self.state.sessions[&id].target == operation.source_target,
2087 "retained Move target is missing or changed; refusing to recreate the environment"
2088 );
2089 }
2090 if self.state.sessions[&id].state == SessionState::Error
2091 && let Some(previous) = operation.recovery_session.as_ref()
2092 {
2093 let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
2096 if operation.in_place {
2097 self.retain_failed_in_place_move(&id, previous, cause)?;
2103 } else if operation.workspace_transfer.is_some() {
2104 self.rollback_move_destination(operation, cause, executor)?;
2105 let record = self.state.sessions.get_mut(&id).unwrap();
2106 record.state = SessionState::Closing;
2107 crate::database::save_resumed_session(record, None)?;
2108 } else {
2109 let failure =
2110 self.rollback_failed_resume(&id, previous, false, cause, executor)?;
2111 ensure!(
2112 self.state.sessions[&id].state == SessionState::Stopped,
2113 "{failure:#}"
2114 );
2115 }
2116 }
2117 let state = self.state.sessions[&id].state;
2118 if matches!(state, SessionState::Closing | SessionState::Destroying)
2119 && !(operation.retains_source_environment() && operation.recovery_session.is_some())
2120 {
2121 Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
2122 operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
2123 }
2124 if matches!(
2125 self.state.sessions[&id].state,
2126 SessionState::Running | SessionState::Disconnected
2127 ) && operation.destination_target.is_none()
2128 {
2129 let source = self.state.sessions[&id].clone();
2130 let source_has_parent_role = parent_tools_enabled(
2131 &source.subagents.clone().unwrap_or_default(),
2132 source.harness_kind,
2133 );
2134 let destination_has_parent_role = parent_tools_enabled(
2135 &move_subagent_policy(&source, &operation.selection),
2136 operation
2137 .selection
2138 .profile_id
2139 .as_deref()
2140 .and_then(|profile| self.config.profiles.get(profile))
2141 .context("destination profile is missing")?
2142 .kind,
2143 );
2144 if operation.in_place {
2145 operation.phase = MovePhase::ClosingSource;
2149 operation.updated_at = now();
2150 crate::database::save_move_operation(operation)?;
2151 let mut stopped_subagents_for_legacy_worker = false;
2152 if source_has_parent_role {
2153 if !source_relay.is_held() {
2154 source_relay = MoveSourceRelay::lease(manager, &id).await?;
2155 }
2156 match source_relay.drain_subagent_mutations(&id, executor).await? {
2157 SubagentMutationDrain::Drained => {}
2158 SubagentMutationDrain::UnsupportedWorkerProtocol(version) => {
2159 executor.notify_notice(&format!(
2160 "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"
2161 ));
2162 executor.before_move_source_stop().await?;
2163 stopped_subagents_for_legacy_worker = true;
2164 }
2165 }
2166 }
2167 if should_stop_move_subagents(true, destination_has_parent_role)
2168 && !stopped_subagents_for_legacy_worker
2169 {
2170 let current = crate::database::load_state()?;
2171 if let Some(error) =
2172 roleless_move_children_error(&live_move_children(¤t, &id))
2173 {
2174 if source_has_parent_role {
2175 source_relay
2176 .reopen_subagent_mutations(&id)
2177 .await
2178 .context("could not reopen sub-agent requests after refusing Move")?;
2179 }
2180 bail!("{error}; close or finish those children before moving to this harness");
2181 }
2182 if let Err(error) = executor.before_move_source_stop().await {
2183 if source_has_parent_role {
2184 source_relay
2185 .reopen_subagent_mutations(&id)
2186 .await
2187 .context("could not reopen sub-agent requests after stopping children failed")?;
2188 }
2189 return Err(error);
2190 }
2191 }
2192 } else {
2193 executor.before_move_source_stop().await?;
2196 }
2197 if executor.cancellation_requested() {
2198 if operation.in_place && source_has_parent_role {
2199 source_relay.reopen_subagent_mutations(&id).await?;
2200 }
2201 bail!("Move cancelled before source interruption");
2202 }
2203 executor.notify_notice("Stopping source");
2204 let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
2205 if !operation.in_place {
2206 operation.phase = MovePhase::ClosingSource;
2207 operation.updated_at = now();
2208 crate::database::save_move_operation(operation)?;
2209 }
2210 let disposition = if operation.in_place {
2214 SourceTargetDisposition::RetainForInPlaceSwap
2215 } else {
2216 SourceTargetDisposition::Destroy
2217 };
2218 let suspended = Box::pin(self.suspend_session_for_move(
2219 &id,
2220 executor,
2221 manager,
2222 operation,
2223 preparation.as_ref(),
2224 disposition,
2225 std::mem::take(&mut source_relay),
2226 ))
2227 .await;
2228 if let Err(error) = suspended {
2229 if operation.in_place
2230 && source_has_parent_role
2231 && matches!(
2232 self.state.sessions[&id].state,
2233 SessionState::Running | SessionState::Disconnected
2234 )
2235 && let Err(reopen) = MoveSourceRelay::set_subagent_admission_via_manager(
2236 manager,
2237 &id,
2238 true,
2239 )
2240 .await
2241 {
2242 return Err(error.context(format!(
2243 "could not reopen sub-agent requests after source checkpoint failed: {reopen:#}"
2244 )));
2245 }
2246 return Err(error);
2247 }
2248 }
2249 drop(source_relay);
2251 if !operation.in_place
2254 && operation.workspace_transfer.is_none()
2255 && self.state.sessions[&id].state == SessionState::Stopped
2256 && self.state.sessions[&id].target.is_some()
2257 {
2258 executor.notify_notice("Cleaning up source");
2259 let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
2260 self.cleanup_stopped_target(&id, executor)?;
2261 }
2262 operation.checkpoint = operation
2263 .checkpoint
2264 .clone()
2265 .or_else(|| self.state.sessions[&id].checkpoint.clone());
2266 ensure!(
2267 operation.restore_artifact().is_some(),
2268 "move has no verified checkpoint"
2269 );
2270 ensure!(
2271 !executor.cancellation_requested(),
2272 "move cancelled after source sealing; checkpoint and remaining environment retained"
2273 );
2274 operation.phase = MovePhase::ResumingDestination;
2275 operation.recovery_session = Some(self.state.sessions[&id].clone());
2276 operation.updated_at = now();
2277 crate::database::save_move_operation(operation)?;
2278 executor.reserve_move_destination();
2279 executor.notify_notice("Preparing destination");
2280 if let Some(policy) = &operation.selection.subagents {
2284 let session = self.state.sessions.get_mut(&id).unwrap();
2285 session.subagents = Some(policy.clone());
2286 crate::database::save_resumed_session(session, None)?;
2287 }
2288 if operation.in_place {
2289 Box::pin(self.restore_session_in_place(
2293 &id,
2294 operation.selection.profile_id.as_deref().unwrap(),
2295 operation.selection.target_template_id.as_deref().unwrap(),
2296 executor,
2297 ))
2298 .await?;
2299 } else {
2300 if operation.selection.clear_resource_allocation {
2301 let session = self.state.sessions.get_mut(&id).unwrap();
2302 session.resource_allocation = None;
2303 session.container_cpus = None;
2304 session.container_memory = None;
2305 crate::database::save_resumed_session(session, None)?;
2306 }
2307 if operation.workspace_transfer.is_some() {
2308 self.capture_move_workspace(operation, executor)?;
2309 Box::pin(self.resume_session_for_move(operation, executor)).await?;
2310 let saved = crate::database::load_move_operation(&id)?
2311 .context("Move intent disappeared during restore")?;
2312 operation.workspace_transfer = saved.workspace_transfer;
2313 operation.prepared_destination = saved.prepared_destination;
2314 } else {
2315 Box::pin(self.resume_session_controlled(
2316 &id,
2317 operation.selection.profile_id.as_deref().unwrap(),
2318 operation.selection.target_template_id.as_deref().unwrap(),
2319 SessionResumeOptions {
2320 additional_mounts: operation.selection.additional_mounts.clone(),
2321 resource_allocation: operation.selection.resource_allocation.clone(),
2322 discard_queue: true,
2323 },
2324 executor,
2325 ))
2326 .await?;
2327 }
2328 }
2329 let destination = &self.state.sessions[&id];
2330 operation.destination_target = destination.target.clone();
2331 operation.destination_native_session_id = destination.native_session_id.clone();
2332 operation.phase = MovePhase::StartingQueue;
2334 operation.queue_admission_started = true;
2335 operation.updated_at = now();
2336 crate::database::save_move_operation(operation)?;
2337 }
2338 restore_move_queue_hold(operation);
2339 self.finish_workspace_transfer(operation, executor)?;
2340 self.admit_move_queue(operation, executor).await
2341
2342 }).await
2343 }
2344
2345 async fn admit_move_queue(
2346 &self,
2347 operation: &mut MoveOperation,
2348 executor: &(impl CommandExecutor + Sync),
2349 ) -> Result<()> {
2350 let timing_id = operation.selection.session_id.clone();
2351 let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
2352 let id = &operation.selection.session_id;
2353 let mut relay = {
2354 let _checking_destination =
2355 ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2356 let destination = &self.state.sessions[id];
2357 ensure!(
2358 destination.state == SessionState::Running
2359 && destination.target == operation.destination_target
2360 && destination.native_session_id == operation.destination_native_session_id,
2361 "cannot prove the same ready destination; refusing to replay potentially executed work"
2362 );
2363 let spec = self.reconnect_command(id)?;
2364 let relay = StandaloneSession::connect_command(&spec, id).await?;
2365 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")?;
2366 if let Some(expected) = &operation.destination_store_id {
2367 ensure!(
2368 *expected == store_id,
2369 "destination relay storage was replaced; refusing to replay potentially executed work"
2370 );
2371 } else {
2372 operation.destination_store_id = Some(store_id);
2374 crate::database::save_move_operation(operation)?;
2375 }
2376 ensure!(
2377 relay.snapshot().operational.native_session_id
2378 == operation.destination_native_session_id,
2379 "destination relay native identity changed; refusing queue replay"
2380 );
2381 relay
2382 };
2383 if operation.queue == ResumeQueueDisposition::Start {
2384 let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
2385 executor.notify_notice("Starting queued work");
2386 let checkpoint = operation
2387 .restore_artifact()
2388 .context("move queue archive is missing")?;
2389 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2390 ensure!(
2391 verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
2392 "move queue checkpoint verification failed"
2393 );
2394 for queued in verified.canonical_session.queued_prompts {
2395 if operation.queue_admission_finished {
2396 continue;
2397 }
2398 ensure!(
2399 !executor.cancellation_requested(),
2400 "move cancelled during queue admission; destination retained"
2401 );
2402 let command = match queued.kind {
2403 CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
2404 prompt: queued
2405 .content
2406 .into_iter()
2407 .map(serde_json::from_value)
2408 .collect::<serde_json::Result<_>>()?,
2409 },
2410 CanonicalQueuedCommandKind::SetConfig { key, value } => {
2411 RelayCommand::SetConfig { key, value }
2412 }
2413 };
2414 relay.submit_accepted(queued.command_id, command).await?;
2415 }
2416 }
2417 operation.queue_admission_finished = true;
2418 crate::database::save_move_operation(operation)?;
2419 restore_move_queue_hold(operation);
2420 let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
2421 "Queued work was discarded; ready and idle."
2422 } else {
2423 "Queued work was accepted."
2424 };
2425 let source_profile = &operation.source_profile_id;
2426 let source_target = &operation.source_target_template_id;
2427 let destination_profile = operation.selection.profile_id.as_deref().unwrap();
2428 let destination_target = operation.selection.target_template_id.as_deref().unwrap();
2429 let text = if operation.in_place {
2432 format!(
2433 "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."
2434 )
2435 } else {
2436 format!(
2437 "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
2438 )
2439 };
2440 relay
2441 .submit(
2442 format!("{}-notice", operation.operation_id),
2443 RelayCommand::RecordNotice { text },
2444 )
2445 .await?;
2446 Ok(())
2447 }
2448
2449 pub(super) fn validate_move_checkpoint(
2450 &self,
2451 operation: &MoveOperation,
2452 preparation: Option<&MovePreparation>,
2453 github_token: Option<&str>,
2454 executor: &(impl CommandExecutor + Sync),
2455 ) -> Result<()> {
2456 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2457 let current = Controller {
2458 config: mj_core::config::Config::load()?,
2459 state: self.state.clone(),
2460 };
2461 ensure!(
2462 current.move_configuration_fingerprint(&operation.selection)?
2463 == operation.configuration_fingerprint,
2464 "destination configuration changed during move"
2465 );
2466 let id = &operation.selection.session_id;
2467 current.validate_move_destination_paths(
2468 &self.state.sessions[id],
2469 operation.selection.target_template_id.as_deref().unwrap(),
2470 executor,
2471 )?;
2472 if !operation.in_place
2473 && operation.workspace_transfer.is_none()
2474 && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
2475 .preflight_resume_repository_sources_with_token(
2476 id,
2477 operation.selection.target_template_id.as_deref().unwrap(),
2478 github_token,
2479 executor,
2480 )?
2481 {
2482 bail!(
2483 "destination repository source is missing checkpoint commit {}; source retained",
2484 mismatch.missing_commit
2485 );
2486 }
2487 if let Some(prepared) = preparation {
2488 let checkpoint = operation
2489 .restore_artifact()
2490 .or(self.state.sessions[id].checkpoint.as_ref())
2491 .context("no move checkpoint")?;
2492 let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2493 let actual: Vec<_> = verified
2494 .canonical_session
2495 .queued_prompts
2496 .iter()
2497 .map(|p| p.command_id.as_str())
2498 .collect();
2499 let expected: Vec<_> = prepared
2500 .queued_commands
2501 .iter()
2502 .map(|p| p.command_id.as_str())
2503 .collect();
2504 ensure!(
2505 actual == expected,
2506 "pending queue changed before checkpoint capture; source retained, confirm Move again"
2507 );
2508 }
2509 Ok(())
2510 }
2511}
2512
2513fn validate_preserved_configuration(
2514 profile_id: &str,
2515 accepted: &mj_core::acp::AcceptedSessionConfig,
2516 choices: &mj_core::worker_launch::ProfileConfig,
2517) -> Result<()> {
2518 for (key, value, offered) in [
2519 ("model", accepted.model.as_deref(), &choices.models),
2520 ("effort", accepted.effort.as_deref(), &choices.efforts),
2521 ] {
2522 let Some(value) = value else { continue };
2523 ensure!(
2524 offered.iter().any(|choice| choice.value == value),
2525 "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
2526 offered
2527 .iter()
2528 .map(|choice| choice.value.as_str())
2529 .collect::<Vec<_>>()
2530 .join(", ")
2531 );
2532 }
2533 Ok(())
2534}
2535
2536fn move_subagent_policy(
2537 source: &mj_core::state::SessionRecord,
2538 selection: &MoveSelection,
2539) -> mj_core::subagent::SubagentPolicy {
2540 selection
2541 .subagents
2542 .clone()
2543 .or_else(|| source.subagents.clone())
2544 .unwrap_or_default()
2545}
2546
2547fn move_changes_nothing(source: &mj_core::state::SessionRecord, selection: &MoveSelection) -> bool {
2552 source.state == SessionState::Running
2553 && selection.profile_id.as_deref() == Some(source.last_profile.as_str())
2554 && move_environment_change(source, selection, false).is_none()
2555 && move_subagent_policy(source, selection) == source.subagents.clone().unwrap_or_default()
2556}
2557
2558pub(crate) fn parent_tools_enabled(
2559 policy: &mj_core::subagent::SubagentPolicy,
2560 harness: mj_core::config::HarnessKind,
2561) -> bool {
2562 policy.for_launch(harness, false).parent_role().is_some()
2563}
2564
2565pub(crate) fn worker_subagent_queue_enabled(
2569 policy: &mj_core::subagent::SubagentPolicy,
2570 harness: mj_core::config::HarnessKind,
2571 subagent: Option<&mj_core::subagent::SubagentRecord>,
2572 agent_mailboxes_enabled: bool,
2573) -> bool {
2574 let managed_subagent = subagent.is_some();
2575 let launch_policy = policy.for_launch(harness, managed_subagent);
2576 mj_core::worker_launch::resolve_subagent_mcp_role(
2577 &launch_policy,
2578 managed_subagent,
2579 subagent.is_some_and(|record| record.handback_tool),
2580 agent_mailboxes_enabled,
2581 !managed_subagent,
2582 )
2583 .is_some()
2584}
2585
2586pub(in crate::controller) fn move_children(
2587 state: &mj_core::state::State,
2588 parent_session_id: &str,
2589) -> Vec<mj_core::subagent::InPlaceSubagent> {
2590 state
2591 .subagents
2592 .values()
2593 .filter(|relation| relation.parent_session_id == parent_session_id)
2594 .filter_map(|relation| {
2595 let child = state.sessions.get(&relation.child_session_id)?;
2596 let state = if child.state.has_live_worker() {
2597 mj_core::subagent::InPlaceSubagentState::Running
2598 } else if child.state == SessionState::Parked {
2599 mj_core::subagent::InPlaceSubagentState::Parked
2600 } else {
2601 return None;
2602 };
2603 Some(mj_core::subagent::InPlaceSubagent {
2604 child_session_id: relation.child_session_id.clone(),
2605 task_name: relation.task_name.clone(),
2606 state,
2607 })
2608 })
2609 .collect()
2610}
2611
2612fn live_move_children(state: &mj_core::state::State, parent_session_id: &str) -> Vec<String> {
2613 move_children(state, parent_session_id)
2614 .into_iter()
2615 .filter(|child| child.state == mj_core::subagent::InPlaceSubagentState::Running)
2616 .map(|child| {
2617 format!(
2618 "{} (child_session_id {})",
2619 child.task_name, child.child_session_id
2620 )
2621 })
2622 .collect()
2623}
2624
2625fn roleless_move_children_error(children: &[String]) -> Option<String> {
2626 (!children.is_empty()).then(|| {
2627 format!(
2628 "cannot in-place Move to a harness without mj-agents parent tools while live sub-agents exist: {}",
2629 children.join(", ")
2630 )
2631 })
2632}
2633
2634fn should_stop_move_subagents(in_place: bool, destination_has_parent_role: bool) -> bool {
2635 !in_place || !destination_has_parent_role
2636}
2637
2638fn subagent_mutations_pending(
2639 requests: &[mj_core::subagent::SubagentToolRequest],
2640 durable_effect_pending: bool,
2641) -> bool {
2642 durable_effect_pending
2643 || requests
2644 .iter()
2645 .any(|request| request.action.mutates_child_state())
2646}
2647
2648pub(super) fn in_place_move_eligible(
2655 source: &mj_core::state::SessionRecord,
2656 selection: &MoveSelection,
2657 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2658 is_subagent: bool,
2659 retry: bool,
2660) -> bool {
2661 !retry
2662 && !is_subagent
2663 && source.target.is_some()
2664 && matches!(
2665 source.state,
2666 SessionState::Running | SessionState::Disconnected
2667 )
2668 && move_environment_change(
2669 source,
2670 selection,
2671 bare_targets_share_environment(
2672 targets,
2673 &source.target_template_id,
2674 selection.target_template_id.as_deref().unwrap_or_default(),
2675 ),
2676 )
2677 .is_none()
2678}
2679
2680fn bare_targets_share_environment(
2685 targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2686 source_id: &str,
2687 destination_id: &str,
2688) -> bool {
2689 use mj_core::config::TargetTemplate::{LocalBare, SshBare};
2690 match (targets.get(source_id), targets.get(destination_id)) {
2691 (Some(LocalBare), Some(LocalBare)) => true,
2692 (
2693 Some(SshBare { ssh: source, .. }),
2694 Some(SshBare {
2695 ssh: destination, ..
2696 }),
2697 ) => {
2698 crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
2699 }
2700 _ => false,
2701 }
2702}
2703fn outcome(
2704 operation_id: &str,
2705 selection: &MoveSelection,
2706 status: &str,
2707 error: Option<String>,
2708 recovery: Option<String>,
2709) -> MoveOutcome {
2710 MoveOutcome {
2711 operation_id: operation_id.into(),
2712 session_id: selection.session_id.clone(),
2713 profile_id: selection.profile_id.clone().unwrap_or_default(),
2714 target_template_id: selection.target_template_id.clone().unwrap_or_default(),
2715 outcome: status.into(),
2716 error,
2717 recovery,
2718 }
2719}
2720
2721fn move_environment_change(
2725 source: &mj_core::state::SessionRecord,
2726 selection: &MoveSelection,
2727 same_environment: bool,
2728) -> Option<&'static str> {
2729 if Some(&source.target_template_id) != selection.target_template_id.as_ref()
2730 && !same_environment
2731 {
2732 Some("target changed")
2733 } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
2734 Some("attached mounts changed")
2735 } else if source.resource_allocation != selection.resource_allocation
2736 || (selection.clear_resource_allocation
2737 && (source.container_cpus.is_some() || source.container_memory.is_some()))
2738 {
2739 Some("resource allocation changed")
2740 } else {
2741 None
2742 }
2743}