use crate::catalog::entity_for_table;
use crate::queue::cache::CachedEntityInfo;
use crate::queue::{
enqueue_refresh, enqueue_refresh_bulk, enqueue_refresh_dedup, enqueue_refresh_patched,
};
use crate::utils::{IntExtraction, quote_identifier, tuple_get_i64};
use pgrx::PgTupleDesc;
use pgrx::prelude::*;
use pgrx::spi;
enum KeyExtraction {
Value(String),
Null,
TypeMismatch,
}
fn extract_distinct_on_key(
tuple: &PgHeapTuple<'_, AllocatedByPostgres>,
key_col: &str,
) -> KeyExtraction {
match tuple.get_by_name::<String>(key_col) {
Ok(Some(val)) => return KeyExtraction::Value(val),
Ok(None) => return KeyExtraction::Null,
Err(_) => {} }
match tuple.get_by_name::<pgrx::Uuid>(key_col) {
Ok(Some(val)) => return KeyExtraction::Value(val.to_string()),
Ok(None) => return KeyExtraction::Null,
Err(_) => {}
}
match tuple.get_by_name::<i64>(key_col) {
Ok(Some(val)) => return KeyExtraction::Value(val.to_string()),
Ok(None) => return KeyExtraction::Null,
Err(_) => {}
}
match tuple.get_by_name::<i32>(key_col) {
Ok(Some(val)) => return KeyExtraction::Value(val.to_string()),
Ok(None) => return KeyExtraction::Null,
Err(_) => {}
}
KeyExtraction::TypeMismatch
}
#[pg_trigger]
#[allow(clippy::unnecessary_wraps)] fn pg_tview_trigger_handler<'a>(
trigger: &'a PgTrigger<'a>,
) -> Result<Option<PgHeapTuple<'a, AllocatedByPostgres>>, spi::Error> {
let table_oid = match trigger.relation() {
Ok(rel) => rel.oid(),
Err(e) => {
warning!("Failed to get trigger relation: {:?}", e);
return Ok(None);
}
};
if crate::config::suspend_triggers() {
if let Ok(Some(entity_info)) =
crate::queue::cache::table_cache::entity_info_cached(table_oid)
{
crate::suspend::record_change(&entity_info.name);
}
let paths: Vec<crate::cascade_path::CascadePath> =
match crate::queue::cache::cascade_cache::cascade_paths_for_table(table_oid) {
Ok(p) => p,
Err(e) => {
warning!(
"Failed to load cascade paths for suspended trigger: {:?}",
e
);
vec![]
}
};
for path in paths {
crate::suspend::record_change(&path.entity_name);
}
return Ok(None);
}
match crate::queue::cache::table_cache::entity_info_cached(table_oid) {
Ok(Some(entity_info)) => {
let entity = &entity_info.name;
if let Some(key_col) = &entity_info.distinct_on_key {
let tuple = if let Some(t) = trigger.new().or_else(|| trigger.old()) {
t
} else {
warning!("No tuple in trigger context for DISTINCT ON TVIEW '{entity}'");
return Ok(None);
};
match extract_distinct_on_key(&tuple, key_col) {
KeyExtraction::Value(key_val) => {
enqueue_refresh_dedup(entity, &key_val);
}
KeyExtraction::Null => {
warning!(
"DISTINCT ON key '{key_col}' is NULL for entity '{entity}' — skipping refresh"
);
}
KeyExtraction::TypeMismatch => {
warning!(
"Cannot extract DISTINCT ON key '{key_col}' for '{entity}': \
unsupported column type — skipping refresh for this row"
);
}
}
} else {
let pk_value = match crate::utils::extract_pk(trigger) {
Ok(pk) => pk,
Err(e) => {
warning!("Failed to extract primary key from trigger: {:?}", e);
return Ok(None);
}
};
if let Some(fields) = try_capture_direct_patch(trigger, &entity_info) {
enqueue_refresh_patched(entity, pk_value, fields);
} else {
enqueue_refresh(entity, pk_value);
}
}
return Ok(None);
}
Ok(None) => { }
Err(e) => {
warning!(
"Failed to resolve entity for table OID {:?}: {:?}",
table_oid,
e
);
return Ok(None);
}
}
enqueue_cascade_parents(trigger, table_oid);
Ok(None)
}
fn enqueue_cascade_parents(trigger: &PgTrigger, table_oid: pg_sys::Oid) {
let paths: Vec<crate::cascade_path::CascadePath> =
match crate::queue::cache::cascade_cache::cascade_paths_for_table(table_oid) {
Ok(p) => p,
Err(e) => {
warning!(
"Failed to load cascade paths for table {:?}: {:?}",
table_oid,
e
);
return;
}
};
if paths.is_empty() {
return;
}
let Some(tuple) = trigger.new().or_else(|| trigger.old()) else {
warning!("No tuple available in trigger context");
return;
};
for path in &paths {
if let Err(e) = follow_cascade_path(path, &tuple) {
warning!(
"Cascade refresh failed for path {} → {}: {:?}",
path.source_table,
path.entity_name,
e
);
}
}
}
fn follow_cascade_path(
path: &crate::cascade_path::CascadePath,
tuple: &PgHeapTuple<AllocatedByPostgres>,
) -> crate::TViewResult<()> {
if path.unresolvable {
warning!(
"Unresolvable cascade path for entity '{}' — full refresh needed",
path.entity_name
);
return Ok(());
}
let mut current_ids = match tuple_get_i64(tuple, &path.initial_col) {
IntExtraction::Value(pk) => vec![pk],
IntExtraction::Null => return Ok(()), IntExtraction::Missing => {
warning!(
"Initial column '{}' not found on tuple for cascade to '{}'",
path.initial_col,
path.entity_name
);
return Ok(());
}
};
for hop in &path.hops {
if current_ids.is_empty() {
return Ok(());
}
current_ids = crate::queue::spi_batch_lookup(
hop.table_oid,
&hop.lookup_col,
&hop.carry_col,
¤t_ids,
)?;
}
for pk in current_ids {
enqueue_refresh(&path.entity_name, pk);
}
Ok(())
}
unsafe extern "C" {
#[link_name = "datumIsEqual"]
fn datum_is_equal(
value1: pg_sys::Datum,
value2: pg_sys::Datum,
typ_by_val: bool,
typ_len: i32,
) -> bool;
}
fn try_capture_direct_patch(
trigger: &PgTrigger,
entity_info: &CachedEntityInfo,
) -> Option<serde_json::Map<String, serde_json::Value>> {
use std::collections::HashSet;
if !crate::config::direct_patch_enabled()
|| entity_info.direct_map.is_empty()
|| entity_info.distinct_on_key.is_some()
|| entity_info.is_union
{
return None;
}
if !crate::lifecycle::check_jsonb_delta_available() {
return None;
}
let changed = changed_columns(trigger)?;
if changed.is_empty() {
return None;
}
let pk_col = format!("pk_{}", entity_info.name);
let fk_set: HashSet<&str> = entity_info.fk_columns.iter().map(String::as_str).collect();
let uuid_fk_set: HashSet<&str> = entity_info
.uuid_fk_columns
.iter()
.map(String::as_str)
.collect();
let output_set: HashSet<&str> = entity_info
.output_columns
.iter()
.map(String::as_str)
.collect();
for col in &changed {
if col == &pk_col
|| fk_set.contains(col.as_str())
|| uuid_fk_set.contains(col.as_str())
|| output_set.contains(col.as_str())
|| !entity_info.direct_map.contains_key(col.as_str())
{
return None;
}
}
let new_tuple = trigger.new()?;
let mut fields = serde_json::Map::with_capacity(changed.len());
for col in &changed {
let key = entity_info.direct_map.get(col.as_str())?;
let value = capture_value(&new_tuple, col)?;
fields.insert(key.clone(), value);
}
Some(fields)
}
fn changed_columns(trigger: &PgTrigger) -> Option<Vec<String>> {
unsafe {
let td: &pg_sys::TriggerData = trigger.trigger_data();
let old_tuple = td.tg_trigtuple;
let new_tuple = td.tg_newtuple;
if old_tuple.is_null() || new_tuple.is_null() {
return None; }
let rel = td.tg_relation;
if rel.is_null() {
return None;
}
let tupdesc_ptr = (*rel).rd_att;
if tupdesc_ptr.is_null() {
return None;
}
let tupdesc = PgTupleDesc::from_pg_unchecked(tupdesc_ptr);
let mut changed = Vec::new();
for i in 0..tupdesc.len() {
let Some(att) = tupdesc.get(i) else { continue };
if att.attisdropped {
continue;
}
let attnum = i32::try_from(i + 1).ok()?;
let mut old_isnull = false;
let mut new_isnull = false;
let old_datum = pg_sys::heap_getattr(old_tuple, attnum, tupdesc_ptr, &mut old_isnull);
let new_datum = pg_sys::heap_getattr(new_tuple, attnum, tupdesc_ptr, &mut new_isnull);
let differs = if old_isnull != new_isnull {
true
} else if old_isnull {
false } else {
!datum_is_equal(old_datum, new_datum, att.attbyval, i32::from(att.attlen))
};
if differs {
changed.push(att.name().to_string());
}
}
Some(changed)
}
}
fn capture_value(
tuple: &PgHeapTuple<'_, AllocatedByPostgres>,
col: &str,
) -> Option<serde_json::Value> {
use serde_json::Value;
match tuple.get_by_name::<String>(col) {
Ok(Some(v)) => return Some(Value::String(v)),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
match tuple.get_by_name::<bool>(col) {
Ok(Some(v)) => return Some(Value::Bool(v)),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
match tuple.get_by_name::<i64>(col) {
Ok(Some(v)) => return Some(Value::Number(v.into())),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
match tuple.get_by_name::<i32>(col) {
Ok(Some(v)) => return Some(Value::Number(v.into())),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
match tuple.get_by_name::<i16>(col) {
Ok(Some(v)) => return Some(Value::Number(v.into())),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
match tuple.get_by_name::<pgrx::Uuid>(col) {
Ok(Some(v)) => return Some(Value::String(v.to_string())),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
match tuple.get_by_name::<pgrx::JsonB>(col) {
Ok(Some(v)) => return Some(v.0),
Ok(None) => return Some(Value::Null),
Err(_) => {}
}
None
}
#[pg_trigger]
#[allow(clippy::unnecessary_wraps)] fn pg_tview_flush_trigger<'a>(
_trigger: &'a PgTrigger<'a>,
) -> Result<Option<PgHeapTuple<'a, AllocatedByPostgres>>, spi::Error> {
if crate::config::suspend_triggers() {
return Ok(None);
}
if let Err(e) = crate::queue::flush_refresh_queue() {
warning!("TVIEW refresh failed in statement trigger: {:?}", e);
}
if let Err(e) = crate::audit::flush_audit_buffer() {
warning!("Audit flush failed in statement trigger: {:?}", e);
}
Ok(None)
}
#[pg_trigger]
#[allow(clippy::unnecessary_wraps)] fn pg_tview_stmt_trigger_handler<'a>(
trigger: &'a PgTrigger<'a>,
) -> Result<Option<PgHeapTuple<'a, AllocatedByPostgres>>, spi::Error> {
let table_oid = match trigger.relation() {
Ok(rel) => rel.oid(),
Err(e) => {
warning!("Failed to get trigger relation: {:?}", e);
return Ok(None);
}
};
let entity = match entity_for_table(table_oid) {
Ok(Some(e)) => e,
Ok(None) => {
return Ok(None);
}
Err(e) => {
warning!(
"Failed to resolve entity for table OID {:?}: {:?}",
table_oid,
e
);
return Ok(None);
}
};
let changed_pks = match extract_pks_from_transition_table(trigger) {
Ok(pks) => pks,
Err(e) => {
warning!("Failed to extract PKs from transition table: {:?}", e);
return Ok(None);
}
};
if changed_pks.is_empty() {
return Ok(None);
}
enqueue_refresh_bulk(&entity, changed_pks);
Ok(None)
}
fn extract_pks_from_transition_table(trigger: &PgTrigger) -> spi::Result<Vec<i64>> {
let transition_table_name = if trigger.new().is_some() && trigger.old().is_none() {
"new_table" } else if trigger.new().is_none() && trigger.old().is_some() {
"old_table" } else if trigger.new().is_some() && trigger.old().is_some() {
"new_table" } else {
return Ok(Vec::new()); };
let pk_column = get_pk_column_name(
trigger
.relation()
.map_err(|_| crate::TViewError::SpiError {
query: "get relation".to_string(),
error: "Failed to get trigger relation".to_string(),
})?
.oid(),
)?;
let query = format!(
"SELECT DISTINCT {} FROM {}",
quote_identifier(&pk_column),
transition_table_name );
Spi::connect(|client| {
let rows = client.select(&query, None, &[])?;
let mut pks = Vec::new();
for row in rows {
if let Some(pk) = row[&pk_column as &str].value::<i64>()? {
pks.push(pk);
}
}
Ok(pks)
})
}
fn get_pk_column_name(table_oid: pg_sys::Oid) -> spi::Result<String> {
let entity = match entity_for_table(table_oid) {
Ok(Some(e)) => e,
Ok(None) => {
return Err(crate::TViewError::SpiError {
query: "entity_for_table".to_string(),
error: "Table not managed by pg_tviews".to_string(),
}
.into());
}
Err(e) => {
return Err(crate::TViewError::SpiError {
query: "entity_for_table".to_string(),
error: format!("Failed to get entity: {e:?}"),
}
.into());
}
};
Ok(format!("pk_{entity}"))
}