use super::ops::{
clear_queue, is_crash_recovery_checked, mark_crash_recovery_checked, take_queue_snapshot,
};
use crate::TViewResult;
use pgrx::pg_sys;
use pgrx::prelude::*;
use std::collections::HashSet;
use std::os::raw::c_void;
use std::panic::AssertUnwindSafe;
thread_local! {
static SAVEPOINT_DEPTH: std::cell::RefCell<usize> = const { std::cell::RefCell::new(0) };
static QUEUE_SNAPSHOTS: std::cell::RefCell<Vec<HashSet<super::key::RefreshKey>>> =
const { std::cell::RefCell::new(Vec::new()) };
static PATCH_SNAPSHOTS: std::cell::RefCell<
Vec<std::collections::HashMap<super::key::RefreshKey, super::patch::PatchState>>,
> = const { std::cell::RefCell::new(Vec::new()) };
static FANOUT_SNAPSHOTS: std::cell::RefCell<Vec<super::patch::FanoutMap>> =
const { std::cell::RefCell::new(Vec::new()) };
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum XactEvent {
Commit,
Abort,
PreCommit,
Prepare, }
pub unsafe fn register_xact_callback() {
unsafe {
pg_sys::RegisterXactCallback(Some(tview_xact_callback), std::ptr::null_mut());
}
}
pub unsafe fn register_subxact_callback() {
unsafe {
pg_sys::RegisterSubXactCallback(Some(tview_subxact_callback), std::ptr::null_mut());
}
let nest_level = unsafe { pg_sys::GetCurrentTransactionNestLevel() };
SAVEPOINT_DEPTH.with(|d| {
*d.borrow_mut() = (nest_level as usize).saturating_sub(1);
});
QUEUE_SNAPSHOTS.with(|s| {
let mut snapshots = s.borrow_mut();
for _ in 0..(nest_level as usize).saturating_sub(1) {
snapshots.push(HashSet::new());
}
});
PATCH_SNAPSHOTS.with(|s| {
let mut snapshots = s.borrow_mut();
for _ in 0..(nest_level as usize).saturating_sub(1) {
snapshots.push(std::collections::HashMap::new());
}
});
FANOUT_SNAPSHOTS.with(|s| {
let mut snapshots = s.borrow_mut();
for _ in 0..(nest_level as usize).saturating_sub(1) {
snapshots.push(std::collections::HashMap::new());
}
});
}
#[unsafe(no_mangle)]
unsafe extern "C-unwind" fn tview_xact_callback(event: u32, _arg: *mut c_void) {
#[allow(non_upper_case_globals)] let xact_event = match event {
pg_sys::XactEvent::XACT_EVENT_COMMIT => XactEvent::Commit,
pg_sys::XactEvent::XACT_EVENT_PRE_COMMIT => XactEvent::PreCommit,
pg_sys::XactEvent::XACT_EVENT_ABORT => XactEvent::Abort,
pg_sys::XactEvent::XACT_EVENT_PREPARE => XactEvent::Prepare,
_ => return, };
match xact_event {
XactEvent::PreCommit | XactEvent::Commit => {
#[allow(clippy::collapsible_if)]
if crate::suspend::is_suspended() {
let stale = crate::suspend::get_changed_entities();
if !stale.is_empty() {
warning!(
"pg_tviews: transaction committed with refresh suspended; TVIEWs {:?} \
are stale until pg_tviews_refresh() is run for each",
stale
);
}
crate::suspend::clear_changed_entities();
}
crate::suspend::force_resume();
warn_unflushed();
clear_transaction_state();
}
XactEvent::Prepare => {
crate::suspend::force_resume();
super::ops::clear_crash_recovery_cache();
clear_transaction_state();
}
XactEvent::Abort => {
crate::suspend::force_resume();
crate::revision::reset();
crate::hooks::release_hook_guard_on_abort(true);
super::ops::clear_crash_recovery_cache();
clear_transaction_state();
}
}
}
fn clear_transaction_state() {
clear_queue();
super::patch::clear_patch_map();
super::patch::clear_fanout_map();
super::cache::cascade_cache::clear_cache();
crate::audit::clear_audit_buffer();
crate::metrics::metrics_api::reset_metrics();
super::affected::clear();
}
fn warn_unflushed() {
let queued = super::state::get_queue_contents();
if queued.is_empty() || !crate::utils::first_time("commit with queued refreshes") {
return;
}
let entities: std::collections::BTreeSet<&str> =
queued.iter().map(|k| k.entity.as_str()).collect();
pg_sys::panic::ErrorReport::new(
PgSqlErrorCode::ERRCODE_WARNING,
format!(
"pg_tviews: transaction committed with {} queued refreshes for {:?}; they were \
not applied (missing flush trigger?)",
queued.len(),
entities
),
function_name!(),
)
.set_hint(
"Run tviews.pg_tviews_health_check() to find missing triggers, and \
tviews.pg_tviews_refresh(entity) to rebuild the TVIEWs named.",
)
.report(PgLogLevel::WARNING);
}
#[unsafe(no_mangle)]
unsafe extern "C-unwind" fn tview_subxact_callback(
event: u32,
_subxid: pg_sys::SubTransactionId,
_parent_subid: pg_sys::SubTransactionId,
_arg: *mut c_void,
) {
let result = std::panic::catch_unwind(AssertUnwindSafe(|| {
match event {
pg_sys::SubXactEvent::SUBXACT_EVENT_START_SUB => {
SAVEPOINT_DEPTH.with(|d| {
let mut depth = d.borrow_mut();
*depth += 1;
});
let snapshot = take_queue_snapshot();
QUEUE_SNAPSHOTS.with(|s| {
s.borrow_mut().push(snapshot);
});
super::affected::savepoint_start();
let patch_snapshot = super::patch::take_patch_snapshot();
PATCH_SNAPSHOTS.with(|s| {
s.borrow_mut().push(patch_snapshot);
});
let fanout_snapshot = super::patch::take_fanout_snapshot();
FANOUT_SNAPSHOTS.with(|s| {
s.borrow_mut().push(fanout_snapshot);
});
}
pg_sys::SubXactEvent::SUBXACT_EVENT_ABORT_SUB => {
decrement_savepoint_depth();
crate::hooks::release_hook_guard_on_abort(false);
if let Some(snapshot) = QUEUE_SNAPSHOTS.with(|s| s.borrow_mut().pop()) {
super::state::replace_queue(snapshot);
}
if let Some(patch_snapshot) = PATCH_SNAPSHOTS.with(|s| s.borrow_mut().pop()) {
super::patch::replace_patch_map(patch_snapshot);
}
if let Some(fanout_snapshot) = FANOUT_SNAPSHOTS.with(|s| s.borrow_mut().pop()) {
super::patch::replace_fanout_map(fanout_snapshot);
}
super::affected::savepoint_abort();
}
pg_sys::SubXactEvent::SUBXACT_EVENT_COMMIT_SUB => {
decrement_savepoint_depth();
QUEUE_SNAPSHOTS.with(|s| {
s.borrow_mut().pop();
});
PATCH_SNAPSHOTS.with(|s| {
s.borrow_mut().pop();
});
FANOUT_SNAPSHOTS.with(|s| {
s.borrow_mut().pop();
});
super::affected::savepoint_commit();
}
_ => {
}
}
}));
if result.is_err() {
warning!("PANIC in subtransaction callback - this is a bug!");
}
}
fn decrement_savepoint_depth() {
SAVEPOINT_DEPTH.with(|d| {
let mut depth = d.borrow_mut();
if *depth == 0 {
warning!("pg_tviews: subxact depth underflow — event ordering unexpected");
}
*depth = depth.saturating_sub(1);
});
}
pub fn flush_refresh_queue() -> TViewResult<()> {
let mut pending = take_queue_snapshot();
let fanouts = super::patch::take_fanout_snapshot();
if pending.is_empty() && fanouts.is_empty() {
return Ok(());
}
crate::revision::check();
crate::config::warn_deprecated_settings();
super::affected::begin_flush();
let mut patches = super::patch::take_patch_snapshot();
let refresh_timer = crate::metrics::metrics_api::record_refresh_start();
let graph = super::cache::graph_cache::load_cached()?;
let mut processed: std::collections::HashSet<super::key::RefreshKey> =
std::collections::HashSet::with_capacity(pending.len().max(16));
let mut parent_meta_cache: std::collections::HashMap<
String,
Option<crate::catalog::TviewMeta>,
> = std::collections::HashMap::new();
apply_fanouts(fanouts, &graph, &mut patches, &mut pending, &processed)?;
let mut iteration = 1;
loop {
while !pending.is_empty() {
let sorted_keys = graph.sort_keys(pending.drain().collect());
let entity = sorted_keys[0].entity.clone();
let mut entity_keys = Vec::new();
for key in sorted_keys {
if key.entity != entity {
pending.insert(key);
} else if processed.insert(key.clone()) {
entity_keys.push(key);
}
}
let _owner = crate::owner::AsOwner::of_entity(&entity)?;
if !is_crash_recovery_checked(&entity) {
mark_crash_recovery_checked(&entity);
if crate::lifecycle::detect_post_crash_truncation(&entity)? {
crate::admin::fill_empty_tview(&entity)?;
}
}
if entity_keys.iter().any(super::key::RefreshKey::is_all) {
let meta =
crate::catalog::TviewMeta::load_by_entity(&entity)?.ok_or_else(|| {
crate::TViewError::MetadataNotFound {
entity: entity.clone(),
}
})?;
let changed: Vec<i64> = crate::ddl::replace::reconcile(&entity, &meta)?
.iter()
.filter_map(|k| k.parse::<i64>().ok())
.collect();
for parent_key in
crate::propagate::find_parents_batch(&entity, &changed, &changed, &graph)?
.into_values()
.flatten()
{
super::patch::poison_into(&mut patches, parent_key.clone());
if !processed.contains(&parent_key) {
pending.insert(parent_key);
}
}
if meta.identity.is_pk(&entity) {
processed.extend(
changed
.into_iter()
.map(|pk| super::key::RefreshKey::pk(&entity, pk)),
);
}
iteration += 1;
continue;
}
let apply_enabled = crate::config::direct_patch_enabled();
let mut patched: Vec<(i64, Vec<super::patch::PatchEntry>)> = Vec::new();
let mut applied_pks: HashSet<i64> = HashSet::new();
let mut recompute_keys: Vec<super::key::RefreshKey> = Vec::new();
for key in entity_keys {
if apply_enabled
&& let Some(pk) = key.key.as_int()
&& let Some(super::patch::PatchState::Direct(chain)) = patches.get(&key)
{
patched.push((pk, chain.clone()));
applied_pks.insert(pk);
continue;
}
recompute_keys.push(key);
}
if !patched.is_empty() {
let meta =
crate::catalog::TviewMeta::load_by_entity(&entity)?.ok_or_else(|| {
crate::TViewError::MetadataNotFound {
entity: entity.clone(),
}
})?;
let fallback = crate::refresh::direct::apply_entity_patches(&meta, patched)?;
for pk in fallback {
applied_pks.remove(&pk);
recompute_keys.push(super::key::RefreshKey::pk(&entity, pk));
}
}
if !recompute_keys.is_empty() {
let keys: Vec<super::key::KeyValue> =
recompute_keys.into_iter().map(|k| k.key).collect();
let touched = if let [key] = keys.as_slice() {
let meta =
crate::catalog::TviewMeta::load_by_entity(&entity)?.ok_or_else(|| {
crate::TViewError::MetadataNotFound {
entity: entity.clone(),
}
})?;
crate::refresh::refresh_key(&meta, key)?
} else {
crate::refresh::refresh_bulk(&entity, &keys)?
};
for parent_key in crate::propagate::find_parents_batch(
&entity,
&touched.pks,
&touched.appeared,
&graph,
)?
.into_values()
.flatten()
{
super::patch::poison_into(&mut patches, parent_key.clone());
if !processed.contains(&parent_key) {
pending.insert(parent_key);
}
}
}
if !applied_pks.is_empty() {
let applied: Vec<i64> = applied_pks.into_iter().collect();
let parent_map =
crate::propagate::find_parents_batch(&entity, &applied, &[], &graph)?;
for (child_pk, parent_keys) in &parent_map {
let child_key = &super::key::RefreshKey::pk(&entity, *child_pk);
let child_chain = match patches.get(child_key) {
Some(super::patch::PatchState::Direct(chain)) => Some(chain.clone()),
_ => None,
};
for parent_key in parent_keys {
let derived = match &child_chain {
Some(chain) => {
load_meta_cached(&parent_key.entity, &mut parent_meta_cache)?
.and_then(|m| {
crate::refresh::direct::derive_parent_chain(
&m,
&child_key.entity,
chain,
)
})
}
None => None,
};
match derived {
Some(chain) => super::patch::merge_chain_into(
&mut patches,
parent_key.clone(),
chain,
),
None => {
super::patch::poison_into(&mut patches, parent_key.clone());
}
}
if !processed.contains(parent_key) {
pending.insert(parent_key.clone());
}
}
}
}
iteration += 1;
let max_depth = crate::config::max_propagation_depth();
if iteration > max_depth {
return Err(crate::TViewError::PropagationDepthExceeded {
max_depth,
processed: processed.len(),
});
}
}
let late = take_queue_snapshot();
let late_fanouts = super::patch::take_fanout_snapshot();
if late.is_empty() && late_fanouts.is_empty() {
break;
}
for (k, v) in super::patch::take_patch_snapshot() {
patches.insert(k, v);
}
pending = late;
apply_fanouts(late_fanouts, &graph, &mut patches, &mut pending, &processed)?;
}
{
let mut entity_counts: std::collections::HashMap<&str, i64> =
std::collections::HashMap::new();
for key in &processed {
*entity_counts.entry(&key.entity).or_insert(0) += 1;
}
for (entity, count) in entity_counts {
crate::audit::log_refresh(entity, count);
}
}
crate::metrics::metrics_api::record_refresh_complete(
processed.len(),
iteration - 1,
&refresh_timer,
);
Ok(())
}
fn load_meta_cached(
entity: &str,
cache: &mut std::collections::HashMap<String, Option<crate::catalog::TviewMeta>>,
) -> spi::Result<Option<crate::catalog::TviewMeta>> {
if let Some(meta) = cache.get(entity) {
return Ok(meta.clone());
}
let meta = crate::catalog::TviewMeta::load_by_entity(entity)?;
cache.insert(entity.to_string(), meta.clone());
Ok(meta)
}
fn apply_fanouts(
fanouts: super::patch::FanoutMap,
graph: &super::graph::EntityDepGraph,
patches: &mut std::collections::HashMap<super::key::RefreshKey, super::patch::PatchState>,
pending: &mut HashSet<super::key::RefreshKey>,
processed: &HashSet<super::key::RefreshKey>,
) -> TViewResult<()> {
let mut groups: std::collections::BTreeMap<(String, String), Vec<_>> =
std::collections::BTreeMap::new();
for ((entity, lookup_col, key), fields) in fanouts {
groups
.entry((entity, lookup_col))
.or_default()
.push((key, fields));
}
for ((entity, lookup_col), rows) in groups {
let meta = crate::catalog::TviewMeta::load_by_entity(&entity)?.ok_or_else(|| {
crate::TViewError::MetadataNotFound {
entity: entity.clone(),
}
})?;
let owner = crate::owner::AsOwner::of_table(meta.tview_oid)?;
let changed = crate::refresh::direct::apply_fanout_patch(&meta, &lookup_col, &rows)?;
drop(owner);
for parent_key in crate::propagate::find_parents_batch(&entity, &changed, &[], graph)?
.into_values()
.flatten()
{
super::patch::poison_into(patches, parent_key.clone());
if !processed.contains(&parent_key) {
pending.insert(parent_key);
}
}
}
Ok(())
}