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