Skip to main content

boatramp_core/
migrate.rs

1//! The control-plane **store migration mechanism**: a versioned, ordered registry of
2//! forward-only migrations the engine walks to bring a store up to the layout this
3//! binary requires.
4//!
5//! ## The mechanism
6//!
7//! - A single monotonic **schema version** ([`SchemaState::version`]) is recorded in
8//!   the [`SCHEMA_KEY`] marker, alongside a crash-resume cursor
9//!   ([`SchemaState::in_progress`]), the set of migrations whose destructive cleanup is
10//!   still deferred ([`SchemaState::unfinalized`]), and an applied-[`AppliedRecord`]
11//!   history for audit.
12//! - Each [`Migration`] declares its target [`version`](Migration::version), a cheap
13//!   [`is_applicable`](Migration::is_applicable) check, an ordered list of
14//!   non-destructive [`forward_steps`](Migration::forward_steps), and (for an *online*
15//!   migration) a list of destructive [`cleanup_steps`](Migration::cleanup_steps)
16//!   deferred to `--finalize`.
17//! - The [`registry`] lists every migration in ascending version order. A store at
18//!   version `N` applies every migration with `version > N`, in order.
19//! - Each [`Step`] is an idempotent, re-verifying unit; the engine records each
20//!   completed step in the resume cursor and persists after every step, so a crash
21//!   re-runs only the interrupted step. The reusable steps ([`RekeyFamily`],
22//!   [`RewriteValues`], [`DeleteOldFamily`], …) are the building blocks a new migration
23//!   composes rather than re-implementing copy/verify/resume.
24//!
25//! ## Online migrations + the dual soak
26//!
27//! An *online* migration writes its new state and leaves the old readable, deferring
28//! the destructive delete to a `finalize` pass — so an operator can `--stage` (copy +
29//! verify, serve off the new state with an old-state fallback), soak, then `--finalize`
30//! (delete). A migration with no cleanup steps is *simple* (its forward is the whole
31//! thing). `MigrateOptions::finalize` (the default / `one_shot`) runs cleanup inline;
32//! `--stage` (`finalize = false`) leaves the store in the `dual` soak.
33//!
34//! ## Safety
35//!
36//! Forward-only (no down-migrations; rollback is a backup, or a revert within the
37//! pre-finalize dual window while the old keys still exist). Copy-before-delete:
38//! [`RekeyFamily`] copies + read-back verifies and never deletes; [`DeleteOldFamily`]
39//! re-reads and byte-compares before each delete and refuses to delete an old key whose
40//! new key is absent/mismatched — so a partially-copied family never loses data. Every
41//! step is idempotent, so a re-run (or a concurrent racer in a cluster) converges
42//! rather than corrupts. The startup guard ([`status`]) refuses to serve a store below
43//! the current version.
44//!
45//! ## Cluster
46//!
47//! The migration runs through the same [`KvStore`] a node serves from; in a Raft
48//! cluster that is the replicated store, so a single leader-run migration replicates
49//! for free. Prefer running `boatramp migrate` **once** (against the leader) before a
50//! rolling upgrade. The startup guard does not elect a leader or block followers, so if
51//! several nodes start with `--auto-migrate` against a still-out-of-date store they may
52//! run concurrently — which is **safe, only redundant**: every step is idempotent and
53//! re-verifying, so racers converge. The marker cursor is a last-writer-wins `put` (no
54//! CAS), so a race can at worst redo a step, never skip a delete-guard.
55
56use 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
63/// The global marker key recording the store's schema version + migration progress.
64pub const SCHEMA_KEY: &str = "schema/version";
65
66/// A pre-migration snapshot of the full key list, written once before any change so an
67/// operator can diff/audit what existed before the first migration ran.
68pub const PREMIGRATION_INDEX_KEY: &str = "schema/premigration-index";
69
70/// The latest schema version this binary knows how to reach — the highest-versioned
71/// entry in [`registry`]. A store below this must be migrated before it will serve.
72pub const CURRENT_VERSION: u32 = 1;
73
74/// The mutable per-name families the project re-key (migration 1) moves under
75/// `project/<default>/…`. Order is irrelevant to correctness (families are independent)
76/// but fixed for a legible, resumable progress record.
77pub 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
92/// The domain-routing families: the **key stays global**, the **value** is rewritten
93/// from a bare site name to a `{project, site}` [`DomainOwner`].
94pub const DOMAIN_FAMILIES: &[&str] = &["domain/", "wildcard/", "httpchallenge/"];
95
96/// Families that are **never** migrated by the project re-key: content-addressed bodies
97/// (dedup-shared, stable across layouts) and control-plane singletons. Listed for
98/// documentation + the migration's own guard against clobbering them. (Blob bodies live
99/// under a two-hex-char shard prefix in the *blob* store, a different backend entirely.)
100pub const GLOBAL_FAMILIES: &[&str] = &[
101    "manifests/",
102    "meta/",
103    "siteconfig/",
104    "computever/",
105    "daemonconfig/",
106    "authz/",
107    "daemon/",
108    "cert/",
109    // The 0.2.0 namespaces themselves — already the project layout, must not be moved.
110    "project/",
111    "projectmeta/",
112    "projectver/",
113    "project-history/",
114    "owner/",
115    "schema/",
116];
117
118// ---- the marker ------------------------------------------------------------------
119
120/// The persisted schema-version marker + migration progress.
121#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
122#[serde(default, deny_unknown_fields)]
123pub struct SchemaState {
124    /// The highest **forward-applied** migration version — the layout the store serves
125    /// at. `0` = never migrated (a fresh store born at the current layout, or a
126    /// pre-migration legacy store).
127    pub version: u32,
128    /// Migrations whose forward is applied but whose destructive cleanup is deferred (a
129    /// `dual` soak, ascending). Empty = nothing awaiting `--finalize`.
130    #[serde(skip_serializing_if = "Vec::is_empty")]
131    pub unfinalized: Vec<u32>,
132    /// A forward phase interrupted mid-flight — the crash-resume cursor. Present only
133    /// while a migration's forward steps are running.
134    #[serde(skip_serializing_if = "Option::is_none")]
135    pub in_progress: Option<InProgress>,
136    /// Applied migrations, for audit.
137    #[serde(skip_serializing_if = "Vec::is_empty")]
138    pub history: Vec<AppliedRecord>,
139    /// Unix time of the last marker write.
140    pub updated_at: u64,
141}
142
143/// A forward phase mid-flight: the migration being applied + the ids of its completed
144/// steps (the resume cursor).
145#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
146pub struct InProgress {
147    /// The version whose forward is running.
148    pub target: u32,
149    /// Ids of the forward steps already completed (skipped on resume).
150    pub steps_done: Vec<String>,
151}
152
153/// One applied migration, recorded in [`SchemaState::history`].
154#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
155pub struct AppliedRecord {
156    /// The migration's target version.
157    pub version: u32,
158    /// The migration's stable id.
159    pub id: String,
160    /// Unix time it was fully applied.
161    pub at: u64,
162}
163
164/// The pre-mechanism (0.2.0-preview) marker shape, mapped onto [`SchemaState`] by the
165/// tolerant reader so a store written by the original one-shot engine is understood.
166#[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            // The project re-key (v1) forward is done; `dual` means cleanup pending.
180            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            // Layout 1: unmigrated. A partial `families_done` resumes v1's forward — its
189            // family names map onto the new per-family step ids.
190            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// ---- public API (stable: consumed by the CLI + serve startup guard) --------------
210
211/// How a store needs to be treated on startup.
212#[derive(Debug, Clone, Copy, PartialEq, Eq)]
213pub enum Status {
214    /// At the current version with nothing pending (or empty) — serve as-is.
215    Ready,
216    /// Below the current version with applicable work (or a forward interrupted
217    /// mid-flight) — refuse to serve until migrated.
218    NeedsMigration,
219    /// Forward done, destructive cleanup deferred (`dual` soak) — serve OK; a
220    /// `finalize` pass is the only remaining work.
221    Dual,
222}
223
224/// Options controlling a migration pass.
225#[derive(Debug, Clone, Copy, Default)]
226pub struct MigrateOptions {
227    /// Scan and report what *would* change, writing nothing.
228    pub dry_run: bool,
229    /// Run each online migration's destructive cleanup inline (a one-shot). With
230    /// `false` the pass stops after the forward phase, leaving old state for a
231    /// soak/rollback window (`dual`); a later `finalize` runs the cleanups.
232    pub finalize: bool,
233}
234
235impl MigrateOptions {
236    /// The default one-shot migration: forward + cleanup, to the current version.
237    pub fn one_shot() -> Self {
238        Self {
239            dry_run: false,
240            finalize: true,
241        }
242    }
243}
244
245/// A per-pass tally of what a migration run moved (or, on a dry run, would move).
246/// Aggregated across every migration applied in the pass.
247#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
248pub struct MigrationReport {
249    /// Keys re-keyed per family (`family` → count).
250    pub rekeyed: Vec<(String, usize)>,
251    /// Index values rewritten per family.
252    pub values_rewritten: Vec<(String, usize)>,
253    /// Owner reverse-index entries written.
254    pub owner_entries: usize,
255    /// Whether the pass created the `default` project pointer.
256    pub created_default_project: bool,
257    /// The store was already at the current version — the pass was a no-op.
258    pub already_migrated: bool,
259    /// The store is left in the `dual` soak (forward done, cleanup deferred).
260    pub dual: bool,
261}
262
263impl MigrationReport {
264    /// Total keys re-keyed across all families.
265    pub fn total_rekeyed(&self) -> usize {
266        self.rekeyed.iter().map(|(_, n)| n).sum()
267    }
268}
269
270/// A migration failure.
271#[derive(Debug, thiserror::Error)]
272pub enum MigrateError {
273    /// The underlying KV store failed.
274    #[error(transparent)]
275    Kv(#[from] crate::error::KvError),
276    /// A copied datum failed read-back verification — the pass aborts rather than risk
277    /// deleting the source of a datum that did not land.
278    #[error("verification failed for migrated key {0}")]
279    Verify(String),
280    /// (De)serializing the marker / a record failed.
281    #[error("migration serde error: {0}")]
282    Serde(String),
283}
284
285/// Read the schema-version marker, tolerantly mapping the pre-mechanism marker shape.
286/// A fresh store (no marker) reads as version 0.
287pub 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    // Not the current shape → the pre-mechanism `{layout, dual, families_done}` marker.
295    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
306/// Classify a store for the serve-startup guard.
307pub async fn status(kv: &dyn KvStore) -> Result<Status, MigrateError> {
308    status_of(kv, &registry()).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    // A forward interrupted mid-flight must be resumed before serving.
317    if state.in_progress.is_some() {
318        return Ok(Status::NeedsMigration);
319    }
320    // Forward done, cleanup deferred: serve OK, finalize pending.
321    if !state.unfinalized.is_empty() {
322        return Ok(Status::Dual);
323    }
324    // Any pending migration with real work refuses serving; a fresh store whose pending
325    // migrations are all inapplicable is Ready (new writes already land at the current
326    // layout).
327    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
335/// Migrate `kv` up to the current version, resumable + idempotent. Applies every
336/// registered migration whose version exceeds the store's, in order; forward phases
337/// are non-destructive, and (unless `finalize`) destructive cleanups are left for a
338/// `dual` soak. Safe to call on an up-to-date store (no-op) and safe to re-run after
339/// an interruption (resumes from the marker).
340pub async fn migrate(
341    kv: &dyn KvStore,
342    opts: MigrateOptions,
343) -> Result<MigrationReport, MigrateError> {
344    run(kv, opts, &registry()).await
345}
346
347/// Complete a store staged into the `dual` soak (and apply any pending forward): a
348/// one-shot pass that runs every deferred cleanup and flips to the current version.
349pub async fn finalize(kv: &dyn KvStore) -> Result<MigrationReport, MigrateError> {
350    migrate(kv, MigrateOptions::one_shot()).await
351}
352
353/// The engine core, parameterized by the migration list so tests can drive it with a
354/// synthetic chain.
355async 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    // A dry run reports the pending forward work without persisting anything.
370    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    // Snapshot the pre-migration key list once (best-effort, never overwriting).
383    if kv.get(PREMIGRATION_INDEX_KEY).await?.is_none() {
384        write_premigration_index(kv).await?;
385    }
386
387    // Forward phase: apply each pending migration in order.
388    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        // A migration with no work on this store still advances the version (records it
398        // as vacuously applied), so a fresh store reaches the current version.
399        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        // Run the forward steps, skipping those already recorded, persisting the cursor
412        // after each so a crash resumes at the next step.
413        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        // Forward complete: the version advances; the destructive cleanup is deferred.
438        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    // Finalize: run the deferred cleanups (this pass's + any pre-staged), oldest first.
453    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
472/// Write the pre-migration key snapshot (a newline-joined list of every key that a
473/// migration might touch, for audit/rollback). Best-effort.
474async 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
485// ---- the migration registry ------------------------------------------------------
486
487/// Every migration this binary can apply, in **ascending version order**. A store at
488/// version `N` applies each entry with `version > N`. Append the next breaking store
489/// change here as a new [`Migration`] with the next version; never renumber or reorder
490/// an existing one.
491fn registry() -> Vec<Box<dyn Migration>> {
492    vec![Box::new(ProjectRekeyV1)]
493}
494
495/// One ordered, versioned, forward-only migration.
496#[async_trait]
497trait Migration: Send + Sync {
498    /// The schema version this migration produces (unique + ascending across the
499    /// registry).
500    fn version(&self) -> u32;
501    /// A stable id for logs + the marker history.
502    fn id(&self) -> &'static str;
503    /// A one-line description.
504    #[allow(dead_code)]
505    fn description(&self) -> &'static str;
506    /// Whether this store has anything for this migration to do. A cheap detector, so a
507    /// fresh store (born at the current layout) is `Ready` without running the
508    /// migration.
509    async fn is_applicable(&self, kv: &dyn KvStore) -> Result<bool, MigrateError>;
510    /// The ordered, non-destructive forward steps that make the store serve at this
511    /// version.
512    fn forward_steps(&self) -> Vec<Box<dyn Step>>;
513    /// The destructive cleanup steps (delete old state), deferred to `finalize`. Empty
514    /// ⇒ a *simple* migration.
515    fn cleanup_steps(&self) -> Vec<Box<dyn Step>> {
516        Vec::new()
517    }
518    /// Whether this migration has a deferrable destructive phase (supports the dual
519    /// soak).
520    fn online(&self) -> bool {
521        !self.cleanup_steps().is_empty()
522    }
523}
524
525/// Migration **1** — re-key the flat, pre-0.2.0 store to the project-scoped layout: move
526/// each mutable per-name family under `project/<default>/…`, rewrite the domain-routing
527/// index values to the `{project, site}` form, create the `default` project pointer, and
528/// build the `owner/*` reverse index. Online (the old-key delete is deferrable).
529struct 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        // Any layout-1 datum in a mutable or domain family means there is work to do.
544        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// ---- reusable migration steps ----------------------------------------------------
583
584/// One idempotent, resumable unit of a migration. `id` is stable (recorded in the
585/// resume cursor). `run` must be safe to re-run and must fail closed (verify before any
586/// destructive action). On `dry_run` it counts without writing.
587#[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
602/// Copy one mutable family to its `project/<project>/…` keys: copy → read-back verify,
603/// **without** deleting the source (that is [`DeleteOldFamily`], the finalize half).
604struct 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; // vanished between listing and read
629            };
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
643/// Delete the old keys of one mutable family whose `project/<project>/…` counterpart is
644/// present and byte-identical — the finalize half of copy-verify-delete. Idempotent (a
645/// missing old key is skipped); refuses to delete an old key whose new key is absent or
646/// mismatched, so a partially-copied family never loses data.
647struct 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
678/// Rewrite the values of one index family through `transform` (idempotent — an
679/// already-canonical value round-trips unchanged, so it is a no-op).
680struct 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
718/// Canonicalize a domain-index value to the `{project, site}` [`DomainOwner`] form (a
719/// bare legacy site string becomes `(default, <site>)`).
720fn domain_owner_canonical(value: &[u8]) -> Vec<u8> {
721    DomainOwner::from_bytes(value).to_bytes()
722}
723
724/// Create the `default` project pointer + content-addressed body if absent.
725struct 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        // Single-source the record shape with the boot-time ensure (fresh installs
747        // materialize `default` the same way a migration does).
748        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
760/// Build the `owner/<kind>/<name>` → project reverse index over the migrated
761/// site/function/compute records. Idempotent.
762struct 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        // Sites: the site-config pointer `project/default/site/<site>`.
777        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        // Functions: the meta key `project/default/functions/<name>` (no further `/`).
789        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        // Compute workloads: `project/default/compute/<name>`.
801        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    /// Seed a synthetic layout-1 store: sites, functions, compute, aliases, domain
826    /// index (raw bare-site values), and a content-addressed body that must NOT move.
827    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 &reg {
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        // The version advanced + is recorded in history.
896        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        // Content-addressed body + singleton untouched.
948        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        // Default project pointer + owner reverse index.
958        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        // The marker records the deferred cleanup.
1044        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        // Simulate a crash after only two forward steps: copy `current/` and `site/`
1058        // by hand and record their step ids in the resume cursor.
1059        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        // No double-nesting from re-processing an already-done family.
1097        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        // A store finalized by the original one-shot engine wrote `{layout:2, ...}`.
1107        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        // A `2-dual` legacy marker maps to a deferred cleanup.
1120        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    // ---- framework tests: prove the engine chains an arbitrary migration list -----
1129
1130    /// A synthetic **simple** migration (no cleanup) that writes one sentinel key,
1131    /// applicable only while the sentinel is absent.
1132    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        // Both migrations applied: v1 re-keyed the store, v2 wrote its sentinel.
1186        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        // History records them in order.
1199        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        // Idempotent re-run.
1203        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        // A store already at v1 (project layout) but not v2: only v2 runs.
1212        let kv = MemoryKv::new();
1213        // Write new-layout site data + a v1 marker directly (as if migrated by v1).
1214        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        // v1 was not re-run (no rekeyed families), v2 applied.
1231        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}