Skip to main content

mj_controller/controller/move_session/
transfer.rs

1//! Move-only workspace transport. It never reads or writes a checkpoint payload.
2use super::*;
3use crate::targets::{self, CommandSpec};
4use mj_core::move_workspace::*;
5use std::path::{Path, PathBuf};
6
7fn repositories(
8    layout: &super::super::checkpoint::SessionExportLayout,
9) -> Vec<WorkspaceRepository> {
10    layout
11        .repositories
12        .iter()
13        .map(|repo| WorkspaceRepository {
14            id: repo.id.clone(),
15            root: Path::new(&layout.workspace_root).join(&repo.relative_destination),
16        })
17        .collect()
18}
19
20fn utility(
21    executor: &(impl CommandExecutor + Sync),
22    backend: &targets::TargetLocator,
23    id: &str,
24    command: &WorkspaceCommand,
25) -> Result<serde_json::Value> {
26    let binary = super::super::worker_binary::worker_binary_for(backend, executor)?;
27    let executable = if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
28        binary.to_string_lossy().into_owned()
29    } else {
30        let root = targets::worker_root(backend, id)?;
31        let helper = format!(
32            "{root}/move-helper-{}",
33            mj_core::worker_build::WORKER_BUILD_ID
34        );
35        let probe = executor.execute(&targets::command_on_locator(
36            backend,
37            id,
38            vec!["test".into(), "-x".into(), helper.clone()],
39            "check Move helper",
40        )?)?;
41        if probe.status != 0 {
42            install_helper(executor, backend, id, &binary, &helper)?;
43        }
44        helper
45    };
46    let command_spec = targets::command_on_locator(
47        backend,
48        id,
49        vec![executable, "worker".into(), "move-workspace".into()],
50        "prepare or verify Move workspace",
51    )?;
52    let output = executor.execute_with_stdin(
53        &command_spec,
54        &mut std::io::Cursor::new(serde_json::to_vec(command)?),
55    )?;
56    ensure!(
57        output.status == 0,
58        "Move workspace operation failed: {}",
59        String::from_utf8_lossy(&output.stderr)
60    );
61    serde_json::from_slice(&output.stdout).context("read Move workspace result")
62}
63
64fn install_helper(
65    executor: &(impl CommandExecutor + Sync),
66    backend: &targets::TargetLocator,
67    id: &str,
68    binary: &Path,
69    helper: &str,
70) -> Result<()> {
71    // Each uploader owns its temporary inode. Rename publishes only a complete
72    // immutable helper; neither checkpoint staging nor a running binary is overwritten.
73    let script = r#"set -eu
74temporary=$(mktemp "$1.XXXXXX")
75trap 'rm -f -- "$temporary"' EXIT
76cat > "$temporary"
77chmod 700 "$temporary"
78mv -f -- "$temporary" "$1"
79"#;
80    let command = targets::command_on_locator(
81        backend,
82        id,
83        vec![
84            "sh".into(),
85            "-c".into(),
86            script.into(),
87            "install-move-helper".into(),
88            helper.into(),
89        ],
90        "install immutable Move helper",
91    )?;
92    let output = executor.execute_with_stdin(&command, &mut std::fs::File::open(binary)?)?;
93    ensure!(
94        output.status == 0,
95        "Move helper installation failed: {}",
96        String::from_utf8_lossy(&output.stderr)
97    );
98    Ok(())
99}
100
101/// Options for every Move rsync. `--protect-args` sends paths inside the rsync
102/// protocol, so the namespace wrapper passes paths with spaces unchanged. The
103/// preflight probes these options, so an rsync that rejects one blocks Move.
104const RSYNC_TRANSFER_OPTIONS: [&str; 6] = [
105    "--recursive",
106    "--times",
107    "--perms",
108    "--protect-args",
109    "--partial",
110    "--partial-dir=.move-partial",
111];
112
113/// Why an rsync failed the probe. macOS ships openrsync, or rsync 2.6.9 on
114/// older releases, and neither accepts `--protect-args`.
115const RSYNC_REQUIREMENT: &str = "Move needs rsync 3.0 or newer, which accepts --protect-args; \
116the rsync that ships with macOS does not (install one with `brew install rsync`)";
117
118/// Rsync arguments that succeed only if rsync accepts every transfer option.
119/// `--version` makes rsync exit once it has parsed them.
120fn rsync_probe_arguments() -> impl Iterator<Item = String> {
121    RSYNC_TRANSFER_OPTIONS
122        .into_iter()
123        .chain(["--version"])
124        .map(String::from)
125}
126
127fn rsync_probe() -> Vec<String> {
128    std::iter::once("rsync".into())
129        .chain(rsync_probe_arguments())
130        .collect()
131}
132
133/// Rsync's remote shell is the existing locator command, including its SSH
134/// session lease and container namespace. The remote endpoint sees only Move
135/// staging. Partial files belong to this operation and survive retry.
136fn copy_workspace(
137    executor: &(impl CommandExecutor + Sync),
138    backend: &targets::TargetLocator,
139    remote: &Path,
140    local: &Path,
141    upload: bool,
142) -> Result<()> {
143    std::fs::create_dir_all(local)?;
144    let source = format!("{}/", if upload { local } else { remote }.display());
145    let destination = format!("{}/", if upload { remote } else { local }.display());
146    let mut args: Vec<String> = RSYNC_TRANSFER_OPTIONS.map(String::from).into();
147    let base = targets::locator_command(backend, vec!["rsync".into()]);
148    let lease = base.open_ssh_session(executor)?;
149    let wrapper = if matches!(backend, targets::TargetLocator::LocalBare { .. }) {
150        args.extend([source, destination]);
151        None
152    } else {
153        let directory = rsync_shell(lease.command())?;
154        let path = directory.path().join("move-rsync-shell");
155        args.extend(["--rsh".into(), path.to_string_lossy().into_owned()]);
156        args.extend(if upload {
157            [source, format!("move:{destination}")]
158        } else {
159            [format!("move:{source}"), destination]
160        });
161        Some(directory)
162    };
163    let result = super::super::execute_checked(
164        executor,
165        CommandSpec::new("rsync", args).purpose("transfer Move workspace with resumable files"),
166    );
167    drop(wrapper);
168    drop(lease);
169    result.map(|_| ())
170}
171
172fn rsync_shell(command: &CommandSpec) -> Result<tempfile::TempDir> {
173    let directory = tempfile::tempdir()?;
174    let path = directory.path().join("move-rsync-shell");
175    let quoted = std::iter::once(command.program.as_str())
176        .chain(command.args.iter().map(String::as_str))
177        .map(targets::posix_quote)
178        .collect::<Vec<_>>()
179        .join(" ");
180    std::fs::write(
181        &path,
182        format!("#!/bin/sh\nset -eu\nshift\nshift\nexec {quoted} \"$@\"\n"),
183    )?;
184    #[cfg(unix)]
185    {
186        use std::os::unix::fs::PermissionsExt;
187        std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o700))?;
188    }
189    Ok(directory)
190}
191
192fn storage_probe(
193    executor: &(impl CommandExecutor + Sync),
194    command: CommandSpec,
195    location: &str,
196    allocation: &str,
197    copies: u64,
198    assessment: &mut WorkspaceAssessment,
199) {
200    let result = executor.execute(&command).and_then(|output| {
201        ensure!(
202            output.status == 0,
203            "{}",
204            String::from_utf8_lossy(&output.stderr)
205        );
206        let report = String::from_utf8_lossy(&output.stdout);
207        let available_bytes = super::super::mbx::available_bytes(&report)
208            .map(|blocks| blocks.saturating_mul(1024))
209            .context("cannot read available disk space")?;
210        let device = report
211            .lines()
212            .nth(1)
213            .and_then(|line| line.split_whitespace().next())
214            .context("cannot read storage identity")?;
215        Ok((
216            format!(
217                "{}:{device}",
218                command.ssh_destination.as_deref().unwrap_or("local")
219            ),
220            available_bytes,
221        ))
222    });
223    match result {
224        Ok((filesystem_key, available_bytes)) => {
225            if let Some(storage) = assessment
226                .storage
227                .iter_mut()
228                .find(|storage| storage.filesystem_key == filesystem_key)
229            {
230                storage.location.push_str(&format!(" + {location}"));
231                if storage.allocations.insert(allocation.into()) {
232                    storage.copies = storage.copies.saturating_add(copies);
233                }
234                storage.available_bytes = storage.available_bytes.min(available_bytes);
235            } else {
236                assessment.storage.push(WorkspaceStorage {
237                    filesystem_key,
238                    allocations: [allocation.into()].into_iter().collect(),
239                    location: location.into(),
240                    available_bytes,
241                    copies,
242                });
243            }
244        }
245        Err(error) => assessment
246            .blockers
247            .push(format!("Check {location} storage: {error:#}")),
248    }
249}
250
251fn at_host(ssh: Option<&targets::SshTarget>, mut arguments: Vec<String>) -> CommandSpec {
252    if let Some(ssh) = ssh {
253        targets::ssh_command_owned(ssh, arguments)
254    } else {
255        let program = arguments.remove(0);
256        CommandSpec::new(program, arguments)
257    }
258}
259
260fn container_storage_arguments(engine: &str) -> Option<Vec<String>> {
261    let format = match engine {
262        "podman" => "{{.Store.GraphRoot}}",
263        "docker" => "{{.DockerRootDir}}",
264        _ => return None,
265    };
266    Some(vec![
267        "sh".into(),
268        "-c".into(),
269        "set -eu; p=$(\"$1\" info --format \"$2\"); test -n \"$p\"; df -Pk -- \"$p\"".into(),
270        "move-container-storage".into(),
271        engine.into(),
272        format.into(),
273    ])
274}
275
276fn storage_arguments(path: &Path) -> Vec<String> {
277    vec![
278        "sh".into(),
279        "-c".into(),
280        "set -eu; p=$1; while [ ! -e \"$p\" ]; do p=$(dirname -- \"$p\"); done; df -Pk -- \"$p\""
281            .into(),
282        "move-storage".into(),
283        path.to_string_lossy().into_owned(),
284    ]
285}
286
287impl Controller {
288    pub fn cleanup_retained_move_source(
289        &self,
290        session_id: &str,
291        operation_id: &str,
292        executor: &(impl CommandExecutor + Sync),
293    ) -> Result<()> {
294        crate::worker_lifecycle::run_blocking(
295            session_id,
296            "cleanup retained move source",
297            executor,
298            || {
299                let retained = crate::database::retained_move_sources(session_id)?
300                    .into_iter()
301                    .find(|source| source.operation_id == operation_id)
302                    .context("retained Move source is missing")?;
303                let source = &retained.source;
304                ensure!(
305                    !self
306                        .state
307                        .sessions
308                        .values()
309                        .any(|session| session.target.is_some() && session.target == source.target),
310                    "retained source is still referenced by an active session"
311                );
312                if let Some(locator) = &source.target {
313                    let backend =
314                        super::super::backend::backend_locator(locator, source, &self.config)?;
315                    targets::retire_move_target_plan(&backend, session_id)?.execute(executor)?;
316                    self.release_retired_move_target(
317                        source,
318                        &backend,
319                        self.state.sessions.get(session_id),
320                        executor,
321                    );
322                }
323                if let mj_core::state::Checkout::ManagedWorktree {
324                    worktree: checkout, ..
325                } = source.checkout()
326                {
327                    ensure!(
328                        !self.state.sessions.keys().any(|session_id| {
329                            matches!(
330                                self.state.checkout(session_id),
331                                Ok(mj_core::state::Checkout::ManagedWorktree {
332                                    worktree: other,
333                                    ..
334                                }) if other == checkout
335                            )
336                        }),
337                        "retained checkout is still referenced by a session"
338                    );
339                    super::super::worktree::retire_managed_worktree(executor, checkout)?;
340                }
341                crate::database::forget_retained_move_source(operation_id)
342            },
343        )
344    }
345
346    /// Release the mbx build state of a target a Move has just retired,
347    /// keeping whatever the placement the session goes on running in,
348    /// `live`, can still build in. When that placement cannot be resolved,
349    /// nothing is released: a stale target costs disk until mbx ages it out,
350    /// while a wrong release costs a live session its build output.
351    fn release_retired_move_target(
352        &self,
353        retired: &mj_core::state::SessionRecord,
354        backend: &targets::TargetLocator,
355        live: Option<&mj_core::state::SessionRecord>,
356        executor: &impl CommandExecutor,
357    ) {
358        let Some(release) = super::super::mbx::release::BuildStateRelease::for_target(
359            retired,
360            backend,
361            &self.config,
362        ) else {
363            return;
364        };
365        let live_backend = match live
366            .and_then(|live| live.target.as_ref().map(|target| (live, target)))
367        {
368            None => None,
369            Some((live, target)) => {
370                match super::super::backend::backend_locator(target, live, &self.config) {
371                    Ok(backend) => Some((live, backend)),
372                    Err(error) => {
373                        tracing::warn!(
374                            session_id = %retired.id,
375                            error = format!("{error:#}"),
376                            "the Move's live target could not be resolved, so the retired target's build state is kept"
377                        );
378                        return;
379                    }
380                }
381            }
382        };
383        if let Some(release) = release.excluding(
384            live_backend
385                .as_ref()
386                .map(|(live, backend)| (*live, backend)),
387        ) {
388            release.run(executor);
389        }
390    }
391
392    pub(in crate::controller) fn move_destination_bundle(
393        &self,
394        id: &str,
395    ) -> Result<Option<targets::ProjectBundleSpec>> {
396        let Some(operation) = crate::database::load_move_operation(id)? else {
397            return Ok(None);
398        };
399        if operation.workspace_transfer.is_none()
400            || operation.phase != MovePhase::ResumingDestination
401        {
402            return Ok(None);
403        }
404        let handoff = operation.handoff.context("Move handoff missing")?;
405        let archive = super::super::resume::verify_resume_checkpoint(id, &handoff)?;
406        let primary = archive
407            .manifest
408            .repositories
409            .iter()
410            .find(|repo| repo.metadata.id == archive.manifest.bundle.primary_repository)
411            .context("Move primary repository missing")?
412            .metadata
413            .relative_destination
414            .to_string_lossy()
415            .into_owned();
416        let repositories = archive
417            .manifest
418            .repositories
419            .iter()
420            .map(|repo| {
421                let metadata = &repo.metadata;
422                mj_core::remote_git::validate_network_url(&metadata.origin)?;
423                Ok(targets::RepositorySpec {
424                    url: Some(metadata.origin.clone()),
425                    push_urls: metadata.push_urls.clone(),
426                    destination: metadata.relative_destination.to_string_lossy().into_owned(),
427                    git_ref: None,
428                    reference: None,
429                })
430            })
431            .collect::<Result<Vec<_>>>()?;
432        Ok(Some(targets::ProjectBundleSpec {
433            primary,
434            repositories,
435        }))
436    }
437
438    pub(super) fn assess_move_workspace(
439        &self,
440        selection: &MoveSelection,
441        executor: &(impl CommandExecutor + Sync),
442    ) -> Result<WorkspaceAssessment> {
443        let id = &selection.session_id;
444        if let Some(operation) = crate::database::load_move_operation(id)?
445            && operation.holds_source_environment()
446            && let Some(transfer) = operation.workspace_transfer
447        {
448            return Ok(transfer.assessment);
449        }
450        let layout = self.session_export_layout(id, executor)?;
451        let mut assessment: WorkspaceAssessment = serde_json::from_value(utility(
452            executor,
453            &layout.backend,
454            id,
455            &WorkspaceCommand::Inspect {
456                repositories: repositories(&layout),
457            },
458        )?)?;
459        for command in [
460            CommandSpec::new("rsync", rsync_probe_arguments())
461                .purpose("check controller Move transport"),
462            targets::command_on_locator(
463                &layout.backend,
464                id,
465                rsync_probe(),
466                "check source Move transport",
467            )?,
468        ] {
469            match executor.execute(&command) {
470                Ok(output) if output.status == 0 => {}
471                Ok(output) => assessment.blockers.push(format!(
472                    "{}: {RSYNC_REQUIREMENT}: {}",
473                    command.purpose,
474                    String::from_utf8_lossy(&output.stderr)
475                )),
476                Err(error) => assessment
477                    .blockers
478                    .push(format!("{}: {error:#}", command.purpose)),
479            }
480        }
481        let source_root = targets::worker_root(&layout.backend, id)?;
482        storage_probe(
483            executor,
484            targets::command_on_locator(
485                &layout.backend,
486                id,
487                storage_arguments(Path::new(&source_root)),
488                "check source Move staging space",
489            )?,
490            "Source",
491            "source",
492            1,
493            &mut assessment,
494        );
495        if let Some(arguments) = layout
496            .backend
497            .container_engine()
498            .and_then(container_storage_arguments)
499        {
500            let host = targets::locator_command(&layout.backend, vec!["true".into()]).ssh_session;
501            storage_probe(
502                executor,
503                at_host(host.as_ref(), arguments).purpose("check source container backing storage"),
504                "Source backing storage",
505                "source",
506                1,
507                &mut assessment,
508            );
509        }
510        let mut local_args = storage_arguments(&mj_core::config::data_dir());
511        let program = local_args.remove(0);
512        storage_probe(
513            executor,
514            CommandSpec::new(program, local_args).purpose("check controller Move staging space"),
515            "Controller",
516            "controller",
517            1,
518            &mut assessment,
519        );
520        let target = &self.config.targets[selection
521            .target_template_id
522            .as_ref()
523            .context("Move target missing")?];
524        if matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. }) {
525            return Ok(assessment);
526        }
527        if matches!(target, mj_core::config::TargetTemplate::LocalBare)
528            && matches!(layout.backend, targets::TargetLocator::LocalBare { .. })
529        {
530            assessment.blockers.push("These local targets share the source worker storage, so no files can be transferred between them.".into());
531        }
532        if let mj_core::config::TargetTemplate::SshBare { ssh, .. } = target
533            && let targets::TargetLocator::SshBare {
534                ssh: source_ssh, ..
535            } = &layout.backend
536            && *source_ssh == targets::SshTarget::from(ssh)
537        {
538            assessment.blockers.push("These SSH targets share the source worker storage, so no files can be transferred between them.".into());
539        }
540        let (destination_ssh, destination_engine, destination_image) = match target {
541            mj_core::config::TargetTemplate::LocalPodman { container } => {
542                (None, Some("podman"), Some(container.image.as_str()))
543            }
544            mj_core::config::TargetTemplate::LocalDocker { container } => {
545                (None, Some("docker"), Some(container.image.as_str()))
546            }
547            mj_core::config::TargetTemplate::AppleContainer { container } => {
548                (None, Some("container"), Some(container.image.as_str()))
549            }
550            mj_core::config::TargetTemplate::SshPodman { ssh, container } => {
551                (Some(ssh), Some("podman"), Some(container.image.as_str()))
552            }
553            mj_core::config::TargetTemplate::SshDocker { ssh, container } => {
554                (Some(ssh), Some("docker"), Some(container.image.as_str()))
555            }
556            mj_core::config::TargetTemplate::SshBare { ssh, .. } => (Some(ssh), None, None),
557            _ => (None, None, None),
558        };
559        let mut arguments =
560            if let (Some(engine), Some(image)) = (destination_engine, destination_image) {
561                let mut arguments = vec![
562                    engine.into(),
563                    "run".into(),
564                    "--rm".into(),
565                    "--entrypoint".into(),
566                    "rsync".into(),
567                    image.into(),
568                ];
569                arguments.extend(rsync_probe_arguments());
570                arguments
571            } else {
572                rsync_probe()
573            };
574        let command = if let Some(ssh) = destination_ssh {
575            targets::ssh_command_owned(&targets::SshTarget::from(ssh), arguments)
576        } else {
577            let program = arguments.remove(0);
578            CommandSpec::new(program, arguments)
579        };
580        match executor.execute(&command.purpose("check destination Move transport")) {
581            Ok(output) if output.status == 0 => {},
582            Ok(output) => assessment.blockers.push(format!("Destination: {RSYNC_REQUIREMENT}. Update the target image or install rsync there: {}", String::from_utf8_lossy(&output.stderr))),
583            Err(error) => assessment.blockers.push(format!("Destination transport check failed: {error:#}")),
584        }
585        let destination_ssh = destination_ssh.map(targets::SshTarget::from);
586        if let mj_core::config::TargetTemplate::LocalPodman { container }
587        | mj_core::config::TargetTemplate::SshPodman { container, .. } = target
588            && let mj_core::config::PodmanWorkspaceStorage::HostHelper { root, .. } =
589                &container.workspace_storage
590        {
591            storage_probe(
592                executor,
593                at_host(destination_ssh.as_ref(), storage_arguments(Path::new(root)))
594                    .purpose("check destination workspace storage"),
595                "Destination workspace",
596                "destination",
597                2,
598                &mut assessment,
599            );
600        }
601        if let Some(engine) = destination_engine {
602            let image = destination_image.context("Move image missing")?;
603            let arguments = vec![
604                engine.into(),
605                "run".into(),
606                "--rm".into(),
607                "--entrypoint".into(),
608                "sh".into(),
609                image.into(),
610                "-c".into(),
611                "df -Pk /".into(),
612            ];
613            storage_probe(
614                executor,
615                at_host(destination_ssh.as_ref(), arguments)
616                    .purpose("check destination container space"),
617                "Destination container",
618                "destination",
619                2,
620                &mut assessment,
621            );
622            if let Some(arguments) = container_storage_arguments(engine) {
623                storage_probe(
624                    executor,
625                    at_host(destination_ssh.as_ref(), arguments)
626                        .purpose("check destination container backing storage"),
627                    "Destination backing storage",
628                    "destination",
629                    2,
630                    &mut assessment,
631                );
632            }
633        } else if !matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. }) {
634            let path = match target {
635                mj_core::config::TargetTemplate::SshBare {
636                    workspace_prefix, ..
637                } => workspace_prefix.clone(),
638                _ => mj_core::config::data_dir(),
639            };
640            storage_probe(
641                executor,
642                at_host(destination_ssh.as_ref(), storage_arguments(&path))
643                    .purpose("check destination Move staging and workspace space"),
644                "Destination",
645                "destination",
646                2,
647                &mut assessment,
648            );
649        }
650        Ok(assessment)
651    }
652
653    pub(super) fn assess_prepared_destination(
654        &self,
655        operation: &MoveOperation,
656        assessment: &mut WorkspaceAssessment,
657        executor: &(impl CommandExecutor + Sync),
658    ) -> Result<()> {
659        let destination = operation
660            .prepared_destination
661            .as_ref()
662            .context("EC2 Move destination missing")?;
663        let backend = super::destination::prepared_backend(
664            destination
665                .target()
666                .context("EC2 Move target not checked")?,
667            &destination.runtime,
668            &operation.selection.session_id,
669        )?;
670        let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
671        super::super::execute_checked(
672            executor,
673            targets::locator_command(&backend, rsync_probe())
674                .purpose("check prepared EC2 Move transport"),
675        )?;
676        let targets::TargetLocator::AwsEc2 { workspace, .. } = &backend else {
677            bail!("prepared target is not EC2");
678        };
679        storage_probe(
680            executor,
681            targets::locator_command(&backend, storage_arguments(Path::new(workspace)))
682                .purpose("check prepared EC2 Move staging and workspace space"),
683            "Destination",
684            "destination",
685            2,
686            assessment,
687        );
688        operation.selection.workspace.validate(assessment)
689    }
690
691    pub(super) fn new_workspace_transfer(
692        &self,
693        id: &str,
694        operation_id: &str,
695        assessment: WorkspaceAssessment,
696        executor: &(impl CommandExecutor + Sync),
697    ) -> Result<WorkspaceTransfer> {
698        let layout = self.session_export_layout(id, executor)?;
699        let root = targets::worker_root(&layout.backend, id)?;
700        Ok(WorkspaceTransfer {
701            assessment,
702            source: Box::new(self.state.sessions[id].clone()),
703            source_stage: Path::new(&root).join(format!("move-{operation_id}")),
704            controller_stage: mj_core::config::data_dir()
705                .join("moves")
706                .join(operation_id)
707                .join("workspace"),
708            phase: WorkspaceTransferPhase::Planned,
709        })
710    }
711
712    pub(super) fn capture_move_workspace(
713        &self,
714        operation: &mut MoveOperation,
715        executor: &(impl CommandExecutor + Sync),
716    ) -> Result<()> {
717        executor.begin_resumable_move_work()?;
718        let result = self.capture_move_workspace_inner(operation, executor);
719        executor.end_resumable_move_work()?;
720        result
721    }
722
723    fn capture_move_workspace_inner(
724        &self,
725        operation: &mut MoveOperation,
726        executor: &(impl CommandExecutor + Sync),
727    ) -> Result<()> {
728        let id = &operation.selection.session_id;
729        let transfer = operation
730            .workspace_transfer
731            .as_ref()
732            .context("Move transfer missing")?
733            .clone();
734        let layout = self.session_export_layout(id, executor)?;
735        if transfer.phase == WorkspaceTransferPhase::Planned {
736            // The target utility holds the stage's file lock before reusing
737            // or replacing an incomplete capture, including after a lost ACK.
738            utility(
739                executor,
740                &layout.backend,
741                id,
742                &WorkspaceCommand::Capture {
743                    repositories: repositories(&layout),
744                    selection: operation.selection.workspace.clone(),
745                    destination: transfer.source_stage.clone(),
746                },
747            )?;
748            operation.workspace_transfer.as_mut().unwrap().phase = WorkspaceTransferPhase::Captured;
749            crate::database::save_move_operation(operation)?;
750        }
751        if operation.workspace_transfer.as_ref().unwrap().phase == WorkspaceTransferPhase::Captured
752        {
753            executor.notify_notice("Transferring workspace to temporary Move storage");
754            copy_workspace(
755                executor,
756                &layout.backend,
757                &transfer.source_stage,
758                &transfer.controller_stage,
759                false,
760            )?;
761            // Verification uses the same installed utility on the controller.
762            let local = targets::TargetLocator::LocalBare {
763                worker_root: mj_core::config::data_dir()
764                    .join("workers")
765                    .join(id)
766                    .to_string_lossy()
767                    .into_owned(),
768            };
769            utility(
770                executor,
771                &local,
772                id,
773                &WorkspaceCommand::Verify {
774                    source: transfer.controller_stage.clone(),
775                },
776            )?;
777            operation.workspace_transfer.as_mut().unwrap().phase =
778                WorkspaceTransferPhase::Downloaded;
779            crate::database::save_move_operation(operation)?;
780        }
781        if operation.workspace_transfer.as_ref().unwrap().phase
782            == WorkspaceTransferPhase::Downloaded
783        {
784            let root = targets::worker_root(&layout.backend, id)?;
785            super::super::execute_checked(
786                executor,
787                targets::command_on_locator(
788                    &layout.backend,
789                    id,
790                    vec![
791                        "sh".into(),
792                        "-c".into(),
793                        targets::stop_worker_daemon_script(&root),
794                    ],
795                    "stop sealed Move source worker, retaining workspace",
796                )?,
797            )?;
798            operation.workspace_transfer.as_mut().unwrap().phase =
799                WorkspaceTransferPhase::SourceStopped;
800            crate::database::save_move_operation(operation)?;
801        }
802        Ok(())
803    }
804
805    pub(in crate::controller) fn restore_move_workspace(
806        &self,
807        id: &str,
808        executor: &(impl CommandExecutor + Sync),
809    ) -> Result<()> {
810        executor.begin_resumable_move_work()?;
811        let result = self.restore_move_workspace_inner(id, executor);
812        executor.end_resumable_move_work()?;
813        result
814    }
815
816    fn restore_move_workspace_inner(
817        &self,
818        id: &str,
819        executor: &(impl CommandExecutor + Sync),
820    ) -> Result<()> {
821        let mut operation =
822            crate::database::load_move_operation(id)?.context("Move intent missing")?;
823        let transfer = operation
824            .workspace_transfer
825            .as_ref()
826            .context("Move transfer missing")?
827            .clone();
828        let layout = self.session_export_layout(id, executor)?;
829        ensure!(
830            self.state.sessions[id].target != transfer.source.target,
831            "Move destination reuses the source environment; refusing to overwrite it"
832        );
833        let stage = PathBuf::from(targets::worker_root(&layout.backend, id)?)
834            .join(format!("move-{}", operation.operation_id));
835        super::super::execute_checked(
836            executor,
837            targets::command_on_locator(
838                &layout.backend,
839                id,
840                vec![
841                    "mkdir".into(),
842                    "-p".into(),
843                    stage.to_string_lossy().into_owned(),
844                ],
845                "create destination Move staging",
846            )?,
847        )?;
848        copy_workspace(
849            executor,
850            &layout.backend,
851            &stage,
852            &transfer.controller_stage,
853            true,
854        )?;
855        utility(
856            executor,
857            &layout.backend,
858            id,
859            &WorkspaceCommand::Restore {
860                source: stage,
861                repositories: repositories(&layout),
862            },
863        )?;
864        operation.workspace_transfer.as_mut().unwrap().phase = WorkspaceTransferPhase::Restored;
865        crate::database::save_move_operation(&operation)?;
866        Ok(())
867    }
868
869    pub(super) fn finish_workspace_transfer(
870        &mut self,
871        operation: &mut MoveOperation,
872        executor: &(impl CommandExecutor + Sync),
873    ) -> Result<()> {
874        let Some(transfer) = operation.workspace_transfer.as_ref().cloned() else {
875            return Ok(());
876        };
877        if transfer.phase == WorkspaceTransferPhase::Ready {
878            return Ok(());
879        }
880        let id = &operation.selection.session_id;
881        ensure!(
882            self.state.sessions[id].state == SessionState::Running
883                && self.state.sessions[id].target != transfer.source.target,
884            "Move destination is not independently ready; source retained"
885        );
886        if !operation.selection.workspace.exclusions.is_empty() {
887            crate::database::retain_move_source(operation)?;
888            if let Some(locator) = &transfer.source.target {
889                let backend = super::super::backend::backend_locator(
890                    locator,
891                    &transfer.source,
892                    &self.config,
893                )?;
894                super::super::execute_checked(
895                    executor,
896                    targets::command_on_locator(
897                        &backend,
898                        id,
899                        vec![
900                            "rm".into(),
901                            "-rf".into(),
902                            "--".into(),
903                            transfer.source_stage.to_string_lossy().into_owned(),
904                            transfer
905                                .source_stage
906                                .with_extension("move-lock")
907                                .to_string_lossy()
908                                .into_owned(),
909                        ],
910                        "remove completed Move source staging",
911                    )?,
912                )?;
913            }
914            executor.notify_notice(&format!("Excluded files remain in the stopped source. Inspect with mj move-sources --session {id}; remove with mj move-sources --session {id} --cleanup {} --yes", operation.operation_id));
915        } else {
916            // Operate on the saved resource identity, never on the destination.
917            // cleanup_stopped_target persists, so use the target/checkout
918            // primitives directly and leave the active session row untouched.
919            let source = &transfer.source;
920            if let Some(locator) = &source.target {
921                let backend =
922                    super::super::backend::backend_locator(locator, source, &self.config)?;
923                targets::retire_move_target_plan(&backend, id)?.execute(executor)?;
924                self.release_retired_move_target(
925                    source,
926                    &backend,
927                    self.state.sessions.get(id),
928                    executor,
929                );
930            }
931            if let mj_core::state::Checkout::ManagedWorktree {
932                worktree: checkout, ..
933            } = source.checkout()
934            {
935                super::super::worktree::retire_managed_worktree(executor, checkout)?;
936            }
937        }
938        if transfer.controller_stage.exists() {
939            std::fs::remove_dir_all(&transfer.controller_stage)?;
940        }
941        operation.workspace_transfer.as_mut().unwrap().phase = WorkspaceTransferPhase::Ready;
942        crate::database::save_move_operation(operation)?;
943        Ok(())
944    }
945
946    pub(in crate::controller) fn rollback_move_destination(
947        &mut self,
948        operation: &MoveOperation,
949        error: anyhow::Error,
950        executor: &(impl CommandExecutor + Sync),
951    ) -> Result<anyhow::Error> {
952        crate::worker_lifecycle::run_blocking(
953            &operation.selection.session_id,
954            "rollback move destination",
955            executor,
956            || {
957                let transfer = operation
958                    .workspace_transfer
959                    .as_ref()
960                    .context("Move transfer missing")?;
961                let id = &operation.selection.session_id;
962                let current = &self.state.sessions[id];
963                if current.target == transfer.source.target {
964                    return self.retain_failed_in_place_move(id, &transfer.source, error);
965                }
966                let prepared = operation
967                    .prepared_destination
968                    .as_ref()
969                    .and_then(|d| d.target())
970                    .cloned();
971                if prepared.is_some() {
972                    let mut saved = crate::database::load_move_operation(id)?
973                        .context("Move cleanup intent missing")?;
974                    self.cleanup_prepared_move_destination(&mut saved, executor)?;
975                }
976                if let Some(locator) = &current.target
977                    && Some(locator) != prepared.as_ref()
978                {
979                    let backend =
980                        super::super::backend::backend_locator(locator, current, &self.config)?;
981                    targets::retire_move_target_plan(&backend, id)?.execute(executor)?;
982                    // The source becomes the session again, so it is the
983                    // placement whose build state must survive.
984                    self.release_retired_move_target(
985                        current,
986                        &backend,
987                        Some(&transfer.source),
988                        executor,
989                    );
990                }
991                let source_checkout = transfer.source.checkout();
992                let source_worktree = match source_checkout {
993                    mj_core::state::Checkout::ManagedWorktree { worktree, .. } => Some(worktree),
994                    _ => None,
995                };
996                if let Ok(mj_core::state::Checkout::ManagedWorktree { worktree, .. }) =
997                    self.state.checkout(id)
998                    && Some(worktree) != source_worktree
999                {
1000                    super::super::worktree::retire_managed_worktree(executor, worktree)?;
1001                }
1002                let mut source = (*transfer.source).clone();
1003                source.state = SessionState::Error;
1004                source.last_error = Some(format!(
1005                    "{error:#}; source environment retained for Move retry"
1006                ));
1007                source.updated_at = now();
1008                crate::database::save_resumed_session(&source, None)?;
1009                self.state.sessions.insert(id.clone(), source);
1010                Ok(error)
1011            },
1012        )
1013    }
1014}
1015
1016#[cfg(test)]
1017mod tests {
1018    use super::*;
1019
1020    // Needs a POSIX shell.
1021    #[cfg(unix)]
1022    #[test]
1023    fn concurrent_move_helper_uploads_publish_complete_files_without_shared_staging() {
1024        let directory = tempfile::tempdir().unwrap();
1025        let id = "move-helper-upload";
1026        let root = directory.path().join(id);
1027        std::fs::create_dir(&root).unwrap();
1028        let binary = directory.path().join("binary");
1029        let body = vec![0x5a; 512 * 1024 + 17];
1030        std::fs::write(&binary, &body).unwrap();
1031        let helper = root
1032            .join("move-helper-build")
1033            .to_string_lossy()
1034            .into_owned();
1035        let backend = targets::TargetLocator::LocalBare {
1036            worker_root: root.to_string_lossy().into_owned(),
1037        };
1038        std::thread::scope(|scope| {
1039            let first = scope.spawn(|| {
1040                install_helper(&targets::ProcessExecutor, &backend, id, &binary, &helper)
1041            });
1042            let second = scope.spawn(|| {
1043                install_helper(&targets::ProcessExecutor, &backend, id, &binary, &helper)
1044            });
1045            first.join().unwrap().unwrap();
1046            second.join().unwrap().unwrap();
1047        });
1048        assert_eq!(std::fs::read(&helper).unwrap(), body);
1049        assert_eq!(std::fs::read_dir(root).unwrap().count(), 1);
1050    }
1051
1052    #[test]
1053    fn shared_filesystem_capacity_counts_allocations_once_and_changes_with_selection() {
1054        struct Disk;
1055        impl CommandExecutor for Disk {
1056            fn execute(&self, _: &CommandSpec) -> Result<targets::CommandOutput> {
1057                Ok(targets::CommandOutput {
1058                    status: 0,
1059                    stdout: b"Filesystem 1024-blocks Used Available Capacity Mounted on\n/dev/test 2000000 0 1500000 0% /\n".to_vec(),
1060                    stderr: Vec::new(),
1061                })
1062            }
1063        }
1064        let path = WorkspacePath {
1065            repository: "project".into(),
1066            path: "large".into(),
1067        };
1068        let mut assessment = WorkspaceAssessment {
1069            files: vec![WorkspaceFile {
1070                location: path.clone(),
1071                bytes: 400_000_000,
1072            }],
1073            ..Default::default()
1074        };
1075        for (location, allocation, copies) in [
1076            ("Source", "source", 1),
1077            ("Controller", "controller", 1),
1078            ("Destination", "destination", 2),
1079            ("Destination backing", "destination", 2),
1080        ] {
1081            storage_probe(
1082                &Disk,
1083                CommandSpec::new("df", ["-Pk"]),
1084                location,
1085                allocation,
1086                copies,
1087                &mut assessment,
1088            );
1089        }
1090        assert_eq!(assessment.storage.len(), 1);
1091        assert_eq!(assessment.storage[0].copies, 4);
1092        let mut selection = WorkspaceSelection::default();
1093        assert!(
1094            assessment
1095                .selection_problem(&selection)
1096                .unwrap()
1097                .contains("free")
1098        );
1099        selection.set_included(&assessment, &path, false);
1100        assert!(assessment.selection_problem(&selection).is_none());
1101    }
1102
1103    // Needs a POSIX shell and rsync.
1104    #[cfg(unix)]
1105    #[test]
1106    fn interrupted_copy_retries_and_exclusions_retain_only_the_source() {
1107        use crate::controller::test_support::{
1108            IsolatedTest, checkpoint_test_session, committed_repository,
1109            resume_compatibility_config,
1110        };
1111        if std::env::var_os("MJ_MOVE_TRANSFER_LIFECYCLE_TEST").is_none() {
1112            let directory = tempfile::tempdir().unwrap();
1113            let name = crate::controller::test_support::test_name(
1114                module_path!(),
1115                "interrupted_copy_retries_and_exclusions_retain_only_the_source",
1116            );
1117            IsolatedTest::new(name)
1118                .isolated_store(directory.path())
1119                .env("MJ_MOVE_TRANSFER_LIFECYCLE_TEST", "1")
1120                .env("MJ_WORKER_BINARY", std::env::current_exe().unwrap())
1121                .run();
1122            return;
1123        }
1124        let _writer = crate::database::install_isolated_test_writer();
1125        struct Executor(std::sync::atomic::AtomicBool);
1126        impl CommandExecutor for Executor {
1127            fn execute(&self, command: &CommandSpec) -> Result<targets::CommandOutput> {
1128                if command.purpose == "transfer Move workspace with resumable files"
1129                    && self.0.swap(false, std::sync::atomic::Ordering::SeqCst)
1130                {
1131                    anyhow::bail!("interrupted copy");
1132                }
1133                targets::ProcessExecutor.execute(command)
1134            }
1135            fn execute_with_stdin(
1136                &self,
1137                _: &CommandSpec,
1138                input: &mut (dyn std::io::Read + Send),
1139            ) -> Result<targets::CommandOutput> {
1140                let request = serde_json::from_reader(input)?;
1141                Ok(targets::CommandOutput {
1142                    status: 0,
1143                    stdout: serde_json::to_vec(&mj_worker::move_workspace::execute(request)?)?,
1144                    stderr: Vec::new(),
1145                })
1146            }
1147        }
1148        let source = committed_repository();
1149        let destination = committed_repository();
1150        std::fs::write(source.path().join("selected"), vec![7; 256 * 1024 + 1]).unwrap();
1151        std::fs::write(source.path().join("excluded"), b"retain me").unwrap();
1152        let id = "move-transfer-lifecycle";
1153        let mut session = checkpoint_test_session(id);
1154        session.project_directory = Some(source.path().into());
1155        session.target_template_id = "local-bare".into();
1156        session.state = SessionState::Closing;
1157        let source_root = mj_core::config::data_dir().join("source-workers").join(id);
1158        std::fs::create_dir_all(&source_root).unwrap();
1159        session.target = Some(mj_core::state::TargetLocator::LocalBare {
1160            worker_root: source_root.clone(),
1161        });
1162        let mut operation = super::super::tests::source_recovery_operation(&session);
1163        operation
1164            .selection
1165            .workspace
1166            .exclusions
1167            .push(WorkspacePath {
1168                repository: "project".into(),
1169                path: "excluded".into(),
1170            });
1171        operation.phase = MovePhase::ResumingDestination;
1172        let executor = Executor(std::sync::atomic::AtomicBool::new(true));
1173        let mut controller = Controller {
1174            config: resume_compatibility_config(),
1175            state: mj_core::state::State::default(),
1176        };
1177        controller.state.sessions.insert(id.into(), session.clone());
1178        crate::database::save_session(&session).unwrap();
1179        let assessment = mj_worker::move_workspace::inspect(&[WorkspaceRepository {
1180            id: "project".into(),
1181            root: source.path().into(),
1182        }])
1183        .unwrap();
1184        operation.workspace_transfer = Some(
1185            controller
1186                .new_workspace_transfer(id, &operation.operation_id, assessment, &executor)
1187                .unwrap(),
1188        );
1189        crate::database::save_move_operation(&operation).unwrap();
1190        assert!(
1191            controller
1192                .capture_move_workspace(&mut operation, &executor)
1193                .is_err()
1194        );
1195        operation = crate::database::load_move_operation(id).unwrap().unwrap();
1196        assert_eq!(
1197            operation.workspace_transfer.as_ref().unwrap().phase,
1198            WorkspaceTransferPhase::Captured
1199        );
1200        controller
1201            .capture_move_workspace(&mut operation, &executor)
1202            .unwrap();
1203        let destination_root = mj_core::config::data_dir()
1204            .join("destination-workers")
1205            .join(id);
1206        std::fs::create_dir_all(&destination_root).unwrap();
1207        let destination_session = controller.state.sessions.get_mut(id).unwrap();
1208        destination_session.target = Some(mj_core::state::TargetLocator::LocalBare {
1209            worker_root: destination_root.clone(),
1210        });
1211        destination_session.project_directory = Some(destination.path().into());
1212        controller.restore_move_workspace(id, &executor).unwrap();
1213        operation = crate::database::load_move_operation(id).unwrap().unwrap();
1214        controller.state.sessions.get_mut(id).unwrap().state = SessionState::Running;
1215        controller
1216            .finish_workspace_transfer(&mut operation, &executor)
1217            .unwrap();
1218        assert!(destination.path().join("selected").exists());
1219        assert!(!destination.path().join("excluded").exists());
1220        assert!(source.path().join("excluded").exists());
1221        assert_eq!(crate::database::retained_move_sources(id).unwrap().len(), 1);
1222        assert!(
1223            !operation
1224                .workspace_transfer
1225                .as_ref()
1226                .unwrap()
1227                .controller_stage
1228                .exists()
1229        );
1230        controller
1231            .cleanup_retained_move_source(id, &operation.operation_id, &executor)
1232            .unwrap();
1233        assert!(
1234            crate::database::retained_move_sources(id)
1235                .unwrap()
1236                .is_empty()
1237        );
1238        assert!(!source_root.exists());
1239        assert!(destination_root.exists());
1240        assert!(
1241            source.path().join("excluded").exists(),
1242            "user-owned checkout must survive cleanup"
1243        );
1244    }
1245
1246    // Needs a POSIX shell and rsync.
1247    #[cfg(unix)]
1248    #[test]
1249    fn rsync_namespace_adapter_streams_large_files_with_spaces_in_paths() {
1250        let directory = tempfile::tempdir().unwrap();
1251        let source = directory.path().join("source with spaces");
1252        let destination = directory.path().join("destination");
1253        std::fs::create_dir(&source).unwrap();
1254        std::fs::create_dir(&destination).unwrap();
1255        let payload = vec![0x71; 256 * 1024 + 13];
1256        std::fs::write(source.join("large file"), &payload).unwrap();
1257        let wrapper = rsync_shell(&CommandSpec::new(
1258            "sh",
1259            ["-c", "exec \"$@\"", "namespace", "rsync"],
1260        ))
1261        .unwrap();
1262        super::super::super::execute_checked(
1263            &targets::ProcessExecutor,
1264            CommandSpec::new(
1265                "rsync",
1266                vec![
1267                    "--recursive".into(),
1268                    "--protect-args".into(),
1269                    "--rsh".into(),
1270                    wrapper
1271                        .path()
1272                        .join("move-rsync-shell")
1273                        .display()
1274                        .to_string(),
1275                    format!("move:{}/", source.display()),
1276                    format!("{}/", destination.display()),
1277                ],
1278            ),
1279        )
1280        .unwrap();
1281        assert_eq!(
1282            std::fs::read(destination.join("large file")).unwrap(),
1283            payload
1284        );
1285    }
1286
1287    // Needs a POSIX shell and rsync.
1288    #[cfg(unix)]
1289    #[test]
1290    fn resumable_local_copy_preserves_large_payloads_and_repairs_partial_files() {
1291        let directory = tempfile::tempdir().unwrap();
1292        let source = directory.path().join("source");
1293        let destination = directory.path().join("destination");
1294        std::fs::create_dir_all(&source).unwrap();
1295        std::fs::create_dir_all(&destination).unwrap();
1296        let bytes = vec![0x35; 192 * 1024 + 17];
1297        std::fs::write(source.join("payload"), &bytes).unwrap();
1298        std::fs::write(destination.join("payload"), &bytes[..70_000]).unwrap();
1299        let backend = targets::TargetLocator::LocalBare {
1300            worker_root: directory.path().display().to_string(),
1301        };
1302        copy_workspace(
1303            &targets::ProcessExecutor,
1304            &backend,
1305            &source,
1306            &destination,
1307            false,
1308        )
1309        .unwrap();
1310        assert_eq!(std::fs::read(destination.join("payload")).unwrap(), bytes);
1311        copy_workspace(
1312            &targets::ProcessExecutor,
1313            &backend,
1314            &source,
1315            &destination,
1316            false,
1317        )
1318        .unwrap();
1319        assert_eq!(std::fs::read(destination.join("payload")).unwrap(), bytes);
1320    }
1321}