Skip to main content

mj_controller/controller/
move_session.rs

1//! One recoverable stop/restore operation, independent of its initiating viewer.
2
3#[cfg(test)]
4mod tests;
5
6use anyhow::{Context, Result, bail, ensure};
7use mj_core::hex::lower_hex;
8use sha2::{Digest, Sha256};
9
10use super::lifecycle::SourceTargetDisposition;
11use super::{Controller, SessionResumeOptions, now};
12
13fn mutation_holds() -> &'static std::sync::Mutex<std::collections::BTreeSet<String>> {
14    static HOLDS: std::sync::OnceLock<std::sync::Mutex<std::collections::BTreeSet<String>>> =
15        std::sync::OnceLock::new();
16    HOLDS.get_or_init(Default::default)
17}
18
19pub fn move_owns_session(session_id: &str) -> bool {
20    move_has_pending_queue(session_id)
21        || mutation_holds()
22            .lock()
23            .unwrap_or_else(std::sync::PoisonError::into_inner)
24            .contains(session_id)
25}
26
27fn queue_holds() -> &'static std::sync::Mutex<std::collections::BTreeSet<String>> {
28    static HOLDS: std::sync::OnceLock<std::sync::Mutex<std::collections::BTreeSet<String>>> =
29        std::sync::OnceLock::new();
30    HOLDS.get_or_init(Default::default)
31}
32
33pub fn move_has_pending_queue(session_id: &str) -> bool {
34    queue_holds()
35        .lock()
36        .unwrap_or_else(std::sync::PoisonError::into_inner)
37        .contains(session_id)
38}
39
40pub fn release_move_queue_hold(session_id: &str) {
41    queue_holds()
42        .lock()
43        .unwrap_or_else(std::sync::PoisonError::into_inner)
44        .remove(session_id);
45}
46
47pub fn restore_move_queue_hold(operation: &MoveOperation) {
48    let mut holds = queue_holds()
49        .lock()
50        .unwrap_or_else(std::sync::PoisonError::into_inner);
51    if operation.queue_admission_started && !operation.queue_admission_finished {
52        holds.insert(operation.selection.session_id.clone());
53    } else {
54        holds.remove(&operation.selection.session_id);
55    }
56}
57
58/// What a recovered, interrupted source stop tells the person.
59///
60/// An in-place move promised to keep the environment; a recovery that tears it
61/// down has broken that promise, so the message says so instead of leaving the
62/// person to infer it from an unchanged "retry" line.
63fn interrupted_source_stop_message(recovered: &str, in_place: bool) -> String {
64    if in_place {
65        format!("{recovered}; the in-place swap was interrupted; the environment was released")
66    } else {
67        recovered.to_owned()
68    }
69}
70
71/// Whether the persisted record shows a source that has finished stopping with
72/// a verified checkpoint to restore.
73///
74/// Only `Stopped` proves the source stop finished. Resume also admits recovery
75/// states, but those do not establish that the source is already stopped.
76fn source_stopped_with_verified_checkpoint(record: &mj_core::state::SessionRecord) -> bool {
77    record.state == SessionState::Stopped && record.checkpoint.is_some()
78}
79
80/// The recovery guidance for a failed Move whose source is already stopped with
81/// a verified checkpoint.
82///
83/// Resume restores the retained checkpoint on the destination profile and
84/// target the Move selected.
85fn stopped_source_recovery(
86    session_id: &str,
87    destination_profile: Option<&str>,
88    destination_target: Option<&str>,
89) -> String {
90    let flag = |name: &str, value: Option<&str>| {
91        value
92            .filter(|value| !value.is_empty())
93            .map(|value| format!(" --{name} {value}"))
94            .unwrap_or_default()
95    };
96    format!(
97        "Source is stopped with a verified checkpoint. Bring it back with \
98         `mj resume --session {session_id}{}{} --queue start`.",
99        flag("profile", destination_profile),
100        flag("target", destination_target),
101    )
102}
103
104/// The recovery guidance a failed or cancelled Move shows, from its persisted
105/// state.
106fn failed_move_recovery(
107    operation: &MoveOperation,
108    record: Option<&mj_core::state::SessionRecord>,
109) -> String {
110    if operation.queue_admission_started {
111        return "Destination is live; retry queue admission on this same destination. Already accepted work may have effects.".to_owned();
112    }
113    match record {
114        Some(record) if source_stopped_with_verified_checkpoint(record) => {
115            stopped_source_recovery(
116                &operation.selection.session_id,
117                operation.selection.profile_id.as_deref(),
118                operation.selection.target_template_id.as_deref(),
119            )
120        }
121        _ => "Source or partial destination is retained. Retry move after resolving the reported error.".to_owned(),
122    }
123}
124
125pub struct MoveMutationGuard(String);
126
127impl MoveMutationGuard {
128    pub fn reserve(session_id: &str) -> Result<Self> {
129        ensure!(
130            mutation_holds()
131                .lock()
132                .unwrap_or_else(std::sync::PoisonError::into_inner)
133                .insert(session_id.to_owned()),
134            "session already has a move owner"
135        );
136        Ok(Self(session_id.to_owned()))
137    }
138}
139
140impl Drop for MoveMutationGuard {
141    fn drop(&mut self) {
142        mutation_holds()
143            .lock()
144            .unwrap_or_else(std::sync::PoisonError::into_inner)
145            .remove(&self.0);
146    }
147}
148
149pub(crate) fn move_refuses_command(session_id: &str, command: &RelayCommand) -> bool {
150    move_owns_session(session_id)
151        && matches!(
152            command,
153            RelayCommand::Prompt { .. }
154                | RelayCommand::SetConfig { .. }
155                | RelayCommand::SetSessionMode { .. }
156                | RelayCommand::RunUserShell { .. }
157                | RelayCommand::CancelUserShell { .. }
158                | RelayCommand::Cancel
159                | RelayCommand::RemoveQueuedPrompt { .. }
160                | RelayCommand::ClearQueuedPrompts
161        )
162}
163use crate::session_manager::{SessionManagerControl, StandaloneSession, new_command_id};
164use mj_checkpoint::archive::{CanonicalQueuedCommandKind, verify_archive_streaming};
165use mj_core::state::{MoveOperation, MovePhase, ResumeQueueDisposition, SessionState};
166
167pub use mj_core::state::{MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest};
168
169use crate::targets::{CommandExecutor, ProvisionStage, ProvisionStageGuard};
170use mj_core::relay::RelayCommand;
171
172/// Refresh source state without turning a dead source harness into a Move prerequisite.
173/// Leasing preserves typed transport failures, unlike the UI's string-valued sync reply.
174pub async fn refresh_move_source(
175    manager: &SessionManagerControl,
176    id: &str,
177) -> Result<Option<mj_core::state::ManagedSessionSnapshot>> {
178    let handle = manager
179        .wait_for_session(id, std::time::Duration::from_secs(5))
180        .await?;
181    let result = async {
182        let mut lease = handle.lease_connection().await?;
183        let snapshot = lease.connection_mut().sync().await?;
184        lease.release();
185        Ok(snapshot)
186    }
187    .await;
188    match result {
189        Ok(snapshot) => Ok(Some(snapshot)),
190        Err(error) if crate::worker_client::RelayTransportDead::marks(&error) => {
191            tracing::warn!(session_id = id, error = %error, "Move will recover the unavailable source without its harness");
192            Ok(None)
193        }
194        Err(error) => Err(error),
195    }
196}
197
198fn digest(value: &impl serde::Serialize) -> Result<String> {
199    Ok(lower_hex(Sha256::digest(serde_json::to_vec(value)?)))
200}
201
202struct MovePhaseTimer<'a> {
203    session_id: &'a str,
204    phase: &'static str,
205    started: std::time::Instant,
206}
207
208impl<'a> MovePhaseTimer<'a> {
209    fn new(session_id: &'a str, phase: &'static str) -> Self {
210        Self {
211            session_id,
212            phase,
213            started: std::time::Instant::now(),
214        }
215    }
216}
217
218impl Drop for MovePhaseTimer<'_> {
219    fn drop(&mut self) {
220        tracing::info!(
221            session_id = self.session_id,
222            phase = self.phase,
223            elapsed_ms = self.started.elapsed().as_millis() as u64,
224            "move phase finished"
225        );
226    }
227}
228
229/// Confirmation is an inspector, not transport for attachment bytes: every
230/// surface that shows a move preparation sees a placeholder for each queued
231/// image. Replay always reads the verified archive after destination readiness.
232fn replace_queued_images_with_placeholders(
233    queued_commands: &mut [mj_core::state::MaterializedQueuedPrompt],
234) {
235    for block in queued_commands
236        .iter_mut()
237        .flat_map(|command| command.content.iter_mut())
238    {
239        if block.get("type").and_then(serde_json::Value::as_str) != Some("image") {
240            continue;
241        }
242        let mime = block
243            .get("mimeType")
244            .or_else(|| block.get("mime_type"))
245            .and_then(serde_json::Value::as_str)
246            .unwrap_or("image");
247        *block = serde_json::json!({"type": "text", "text": format!("[Image attachment: {mime}]")});
248    }
249}
250
251impl Controller {
252    /// Returns the planned conversion when this move turns a local checkout
253    /// into an isolated workspace, so the caller can describe it without
254    /// reading Git a second time.
255    fn validate_move_destination_paths(
256        &self,
257        source: &mj_core::state::SessionRecord,
258        target_id: &str,
259        executor: &(impl CommandExecutor + Sync),
260    ) -> Result<Option<super::worktree::RawToWorkspaceConversion>> {
261        use super::worktree::ResumePlan;
262        match super::worktree::resume_compatibility(source, &self.config, target_id)
263            .map_err(anyhow::Error::msg)?
264        {
265            ResumePlan::RawToWorkspace => {
266                return Ok(Some(super::worktree::plan_raw_to_workspace(
267                    source,
268                    &self.config,
269                    executor,
270                )?));
271            }
272            ResumePlan::WorkspaceToRaw => {
273                self.plan_workspace_to_raw(source, target_id, executor)?;
274            }
275            ResumePlan::InPlace if source.managed_worktree.is_none() => {
276                if let Some(path) = &source.project_directory {
277                    self.validate_project_directory(target_id, path, executor)?;
278                }
279            }
280            ResumePlan::InPlace => {}
281        }
282        Ok(None)
283    }
284    fn move_confirmation(
285        &self,
286        selection: &MoveSelection,
287        conversion: Option<&mj_core::state::RawConversionPreview>,
288    ) -> Result<(bool, Vec<mj_core::state::MaterializedQueuedPrompt>, String)> {
289        let source = self
290            .state
291            .sessions
292            .get(&selection.session_id)
293            .context("unknown move session")?;
294        let (mut active, mut queued) = crate::database::move_pending_work(&source.id)?;
295        if let Some(operation) = crate::database::load_move_operation(&source.id)?
296            && operation.queue_admission_started
297            && !operation.queue_admission_finished
298        {
299            ensure!(
300                operation.selection == *selection,
301                "queue admission is incomplete on the live destination; retry that move before selecting another destination"
302            );
303            let checkpoint = operation
304                .checkpoint
305                .as_ref()
306                .context("retained queue checkpoint is missing")?;
307            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
308            ensure!(
309                verified.archive_sha256 == checkpoint.sha256
310                    && verified.manifest.session.id == source.id,
311                "retained move checkpoint verification failed"
312            );
313            queued = verified
314                .canonical_session
315                .queued_prompts
316                .into_iter()
317                .map(|entry| mj_core::state::MaterializedQueuedPrompt {
318                    accepted_ordinal: None,
319                    command_id: entry.command_id,
320                    kind: match entry.kind {
321                        CanonicalQueuedCommandKind::Prompt => {
322                            mj_core::state::QueuedCommandKind::Prompt
323                        }
324                        CanonicalQueuedCommandKind::SetConfig { key, value } => {
325                            mj_core::state::QueuedCommandKind::SetConfig { key, value }
326                        }
327                    },
328                    content: entry.content,
329                    queued_at_ms: entry.queued_at_ms,
330                })
331                .collect();
332            active = false;
333        }
334        let fingerprint = digest(&(
335            &source.last_profile,
336            &source.target_template_id,
337            &source.target,
338            &source.target_runtime,
339            &source.native_session_id,
340            &source.resource_allocation,
341            &source.additional_mounts,
342            &source.container_cpus,
343            &source.container_memory,
344            selection,
345            self.move_configuration_fingerprint(selection)?,
346            &queued,
347            // Only what the destination is built from. The dirty counts move
348            // with every keystroke of a live agent, and hashing them would
349            // invalidate the confirmation the person is reading.
350            conversion.map(|preview| {
351                (
352                    &preview.fetch_url,
353                    &preview.push_urls,
354                    &preview.branch,
355                    &preview.destination,
356                )
357            }),
358        ))?;
359        Ok((active, queued, fingerprint))
360    }
361    pub(super) fn move_configuration_fingerprint(
362        &self,
363        selection: &MoveSelection,
364    ) -> Result<String> {
365        let profile = self
366            .config
367            .profiles
368            .get(
369                selection
370                    .profile_id
371                    .as_deref()
372                    .context("move profile is unresolved")?,
373            )
374            .context("move profile no longer exists")?;
375        let target = self
376            .config
377            .targets
378            .get(
379                selection
380                    .target_template_id
381                    .as_deref()
382                    .context("move target is unresolved")?,
383            )
384            .context("move target no longer exists")?;
385        // Only a digest crosses IPC or enters the move record. Configuration
386        // may contain credential-bearing values and must never be copied there.
387        let source = self
388            .state
389            .sessions
390            .get(&selection.session_id)
391            .context("unknown move session")?;
392        digest(&(
393            profile,
394            target,
395            &self.config.bundles,
396            self.config.targets.get(&source.target_template_id),
397            self.config.profiles.get(&source.last_profile),
398        ))
399    }
400
401    async fn validate_move_destination_configuration(
402        &self,
403        selection: &MoveSelection,
404        source_harness: mj_core::config::HarnessKind,
405        operational: &mj_core::relay::RelayOperationalState,
406    ) -> Result<()> {
407        let profile_id = selection
408            .profile_id
409            .as_deref()
410            .context("move profile is unresolved")?;
411        let profile = self
412            .config
413            .profiles
414            .get(profile_id)
415            .context("move profile no longer exists")?;
416        let source = self
417            .state
418            .sessions
419            .get(&selection.session_id)
420            .context("unknown move session")?;
421        if profile.kind != source_harness || profile_id == source.last_profile {
422            return Ok(());
423        }
424        let accepted = mj_core::acp::AcceptedSessionConfig::from_configuration(
425            &operational.config,
426            &operational.config_options,
427        );
428        if accepted.model.is_none() && accepted.effort.is_none() {
429            return Ok(());
430        }
431        // A fresh probe prevents a stale local catalogue from approving a
432        // move that the destination profile will immediately reset.
433        let choices =
434            super::profile_config::discover(profile_id.to_owned(), accepted.model.clone(), true)
435                .await
436                .with_context(|| {
437                    format!("discover destination profile {profile_id:?} configuration")
438                })?;
439        validate_preserved_configuration(profile_id, &accepted, &choices)
440    }
441
442    pub async fn prepare_move_session_controlled(
443        &self,
444        mut selection: MoveSelection,
445        executor: &(impl CommandExecutor + Sync),
446    ) -> Result<MovePreparation> {
447        ensure!(
448            selection.profile_id.is_some() || selection.target_template_id.is_some(),
449            "move requires a target or profile selection"
450        );
451        let source = self
452            .state
453            .sessions
454            .get(&selection.session_id)
455            .context("unknown session")?;
456        ensure!(
457            !self.state.subagents.contains_key(&source.id),
458            "sub-agent sessions cannot move independently of their parent"
459        );
460        ensure!(
461            !self.state.subagents.values().any(|child| {
462                child.parent_session_id == source.id
463                    && self
464                        .state
465                        .sessions
466                        .get(&child.child_session_id)
467                        .is_some_and(|session| session.state.is_active())
468            }),
469            "stop active sub-agents before moving their parent session"
470        );
471        let previous = crate::database::load_move_operation(&source.id)?;
472        let retry = previous.as_ref().is_some_and(|op| {
473            !matches!(op.phase, MovePhase::Completed) && source.checkpoint.is_some()
474        });
475        ensure!(
476            matches!(
477                source.state,
478                SessionState::Running | SessionState::Disconnected
479            ) || retry,
480            "only active sessions can move; run `mj resume` (or POST /api/v1/sessions/{}/resume) for a stopped or lost session",
481            source.id
482        );
483        selection
484            .profile_id
485            .get_or_insert_with(|| source.last_profile.clone());
486        selection
487            .target_template_id
488            .get_or_insert_with(|| source.target_template_id.clone());
489        selection
490            .additional_mounts
491            .get_or_insert_with(|| source.additional_mounts.clone());
492        ensure!(
493            !selection.clear_resource_allocation || selection.resource_allocation.is_none(),
494            "select resource allocation or explicitly clear it, not both"
495        );
496        if !selection.clear_resource_allocation {
497            selection.resource_allocation = selection
498                .resource_allocation
499                .or_else(|| source.resource_allocation.clone());
500        }
501        let profile_id = selection.profile_id.as_deref().unwrap();
502        let target_id = selection.target_template_id.as_deref().unwrap();
503        let profile = self
504            .config
505            .profiles
506            .get(profile_id)
507            .context("unknown destination profile")?;
508        ensure!(
509            profile.enabled,
510            "destination profile {profile_id:?} is disabled"
511        );
512        let target = self
513            .config
514            .targets
515            .get(target_id)
516            .context("unknown destination target")?;
517        self.validate_muse_resume_destination(source, profile.kind, target_id)?;
518        super::worktree::resume_compatibility(source, &self.config, target_id)
519            .map_err(anyhow::Error::msg)?;
520        super::backend::validate_resource_allocation(
521            target,
522            selection.resource_allocation.as_ref(),
523        )?;
524        let mounts = selection.additional_mounts.as_deref().unwrap_or_default();
525        ensure!(
526            profile.kind != mj_core::config::HarnessKind::Muse || mounts.is_empty(),
527            "Muse Code ACP supports one workspace root; attached directories are unsupported"
528        );
529        ensure!(
530            mounts.is_empty() || mj_core::config::mount_history_host(target).is_some(),
531            "attached resources are unsupported for this target; select compatible resources explicitly"
532        );
533        crate::targets::validate_additional_mounts(mounts)?;
534        for mount in mounts {
535            self.validate_mount_source(target_id, &mount.source, executor)?;
536        }
537        let planned_conversion =
538            self.validate_move_destination_paths(source, target_id, executor)?;
539        ensure!(
540            profile.home.is_dir(),
541            "destination profile home is unavailable; configure the profile before moving"
542        );
543        super::worker_binary::preflight_worker_binary(target)?;
544        super::backend::preflight_target(target, executor, super::backend::TargetCheck::Launch)?;
545        let source_harness = previous
546            .as_ref()
547            .filter(|operation| {
548                !operation.queue_admission_started && operation.phase != MovePhase::Completed
549            })
550            .and_then(|operation| operation.recovery_session.as_ref())
551            .map_or(source.harness_kind, |record| record.harness_kind);
552        let cross_harness = profile.kind != source_harness;
553        if cross_harness {
554            let cancel = tokio_util::sync::CancellationToken::new();
555            let resolve =
556                crate::utility_llm::UtilityLlmRuntime::shared().resolve(&self.config, &cancel);
557            tokio::pin!(resolve);
558            loop {
559                tokio::select! {
560                    result = &mut resolve => { result.context("cross-harness move needs an available utility model")?; break; }
561                    _ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {
562                        if executor.cancellation_requested() { cancel.cancel(); bail!("move preparation cancelled"); }
563                    }
564                }
565            }
566        }
567        // What moving a local checkout into a target really does, computed
568        // before anything is stopped so a person can confirm it.
569        let conversion = planned_conversion
570            .map(|conversion| {
571                super::worktree::raw_conversion_preview(source, &conversion, executor)
572                    .context("describe the move of this checkout into the target")
573            })
574            .transpose()?
575            .map(Box::new);
576        let (active, mut queued_commands, fingerprint) =
577            self.move_confirmation(&selection, conversion.as_deref())?;
578        replace_queued_images_with_placeholders(&mut queued_commands);
579        let operation_id = previous
580            .as_ref()
581            .filter(|operation| {
582                operation.selection == selection
583                    && operation.phase != MovePhase::Completed
584                    && operation.checkpoint.is_some()
585            })
586            .map(|operation| operation.operation_id.clone())
587            .unwrap_or(new_command_id("move")?);
588        Ok(MovePreparation {
589            source_unavailable: false,
590            in_place: in_place_move_eligible(
591                source,
592                &selection,
593                self.state.subagents.contains_key(&source.id),
594                retry,
595            ),
596            conversion,
597            selection,
598            source_profile_id: source.last_profile.clone(),
599            source_target_template_id: source.target_template_id.clone(),
600            cross_harness,
601            active,
602            queued_commands,
603            fingerprint,
604            operation_id,
605        })
606    }
607
608    /// Called with one daemon lifecycle and recovery reservation already held.
609    pub async fn move_session_managed_controlled(
610        &mut self,
611        request: MoveSessionRequest,
612        executor: &(impl CommandExecutor + Sync),
613        manager: &SessionManagerControl,
614    ) -> Result<MoveOutcome> {
615        let prepared = &request.preparation;
616        let id = prepared.selection.session_id.clone();
617        let started = std::time::Instant::now();
618        executor.notify_notice("Checking destination");
619        let checked = {
620            let _checking_destination =
621                ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
622            let mut checked = self
623                .prepare_move_session_controlled(prepared.selection.clone(), executor)
624                .await?;
625            // Destination checks can outlast a turn. Refresh confirmation from the
626            // relay immediately before interruption, while new submissions are held.
627            if matches!(
628                self.state.sessions[&id].state,
629                SessionState::Running | SessionState::Disconnected
630            ) {
631                let source_harness = self.state.sessions[&id].harness_kind;
632                let snapshot = refresh_move_source(manager, &id).await?;
633                let (active, queue, fingerprint) =
634                    self.move_confirmation(&checked.selection, checked.conversion.as_deref())?;
635                checked.source_unavailable = snapshot
636                    .as_ref()
637                    .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
638                checked.active = active
639                    || checked.source_unavailable
640                    || snapshot.as_ref().is_some_and(|snapshot| {
641                        let mut operational = snapshot.operational.clone();
642                        operational.queued_prompts.clear();
643                        operational.checkpoint_barrier = None;
644                        !operational.safe_to_replace(source_harness)
645                    });
646                if let Some(snapshot) = &snapshot {
647                    self.validate_move_destination_configuration(
648                        &checked.selection,
649                        source_harness,
650                        &snapshot.operational,
651                    )
652                    .await?;
653                }
654                checked.queued_commands = queue;
655                checked.fingerprint = fingerprint;
656            }
657            checked
658        };
659        ensure!(
660            checked.fingerprint == prepared.fingerprint,
661            "session, pending work, or destination configuration changed; prepare and confirm Move again"
662        );
663        let queue = request.queue.unwrap_or(ResumeQueueDisposition::Discard);
664        let source = self.state.sessions[&id].clone();
665        let old_operation = crate::database::load_move_operation(&id)?;
666        let retry = old_operation.filter(|op| {
667            op.selection == prepared.selection
668                && op.phase != MovePhase::Completed
669                && op.checkpoint.is_some()
670        });
671        if retry.is_none()
672            && source.state == SessionState::Running
673            && source.last_profile == checked.selection.profile_id.as_deref().unwrap()
674            && source.target_template_id == checked.selection.target_template_id.as_deref().unwrap()
675            && Some(&source.additional_mounts) == checked.selection.additional_mounts.as_ref()
676            && source.resource_allocation == checked.selection.resource_allocation
677            && (!checked.selection.clear_resource_allocation
678                || (source.container_cpus.is_none() && source.container_memory.is_none()))
679        {
680            return Ok(outcome(
681                &prepared.operation_id,
682                &prepared.selection,
683                "unchanged",
684                None,
685                None,
686            ));
687        }
688        ensure!(
689            !checked.active || request.acknowledge_interruption,
690            "active work will be interrupted; confirm Move again with interruption acknowledgement"
691        );
692        ensure!(
693            checked.queued_commands.is_empty() || request.queue.is_some(),
694            "pending work requires an explicit queue choice: discard or start"
695        );
696        let timestamp = now();
697        let mut operation = match retry {
698            Some(mut op) => {
699                ensure!(
700                    !op.queue_admission_started || op.queue == queue,
701                    "queued work may already have run; retry with the original queue choice on the same destination"
702                );
703                crate::database::clear_move_cancellation_for_retry(&id)?;
704                op.configuration_fingerprint =
705                    self.move_configuration_fingerprint(&checked.selection)?;
706                op.queue = queue;
707                op.cancellation_requested = false;
708                op.error = None;
709                op
710            }
711            None => MoveOperation {
712                source_checkpoint_only: false,
713                in_place: in_place_move_eligible(
714                    &source,
715                    &checked.selection,
716                    self.state.subagents.contains_key(&id),
717                    false,
718                ),
719                operation_id: prepared.operation_id.clone(),
720                selection: checked.selection.clone(),
721                source_profile_id: source.last_profile.clone(),
722                source_target_template_id: source.target_template_id.clone(),
723                source_target: source.target.clone(),
724                source_native_session_id: source.native_session_id.clone(),
725                source_additional_mounts: source.additional_mounts.clone(),
726                source_resource_allocation: source.resource_allocation.clone(),
727                destination_target: None,
728                destination_native_session_id: None,
729                destination_store_id: None,
730                configuration_fingerprint: self
731                    .move_configuration_fingerprint(&checked.selection)?,
732                checkpoint: None,
733                recovery_session: None,
734                queue,
735                phase: MovePhase::Preparing,
736                queue_admission_started: false,
737                queue_admission_finished: false,
738                cancellation_requested: false,
739                created_at: timestamp.clone(),
740                updated_at: timestamp,
741                error: None,
742            },
743        };
744        crate::database::save_move_operation(&operation)?;
745        tracing::info!(
746            session_id = id,
747            phase = "preflight",
748            elapsed_ms = started.elapsed().as_millis() as u64,
749            "move phase completed"
750        );
751        let result = self
752            .execute_move(&mut operation, Some(&checked), executor, manager)
753            .await;
754        self.finish_move_result(&mut operation, result, executor)
755    }
756
757    fn finish_move_result(
758        &self,
759        operation: &mut MoveOperation,
760        result: Result<()>,
761        executor: &impl CommandExecutor,
762    ) -> Result<MoveOutcome> {
763        let (status, error, recovery) = match result {
764            Ok(()) => {
765                operation.phase = MovePhase::Completed;
766                ("completed", None, None)
767            }
768            Err(error) => {
769                let cancelled =
770                    executor.cancellation_requested() || operation.cancellation_requested;
771                operation.phase = if cancelled {
772                    MovePhase::Cancelled
773                } else {
774                    MovePhase::Failed
775                };
776                operation.cancellation_requested = cancelled;
777                let recovery = failed_move_recovery(
778                    operation,
779                    self.state.sessions.get(&operation.selection.session_id),
780                );
781                let error = format!("{error:#}");
782                operation.error = Some(error.clone());
783                (
784                    if cancelled { "cancelled" } else { "failed" },
785                    Some(error),
786                    Some(recovery),
787                )
788            }
789        };
790        operation.updated_at = now();
791        crate::database::save_move_operation(operation)?;
792        Ok(outcome(
793            &operation.operation_id,
794            &operation.selection,
795            status,
796            error,
797            recovery,
798        ))
799    }
800
801    pub async fn recover_move_managed_controlled(
802        &mut self,
803        mut operation: MoveOperation,
804        executor: &(impl CommandExecutor + Sync),
805        manager: &SessionManagerControl,
806    ) -> Result<MoveOutcome> {
807        let id = operation.selection.session_id.clone();
808        let result = async {
809            let session = self.state.sessions.get(&id).context("move session is missing")?.clone();
810            if operation.queue_admission_started {
811                ensure!(!operation.cancellation_requested, "Move was cancelled; destination retained without further queue admission");
812                return self.admit_move_queue(&mut operation, executor).await;
813            }
814            if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
815                let cleanup = crate::targets::CancellableProcessExecutor::with_timeout(std::time::Duration::from_secs(15));
816                self.recover_move_source_stop(&mut operation, &cleanup, manager).await?;
817                operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
818                if self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
819                    self.cleanup_stopped_target(&id, &cleanup)?;
820                }
821                bail!("{}", interrupted_source_stop_message(
822                    "Move source stop recovered; no destination work was started. Retry Move or Resume with previous settings",
823                    operation.in_place,
824                ));
825            }
826            match operation.phase {
827                MovePhase::Preparing => bail!("Move preparation was interrupted; source retained. Prepare Move again."),
828                MovePhase::ResumingDestination if session.state == SessionState::Running => {
829                    // Running is installed only after native readiness and the
830                    // handoff. A crash before the next intent write is safe to
831                    // advance, since restore started with an empty queue.
832                    ensure!(Some(&session.last_profile) == operation.selection.profile_id.as_ref()
833                        && Some(&session.target_template_id) == operation.selection.target_template_id.as_ref(),
834                        "ready destination does not match the move intent");
835                    operation.destination_target = session.target.clone();
836                    operation.destination_native_session_id = session.native_session_id.clone();
837                    operation.queue_admission_started = true;
838                    operation.phase = MovePhase::StartingQueue;
839                    crate::database::save_move_operation(&operation)?;
840                    restore_move_queue_hold(&operation);
841                    ensure!(!operation.cancellation_requested, "Move was cancelled; ready destination retained");
842                    self.admit_move_queue(&mut operation, executor).await
843                }
844                MovePhase::ResumingDestination => {
845                    let previous = operation.recovery_session.as_ref().context("move lacks its stopped recovery identity; retain resources for inspection")?;
846                    // An in-place swap still owns the checkout the source was
847                    // using, so the rollback retires it exactly as a resume
848                    // that recreated one.
849                    let error = self.rollback_failed_resume(&id, previous, operation.in_place,
850                        anyhow::anyhow!("destination restoration was interrupted; checkpoint retained for an explicit retry"), executor)?;
851                    Err(error)
852                }
853                MovePhase::ClosingSource => {
854                    if matches!(session.state, SessionState::Closing | SessionState::Destroying) {
855                        self.recover_move_source_stop(&mut operation, executor, manager).await?;
856                    }
857                    operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
858                    if self.state.sessions[&id].state == SessionState::Stopped && self.state.sessions[&id].target.is_some() {
859                        self.cleanup_stopped_target(&id, executor)?;
860                    }
861                    bail!("{}", interrupted_source_stop_message(
862                        "source stop was recovered; verified checkpoint retained. Retry move or Resume with previous settings",
863                        operation.in_place,
864                    ))
865                }
866                _ => bail!("Move requires an explicit retry after the daemon restarted"),
867            }
868        }.await;
869        self.finish_move_result(&mut operation, result, executor)
870    }
871
872    async fn recover_move_source_stop(
873        &mut self,
874        operation: &mut MoveOperation,
875        executor: &(impl CommandExecutor + Sync),
876        manager: &SessionManagerControl,
877    ) -> Result<()> {
878        let id = operation.selection.session_id.clone();
879        if self.state.sessions[&id].state == SessionState::Closing {
880            self.prepare_move_source_checkpoint(&id, executor, manager, operation)
881                .await?;
882            let handle = manager
883                .wait_for_session(&id, std::time::Duration::from_secs(5))
884                .await?;
885            let mut lease = handle.lease_connection().await?;
886            let execution = lease.connection_mut().sync().await?.operational.execution;
887            lease.release();
888            if matches!(
889                execution,
890                mj_core::relay::RelayExecutionState::Idle
891                    | mj_core::relay::RelayExecutionState::Running
892            ) {
893                if operation.cancellation_requested {
894                    let record = self.state.sessions.get_mut(&id).unwrap();
895                    record.state = SessionState::Running;
896                    record.updated_at = now();
897                    record.last_error =
898                        Some("Move was cancelled before the source was sealed".into());
899                    crate::database::save_lifecycle_session(record)?;
900                } else {
901                    // A recovered source stop always tears its target down: an
902                    // interrupted in-place swap falls back to the fresh path.
903                    self.suspend_session_for_move(
904                        &id,
905                        executor,
906                        manager,
907                        operation,
908                        None,
909                        SourceTargetDisposition::Destroy,
910                    )
911                    .await?;
912                }
913                return Ok(());
914            }
915        }
916        // A move's source stop was admitted by the move itself.
917        self.recover_interrupted_close_managed(&id, executor, manager, true, None)
918            .await?;
919        Ok(())
920    }
921
922    async fn execute_move(
923        &mut self,
924        operation: &mut MoveOperation,
925        preparation: Option<&MovePreparation>,
926        executor: &(impl CommandExecutor + Sync),
927        manager: &SessionManagerControl,
928    ) -> Result<()> {
929        let id = operation.selection.session_id.clone();
930        ensure!(
931            !executor.cancellation_requested(),
932            "move cancelled before source interruption"
933        );
934        if !operation.queue_admission_started {
935            if self.state.sessions[&id].state == SessionState::Error
936                && let Some(previous) = operation.recovery_session.as_ref()
937            {
938                // A prior failed teardown retains its exact partial checkout.
939                // Restore source identity only after that owner is stopped.
940                let failure = self.rollback_failed_resume(
941                    &id,
942                    previous,
943                    false,
944                    anyhow::anyhow!("clean up the partial Move destination before retry"),
945                    executor,
946                )?;
947                ensure!(
948                    self.state.sessions[&id].state == SessionState::Stopped,
949                    "{failure:#}"
950                );
951            }
952            let state = self.state.sessions[&id].state;
953            if matches!(state, SessionState::Closing | SessionState::Destroying) {
954                self.recover_move_source_stop(operation, executor, manager)
955                    .await?;
956                operation.checkpoint = self.state.sessions[&id].checkpoint.clone();
957            } else if matches!(state, SessionState::Running | SessionState::Disconnected)
958                && operation.destination_target.is_none()
959            {
960                executor.notify_notice("Stopping source");
961                let _timing = MovePhaseTimer::new(&id, "checkpoint and source stop");
962                operation.phase = MovePhase::ClosingSource;
963                operation.updated_at = now();
964                crate::database::save_move_operation(operation)?;
965                // An eligible move keeps its environment: the close stops at
966                // the sealed relay and the record keeps its target for
967                // `restore_session_in_place`.
968                let disposition = if operation.in_place {
969                    SourceTargetDisposition::RetainForInPlaceSwap
970                } else {
971                    SourceTargetDisposition::Destroy
972                };
973                self.suspend_session_for_move(
974                    &id,
975                    executor,
976                    manager,
977                    operation,
978                    preparation,
979                    disposition,
980                )
981                .await?;
982            }
983            // An in-place move never reaches `Stopped` with a target: its
984            // source stays `Closing` in the environment the destination reuses.
985            if !operation.in_place
986                && self.state.sessions[&id].state == SessionState::Stopped
987                && self.state.sessions[&id].target.is_some()
988            {
989                executor.notify_notice("Cleaning up source");
990                let _timing = MovePhaseTimer::new(&id, "source storage cleanup");
991                self.cleanup_stopped_target(&id, executor)?;
992            }
993            operation.checkpoint = operation
994                .checkpoint
995                .clone()
996                .or_else(|| self.state.sessions[&id].checkpoint.clone());
997            ensure!(
998                operation.checkpoint.is_some(),
999                "move has no verified checkpoint"
1000            );
1001            ensure!(
1002                !executor.cancellation_requested(),
1003                "move cancelled after source teardown; session is stopped"
1004            );
1005            operation.phase = MovePhase::ResumingDestination;
1006            operation.recovery_session = Some(self.state.sessions[&id].clone());
1007            operation.updated_at = now();
1008            crate::database::save_move_operation(operation)?;
1009            executor.notify_notice("Preparing destination");
1010            if operation.in_place {
1011                // The environment is kept, so there is nothing to provision and
1012                // no allocation to clear: eligibility already required the
1013                // allocation to be unchanged.
1014                self.restore_session_in_place(
1015                    &id,
1016                    operation.selection.profile_id.as_deref().unwrap(),
1017                    executor,
1018                )
1019                .await?;
1020            } else {
1021                if operation.selection.clear_resource_allocation {
1022                    let session = self.state.sessions.get_mut(&id).unwrap();
1023                    session.resource_allocation = None;
1024                    session.container_cpus = None;
1025                    session.container_memory = None;
1026                    crate::database::save_session(session)?;
1027                }
1028                self.resume_session_controlled(
1029                    &id,
1030                    operation.selection.profile_id.as_deref().unwrap(),
1031                    operation.selection.target_template_id.as_deref().unwrap(),
1032                    SessionResumeOptions {
1033                        additional_mounts: operation.selection.additional_mounts.clone(),
1034                        resource_allocation: operation.selection.resource_allocation.clone(),
1035                        discard_queue: true,
1036                    },
1037                    executor,
1038                )
1039                .await?;
1040            }
1041            let destination = &self.state.sessions[&id];
1042            operation.destination_target = destination.target.clone();
1043            operation.destination_native_session_id = destination.native_session_id.clone();
1044            // Persist readiness before submitting even the first queued command.
1045            operation.phase = MovePhase::StartingQueue;
1046            operation.queue_admission_started = true;
1047            operation.updated_at = now();
1048            crate::database::save_move_operation(operation)?;
1049        }
1050        restore_move_queue_hold(operation);
1051        self.admit_move_queue(operation, executor).await
1052    }
1053
1054    async fn admit_move_queue(
1055        &self,
1056        operation: &mut MoveOperation,
1057        executor: &(impl CommandExecutor + Sync),
1058    ) -> Result<()> {
1059        let timing_id = operation.selection.session_id.clone();
1060        let _timing = MovePhaseTimer::new(&timing_id, "queue admission");
1061        let id = &operation.selection.session_id;
1062        let mut relay = {
1063            let _checking_destination =
1064                ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1065            let destination = &self.state.sessions[id];
1066            ensure!(
1067                destination.state == SessionState::Running
1068                    && destination.target == operation.destination_target
1069                    && destination.native_session_id == operation.destination_native_session_id,
1070                "cannot prove the same ready destination; refusing to replay potentially executed work"
1071            );
1072            let spec = self.reconnect_command(id)?;
1073            let relay = StandaloneSession::connect_command(&spec, id).await?;
1074            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")?;
1075            if let Some(expected) = &operation.destination_store_id {
1076                ensure!(
1077                    *expected == store_id,
1078                    "destination relay storage was replaced; refusing to replay potentially executed work"
1079                );
1080            } else {
1081                // No command can be admitted until this identity is durable.
1082                operation.destination_store_id = Some(store_id);
1083                crate::database::save_move_operation(operation)?;
1084            }
1085            ensure!(
1086                relay.snapshot().operational.native_session_id
1087                    == operation.destination_native_session_id,
1088                "destination relay native identity changed; refusing queue replay"
1089            );
1090            relay
1091        };
1092        if operation.queue == ResumeQueueDisposition::Start && !operation.queue_admission_finished {
1093            let _starting_queue = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
1094            executor.notify_notice("Starting queued work");
1095            let checkpoint = operation
1096                .checkpoint
1097                .as_ref()
1098                .context("move queue archive is missing")?;
1099            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1100            ensure!(
1101                verified.archive_sha256 == checkpoint.sha256 && verified.manifest.session.id == *id,
1102                "move queue checkpoint verification failed"
1103            );
1104            for queued in verified.canonical_session.queued_prompts {
1105                ensure!(
1106                    !executor.cancellation_requested(),
1107                    "move cancelled during queue admission; destination retained"
1108                );
1109                let command = match queued.kind {
1110                    CanonicalQueuedCommandKind::Prompt => RelayCommand::Prompt {
1111                        prompt: queued
1112                            .content
1113                            .into_iter()
1114                            .map(serde_json::from_value)
1115                            .collect::<serde_json::Result<_>>()?,
1116                    },
1117                    CanonicalQueuedCommandKind::SetConfig { key, value } => {
1118                        RelayCommand::SetConfig { key, value }
1119                    }
1120                };
1121                relay.submit(queued.command_id, command).await?;
1122            }
1123        }
1124        operation.queue_admission_finished = true;
1125        crate::database::save_move_operation(operation)?;
1126        restore_move_queue_hold(operation);
1127        let queue_sentence = if operation.queue == ResumeQueueDisposition::Discard {
1128            "Queued work was discarded; ready and idle."
1129        } else {
1130            "Queued work was accepted."
1131        };
1132        let source_profile = &operation.source_profile_id;
1133        let source_target = &operation.source_target_template_id;
1134        let destination_profile = operation.selection.profile_id.as_deref().unwrap();
1135        let destination_target = operation.selection.target_template_id.as_deref().unwrap();
1136        // The two moves are different events for the person reading the
1137        // conversation: one rebuilt the environment, the other kept it.
1138        let text = if operation.in_place {
1139            format!(
1140                "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."
1141            )
1142        } else {
1143            format!(
1144                "Moved from {source_profile} / {source_target} to {destination_profile} / {destination_target} in a fresh environment. {queue_sentence} The interrupted prompt was not replayed."
1145            )
1146        };
1147        relay
1148            .submit(
1149                format!("{}-notice", operation.operation_id),
1150                RelayCommand::RecordNotice { text },
1151            )
1152            .await?;
1153        Ok(())
1154    }
1155
1156    pub(super) fn validate_move_checkpoint(
1157        &self,
1158        operation: &MoveOperation,
1159        preparation: Option<&MovePreparation>,
1160        executor: &(impl CommandExecutor + Sync),
1161    ) -> Result<()> {
1162        let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1163        let current = Controller {
1164            config: mj_core::config::Config::load()?,
1165            state: self.state.clone(),
1166        };
1167        ensure!(
1168            current.move_configuration_fingerprint(&operation.selection)?
1169                == operation.configuration_fingerprint,
1170            "destination configuration changed during move"
1171        );
1172        let id = &operation.selection.session_id;
1173        current.validate_move_destination_paths(
1174            &self.state.sessions[id],
1175            operation.selection.target_template_id.as_deref().unwrap(),
1176            executor,
1177        )?;
1178        if let super::ResumeRepositorySourcePreflight::RepositoryMoved(mismatch) = self
1179            .preflight_resume_repository_sources(
1180                id,
1181                operation.selection.target_template_id.as_deref().unwrap(),
1182                executor,
1183            )?
1184        {
1185            bail!(
1186                "destination repository source is missing checkpoint commit {}; source retained",
1187                mismatch.missing_commit
1188            );
1189        }
1190        if let Some(prepared) = preparation {
1191            let checkpoint = self.state.sessions[id]
1192                .checkpoint
1193                .as_ref()
1194                .context("no move checkpoint")?;
1195            let verified = verify_archive_streaming(&checkpoint.archive_path)?;
1196            let actual: Vec<_> = verified
1197                .canonical_session
1198                .queued_prompts
1199                .iter()
1200                .map(|p| p.command_id.as_str())
1201                .collect();
1202            let expected: Vec<_> = prepared
1203                .queued_commands
1204                .iter()
1205                .map(|p| p.command_id.as_str())
1206                .collect();
1207            ensure!(
1208                actual == expected,
1209                "pending queue changed before checkpoint capture; source retained, confirm Move again"
1210            );
1211        }
1212        Ok(())
1213    }
1214}
1215
1216fn validate_preserved_configuration(
1217    profile_id: &str,
1218    accepted: &mj_core::acp::AcceptedSessionConfig,
1219    choices: &mj_core::worker_launch::ProfileConfig,
1220) -> Result<()> {
1221    for (key, value, offered) in [
1222        ("model", accepted.model.as_deref(), &choices.models),
1223        ("effort", accepted.effort.as_deref(), &choices.efforts),
1224    ] {
1225        let Some(value) = value else { continue };
1226        ensure!(
1227            offered.iter().any(|choice| choice.value == value),
1228            "destination profile {profile_id:?} does not offer the session's accepted {key} {value:?}; choices: {}",
1229            offered
1230                .iter()
1231                .map(|choice| choice.value.as_str())
1232                .collect::<Vec<_>>()
1233                .join(", ")
1234        );
1235    }
1236    Ok(())
1237}
1238
1239/// Whether this move can replace only the harness inside the source target.
1240///
1241/// The environment is kept only when nothing about it changes: the same target
1242/// template, the same attached mounts, and the same resource allocation. A
1243/// retry starts from a torn-down or unknown source, and a sub-agent never owns
1244/// its own target, so both take the full fresh-environment path.
1245pub(super) fn in_place_move_eligible(
1246    source: &mj_core::state::SessionRecord,
1247    selection: &MoveSelection,
1248    is_subagent: bool,
1249    retry: bool,
1250) -> bool {
1251    !retry
1252        && !is_subagent
1253        && source.target.is_some()
1254        && matches!(
1255            source.state,
1256            SessionState::Running | SessionState::Disconnected
1257        )
1258        && Some(&source.target_template_id) == selection.target_template_id.as_ref()
1259        && Some(&source.additional_mounts) == selection.additional_mounts.as_ref()
1260        && source.resource_allocation == selection.resource_allocation
1261        && !selection.clear_resource_allocation
1262}
1263
1264fn outcome(
1265    operation_id: &str,
1266    selection: &MoveSelection,
1267    status: &str,
1268    error: Option<String>,
1269    recovery: Option<String>,
1270) -> MoveOutcome {
1271    MoveOutcome {
1272        operation_id: operation_id.into(),
1273        session_id: selection.session_id.clone(),
1274        profile_id: selection.profile_id.clone().unwrap_or_default(),
1275        target_template_id: selection.target_template_id.clone().unwrap_or_default(),
1276        outcome: status.into(),
1277        error,
1278        recovery,
1279    }
1280}