use pgrx::prelude::*;
use crate::catalog::{DependencyDetail, DependencyType, TviewMeta};
use crate::queue::key::KeyValue;
use crate::lifecycle::jsonb_delta_schema;
use crate::utils::{qualified_relname_from_oid, quote_identifier};
pub fn refresh_key(meta: &TviewMeta, key: &KeyValue) -> spi::Result<super::Touched> {
let keys = std::slice::from_ref(key);
let before = super::lock_rows(
meta,
&qualified_relname_from_oid(meta.tview_oid)?,
keys,
!meta.identity.is_pk(&meta.entity_name),
)?;
let (written, deleted) = if meta.is_union && !view_row_exists(meta, key)? {
(super::Written::default(), delete_tview_row(meta, key)?)
} else {
crate::metrics::metrics_api::record_view_recomputes(1);
let (produced, written) = apply_patch(meta, key)?;
let deleted = if produced == 0 {
delete_tview_row(meta, key)?
} else {
Vec::new()
};
(written, deleted)
};
Ok(super::touched(meta, keys, before, written, deleted))
}
fn delete_tview_row(meta: &TviewMeta, key: &KeyValue) -> spi::Result<Vec<i64>> {
let key_type = meta.key_type()?;
let qi_tv = qualified_relname_from_oid(meta.tview_oid)?;
let qi_key = quote_identifier(&meta.identity.column);
let qi_pk = quote_identifier(&format!("pk_{}", meta.entity_name));
let sql = format!(
"DELETE FROM {qi_tv} WHERE {qi_key} = {} \
RETURNING {qi_pk}::text, to_jsonb({qi_tv}.*)->>'id'",
super::key_cast(&key_type, "$1", false)
);
super::run_journaled_delete(
&meta.entity_name,
&sql,
&[super::key_scalar(&key_type, key)?],
)
}
fn view_row_exists(meta: &TviewMeta, key: &KeyValue) -> spi::Result<bool> {
let key_type = meta.key_type()?;
let qi_view = qualified_relname_from_oid(meta.view_oid)?;
let sql = format!(
"SELECT 1 FROM {qi_view} WHERE {} = {} LIMIT 2",
quote_identifier(&meta.identity.column),
super::key_cast(&key_type, "$1", false)
);
Spi::connect(|client| {
let args = [super::key_scalar(&key_type, key)?];
let mut rows = client.select(&sql, None, &args)?;
if rows.next().is_none() {
return Ok(false);
}
if meta.is_union && rows.next().is_some() {
let policy = crate::config::union_duplicate_policy();
if policy == "first" {
crate::utils::log_once(
&format!("union_duplicate:{}", meta.entity_name),
&format!(
"TVIEW '{}': UNION ALL backing view returned multiple rows for {}={key}; \
taking the first row (union_duplicate_policy=first). Reported once \
per backend.",
meta.entity_name, meta.identity.column
),
);
} else {
return Err(spi::Error::from(crate::TViewError::SpiError {
query: sql.clone(),
error: format!(
"TVIEW '{}': UNION ALL backing view returned multiple rows for {}={key}. \
Ensure UNION ALL branches are mutually exclusive, or set \
pg_tviews.union_duplicate_policy='first' to suppress this error.",
meta.entity_name, meta.identity.column
),
}));
}
}
Ok(true)
})
}
fn apply_patch(meta: &TviewMeta, key: &KeyValue) -> spi::Result<(i64, super::Written)> {
let key_type = meta.key_type()?;
let key_col = &meta.identity.column;
let Some(delta_schema) = jsonb_delta_schema() else {
crate::utils::log_once(
crate::lifecycle::JSONB_DELTA_MISSING,
"jsonb_delta is not installed: smart JSONB patching is disabled and cascades \
replace whole documents (about 2x slower). CREATE EXTENSION jsonb_delta to enable it.",
);
return apply_full_replacement(meta, key);
};
let deps = meta.parse_dependencies();
if deps.is_empty() {
return apply_full_replacement(meta, key);
}
if deps.iter().any(|d| {
matches!(
d.dep_type,
DependencyType::Array | DependencyType::NestedObject
)
}) {
return apply_full_replacement(meta, key);
}
let col_names = crate::utils::get_view_columns_by_oid(meta.view_oid)?;
if col_names.is_empty() {
return apply_full_replacement(meta, key);
}
let col_list = super::column_list(&col_names);
let qi_tv = qualified_relname_from_oid(meta.tview_oid)?;
let qi_view = qualified_relname_from_oid(meta.view_oid)?;
let qi_key = quote_identifier(key_col);
let patch_expr = build_smart_patch_expr(
&delta_schema,
&deps,
&format!("{qi_tv}.data"),
"EXCLUDED.\"data\"",
);
let conflict = format!(
"ON CONFLICT ({qi_key}) {}",
super::upsert_conflict_action(&qi_tv, &col_names, key_col, Some(&patch_expr))
);
super::run_counted_upsert(
&meta.entity_name,
&qi_tv,
&col_list,
&format!(
"SELECT {col_list} FROM {qi_view} WHERE {qi_key} = {}",
super::key_cast(&key_type, "$1", false)
),
&conflict,
&[super::key_scalar(&key_type, key)?],
)
}
fn build_smart_patch_expr(
schema: &str,
deps: &[DependencyDetail],
base_data_expr: &str,
source: &str,
) -> String {
let mut patch_expr = base_data_expr.to_string();
for dep in deps {
patch_expr = match dep.dep_type {
DependencyType::NestedObject => {
if let Some(path) = &dep.path {
let path_str = path.join(",");
format!(
"{schema}.jsonb_smart_patch_nested({patch_expr}, {source}, ARRAY['{path_str}'])"
)
} else {
warning!("NestedObject dependency missing path, skipping");
patch_expr
}
}
DependencyType::Array => {
patch_expr
}
DependencyType::Scalar => {
format!("{schema}.jsonb_smart_patch_scalar({patch_expr}, {source})")
}
};
}
patch_expr
}
fn apply_full_replacement(meta: &TviewMeta, key: &KeyValue) -> spi::Result<(i64, super::Written)> {
let key_type = meta.key_type()?;
let qi_tv = qualified_relname_from_oid(meta.tview_oid)?;
let key_col = &meta.identity.column;
let qi_key = quote_identifier(key_col);
let qi_view = qualified_relname_from_oid(meta.view_oid)?;
let col_names = crate::utils::get_view_columns_by_oid(meta.view_oid)?;
let col_list = super::column_list(&col_names);
super::run_counted_upsert(
&meta.entity_name,
&qi_tv,
&col_list,
&format!(
"SELECT {col_list} FROM {qi_view} WHERE {qi_key} = {}",
super::key_cast(&key_type, "$1", false)
),
&format!(
"ON CONFLICT ({qi_key}) {}",
super::upsert_conflict_action(&qi_tv, &col_names, key_col, None)
),
&[super::key_scalar(&key_type, key)?],
)
}
#[cfg(any(test, feature = "pg_test"))]
#[pg_schema]
mod tests {
use crate::queue::key::KeyValue;
use pgrx::JsonB;
use pgrx::prelude::*;
fn refresh_row(source: pg_sys::Oid, key: &KeyValue) -> spi::Result<crate::refresh::Touched> {
let entity = Spi::get_one::<String>(&format!(
"SELECT entity FROM {} WHERE view_oid::oid = {1} OR table_oid::oid = {1}",
crate::utils::meta_table(),
source.to_u32()
))?
.expect("a registered TVIEW");
let meta = crate::catalog::TviewMeta::load_by_entity(&entity)?.expect("its metadata");
crate::refresh::refresh_key(&meta, key)
}
#[pg_test]
fn test_apply_patch_nested_object() {
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_user BIGINT REFERENCES tb_user(pk_user),
title TEXT
)",
)
.unwrap();
Spi::run("INSERT INTO tb_user (pk_user, name) VALUES (1, 'Alice')").unwrap();
Spi::run("INSERT INTO tb_post (pk_post, fk_user, title) VALUES (1, 1, 'Hello')").unwrap();
Spi::run(
"
SELECT pg_tviews_create('user', $$
SELECT pk_user, jsonb_build_object('name', name) AS data
FROM tb_user
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create(
'post',
$$
SELECT pk_post, fk_user,
jsonb_build_object(
'title', title,
'author', v_user.data
) AS data
FROM tb_post
LEFT JOIN v_user ON v_user.pk_user = tb_post.fk_user
$$
)
",
)
.unwrap();
let meta = crate::utils::spi_get_string(
"
SELECT dependency_types::text FROM pg_tview_meta
WHERE entity = 'post'
",
)
.unwrap()
.unwrap();
assert!(
meta.contains("nested_object"),
"Expected nested_object dependency, got: {meta}"
);
let initial_data = Spi::get_one::<JsonB>(
"
SELECT data FROM tv_post WHERE pk_post = 1
",
)
.unwrap()
.unwrap();
let initial_json = &initial_data.0;
assert_eq!(initial_json["title"], "Hello");
assert_eq!(initial_json["author"]["name"], "Alice");
Spi::run("UPDATE tb_user SET name = 'Alice Updated' WHERE pk_user = 1").unwrap();
let user_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_user'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(user_oid, &KeyValue::Int(1)).unwrap();
let post_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_post'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(post_oid, &KeyValue::Int(1)).unwrap();
let updated_data = Spi::get_one::<JsonB>(
"
SELECT data FROM tv_post WHERE pk_post = 1
",
)
.unwrap()
.unwrap();
let updated_json = &updated_data.0;
assert_eq!(
updated_json["title"], "Hello",
"Title should be recomputed unchanged (own column not modified)"
);
assert_eq!(
updated_json["author"]["name"], "Alice Updated",
"Author name should be updated via smart patch"
);
}
#[pg_test]
fn test_apply_patch_array() {
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_user BIGINT REFERENCES tb_user(pk_user),
title TEXT
)",
)
.unwrap();
Spi::run(
"CREATE TABLE tb_comment (
pk_comment BIGSERIAL PRIMARY KEY,
fk_post BIGINT REFERENCES tb_post(pk_post),
fk_user BIGINT REFERENCES tb_user(pk_user),
text TEXT
)",
)
.unwrap();
Spi::run("INSERT INTO tb_user (pk_user, name) VALUES (1, 'Alice')").unwrap();
Spi::run("INSERT INTO tb_post (pk_post, fk_user, title) VALUES (1, 1, 'Hello')").unwrap();
Spi::run(
"INSERT INTO tb_comment (pk_comment, fk_post, fk_user, text)
VALUES (1, 1, 1, 'Great post!')",
)
.unwrap();
Spi::run(
"INSERT INTO tb_comment (pk_comment, fk_post, fk_user, text)
VALUES (2, 1, 1, 'Thanks!')",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('user', $$
SELECT pk_user, jsonb_build_object('name', name) AS data
FROM tb_user
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('comment', $$
SELECT pk_comment, fk_post, fk_user,
jsonb_build_object('text', text) AS data
FROM tb_comment
$$)
",
)
.unwrap();
Spi::run("
SELECT pg_tviews_create(
'post',
$$
SELECT pk_post, fk_user,
jsonb_build_object(
'title', title,
'author', v_user.data,
'comments', COALESCE(jsonb_agg(v_comment.data ORDER BY v_comment.pk_comment), '[]'::jsonb)
) AS data
FROM tb_post
LEFT JOIN v_user ON v_user.pk_user = tb_post.fk_user
LEFT JOIN v_comment ON v_comment.fk_post = tb_post.pk_post
GROUP BY pk_post, fk_user, title, v_user.data
$$
)
").unwrap();
let meta = crate::utils::spi_get_string(
"
SELECT dependency_types::text FROM pg_tview_meta
WHERE entity = 'post'
",
)
.unwrap()
.unwrap();
assert!(
meta.contains("array"),
"Expected array dependency, got: {meta}"
);
let initial_data = Spi::get_one::<JsonB>(
"
SELECT data FROM tv_post WHERE pk_post = 1
",
)
.unwrap()
.unwrap();
let initial_comments = initial_data.0["comments"].as_array().unwrap();
assert_eq!(
initial_comments.len(),
2,
"Should have 2 comments initially"
);
Spi::run("UPDATE tb_comment SET text = 'Updated!' WHERE pk_comment = 1").unwrap();
let comment_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_comment'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(comment_oid, &KeyValue::Int(1)).unwrap();
let post_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_post'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(post_oid, &KeyValue::Int(1)).unwrap();
let updated_data = Spi::get_one::<JsonB>(
"
SELECT data FROM tv_post WHERE pk_post = 1
",
)
.unwrap()
.unwrap();
let comments = updated_data.0["comments"].as_array().unwrap();
assert_eq!(comments.len(), 2, "Should still have 2 comments");
let comment_1 = comments
.iter()
.find(|c| c["id"].as_i64() == Some(1))
.expect("Should find comment with id=1");
let comment_2 = comments
.iter()
.find(|c| c["id"].as_i64() == Some(2))
.expect("Should find comment with id=2");
assert_eq!(comment_1["text"], "Updated!", "Comment 1 should be updated");
assert_eq!(
comment_2["text"], "Thanks!",
"Comment 2 should be unchanged"
);
}
#[pg_test]
fn test_apply_patch_scalar() {
Spi::run("CREATE TABLE tb_category (pk_category BIGSERIAL PRIMARY KEY, name TEXT)")
.unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_category BIGINT REFERENCES tb_category(pk_category),
title TEXT
)",
)
.unwrap();
Spi::run("INSERT INTO tb_category (pk_category, name) VALUES (1, 'Tech')").unwrap();
Spi::run("INSERT INTO tb_post (pk_post, fk_category, title) VALUES (1, 1, 'Hello')")
.unwrap();
Spi::run(
"
SELECT pg_tviews_create(
'post',
$$
SELECT pk_post, fk_category,
jsonb_build_object('title', title) AS data
FROM tb_post
$$
)
",
)
.unwrap();
let meta = crate::utils::spi_get_string(
"
SELECT dependency_types::text FROM pg_tview_meta
WHERE entity ='post'
",
)
.unwrap()
.unwrap();
assert!(
meta.contains("scalar"),
"Expected scalar dependency, got: {meta}"
);
let initial_data = Spi::get_one::<JsonB>(
"
SELECT data FROM tv_post WHERE pk_post = 1
",
)
.unwrap()
.unwrap();
assert_eq!(initial_data.0["title"], "Hello");
assert!(
initial_data.0.get("category").is_none(),
"Should not have category in data"
);
Spi::run("UPDATE tb_category SET name = 'Technology' WHERE pk_category = 1").unwrap();
let post_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_post'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(post_oid, &KeyValue::Int(1)).unwrap();
let updated_data = Spi::get_one::<JsonB>(
"
SELECT data FROM tv_post WHERE pk_post = 1
",
)
.unwrap()
.unwrap();
assert_eq!(
updated_data.0["title"], "Hello",
"Title should be unchanged"
);
assert!(
updated_data.0.get("category").is_none(),
"Still no category in data"
);
}
#[pg_test]
fn test_smart_patch_full_integration() {
let _ = Spi::run("CREATE EXTENSION IF NOT EXISTS jsonb_delta");
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT, email TEXT)")
.unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_user BIGINT REFERENCES tb_user(pk_user),
title TEXT,
content TEXT
)",
)
.unwrap();
Spi::run(
"CREATE TABLE tb_comment (
pk_comment BIGSERIAL PRIMARY KEY,
fk_post BIGINT REFERENCES tb_post(pk_post),
fk_user BIGINT REFERENCES tb_user(pk_user),
text TEXT
)",
)
.unwrap();
Spi::run(
"INSERT INTO tb_user (pk_user, name, email) VALUES (1, 'Alice', 'alice@example.com')",
)
.unwrap();
Spi::run("INSERT INTO tb_user (pk_user, name, email) VALUES (2, 'Bob', 'bob@example.com')")
.unwrap();
Spi::run(
"INSERT INTO tb_post (pk_post, fk_user, title, content)
VALUES (1, 1, 'First Post', 'Hello World')",
)
.unwrap();
Spi::run(
"INSERT INTO tb_comment (pk_comment, fk_post, fk_user, text)
VALUES (1, 1, 1, 'Great post!')",
)
.unwrap();
Spi::run(
"INSERT INTO tb_comment (pk_comment, fk_post, fk_user, text)
VALUES (2, 1, 2, 'Thanks for sharing!')",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('user', $$
SELECT pk_user, jsonb_build_object('name', name, 'email', email) AS data
FROM tb_user
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('comment', $$
SELECT pk_comment, fk_post, fk_user,
jsonb_build_object('text', text) AS data
FROM tb_comment
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('post', $$
SELECT pk_post, fk_user,
jsonb_build_object(
'title', title,
'content', content,
'author', v_user.data,
'comments', COALESCE(
jsonb_agg(
v_comment.data
ORDER BY v_comment.pk_comment
),
'[]'::jsonb
)
) AS data
FROM tb_post
LEFT JOIN v_user ON v_user.pk_user = tb_post.fk_user
LEFT JOIN v_comment ON v_comment.fk_post = tb_post.pk_post
GROUP BY pk_post, fk_user, title, content, v_user.data
$$)
",
)
.unwrap();
let initial = Spi::get_one::<JsonB>("SELECT data FROM tv_post WHERE pk_post = 1")
.unwrap()
.unwrap();
assert_eq!(initial.0["title"], "First Post");
assert_eq!(initial.0["author"]["name"], "Alice");
assert_eq!(initial.0["comments"].as_array().unwrap().len(), 2);
Spi::run(
"UPDATE tb_user SET name = 'Alice Updated', email = 'alice.new@example.com'
WHERE pk_user = 1",
)
.unwrap();
let user_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_user'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(user_oid, &KeyValue::Int(1)).unwrap();
let post_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_post'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(post_oid, &KeyValue::Int(1)).unwrap();
let after_author_update =
Spi::get_one::<JsonB>("SELECT data FROM tv_post WHERE pk_post = 1")
.unwrap()
.unwrap();
assert_eq!(after_author_update.0["author"]["name"], "Alice Updated");
assert_eq!(
after_author_update.0["author"]["email"],
"alice.new@example.com"
);
assert_eq!(after_author_update.0["title"], "First Post");
assert_eq!(after_author_update.0["content"], "Hello World");
assert_eq!(
after_author_update.0["comments"].as_array().unwrap().len(),
2
);
Spi::run("UPDATE tb_comment SET text = 'Updated comment!' WHERE pk_comment = 1").unwrap();
let comment_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_comment'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(comment_oid, &KeyValue::Int(1)).unwrap();
refresh_row(post_oid, &KeyValue::Int(1)).unwrap();
let after_comment_update =
Spi::get_one::<JsonB>("SELECT data FROM tv_post WHERE pk_post = 1")
.unwrap()
.unwrap();
let comments = after_comment_update.0["comments"].as_array().unwrap();
assert_eq!(comments.len(), 2, "Should still have 2 comments");
let comment_1 = comments
.iter()
.find(|c| c["id"].as_i64() == Some(1))
.expect("Should find comment 1");
assert_eq!(comment_1["text"], "Updated comment!");
let comment_2 = comments
.iter()
.find(|c| c["id"].as_i64() == Some(2))
.expect("Should find comment 2");
assert_eq!(comment_2["text"], "Thanks for sharing!");
}
#[pg_test]
fn test_fallback_without_jsonb_delta() {
let _ = Spi::run("DROP EXTENSION IF EXISTS jsonb_delta CASCADE");
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_user BIGINT REFERENCES tb_user(pk_user),
title TEXT
)",
)
.unwrap();
Spi::run("INSERT INTO tb_user VALUES (1, 'Alice')").unwrap();
Spi::run("INSERT INTO tb_post VALUES (1, 1, 'Hello')").unwrap();
Spi::run(
"
SELECT pg_tviews_create('user', $$
SELECT pk_user, jsonb_build_object('name', name) AS data
FROM tb_user
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('post', $$
SELECT pk_post, fk_user,
jsonb_build_object('title', title, 'author', v_user.data) AS data
FROM tb_post
LEFT JOIN v_user ON v_user.pk_user = tb_post.fk_user
$$)
",
)
.unwrap();
let meta = crate::utils::spi_get_string(
"
SELECT dependency_types::text FROM pg_tview_meta WHERE entity = 'post'
",
);
assert!(
meta.is_ok(),
"Metadata should be captured even without jsonb_delta"
);
Spi::run("UPDATE tb_user SET name = 'Alice Fallback' WHERE pk_user = 1").unwrap();
let user_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_user'::regclass::oid")
.unwrap()
.unwrap();
let result = refresh_row(user_oid, &KeyValue::Int(1));
assert!(result.is_ok(), "Fallback should work without jsonb_delta");
let post_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_post'::regclass::oid")
.unwrap()
.unwrap();
refresh_row(post_oid, &KeyValue::Int(1)).unwrap();
let updated = Spi::get_one::<JsonB>("SELECT data FROM tv_post WHERE pk_post = 1")
.unwrap()
.unwrap();
assert_eq!(updated.0["author"]["name"], "Alice Fallback");
assert_eq!(updated.0["title"], "Hello");
}
#[pg_test]
fn test_legacy_tview_fallback() {
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run("INSERT INTO tb_user VALUES (1, 'Alice')").unwrap();
Spi::run(
"
SELECT pg_tviews_create('user', $$
SELECT pk_user, jsonb_build_object('name', name) AS data
FROM tb_user
$$)
",
)
.unwrap();
Spi::run(
"
UPDATE pg_tview_meta
SET dependency_types = NULL,
dependency_paths = NULL,
array_match_keys = NULL
WHERE entity ='user'
",
)
.unwrap();
Spi::run("UPDATE tb_user SET name = 'Alice Legacy' WHERE pk_user = 1").unwrap();
let user_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_user'::regclass::oid")
.unwrap()
.unwrap();
let result = refresh_row(user_oid, &KeyValue::Int(1));
assert!(result.is_ok(), "Should handle legacy TVIEW gracefully");
let updated = Spi::get_one::<JsonB>("SELECT data FROM tv_user WHERE pk_user = 1")
.unwrap()
.unwrap();
assert_eq!(updated.0["name"], "Alice Legacy");
}
#[pg_test]
fn test_refresh_distinct_on_basic() {
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_user BIGINT REFERENCES tb_user(pk_user),
title TEXT,
created_at TIMESTAMP DEFAULT NOW()
)",
)
.unwrap();
Spi::run("INSERT INTO tb_user VALUES (1, 'Alice')").unwrap();
Spi::run(
"INSERT INTO tb_post (pk_post, fk_user, title) VALUES
(1, 1, 'First Post'),
(2, 1, 'Second Post'),
(3, 1, 'Third Post')",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('user', $$
SELECT pk_user, jsonb_build_object('name', name) AS data
FROM tb_user
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('post_by_user', $$
SELECT DISTINCT ON (fk_user)
pk_post, fk_user,
jsonb_build_object('title', title) AS data
FROM tb_post
ORDER BY fk_user, pk_post
$$, 'fk_user')
",
)
.unwrap();
let identity = crate::utils::spi_get_string(
"
SELECT identity::text FROM pg_tview_meta
WHERE entity = 'post_by_user'
",
)
.unwrap()
.unwrap();
assert!(
identity.contains("fk_user"),
"Should record the DISTINCT ON key as the identity"
);
let initial_count: i64 = Spi::get_one(
"
SELECT COUNT(*) FROM tv_post_by_user WHERE fk_user = 1
",
)
.unwrap()
.unwrap();
assert_eq!(initial_count, 1, "Should have exactly 1 row for fk_user=1");
let initial_title: String = Spi::get_one(
"
SELECT data->>'title' FROM tv_post_by_user WHERE fk_user = 1
",
)
.unwrap()
.unwrap();
assert_eq!(
initial_title, "First Post",
"Should be first post initially"
);
}
#[pg_test]
fn test_refresh_distinct_on_multiple_keys() {
Spi::run("CREATE TABLE tb_category (pk_category BIGSERIAL PRIMARY KEY, name TEXT)")
.unwrap();
Spi::run(
"CREATE TABLE tb_item (
pk_item BIGSERIAL PRIMARY KEY,
fk_category BIGINT REFERENCES tb_category(pk_category),
title TEXT
)",
)
.unwrap();
Spi::run("INSERT INTO tb_category VALUES (1, 'Tech'), (2, 'News')").unwrap();
Spi::run(
"INSERT INTO tb_item (pk_item, fk_category, title) VALUES
(1, 1, 'Item 1A'),
(2, 1, 'Item 1B'),
(3, 1, 'Item 1C'),
(4, 2, 'Item 2A'),
(5, 2, 'Item 2B')",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('category', $$
SELECT pk_category, jsonb_build_object('name', name) AS data
FROM tb_category
$$)
",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('item_by_cat', $$
SELECT DISTINCT ON (fk_category)
pk_item, fk_category,
jsonb_build_object('title', title) AS data
FROM tb_item
ORDER BY fk_category, pk_item
$$, 'fk_category')
",
)
.unwrap();
let cat1_count: i64 = Spi::get_one(
"
SELECT COUNT(*) FROM tv_item_by_cat WHERE fk_category = 1
",
)
.unwrap()
.unwrap();
assert_eq!(cat1_count, 1, "Should have 1 row for category 1");
let cat1_title: String = Spi::get_one(
"
SELECT data->>'title' FROM tv_item_by_cat WHERE fk_category = 1
",
)
.unwrap()
.unwrap();
assert_eq!(cat1_title, "Item 1A", "Category 1 should show Item 1A");
Spi::run("DELETE FROM tb_item WHERE pk_item = 1").unwrap();
let view_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'v_item_by_cat'::regclass::oid")
.unwrap()
.unwrap();
let result = refresh_row(view_oid, &KeyValue::Int(1));
assert!(result.is_ok(), "First refresh should succeed");
let cat1_new_title: String = Spi::get_one(
"
SELECT data->>'title' FROM tv_item_by_cat WHERE fk_category = 1
",
)
.unwrap()
.unwrap();
assert_eq!(
cat1_new_title, "Item 1B",
"Category 1 should now show Item 1B"
);
Spi::run("DELETE FROM tb_item WHERE pk_item = 2").unwrap();
let result2 = refresh_row(view_oid, &KeyValue::Int(1));
assert!(result2.is_ok(), "Second refresh should succeed");
let cat1_final_title: String = Spi::get_one(
"
SELECT data->>'title' FROM tv_item_by_cat WHERE fk_category = 1
",
)
.unwrap()
.unwrap();
assert_eq!(
cat1_final_title, "Item 1C",
"Category 1 should now show Item 1C"
);
}
#[pg_test]
fn test_audit_buffer_and_flush() {
Spi::run("SET pg_tviews.audit_enabled = true").unwrap();
crate::audit::log_refresh("user", 5);
crate::audit::log_refresh("post", 3);
crate::audit::log_create("comment", "SELECT ...");
let count: i64 = Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log")
.unwrap()
.unwrap_or(0);
assert_eq!(count, 0, "Buffer should not write to DB before flush");
crate::audit::flush_audit_buffer().unwrap();
let count: i64 = Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log")
.unwrap()
.unwrap_or(0);
assert_eq!(count, 3, "Flush should write all buffered entries");
let refresh_count: i64 =
Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log WHERE operation = 'REFRESH'")
.unwrap()
.unwrap_or(0);
assert_eq!(refresh_count, 2, "Should have 2 REFRESH entries");
let create_count: i64 =
Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log WHERE operation = 'CREATE'")
.unwrap()
.unwrap_or(0);
assert_eq!(create_count, 1, "Should have 1 CREATE entry");
crate::audit::flush_audit_buffer().unwrap();
let count_after: i64 = Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log")
.unwrap()
.unwrap_or(0);
assert_eq!(count_after, 3, "Second flush should be no-op");
}
#[pg_test]
fn test_audit_buffer_clear() {
Spi::run("SET pg_tviews.audit_enabled = true").unwrap();
crate::audit::log_refresh("user", 10);
crate::audit::log_drop("post");
crate::audit::clear_audit_buffer();
crate::audit::flush_audit_buffer().unwrap();
let count: i64 = Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log")
.unwrap()
.unwrap_or(0);
assert_eq!(count, 0, "Cleared buffer should not produce any rows");
}
#[pg_test]
fn test_audit_disabled_skips_flush() {
Spi::run("SET pg_tviews.audit_enabled = false").unwrap();
crate::audit::log_refresh("user", 5);
crate::audit::flush_audit_buffer().unwrap();
let count: i64 = Spi::get_one("SELECT COUNT(*) FROM pg_tview_audit_log")
.unwrap()
.unwrap_or(0);
assert_eq!(count, 0, "Disabled audit should not write any rows");
}
#[pg_test]
fn test_missing_row_deletes_tview_row() {
Spi::run("CREATE TABLE tb_user (pk_user BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
fk_user BIGINT REFERENCES tb_user(pk_user),
title TEXT
)",
)
.unwrap();
Spi::run("INSERT INTO tb_user VALUES (1, 'Alice')").unwrap();
Spi::run("INSERT INTO tb_post VALUES (1, 1, 'Hello')").unwrap();
Spi::run(
"
SELECT pg_tviews_create('post', $$
SELECT pk_post, fk_user,
jsonb_build_object('title', title) AS data
FROM tb_post
$$)
",
)
.unwrap();
let before: i64 = Spi::get_one("SELECT count(*) FROM tv_post WHERE pk_post = 1")
.unwrap()
.unwrap();
assert_eq!(before, 1, "tview row should exist before delete");
Spi::run("DELETE FROM tb_post WHERE pk_post = 1").unwrap();
let post_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_post'::regclass::oid")
.unwrap()
.unwrap();
let result = refresh_row(post_oid, &KeyValue::Int(1));
assert!(
result.is_ok(),
"Refresh of a deleted row should succeed by removing the tview row, got {result:?}"
);
let after: i64 = Spi::get_one("SELECT count(*) FROM tv_post WHERE pk_post = 1")
.unwrap()
.unwrap();
assert_eq!(
after, 0,
"tview row should be gone after refreshing a deleted pk"
);
}
#[pg_test]
fn test_null_data_column_error_handling() {
Spi::run("CREATE TABLE tb_item (pk_item BIGSERIAL PRIMARY KEY, name TEXT)").unwrap();
Spi::run("INSERT INTO tb_item VALUES (1, 'Widget')").unwrap();
Spi::run(
"
SELECT pg_tviews_create('item', $$
SELECT pk_item,
CASE WHEN name = 'Widget' THEN jsonb_build_object('name', name)
ELSE NULL
END AS data
FROM tb_item
$$)
",
)
.unwrap();
let initial_data: Option<String> =
Spi::get_one("SELECT data::text FROM tv_item WHERE pk_item = 1").unwrap();
assert!(initial_data.is_some(), "Should have valid data initially");
Spi::run("UPDATE tb_item SET name = 'Widget-Modified' WHERE pk_item = 1").unwrap();
let item_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'tv_item'::regclass::oid")
.unwrap()
.unwrap();
let result = refresh_row(item_oid, &KeyValue::Int(1));
assert!(
result.is_err(),
"Refresh should fail when data column is NULL"
);
let error_msg = format!("{:?}", result.unwrap_err());
assert!(
error_msg.to_lowercase().contains("null") || error_msg.contains("data"),
"Error should mention NULL or data column issue"
);
}
#[pg_test]
fn test_refresh_distinct_on_text_keys() {
Spi::run(
"CREATE TABLE tb_post (
pk_post BIGSERIAL PRIMARY KEY,
title TEXT
)",
)
.unwrap();
Spi::run(
"INSERT INTO tb_post (pk_post, title) VALUES
(1, 'Post 1'),
(2, 'Post 2')",
)
.unwrap();
Spi::run(
"
SELECT pg_tviews_create('post_by_title', $$
SELECT DISTINCT ON (title)
pk_post,
jsonb_build_object('title', title) AS data
FROM tb_post
ORDER BY title, pk_post
$$, 'title')
",
)
.unwrap();
let view_oid: pgrx::pg_sys::Oid = Spi::get_one("SELECT 'v_post_by_title'::regclass::oid")
.unwrap()
.unwrap();
let result1 = refresh_row(view_oid, &KeyValue::Text("Post 1".into()));
assert!(result1.is_ok(), "Initial refresh should succeed");
let result2 = refresh_row(view_oid, &KeyValue::Text("Post 2".into()));
assert!(result2.is_ok(), "Second refresh should succeed");
}
}