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