use crate::mutation::candidate::PhysicalDocumentIdentity;
use crate::mutation::publication::MutationPublicationBatch;
use crate::mutation::statement::context::MutationStatementContext;
use crate::mutation::triggers::context::TriggerContext;
use crate::mutation::triggers::queue::StatementEvent;
use crate::mutation::triggers::AfterRowTriggerEvent;
use crate::query::scope::StatementCommands;
use crate::query::CteScope;
use std::cell::RefCell;
use std::sync::Arc;
use uqa_sql::{plan::CtePlan, SQLError, SQLParam};
thread_local! {
static RUNNING_STATEMENTS: RefCell<Vec<Arc<StatementCommands>>> = const { RefCell::new(Vec::new()) };
}
pub struct RunningStatement(());
impl Drop for RunningStatement {
fn drop(&mut self) {
RUNNING_STATEMENTS.with(|statements| {
statements.borrow_mut().pop();
});
}
}
pub fn enter_statement(statement: Arc<StatementCommands>) -> RunningStatement {
RUNNING_STATEMENTS.with(|statements| statements.borrow_mut().push(statement));
RunningStatement(())
}
pub fn statement_commands<S: Clone>(
inherited: Option<&CteScope<S>>,
) -> (Arc<StatementCommands>, Option<RunningStatement>) {
if let Some(statement) = inherited.and_then(CteScope::statement_commands) {
return (Arc::clone(statement), None);
}
let statement = Arc::<StatementCommands>::default();
let running = enter_statement(Arc::clone(&statement));
(statement, Some(running))
}
fn started_by_another_statement() -> bool {
RUNNING_STATEMENTS.with(|statements| statements.borrow().len() > 1)
}
fn note_triggered_rows(rows: &[PhysicalDocumentIdentity]) {
RUNNING_STATEMENTS.with(|statements| {
let statements = statements.borrow();
if let Some((_, starters)) = statements.split_last() {
for starter in starters {
starter.note_triggered(rows.iter().cloned());
}
}
});
}
pub fn fire_before_statements(
statement: &StatementCommands,
context: &TriggerContext<'_>,
statements: &[StatementEvent],
) -> Result<(), SQLError> {
statements.iter().try_for_each(|event| {
statement
.after_triggers()
.fire_before_statement(context, event)
})
}
pub fn end_command(
statement: &StatementCommands,
context: &TriggerContext<'_>,
statements: &[StatementEvent],
rows: Vec<AfterRowTriggerEvent>,
) -> Result<(), SQLError> {
statement
.after_triggers()
.queue_command(context, statements, rows)?;
if statement.modifies_with() {
return Ok(());
}
statement.after_triggers().fire(context)
}
pub fn publication_batch(statement: &StatementCommands) -> MutationPublicationBatch {
MutationPublicationBatch::recording_writes(
statement.modifies_with() || started_by_another_statement(),
)
}
pub fn note_written_rows(
statement: &StatementCommands,
publication: &mut MutationPublicationBatch,
) {
let rows = publication.take_written();
if rows.is_empty() {
return;
}
note_triggered_rows(&rows);
if statement.modifies_with() {
statement.note_written(rows);
}
}
pub fn note_written_row(table: &str, doc_id: uqa_core::DocId) {
if started_by_another_statement() {
note_triggered_rows(&[PhysicalDocumentIdentity {
table: table.to_string(),
doc_id,
}]);
}
}
pub fn action_publication_batch() -> MutationPublicationBatch {
MutationPublicationBatch::recording_writes(started_by_another_statement())
}
pub fn note_action_rows(publication: &mut MutationPublicationBatch) {
note_triggered_rows(&publication.take_written());
}
pub fn finish_statement<S: Clone + Send + Sync + 'static>(
context: &MutationStatementContext<'_, S>,
params: &[SQLParam],
ctes: &[CtePlan],
scope: Option<&mut CteScope<S>>,
) -> Result<(), SQLError> {
let Some(scope) = scope.filter(|_| ctes.iter().any(|cte| cte.body.modifies_data())) else {
return Ok(());
};
crate::query::cte::finish_statement_ctes(context.query.source.ctes, params, scope)?;
let Some(commands) = scope.statement_commands() else {
return Ok(());
};
commands
.after_triggers()
.fire(&context.mutation.preparation.referential.triggers)
}