use crate::queue::cache::CachedEntityInfo;
use crate::queue::key::KeyValue;
use crate::queue::{enqueue_refresh, enqueue_refresh_patched};
use crate::utils::{IntExtraction, tuple_get_i64};
use pgrx::PgTupleDesc;
use pgrx::prelude::*;
use pgrx::spi;
enum KeyExtraction {
Value(KeyValue),
Null,
Missing,
}
fn attnum_of(tupdesc: &PgTupleDesc<'_>, name: &str, hint: Option<i16>) -> Option<usize> {
let named = |i: usize| {
tupdesc
.get(i)
.is_some_and(|att| !att.attisdropped && att.name() == name)
};
hint.and_then(|n| usize::try_from(n).ok())
.and_then(|n| n.checked_sub(1))
.filter(|&i| named(i))
.or_else(|| (0..tupdesc.len()).find(|&i| named(i)))
.map(|i| i + 1)
}
unsafe fn tuple_key(
tuple: *mut pg_sys::HeapTupleData,
tupdesc: &PgTupleDesc<'_>,
name: &str,
hint: Option<i16>,
) -> KeyExtraction {
let Some(attnum) = attnum_of(tupdesc, name, hint) else {
return KeyExtraction::Missing;
};
let Some(att) = tupdesc.get(attnum - 1) else {
return KeyExtraction::Missing;
};
let typid = unsafe { pg_sys::getBaseType(att.atttypid) };
unsafe {
let mut isnull = false;
let datum = pg_sys::heap_getattr(
tuple,
i32::try_from(attnum).unwrap_or(i32::MAX),
tupdesc.as_ptr(),
&mut isnull,
);
if isnull {
return KeyExtraction::Null;
}
let int = match typid {
pg_sys::INT2OID => i16::from_datum(datum, false).map(i64::from),
pg_sys::INT4OID => i32::from_datum(datum, false).map(i64::from),
pg_sys::INT8OID => i64::from_datum(datum, false),
_ => {
let mut output = pg_sys::InvalidOid;
let mut varlena = false;
pg_sys::getTypeOutputInfo(typid, &raw mut output, &raw mut varlena);
let text = pg_sys::OidOutputFunctionCall(output, datum);
let value = std::ffi::CStr::from_ptr(text)
.to_string_lossy()
.into_owned();
pg_sys::pfree(text.cast());
return KeyExtraction::Value(KeyValue::Text(value));
}
};
int.map_or(KeyExtraction::Null, |v| {
KeyExtraction::Value(KeyValue::Int(v))
})
}
}
fn row_images(trigger: &PgTrigger<'_>) -> Vec<*mut pg_sys::HeapTupleData> {
let td = trigger.trigger_data();
[td.tg_trigtuple, td.tg_newtuple]
.into_iter()
.filter(|t| !t.is_null())
.collect()
}
fn trigger_tupdesc<'a>(trigger: &'a PgTrigger<'a>) -> Option<PgTupleDesc<'a>> {
let rel = trigger.trigger_data().tg_relation;
unsafe {
(!rel.is_null() && !(*rel).rd_att.is_null())
.then(|| PgTupleDesc::from_pg_unchecked((*rel).rd_att))
}
}
#[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
}
};
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
);
vec![]
}
};
let own = match crate::queue::cache::table_cache::entity_info_cached(table_oid) {
Ok(info) => info.filter(|i| serves(&i.name)),
Err(e) => {
warning!(
"Failed to resolve entity for table OID {:?}: {:?}",
table_oid,
e
);
None
}
};
if crate::config::suspend_triggers() || crate::suspend::is_suspended() {
if let Some(entity) = &served {
crate::suspend::record_change(entity);
return Ok(None);
}
if let Some(info) = &own {
crate::suspend::record_change(&info.name);
}
for path in &paths {
crate::suspend::record_change(&path.entity_name);
}
return Ok(None);
}
if let Some((info, legacy)) = own.as_ref().and_then(|i| i.legacy_root.map(|l| (i, l))) {
enqueue_legacy_root(trigger, info, legacy);
}
enqueue_cascade_parents(trigger, own.as_ref(), &paths);
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_legacy_root(
trigger: &PgTrigger,
info: &CachedEntityInfo,
legacy: crate::queue::cache::LegacyRoot,
) {
let entity = &info.name;
let reregister = format!(
"tv_{entity} was registered by an older release: SELECT * FROM \
tviews.pg_tviews_reregister_all() re-registers it"
);
if legacy == crate::queue::cache::LegacyRoot::DistinctOn {
crate::utils::log_once(
&format!("legacy_distinct_on:{entity}"),
&format!("{reregister}; until then it is refreshed in full on writes"),
);
crate::queue::enqueue_refresh_all(entity);
return;
}
let pk_col = format!("pk_{entity}");
let Some(tupdesc) = trigger_tupdesc(trigger) else {
return;
};
if let Some(fields) =
try_capture_direct_patch(trigger, info, &pk_col, changed_columns(trigger).as_deref())
&& let Some(new) = new_image(trigger)
&& let KeyExtraction::Value(KeyValue::Int(pk)) =
unsafe { tuple_key(new, &tupdesc, &pk_col, None) }
{
enqueue_refresh_patched(entity, pk, fields);
return;
}
for image in row_images(trigger) {
match unsafe { tuple_key(image, &tupdesc, &pk_col, None) } {
KeyExtraction::Value(key) => enqueue_refresh(entity, key),
KeyExtraction::Null => {}
KeyExtraction::Missing => warning!("{pk_col} not found on the changed row"),
}
}
}
fn new_image(trigger: &PgTrigger<'_>) -> Option<*mut pg_sys::HeapTupleData> {
let new = trigger.trigger_data().tg_newtuple;
(!new.is_null()).then_some(new)
}
fn enqueue_cascade_parents(
trigger: &PgTrigger,
own: Option<&CachedEntityInfo>,
paths: &[crate::cascade_path::CascadePath],
) {
if paths.is_empty() {
return;
}
let Some(tupdesc) = trigger_tupdesc(trigger) else {
warning!("No relation in trigger context");
return;
};
let images = row_images(trigger);
if images.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.root
&& !path.source_columns.is_empty()
&& !path.source_columns.iter().any(|c| changed.contains(c))
{
continue;
}
if path.root
&& let Some(info) = own.filter(|i| i.name == path.entity_name)
&& let Some(fields) =
try_capture_direct_patch(trigger, info, &path.initial_col, changed.as_deref())
&& let Some(new) = new_image(trigger)
&& let KeyExtraction::Value(KeyValue::Int(pk)) =
unsafe { tuple_key(new, &tupdesc, &path.initial_col, path.initial_attnum) }
{
enqueue_refresh_patched(&path.entity_name, pk, fields);
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 &image in &images {
follow_cascade_path(path, image, &tupdesc);
}
}
}
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,
image: *mut pg_sys::HeapTupleData,
tupdesc: &PgTupleDesc<'_>,
) {
if path.unresolvable {
return;
}
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;
}
match unsafe { tuple_key(image, tupdesc, &path.initial_col, path.initial_attnum) } {
KeyExtraction::Value(key) => enqueue_refresh(&path.entity_name, key),
KeyExtraction::Null => {} KeyExtraction::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
),
);
}
}
}
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,
key_col: &str,
changed: Option<&[String]>,
) -> 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
|| entity_info.is_union
{
return None;
}
if !crate::lifecycle::check_jsonb_delta_available() {
return None;
}
let changed = changed?;
if changed.is_empty() {
return None;
}
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 == key_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)
}