uqa-execution 0.5.0

Volcano physical operators with row-batch pipelines
//
// Unified Query Algebra
//
// Copyright (c) 2023-2026 Cognica, Inc.
//

//! The AFTER trigger events of one statement, queued and fired as `PostgreSQL`'s `trigger.c` queues and fires the events of one query level.
//!
//! Every command of a statement queues into one queue: the primary command, its data-modifying WITH items, and the referential actions that the queued events run as they fire. A relation and an operation that the statement writes keep one state (`AfterTriggersTableData`): the rows its transition tables collect, whether its BEFORE STATEMENT triggers fired, and where its AFTER STATEMENT events wait in the queue. A later command on the same relation and operation therefore fires no second BEFORE STATEMENT trigger and moves the AFTER STATEMENT events to the end of the queue (`cancel_prior_stmt_triggers`), until a trigger reads the transition tables and closes the state, after which the rows start a new one.

use super::context::TriggerContext;
use super::{AfterRowTriggerEvent, TransitionTables};
use std::collections::BTreeSet;
use std::sync::Arc;
use uqa_core::Value;
use uqa_sql::ast::{TriggerEvent, TriggerTiming};
use uqa_sql::error::Result;
use uqa_sql::SQLError;

/// A relation and operation that a command writes, with the columns an UPDATE sets, which its statement triggers fire for.
#[derive(Debug, Clone)]
pub struct StatementEvent {
    pub relation: String,
    pub event: TriggerEvent,
    pub columns: Vec<String>,
}

impl StatementEvent {
    pub fn new(relation: &str, event: TriggerEvent, columns: &[String]) -> Self {
        Self {
            relation: relation.to_string(),
            event,
            columns: columns.to_vec(),
        }
    }
}

/// The queued AFTER events of one statement.
#[derive(Default)]
pub struct AfterTriggerQueue {
    state: parking_lot::Mutex<QueueState>,
}

#[derive(Default)]
struct QueueState {
    entries: Vec<QueueEntry>,
    /// The first entry that has not fired.
    next: usize,
    tables: Vec<TableState>,
}

enum QueueEntry {
    Row {
        event: AfterRowTriggerEvent,
        /// The state whose transition tables collected the row.
        table: Option<usize>,
    },
    Statement {
        statement: StatementEvent,
        table: usize,
    },
    /// An entry that fired or that a later command cancelled.
    Done,
}

/// The state of one relation and operation within the statement.
struct TableState {
    relation: String,
    event: TriggerEvent,
    /// The relation and its partitions, whose rows the transition tables collect.
    sources: BTreeSet<String>,
    capture: TransitionCapture,
    /// A trigger has read the transition tables, which take no further rows.
    closed: bool,
    before_fired: bool,
    /// Where the AFTER STATEMENT events of the state were queued, from which a later command cancels them.
    statements_from: Option<usize>,
    old_rows: Vec<Value>,
    new_rows: Vec<Value>,
    transitions: Option<Arc<TransitionTables>>,
}

/// Which images of its rows a relation's triggers read as transition tables.
#[derive(Clone, Copy, Default)]
struct TransitionCapture {
    old: bool,
    new: bool,
}

impl QueueState {
    /// The open state of `relation` and `event`, created when every earlier one was closed (`GetAfterTriggersTableData`).
    fn open_table(
        &mut self,
        context: &TriggerContext<'_>,
        relation: &str,
        event: TriggerEvent,
    ) -> Result<usize> {
        let relation = context
            .relations
            .try_resolve_table_name(relation)
            .map_err(|error| SQLError::Internal(format!("resolve trigger relation: {error}")))?
            .unwrap_or_else(|| relation.to_string());
        if let Some(index) = self
            .tables
            .iter()
            .position(|table| !table.closed && table.event == event && table.relation == relation)
        {
            return Ok(index);
        }
        let mut capture = TransitionCapture::default();
        for row in [false, true] {
            for trigger in
                context
                    .catalog
                    .triggers_for(&relation, TriggerTiming::After, event, row, &[])?
            {
                capture.old |= trigger.definition.old_transition_table().is_some();
                capture.new |= trigger.definition.new_transition_table().is_some();
            }
        }
        // Only a relation whose triggers read transition tables collects rows, and a view has no partitions to collect them from.
        let sources = if capture.old || capture.new {
            context
                .catalog
                .hierarchy_scan_tables(&relation, true)?
                .into_iter()
                .collect()
        } else {
            BTreeSet::new()
        };
        self.tables.push(TableState {
            relation,
            event,
            sources,
            capture,
            closed: false,
            before_fired: false,
            statements_from: None,
            old_rows: Vec::new(),
            new_rows: Vec::new(),
            transitions: None,
        });
        Ok(self.tables.len() - 1)
    }

    /// Cancel the AFTER STATEMENT events that an earlier command queued for `table` and that have not fired.
    fn cancel_prior_statements(&mut self, table: usize) {
        let Some(from) = self.tables[table].statements_from else {
            return;
        };
        for entry in self.entries.iter_mut().skip(from.max(self.next)) {
            if matches!(entry, QueueEntry::Statement { table: queued, .. } if *queued == table) {
                *entry = QueueEntry::Done;
            }
        }
    }

