use nodedb_types::DatabaseId;
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;
use super::RETENTION_POLICIES_CRDT_COLLECTION;
use super::parse::parse_create_retention_policy;
fn err(sqlstate: &str, message: String) -> DdlError {
DdlError {
sqlstate: sqlstate.to_string(),
message,
}
}
pub async fn create_retention_policy(
state: &SharedState,
identity: &AuthenticatedIdentity,
database_id: DatabaseId,
name: &str,
collection: &str,
body_raw: &str,
eval_interval_raw: Option<&str>,
) -> Result<Vec<DdlResult>, DdlError> {
require_tenant_admin(identity, "create retention policies")?;
let reconstructed = if let Some(eval) = eval_interval_raw {
format!(
"CREATE RETENTION POLICY {name} ON {collection} ({body_raw}) WITH (EVAL_INTERVAL = '{eval}')"
)
} else {
format!("CREATE RETENTION POLICY {name} ON {collection} ({body_raw})")
};
let parsed = parse_create_retention_policy(&reconstructed)?;
let tenant_id = identity.tenant_id.as_u64();
{
let catalog = state.credentials.catalog();
match catalog.get_collection(database_id, tenant_id, &parsed.collection) {
Ok(Some(coll)) if coll.collection_type.is_timeseries() => {}
Ok(Some(_)) => {
return Err(err(
"42809",
format!("'{}' is not a timeseries collection", parsed.collection),
));
}
_ => {
return Err(err(
"42P01",
format!("collection '{}' does not exist", parsed.collection),
));
}
}
}
if state
.retention_policy_registry
.get(database_id.as_u64(), tenant_id, &parsed.name)
.is_some()
{
return Err(err(
"42710",
format!("retention policy '{}' already exists", parsed.name),
));
}
if state
.retention_policy_registry
.get_for_collection(database_id.as_u64(), tenant_id, &parsed.collection)
.is_some()
{
return Err(err(
"42710",
format!(
"collection '{}' already has a retention policy",
parsed.collection
),
));
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_err(|_| err("XX000", "system clock error".to_string()))?
.as_secs();
let def = crate::engine::timeseries::retention_policy::types::RetentionPolicyDef {
database_id: database_id.as_u64(),
tenant_id,
name: parsed.name.clone(),
collection: parsed.collection.clone(),
tiers: parsed.tiers,
auto_tier: false,
enabled: true,
eval_interval_ms: parsed.eval_interval_ms,
owner: identity.username.clone(),
created_at: now,
};
let catalog = state.credentials.catalog();
catalog
.put_retention_policy(&def)
.map_err(|e| err("XX000", format!("catalog write: {e}")))?;
{
let delta_payload = zerompk::to_msgpack_vec(&def).unwrap_or_default();
let delta = crate::event::crdt_sync::types::OutboundDelta {
collection: RETENTION_POLICIES_CRDT_COLLECTION.into(),
document_id: def.name.clone(),
payload: delta_payload,
op: crate::event::crdt_sync::types::DeltaOp::Upsert,
lsn: 0,
tenant_id,
peer_id: state.node_id,
sequence: 0,
};
state.crdt_sync_delivery.enqueue(tenant_id, delta);
}
state.retention_policy_registry.register(def.clone());
if !def.downsample_tiers().is_empty() {
crate::engine::timeseries::retention_policy::autowire::register_tiers(state, &def)
.await
.map_err(|e| {
state.retention_policy_registry.unregister(
database_id.as_u64(),
tenant_id,
&def.name,
);
let _ = catalog.delete_retention_policy(database_id.as_u64(), tenant_id, &def.name);
err("XX000", format!("failed to auto-wire aggregates: {e}"))
})?;
}
state.audit_record(
crate::control::security::audit::AuditEvent::AdminAction,
Some(identity.tenant_id),
&identity.username,
&format!(
"CREATE RETENTION POLICY {} ON {}",
parsed.name, parsed.collection
),
);
tracing::info!(
name = parsed.name,
collection = parsed.collection,
tiers = parsed.tier_count,
"retention policy created"
);
Ok(vec![DdlResult::Status {
command: "CREATE RETENTION POLICY".to_string(),
rows_affected: None,
}])
}