use std::borrow::Cow;
use std::time::Duration;
use anyhow::Result;
use tracing::{debug, trace, warn};
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider, TableProvider};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::schema::{DocKeyKey, DocKeyPrefix, DocLookupIdentityKey};
use crate::key::{KVKeyDecode, KVValue, Resumable};
use crate::kvs::ds::Datastore;
use crate::kvs::tasklease::LeaseHandler;
#[cfg(test)]
use crate::kvs::testing::{RetryableConflictSite, maybe_inject_retryable_conflict};
use crate::kvs::{Direction, TransactionType};
use crate::val::{RecordIdKey, RecordIdentity, TableName};
const TARGET: &str = "surrealdb::core::kvs::migration";
const BATCH: usize = 500;
pub(super) async fn rebuild_forward_doc_ids_under_an_injective_key(ds: &Datastore) -> Result<()> {
let txn = ds.transaction(TransactionType::Read).await?;
let tables = async {
let mut tables = Vec::new();
for ns in txn.all_ns(None).await?.iter() {
for db in txn.all_db(ns.namespace_id, None).await?.iter() {
for tb in txn.all_tb(ns.namespace_id, db.database_id, None).await?.iter() {
tables.push((ns.namespace_id, db.database_id, tb.name.clone()));
}
}
}
Ok::<_, anyhow::Error>(tables)
}
.await;
txn.cancel().await?;
let lease = super::lease(ds)?;
let mut copied = 0usize;
for (ns, db, tb) in tables? {
copied += migrate_table(ds, &lease, ns, db, &tb).await?;
}
if copied > 0 {
debug!(target: TARGET, count = copied, "Rebuilt forward doc-ID mappings under an injective key");
}
Ok(())
}
async fn migrate_table(
ds: &Datastore,
lease: &LeaseHandler,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
) -> Result<usize> {
let mut copied = 0usize;
let mut resume: Option<Vec<u8>> = None;
loop {
let Some((batch, last)) = read_batch(ds, ns, db, tb, resume.as_deref()).await? else {
return Ok(copied);
};
resume = Some(last);
match write_batch_with_retry(ds, ns, db, tb, &batch).await {
Ok(written) => copied += written,
Err(e) if is_write_conflict(&e) => {
warn!(
target: TARGET,
table = %tb,
mappings = batch.len(),
error = %e,
"Skipped a batch of doc-ID mappings after repeated write conflicts; \
they resolve through the verified fallback until written again"
);
}
Err(e) => return Err(e),
}
let _ = lease.try_maintain_lease().await;
}
}
fn is_write_conflict(e: &anyhow::Error) -> bool {
surrealdb_datastore::is_retryable_transaction_conflict(e)
|| crate::kvs::is_conditional_write_conflict(e)
}
const WRITE_ATTEMPTS: usize = 5;
const RETRY_BACKOFF: Duration = Duration::from_millis(10);
async fn write_batch_with_retry(
ds: &Datastore,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
batch: &[(Vec<u8>, u64, RecordIdKey)],
) -> Result<usize> {
for attempt in 1..=WRITE_ATTEMPTS {
match write_batch(ds, ns, db, tb, batch).await {
Err(e) if attempt < WRITE_ATTEMPTS && is_write_conflict(&e) => {
trace!(
target: "surrealdb::core::kvs::migration",
"Retrying a doc-ID mapping batch after a write conflict (attempt {attempt})"
);
let base = RETRY_BACKOFF * (1 << (attempt - 1));
let jitter = Duration::from_millis(rand::random::<u64>() % 20);
crate::common::time::sleep(base + jitter).await;
}
other => return other,
}
}
unreachable!("the final attempt returns rather than looping")
}
async fn read_batch(
ds: &Datastore,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
resume: Option<&[u8]>,
) -> Result<Option<(Vec<(Vec<u8>, u64, RecordIdKey)>, Vec<u8>)>> {
let txn = ds.transaction(TransactionType::Read).await?;
let collected = async {
let range = DocKeyPrefix::new(ns, db, Cow::Borrowed(tb)).range()?;
let range = match resume {
Some(resume) => range.resume_after(resume, Direction::Forward),
None => range,
};
let mut cursor = txn.open_vals_cursor(range, Direction::Forward, 0, None).await?;
let page = cursor.next_batch(BATCH as u32).await?;
let Some((last, _)) = page.iter().last() else {
return Ok(None);
};
let last = last.to_vec();
let mut out = Vec::new();
for (key, value) in page.iter() {
let decoded = DocKeyKey::decode_key(key)
.and_then(|k| RecordIdKey::kv_decode_value(value, ()).map(|id| (k.doc_id, id)));
match decoded {
Ok((doc_id, id)) => out.push((key.to_vec(), doc_id, id)),
Err(e) => warn!(
target: TARGET,
table = %tb,
error = %e,
"Skipped a doc-ID mapping that does not decode"
),
}
}
Ok::<_, anyhow::Error>(Some((out, last)))
}
.await;
txn.cancel().await?;
collected
}
async fn write_batch(
ds: &Datastore,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
batch: &[(Vec<u8>, u64, RecordIdKey)],
) -> Result<usize> {
let txn = ds.transaction(TransactionType::Write).await?;
let written = async {
#[cfg(test)]
maybe_inject_retryable_conflict(
RetryableConflictSite::DocLookupIdentityMigration,
ds.id(),
)?;
let forward = |id: &RecordIdKey| {
DocLookupIdentityKey::new(ns, db, Cow::Borrowed(tb), RecordIdentity(id.clone()))
};
let reverse = batch
.iter()
.map(|(_, doc_id, _)| DocKeyKey::new(ns, db, Cow::Borrowed(tb), *doc_id))
.collect();
let current = txn.get_many_key(reverse, None).await?;
let published =
txn.get_many_key(batch.iter().map(|(_, _, id)| forward(id)).collect(), None).await?;
let mut written = 0usize;
for (((_, doc_id, id), current), published) in batch.iter().zip(current).zip(published) {
if !current.is_some_and(|current| current.addresses_same_record(id)) {
continue;
}
if published.is_some() {
continue;
}
match txn.put_key(&forward(id), doc_id).await {
Ok(()) => written += 1,
Err(e) if super::already_exists(&e) => {}
Err(e) => return Err(e),
}
}
Ok::<_, anyhow::Error>(written)
}
.await;
match written {
Ok(written) => {
txn.commit().await?;
Ok(written)
}
Err(e) => {
txn.cancel().await?;
Err(e)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::idx::docids::TableDocIds;
use crate::val::{Number, Value};
fn arr(n: Number) -> RecordIdKey {
RecordIdKey::Array(vec![Value::Number(n)].into())
}
#[tokio::test]
async fn the_migration_rebuilds_the_injective_mapping_from_the_reverse_one() {
let ds = crate::kvs::Datastore::new("memory").await.unwrap();
ds.execute("DEFINE NAMESPACE n; USE NS n; DEFINE DATABASE d;", &Default::default(), None)
.await
.unwrap();
let (ns, db) = {
let txn = ds.transaction(TransactionType::Read).await.unwrap();
let ns = txn.all_ns(None).await.unwrap()[0].namespace_id;
let db = txn.all_db(ns, None).await.unwrap()[0].database_id;
txn.cancel().await.unwrap();
(ns, db)
};
let tb: TableName = "t".into();
let docids = TableDocIds::new(ns, db, tb.clone());
let int = arr(Number::Int(1));
ds.execute("USE NS n DB d; DEFINE TABLE t;", &Default::default(), None).await.unwrap();
{
let txn = ds.transaction(TransactionType::Write).await.unwrap();
txn.set_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&tb), 7), &int).await.unwrap();
txn.commit().await.unwrap();
}
{
let txn = ds.transaction(TransactionType::Read).await.unwrap();
assert_eq!(
txn.get_key(
&DocLookupIdentityKey::new(
ns,
db,
Cow::Borrowed(&tb),
RecordIdentity(int.clone())
),
None
)
.await
.unwrap(),
None,
"the injective mapping is what the migration is there to write"
);
txn.cancel().await.unwrap();
}
rebuild_forward_doc_ids_under_an_injective_key(&ds).await.unwrap();
let dj = DocLookupIdentityKey::new(ns, db, Cow::Borrowed(&tb), RecordIdentity(int.clone()));
let txn = ds.transaction(TransactionType::Read).await.unwrap();
assert_eq!(
txn.get_key(&dj, None).await.unwrap(),
Some(7),
"the migration must populate the injective key"
);
assert_eq!(docids.get_doc_id(&txn, &int).await.unwrap(), Some(7));
txn.cancel().await.unwrap();
rebuild_forward_doc_ids_under_an_injective_key(&ds).await.unwrap();
let txn = ds.transaction(TransactionType::Read).await.unwrap();
assert_eq!(txn.get_key(&dj, None).await.unwrap(), Some(7));
assert_eq!(docids.get_doc_id(&txn, &int).await.unwrap(), Some(7));
txn.cancel().await.unwrap();
}
#[tokio::test]
async fn a_stale_batch_does_not_overwrite_what_moved_on_under_it() {
let ds = crate::kvs::Datastore::new("memory").await.unwrap();
ds.execute(
"DEFINE NAMESPACE n; USE NS n; DEFINE DATABASE d; USE NS n DB d; DEFINE TABLE t;",
&Default::default(),
None,
)
.await
.unwrap();
let (ns, db) = {
let txn = ds.transaction(TransactionType::Read).await.unwrap();
let ns = txn.all_ns(None).await.unwrap()[0].namespace_id;
let db = txn.all_db(ns, None).await.unwrap()[0].database_id;
txn.cancel().await.unwrap();
(ns, db)
};
let tb: TableName = "t".into();
let deleted = arr(Number::Int(1));
let recreated = arr(Number::Int(2));
let key = |id: &RecordIdKey| {
DocLookupIdentityKey::new(ns, db, Cow::Borrowed(&tb), RecordIdentity(id.clone()))
};
{
let txn = ds.transaction(TransactionType::Write).await.unwrap();
txn.set_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&tb), 9), &recreated).await.unwrap();
txn.set_key(&key(&recreated), &9u64).await.unwrap();
txn.commit().await.unwrap();
}
let stale =
vec![(Vec::new(), 7u64, deleted.clone()), (Vec::new(), 8u64, recreated.clone())];
assert_eq!(write_batch(&ds, ns, db, &tb, &stale).await.unwrap(), 0);
let txn = ds.transaction(TransactionType::Read).await.unwrap();
assert_eq!(
txn.get_key(&key(&deleted), None).await.unwrap(),
None,
"a deleted record's mapping must not be resurrected"
);
assert_eq!(
txn.get_key(&key(&recreated), None).await.unwrap(),
Some(9),
"a re-created record keeps the doc-ID it was given, not the stale one"
);
txn.cancel().await.unwrap();
}
async fn only_ns_db(ds: &Datastore) -> (NamespaceId, DatabaseId) {
let txn = ds.transaction(TransactionType::Read).await.unwrap();
let ns = txn.all_ns(None).await.unwrap()[0].namespace_id;
let db = txn.all_db(ns, None).await.unwrap()[0].database_id;
txn.cancel().await.unwrap();
(ns, db)
}
async fn injective(
ds: &Datastore,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
id: &RecordIdKey,
) -> Option<u64> {
let txn = ds.transaction(TransactionType::Read).await.unwrap();
let dj = DocLookupIdentityKey::new(ns, db, Cow::Borrowed(tb), RecordIdentity(id.clone()));
let found = txn.get_key(&dj, None).await.unwrap();
txn.cancel().await.unwrap();
found
}
#[tokio::test]
async fn a_batch_that_keeps_losing_its_commit_is_skipped_not_fatal() {
let ds = crate::kvs::Datastore::new("memory").await.unwrap();
ds.execute(
"DEFINE NAMESPACE n; USE NS n; DEFINE DATABASE d; USE NS n DB d; DEFINE TABLE a; DEFINE TABLE b;",
&Default::default(),
None,
)
.await
.unwrap();
let (ns, db) = only_ns_db(&ds).await;
let (a, b): (TableName, TableName) = ("a".into(), "b".into());
let id = arr(Number::Int(1));
{
let txn = ds.transaction(TransactionType::Write).await.unwrap();
txn.set_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&a), 7), &id).await.unwrap();
txn.set_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&b), 8), &id).await.unwrap();
txn.commit().await.unwrap();
}
let _guard = crate::kvs::testing::inject_retryable_conflicts(
RetryableConflictSite::DocLookupIdentityMigration,
ds.id(),
WRITE_ATTEMPTS,
);
rebuild_forward_doc_ids_under_an_injective_key(&ds)
.await
.expect("a batch lost to conflicts must not fail the migration");
assert_eq!(injective(&ds, ns, db, &a, &id).await, None, "the lost batch is skipped");
assert_eq!(injective(&ds, ns, db, &b, &id).await, Some(8), "and the walk goes on past it");
rebuild_forward_doc_ids_under_an_injective_key(&ds).await.unwrap();
assert_eq!(injective(&ds, ns, db, &a, &id).await, Some(7), "a later run copies it");
}
#[tokio::test]
async fn a_reverse_mapping_that_does_not_decode_is_skipped_not_fatal() {
let ds = crate::kvs::Datastore::new("memory").await.unwrap();
ds.execute(
"DEFINE NAMESPACE n; USE NS n; DEFINE DATABASE d; USE NS n DB d; DEFINE TABLE t;",
&Default::default(),
None,
)
.await
.unwrap();
let (ns, db) = only_ns_db(&ds).await;
let tb: TableName = "t".into();
let (before, after) = (arr(Number::Int(1)), arr(Number::Int(2)));
{
let txn = ds.transaction(TransactionType::Write).await.unwrap();
txn.set_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&tb), 4), &before).await.unwrap();
let garbage =
crate::key::KVKey::encode_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&tb), 5))
.unwrap();
txn.set(garbage, vec![0xff; 3]).await.unwrap();
txn.set_key(&DocKeyKey::new(ns, db, Cow::Borrowed(&tb), 6), &after).await.unwrap();
txn.commit().await.unwrap();
}
rebuild_forward_doc_ids_under_an_injective_key(&ds)
.await
.expect("an undecodable entry must not fail the migration");
assert_eq!(injective(&ds, ns, db, &tb, &before).await, Some(4));
assert_eq!(injective(&ds, ns, db, &tb, &after).await, Some(6));
}
}