1use 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 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
98const RSYNC_TRANSFER_OPTIONS: [&str; 6] = [
102 "--recursive",
103 "--times",
104 "--perms",
105 "--protect-args",
106 "--partial",
107 "--partial-dir=.move-partial",
108];
109
110const 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
115fn 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
130fn 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 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 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 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) = ¤t.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) = ¤t.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 #[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 #[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 #[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 #[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}