1use super::*;
9use crate::catalog::CollectionModel;
10use crate::runtime::audit_log::{AuditAuthSource, AuditEvent, AuditFieldEscaper, Outcome};
11use crate::runtime::ddl::polymorphic_resolver;
12use crate::storage::query::ast::{
13 CreateVcsRefQuery, DropForkQuery, DropVcsRefQuery, ForkStoreQuery, VcsRefKind,
14};
15use crate::storage::query::{analyze_create_table, resolve_declared_data_type, CreateColumnDef};
16use std::collections::{BTreeSet, HashMap, HashSet};
17
18fn operational_wal_retention_floor(durable_lsn: u64) -> u64 {
19 if durable_lsn == 0 {
20 0
21 } else {
22 1
23 }
24}
25
26fn validate_store_fork_lsn(
27 fork_lsn: u64,
28 retention_floor: u64,
29 durable_lsn: u64,
30) -> RedDBResult<()> {
31 if fork_lsn < retention_floor {
32 return Err(RedDBError::Query(format!(
33 "cannot fork store at LSN {fork_lsn}: below operational WAL retention floor \
34 {retention_floor}; restore from backup for older recovery points"
35 )));
36 }
37 if fork_lsn > durable_lsn {
38 return Err(RedDBError::Query(format!(
39 "cannot fork store at LSN {fork_lsn}: current durable LSN is {durable_lsn}; \
40 restore from backup for recovery points outside the retained WAL window"
41 )));
42 }
43 Ok(())
44}
45
46fn vault_master_key_ref(collection: &str) -> String {
47 format!("red.vault.{collection}.master_key")
48}
49
50impl RedDBRuntime {
51 pub fn execute_create_table(
57 &self,
58 raw_query: &str,
59 query: &CreateTableQuery,
60 ) -> RedDBResult<RuntimeQueryResult> {
61 if query.collection_model != CollectionModel::Table {
62 return self.execute_create_keyed_collection(raw_query, query);
63 }
64 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
65 let store = self.inner.db.store();
66 analyze_create_table(query).map_err(|err| RedDBError::Query(err.to_string()))?;
67 if let Some(policy) = &query.ai_policy {
71 validate_ai_policy_modalities(policy)?;
72 }
73 crate::reserved_fields::ensure_no_reserved_public_item_fields(
74 query.columns.iter().map(|column| column.name.as_str()),
75 &format!("table '{}'", query.name),
76 )?;
77 let exists = store.get_collection(&query.name).is_some();
79 if exists {
80 if query.if_not_exists {
81 return Ok(RuntimeQueryResult::ok_message(
82 raw_query.to_string(),
83 &format!("table '{}' already exists", query.name),
84 "create",
85 ));
86 }
87 return Err(RedDBError::Query(format!(
88 "table '{}' already exists",
89 query.name
90 )));
91 }
92
93 let contract = collection_contract_from_create_table(query)?;
96 validate_event_subscriptions(self, &query.name, &contract.subscriptions)?;
97 store
99 .create_collection(&query.name)
100 .map_err(|err| RedDBError::Internal(err.to_string()))?;
101 for subscription in &contract.subscriptions {
102 ensure_event_target_queue(self, &subscription.target_queue)?;
103 }
104 if let Some(default_ttl_ms) = query.default_ttl_ms {
105 self.inner
106 .db
107 .set_collection_default_ttl_ms(&query.name, default_ttl_ms);
108 }
109 self.inner
110 .db
111 .save_collection_contract(contract)
112 .map_err(|err| RedDBError::Internal(err.to_string()))?;
113 if let Some(tenant_id) = crate::runtime::impl_core::current_tenant() {
114 store.set_config_tree(
115 &format!("red.collection_tenants.{}", query.name),
116 &crate::serde_json::Value::String(tenant_id),
117 );
118 }
119 self.inner
120 .db
121 .persist_metadata()
122 .map_err(|err| RedDBError::Internal(err.to_string()))?;
123 self.refresh_table_planner_stats(&query.name);
124 self.invalidate_result_cache();
125 let columns: Vec<String> = query.columns.iter().map(|col| col.name.clone()).collect();
128 self.schema_vocabulary_apply(
129 crate::runtime::schema_vocabulary::DdlEvent::CreateCollection {
130 collection: query.name.clone(),
131 columns,
132 type_tags: Vec::new(),
133 description: None,
134 },
135 );
136 if let Some(spec) = &query.partition_by {
143 let kind_str = match spec.kind {
144 crate::storage::query::ast::PartitionKind::Range => "range",
145 crate::storage::query::ast::PartitionKind::List => "list",
146 crate::storage::query::ast::PartitionKind::Hash => "hash",
147 };
148 store.set_config_tree(
149 &format!("partition.{}.by", query.name),
150 &crate::serde_json::Value::String(kind_str.to_string()),
151 );
152 store.set_config_tree(
153 &format!("partition.{}.column", query.name),
154 &crate::serde_json::Value::String(spec.column.clone()),
155 );
156 }
157
158 if let Some(col) = &query.tenant_by {
169 store.set_config_tree(
170 &format!("tenant_tables.{}.column", query.name),
171 &crate::serde_json::Value::String(col.clone()),
172 );
173 self.register_tenant_table(&query.name, col);
174 }
175
176 let ttl_suffix = query
177 .default_ttl_ms
178 .map(|ttl_ms| format!(" with default TTL {}ms", ttl_ms))
179 .unwrap_or_default();
180
181 let tenant_suffix = query
182 .tenant_by
183 .as_ref()
184 .map(|col| format!(" (tenant-scoped by {col})"))
185 .unwrap_or_default();
186
187 Ok(RuntimeQueryResult::ok_message(
188 raw_query.to_string(),
189 &format!(
190 "table '{}' created{}{}",
191 query.name, ttl_suffix, tenant_suffix
192 ),
193 "create",
194 ))
195 }
196
197 fn execute_create_keyed_collection(
198 &self,
199 raw_query: &str,
200 query: &CreateTableQuery,
201 ) -> RedDBResult<RuntimeQueryResult> {
202 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
203 if is_system_schema_name(&query.name) {
204 return Err(RedDBError::Query("system schema is read-only".to_string()));
205 }
206 let store = self.inner.db.store();
207 let label = polymorphic_resolver::model_name(query.collection_model);
208 if store.get_collection(&query.name).is_some() {
209 if query.if_not_exists {
210 return Ok(RuntimeQueryResult::ok_message(
211 raw_query.to_string(),
212 &format!("{label} '{}' already exists", query.name),
213 "create",
214 ));
215 }
216 return Err(RedDBError::Query(format!(
217 "{label} '{}' already exists",
218 query.name
219 )));
220 }
221
222 store
223 .create_collection(&query.name)
224 .map_err(|err| RedDBError::Internal(err.to_string()))?;
225 if query.collection_model == CollectionModel::Vault {
226 self.provision_vault_key_material(&query.name, query.vault_own_master_key)?;
227 let key_scope = if query.vault_own_master_key {
228 "own"
229 } else {
230 "cluster"
231 };
232 store.set_config_tree(
233 &format!("red.vault.{}.key_scope", query.name),
234 &crate::serde_json::Value::String(key_scope.to_string()),
235 );
236 store.set_config_tree(
237 &format!("red.vault.{}.status", query.name),
238 &crate::serde_json::Value::String("sealed".to_string()),
239 );
240 }
241 if query.collection_model == CollectionModel::Metrics {
242 for spec in &query.metrics_rollup_policies {
243 let policy = crate::storage::timeseries::retention::DownsamplePolicy::parse(spec)
244 .ok_or_else(|| {
245 RedDBError::Query(format!("invalid metrics rollup policy '{}'", spec))
246 })?;
247 if policy.source != "raw" {
248 return Err(RedDBError::Query(format!(
249 "invalid metrics rollup policy '{}': metrics v0 rollups must use raw as source",
250 spec
251 )));
252 }
253 if !matches!(
254 policy.aggregation.as_str(),
255 "avg" | "sum" | "min" | "max" | "count"
256 ) {
257 return Err(RedDBError::Query(format!(
258 "invalid metrics rollup policy '{}': supported aggregations are avg, sum, min, max, count",
259 spec
260 )));
261 }
262 }
263 if let Some(raw_retention_ms) = query.default_ttl_ms {
264 self.inner
265 .db
266 .set_collection_default_ttl_ms(&query.name, raw_retention_ms);
267 store.set_config_tree(
268 &format!("red.metrics.{}.raw_retention_ms", query.name),
269 &crate::serde_json::Value::Number(raw_retention_ms as f64),
270 );
271 }
272 let tenant_identity = query
273 .tenant_by
274 .clone()
275 .unwrap_or_else(|| "current_tenant".to_string());
276 store.set_config_tree(
277 &format!("red.metrics.{}.tenant_identity", query.name),
278 &crate::serde_json::Value::String(tenant_identity),
279 );
280 store.set_config_tree(
281 &format!("red.metrics.{}.namespace", query.name),
282 &crate::serde_json::Value::String("default".to_string()),
283 );
284 if !query.metrics_rollup_policies.is_empty() {
285 store.set_config_tree(
286 &format!("red.metrics.{}.rollup_policies", query.name),
287 &crate::serde_json::Value::Array(
288 query
289 .metrics_rollup_policies
290 .iter()
291 .cloned()
292 .map(crate::serde_json::Value::String)
293 .collect(),
294 ),
295 );
296 }
297 }
298 let contract = if query.collection_model == CollectionModel::Metrics {
299 metrics_collection_contract(query)
300 } else {
301 keyed_collection_contract(
302 &query.name,
303 query.collection_model,
304 query.analytics_config.clone(),
305 )
306 };
307 self.inner
308 .db
309 .save_collection_contract(contract)
310 .map_err(|err| RedDBError::Internal(err.to_string()))?;
311 if let Some(tenant_id) = crate::runtime::impl_core::current_tenant() {
312 store.set_config_tree(
313 &format!("red.collection_tenants.{}", query.name),
314 &crate::serde_json::Value::String(tenant_id),
315 );
316 }
317 self.inner
318 .db
319 .persist_metadata()
320 .map_err(|err| RedDBError::Internal(err.to_string()))?;
321 self.invalidate_result_cache();
322
323 Ok(RuntimeQueryResult::ok_message(
324 raw_query.to_string(),
325 &format!("{label} '{}' created", query.name),
326 "create",
327 ))
328 }
329
330 pub fn ensure_system_graph_with_analytics(
339 &self,
340 name: &str,
341 outputs: &[crate::catalog::AnalyticsOutput],
342 ) -> RedDBResult<()> {
343 let store = self.inner.db.store();
344 if store.get_collection(name).is_some() {
345 return Ok(());
346 }
347 let analytics_config = outputs
348 .iter()
349 .map(|output| crate::catalog::AnalyticsViewDescriptor {
350 output: *output,
351 algorithm: None,
352 resolution: None,
353 max_iterations: None,
354 tolerance: None,
355 })
356 .collect();
357 store
358 .create_collection(name)
359 .map_err(|err| RedDBError::Internal(err.to_string()))?;
360 let contract = keyed_collection_contract(
361 name,
362 crate::catalog::CollectionModel::Graph,
363 analytics_config,
364 );
365 self.inner
366 .db
367 .save_collection_contract(contract)
368 .map_err(|err| RedDBError::Internal(err.to_string()))?;
369 self.inner
370 .db
371 .persist_metadata()
372 .map_err(|err| RedDBError::Internal(err.to_string()))?;
373 self.invalidate_result_cache_process_only();
374 Ok(())
375 }
376
377 pub fn execute_create_collection(
378 &self,
379 raw_query: &str,
380 query: &CreateCollectionQuery,
381 ) -> RedDBResult<RuntimeQueryResult> {
382 let model = match query.kind.as_str() {
383 "graph" => CollectionModel::Graph,
384 "document" => CollectionModel::Document,
385 "metrics" => CollectionModel::Metrics,
386 "vector.turbo" => {
387 let dimension = query.vector_dimension.ok_or_else(|| {
388 RedDBError::Query(
389 "CREATE COLLECTION KIND vector.turbo requires DIM".to_string(),
390 )
391 })?;
392 let create = CreateVectorQuery {
393 name: query.name.clone(),
394 dimension,
395 metric: query
396 .vector_metric
397 .unwrap_or(crate::storage::engine::distance::DistanceMetric::Cosine),
398 if_not_exists: query.if_not_exists,
399 };
400 let result = self.execute_create_vector(raw_query, &create)?;
401 let store = self.inner.db.store();
408 crate::runtime::vector_turbo_kind::mark_as_turbo(&store, &query.name);
409 self.inner
410 .db
411 .persist_metadata()
412 .map_err(|err| RedDBError::Internal(err.to_string()))?;
413 let _ = self.inner.db.turbo_state(&query.name);
418 return Ok(result);
419 }
420 "blockchain" => CollectionModel::Table,
425 other => {
426 return Err(RedDBError::Query(format!(
427 "NOT_YET_SUPPORTED: CREATE COLLECTION KIND {other} is not implemented"
428 )));
429 }
430 };
431 let create = CreateTableQuery {
432 collection_model: model,
433 name: query.name.clone(),
434 columns: Vec::new(),
435 if_not_exists: query.if_not_exists,
436 default_ttl_ms: None,
437 metrics_rollup_policies: Vec::new(),
438 context_index_fields: Vec::new(),
439 context_index_enabled: false,
440 timestamps: false,
441 partition_by: None,
442 tenant_by: None,
443 append_only: false,
444 subscriptions: Vec::new(),
445 analytics_config: Vec::new(),
446 vault_own_master_key: false,
447 ai_policy: None,
448 };
449 let result = self.execute_create_table(raw_query, &create)?;
450 if query.kind == "blockchain" {
451 self.install_blockchain_kind(&query.name)?;
452 }
453 if !query.allowed_signers.is_empty() {
458 let actor = crate::runtime::impl_core::current_user_projected()
459 .unwrap_or_else(|| "@system/create-collection".to_string());
460 crate::runtime::signed_writes_kind::install(
461 &self.inner.db.store(),
462 &query.name,
463 &query.allowed_signers,
464 &actor,
465 );
466 }
467 Ok(result)
468 }
469
470 fn install_blockchain_kind(&self, name: &str) -> RedDBResult<()> {
474 use crate::runtime::blockchain_kind;
475 use crate::storage::unified::{EntityData, EntityId, EntityKind, RowData, UnifiedEntity};
476 use std::sync::Arc;
477
478 let store = self.inner.db.store();
479 blockchain_kind::mark_as_chain(&store, name);
480
481 let existing_tip = blockchain_kind::chain_tip(&store, name);
482 if existing_tip.height.is_some() {
483 return Ok(());
484 }
485
486 let fields = blockchain_kind::genesis_fields(blockchain_kind::now_ms());
487 let named: std::collections::HashMap<String, crate::storage::schema::Value> =
488 fields.into_iter().collect();
489 let entity = UnifiedEntity::new(
490 EntityId::new(0),
491 EntityKind::TableRow {
492 table: Arc::from(name),
493 row_id: 0,
494 },
495 EntityData::Row(RowData {
496 columns: Vec::new(),
497 named: Some(named),
498 schema: None,
499 }),
500 );
501 store
502 .insert_auto(name, entity)
503 .map_err(|err| RedDBError::Internal(err.to_string()))?;
504 if let Some(tip) = blockchain_kind::chain_tip_full(&store, name) {
507 self.inner
508 .chain_tip_cache
509 .lock()
510 .insert(name.to_string(), tip);
511 }
512 Ok(())
513 }
514
515 pub fn execute_create_vector(
516 &self,
517 raw_query: &str,
518 query: &CreateVectorQuery,
519 ) -> RedDBResult<RuntimeQueryResult> {
520 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
521 if is_system_schema_name(&query.name) {
522 return Err(RedDBError::Query("system schema is read-only".to_string()));
523 }
524 let store = self.inner.db.store();
525 if store.get_collection(&query.name).is_some() {
526 if query.if_not_exists {
527 return Ok(RuntimeQueryResult::ok_message(
528 raw_query.to_string(),
529 &format!("vector '{}' already exists", query.name),
530 "create",
531 ));
532 }
533 return Err(RedDBError::Query(format!(
534 "vector '{}' already exists",
535 query.name
536 )));
537 }
538
539 store
540 .create_collection(&query.name)
541 .map_err(|err| RedDBError::Internal(err.to_string()))?;
542 self.inner
543 .db
544 .save_collection_contract(vector_collection_contract(query))
545 .map_err(|err| RedDBError::Internal(err.to_string()))?;
546 if let Some(tenant_id) = crate::runtime::impl_core::current_tenant() {
547 store.set_config_tree(
548 &format!("red.collection_tenants.{}", query.name),
549 &crate::serde_json::Value::String(tenant_id),
550 );
551 }
552 crate::runtime::vector_turbo_kind::mark_as_turbo(&store, &query.name);
561 self.inner
562 .db
563 .persist_metadata()
564 .map_err(|err| RedDBError::Internal(err.to_string()))?;
565 let _ = self.inner.db.turbo_state(&query.name);
570 self.invalidate_result_cache();
571
572 Ok(RuntimeQueryResult::ok_message(
573 raw_query.to_string(),
574 &format!("vector '{}' created", query.name),
575 "create",
576 ))
577 }
578
579 fn provision_vault_key_material(
580 &self,
581 collection: &str,
582 own_master_key: bool,
583 ) -> RedDBResult<()> {
584 let auth_store = self.inner.auth_store.read().clone().ok_or_else(|| {
585 RedDBError::Query("CREATE VAULT requires an enabled, unsealed vault".to_string())
586 })?;
587 if !auth_store.is_vault_backed() {
588 return Err(RedDBError::Query(
589 "CREATE VAULT requires an enabled, unsealed vault".to_string(),
590 ));
591 }
592
593 if auth_store.vault_secret_key().is_none() {
594 let key = crate::auth::store::random_bytes(32);
595 auth_store
596 .vault_kv_try_set("red.secret.aes_key".to_string(), hex::encode(key))
597 .map_err(|err| RedDBError::Query(err.to_string()))?;
598 }
599
600 if own_master_key {
601 let key = crate::auth::store::random_bytes(32);
602 auth_store
603 .vault_kv_try_set(vault_master_key_ref(collection), hex::encode(key))
604 .map_err(|err| RedDBError::Query(err.to_string()))?;
605 }
606
607 Ok(())
608 }
609
610 pub fn execute_drop_table(
614 &self,
615 raw_query: &str,
616 query: &DropTableQuery,
617 ) -> RedDBResult<RuntimeQueryResult> {
618 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
619 let store = self.inner.db.store();
620
621 if is_system_schema_name(&query.name) {
622 return Err(RedDBError::Query("system schema is read-only".to_string()));
623 }
624
625 let exists = store.get_collection(&query.name).is_some();
626 if !exists {
627 if query.if_exists {
628 return Ok(RuntimeQueryResult::ok_message(
629 raw_query.to_string(),
630 &format!("table '{}' does not exist", query.name),
631 "drop",
632 ));
633 }
634 return Err(RedDBError::NotFound(format!(
635 "table '{}' not found",
636 query.name
637 )));
638 }
639 let actual =
640 polymorphic_resolver::resolve(&query.name, &self.inner.db.catalog_model_snapshot())?;
641 polymorphic_resolver::ensure_model_match(CollectionModel::Table, actual)?;
642
643 let final_count = store
646 .get_collection(&query.name)
647 .map(|manager| manager.query_all(|_| true).len() as u64)
648 .unwrap_or(0);
649 crate::runtime::mutation::emit_collection_dropped_event_for_collection(
650 self,
651 &query.name,
652 final_count,
653 )?;
654
655 let orphaned_indices: Vec<String> = self
656 .inner
657 .index_store
658 .list_indices(&query.name)
659 .into_iter()
660 .map(|index| index.name)
661 .collect();
662 for name in &orphaned_indices {
663 self.inner.index_store.drop_index(name, &query.name);
664 }
665
666 store
667 .drop_collection(&query.name)
668 .map_err(|err| RedDBError::Internal(err.to_string()))?;
669 self.inner.db.invalidate_vector_index(&query.name);
670 self.inner.db.clear_collection_default_ttl_ms(&query.name);
671 self.inner
672 .db
673 .remove_collection_contract(&query.name)
674 .map_err(|err| RedDBError::Internal(err.to_string()))?;
675 self.clear_table_planner_stats(&query.name);
676 self.invalidate_result_cache();
677 if let Some(store) = self.inner.auth_store.read().clone() {
681 store.invalidate_visible_collections_cache();
682 }
683 self.inner
684 .db
685 .persist_metadata()
686 .map_err(|err| RedDBError::Internal(err.to_string()))?;
687 self.schema_vocabulary_apply(
693 crate::runtime::schema_vocabulary::DdlEvent::DropCollection {
694 collection: query.name.clone(),
695 },
696 );
697
698 Ok(RuntimeQueryResult::ok_message(
699 raw_query.to_string(),
700 &format!("table '{}' dropped", query.name),
701 "drop",
702 ))
703 }
704
705 pub fn execute_drop_graph(
706 &self,
707 raw_query: &str,
708 query: &DropGraphQuery,
709 ) -> RedDBResult<RuntimeQueryResult> {
710 self.execute_drop_typed_collection(
711 raw_query,
712 &query.name,
713 query.if_exists,
714 CollectionModel::Graph,
715 "graph",
716 )
717 }
718
719 pub fn execute_drop_vector(
720 &self,
721 raw_query: &str,
722 query: &DropVectorQuery,
723 ) -> RedDBResult<RuntimeQueryResult> {
724 self.execute_drop_typed_collection(
725 raw_query,
726 &query.name,
727 query.if_exists,
728 CollectionModel::Vector,
729 "vector",
730 )
731 }
732
733 pub fn execute_drop_document(
734 &self,
735 raw_query: &str,
736 query: &DropDocumentQuery,
737 ) -> RedDBResult<RuntimeQueryResult> {
738 self.execute_drop_typed_collection(
739 raw_query,
740 &query.name,
741 query.if_exists,
742 CollectionModel::Document,
743 "document",
744 )
745 }
746
747 pub fn execute_drop_kv(
748 &self,
749 raw_query: &str,
750 query: &DropKvQuery,
751 ) -> RedDBResult<RuntimeQueryResult> {
752 let label = polymorphic_resolver::model_name(query.model);
753 self.execute_drop_typed_collection(
754 raw_query,
755 &query.name,
756 query.if_exists,
757 query.model,
758 label,
759 )
760 }
761
762 pub fn execute_drop_collection(
763 &self,
764 raw_query: &str,
765 query: &DropCollectionQuery,
766 ) -> RedDBResult<RuntimeQueryResult> {
767 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
768 if is_system_schema_name(&query.name) {
769 return Err(RedDBError::Query("system schema is read-only".to_string()));
770 }
771 let store = self.inner.db.store();
772 if store.get_collection(&query.name).is_none() {
773 if query.if_exists {
774 return Ok(RuntimeQueryResult::ok_message(
775 raw_query.to_string(),
776 &format!("collection '{}' does not exist", query.name),
777 "drop",
778 ));
779 }
780 return Err(RedDBError::NotFound(format!(
781 "collection '{}' not found",
782 query.name
783 )));
784 }
785
786 let actual =
787 polymorphic_resolver::resolve(&query.name, &self.inner.db.catalog_model_snapshot())?;
788 if let Some(expected) = query.model {
789 polymorphic_resolver::ensure_model_match(expected, actual)?;
790 }
791
792 match actual {
793 CollectionModel::Table => self.execute_drop_table(
794 raw_query,
795 &DropTableQuery {
796 name: query.name.clone(),
797 if_exists: query.if_exists,
798 },
799 ),
800 CollectionModel::TimeSeries => self.execute_drop_timeseries(
801 raw_query,
802 &DropTimeSeriesQuery {
803 name: query.name.clone(),
804 if_exists: query.if_exists,
805 },
806 ),
807 CollectionModel::Queue => self.execute_drop_queue(
808 raw_query,
809 &DropQueueQuery {
810 name: query.name.clone(),
811 if_exists: query.if_exists,
812 },
813 ),
814 CollectionModel::Graph => self.execute_drop_graph(
815 raw_query,
816 &DropGraphQuery {
817 name: query.name.clone(),
818 if_exists: query.if_exists,
819 },
820 ),
821 CollectionModel::Vector => self.execute_drop_vector(
822 raw_query,
823 &DropVectorQuery {
824 name: query.name.clone(),
825 if_exists: query.if_exists,
826 },
827 ),
828 CollectionModel::Document => self.execute_drop_document(
829 raw_query,
830 &DropDocumentQuery {
831 name: query.name.clone(),
832 if_exists: query.if_exists,
833 },
834 ),
835 CollectionModel::Kv => self.execute_drop_kv(
836 raw_query,
837 &DropKvQuery {
838 name: query.name.clone(),
839 if_exists: query.if_exists,
840 model: CollectionModel::Kv,
841 },
842 ),
843 CollectionModel::Config => self.execute_drop_kv(
844 raw_query,
845 &DropKvQuery {
846 name: query.name.clone(),
847 if_exists: query.if_exists,
848 model: CollectionModel::Config,
849 },
850 ),
851 CollectionModel::Vault => self.execute_drop_kv(
852 raw_query,
853 &DropKvQuery {
854 name: query.name.clone(),
855 if_exists: query.if_exists,
856 model: CollectionModel::Vault,
857 },
858 ),
859 CollectionModel::Hll => self.execute_probabilistic_command(
860 raw_query,
861 &ProbabilisticCommand::DropHll {
862 name: query.name.clone(),
863 if_exists: query.if_exists,
864 },
865 ),
866 CollectionModel::Sketch => self.execute_probabilistic_command(
867 raw_query,
868 &ProbabilisticCommand::DropSketch {
869 name: query.name.clone(),
870 if_exists: query.if_exists,
871 },
872 ),
873 CollectionModel::Filter => self.execute_probabilistic_command(
874 raw_query,
875 &ProbabilisticCommand::DropFilter {
876 name: query.name.clone(),
877 if_exists: query.if_exists,
878 },
879 ),
880 CollectionModel::Metrics => self.execute_drop_typed_collection(
881 raw_query,
882 &query.name,
883 query.if_exists,
884 CollectionModel::Metrics,
885 "metrics",
886 ),
887 CollectionModel::Mixed => self.execute_drop_typed_collection(
888 raw_query,
889 &query.name,
890 query.if_exists,
891 CollectionModel::Mixed,
892 "collection",
893 ),
894 }
895 }
896
897 fn require_graph_for_analytics(&self, name: &str) -> RedDBResult<()> {
903 let is_graph = self
904 .inner
905 .db
906 .collection_contract(name)
907 .map(|c| c.declared_model == crate::catalog::CollectionModel::Graph)
908 .unwrap_or(false);
909 if !is_graph {
910 return Err(RedDBError::Query(format!(
911 "ALTER GRAPH ... ANALYTICS: '{name}' is not a graph collection"
912 )));
913 }
914 Ok(())
915 }
916
917 pub fn execute_alter_table(
923 &self,
924 raw_query: &str,
925 query: &AlterTableQuery,
926 ) -> RedDBResult<RuntimeQueryResult> {
927 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
928 let store = self.inner.db.store();
929
930 if store.get_collection(&query.name).is_none() {
932 return Err(RedDBError::NotFound(format!(
933 "table '{}' not found",
934 query.name
935 )));
936 }
937
938 let mut messages = Vec::new();
939
940 let fields_added: Vec<String> = query
942 .operations
943 .iter()
944 .filter_map(|op| {
945 if let AlterOperation::AddColumn(col) = op {
946 Some(col.name.clone())
947 } else {
948 None
949 }
950 })
951 .collect();
952 let fields_removed: Vec<String> = query
953 .operations
954 .iter()
955 .filter_map(|op| {
956 if let AlterOperation::DropColumn(name) = op {
957 Some(name.clone())
958 } else {
959 None
960 }
961 })
962 .collect();
963
964 for op in &query.operations {
965 match op {
966 AlterOperation::AddColumn(col) => {
967 messages.push(format!("column '{}' added", col.name));
969 }
970 AlterOperation::DropColumn(name) => {
971 messages.push(format!("column '{}' dropped", name));
972 }
973 AlterOperation::RenameColumn { from, to } => {
974 messages.push(format!("column '{}' renamed to '{}'", from, to));
975 }
976 AlterOperation::AttachPartition { child, bound } => {
977 store.set_config_tree(
981 &format!("partition.{}.children.{}", query.name, child),
982 &crate::serde_json::Value::String(bound.clone()),
983 );
984 messages.push(format!(
985 "partition '{child}' attached to '{}' ({bound})",
986 query.name
987 ));
988 }
989 AlterOperation::DetachPartition { child } => {
990 store.set_config_tree(
991 &format!("partition.{}.children.{}", query.name, child),
992 &crate::serde_json::Value::Null,
993 );
994 messages.push(format!(
995 "partition '{child}' detached from '{}'",
996 query.name
997 ));
998 }
999 AlterOperation::EnableRowLevelSecurity => {
1000 self.inner
1001 .rls_enabled_tables
1002 .write()
1003 .insert(query.name.clone());
1004 store.set_config_tree(
1006 &format!("rls.enabled.{}", query.name),
1007 &crate::serde_json::Value::Bool(true),
1008 );
1009 self.invalidate_plan_cache();
1010 messages.push(format!("row level security enabled on '{}'", query.name));
1011 }
1012 AlterOperation::DisableRowLevelSecurity => {
1013 self.inner.rls_enabled_tables.write().remove(&query.name);
1014 store.set_config_tree(
1015 &format!("rls.enabled.{}", query.name),
1016 &crate::serde_json::Value::Null,
1017 );
1018 self.invalidate_plan_cache();
1019 messages.push(format!("row level security disabled on '{}'", query.name));
1020 }
1021 AlterOperation::EnableTenancy { column } => {
1023 store.set_config_tree(
1024 &format!("tenant_tables.{}.column", query.name),
1025 &crate::serde_json::Value::String(column.clone()),
1026 );
1027 self.register_tenant_table(&query.name, column);
1028 self.invalidate_plan_cache();
1029 messages.push(format!(
1030 "tenancy enabled on '{}' by column '{column}'",
1031 query.name
1032 ));
1033 }
1034 AlterOperation::DisableTenancy => {
1035 store.set_config_tree(
1036 &format!("tenant_tables.{}.column", query.name),
1037 &crate::serde_json::Value::Null,
1038 );
1039 self.unregister_tenant_table(&query.name);
1040 self.invalidate_plan_cache();
1041 messages.push(format!("tenancy disabled on '{}'", query.name));
1042 }
1043 AlterOperation::SetAppendOnly(on) => {
1044 messages.push(format!(
1049 "append_only {} on '{}'",
1050 if *on { "enabled" } else { "disabled" },
1051 query.name
1052 ));
1053 }
1054 AlterOperation::SetVersioned(on) => {
1055 self.vcs_set_versioned(&query.name, *on)?;
1062 messages.push(format!(
1063 "versioned {} on '{}'",
1064 if *on { "enabled" } else { "disabled" },
1065 query.name
1066 ));
1067 }
1068 AlterOperation::EnableEvents(subscription) => {
1069 let mut subscription = subscription.clone();
1070 subscription.source = query.name.clone();
1071 validate_event_subscriptions(
1072 self,
1073 &query.name,
1074 std::slice::from_ref(&subscription),
1075 )?;
1076 ensure_event_target_queue(self, &subscription.target_queue)?;
1077 messages.push(format!(
1078 "events enabled on '{}' to '{}'",
1079 query.name, subscription.target_queue
1080 ));
1081 }
1082 AlterOperation::DisableEvents => {
1083 messages.push(format!("events disabled on '{}'", query.name));
1084 }
1085 AlterOperation::AddSubscription { name, descriptor } => {
1086 let mut sub = descriptor.clone();
1087 sub.name = name.clone();
1088 sub.source = query.name.clone();
1089 validate_event_subscriptions(self, &query.name, std::slice::from_ref(&sub))?;
1090 ensure_event_target_queue(self, &sub.target_queue)?;
1091 messages.push(format!(
1092 "subscription '{}' added on '{}' to '{}'",
1093 name, query.name, sub.target_queue
1094 ));
1095 }
1096 AlterOperation::DropSubscription { name } => {
1097 messages.push(format!(
1098 "subscription '{}' dropped on '{}'",
1099 name, query.name
1100 ));
1101 }
1102 AlterOperation::AddSigner { pubkey } => {
1103 if !crate::runtime::signed_writes_kind::is_signed(&store, &query.name) {
1110 return Err(RedDBError::Query(format!(
1111 "ALTER COLLECTION ADD SIGNER: '{}' has no signer registry; \
1112 recreate it with CREATE COLLECTION ... SIGNED_BY (...)",
1113 query.name
1114 )));
1115 }
1116 let actor = crate::runtime::impl_core::current_user_projected()
1117 .unwrap_or_else(|| "@system/alter".to_string());
1118 let changed = crate::runtime::signed_writes_kind::add_signer(
1119 &store,
1120 &query.name,
1121 *pubkey,
1122 &actor,
1123 );
1124 messages.push(format!(
1125 "signer {} on '{}'",
1126 if changed { "added" } else { "already present" },
1127 query.name
1128 ));
1129 }
1130 AlterOperation::RevokeSigner { pubkey } => {
1131 if !crate::runtime::signed_writes_kind::is_signed(&store, &query.name) {
1132 return Err(RedDBError::Query(format!(
1133 "ALTER COLLECTION REVOKE SIGNER: '{}' has no signer registry",
1134 query.name
1135 )));
1136 }
1137 let actor = crate::runtime::impl_core::current_user_projected()
1138 .unwrap_or_else(|| "@system/alter".to_string());
1139 let changed = crate::runtime::signed_writes_kind::revoke_signer(
1140 &store,
1141 &query.name,
1142 pubkey,
1143 &actor,
1144 );
1145 messages.push(format!(
1146 "signer {} on '{}'",
1147 if changed {
1148 "revoked"
1149 } else {
1150 "already revoked"
1151 },
1152 query.name
1153 ));
1154 }
1155 AlterOperation::SetRetention { duration_ms } => {
1156 let existing = self.inner.db.collection_contract(&query.name);
1162 let has_ts_column = existing
1163 .as_ref()
1164 .map(retention_timestamp_column_exists)
1165 .unwrap_or(false);
1166 if !has_ts_column {
1167 return Err(RedDBError::Query(format!(
1168 "ALTER COLLECTION SET RETENTION: '{}' has no timestamp \
1169 column — declare a TIMESTAMP/TIMESTAMPMS/DATETIME column \
1170 or enable WITH timestamps = true before setting a \
1171 retention policy",
1172 query.name
1173 )));
1174 }
1175 messages.push(format!(
1176 "retention set to {duration_ms} ms on '{}'",
1177 query.name
1178 ));
1179 }
1180 AlterOperation::UnsetRetention => {
1181 messages.push(format!("retention cleared on '{}'", query.name));
1182 }
1183 AlterOperation::AddAnalytics(views) => {
1184 self.require_graph_for_analytics(&query.name)?;
1189 let existing = self.inner.db.collection_contract(&query.name);
1190 let enabled: std::collections::BTreeSet<crate::catalog::AnalyticsOutput> =
1191 existing
1192 .as_ref()
1193 .map(|c| c.analytics_config.iter().map(|v| v.output).collect())
1194 .unwrap_or_default();
1195 for view in views {
1196 if enabled.contains(&view.output) {
1197 messages.push(format!(
1200 "analytics '{}' already enabled on '{}'",
1201 view.output.as_str(),
1202 query.name
1203 ));
1204 } else {
1205 messages.push(format!(
1206 "analytics '{}' enabled on '{}'",
1207 view.output.as_str(),
1208 query.name
1209 ));
1210 }
1211 }
1212 }
1213 AlterOperation::DropAnalytics(output) => {
1214 self.require_graph_for_analytics(&query.name)?;
1215 let enabled = self
1216 .inner
1217 .db
1218 .collection_contract(&query.name)
1219 .map(|c| c.analytics_config.iter().any(|v| v.output == *output))
1220 .unwrap_or(false);
1221 if !enabled {
1222 return Err(RedDBError::Query(format!(
1225 "ALTER GRAPH DROP ANALYTICS: analytics output '{}' is not enabled on graph '{}'",
1226 output.as_str(),
1227 query.name
1228 )));
1229 }
1230 messages.push(format!(
1231 "analytics '{}' disabled on '{}'",
1232 output.as_str(),
1233 query.name
1234 ));
1235 }
1236 }
1237 }
1238
1239 let mut contract = self
1240 .inner
1241 .db
1242 .collection_contract(&query.name)
1243 .unwrap_or_else(|| default_collection_contract_for_existing_table(&query.name));
1244 apply_alter_operations_to_contract(&mut contract, &query.operations);
1245 contract.version = contract.version.saturating_add(1);
1246 contract.updated_at_unix_ms = current_unix_ms();
1247 self.inner
1248 .db
1249 .save_collection_contract(contract)
1250 .map_err(|err| RedDBError::Internal(err.to_string()))?;
1251 if !fields_added.is_empty() || !fields_removed.is_empty() {
1255 let sub_names: Vec<String> = self
1256 .inner
1257 .db
1258 .collection_contract(&query.name)
1259 .map(|c| {
1260 c.subscriptions
1261 .iter()
1262 .filter(|s| s.enabled)
1263 .map(|s| s.name.clone())
1264 .collect()
1265 })
1266 .unwrap_or_default();
1267 if !sub_names.is_empty() {
1268 crate::telemetry::operator_event::OperatorEvent::SubscriptionSchemaChange {
1269 collection: query.name.clone(),
1270 subscription_names: sub_names.join(", "),
1271 fields_added: fields_added.join(", "),
1272 fields_removed: fields_removed.join(", "),
1273 lsn: self.cdc_current_lsn(),
1274 }
1275 .emit_global();
1276 }
1277 }
1278
1279 self.clear_table_planner_stats(&query.name);
1280 self.invalidate_result_cache();
1281 let post_alter_columns: Vec<String> = self
1286 .inner
1287 .db
1288 .collection_contract(&query.name)
1289 .map(|contract| {
1290 contract
1291 .declared_columns
1292 .iter()
1293 .map(|col| col.name.clone())
1294 .collect()
1295 })
1296 .unwrap_or_default();
1297 self.schema_vocabulary_apply(
1298 crate::runtime::schema_vocabulary::DdlEvent::AlterCollection {
1299 collection: query.name.clone(),
1300 columns: post_alter_columns,
1301 type_tags: Vec::new(),
1302 description: None,
1303 },
1304 );
1305
1306 let message = if messages.is_empty() {
1307 format!("table '{}' altered (no operations)", query.name)
1308 } else {
1309 format!("table '{}' altered: {}", query.name, messages.join(", "))
1310 };
1311
1312 Ok(RuntimeQueryResult::ok_message(
1313 raw_query.to_string(),
1314 &message,
1315 "alter",
1316 ))
1317 }
1318
1319 pub fn execute_create_vcs_ref(
1320 &self,
1321 raw_query: &str,
1322 query: &CreateVcsRefQuery,
1323 ) -> RedDBResult<RuntimeQueryResult> {
1324 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1325 let connection_id = crate::runtime::impl_core::current_connection_id();
1326 match query.kind {
1327 VcsRefKind::Branch => {
1328 self.vcs_branch_create(crate::application::vcs::CreateBranchInput {
1329 name: query.name.clone(),
1330 from: query.target.clone(),
1331 connection_id,
1332 })?;
1333 }
1334 VcsRefKind::Tag => {
1335 let target = match &query.target {
1336 Some(target) => target.clone(),
1337 None => self
1338 .vcs_status(crate::application::vcs::StatusInput { connection_id })?
1339 .head_commit
1340 .filter(|head| !head.is_empty())
1341 .ok_or_else(|| {
1342 RedDBError::InvalidConfig(
1343 "cannot create tag: HEAD has no commits".to_string(),
1344 )
1345 })?,
1346 };
1347 self.vcs_tag_create(crate::application::vcs::CreateTagInput {
1348 name: query.name.clone(),
1349 target,
1350 annotation: None,
1351 })?;
1352 }
1353 }
1354 let kind = match query.kind {
1355 VcsRefKind::Branch => "branch",
1356 VcsRefKind::Tag => "tag",
1357 };
1358 Ok(RuntimeQueryResult::ok_message(
1359 raw_query.to_string(),
1360 &format!("{kind} '{}' created", query.name),
1361 "create",
1362 ))
1363 }
1364
1365 pub fn execute_drop_vcs_ref(
1366 &self,
1367 raw_query: &str,
1368 query: &DropVcsRefQuery,
1369 ) -> RedDBResult<RuntimeQueryResult> {
1370 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1371 match query.kind {
1372 VcsRefKind::Branch => self.vcs_branch_delete(&query.name)?,
1373 VcsRefKind::Tag => self.vcs_tag_delete(&query.name)?,
1374 }
1375 let kind = match query.kind {
1376 VcsRefKind::Branch => "branch",
1377 VcsRefKind::Tag => "tag",
1378 };
1379 Ok(RuntimeQueryResult::ok_message(
1380 raw_query.to_string(),
1381 &format!("{kind} '{}' dropped", query.name),
1382 "drop",
1383 ))
1384 }
1385
1386 pub fn execute_fork_store(
1387 &self,
1388 raw_query: &str,
1389 query: &ForkStoreQuery,
1390 ) -> RedDBResult<RuntimeQueryResult> {
1391 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1392 let path = self.inner.db.path().ok_or_else(|| {
1393 RedDBError::Query("FORK STORE requires a persistent store".to_string())
1394 })?;
1395 let durable_lsn = self.cdc_current_lsn();
1396 let fork_lsn = query.at_lsn.unwrap_or(durable_lsn);
1397 let retention_floor = operational_wal_retention_floor(durable_lsn);
1398 validate_store_fork_lsn(fork_lsn, retention_floor, durable_lsn)?;
1399
1400 let manifest = reddb_file::OperationalManifest::for_db_path(path);
1401 let mut collections = self.inner.db.store().list_collections();
1402 collections.sort();
1403 manifest.recover_or_bootstrap(&collections).map_err(|err| {
1404 RedDBError::Query(format!("failed to prepare store fork manifest: {err}"))
1405 })?;
1406 manifest
1407 .create_fork(&query.name, fork_lsn)
1408 .map_err(|err| RedDBError::Query(format!("failed to create store fork: {err}")))?;
1409 Ok(RuntimeQueryResult::ok_message(
1410 raw_query.to_string(),
1411 &format!("store fork '{}' created at LSN {fork_lsn}", query.name),
1412 "fork_store",
1413 ))
1414 }
1415
1416 pub fn execute_drop_fork(
1417 &self,
1418 raw_query: &str,
1419 query: &DropForkQuery,
1420 ) -> RedDBResult<RuntimeQueryResult> {
1421 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1422 let path = self.inner.db.path().ok_or_else(|| {
1423 RedDBError::Query("DROP FORK requires a persistent store".to_string())
1424 })?;
1425 let manifest = reddb_file::OperationalManifest::for_db_path(path);
1426 let dropped = manifest
1427 .drop_fork(&query.name)
1428 .map_err(|err| RedDBError::Query(format!("failed to drop store fork: {err}")))?;
1429 if !dropped && !query.if_exists {
1430 return Err(RedDBError::NotFound(format!("store fork '{}'", query.name)));
1431 }
1432 Ok(RuntimeQueryResult::ok_message(
1433 raw_query.to_string(),
1434 &format!("store fork '{}' dropped", query.name),
1435 "drop_fork",
1436 ))
1437 }
1438
1439 pub fn execute_promote_fork(
1440 &self,
1441 raw_query: &str,
1442 query: &PromoteForkQuery,
1443 ) -> RedDBResult<RuntimeQueryResult> {
1444 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1445 let path = self.inner.db.path().ok_or_else(|| {
1446 RedDBError::Query("PROMOTE FORK requires a persistent store".to_string())
1447 })?;
1448 self.flush()?;
1449 let manifest = reddb_file::OperationalManifest::for_db_path(path);
1450 let outcome = manifest
1451 .promote_fork(&query.name)
1452 .map_err(|err| RedDBError::Query(format!("failed to promote store fork: {err}")))?;
1453 let outcome =
1454 outcome.ok_or_else(|| RedDBError::NotFound(format!("store fork '{}'", query.name)))?;
1455 Ok(RuntimeQueryResult::ok_message(
1456 raw_query.to_string(),
1457 &format!(
1458 "store fork '{}' promoted at LSN {}; retired parent archived at {}",
1459 outcome.name,
1460 outcome.fork_lsn,
1461 outcome.archived_parent.store_identity()
1462 ),
1463 "promote_fork",
1464 ))
1465 }
1466
1467 pub fn execute_explain_alter(
1474 &self,
1475 raw_query: &str,
1476 query: &ExplainAlterQuery,
1477 ) -> RedDBResult<RuntimeQueryResult> {
1478 analyze_create_table(&query.target).map_err(|err| RedDBError::Query(err.to_string()))?;
1482
1483 let current_contract = self.inner.db.collection_contract(&query.target.name);
1484
1485 let current_columns: Vec<crate::physical::DeclaredColumnContract> = current_contract
1486 .as_ref()
1487 .map(|c| c.declared_columns.clone())
1488 .unwrap_or_default();
1489
1490 let diff = super::schema_diff::compute_column_diff(
1491 &query.target.name,
1492 ¤t_columns,
1493 &query.target.columns,
1494 );
1495
1496 let rendered = match query.format {
1497 ExplainFormat::Sql => super::schema_diff::format_as_sql(&diff),
1498 ExplainFormat::Json => super::schema_diff::format_as_json(&diff),
1499 };
1500
1501 let format_label = match query.format {
1502 ExplainFormat::Sql => "sql",
1503 ExplainFormat::Json => "json",
1504 };
1505
1506 let columns = vec![
1507 "table".to_string(),
1508 "format".to_string(),
1509 "diff".to_string(),
1510 ];
1511 let row = vec![
1512 ("table".to_string(), Value::text(query.target.name.clone())),
1513 ("format".to_string(), Value::text(format_label.to_string())),
1514 ("diff".to_string(), Value::text(rendered)),
1515 ];
1516
1517 Ok(RuntimeQueryResult::ok_records(
1518 raw_query.to_string(),
1519 columns,
1520 vec![row],
1521 "explain",
1522 ))
1523 }
1524
1525 pub fn execute_create_index(
1530 &self,
1531 raw_query: &str,
1532 query: &CreateIndexQuery,
1533 ) -> RedDBResult<RuntimeQueryResult> {
1534 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1535 let store = self.inner.db.store();
1536
1537 let manager = store
1539 .get_collection(&query.table)
1540 .ok_or_else(|| RedDBError::NotFound(format!("table '{}' not found", query.table)))?;
1541
1542 let method_kind = match query.method {
1543 IndexMethod::Hash => super::index_store::IndexMethodKind::Hash,
1544 IndexMethod::BTree => super::index_store::IndexMethodKind::BTree,
1545 IndexMethod::Bitmap => super::index_store::IndexMethodKind::Bitmap,
1546 IndexMethod::Spatial => super::index_store::IndexMethodKind::default_spatial(),
1549 IndexMethod::H3 { resolution } => {
1550 super::index_store::IndexMethodKind::H3 { resolution }
1551 }
1552 };
1553
1554 let entities = manager.query_all(|_| true);
1564 let entity_fields: Vec<(crate::storage::unified::EntityId, Vec<(String, Value)>)> =
1565 entities
1566 .iter()
1567 .map(|e| {
1568 let fields = match &e.data {
1569 crate::storage::EntityData::Row(row) => {
1570 if let Some(ref named) = row.named {
1571 named.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
1572 } else if let Some(ref schema) = row.schema {
1573 schema
1577 .iter()
1578 .zip(row.columns.iter())
1579 .map(|(k, v)| (k.clone(), v.clone()))
1580 .collect()
1581 } else {
1582 Vec::new()
1583 }
1584 }
1585 crate::storage::EntityData::Node(node) => node
1586 .properties
1587 .iter()
1588 .map(|(k, v)| (k.clone(), v.clone()))
1589 .collect(),
1590 crate::storage::EntityData::TimeSeries(ts) => ts
1591 .fields
1592 .iter()
1593 .map(|(k, v)| (k.clone(), v.clone()))
1594 .collect(),
1595 _ => Vec::new(),
1596 };
1597 (e.id, fields)
1598 })
1599 .collect();
1600
1601 let index_growth_rows = entity_fields
1602 .iter()
1603 .map(|(_, fields)| fields.clone())
1604 .collect::<Vec<_>>();
1605 self.admit_non_evictable_growth(
1606 crate::storage::memory_pools::MemoryPool::IndexMemory,
1607 &format!("create index {}", query.name),
1608 crate::runtime::memory_admission::estimate_index_growth(
1609 &index_growth_rows,
1610 &query.columns,
1611 ),
1612 )?;
1613
1614 let indexed_count = self
1616 .inner
1617 .index_store
1618 .create_index(
1619 &query.name,
1620 &query.table,
1621 &query.columns,
1622 method_kind,
1623 query.unique,
1624 &entity_fields,
1625 )
1626 .map_err(RedDBError::Internal)?;
1627
1628 let analyzed = crate::storage::query::planner::stats_catalog::analyze_entity_fields(
1629 &query.table,
1630 &entity_fields,
1631 );
1632 crate::storage::query::planner::stats_catalog::persist_table_stats(&store, &analyzed);
1633 self.invalidate_plan_cache();
1634
1635 self.inner
1637 .index_store
1638 .register(super::index_store::RegisteredIndex {
1639 name: query.name.clone(),
1640 collection: query.table.clone(),
1641 columns: query.columns.clone(),
1642 method: method_kind,
1643 unique: query.unique,
1644 });
1645 self.persist_runtime_index_descriptor(super::index_store::RegisteredIndex {
1646 name: query.name.clone(),
1647 collection: query.table.clone(),
1648 columns: query.columns.clone(),
1649 method: method_kind,
1650 unique: query.unique,
1651 })?;
1652 self.schema_vocabulary_apply(crate::runtime::schema_vocabulary::DdlEvent::CreateIndex {
1656 collection: query.table.clone(),
1657 index: query.name.clone(),
1658 columns: query.columns.clone(),
1659 });
1660
1661 let method_str = format!("{}", query.method);
1662 let unique_str = if query.unique { "unique " } else { "" };
1663 let cols = query.columns.join(", ");
1664 let total_count = entity_fields.len();
1665 let coverage_detail =
1666 if matches!(method_kind, super::index_store::IndexMethodKind::H3 { .. })
1667 && total_count > 0
1668 && indexed_count == 0
1669 {
1670 format!(
1671 "{} of {} entities indexed — no indexable geo value in '{}'; expected {}",
1672 indexed_count,
1673 total_count,
1674 cols,
1675 crate::geo::RECOGNIZED_GEO_SHAPES
1676 )
1677 } else if matches!(method_kind, super::index_store::IndexMethodKind::H3 { .. }) {
1678 format!("{} of {} entities indexed", indexed_count, total_count)
1679 } else {
1680 format!("{} entities indexed", indexed_count)
1681 };
1682
1683 Ok(RuntimeQueryResult::ok_message(
1684 raw_query.to_string(),
1685 &format!(
1686 "{}index '{}' created on '{}' ({}) using {} ({})",
1687 unique_str, query.name, query.table, cols, method_str, coverage_detail
1688 ),
1689 "create",
1690 ))
1691 }
1692
1693 pub fn execute_drop_index(
1697 &self,
1698 raw_query: &str,
1699 query: &DropIndexQuery,
1700 ) -> RedDBResult<RuntimeQueryResult> {
1701 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1702 let store = self.inner.db.store();
1703
1704 if store.get_collection(&query.table).is_none() {
1706 if query.if_exists {
1707 return Ok(RuntimeQueryResult::ok_message(
1708 raw_query.to_string(),
1709 &format!("table '{}' does not exist", query.table),
1710 "drop",
1711 ));
1712 }
1713 return Err(RedDBError::NotFound(format!(
1714 "table '{}' not found",
1715 query.table
1716 )));
1717 }
1718
1719 self.inner.index_store.drop_index(&query.name, &query.table);
1721 self.persist_runtime_index_drop(&query.table, &query.name)?;
1722 self.invalidate_plan_cache();
1723 self.schema_vocabulary_apply(crate::runtime::schema_vocabulary::DdlEvent::DropIndex {
1725 collection: query.table.clone(),
1726 index: query.name.clone(),
1727 });
1728
1729 Ok(RuntimeQueryResult::ok_message(
1730 raw_query.to_string(),
1731 &format!("index '{}' dropped from '{}'", query.name, query.table),
1732 "drop",
1733 ))
1734 }
1735
1736 fn execute_drop_typed_collection(
1737 &self,
1738 raw_query: &str,
1739 name: &str,
1740 if_exists: bool,
1741 expected_model: CollectionModel,
1742 label: &str,
1743 ) -> RedDBResult<RuntimeQueryResult> {
1744 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1745 if is_system_schema_name(name) {
1746 return Err(RedDBError::Query("system schema is read-only".to_string()));
1747 }
1748 let store = self.inner.db.store();
1749 if store.get_collection(name).is_none() {
1750 if if_exists {
1751 return Ok(RuntimeQueryResult::ok_message(
1752 raw_query.to_string(),
1753 &format!("{label} '{name}' does not exist"),
1754 "drop",
1755 ));
1756 }
1757 return Err(RedDBError::NotFound(format!("{label} '{name}' not found")));
1758 }
1759
1760 let actual = self
1761 .inner
1762 .db
1763 .collection_contract(name)
1764 .map(|contract| contract.declared_model)
1765 .map(Ok)
1766 .unwrap_or_else(|| {
1767 polymorphic_resolver::resolve(name, &self.inner.db.catalog_model_snapshot())
1768 })?;
1769 polymorphic_resolver::ensure_model_match(expected_model, actual)?;
1770 self.drop_collection_storage(raw_query, name, label)
1771 }
1772
1773 pub fn execute_truncate(
1774 &self,
1775 raw_query: &str,
1776 query: &TruncateQuery,
1777 ) -> RedDBResult<RuntimeQueryResult> {
1778 self.check_write(crate::runtime::write_gate::WriteKind::Ddl)?;
1779 if is_system_schema_name(&query.name) {
1780 return Err(RedDBError::Query("system schema is read-only".to_string()));
1781 }
1782
1783 let label = query
1784 .model
1785 .map(polymorphic_resolver::model_name)
1786 .unwrap_or("collection");
1787 let store = self.inner.db.store();
1788 if store.get_collection(&query.name).is_none() {
1789 if query.if_exists {
1790 return Ok(RuntimeQueryResult::ok_message(
1791 raw_query.to_string(),
1792 &format!("{label} '{}' does not exist", query.name),
1793 "truncate",
1794 ));
1795 }
1796 return Err(RedDBError::NotFound(format!(
1797 "{label} '{}' not found",
1798 query.name
1799 )));
1800 }
1801
1802 let actual =
1803 polymorphic_resolver::resolve(&query.name, &self.inner.db.catalog_model_snapshot())?;
1804 if let Some(expected) = query.model {
1805 polymorphic_resolver::ensure_model_match(expected, actual)?;
1806 }
1807
1808 if actual == CollectionModel::Queue {
1809 return self.execute_queue_command(
1810 raw_query,
1811 &QueueCommand::Purge {
1812 queue: query.name.clone(),
1813 },
1814 );
1815 }
1816
1817 let affected = self.truncate_collection_entities(&query.name)?;
1819 crate::runtime::mutation::emit_truncate_event_for_collection(self, &query.name, affected)?;
1821 self.inner.db.invalidate_vector_index(&query.name);
1822 self.clear_table_planner_stats(&query.name);
1823 self.invalidate_result_cache();
1824
1825 Ok(RuntimeQueryResult::ok_message(
1826 raw_query.to_string(),
1827 &format!(
1828 "{affected} entities truncated from {label} '{}'",
1829 query.name
1830 ),
1831 "truncate",
1832 ))
1833 }
1834
1835 fn truncate_collection_entities(&self, name: &str) -> RedDBResult<u64> {
1836 let store = self.inner.db.store();
1837 let Some(manager) = store.get_collection(name) else {
1838 return Ok(0);
1839 };
1840 let entities = manager.query_all(|_| true);
1841 if entities.is_empty() {
1842 return Ok(0);
1843 }
1844
1845 for entity in &entities {
1846 let fields = entity_index_fields(&entity.data);
1847 self.inner
1848 .index_store
1849 .index_entity_delete(name, entity.id, &fields)
1850 .map_err(RedDBError::Internal)?;
1851 }
1852
1853 let ids = entities.iter().map(|entity| entity.id).collect::<Vec<_>>();
1854 let deleted_ids = store
1855 .delete_batch(name, &ids)
1856 .map_err(|err| RedDBError::Internal(err.to_string()))?;
1857 for id in &deleted_ids {
1858 store.context_index().remove_entity(*id);
1859 }
1860 Ok(deleted_ids.len() as u64)
1861 }
1862
1863 fn drop_collection_storage(
1864 &self,
1865 raw_query: &str,
1866 name: &str,
1867 label: &str,
1868 ) -> RedDBResult<RuntimeQueryResult> {
1869 let store = self.inner.db.store();
1870
1871 let final_count = store
1874 .get_collection(name)
1875 .map(|manager| manager.query_all(|_| true).len() as u64)
1876 .unwrap_or(0);
1877 crate::runtime::mutation::emit_collection_dropped_event_for_collection(
1878 self,
1879 name,
1880 final_count,
1881 )?;
1882
1883 let orphaned_indices: Vec<String> = self
1884 .inner
1885 .index_store
1886 .list_indices(name)
1887 .into_iter()
1888 .map(|index| index.name)
1889 .collect();
1890 for index_name in &orphaned_indices {
1891 self.inner.index_store.drop_index(index_name, name);
1892 }
1893
1894 store
1895 .drop_collection(name)
1896 .map_err(|err| RedDBError::Internal(err.to_string()))?;
1897 self.inner.db.invalidate_vector_index(name);
1898 self.inner.db.clear_collection_default_ttl_ms(name);
1899 self.inner
1900 .db
1901 .remove_collection_contract(name)
1902 .map_err(|err| RedDBError::Internal(err.to_string()))?;
1903 self.clear_table_planner_stats(name);
1904 self.invalidate_result_cache();
1905 if let Some(store) = self.inner.auth_store.read().clone() {
1906 store.invalidate_visible_collections_cache();
1907 }
1908 self.inner
1909 .db
1910 .persist_metadata()
1911 .map_err(|err| RedDBError::Internal(err.to_string()))?;
1912 self.schema_vocabulary_apply(
1913 crate::runtime::schema_vocabulary::DdlEvent::DropCollection {
1914 collection: name.to_string(),
1915 },
1916 );
1917
1918 Ok(RuntimeQueryResult::ok_message(
1919 raw_query.to_string(),
1920 &format!("{label} '{name}' dropped"),
1921 "drop",
1922 ))
1923 }
1924}
1925
1926pub(crate) fn is_system_schema_name(name: &str) -> bool {
1927 name == "red" || name.starts_with("red.") || name.starts_with("__red_schema_")
1928}
1929
1930fn entity_index_fields(data: &EntityData) -> Vec<(String, Value)> {
1931 match data {
1932 EntityData::Row(row) => {
1933 if let Some(ref named) = row.named {
1934 named
1935 .iter()
1936 .map(|(key, value)| (key.clone(), value.clone()))
1937 .collect()
1938 } else if let Some(ref schema) = row.schema {
1939 schema
1940 .iter()
1941 .zip(row.columns.iter())
1942 .map(|(key, value)| (key.clone(), value.clone()))
1943 .collect()
1944 } else {
1945 Vec::new()
1946 }
1947 }
1948 EntityData::Node(node) => node
1949 .properties
1950 .iter()
1951 .map(|(key, value)| (key.clone(), value.clone()))
1952 .collect(),
1953 _ => Vec::new(),
1954 }
1955}
1956
1957fn validate_ai_policy_modalities(policy: &crate::catalog::AiPolicy) -> RedDBResult<()> {
1963 use crate::runtime::ai::provider_capabilities::{Modality, Registry};
1964 let registry = Registry::new();
1965 let check = |provider: &str, model: &str, modality: Modality| -> RedDBResult<()> {
1966 registry
1967 .validate_policy_modality(provider, model, modality)
1968 .map_err(|err| RedDBError::Query(err.to_string()))
1969 };
1970 if let Some(embed) = &policy.embed {
1971 check(&embed.provider, &embed.model, Modality::Embed)?;
1972 }
1973 if let Some(moderate) = &policy.moderate {
1974 check(&moderate.provider, &moderate.model, Modality::Moderate)?;
1975 }
1976 if let Some(vision) = &policy.vision {
1977 check(&vision.provider, &vision.model, Modality::Vision)?;
1978 }
1979 Ok(())
1980}
1981
1982fn collection_contract_from_create_table(
1983 query: &CreateTableQuery,
1984) -> RedDBResult<crate::physical::CollectionContract> {
1985 let now = current_unix_ms();
1986 let mut declared_columns: Vec<crate::physical::DeclaredColumnContract> = query
1987 .columns
1988 .iter()
1989 .map(declared_column_contract_from_ddl)
1990 .collect();
1991 if query.timestamps {
1992 declared_columns.push(crate::physical::DeclaredColumnContract {
1996 name: "created_at".to_string(),
1997 data_type: "BIGINT".to_string(),
1998 sql_type: Some(crate::storage::schema::SqlTypeName::simple("BIGINT")),
1999 not_null: true,
2000 default: None,
2001 compress: None,
2002 unique: false,
2003 primary_key: false,
2004 enum_variants: Vec::new(),
2005 array_element: None,
2006 decimal_precision: None,
2007 });
2008 declared_columns.push(crate::physical::DeclaredColumnContract {
2009 name: "updated_at".to_string(),
2010 data_type: "BIGINT".to_string(),
2011 sql_type: Some(crate::storage::schema::SqlTypeName::simple("BIGINT")),
2012 not_null: true,
2013 default: None,
2014 compress: None,
2015 unique: false,
2016 primary_key: false,
2017 enum_variants: Vec::new(),
2018 array_element: None,
2019 decimal_precision: None,
2020 });
2021 }
2022 Ok(crate::physical::CollectionContract {
2023 name: query.name.clone(),
2024 declared_model: crate::catalog::CollectionModel::Table,
2025 schema_mode: crate::catalog::SchemaMode::SemiStructured,
2026 origin: crate::physical::ContractOrigin::Explicit,
2027 version: 1,
2028 created_at_unix_ms: now,
2029 updated_at_unix_ms: now,
2030 default_ttl_ms: query.default_ttl_ms,
2031 vector_dimension: None,
2032 vector_metric: None,
2033 context_index_fields: query.context_index_fields.clone(),
2034 declared_columns,
2035 table_def: Some(build_table_def_from_create_table(query)?),
2036 timestamps_enabled: query.timestamps,
2037 context_index_enabled: query.context_index_enabled
2038 || !query.context_index_fields.is_empty(),
2039 metrics_raw_retention_ms: None,
2040 metrics_rollup_policies: Vec::new(),
2041 metrics_tenant_identity: None,
2042 metrics_namespace: None,
2043 append_only: query.append_only,
2044 subscriptions: query.subscriptions.clone(),
2045 analytics_config: Vec::new(),
2046 session_key: None,
2047 session_gap_ms: None,
2048 retention_duration_ms: None,
2049 analytical_storage: None,
2050 ai_policy: query.ai_policy.clone(),
2051 })
2052}
2053
2054fn default_collection_contract_for_existing_table(
2055 name: &str,
2056) -> crate::physical::CollectionContract {
2057 let now = current_unix_ms();
2058 crate::physical::CollectionContract {
2059 name: name.to_string(),
2060 declared_model: crate::catalog::CollectionModel::Table,
2061 schema_mode: crate::catalog::SchemaMode::SemiStructured,
2062 origin: crate::physical::ContractOrigin::Explicit,
2063 version: 0,
2064 created_at_unix_ms: now,
2065 updated_at_unix_ms: now,
2066 default_ttl_ms: None,
2067 vector_dimension: None,
2068 vector_metric: None,
2069 context_index_fields: Vec::new(),
2070 declared_columns: Vec::new(),
2071 table_def: Some(crate::storage::schema::TableDef::new(name.to_string())),
2072 timestamps_enabled: false,
2073 context_index_enabled: false,
2074 metrics_raw_retention_ms: None,
2075 metrics_rollup_policies: Vec::new(),
2076 metrics_tenant_identity: None,
2077 metrics_namespace: None,
2078 append_only: false,
2079 subscriptions: Vec::new(),
2080 analytics_config: Vec::new(),
2081 session_key: None,
2082 session_gap_ms: None,
2083 retention_duration_ms: None,
2084 analytical_storage: None,
2085
2086 ai_policy: None,
2087 }
2088}
2089
2090fn keyed_collection_contract(
2091 name: &str,
2092 model: crate::catalog::CollectionModel,
2093 analytics_config: Vec<crate::catalog::AnalyticsViewDescriptor>,
2094) -> crate::physical::CollectionContract {
2095 let now = current_unix_ms();
2096 crate::physical::CollectionContract {
2097 name: name.to_string(),
2098 declared_model: model,
2099 schema_mode: crate::catalog::SchemaMode::Dynamic,
2100 origin: crate::physical::ContractOrigin::Explicit,
2101 version: 1,
2102 created_at_unix_ms: now,
2103 updated_at_unix_ms: now,
2104 default_ttl_ms: None,
2105 vector_dimension: None,
2106 vector_metric: None,
2107 context_index_fields: Vec::new(),
2108 declared_columns: Vec::new(),
2109 table_def: None,
2110 timestamps_enabled: false,
2111 context_index_enabled: false,
2112 metrics_raw_retention_ms: None,
2113 metrics_rollup_policies: Vec::new(),
2114 metrics_tenant_identity: None,
2115 metrics_namespace: None,
2116 append_only: false,
2117 subscriptions: Vec::new(),
2118 analytics_config,
2119 session_key: None,
2120 session_gap_ms: None,
2121 retention_duration_ms: None,
2122 analytical_storage: None,
2123
2124 ai_policy: None,
2125 }
2126}
2127
2128fn metrics_collection_contract(query: &CreateTableQuery) -> crate::physical::CollectionContract {
2129 let now = current_unix_ms();
2130 crate::physical::CollectionContract {
2131 name: query.name.clone(),
2132 declared_model: crate::catalog::CollectionModel::Metrics,
2133 schema_mode: crate::catalog::SchemaMode::SemiStructured,
2134 origin: crate::physical::ContractOrigin::Explicit,
2135 version: 1,
2136 created_at_unix_ms: now,
2137 updated_at_unix_ms: now,
2138 default_ttl_ms: query.default_ttl_ms,
2139 vector_dimension: None,
2140 vector_metric: None,
2141 context_index_fields: Vec::new(),
2142 declared_columns: Vec::new(),
2143 table_def: None,
2144 timestamps_enabled: false,
2145 context_index_enabled: false,
2146 metrics_raw_retention_ms: query.default_ttl_ms,
2147 metrics_rollup_policies: query.metrics_rollup_policies.clone(),
2148 metrics_tenant_identity: Some(
2149 query
2150 .tenant_by
2151 .clone()
2152 .unwrap_or_else(|| "current_tenant".to_string()),
2153 ),
2154 metrics_namespace: Some("default".to_string()),
2155 append_only: true,
2156 subscriptions: Vec::new(),
2157 analytics_config: Vec::new(),
2158 session_key: None,
2159 session_gap_ms: None,
2160 retention_duration_ms: None,
2161 analytical_storage: None,
2162
2163 ai_policy: None,
2164 }
2165}
2166
2167fn vector_collection_contract(query: &CreateVectorQuery) -> crate::physical::CollectionContract {
2168 let now = current_unix_ms();
2169 crate::physical::CollectionContract {
2170 name: query.name.clone(),
2171 declared_model: crate::catalog::CollectionModel::Vector,
2172 schema_mode: crate::catalog::SchemaMode::Dynamic,
2173 origin: crate::physical::ContractOrigin::Explicit,
2174 version: 1,
2175 created_at_unix_ms: now,
2176 updated_at_unix_ms: now,
2177 default_ttl_ms: None,
2178 vector_dimension: Some(query.dimension),
2179 vector_metric: Some(query.metric),
2180 context_index_fields: Vec::new(),
2181 declared_columns: Vec::new(),
2182 table_def: None,
2183 timestamps_enabled: false,
2184 context_index_enabled: false,
2185 metrics_raw_retention_ms: None,
2186 metrics_rollup_policies: Vec::new(),
2187 metrics_tenant_identity: None,
2188 metrics_namespace: None,
2189 append_only: false,
2190 subscriptions: Vec::new(),
2191 analytics_config: Vec::new(),
2192 session_key: None,
2193 session_gap_ms: None,
2194 retention_duration_ms: None,
2195 analytical_storage: None,
2196
2197 ai_policy: None,
2198 }
2199}
2200
2201fn declared_column_contract_from_ddl(
2202 column: &CreateColumnDef,
2203) -> crate::physical::DeclaredColumnContract {
2204 crate::physical::DeclaredColumnContract {
2205 name: column.name.clone(),
2206 data_type: column.data_type.clone(),
2207 sql_type: Some(column.sql_type.clone()),
2208 not_null: column.not_null,
2209 default: column.default.clone(),
2210 compress: column.compress,
2211 unique: column.unique,
2212 primary_key: column.primary_key,
2213 enum_variants: column.enum_variants.clone(),
2214 array_element: column.array_element.clone(),
2215 decimal_precision: column.decimal_precision,
2216 }
2217}
2218
2219fn apply_alter_operations_to_contract(
2220 contract: &mut crate::physical::CollectionContract,
2221 operations: &[AlterOperation],
2222) {
2223 if contract.table_def.is_none() {
2224 contract.table_def = Some(crate::storage::schema::TableDef::new(contract.name.clone()));
2225 }
2226 for operation in operations {
2227 match operation {
2228 AlterOperation::AddColumn(column) => {
2229 if !contract
2230 .declared_columns
2231 .iter()
2232 .any(|existing| existing.name == column.name)
2233 {
2234 contract
2235 .declared_columns
2236 .push(declared_column_contract_from_ddl(column));
2237 }
2238 if let Some(table_def) = contract.table_def.as_mut() {
2239 if table_def.get_column(&column.name).is_none() {
2240 if let Ok(column_def) = column_def_from_ddl(column) {
2241 if column.primary_key {
2242 table_def.primary_key.push(column.name.clone());
2243 table_def.constraints.push(
2244 crate::storage::schema::Constraint::new(
2245 format!("pk_{}", column.name),
2246 crate::storage::schema::ConstraintType::PrimaryKey,
2247 )
2248 .on_columns(vec![column.name.clone()]),
2249 );
2250 }
2251 if column.unique {
2252 table_def.constraints.push(
2253 crate::storage::schema::Constraint::new(
2254 format!("uniq_{}", column.name),
2255 crate::storage::schema::ConstraintType::Unique,
2256 )
2257 .on_columns(vec![column.name.clone()]),
2258 );
2259 }
2260 if column.not_null {
2261 table_def.constraints.push(
2262 crate::storage::schema::Constraint::new(
2263 format!("not_null_{}", column.name),
2264 crate::storage::schema::ConstraintType::NotNull,
2265 )
2266 .on_columns(vec![column.name.clone()]),
2267 );
2268 }
2269 table_def.columns.push(column_def);
2270 }
2271 }
2272 }
2273 }
2274 AlterOperation::DropColumn(name) => {
2275 contract
2276 .declared_columns
2277 .retain(|column| column.name != *name);
2278 if let Some(table_def) = contract.table_def.as_mut() {
2279 if let Some(index) = table_def.column_index(name) {
2280 table_def.columns.remove(index);
2281 }
2282 table_def.primary_key.retain(|column| column != name);
2283 table_def.constraints.retain(|constraint| {
2284 !constraint.columns.iter().any(|column| column == name)
2285 });
2286 table_def
2287 .indexes
2288 .retain(|index| !index.columns.iter().any(|column| column == name));
2289 }
2290 }
2291 AlterOperation::RenameColumn { from, to } => {
2292 if contract
2293 .declared_columns
2294 .iter()
2295 .any(|column| column.name == *to)
2296 {
2297 continue;
2298 }
2299 if let Some(column) = contract
2300 .declared_columns
2301 .iter_mut()
2302 .find(|column| column.name == *from)
2303 {
2304 column.name = to.clone();
2305 }
2306 if let Some(table_def) = contract.table_def.as_mut() {
2307 if let Some(column) = table_def
2308 .columns
2309 .iter_mut()
2310 .find(|column| column.name == *from)
2311 {
2312 column.name = to.clone();
2313 }
2314 for primary_key in &mut table_def.primary_key {
2315 if *primary_key == *from {
2316 *primary_key = to.clone();
2317 }
2318 }
2319 for constraint in &mut table_def.constraints {
2320 for column in &mut constraint.columns {
2321 if *column == *from {
2322 *column = to.clone();
2323 }
2324 }
2325 if let Some(ref_columns) = constraint.ref_columns.as_mut() {
2326 for column in ref_columns {
2327 if *column == *from {
2328 *column = to.clone();
2329 }
2330 }
2331 }
2332 }
2333 for index in &mut table_def.indexes {
2334 for column in &mut index.columns {
2335 if *column == *from {
2336 *column = to.clone();
2337 }
2338 }
2339 }
2340 }
2341 }
2342 AlterOperation::AttachPartition { .. } | AlterOperation::DetachPartition { .. } => {}
2345 AlterOperation::EnableRowLevelSecurity | AlterOperation::DisableRowLevelSecurity => {}
2349 AlterOperation::EnableTenancy { .. } | AlterOperation::DisableTenancy => {}
2352 AlterOperation::SetAppendOnly(on) => {
2353 contract.append_only = *on;
2354 }
2355 AlterOperation::SetVersioned(_) => {}
2358 AlterOperation::EnableEvents(subscription) => {
2359 let mut subscription = subscription.clone();
2360 subscription.source = contract.name.clone();
2361 subscription.enabled = true;
2362 if let Some(existing) = contract
2363 .subscriptions
2364 .iter_mut()
2365 .find(|existing| existing.target_queue == subscription.target_queue)
2366 {
2367 *existing = subscription;
2368 } else {
2369 contract.subscriptions.push(subscription);
2370 }
2371 }
2372 AlterOperation::DisableEvents => {
2373 for subscription in &mut contract.subscriptions {
2374 subscription.enabled = false;
2375 }
2376 }
2377 AlterOperation::AddSubscription { name, descriptor } => {
2378 let mut sub = descriptor.clone();
2379 sub.name = name.clone();
2380 sub.source = contract.name.clone();
2381 sub.enabled = true;
2382 if let Some(existing) = contract.subscriptions.iter_mut().find(|s| s.name == *name)
2383 {
2384 *existing = sub;
2385 } else {
2386 contract.subscriptions.push(sub);
2387 }
2388 }
2389 AlterOperation::DropSubscription { name } => {
2390 contract.subscriptions.retain(|s| s.name != *name);
2391 }
2392 AlterOperation::AddSigner { .. } | AlterOperation::RevokeSigner { .. } => {}
2397 AlterOperation::SetRetention { duration_ms } => {
2398 contract.retention_duration_ms = Some(*duration_ms);
2399 }
2400 AlterOperation::UnsetRetention => {
2401 contract.retention_duration_ms = None;
2402 }
2403 AlterOperation::AddAnalytics(views) => {
2407 for view in views {
2408 if !contract
2412 .analytics_config
2413 .iter()
2414 .any(|existing| existing.output == view.output)
2415 {
2416 contract.analytics_config.push(view.clone());
2417 }
2418 }
2419 }
2420 AlterOperation::DropAnalytics(output) => {
2421 contract
2422 .analytics_config
2423 .retain(|view| view.output != *output);
2424 }
2425 }
2426 }
2427}
2428
2429pub(crate) fn retention_timestamp_column_exists(
2434 contract: &crate::physical::CollectionContract,
2435) -> bool {
2436 if contract.timestamps_enabled {
2437 return true;
2438 }
2439 if matches!(
2440 contract.declared_model,
2441 crate::catalog::CollectionModel::TimeSeries | crate::catalog::CollectionModel::Metrics
2442 ) {
2443 return true;
2447 }
2448 contract
2449 .declared_columns
2450 .iter()
2451 .any(|column| is_temporal_data_type(&column.data_type))
2452}
2453
2454fn is_temporal_data_type(data_type: &str) -> bool {
2455 let upper = data_type.to_ascii_uppercase();
2456 matches!(
2457 upper.as_str(),
2458 "TIMESTAMP" | "TIMESTAMPMS" | "TIMESTAMP_MS" | "DATETIME" | "DATE"
2459 )
2460}
2461
2462fn validate_event_subscriptions(
2463 runtime: &RedDBRuntime,
2464 source: &str,
2465 subscriptions: &[crate::catalog::SubscriptionDescriptor],
2466) -> RedDBResult<()> {
2467 for subscription in subscriptions
2468 .iter()
2469 .filter(|subscription| subscription.enabled)
2470 {
2471 if subscription.all_tenants && crate::runtime::impl_core::current_tenant().is_some() {
2472 return Err(RedDBError::Query(
2473 "cross-tenant subscription requires cluster-admin capability (events:cluster_subscribe)".to_string(),
2474 ));
2475 }
2476 validate_subscription_auth(runtime, source, subscription)?;
2477 if subscription.target_queue == source
2478 || subscription_would_create_cycle(
2479 &runtime.inner.db,
2480 source,
2481 &subscription.target_queue,
2482 )
2483 {
2484 return Err(RedDBError::Query(
2485 "subscription would create cycle".to_string(),
2486 ));
2487 }
2488 audit_subscription_redact_gap(runtime, source, subscription);
2489 }
2490 Ok(())
2491}
2492
2493fn validate_subscription_auth(
2494 runtime: &RedDBRuntime,
2495 source: &str,
2496 subscription: &crate::catalog::SubscriptionDescriptor,
2497) -> RedDBResult<()> {
2498 let auth_store = match runtime.inner.auth_store.read().clone() {
2499 Some(store) => store,
2500 None => return Ok(()),
2501 };
2502 let (username, role) = match crate::runtime::impl_core::current_auth_identity() {
2503 Some(identity) => identity,
2504 None => return Ok(()),
2505 };
2506 let tenant = crate::runtime::impl_core::current_tenant();
2507 let principal = crate::auth::UserId::from_parts(tenant.as_deref(), &username);
2508
2509 if auth_store.iam_authorization_enabled() {
2510 let ctx = crate::auth::policies::EvalContext {
2511 principal_tenant: tenant.clone(),
2512 current_tenant: tenant.clone(),
2513 peer_ip: None,
2514 mfa_present: false,
2515 now_ms: crate::auth::now_ms(),
2516 principal_is_admin_role: role == crate::auth::Role::Admin,
2517 principal_is_platform_scoped: principal.tenant.is_none(),
2518 };
2519 let mut source_resource = crate::auth::policies::ResourceRef::new("table", source);
2520 if let Some(t) = tenant.as_deref() {
2521 source_resource = source_resource.with_tenant(t.to_string());
2522 }
2523 if !auth_store.check_policy_authz_with_role(
2524 &principal,
2525 "select",
2526 &source_resource,
2527 &ctx,
2528 role,
2529 ) {
2530 return Err(RedDBError::Query(format!(
2531 "permission denied: principal=`{}` action=`select` resource=`{}:{}` denied by IAM policy",
2532 principal, source_resource.kind, source_resource.name
2533 )));
2534 }
2535
2536 let mut target_resource =
2537 crate::auth::policies::ResourceRef::new("queue", subscription.target_queue.clone());
2538 if let Some(t) = tenant.as_deref() {
2539 target_resource = target_resource.with_tenant(t.to_string());
2540 }
2541 if !auth_store.check_policy_authz_with_role(
2542 &principal,
2543 "write",
2544 &target_resource,
2545 &ctx,
2546 role,
2547 ) {
2548 return Err(RedDBError::Query(format!(
2549 "permission denied: principal=`{}` action=`write` resource=`{}:{}` denied by IAM policy",
2550 principal, target_resource.kind, target_resource.name
2551 )));
2552 }
2553 return Ok(());
2554 }
2555
2556 let ctx = crate::auth::privileges::AuthzContext {
2557 principal: &username,
2558 effective_role: role,
2559 tenant: tenant.as_deref(),
2560 };
2561 auth_store
2562 .check_grant(
2563 &ctx,
2564 crate::auth::privileges::Action::Select,
2565 &crate::auth::privileges::Resource::table_from_name(source),
2566 )
2567 .map_err(|err| RedDBError::Query(format!("permission denied: {err}")))?;
2568 auth_store
2569 .check_grant(
2570 &ctx,
2571 crate::auth::privileges::Action::Insert,
2572 &crate::auth::privileges::Resource::table_from_name(&subscription.target_queue),
2573 )
2574 .map_err(|err| RedDBError::Query(format!("permission denied: {err}")))?;
2575 Ok(())
2576}
2577
2578fn audit_subscription_redact_gap(
2579 runtime: &RedDBRuntime,
2580 source: &str,
2581 subscription: &crate::catalog::SubscriptionDescriptor,
2582) {
2583 let auth_store = match runtime.inner.auth_store.read().clone() {
2584 Some(store) if store.iam_authorization_enabled() => store,
2585 _ => return,
2586 };
2587 let (username, role) = match crate::runtime::impl_core::current_auth_identity() {
2588 Some(identity) => identity,
2589 None => return,
2590 };
2591 let tenant = crate::runtime::impl_core::current_tenant();
2592 let principal = crate::auth::UserId::from_parts(tenant.as_deref(), &username);
2593 let missing = subscription_redact_gap_columns(&auth_store, &principal, source, subscription);
2594 if missing.is_empty() {
2595 return;
2596 }
2597
2598 let columns = missing.into_iter().collect::<Vec<_>>().join(", ");
2599 tracing::warn!(
2600 target: "reddb::operator",
2601 "subscription_redact_gap: source={} target_queue={} columns=[{}]",
2602 source,
2603 subscription.target_queue,
2604 columns
2605 );
2606 let mut event = AuditEvent::builder("subscription_redact_gap")
2607 .principal(username)
2608 .source(AuditAuthSource::System)
2609 .resource(format!(
2610 "subscription:{}->{}",
2611 source, subscription.target_queue
2612 ))
2613 .outcome(Outcome::Success)
2614 .field(AuditFieldEscaper::field("source", source))
2615 .field(AuditFieldEscaper::field(
2616 "target_queue",
2617 subscription.target_queue.clone(),
2618 ))
2619 .field(AuditFieldEscaper::field(
2620 "subscription",
2621 subscription.name.clone(),
2622 ))
2623 .field(AuditFieldEscaper::field("columns", columns))
2624 .field(AuditFieldEscaper::field("role", role.as_str()));
2625 if let Some(t) = tenant {
2626 event = event.tenant(t);
2627 }
2628 runtime.inner.audit_log.record_event(event.build());
2629}
2630
2631fn subscription_redact_gap_columns(
2632 auth_store: &crate::auth::store::AuthStore,
2633 principal: &crate::auth::UserId,
2634 source: &str,
2635 subscription: &crate::catalog::SubscriptionDescriptor,
2636) -> BTreeSet<String> {
2637 let redacted: HashSet<String> = subscription
2638 .redact_fields
2639 .iter()
2640 .map(|field| field.to_ascii_lowercase())
2641 .collect();
2642 auth_store
2643 .effective_policies(principal)
2644 .iter()
2645 .flat_map(|policy| policy.statements.iter())
2646 .filter(|statement| statement.effect == crate::auth::policies::Effect::Deny)
2647 .filter(|statement| statement.actions.iter().any(action_pattern_matches_select))
2648 .flat_map(|statement| statement.resources.iter())
2649 .filter_map(|resource| denied_column_for_source(resource, source))
2650 .filter(|column| !redact_covers_column(&redacted, source, column))
2651 .collect()
2652}
2653
2654fn action_pattern_matches_select(pattern: &crate::auth::policies::ActionPattern) -> bool {
2655 match pattern {
2656 crate::auth::policies::ActionPattern::Wildcard => true,
2657 crate::auth::policies::ActionPattern::Exact(action) => action == "select",
2658 crate::auth::policies::ActionPattern::Prefix(prefix) => {
2659 "select".len() > prefix.len() + 1
2660 && "select".starts_with(prefix)
2661 && "select".as_bytes()[prefix.len()] == b':'
2662 }
2663 }
2664}
2665
2666fn denied_column_for_source(
2667 resource: &crate::auth::policies::ResourcePattern,
2668 source: &str,
2669) -> Option<String> {
2670 let crate::auth::policies::ResourcePattern::Exact { kind, name } = resource else {
2671 return None;
2672 };
2673 if kind != "column" {
2674 return None;
2675 }
2676 let column = crate::auth::ColumnRef::parse_resource_name(name).ok()?;
2677 (column.table_resource_name() == source).then_some(column.column)
2678}
2679
2680fn redact_covers_column(redacted: &HashSet<String>, source: &str, column: &str) -> bool {
2681 let column = column.to_ascii_lowercase();
2682 let qualified = format!("{}.{}", source.to_ascii_lowercase(), column);
2683 redacted.contains("*") || redacted.contains(&column) || redacted.contains(&qualified)
2684}
2685
2686fn subscription_would_create_cycle(
2687 db: &crate::storage::unified::devx::RedDB,
2688 source: &str,
2689 target: &str,
2690) -> bool {
2691 let mut graph: HashMap<String, Vec<String>> = HashMap::new();
2692 for contract in db.collection_contracts() {
2693 for subscription in contract
2694 .subscriptions
2695 .into_iter()
2696 .filter(|subscription| subscription.enabled)
2697 {
2698 graph
2699 .entry(subscription.source)
2700 .or_default()
2701 .push(subscription.target_queue);
2702 }
2703 }
2704 graph
2705 .entry(source.to_string())
2706 .or_default()
2707 .push(target.to_string());
2708
2709 let mut stack = vec![target.to_string()];
2710 let mut seen = HashSet::new();
2711 while let Some(node) = stack.pop() {
2712 if node == source {
2713 return true;
2714 }
2715 if !seen.insert(node.clone()) {
2716 continue;
2717 }
2718 if let Some(next) = graph.get(&node) {
2719 stack.extend(next.iter().cloned());
2720 }
2721 }
2722 false
2723}
2724
2725pub(crate) fn ensure_event_target_queue_pub(
2726 runtime: &RedDBRuntime,
2727 queue: &str,
2728) -> RedDBResult<()> {
2729 ensure_event_target_queue(runtime, queue)
2730}
2731
2732fn ensure_event_target_queue(runtime: &RedDBRuntime, queue: &str) -> RedDBResult<()> {
2733 let store = runtime.inner.db.store();
2734 if store.get_collection(queue).is_some() {
2735 return Ok(());
2736 }
2737 store
2738 .create_collection(queue)
2739 .map_err(|err| RedDBError::Internal(err.to_string()))?;
2740 runtime
2741 .inner
2742 .db
2743 .save_collection_contract(event_queue_collection_contract(queue))
2744 .map_err(|err| RedDBError::Internal(err.to_string()))?;
2745 store.set_config_tree(
2746 &format!("queue.{queue}.mode"),
2747 &crate::serde_json::Value::String("fanout".to_string()),
2748 );
2749 Ok(())
2750}
2751
2752fn event_queue_collection_contract(queue: &str) -> crate::physical::CollectionContract {
2753 let now = current_unix_ms();
2754 crate::physical::CollectionContract {
2755 name: queue.to_string(),
2756 declared_model: crate::catalog::CollectionModel::Queue,
2757 schema_mode: crate::catalog::SchemaMode::Dynamic,
2758 origin: crate::physical::ContractOrigin::Implicit,
2759 version: 1,
2760 created_at_unix_ms: now,
2761 updated_at_unix_ms: now,
2762 default_ttl_ms: None,
2763 vector_dimension: None,
2764 vector_metric: None,
2765 context_index_fields: Vec::new(),
2766 declared_columns: Vec::new(),
2767 table_def: None,
2768 timestamps_enabled: false,
2769 context_index_enabled: false,
2770 metrics_raw_retention_ms: None,
2771 metrics_rollup_policies: Vec::new(),
2772 metrics_tenant_identity: None,
2773 metrics_namespace: None,
2774 append_only: true,
2775 subscriptions: Vec::new(),
2776 analytics_config: Vec::new(),
2777 session_key: None,
2778 session_gap_ms: None,
2779 retention_duration_ms: None,
2780 analytical_storage: None,
2781
2782 ai_policy: None,
2783 }
2784}
2785
2786fn build_table_def_from_create_table(
2787 query: &CreateTableQuery,
2788) -> RedDBResult<crate::storage::schema::TableDef> {
2789 let mut table = crate::storage::schema::TableDef::new(query.name.clone());
2790 for column in &query.columns {
2791 if column.primary_key {
2792 table.primary_key.push(column.name.clone());
2793 table.constraints.push(
2794 crate::storage::schema::Constraint::new(
2795 format!("pk_{}", column.name),
2796 crate::storage::schema::ConstraintType::PrimaryKey,
2797 )
2798 .on_columns(vec![column.name.clone()]),
2799 );
2800 }
2801 if column.unique {
2802 table.constraints.push(
2803 crate::storage::schema::Constraint::new(
2804 format!("uniq_{}", column.name),
2805 crate::storage::schema::ConstraintType::Unique,
2806 )
2807 .on_columns(vec![column.name.clone()]),
2808 );
2809 }
2810 if column.not_null {
2811 table.constraints.push(
2812 crate::storage::schema::Constraint::new(
2813 format!("not_null_{}", column.name),
2814 crate::storage::schema::ConstraintType::NotNull,
2815 )
2816 .on_columns(vec![column.name.clone()]),
2817 );
2818 }
2819 table.columns.push(column_def_from_ddl(column)?);
2820 }
2821 if query.timestamps {
2826 table.columns.push(
2827 crate::storage::schema::ColumnDef::new(
2828 "created_at".to_string(),
2829 crate::storage::schema::DataType::UnsignedInteger,
2830 )
2831 .not_null(),
2832 );
2833 table.columns.push(
2834 crate::storage::schema::ColumnDef::new(
2835 "updated_at".to_string(),
2836 crate::storage::schema::DataType::UnsignedInteger,
2837 )
2838 .not_null(),
2839 );
2840 table.constraints.push(
2841 crate::storage::schema::Constraint::new(
2842 "not_null_created_at".to_string(),
2843 crate::storage::schema::ConstraintType::NotNull,
2844 )
2845 .on_columns(vec!["created_at".to_string()]),
2846 );
2847 table.constraints.push(
2848 crate::storage::schema::Constraint::new(
2849 "not_null_updated_at".to_string(),
2850 crate::storage::schema::ConstraintType::NotNull,
2851 )
2852 .on_columns(vec!["updated_at".to_string()]),
2853 );
2854 }
2855 table
2856 .validate()
2857 .map_err(|err| RedDBError::Query(format!("invalid table definition: {err}")))?;
2858 Ok(table)
2859}
2860
2861fn column_def_from_ddl(column: &CreateColumnDef) -> RedDBResult<crate::storage::schema::ColumnDef> {
2862 let data_type = resolve_declared_data_type(&column.data_type)
2863 .map_err(|err| RedDBError::Query(err.to_string()))?;
2864 let mut column_def = crate::storage::schema::ColumnDef::new(column.name.clone(), data_type);
2865 if column.not_null {
2866 column_def = column_def.not_null();
2867 }
2868 if let Some(default) = &column.default {
2869 column_def = column_def.with_default(default.as_bytes().to_vec());
2870 }
2871 if column.compress.unwrap_or(0) > 0 {
2872 column_def = column_def.compressed();
2873 }
2874 if !column.enum_variants.is_empty() {
2875 column_def = column_def.with_variants(column.enum_variants.clone());
2876 }
2877 if let Some(precision) = column.decimal_precision {
2878 column_def = column_def.with_precision(precision);
2879 }
2880 if let Some(element_type) = &column.array_element {
2881 column_def = column_def.with_element_type(
2882 resolve_declared_data_type(element_type)
2883 .map_err(|err| RedDBError::Query(err.to_string()))?,
2884 );
2885 }
2886 column_def = column_def.with_metadata("ddl_data_type", column.data_type.clone());
2887 if column.unique {
2888 column_def = column_def.with_metadata("unique", "true");
2889 }
2890 if column.primary_key {
2891 column_def = column_def.with_metadata("primary_key", "true");
2892 }
2893 Ok(column_def)
2894}
2895
2896fn current_unix_ms() -> u128 {
2897 std::time::SystemTime::now()
2898 .duration_since(std::time::UNIX_EPOCH)
2899 .unwrap_or_default()
2900 .as_millis()
2901}
2902
2903#[cfg(test)]
2904mod tests {
2905 use crate::auth::policies::{ActionPattern, Effect, Policy, ResourcePattern, Statement};
2906 use crate::auth::store::{AuthStore, PrincipalRef};
2907 use crate::auth::UserId;
2908 use crate::auth::{AuthConfig, Role};
2909 use crate::runtime::impl_core::{clear_current_auth_identity, set_current_auth_identity};
2910 use crate::storage::schema::Value;
2911 use crate::{RedDBOptions, RedDBRuntime};
2912 use std::sync::Arc;
2913
2914 fn make_allow_policy(id: &str, action: &str, collection: &str) -> Policy {
2915 Policy {
2916 id: id.to_string(),
2917 version: 1,
2918 tenant: None,
2919 created_at: 0,
2920 updated_at: 0,
2921 statements: vec![Statement {
2922 sid: None,
2923 effect: Effect::Allow,
2924 actions: vec![ActionPattern::Exact(action.to_string())],
2925 resources: vec![ResourcePattern::Exact {
2926 kind: "collection".to_string(),
2927 name: collection.to_string(),
2928 }],
2929 condition: None,
2930 }],
2931 }
2932 }
2933
2934 fn wire_auth_store(rt: &RedDBRuntime) -> Arc<AuthStore> {
2935 let store = Arc::new(AuthStore::new(AuthConfig::default()));
2936 *rt.inner.auth_store.write() = Some(store.clone());
2937 store
2938 }
2939
2940 #[test]
2941 fn drop_denied_without_iam_policy() {
2942 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
2943 rt.execute_query("CREATE TABLE foo (id INT)").unwrap();
2944 let store = wire_auth_store(&rt);
2945 let select_only = Policy {
2947 id: "select-only".to_string(),
2948 version: 1,
2949 tenant: None,
2950 created_at: 0,
2951 updated_at: 0,
2952 statements: vec![Statement {
2953 sid: None,
2954 effect: Effect::Allow,
2955 actions: vec![ActionPattern::Exact("select".to_string())],
2956 resources: vec![ResourcePattern::Wildcard],
2957 condition: None,
2958 }],
2959 };
2960 store.put_policy_internal(select_only).unwrap();
2961 let alice = UserId::from_parts(None, "alice");
2962 store
2963 .attach_policy(PrincipalRef::User(alice), "select-only")
2964 .unwrap();
2965 set_current_auth_identity("alice".to_string(), Role::Write);
2966 let err = rt.execute_query("DROP TABLE foo").unwrap_err();
2967 clear_current_auth_identity();
2968 assert!(
2969 format!("{err}").contains("denied by IAM policy"),
2970 "got: {err}"
2971 );
2972 }
2973
2974 #[test]
2975 fn drop_allowed_with_explicit_iam_policy() {
2976 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
2977 rt.execute_query("CREATE TABLE bar (id INT)").unwrap();
2978 let store = wire_auth_store(&rt);
2979 let policy = make_allow_policy("allow-drop-bar", "drop", "bar");
2980 store.put_policy_internal(policy).unwrap();
2981 let bob = UserId::from_parts(None, "bob");
2982 store
2983 .attach_policy(PrincipalRef::User(bob), "allow-drop-bar")
2984 .unwrap();
2985 set_current_auth_identity("bob".to_string(), Role::Write);
2986 rt.execute_query("DROP TABLE bar").unwrap();
2987 clear_current_auth_identity();
2988 }
2989
2990 #[test]
2991 fn drop_allowed_with_wildcard_iam_policy() {
2992 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
2993 rt.execute_query("CREATE TABLE baz (id INT)").unwrap();
2994 let store = wire_auth_store(&rt);
2995 let policy = Policy {
2996 id: "allow-drop-all".to_string(),
2997 version: 1,
2998 tenant: None,
2999 created_at: 0,
3000 updated_at: 0,
3001 statements: vec![Statement {
3002 sid: None,
3003 effect: Effect::Allow,
3004 actions: vec![ActionPattern::Exact("drop".to_string())],
3005 resources: vec![ResourcePattern::Wildcard],
3006 condition: None,
3007 }],
3008 };
3009 store.put_policy_internal(policy).unwrap();
3010 let carl = UserId::from_parts(None, "carl");
3011 store
3012 .attach_policy(PrincipalRef::User(carl), "allow-drop-all")
3013 .unwrap();
3014 set_current_auth_identity("carl".to_string(), Role::Write);
3015 rt.execute_query("DROP TABLE baz").unwrap();
3016 clear_current_auth_identity();
3017 }
3018
3019 #[test]
3020 fn truncate_denied_without_iam_policy() {
3021 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3022 rt.execute_query("CREATE TABLE qux (id INT)").unwrap();
3023 let store = wire_auth_store(&rt);
3024 store
3031 .set_enforcement_mode(crate::auth::enforcement_mode::PolicyEnforcementMode::PolicyOnly);
3032 let select_only = Policy {
3034 id: "select-only-2".to_string(),
3035 version: 1,
3036 tenant: None,
3037 created_at: 0,
3038 updated_at: 0,
3039 statements: vec![Statement {
3040 sid: None,
3041 effect: Effect::Allow,
3042 actions: vec![ActionPattern::Exact("select".to_string())],
3043 resources: vec![ResourcePattern::Wildcard],
3044 condition: None,
3045 }],
3046 };
3047 store.put_policy_internal(select_only).unwrap();
3048 let dana = UserId::from_parts(None, "dana");
3049 store
3050 .attach_policy(PrincipalRef::User(dana), "select-only-2")
3051 .unwrap();
3052 set_current_auth_identity("dana".to_string(), Role::Write);
3053 let err = rt.execute_query("TRUNCATE TABLE qux").unwrap_err();
3054 clear_current_auth_identity();
3055 assert!(
3056 format!("{err}").contains("denied by IAM policy"),
3057 "got: {err}"
3058 );
3059 }
3060
3061 #[test]
3062 fn truncate_table_clears_rows_and_preserves_schema_and_indexes() {
3063 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3064 rt.execute_query("CREATE TABLE users (id INT, name TEXT)")
3065 .unwrap();
3066 rt.execute_query("INSERT INTO users (id, name) VALUES (1, 'ana'), (2, 'bob')")
3067 .unwrap();
3068 rt.execute_query("CREATE INDEX idx_users_id ON users (id) USING HASH")
3069 .unwrap();
3070
3071 let truncated = rt.execute_query("TRUNCATE TABLE users").unwrap();
3072 assert_eq!(truncated.statement_type, "truncate");
3073 assert_eq!(truncated.affected_rows, 0);
3074
3075 let empty = rt.execute_query("SELECT id FROM users").unwrap();
3076 assert!(empty.result.records.is_empty());
3077
3078 rt.execute_query("INSERT INTO users (id, name) VALUES (3, 'cy')")
3079 .unwrap();
3080 let selected = rt
3081 .execute_query("SELECT name FROM users WHERE id = 3")
3082 .unwrap();
3083 let name = selected.result.records[0].get("name").unwrap();
3084 assert_eq!(name, &Value::text("cy"));
3085 assert!(rt.db().collection_contract("users").is_some());
3086 assert!(rt
3087 .inner
3088 .index_store
3089 .list_indices("users")
3090 .iter()
3091 .any(|index| index.name == "idx_users_id"));
3092 }
3093
3094 #[test]
3095 fn truncate_collection_is_polymorphic_and_typed_mismatch_fails() {
3096 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3097 rt.execute_query("CREATE QUEUE tasks").unwrap();
3098 rt.execute_query("QUEUE PUSH tasks {'job':'a'}").unwrap();
3099
3100 let err = rt.execute_query("TRUNCATE TABLE tasks").unwrap_err();
3101 assert!(format!("{err}").contains("model mismatch: expected table, got queue"));
3102
3103 rt.execute_query("TRUNCATE COLLECTION tasks").unwrap();
3104 let len = rt.execute_query("QUEUE LEN tasks").unwrap();
3105 assert_eq!(
3106 len.result.records[0].get("len"),
3107 Some(&Value::UnsignedInteger(0))
3108 );
3109 }
3110
3111 #[test]
3112 fn truncate_system_schema_is_read_only() {
3113 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3114 let err = rt
3115 .execute_query("TRUNCATE COLLECTION red.collections")
3116 .unwrap_err();
3117 assert!(format!("{err}").contains("system schema is read-only"));
3118 }
3119
3120 fn queue_payloads(rt: &RedDBRuntime, queue: &str) -> Vec<crate::json::Value> {
3123 let result = rt
3124 .execute_query(&format!("QUEUE PEEK {queue} 100"))
3125 .expect("peek queue");
3126 result
3127 .result
3128 .records
3129 .iter()
3130 .map(
3131 |record| match record.get("payload").expect("payload column") {
3132 Value::Json(bytes) => crate::json::from_slice(bytes).expect("json payload"),
3133 other => panic!("expected JSON queue payload, got {other:?}"),
3134 },
3135 )
3136 .collect()
3137 }
3138
3139 #[test]
3142 fn truncate_event_enabled_table_emits_single_truncate_event() {
3143 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3144 rt.execute_query("CREATE TABLE users (id INT, name TEXT) WITH EVENTS TO users_events")
3145 .unwrap();
3146 rt.execute_query(
3147 "INSERT INTO users (id, name) VALUES (1, 'alice'), (2, 'bob'), (3, 'carol')",
3148 )
3149 .unwrap();
3150
3151 rt.execute_query("QUEUE POP users_events COUNT 10").unwrap();
3153
3154 rt.execute_query("TRUNCATE TABLE users").unwrap();
3155
3156 let events = queue_payloads(&rt, "users_events");
3157 assert_eq!(
3159 events.len(),
3160 1,
3161 "expected 1 truncate event, got {}",
3162 events.len()
3163 );
3164 let ev = events[0].as_object().expect("event is object");
3165 assert_eq!(
3166 ev.get("op").and_then(crate::json::Value::as_str),
3167 Some("truncate")
3168 );
3169 assert_eq!(
3170 ev.get("collection").and_then(crate::json::Value::as_str),
3171 Some("users")
3172 );
3173 assert_eq!(
3174 ev.get("entities_count")
3175 .and_then(crate::json::Value::as_u64),
3176 Some(3)
3177 );
3178 assert!(ev.get("ts").and_then(crate::json::Value::as_u64).is_some());
3179 assert!(ev.get("lsn").and_then(crate::json::Value::as_u64).is_some());
3180 assert!(ev
3181 .get("event_id")
3182 .and_then(crate::json::Value::as_str)
3183 .is_some_and(|s| !s.is_empty()));
3184 }
3185
3186 #[test]
3188 fn truncate_no_events_collection_emits_nothing() {
3189 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3190 rt.execute_query("CREATE TABLE plain (id INT, val TEXT)")
3191 .unwrap();
3192 rt.execute_query("INSERT INTO plain (id, val) VALUES (1, 'a'), (2, 'b')")
3193 .unwrap();
3194 rt.execute_query("TRUNCATE TABLE plain").unwrap();
3196 let rows = rt.execute_query("SELECT id FROM plain").unwrap();
3198 assert!(rows.result.records.is_empty());
3199 }
3200
3201 #[test]
3205 fn drop_event_enabled_table_emits_single_collection_dropped_event() {
3206 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3207 rt.execute_query("CREATE TABLE users (id INT, name TEXT) WITH EVENTS TO users_events")
3208 .unwrap();
3209 rt.execute_query("INSERT INTO users (id, name) VALUES (1, 'alice'), (2, 'bob')")
3210 .unwrap();
3211
3212 rt.execute_query("QUEUE POP users_events COUNT 10").unwrap();
3214
3215 rt.execute_query("DROP TABLE users").unwrap();
3216
3217 let events = queue_payloads(&rt, "users_events");
3219 assert_eq!(
3220 events.len(),
3221 1,
3222 "expected 1 collection_dropped event, got {}",
3223 events.len()
3224 );
3225 let ev = events[0].as_object().expect("event is object");
3226 assert_eq!(
3227 ev.get("op").and_then(crate::json::Value::as_str),
3228 Some("collection_dropped")
3229 );
3230 assert_eq!(
3231 ev.get("collection").and_then(crate::json::Value::as_str),
3232 Some("users")
3233 );
3234 assert_eq!(
3235 ev.get("final_entities_count")
3236 .and_then(crate::json::Value::as_u64),
3237 Some(2)
3238 );
3239 assert!(ev.get("ts").and_then(crate::json::Value::as_u64).is_some());
3240 assert!(ev.get("lsn").and_then(crate::json::Value::as_u64).is_some());
3241 assert!(ev
3242 .get("event_id")
3243 .and_then(crate::json::Value::as_str)
3244 .is_some_and(|s| !s.is_empty()));
3245
3246 let err = rt.execute_query("SELECT id FROM users").unwrap_err();
3248 assert!(
3249 format!("{err}").contains("users"),
3250 "expected not-found error"
3251 );
3252 }
3253
3254 #[test]
3257 fn drop_no_events_collection_emits_nothing() {
3258 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3259 rt.execute_query("CREATE TABLE plain (id INT, val TEXT)")
3260 .unwrap();
3261 rt.execute_query("INSERT INTO plain (id, val) VALUES (1, 'a')")
3262 .unwrap();
3263 rt.execute_query("DROP TABLE plain").unwrap();
3264 let err = rt.execute_query("SELECT id FROM plain").unwrap_err();
3266 assert!(format!("{err}").contains("plain"));
3267 }
3268
3269 #[test]
3273 fn ops_filter_insert_only_ignores_update_and_delete() {
3274 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3275 rt.execute_query(
3276 "CREATE TABLE items (id INT, val TEXT) WITH EVENTS (INSERT) TO items_events",
3277 )
3278 .unwrap();
3279 rt.execute_query("INSERT INTO items (id, val) VALUES (1, 'a')")
3280 .unwrap();
3281 rt.execute_query("UPDATE items SET val = 'b' WHERE id = 1")
3282 .unwrap();
3283 rt.execute_query("DELETE FROM items WHERE id = 1").unwrap();
3284
3285 let events = queue_payloads(&rt, "items_events");
3286 assert_eq!(
3288 events.len(),
3289 1,
3290 "expected 1 insert event, got {}",
3291 events.len()
3292 );
3293 assert_eq!(
3294 events[0]
3295 .as_object()
3296 .unwrap()
3297 .get("op")
3298 .and_then(crate::json::Value::as_str),
3299 Some("insert")
3300 );
3301 }
3302
3303 #[test]
3305 fn where_filter_skips_rows_that_do_not_match() {
3306 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3307 rt.execute_query(
3308 "CREATE TABLE users (id INT, status TEXT) WITH EVENTS WHERE status = 'active' TO users_events",
3309 )
3310 .unwrap();
3311
3312 rt.execute_query("INSERT INTO users (id, status) VALUES (1, 'active')")
3314 .unwrap();
3315 rt.execute_query("INSERT INTO users (id, status) VALUES (2, 'inactive')")
3317 .unwrap();
3318
3319 let events = queue_payloads(&rt, "users_events");
3320 assert_eq!(
3321 events.len(),
3322 1,
3323 "expected 1 event (only active), got {}",
3324 events.len()
3325 );
3326 let ev = events[0].as_object().unwrap();
3327 assert_eq!(
3328 ev.get("op").and_then(crate::json::Value::as_str),
3329 Some("insert")
3330 );
3331 let after = ev.get("after").unwrap().as_object().unwrap();
3332 assert_eq!(
3333 after.get("status").and_then(crate::json::Value::as_str),
3334 Some("active")
3335 );
3336 }
3337
3338 #[test]
3340 fn ops_filter_and_where_filter_combined() {
3341 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3342 rt.execute_query(
3343 "CREATE TABLE items (id INT, status TEXT) WITH EVENTS (INSERT, UPDATE) WHERE status = 'active' TO items_events",
3344 )
3345 .unwrap();
3346
3347 rt.execute_query("INSERT INTO items (id, status) VALUES (1, 'active')")
3349 .unwrap();
3350 rt.execute_query("INSERT INTO items (id, status) VALUES (2, 'inactive')")
3352 .unwrap();
3353 rt.execute_query("UPDATE items SET status = 'inactive' WHERE id = 1")
3355 .unwrap();
3356 rt.execute_query("DELETE FROM items WHERE id = 2").unwrap();
3358
3359 let events = queue_payloads(&rt, "items_events");
3360 assert_eq!(
3362 events.len(),
3363 1,
3364 "expected 1 event, got {}: {events:?}",
3365 events.len()
3366 );
3367 assert_eq!(
3368 events[0]
3369 .as_object()
3370 .unwrap()
3371 .get("op")
3372 .and_then(crate::json::Value::as_str),
3373 Some("insert")
3374 );
3375 }
3376
3377 #[test]
3379 fn where_filter_on_delete_checks_before_state() {
3380 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3381 rt.execute_query(
3382 "CREATE TABLE users (id INT, status TEXT) WITH EVENTS (DELETE) WHERE status = 'active' TO users_events",
3383 )
3384 .unwrap();
3385
3386 rt.execute_query("INSERT INTO users (id, status) VALUES (1, 'active'), (2, 'inactive')")
3387 .unwrap();
3388
3389 rt.execute_query("DELETE FROM users WHERE id = 1").unwrap();
3391 rt.execute_query("DELETE FROM users WHERE id = 2").unwrap();
3393
3394 let events = queue_payloads(&rt, "users_events");
3395 assert_eq!(
3396 events.len(),
3397 1,
3398 "expected 1 delete event, got {}",
3399 events.len()
3400 );
3401 let ev = events[0].as_object().unwrap();
3402 assert_eq!(
3403 ev.get("op").and_then(crate::json::Value::as_str),
3404 Some("delete")
3405 );
3406 }
3407
3408 #[test]
3412 fn alter_add_column_on_event_enabled_table_succeeds() {
3413 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3414 rt.execute_query("CREATE TABLE users (id INT, name TEXT) WITH EVENTS TO users_events")
3415 .unwrap();
3416 rt.execute_query("ALTER TABLE users ADD COLUMN phone TEXT")
3418 .unwrap();
3419 let contract = rt.db().collection_contract("users").unwrap();
3421 assert!(
3422 contract.declared_columns.iter().any(|c| c.name == "phone"),
3423 "phone column should be in contract"
3424 );
3425 assert!(
3427 contract.subscriptions.iter().any(|s| s.enabled),
3428 "subscription should remain enabled"
3429 );
3430 }
3431
3432 #[test]
3435 fn alter_drop_column_and_rls_on_event_enabled_table_succeeds() {
3436 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3437 rt.execute_query(
3438 "CREATE TABLE items (id INT, secret TEXT, status TEXT) WITH EVENTS TO items_events",
3439 )
3440 .unwrap();
3441 rt.execute_query("ALTER TABLE items DROP COLUMN secret")
3443 .unwrap();
3444 let contract = rt.db().collection_contract("items").unwrap();
3445 assert!(
3446 !contract.declared_columns.iter().any(|c| c.name == "secret"),
3447 "secret column should be removed"
3448 );
3449 rt.execute_query("ALTER TABLE items ENABLE ROW LEVEL SECURITY")
3451 .unwrap();
3452 assert!(
3454 contract.subscriptions.iter().any(|s| s.enabled),
3455 "subscription should remain enabled"
3456 );
3457 }
3458
3459 #[test]
3465 fn create_vector_marks_collection_as_turbo_baseline() {
3466 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3467 rt.execute_query("CREATE VECTOR embeddings DIM 4").unwrap();
3468 let store = rt.db().store();
3469 assert!(
3470 crate::runtime::vector_turbo_kind::is_turbo(&store, "embeddings"),
3471 "new vector collections must be turbo-marked baseline"
3472 );
3473 assert!(
3474 rt.db().turbo_state("embeddings").is_some(),
3475 "turbo_state must materialise after CREATE VECTOR"
3476 );
3477 }
3478
3479 #[test]
3483 fn create_table_does_not_mark_turbo() {
3484 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3485 rt.execute_query("CREATE TABLE plain (id INT)").unwrap();
3486 let store = rt.db().store();
3487 assert!(
3488 !crate::runtime::vector_turbo_kind::is_turbo(&store, "plain"),
3489 "non-vector collections must not gain the turbo marker"
3490 );
3491 assert!(rt.db().turbo_state("plain").is_none());
3492 }
3493
3494 #[test]
3497 fn create_collection_kind_vector_turbo_still_marked() {
3498 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3499 rt.execute_query("CREATE COLLECTION turbo_v KIND vector.turbo DIM 4")
3500 .unwrap();
3501 let store = rt.db().store();
3502 assert!(crate::runtime::vector_turbo_kind::is_turbo(
3503 &store, "turbo_v"
3504 ));
3505 assert!(rt.db().turbo_state("turbo_v").is_some());
3506 }
3507
3508 #[test]
3511 fn create_table_persists_and_introspects_ai_policy() {
3512 use crate::catalog::{ModerateDegradedMode, ModerateRejectAction};
3513 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3514 rt.execute_query(
3515 "CREATE TABLE posts (id INT, title TEXT, body TEXT, photo TEXT) WITH ( \
3516 EMBED (fields = ('title', 'body'), provider = 'openai', model = 'text-embedding-3-small'), \
3517 MODERATE (fields = ('body'), provider = 'openai', model = 'omni-moderation-latest', sync = true, degraded = closed, on_reject = flag), \
3518 VISION (image_field = 'photo', outputs = ('caption'), provider = 'openai', model = 'gpt-4o') \
3519 )",
3520 )
3521 .expect("create table with ai policy");
3522
3523 let contracts = rt.db().collection_contracts();
3524 let contract = contracts
3525 .iter()
3526 .find(|c| c.name == "posts")
3527 .expect("contract persisted");
3528 let policy = contract.ai_policy.as_ref().expect("ai policy persisted");
3529
3530 let embed = policy.embed.as_ref().expect("embed");
3531 assert_eq!(embed.fields, vec!["title".to_string(), "body".to_string()]);
3532 assert_eq!(embed.model, "text-embedding-3-small");
3533
3534 let moderate = policy.moderate.as_ref().expect("moderate");
3535 assert!(moderate.sync_gate);
3536 assert_eq!(moderate.degraded_mode, ModerateDegradedMode::Closed);
3537 assert_eq!(moderate.reject_action, ModerateRejectAction::Flag);
3538
3539 let vision = policy.vision.as_ref().expect("vision");
3540 assert_eq!(vision.image_field, "photo");
3541 assert_eq!(vision.model, "gpt-4o");
3542 }
3543
3544 #[test]
3548 fn create_table_rejects_ai_policy_incapable_provider() {
3549 let rt = RedDBRuntime::with_options(RedDBOptions::in_memory()).unwrap();
3550 let err = rt
3551 .execute_query(
3552 "CREATE TABLE bad (id INT, body TEXT) WITH ( \
3553 EMBED (fields = ('body'), provider = 'anthropic', model = 'claude-3-5-sonnet') \
3554 )",
3555 )
3556 .unwrap_err();
3557 assert!(
3558 err.to_string()
3559 .contains("cannot serve the 'embed' modality"),
3560 "{err}"
3561 );
3562 assert!(
3563 rt.db().store().get_collection("bad").is_none(),
3564 "rejected DDL must not create the collection"
3565 );
3566 }
3567}