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