1use async_trait::async_trait;
57use serde::{Deserialize, Serialize};
58
59use crate::kv::{KvStore, WriteOp};
60use crate::project::{self, owner_kind, DomainOwner, DEFAULT_PROJECT};
61use crate::time::now_unix;
62
63pub const SCHEMA_KEY: &str = "schema/version";
65
66pub const PREMIGRATION_INDEX_KEY: &str = "schema/premigration-index";
69
70pub const CURRENT_VERSION: u32 = 1;
73
74pub const MUTABLE_FAMILIES: &[&str] = &[
78 "current/",
79 "site/",
80 "history/",
81 "alias/",
82 "domainverify/",
83 "dnsmanaged/",
84 "functions/",
85 "metering/",
86 "blobnotify/",
87 "workflows/",
88 "compute/",
89 "compute_state/",
90];
91
92pub const DOMAIN_FAMILIES: &[&str] = &["domain/", "wildcard/", "httpchallenge/"];
95
96pub const GLOBAL_FAMILIES: &[&str] = &[
101 "manifests/",
102 "meta/",
103 "siteconfig/",
104 "computever/",
105 "daemonconfig/",
106 "authz/",
107 "daemon/",
108 "cert/",
109 "project/",
111 "projectmeta/",
112 "projectver/",
113 "project-history/",
114 "owner/",
115 "schema/",
116];
117
118#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
122#[serde(default, deny_unknown_fields)]
123pub struct SchemaState {
124 pub version: u32,
128 #[serde(skip_serializing_if = "Vec::is_empty")]
131 pub unfinalized: Vec<u32>,
132 #[serde(skip_serializing_if = "Option::is_none")]
135 pub in_progress: Option<InProgress>,
136 #[serde(skip_serializing_if = "Vec::is_empty")]
138 pub history: Vec<AppliedRecord>,
139 pub updated_at: u64,
141}
142
143#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
146pub struct InProgress {
147 pub target: u32,
149 pub steps_done: Vec<String>,
151}
152
153#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
155pub struct AppliedRecord {
156 pub version: u32,
158 pub id: String,
160 pub at: u64,
162}
163
164#[derive(Debug, Deserialize)]
167struct LegacyMarker {
168 #[serde(default)]
169 layout: u32,
170 #[serde(default)]
171 dual: bool,
172 #[serde(default)]
173 families_done: Vec<String>,
174}
175
176impl LegacyMarker {
177 fn into_state(self) -> SchemaState {
178 if self.layout >= 2 {
179 SchemaState {
181 version: 1,
182 unfinalized: if self.dual { vec![1] } else { Vec::new() },
183 in_progress: None,
184 history: Vec::new(),
185 updated_at: now_unix(),
186 }
187 } else {
188 let steps_done: Vec<String> = self
191 .families_done
192 .iter()
193 .map(|f| rekey_step_id(DEFAULT_PROJECT, f))
194 .collect();
195 SchemaState {
196 version: 0,
197 unfinalized: Vec::new(),
198 in_progress: (!steps_done.is_empty()).then_some(InProgress {
199 target: 1,
200 steps_done,
201 }),
202 history: Vec::new(),
203 updated_at: now_unix(),
204 }
205 }
206 }
207}
208
209#[derive(Debug, Clone, Copy, PartialEq, Eq)]
213pub enum Status {
214 Ready,
216 NeedsMigration,
219 Dual,
222}
223
224#[derive(Debug, Clone, Copy, Default)]
226pub struct MigrateOptions {
227 pub dry_run: bool,
229 pub finalize: bool,
233}
234
235impl MigrateOptions {
236 pub fn one_shot() -> Self {
238 Self {
239 dry_run: false,
240 finalize: true,
241 }
242 }
243}
244
245#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
248pub struct MigrationReport {
249 pub rekeyed: Vec<(String, usize)>,
251 pub values_rewritten: Vec<(String, usize)>,
253 pub owner_entries: usize,
255 pub created_default_project: bool,
257 pub already_migrated: bool,
259 pub dual: bool,
261}
262
263impl MigrationReport {
264 pub fn total_rekeyed(&self) -> usize {
266 self.rekeyed.iter().map(|(_, n)| n).sum()
267 }
268}
269
270#[derive(Debug, thiserror::Error)]
272pub enum MigrateError {
273 #[error(transparent)]
275 Kv(#[from] crate::error::KvError),
276 #[error("verification failed for migrated key {0}")]
279 Verify(String),
280 #[error("migration serde error: {0}")]
282 Serde(String),
283}
284
285pub async fn read_state(kv: &dyn KvStore) -> Result<SchemaState, MigrateError> {
288 let Some(bytes) = kv.get(SCHEMA_KEY).await? else {
289 return Ok(SchemaState::default());
290 };
291 if let Ok(state) = serde_json::from_slice::<SchemaState>(&bytes) {
292 return Ok(state);
293 }
294 let legacy: LegacyMarker =
296 serde_json::from_slice(&bytes).map_err(|e| MigrateError::Serde(e.to_string()))?;
297 Ok(legacy.into_state())
298}
299
300async fn persist_state(kv: &dyn KvStore, state: &SchemaState) -> Result<(), MigrateError> {
301 let bytes = serde_json::to_vec(state).map_err(|e| MigrateError::Serde(e.to_string()))?;
302 kv.put(SCHEMA_KEY, bytes).await?;
303 Ok(())
304}
305
306pub async fn status(kv: &dyn KvStore) -> Result<Status, MigrateError> {
308 status_of(kv, ®istry()).await
309}
310
311async fn status_of(
312 kv: &dyn KvStore,
313 migrations: &[Box<dyn Migration>],
314) -> Result<Status, MigrateError> {
315 let state = read_state(kv).await?;
316 if state.in_progress.is_some() {
318 return Ok(Status::NeedsMigration);
319 }
320 if !state.unfinalized.is_empty() {
322 return Ok(Status::Dual);
323 }
324 for m in migrations.iter().filter(|m| m.version() > state.version) {
328 if m.is_applicable(kv).await? {
329 return Ok(Status::NeedsMigration);
330 }
331 }
332 Ok(Status::Ready)
333}
334
335pub async fn migrate(
341 kv: &dyn KvStore,
342 opts: MigrateOptions,
343) -> Result<MigrationReport, MigrateError> {
344 run(kv, opts, ®istry()).await
345}
346
347pub async fn finalize(kv: &dyn KvStore) -> Result<MigrationReport, MigrateError> {
350 migrate(kv, MigrateOptions::one_shot()).await
351}
352
353async fn run(
356 kv: &dyn KvStore,
357 opts: MigrateOptions,
358 migrations: &[Box<dyn Migration>],
359) -> Result<MigrationReport, MigrateError> {
360 let mut state = read_state(kv).await?;
361 let mut report = MigrationReport::default();
362
363 let has_pending_forward = migrations.iter().any(|m| m.version() > state.version);
364 if !has_pending_forward && state.unfinalized.is_empty() && state.in_progress.is_none() {
365 report.already_migrated = true;
366 return Ok(report);
367 }
368
369 if opts.dry_run {
371 for m in migrations.iter().filter(|m| m.version() > state.version) {
372 if m.is_applicable(kv).await? {
373 for step in m.forward_steps() {
374 step.run(kv, true, &mut report).await?;
375 }
376 }
377 }
378 report.dual = !opts.finalize;
379 return Ok(report);
380 }
381
382 if kv.get(PREMIGRATION_INDEX_KEY).await?.is_none() {
384 write_premigration_index(kv).await?;
385 }
386
387 for m in migrations {
389 if m.version() <= state.version {
390 continue;
391 }
392 let resuming = state
393 .in_progress
394 .as_ref()
395 .is_some_and(|ip| ip.target == m.version());
396
397 if !resuming && !m.is_applicable(kv).await? {
400 state.version = m.version();
401 state.history.push(AppliedRecord {
402 version: m.version(),
403 id: m.id().to_string(),
404 at: now_unix(),
405 });
406 state.updated_at = now_unix();
407 persist_state(kv, &state).await?;
408 continue;
409 }
410
411 let mut done = if resuming {
414 state
415 .in_progress
416 .take()
417 .map(|ip| ip.steps_done)
418 .unwrap_or_default()
419 } else {
420 Vec::new()
421 };
422 for step in m.forward_steps() {
423 let sid = step.id();
424 if done.contains(&sid) {
425 continue;
426 }
427 step.run(kv, false, &mut report).await?;
428 done.push(sid);
429 state.in_progress = Some(InProgress {
430 target: m.version(),
431 steps_done: done.clone(),
432 });
433 state.updated_at = now_unix();
434 persist_state(kv, &state).await?;
435 }
436
437 state.version = m.version();
439 state.in_progress = None;
440 state.history.push(AppliedRecord {
441 version: m.version(),
442 id: m.id().to_string(),
443 at: now_unix(),
444 });
445 if m.online() {
446 state.unfinalized.push(m.version());
447 }
448 state.updated_at = now_unix();
449 persist_state(kv, &state).await?;
450 }
451
452 if opts.finalize && !state.unfinalized.is_empty() {
454 let mut pending = state.unfinalized.clone();
455 pending.sort_unstable();
456 for v in pending {
457 if let Some(m) = migrations.iter().find(|m| m.version() == v) {
458 for step in m.cleanup_steps() {
459 step.run(kv, false, &mut report).await?;
460 }
461 }
462 }
463 state.unfinalized.clear();
464 state.updated_at = now_unix();
465 persist_state(kv, &state).await?;
466 }
467
468 report.dual = !state.unfinalized.is_empty();
469 Ok(report)
470}
471
472async fn write_premigration_index(kv: &dyn KvStore) -> Result<(), MigrateError> {
475 let mut keys = Vec::new();
476 for family in MUTABLE_FAMILIES.iter().chain(DOMAIN_FAMILIES) {
477 keys.extend(kv.list_prefix(family).await?);
478 }
479 keys.sort();
480 kv.put(PREMIGRATION_INDEX_KEY, keys.join("\n").into_bytes())
481 .await?;
482 Ok(())
483}
484
485fn registry() -> Vec<Box<dyn Migration>> {
492 vec![Box::new(ProjectRekeyV1)]
493}
494
495#[async_trait]
497trait Migration: Send + Sync {
498 fn version(&self) -> u32;
501 fn id(&self) -> &'static str;
503 #[allow(dead_code)]
505 fn description(&self) -> &'static str;
506 async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError>;
510 fn forward_steps(&self) -> Vec<Box<dyn Step>>;
513 fn cleanup_steps(&self) -> Vec<Box<dyn Step>> {
516 Vec::new()
517 }
518 fn online(&self) -> bool {
521 !self.cleanup_steps().is_empty()
522 }
523}
524
525struct ProjectRekeyV1;
530
531#[async_trait]
532impl Migration for ProjectRekeyV1 {
533 fn version(&self) -> u32 {
534 1
535 }
536 fn id(&self) -> &'static str {
537 "project-rekey"
538 }
539 fn description(&self) -> &'static str {
540 "re-key the pre-0.2.0 store under project/<default>/ and add the project-scoped index"
541 }
542 async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError> {
543 for family in MUTABLE_FAMILIES.iter().chain(DOMAIN_FAMILIES) {
545 if !kv.list_prefix(family).await?.is_empty() {
546 return Ok(true);
547 }
548 }
549 Ok(false)
550 }
551 fn forward_steps(&self) -> Vec<Box<dyn Step>> {
552 let mut steps: Vec<Box<dyn Step>> = Vec::new();
553 for family in MUTABLE_FAMILIES {
554 steps.push(Box::new(RekeyFamily {
555 family,
556 project: DEFAULT_PROJECT,
557 }));
558 }
559 for family in DOMAIN_FAMILIES {
560 steps.push(Box::new(RewriteValues {
561 family,
562 transform: domain_owner_canonical,
563 }));
564 }
565 steps.push(Box::new(EnsureDefaultProject));
566 steps.push(Box::new(BuildOwnerIndex));
567 steps
568 }
569 fn cleanup_steps(&self) -> Vec<Box<dyn Step>> {
570 MUTABLE_FAMILIES
571 .iter()
572 .map(|family| {
573 Box::new(DeleteOldFamily {
574 family,
575 project: DEFAULT_PROJECT,
576 }) as Box<dyn Step>
577 })
578 .collect()
579 }
580}
581
582#[async_trait]
588trait Step: Send + Sync {
589 fn id(&self) -> String;
590 async fn run(
591 &self,
592 kv: &dyn KvStore,
593 dry_run: bool,
594 report: &mut MigrationReport,
595 ) -> Result<(), MigrateError>;
596}
597
598fn rekey_step_id(project: &str, family: &str) -> String {
599 format!("rekey:{project}:{family}")
600}
601
602struct RekeyFamily {
605 family: &'static str,
606 project: &'static str,
607}
608
609#[async_trait]
610impl Step for RekeyFamily {
611 fn id(&self) -> String {
612 rekey_step_id(self.project, self.family)
613 }
614 async fn run(
615 &self,
616 kv: &dyn KvStore,
617 dry_run: bool,
618 report: &mut MigrationReport,
619 ) -> Result<(), MigrateError> {
620 let mut moved = 0;
621 for old_key in kv.list_prefix(self.family).await? {
622 let new_key = format!("project/{}/{}", self.project, old_key);
623 if dry_run {
624 moved += 1;
625 continue;
626 }
627 let Some(value) = kv.get(&old_key).await? else {
628 continue; };
630 kv.put(&new_key, value.clone()).await?;
631 if kv.get(&new_key).await?.as_deref() != Some(value.as_slice()) {
632 return Err(MigrateError::Verify(new_key));
633 }
634 moved += 1;
635 }
636 if moved > 0 {
637 report.rekeyed.push((self.family.to_string(), moved));
638 }
639 Ok(())
640 }
641}
642
643struct DeleteOldFamily {
648 family: &'static str,
649 project: &'static str,
650}
651
652#[async_trait]
653impl Step for DeleteOldFamily {
654 fn id(&self) -> String {
655 format!("delete:{}:{}", self.project, self.family)
656 }
657 async fn run(
658 &self,
659 kv: &dyn KvStore,
660 dry_run: bool,
661 _report: &mut MigrationReport,
662 ) -> Result<(), MigrateError> {
663 for old_key in kv.list_prefix(self.family).await? {
664 let new_key = format!("project/{}/{}", self.project, old_key);
665 if dry_run {
666 continue;
667 }
668 match (kv.get(&old_key).await?, kv.get(&new_key).await?) {
669 (Some(o), Some(n)) if o == n => kv.delete(&old_key).await?,
670 (Some(_), _) => return Err(MigrateError::Verify(new_key)),
671 (None, _) => {}
672 }
673 }
674 Ok(())
675 }
676}
677
678struct RewriteValues {
681 family: &'static str,
682 transform: fn(&[u8]) -> Vec<u8>,
683}
684
685#[async_trait]
686impl Step for RewriteValues {
687 fn id(&self) -> String {
688 format!("rewrite:{}", self.family)
689 }
690 async fn run(
691 &self,
692 kv: &dyn KvStore,
693 dry_run: bool,
694 report: &mut MigrationReport,
695 ) -> Result<(), MigrateError> {
696 let mut rewritten = 0;
697 for key in kv.list_prefix(self.family).await? {
698 let Some(value) = kv.get(&key).await? else {
699 continue;
700 };
701 let canonical = (self.transform)(&value);
702 if canonical != value {
703 if !dry_run {
704 kv.put(&key, canonical).await?;
705 }
706 rewritten += 1;
707 }
708 }
709 if rewritten > 0 {
710 report
711 .values_rewritten
712 .push((self.family.to_string(), rewritten));
713 }
714 Ok(())
715 }
716}
717
718fn domain_owner_canonical(value: &[u8]) -> Vec<u8> {
721 DomainOwner::from_bytes(value).to_bytes()
722}
723
724struct EnsureDefaultProject;
726
727#[async_trait]
728impl Step for EnsureDefaultProject {
729 fn id(&self) -> String {
730 "ensure-default-project".to_string()
731 }
732 async fn run(
733 &self,
734 kv: &dyn KvStore,
735 dry_run: bool,
736 report: &mut MigrationReport,
737 ) -> Result<(), MigrateError> {
738 let pointer = project::pointer_key(DEFAULT_PROJECT);
739 if kv.get(&pointer).await?.is_some() {
740 return Ok(());
741 }
742 report.created_default_project = true;
743 if dry_run {
744 return Ok(());
745 }
746 let default = crate::deploy::DeployStore::default_project_record();
749 let hash = default.id();
750 let body = serde_json::to_vec(&default).map_err(|e| MigrateError::Serde(e.to_string()))?;
751 kv.write_batch(vec![
752 WriteOp::Put(project::spec_key(&hash), body),
753 WriteOp::Put(pointer, hash.into_bytes()),
754 ])
755 .await?;
756 Ok(())
757 }
758}
759
760struct BuildOwnerIndex;
763
764#[async_trait]
765impl Step for BuildOwnerIndex {
766 fn id(&self) -> String {
767 "build-owner-index".to_string()
768 }
769 async fn run(
770 &self,
771 kv: &dyn KvStore,
772 dry_run: bool,
773 report: &mut MigrationReport,
774 ) -> Result<(), MigrateError> {
775 let mut ops = Vec::new();
776 let site_prefix = format!("project/{DEFAULT_PROJECT}/site/");
778 for key in kv.list_prefix(&site_prefix).await? {
779 if let Some(site) = key.strip_prefix(&site_prefix) {
780 if !site.is_empty() {
781 ops.push(WriteOp::Put(
782 project::owner_key(owner_kind::SITE, site),
783 DEFAULT_PROJECT.as_bytes().to_vec(),
784 ));
785 }
786 }
787 }
788 let fn_prefix = format!("project/{DEFAULT_PROJECT}/functions/");
790 for key in kv.list_prefix(&fn_prefix).await? {
791 if let Some(rest) = key.strip_prefix(&fn_prefix) {
792 if !rest.is_empty() && !rest.contains('/') {
793 ops.push(WriteOp::Put(
794 project::owner_key(owner_kind::FUNCTION, rest),
795 DEFAULT_PROJECT.as_bytes().to_vec(),
796 ));
797 }
798 }
799 }
800 let compute_prefix = format!("project/{DEFAULT_PROJECT}/compute/");
802 for key in kv.list_prefix(&compute_prefix).await? {
803 if let Some(name) = key.strip_prefix(&compute_prefix) {
804 if !name.is_empty() {
805 ops.push(WriteOp::Put(
806 project::owner_key(owner_kind::COMPUTE, name),
807 DEFAULT_PROJECT.as_bytes().to_vec(),
808 ));
809 }
810 }
811 }
812 report.owner_entries += ops.len();
813 if !dry_run && !ops.is_empty() {
814 kv.write_batch(ops).await?;
815 }
816 Ok(())
817 }
818}
819
820#[cfg(test)]
821mod tests {
822 use super::*;
823 use crate::kv::MemoryKv;
824
825 async fn seed_legacy(kv: &MemoryKv) {
828 kv.put("current/blog", b"dep-1".to_vec()).await.unwrap();
829 kv.put("site/blog", b"cfghash".to_vec()).await.unwrap();
830 kv.put("history/blog", b"[]".to_vec()).await.unwrap();
831 kv.put("alias/blog/staging", b"dep-1".to_vec())
832 .await
833 .unwrap();
834 kv.put("domainverify/blog/www.example", b"{}".to_vec())
835 .await
836 .unwrap();
837 kv.put("dnsmanaged/blog/www.example", b"{}".to_vec())
838 .await
839 .unwrap();
840 kv.put("functions/resize", b"{}".to_vec()).await.unwrap();
841 kv.put("functions/resize/versions/v1", b"{}".to_vec())
842 .await
843 .unwrap();
844 kv.put("metering/resize", b"{}".to_vec()).await.unwrap();
845 kv.put("blobnotify/resize/uploads", b"{}".to_vec())
846 .await
847 .unwrap();
848 kv.put("workflows/etl", b"{}".to_vec()).await.unwrap();
849 kv.put("compute/api", b"{}".to_vec()).await.unwrap();
850 kv.put("compute_state/api/0", b"{}".to_vec()).await.unwrap();
851 kv.put("domain/www.example", b"blog".to_vec())
852 .await
853 .unwrap();
854 kv.put("wildcard/preview.example", b"blog".to_vec())
855 .await
856 .unwrap();
857 kv.put("httpchallenge/www.example/tok", b"blog".to_vec())
858 .await
859 .unwrap();
860 kv.put("siteconfig/cfghash", b"the-config".to_vec())
861 .await
862 .unwrap();
863 kv.put("manifests/dep-1", b"the-manifest".to_vec())
864 .await
865 .unwrap();
866 kv.put("authz/tokens/t1", b"tok".to_vec()).await.unwrap();
867 }
868
869 #[test]
870 fn registry_versions_are_strictly_ascending_and_reach_current() {
871 let reg = registry();
872 assert!(!reg.is_empty());
873 let mut last = 0;
874 for m in ® {
875 assert!(m.version() > last, "versions must strictly ascend");
876 last = m.version();
877 }
878 assert_eq!(
879 last, CURRENT_VERSION,
880 "CURRENT_VERSION == the top migration"
881 );
882 }
883
884 #[tokio::test]
885 async fn status_detects_legacy_fresh_and_migrated() {
886 let fresh = MemoryKv::new();
887 assert_eq!(status(&fresh).await.unwrap(), Status::Ready);
888
889 let legacy = MemoryKv::new();
890 seed_legacy(&legacy).await;
891 assert_eq!(status(&legacy).await.unwrap(), Status::NeedsMigration);
892
893 migrate(&legacy, MigrateOptions::one_shot()).await.unwrap();
894 assert_eq!(status(&legacy).await.unwrap(), Status::Ready);
895 let state = read_state(&legacy).await.unwrap();
897 assert_eq!(state.version, CURRENT_VERSION);
898 assert!(state.history.iter().any(|a| a.id == "project-rekey"));
899 }
900
901 #[tokio::test]
902 async fn migrate_rekeys_mutable_families_and_rewrites_domain_values() {
903 let kv = MemoryKv::new();
904 seed_legacy(&kv).await;
905 let report = migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
906
907 assert_eq!(
908 kv.get("project/default/current/blog")
909 .await
910 .unwrap()
911 .as_deref(),
912 Some(&b"dep-1"[..])
913 );
914 assert_eq!(
915 kv.get("project/default/functions/resize/versions/v1")
916 .await
917 .unwrap()
918 .as_deref(),
919 Some(&b"{}"[..])
920 );
921 assert_eq!(
922 kv.get("project/default/compute_state/api/0")
923 .await
924 .unwrap()
925 .as_deref(),
926 Some(&b"{}"[..])
927 );
928 assert!(kv.get("current/blog").await.unwrap().is_none());
929 assert!(kv.get("compute/api").await.unwrap().is_none());
930
931 let dv = kv.get("domain/www.example").await.unwrap().unwrap();
932 assert_eq!(
933 DomainOwner::from_bytes(&dv),
934 DomainOwner::new("default", "blog")
935 );
936 assert!(String::from_utf8_lossy(&dv).contains("\"project\":\"default\""));
937 assert_eq!(
938 DomainOwner::from_bytes(
939 &kv.get("httpchallenge/www.example/tok")
940 .await
941 .unwrap()
942 .unwrap()
943 ),
944 DomainOwner::new("default", "blog")
945 );
946
947 assert_eq!(
949 kv.get("siteconfig/cfghash").await.unwrap().as_deref(),
950 Some(&b"the-config"[..])
951 );
952 assert_eq!(
953 kv.get("authz/tokens/t1").await.unwrap().as_deref(),
954 Some(&b"tok"[..])
955 );
956
957 assert!(kv.get("projectmeta/default").await.unwrap().is_some());
959 assert!(report.created_default_project);
960 assert_eq!(
961 kv.get("owner/site/blog").await.unwrap().as_deref(),
962 Some(&b"default"[..])
963 );
964 assert_eq!(
965 kv.get("owner/function/resize").await.unwrap().as_deref(),
966 Some(&b"default"[..])
967 );
968 assert!(kv
969 .get("owner/function/resize/versions/v1")
970 .await
971 .unwrap()
972 .is_none());
973
974 assert_eq!(status(&kv).await.unwrap(), Status::Ready);
975 assert!(!report.already_migrated);
976 }
977
978 #[tokio::test]
979 async fn migrate_is_idempotent() {
980 let kv = MemoryKv::new();
981 seed_legacy(&kv).await;
982 let first = migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
983 assert!(first.total_rekeyed() > 0);
984
985 let before: Vec<String> = kv.list_prefix("project/").await.unwrap();
986 let second = migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
987 assert!(second.already_migrated);
988 assert_eq!(second.total_rekeyed(), 0);
989 assert_eq!(kv.list_prefix("project/").await.unwrap(), before);
990 }
991
992 #[tokio::test]
993 async fn dry_run_writes_nothing() {
994 let kv = MemoryKv::new();
995 seed_legacy(&kv).await;
996 let report = migrate(
997 &kv,
998 MigrateOptions {
999 dry_run: true,
1000 finalize: true,
1001 },
1002 )
1003 .await
1004 .unwrap();
1005
1006 assert!(report.total_rekeyed() > 0);
1007 assert!(report.created_default_project);
1008 assert!(kv
1009 .get("project/default/current/blog")
1010 .await
1011 .unwrap()
1012 .is_none());
1013 assert!(kv.get("current/blog").await.unwrap().is_some());
1014 assert!(kv.get("projectmeta/default").await.unwrap().is_none());
1015 assert_eq!(status(&kv).await.unwrap(), Status::NeedsMigration);
1016 }
1017
1018 #[tokio::test]
1019 async fn dual_stage_then_finalize() {
1020 let kv = MemoryKv::new();
1021 seed_legacy(&kv).await;
1022
1023 let staged = migrate(
1024 &kv,
1025 MigrateOptions {
1026 dry_run: false,
1027 finalize: false,
1028 },
1029 )
1030 .await
1031 .unwrap();
1032 assert!(staged.dual);
1033 assert_eq!(status(&kv).await.unwrap(), Status::Dual);
1034 assert!(kv
1035 .get("project/default/current/blog")
1036 .await
1037 .unwrap()
1038 .is_some());
1039 assert!(
1040 kv.get("current/blog").await.unwrap().is_some(),
1041 "old key kept during dual soak"
1042 );
1043 assert_eq!(read_state(&kv).await.unwrap().unfinalized, vec![1]);
1045
1046 let done = finalize(&kv).await.unwrap();
1047 assert!(!done.dual);
1048 assert!(kv.get("current/blog").await.unwrap().is_none());
1049 assert_eq!(status(&kv).await.unwrap(), Status::Ready);
1050 }
1051
1052 #[tokio::test]
1053 async fn resumes_after_a_crash_mid_migration() {
1054 let kv = MemoryKv::new();
1055 seed_legacy(&kv).await;
1056
1057 for old in ["current/blog", "site/blog"] {
1060 let v = kv.get(old).await.unwrap().unwrap();
1061 kv.put(&format!("project/default/{old}"), v).await.unwrap();
1062 kv.delete(old).await.unwrap();
1063 }
1064 let partial = SchemaState {
1065 version: 0,
1066 unfinalized: Vec::new(),
1067 in_progress: Some(InProgress {
1068 target: 1,
1069 steps_done: vec![
1070 rekey_step_id("default", "current/"),
1071 rekey_step_id("default", "site/"),
1072 ],
1073 }),
1074 history: Vec::new(),
1075 updated_at: 1,
1076 };
1077 persist_state(&kv, &partial).await.unwrap();
1078
1079 migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
1080 assert_eq!(status(&kv).await.unwrap(), Status::Ready);
1081 assert_eq!(
1082 kv.get("project/default/compute/api")
1083 .await
1084 .unwrap()
1085 .as_deref(),
1086 Some(&b"{}"[..])
1087 );
1088 assert_eq!(
1089 kv.get("project/default/current/blog")
1090 .await
1091 .unwrap()
1092 .as_deref(),
1093 Some(&b"dep-1"[..])
1094 );
1095 assert!(kv.get("compute/api").await.unwrap().is_none());
1096 assert!(kv
1098 .get("project/default/project/default/current/blog")
1099 .await
1100 .unwrap()
1101 .is_none());
1102 }
1103
1104 #[tokio::test]
1105 async fn reads_the_pre_mechanism_legacy_marker() {
1106 let kv = MemoryKv::new();
1108 kv.put(
1109 SCHEMA_KEY,
1110 br#"{"layout":2,"dual":false,"migrated_at":9,"families_done":["current/"]}"#.to_vec(),
1111 )
1112 .await
1113 .unwrap();
1114 let state = read_state(&kv).await.unwrap();
1115 assert_eq!(state.version, 1);
1116 assert!(state.unfinalized.is_empty());
1117 assert_eq!(status(&kv).await.unwrap(), Status::Ready);
1118
1119 let dual = MemoryKv::new();
1121 dual.put(SCHEMA_KEY, br#"{"layout":2,"dual":true}"#.to_vec())
1122 .await
1123 .unwrap();
1124 assert_eq!(read_state(&dual).await.unwrap().unfinalized, vec![1]);
1125 assert_eq!(status(&dual).await.unwrap(), Status::Dual);
1126 }
1127
1128 struct AddSentinelV2;
1133
1134 #[async_trait]
1135 impl Migration for AddSentinelV2 {
1136 fn version(&self) -> u32 {
1137 2
1138 }
1139 fn id(&self) -> &'static str {
1140 "add-sentinel"
1141 }
1142 fn description(&self) -> &'static str {
1143 "test migration: write a sentinel key"
1144 }
1145 async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError> {
1146 Ok(kv.get("demo/v2").await?.is_none())
1147 }
1148 fn forward_steps(&self) -> Vec<Box<dyn Step>> {
1149 vec![Box::new(WriteSentinel)]
1150 }
1151 }
1152
1153 struct WriteSentinel;
1154 #[async_trait]
1155 impl Step for WriteSentinel {
1156 fn id(&self) -> String {
1157 "write-sentinel".to_string()
1158 }
1159 async fn run(
1160 &self,
1161 kv: &dyn KvStore,
1162 dry_run: bool,
1163 _report: &mut MigrationReport,
1164 ) -> Result<(), MigrateError> {
1165 if !dry_run {
1166 kv.put("demo/v2", b"ok".to_vec()).await?;
1167 }
1168 Ok(())
1169 }
1170 }
1171
1172 fn chain() -> Vec<Box<dyn Migration>> {
1173 vec![Box::new(ProjectRekeyV1), Box::new(AddSentinelV2)]
1174 }
1175
1176 #[tokio::test]
1177 async fn engine_applies_a_multi_migration_chain_in_order() {
1178 let kv = MemoryKv::new();
1179 seed_legacy(&kv).await;
1180
1181 run(&kv, MigrateOptions::one_shot(), &chain())
1182 .await
1183 .unwrap();
1184
1185 assert!(kv
1187 .get("project/default/current/blog")
1188 .await
1189 .unwrap()
1190 .is_some());
1191 assert_eq!(
1192 kv.get("demo/v2").await.unwrap().as_deref(),
1193 Some(&b"ok"[..])
1194 );
1195
1196 let state = read_state(&kv).await.unwrap();
1197 assert_eq!(state.version, 2);
1198 let ids: Vec<&str> = state.history.iter().map(|a| a.id.as_str()).collect();
1200 assert_eq!(ids, vec!["project-rekey", "add-sentinel"]);
1201
1202 let again = run(&kv, MigrateOptions::one_shot(), &chain())
1204 .await
1205 .unwrap();
1206 assert!(again.already_migrated);
1207 }
1208
1209 #[tokio::test]
1210 async fn engine_applies_only_pending_migrations_from_a_version() {
1211 let kv = MemoryKv::new();
1213 kv.put("project/default/site/blog", b"cfg".to_vec())
1215 .await
1216 .unwrap();
1217 persist_state(
1218 &kv,
1219 &SchemaState {
1220 version: 1,
1221 ..Default::default()
1222 },
1223 )
1224 .await
1225 .unwrap();
1226
1227 let report = run(&kv, MigrateOptions::one_shot(), &chain())
1228 .await
1229 .unwrap();
1230 assert!(report.rekeyed.is_empty());
1232 assert_eq!(
1233 kv.get("demo/v2").await.unwrap().as_deref(),
1234 Some(&b"ok"[..])
1235 );
1236 assert_eq!(read_state(&kv).await.unwrap().version, 2);
1237 }
1238}