Skip to main content

mj_controller/controller/
move_session.rs

1//! One recoverable stop/restore operation, independent of its initiating viewer.
2
3mod 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/// Execution and retained queue admission are phases of one ownership record.
17#[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
86/// What a recovered, interrupted source stop tells the person.
87///
88/// Recovery keeps the environment promised by an in-place move.
89fn 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
97/// Whether the persisted record shows a source that has finished stopping with
98/// a verified checkpoint to restore.
99///
100/// Only `Stopped` proves the source stop finished. Resume also admits recovery
101/// states, but those do not establish that the source is already stopped.
102fn source_stopped_with_verified_checkpoint(record: &mj_core::state::SessionRecord) -> bool {
103    record.state == SessionState::Stopped && record.checkpoint.is_some()
104}
105
106/// The recovery guidance for a failed Move whose source is already stopped with
107/// a verified checkpoint.
108///
109/// Resume restores the retained checkpoint on the destination profile and
110/// target the Move selected.
111fn 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
130/// The recovery guidance a failed or cancelled Move shows, from its persisted
131/// state.
132///
133/// What is retained is read from the Move record alone
134/// ([`MoveOperation::checkpoint_retained`] and
135/// [`MoveOperation::holds_source_environment`]), the same facts the API's
136/// `move_recovery` and the web page read, so the text never promises a retry
137/// the record cannot support.
138fn 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
198/// The guidance for a Move that holds its source environment but has lost the
199/// checkpoint a retry would restore. Neither a retry nor Resume can bring the
200/// conversation back, and Destroy removes the checkout, so the text says where
201/// the files are first.
202fn 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
219/// The published message of a Move that did not finish. `phase` is where it
220/// stopped, when that is known.
221fn 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
237/// Drop every archive reference whose file is gone from this Move record, and
238/// say whether anything changed.
239///
240/// A Move that took a handoff restores only from it: its `checkpoint` is the
241/// source's older full checkpoint, and restoring that would silently lose the
242/// conversation since. So a missing handoff drops both references, and the
243/// record then says plainly that nothing is retained to restore.
244pub(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
271/// Record, at daemon startup, every Move whose archive is missing. An active
272/// Move is then recovered from the corrected record; a finished one gets its
273/// published guidance rewritten from it, since nothing else publishes it again.
274pub(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
296/// Refresh recovery guidance when a durable background cleanup changes it.
297pub(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
320/// Why a retried sealed Move's selection is refused: each part that differs
321/// from the selection the Move was sealed with, and what to send to match it.
322fn 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
458/// Refresh source state without turning a dead source harness into a Move prerequisite.
459pub 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/// A Move's hold on its source's relay connection, from the confirmation
467/// until the checkpoint that seals the source takes it over. While it is
468/// held, the session actor admits nothing to the relay, so the state the Move
469/// was confirmed against is the state it interrupts. Dropping it hands the
470/// connection back to the actor.
471#[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    /// Take the source's connection from its session actor, which syncs it
483    /// before handing it over. Holds nothing when the source's worker cannot
484    /// be reached; the Move then recovers it without its harness. Leasing
485    /// preserves typed transport failures, unlike the UI's string-valued sync.
486    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    /// The source's state as of its last sync, or `None` when it is unreachable.
536    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    /// Sync the held connection again. A source whose transport died since it
545    /// was leased is let go and reported unreachable.
546    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    /// Close mutation admission at the worker, let its actor continue polling
565    /// accepted requests until they finish, then take the relay back for the
566    /// checkpoint. The worker queue lock is the admission/drain boundary.
567    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            // The worker may have persisted the close before its reply was
598            // lost. Drop this connection and make a best-effort idempotent
599            // reopen on the actor's next relay connection.
600            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    /// Hand the connection to the checkpoint that seals the source.
724    pub(in crate::controller) fn take(
725        &mut self,
726    ) -> Option<super::checkpoint::ControllerRelayLease> {
727        self.0.take()
728    }
729
730    /// Keep holding the source across a worker restart, on the new worker's
731    /// connection.
732    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
782/// Confirmation is an inspector, not transport for attachment bytes: every
783/// surface that shows a move preparation sees a placeholder for each queued
784/// image. Replay always reads the verified archive after destination readiness.
785fn 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    /// Returns the planned conversion when this move turns a local checkout
806    /// into an isolated workspace, so the caller can describe it without
807    /// reading Git a second time.
808    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            // Only what the destination is built from. The dirty counts move
900            // with every keystroke of a live agent, and hashing them would
901            // invalidate the confirmation the person is reading.
902            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        // Only a digest crosses IPC or enters the move record. Configuration
951        // may contain credential-bearing values and must never be copied there.
952        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        // A fresh probe prevents a stale local catalogue from approving a
997        // move that the destination profile will immediately reset.
998        let choices =
999            super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
1000                .await
1001                .with_context(|| {
1002                    format!("discover destination profile {profile_id:?} configuration")
1003                })?;
1004        validate_preserved_configuration(profile_id, &accepted, &choices)
1005    }
1006
1007    pub async fn prepare_move_session_controlled(
1008        &self,
1009        mut selection: MoveSelection,
1010        executor: &(impl CommandExecutor + Sync),
1011    ) -> Result<MovePreparation> {
1012        ensure!(
1013            selection.profile_id.is_some() || selection.target_template_id.is_some(),
1014            "move requires a target or profile selection"
1015        );
1016        let source = self
1017            .state
1018            .sessions
1019            .get(&selection.session_id)
1020            .context("unknown session")?;
1021        ensure!(
1022            !self.state.subagents.contains_key(&source.id),
1023            "sub-agent sessions cannot move independently of their parent"
1024        );
1025        let previous = crate::database::load_move_operation(&source.id)?;
1026        ensure!(
1027            previous
1028                .as_ref()
1029                .and_then(|op| op.prepared_destination.as_ref())
1030                .is_none_or(|d| !matches!(
1031                    d.state,
1032                    PreparedDestinationState::CleanupPending { .. }
1033                )),
1034            "EC2 destination cleanup is pending; source retained. Automatic cleanup will retry before another Move"
1035        );
1036
1037        if previous
1038            .as_ref()
1039            .is_some_and(|op| op.holds_source_environment() && !op.checkpoint_retained())
1040        {
1041            bail!(
1042                "this Move cannot be retried. {}",
1043                missing_move_checkpoint_recovery(Some(source))
1044            );
1045        }
1046        let retry = previous.as_ref().is_some_and(|op| {
1047            !matches!(op.phase, MovePhase::Completed)
1048                && (op.restore_artifact().is_some() || source.checkpoint.is_some())
1049        });
1050        ensure!(
1051            matches!(
1052                source.state,
1053                SessionState::Running | SessionState::Disconnected
1054            ) || retry,
1055            "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
1056            source.id
1057        );
1058        selection
1059            .profile_id
1060            .get_or_insert_with(|| source.last_profile.clone());
1061        selection
1062            .target_template_id
1063            .get_or_insert_with(|| source.target_template_id.clone());
1064        selection
1065            .additional_mounts
1066            .get_or_insert_with(|| source.additional_mounts.clone());
1067        ensure!(
1068            !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
1069            "select resource allocation or explicitly clear it, not both"
1070        );
1071        if selection.resource_allocation.is_none()
1072            && matches!(
1073                self.config
1074                    .targets
1075                    .get(selection.target_template_id.as_deref().unwrap()),
1076                Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1077            )
1078            && (matches!(
1079                source.resource_allocation,
1080                Some(mj_core::state::SessionResourceAllocation::Container { .. })
1081            ) || source.container_cpus.is_some()
1082                || source.container_memory.is_some())
1083        {
1084            selection.clear_resource_allocation = true;
1085        }
1086        if !selection.clear_resource_allocation {
1087            selection.resource_allocation = selection
1088                .resource_allocation
1089                .or_else(|| source.resource_allocation.clone());
1090        }
1091        if selection.clear_resource_allocation
1092            && source.resource_allocation.is_none()
1093            && source.container_cpus.is_none()
1094            && source.container_memory.is_none()
1095        {
1096            selection.clear_resource_allocation = false;
1097        }
1098        if let Some(retained) = previous.as_ref().filter(|op| op.holds_source_environment()) {
1099            // An in-place Move transfers no files, so its file selection is
1100            // not part of what a retry has to repeat. Every surface (the web
1101            // retry has no large-transfer flag at all) retries it the same way.
1102            if retained.in_place {
1103                selection.workspace = retained.selection.workspace.clone();
1104            }
1105            ensure!(
1106                retained.selection == selection,
1107                "{}",
1108                sealed_selection_difference(&retained.selection, &selection)
1109            );
1110            ensure!(
1111                !retained.in_place
1112                    || (retained.source_target.is_some()
1113                        && source.target == retained.source_target),
1114                "retained Move target is missing or changed; refusing to recreate it"
1115            );
1116        }
1117        let profile_id = selection.profile_id.as_deref().unwrap();
1118        let target_id = selection.target_template_id.as_deref().unwrap();
1119        let profile = self
1120            .config
1121            .profiles
1122            .get(profile_id)
1123            .context("unknown destination profile")?;
1124        ensure!(
1125            profile.enabled,
1126            "destination profile {profile_id:?} is disabled"
1127        );
1128        if let Some(policy) = &selection.subagents {
1129            super::profile_config::validate_session_subagent_policy(
1130                &self.config,
1131                profile_id,
1132                policy,
1133            )
1134            .await?;
1135        }
1136        let target = self
1137            .config
1138            .targets
1139            .get(target_id)
1140            .context("unknown destination target")?;
1141        self.validate_muse_resume_destination(source, profile.kind, target_id)?;
1142        super::worktree::resume_compatibility(source, &self.config, target_id)
1143            .map_err(anyhow::Error::msg)?;
1144        super::backend::validate_resource_allocation(
1145            target,
1146            selection.resource_allocation.as_ref(),
1147        )?;
1148        let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
1149        ensure!(
1150            profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
1151            "Muse Code ACP supports one workspace root; attached directories are unsupported"
1152        );
1153        ensure!(
1154            mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
1155            "attached resources are unsupported for this target; select compatible resources explicitly"
1156        );
1157        crate::targets::validate_additional_mounts(mounts)?;
1158        for mount in mounts {
1159            self.validate_mount_source(target_id, &mount.source, executor)?;
1160        }
1161        let planned_conversion =
1162            self.validate_move_destination_paths(source, target_id, executor)?;
1163        ensure!(
1164            profile.home.is_dir(),
1165            "destination profile home is unavailable; configure the profile before moving"
1166        );
1167        super::worker_binary::preflight_worker_binary(target, executor)?;
1168        super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
1169        let source_harness = previous
1170            .as_ref()
1171            .filter(|operation| {
1172                !operation.queue_admission_started && operation.phase != MovePhase::Completed
1173            })
1174            .and_then(|operation| operation.recovery_session.as_ref())
1175            .map_or(source.harness_kind, |record| record.harness_kind);
1176        let cross_harness = profile.kind != source_harness;
1177        if cross_harness {
1178            let cancel = tokio_util::sync::CancellationToken::new();
1179            let resolve =
1180                crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
1181            tokio::pin!(resolve);
1182            loop {
1183                tokio::select! {
1184                    result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
1185                    _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
1186                        if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
1187                    }
1188                }
1189            }
1190        }
1191        // What moving a local checkout into a target really does, computed
1192        // before anything is stopped so a person can confirm it.
1193        let conversion = planned_conversion
1194            .map(|conversion| {
1195                super::worktree::raw_conversion_preview(source, &conversion, executor)
1196                    .context("describe the move of this checkout into the target")
1197            })
1198            .transpose()?
1199            .map(Box::new);
1200        let (active, mut queued_commands, fingerprint) =
1201            self.move_confirmation(&selection, conversion.as_deref())?;
1202        replace_queued_images_with_placeholders(&mut queued_commands);
1203        let operation_id = previous
1204            .as_ref()
1205            .filter(|operation| {
1206                operation.selection == selection
1207                    && operation.phase != MovePhase::Completed
1208                    && operation.restore_artifact().is_some()
1209            })
1210            .map(|operation| operation.operation_id.clone())
1211            .unwrap_or(new_command_id("move")?);
1212        let in_place = previous
1213            .as_ref()
1214            .is_some_and(|op| retry && op.in_place && op.selection == selection)
1215            || in_place_move_eligible(
1216                source,
1217                &selection,
1218                &self.config.targets,
1219                self.state.subagents.contains_key(&source.id),
1220                retry,
1221            );
1222        let destination_has_parent_role =
1223            parent_tools_enabled(&move_subagent_policy(source, &selection), profile.kind);
1224        if in_place
1225            && !destination_has_parent_role
1226            && let Some(error) =
1227                roleless_move_children_error(&live_move_children(&self.state, &source.id))
1228        {
1229            bail!("{error}; close or finish those children before moving to this harness");
1230        }
1231        let workspace = if in_place {
1232            None
1233        } else {
1234            Some(self.assess_move_workspace(&selection, executor)?)
1235        };
1236        Ok(MovePreparation {
1237            destination_checks: if !in_place
1238                && matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. })
1239            {
1240                DestinationChecks::AfterProvisioning
1241            } else {
1242                DestinationChecks::Checked
1243            },
1244            workspace,
1245            source_unavailable: false,
1246            in_place,
1247            conversion,
1248            selection,
1249            source_profile_id: source.last_profile.clone(),
1250            source_target_template_id: source.target_template_id.clone(),
1251            cross_harness,
1252            active,
1253            queued_commands,
1254            fingerprint,
1255            operation_id,
1256        })
1257    }
1258
1259    /// Called with one daemon lifecycle and recovery reservation already held.
1260    pub async fn move_session_managed_controlled(
1261        &mut self,
1262        request: MoveSessionRequest,
1263        executor: &(impl CommandExecutor + Sync),
1264        manager: &SessionManagerControl,
1265    ) -> Result<MoveOutcome> {
1266        crate::worker_lifecycle::run(&request.preparation.selection.session_id.clone(), "move session managed controlled", executor, async {
1267        let prepared = &request.preparation;
1268        let id = prepared.selection.session_id.clone();
1269        let started = std::time::Instant::now();
1270        executor.notify_notice("Checking destination");
1271        let mut source_relay = MoveSourceRelay::default();
1272        let checked = {
1273            let _checking_destination =
1274                ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1275            let mut checked = {
1276                let _timing = MovePhaseTimer::new(&id, "preflight destination checks");
1277                self.prepare_move_session_controlled(prepared.selection.clone(), executor)
1278                    .await?
1279            };
1280            // Destination checks can outlast a turn. Confirm against the relay
1281            // immediately before interruption, and hold its connection from
1282            // here until the checkpoint seals the source, so nothing reaches
1283            // the relay between this confirmation and the interruption.
1284            if matches!(
1285                self.state.sessions[&id].state,
1286                SessionState::Running | SessionState::Disconnected
1287            ) {
1288                let source_harness = self.state.sessions[&id].harness_kind;
1289                source_relay = {
1290                    let _timing = MovePhaseTimer::new(&id, "preflight source lease");
1291                    MoveSourceRelay::lease(manager, &id).await?
1292                };
1293                let snapshot = source_relay.snapshot();
1294                let (active, queue, fingerprint) =
1295                    self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
1296                checked.source_unavailable = snapshot
1297                    .as_ref()
1298                    .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
1299                checked.active = active
1300                    || checked.source_unavailable
1301                    || snapshot.as_ref().is_some_and(|snapshot| {
1302                        let mut operational = snapshot.operational.clone();
1303                        operational.queued_prompts.clear();
1304                        operational.checkpoint_barrier = None;
1305                        !operational.safe_to_replace(source_harness)
1306                    });
1307                if let Some(snapshot) = &snapshot {
1308                    let _timing = MovePhaseTimer::new(&id, "preflight destination configuration");
1309                    self.validate_move_destination_configuration(
1310                        &checked.selection,
1311                        source_harness,
1312                        &snapshot.operational,
1313                    )
1314                    .await?;
1315                }
1316                checked.queued_commands = queue;
1317                checked.fingerprint = fingerprint;
1318            }
1319            checked
1320        };
1321        ensure!(
1322            checked.fingerprint == prepared.fingerprint,
1323            "session, pending work, or destination configuration changed; prepare and confirm Move again"
1324        );
1325        let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
1326        let source = self.state.sessions[&id].clone();
1327        if let Some(assessment) = &checked.workspace {
1328            checked.selection.workspace.validate(assessment)?;
1329        }
1330        let old_operation = crate::database::load_move_operation(&id)?;
1331        let retry = old_operation.filter(|op| {
1332            op.selection == prepared.selection
1333                && op.phase != MovePhase::Completed
1334                && op.restore_artifact().is_some()
1335        });
1336        if retry.is_none()
1337            && source.state == SessionState::Running
1338            && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
1339            && move_environment_change(&source, &checked.selection, false).is_none()
1340        {
1341            return Ok(outcome(
1342                &prepared.operation_id,
1343                &prepared.selection,
1344                "unchanged",
1345                None,
1346                None,
1347            ));
1348        }
1349        ensure!(
1350            !checked.active || request.acknowledge_interruption,
1351            "active work will be interrupted; confirm Move again with interruption acknowledgement"
1352        );
1353        ensure!(
1354            checked.queued_commands.is_empty() || request.queue.is_some(),
1355            "pending work requires an explicit queue choice: discard or start"
1356        );
1357        let timestamp = now();
1358        let mut operation = match retry {
1359            Some(mut op) => {
1360                ensure!(
1361                    !op.queue_admission_started || op.queue == queue,
1362                    "queued work may already have run; retry with the original queue choice on the same destination"
1363                );
1364                crate::database::clear_move_cancellation_for_retry(&id)?;
1365                op.configuration_fingerprint =
1366                    self.move_configuration_fingerprint(&checked.selection)?;
1367                op.queue = queue;
1368                op.cancellation_requested = false;
1369                op.error = None;
1370                op
1371            }
1372            None => MoveOperation {
1373                prepared_destination: None,
1374                accepted_preparation: None,
1375                acknowledge_interruption: false,
1376                workspace_transfer: if checked.in_place {
1377                    None
1378                } else {
1379                    Some(
1380                        self.new_workspace_transfer(
1381                            &id,
1382                            &prepared.operation_id,
1383                            checked
1384                                .workspace
1385                                .clone()
1386                                .context("Move workspace assessment missing")?,
1387                            executor,
1388                        )?,
1389                    )
1390                },
1391                handoff: None,
1392                source_checkpoint_only: false,
1393                in_place: in_place_move_eligible(
1394                    &source,
1395                    &checked.selection,
1396                    &self.config.targets,
1397                    self.state.subagents.contains_key(&id),
1398                    false,
1399                ),
1400                operation_id: prepared.operation_id.clone(),
1401                selection: checked.selection.clone(),
1402                source_profile_id: source.last_profile.clone(),
1403                source_target_template_id: source.target_template_id.clone(),
1404                source_target: source.target.clone(),
1405                source_native_session_id: source.native_session_id.clone(),
1406                source_additional_mounts: source.additional_mounts.clone(),
1407                source_resource_allocation: source.resource_allocation.clone(),
1408                destination_target: None,
1409                destination_native_session_id: None,
1410                destination_store_id: None,
1411                configuration_fingerprint: self
1412                    .move_configuration_fingerprint(&checked.selection)?,
1413                checkpoint: None,
1414                recovery_session: None,
1415                queue,
1416                phase: MovePhase::Preparing,
1417                queue_admission_started: false,
1418                queue_admission_finished: false,
1419                cancellation_requested: false,
1420                created_at: timestamp.clone(),
1421                updated_at: timestamp,
1422                error: None,
1423            },
1424        };
1425        operation.accepted_preparation = Some(Box::new(checked.clone()));
1426        operation.acknowledge_interruption = request.acknowledge_interruption;
1427        crate::database::save_move_operation(&operation)?;
1428        tracing::info!(
1429            session_id = id,
1430            in_place = operation.in_place,
1431            reason = move_environment_change(
1432                &source,
1433                &checked.selection,
1434                bare_targets_share_environment(
1435                    &self.config.targets,
1436                    &source.target_template_id,
1437                    checked
1438                        .selection
1439                        .target_template_id
1440                        .as_deref()
1441                        .unwrap_or_default(),
1442                ),
1443            )
1444            .unwrap_or(if operation.in_place {
1445                "environment unchanged"
1446            } else {
1447                "source environment unavailable or previously released"
1448            }),
1449            "move environment decision"
1450        );
1451        tracing::info!(
1452            session_id = id,
1453            phase = "preflight",
1454            elapsed_ms = started.elapsed().as_millis() as u64,
1455            "move phase completed"
1456        );
1457        // Move nests checkpoint and restore state machines. Heap-own their
1458        // futures so dev builds fit the runtime's ordinary thread stacks.
1459        let result = Box::pin(self.execute_move(
1460            &mut operation,
1461            Some(&checked),
1462            executor,
1463            manager,
1464            source_relay,
1465        ))
1466        .await;
1467        self.finish_move_result_with_subagent_recovery(
1468            &mut operation,
1469            result,
1470            executor,
1471            manager,
1472        )
1473        .await
1474
1475        }).await
1476    }
1477
1478    fn finish_move_result(
1479        &mut self,
1480        operation: &mut MoveOperation,
1481        result: Result<()>,
1482        executor: &impl CommandExecutor,
1483    ) -> Result<MoveOutcome> {
1484        if result.is_err()
1485            && operation.workspace_transfer.is_some()
1486            && !crate::upgrade::gate().is_open()
1487        {
1488            // Handoff cancellation is not user cancellation. Leave the durable
1489            // phase active so the next daemon resumes the accepted Move.
1490            return Ok(outcome(
1491                &operation.operation_id,
1492                &operation.selection,
1493                "interrupted",
1494                None,
1495                Some("Move will continue after the daemon upgrade".into()),
1496            ));
1497        }
1498        if let Some(saved) = crate::database::load_move_operation(&operation.selection.session_id)?
1499            && saved.operation_id == operation.operation_id
1500        {
1501            operation.prepared_destination = saved.prepared_destination;
1502        }
1503        let result = match result {
1504            Err(error)
1505                if !operation.queue_admission_started
1506                    && operation
1507                        .prepared_destination
1508                        .as_ref()
1509                        .is_some_and(|d| d.owns_resource()) =>
1510            {
1511                executor.begin_resumable_move_work()?;
1512                let cleaned = self.cleanup_prepared_move_destination(operation, executor);
1513                executor.end_resumable_move_work()?;
1514                Err(match cleaned {
1515                    Ok(()) => error,
1516                    Err(cleanup_error) => error.context(format!("EC2 destination cleanup is pending and will retry automatically: {cleanup_error:#}")),
1517                })
1518            }
1519            result => result,
1520        };
1521        let session_id = operation.selection.session_id.clone();
1522        let mut last_error = self
1523            .state
1524            .sessions
1525            .get(&session_id)
1526            .context("Move outcome session is missing")?
1527            .last_error
1528            .clone();
1529        let (status, error, recovery) = match result {
1530            Ok(()) => {
1531                operation.phase = MovePhase::Completed;
1532                operation.error = None;
1533                if last_error
1534                    .as_deref()
1535                    .is_some_and(|error| error.starts_with(mj_core::state::MOVE_FAILURE_PREFIX))
1536                {
1537                    last_error = None;
1538                }
1539                ("completed", None, None)
1540            }
1541            Err(error) => {
1542                let phase = match operation.phase {
1543                    MovePhase::Preparing => "preparing the destination",
1544                    MovePhase::ClosingSource => "checkpointing and suspending the source",
1545                    MovePhase::ResumingDestination => "resuming the destination",
1546                    MovePhase::StartingQueue => "starting the destination queue",
1547                    MovePhase::Completed | MovePhase::Failed | MovePhase::Cancelled => {
1548                        "recovering the move"
1549                    }
1550                };
1551                let cancelled =
1552                    executor.cancellation_requested() || operation.cancellation_requested;
1553                operation.phase = if cancelled {
1554                    MovePhase::Cancelled
1555                } else {
1556                    MovePhase::Failed
1557                };
1558                operation.cancellation_requested = cancelled;
1559                let recovery = failed_move_recovery(
1560                    operation,
1561                    self.state.sessions.get(&operation.selection.session_id),
1562                );
1563                let error = format!("{error:#}");
1564                tracing::warn!(
1565                    %session_id,
1566                    reference = %operation.operation_id,
1567                    phase,
1568                    cancelled,
1569                    %error,
1570                    "session move did not finish"
1571                );
1572                last_error = Some(failed_move_message(
1573                    Some(phase),
1574                    cancelled,
1575                    &recovery,
1576                    &operation.operation_id,
1577                ));
1578                operation.error = Some(error.clone());
1579                (
1580                    if cancelled { "cancelled" } else { "failed" },
1581                    Some(error),
1582                    Some(recovery),
1583                )
1584            }
1585        };
1586        operation.updated_at = now();
1587        crate::database::save_move_outcome(operation, last_error.as_deref())?;
1588        let record = self
1589            .state
1590            .sessions
1591            .get_mut(&session_id)
1592            .expect("Move outcome session");
1593        record.last_error = last_error;
1594        record.updated_at = operation.updated_at.clone();
1595        Ok(outcome(
1596            &operation.operation_id,
1597            &operation.selection,
1598            status,
1599            error,
1600            recovery,
1601        ))
1602    }
1603
1604    async fn finish_move_result_with_subagent_recovery(
1605        &mut self,
1606        operation: &mut MoveOperation,
1607        result: Result<()>,
1608        executor: &(impl CommandExecutor + Sync),
1609        manager: &SessionManagerControl,
1610    ) -> Result<MoveOutcome> {
1611        let outcome = self.finish_move_result(operation, result, executor)?;
1612        if !operation.in_place
1613            || !matches!(operation.phase, MovePhase::Failed | MovePhase::Cancelled)
1614        {
1615            return Ok(outcome);
1616        }
1617        let session_id = &operation.selection.session_id;
1618        let Some(session) = self.state.sessions.get(session_id) else {
1619            return Ok(outcome);
1620        };
1621        if !matches!(
1622            session.state,
1623            SessionState::Running | SessionState::Disconnected
1624        ) || !parent_tools_enabled(
1625            &session.subagents.clone().unwrap_or_default(),
1626            session.harness_kind,
1627        ) {
1628            return Ok(outcome);
1629        }
1630
1631        // A failed reply is ambiguous: the worker may already have persisted
1632        // `open=true`. Retry through a fresh actor lease once so the retry uses
1633        // the next relay connection if the first one was abandoned.
1634        for attempt in 1..=2 {
1635            match MoveSourceRelay::set_subagent_admission_via_manager(manager, session_id, true)
1636                .await
1637            {
1638                Ok(true) => {
1639                    if attempt > 1 {
1640                        tracing::info!(
1641                            session_id,
1642                            "reopened sub-agent requests on a replacement relay connection"
1643                        );
1644                    }
1645                    break;
1646                }
1647                Ok(false) => break,
1648                Err(error) => {
1649                    tracing::warn!(
1650                        session_id,
1651                        attempt,
1652                        error = %error,
1653                        "could not reopen sub-agent requests after the in-place Move left its source running"
1654                    );
1655                    if attempt == 1 {
1656                        tokio::task::yield_now().await;
1657                    }
1658                }
1659            }
1660        }
1661        Ok(outcome)
1662    }
1663
1664    pub async fn recover_move_managed_controlled(
1665        &mut self,
1666        mut operation: MoveOperation,
1667        executor: &(impl CommandExecutor + Sync),
1668        manager: &SessionManagerControl,
1669    ) -> Result<MoveOutcome> {
1670        crate::worker_lifecycle::run(&operation.selection.session_id.clone(), "recover move managed controlled", executor, async {
1671        let id = operation.selection.session_id.clone();
1672        let result = async {
1673            let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1674            if operation.phase == MovePhase::Preparing
1675                && operation.accepted_preparation.is_some()
1676                && matches!(self.config.targets.get(operation.selection.target_template_id.as_deref().unwrap_or_default()),
1677                    Some(mj_core::config::TargetTemplate::AwsEc2 { .. }))
1678                && !operation.cancellation_requested {
1679                return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1680            }
1681            if operation.queue_admission_started {
1682                ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1683                self.finish_workspace_transfer(&mut operation, executor)?;
1684                return self.admit_move_queue(&mut operation, executor).await;
1685            }
1686            if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1687                && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1688                if session.state == SessionState::Provisioning || session.target != operation.source_target {
1689                    self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1690                }
1691                return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1692            }
1693            if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1694                let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1695                // An in-place Move writes `recovery_session` only after its
1696                // relay reported Closed, so the seal is a durable fact of the
1697                // Move. Asking the sealed relay again cannot change it and
1698                // fails when its actor is gone, which left the record
1699                // `Closing` (suspending) with nothing to finish it.
1700                let sealed = operation.in_place && operation.recovery_session.is_some() && session.state == SessionState::Closing;
1701                if !sealed {
1702                    Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1703                }
1704                // A Move that keeps the environment restores only from the
1705                // handoff its seal recorded. Taking the source's older full
1706                // checkpoint here would make a lost handoff look retained.
1707                if !operation.retains_source_environment() {
1708                    operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1709                }
1710                if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1711                    let previous = self.state.sessions[&id].clone();
1712                    operation.recovery_session = Some(previous.clone());
1713                    let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1714                    return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1715                }
1716                if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1717                    self.cleanup_stopped_target(&id, &cleanup)?;
1718                }
1719                bail!("{}", interrupted_source_stop_message(
1720                    "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1721                    operation.in_place,
1722                ));
1723            }
1724            match operation.phase {
1725                MovePhase::ClosingSource
1726                    if operation.in_place
1727                        && matches!(
1728                            session.state,
1729                            SessionState::Running | SessionState::Disconnected
1730                        )
1731                        && !operation.cancellation_requested =>
1732                {
1733                    // The phase is durable before the worker admission gate
1734                    // closes. If the daemon died before the session reached
1735                    // Closing, resume the gate/drain/stop sequence on startup.
1736                    let relay = MoveSourceRelay::lease(manager, &id).await?;
1737                    return Box::pin(self.execute_move(
1738                        &mut operation,
1739                        None,
1740                        executor,
1741                        manager,
1742                        relay,
1743                    ))
1744                    .await;
1745                }
1746                MovePhase::ClosingSource
1747                    if operation.in_place
1748                        && matches!(
1749                            session.state,
1750                            SessionState::Running | SessionState::Disconnected
1751                        )
1752                        && operation.cancellation_requested =>
1753                {
1754                    let source = self.state.sessions.get(&id).context("move source is missing")?;
1755                    if parent_tools_enabled(
1756                        &source.subagents.clone().unwrap_or_default(),
1757                        source.harness_kind,
1758                    ) {
1759                        MoveSourceRelay::set_subagent_admission_via_manager(
1760                            manager, &id, true,
1761                        )
1762                        .await?;
1763                    }
1764                    bail!("Move was cancelled before source interruption; source retained and sub-agent requests reopened")
1765                }
1766                MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1767                MovePhase::ResumingDestination if session.state == SessionState::Running => {
1768                    // Running is installed only after native readiness and the
1769                    // handoff. A crash before the next intent write is safe to
1770                    // advance, since restore started with an empty queue.
1771                    ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1772                        && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1773                        "ready destination does not match the move intent");
1774                    operation.destination_target = session.target.clone();
1775                    operation.destination_native_session_id = session.native_session_id.clone();
1776                    operation.queue_admission_started = true;
1777                    operation.phase = MovePhase::StartingQueue;
1778                    crate::database::save_move_operation(&operation)?;
1779                    restore_move_queue_hold(&operation);
1780                    ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1781                    self.finish_workspace_transfer(&mut operation, executor)?;
1782                    self.admit_move_queue(&mut operation, executor).await
1783                }
1784                MovePhase::ResumingDestination => {
1785                    let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1786                    // The retained environment belongs to this Move, even
1787                    // when only part of its worker installation completed.
1788                    let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1789                    let error = if operation.in_place {
1790                        self.retain_failed_in_place_move(&id, previous, cause)?
1791                    } else if operation.workspace_transfer.is_some() {
1792                        self.rollback_move_destination(&operation, cause, executor)?
1793                    } else {
1794                        self.rollback_failed_resume(&id, previous, false, cause, executor)?
1795                    };
1796                    Err(error)
1797                }
1798                MovePhase::ClosingSource => {
1799                    if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1800                        Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1801                    }
1802                    operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1803                    if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1804                        self.cleanup_stopped_target(&id, executor)?;
1805                    }
1806                    bail!("{}", interrupted_source_stop_message(
1807                        "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1808                        operation.in_place,
1809                    ))
1810                }
1811                _ => bail!("Move requires an explicit retry after the daemon restarted"),
1812            }
1813        }.await;
1814        self.finish_move_result_with_subagent_recovery(
1815            &mut operation,
1816            result,
1817            executor,
1818            manager,
1819        )
1820        .await
1821
1822        }).await
1823    }
1824
1825    async fn recover_move_source_stop(
1826        &mut self,
1827        operation: &mut MoveOperation,
1828        executor: &(impl CommandExecutor + Sync),
1829        manager: &SessionManagerControl,
1830    ) -> Result<()> {
1831        let id = operation.selection.session_id.clone();
1832        if self.state.sessions[&id].state == SessionState::Closing {
1833            self.prepare_move_source_checkpoint(
1834                &id,
1835                executor,
1836                manager,
1837                operation,
1838                &mut MoveSourceRelay::default(),
1839            )
1840            .await?;
1841            let handle = manager
1842                .wait_for_session(&id, std::time::Duration::from_secs(5))
1843                .await?;
1844            let mut lease = handle.lease_connection().await?;
1845            let execution = lease.connection_mut().sync().await?.operational.execution;
1846            if operation.retains_source_environment()
1847                && matches!(
1848                    execution,
1849                    mj_core::relay::RelayExecutionState::Closing
1850                        | mj_core::relay::RelayExecutionState::Closed
1851                )
1852            {
1853                super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1854                lease.release();
1855                // The source is sealed whether or not its handoff survived.
1856                // A missing handoff is already recorded on the Move, whose
1857                // owner retains the environment and publishes the loss;
1858                // refusing here would leave the session suspending instead.
1859                return Ok(());
1860            }
1861            lease.release();
1862            if matches!(
1863                execution,
1864                mj_core::relay::RelayExecutionState::Idle
1865                    | mj_core::relay::RelayExecutionState::Running
1866            ) {
1867                if operation.cancellation_requested || operation.retains_source_environment() {
1868                    let record = self.state.sessions.get_mut(&id).unwrap();
1869                    record.state = SessionState::Running;
1870                    record.updated_at = now();
1871                    record.last_error = Some(
1872                        "Move was interrupted before the source was sealed; source retained".into(),
1873                    );
1874                    crate::database::save_lifecycle_session(record)?;
1875                } else {
1876                    // Fresh-environment moves finish their admitted source stop.
1877                    Box::pin(self.suspend_session_for_move(
1878                        &id,
1879                        executor,
1880                        manager,
1881                        operation,
1882                        None,
1883                        SourceTargetDisposition::Destroy,
1884                        MoveSourceRelay::default(),
1885                    ))
1886                    .await?;
1887                }
1888                return Ok(());
1889            }
1890        }
1891        ensure!(
1892            !operation.retains_source_environment(),
1893            "retained Move source cannot be proven; refusing target teardown"
1894        );
1895        // A move's source stop was admitted by the move itself.
1896        self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1897            .await?;
1898        Ok(())
1899    }
1900
1901    /// `source_relay` is the source connection the Move confirmed against,
1902    /// held until the checkpoint that seals the source takes it over.
1903    async fn execute_move(
1904        &mut self,
1905        operation: &mut MoveOperation,
1906        preparation: Option<&MovePreparation>,
1907        executor: &(impl CommandExecutor + Sync),
1908        manager: &SessionManagerControl,
1909        mut source_relay: MoveSourceRelay,
1910    ) -> Result<()> {
1911        let retained = source_relay.owner();
1912        let admission_id = operation.selection.session_id.clone();
1913        crate::worker_lifecycle::run_with_owner(&admission_id, "execute move", executor, retained, async {
1914        let id = operation.selection.session_id.clone();
1915        let mut preparation = preparation.cloned();
1916        if let Some(saved) = crate::database::load_move_operation(&id)?
1917            && saved.operation_id == operation.operation_id
1918        {
1919            operation.prepared_destination = saved.prepared_destination;
1920        }
1921
1922        if !operation.in_place
1923            && !operation.queue_admission_started
1924            && matches!(
1925                self.config.targets.get(
1926                    operation
1927                        .selection
1928                        .target_template_id
1929                        .as_deref()
1930                        .unwrap_or_default()
1931                ),
1932                Some(mj_core::config::TargetTemplate::AwsEc2 { .. })
1933            )
1934        {
1935            drop(source_relay);
1936            source_relay = MoveSourceRelay::default();
1937            self.verify_current_move_configuration(operation)?;
1938            // Resolve the previous destination before acquiring another one. The
1939            // rollback must never interpret a newly prepared instance as the old
1940            // destination whose session record still needs restoring.
1941            if self.state.sessions[&id].state == SessionState::Error
1942                && operation.recovery_session.is_some()
1943            {
1944                ensure!(
1945                    !forget_missing_move_archives(operation),
1946                    "the Move's checkpoint archive is missing; nothing is left to restore"
1947                );
1948                self.rollback_move_destination(
1949                    operation,
1950                    anyhow::anyhow!("clean up the partial Move destination before retry"),
1951                    executor,
1952                )?;
1953                let record = self.state.sessions.get_mut(&id).unwrap();
1954                record.state = SessionState::Closing;
1955                crate::database::save_resumed_session(record, None)?;
1956                operation.prepared_destination = crate::database::load_move_operation(&id)?
1957                    .context("Move cleanup intent missing")?
1958                    .prepared_destination;
1959            }
1960            self.prepare_ec2_move_destination(operation, executor)?;
1961            self.verify_current_move_configuration(operation)?;
1962            if matches!(
1963                self.state.sessions[&id].state,
1964                SessionState::Running | SessionState::Disconnected
1965            ) && operation.recovery_session.is_none()
1966            {
1967                let accepted = operation
1968                    .accepted_preparation
1969                    .as_ref()
1970                    .context("EC2 Move confirmation missing")?;
1971                let mut checked = self
1972                    .prepare_move_session_controlled(operation.selection.clone(), executor)
1973                    .await?;
1974                let harness = self.state.sessions[&id].harness_kind;
1975                source_relay = MoveSourceRelay::lease(manager, &id).await?;
1976                let snapshot = source_relay.snapshot();
1977                let (active, queue, fingerprint) =
1978                    self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
1979                let unavailable = snapshot
1980                    .as_ref()
1981                    .is_none_or(|s| !s.operational.native_session_is_ready());
1982                let active = active
1983                    || unavailable
1984                    || snapshot.as_ref().is_some_and(|s| {
1985                        let mut state = s.operational.clone();
1986                        state.queued_prompts.clear();
1987                        state.checkpoint_barrier = None;
1988                        !state.safe_to_replace(harness)
1989                    });
1990                ensure!(
1991                    fingerprint == accepted.fingerprint,
1992                    "session, pending work, or destination configuration changed; source retained, prepare and confirm Move again"
1993                );
1994                ensure!(
1995                    !active || operation.acknowledge_interruption,
1996                    "active work will be interrupted; source retained, confirm Move again with interruption acknowledgement"
1997                );
1998                if let Some(snapshot) = snapshot {
1999                    self.validate_move_destination_configuration(
2000                        &checked.selection,
2001                        harness,
2002                        &snapshot.operational,
2003                    )
2004                    .await?;
2005                }
2006                checked.active = active;
2007                checked.queued_commands = queue;
2008                checked.source_unavailable = unavailable;
2009                let assessment = checked
2010                    .workspace
2011                    .as_mut()
2012                    .context("EC2 Move workspace missing")?;
2013                self.assess_prepared_destination(operation, assessment, executor)?;
2014                operation
2015                    .workspace_transfer
2016                    .as_mut()
2017                    .context("EC2 Move transfer missing")?
2018                    .assessment = assessment.clone();
2019                crate::database::save_move_operation(operation)?;
2020                preparation = Some(checked);
2021            } else {
2022                let mut assessment = operation
2023                    .workspace_transfer
2024                    .as_ref()
2025                    .context("EC2 Move transfer missing")?
2026                    .assessment
2027                    .clone();
2028                assessment
2029                    .storage
2030                    .retain(|s| !s.allocations.contains("destination"));
2031                self.assess_prepared_destination(operation, &mut assessment, executor)?;
2032            }
2033        }
2034        ensure!(
2035            !executor.cancellation_requested(),
2036            "move cancelled before source interruption"
2037        );
2038        if !operation.queue_admission_started {
2039            // A retry restores from the archive this record names. When the
2040            // file is gone, record that before touching the session, so the
2041            // failure publishes what is really retained and the session stays
2042            // in the state Destroy can act on.
2043            if forget_missing_move_archives(operation) {
2044                operation.updated_at = now();
2045                crate::database::save_move_operation(operation)?;
2046                bail!("the Move's checkpoint archive is missing; nothing is left to restore");
2047            }
2048            if operation.in_place {
2049                ensure!(
2050                    operation.source_target.is_some()
2051                        && self.state.sessions[&id].target == operation.source_target,
2052                    "retained Move target is missing or changed; refusing to recreate the environment"
2053                );
2054            }
2055            if self.state.sessions[&id].state == SessionState::Error
2056                && let Some(previous) = operation.recovery_session.as_ref()
2057            {
2058                // A prior failed teardown retains its exact partial checkout.
2059                // Restore source identity only after that owner is stopped.
2060                let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
2061                if operation.in_place {
2062                    // The record stays `Error`: `restore_session_in_place`
2063                    // accepts a retained environment in that state, and every
2064                    // failure before the swap leaves it there for Destroy or
2065                    // another retry. Marking it `Closing` here left a failed
2066                    // retry suspending with nothing to finish it.
2067                    self.retain_failed_in_place_move(&id, previous, cause)?;
2068                } else if operation.workspace_transfer.is_some() {
2069                    self.rollback_move_destination(operation, cause, executor)?;
2070                    let record = self.state.sessions.get_mut(&id).unwrap();
2071                    record.state = SessionState::Closing;
2072                    crate::database::save_resumed_session(record, None)?;
2073                } else {
2074                    let failure =
2075                        self.rollback_failed_resume(&id, previous, false, cause, executor)?;
2076                    ensure!(
2077                        self.state.sessions[&id].state == SessionState::Stopped,
2078                        "{failure:#}"
2079                    );
2080                }
2081            }
2082            let state = self.state.sessions[&id].state;
2083            if matches!(state, SessionState::Closing | SessionState::Destroying)
2084                && !(operation.retains_source_environment() && operation.recovery_session.is_some())
2085            {
2086                Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
2087                operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
2088            }
2089            if matches!(
2090                self.state.sessions[&id].state,
2091                SessionState::Running | SessionState::Disconnected
2092            ) && operation.destination_target.is_none()
2093            {
2094                let source = self.state.sessions[&id].clone();
2095                let source_has_parent_role = parent_tools_enabled(
2096                    &source.subagents.clone().unwrap_or_default(),
2097                    source.harness_kind,
2098                );
2099                let destination_has_parent_role = parent_tools_enabled(
2100                    &move_subagent_policy(&source, &operation.selection),
2101                    operation
2102                        .selection
2103                        .profile_id
2104                        .as_deref()
2105                        .and_then(|profile| self.config.profiles.get(profile))
2106                        .context("destination profile is missing")?
2107                        .kind,
2108                );
2109                if operation.in_place {
2110                    // Persist the phase before closing worker admission. A daemon
2111                    // crash from this point can reopen the relay, finish the
2112                    // drain, or observe that source sealing already began.
2113                    operation.phase = MovePhase::ClosingSource;
2114                    operation.updated_at = now();
2115                    crate::database::save_move_operation(operation)?;
2116                    let mut stopped_subagents_for_legacy_worker = false;
2117                    if source_has_parent_role {
2118                        if !source_relay.is_held() {
2119                            source_relay = MoveSourceRelay::lease(manager, &id).await?;
2120                        }
2121                        match source_relay.drain_subagent_mutations(&id, executor).await? {
2122                            SubagentMutationDrain::Drained => {}
2123                            SubagentMutationDrain::UnsupportedWorkerProtocol(version) => {
2124                                executor.notify_notice(&format!(
2125                                    "The source worker uses relay protocol {version}; protocol 32 is required to safely drain sub-agent requests, so stopping sub-agents before this in-place Move"
2126                                ));
2127                                executor.before_move_source_stop().await?;
2128                                stopped_subagents_for_legacy_worker = true;
2129                            }
2130                        }
2131                    }
2132                    if should_stop_move_subagents(true, destination_has_parent_role)
2133                        && !stopped_subagents_for_legacy_worker
2134                    {
2135                        let current = crate::database::load_state()?;
2136                        if let Some(error) =
2137                            roleless_move_children_error(&live_move_children(&current, &id))
2138                        {
2139                            if source_has_parent_role {
2140                                source_relay
2141                                    .reopen_subagent_mutations(&id)
2142                                    .await
2143                                    .context("could not reopen sub-agent requests after refusing Move")?;
2144                            }
2145                            bail!("{error}; close or finish those children before moving to this harness");
2146                        }
2147                        if let Err(error) = executor.before_move_source_stop().await {
2148                            if source_has_parent_role {
2149                                source_relay
2150                                    .reopen_subagent_mutations(&id)
2151                                    .await
2152                                    .context("could not reopen sub-agent requests after stopping children failed")?;
2153                            }
2154                            return Err(error);
2155                        }
2156                    }
2157                } else {
2158                    // Non-in-place Moves retain their existing child-stop
2159                    // behavior before the source checkpoint begins.
2160                    executor.before_move_source_stop().await?;
2161                }
2162                if executor.cancellation_requested() {
2163                    if operation.in_place && source_has_parent_role {
2164                        source_relay.reopen_subagent_mutations(&id).await?;
2165                    }
2166                    bail!("Move cancelled before source interruption");
2167                }
2168                executor.notify_notice("Stopping source");
2169                let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
2170                if !operation.in_place {
2171                    operation.phase = MovePhase::ClosingSource;
2172                    operation.updated_at = now();
2173                    crate::database::save_move_operation(operation)?;
2174                }
2175                // An eligible move keeps its environment: the close stops at
2176                // the sealed relay and the record keeps its target for
2177                // `restore_session_in_place`.
2178                let disposition = if operation.in_place {
2179                    SourceTargetDisposition::RetainForInPlaceSwap
2180                } else {
2181                    SourceTargetDisposition::Destroy
2182                };
2183                let suspended = Box::pin(self.suspend_session_for_move(
2184                    &id,
2185                    executor,
2186                    manager,
2187                    operation,
2188                    preparation.as_ref(),
2189                    disposition,
2190                    std::mem::take(&mut source_relay),
2191                ))
2192                .await;
2193                if let Err(error) = suspended {
2194                    if operation.in_place
2195                        && source_has_parent_role
2196                        && matches!(
2197                            self.state.sessions[&id].state,
2198                            SessionState::Running | SessionState::Disconnected
2199                        )
2200                        && let Err(reopen) = MoveSourceRelay::set_subagent_admission_via_manager(
2201                            manager,
2202                            &id,
2203                            true,
2204                        )
2205                        .await
2206                    {
2207                        return Err(error.context(format!(
2208                            "could not reopen sub-agent requests after source checkpoint failed: {reopen:#}"
2209                        )));
2210                    }
2211                    return Err(error);
2212                }
2213            }
2214            // Past the source stop, nothing else needs the source's connection.
2215            drop(source_relay);
2216            // An in-place move never reaches `Stopped` with a target: its
2217            // source stays `Closing` in the environment the destination reuses.
2218            if !operation.in_place
2219                && operation.workspace_transfer.is_none()
2220                && self.state.sessions[&id].state == SessionState::Stopped
2221                && self.state.sessions[&id].target.is_some()
2222            {
2223                executor.notify_notice("Cleaning up source");
2224                let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
2225                self.cleanup_stopped_target(&id, executor)?;
2226            }
2227            operation.checkpoint = operation
2228                .checkpoint
2229                .clone()
2230                .or_else(|| self.state.sessions[&id].checkpoint.clone());
2231            ensure!(
2232                operation.restore_artifact().is_some(),
2233                "move has no verified checkpoint"
2234            );
2235            ensure!(
2236                !executor.cancellation_requested(),
2237                "move cancelled after source sealing; checkpoint and remaining environment retained"
2238            );
2239            operation.phase = MovePhase::ResumingDestination;
2240            operation.recovery_session = Some(self.state.sessions[&id].clone());
2241            operation.updated_at = now();
2242            crate::database::save_move_operation(operation)?;
2243            executor.reserve_move_destination();
2244            executor.notify_notice("Preparing destination");
2245            // The destination worker reads its delegation policy from the
2246            // record at launch. `recovery_session` above keeps the old one
2247            // for a rollback.
2248            if let Some(policy) = &operation.selection.subagents {
2249                let session = self.state.sessions.get_mut(&id).unwrap();
2250                session.subagents = Some(policy.clone());
2251                crate::database::save_resumed_session(session, None)?;
2252            }
2253            if operation.in_place {
2254                // The environment is kept, so there is nothing to provision and
2255                // no allocation to clear: eligibility already required the
2256                // allocation to be unchanged.
2257                Box::pin(self.restore_session_in_place(
2258                    &id,
2259                    operation.selection.profile_id.as_deref().unwrap(),
2260                    operation.selection.target_template_id.as_deref().unwrap(),
2261                    executor,
2262                ))
2263                .await?;
2264            } else {
2265                if operation.selection.clear_resource_allocation {
2266                    let session = self.state.sessions.get_mut(&id).unwrap();
2267                    session.resource_allocation = None;
2268                    session.container_cpus = None;
2269                    session.container_memory = None;
2270                    crate::database::save_resumed_session(session, None)?;
2271                }
2272                if operation.workspace_transfer.is_some() {
2273                    self.capture_move_workspace(operation, executor)?;
2274                    Box::pin(self.resume_session_for_move(operation, executor)).await?;
2275                    let saved = crate::database::load_move_operation(&id)?
2276                        .context("Move intent disappeared during restore")?;
2277                    operation.workspace_transfer = saved.workspace_transfer;
2278                    operation.prepared_destination = saved.prepared_destination;
2279                } else {
2280                    Box::pin(self.resume_session_controlled(
2281                        &id,
2282                        operation.selection.profile_id.as_deref().unwrap(),
2283                        operation.selection.target_template_id.as_deref().unwrap(),
2284                        SessionResumeOptions {
2285                            additional_mounts: operation.selection.additional_mounts.clone(),
2286                            resource_allocation: operation.selection.resource_allocation.clone(),
2287                            discard_queue: true,
2288                        },
2289                        executor,
2290                    ))
2291                    .await?;
2292                }
2293            }
2294            let destination = &self.state.sessions[&id];
2295            operation.destination_target = destination.target.clone();
2296            operation.destination_native_session_id = destination.native_session_id.clone();
2297            // Persist readiness before submitting even the first queued command.
2298            operation.phase = MovePhase::StartingQueue;
2299            operation.queue_admission_started = true;
2300            operation.updated_at = now();
2301            crate::database::save_move_operation(operation)?;
2302        }
2303        restore_move_queue_hold(operation);
2304        self.finish_workspace_transfer(operation, executor)?;
2305        self.admit_move_queue(operation, executor).await
2306
2307        }).await
2308    }
2309
2310    async fn admit_move_queue(
2311        &self,
2312        operation: &mut MoveOperation,
2313        executor: &(impl CommandExecutor + Sync),
2314    ) -> Result<()> {
2315        let timing_id = operation.selection.session_id.clone();
2316        let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
2317        let id = &operation.selection.session_id;
2318        let mut relay = {
2319            let _checking_destination =
2320                ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2321            let destination = &self.state.sessions[id];
2322            ensure!(
2323                destination.state == SessionState::Running
2324                    && destination.target == operation.destination_target
2325                    && destination.native_session_id == operation.destination_native_session_id,
2326                "cannot prove the same ready destination; refusing to replay potentially executed work"
2327            );
2328            let spec = self.reconnect_command(id)?;
2329            let relay = StandaloneSession::connect_command(&spec, id).await?;
2330            let store_id = relay.snapshot().operational.store_id.context("destination worker does not expose its durable store identity; upgrade the worker before admitting queued work")?;
2331            if let Some(expected) = &operation.destination_store_id {
2332                ensure!(
2333                    *expected == store_id,
2334                    "destination relay storage was replaced; refusing to replay potentially executed work"
2335                );
2336            } else {
2337                // No command can be admitted until this identity is durable.
2338                operation.destination_store_id = Some(store_id);
2339                crate::database::save_move_operation(operation)?;
2340            }
2341            ensure!(
2342                relay.snapshot().operational.native_session_id
2343                    == operation.destination_native_session_id,
2344                "destination relay native identity changed; refusing queue replay"
2345            );
2346            relay
2347        };
2348        if operation.queue == ResumeQueueDisposition::Start {
2349            let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
2350            executor.notify_notice("Starting queued work");
2351            let checkpoint = operation
2352                .restore_artifact()
2353                .context("move queue archive is missing")?;
2354            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2355            ensure!(
2356                verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
2357                "move queue checkpoint verification failed"
2358            );
2359            for queued in verified.canonical_session.queued_prompts {
2360                if operation.queue_admission_finished {
2361                    continue;
2362                }
2363                ensure!(
2364                    !executor.cancellation_requested(),
2365                    "move cancelled during queue admission; destination retained"
2366                );
2367                let command = match queued.kind {
2368                    CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
2369                        prompt: queued
2370                            .content
2371                            .into_iter()
2372                            .map(serde_json::from_value)
2373                            .collect::<serde_json::Result<_>>()?,
2374                    },
2375                    CanonicalQueuedCommandKind::SetConfig { key, value } => {
2376                        RelayCommand::SetConfig { key, value }
2377                    }
2378                };
2379                relay.submit_accepted(queued.command_id, command).await?;
2380            }
2381        }
2382        operation.queue_admission_finished = true;
2383        crate::database::save_move_operation(operation)?;
2384        restore_move_queue_hold(operation);
2385        let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
2386            "Queued work was discarded; ready and idle."
2387        } else {
2388            "Queued work was accepted."
2389        };
2390        let source_profile = &operation.source_profile_id;
2391        let source_target = &operation.source_target_template_id;
2392        let destination_profile = operation.selection.profile_id.as_deref().unwrap();
2393        let destination_target = operation.selection.target_template_id.as_deref().unwrap();
2394        // The two moves are different events for the person reading the
2395        // conversation: one rebuilt the environment, the other kept it.
2396        let text = if operation.in_place {
2397            format!(
2398                "Switched from {source_profile} / {source_target} to {destination_profile} / {destination_target} in place; the workspace and environment were kept. {queue_sentence} The interrupted prompt was not replayed."
2399            )
2400        } else {
2401            format!(
2402                "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
2403            )
2404        };
2405        relay
2406            .submit(
2407                format!("{}-notice", operation.operation_id),
2408                RelayCommand::RecordNotice { text },
2409            )
2410            .await?;
2411        Ok(())
2412    }
2413
2414    pub(super) fn validate_move_checkpoint(
2415        &self,
2416        operation: &MoveOperation,
2417        preparation: Option<&MovePreparation>,
2418        executor: &(impl CommandExecutor + Sync),
2419    ) -> Result<()> {
2420        let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
2421        let current = Controller {
2422            config: mj_core::config::Config::load()?,
2423            state: self.state.clone(),
2424        };
2425        ensure!(
2426            current.move_configuration_fingerprint(&operation.selection)?
2427                == operation.configuration_fingerprint,
2428            "destination configuration changed during move"
2429        );
2430        let id = &operation.selection.session_id;
2431        current.validate_move_destination_paths(
2432            &self.state.sessions[id],
2433            operation.selection.target_template_id.as_deref().unwrap(),
2434            executor,
2435        )?;
2436        if !operation.in_place
2437            && operation.workspace_transfer.is_none()
2438            && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
2439                .preflight_resume_repository_sources(
2440                    id,
2441                    operation.selection.target_template_id.as_deref().unwrap(),
2442                    executor,
2443                )?
2444        {
2445            bail!(
2446                "destination repository source is missing checkpoint commit {}; source retained",
2447                mismatch.missing_commit
2448            );
2449        }
2450        if let Some(prepared) = preparation {
2451            let checkpoint = operation
2452                .restore_artifact()
2453                .or(self.state.sessions[id].checkpoint.as_ref())
2454                .context("no move checkpoint")?;
2455            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
2456            let actual: Vec<_> = verified
2457                .canonical_session
2458                .queued_prompts
2459                .iter()
2460                .map(|p| p.command_id.as_str())
2461                .collect();
2462            let expected: Vec<_> = prepared
2463                .queued_commands
2464                .iter()
2465                .map(|p| p.command_id.as_str())
2466                .collect();
2467            ensure!(
2468                actual == expected,
2469                "pending queue changed before checkpoint capture; source retained, confirm Move again"
2470            );
2471        }
2472        Ok(())
2473    }
2474}
2475
2476fn validate_preserved_configuration(
2477    profile_id: &str,
2478    accepted: &mj_core::acp::AcceptedSessionConfig,
2479    choices: &mj_core::worker_launch::ProfileConfig,
2480) -> Result<()> {
2481    for (key, value, offered) in [
2482        ("model", accepted.model.as_deref(), &choices.models),
2483        ("effort", accepted.effort.as_deref(), &choices.efforts),
2484    ] {
2485        let Some(value) = value else { continue };
2486        ensure!(
2487            offered.iter().any(|choice| choice.value == value),
2488            "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
2489            offered
2490                .iter()
2491                .map(|choice| choice.value.as_str())
2492                .collect::<Vec<_>>()
2493                .join(", ")
2494        );
2495    }
2496    Ok(())
2497}
2498
2499fn move_subagent_policy(
2500    source: &mj_core::state::SessionRecord,
2501    selection: &MoveSelection,
2502) -> mj_core::subagent::SubagentPolicy {
2503    selection
2504        .subagents
2505        .clone()
2506        .or_else(|| source.subagents.clone())
2507        .unwrap_or_default()
2508}
2509
2510pub(crate) fn parent_tools_enabled(
2511    policy: &mj_core::subagent::SubagentPolicy,
2512    harness: mj_core::config::HarnessKind,
2513) -> bool {
2514    policy.for_launch(harness, false).parent_role().is_some()
2515}
2516
2517pub(in crate::controller) fn move_children(
2518    state: &mj_core::state::State,
2519    parent_session_id: &str,
2520) -> Vec<mj_core::subagent::InPlaceSubagent> {
2521    state
2522        .subagents
2523        .values()
2524        .filter(|relation| relation.parent_session_id == parent_session_id)
2525        .filter_map(|relation| {
2526            let child = state.sessions.get(&relation.child_session_id)?;
2527            let state = if child.state.has_live_worker() {
2528                mj_core::subagent::InPlaceSubagentState::Running
2529            } else if child.state == SessionState::Parked {
2530                mj_core::subagent::InPlaceSubagentState::Parked
2531            } else {
2532                return None;
2533            };
2534            Some(mj_core::subagent::InPlaceSubagent {
2535                child_session_id: relation.child_session_id.clone(),
2536                task_name: relation.task_name.clone(),
2537                state,
2538            })
2539        })
2540        .collect()
2541}
2542
2543fn live_move_children(state: &mj_core::state::State, parent_session_id: &str) -> Vec<String> {
2544    move_children(state, parent_session_id)
2545        .into_iter()
2546        .filter(|child| child.state == mj_core::subagent::InPlaceSubagentState::Running)
2547        .map(|child| {
2548            format!(
2549                "{} (child_session_id {})",
2550                child.task_name, child.child_session_id
2551            )
2552        })
2553        .collect()
2554}
2555
2556fn roleless_move_children_error(children: &[String]) -> Option<String> {
2557    (!children.is_empty()).then(|| {
2558        format!(
2559            "cannot in-place Move to a harness without mj-agents parent tools while live sub-agents exist: {}",
2560            children.join(", ")
2561        )
2562    })
2563}
2564
2565fn should_stop_move_subagents(in_place: bool, destination_has_parent_role: bool) -> bool {
2566    !in_place || !destination_has_parent_role
2567}
2568
2569fn subagent_mutations_pending(
2570    requests: &[mj_core::subagent::SubagentToolRequest],
2571    durable_effect_pending: bool,
2572) -> bool {
2573    durable_effect_pending
2574        || requests
2575            .iter()
2576            .any(|request| request.action.mutates_child_state())
2577}
2578
2579/// Whether this move can replace only the harness inside the source target.
2580///
2581/// The environment is kept only when nothing about it changes: the same target
2582/// template (or another one that names the same bare machine), the same
2583/// attached mounts, and the same resource allocation. A retry takes its retention decision from the durable Move intent instead;
2584/// a sub-agent never owns its own target.
2585pub(super) fn in_place_move_eligible(
2586    source: &mj_core::state::SessionRecord,
2587    selection: &MoveSelection,
2588    targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2589    is_subagent: bool,
2590    retry: bool,
2591) -> bool {
2592    !retry
2593        && !is_subagent
2594        && source.target.is_some()
2595        && matches!(
2596            source.state,
2597            SessionState::Running | SessionState::Disconnected
2598        )
2599        && move_environment_change(
2600            source,
2601            selection,
2602            bare_targets_share_environment(
2603                targets,
2604                &source.target_template_id,
2605                selection.target_template_id.as_deref().unwrap_or_default(),
2606            ),
2607        )
2608        .is_none()
2609}
2610
2611/// Whether two target templates are bare targets on the same machine, so the
2612/// session's worker root and workspace are already where the destination wants
2613/// them. Two local bare targets qualify, and so do two SSH bare targets with
2614/// the same connection; only the template the session names changes.
2615fn bare_targets_share_environment(
2616    targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
2617    source_id: &str,
2618    destination_id: &str,
2619) -> bool {
2620    use mj_core::config::TargetTemplate::{LocalBare, SshBare};
2621    match (targets.get(source_id), targets.get(destination_id)) {
2622        (Some(LocalBare), Some(LocalBare)) => true,
2623        (
2624            Some(SshBare { ssh: source, .. }),
2625            Some(SshBare {
2626                ssh: destination, ..
2627            }),
2628        ) => {
2629            crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
2630        }
2631        _ => false,
2632    }
2633}
2634fn outcome(
2635    operation_id: &str,
2636    selection: &MoveSelection,
2637    status: &str,
2638    error: Option<String>,
2639    recovery: Option<String>,
2640) -> MoveOutcome {
2641    MoveOutcome {
2642        operation_id: operation_id.into(),
2643        session_id: selection.session_id.clone(),
2644        profile_id: selection.profile_id.clone().unwrap_or_default(),
2645        target_template_id: selection.target_template_id.clone().unwrap_or_default(),
2646        outcome: status.into(),
2647        error,
2648        recovery,
2649    }
2650}
2651
2652/// One decision shared by no-op detection and environment retention. A
2653/// different target template counts as a change unless `same_environment` says
2654/// the two name the same bare machine.
2655fn move_environment_change(
2656    source: &mj_core::state::SessionRecord,
2657    selection: &MoveSelection,
2658    same_environment: bool,
2659) -> Option<&'static str> {
2660    if Some(&source.target_template_id) != selection.target_template_id.as_ref()
2661        && !same_environment
2662    {
2663        Some("target changed")
2664    } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
2665        Some("attached mounts changed")
2666    } else if source.resource_allocation != selection.resource_allocation
2667        || (selection.clear_resource_allocation
2668            && (source.container_cpus.is_some() || source.container_memory.is_some()))
2669    {
2670        Some("resource allocation changed")
2671    } else {
2672        None
2673    }
2674}