use std::collections::HashMap;
use std::collections::hash_map::Entry;
use std::ops::Range;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use anyhow::Result;
use rand::Rng;
use revision::revisioned;
use serde::{Deserialize, Serialize};
use tokio::sync::{Mutex, RwLock};
use tokio::time::sleep;
use uuid::Uuid;
use web_time::Instant;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider, TableProvider};
use crate::catalog::{DatabaseId, IndexId, NamespaceId, SequenceDefinition, TableId};
use crate::ctx::Context;
use crate::err::Error;
use crate::idx::IndexKeyBase;
use crate::idx::seqdocids::DocId;
use crate::key::database::sq::Sq;
use crate::key::database::th::TableIdGeneratorBatchKey;
use crate::key::database::ti::TableIdGeneratorStateKey;
use crate::key::namespace::dh::DatabaseIdGeneratorBatchKey;
use crate::key::namespace::di::DatabaseIdGeneratorStateKey;
use crate::key::root::nh::NamespaceIdGeneratorBatchKey;
use crate::key::root::ni::NamespaceIdGeneratorStateKey;
use crate::key::sequence::Prefix;
use crate::key::sequence::ba::Ba;
use crate::key::sequence::st::St;
use crate::key::table::ih::IndexIdGeneratorBatchKey;
use crate::key::table::is::IndexIdGeneratorStateKey;
use crate::kvs::ds::TransactionFactory;
use crate::kvs::tx::ProvisionalSequence;
use crate::kvs::{KVKey, LockType, Transaction, TransactionType, impl_kv_value_revisioned};
use crate::val::TableName;
type SequencesMap = Arc<RwLock<HashMap<Arc<SequenceDomain>, Arc<CachedSequence>>>>;
struct CachedSequence {
evicted: AtomicBool,
sequence: Mutex<Sequence>,
}
impl CachedSequence {
fn new(sequence: Sequence) -> Self {
Self {
evicted: AtomicBool::new(false),
sequence: Mutex::new(sequence),
}
}
fn evict(&self) {
self.evicted.store(true, Ordering::Release);
}
}
pub(crate) async fn next_unissued_value(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
sq: &str,
start: i64,
version: Option<u64>,
) -> Result<i64> {
let range = Prefix::new_st_range(ns, db, sq)?;
let mut next = start;
for (_, v) in tx.getr(range, version).await? {
next = next.max(SequenceState::decode(&v)?.next);
}
Ok(next)
}
#[derive(Clone, Copy)]
struct ClaimGuard<'a> {
definition: Option<&'a SequenceDefinition>,
cursor: Option<&'a SequenceState>,
caller: Option<&'a Transaction>,
}
struct BatchAllocation {
from: i64,
to: i64,
window: Vec<u8>,
}
enum AllocatorTx<'a> {
Caller(&'a Transaction),
Own(Box<Transaction>),
}
impl<'a> AllocatorTx<'a> {
async fn open(
caller: Option<&'a Transaction>,
tf: &TransactionFactory,
kind: TransactionType,
sqs: &Sequences,
) -> SequenceResult<Self> {
Ok(match caller {
Some(tx) => Self::Caller(tx),
None => {
Self::Own(Box::new(tf.transaction(kind, LockType::Optimistic, sqs.clone()).await?))
}
})
}
fn tx(&self) -> &Transaction {
match self {
Self::Caller(tx) => tx,
Self::Own(tx) => tx,
}
}
async fn commit(self) -> SequenceResult<()> {
if let Self::Own(tx) = self {
tx.commit().await?;
}
Ok(())
}
async fn cancel(self) -> SequenceResult<()> {
if let Self::Own(tx) = self {
tx.cancel().await?;
}
Ok(())
}
}
#[derive(Debug)]
enum SequenceError {
Reset,
Other(anyhow::Error),
}
type SequenceResult<T> = std::result::Result<T, SequenceError>;
impl std::fmt::Display for SequenceError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Reset => f.write_str(
"the sequence was redefined or removed while this allocator was serving it",
),
Self::Other(e) => write!(f, "{e}"),
}
}
}
impl From<anyhow::Error> for SequenceError {
fn from(err: anyhow::Error) -> Self {
Self::Other(err)
}
}
impl SequenceError {
fn into_public(self) -> anyhow::Error {
match self {
Self::Other(err) => err,
Self::Reset => anyhow::Error::new(Error::SequenceReset),
}
}
}
#[cfg(all(test, feature = "kv-rocksdb"))]
fn is_sequence_reset(err: &anyhow::Error) -> bool {
err.chain().any(|e| matches!(e.downcast_ref::<Error>(), Some(Error::SequenceReset)))
}
#[derive(Clone)]
pub struct Sequences {
tf: TransactionFactory,
nid: Uuid,
sequences: SequencesMap,
}
#[derive(Hash, PartialEq, Eq)]
enum SequenceDomain {
UserName(NamespaceId, DatabaseId, String),
FullTextDocIds(IndexKeyBase),
NameSpacesIds,
DatabasesIds(NamespaceId),
TablesIds(NamespaceId, DatabaseId),
IndexIds(NamespaceId, DatabaseId, TableName),
}
impl SequenceDomain {
fn new_user(ns: NamespaceId, db: DatabaseId, sq: &str) -> Self {
Self::UserName(ns, db, sq.to_string())
}
pub(crate) fn new_ft_doc_ids(ikb: IndexKeyBase) -> Self {
Self::FullTextDocIds(ikb)
}
pub(crate) fn new_namespace_ids() -> Self {
Self::NameSpacesIds
}
pub(crate) fn new_database_ids(ns: NamespaceId) -> Self {
Self::DatabasesIds(ns)
}
pub(crate) fn new_table_ids(ns: NamespaceId, db: DatabaseId) -> Self {
Self::TablesIds(ns, db)
}
pub(crate) fn new_index_ids(ns: NamespaceId, db: DatabaseId, tb: TableName) -> Self {
Self::IndexIds(ns, db, tb)
}
fn is_user(&self) -> bool {
matches!(self, Self::UserName(..))
}
fn new_batch_range_keys(&self) -> Result<Range<Vec<u8>>> {
match self {
Self::UserName(ns, db, sq) => Prefix::new_ba_range(*ns, *db, sq),
Self::FullTextDocIds(ibk) => ibk.new_ib_range(),
Self::NameSpacesIds => NamespaceIdGeneratorBatchKey::range(),
Self::DatabasesIds(ns) => DatabaseIdGeneratorBatchKey::range(*ns),
Self::TablesIds(ns, db) => TableIdGeneratorBatchKey::range(*ns, *db),
Self::IndexIds(ns, db, tb) => IndexIdGeneratorBatchKey::range(*ns, *db, tb),
}
}
fn new_batch_key(&self, start: i64) -> Result<Vec<u8>> {
match &self {
Self::UserName(ns, db, sq) => Ba::new(*ns, *db, sq, start).encode_key(),
Self::FullTextDocIds(ikb) => ikb.new_ib_key(start).encode_key(),
Self::NameSpacesIds => NamespaceIdGeneratorBatchKey::new(start).encode_key(),
Self::DatabasesIds(ns) => DatabaseIdGeneratorBatchKey::new(*ns, start).encode_key(),
Self::TablesIds(ns, db) => TableIdGeneratorBatchKey::new(*ns, *db, start).encode_key(),
Self::IndexIds(ns, db, tb) => {
IndexIdGeneratorBatchKey::new(*ns, *db, tb, start).encode_key()
}
}
}
fn new_state_key(&self, nid: Uuid) -> Result<Vec<u8>> {
match &self {
Self::UserName(ns, db, sq) => St::new(*ns, *db, sq, nid).encode_key(),
Self::FullTextDocIds(ikb) => ikb.new_is_key(nid).encode_key(),
Self::NameSpacesIds => NamespaceIdGeneratorStateKey::new(nid).encode_key(),
Self::DatabasesIds(ns) => DatabaseIdGeneratorStateKey::new(*ns, nid).encode_key(),
Self::TablesIds(ns, db) => TableIdGeneratorStateKey::new(*ns, *db, nid).encode_key(),
Self::IndexIds(ns, db, tb) => {
IndexIdGeneratorStateKey::new(*ns, *db, tb, nid).encode_key()
}
}
}
}
#[revisioned(revision = 1)]
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Serialize, Deserialize, Hash)]
pub(crate) struct BatchValue {
to: i64,
owner: Uuid,
}
impl_kv_value_revisioned!(BatchValue);
impl BatchValue {
#[cfg(test)]
pub(crate) fn new(to: i64, owner: Uuid) -> Self {
Self {
to,
owner,
}
}
}
#[revisioned(revision = 1)]
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Serialize, Deserialize, Hash)]
pub(crate) struct SequenceState {
next: i64,
}
impl_kv_value_revisioned!(SequenceState);
impl SequenceState {
#[cfg(test)]
pub(crate) fn new(next: i64) -> Self {
Self {
next,
}
}
fn decode(v: &[u8]) -> Result<Self> {
Ok(revision::from_slice(v)?)
}
fn encode(&self) -> Result<Vec<u8>> {
Ok(revision::to_vec(self)?)
}
}
impl Sequences {
pub(super) fn new(tf: TransactionFactory, nid: Uuid) -> Self {
Self {
tf,
sequences: Arc::new(Default::default()),
nid,
}
}
pub(crate) async fn namespace_removed(&self, tx: &Transaction, ns: NamespaceId) -> Result<()> {
for db in tx.all_db(ns, None).await?.iter() {
self.database_removed(tx, ns, db.database_id).await?;
}
Ok(())
}
pub(crate) async fn database_removed(
&self,
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
) -> Result<()> {
for sqs in tx.all_db_sequences(ns, db, None).await?.iter() {
self.sequence_removed(ns, db, &sqs.name).await;
}
Ok(())
}
pub(crate) async fn sequence_removed(&self, ns: NamespaceId, db: DatabaseId, sq: &str) {
let key = SequenceDomain::new_user(ns, db, sq);
let mut cached = self.sequences.write().await;
if let Some(s) = cached.remove(&key) {
s.evict();
}
}
async fn next_val(
&self,
ctx: Option<&Context>,
seq: Arc<SequenceDomain>,
start: i64,
batch: u32,
timeout: Option<Duration>,
) -> Result<i64> {
let sequence = self.sequences.read().await.get(&seq).cloned();
if let Some(s) = sequence {
let res =
s.sequence.lock().await.next(self, ctx, &seq, batch, None, Some(&s.evicted)).await;
match res {
Err(SequenceError::Reset) => {
self.evict_this_allocator(&seq, &s).await;
}
res => return res.map_err(SequenceError::into_public),
}
return self
.build_and_next(ctx, &seq, start, batch, timeout)
.await
.map_err(SequenceError::into_public);
}
match self.build_and_next(ctx, &seq, start, batch, timeout).await {
Err(SequenceError::Reset) => self
.build_and_next(ctx, &seq, start, batch, timeout)
.await
.map_err(SequenceError::into_public),
res => res.map_err(SequenceError::into_public),
}
}
async fn build_and_next(
&self,
ctx: Option<&Context>,
seq: &Arc<SequenceDomain>,
start: i64,
batch: u32,
timeout: Option<Duration>,
) -> SequenceResult<i64> {
let s = match self.sequences.write().await.entry(Arc::clone(seq)) {
Entry::Occupied(e) => Arc::clone(e.get()),
Entry::Vacant(e) => {
let s = Arc::new(CachedSequence::new(
Sequence::load(ctx, self, seq, start, batch, timeout).await?,
));
Arc::clone(e.insert(s))
}
};
let res = s.sequence.lock().await.next(self, ctx, seq, batch, None, Some(&s.evicted)).await;
if matches!(res, Err(SequenceError::Reset)) {
self.evict_this_allocator(seq, &s).await;
}
res
}
async fn evict_this_allocator(&self, seq: &Arc<SequenceDomain>, s: &Arc<CachedSequence>) {
let mut cached = self.sequences.write().await;
s.evict();
if let Entry::Occupied(e) = cached.entry(Arc::clone(seq))
&& Arc::ptr_eq(e.get(), s)
{
e.remove();
}
}
pub(crate) async fn next_namespace_id(&self, ctx: Option<&Context>) -> Result<NamespaceId> {
let domain = Arc::new(SequenceDomain::new_namespace_ids());
let id = self.next_val(ctx, domain, 0, 100, None).await?;
Ok(NamespaceId(id as u32))
}
pub(crate) async fn next_database_id(
&self,
ctx: Option<&Context>,
ns: NamespaceId,
) -> Result<DatabaseId> {
let domain = Arc::new(SequenceDomain::new_database_ids(ns));
let id = self.next_val(ctx, domain, 0, 100, None).await?;
Ok(DatabaseId(id as u32))
}
pub(crate) async fn next_table_id(
&self,
ctx: Option<&Context>,
ns: NamespaceId,
db: DatabaseId,
) -> Result<TableId> {
let domain = Arc::new(SequenceDomain::new_table_ids(ns, db));
let id = self.next_val(ctx, domain, 0, 100, None).await?;
Ok(TableId(id as u32))
}
pub(crate) async fn next_index_id(
&self,
ctx: Option<&Context>,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
) -> Result<IndexId> {
let domain = Arc::new(SequenceDomain::new_index_ids(ns, db, tb));
let id = self.next_val(ctx, domain, 0, 100, None).await?;
Ok(IndexId(id as u32))
}
pub(crate) async fn next_user_sequence_id(
&self,
ctx: Option<&Context>,
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
sq: &str,
) -> Result<i64> {
let seq_def = tx.get_db_sequence(ns, db, sq, None).await?;
if let Some(slot) = tx.provisional_sequence(ns, db, sq).await {
let domain = SequenceDomain::new_user(ns, db, sq);
return self.next_provisional_val(ctx, tx, &slot, &domain, &seq_def).await;
}
let domain = Arc::new(SequenceDomain::new_user(ns, db, sq));
self.next_val(ctx, domain, seq_def.start, seq_def.batch, seq_def.timeout).await
}
async fn next_provisional_val(
&self,
ctx: Option<&Context>,
tx: &Transaction,
slot: &ProvisionalSequence,
domain: &SequenceDomain,
seq_def: &SequenceDefinition,
) -> Result<i64> {
let mut slot = slot.lock().await;
let allocator = match &mut *slot {
Some(allocator) => allocator,
none => none.insert(
Sequence::load_provisional(ctx, self, domain, seq_def, tx)
.await
.map_err(SequenceError::into_public)?,
),
};
allocator
.next(self, ctx, domain, seq_def.batch, Some(tx), None)
.await
.map_err(SequenceError::into_public)
}
pub(crate) async fn next_fts_doc_id(
&self,
ctx: Option<&Context>,
ikb: IndexKeyBase,
batch: u32,
) -> Result<DocId> {
let domain = Arc::new(SequenceDomain::new_ft_doc_ids(ikb));
let id = self.next_val(ctx, domain, 0, batch, None).await?;
Ok(id as DocId)
}
}
#[derive(Clone)]
pub(crate) struct Sequence {
tf: TransactionFactory,
st: SequenceState,
timeout: Option<Duration>,
to: i64,
state_key: Vec<u8>,
durable: Option<SequenceState>,
definition: Option<SequenceDefinition>,
window: Vec<u8>,
}
impl Sequence {
async fn load(
ctx: Option<&Context>,
sqs: &Sequences,
seq: &SequenceDomain,
start: i64,
batch: u32,
timeout: Option<Duration>,
) -> SequenceResult<Self> {
Self::build(ctx, sqs, seq, start, batch, timeout, None).await
}
async fn load_provisional(
ctx: Option<&Context>,
sqs: &Sequences,
seq: &SequenceDomain,
def: &SequenceDefinition,
caller: &Transaction,
) -> SequenceResult<Self> {
Self::build(ctx, sqs, seq, def.start, def.batch, def.timeout, Some(caller)).await
}
async fn build(
ctx: Option<&Context>,
sqs: &Sequences,
seq: &SequenceDomain,
start: i64,
batch: u32,
timeout: Option<Duration>,
caller: Option<&Transaction>,
) -> SequenceResult<Self> {
let state_key = seq.new_state_key(sqs.nid)?;
let own = AllocatorTx::open(caller, &sqs.tf, TransactionType::Read, sqs).await?;
let tx = own.tx();
let durable = match tx.get(&state_key, None).await? {
Some(v) => Some(SequenceState::decode(&v)?),
None => None,
};
let (start, batch, timeout, definition) = match seq {
SequenceDomain::UserName(..) if caller.is_some() => (start, batch, timeout, None),
SequenceDomain::UserName(ns, db, sq) => {
match tx.get(&Sq::new(*ns, *db, sq), None).await? {
Some(def) => (def.start, def.batch, def.timeout, Some(def)),
None => {
return Err(anyhow::Error::new(Error::SeqNotFound {
name: sq.clone(),
})
.into());
}
}
}
_ => (start, batch, timeout, None),
};
let mut st = if let Some(st) = durable.clone() {
st
} else {
let start = Self::seed_start_from_catalog(tx, seq, start).await?;
SequenceState {
next: start,
}
};
own.cancel().await?;
let BatchAllocation {
from,
to,
window,
} = Self::find_batch_allocation(
sqs,
ctx,
seq,
st.next,
batch,
timeout,
ClaimGuard {
definition: definition.as_ref(),
cursor: durable.as_ref(),
caller,
},
)
.await?;
st.next = from;
Ok(Self {
tf: sqs.tf.clone(),
state_key,
to,
st,
timeout,
durable,
definition,
window,
})
}
async fn seed_start_from_catalog(
tx: &Transaction,
seq: &SequenceDomain,
start: i64,
) -> Result<i64> {
let mut seeded = start;
match seq {
SequenceDomain::NameSpacesIds => {
for ns in tx.all_ns(None).await?.iter() {
seeded = seeded.max(ns.namespace_id.0 as i64 + 1);
}
}
SequenceDomain::DatabasesIds(ns) => {
for db in tx.all_db(*ns, None).await?.iter() {
seeded = seeded.max(db.database_id.0 as i64 + 1);
}
}
SequenceDomain::TablesIds(ns, db) => {
for tb in tx.all_tb(*ns, *db, None).await?.iter() {
seeded = seeded.max(tb.table_id.0 as i64 + 1);
}
}
SequenceDomain::IndexIds(ns, db, tb) => {
for ix in tx.all_tb_indexes(*ns, *db, tb, None).await?.iter() {
seeded = seeded.max(ix.index_id.0 as i64 + 1);
}
}
SequenceDomain::FullTextDocIds(_) | SequenceDomain::UserName(..) => {}
}
Ok(seeded)
}
async fn next(
&mut self,
sqs: &Sequences,
ctx: Option<&Context>,
seq: &SequenceDomain,
batch: u32,
caller: Option<&Transaction>,
evicted: Option<&AtomicBool>,
) -> SequenceResult<i64> {
let batch = self.definition.as_ref().map(|d| d.batch).unwrap_or(batch);
if self.st.next >= self.to {
BatchAllocation {
from: self.st.next,
to: self.to,
window: self.window,
} = Self::find_batch_allocation(
sqs,
ctx,
seq,
self.st.next,
batch,
self.timeout,
ClaimGuard {
definition: self.definition.as_ref(),
cursor: self.durable.as_ref(),
caller,
},
)
.await?;
}
let v = self.st.next;
self.st.next += 1;
let written = self.st.encode()?;
let expected = self.durable.as_ref().map(SequenceState::encode).transpose()?;
let own = AllocatorTx::open(caller, &self.tf, TransactionType::Write, sqs).await?;
let tx = own.tx();
if caller.is_none() && self.durable.is_none() && seq.is_user() {
let claim = match tx.get(&self.window, None).await {
Ok(Some(claim)) => claim,
Ok(None) => {
own.cancel().await?;
return Err(SequenceError::Reset);
}
Err(e) => {
own.cancel().await?;
return Err(e.into());
}
};
match revision::from_slice::<BatchValue>(&claim) {
Ok(ba) if ba.owner == sqs.nid && ba.to == self.to => {}
Ok(_) => {
own.cancel().await?;
return Err(SequenceError::Reset);
}
Err(e) => {
own.cancel().await?;
return Err(anyhow::Error::from(e).into());
}
}
if let Err(e) = tx.set(&self.window, &claim).await {
own.cancel().await?;
return Err(e.into());
}
}
let res = if caller.is_some() || !seq.is_user() {
tx.set(&self.state_key, &written).await
} else {
tx.putc(&self.state_key, &written, expected.as_ref()).await
};
match res {
Ok(_) if evicted.is_some_and(|e| e.load(Ordering::Acquire)) => {
own.cancel().await?;
Err(SequenceError::Reset)
}
Ok(_) => match own.commit().await {
Ok(()) => {
self.durable = Some(self.st.clone());
Ok(v)
}
Err(SequenceError::Other(e))
if seq.is_user() && crate::kvs::is_retryable_transaction_conflict(&e) =>
{
Err(SequenceError::Reset)
}
Err(e) => Err(e),
},
Err(e) => {
own.cancel().await?;
if crate::kvs::is_condition_not_met(&e) {
return Err(SequenceError::Reset);
}
Err(e.into())
}
}
}
async fn find_batch_allocation(
sqs: &Sequences,
ctx: Option<&Context>,
seq: &SequenceDomain,
next: i64,
batch: u32,
to: Option<Duration>,
guard: ClaimGuard<'_>,
) -> SequenceResult<BatchAllocation> {
let mut tempo = 4;
const MAX_BACKOFF: u64 = 32_768;
let start = if to.is_some() {
Some(Instant::now())
} else {
None
};
loop {
if let Some(ctx) = ctx {
ctx.expect_not_timedout().await?;
} else {
yield_now!();
}
if let (Some(ref start), Some(ref to)) = (start, to) {
if start.elapsed().ge(to) {
let timeout = (*to).into();
return Err(anyhow::Error::new(Error::QueryTimedout(timeout)).into());
}
}
match Self::check_batch_allocation(sqs, seq, next, batch, guard).await {
Ok(r) => return Ok(r),
Err(SequenceError::Reset) => return Err(SequenceError::Reset),
Err(e) if guard.caller.is_some() => return Err(e),
Err(_) => {}
}
let sleep_ms = rand::rng().random_range(1..=tempo);
sleep(Duration::from_millis(sleep_ms)).await;
if tempo < MAX_BACKOFF {
tempo *= 2;
}
}
}
async fn check_batch_allocation(
sqs: &Sequences,
seq: &SequenceDomain,
next: i64,
batch: u32,
guard: ClaimGuard<'_>,
) -> SequenceResult<BatchAllocation> {
let own = AllocatorTx::open(guard.caller, &sqs.tf, TransactionType::Write, sqs).await?;
let tx = own.tx();
let result = async {
if let (Some(expected), SequenceDomain::UserName(ns, db, sq)) = (guard.definition, seq)
{
let key = Sq::new(*ns, *db, sq);
let current = tx.get(&key, None).await?;
if current.as_ref() != Some(expected) {
return Err(SequenceError::Reset);
}
tx.set(&key, expected).await?;
}
if seq.is_user()
&& guard.caller.is_none()
&& let Some(cursor) = guard.cursor
{
let current = match tx.get(&seq.new_state_key(sqs.nid)?, None).await? {
Some(v) => Some(SequenceState::decode(&v)?),
None => None,
};
if current.as_ref() != Some(cursor) {
return Err(SequenceError::Reset);
}
}
let batch_range = seq.new_batch_range_keys()?;
let val = tx.getr(batch_range, None).await?;
let mut next_start = next;
for (key, val) in val.iter() {
let ba: BatchValue = revision::from_slice(val).map_err(anyhow::Error::from)?;
next_start = next_start.max(ba.to);
if ba.owner == sqs.nid {
if next < ba.to {
return Ok(BatchAllocation {
from: next,
to: ba.to,
window: key.clone(),
});
}
tx.del(key).await?;
}
}
let next_to = next_start + batch as i64;
let bv = revision::to_vec(&BatchValue {
to: next_to,
owner: sqs.nid,
})
.map_err(anyhow::Error::from)?;
let batch_key = seq.new_batch_key(next_start)?;
tx.set(&batch_key, &bv).await?;
Ok::<BatchAllocation, SequenceError>(BatchAllocation {
from: next_start,
to: next_to,
window: batch_key,
})
}
.await;
match result {
Ok(res) => {
own.commit().await?;
Ok(res)
}
Err(e) => {
own.cancel().await?;
Err(e)
}
}
}
}
#[cfg(test)]
mod tests {
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider, TableProvider};
use crate::catalog::{
DatabaseDefinition, DatabaseId, Index, IndexDefinition, IndexId, NamespaceDefinition,
NamespaceId, TableDefinition, TableId,
};
use crate::kvs::sequences::{Sequence, SequenceDomain};
use crate::kvs::{Datastore, LockType, TransactionType};
use crate::val::TableName;
#[tokio::test]
async fn seed_start_from_catalog_uses_max_existing_id() {
let ds = Datastore::new("memory").await.unwrap();
let ns_id = NamespaceId(7);
let db_id = DatabaseId(11);
let tb_name: TableName = "tb".into();
let tx = ds.transaction(TransactionType::Write, LockType::Optimistic).await.unwrap();
tx.put_ns(NamespaceDefinition {
namespace_id: ns_id,
name: "ns".into(),
comment: None,
})
.await
.unwrap();
tx.put_db(
"ns",
DatabaseDefinition {
namespace_id: ns_id,
database_id: db_id,
name: "db".into(),
comment: None,
changefeed: None,
strict: false,
},
)
.await
.unwrap();
tx.put_tb("ns", "db", &TableDefinition::new(ns_id, db_id, TableId(13), tb_name.clone()))
.await
.unwrap();
tx.put_tb_index(
ns_id,
db_id,
&tb_name,
&IndexDefinition {
index_id: IndexId(17),
name: "ix".into(),
table_name: tb_name.clone(),
cols: vec![],
index: Index::Idx,
comment: None,
prepare_remove: false,
},
)
.await
.unwrap();
tx.commit().await.unwrap();
let tx = ds.transaction(TransactionType::Read, LockType::Optimistic).await.unwrap();
assert_eq!(
Sequence::seed_start_from_catalog(&tx, &SequenceDomain::NameSpacesIds, 0)
.await
.unwrap(),
8
);
assert_eq!(
Sequence::seed_start_from_catalog(&tx, &SequenceDomain::DatabasesIds(ns_id), 0)
.await
.unwrap(),
12
);
assert_eq!(
Sequence::seed_start_from_catalog(&tx, &SequenceDomain::TablesIds(ns_id, db_id), 0)
.await
.unwrap(),
14
);
assert_eq!(
Sequence::seed_start_from_catalog(
&tx,
&SequenceDomain::IndexIds(ns_id, db_id, tb_name.clone()),
0,
)
.await
.unwrap(),
18
);
assert_eq!(
Sequence::seed_start_from_catalog(&tx, &SequenceDomain::NameSpacesIds, 100)
.await
.unwrap(),
100
);
assert_eq!(
Sequence::seed_start_from_catalog(
&tx,
&SequenceDomain::DatabasesIds(NamespaceId(999)),
3,
)
.await
.unwrap(),
3
);
tx.cancel().await.unwrap();
}
#[tokio::test]
async fn seed_start_from_catalog_returns_start_on_empty_store() {
let ds = Datastore::new("memory").await.unwrap();
let tx = ds.transaction(TransactionType::Read, LockType::Optimistic).await.unwrap();
assert_eq!(
Sequence::seed_start_from_catalog(&tx, &SequenceDomain::NameSpacesIds, 0)
.await
.unwrap(),
0
);
assert_eq!(
Sequence::seed_start_from_catalog(
&tx,
&SequenceDomain::DatabasesIds(NamespaceId(0)),
0,
)
.await
.unwrap(),
0
);
tx.cancel().await.unwrap();
}
#[cfg(any(feature = "kv-mem", feature = "kv-rocksdb"))]
mod support {
use anyhow::Result;
use uuid::Uuid;
use crate::catalog::providers::DatabaseProvider;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::database::sq::Sq;
use crate::key::sequence::Prefix;
use crate::kvs::ds::TransactionFactory;
use crate::kvs::sequences::Sequences;
use crate::kvs::{LockType, TransactionType};
pub(super) fn tf_sequences(tf: &TransactionFactory) -> Sequences {
Sequences::new(tf.clone(), Uuid::from_u128(99))
}
pub(super) async fn define_sequence_in(
tx: &crate::kvs::Transaction,
ns: NamespaceId,
db: DatabaseId,
name: &str,
start: i64,
batch: u32,
) -> Result<()> {
tx.set(
&Sq::new(ns, db, name),
&crate::catalog::SequenceDefinition {
name: name.into(),
batch,
start,
timeout: None,
},
)
.await?;
tx.delr(Prefix::new_ba_range(ns, db, name)?).await?;
tx.delr(Prefix::new_st_range(ns, db, name)?).await?;
Ok(())
}
pub(super) async fn define_sequence(
tf: &TransactionFactory,
sqs: &Sequences,
ns: NamespaceId,
db: DatabaseId,
name: &str,
start: i64,
batch: u32,
) -> Result<()> {
let tx =
tf.transaction(TransactionType::Write, LockType::Optimistic, sqs.clone()).await?;
match define_sequence_in(&tx, ns, db, name, start, batch).await {
Ok(()) => tx.commit().await,
Err(e) => {
tx.cancel().await?;
Err(e)
}
}
}
pub(super) async fn reposition_to_resume(
tf: &TransactionFactory,
sqs: &Sequences,
ns: NamespaceId,
db: DatabaseId,
name: &str,
batch: u32,
) -> Result<()> {
let tx =
tf.transaction(TransactionType::Write, LockType::Optimistic, sqs.clone()).await?;
let res = async {
let def = tx.get_db_sequence(ns, db, name, None).await?;
let start =
crate::kvs::sequences::next_unissued_value(&tx, ns, db, name, def.start, None)
.await?;
define_sequence_in(&tx, ns, db, name, start, batch).await
}
.await;
match res {
Ok(()) => tx.commit().await,
Err(e) => {
tx.cancel().await?;
Err(e)
}
}
}
}
#[cfg(feature = "kv-mem")]
mod over_memory {
use std::sync::Arc;
use tokio::sync::Notify;
use uuid::Uuid;
use super::support::{
define_sequence, define_sequence_in, reposition_to_resume, tf_sequences,
};
use crate::catalog::providers::DatabaseProvider;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::database::sq::Sq;
use crate::key::sequence::Prefix;
use crate::kvs::ds::{DatastoreFlavor, TransactionFactory};
use crate::kvs::sequences::{
Sequence, SequenceDomain, SequenceError, SequenceResult, Sequences,
};
use crate::kvs::{LockType, TransactionType};
async fn factory() -> (TransactionFactory, Sequences) {
let flavor = crate::kvs::mem::Datastore::new(crate::kvs::mem::MemoryConfig::default())
.await
.map(DatastoreFlavor::Mem)
.unwrap();
let tf = TransactionFactory::new(
Arc::new(Notify::new()),
Box::new(flavor),
Arc::new(Default::default()),
);
let sequences = Sequences::new(tf.clone(), Uuid::new_v4());
(tf, sequences)
}
#[tokio::test]
async fn a_reset_between_the_claim_and_the_first_write_is_caught() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let seqs = Sequences::new(tf.clone(), Uuid::from_u128(21));
define_sequence(&tf, &seqs, ns, db, name, 1, 1000).await.unwrap();
let domain = SequenceDomain::new_user(ns, db, name);
let mut seq = Sequence::load(None, &seqs, &domain, 1, 1000, None).await.unwrap();
define_sequence(&tf, &seqs, ns, db, name, 500, 1000).await.unwrap();
let err = seq
.next(&seqs, None, &domain, 1000, None, None)
.await
.expect_err("the first write served a window the reposition had replaced");
assert!(matches!(err, SequenceError::Reset), "expected a reset, got: {err}");
}
#[tokio::test]
async fn an_identical_overwrite_still_resets_an_allocator() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let seqs = Sequences::new(tf.clone(), Uuid::from_u128(31));
define_sequence(&tf, &seqs, ns, db, name, 1, 1000).await.unwrap();
let domain = SequenceDomain::new_user(ns, db, name);
let mut seq = Sequence::load(None, &seqs, &domain, 1, 1000, None).await.unwrap();
define_sequence(&tf, &seqs, ns, db, name, 1, 1000).await.unwrap();
let err = seq
.next(&seqs, None, &domain, 1000, None, None)
.await
.expect_err("an allocator survived a reset that rewrote the same definition");
assert!(matches!(err, SequenceError::Reset), "expected a reset, got: {err}");
}
#[tokio::test]
async fn an_allocator_is_not_built_over_a_removed_sequence() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let seqs = Sequences::new(tf.clone(), Uuid::from_u128(32));
define_sequence(&tf, &seqs, ns, db, name, 1, 1000).await.unwrap();
let tx = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
tx.del(&Sq::new(ns, db, name)).await.unwrap();
tx.commit().await.unwrap();
let domain = SequenceDomain::new_user(ns, db, name);
let err = match Sequence::load(None, &seqs, &domain, 1, 1000, None).await {
Ok(_) => panic!("an allocator was built over a removed sequence"),
Err(e) => e,
};
assert!(
err.to_string().contains("does not exist"),
"expected the sequence to be reported missing, got: {err}"
);
}
#[tokio::test]
async fn a_peer_with_an_exhausted_window_picks_up_an_identical_overwrite() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let peer = Sequences::new(tf.clone(), Uuid::from_u128(41));
define_sequence(&tf, &peer, ns, db, name, 1, 1).await.unwrap();
let serve = async |tf: TransactionFactory, peer: Sequences| -> i64 {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, peer.clone())
.await
.unwrap();
let v = peer.next_user_sequence_id(None, &tx, ns, db, name).await;
tx.cancel().await.unwrap();
v.unwrap()
};
assert_eq!(serve(tf.clone(), peer.clone()).await, 1);
define_sequence(&tf, &peer, ns, db, name, 1, 1).await.unwrap();
assert_eq!(
serve(tf.clone(), peer.clone()).await,
1,
"the peer carried on from its old window instead of resetting"
);
}
mod hooked {
use std::fmt;
use std::future::Future;
use std::ops::Range;
use std::pin::Pin;
use std::sync::{Arc, Mutex as StdMutex};
use tokio::sync::Notify;
use crate::kvs::api::{BoxFut, KeysResult, ScanResult, Transactable};
use crate::kvs::ds::{
DatastoreFlavor, Metrics, TransactionBuilder, TransactionFactory,
};
use crate::kvs::err::Result as KvsResult;
use crate::kvs::{Key, TransactionBuilderRequirements, Val};
pub(super) type Hook =
Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()> + Send>> + Send>;
pub(super) type HookSlot = Arc<StdMutex<Option<Hook>>>;
struct HookedBuilder {
inner: Box<dyn TransactionBuilder>,
hook: HookSlot,
}
struct HookedTx {
inner: Box<dyn Transactable>,
hook: HookSlot,
}
impl fmt::Display for HookedBuilder {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "hooked {}", self.inner)
}
}
impl TransactionBuilderRequirements for HookedBuilder {}
impl TransactionBuilder for HookedBuilder {
fn new_transaction(
&self,
write: bool,
lock: bool,
) -> BoxFut<'_, anyhow::Result<(Box<dyn Transactable>, bool)>> {
Box::pin(async move {
let (inner, local) = self.inner.new_transaction(write, lock).await?;
let tx: Box<dyn Transactable> = Box::new(HookedTx {
inner,
hook: Arc::clone(&self.hook),
});
Ok((tx, local))
})
}
fn shutdown(&self) -> BoxFut<'_, anyhow::Result<()>> {
self.inner.shutdown()
}
fn register_metrics(&self) -> Option<Metrics> {
self.inner.register_metrics()
}
fn collect_u64_metric(&self, metric: &str) -> Option<u64> {
self.inner.collect_u64_metric(metric)
}
}
impl Transactable for HookedTx {
fn kind(&self) -> &'static str {
self.inner.kind()
}
fn closed(&self) -> bool {
self.inner.closed()
}
fn writeable(&self) -> bool {
self.inner.writeable()
}
fn cancel(&self) -> BoxFut<'_, KvsResult<()>> {
self.inner.cancel()
}
fn commit(&self) -> BoxFut<'_, KvsResult<()>> {
Box::pin(async move {
let hook = self.hook.lock().unwrap().take();
if let Some(hook) = hook {
hook().await;
}
self.inner.commit().await
})
}
fn exists(&self, key: Key, version: Option<u64>) -> BoxFut<'_, KvsResult<bool>> {
self.inner.exists(key, version)
}
fn get(
&self,
key: Key,
version: Option<u64>,
) -> BoxFut<'_, KvsResult<Option<Val>>> {
self.inner.get(key, version)
}
fn set(&self, key: Key, val: Val) -> BoxFut<'_, KvsResult<()>> {
self.inner.set(key, val)
}
fn put(&self, key: Key, val: Val) -> BoxFut<'_, KvsResult<()>> {
self.inner.put(key, val)
}
fn putc(&self, key: Key, val: Val, chk: Option<Val>) -> BoxFut<'_, KvsResult<()>> {
self.inner.putc(key, val, chk)
}
fn del(&self, key: Key) -> BoxFut<'_, KvsResult<()>> {
self.inner.del(key)
}
fn delc(&self, key: Key, chk: Option<Val>) -> BoxFut<'_, KvsResult<()>> {
self.inner.delc(key, chk)
}
fn delr(&self, rng: Range<Key>) -> BoxFut<'_, KvsResult<()>> {
self.inner.delr(rng)
}
fn getr(
&self,
rng: Range<Key>,
version: Option<u64>,
) -> BoxFut<'_, KvsResult<ScanResult>> {
self.inner.getr(rng, version)
}
fn keys(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, KvsResult<KeysResult>> {
self.inner.keys(rng, limit, skip, version)
}
fn keysr(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, KvsResult<KeysResult>> {
self.inner.keysr(rng, limit, skip, version)
}
fn scan(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, KvsResult<ScanResult>> {
self.inner.scan(rng, limit, skip, version)
}
fn scanr(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, KvsResult<ScanResult>> {
self.inner.scanr(rng, limit, skip, version)
}
fn new_save_point(&self) -> BoxFut<'_, KvsResult<()>> {
self.inner.new_save_point()
}
fn release_last_save_point(&self) -> BoxFut<'_, KvsResult<()>> {
self.inner.release_last_save_point()
}
fn rollback_to_save_point(&self) -> BoxFut<'_, KvsResult<()>> {
self.inner.rollback_to_save_point()
}
}
pub(super) async fn hooked_factory() -> (TransactionFactory, HookSlot) {
let hook: HookSlot = Arc::new(StdMutex::new(None));
let flavor =
crate::kvs::mem::Datastore::new(crate::kvs::mem::MemoryConfig::default())
.await
.map(DatastoreFlavor::Mem)
.unwrap();
let tf = TransactionFactory::new(
Arc::new(Notify::new()),
Box::new(HookedBuilder {
inner: Box::new(flavor),
hook: Arc::clone(&hook),
}),
Arc::new(Default::default()),
);
(tf, hook)
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_provisional_draw_is_invisible_to_other_transactions() {
let (tf, seqs) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
define_sequence(&tf, &seqs, ns, db, name, 1, 10).await.unwrap();
let repositioning = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
define_sequence_in(&repositioning, ns, db, name, 900, 10).await.unwrap();
repositioning.sequence_defined_here(ns, db, name.to_string()).await;
let drawn =
seqs.next_user_sequence_id(None, &repositioning, ns, db, name).await.unwrap();
assert_eq!(drawn, 900, "the defining transaction draws from what it wrote");
let other = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
let elsewhere = seqs.next_user_sequence_id(None, &other, ns, db, name).await.unwrap();
other.cancel().await.unwrap();
assert_eq!(
elsewhere, 1,
"a concurrent transaction was served a value derived from an uncommitted `START`"
);
repositioning.cancel().await.unwrap();
}
#[tokio::test(flavor = "multi_thread")]
async fn a_cancelled_provisional_draw_leaves_no_rows_behind() {
let (tf, seqs) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let defining = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
define_sequence_in(&defining, ns, db, name, 7, 10).await.unwrap();
defining.sequence_defined_here(ns, db, name.to_string()).await;
assert_eq!(
seqs.next_user_sequence_id(None, &defining, ns, db, name).await.unwrap(),
7,
"the defining transaction draws from what it wrote"
);
defining.cancel().await.unwrap();
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
let batches = tx.getr(Prefix::new_ba_range(ns, db, name).unwrap(), None).await.unwrap();
let states = tx.getr(Prefix::new_st_range(ns, db, name).unwrap(), None).await.unwrap();
tx.cancel().await.unwrap();
assert!(
batches.is_empty(),
"a rolled-back definition left {} claimed window(s) behind it",
batches.len()
);
assert!(
states.is_empty(),
"a rolled-back definition left {} node cursor(s) behind it",
states.len()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn an_allocator_loaded_during_a_reposition_cannot_outlive_it() {
let (tf, seqs) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
define_sequence(&tf, &seqs, ns, db, name, 1, 100).await.unwrap();
{
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
assert_eq!(seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
}
let repositioning = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
define_sequence_in(&repositioning, ns, db, name, 1, 1).await.unwrap();
repositioning.sequence_defined_here(ns, db, name.to_string()).await;
seqs.sequence_removed(ns, db, name).await;
let domain = SequenceDomain::new_user(ns, db, name);
let stale = Sequence::load(None, &seqs, &domain, 1, 100, None).await;
let drawn =
seqs.next_user_sequence_id(None, &repositioning, ns, db, name).await.unwrap();
assert_eq!(drawn, 1, "the reposition draws from what it wrote");
match repositioning.commit().await {
Err(e) => assert!(
crate::kvs::is_retryable_transaction_conflict(&e),
"the reposition failed for some reason other than the contention \
that protects it: {e}"
),
Ok(()) => {
let from_stale = match stale {
Err(_) => None,
Ok(mut stale) => {
match stale.next(&seqs, None, &domain, 100, None, None).await {
Err(SequenceError::Reset) => None,
Err(e) => {
panic!("unexpected error from the stale allocator: {e:?}")
}
Ok(v) => Some(v),
}
}
};
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
let fresh = seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap();
tx.cancel().await.unwrap();
if let Some(from_stale) = from_stale {
assert_ne!(
from_stale, fresh,
"an allocator loaded during the reposition served a value the new \
run also served"
);
}
}
}
}
async fn draw_held_across_a_reposition(
reposition_draws: bool,
) -> (SequenceResult<i64>, Vec<i64>) {
let (tf, seqs) = factory().await;
let peer = Sequences::new(tf.clone(), Uuid::from_u128(101));
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let domain = SequenceDomain::new_user(ns, db, name);
define_sequence(&tf, &seqs, ns, db, name, 1, 100).await.unwrap();
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
assert_eq!(seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
let held =
seqs.sequences.read().await.get(&domain).cloned().expect("a cached allocator");
let mut issued = Vec::new();
let repositioning = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
define_sequence_in(&repositioning, ns, db, name, 1, 1).await.unwrap();
repositioning.sequence_defined_here(ns, db, name.to_string()).await;
seqs.sequence_removed(ns, db, name).await;
if reposition_draws {
issued.push(
seqs.next_user_sequence_id(None, &repositioning, ns, db, name).await.unwrap(),
);
}
repositioning.commit().await.unwrap();
if !reposition_draws {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
issued.push(seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap());
tx.cancel().await.unwrap();
}
let from_held = held
.sequence
.lock()
.await
.next(&seqs, None, &domain, 100, None, Some(&held.evicted))
.await;
for node in [&peer, &seqs] {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, node.clone())
.await
.unwrap();
issued.push(node.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap());
tx.cancel().await.unwrap();
}
(from_held, issued)
}
#[tokio::test(flavor = "multi_thread")]
async fn a_draw_held_across_a_reposition_is_refused() {
let (from_held, issued) = draw_held_across_a_reposition(true).await;
assert!(
matches!(from_held, Err(SequenceError::Reset)),
"an allocator evicted by the reposition served {from_held:?}; the new run issued \
{issued:?}"
);
let mut distinct = issued.clone();
distinct.sort_unstable();
distinct.dedup();
assert_eq!(
distinct.len(),
issued.len(),
"the new run issued a value twice: {issued:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_draw_held_across_a_reposition_and_a_rebuild_is_refused() {
let (from_held, issued) = draw_held_across_a_reposition(false).await;
assert!(
matches!(from_held, Err(SequenceError::Reset)),
"an allocator evicted by the reposition served {from_held:?}; the new run issued \
{issued:?}"
);
let mut distinct = issued.clone();
distinct.sort_unstable();
distinct.dedup();
assert_eq!(
distinct.len(),
issued.len(),
"the new run issued a value twice: {issued:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn draws_queued_on_a_reset_allocator_do_not_resume_after_its_rebuild() {
let (tf, repositioner) = factory().await;
let peer = Sequences::new(tf.clone(), Uuid::from_u128(101));
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let domain = SequenceDomain::new_user(ns, db, name);
define_sequence(&tf, &peer, ns, db, name, 1, 100).await.unwrap();
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, peer.clone())
.await
.unwrap();
assert_eq!(peer.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
let queued =
peer.sequences.read().await.get(&domain).cloned().expect("a cached allocator");
let tx = tf
.transaction(TransactionType::Write, LockType::Optimistic, repositioner.clone())
.await
.unwrap();
define_sequence_in(&tx, ns, db, name, 1, 1).await.unwrap();
tx.sequence_defined_here(ns, db, name.to_string()).await;
repositioner.sequence_removed(ns, db, name).await;
tx.commit().await.unwrap();
let first = queued
.sequence
.lock()
.await
.next(&peer, None, &domain, 100, None, Some(&queued.evicted))
.await;
assert!(
matches!(first, Err(SequenceError::Reset)),
"the reset went unnoticed: {first:?}"
);
peer.evict_this_allocator(&Arc::new(SequenceDomain::new_user(ns, db, name)), &queued)
.await;
let mut issued = Vec::new();
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, peer.clone())
.await
.unwrap();
issued.push(peer.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap());
tx.cancel().await.unwrap();
let second = queued
.sequence
.lock()
.await
.next(&peer, None, &domain, 100, None, Some(&queued.evicted))
.await;
for node in [&repositioner, &repositioner, &peer] {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, node.clone())
.await
.unwrap();
issued.push(node.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap());
tx.cancel().await.unwrap();
}
assert!(
matches!(second, Err(SequenceError::Reset)),
"a draw queued on the reset allocator served {second:?} after its rebuild; the new \
run issued {issued:?}"
);
let mut distinct = issued.clone();
distinct.sort_unstable();
distinct.dedup();
assert_eq!(
distinct.len(),
issued.len(),
"the new run issued a value twice: {issued:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_peers_allocator_cannot_outlive_a_reposition_either() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let peer = Sequences::new(tf.clone(), Uuid::from_u128(101));
let repositioner = Sequences::new(tf.clone(), Uuid::from_u128(102));
define_sequence(&tf, &peer, ns, db, name, 1, 100).await.unwrap();
{
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, peer.clone())
.await
.unwrap();
assert_eq!(peer.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
}
let tx = tf
.transaction(TransactionType::Write, LockType::Optimistic, repositioner.clone())
.await
.unwrap();
define_sequence_in(&tx, ns, db, name, 1, 1).await.unwrap();
tx.sequence_defined_here(ns, db, name.to_string()).await;
repositioner.sequence_removed(ns, db, name).await;
let domain = SequenceDomain::new_user(ns, db, name);
let loaded = Sequence::load(None, &peer, &domain, 1, 100, None).await;
let drawn = repositioner.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap();
assert_eq!(drawn, 1, "the reposition draws from what it wrote");
let committed = tx.commit().await;
assert!(
committed.is_err() || loaded.is_err(),
"a peer loaded a window inside the reposition and both were allowed to stand"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_definition_rolled_back_by_a_save_point_is_not_served_from() {
for rolled_back_start in [900, 1] {
let (tf, seqs) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
define_sequence(&tf, &seqs, ns, db, name, 1, 10).await.unwrap();
for expected in 1..=3 {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
assert_eq!(
seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(),
expected
);
tx.cancel().await.unwrap();
}
let tx = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
tx.new_save_point().await.unwrap();
define_sequence_in(&tx, ns, db, name, rolled_back_start, 10).await.unwrap();
tx.sequence_defined_here(ns, db, name.to_string()).await;
seqs.sequence_removed(ns, db, name).await;
assert_eq!(
seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(),
rolled_back_start,
"inside the save point the redefinition is in force"
);
tx.rollback_to_save_point().await.unwrap();
assert_eq!(
seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(),
4,
"after rolling back a redefinition to START {rolled_back_start}, the transaction \
was served from it rather than from the committed run"
);
tx.commit().await.unwrap();
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
assert_eq!(
seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(),
5,
"the committed run did not continue after the transaction"
);
tx.cancel().await.unwrap();
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_window_claimed_inside_a_rolled_back_save_point_is_not_served_from() {
let (tf, seqs) = factory().await;
let peer = Sequences::new(tf.clone(), Uuid::from_u128(101));
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let tx = tf
.transaction(TransactionType::Write, LockType::Optimistic, seqs.clone())
.await
.unwrap();
define_sequence_in(&tx, ns, db, name, 1, 10).await.unwrap();
tx.sequence_defined_here(ns, db, name.to_string()).await;
tx.new_save_point().await.unwrap();
seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap();
tx.rollback_to_save_point().await.unwrap();
let kept = seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap();
tx.commit().await.unwrap();
let mut from_peer = Vec::new();
for _ in 0..3 {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, peer.clone())
.await
.unwrap();
from_peer.push(peer.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap());
tx.cancel().await.unwrap();
}
assert!(
!from_peer.contains(&kept),
"a peer issued {kept} again, which the defining transaction committed; the peer \
issued {from_peer:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_claim_straddling_a_reposition_does_not_survive_it() {
use std::collections::HashSet;
use std::sync::Mutex as StdMutex;
let (tf, hook) = hooked::hooked_factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
const BATCH: u32 = 1000;
let n1 = Sequences::new(tf.clone(), Uuid::from_u128(101));
let n2 = Sequences::new(tf.clone(), Uuid::from_u128(102));
define_sequence(&tf, &n1, ns, db, name, 1, BATCH).await.unwrap();
{
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, n1.clone())
.await
.unwrap();
assert_eq!(n1.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
}
let a1_slot: Arc<StdMutex<Option<Sequence>>> = Arc::new(StdMutex::new(None));
{
let tf = tf.clone();
let n1 = n1.clone();
let a1_slot = Arc::clone(&a1_slot);
*hook.lock().unwrap() = Some(Box::new(move || {
Box::pin(async move {
define_sequence(&tf, &n1, ns, db, name, 500, BATCH).await.unwrap();
let domain = SequenceDomain::new_user(ns, db, name);
let a1 = Sequence::load(None, &n1, &domain, 500, BATCH, None)
.await
.expect("N1's claim after the reposition");
*a1_slot.lock().unwrap() = Some(a1);
})
}));
}
let domain = SequenceDomain::new_user(ns, db, name);
let a2 = match Sequence::load(None, &n2, &domain, 1, BATCH, None).await {
Err(SequenceError::Reset) => {
Sequence::load(None, &n2, &domain, 1, BATCH, None).await.unwrap()
}
Ok(a2) => a2,
Err(e) => panic!("N2's load failed: {e}"),
};
assert!(hook.lock().unwrap().is_none(), "the straddle did not run");
let a1 = a1_slot.lock().unwrap().take().expect("N1's allocator");
let (w1, w2) = ((a1.st.next, a1.to), (a2.st.next, a2.to));
let mut a1 = a1;
let mut a2 = a2;
let mut issued = HashSet::new();
let mut twice = Vec::new();
for _ in 0..502 {
let v = a1.next(&n1, None, &domain, BATCH, None, None).await.unwrap();
if !issued.insert(v) {
twice.push(v);
}
}
for _ in 0..3 {
let v = a2.next(&n2, None, &domain, BATCH, None, None).await.unwrap();
if !issued.insert(v) {
twice.push(v);
}
}
assert!(twice.is_empty(), "ids issued twice after a reposition: {twice:?}");
assert!(
w1.1 <= w2.0 || w2.1 <= w1.0,
"the windows overlap: N1 [{}, {}) and N2 [{}, {})",
w1.0,
w1.1,
w2.0,
w2.1
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_reset_while_building_is_rebuilt_not_reported() {
let (tf, hook) = hooked::hooked_factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
const BATCH: u32 = 1000;
let n1 = Sequences::new(tf.clone(), Uuid::from_u128(101));
let n2 = Sequences::new(tf.clone(), Uuid::from_u128(102));
define_sequence(&tf, &n1, ns, db, name, 1, BATCH).await.unwrap();
{
let tf = tf.clone();
let n1 = n1.clone();
*hook.lock().unwrap() = Some(Box::new(move || {
Box::pin(async move {
define_sequence(&tf, &n1, ns, db, name, 500, BATCH).await.unwrap();
})
}));
}
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, n2.clone())
.await
.unwrap();
let got = n2
.next_user_sequence_id(None, &tx, ns, db, name)
.await
.expect("a reset seen while building must be rebuilt through, not reported");
tx.cancel().await.unwrap();
assert!(hook.lock().unwrap().is_none(), "the straddle did not run");
assert_eq!(got, 500, "the rebuild allocates from the committed `START`");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_first_cursor_write_straddling_a_reposition_does_not_survive_it() {
use std::collections::HashSet;
use std::sync::Mutex as StdMutex;
let (tf, hook) = hooked::hooked_factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
const BATCH: u32 = 10;
let n1 = Sequences::new(tf.clone(), Uuid::from_u128(101));
let n2 = Sequences::new(tf.clone(), Uuid::from_u128(102));
define_sequence(&tf, &n1, ns, db, name, 1, BATCH).await.unwrap();
let domain = SequenceDomain::new_user(ns, db, name);
let mut a1 = Sequence::load(None, &n1, &domain, 1, BATCH, None).await.unwrap();
assert_eq!((a1.st.next, a1.to), (1, 11));
let a2_slot: Arc<StdMutex<Option<Sequence>>> = Arc::new(StdMutex::new(None));
{
let tf = tf.clone();
let n2 = n2.clone();
let a2_slot = Arc::clone(&a2_slot);
*hook.lock().unwrap() = Some(Box::new(move || {
Box::pin(async move {
reposition_to_resume(&tf, &n2, ns, db, name, BATCH).await.unwrap();
let domain = SequenceDomain::new_user(ns, db, name);
let a2 = Sequence::load(None, &n2, &domain, 1, BATCH, None)
.await
.expect("N2's claim after the reposition");
*a2_slot.lock().unwrap() = Some(a2);
})
}));
}
let mut issued = HashSet::new();
let mut twice = Vec::new();
match a1.next(&n1, None, &domain, BATCH, None, None).await {
Ok(v) => {
issued.insert(v);
}
Err(SequenceError::Reset) => {
a1 = Sequence::load(None, &n1, &domain, 1, BATCH, None).await.unwrap();
}
Err(e) => panic!("N1's first write failed: {e}"),
}
assert!(hook.lock().unwrap().is_none(), "the straddle did not run");
let mut a2 = a2_slot.lock().unwrap().take().expect("N2's allocator");
let (w1, w2) = ((a1.st.next, a1.to), (a2.st.next, a2.to));
for _ in 0..(BATCH as usize - 1) {
let v = a1.next(&n1, None, &domain, BATCH, None, None).await.unwrap();
if !issued.insert(v) {
twice.push(v);
}
}
for _ in 0..BATCH as usize {
let v = a2.next(&n2, None, &domain, BATCH, None, None).await.unwrap();
if !issued.insert(v) {
twice.push(v);
}
}
assert!(twice.is_empty(), "ids issued twice after a reposition: {twice:?}");
assert!(
w1.1 <= w2.0 || w2.1 <= w1.0,
"the windows overlap: N1 [{}, {}) and N2 [{}, {})",
w1.0,
w1.1,
w2.0,
w2.1
);
}
#[tokio::test]
async fn a_peer_reloads_a_sequence_reset_elsewhere() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let peer = Sequences::new(tf.clone(), Uuid::from_u128(1));
const BATCH: u32 = 1000;
define_sequence(&tf, &peer, ns, db, name, 1, BATCH).await.unwrap();
let serve = async |tf: TransactionFactory, peer: Sequences| -> i64 {
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, peer.clone())
.await
.unwrap();
let v = peer.next_user_sequence_id(None, &tx, ns, db, name).await;
tx.cancel().await.unwrap();
v.unwrap()
};
for expected in 1..=2 {
assert_eq!(serve(tf.clone(), peer.clone()).await, expected);
}
define_sequence(&tf, &peer, ns, db, name, 500, BATCH).await.unwrap();
assert_eq!(
serve(tf.clone(), peer.clone()).await,
500,
"a peer kept serving a sequence that was repositioned elsewhere"
);
}
#[tokio::test]
async fn a_first_load_ignores_a_start_the_caller_read_before_a_reset() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
const BATCH: u32 = 1000;
define_sequence(&tf, &tf_sequences(&tf), ns, db, name, 1, BATCH).await.unwrap();
let first = Sequences::new(tf.clone(), Uuid::from_u128(11));
{
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, first.clone())
.await
.unwrap();
assert_eq!(first.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
}
let fresh = Sequences::new(tf.clone(), Uuid::from_u128(12));
let stale = tf
.transaction(TransactionType::Read, LockType::Optimistic, fresh.clone())
.await
.unwrap();
let _ = stale.get_db_sequence(ns, db, name, None).await.unwrap();
define_sequence(&tf, &fresh, ns, db, name, 500, BATCH).await.unwrap();
let got = fresh.next_user_sequence_id(None, &stale, ns, db, name).await.unwrap();
stale.cancel().await.unwrap();
assert_eq!(
got, 500,
"a first load reissued from a start the caller read before the reset"
);
}
#[tokio::test]
async fn a_reset_mid_statement_allocates_from_the_new_definition() {
let (tf, _) = factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
let node = Sequences::new(tf.clone(), Uuid::from_u128(7));
const BATCH: u32 = 1000;
define_sequence(&tf, &node, ns, db, name, 1, BATCH).await.unwrap();
{
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, node.clone())
.await
.unwrap();
assert_eq!(node.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap(), 1);
tx.cancel().await.unwrap();
}
let stale = tf
.transaction(TransactionType::Read, LockType::Optimistic, node.clone())
.await
.unwrap();
let _ = stale.get_db_sequence(ns, db, name, None).await.unwrap();
define_sequence(&tf, &node, ns, db, name, 500, BATCH).await.unwrap();
let got = node.next_user_sequence_id(None, &stale, ns, db, name).await.unwrap();
stale.cancel().await.unwrap();
assert_eq!(
got, 500,
"the retry allocated from the definition the caller had already read"
);
}
}
#[cfg(feature = "kv-rocksdb")]
mod over_rocksdb {
use std::sync::Arc;
use tokio::sync::Notify;
use uuid::Uuid;
use super::support::{define_sequence, reposition_to_resume, tf_sequences};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::kvs::ds::{DatastoreFlavor, TransactionFactory};
use crate::kvs::sequences::Sequences;
use crate::kvs::{LockType, TransactionType};
#[cfg(feature = "kv-rocksdb")]
async fn rocksdb_factory() -> (TransactionFactory, Sequences, tempfile::TempDir) {
let dir = tempfile::TempDir::new().unwrap();
let flavor = crate::kvs::rocksdb::Datastore::new(
&dir.path().to_string_lossy(),
crate::kvs::rocksdb::RocksDbConfig::default(),
)
.await
.map(DatastoreFlavor::RocksDB)
.unwrap();
let tf = TransactionFactory::new(
Arc::new(Notify::new()),
Box::new(flavor),
Arc::new(Default::default()),
);
let sequences = Sequences::new(tf.clone(), Uuid::new_v4());
(tf, sequences, dir)
}
#[cfg(feature = "kv-rocksdb")]
#[tokio::test(flavor = "multi_thread")]
async fn concurrent_allocation_never_reissues_across_repositions() {
use std::collections::HashSet;
use std::time::Duration;
let (tf, _, _dir) = rocksdb_factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
const BATCH: u32 = 10;
const NODES: u128 = 3;
const TASKS_PER_NODE: usize = 3;
const ALLOCATIONS: usize = 40;
const STRIDE: i64 = 10_000;
const EPOCHS: i64 = 3;
define_sequence(&tf, &tf_sequences(&tf), ns, db, name, 1, BATCH).await.unwrap();
let mut workers = Vec::new();
for node in 0..NODES {
let seqs = Sequences::new(tf.clone(), Uuid::from_u128(node));
for _ in 0..TASKS_PER_NODE {
let seqs = seqs.clone();
let tf = tf.clone();
workers.push(tokio::spawn(async move {
let mut got = Vec::new();
for _ in 0..ALLOCATIONS {
let tx = tf
.transaction(
TransactionType::Read,
LockType::Optimistic,
seqs.clone(),
)
.await
.unwrap();
let res = seqs.next_user_sequence_id(None, &tx, ns, db, name).await;
tx.cancel().await.unwrap();
match res {
Ok(v) => got.push(v),
Err(e)
if crate::kvs::sequences::is_sequence_reset(&e)
|| crate::kvs::is_retryable_transaction_conflict(&e) => {}
Err(e) => panic!("allocation failed: {e}"),
}
}
got
}));
}
}
let repositioner = {
let tf = tf.clone();
tokio::spawn(async move {
let sqs = tf_sequences(&tf);
for epoch in 1..=EPOCHS {
tokio::time::sleep(Duration::from_millis(5)).await;
let start = epoch * STRIDE;
loop {
match define_sequence(&tf, &sqs, ns, db, name, start, BATCH).await {
Ok(()) => break,
Err(e) if crate::kvs::is_retryable_transaction_conflict(&e) => {
continue;
}
Err(e) => panic!("reposition failed: {e}"),
}
}
}
})
};
let mut issued = Vec::new();
for w in workers {
issued.extend(
tokio::time::timeout(Duration::from_secs(60), w)
.await
.expect("allocation stalled")
.expect("a worker panicked"),
);
}
tokio::time::timeout(Duration::from_secs(60), repositioner)
.await
.expect("repositioning stalled")
.expect("the repositioner panicked");
assert!(!issued.is_empty(), "the fixture issued nothing, so it asserts nothing");
let unique: HashSet<i64> = issued.iter().copied().collect();
assert_eq!(
unique.len(),
issued.len(),
"an id was issued twice across a reposition ({} of {} distinct)",
unique.len(),
issued.len()
);
let seqs = tf_sequences(&tf);
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
let after = seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap();
tx.cancel().await.unwrap();
assert!(
after >= EPOCHS * STRIDE,
"allocation resumed below the last reposition: {after} < {}",
EPOCHS * STRIDE
);
}
#[cfg(feature = "kv-rocksdb")]
#[tokio::test(flavor = "multi_thread")]
async fn repositioning_a_live_sequence_to_resume_never_reissues() {
use std::collections::HashSet;
use std::time::Duration;
let (tf, _, _dir) = rocksdb_factory().await;
let ns = NamespaceId(1);
let db = DatabaseId(2);
let name = "sq";
const BATCH: u32 = 10;
const NODES: u128 = 3;
const TASKS_PER_NODE: usize = 3;
const ALLOCATIONS: usize = 40;
const EPOCHS: usize = 3;
define_sequence(&tf, &tf_sequences(&tf), ns, db, name, 1, BATCH).await.unwrap();
let mut workers = Vec::new();
for node in 0..NODES {
let seqs = Sequences::new(tf.clone(), Uuid::from_u128(node));
for _ in 0..TASKS_PER_NODE {
let seqs = seqs.clone();
let tf = tf.clone();
workers.push(tokio::spawn(async move {
let mut got = Vec::new();
for _ in 0..ALLOCATIONS {
let tx = tf
.transaction(
TransactionType::Read,
LockType::Optimistic,
seqs.clone(),
)
.await
.unwrap();
let res = seqs.next_user_sequence_id(None, &tx, ns, db, name).await;
tx.cancel().await.unwrap();
match res {
Ok(v) => got.push(v),
Err(e)
if crate::kvs::sequences::is_sequence_reset(&e)
|| crate::kvs::is_retryable_transaction_conflict(&e) => {}
Err(e) => panic!("allocation failed: {e}"),
}
}
got
}));
}
}
let repositioner = {
let tf = tf.clone();
tokio::spawn(async move {
let sqs = tf_sequences(&tf);
for _ in 0..EPOCHS {
tokio::time::sleep(Duration::from_millis(5)).await;
loop {
match reposition_to_resume(&tf, &sqs, ns, db, name, BATCH).await {
Ok(()) => break,
Err(e) if crate::kvs::is_retryable_transaction_conflict(&e) => {
continue;
}
Err(e) => panic!("reposition failed: {e}"),
}
}
}
})
};
let mut issued = Vec::new();
for w in workers {
issued.extend(
tokio::time::timeout(Duration::from_secs(60), w)
.await
.expect("allocation stalled")
.expect("a worker panicked"),
);
}
tokio::time::timeout(Duration::from_secs(60), repositioner)
.await
.expect("repositioning stalled")
.expect("the repositioner panicked");
assert!(!issued.is_empty(), "the fixture issued nothing, so it asserts nothing");
let unique: HashSet<i64> = issued.iter().copied().collect();
assert_eq!(
unique.len(),
issued.len(),
"an id was issued twice across a reposition ({} of {} distinct)",
unique.len(),
issued.len()
);
let top = issued.iter().copied().max().unwrap();
let seqs = tf_sequences(&tf);
let tx = tf
.transaction(TransactionType::Read, LockType::Optimistic, seqs.clone())
.await
.unwrap();
let after = seqs.next_user_sequence_id(None, &tx, ns, db, name).await.unwrap();
tx.cancel().await.unwrap();
assert!(
after > top,
"allocation resumed at {after}, below an id already issued ({top})"
);
}
}
}