    /// Take the next entry to fire, leaving it done.
    fn take_next(&mut self) -> Option<QueueEntry> {
        while self.next < self.entries.len() {
            let entry = std::mem::replace(&mut self.entries[self.next], QueueEntry::Done);
            self.next += 1;
            if !matches!(entry, QueueEntry::Done) {
                return Some(entry);
            }
        }
        None
    }
}

impl AfterTriggerQueue {
    /// Fire the BEFORE STATEMENT triggers of `statement` unless the statement already fired them for its open state (`before_stmt_triggers_fired`).
    pub fn fire_before_statement(
        &self,
        context: &TriggerContext<'_>,
        statement: &StatementEvent,
    ) -> Result<()> {
        let fire = {
            let mut state = self.state.lock();
            let table = state.open_table(context, &statement.relation, statement.event)?;
            !std::mem::replace(&mut state.tables[table].before_fired, true)
        };
        if fire {
            super::fire_statement_triggers(
                context,
                &statement.relation,
                TriggerTiming::Before,
                statement.event,
                &statement.columns,
            )?;
        }
        Ok(())
    }

    /// Queue the AFTER events of one command: the events of its rows, whose images join the transition tables of the statement event they belong to, then an AFTER STATEMENT event for each of `statements`, which replaces the one an earlier command queued for the same state.
    pub fn queue_command(
        &self,
        context: &TriggerContext<'_>,
        statements: &[StatementEvent],
        rows: Vec<AfterRowTriggerEvent>,
    ) -> Result<()> {
        let mut state = self.state.lock();
        let tables = statements
            .iter()
            .map(|statement| state.open_table(context, &statement.relation, statement.event))
            .collect::<Result<Vec<_>>>()?;
        for event in rows {
            let table = if event.captured {
                statements
                    .iter()
                    .zip(&tables)
                    .find(|(statement, table)| {
                        statement.event == event.event
                            && state.tables[**table].sources.contains(&event.table)
                    })
                    .map(|(_, table)| *table)
            } else {
                None
            };
            if let Some(table) = table {
                let target = &mut state.tables[table];
                if target.capture.old && !matches!(event.old, Value::Null) {
                    target.old_rows.push(event.old.clone());
                }
                if target.capture.new && !matches!(event.new, Value::Null) {
                    target.new_rows.push(event.new.clone());
                }
            }
            state.entries.push(QueueEntry::Row { event, table });
        }
        for (statement, table) in statements.iter().zip(tables) {
            state.cancel_prior_statements(table);
            let position = state.entries.len();
            state.tables[table].statements_from = Some(position);
            state.entries.push(QueueEntry::Statement {
                statement: statement.clone(),
                table,
            });
        }
        Ok(())
    }

    /// Fire the queued events in order, with the events that the referential actions among them queue as they run (`afterTriggerInvokeEvents`).
    pub fn fire(&self, context: &TriggerContext<'_>) -> Result<()> {
        loop {
            let Some(entry) = self.state.lock().take_next() else {
                return Ok(());
            };
            match entry {
                QueueEntry::Row { event, table } => {
                    super::fire_after_row_trigger(context, &event, table, self)?;
                }
                QueueEntry::Statement { statement, table } => {
                    self.fire_after_statement(context, &statement, table)?;
                }
                QueueEntry::Done => {}
            }
        }
    }

    fn fire_after_statement(
        &self,
        context: &TriggerContext<'_>,
        statement: &StatementEvent,
        table: usize,
    ) -> Result<()> {
        let triggers = context.catalog.triggers_for(
            &statement.relation,
            TriggerTiming::After,
            statement.event,
            false,
            &statement.columns,
        )?;
        if triggers.is_empty() {
            return Ok(());
        }
        let transitions = if triggers
            .iter()
            .any(|trigger| !trigger.definition.transition_relations.is_empty())
        {
            Some(self.transitions(context, table)?)
        } else {
            None
        };
        super::fire_statement_triggers_with_transition(
            context,
            &statement.relation,
            TriggerTiming::After,
            statement.event,
            &statement.columns,
            transitions.as_deref(),
        )
    }

    /// The transition tables of `table`, which reading closes to further rows.
    pub(super) fn transitions(
        &self,
        context: &TriggerContext<'_>,
        table: usize,
    ) -> Result<Arc<TransitionTables>> {
        let mut state = self.state.lock();
        let target = &mut state.tables[table];
        target.closed = true;
        if let Some(transitions) = &target.transitions {
            return Ok(Arc::clone(transitions));
        }
        let old = target
            .capture
            .old
            .then(|| {
                super::transitions::materialize_transition_rows(
                    context,
                    &target.relation,
                    std::mem::take(&mut target.old_rows),
                )
            })
            .transpose()?;
        let new = target
            .capture
            .new
            .then(|| {
                super::transitions::materialize_transition_rows(
                    context,
                    &target.relation,
                    std::mem::take(&mut target.new_rows),
                )
            })
            .transpose()?;
        let transitions = Arc::new(TransitionTables { old, new });
        target.transitions = Some(Arc::clone(&transitions));
        Ok(transitions)
    }
}