use crate::queue::cache::CachedEntityInfo;
use crate::queue::{enqueue_refresh, enqueue_refresh_dedup, enqueue_refresh_patched};
use crate::utils::{IntExtraction, 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> {
crate::revision::check();
let table_oid = match trigger.relation() {
Ok(rel) => rel.oid(),
Err(e) => {
warning!("Failed to get trigger relation: {:?}", e);
return Ok(None);
}
};
let served = crate::delta::trigger_entity(trigger);
let serves = |entity: &str| served.as_deref().is_none_or(|e| e == entity);
let table_oid = match crate::delta::partition_root(table_oid) {
Ok(root) => root,
Err(e) => {
warning!(
"Failed to resolve the partition root of {:?}: {}",
table_oid,
e
);
table_oid
}
};
if crate::config::suspend_triggers() || crate::suspend::is_suspended() {
if let Some(entity) = &served {
crate::suspend::record_change(entity);
return Ok(None);
}
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)) if serves(&entity_info.name) => {
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, entity) {
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);
}
}
}
Ok(_) => { }
Err(e) => {
warning!(
"Failed to resolve entity for table OID {:?}: {:?}",
table_oid,
e
);
return Ok(None);
}
}
enqueue_cascade_parents(trigger, table_oid, &serves);
if let Some(entity) = &served
&& let Err(e) = crate::delta::map_row(trigger, entity, table_oid)
{
error!("pg_tviews: could not map the changed row to tv_{entity} keys: {e}");
}
Ok(None)
}
fn enqueue_cascade_parents(
trigger: &PgTrigger,
table_oid: pg_sys::Oid,
serves: &dyn Fn(&str) -> bool,
) {
let paths: Vec<crate::cascade_path::CascadePath> =
match crate::queue::cache::cascade_cache::cascade_paths_for_table(table_oid) {
Ok(p) => p.into_iter().filter(|p| serves(&p.entity_name)).collect(),
Err(e) => {
warning!(
"Failed to load cascade paths for table {:?}: {:?}",
table_oid,
e
);
return;
}
};
if paths.is_empty() {
return;
}
let tuples: Vec<_> = [trigger.old(), trigger.new()]
.into_iter()
.flatten()
.collect();
if tuples.is_empty() {
warning!("No tuple available in trigger context");
return;
}
let changed = changed_columns(trigger);
for path in &paths {
if let Some(changed) = &changed
&& !path.source_columns.is_empty()
&& !path.source_columns.iter().any(|c| changed.contains(c))
{
continue;
}
if let Some(changed) = &changed
&& let Some((fanout, key, fields)) = try_capture_fanout(trigger, path, changed)
{
crate::queue::patch::record_fanout(
(path.entity_name.clone(), fanout.lookup_col.clone(), key),
fields,
);
continue;
}
for tuple in &tuples {
if let Err(e) = follow_cascade_path(path, tuple) {
warning!(
"Cascade refresh failed for path {} → {}: {:?}",
path.source_table,
path.entity_name,
e
);
}
}
}
}
fn try_capture_fanout<'p>(
trigger: &PgTrigger,
path: &'p crate::cascade_path::CascadePath,
changed: &[String],
) -> Option<(
&'p crate::cascade_path::FanoutPatch,
i64,
serde_json::Map<String, serde_json::Value>,
)> {
let fanout = path.fanout.as_ref()?;
if !crate::config::direct_patch_enabled()
|| !crate::lifecycle::check_jsonb_delta_available()
|| changed.contains(&path.initial_col)
{
return None;
}
let new_tuple = trigger.new()?;
let IntExtraction::Value(key) = tuple_get_i64(&new_tuple, &path.initial_col) else {
return None;
};
let mut fields = serde_json::Map::new();
for col in changed.iter().filter(|c| path.source_columns.contains(c)) {
let (_, data_key) = fanout.fields.iter().find(|(c, _)| c == col)?;
fields.insert(data_key.clone(), capture_value(&new_tuple, col)?);
}
(!fields.is_empty()).then_some((fanout, key, fields))
}
fn follow_cascade_path(
path: &crate::cascade_path::CascadePath,
tuple: &PgHeapTuple<AllocatedByPostgres>,
) -> crate::TViewResult<()> {
if path.unresolvable {
return Ok(());
}
if !path.hops.is_empty() {
crate::utils::log_once(
&format!("legacy_hops:{}", path.entity_name),
&format!(
"tv_{0} was registered by an older release: writes to {1} refresh it in full \
until SELECT * FROM tviews.pg_tviews_reregister_all() re-registers it",
path.entity_name, path.source_table
),
);
crate::queue::enqueue_refresh_all(&path.entity_name);
return Ok(());
}
match tuple_get_i64(tuple, &path.initial_col) {
IntExtraction::Value(pk) => enqueue_refresh(&path.entity_name, pk),
IntExtraction::Null => {} IntExtraction::Missing => {
crate::utils::log_once(
&format!("initial_col:{}:{}", path.source_table, path.initial_col),
&format!(
"column '{}' of {} is gone: its writes no longer cascade to tv_{}; \
re-register it with pg_tviews_reregister('{}')",
path.initial_col, path.source_table, path.entity_name, path.entity_name
),
);
}
}
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() || crate::suspend::is_suspended() {
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)
}