Skip to main content

mj_controller/controller/move_session/
transfer.rs

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