pub mod aggregate;
pub mod create;
pub mod drop;
pub(crate) mod privileges;
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::*;
pub(crate) fn backing_view_name(schema: &str, table: &str) -> (String, String) {
(
crate::utils::ext_schema().to_string(),
crate::utils::fit_identifier(format!("{schema}__{table}")),
)
}
pub(crate) fn follow_table_move(table: pg_sys::Oid) -> TViewResult<()> {
let args = [unsafe { DatumWithOid::new(table, PgOid::BuiltIn(PgBuiltInOids::OIDOID).value()) }];
let view = Spi::get_one_with_args::<pg_sys::Oid>(
&format!(
"SELECT (SELECT m.view_oid::pg_catalog.oid FROM {} m \
WHERE m.table_oid::pg_catalog.oid = $1)",
crate::utils::meta_table()
),
&args,
)
.map_err(|e| TViewError::CatalogError {
operation: "Find the TVIEW of a moved table".to_string(),
pg_error: e.to_string(),
})?;
let Some(view) = view else {
return Ok(());
};
let (schema, name) = relation_name(table)?;
let (view_schema, wanted) = backing_view_name(&schema, &name);
let (_, current) = relation_name(view)?;
if current == wanted {
return Ok(());
}
let qualified = format!(
"{}.{}",
crate::utils::quote_identifier(&view_schema),
crate::utils::quote_identifier(&wanted)
);
let taken = Spi::get_one_with_args::<bool>(
"SELECT pg_catalog.to_regclass($1) IS NOT NULL",
&[unsafe {
DatumWithOid::new(
qualified.as_str(),
PgOid::BuiltIn(PgBuiltInOids::TEXTOID).value(),
)
}],
)
.map_err(|e| TViewError::CatalogError {
operation: format!("Check {view_schema}.{wanted}"),
pg_error: e.to_string(),
})?
.unwrap_or(false);
if taken {
return Err(TViewError::InvalidInput {
parameter: "table name".to_string(),
reason: format!(
"the backing view of {schema}.{name} would be {view_schema}.{wanted}, which is \
already taken by another relation"
),
});
}
let sql = format!(
"ALTER VIEW {} RENAME TO {}",
crate::utils::qualified_relname_from_oid(view)?,
crate::utils::quote_identifier(&wanted)
);
{
let _owner = crate::owner::AsOwner::of_table(view)?;
crate::utils::spi_run_ddl(&sql)
.map_err(|error| TViewError::SpiError { query: sql, error })?;
}
crate::queue::cache::invalidate_all_caches();
unsafe { pg_sys::CacheInvalidateRelcacheByRelid(table) };
Ok(())
}
pub(crate) fn in_extension_schema<T>(ddl: impl FnOnce() -> TViewResult<T>) -> TViewResult<T> {
in_extension_schema_for(unsafe { pg_sys::GetUserId() }, ddl)
}
pub(crate) fn in_extension_schema_for<T>(
role: pg_sys::Oid,
ddl: impl FnOnce() -> TViewResult<T>,
) -> TViewResult<T> {
let schema = crate::utils::ext_schema();
let args = [unsafe { DatumWithOid::new(role, PgOid::BuiltIn(PgBuiltInOids::OIDOID).value()) }];
let (can_create, name) = Spi::get_two_with_args::<bool, String>(
&format!(
"SELECT pg_catalog.has_schema_privilege($1, '{schema}', 'CREATE'), \
pg_catalog.quote_ident(pg_catalog.pg_get_userbyid($1))"
),
&args,
)
.map_err(|e| TViewError::CatalogError {
operation: format!("Check CREATE on schema {schema}"),
pg_error: e.to_string(),
})?;
if can_create.unwrap_or(false) {
return ddl();
}
let role = name.unwrap_or_default();
let privilege = |verb: &str| -> TViewResult<()> {
let sql = if verb == "GRANT" {
format!("GRANT CREATE ON SCHEMA {schema} TO {role}")
} else {
format!("REVOKE CREATE ON SCHEMA {schema} FROM {role}")
};
let _owner = crate::owner::AsOwner::of_extension()?;
crate::utils::spi_run_ddl(&sql).map_err(|error| TViewError::SpiError { query: sql, error })
};
privilege("GRANT")?;
let result = ddl()?;
privilege("REVOKE")?;
Ok(result)
}
pub(crate) fn relation_name(oid: pg_sys::Oid) -> TViewResult<(String, String)> {
let args = [unsafe { DatumWithOid::new(oid, PgOid::BuiltIn(PgBuiltInOids::OIDOID).value()) }];
let (name, schema) = Spi::get_two_with_args::<String, String>(
"SELECT c.relname::text, n.nspname::text FROM pg_catalog.pg_class c \
JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace WHERE c.oid = $1",
&args,
)
.map_err(|e| TViewError::CatalogError {
operation: format!("Name relation {oid:?}"),
pg_error: e.to_string(),
})?;
match (schema, name) {
(Some(schema), Some(name)) => Ok((schema, name)),
_ => Err(TViewError::CatalogError {
operation: format!("Name relation {oid:?}"),
pg_error: "relation not found".to_string(),
}),
}
}
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"
))
}