pub mod aggregate;
pub mod create;
pub mod drop;
pub mod rename;
pub mod replace;
pub(crate) mod uncascaded;
pub use drop::drop_tview;
use crate::error::{TViewError, TViewResult};
use pgrx::datum::DatumWithOid;
use pgrx::prelude::*;
const REGISTRATION_LOCK_CLASS: i32 = 0x7476_6965;
pub(crate) fn lock_entity(entity: &str) -> TViewResult<()> {
Spi::run_with_args(
"SELECT pg_catalog.pg_advisory_xact_lock($1, pg_catalog.hashtext($2))",
&[
unsafe {
DatumWithOid::new(
REGISTRATION_LOCK_CLASS,
PgOid::BuiltIn(PgBuiltInOids::INT4OID).value(),
)
},
unsafe { DatumWithOid::new(entity, PgOid::BuiltIn(PgBuiltInOids::TEXTOID).value()) },
],
)
.map_err(|e| TViewError::CatalogError {
operation: format!("Lock the registration of TVIEW {entity}"),
pg_error: e.to_string(),
})
}
#[pg_extern]
fn pg_tviews_create(tview_name: &str, select_sql: &str) -> Result<String, String> {
crate::revision::check();
create_reported(tview_name, select_sql, replace::Options::default())
}
#[pg_extern]
#[allow(clippy::needless_pass_by_value)] fn pg_tviews_create_aggregate(
tview_name: &str,
select_sql: &str,
group_keys: pgrx::JsonB,
) -> Result<String, String> {
crate::revision::check();
let keys: aggregate::GroupKeys = serde_json::from_value(group_keys.0)
.ok()
.filter(|keys: &aggregate::GroupKeys| !keys.is_empty())
.ok_or_else(|| {
"group_keys must be a JSON object mapping source table names to column names, e.g. \
'{\"tb_order\": \"fk_user\"}'"
.to_string()
})?;
create_reported(tview_name, select_sql, replace::Options::aggregate(keys))
}
fn create_reported(
tview_name: &str,
select_sql: &str,
options: replace::Options,
) -> Result<String, String> {
unsafe {
crate::hooks::ensure_hook_installed();
}
match replace::create_only(tview_name, select_sql, options, false) {
Ok(replace::Created::Rows(_) | replace::Created::Skipped) => {
Ok(format!("TVIEW '{tview_name}' created successfully"))
}
Ok(replace::Created::Exists(name)) => Err(format!(
"TVIEW {name} already exists; pg_tviews_create_or_replace() changes an existing TVIEW"
)),
Err(e) => Err(format!("Failed to create TVIEW: {e}")),
}
}
#[pg_extern]
fn pg_tviews_handle_dropped(entity: &str) -> Result<(), String> {
drop::handle_dropped(entity).map_err(|e| format!("Failed to deregister TVIEW '{entity}': {e}"))
}
#[pg_extern]
#[allow(clippy::needless_pass_by_value)] fn pg_tviews_create_or_replace(
tview_name: &str,
query: &str,
options: default!(pgrx::JsonB, "'{}'"),
) -> Result<String, String> {
crate::revision::check();
replace::create_or_replace(tview_name, query, &options.0)
.map(str::to_string)
.map_err(|e| format!("Failed to create or replace TVIEW '{tview_name}': {e}"))
}
#[pg_extern]
fn pg_tviews_drop(
tview_name: &str,
if_exists: default!(bool, false),
cascade: default!(bool, false),
) -> Result<String, String> {
crate::revision::check();
match drop_tview(tview_name, if_exists, cascade) {
Ok(true) => Ok(format!("TVIEW '{tview_name}' dropped successfully")),
Ok(false) => Ok(format!(
"TVIEW '{tview_name}' does not exist, nothing dropped"
)),
Err(e) => Err(format!("Failed to drop TVIEW: {e}")),
}
}
#[pg_extern]
fn pg_tviews_reregister(tview_name: &str) -> Result<String, String> {
crate::revision::check();
crate::validation::validate_sql_identifier(tview_name, "tview_name")
.map_err(|e| format!("Invalid TVIEW name: {e}"))?;
let entity = tview_name.strip_prefix("tv_").unwrap_or(tview_name);
create::reregister_tview(entity)
.map(|()| "reregistered".to_string())
.map_err(|e| format!("Failed to re-register TVIEW '{entity}': {e}"))
}
extension_sql!(
r"
CREATE FUNCTION @extschema@.pg_tviews_reregister_all(strict BOOLEAN DEFAULT false)
RETURNS TABLE (entity TEXT, status TEXT)
LANGUAGE plpgsql
AS $$
#variable_conflict use_column
DECLARE
next_entity TEXT;
failures INTEGER := 0;
BEGIN
FOR next_entity IN
WITH RECURSIVE edges(entity, dependency) AS (
SELECT DISTINCT r.entity, m.entity
FROM @extschema@.pg_tview_reads r
JOIN @extschema@.pg_tview_meta m
ON r.relid IN (m.view_oid::oid, m.table_oid::oid)
WHERE m.entity <> r.entity
),
depth(entity, level) AS (
SELECT m.entity, 0 FROM @extschema@.pg_tview_meta m
UNION
SELECT e.entity, d.level + 1
FROM depth d JOIN edges e ON e.dependency = d.entity
WHERE d.level < 100
)
SELECT d.entity FROM depth d GROUP BY d.entity ORDER BY max(d.level), d.entity
LOOP
entity := next_entity;
BEGIN
PERFORM @extschema@.pg_tviews_reregister(next_entity);
status := 'reregistered';
EXCEPTION WHEN OTHERS THEN
status := SQLERRM;
failures := failures + 1;
END;
RETURN NEXT;
END LOOP;
IF strict AND failures > 0 THEN
RAISE EXCEPTION 'pg_tviews: % TVIEW(s) could not be re-registered', failures
USING HINT = 'SELECT * FROM tviews.pg_tviews_reregister_all() lists them';
END IF;
END;
$$;
",
name = "reregister_all",
requires = [
pg_tviews_reregister,
"create_metadata_tables",
"tview_reads"
],
);
#[pg_extern]
#[allow(clippy::needless_pass_by_value)] fn pg_tviews_rebind_cascade_paths(
view_oid: pg_sys::Oid,
cascade_paths: Vec<String>,
) -> Result<Vec<String>, String> {
crate::revision::check();
create::rebind_cascade_paths(view_oid, &cascade_paths)
.map_err(|e| format!("Failed to rebind cascade paths: {e}"))
}
#[pg_extern]
fn pg_tviews_convert_existing_table(table_name: &str) -> Result<String, String> {
crate::validation::validate_sql_identifier(table_name, "table_name")
.map_err(|e| format!("Invalid table name: {e}"))?;
Err(format!(
"pg_tviews_convert_existing_table() is deprecated and no longer converts '{table_name}'; \
use pg_tviews_create() or CREATE TABLE tv_<entity> AS SELECT ... instead"
))
}