use std::borrow::Cow;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Result, bail};
use common::time::sleep;
use reblessive::TreeStack;
use reblessive::tree::Stk;
pub(crate) use surrealdb_datastore::values::event_queue::AsyncEventRecord;
use surrealdb_kvs::TransactionType::Write;
use surrealdb_kvs::timestamp::HlcTimeStamp;
#[cfg(not(target_family = "wasm"))]
use tokio::spawn;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::catalog::{EventDefinition, FromStored, StoredEventDefinition};
use crate::ctx::{Context, FrozenContext};
use crate::dbs::{Options, Session};
use crate::doc::{Action, CursorDoc, Document, DocumentContext, Error};
use crate::exe::FlowResultExt as _;
use crate::iam::AuthLimit;
use crate::key::schema::{EventQueueKey, EventQueuePrefix};
use crate::key::{KVKeyDecode, KVValue};
use crate::kvs::sequences::Sequences;
use crate::kvs::tasklease::LeaseHandler;
#[cfg(test)]
use crate::kvs::testing::{
NonRetryableErrorSite, RetryableConflictSite, maybe_inject_non_retryable_error,
maybe_inject_retryable_conflict,
};
use crate::kvs::{
Datastore, NORMAL_BATCH_SIZE, Transaction, TransactionFactory, TransactionType,
is_indeterminate_commit, is_retryable_transaction_conflict,
};
use crate::val::Value;
const EVENT_CONFLICT_RETRIES: u32 = 3;
const EVENT_CONFLICT_RETRY_SLEEP: Duration = Duration::from_millis(50);
impl Document {
pub(super) async fn process_table_events(
&mut self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
action: Action,
) -> Result<()> {
if opt.import {
return Ok(());
}
if !self.is_modified() {
return Ok(());
}
let opt = &opt.new_with_perms(false);
if self.doc_ctx.ev()?.is_empty() {
return Ok(());
}
let input = self.materialize_input_value(stk, ctx, opt).await?;
self.process_events(stk, ctx, opt, action, input).await
}
pub(super) async fn process_events(
&mut self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
action: Action,
input: Option<Arc<Value>>,
) -> Result<()> {
if opt.import {
return Ok(());
}
if !self.is_modified() {
return Ok(());
}
let opt = &opt.new_with_perms(false);
for ev in self.doc_ctx.ev()?.iter() {
let opt = opt.limited_by(&AuthLimit::try_from(&ev.auth_limit)?);
let evt = match action {
Action::Create => Value::from("CREATE"),
Action::Update => Value::from("UPDATE"),
Action::Delete => Value::from("DELETE"),
};
let after = self.current.doc.as_arc();
let before = self.initial.doc.as_arc();
let doc = if action == Action::Delete {
&mut self.initial
} else {
&mut self.current
};
let mut ctx = Context::new_child(ctx);
ctx.add_value("after", after);
ctx.add_value("before", before);
ctx.add_value("event", evt.into());
ctx.add_value("value", doc.doc.as_arc());
ctx.add_value("input", input.clone().unwrap_or_default());
let ctx = ctx.freeze();
let val = stk
.run(|stk| crate::legacy::expr_compute(&ev.when, stk, &ctx, &opt, Some(doc)))
.await
.catch_return()
.map_err(|e| anyhow::anyhow!("Error while processing event {}: {e}", ev.name))?;
if val.is_truthy() {
if ev.is_async() {
Self::process_event_async(ctx, opt, ev, &self.doc_ctx, doc).await?;
} else {
Self::process_event_sync(stk, ctx, opt, None, ev, doc).await?;
}
}
}
Ok(())
}
async fn process_event_sync(
stk: &mut Stk,
ctx: FrozenContext,
opt: Options,
_lh: Option<&LeaseHandler>,
ev: &EventDefinition,
doc: &CursorDoc,
) -> Result<()> {
for then in ev.then.iter() {
stk.run(|stk| crate::legacy::expr_compute(then, stk, &ctx, &opt, Some(doc)))
.await
.catch_return()
.map_err(|e| anyhow::anyhow!("Error while processing event {}: {e}", ev.name))?;
}
Ok(())
}
async fn process_event_async(
ctx: FrozenContext,
opt: Options,
ev: &EventDefinition,
doc_ctx: &DocumentContext,
cursor_doc: &mut CursorDoc,
) -> Result<()> {
let node_id = ctx.node_id();
let ts = HlcTimeStamp::next();
let db = doc_ctx.db();
let tx = ctx.tx();
let key = EventQueueKey {
ns: db.namespace_id,
db: db.database_id,
tb: Cow::Borrowed(&ev.target_table),
ev: Cow::Borrowed(&ev.name),
ts: ts.0,
node_id,
};
let event_record = queue_async_event(&opt, &ctx, ev.stored(), cursor_doc)?;
tx.put_key(&key, &event_record).await?;
tx.trigger_async_event();
Ok(())
}
}
fn queue_async_event(
opt: &Options,
ctx: &FrozenContext,
event_definition: &StoredEventDefinition,
cursor_doc: &CursorDoc,
) -> Result<AsyncEventRecord> {
let (ns, db) = opt.arc_ns_db()?;
if let Some(d) = opt.async_event_depth()
&& d >= event_definition.max_depth()
{
bail!(Error::EvReachMaxDepth(event_definition.name.to_string(), d))
}
Ok(AsyncEventRecord {
attempt: 0,
event_depth: opt.async_event_depth().map(|d| d + 1).unwrap_or(0),
rid: cursor_doc.rid.clone(),
cursor_record: cursor_doc.doc.clone().into_read_only(),
fields_computed: cursor_doc.fields_computed,
ns,
db,
perms: opt.perms,
auth_enabled: ctx.auth_enabled(),
values: ctx.collect_values(HashMap::new()),
auth_with_limit: Arc::clone(&opt.auth),
event_definition: event_definition.clone(),
})
}
fn build_event_context(record: &AsyncEventRecord, ctx: &FrozenContext) -> FrozenContext {
let mut ctx = Context::new_child(ctx);
ctx.add_values(record.values.clone());
ctx.auth_enabled = record.auth_enabled;
ctx.freeze()
}
async fn build_event_options(
record: &AsyncEventRecord,
tx: &Transaction,
parent_opts: &Options,
eq: &EventQueueKey<'_>,
) -> Result<Options> {
let ns = tx.expect_ns_by_name(&record.ns).await?;
if ns.namespace_id != eq.ns {
bail!(Error::EvNamespaceMismatch(
record.event_definition.name.to_string(),
ns.name.to_string(),
));
}
let db = tx.expect_db_by_name(&record.ns, &record.db).await?;
if db.database_id != eq.db {
bail!(Error::EvDatabaseMismatch(
record.event_definition.name.to_string(),
db.name.to_string(),
));
}
let opt = parent_opts.clone();
let opt = opt
.with_perms(record.perms)
.with_auth(Arc::clone(&record.auth_with_limit))
.with_async_event_depth(record.event_depth)
.with_ns(Some(Arc::clone(&record.ns)))
.with_db(Some(Arc::clone(&record.db)));
Ok(opt)
}
fn build_event_cursor_doc(record: &AsyncEventRecord) -> CursorDoc {
CursorDoc {
rid: record.rid.clone(),
ir: None,
doc: Arc::clone(&record.cursor_record).into(),
fields_computed: record.fields_computed,
}
}
pub async fn process_next_events_batch(ds: &Datastore, lh: Option<&LeaseHandler>) -> Result<usize> {
let res = {
if let Some(lh) = lh.as_ref() {
lh.try_maintain_lease().await?;
}
let tx = ds.transaction(TransactionType::Read).await?;
let range = EventQueuePrefix {}.range()?;
let res = catch!(tx, tx.scan_raw(range, NORMAL_BATCH_SIZE, 0, None).await);
tx.cancel().await?;
res
};
let count = res.len();
process_events_batch(ds, res, lh).await?;
Ok(count)
}
#[cfg(not(target_family = "wasm"))]
pub(crate) async fn process_events_batch(
ds: &Datastore,
res: Vec<(Vec<u8>, Vec<u8>)>,
lh: Option<&LeaseHandler>,
) -> Result<()> {
if res.is_empty() {
return Ok(());
}
let concurrency: usize = num_cpus::get().max(4);
let workers = res.len().min(concurrency);
let mut join_handles = Vec::with_capacity(workers);
let (sender, receiver) = async_channel::bounded::<AsyncEventContext>(workers);
for _ in 0..workers {
let receiver = receiver.clone();
let jh = spawn(async move {
let mut stack = TreeStack::new();
while let Ok(event_context) = receiver.recv().await {
stack
.enter(|stk| stk.run(|stk| event_context.run_event_checked(stk)))
.finish()
.await;
}
});
join_handles.push(jh);
}
for (k, v) in res {
let Some(v) = decode_queued(&k, &v) else {
continue;
};
match AsyncEventContext::new(ds, lh.cloned(), k, v) {
Ok(event_context) => {
sender.send(event_context).await?;
}
Err(e) => {
error!("Unexpected Error while processing event: {e}");
}
};
if let Some(lh) = lh {
lh.try_maintain_lease().await?;
}
}
sender.close();
for jh in join_handles {
if let Err(e) = jh.await {
error!("Error while processing an event: {e}");
}
}
Ok(())
}
#[cfg(target_family = "wasm")]
pub(crate) async fn process_events_batch(
ds: &Datastore,
res: Vec<(Vec<u8>, Vec<u8>)>,
lh: Option<&LeaseHandler>,
) -> Result<()> {
let mut stack = TreeStack::new();
for (k, v) in res {
if let Some(lh) = lh {
lh.try_maintain_lease().await?;
}
let Some(v) = decode_queued(&k, &v) else {
continue;
};
let event_context = AsyncEventContext::new(ds, lh.cloned(), k, v)?;
stack.enter(|stk| stk.run(|stk| event_context.run_event_checked(stk))).finish().await;
}
Ok(())
}
fn decode_queued(k: &[u8], v: &[u8]) -> Option<AsyncEventRecord> {
if let Err(e) = EventQueueKey::decode_key(k) {
error!("Skipping async event queue entry with an undecodable key: {e} - Key: {k:?}");
return None;
}
match KVValue::kv_decode_value(v, ()) {
Ok(ev) => Some(ev),
Err(e) => {
error!("Skipping undecodable async event queue entry: {e} - Key: {k:?}");
None
}
}
}
struct AsyncEventContext {
ctx: Option<Context>,
opt: Options,
tf: TransactionFactory,
sequences: Sequences,
lh: Option<LeaseHandler>,
k: Vec<u8>,
v: Option<AsyncEventRecord>,
}
impl AsyncEventContext {
fn new(
ds: &Datastore,
lh: Option<LeaseHandler>,
k: Vec<u8>,
v: AsyncEventRecord,
) -> Result<Self> {
Ok(Self {
ctx: Some(ds.setup_ctx()?),
opt: ds.setup_options(&Session::default()),
tf: ds.transaction_factory().clone(),
sequences: ds.sequences().clone(),
lh,
k,
v: Some(v),
})
}
async fn run_event_checked(mut self, stk: &mut Stk) {
if let Some(ctx) = self.ctx.take()
&& let Some(v) = self.v.take()
&& let Err(e) = self.run_event(stk, ctx, v).await
{
error!("Unexpected error while processing an event. Error: {e} - Key: {:?}", self.k);
}
}
async fn new_write_tx(&self) -> Result<Transaction> {
self.tf.transaction(Write, self.sequences.clone()).await
}
async fn run_event(
&mut self,
stk: &mut Stk,
ctx: Context,
mut ev: AsyncEventRecord,
) -> Result<()> {
let eq = EventQueueKey::decode_key(&self.k)?;
let base = ctx.freeze();
let mut conflicts = 0;
let err = loop {
let tx = self.new_write_tx().await?;
let mut ctx = Context::new_child(&base);
ctx.set_transaction(Arc::new(tx));
let ctx = ctx.freeze();
let Err(e) = self.run_event_once(stk, &ctx, &eq, &ev).await else {
return Ok(());
};
let transient = is_retryable_transaction_conflict(&e) || is_indeterminate_commit(&e);
if transient && conflicts < EVENT_CONFLICT_RETRIES {
conflicts += 1;
debug!("Re-running the event `{}` in place: {e}", eq.ev);
sleep(EVENT_CONFLICT_RETRY_SLEEP).await;
continue;
}
break e;
};
if is_indeterminate_commit(&err) {
return Err(err);
}
if let Some(final_error) = Self::is_final_error(&err).await? {
let tx = self.new_write_tx().await?;
return Self::final_error(tx, &eq, final_error).await;
}
let tx = self.new_write_tx().await?;
Self::retry_attempt(tx, err, &eq, &mut ev).await
}
async fn run_event_once(
&self,
stk: &mut Stk,
ctx: &FrozenContext,
eq: &EventQueueKey<'_>,
ev: &AsyncEventRecord,
) -> Result<()> {
let tx = ctx.tx();
match tx.exists_key(eq, None).await {
Ok(true) => {}
Ok(false) => return tx.cancel().await,
Err(e) => {
let _ = tx.cancel().await;
return Err(e);
}
}
if let Err(e) = Self::process_event(stk, ctx, &self.opt, self.lh.as_ref(), eq, ev).await {
let _ = tx.cancel().await;
return Err(e);
}
if let Err(e) = tx.del_key(eq).await {
let _ = tx.cancel().await;
return Err(e);
}
#[cfg(test)]
if let Err(e) =
maybe_inject_retryable_conflict(RetryableConflictSite::AsyncEventCommit, ctx.node_id())
{
let _ = tx.cancel().await;
return Err(e);
}
#[cfg(test)]
if let Err(e) = maybe_inject_non_retryable_error(
NonRetryableErrorSite::AsyncEventCommitDiscardedUnknown,
ctx.node_id(),
) {
let _ = tx.cancel().await;
return Err(e);
}
let res = tx.commit().await;
#[cfg(test)]
if res.is_ok() {
maybe_inject_non_retryable_error(
NonRetryableErrorSite::AsyncEventCommitAppliedUnknown,
ctx.node_id(),
)?;
}
res
}
async fn retry_attempt(
tx: Transaction,
e: anyhow::Error,
eq: &EventQueueKey<'_>,
ev: &mut AsyncEventRecord,
) -> Result<()> {
if !catch!(tx, tx.exists_key(eq, None).await) {
return tx.cancel().await;
}
ev.attempt += 1;
if ev.attempt <= ev.event_definition.retry() {
catch!(tx, tx.set_key(eq, ev).await);
} else {
warn!(
"Final error after processing the event `{}` on table {} {} times: {e}",
eq.ev, ev.event_definition.target_table, ev.attempt
);
catch!(tx, tx.del_key(eq).await);
}
catch!(tx, tx.commit().await);
Ok(())
}
async fn is_final_error(e: &anyhow::Error) -> Result<Option<&Error>> {
let se: Option<&Error> = e.downcast_ref();
if matches!(
se,
Some(Error::EvNamespaceMismatch(..))
| Some(Error::EvDatabaseMismatch(..))
| Some(Error::EvReachMaxDepth(..))
) {
Ok(se)
} else {
Ok(None)
}
}
async fn final_error(tx: Transaction, eq: &EventQueueKey<'_>, e: &Error) -> Result<()> {
warn!("Event processing failed: {:?}", e);
catch!(tx, tx.del_key(eq).await);
catch!(tx, tx.commit().await);
Ok(())
}
async fn process_event(
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
lh: Option<&LeaseHandler>,
eq: &EventQueueKey<'_>,
ev: &AsyncEventRecord,
) -> Result<()> {
let ctx = build_event_context(ev, ctx);
let opt = build_event_options(ev, &ctx.tx(), opt, eq).await?;
let doc = build_event_cursor_doc(ev);
let compiled = EventDefinition::from_stored(&ev.event_definition)?;
Document::process_event_sync(stk, ctx, opt, lh, &compiled, &doc).await
}
}
#[cfg(test)]
mod tests {
use std::borrow::Cow;
use uuid::Uuid;
use super::{AsyncEventContext, decode_queued};
use crate::catalog::providers::CatalogProvider;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::dbs::Session;
use crate::key::schema::{EventQueueKey, EventQueuePrefix};
use crate::key::{KVKey, KVKeyDecode};
use crate::kvs::Datastore;
use crate::kvs::TransactionType::{Read, Write};
fn valid_key() -> Vec<u8> {
EventQueueKey {
ns: NamespaceId(1),
db: DatabaseId(2),
tb: Cow::Owned("tb".into()),
ev: Cow::Owned("ev".into()),
ts: 42,
node_id: Uuid::from_u128(7),
}
.encode_key()
.expect("a well-formed queue key must encode")
.to_vec()
}
#[test]
fn a_well_formed_queue_key_still_decodes() {
let encoded = valid_key();
let decoded = EventQueueKey::decode_key(&encoded).expect("valid key must decode");
assert_eq!(decoded.ns, NamespaceId(1));
assert_eq!(decoded.db, DatabaseId(2));
assert_eq!(decoded.ts, 42);
assert_eq!(decoded.node_id, Uuid::from_u128(7));
}
#[test]
fn an_undecodable_key_is_skipped() {
assert!(
decode_queued(b"/!eq\xff-not-a-valid-entry", &[]).is_none(),
"an undecodable key must be skipped, not passed on to run_event"
);
}
#[test]
fn an_undecodable_value_behind_a_valid_key_is_skipped() {
assert!(
decode_queued(&valid_key(), b"not-an-async-event-record").is_none(),
"an undecodable value must still be skipped"
);
}
async fn queued_entries(ds: &Datastore) -> anyhow::Result<Vec<(Vec<u8>, Vec<u8>)>> {
let tx = ds.transaction(Read).await?;
let res = tx.scan_raw(EventQueuePrefix {}.range()?, 1000, 0, None).await;
tx.cancel().await?;
res
}
#[tokio::test]
async fn a_retry_does_not_requeue_an_entry_that_already_ran() -> anyhow::Result<()> {
let ds = Datastore::new("memory").await?;
let session = Session::owner().with_ns("test").with_db("test");
let tx = ds.transaction(Write).await?;
tx.ensure_ns_db(None, "test", "test").await?;
tx.commit().await?;
let sql = "DEFINE TABLE person SCHEMALESS;
DEFINE EVENT log ON person ASYNC RETRY 1 THEN (CREATE logged);
CREATE person:1 RETURN NONE;";
for res in ds.execute(sql, &session, None).await? {
res.result?;
}
let (k, v) = queued_entries(&ds).await?.pop().expect("the event must be queued");
let eq = EventQueueKey::decode_key(&k)?;
let mut ev = decode_queued(&k, &v).expect("the queued entry must decode");
let tx = ds.transaction(Write).await?;
tx.del_key(&eq).await?;
tx.commit().await?;
let tx = ds.transaction(Write).await?;
AsyncEventContext::retry_attempt(tx, anyhow::anyhow!("the event failed"), &eq, &mut ev)
.await?;
assert!(queued_entries(&ds).await?.is_empty(), "the retry must not write the entry back");
Ok(())
}
}