use crate::TViewResult;
use crate::catalog::TviewMeta;
use crate::queue::key::KeyValue;
pub fn refresh_bulk(entity: &str, keys: &[KeyValue]) -> TViewResult<super::Touched> {
if keys.is_empty() {
return Ok(super::Touched::default());
}
crate::metrics::metrics_api::record_view_recomputes(keys.len() as u64);
let meta =
TviewMeta::load_by_entity(entity)?.ok_or_else(|| crate::TViewError::MetadataNotFound {
entity: entity.to_string(),
})?;
let qi_view = crate::utils::qualified_relname_from_oid(meta.view_oid)?;
let qi_tv = crate::utils::qualified_relname_from_oid(meta.tview_oid)?;
let key_col = &meta.identity.column;
let key_type = meta.key_type()?;
let col_names = crate::utils::get_view_columns_by_oid(meta.view_oid)?;
if col_names.is_empty() {
return Ok(super::Touched::default());
}
let col_list = super::column_list(&col_names);
let qi_key = crate::utils::quote_identifier(key_col);
let qi_pk = crate::utils::quote_identifier(&format!("pk_{entity}"));
let any_key = format!("ANY({})", super::key_cast(&key_type, "$1", true));
let source_sql = format!("SELECT {col_list} FROM {qi_view} WHERE {qi_key} = {any_key}");
let conflict = format!(
"ON CONFLICT ({qi_key}) {}",
super::upsert_conflict_action(&qi_tv, &col_names, key_col, None)
);
let delete_sql = format!(
"DELETE FROM {qi_tv} t \
WHERE t.{qi_key} = {any_key} \
AND NOT EXISTS (SELECT 1 FROM {qi_view} v WHERE v.{qi_key} = t.{qi_key}) \
RETURNING t.{qi_pk}::text, to_jsonb(t.*)->>'id'"
);
let with_pks = !meta.identity.is_pk(entity);
let mut touched = super::Touched::default();
for chunk in keys.chunks(crate::config::batch_size()) {
let before = super::lock_rows(&meta, &qi_tv, chunk, with_pks)?;
let (_, written) = super::run_counted_upsert(
entity,
&qi_tv,
&col_list,
&source_sql,
&conflict,
&[super::key_array(&key_type, chunk)?],
)?;
let deleted = super::run_journaled_delete(
entity,
&delete_sql,
&[super::key_array(&key_type, chunk)?],
)?;
touched.extend(super::touched(&meta, chunk, before, written, deleted));
}
Ok(touched)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_refresh_bulk_empty() {
assert!(refresh_bulk("test", &[]).is_ok());
}
}