mod doc_lookup_identity;
mod module_keys;
mod sequence_defs;
use std::future::Future;
use std::pin::Pin;
use anyhow::{Result, bail};
use tracing::{info, warn};
use crate::key::schema::{MigrationKey, StorageVersionKey, VersionHistoryKey};
use crate::kvs::ds::Datastore;
use crate::kvs::tasklease::{LeaseHandler, TaskLeaseType};
use crate::kvs::version::{MigrationRecord, StorageVersion, VersionHistoryEntry};
use crate::kvs::{DatastoreError, TransactionType};
const TARGET: &str = "surrealdb::core::kvs::migration";
const LEASE_DURATION: std::time::Duration = std::time::Duration::from_secs(60);
const WAIT_FOR_HOLDER: std::time::Duration = std::time::Duration::from_secs(300);
const WAIT_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_secs(2);
type MigrationFuture<'a> = Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
pub(crate) struct Migration {
pub id: u32,
pub name: &'static str,
pub version: (u64, u64, u64),
pub run: for<'a> fn(&'a Datastore) -> MigrationFuture<'a>,
}
pub(crate) static MIGRATIONS: &[Migration] = &[
Migration {
id: 1,
name: "copy sequence definitions out of the table band",
version: (3, 3, 0),
run: |ds| Box::pin(sequence_defs::copy_definitions_out_of_the_table_band(ds)),
},
Migration {
id: 2,
name: "re-key module definitions to their derived name",
version: (3, 3, 0),
run: |ds| Box::pin(module_keys::rekey_definitions_to_their_derived_name(ds)),
},
Migration {
id: 3,
name: "rebuild forward doc-ID mappings under an injective key",
version: (3, 3, 0),
run: |ds| Box::pin(doc_lookup_identity::rebuild_forward_doc_ids_under_an_injective_key(ds)),
},
];
pub(crate) async fn run(ds: &Datastore, is_new: bool) -> Result<()> {
run_with(ds, MIGRATIONS, &StorageVersion::current(), is_new).await
}
async fn run_with(
ds: &Datastore,
registry: &[Migration],
current: &StorageVersion,
is_new: bool,
) -> Result<()> {
let stored = ds.storage_version().await?;
if is_new && stored.is_none() {
record_baseline(ds, registry).await?;
}
check_downgrade(ds, registry, stored.as_ref(), current).await?;
let pending = unapplied(ds, &candidate_ids(registry, current)).await?;
if !pending.is_empty() {
apply(ds, registry, &pending).await?;
}
advance_stamp(ds, stored, current).await
}
async fn record_baseline(ds: &Datastore, registry: &[Migration]) -> Result<()> {
for migration in registry {
record_applied(ds, migration).await?;
}
Ok(())
}
async fn check_downgrade(
ds: &Datastore,
registry: &[Migration],
stored: Option<&StorageVersion>,
current: &StorageVersion,
) -> Result<()> {
let unknown: Vec<String> = ds
.applied_migration_ids()
.await?
.into_iter()
.filter(|(id, _)| !registry.iter().any(|m| m.id == *id))
.map(|(id, described)| match described {
Some(record) => format!("{} ({}, id {id})", record.name, record.version),
None => format!("id {id}, unreadable record"),
})
.collect();
if !unknown.is_empty() {
bail!(DatastoreError::MigratedBeyondStorageVersion {
stored: stored.map(|s| s.to_string()).unwrap_or_else(|| "unknown".to_owned()),
running: current.to_string(),
migrations: unknown.join(", "),
});
}
let Some(from) = stored else {
return Ok(());
};
if from.triple() <= current.triple() {
return Ok(());
}
warn!(
target: TARGET,
stored = %from,
running = %current,
"This datastore was last written by a newer version of SurrealDB. \
No data migration separates the two, so startup continues, but \
running a mixed set of versions is not supported."
);
Ok(())
}
fn candidate_ids(registry: &[Migration], current: &StorageVersion) -> Vec<u32> {
registry.iter().filter(|m| m.version <= current.triple()).map(|m| m.id).collect()
}
async fn unapplied(ds: &Datastore, ids: &[u32]) -> Result<Vec<u32>> {
if ids.is_empty() {
return Ok(Vec::new());
}
let txn = ds.transaction(TransactionType::Read).await?;
let mut outstanding = Vec::new();
for id in ids {
let applied = catch!(txn, txn.exists_key(&MigrationKey::new(*id), None).await);
if !applied {
outstanding.push(*id);
}
}
txn.cancel().await?;
Ok(outstanding)
}
async fn apply(ds: &Datastore, registry: &[Migration], pending: &[u32]) -> Result<()> {
let lease = lease(ds)?;
let deadline = web_time::Instant::now() + WAIT_FOR_HOLDER;
loop {
let outstanding = unapplied(ds, pending).await?;
if outstanding.is_empty() {
return Ok(());
}
if lease.has_lease().await? {
for id in outstanding {
let migration = registry
.iter()
.find(|m| m.id == id)
.expect("pending ids are taken from the registry");
info!(
target: TARGET,
id = migration.id,
name = migration.name,
"Applying data migration"
);
(migration.run)(ds).await?;
record_applied(ds, migration).await?;
let _ = lease.try_maintain_lease().await;
}
return Ok(());
}
if web_time::Instant::now() >= deadline {
bail!(DatastoreError::MigrationTimedOut {
migrations: describe(registry, &outstanding),
});
}
info!(
target: TARGET,
count = outstanding.len(),
"Another node is applying data migrations; waiting for it to finish"
);
common::time::sleep(WAIT_POLL_INTERVAL).await;
}
}
pub(super) fn lease(ds: &Datastore) -> Result<LeaseHandler> {
LeaseHandler::new(
ds.sequences().clone(),
ds.id(),
ds.transaction_factory().clone(),
TaskLeaseType::DataMigration,
LEASE_DURATION,
)
}
async fn record_applied(ds: &Datastore, migration: &Migration) -> Result<()> {
let record = MigrationRecord {
name: migration.name.to_owned(),
version: StorageVersion::from_triple(migration.version),
node: ds.id(),
timestamp: ds.clock_now().value,
};
let txn = ds.transaction(TransactionType::Write).await?;
let key = MigrationKey::new(migration.id);
match run!(txn, txn.set_key(&key, &record).await) {
Err(e) if lost_the_race(&e) => Ok(()),
other => other,
}
}
async fn advance_stamp(
ds: &Datastore,
stored: Option<StorageVersion>,
current: &StorageVersion,
) -> Result<()> {
if stored.as_ref().is_some_and(|s| s.triple() >= current.triple()) {
return Ok(());
}
let timestamp = ds.clock_now().value;
let entry = VersionHistoryEntry {
from: stored.clone(),
to: current.clone(),
node: ds.id(),
timestamp,
};
let txn = ds.transaction(TransactionType::Write).await?;
let result = async {
txn.put_compare_key(&StorageVersionKey {}, current, stored.as_ref()).await?;
txn.set_key(&VersionHistoryKey::new(timestamp, ds.id()), &entry).await?;
txn.commit().await
}
.await;
match result {
Ok(()) => {
if let Some(from) = &entry.from {
info!(target: TARGET, %from, to = %current, "Datastore version advanced");
} else {
info!(target: TARGET, version = %current, "Datastore version recorded");
}
Ok(())
}
Err(e) if lost_the_race(&e) => {
let _ = txn.cancel().await;
Ok(())
}
Err(e) => {
let _ = txn.cancel().await;
Err(e)
}
}
}
fn describe(registry: &[Migration], ids: &[u32]) -> String {
ids.iter()
.map(|id| match registry.iter().find(|m| m.id == *id) {
Some(m) => format!("{} ({})", m.name, id),
None => id.to_string(),
})
.collect::<Vec<_>>()
.join(", ")
}
pub(super) fn already_exists(err: &anyhow::Error) -> bool {
matches!(
err.downcast_ref::<crate::kvs::Error>(),
Some(crate::kvs::Error::TransactionKeyAlreadyExists)
) || matches!(
err.downcast_ref::<crate::err::Error>(),
Some(crate::err::Error::Kvs(crate::kvs::Error::TransactionKeyAlreadyExists))
)
}
fn lost_the_race(err: &anyhow::Error) -> bool {
crate::kvs::is_conditional_write_conflict(err)
|| crate::kvs::is_retryable_transaction_conflict(err)
}
#[cfg(test)]
mod tests;