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        let retained = crate::database::retained_move_sources(session_id)?
292            .into_iter()
293            .find(|source| source.operation_id == operation_id)
294            .context("retained Move source is missing")?;
295        let source = &retained.source;
296        ensure!(
297            !self
298                .state
299                .sessions
300                .values()
301                .any(|session| session.target.is_some() && session.target == source.target),
302            "retained source is still referenced by an active session"
303        );
304        if let Some(locator) = &source.target {
305            let backend = super::super::backend::backend_locator(locator, source, &self.config)?;
306            targets::retire_move_target_plan(&backend, session_id)?.execute(executor)?;
307        }
308        if let Some(checkout) = &source.managed_worktree {
309            ensure!(
310                !self
311                    .state
312                    .sessions
313                    .values()
314                    .any(|session| session.managed_worktree.as_ref() == Some(checkout)),
315                "retained checkout is still referenced by a session"
316            );
317            super::super::worktree::retire_managed_worktree(executor, checkout)?;
318        }
319        crate::database::forget_retained_move_source(operation_id)
320    }
321
322    pub(in crate::controller) fn move_destination_bundle(
323        &self,
324        id: &str,
325    ) -> Result<Option<targets::ProjectBundleSpec>> {
326        let Some(operation) = crate::database::load_move_operation(id)? else {
327            return Ok(None);
328        };
329        if operation.workspace_transfer.is_none()
330            || operation.phase != MovePhase::ResumingDestination
331        {
332            return Ok(None);
333        }
334        let handoff = operation.handoff.context("Move handoff missing")?;
335        let archive = super::super::resume::verify_resume_checkpoint(id, &handoff)?;
336        let primary = archive
337            .manifest
338            .repositories
339            .iter()
340            .find(|repo| repo.metadata.id == archive.manifest.bundle.primary_repository)
341            .context("Move primary repository missing")?
342            .metadata
343            .relative_destination
344            .to_string_lossy()
345            .into_owned();
346        let repositories = archive
347            .manifest
348            .repositories
349            .iter()
350            .map(|repo| {
351                let metadata = &repo.metadata;
352                mj_core::remote_git::validate_network_url(&metadata.origin)?;
353                Ok(targets::RepositorySpec {
354                    url: Some(metadata.origin.clone()),
355                    push_urls: metadata.push_urls.clone(),
356                    destination: metadata.relative_destination.to_string_lossy().into_owned(),
357                    git_ref: None,
358                    reference: None,
359                })
360            })
361            .collect::<Result<Vec<_>>>()?;
362        Ok(Some(targets::ProjectBundleSpec {
363            primary,
364            repositories,
365        }))
366    }
367
368    pub(super) fn assess_move_workspace(
369        &self,
370        selection: &MoveSelection,
371        executor: &(impl CommandExecutor + Sync),
372    ) -> Result<WorkspaceAssessment> {
373        let id = &selection.session_id;
374        if let Some(operation) = crate::database::load_move_operation(id)?
375            && operation.phase != MovePhase::Completed
376            && let Some(transfer) = operation.workspace_transfer
377        {
378            return Ok(transfer.assessment);
379        }
380        let layout = self.session_export_layout(id, executor)?;
381        let mut assessment: WorkspaceAssessment = serde_json::from_value(utility(
382            executor,
383            &layout.backend,
384            id,
385            &WorkspaceCommand::Inspect {
386                repositories: repositories(&layout),
387            },
388        )?)?;
389        for command in [
390            CommandSpec::new("rsync", rsync_probe_arguments())
391                .purpose("check controller Move transport"),
392            targets::command_on_locator(
393                &layout.backend,
394                id,
395                rsync_probe(),
396                "check source Move transport",
397            )?,
398        ] {
399            match executor.execute(&command) {
400                Ok(output) if output.status == 0 => {}
401                Ok(output) => assessment.blockers.push(format!(
402                    "{}: {RSYNC_REQUIREMENT}: {}",
403                    command.purpose,
404                    String::from_utf8_lossy(&output.stderr)
405                )),
406                Err(error) => assessment
407                    .blockers
408                    .push(format!("{}: {error:#}", command.purpose)),
409            }
410        }
411        let source_root = targets::worker_root(&layout.backend, id)?;
412        storage_probe(
413            executor,
414            targets::command_on_locator(
415                &layout.backend,
416                id,
417                storage_arguments(Path::new(&source_root)),
418                "check source Move staging space",
419            )?,
420            "Source",
421            "source",
422            1,
423            &mut assessment,
424        );
425        if let Some(arguments) = layout
426            .backend
427            .container_engine()
428            .and_then(container_storage_arguments)
429        {
430            let host = targets::locator_command(&layout.backend, vec!["true".into()]).ssh_session;
431            storage_probe(
432                executor,
433                at_host(host.as_ref(), arguments).purpose("check source container backing storage"),
434                "Source backing storage",
435                "source",
436                1,
437                &mut assessment,
438            );
439        }
440        let mut local_args = storage_arguments(&mj_core::config::data_dir());
441        let program = local_args.remove(0);
442        storage_probe(
443            executor,
444            CommandSpec::new(program, local_args).purpose("check controller Move staging space"),
445            "Controller",
446            "controller",
447            1,
448            &mut assessment,
449        );
450        let target = &self.config.targets[selection
451            .target_template_id
452            .as_ref()
453            .context("Move target missing")?];
454        if matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. }) {
455            assessment.blockers.push("Move cannot inspect a not-yet-created EC2 instance's transfer tools or free disk space. Use a configured SSH target on an existing instance.".into());
456            return Ok(assessment);
457        }
458        if matches!(target, mj_core::config::TargetTemplate::LocalBare)
459            && matches!(layout.backend, targets::TargetLocator::LocalBare { .. })
460        {
461            assessment.blockers.push("These local targets share the source worker storage, so no files can be transferred between them.".into());
462        }
463        if let mj_core::config::TargetTemplate::SshBare { ssh, .. } = target
464            && let targets::TargetLocator::SshBare {
465                ssh: source_ssh, ..
466            } = &layout.backend
467            && *source_ssh == targets::SshTarget::from(ssh)
468        {
469            assessment.blockers.push("These SSH targets share the source worker storage, so no files can be transferred between them.".into());
470        }
471        let (destination_ssh, destination_engine, destination_image) = match target {
472            mj_core::config::TargetTemplate::LocalPodman { container } => {
473                (None, Some("podman"), Some(container.image.as_str()))
474            }
475            mj_core::config::TargetTemplate::LocalDocker { container } => {
476                (None, Some("docker"), Some(container.image.as_str()))
477            }
478            mj_core::config::TargetTemplate::AppleContainer { container } => {
479                (None, Some("container"), Some(container.image.as_str()))
480            }
481            mj_core::config::TargetTemplate::SshPodman { ssh, container } => {
482                (Some(ssh), Some("podman"), Some(container.image.as_str()))
483            }
484            mj_core::config::TargetTemplate::SshDocker { ssh, container } => {
485                (Some(ssh), Some("docker"), Some(container.image.as_str()))
486            }
487            mj_core::config::TargetTemplate::SshBare { ssh, .. } => (Some(ssh), None, None),
488            _ => (None, None, None),
489        };
490        let mut arguments =
491            if let (Some(engine), Some(image)) = (destination_engine, destination_image) {
492                let mut arguments = vec![
493                    engine.into(),
494                    "run".into(),
495                    "--rm".into(),
496                    "--entrypoint".into(),
497                    "rsync".into(),
498                    image.into(),
499                ];
500                arguments.extend(rsync_probe_arguments());
501                arguments
502            } else {
503                rsync_probe()
504            };
505        let command = if let Some(ssh) = destination_ssh {
506            targets::ssh_command_owned(&targets::SshTarget::from(ssh), arguments)
507        } else {
508            let program = arguments.remove(0);
509            CommandSpec::new(program, arguments)
510        };
511        match executor.execute(&command.purpose("check destination Move transport")) {
512            Ok(output) if output.status == 0 => {},
513            Ok(output) => assessment.blockers.push(format!("Destination: {RSYNC_REQUIREMENT}. Update the target image or install rsync there: {}", String::from_utf8_lossy(&output.stderr))),
514            Err(error) => assessment.blockers.push(format!("Destination transport check failed: {error:#}")),
515        }
516        let destination_ssh = destination_ssh.map(targets::SshTarget::from);
517        if let mj_core::config::TargetTemplate::LocalPodman { container }
518        | mj_core::config::TargetTemplate::SshPodman { container, .. } = target
519            && let mj_core::config::PodmanWorkspaceStorage::HostHelper { root, .. } =
520                &container.workspace_storage
521        {
522            storage_probe(
523                executor,
524                at_host(destination_ssh.as_ref(), storage_arguments(Path::new(root)))
525                    .purpose("check destination workspace storage"),
526                "Destination workspace",
527                "destination",
528                2,
529                &mut assessment,
530            );
531        }
532        if let Some(engine) = destination_engine {
533            let image = destination_image.context("Move image missing")?;
534            let arguments = vec![
535                engine.into(),
536                "run".into(),
537                "--rm".into(),
538                "--entrypoint".into(),
539                "sh".into(),
540                image.into(),
541                "-c".into(),
542                "df -Pk /".into(),
543            ];
544            storage_probe(
545                executor,
546                at_host(destination_ssh.as_ref(), arguments)
547                    .purpose("check destination container space"),
548                "Destination container",
549                "destination",
550                2,
551                &mut assessment,
552            );
553            if let Some(arguments) = container_storage_arguments(engine) {
554                storage_probe(
555                    executor,
556                    at_host(destination_ssh.as_ref(), arguments)
557                        .purpose("check destination container backing storage"),
558                    "Destination backing storage",
559                    "destination",
560                    2,
561                    &mut assessment,
562                );
563            }
564        } else if !matches!(target, mj_core::config::TargetTemplate::AwsEc2 { .. }) {
565            let path = match target {
566                mj_core::config::TargetTemplate::SshBare {
567                    workspace_prefix, ..
568                } => workspace_prefix.clone(),
569                _ => mj_core::config::data_dir(),
570            };
571            storage_probe(
572                executor,
573                at_host(destination_ssh.as_ref(), storage_arguments(&path))
574                    .purpose("check destination Move staging and workspace space"),
575                "Destination",
576                "destination",
577                2,
578                &mut assessment,
579            );
580        }
581        Ok(assessment)
582    }
583
584    pub(super) fn new_workspace_transfer(
585        &self,
586        id: &str,
587        operation_id: &str,
588        assessment: WorkspaceAssessment,
589        executor: &(impl CommandExecutor + Sync),
590    ) -> Result<WorkspaceTransfer> {
591        let layout = self.session_export_layout(id, executor)?;
592        let root = targets::worker_root(&layout.backend, id)?;
593        Ok(WorkspaceTransfer {
594            assessment,
595            source: Box::new(self.state.sessions[id].clone()),
596            source_stage: Path::new(&root).join(format!("move-{operation_id}")),
597            controller_stage: mj_core::config::data_dir()
598                .join("moves")
599                .join(operation_id)
600                .join("workspace"),
601            phase: WorkspaceTransferPhase::Planned,
602        })
603    }
604
605    pub(super) fn capture_move_workspace(
606        &self,
607        operation: &mut MoveOperation,
608        executor: &(impl CommandExecutor + Sync),
609    ) -> Result<()> {
610        executor.begin_resumable_move_work()?;
611        let result = self.capture_move_workspace_inner(operation, executor);
612        executor.end_resumable_move_work()?;
613        result
614    }
615
616    fn capture_move_workspace_inner(
617        &self,
618        operation: &mut MoveOperation,
619        executor: &(impl CommandExecutor + Sync),
620    ) -> Result<()> {
621        let id = &operation.selection.session_id;
622        let transfer = operation
623            .workspace_transfer
624            .as_ref()
625            .context("Move transfer missing")?
626            .clone();
627        let layout = self.session_export_layout(id, executor)?;
628        if transfer.phase == WorkspaceTransferPhase::Planned {
629            // The target utility holds the stage's file lock before reusing
630            // or replacing an incomplete capture, including after a lost ACK.
631            utility(
632                executor,
633                &layout.backend,
634                id,
635                &WorkspaceCommand::Capture {
636                    repositories: repositories(&layout),
637                    selection: operation.selection.workspace.clone(),
638                    destination: transfer.source_stage.clone(),
639                },
640            )?;
641            operation.workspace_transfer.as_mut().unwrap().phase = WorkspaceTransferPhase::Captured;
642            crate::database::save_move_operation(operation)?;
643        }
644        if operation.workspace_transfer.as_ref().unwrap().phase == WorkspaceTransferPhase::Captured
645        {
646            executor.notify_notice("Transferring workspace to temporary Move storage");
647            copy_workspace(
648                executor,
649                &layout.backend,
650                &transfer.source_stage,
651                &transfer.controller_stage,
652                false,
653            )?;
654            // Verification uses the same installed utility on the controller.
655            let local = targets::TargetLocator::LocalBare {
656                worker_root: mj_core::config::data_dir()
657                    .join("workers")
658                    .join(id)
659                    .to_string_lossy()
660                    .into_owned(),
661            };
662            utility(
663                executor,
664                &local,
665                id,
666                &WorkspaceCommand::Verify {
667                    source: transfer.controller_stage.clone(),
668                },
669            )?;
670            operation.workspace_transfer.as_mut().unwrap().phase =
671                WorkspaceTransferPhase::Downloaded;
672            crate::database::save_move_operation(operation)?;
673        }
674        if operation.workspace_transfer.as_ref().unwrap().phase
675            == WorkspaceTransferPhase::Downloaded
676        {
677            let root = targets::worker_root(&layout.backend, id)?;
678            super::super::execute_checked(
679                executor,
680                targets::command_on_locator(
681                    &layout.backend,
682                    id,
683                    vec![
684                        "sh".into(),
685                        "-c".into(),
686                        targets::stop_worker_daemon_script(&root),
687                    ],
688                    "stop sealed Move source worker, retaining workspace",
689                )?,
690            )?;
691            operation.workspace_transfer.as_mut().unwrap().phase =
692                WorkspaceTransferPhase::SourceStopped;
693            crate::database::save_move_operation(operation)?;
694        }
695        Ok(())
696    }
697
698    pub(in crate::controller) fn restore_move_workspace(
699        &self,
700        id: &str,
701        executor: &(impl CommandExecutor + Sync),
702    ) -> Result<()> {
703        executor.begin_resumable_move_work()?;
704        let result = self.restore_move_workspace_inner(id, executor);
705        executor.end_resumable_move_work()?;
706        result
707    }
708
709    fn restore_move_workspace_inner(
710        &self,
711        id: &str,
712        executor: &(impl CommandExecutor + Sync),
713    ) -> Result<()> {
714        let mut operation =
715            crate::database::load_move_operation(id)?.context("Move intent missing")?;
716        let transfer = operation
717            .workspace_transfer
718            .as_ref()
719            .context("Move transfer missing")?
720            .clone();
721        let layout = self.session_export_layout(id, executor)?;
722        ensure!(
723            self.state.sessions[id].target != transfer.source.target,
724            "Move destination reuses the source environment; refusing to overwrite it"
725        );
726        let stage = PathBuf::from(targets::worker_root(&layout.backend, id)?)
727            .join(format!("move-{}", operation.operation_id));
728        super::super::execute_checked(
729            executor,
730            targets::command_on_locator(
731                &layout.backend,
732                id,
733                vec![
734                    "mkdir".into(),
735                    "-p".into(),
736                    stage.to_string_lossy().into_owned(),
737                ],
738                "create destination Move staging",
739            )?,
740        )?;
741        copy_workspace(
742            executor,
743            &layout.backend,
744            &stage,
745            &transfer.controller_stage,
746            true,
747        )?;
748        utility(
749            executor,
750            &layout.backend,
751            id,
752            &WorkspaceCommand::Restore {
753                source: stage,
754                repositories: repositories(&layout),
755            },
756        )?;
757        operation.workspace_transfer.as_mut().unwrap().phase = WorkspaceTransferPhase::Restored;
758        crate::database::save_move_operation(&operation)?;
759        Ok(())
760    }
761
762    pub(super) fn finish_workspace_transfer(
763        &mut self,
764        operation: &mut MoveOperation,
765        executor: &(impl CommandExecutor + Sync),
766    ) -> Result<()> {
767        let Some(transfer) = operation.workspace_transfer.as_ref().cloned() else {
768            return Ok(());
769        };
770        if transfer.phase == WorkspaceTransferPhase::Ready {
771            return Ok(());
772        }
773        let id = &operation.selection.session_id;
774        ensure!(
775            self.state.sessions[id].state == SessionState::Running
776                && self.state.sessions[id].target != transfer.source.target,
777            "Move destination is not independently ready; source retained"
778        );
779        if !operation.selection.workspace.exclusions.is_empty() {
780            crate::database::retain_move_source(operation)?;
781            if let Some(locator) = &transfer.source.target {
782                let backend = super::super::backend::backend_locator(
783                    locator,
784                    &transfer.source,
785                    &self.config,
786                )?;
787                super::super::execute_checked(
788                    executor,
789                    targets::command_on_locator(
790                        &backend,
791                        id,
792                        vec![
793                            "rm".into(),
794                            "-rf".into(),
795                            "--".into(),
796                            transfer.source_stage.to_string_lossy().into_owned(),
797                            transfer
798                                .source_stage
799                                .with_extension("move-lock")
800                                .to_string_lossy()
801                                .into_owned(),
802                        ],
803                        "remove completed Move source staging",
804                    )?,
805                )?;
806            }
807            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));
808        } else {
809            // Operate on the saved resource identity, never on the destination.
810            // cleanup_stopped_target persists, so use the target/checkout
811            // primitives directly and leave the active session row untouched.
812            let source = &transfer.source;
813            if let Some(locator) = &source.target {
814                let backend =
815                    super::super::backend::backend_locator(locator, source, &self.config)?;
816                targets::retire_move_target_plan(&backend, id)?.execute(executor)?;
817            }
818            if let Some(checkout) = &source.managed_worktree {
819                super::super::worktree::retire_managed_worktree(executor, checkout)?;
820            }
821        }
822        if transfer.controller_stage.exists() {
823            std::fs::remove_dir_all(&transfer.controller_stage)?;
824        }
825        operation.workspace_transfer.as_mut().unwrap().phase = WorkspaceTransferPhase::Ready;
826        crate::database::save_move_operation(operation)?;
827        Ok(())
828    }
829
830    pub(in crate::controller) fn rollback_move_destination(
831        &mut self,
832        operation: &MoveOperation,
833        error: anyhow::Error,
834        executor: &(impl CommandExecutor + Sync),
835    ) -> Result<anyhow::Error> {
836        let transfer = operation
837            .workspace_transfer
838            .as_ref()
839            .context("Move transfer missing")?;
840        let id = &operation.selection.session_id;
841        let current = &self.state.sessions[id];
842        if current.target == transfer.source.target {
843            return self.retain_failed_in_place_move(id, &transfer.source, error);
844        }
845        if let Some(locator) = &current.target {
846            let backend = super::super::backend::backend_locator(locator, current, &self.config)?;
847            targets::retire_move_target_plan(&backend, id)?.execute(executor)?;
848        }
849        if let Some(checkout) = &current.managed_worktree
850            && Some(checkout) != transfer.source.managed_worktree.as_ref()
851        {
852            super::super::worktree::retire_managed_worktree(executor, checkout)?;
853        }
854        let mut source = (*transfer.source).clone();
855        source.state = SessionState::Error;
856        source.last_error = Some(format!(
857            "{error:#}; source environment retained for Move retry"
858        ));
859        source.updated_at = now();
860        crate::database::save_resumed_session(&source, None)?;
861        self.state.sessions.insert(id.clone(), source);
862        Ok(error)
863    }
864}
865
866#[cfg(test)]
867mod tests {
868    use super::*;
869
870    // Needs a POSIX shell.
871    #[cfg(unix)]
872    #[test]
873    fn concurrent_move_helper_uploads_publish_complete_files_without_shared_staging() {
874        let directory = tempfile::tempdir().unwrap();
875        let id = "move-helper-upload";
876        let root = directory.path().join(id);
877        std::fs::create_dir(&root).unwrap();
878        let binary = directory.path().join("binary");
879        let body = vec![0x5a; 512 * 1024 + 17];
880        std::fs::write(&binary, &body).unwrap();
881        let helper = root
882            .join("move-helper-build")
883            .to_string_lossy()
884            .into_owned();
885        let backend = targets::TargetLocator::LocalBare {
886            worker_root: root.to_string_lossy().into_owned(),
887        };
888        std::thread::scope(|scope| {
889            let first = scope.spawn(|| {
890                install_helper(&targets::ProcessExecutor, &backend, id, &binary, &helper)
891            });
892            let second = scope.spawn(|| {
893                install_helper(&targets::ProcessExecutor, &backend, id, &binary, &helper)
894            });
895            first.join().unwrap().unwrap();
896            second.join().unwrap().unwrap();
897        });
898        assert_eq!(std::fs::read(&helper).unwrap(), body);
899        assert_eq!(std::fs::read_dir(root).unwrap().count(), 1);
900    }
901
902    #[test]
903    fn shared_filesystem_capacity_counts_allocations_once_and_changes_with_selection() {
904        struct Disk;
905        impl CommandExecutor for Disk {
906            fn execute(&self, _: &CommandSpec) -> Result<targets::CommandOutput> {
907                Ok(targets::CommandOutput {
908                    status: 0,
909                    stdout: b"Filesystem 1024-blocks Used Available Capacity Mounted on\n/dev/test 2000000 0 1500000 0% /\n".to_vec(),
910                    stderr: Vec::new(),
911                })
912            }
913        }
914        let path = WorkspacePath {
915            repository: "project".into(),
916            path: "large".into(),
917        };
918        let mut assessment = WorkspaceAssessment {
919            files: vec![WorkspaceFile {
920                location: path.clone(),
921                bytes: 400_000_000,
922            }],
923            ..Default::default()
924        };
925        for (location, allocation, copies) in [
926            ("Source", "source", 1),
927            ("Controller", "controller", 1),
928            ("Destination", "destination", 2),
929            ("Destination backing", "destination", 2),
930        ] {
931            storage_probe(
932                &Disk,
933                CommandSpec::new("df", ["-Pk"]),
934                location,
935                allocation,
936                copies,
937                &mut assessment,
938            );
939        }
940        assert_eq!(assessment.storage.len(), 1);
941        assert_eq!(assessment.storage[0].copies, 4);
942        let mut selection = WorkspaceSelection::default();
943        assert!(
944            assessment
945                .selection_problem(&selection)
946                .unwrap()
947                .contains("free")
948        );
949        selection.set_included(&assessment, &path, false);
950        assert!(assessment.selection_problem(&selection).is_none());
951    }
952
953    // Needs a POSIX shell and rsync.
954    #[cfg(unix)]
955    #[test]
956    fn interrupted_copy_retries_and_exclusions_retain_only_the_source() {
957        use crate::controller::test_support::{
958            IsolatedTest, checkpoint_test_session, committed_repository,
959            resume_compatibility_config,
960        };
961        if std::env::var_os("MJ_MOVE_TRANSFER_LIFECYCLE_TEST").is_none() {
962            let directory = tempfile::tempdir().unwrap();
963            let name = crate::controller::test_support::test_name(
964                module_path!(),
965                "interrupted_copy_retries_and_exclusions_retain_only_the_source",
966            );
967            IsolatedTest::new(name)
968                .isolated_store(directory.path())
969                .env("MJ_MOVE_TRANSFER_LIFECYCLE_TEST", "1")
970                .env("MJ_WORKER_BINARY", std::env::current_exe().unwrap())
971                .run();
972            return;
973        }
974        let _writer = crate::database::install_isolated_test_writer();
975        struct Executor(std::sync::atomic::AtomicBool);
976        impl CommandExecutor for Executor {
977            fn execute(&self, command: &CommandSpec) -> Result<targets::CommandOutput> {
978                if command.purpose == "transfer Move workspace with resumable files"
979                    && self.0.swap(false, std::sync::atomic::Ordering::SeqCst)
980                {
981                    anyhow::bail!("interrupted copy");
982                }
983                targets::ProcessExecutor.execute(command)
984            }
985            fn execute_with_stdin(
986                &self,
987                _: &CommandSpec,
988                input: &mut (dyn std::io::Read + Send),
989            ) -> Result<targets::CommandOutput> {
990                let request = serde_json::from_reader(input)?;
991                Ok(targets::CommandOutput {
992                    status: 0,
993                    stdout: serde_json::to_vec(&mj_worker::move_workspace::execute(request)?)?,
994                    stderr: Vec::new(),
995                })
996            }
997        }
998        let source = committed_repository();
999        let destination = committed_repository();
1000        std::fs::write(source.path().join("selected"), vec![7; 256 * 1024 + 1]).unwrap();
1001        std::fs::write(source.path().join("excluded"), b"retain me").unwrap();
1002        let id = "move-transfer-lifecycle";
1003        let mut session = checkpoint_test_session(id);
1004        session.project_directory = Some(source.path().into());
1005        session.target_template_id = "local-bare".into();
1006        session.state = SessionState::Closing;
1007        let source_root = mj_core::config::data_dir().join("source-workers").join(id);
1008        std::fs::create_dir_all(&source_root).unwrap();
1009        session.target = Some(mj_core::state::TargetLocator::LocalBare {
1010            worker_root: source_root.clone(),
1011        });
1012        let mut operation = super::super::tests::source_recovery_operation(&session);
1013        operation
1014            .selection
1015            .workspace
1016            .exclusions
1017            .push(WorkspacePath {
1018                repository: "project".into(),
1019                path: "excluded".into(),
1020            });
1021        operation.phase = MovePhase::ResumingDestination;
1022        let executor = Executor(std::sync::atomic::AtomicBool::new(true));
1023        let mut controller = Controller {
1024            config: resume_compatibility_config(),
1025            state: mj_core::state::State::default(),
1026        };
1027        controller.state.sessions.insert(id.into(), session.clone());
1028        crate::database::save_session(&session).unwrap();
1029        let assessment = mj_worker::move_workspace::inspect(&[WorkspaceRepository {
1030            id: "project".into(),
1031            root: source.path().into(),
1032        }])
1033        .unwrap();
1034        operation.workspace_transfer = Some(
1035            controller
1036                .new_workspace_transfer(id, &operation.operation_id, assessment, &executor)
1037                .unwrap(),
1038        );
1039        crate::database::save_move_operation(&operation).unwrap();
1040        assert!(
1041            controller
1042                .capture_move_workspace(&mut operation, &executor)
1043                .is_err()
1044        );
1045        operation = crate::database::load_move_operation(id).unwrap().unwrap();
1046        assert_eq!(
1047            operation.workspace_transfer.as_ref().unwrap().phase,
1048            WorkspaceTransferPhase::Captured
1049        );
1050        controller
1051            .capture_move_workspace(&mut operation, &executor)
1052            .unwrap();
1053        let destination_root = mj_core::config::data_dir()
1054            .join("destination-workers")
1055            .join(id);
1056        std::fs::create_dir_all(&destination_root).unwrap();
1057        let destination_session = controller.state.sessions.get_mut(id).unwrap();
1058        destination_session.target = Some(mj_core::state::TargetLocator::LocalBare {
1059            worker_root: destination_root.clone(),
1060        });
1061        destination_session.project_directory = Some(destination.path().into());
1062        controller.restore_move_workspace(id, &executor).unwrap();
1063        operation = crate::database::load_move_operation(id).unwrap().unwrap();
1064        controller.state.sessions.get_mut(id).unwrap().state = SessionState::Running;
1065        controller
1066            .finish_workspace_transfer(&mut operation, &executor)
1067            .unwrap();
1068        assert!(destination.path().join("selected").exists());
1069        assert!(!destination.path().join("excluded").exists());
1070        assert!(source.path().join("excluded").exists());
1071        assert_eq!(crate::database::retained_move_sources(id).unwrap().len(), 1);
1072        assert!(
1073            !operation
1074                .workspace_transfer
1075                .as_ref()
1076                .unwrap()
1077                .controller_stage
1078                .exists()
1079        );
1080        controller
1081            .cleanup_retained_move_source(id, &operation.operation_id, &executor)
1082            .unwrap();
1083        assert!(
1084            crate::database::retained_move_sources(id)
1085                .unwrap()
1086                .is_empty()
1087        );
1088        assert!(!source_root.exists());
1089        assert!(destination_root.exists());
1090        assert!(
1091            source.path().join("excluded").exists(),
1092            "user-owned checkout must survive cleanup"
1093        );
1094    }
1095
1096    // Needs a POSIX shell and rsync.
1097    #[cfg(unix)]
1098    #[test]
1099    fn rsync_namespace_adapter_streams_large_files_with_spaces_in_paths() {
1100        let directory = tempfile::tempdir().unwrap();
1101        let source = directory.path().join("source with spaces");
1102        let destination = directory.path().join("destination");
1103        std::fs::create_dir(&source).unwrap();
1104        std::fs::create_dir(&destination).unwrap();
1105        let payload = vec![0x71; 256 * 1024 + 13];
1106        std::fs::write(source.join("large file"), &payload).unwrap();
1107        let wrapper = rsync_shell(&CommandSpec::new(
1108            "sh",
1109            ["-c", "exec \"$@\"", "namespace", "rsync"],
1110        ))
1111        .unwrap();
1112        super::super::super::execute_checked(
1113            &targets::ProcessExecutor,
1114            CommandSpec::new(
1115                "rsync",
1116                vec![
1117                    "--recursive".into(),
1118                    "--protect-args".into(),
1119                    "--rsh".into(),
1120                    wrapper
1121                        .path()
1122                        .join("move-rsync-shell")
1123                        .display()
1124                        .to_string(),
1125                    format!("move:{}/", source.display()),
1126                    format!("{}/", destination.display()),
1127                ],
1128            ),
1129        )
1130        .unwrap();
1131        assert_eq!(
1132            std::fs::read(destination.join("large file")).unwrap(),
1133            payload
1134        );
1135    }
1136
1137    // Needs a POSIX shell and rsync.
1138    #[cfg(unix)]
1139    #[test]
1140    fn resumable_local_copy_preserves_large_payloads_and_repairs_partial_files() {
1141        let directory = tempfile::tempdir().unwrap();
1142        let source = directory.path().join("source");
1143        let destination = directory.path().join("destination");
1144        std::fs::create_dir_all(&source).unwrap();
1145        std::fs::create_dir_all(&destination).unwrap();
1146        let bytes = vec![0x35; 192 * 1024 + 17];
1147        std::fs::write(source.join("payload"), &bytes).unwrap();
1148        std::fs::write(destination.join("payload"), &bytes[..70_000]).unwrap();
1149        let backend = targets::TargetLocator::LocalBare {
1150            worker_root: directory.path().display().to_string(),
1151        };
1152        copy_workspace(
1153            &targets::ProcessExecutor,
1154            &backend,
1155            &source,
1156            &destination,
1157            false,
1158        )
1159        .unwrap();
1160        assert_eq!(std::fs::read(destination.join("payload")).unwrap(), bytes);
1161        copy_workspace(
1162            &targets::ProcessExecutor,
1163            &backend,
1164            &source,
1165            &destination,
1166            false,
1167        )
1168        .unwrap();
1169        assert_eq!(std::fs::read(destination.join("payload")).unwrap(), bytes);
1170    }
1171}