use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use crate::kv::{KvStore, WriteOp};
use crate::project::{self, owner_kind, DomainOwner, DEFAULT_PROJECT};
use crate::time::now_unix;
pub const SCHEMA_KEY: &str = "schema/version";
pub const PREMIGRATION_INDEX_KEY: &str = "schema/premigration-index";
pub const CURRENT_VERSION: u32 = 1;
pub const MUTABLE_FAMILIES: &[&str] = &[
"current/",
"site/",
"history/",
"alias/",
"domainverify/",
"dnsmanaged/",
"functions/",
"metering/",
"blobnotify/",
"workflows/",
"compute/",
"compute_state/",
];
pub const DOMAIN_FAMILIES: &[&str] = &["domain/", "wildcard/", "httpchallenge/"];
pub const GLOBAL_FAMILIES: &[&str] = &[
"manifests/",
"meta/",
"siteconfig/",
"computever/",
"daemonconfig/",
"authz/",
"daemon/",
"cert/",
"project/",
"projectmeta/",
"projectver/",
"project-history/",
"owner/",
"schema/",
];
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct SchemaState {
pub version: u32,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub unfinalized: Vec<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub in_progress: Option<InProgress>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub history: Vec<AppliedRecord>,
pub updated_at: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct InProgress {
pub target: u32,
pub steps_done: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AppliedRecord {
pub version: u32,
pub id: String,
pub at: u64,
}
#[derive(Debug, Deserialize)]
struct LegacyMarker {
#[serde(default)]
layout: u32,
#[serde(default)]
dual: bool,
#[serde(default)]
families_done: Vec<String>,
}
impl LegacyMarker {
fn into_state(self) -> SchemaState {
if self.layout >= 2 {
SchemaState {
version: 1,
unfinalized: if self.dual { vec![1] } else { Vec::new() },
in_progress: None,
history: Vec::new(),
updated_at: now_unix(),
}
} else {
let steps_done: Vec<String> = self
.families_done
.iter()
.map(|f| rekey_step_id(DEFAULT_PROJECT, f))
.collect();
SchemaState {
version: 0,
unfinalized: Vec::new(),
in_progress: (!steps_done.is_empty()).then_some(InProgress {
target: 1,
steps_done,
}),
history: Vec::new(),
updated_at: now_unix(),
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Status {
Ready,
NeedsMigration,
Dual,
}
#[derive(Debug, Clone, Copy, Default)]
pub struct MigrateOptions {
pub dry_run: bool,
pub finalize: bool,
}
impl MigrateOptions {
pub fn one_shot() -> Self {
Self {
dry_run: false,
finalize: true,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct MigrationReport {
pub rekeyed: Vec<(String, usize)>,
pub values_rewritten: Vec<(String, usize)>,
pub owner_entries: usize,
pub created_default_project: bool,
pub already_migrated: bool,
pub dual: bool,
}
impl MigrationReport {
pub fn total_rekeyed(&self) -> usize {
self.rekeyed.iter().map(|(_, n)| n).sum()
}
}
#[derive(Debug, thiserror::Error)]
pub enum MigrateError {
#[error(transparent)]
Kv(#[from] crate::error::KvError),
#[error("verification failed for migrated key {0}")]
Verify(String),
#[error("migration serde error: {0}")]
Serde(String),
}
pub async fn read_state(kv: &dyn KvStore) -> Result<SchemaState, MigrateError> {
let Some(bytes) = kv.get(SCHEMA_KEY).await? else {
return Ok(SchemaState::default());
};
if let Ok(state) = serde_json::from_slice::<SchemaState>(&bytes) {
return Ok(state);
}
let legacy: LegacyMarker =
serde_json::from_slice(&bytes).map_err(|e| MigrateError::Serde(e.to_string()))?;
Ok(legacy.into_state())
}
async fn persist_state(kv: &dyn KvStore, state: &SchemaState) -> Result<(), MigrateError> {
let bytes = serde_json::to_vec(state).map_err(|e| MigrateError::Serde(e.to_string()))?;
kv.put(SCHEMA_KEY, bytes).await?;
Ok(())
}
pub async fn status(kv: &dyn KvStore) -> Result<Status, MigrateError> {
status_of(kv, ®istry()).await
}
async fn status_of(
kv: &dyn KvStore,
migrations: &[Box<dyn Migration>],
) -> Result<Status, MigrateError> {
let state = read_state(kv).await?;
if state.in_progress.is_some() {
return Ok(Status::NeedsMigration);
}
if !state.unfinalized.is_empty() {
return Ok(Status::Dual);
}
for m in migrations.iter().filter(|m| m.version() > state.version) {
if m.is_applicable(kv).await? {
return Ok(Status::NeedsMigration);
}
}
Ok(Status::Ready)
}
pub async fn migrate(
kv: &dyn KvStore,
opts: MigrateOptions,
) -> Result<MigrationReport, MigrateError> {
run(kv, opts, ®istry()).await
}
pub async fn finalize(kv: &dyn KvStore) -> Result<MigrationReport, MigrateError> {
migrate(kv, MigrateOptions::one_shot()).await
}
async fn run(
kv: &dyn KvStore,
opts: MigrateOptions,
migrations: &[Box<dyn Migration>],
) -> Result<MigrationReport, MigrateError> {
let mut state = read_state(kv).await?;
let mut report = MigrationReport::default();
let has_pending_forward = migrations.iter().any(|m| m.version() > state.version);
if !has_pending_forward && state.unfinalized.is_empty() && state.in_progress.is_none() {
report.already_migrated = true;
return Ok(report);
}
if opts.dry_run {
for m in migrations.iter().filter(|m| m.version() > state.version) {
if m.is_applicable(kv).await? {
for step in m.forward_steps() {
step.run(kv, true, &mut report).await?;
}
}
}
report.dual = !opts.finalize;
return Ok(report);
}
if kv.get(PREMIGRATION_INDEX_KEY).await?.is_none() {
write_premigration_index(kv).await?;
}
for m in migrations {
if m.version() <= state.version {
continue;
}
let resuming = state
.in_progress
.as_ref()
.is_some_and(|ip| ip.target == m.version());
if !resuming && !m.is_applicable(kv).await? {
state.version = m.version();
state.history.push(AppliedRecord {
version: m.version(),
id: m.id().to_string(),
at: now_unix(),
});
state.updated_at = now_unix();
persist_state(kv, &state).await?;
continue;
}
let mut done = if resuming {
state
.in_progress
.take()
.map(|ip| ip.steps_done)
.unwrap_or_default()
} else {
Vec::new()
};
for step in m.forward_steps() {
let sid = step.id();
if done.contains(&sid) {
continue;
}
step.run(kv, false, &mut report).await?;
done.push(sid);
state.in_progress = Some(InProgress {
target: m.version(),
steps_done: done.clone(),
});
state.updated_at = now_unix();
persist_state(kv, &state).await?;
}
state.version = m.version();
state.in_progress = None;
state.history.push(AppliedRecord {
version: m.version(),
id: m.id().to_string(),
at: now_unix(),
});
if m.online() {
state.unfinalized.push(m.version());
}
state.updated_at = now_unix();
persist_state(kv, &state).await?;
}
if opts.finalize && !state.unfinalized.is_empty() {
let mut pending = state.unfinalized.clone();
pending.sort_unstable();
for v in pending {
if let Some(m) = migrations.iter().find(|m| m.version() == v) {
for step in m.cleanup_steps() {
step.run(kv, false, &mut report).await?;
}
}
}
state.unfinalized.clear();
state.updated_at = now_unix();
persist_state(kv, &state).await?;
}
report.dual = !state.unfinalized.is_empty();
Ok(report)
}
async fn write_premigration_index(kv: &dyn KvStore) -> Result<(), MigrateError> {
let mut keys = Vec::new();
for family in MUTABLE_FAMILIES.iter().chain(DOMAIN_FAMILIES) {
keys.extend(kv.list_prefix(family).await?);
}
keys.sort();
kv.put(PREMIGRATION_INDEX_KEY, keys.join("\n").into_bytes())
.await?;
Ok(())
}
fn registry() -> Vec<Box<dyn Migration>> {
vec![Box::new(ProjectRekeyV1)]
}
#[async_trait]
trait Migration: Send + Sync {
fn version(&self) -> u32;
fn id(&self) -> &'static str;
#[allow(dead_code)]
fn description(&self) -> &'static str;
async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError>;
fn forward_steps(&self) -> Vec<Box<dyn Step>>;
fn cleanup_steps(&self) -> Vec<Box<dyn Step>> {
Vec::new()
}
fn online(&self) -> bool {
!self.cleanup_steps().is_empty()
}
}
struct ProjectRekeyV1;
#[async_trait]
impl Migration for ProjectRekeyV1 {
fn version(&self) -> u32 {
1
}
fn id(&self) -> &'static str {
"project-rekey"
}
fn description(&self) -> &'static str {
"re-key the pre-0.2.0 store under project/<default>/ and add the project-scoped index"
}
async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError> {
for family in MUTABLE_FAMILIES.iter().chain(DOMAIN_FAMILIES) {
if !kv.list_prefix(family).await?.is_empty() {
return Ok(true);
}
}
Ok(false)
}
fn forward_steps(&self) -> Vec<Box<dyn Step>> {
let mut steps: Vec<Box<dyn Step>> = Vec::new();
for family in MUTABLE_FAMILIES {
steps.push(Box::new(RekeyFamily {
family,
project: DEFAULT_PROJECT,
}));
}
for family in DOMAIN_FAMILIES {
steps.push(Box::new(RewriteValues {
family,
transform: domain_owner_canonical,
}));
}
steps.push(Box::new(EnsureDefaultProject));
steps.push(Box::new(BuildOwnerIndex));
steps
}
fn cleanup_steps(&self) -> Vec<Box<dyn Step>> {
MUTABLE_FAMILIES
.iter()
.map(|family| {
Box::new(DeleteOldFamily {
family,
project: DEFAULT_PROJECT,
}) as Box<dyn Step>
})
.collect()
}
}
#[async_trait]
trait Step: Send + Sync {
fn id(&self) -> String;
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
report: &mut MigrationReport,
) -> Result<(), MigrateError>;
}
fn rekey_step_id(project: &str, family: &str) -> String {
format!("rekey:{project}:{family}")
}
struct RekeyFamily {
family: &'static str,
project: &'static str,
}
#[async_trait]
impl Step for RekeyFamily {
fn id(&self) -> String {
rekey_step_id(self.project, self.family)
}
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
report: &mut MigrationReport,
) -> Result<(), MigrateError> {
let mut moved = 0;
for old_key in kv.list_prefix(self.family).await? {
let new_key = format!("project/{}/{}", self.project, old_key);
if dry_run {
moved += 1;
continue;
}
let Some(value) = kv.get(&old_key).await? else {
continue; };
kv.put(&new_key, value.clone()).await?;
if kv.get(&new_key).await?.as_deref() != Some(value.as_slice()) {
return Err(MigrateError::Verify(new_key));
}
moved += 1;
}
if moved > 0 {
report.rekeyed.push((self.family.to_string(), moved));
}
Ok(())
}
}
struct DeleteOldFamily {
family: &'static str,
project: &'static str,
}
#[async_trait]
impl Step for DeleteOldFamily {
fn id(&self) -> String {
format!("delete:{}:{}", self.project, self.family)
}
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
_report: &mut MigrationReport,
) -> Result<(), MigrateError> {
for old_key in kv.list_prefix(self.family).await? {
let new_key = format!("project/{}/{}", self.project, old_key);
if dry_run {
continue;
}
match (kv.get(&old_key).await?, kv.get(&new_key).await?) {
(Some(o), Some(n)) if o == n => kv.delete(&old_key).await?,
(Some(_), _) => return Err(MigrateError::Verify(new_key)),
(None, _) => {}
}
}
Ok(())
}
}
struct RewriteValues {
family: &'static str,
transform: fn(&[u8]) -> Vec<u8>,
}
#[async_trait]
impl Step for RewriteValues {
fn id(&self) -> String {
format!("rewrite:{}", self.family)
}
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
report: &mut MigrationReport,
) -> Result<(), MigrateError> {
let mut rewritten = 0;
for key in kv.list_prefix(self.family).await? {
let Some(value) = kv.get(&key).await? else {
continue;
};
let canonical = (self.transform)(&value);
if canonical != value {
if !dry_run {
kv.put(&key, canonical).await?;
}
rewritten += 1;
}
}
if rewritten > 0 {
report
.values_rewritten
.push((self.family.to_string(), rewritten));
}
Ok(())
}
}
fn domain_owner_canonical(value: &[u8]) -> Vec<u8> {
DomainOwner::from_bytes(value).to_bytes()
}
struct EnsureDefaultProject;
#[async_trait]
impl Step for EnsureDefaultProject {
fn id(&self) -> String {
"ensure-default-project".to_string()
}
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
report: &mut MigrationReport,
) -> Result<(), MigrateError> {
let pointer = project::pointer_key(DEFAULT_PROJECT);
if kv.get(&pointer).await?.is_some() {
return Ok(());
}
report.created_default_project = true;
if dry_run {
return Ok(());
}
let default = crate::deploy::DeployStore::default_project_record();
let hash = default.id();
let body = serde_json::to_vec(&default).map_err(|e| MigrateError::Serde(e.to_string()))?;
kv.write_batch(vec![
WriteOp::Put(project::spec_key(&hash), body),
WriteOp::Put(pointer, hash.into_bytes()),
])
.await?;
Ok(())
}
}
struct BuildOwnerIndex;
#[async_trait]
impl Step for BuildOwnerIndex {
fn id(&self) -> String {
"build-owner-index".to_string()
}
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
report: &mut MigrationReport,
) -> Result<(), MigrateError> {
let mut ops = Vec::new();
let site_prefix = format!("project/{DEFAULT_PROJECT}/site/");
for key in kv.list_prefix(&site_prefix).await? {
if let Some(site) = key.strip_prefix(&site_prefix) {
if !site.is_empty() {
ops.push(WriteOp::Put(
project::owner_key(owner_kind::SITE, site),
DEFAULT_PROJECT.as_bytes().to_vec(),
));
}
}
}
let fn_prefix = format!("project/{DEFAULT_PROJECT}/functions/");
for key in kv.list_prefix(&fn_prefix).await? {
if let Some(rest) = key.strip_prefix(&fn_prefix) {
if !rest.is_empty() && !rest.contains('/') {
ops.push(WriteOp::Put(
project::owner_key(owner_kind::FUNCTION, rest),
DEFAULT_PROJECT.as_bytes().to_vec(),
));
}
}
}
let compute_prefix = format!("project/{DEFAULT_PROJECT}/compute/");
for key in kv.list_prefix(&compute_prefix).await? {
if let Some(name) = key.strip_prefix(&compute_prefix) {
if !name.is_empty() {
ops.push(WriteOp::Put(
project::owner_key(owner_kind::COMPUTE, name),
DEFAULT_PROJECT.as_bytes().to_vec(),
));
}
}
}
report.owner_entries += ops.len();
if !dry_run && !ops.is_empty() {
kv.write_batch(ops).await?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::kv::MemoryKv;
async fn seed_legacy(kv: &MemoryKv) {
kv.put("current/blog", b"dep-1".to_vec()).await.unwrap();
kv.put("site/blog", b"cfghash".to_vec()).await.unwrap();
kv.put("history/blog", b"[]".to_vec()).await.unwrap();
kv.put("alias/blog/staging", b"dep-1".to_vec())
.await
.unwrap();
kv.put("domainverify/blog/www.example", b"{}".to_vec())
.await
.unwrap();
kv.put("dnsmanaged/blog/www.example", b"{}".to_vec())
.await
.unwrap();
kv.put("functions/resize", b"{}".to_vec()).await.unwrap();
kv.put("functions/resize/versions/v1", b"{}".to_vec())
.await
.unwrap();
kv.put("metering/resize", b"{}".to_vec()).await.unwrap();
kv.put("blobnotify/resize/uploads", b"{}".to_vec())
.await
.unwrap();
kv.put("workflows/etl", b"{}".to_vec()).await.unwrap();
kv.put("compute/api", b"{}".to_vec()).await.unwrap();
kv.put("compute_state/api/0", b"{}".to_vec()).await.unwrap();
kv.put("domain/www.example", b"blog".to_vec())
.await
.unwrap();
kv.put("wildcard/preview.example", b"blog".to_vec())
.await
.unwrap();
kv.put("httpchallenge/www.example/tok", b"blog".to_vec())
.await
.unwrap();
kv.put("siteconfig/cfghash", b"the-config".to_vec())
.await
.unwrap();
kv.put("manifests/dep-1", b"the-manifest".to_vec())
.await
.unwrap();
kv.put("authz/tokens/t1", b"tok".to_vec()).await.unwrap();
}
#[test]
fn registry_versions_are_strictly_ascending_and_reach_current() {
let reg = registry();
assert!(!reg.is_empty());
let mut last = 0;
for m in ® {
assert!(m.version() > last, "versions must strictly ascend");
last = m.version();
}
assert_eq!(
last, CURRENT_VERSION,
"CURRENT_VERSION == the top migration"
);
}
#[tokio::test]
async fn status_detects_legacy_fresh_and_migrated() {
let fresh = MemoryKv::new();
assert_eq!(status(&fresh).await.unwrap(), Status::Ready);
let legacy = MemoryKv::new();
seed_legacy(&legacy).await;
assert_eq!(status(&legacy).await.unwrap(), Status::NeedsMigration);
migrate(&legacy, MigrateOptions::one_shot()).await.unwrap();
assert_eq!(status(&legacy).await.unwrap(), Status::Ready);
let state = read_state(&legacy).await.unwrap();
assert_eq!(state.version, CURRENT_VERSION);
assert!(state.history.iter().any(|a| a.id == "project-rekey"));
}
#[tokio::test]
async fn migrate_rekeys_mutable_families_and_rewrites_domain_values() {
let kv = MemoryKv::new();
seed_legacy(&kv).await;
let report = migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
assert_eq!(
kv.get("project/default/current/blog")
.await
.unwrap()
.as_deref(),
Some(&b"dep-1"[..])
);
assert_eq!(
kv.get("project/default/functions/resize/versions/v1")
.await
.unwrap()
.as_deref(),
Some(&b"{}"[..])
);
assert_eq!(
kv.get("project/default/compute_state/api/0")
.await
.unwrap()
.as_deref(),
Some(&b"{}"[..])
);
assert!(kv.get("current/blog").await.unwrap().is_none());
assert!(kv.get("compute/api").await.unwrap().is_none());
let dv = kv.get("domain/www.example").await.unwrap().unwrap();
assert_eq!(
DomainOwner::from_bytes(&dv),
DomainOwner::new("default", "blog")
);
assert!(String::from_utf8_lossy(&dv).contains("\"project\":\"default\""));
assert_eq!(
DomainOwner::from_bytes(
&kv.get("httpchallenge/www.example/tok")
.await
.unwrap()
.unwrap()
),
DomainOwner::new("default", "blog")
);
assert_eq!(
kv.get("siteconfig/cfghash").await.unwrap().as_deref(),
Some(&b"the-config"[..])
);
assert_eq!(
kv.get("authz/tokens/t1").await.unwrap().as_deref(),
Some(&b"tok"[..])
);
assert!(kv.get("projectmeta/default").await.unwrap().is_some());
assert!(report.created_default_project);
assert_eq!(
kv.get("owner/site/blog").await.unwrap().as_deref(),
Some(&b"default"[..])
);
assert_eq!(
kv.get("owner/function/resize").await.unwrap().as_deref(),
Some(&b"default"[..])
);
assert!(kv
.get("owner/function/resize/versions/v1")
.await
.unwrap()
.is_none());
assert_eq!(status(&kv).await.unwrap(), Status::Ready);
assert!(!report.already_migrated);
}
#[tokio::test]
async fn migrate_is_idempotent() {
let kv = MemoryKv::new();
seed_legacy(&kv).await;
let first = migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
assert!(first.total_rekeyed() > 0);
let before: Vec<String> = kv.list_prefix("project/").await.unwrap();
let second = migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
assert!(second.already_migrated);
assert_eq!(second.total_rekeyed(), 0);
assert_eq!(kv.list_prefix("project/").await.unwrap(), before);
}
#[tokio::test]
async fn dry_run_writes_nothing() {
let kv = MemoryKv::new();
seed_legacy(&kv).await;
let report = migrate(
&kv,
MigrateOptions {
dry_run: true,
finalize: true,
},
)
.await
.unwrap();
assert!(report.total_rekeyed() > 0);
assert!(report.created_default_project);
assert!(kv
.get("project/default/current/blog")
.await
.unwrap()
.is_none());
assert!(kv.get("current/blog").await.unwrap().is_some());
assert!(kv.get("projectmeta/default").await.unwrap().is_none());
assert_eq!(status(&kv).await.unwrap(), Status::NeedsMigration);
}
#[tokio::test]
async fn dual_stage_then_finalize() {
let kv = MemoryKv::new();
seed_legacy(&kv).await;
let staged = migrate(
&kv,
MigrateOptions {
dry_run: false,
finalize: false,
},
)
.await
.unwrap();
assert!(staged.dual);
assert_eq!(status(&kv).await.unwrap(), Status::Dual);
assert!(kv
.get("project/default/current/blog")
.await
.unwrap()
.is_some());
assert!(
kv.get("current/blog").await.unwrap().is_some(),
"old key kept during dual soak"
);
assert_eq!(read_state(&kv).await.unwrap().unfinalized, vec![1]);
let done = finalize(&kv).await.unwrap();
assert!(!done.dual);
assert!(kv.get("current/blog").await.unwrap().is_none());
assert_eq!(status(&kv).await.unwrap(), Status::Ready);
}
#[tokio::test]
async fn resumes_after_a_crash_mid_migration() {
let kv = MemoryKv::new();
seed_legacy(&kv).await;
for old in ["current/blog", "site/blog"] {
let v = kv.get(old).await.unwrap().unwrap();
kv.put(&format!("project/default/{old}"), v).await.unwrap();
kv.delete(old).await.unwrap();
}
let partial = SchemaState {
version: 0,
unfinalized: Vec::new(),
in_progress: Some(InProgress {
target: 1,
steps_done: vec![
rekey_step_id("default", "current/"),
rekey_step_id("default", "site/"),
],
}),
history: Vec::new(),
updated_at: 1,
};
persist_state(&kv, &partial).await.unwrap();
migrate(&kv, MigrateOptions::one_shot()).await.unwrap();
assert_eq!(status(&kv).await.unwrap(), Status::Ready);
assert_eq!(
kv.get("project/default/compute/api")
.await
.unwrap()
.as_deref(),
Some(&b"{}"[..])
);
assert_eq!(
kv.get("project/default/current/blog")
.await
.unwrap()
.as_deref(),
Some(&b"dep-1"[..])
);
assert!(kv.get("compute/api").await.unwrap().is_none());
assert!(kv
.get("project/default/project/default/current/blog")
.await
.unwrap()
.is_none());
}
#[tokio::test]
async fn reads_the_pre_mechanism_legacy_marker() {
let kv = MemoryKv::new();
kv.put(
SCHEMA_KEY,
br#"{"layout":2,"dual":false,"migrated_at":9,"families_done":["current/"]}"#.to_vec(),
)
.await
.unwrap();
let state = read_state(&kv).await.unwrap();
assert_eq!(state.version, 1);
assert!(state.unfinalized.is_empty());
assert_eq!(status(&kv).await.unwrap(), Status::Ready);
let dual = MemoryKv::new();
dual.put(SCHEMA_KEY, br#"{"layout":2,"dual":true}"#.to_vec())
.await
.unwrap();
assert_eq!(read_state(&dual).await.unwrap().unfinalized, vec![1]);
assert_eq!(status(&dual).await.unwrap(), Status::Dual);
}
struct AddSentinelV2;
#[async_trait]
impl Migration for AddSentinelV2 {
fn version(&self) -> u32 {
2
}
fn id(&self) -> &'static str {
"add-sentinel"
}
fn description(&self) -> &'static str {
"test migration: write a sentinel key"
}
async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError> {
Ok(kv.get("demo/v2").await?.is_none())
}
fn forward_steps(&self) -> Vec<Box<dyn Step>> {
vec![Box::new(WriteSentinel)]
}
}
struct WriteSentinel;
#[async_trait]
impl Step for WriteSentinel {
fn id(&self) -> String {
"write-sentinel".to_string()
}
async fn run(
&self,
kv: &dyn KvStore,
dry_run: bool,
_report: &mut MigrationReport,
) -> Result<(), MigrateError> {
if !dry_run {
kv.put("demo/v2", b"ok".to_vec()).await?;
}
Ok(())
}
}
fn chain() -> Vec<Box<dyn Migration>> {
vec![Box::new(ProjectRekeyV1), Box::new(AddSentinelV2)]
}
#[tokio::test]
async fn engine_applies_a_multi_migration_chain_in_order() {
let kv = MemoryKv::new();
seed_legacy(&kv).await;
run(&kv, MigrateOptions::one_shot(), &chain())
.await
.unwrap();
assert!(kv
.get("project/default/current/blog")
.await
.unwrap()
.is_some());
assert_eq!(
kv.get("demo/v2").await.unwrap().as_deref(),
Some(&b"ok"[..])
);
let state = read_state(&kv).await.unwrap();
assert_eq!(state.version, 2);
let ids: Vec<&str> = state.history.iter().map(|a| a.id.as_str()).collect();
assert_eq!(ids, vec!["project-rekey", "add-sentinel"]);
let again = run(&kv, MigrateOptions::one_shot(), &chain())
.await
.unwrap();
assert!(again.already_migrated);
}
#[tokio::test]
async fn engine_applies_only_pending_migrations_from_a_version() {
let kv = MemoryKv::new();
kv.put("project/default/site/blog", b"cfg".to_vec())
.await
.unwrap();
persist_state(
&kv,
&SchemaState {
version: 1,
..Default::default()
},
)
.await
.unwrap();
let report = run(&kv, MigrateOptions::one_shot(), &chain())
.await
.unwrap();
assert!(report.rekeyed.is_empty());
assert_eq!(
kv.get("demo/v2").await.unwrap().as_deref(),
Some(&b"ok"[..])
);
assert_eq!(read_state(&kv).await.unwrap().version, 2);
}
}