use nodedb_types::DatabaseId;
use crate::control::security::audit::AuditEvent;
use crate::control::security::catalog::StoredCollection;
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::state::SharedState;
use super::super::super::super::catalog::propose_and_apply;
use super::super::super::super::result::{DdlError, DdlResult};
use super::super::enforcement::{parse_balanced_clause_from_raw, resolve_custom_type_columns};
use super::engine_option::validate_engine_name;
use super::request::CreateCollectionRequest;
fn err(sqlstate: &str, message: String) -> DdlError {
DdlError {
sqlstate: sqlstate.to_string(),
message,
}
}
fn parse_crdt_flag(value: &str) -> Result<bool, DdlError> {
match value.trim() {
v if v.eq_ignore_ascii_case("true") => Ok(true),
v if v.eq_ignore_ascii_case("false") => Ok(false),
other => Err(err(
"42601",
format!("invalid value for WITH (crdt=...): '{other}'; expected 'true' or 'false'"),
)),
}
}
fn resolve_crdt_flag(
options: &[(String, String)],
collection_type: &nodedb_types::CollectionType,
) -> Result<bool, DdlError> {
let crdt = match options.iter().find(|(k, _)| k.eq_ignore_ascii_case("crdt")) {
Some((_, v)) => parse_crdt_flag(v)?,
None => false,
};
if crdt && !matches!(collection_type, nodedb_types::CollectionType::Document(_)) {
return Err(err(
"42601",
"WITH (crdt=true) is only supported on document collections".to_string(),
));
}
Ok(crdt)
}
pub struct Variant {
pub label: &'static str,
pub response_tag: &'static str,
pub require_columns: bool,
pub default_strict: bool,
}
pub async fn build_and_persist(
state: &SharedState,
identity: &AuthenticatedIdentity,
req: &CreateCollectionRequest<'_>,
database_id: DatabaseId,
variant: &Variant,
) -> Result<Vec<DdlResult>, DdlError> {
let CreateCollectionRequest {
name,
engine,
columns,
options,
flags,
balanced_raw,
} = *req;
validate_name(name, variant.label)?;
if variant.require_columns && columns.is_empty() {
return Err(err(
"42601",
"CREATE TABLE requires a column list; for schemaless collections use CREATE COLLECTION"
.to_string(),
));
}
let tenant_id = identity.tenant_id;
let mut local_lifecycle = if state.metadata_raft.get().is_none() {
Some(
state
.quiesce
.acquire_lifecycle(database_id.as_u64(), tenant_id.as_u64(), name)
.await,
)
} else {
None
};
let catalog = state.credentials.catalog();
if catalog
.get_materialized_view(tenant_id.as_u64(), name)
.map_err(|error| err("XX000", error.to_string()))?
.is_some()
{
return Err(err(
"42P07",
format!("materialized view '{name}' already owns this collection name"),
));
}
let existing = catalog
.get_collection(database_id, tenant_id.as_u64(), name)
.map_err(|error| err("XX000", error.to_string()))?;
if let Some(existing) = existing {
if existing.is_active {
return Err(err(
"42P07",
format!("{} '{name}' already exists", variant.label),
));
}
let purge_lsn = state.wal.next_lsn().as_u64();
let purge_result =
crate::control::server::shared::ddl::neutral::collection::purge::hard_purge_collection(
state,
database_id.as_u64(),
tenant_id.as_u64(),
name,
purge_lsn,
local_lifecycle.is_some(),
)
.await;
if let Err(failure) = purge_result {
if failure.retry_queued
&& let Some(guard) = local_lifecycle.take()
{
guard.disarm();
}
return Err(err("XX000", failure.error.to_string()));
}
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let canonical_engine = validate_engine_name(engine, options)?;
let bitemporal_flag = flags.iter().any(|f| f == "BITEMPORAL");
let resolved_columns: Vec<(String, String)> =
resolve_custom_type_columns(columns, state, tenant_id.as_u64());
let (collection_type, columnar_schema_columns) = nodedb_sql::ddl_ast::build_collection_type(
canonical_engine,
&resolved_columns,
options,
bitemporal_flag,
variant.default_strict,
)
.map_err(|e| err("42601", e.to_string()))?;
let (mut fields, serial_fields) =
crate::control::server::shared::ddl::schema_validation::parse_fields_clause_from_pairs(
columns,
);
if fields.is_empty() && !columnar_schema_columns.is_empty() {
fields = columnar_schema_columns;
}
let schema_json = match &collection_type {
nodedb_types::CollectionType::Document(nodedb_types::DocumentMode::Strict(schema)) => {
sonic_rs::to_string(schema).ok()
}
nodedb_types::CollectionType::KeyValue(config) => sonic_rs::to_string(config).ok(),
_ => None,
};
let (primary, vector_primary) =
resolve_primary_engine(options, columns, &fields, &collection_type)?;
let append_only = flags.iter().any(|f| f == "APPEND_ONLY");
let hash_chain = flags.iter().any(|f| f == "HASH_CHAIN");
let bitemporal = bitemporal_flag;
if hash_chain && !append_only {
return Err(err("42601", "HASH_CHAIN requires APPEND_ONLY".to_string()));
}
let crdt = resolve_crdt_flag(options, &collection_type)?;
let balanced =
parse_balanced_clause_from_raw(balanced_raw.unwrap_or("")).map_err(|e| err("42601", e))?;
let partition_strategy =
nodedb_types::PartitionStrategy::default_for_collection_type(&collection_type);
let declared_primary_key = columns.iter().find_map(|(col_name, type_str)| {
let (_, is_pk, _, _) =
nodedb_sql::ddl_ast::collection_type::parse_column_type_str_full(type_str);
is_pk.then(|| col_name.clone())
});
let coll = StoredCollection {
tenant_id: tenant_id.as_u64(),
name: name.to_string(),
owner: identity.username.clone(),
created_at: now,
descriptor_version: 0,
constraint_version: 0,
modification_hlc: nodedb_types::Hlc::ZERO,
fields,
field_defs: Vec::new(),
event_defs: Vec::new(),
collection_type,
timeseries_config: schema_json,
conflict_policy: None,
is_active: true,
append_only,
hash_chain,
balanced,
last_chain_hash: None,
period_lock: None,
retention_period: None,
legal_holds: Vec::new(),
state_constraints: Vec::new(),
transition_checks: Vec::new(),
type_guards: Vec::new(),
check_constraints: Vec::new(),
materialized_sums: Vec::new(),
lvc_enabled: false,
bitemporal,
crdt,
permission_tree_def: None,
indexes: Vec::new(),
size_bytes_estimate: 0,
primary,
vector_primary,
partition_strategy,
database_id,
cloned_from: None,
clone_status: nodedb_types::CloneStatus::default(),
has_implicit_edges: false,
declared_primary_key,
};
let entry = crate::control::catalog_entry::CatalogEntry::PutCollection(Box::new(coll.clone()));
propose_and_apply(state, &entry)?;
log_vector_fields(name, &coll.fields);
create_serial_sequences(state, identity, name, &serial_fields, now)?;
state.audit_record(
AuditEvent::AdminAction,
Some(tenant_id),
&identity.username,
&format!("created {} '{name}'", variant.label),
);
Ok(vec![DdlResult::Status {
command: variant.response_tag.to_string(),
rows_affected: None,
}])
}
fn validate_name(name: &str, label: &str) -> Result<(), DdlError> {
if name.is_empty()
|| !name
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
{
return Err(err(
"42601",
format!(
"invalid {label} name '{name}': only letters, digits, '-', and '_' are allowed"
),
));
}
Ok(())
}
fn resolve_primary_engine(
options: &[(String, String)],
columns: &[(String, String)],
fields: &[(String, String)],
collection_type: &nodedb_types::CollectionType,
) -> Result<
(
nodedb_types::PrimaryEngine,
Option<nodedb_types::VectorPrimaryConfig>,
),
DdlError,
> {
match nodedb_sql::ddl_ast::parse::vector_primary::parse_vector_primary_options_from_kvs(options)
{
Ok(Some(mut vp_cfg)) => {
let col_list: Vec<(String, String)> = if fields.is_empty() {
columns.to_vec()
} else {
fields.to_vec()
};
nodedb_sql::ddl_ast::parse::vector_primary::validate_vector_field(&vp_cfg, &col_list)
.map_err(|e| err("42601", e.to_string()))?;
nodedb_sql::ddl_ast::parse::vector_primary::validate_payload_indexes(
&mut vp_cfg,
&col_list,
)
.map_err(|e| err("42601", e.to_string()))?;
if let Some((_, type_str)) = col_list
.iter()
.find(|(n, _)| n.eq_ignore_ascii_case(&vp_cfg.vector_field))
{
let upper_t = type_str.to_uppercase();
if let Some(inner) = upper_t
.strip_prefix("VECTOR(")
.and_then(|s| s.strip_suffix(')'))
&& let Ok(d) = inner.trim().parse::<u32>()
{
if vp_cfg.dim == 0 {
vp_cfg.dim = d;
} else if vp_cfg.dim != d {
return Err(err(
"42601",
format!(
"vector dim mismatch: WITH clause specifies {}, column type VECTOR({}) specifies {}",
vp_cfg.dim, d, d
),
));
}
}
}
Ok((nodedb_types::PrimaryEngine::Vector, Some(vp_cfg)))
}
Ok(None) => Ok((
nodedb_types::PrimaryEngine::infer_from_collection_type(collection_type),
None,
)),
Err(e) => Err(err("42601", e.to_string())),
}
}
fn log_vector_fields(collection_name: &str, fields: &[(String, String)]) {
let vector_fields =
crate::control::server::shared::ddl::schema_validation::extract_vector_fields(fields);
for (field_name, _dim, metric) in &vector_fields {
tracing::info!(
name = %collection_name,
field = %field_name,
%metric,
"auto-configuring vector field"
);
}
}
fn create_serial_sequences(
state: &SharedState,
identity: &AuthenticatedIdentity,
collection_name: &str,
serial_fields: &[String],
now: u64,
) -> Result<(), DdlError> {
for field_name in serial_fields {
let seq_name = format!("{collection_name}_{field_name}_seq");
let mut seq_def = crate::control::security::catalog::sequence_types::StoredSequence::new(
identity.tenant_id.as_u64(),
seq_name.clone(),
identity.username.clone(),
);
seq_def.created_at = now;
let seq_entry =
crate::control::catalog_entry::CatalogEntry::PutSequence(Box::new(seq_def.clone()));
propose_and_apply(state, &seq_entry)?;
let _ = state.sequence_registry.create(seq_def);
tracing::info!(
collection = %collection_name,
field = %field_name,
sequence = %seq_name,
"auto-created SERIAL sequence"
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::resolve_crdt_flag;
fn opts(pairs: &[(&str, &str)]) -> Vec<(String, String)> {
pairs
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect()
}
#[test]
fn crdt_true_on_document_collection_resolves_true() {
let options = opts(&[("crdt", "true")]);
let flag = resolve_crdt_flag(&options, &nodedb_types::CollectionType::document())
.expect("crdt=true on a document collection must resolve");
assert!(flag);
}
#[test]
fn crdt_true_on_non_document_collection_rejected() {
let options = opts(&[("crdt", "true")]);
let err = resolve_crdt_flag(&options, &nodedb_types::CollectionType::columnar())
.expect_err("crdt=true on a non-document collection must be rejected");
assert_eq!(err.sqlstate, "42601");
}
#[test]
fn crdt_garbage_value_rejected() {
let options = opts(&[("crdt", "maybe")]);
let err = resolve_crdt_flag(&options, &nodedb_types::CollectionType::document())
.expect_err("a non-boolean crdt value must be rejected");
assert_eq!(err.sqlstate, "42601");
}
#[test]
fn crdt_absent_defaults_false() {
let options = opts(&[("engine", "kv")]);
let flag = resolve_crdt_flag(&options, &nodedb_types::CollectionType::document())
.expect("absent crdt option must resolve to a default");
assert!(!flag);
}
fn validate_name(name: &str) -> bool {
!name.is_empty()
&& name
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
}
#[test]
fn valid_collection_names() {
assert!(validate_name("docs"));
assert!(validate_name("my_collection"));
assert!(validate_name("my-collection"));
assert!(validate_name("Collection123"));
assert!(validate_name("a"));
}
#[test]
fn invalid_collection_names_rejected() {
assert!(!validate_name("docs;"));
assert!(!validate_name("bad;name"));
assert!(!validate_name("bad name"));
assert!(!validate_name("bad.name"));
assert!(!validate_name("bad/name"));
assert!(!validate_name(""));
assert!(!validate_name("events;"));
assert!(!validate_name("orders;"));
}
}