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