Skip to main content

uqa_sql/catalog/events/definition/lookup/
selections.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Rule and trigger selection over pinned query catalogs and live execution registries.
8
9use std::collections::BTreeMap;
10
11use crate::ast::{RuleEvent, TriggerEvent, TriggerTiming};
12use crate::SQLError;
13
14use super::EventLookupContext;
15
16use crate::catalog::events::{StoredRule, StoredTrigger};
17
18impl EventLookupContext<'_> {
19    pub fn rule_definitions_for(
20        &self,
21        table: &str,
22        event: RuleEvent,
23    ) -> Result<Vec<StoredRule>, SQLError> {
24        let relation = self.analysis.resolve_rule_relation(table)?;
25        if let Some(snapshot) = self.state.query_rules() {
26            return Ok(snapshot
27                .get(&relation)
28                .into_iter()
29                .flat_map(BTreeMap::values)
30                .filter(|rule| rule.definition.event == event)
31                .cloned()
32                .collect());
33        }
34        Ok(self
35            .registry
36            .read_rules()
37            .get(&relation)
38            .into_iter()
39            .flat_map(BTreeMap::values)
40            .filter(|rule| rule.definition.event == event)
41            .cloned()
42            .collect())
43    }
44
45    pub fn rules_for(&self, table: &str, event: RuleEvent) -> Result<Vec<StoredRule>, SQLError> {
46        let relation = self.analysis.resolve_rule_relation(table)?;
47        let replica = self.state.session_replication_role_is_replica();
48        Ok(self
49            .registry
50            .read_rules()
51            .get(&relation)
52            .into_iter()
53            .flat_map(BTreeMap::values)
54            .filter(|rule| {
55                (if replica {
56                    rule.enabled.fires_in_replica()
57                } else {
58                    rule.enabled.fires_in_origin()
59                }) && rule.definition.event == event
60            })
61            .cloned()
62            .collect())
63    }
64
65    pub fn relation_has_rules(&self, table: &str) -> Result<bool, SQLError> {
66        let relation = self.analysis.resolve_rule_relation(table)?;
67        Ok(self
68            .registry
69            .read_rules()
70            .get(&relation)
71            .is_some_and(|entries| !entries.is_empty()))
72    }
73
74    pub fn triggers_for(
75        &self,
76        table: &str,
77        timing: TriggerTiming,
78        event: TriggerEvent,
79        row: bool,
80        updated_columns: &[String],
81    ) -> Result<Vec<StoredTrigger>, SQLError> {
82        let relation = self.analysis.resolve_trigger_table(table)?;
83        let replica = self.state.session_replication_role_is_replica();
84        Ok(self
85            .matching_triggers(&relation, timing, event, row, updated_columns)?
86            .into_iter()
87            .filter(|trigger| {
88                if replica {
89                    trigger.enabled.fires_in_replica()
90                } else {
91                    trigger.enabled.fires_in_origin()
92                }
93            })
94            .collect())
95    }
96
97    /// The triggers of `table` for a timing, a level and an event, whatever replication role fires them. `PostgreSQL` copies a relation's trigger descriptor when a statement begins to write the relation and reads the replication role at each firing, so a statement takes its triggers from here once and applies the role itself.
98    pub fn trigger_definitions_for(
99        &self,
100        table: &str,
101        timing: TriggerTiming,
102        event: TriggerEvent,
103        row: bool,
104        updated_columns: &[String],
105    ) -> Result<Vec<StoredTrigger>, SQLError> {
106        let relation = self.analysis.resolve_trigger_table(table)?;
107        self.matching_triggers(&relation, timing, event, row, updated_columns)
108    }
109
110    fn matching_triggers(
111        &self,
112        relation: &uqa_core::RelationIdentity,
113        timing: TriggerTiming,
114        event: TriggerEvent,
115        row: bool,
116        updated_columns: &[String],
117    ) -> Result<Vec<StoredTrigger>, SQLError> {
118        let relations = if row {
119            self.partition_trigger_sources(&relation.qualified_name())?
120        } else {
121            vec![relation.clone()]
122        };
123        let triggers = self.registry.read_triggers();
124        let mut candidates = BTreeMap::new();
125        for source in relations {
126            for trigger in triggers.get(&source).into_iter().flat_map(BTreeMap::values) {
127                let mut trigger = trigger.clone();
128                if source != *relation {
129                    trigger.definition.table = relation.qualified_name();
130                }
131                candidates
132                    .entry(trigger.definition.name.clone())
133                    .or_insert(trigger);
134            }
135        }
136        Ok(candidates
137            .into_values()
138            .filter(|trigger| {
139                trigger.definition.timing == timing
140                    && trigger.definition.row == row
141                    && trigger.definition.events.contains(&event)
142                    && (event != TriggerEvent::Update
143                        || trigger.definition.update_columns.is_empty()
144                        || trigger
145                            .definition
146                            .update_columns
147                            .iter()
148                            .any(|column| updated_columns.contains(column)))
149            })
150            .collect())
151    }
152
153    pub fn has_trigger_definition(
154        &self,
155        table: &str,
156        timing: TriggerTiming,
157        event: TriggerEvent,
158        row: bool,
159    ) -> Result<bool, SQLError> {
160        let relation = self.analysis.resolve_trigger_table(table)?;
161        let matches = |entries: &BTreeMap<String, StoredTrigger>| {
162            entries.values().any(|trigger| {
163                trigger.definition.timing == timing
164                    && trigger.definition.row == row
165                    && trigger.definition.events.contains(&event)
166            })
167        };
168        if let Some(snapshot) = self.state.query_triggers() {
169            return Ok(snapshot.get(&relation).is_some_and(matches));
170        }
171        Ok(self
172            .registry
173            .read_triggers()
174            .get(&relation)
175            .is_some_and(matches))
176    }
177
178    pub fn has_row_triggers(&self, table: &str, event: TriggerEvent) -> Result<bool, SQLError> {
179        let relation = self.analysis.resolve_trigger_table(table)?;
180        let sources = self.partition_trigger_sources(&relation.qualified_name())?;
181        let replica = self.state.session_replication_role_is_replica();
182        let triggers = self.registry.read_triggers();
183        Ok(sources.iter().any(|source| {
184            triggers.get(source).is_some_and(|entries| {
185                entries.values().any(|trigger| {
186                    (if replica {
187                        trigger.enabled.fires_in_replica()
188                    } else {
189                        trigger.enabled.fires_in_origin()
190                    }) && trigger.definition.row
191                        && trigger.definition.events.contains(&event)
192                })
193            })
194        }))
195    }
196
197    pub fn list_triggers(&self) -> Vec<StoredTrigger> {
198        if let Some(snapshot) = self.state.query_triggers() {
199            return snapshot
200                .values()
201                .flat_map(BTreeMap::values)
202                .cloned()
203                .collect();
204        }
205        self.registry
206            .read_triggers()
207            .values()
208            .flat_map(BTreeMap::values)
209            .cloned()
210            .collect()
211    }
212}