use crate::TViewResult;
use crate::catalog::TviewMeta;
use crate::utils::lookup_view_for_source;
use pgrx::datum::DatumWithOid;
use pgrx::prelude::*;
pub fn refresh_bulk(entity: &str, pks: &[i64]) -> TViewResult<()> {
if pks.is_empty() {
return Ok(());
}
let meta =
TviewMeta::load_by_entity(entity)?.ok_or_else(|| crate::TViewError::MetadataNotFound {
entity: entity.to_string(),
})?;
let view_name = lookup_view_for_source(meta.view_oid)?;
let tv_name = crate::utils::relname_from_oid(meta.tview_oid)?;
let pk_col = format!("pk_{entity}");
let col_names = crate::utils::get_view_columns_by_oid(meta.view_oid)?;
if col_names.is_empty() {
return Ok(());
}
let col_list = col_names.join(", ");
let do_update: String = {
let mut parts = Vec::with_capacity(col_names.len());
for c in &col_names {
if c.as_str() != pk_col.as_str() {
parts.push(format!("{c} = EXCLUDED.{c}"));
}
}
parts.push("updated_at = NOW()".to_string());
parts.join(", ")
};
let upsert_sql = format!(
"INSERT INTO {tv_name} ({col_list}) \
SELECT {col_list} FROM {view_name} WHERE {pk_col} = ANY($1) \
ON CONFLICT ({pk_col}) DO UPDATE SET {do_update}"
);
let delete_sql = format!(
"DELETE FROM {tv_name} t \
WHERE t.{pk_col} = ANY($1) \
AND NOT EXISTS (SELECT 1 FROM {view_name} v WHERE v.{pk_col} = t.{pk_col})"
);
let batch = crate::config::batch_size();
for chunk in pks.chunks(batch) {
Spi::run_with_args(
&upsert_sql,
&[unsafe {
DatumWithOid::new(
chunk.to_vec(),
PgOid::BuiltIn(PgBuiltInOids::INT8ARRAYOID).value(),
)
}],
)?;
Spi::run_with_args(
&delete_sql,
&[unsafe {
DatumWithOid::new(
chunk.to_vec(),
PgOid::BuiltIn(PgBuiltInOids::INT8ARRAYOID).value(),
)
}],
)?;
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_refresh_bulk_empty() {
assert!(refresh_bulk("test", &[]).is_ok());
}
}