use tracing::{debug, warn};
use nodedb_cluster::MetadataEntry;
use crate::control::catalog_entry;
use super::audit::emit_ddl_audit;
use super::types::MetadataCommitApplier;
impl MetadataCommitApplier {
pub(super) fn apply_catalog_ddl(
&self,
entry: &MetadataEntry,
raft_index: u64,
) -> Result<(), crate::Error> {
let catalog = self.credentials.catalog();
let (payload, audit) = match entry {
MetadataEntry::CatalogDdl { payload } => (payload, None),
MetadataEntry::CatalogDdlAudited {
payload,
auth_user_id,
auth_user_name,
sql_text,
} => (
payload,
Some((
auth_user_id.clone(),
auth_user_name.clone(),
sql_text.clone(),
)),
),
_ => return Ok(()),
};
let stamped = match catalog_entry::decode(payload) {
Ok(e) => e,
Err(e) => {
warn!(error = %e, "metadata applier: failed to decode CatalogEntry payload");
return Ok(());
}
};
if matches!(
catalog_entry::descriptor_stamp::validate(&stamped, catalog)?,
catalog_entry::descriptor_stamp::ValidationOutcome::AlreadyApplied
) {
debug!(
kind = stamped.kind(),
"catalog_entry: descriptor entry already superseded or applied"
);
return Ok(());
}
debug!(kind = stamped.kind(), "catalog_entry: applying to redb");
if !catalog_entry::apply::apply_to(&stamped, catalog) {
return Ok(());
}
if let Some(weak) = self.shared.get()
&& let Some(shared) = weak.upgrade()
{
if let Some(drained_id) =
crate::control::lease::drain_propose::descriptor_id_for_implicit_clear(&stamped)
{
shared.lease_drain.install_end(&drained_id);
}
catalog_entry::post_apply::apply_post_apply_side_effects_sync(&stamped, &shared);
emit_ddl_audit(&shared, raft_index, &stamped, audit.as_ref());
catalog_entry::post_apply::spawn_post_apply_async_side_effects(
stamped, shared, raft_index,
);
}
Ok(())
}
}