1use std::fmt;
47use std::path::PathBuf;
48
49use serde::Serialize;
50
51use workload_spec::{LifecycleArchetype, SecretRef, VolumeSource, WorkloadSpec};
52
53use crate::config::{CloudConfig, MachineConfig};
54
55pub const KAMAJI_VOLUME_ROOT: &str = "/var/lib/yah/kamaji/volumes";
64
65pub fn named_volume_path(name: &str) -> String {
67 format!("{KAMAJI_VOLUME_ROOT}/{name}")
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
77pub enum MigrationRefusal {
78 UnknownWorkload { name: String, declared: Vec<String> },
80
81 UnknownGroup {
84 group: String,
85 declared: Vec<String>,
86 },
87
88 NotDeployed { workload: String },
91
92 AlreadyInGroup {
94 workload: String,
95 group: String,
96 machines: Vec<String>,
97 },
98
99 MultipleSources {
107 workload: String,
108 archetype: LifecycleArchetype,
109 machines: Vec<String>,
110 },
111
112 FungibleWithDurableVolumes {
122 workload: String,
123 archetype: LifecycleArchetype,
124 volumes: Vec<String>,
125 },
126
127 NoAdmissibleTarget {
135 workload: String,
136 group: String,
137 reason: String,
138 },
139}
140
141impl fmt::Display for MigrationRefusal {
142 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
143 match self {
144 Self::UnknownWorkload { name, declared } => write!(
145 f,
146 "no workload '{name}' in .yah/infra/workloads/ — declared: {}",
147 list_or_none(declared)
148 ),
149 Self::UnknownGroup { group, declared } => write!(
150 f,
151 "no machine declares sovereign_group = \"{group}\" — declared groups: {}. \
152 A group is the set of machines naming it, so migrating to one that \
153 nothing declares would place the workload nowhere; stamp a machine in \
154 .yah/infra/machines/<name>.toml first.",
155 list_or_none(declared)
156 ),
157 Self::NotDeployed { workload } => write!(
158 f,
159 "workload '{workload}' is not running on any declared machine — nothing to \
160 migrate. To place it for the first time use \
161 `yah cloud workload deploy {workload} <machine>`; pass `--from <machine>` \
162 if it is running somewhere this camp cannot reach."
163 ),
164 Self::AlreadyInGroup {
165 workload,
166 group,
167 machines,
168 } => write!(
169 f,
170 "workload '{workload}' is already in sovereign group '{group}' (on {}) — \
171 nothing to do",
172 machines.join(", ")
173 ),
174 Self::MultipleSources {
175 workload,
176 archetype,
177 machines,
178 } => {
179 write!(
180 f,
181 "workload '{workload}' is live on {} machines outside the target group \
182 ({}), and this plans one move at a time — re-run with \
183 `--from <machine>`.",
184 machines.len(),
185 machines.join(", ")
186 )?;
187 if matches!(archetype, LifecycleArchetype::Appliance) {
188 write!(
189 f,
190 " NOTE: '{workload}' is an appliance, whose defining property is at \
191 most one live instance. Two is a violation that predates this \
192 migration — resolve it before moving, or the move carries it \
193 across the cut."
194 )?;
195 }
196 Ok(())
197 }
198 Self::FungibleWithDurableVolumes {
199 workload,
200 archetype,
201 volumes,
202 } => write!(
203 f,
204 "workload '{workload}' declares archetype = \"{}\" but mounts durable \
205 volume(s) [{}]. A fungible workload is dropped and rescheduled, so the \
206 target would come up with empty storage and the data would be silently \
207 left on the source. Declare `archetype = \"appliance\"` if the state \
208 matters, or drop the volume if it does not.",
209 archetype.taint_key(),
210 volumes.join(", ")
211 ),
212 Self::NoAdmissibleTarget {
213 workload,
214 group,
215 reason,
216 } => write!(
217 f,
218 "no machine in sovereign group '{group}' admits workload '{workload}': \
219 {reason}"
220 ),
221 }
222 }
223}
224
225impl std::error::Error for MigrationRefusal {}
226
227fn list_or_none<S: AsRef<str>>(items: &[S]) -> String {
228 if items.is_empty() {
229 "(none)".to_string()
230 } else {
231 items
232 .iter()
233 .map(|s| s.as_ref())
234 .collect::<Vec<_>>()
235 .join(", ")
236 }
237}
238
239#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
241pub enum VolumeDisposition {
242 Copy {
245 name: String,
246 host_path: String,
247 mounted_at: PathBuf,
248 },
249 Precondition {
253 host_path: PathBuf,
254 mounted_at: PathBuf,
255 },
256 Discard { mounted_at: PathBuf, size_mb: u32 },
260}
261
262impl VolumeDisposition {
263 pub fn is_durable(&self) -> bool {
266 !matches!(self, Self::Discard { .. })
267 }
268
269 pub fn label(&self) -> String {
271 match self {
272 Self::Copy { name, .. } => name.clone(),
273 Self::Precondition { host_path, .. } => host_path.display().to_string(),
274 Self::Discard { mounted_at, .. } => format!("tmpfs:{}", mounted_at.display()),
275 }
276 }
277}
278
279#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
282pub struct Precondition {
283 pub what: String,
285 pub because: String,
287}
288
289#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
291pub struct MigrationStep {
292 pub what: String,
294 pub command: Option<String>,
297}
298
299impl MigrationStep {
300 fn narrate(what: impl Into<String>) -> Self {
301 Self {
302 what: what.into(),
303 command: None,
304 }
305 }
306
307 fn run(what: impl Into<String>, command: impl Into<String>) -> Self {
308 Self {
309 what: what.into(),
310 command: Some(command.into()),
311 }
312 }
313}
314
315#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
320pub struct MigrationPlan {
321 pub workload: String,
322 pub archetype: LifecycleArchetype,
323 pub stateful: bool,
326
327 pub source_machine: String,
328 pub source_group: Option<String>,
329 pub source_ssh: Option<String>,
330 pub source_yubaba: Option<String>,
331
332 pub target_machine: String,
333 pub target_group: String,
334 pub target_ssh: Option<String>,
335 pub target_yubaba: Option<String>,
336
337 pub volumes: Vec<VolumeDisposition>,
338 pub preconditions: Vec<Precondition>,
339 pub steps: Vec<MigrationStep>,
340}
341
342impl MigrationPlan {
343 pub fn volumes_to_copy(&self) -> impl Iterator<Item = &VolumeDisposition> {
345 self.volumes
346 .iter()
347 .filter(|v| matches!(v, VolumeDisposition::Copy { .. }))
348 }
349
350 pub fn incurs_downtime(&self) -> bool {
353 self.stateful
354 }
355}
356
357pub fn plan_migration(
369 cfg: &CloudConfig,
370 workload: &str,
371 to_group: &str,
372 observed: &[String],
373) -> Result<MigrationPlan, MigrationRefusal> {
374 let wl = cfg
375 .workload(workload)
376 .ok_or_else(|| MigrationRefusal::UnknownWorkload {
377 name: workload.to_string(),
378 declared: cfg
379 .workloads
380 .iter()
381 .map(|w| w.spec.name.clone())
382 .collect(),
383 })?;
384 let spec = &wl.spec;
385
386 let members = cfg.machines_in_group(to_group);
389 if members.is_empty() {
390 return Err(MigrationRefusal::UnknownGroup {
391 group: to_group.to_string(),
392 declared: cfg
393 .declared_sovereign_groups()
394 .into_iter()
395 .map(String::from)
396 .collect(),
397 });
398 }
399
400 if observed.is_empty() {
401 return Err(MigrationRefusal::NotDeployed {
402 workload: workload.to_string(),
403 });
404 }
405
406 let in_group: Vec<String> = observed
408 .iter()
409 .filter(|name| {
410 cfg.machine(name)
411 .and_then(|m| m.sovereign_group.as_deref())
412 == Some(to_group)
413 })
414 .cloned()
415 .collect();
416 let sources: Vec<String> = observed
417 .iter()
418 .filter(|name| !in_group.contains(name))
419 .cloned()
420 .collect();
421
422 if sources.is_empty() {
423 return Err(MigrationRefusal::AlreadyInGroup {
424 workload: workload.to_string(),
425 group: to_group.to_string(),
426 machines: in_group,
427 });
428 }
429
430 let archetype = spec.effective_archetype();
431
432 if sources.len() > 1 {
433 return Err(MigrationRefusal::MultipleSources {
434 workload: workload.to_string(),
435 archetype,
436 machines: sources,
437 });
438 }
439
440 let volumes = dispositions(spec);
441 let stateful = matches!(archetype, LifecycleArchetype::Appliance);
442
443 if !stateful {
446 let durable: Vec<String> = volumes
447 .iter()
448 .filter(|v| v.is_durable())
449 .map(|v| v.label())
450 .collect();
451 if !durable.is_empty() {
452 return Err(MigrationRefusal::FungibleWithDurableVolumes {
453 workload: workload.to_string(),
454 archetype,
455 volumes: durable,
456 });
457 }
458 }
459
460 let target =
462 cfg.admit_workload_in_group(spec, to_group)
463 .map_err(|e| MigrationRefusal::NoAdmissibleTarget {
464 workload: workload.to_string(),
465 group: to_group.to_string(),
466 reason: e.to_string(),
467 })?;
468
469 let source_name = sources[0].clone();
470 let source = cfg.machine(&source_name);
471
472 let preconditions = preconditions(spec, &volumes, &source_name, &target.name, to_group);
473 let steps = steps(
474 workload,
475 stateful,
476 &volumes,
477 source,
478 &source_name,
479 target,
480 );
481
482 Ok(MigrationPlan {
483 workload: workload.to_string(),
484 archetype,
485 stateful,
486 source_group: source.and_then(|m| m.sovereign_group.clone()),
487 source_ssh: source.and_then(|m| m.connect.as_ref().map(|c| c.ssh.clone())),
488 source_yubaba: source.and_then(|m| m.yubaba_url()),
489 source_machine: source_name,
490 target_machine: target.name.clone(),
491 target_group: to_group.to_string(),
492 target_ssh: target.connect.as_ref().map(|c| c.ssh.clone()),
493 target_yubaba: target.yubaba_url(),
494 volumes,
495 preconditions,
496 steps,
497 })
498}
499
500fn dispositions(spec: &WorkloadSpec) -> Vec<VolumeDisposition> {
503 spec.volumes
504 .iter()
505 .map(|v| match &v.source {
506 VolumeSource::Named { name } => VolumeDisposition::Copy {
507 name: name.clone(),
508 host_path: named_volume_path(name),
509 mounted_at: v.target.clone(),
510 },
511 VolumeSource::Bind { host_path } => VolumeDisposition::Precondition {
512 host_path: host_path.clone(),
513 mounted_at: v.target.clone(),
514 },
515 VolumeSource::Tmpfs { size_mb } => VolumeDisposition::Discard {
516 mounted_at: v.target.clone(),
517 size_mb: *size_mb,
518 },
519 })
520 .collect()
521}
522
523fn preconditions(
526 spec: &WorkloadSpec,
527 volumes: &[VolumeDisposition],
528 source: &str,
529 target: &str,
530 to_group: &str,
531) -> Vec<Precondition> {
532 let mut out = Vec::new();
533
534 let cluster_secrets: Vec<String> = spec
543 .secrets
544 .iter()
545 .filter_map(|s| match &s.source {
546 SecretRef::Cluster { name } => Some(name.clone()),
547 _ => None,
548 })
549 .collect();
550 if !cluster_secrets.is_empty() {
551 out.push(Precondition {
552 what: format!(
553 "re-put cluster secret(s) [{}] into sovereign group '{to_group}': \
554 `yah cloud secret put <name> --machine <raft leader of {to_group}>`",
555 cluster_secrets.join(", ")
556 ),
557 because: "a cluster secret is decrypted from the LOCAL raft replica, and a \
558 sovereign group is a separate raft. The camp KEK is shared, so the \
559 record would decrypt — but it is not replicated across groups, so \
560 it is absent. The workload deploys and then fails to start on a \
561 missing secret."
562 .to_string(),
563 });
564 }
565
566 for v in volumes {
567 if let VolumeDisposition::Precondition {
568 host_path,
569 mounted_at,
570 } = v
571 {
572 out.push(Precondition {
573 what: format!(
574 "ensure {} exists on {target} with the contents {mounted_at:?} expects",
575 host_path.display()
576 ),
577 because: "a bind mount is an operator-managed host path. The camp did not \
578 create it and has no way to know whether it is reproducible on \
579 another box, so it is not copied automatically — an empty \
580 directory would mount cleanly and lose the data silently."
581 .to_string(),
582 });
583 }
584 }
585
586 out.push(Precondition {
591 what: format!(
592 "pre-pull the image on {target} under the exact `repo:tag@digest` ref \
593 (not the bare tag)"
594 ),
595 because: "automatic pulling at admission is unbuilt, and a pull of the plain tag \
596 leaves containerd's image store keyed on a name kamaji never asks for — \
597 the deploy still fails 'image not found in containerd … pre-pull \
598 required'."
599 .to_string(),
600 });
601
602 out.push(Precondition {
605 what: format!(
606 "confirm the image architecture matches {target} (the source is {source})"
607 ),
608 because: "sovereign groups in this fleet differ by hardware — prod is x86 and the \
609 dev group is aarch64 Pis — so an image that runs on the source may have \
610 no matching platform on the target."
611 .to_string(),
612 });
613
614 out
615}
616
617fn steps(
620 workload: &str,
621 stateful: bool,
622 volumes: &[VolumeDisposition],
623 source: Option<&MachineConfig>,
624 source_name: &str,
625 target: &MachineConfig,
626) -> Vec<MigrationStep> {
627 let mut out = Vec::new();
628
629 let source_yubaba = source
630 .and_then(|m| m.yubaba_url())
631 .unwrap_or_else(|| format!("<{source_name} yubaba url>"));
632 let source_ssh = source
633 .and_then(|m| m.connect.as_ref().map(|c| c.ssh.clone()))
634 .unwrap_or_else(|| format!("<{source_name} ssh target>"));
635 let target_ssh = target
636 .connect
637 .as_ref()
638 .map(|c| c.ssh.clone())
639 .unwrap_or_else(|| format!("<{} ssh target>", target.name));
640
641 let find_ident = MigrationStep::run(
647 format!("find the server-assigned ident for '{workload}' on {source_name}"),
648 format!("curl -s {source_yubaba}/workloads"),
649 );
650 let destroy = MigrationStep::run(
651 format!("stop '{workload}' on {source_name}"),
652 format!("curl -X POST {source_yubaba}/workloads/<ident>/destroy"),
653 );
654 let deploy = MigrationStep::run(
655 format!("deploy '{workload}' onto {}", target.name),
656 format!("yah cloud workload deploy {workload} {}", target.name),
657 );
658 let health = MigrationStep::narrate(format!(
659 "confirm '{workload}' is healthy on {} before doing anything else",
660 target.name
661 ));
662
663 if stateful {
664 out.push(MigrationStep::narrate(format!(
668 "NOTE: '{workload}' is stateful — this move takes it DOWN. The source must \
669 stop before the target starts, because an appliance is defined by having at \
670 most one live instance."
671 )));
672 out.push(find_ident);
673 out.push(destroy);
674 out.push(MigrationStep::narrate(
675 "confirm the container is gone on the source before copying — copying a \
676 volume out from under a running writer is how a half-written state file \
677 reaches the target",
678 ));
679
680 for v in volumes {
681 match v {
682 VolumeDisposition::Copy {
683 name, host_path, ..
684 } => {
685 out.push(MigrationStep::run(
686 format!("copy volume '{name}' to {}", target.name),
687 format!(
688 "rsync -aHAX --numeric-ids --delete \
689 {source_ssh}:{host_path}/ {target_ssh}:{host_path}/"
690 ),
691 ));
692 }
693 VolumeDisposition::Precondition { host_path, .. } => {
694 out.push(MigrationStep::narrate(format!(
695 "bind mount {} is operator-managed and is NOT copied — see \
696 preconditions",
697 host_path.display()
698 )));
699 }
700 VolumeDisposition::Discard {
701 mounted_at,
702 size_mb,
703 } => {
704 out.push(MigrationStep::narrate(format!(
705 "tmpfs at {} ({size_mb} MiB) is discarded by design — do not copy it",
706 mounted_at.display()
707 )));
708 }
709 }
710 }
711
712 out.push(deploy);
713 out.push(health);
714 out.push(MigrationStep::narrate(format!(
715 "LAST, and only after the target is verified healthy: remove the stale volume \
716 data on {source_name}. Deliberately manual and deliberately last — it is the \
717 one step with no undo, and keeping it is the whole rollback."
718 )));
719 } else {
720 out.push(MigrationStep::narrate(format!(
723 "NOTE: '{workload}' is fungible — no state moves and there is no downtime. \
724 The target comes up first so it can be proven before the source is dropped."
725 )));
726 out.push(deploy);
727 out.push(health);
728 out.push(find_ident);
729 out.push(destroy);
730 }
731
732 out
733}
734
735#[cfg(test)]
736mod tests {
737 use super::*;
738 use std::path::Path;
739 use tempfile::{tempdir, TempDir};
740
741 struct Camp {
746 dir: TempDir,
747 }
748
749 impl Camp {
750 fn new() -> Self {
751 Self {
752 dir: tempdir().unwrap(),
753 }
754 }
755
756 fn root(&self) -> &Path {
757 self.dir.path()
758 }
759
760 fn machine(self, name: &str, group: Option<&str>) -> Self {
762 self.machine_with(name, group, "")
763 }
764
765 fn machine_with(self, name: &str, group: Option<&str>, extra: &str) -> Self {
766 let dir = self.root().join(".yah/infra/machines");
767 std::fs::create_dir_all(&dir).unwrap();
768 let group_line = group
769 .map(|g| format!("sovereign_group = \"{g}\"\n"))
770 .unwrap_or_default();
771 std::fs::write(
772 dir.join(format!("{name}.toml")),
773 format!(
774 "name = \"{name}\"\nprovider = \"static\"\nmesh_tags = []\n\
775 {group_line}{extra}\n\
776 [connect]\naddress = \"10.0.0.1\"\nssh = \"root@{name}\"\n\
777 yubaba = \"http://{name}:7443\"\n"
778 ),
779 )
780 .unwrap();
781 self
782 }
783
784 fn workload(self, name: &str, extra: &str) -> Self {
787 let dir = self.root().join(".yah/infra/workloads");
788 std::fs::create_dir_all(&dir).unwrap();
789 std::fs::write(
790 dir.join(format!("{name}.toml")),
791 format!(
792 "schema_version = 1\nname = \"{name}\"\ntier = \"infra\"\n\
793 replicas = 1\nrestart_policy = \"always\"\n\
794 {extra}\n\
795 [image]\nregistry = \"cr.yah.dev\"\nrepository = \"{name}\"\n\
796 tag = \"v1\"\ndigest = \"sha256:abc\"\n\
797 [resources]\nmemory_mb = 128\ncpu_millis = 100\n\
798 ephemeral_storage_mb = 64\n\
799 [stop_policy]\nsignal = 15\ngrace_period = 10000\n\
800 [expose.mesh]\nidentity = \"{name}\"\nports = [8080]\nallow_from = []\n"
801 ),
802 )
803 .unwrap();
804 self
805 }
806
807 fn load(&self) -> CloudConfig {
808 CloudConfig::load(self.root()).expect("fixture camp must load")
809 }
810 }
811
812 const NAMED_VOLUME: &str = "[[volumes]]\nsource = { named = { name = \"pgdata\" } }\n\
813 target = \"/var/lib/postgresql\"\nread_only = false\n";
814
815 #[test]
816 fn a_named_volume_renders_the_kamaji_host_path() {
817 assert_eq!(
821 named_volume_path("pgdata"),
822 "/var/lib/yah/kamaji/volumes/pgdata"
823 );
824 }
825
826 #[test]
827 fn an_unknown_group_names_the_declared_vocabulary() {
828 let camp = Camp::new()
829 .machine("a", Some("prod"))
830 .machine("b", Some("dev"))
831 .workload("svc", "archetype = \"server\"");
832 let cfg = camp.load();
833
834 let err = plan_migration(&cfg, "svc", "Dev", &["a".into()]).unwrap_err();
837 let MigrationRefusal::UnknownGroup { declared, .. } = &err else {
838 panic!("expected UnknownGroup, got {err:?}");
839 };
840 assert_eq!(declared, &["dev".to_string(), "prod".to_string()]);
841 let msg = err.to_string();
842 assert!(msg.contains("dev") && msg.contains("prod"), "{msg}");
843 }
844
845 #[test]
846 fn an_undeployed_workload_is_refused_naming_the_deploy_verb() {
847 let camp = Camp::new()
848 .machine("a", Some("dev"))
849 .workload("svc", "archetype = \"server\"");
850 let err = plan_migration(&camp.load(), "svc", "dev", &[]).unwrap_err();
851 assert!(matches!(err, MigrationRefusal::NotDeployed { .. }));
852 assert!(err.to_string().contains("yah cloud workload deploy svc"));
853 }
854
855 #[test]
856 fn a_workload_already_in_the_target_group_is_a_no_op() {
857 let camp = Camp::new()
858 .machine("a", Some("dev"))
859 .machine("b", Some("prod"))
860 .workload("svc", "archetype = \"server\"");
861 let err = plan_migration(&camp.load(), "svc", "dev", &["a".into()]).unwrap_err();
862 assert!(matches!(err, MigrationRefusal::AlreadyInGroup { .. }));
863 }
864
865 #[test]
866 fn two_live_appliance_instances_are_refused_as_a_pre_existing_violation() {
867 let camp = Camp::new()
868 .machine("a", Some("prod"))
869 .machine("b", Some("prod"))
870 .machine("c", Some("dev"))
871 .workload("db", &format!("archetype = \"appliance\"\n{NAMED_VOLUME}"));
872
873 let err =
874 plan_migration(&camp.load(), "db", "dev", &["a".into(), "b".into()]).unwrap_err();
875 assert!(matches!(err, MigrationRefusal::MultipleSources { .. }));
876 let msg = err.to_string();
877 assert!(msg.contains("at most one live instance"), "{msg}");
878 assert!(msg.contains("--from"), "{msg}");
879 }
880
881 #[test]
882 fn a_declared_server_with_a_named_volume_is_refused_not_silently_rescheduled() {
883 let camp = Camp::new()
888 .machine("a", Some("prod"))
889 .machine("b", Some("dev"))
890 .workload("cache", &format!("archetype = \"server\"\n{NAMED_VOLUME}"));
891
892 let err = plan_migration(&camp.load(), "cache", "dev", &["a".into()]).unwrap_err();
893 let MigrationRefusal::FungibleWithDurableVolumes { volumes, .. } = &err else {
894 panic!("expected FungibleWithDurableVolumes, got {err:?}");
895 };
896 assert_eq!(volumes, &["pgdata".to_string()]);
897 assert!(err.to_string().contains("silently left on the source"));
898 }
899
900 #[test]
901 fn a_tmpfs_on_a_fungible_workload_is_not_durable_and_does_not_refuse() {
902 let camp = Camp::new().machine("a", Some("prod")).machine("b", Some("dev")).workload(
903 "svc",
904 "archetype = \"server\"\n[[volumes]]\n\
905 source = { tmpfs = { size_mb = 64 } }\ntarget = \"/scratch\"\nread_only = false\n",
906 );
907
908 let plan =
909 plan_migration(&camp.load(), "svc", "dev", &["a".into()]).expect("tmpfs is not state");
910 assert!(!plan.stateful);
911 assert_eq!(plan.volumes_to_copy().count(), 0);
912
913 assert_eq!(
918 plan.volumes,
919 vec![VolumeDisposition::Discard {
920 mounted_at: PathBuf::from("/scratch"),
921 size_mb: 64,
922 }]
923 );
924 assert!(!plan.volumes[0].is_durable());
925 assert!(plan
926 .steps
927 .iter()
928 .any(|s| s.what.contains("no state moves")));
929 }
930
931 #[test]
932 fn a_repelling_taint_in_the_target_group_refuses_through_the_admission_seam() {
933 let camp = Camp::new()
937 .machine("a", Some("prod"))
938 .machine_with("b", Some("dev"), "taints = [\"no-appliance\"]")
939 .workload("db", &format!("archetype = \"appliance\"\n{NAMED_VOLUME}"));
940
941 let err = plan_migration(&camp.load(), "db", "dev", &["a".into()]).unwrap_err();
942 let MigrationRefusal::NoAdmissibleTarget { reason, .. } = &err else {
943 panic!("expected NoAdmissibleTarget, got {err:?}");
944 };
945 assert!(reason.contains('b'), "must name the candidate: {reason}");
946 }
947
948 #[test]
949 fn the_stateful_half_stops_before_it_copies_and_copies_before_it_starts() {
950 let camp = Camp::new()
951 .machine("a", Some("prod"))
952 .machine("b", Some("dev"))
953 .workload("db", &format!("archetype = \"appliance\"\n{NAMED_VOLUME}"));
954
955 let plan = plan_migration(&camp.load(), "db", "dev", &["a".into()]).unwrap();
956 assert!(plan.stateful && plan.incurs_downtime());
957 assert_eq!(plan.volumes_to_copy().count(), 1);
958 assert_eq!(plan.target_machine, "b");
959 assert_eq!(plan.source_group.as_deref(), Some("prod"));
960
961 let idx = |needle: &str| {
962 plan.steps
963 .iter()
964 .position(|s| {
965 s.command.as_deref().unwrap_or("").contains(needle) || s.what.contains(needle)
966 })
967 .unwrap_or_else(|| panic!("no step matching {needle}: {:#?}", plan.steps))
968 };
969 let (stop, copy, start) = (idx("/destroy"), idx("rsync"), idx("workload deploy"));
970 assert!(
971 stop < copy && copy < start,
972 "stateful ordering must be stop({stop}) → copy({copy}) → start({start})"
973 );
974
975 let rsync = plan.steps[copy].command.as_deref().unwrap();
978 assert!(rsync.contains("root@a:/var/lib/yah/kamaji/volumes/pgdata/"), "{rsync}");
979 assert!(rsync.contains("root@b:/var/lib/yah/kamaji/volumes/pgdata/"), "{rsync}");
980 }
981
982 #[test]
983 fn the_fungible_half_starts_before_it_stops_and_never_copies() {
984 let camp = Camp::new()
985 .machine("a", Some("prod"))
986 .machine("b", Some("dev"))
987 .workload("svc", "archetype = \"server\"");
988
989 let plan = plan_migration(&camp.load(), "svc", "dev", &["a".into()]).unwrap();
990 assert!(!plan.stateful && !plan.incurs_downtime());
991 assert!(
992 !plan
993 .steps
994 .iter()
995 .any(|s| s.command.as_deref().unwrap_or("").contains("rsync")),
996 "the fungible half has nothing to copy"
997 );
998
999 let idx = |needle: &str| {
1000 plan.steps
1001 .iter()
1002 .position(|s| s.command.as_deref().unwrap_or("").contains(needle))
1003 .unwrap_or_else(|| panic!("no step matching {needle}: {:#?}", plan.steps))
1004 };
1005 assert!(
1006 idx("workload deploy") < idx("/destroy"),
1007 "fungible ordering must be start → stop, so the target is proven first"
1008 );
1009 }
1010
1011 #[test]
1012 fn a_cluster_secret_raises_the_cross_group_raft_precondition() {
1013 let camp = Camp::new().machine("a", Some("prod")).machine("b", Some("dev")).workload(
1016 "svc",
1017 "archetype = \"server\"\n[[secrets]]\n\
1018 source = { cluster = { name = \"cheers/cloud-admin/verify-key\" } }\n\
1019 target = { file = { path = \"/run/secrets/k\", mode = 256 } }\n",
1020 );
1021
1022 let plan = plan_migration(&camp.load(), "svc", "dev", &["a".into()]).unwrap();
1023 let p = plan
1024 .preconditions
1025 .iter()
1026 .find(|p| p.what.contains("cluster secret"))
1027 .expect("cluster secret precondition must be raised");
1028 assert!(p.what.contains("cheers/cloud-admin/verify-key"), "{p:?}");
1029 assert!(p.because.contains("separate raft"), "{p:?}");
1030 }
1031
1032 #[test]
1033 fn a_bind_mount_is_a_precondition_and_is_never_rsynced() {
1034 let camp = Camp::new().machine("a", Some("prod")).machine("b", Some("dev")).workload(
1035 "db",
1036 "archetype = \"appliance\"\n[[volumes]]\n\
1037 source = { bind = { host_path = \"/srv/data\" } }\n\
1038 target = \"/data\"\nread_only = false\n",
1039 );
1040
1041 let plan = plan_migration(&camp.load(), "db", "dev", &["a".into()]).unwrap();
1042 assert_eq!(plan.volumes_to_copy().count(), 0);
1043 assert!(plan
1044 .preconditions
1045 .iter()
1046 .any(|p| p.what.contains("/srv/data")));
1047 assert!(
1048 !plan
1049 .steps
1050 .iter()
1051 .any(|s| s.command.as_deref().unwrap_or("").contains("/srv/data")),
1052 "an operator-managed host path must never be rsynced blindly"
1053 );
1054 }
1055}