use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::state::SharedState;
use super::super::super::result::{DdlError, DdlResult};
use super::super::auth_support::{require_tenant_admin, status};
pub fn change_stream_exists(
state: &SharedState,
identity: &AuthenticatedIdentity,
name: &str,
) -> bool {
let tid = identity.tenant_id.as_u64();
state.stream_registry.get(tid, name).is_some()
}
pub fn drop_change_stream(
state: &SharedState,
identity: &AuthenticatedIdentity,
parts: &[&str],
) -> Result<Vec<DdlResult>, DdlError> {
require_tenant_admin(identity, "drop change streams")?;
let (if_exists, name) = if parts.len() >= 6
&& parts[3].eq_ignore_ascii_case("IF")
&& parts[4].eq_ignore_ascii_case("EXISTS")
{
(true, parts[5].to_lowercase())
} else if parts.len() >= 4 {
(false, parts[3].to_lowercase())
} else {
return Err(DdlError {
sqlstate: "42601".to_string(),
message: "expected DROP CHANGE STREAM [IF EXISTS] <name>".to_string(),
});
};
let tenant_id = identity.tenant_id.as_u64();
let catalog = state.credentials.catalog();
let existed_before = catalog
.get_change_stream(tenant_id, &name)
.map(|opt| opt.is_some())
.unwrap_or(false);
if !existed_before && !if_exists {
return Err(DdlError {
sqlstate: "42704".to_string(),
message: format!("change stream '{name}' does not exist"),
});
}
if !existed_before {
return Ok(status("DROP CHANGE STREAM"));
}
let entry = crate::control::catalog_entry::CatalogEntry::DeleteChangeStream {
tenant_id,
name: name.clone(),
};
let log_index = crate::control::metadata_proposer::propose_catalog_entry(state, &entry)
.map_err(|e| DdlError {
sqlstate: "XX000".to_string(),
message: format!("metadata propose: {e}"),
})?;
if log_index == 0 {
let _ = catalog
.delete_change_stream(tenant_id, &name)
.map_err(|e| DdlError {
sqlstate: "XX000".to_string(),
message: format!("catalog delete: {e}"),
})?;
state.stream_registry.unregister(tenant_id, &name);
state.cdc_router.remove_buffer(tenant_id, &name);
}
state.webhook_manager.stop_task(tenant_id, &name);
state.audit_record(
crate::control::security::audit::AuditEvent::AdminAction,
Some(identity.tenant_id),
&identity.username,
&format!("DROP CHANGE STREAM {name}"),
);
Ok(status("DROP CHANGE STREAM"))
}