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