Skip to main content

mj_controller/controller/
move_session.rs

1//! One recoverable stop/restore operation, independent of its initiating viewer.
2
3mod handoff;
4#[cfg(test)]
5mod tests;
6mod transfer;
7
8use anyhow::{Context, Result, bail, ensure};
9use mj_core::hex::lower_hex;
10use sha2::{Digest, Sha256};
11
12use super::lifecycle::SourceTargetDisposition;
13use super::{Controller, SessionResumeOptions, now};
14
15/// Execution and retained queue admission are phases of one ownership record.
16#[derive(Clone, Copy)]
17enum MoveOwnership {
18    Executing,
19    ExecutingQueue,
20    PendingQueue,
21}
22
23fn move_ownership() -> &'static std::sync::Mutex<std::collections::BTreeMap<String, MoveOwnership>>
24{
25    static OWNER: std::sync::OnceLock<
26        std::sync::Mutex<std::collections::BTreeMap<String, MoveOwnership>>,
27    > = std::sync::OnceLock::new();
28    OWNER.get_or_init(Default::default)
29}
30
31pub fn move_owns_session(session_id: &str) -> bool {
32    move_ownership()
33        .lock()
34        .unwrap_or_else(std::sync::PoisonError::into_inner)
35        .contains_key(session_id)
36}
37
38pub fn move_has_pending_queue(session_id: &str) -> bool {
39    matches!(
40        move_ownership()
41            .lock()
42            .unwrap_or_else(std::sync::PoisonError::into_inner)
43            .get(session_id),
44        Some(MoveOwnership::ExecutingQueue | MoveOwnership::PendingQueue)
45    )
46}
47
48pub fn release_move_queue_hold(session_id: &str) {
49    set_move_queue_hold(session_id, false);
50}
51
52fn set_move_queue_hold(session_id: &str, pending: bool) {
53    let mut owner = move_ownership()
54        .lock()
55        .unwrap_or_else(std::sync::PoisonError::into_inner);
56    let next = match (owner.get(session_id), pending) {
57        (Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue), true) => {
58            Some(MoveOwnership::ExecutingQueue)
59        }
60        (Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue), false) => {
61            Some(MoveOwnership::Executing)
62        }
63        (_, true) => Some(MoveOwnership::PendingQueue),
64        (_, false) => None,
65    };
66    if let Some(next) = next {
67        owner.insert(session_id.to_owned(), next);
68    } else {
69        owner.remove(session_id);
70    }
71}
72
73pub fn restore_move_queue_hold(operation: &MoveOperation) {
74    set_move_queue_hold(
75        &operation.selection.session_id,
76        operation.queue_admission_started && !operation.queue_admission_finished,
77    );
78}
79
80/// What a recovered, interrupted source stop tells the person.
81///
82/// Recovery keeps the environment promised by an in-place move.
83fn interrupted_source_stop_message(recovered: &str, in_place: bool) -> String {
84    if in_place {
85        format!("{recovered}; the in-place swap was interrupted; the environment was retained")
86    } else {
87        recovered.to_owned()
88    }
89}
90
91/// Whether the persisted record shows a source that has finished stopping with
92/// a verified checkpoint to restore.
93///
94/// Only `Stopped` proves the source stop finished. Resume also admits recovery
95/// states, but those do not establish that the source is already stopped.
96fn source_stopped_with_verified_checkpoint(record: &mj_core::state::SessionRecord) -> bool {
97    record.state == SessionState::Stopped && record.checkpoint.is_some()
98}
99
100/// The recovery guidance for a failed Move whose source is already stopped with
101/// a verified checkpoint.
102///
103/// Resume restores the retained checkpoint on the destination profile and
104/// target the Move selected.
105fn stopped_source_recovery(
106    session_id: &str,
107    destination_profile: Option<&str>,
108    destination_target: Option<&str>,
109) -> String {
110    let flag = |name: &str, value: Option<&str>| {
111        value
112            .filter(|value| !value.is_empty())
113            .map(|value| format!(" --{name} {value}"))
114            .unwrap_or_default()
115    };
116    format!(
117        "Source is stopped with a verified checkpoint. Bring it back with \
118         `mj resume --session {session_id}{}{} --queue start`.",
119        flag("profile", destination_profile),
120        flag("target", destination_target),
121    )
122}
123
124/// The recovery guidance a failed or cancelled Move shows, from its persisted
125/// state.
126fn failed_move_recovery(
127    operation: &MoveOperation,
128    record: Option<&mj_core::state::SessionRecord>,
129) -> String {
130    if operation.queue_admission_started {
131        return "Destination is live; retry queue admission on this same destination. Already accepted work may have effects.".to_owned();
132    }
133    if operation.in_place && record.is_some_and(|record| record.target.is_some()) {
134        return "Environment and checkpoint retained. Retry Move on the same target; the checkout will not be recreated.".to_owned();
135    }
136    match record {
137        Some(record) if source_stopped_with_verified_checkpoint(record) => {
138            stopped_source_recovery(
139                &operation.selection.session_id,
140                operation.selection.profile_id.as_deref(),
141                operation.selection.target_template_id.as_deref(),
142            )
143        }
144        _ => "Source or partial destination is retained. Retry move after resolving the reported error.".to_owned(),
145    }
146}
147
148pub struct MoveMutationGuard(String);
149
150impl MoveMutationGuard {
151    pub fn reserve(session_id: &str) -> Result<Self> {
152        let mut owner = move_ownership()
153            .lock()
154            .unwrap_or_else(std::sync::PoisonError::into_inner);
155        ensure!(
156            !matches!(
157                owner.get(session_id),
158                Some(MoveOwnership::Executing | MoveOwnership::ExecutingQueue)
159            ),
160            "session already has a move owner"
161        );
162        let next = if owner.contains_key(session_id) {
163            MoveOwnership::ExecutingQueue
164        } else {
165            MoveOwnership::Executing
166        };
167        owner.insert(session_id.to_owned(), next);
168        Ok(Self(session_id.to_owned()))
169    }
170}
171
172impl Drop for MoveMutationGuard {
173    fn drop(&mut self) {
174        let mut owner = move_ownership()
175            .lock()
176            .unwrap_or_else(std::sync::PoisonError::into_inner);
177        match owner.get(&self.0) {
178            Some(MoveOwnership::ExecutingQueue) => {
179                owner.insert(self.0.clone(), MoveOwnership::PendingQueue);
180            }
181            Some(MoveOwnership::Executing) => {
182                owner.remove(&self.0);
183            }
184            _ => {}
185        }
186    }
187}
188
189pub(crate) fn move_refuses_command(session_id: &str, command: &RelayCommand) -> bool {
190    move_owns_session(session_id)
191        && matches!(
192            command,
193            RelayCommand::Prompt { .. }
194                | RelayCommand::SetConfig { .. }
195                | RelayCommand::SetSessionMode { .. }
196                | RelayCommand::RestoreExecutionMode
197                | RelayCommand::RunUserShell { .. }
198                | RelayCommand::CancelUserShell { .. }
199                | RelayCommand::Cancel
200                | RelayCommand::RemoveQueuedPrompt { .. }
201                | RelayCommand::ClearQueuedPrompts
202        )
203}
204use crate::session_manager::{SessionManagerControl, StandaloneSession, new_command_id};
205use mj_checkpoint::archive::{CanonicalQueuedCommandKind, verify_archive_streaming};
206use mj_core::state::{MoveOperation, MovePhase, ResumeQueueDisposition, SessionState};
207
208pub use mj_core::state::{MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest};
209
210use crate::targets::{CommandExecutor, ProvisionStage, ProvisionStageGuard};
211use mj_core::relay::RelayCommand;
212
213/// Refresh source state without turning a dead source harness into a Move prerequisite.
214pub async fn refresh_move_source(
215    manager: &SessionManagerControl,
216    id: &str,
217) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
218    Ok(MoveSourceRelay::lease(manager, id).await?.snapshot())
219}
220
221/// A Move's hold on its source's relay connection, from the confirmation
222/// until the checkpoint that seals the source takes it over. While it is
223/// held, the session actor admits nothing to the relay, so the state the Move
224/// was confirmed against is the state it interrupts. Dropping it hands the
225/// connection back to the actor.
226#[derive(Default)]
227pub(in crate::controller) struct MoveSourceRelay(Option<super::checkpoint::ControllerRelayLease>);
228
229impl MoveSourceRelay {
230    /// Take the source's connection from its session actor, which syncs it
231    /// before handing it over. Holds nothing when the source's worker cannot
232    /// be reached; the Move then recovers it without its harness. Leasing
233    /// preserves typed transport failures, unlike the UI's string-valued sync.
234    pub(in crate::controller) async fn lease(
235        manager: &SessionManagerControl,
236        id: &str,
237    ) -> Result<Self> {
238        let handle = manager
239            .wait_for_session(id, std::time::Duration::from_secs(5))
240            .await?;
241        match handle.lease_connection().await {
242            Ok(lease) => Ok(Self(Some(
243                super::checkpoint::ControllerRelayLease::Managed {
244                    handle,
245                    lease: Some(lease),
246                },
247            ))),
248            Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
249                tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
250                Ok(Self(None))
251            }
252            Err(error) => Err(error),
253        }
254    }
255
256    /// The source's state as of its last sync, or `None` when it is unreachable.
257    pub(in crate::controller) fn snapshot(
258        &mut self,
259    ) -> Option<mj_core::state::ManagedSessionSnapshot> {
260        self.0
261            .as_mut()
262            .map(|relay| relay.connection_mut().snapshot())
263    }
264
265    /// Sync the held connection again. A source whose transport died since it
266    /// was leased is let go and reported unreachable.
267    pub(in crate::controller) async fn sync(
268        &mut self,
269        id: &str,
270    ) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
271        let Some(relay) = self.0.as_mut() else {
272            return Ok(None);
273        };
274        match relay.connection_mut().sync().await {
275            Ok(snapshot) => Ok(Some(snapshot)),
276            Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
277                tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
278                self.0 = None;
279                Ok(None)
280            }
281            Err(error) => Err(error),
282        }
283    }
284
285    pub(in crate::controller) fn is_held(&self) -> bool {
286        self.0.is_some()
287    }
288
289    /// Hand the connection to the checkpoint that seals the source.
290    pub(in crate::controller) fn take(
291        &mut self,
292    ) -> Option<super::checkpoint::ControllerRelayLease> {
293        self.0.take()
294    }
295
296    /// Keep holding the source across a worker restart, on the new worker's
297    /// connection.
298    pub(in crate::controller) fn replace_connection(
299        &mut self,
300        connection: crate::session_manager::StandaloneSession,
301    ) {
302        if let Some(relay) = self.0.as_mut() {
303            relay.replace_connection(connection);
304        }
305    }
306}
307
308impl Drop for MoveSourceRelay {
309    fn drop(&mut self) {
310        if let Some(relay) = self.0.take() {
311            relay.release();
312        }
313    }
314}
315
316fn digest(value: &impl serde::Serialize) -> Result<String> {
317    Ok(lower_hex(Sha256::digest(serde_json::to_vec(value)?)))
318}
319
320struct MovePhaseTimer<'a> {
321    session_id: &'a str,
322    phase: &'static str,
323    started: std::time::Instant,
324}
325
326impl<'a> MovePhaseTimer<'a> {
327    fn new(session_id: &'a str, phase: &'static str) -> Self {
328        Self {
329            session_id,
330            phase,
331            started: std::time::Instant::now(),
332        }
333    }
334}
335
336impl Drop for MovePhaseTimer<'_> {
337    fn drop(&mut self) {
338        tracing::info!(
339            session_id = self.session_id,
340            phase = self.phase,
341            elapsed_ms = self.started.elapsed().as_millis() as u64,
342            "move phase finished"
343        );
344    }
345}
346
347/// Confirmation is an inspector, not transport for attachment bytes: every
348/// surface that shows a move preparation sees a placeholder for each queued
349/// image. Replay always reads the verified archive after destination readiness.
350fn replace_queued_images_with_placeholders(
351    queued_commands: &mut [mj_core::state::MaterializedQueuedPrompt],
352) {
353    for block in queued_commands
354        .iter_mut()
355        .flat_map(|command| command.content.iter_mut())
356    {
357        if block.get("type").and_then(serde_json::Value::as_str) != Some("image") {
358            continue;
359        }
360        let mime = block
361            .get("mimeType")
362            .or_else(|| block.get("mime_type"))
363            .and_then(serde_json::Value::as_str)
364            .unwrap_or("image");
365        *block = serde_json::json!({"type": "text", "text": format!("[Image attachment: {mime}]")});
366    }
367}
368
369impl Controller {
370    /// Returns the planned conversion when this move turns a local checkout
371    /// into an isolated workspace, so the caller can describe it without
372    /// reading Git a second time.
373    fn validate_move_destination_paths(
374        &self,
375        source: &mj_core::state::SessionRecord,
376        target_id: &str,
377        executor: &(impl CommandExecutor + Sync),
378    ) -> Result<Option<super::worktree::RawToWorkspaceConversion>> {
379        use super::worktree::ResumePlan;
380        match super::worktree::resume_compatibility(source, &self.config, target_id)
381            .map_err(anyhow::Error::msg)?
382        {
383            ResumePlan::RawToWorkspace => {
384                return Ok(Some(super::worktree::plan_raw_to_workspace(
385                    source,
386                    &self.config,
387                    executor,
388                )?));
389            }
390            ResumePlan::WorkspaceToRaw => {
391                self.plan_workspace_to_raw(source, target_id, executor)?;
392            }
393            ResumePlan::InPlace if source.managed_worktree.is_none() => {
394                if let Some(path) = &source.project_directory {
395                    self.validate_project_directory(target_id, path, executor)?;
396                }
397            }
398            ResumePlan::InPlace => {}
399        }
400        Ok(None)
401    }
402    fn move_confirmation(
403        &self,
404        selection: &MoveSelection,
405        conversion: Option<&mj_core::state::RawConversionPreview>,
406    ) -> Result<(bool, Vec<mj_core::state::MaterializedQueuedPrompt>, String)> {
407        let source = self
408            .state
409            .sessions
410            .get(&selection.session_id)
411            .context("unknown move session")?;
412        let (mut active, mut queued) = crate::database::move_pending_work(&source.id)?;
413        if let Some(operation) = crate::database::load_move_operation(&source.id)?
414            && operation.queue_admission_started
415            && !operation.queue_admission_finished
416        {
417            ensure!(
418                operation.selection == *selection,
419                "queue admission is incomplete on the live destination; retry that move before selecting another destination"
420            );
421            let checkpoint = operation
422                .restore_artifact()
423                .context("retained queue checkpoint is missing")?;
424            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
425            ensure!(
426                verified.archive_sha256 == checkpoint.sha256
427                    && verified.manifest.session.id == source.id,
428                "retained move checkpoint verification failed"
429            );
430            queued = verified
431                .canonical_session
432                .queued_prompts
433                .into_iter()
434                .map(|entry| mj_core::state::MaterializedQueuedPrompt {
435                    accepted_ordinal: None,
436                    command_id: entry.command_id,
437                    kind: match entry.kind {
438                        CanonicalQueuedCommandKind::Prompt => {
439                            mj_core::state::QueuedCommandKind::Prompt
440                        }
441                        CanonicalQueuedCommandKind::SetConfig { key, value } => {
442                            mj_core::state::QueuedCommandKind::SetConfig { key, value }
443                        }
444                    },
445                    content: entry.content,
446                    queued_at_ms: entry.queued_at_ms,
447                })
448                .collect();
449            active = false;
450        }
451        let fingerprint = digest(&(
452            &source.last_profile,
453            &source.target_template_id,
454            &source.target,
455            &source.target_runtime,
456            &source.native_session_id,
457            &source.resource_allocation,
458            &source.additional_mounts,
459            &source.container_cpus,
460            &source.container_memory,
461            selection,
462            self.move_configuration_fingerprint(selection)?,
463            &queued,
464            // Only what the destination is built from. The dirty counts move
465            // with every keystroke of a live agent, and hashing them would
466            // invalidate the confirmation the person is reading.
467            conversion.map(|preview| {
468                (
469                    &preview.fetch_url,
470                    &preview.push_urls,
471                    &preview.branch,
472                    &preview.destination,
473                )
474            }),
475        ))?;
476        Ok((active, queued, fingerprint))
477    }
478    pub(super) fn move_configuration_fingerprint(
479        &self,
480        selection: &MoveSelection,
481    ) -> Result<String> {
482        let profile = self
483            .config
484            .profiles
485            .get(
486                selection
487                    .profile_id
488                    .as_deref()
489                    .context("move profile is unresolved")?,
490            )
491            .context("move profile no longer exists")?;
492        let target = self
493            .config
494            .targets
495            .get(
496                selection
497                    .target_template_id
498                    .as_deref()
499                    .context("move target is unresolved")?,
500            )
501            .context("move target no longer exists")?;
502        // Only a digest crosses IPC or enters the move record. Configuration
503        // may contain credential-bearing values and must never be copied there.
504        let source = self
505            .state
506            .sessions
507            .get(&selection.session_id)
508            .context("unknown move session")?;
509        digest(&(
510            profile,
511            target,
512            &self.config.bundles,
513            self.config.targets.get(&source.target_template_id),
514            self.config.profiles.get(&source.last_profile),
515        ))
516    }
517
518    async fn validate_move_destination_configuration(
519        &self,
520        selection: &MoveSelection,
521        source_harness: mj_core::config::HarnessKind,
522        operational: &mj_core::relay::RelayOperationalState,
523    ) -> Result<()> {
524        let profile_id = selection
525            .profile_id
526            .as_deref()
527            .context("move profile is unresolved")?;
528        let profile = self
529            .config
530            .profiles
531            .get(profile_id)
532            .context("move profile no longer exists")?;
533        let source = self
534            .state
535            .sessions
536            .get(&selection.session_id)
537            .context("unknown move session")?;
538        if profile.kind != source_harness || profile_id == source.last_profile {
539            return Ok(());
540        }
541        let accepted = mj_core::acp::AcceptedSessionConfig::from_configuration(
542            &operational.config,
543            &operational.config_options,
544        );
545        if accepted.model.is_none() && accepted.effort.is_none() {
546            return Ok(());
547        }
548        // A fresh probe prevents a stale local catalogue from approving a
549        // move that the destination profile will immediately reset.
550        let choices =
551            super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
552                .await
553                .with_context(|| {
554                    format!("discover destination profile {profile_id:?} configuration")
555                })?;
556        validate_preserved_configuration(profile_id, &accepted, &choices)
557    }
558
559    pub async fn prepare_move_session_controlled(
560        &self,
561        mut selection: MoveSelection,
562        executor: &(impl CommandExecutor + Sync),
563    ) -> Result<MovePreparation> {
564        ensure!(
565            selection.profile_id.is_some() || selection.target_template_id.is_some(),
566            "move requires a target or profile selection"
567        );
568        let source = self
569            .state
570            .sessions
571            .get(&selection.session_id)
572            .context("unknown session")?;
573        ensure!(
574            !self.state.subagents.contains_key(&source.id),
575            "sub-agent sessions cannot move independently of their parent"
576        );
577        let previous = crate::database::load_move_operation(&source.id)?;
578        let retry = previous.as_ref().is_some_and(|op| {
579            !matches!(op.phase, MovePhase::Completed)
580                && (op.restore_artifact().is_some() || source.checkpoint.is_some())
581        });
582        ensure!(
583            matches!(
584                source.state,
585                SessionState::Running | SessionState::Disconnected
586            ) || retry,
587            "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
588            source.id
589        );
590        selection
591            .profile_id
592            .get_or_insert_with(|| source.last_profile.clone());
593        selection
594            .target_template_id
595            .get_or_insert_with(|| source.target_template_id.clone());
596        selection
597            .additional_mounts
598            .get_or_insert_with(|| source.additional_mounts.clone());
599        ensure!(
600            !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
601            "select resource allocation or explicitly clear it, not both"
602        );
603        if !selection.clear_resource_allocation {
604            selection.resource_allocation = selection
605                .resource_allocation
606                .or_else(|| source.resource_allocation.clone());
607        }
608        if selection.clear_resource_allocation
609            && source.resource_allocation.is_none()
610            && source.container_cpus.is_none()
611            && source.container_memory.is_none()
612        {
613            selection.clear_resource_allocation = false;
614        }
615        if let Some(retained) = previous.as_ref().filter(|op| {
616            op.retains_source_environment()
617                && op.phase != MovePhase::Completed
618                && op.recovery_session.is_some()
619        }) {
620            ensure!(
621                retained.selection == selection,
622                "a sealed Move must be retried with its prepared destination and file selection"
623            );
624            ensure!(
625                !retained.in_place
626                    || (retained.source_target.is_some()
627                        && source.target == retained.source_target),
628                "retained Move target is missing or changed; refusing to recreate it"
629            );
630        }
631        let profile_id = selection.profile_id.as_deref().unwrap();
632        let target_id = selection.target_template_id.as_deref().unwrap();
633        let profile = self
634            .config
635            .profiles
636            .get(profile_id)
637            .context("unknown destination profile")?;
638        ensure!(
639            profile.enabled,
640            "destination profile {profile_id:?} is disabled"
641        );
642        if let Some(policy) = &selection.subagents {
643            super::profile_config::validate_session_subagent_policy(
644                &self.config,
645                profile_id,
646                policy,
647            )
648            .await?;
649        }
650        let target = self
651            .config
652            .targets
653            .get(target_id)
654            .context("unknown destination target")?;
655        self.validate_muse_resume_destination(source, profile.kind, target_id)?;
656        super::worktree::resume_compatibility(source, &self.config, target_id)
657            .map_err(anyhow::Error::msg)?;
658        super::backend::validate_resource_allocation(
659            target,
660            selection.resource_allocation.as_ref(),
661        )?;
662        let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
663        ensure!(
664            profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
665            "Muse Code ACP supports one workspace root; attached directories are unsupported"
666        );
667        ensure!(
668            mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
669            "attached resources are unsupported for this target; select compatible resources explicitly"
670        );
671        crate::targets::validate_additional_mounts(mounts)?;
672        for mount in mounts {
673            self.validate_mount_source(target_id, &mount.source, executor)?;
674        }
675        let planned_conversion =
676            self.validate_move_destination_paths(source, target_id, executor)?;
677        ensure!(
678            profile.home.is_dir(),
679            "destination profile home is unavailable; configure the profile before moving"
680        );
681        super::worker_binary::preflight_worker_binary(target, executor)?;
682        super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
683        let source_harness = previous
684            .as_ref()
685            .filter(|operation| {
686                !operation.queue_admission_started && operation.phase != MovePhase::Completed
687            })
688            .and_then(|operation| operation.recovery_session.as_ref())
689            .map_or(source.harness_kind, |record| record.harness_kind);
690        let cross_harness = profile.kind != source_harness;
691        if cross_harness {
692            let cancel = tokio_util::sync::CancellationToken::new();
693            let resolve =
694                crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
695            tokio::pin!(resolve);
696            loop {
697                tokio::select! {
698                    result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
699                    _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
700                        if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
701                    }
702                }
703            }
704        }
705        // What moving a local checkout into a target really does, computed
706        // before anything is stopped so a person can confirm it.
707        let conversion = planned_conversion
708            .map(|conversion| {
709                super::worktree::raw_conversion_preview(source, &conversion, executor)
710                    .context("describe the move of this checkout into the target")
711            })
712            .transpose()?
713            .map(Box::new);
714        let (active, mut queued_commands, fingerprint) =
715            self.move_confirmation(&selection, conversion.as_deref())?;
716        replace_queued_images_with_placeholders(&mut queued_commands);
717        let operation_id = previous
718            .as_ref()
719            .filter(|operation| {
720                operation.selection == selection
721                    && operation.phase != MovePhase::Completed
722                    && operation.restore_artifact().is_some()
723            })
724            .map(|operation| operation.operation_id.clone())
725            .unwrap_or(new_command_id("move")?);
726        let in_place = previous
727            .as_ref()
728            .is_some_and(|op| retry && op.in_place && op.selection == selection)
729            || in_place_move_eligible(
730                source,
731                &selection,
732                &self.config.targets,
733                self.state.subagents.contains_key(&source.id),
734                retry,
735            );
736        let workspace = if in_place {
737            None
738        } else {
739            Some(self.assess_move_workspace(&selection, executor)?)
740        };
741        Ok(MovePreparation {
742            workspace,
743            source_unavailable: false,
744            in_place,
745            conversion,
746            selection,
747            source_profile_id: source.last_profile.clone(),
748            source_target_template_id: source.target_template_id.clone(),
749            cross_harness,
750            active,
751            queued_commands,
752            fingerprint,
753            operation_id,
754        })
755    }
756
757    /// Called with one daemon lifecycle and recovery reservation already held.
758    pub async fn move_session_managed_controlled(
759        &mut self,
760        request: MoveSessionRequest,
761        executor: &(impl CommandExecutor + Sync),
762        manager: &SessionManagerControl,
763    ) -> Result<MoveOutcome> {
764        let prepared = &request.preparation;
765        let id = prepared.selection.session_id.clone();
766        let started = std::time::Instant::now();
767        executor.notify_notice("Checking destination");
768        let mut source_relay = MoveSourceRelay::default();
769        let checked = {
770            let _checking_destination =
771                ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
772            let mut checked = self
773                .prepare_move_session_controlled(prepared.selection.clone(), executor)
774                .await?;
775            // Destination checks can outlast a turn. Confirm against the relay
776            // immediately before interruption, and hold its connection from
777            // here until the checkpoint seals the source, so nothing reaches
778            // the relay between this confirmation and the interruption.
779            if matches!(
780                self.state.sessions[&id].state,
781                SessionState::Running | SessionState::Disconnected
782            ) {
783                let source_harness = self.state.sessions[&id].harness_kind;
784                source_relay = MoveSourceRelay::lease(manager, &id).await?;
785                let snapshot = source_relay.snapshot();
786                let (active, queue, fingerprint) =
787                    self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
788                checked.source_unavailable = snapshot
789                    .as_ref()
790                    .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
791                checked.active = active
792                    || checked.source_unavailable
793                    || snapshot.as_ref().is_some_and(|snapshot| {
794                        let mut operational = snapshot.operational.clone();
795                        operational.queued_prompts.clear();
796                        operational.checkpoint_barrier = None;
797                        !operational.safe_to_replace(source_harness)
798                    });
799                if let Some(snapshot) = &snapshot {
800                    self.validate_move_destination_configuration(
801                        &checked.selection,
802                        source_harness,
803                        &snapshot.operational,
804                    )
805                    .await?;
806                }
807                checked.queued_commands = queue;
808                checked.fingerprint = fingerprint;
809            }
810            checked
811        };
812        ensure!(
813            checked.fingerprint == prepared.fingerprint,
814            "session, pending work, or destination configuration changed; prepare and confirm Move again"
815        );
816        let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
817        let source = self.state.sessions[&id].clone();
818        if let Some(assessment) = &checked.workspace {
819            checked.selection.workspace.validate(assessment)?;
820        }
821        let old_operation = crate::database::load_move_operation(&id)?;
822        let retry = old_operation.filter(|op| {
823            op.selection == prepared.selection
824                && op.phase != MovePhase::Completed
825                && op.restore_artifact().is_some()
826        });
827        if retry.is_none()
828            && source.state == SessionState::Running
829            && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
830            && move_environment_change(&source, &checked.selection, false).is_none()
831        {
832            return Ok(outcome(
833                &prepared.operation_id,
834                &prepared.selection,
835                "unchanged",
836                None,
837                None,
838            ));
839        }
840        ensure!(
841            !checked.active || request.acknowledge_interruption,
842            "active work will be interrupted; confirm Move again with interruption acknowledgement"
843        );
844        ensure!(
845            checked.queued_commands.is_empty() || request.queue.is_some(),
846            "pending work requires an explicit queue choice: discard or start"
847        );
848        let timestamp = now();
849        let mut operation = match retry {
850            Some(mut op) => {
851                ensure!(
852                    !op.queue_admission_started || op.queue == queue,
853                    "queued work may already have run; retry with the original queue choice on the same destination"
854                );
855                crate::database::clear_move_cancellation_for_retry(&id)?;
856                op.configuration_fingerprint =
857                    self.move_configuration_fingerprint(&checked.selection)?;
858                op.queue = queue;
859                op.cancellation_requested = false;
860                op.error = None;
861                op
862            }
863            None => MoveOperation {
864                workspace_transfer: if checked.in_place {
865                    None
866                } else {
867                    Some(
868                        self.new_workspace_transfer(
869                            &id,
870                            &prepared.operation_id,
871                            checked
872                                .workspace
873                                .clone()
874                                .context("Move workspace assessment missing")?,
875                            executor,
876                        )?,
877                    )
878                },
879                handoff: None,
880                source_checkpoint_only: false,
881                in_place: in_place_move_eligible(
882                    &source,
883                    &checked.selection,
884                    &self.config.targets,
885                    self.state.subagents.contains_key(&id),
886                    false,
887                ),
888                operation_id: prepared.operation_id.clone(),
889                selection: checked.selection.clone(),
890                source_profile_id: source.last_profile.clone(),
891                source_target_template_id: source.target_template_id.clone(),
892                source_target: source.target.clone(),
893                source_native_session_id: source.native_session_id.clone(),
894                source_additional_mounts: source.additional_mounts.clone(),
895                source_resource_allocation: source.resource_allocation.clone(),
896                destination_target: None,
897                destination_native_session_id: None,
898                destination_store_id: None,
899                configuration_fingerprint: self
900                    .move_configuration_fingerprint(&checked.selection)?,
901                checkpoint: None,
902                recovery_session: None,
903                queue,
904                phase: MovePhase::Preparing,
905                queue_admission_started: false,
906                queue_admission_finished: false,
907                cancellation_requested: false,
908                created_at: timestamp.clone(),
909                updated_at: timestamp,
910                error: None,
911            },
912        };
913        crate::database::save_move_operation(&operation)?;
914        tracing::info!(
915            session_id = id,
916            in_place = operation.in_place,
917            reason = move_environment_change(
918                &source,
919                &checked.selection,
920                bare_targets_share_environment(
921                    &self.config.targets,
922                    &source.target_template_id,
923                    checked
924                        .selection
925                        .target_template_id
926                        .as_deref()
927                        .unwrap_or_default(),
928                ),
929            )
930            .unwrap_or(if operation.in_place {
931                "environment unchanged"
932            } else {
933                "source environment unavailable or previously released"
934            }),
935            "move environment decision"
936        );
937        tracing::info!(
938            session_id = id,
939            phase = "preflight",
940            elapsed_ms = started.elapsed().as_millis() as u64,
941            "move phase completed"
942        );
943        // Move nests checkpoint and restore state machines. Heap-own their
944        // futures so dev builds fit the runtime's ordinary thread stacks.
945        let result = Box::pin(self.execute_move(
946            &mut operation,
947            Some(&checked),
948            executor,
949            manager,
950            source_relay,
951        ))
952        .await;
953        self.finish_move_result(&mut operation, result, executor)
954    }
955
956    fn finish_move_result(
957        &self,
958        operation: &mut MoveOperation,
959        result: Result<()>,
960        executor: &impl CommandExecutor,
961    ) -> Result<MoveOutcome> {
962        if result.is_err()
963            && operation.workspace_transfer.is_some()
964            && !crate::upgrade::gate().is_open()
965        {
966            // Handoff cancellation is not user cancellation. Leave the durable
967            // phase active so the next daemon resumes the accepted Move.
968            return Ok(outcome(
969                &operation.operation_id,
970                &operation.selection,
971                "interrupted",
972                None,
973                Some("Move will continue after the daemon upgrade".into()),
974            ));
975        }
976        let (status, error, recovery) = match result {
977            Ok(()) => {
978                operation.phase = MovePhase::Completed;
979                ("completed", None, None)
980            }
981            Err(error) => {
982                let cancelled =
983                    executor.cancellation_requested() || operation.cancellation_requested;
984                operation.phase = if cancelled {
985                    MovePhase::Cancelled
986                } else {
987                    MovePhase::Failed
988                };
989                operation.cancellation_requested = cancelled;
990                let recovery = failed_move_recovery(
991                    operation,
992                    self.state.sessions.get(&operation.selection.session_id),
993                );
994                let error = format!("{error:#}");
995                operation.error = Some(error.clone());
996                (
997                    if cancelled { "cancelled" } else { "failed" },
998                    Some(error),
999                    Some(recovery),
1000                )
1001            }
1002        };
1003        operation.updated_at = now();
1004        crate::database::save_move_operation(operation)?;
1005        Ok(outcome(
1006            &operation.operation_id,
1007            &operation.selection,
1008            status,
1009            error,
1010            recovery,
1011        ))
1012    }
1013
1014    pub async fn recover_move_managed_controlled(
1015        &mut self,
1016        mut operation: MoveOperation,
1017        executor: &(impl CommandExecutor + Sync),
1018        manager: &SessionManagerControl,
1019    ) -> Result<MoveOutcome> {
1020        let id = operation.selection.session_id.clone();
1021        let result = async {
1022            let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
1023            if operation.queue_admission_started {
1024                ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
1025                self.finish_workspace_transfer(&mut operation, executor)?;
1026                return self.admit_move_queue(&mut operation, executor).await;
1027            }
1028            if operation.workspace_transfer.is_some() && operation.recovery_session.is_some() && !operation.cancellation_requested
1029                && !(operation.phase == MovePhase::ResumingDestination && session.state == SessionState::Running) {
1030                if session.state == SessionState::Provisioning || session.target != operation.source_target {
1031                    self.rollback_move_destination(&operation, anyhow::anyhow!("resume interrupted Move transfer"), executor)?;
1032                }
1033                return Box::pin(self.execute_move(&mut operation, None, executor, manager, MoveSourceRelay::default())).await;
1034            }
1035            if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1036                let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
1037                Box::pin(self.recover_move_source_stop(&mut operation, &cleanup, manager)).await?;
1038                operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1039                if operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Closing {
1040                    let previous = self.state.sessions[&id].clone();
1041                    operation.recovery_session = Some(previous.clone());
1042                    let cause = anyhow::anyhow!("Move source sealed; environment retained for explicit retry");
1043                    return Err(self.retain_failed_in_place_move(&id, &previous, cause)?);
1044                }
1045                if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1046                    self.cleanup_stopped_target(&id, &cleanup)?;
1047                }
1048                bail!("{}", interrupted_source_stop_message(
1049                    "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
1050                    operation.in_place,
1051                ));
1052            }
1053            match operation.phase {
1054                MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
1055                MovePhase::ResumingDestination if session.state == SessionState::Running => {
1056                    // Running is installed only after native readiness and the
1057                    // handoff. A crash before the next intent write is safe to
1058                    // advance, since restore started with an empty queue.
1059                    ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
1060                        && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
1061                        "ready destination does not match the move intent");
1062                    operation.destination_target = session.target.clone();
1063                    operation.destination_native_session_id = session.native_session_id.clone();
1064                    operation.queue_admission_started = true;
1065                    operation.phase = MovePhase::StartingQueue;
1066                    crate::database::save_move_operation(&operation)?;
1067                    restore_move_queue_hold(&operation);
1068                    ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
1069                    self.finish_workspace_transfer(&mut operation, executor)?;
1070                    self.admit_move_queue(&mut operation, executor).await
1071                }
1072                MovePhase::ResumingDestination => {
1073                    let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
1074                    // The retained environment belongs to this Move, even
1075                    // when only part of its worker installation completed.
1076                    let cause = anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry");
1077                    let error = if operation.in_place {
1078                        self.retain_failed_in_place_move(&id, previous, cause)?
1079                    } else if operation.workspace_transfer.is_some() {
1080                        self.rollback_move_destination(&operation, cause, executor)?
1081                    } else {
1082                        self.rollback_failed_resume(&id, previous, false, cause, executor)?
1083                    };
1084                    Err(error)
1085                }
1086                MovePhase::ClosingSource => {
1087                    if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
1088                        Box::pin(self.recover_move_source_stop(&mut operation, executor, manager)).await?;
1089                    }
1090                    operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1091                    if !operation.retains_source_environment() && self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
1092                        self.cleanup_stopped_target(&id, executor)?;
1093                    }
1094                    bail!("{}", interrupted_source_stop_message(
1095                        "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
1096                        operation.in_place,
1097                    ))
1098                }
1099                _ => bail!("Move requires an explicit retry after the daemon restarted"),
1100            }
1101        }.await;
1102        self.finish_move_result(&mut operation, result, executor)
1103    }
1104
1105    async fn recover_move_source_stop(
1106        &mut self,
1107        operation: &mut MoveOperation,
1108        executor: &(impl CommandExecutor + Sync),
1109        manager: &SessionManagerControl,
1110    ) -> Result<()> {
1111        let id = operation.selection.session_id.clone();
1112        if self.state.sessions[&id].state == SessionState::Closing {
1113            self.prepare_move_source_checkpoint(
1114                &id,
1115                executor,
1116                manager,
1117                operation,
1118                &mut MoveSourceRelay::default(),
1119            )
1120            .await?;
1121            let handle = manager
1122                .wait_for_session(&id, std::time::Duration::from_secs(5))
1123                .await?;
1124            let mut lease = handle.lease_connection().await?;
1125            let execution = lease.connection_mut().sync().await?.operational.execution;
1126            if operation.retains_source_environment()
1127                && matches!(
1128                    execution,
1129                    mj_core::relay::RelayExecutionState::Closing
1130                        | mj_core::relay::RelayExecutionState::Closed
1131                )
1132            {
1133                super::checkpoint::wait_for_relay_closed(lease.connection_mut()).await?;
1134                lease.release();
1135                ensure!(
1136                    operation.restore_artifact().is_some()
1137                        || self.state.sessions[&id].checkpoint.is_some(),
1138                    "sealed Move source has no checkpoint"
1139                );
1140                return Ok(());
1141            }
1142            lease.release();
1143            if matches!(
1144                execution,
1145                mj_core::relay::RelayExecutionState::Idle
1146                    | mj_core::relay::RelayExecutionState::Running
1147            ) {
1148                if operation.cancellation_requested || operation.retains_source_environment() {
1149                    let record = self.state.sessions.get_mut(&id).unwrap();
1150                    record.state = SessionState::Running;
1151                    record.updated_at = now();
1152                    record.last_error = Some(
1153                        "Move was interrupted before the source was sealed; source retained".into(),
1154                    );
1155                    crate::database::save_lifecycle_session(record)?;
1156                } else {
1157                    // Fresh-environment moves finish their admitted source stop.
1158                    Box::pin(self.suspend_session_for_move(
1159                        &id,
1160                        executor,
1161                        manager,
1162                        operation,
1163                        None,
1164                        SourceTargetDisposition::Destroy,
1165                        MoveSourceRelay::default(),
1166                    ))
1167                    .await?;
1168                }
1169                return Ok(());
1170            }
1171        }
1172        ensure!(
1173            !operation.retains_source_environment(),
1174            "retained Move source cannot be proven; refusing target teardown"
1175        );
1176        // A move's source stop was admitted by the move itself.
1177        self.recover_interrupted_close_managed(&id, executor, manager, true, None)
1178            .await?;
1179        Ok(())
1180    }
1181
1182    /// `source_relay` is the source connection the Move confirmed against,
1183    /// held until the checkpoint that seals the source takes it over.
1184    async fn execute_move(
1185        &mut self,
1186        operation: &mut MoveOperation,
1187        preparation: Option<&MovePreparation>,
1188        executor: &(impl CommandExecutor + Sync),
1189        manager: &SessionManagerControl,
1190        mut source_relay: MoveSourceRelay,
1191    ) -> Result<()> {
1192        let id = operation.selection.session_id.clone();
1193        ensure!(
1194            !executor.cancellation_requested(),
1195            "move cancelled before source interruption"
1196        );
1197        if !operation.queue_admission_started {
1198            if operation.in_place {
1199                ensure!(
1200                    operation.source_target.is_some()
1201                        && self.state.sessions[&id].target == operation.source_target,
1202                    "retained Move target is missing or changed; refusing to recreate the environment"
1203                );
1204            }
1205            if self.state.sessions[&id].state == SessionState::Error
1206                && let Some(previous) = operation.recovery_session.as_ref()
1207            {
1208                // A prior failed teardown retains its exact partial checkout.
1209                // Restore source identity only after that owner is stopped.
1210                let cause = anyhow::anyhow!("clean up the partial Move destination before retry");
1211                if operation.in_place {
1212                    self.retain_failed_in_place_move(&id, previous, cause)?;
1213                    let record = self.state.sessions.get_mut(&id).unwrap();
1214                    record.state = SessionState::Closing;
1215                    crate::database::save_resumed_session(record, None)?;
1216                } else if operation.workspace_transfer.is_some() {
1217                    self.rollback_move_destination(operation, cause, executor)?;
1218                    let record = self.state.sessions.get_mut(&id).unwrap();
1219                    record.state = SessionState::Closing;
1220                    crate::database::save_resumed_session(record, None)?;
1221                } else {
1222                    let failure =
1223                        self.rollback_failed_resume(&id, previous, false, cause, executor)?;
1224                    ensure!(
1225                        self.state.sessions[&id].state == SessionState::Stopped,
1226                        "{failure:#}"
1227                    );
1228                }
1229            }
1230            let state = self.state.sessions[&id].state;
1231            if matches!(state, SessionState::Closing | SessionState::Destroying)
1232                && !(operation.retains_source_environment() && operation.recovery_session.is_some())
1233            {
1234                Box::pin(self.recover_move_source_stop(operation, executor, manager)).await?;
1235                operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
1236            }
1237            if matches!(
1238                self.state.sessions[&id].state,
1239                SessionState::Running | SessionState::Disconnected
1240            ) && operation.destination_target.is_none()
1241            {
1242                executor.notify_notice("Stopping source");
1243                let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
1244                operation.phase = MovePhase::ClosingSource;
1245                operation.updated_at = now();
1246                crate::database::save_move_operation(operation)?;
1247                // An eligible move keeps its environment: the close stops at
1248                // the sealed relay and the record keeps its target for
1249                // `restore_session_in_place`.
1250                let disposition = if operation.in_place {
1251                    SourceTargetDisposition::RetainForInPlaceSwap
1252                } else {
1253                    SourceTargetDisposition::Destroy
1254                };
1255                Box::pin(self.suspend_session_for_move(
1256                    &id,
1257                    executor,
1258                    manager,
1259                    operation,
1260                    preparation,
1261                    disposition,
1262                    std::mem::take(&mut source_relay),
1263                ))
1264                .await?;
1265            }
1266            // Past the source stop, nothing else needs the source's connection.
1267            drop(source_relay);
1268            // An in-place move never reaches `Stopped` with a target: its
1269            // source stays `Closing` in the environment the destination reuses.
1270            if !operation.in_place
1271                && operation.workspace_transfer.is_none()
1272                && self.state.sessions[&id].state == SessionState::Stopped
1273                && self.state.sessions[&id].target.is_some()
1274            {
1275                executor.notify_notice("Cleaning up source");
1276                let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
1277                self.cleanup_stopped_target(&id, executor)?;
1278            }
1279            operation.checkpoint = operation
1280                .checkpoint
1281                .clone()
1282                .or_else(|| self.state.sessions[&id].checkpoint.clone());
1283            ensure!(
1284                operation.restore_artifact().is_some(),
1285                "move has no verified checkpoint"
1286            );
1287            ensure!(
1288                !executor.cancellation_requested(),
1289                "move cancelled after source sealing; checkpoint and remaining environment retained"
1290            );
1291            operation.phase = MovePhase::ResumingDestination;
1292            operation.recovery_session = Some(self.state.sessions[&id].clone());
1293            operation.updated_at = now();
1294            crate::database::save_move_operation(operation)?;
1295            executor.reserve_move_destination();
1296            executor.notify_notice("Preparing destination");
1297            // The destination worker reads its delegation policy from the
1298            // record at launch. `recovery_session` above keeps the old one
1299            // for a rollback.
1300            if let Some(policy) = &operation.selection.subagents {
1301                let session = self.state.sessions.get_mut(&id).unwrap();
1302                session.subagents = Some(policy.clone());
1303                crate::database::save_resumed_session(session, None)?;
1304            }
1305            if operation.in_place {
1306                // The environment is kept, so there is nothing to provision and
1307                // no allocation to clear: eligibility already required the
1308                // allocation to be unchanged.
1309                Box::pin(self.restore_session_in_place(
1310                    &id,
1311                    operation.selection.profile_id.as_deref().unwrap(),
1312                    operation.selection.target_template_id.as_deref().unwrap(),
1313                    executor,
1314                ))
1315                .await?;
1316            } else {
1317                if operation.selection.clear_resource_allocation {
1318                    let session = self.state.sessions.get_mut(&id).unwrap();
1319                    session.resource_allocation = None;
1320                    session.container_cpus = None;
1321                    session.container_memory = None;
1322                    crate::database::save_resumed_session(session, None)?;
1323                }
1324                if operation.workspace_transfer.is_some() {
1325                    self.capture_move_workspace(operation, executor)?;
1326                    Box::pin(self.resume_session_for_move(operation, executor)).await?;
1327                    operation.workspace_transfer = crate::database::load_move_operation(&id)?
1328                        .context("Move intent disappeared during restore")?
1329                        .workspace_transfer;
1330                } else {
1331                    Box::pin(self.resume_session_controlled(
1332                        &id,
1333                        operation.selection.profile_id.as_deref().unwrap(),
1334                        operation.selection.target_template_id.as_deref().unwrap(),
1335                        SessionResumeOptions {
1336                            additional_mounts: operation.selection.additional_mounts.clone(),
1337                            resource_allocation: operation.selection.resource_allocation.clone(),
1338                            discard_queue: true,
1339                        },
1340                        executor,
1341                    ))
1342                    .await?;
1343                }
1344            }
1345            let destination = &self.state.sessions[&id];
1346            operation.destination_target = destination.target.clone();
1347            operation.destination_native_session_id = destination.native_session_id.clone();
1348            // Persist readiness before submitting even the first queued command.
1349            operation.phase = MovePhase::StartingQueue;
1350            operation.queue_admission_started = true;
1351            operation.updated_at = now();
1352            crate::database::save_move_operation(operation)?;
1353        }
1354        restore_move_queue_hold(operation);
1355        self.finish_workspace_transfer(operation, executor)?;
1356        self.admit_move_queue(operation, executor).await
1357    }
1358
1359    async fn admit_move_queue(
1360        &self,
1361        operation: &mut MoveOperation,
1362        executor: &(impl CommandExecutor + Sync),
1363    ) -> Result<()> {
1364        let timing_id = operation.selection.session_id.clone();
1365        let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
1366        let id = &operation.selection.session_id;
1367        let mut relay = {
1368            let _checking_destination =
1369                ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1370            let destination = &self.state.sessions[id];
1371            ensure!(
1372                destination.state == SessionState::Running
1373                    && destination.target == operation.destination_target
1374                    && destination.native_session_id == operation.destination_native_session_id,
1375                "cannot prove the same ready destination; refusing to replay potentially executed work"
1376            );
1377            let spec = self.reconnect_command(id)?;
1378            let relay = StandaloneSession::connect_command(&spec, id).await?;
1379            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")?;
1380            if let Some(expected) = &operation.destination_store_id {
1381                ensure!(
1382                    *expected == store_id,
1383                    "destination relay storage was replaced; refusing to replay potentially executed work"
1384                );
1385            } else {
1386                // No command can be admitted until this identity is durable.
1387                operation.destination_store_id = Some(store_id);
1388                crate::database::save_move_operation(operation)?;
1389            }
1390            ensure!(
1391                relay.snapshot().operational.native_session_id
1392                    == operation.destination_native_session_id,
1393                "destination relay native identity changed; refusing queue replay"
1394            );
1395            relay
1396        };
1397        if operation.queue == ResumeQueueDisposition::Start {
1398            let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
1399            executor.notify_notice("Starting queued work");
1400            let checkpoint = operation
1401                .restore_artifact()
1402                .context("move queue archive is missing")?;
1403            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1404            ensure!(
1405                verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
1406                "move queue checkpoint verification failed"
1407            );
1408            for queued in verified.canonical_session.queued_prompts {
1409                if operation.queue_admission_finished {
1410                    continue;
1411                }
1412                ensure!(
1413                    !executor.cancellation_requested(),
1414                    "move cancelled during queue admission; destination retained"
1415                );
1416                let command = match queued.kind {
1417                    CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
1418                        prompt: queued
1419                            .content
1420                            .into_iter()
1421                            .map(serde_json::from_value)
1422                            .collect::<serde_json::Result<_>>()?,
1423                    },
1424                    CanonicalQueuedCommandKind::SetConfig { key, value } => {
1425                        RelayCommand::SetConfig { key, value }
1426                    }
1427                };
1428                relay.submit_accepted(queued.command_id, command).await?;
1429            }
1430        }
1431        operation.queue_admission_finished = true;
1432        crate::database::save_move_operation(operation)?;
1433        restore_move_queue_hold(operation);
1434        let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
1435            "Queued work was discarded; ready and idle."
1436        } else {
1437            "Queued work was accepted."
1438        };
1439        let source_profile = &operation.source_profile_id;
1440        let source_target = &operation.source_target_template_id;
1441        let destination_profile = operation.selection.profile_id.as_deref().unwrap();
1442        let destination_target = operation.selection.target_template_id.as_deref().unwrap();
1443        // The two moves are different events for the person reading the
1444        // conversation: one rebuilt the environment, the other kept it.
1445        let text = if operation.in_place {
1446            format!(
1447                "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."
1448            )
1449        } else {
1450            format!(
1451                "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
1452            )
1453        };
1454        relay
1455            .submit(
1456                format!("{}-notice", operation.operation_id),
1457                RelayCommand::RecordNotice { text },
1458            )
1459            .await?;
1460        Ok(())
1461    }
1462
1463    pub(super) fn validate_move_checkpoint(
1464        &self,
1465        operation: &MoveOperation,
1466        preparation: Option<&MovePreparation>,
1467        executor: &(impl CommandExecutor + Sync),
1468    ) -> Result<()> {
1469        let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1470        let current = Controller {
1471            config: mj_core::config::Config::load()?,
1472            state: self.state.clone(),
1473        };
1474        ensure!(
1475            current.move_configuration_fingerprint(&operation.selection)?
1476                == operation.configuration_fingerprint,
1477            "destination configuration changed during move"
1478        );
1479        let id = &operation.selection.session_id;
1480        current.validate_move_destination_paths(
1481            &self.state.sessions[id],
1482            operation.selection.target_template_id.as_deref().unwrap(),
1483            executor,
1484        )?;
1485        if !operation.in_place
1486            && operation.workspace_transfer.is_none()
1487            && let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
1488                .preflight_resume_repository_sources(
1489                    id,
1490                    operation.selection.target_template_id.as_deref().unwrap(),
1491                    executor,
1492                )?
1493        {
1494            bail!(
1495                "destination repository source is missing checkpoint commit {}; source retained",
1496                mismatch.missing_commit
1497            );
1498        }
1499        if let Some(prepared) = preparation {
1500            let checkpoint = operation
1501                .restore_artifact()
1502                .or(self.state.sessions[id].checkpoint.as_ref())
1503                .context("no move checkpoint")?;
1504            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1505            let actual: Vec<_> = verified
1506                .canonical_session
1507                .queued_prompts
1508                .iter()
1509                .map(|p| p.command_id.as_str())
1510                .collect();
1511            let expected: Vec<_> = prepared
1512                .queued_commands
1513                .iter()
1514                .map(|p| p.command_id.as_str())
1515                .collect();
1516            ensure!(
1517                actual == expected,
1518                "pending queue changed before checkpoint capture; source retained, confirm Move again"
1519            );
1520        }
1521        Ok(())
1522    }
1523}
1524
1525fn validate_preserved_configuration(
1526    profile_id: &str,
1527    accepted: &mj_core::acp::AcceptedSessionConfig,
1528    choices: &mj_core::worker_launch::ProfileConfig,
1529) -> Result<()> {
1530    for (key, value, offered) in [
1531        ("model", accepted.model.as_deref(), &choices.models),
1532        ("effort", accepted.effort.as_deref(), &choices.efforts),
1533    ] {
1534        let Some(value) = value else { continue };
1535        ensure!(
1536            offered.iter().any(|choice| choice.value == value),
1537            "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
1538            offered
1539                .iter()
1540                .map(|choice| choice.value.as_str())
1541                .collect::<Vec<_>>()
1542                .join(", ")
1543        );
1544    }
1545    Ok(())
1546}
1547
1548/// Whether this move can replace only the harness inside the source target.
1549///
1550/// The environment is kept only when nothing about it changes: the same target
1551/// template (or another one that names the same bare machine), the same
1552/// attached mounts, and the same resource allocation. A retry takes its retention decision from the durable Move intent instead;
1553/// a sub-agent never owns its own target.
1554pub(super) fn in_place_move_eligible(
1555    source: &mj_core::state::SessionRecord,
1556    selection: &MoveSelection,
1557    targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
1558    is_subagent: bool,
1559    retry: bool,
1560) -> bool {
1561    !retry
1562        && !is_subagent
1563        && source.target.is_some()
1564        && matches!(
1565            source.state,
1566            SessionState::Running | SessionState::Disconnected
1567        )
1568        && move_environment_change(
1569            source,
1570            selection,
1571            bare_targets_share_environment(
1572                targets,
1573                &source.target_template_id,
1574                selection.target_template_id.as_deref().unwrap_or_default(),
1575            ),
1576        )
1577        .is_none()
1578}
1579
1580/// Whether two target templates are bare targets on the same machine, so the
1581/// session's worker root and workspace are already where the destination wants
1582/// them. Two local bare targets qualify, and so do two SSH bare targets with
1583/// the same connection; only the template the session names changes.
1584fn bare_targets_share_environment(
1585    targets: &std::collections::BTreeMap<String, mj_core::config::TargetTemplate>,
1586    source_id: &str,
1587    destination_id: &str,
1588) -> bool {
1589    use mj_core::config::TargetTemplate::{LocalBare, SshBare};
1590    match (targets.get(source_id), targets.get(destination_id)) {
1591        (Some(LocalBare), Some(LocalBare)) => true,
1592        (
1593            Some(SshBare { ssh: source, .. }),
1594            Some(SshBare {
1595                ssh: destination, ..
1596            }),
1597        ) => {
1598            crate::targets::SshTarget::from(source) == crate::targets::SshTarget::from(destination)
1599        }
1600        _ => false,
1601    }
1602}
1603fn outcome(
1604    operation_id: &str,
1605    selection: &MoveSelection,
1606    status: &str,
1607    error: Option<String>,
1608    recovery: Option<String>,
1609) -> MoveOutcome {
1610    MoveOutcome {
1611        operation_id: operation_id.into(),
1612        session_id: selection.session_id.clone(),
1613        profile_id: selection.profile_id.clone().unwrap_or_default(),
1614        target_template_id: selection.target_template_id.clone().unwrap_or_default(),
1615        outcome: status.into(),
1616        error,
1617        recovery,
1618    }
1619}
1620
1621/// One decision shared by no-op detection and environment retention. A
1622/// different target template counts as a change unless `same_environment` says
1623/// the two name the same bare machine.
1624fn move_environment_change(
1625    source: &mj_core::state::SessionRecord,
1626    selection: &MoveSelection,
1627    same_environment: bool,
1628) -> Option<&'static str> {
1629    if Some(&source.target_template_id) != selection.target_template_id.as_ref()
1630        && !same_environment
1631    {
1632        Some("target changed")
1633    } else if Some(&source.additional_mounts) != selection.additional_mounts.as_ref() {
1634        Some("attached mounts changed")
1635    } else if source.resource_allocation != selection.resource_allocation
1636        || (selection.clear_resource_allocation
1637            && (source.container_cpus.is_some() || source.container_memory.is_some()))
1638    {
1639        Some("resource allocation changed")
1640    } else {
1641        None
1642    }
1643}