use std::sync::Arc;
use anyhow::Result;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider, TableProvider};
use crate::catalog::{self, DatabaseId, NamespaceId};
use crate::ctx::FrozenContext;
use crate::dbs::{RoutedNotification, SendKill};
use crate::key::schema::{NodeLiveQueryKey, SubscriptionKey};
use crate::kvs::Transaction;
use crate::types::{PublicAction, PublicNotification, PublicValue};
use crate::val::TableName;
pub(crate) async fn kill_table_subscriptions_where(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
owned_by: &dyn Fn(&catalog::SubscriptionDefinition) -> bool,
) -> Result<()> {
let Some(sender) = ctx.broker() else {
return Ok(());
};
let lvs = match txn.all_tb_lives(ns, db, tb, None).await {
Ok(lvs) => lvs,
Err(e) => {
warn!(
target: "surrealdb::core::expr",
table = %tb,
error = %e,
"Could not read the subscriptions on a table being removed; its \
subscribers will not be told that it is gone"
);
return Ok(());
}
};
for lv in lvs.iter().filter(|lv| owned_by(lv)) {
txn.on_commit(SendKill::boxed(
Arc::clone(sender),
RoutedNotification::new(
lv.node,
PublicNotification::new(
lv.id.into(),
None,
PublicAction::Killed,
PublicValue::None,
PublicValue::None,
),
),
))
.await;
}
Ok(())
}
pub(crate) async fn kill_table_subscriptions(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
) -> Result<()> {
kill_table_subscriptions_where(ctx, txn, ns, db, tb, &|_| true).await
}
pub(crate) async fn kill_database_subscriptions(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
) -> Result<()> {
kill_database_subscriptions_where(ctx, txn, ns, db, &|_| true).await
}
pub(crate) async fn kill_database_subscriptions_where(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
owned_by: &dyn Fn(&catalog::SubscriptionDefinition) -> bool,
) -> Result<()> {
for tb in txn.all_tb(ns, db, None).await?.iter() {
kill_table_subscriptions_where(ctx, txn, ns, db, &tb.name, owned_by).await?;
}
Ok(())
}
pub(crate) async fn kill_principal_subscriptions(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
actor: &str,
) -> Result<()> {
for tb in txn.all_tb(ns, db, None).await?.iter() {
let lvs = match txn.all_tb_lives(ns, db, &tb.name, None).await {
Ok(lvs) => lvs,
Err(e) => {
warn!(
target: "surrealdb::core::expr",
table = %tb.name,
error = %e,
"Could not read the subscriptions on a table while revoking a principal; any it owns there keep running"
);
continue;
}
};
let mut removed_any = false;
for lv in lvs.iter().filter(|lv| lv.auth.as_ref().is_some_and(|a| a.id() == actor)) {
txn.clr_key(&SubscriptionKey {
ns,
db,
tb: std::borrow::Cow::Borrowed(&tb.name),
lq: lv.id,
})
.await?;
txn.clr_key(&NodeLiveQueryKey {
nd: lv.node,
lq: lv.id,
})
.await?;
removed_any = true;
if let Some(sender) = ctx.broker() {
txn.on_commit(SendKill::boxed(
Arc::clone(sender),
RoutedNotification::new(
lv.node,
PublicNotification::new(
lv.id.into(),
None,
PublicAction::Killed,
PublicValue::None,
PublicValue::None,
),
),
))
.await;
}
}
if removed_any {
txn.bump_table_lives_cache(ns, db, &tb.name).await?;
}
}
Ok(())
}
pub(crate) async fn kill_namespace_subscriptions(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
) -> Result<()> {
for db in txn.all_db(ns, None).await?.iter() {
kill_database_subscriptions(ctx, txn, ns, db.database_id).await?;
}
Ok(())
}
pub(crate) async fn kill_namespace_principal_subscriptions(
ctx: &FrozenContext,
txn: &Transaction,
ns: NamespaceId,
actor: &str,
) -> Result<()> {
for db in txn.all_db(ns, None).await?.iter() {
kill_principal_subscriptions(ctx, txn, ns, db.database_id, actor).await?;
}
Ok(())
}
pub(crate) async fn kill_root_principal_subscriptions(
ctx: &FrozenContext,
txn: &Transaction,
actor: &str,
) -> Result<()> {
for ns in txn.all_ns(None).await?.iter() {
kill_namespace_principal_subscriptions(ctx, txn, ns.namespace_id, actor).await?;
}
Ok(())
}