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