1use super::{
13 executor::{QueryExecutor, QueryResult},
14 parser::QueryParser,
15 planner::QueryPlanner,
16 prepared::PreparedQuery,
17 result::{QueryResultIterator, StreamingConfig},
18 QueryStats,
19};
20
21use super::engine_stats::AtomicQueryStats;
22#[cfg(feature = "state_machine")]
23use super::{
24 select_executor::SelectExecutor,
25 select_optimizer::{OptimizedQueryPlan, SelectOptimizer},
26 select_parser,
27};
28use crate::{
29 memory::MemoryManager, schema::SchemaManager, storage::StorageEngine, Config, Error, Result,
30 Value,
31};
32use dashmap::DashMap;
33use std::sync::atomic::{AtomicU64, Ordering};
34use std::sync::Arc;
35use std::time::Instant;
36
37#[derive(Debug)]
43pub(crate) struct QueryCacheEntry {
44 pub parsed_query: super::ParsedQuery,
46 pub plan: super::planner::QueryPlan,
48 pub cached_at: Instant,
50 pub hit_count: AtomicU64,
54}
55
56impl Clone for QueryCacheEntry {
60 fn clone(&self) -> Self {
61 Self {
62 parsed_query: self.parsed_query.clone(),
63 plan: self.plan.clone(),
64 cached_at: self.cached_at,
65 hit_count: AtomicU64::new(self.hit_count.load(Ordering::Relaxed)),
66 }
67 }
68}
69
70#[derive(Debug, Clone)]
72pub enum SchemaStatus {
73 Available { keyspace: String, table: String },
75 Missing { table: String, reason: String },
77 ExtractionFailed {
79 table: String,
80 cause: String,
81 suggestion: String,
82 },
83}
84
85#[derive(Debug)]
87pub struct QueryEngine {
88 parser: QueryParser,
90 planner: QueryPlanner,
92 executor: QueryExecutor,
94 schema_manager: Arc<SchemaManager>,
96 #[cfg(feature = "state_machine")]
99 select_optimizer: Arc<SelectOptimizer>,
100 #[cfg(feature = "state_machine")]
102 select_executor: Arc<SelectExecutor>,
103 prepared_cache: DashMap<String, Arc<PreparedQuery>>,
105 plan_cache: DashMap<String, QueryCacheEntry>,
107 #[cfg(feature = "state_machine")]
114 select_plan_cache: DashMap<String, (Instant, Arc<OptimizedQueryPlan>)>,
115 stats: AtomicQueryStats,
118 config: Config,
120}
121
122impl QueryEngine {
123 pub fn new(
125 storage: Arc<StorageEngine>,
126 schema: Arc<SchemaManager>,
127 _memory: Arc<MemoryManager>,
128 config: &Config,
129 ) -> Result<Self> {
130 let parser = QueryParser::new(config);
131 let planner = QueryPlanner::new(schema.clone(), config);
132 let executor = QueryExecutor::new(storage.clone(), schema.clone(), config);
133
134 #[cfg(feature = "state_machine")]
136 let select_optimizer = Arc::new(SelectOptimizer::new(schema.clone(), storage.clone()));
137 #[cfg(feature = "state_machine")]
140 let select_executor = Arc::new(
141 SelectExecutor::new(schema.clone(), storage)
142 .with_max_result_bytes(
143 usize::try_from(config.query.max_result_bytes).unwrap_or(usize::MAX),
144 )
145 .with_max_result_rows(
146 usize::try_from(config.query.max_result_rows).unwrap_or(usize::MAX),
147 ),
148 );
149
150 Ok(Self {
151 parser,
152 planner,
153 executor,
154 schema_manager: schema,
155 #[cfg(feature = "state_machine")]
156 select_optimizer,
157 #[cfg(feature = "state_machine")]
158 select_executor,
159 prepared_cache: DashMap::new(),
160 plan_cache: DashMap::new(),
161 #[cfg(feature = "state_machine")]
162 select_plan_cache: DashMap::new(),
163 stats: AtomicQueryStats::default(),
164 config: config.clone(),
165 })
166 }
167
168 fn enforce_legacy_result_budget(&self, result: &QueryResult) -> Result<()> {
176 super::result_budget::enforce_materialized_rows(
177 &result.rows,
178 usize::try_from(self.config.query.max_result_bytes).unwrap_or(usize::MAX),
179 usize::try_from(self.config.query.max_result_rows).unwrap_or(usize::MAX),
180 )
181 .inspect_err(|e| {
182 self.inc_error_queries();
183 crate::observability::record_error(e, "query");
184 })
185 }
186
187 fn inc_total_queries(&self) {
189 self.stats.record_query();
190 }
191
192 fn inc_error_queries(&self) {
194 self.stats.record_error();
195 }
196
197 fn record_cache_hit(&self) {
200 self.stats.record_cache_hit();
201 }
202
203 #[tracing::instrument(
212 name = "query.execute",
213 skip(self, cql),
214 fields(
215 cqlite.query.plan_type = tracing::field::Empty,
216 cqlite.query.access_path = tracing::field::Empty,
217 cqlite.query.rows = tracing::field::Empty,
218 )
219 )]
220 pub async fn execute(&self, cql: &str) -> Result<QueryResult> {
221 let start_time = Instant::now();
222 self.inc_total_queries();
223
224 let trimmed_cql = cql.trim().to_uppercase();
234 if trimmed_cql.starts_with("SELECT") {
235 return self.execute_select_query(cql, start_time).await;
236 }
237
238 let cached_plan = self.plan_cache.get(cql).map(|entry| {
244 entry.hit_count.fetch_add(1, Ordering::Relaxed);
245 entry.plan.clone()
246 });
247 if let Some(plan) = cached_plan {
248 self.record_cache_hit();
249
250 let mut result =
251 crate::observability::record_result("query", self.executor.execute(&plan).await)?;
252 self.enforce_legacy_result_budget(&result)?;
257 self.update_execution_stats(&mut result, start_time);
258 return Ok(result);
259 }
260
261 let parsed_query = self.parser.parse(cql).inspect_err(|e| {
262 self.inc_error_queries();
263 crate::observability::record_error(e, "query");
264 })?;
265 let plan =
266 crate::observability::record_result("query", self.planner.plan(&parsed_query).await)?;
267
268 if self.config.query.query_cache_size.unwrap_or(0) > 0 {
269 self.cache_query_plan(cql, parsed_query, plan.clone());
270 }
271
272 let mut result =
273 crate::observability::record_result("query", self.executor.execute(&plan).await)?;
274 self.enforce_legacy_result_budget(&result)?;
278 self.update_execution_stats(&mut result, start_time);
279 Ok(result)
280 }
281
282 #[cfg(feature = "state_machine")]
304 pub async fn execute_streaming(
305 &self,
306 cql: &str,
307 config: StreamingConfig,
308 ) -> Result<QueryResultIterator> {
309 self.inc_total_queries();
310
311 if !cql.trim().to_uppercase().starts_with("SELECT") {
312 return Err(Error::query_execution(
313 "Streaming execution only supports SELECT queries",
314 ));
315 }
316
317 let select_statement =
318 select_parser::parse_select(cql).inspect_err(|_| self.inc_error_queries())?;
319 let optimized_plan = self.select_optimizer.optimize(select_statement).await?;
320
321 self.select_executor
322 .execute_streaming(optimized_plan, config)
323 .await
324 }
325
326 async fn execute_select_query(&self, cql: &str, start_time: Instant) -> Result<QueryResult> {
328 enum HitOutcome {
334 Reuse(super::planner::QueryPlan),
335 Placeholder,
336 Miss,
337 }
338 let outcome = match self.plan_cache.get(cql) {
339 Some(entry) if entry.plan.table.is_some() => {
340 entry.hit_count.fetch_add(1, Ordering::Relaxed);
341 HitOutcome::Reuse(entry.plan.clone())
342 }
343 Some(_) => HitOutcome::Placeholder,
344 None => HitOutcome::Miss,
345 };
346 match outcome {
347 HitOutcome::Reuse(plan) => {
348 self.record_cache_hit();
349
350 let mut result = crate::observability::record_result(
351 "query",
352 self.executor.execute(&plan).await,
353 )?;
354 self.enforce_legacy_result_budget(&result)?;
358 self.update_execution_stats(&mut result, start_time);
359 return Ok(result);
360 }
361 HitOutcome::Placeholder => {
362 self.plan_cache.remove(cql);
364 }
365 HitOutcome::Miss => {}
366 }
367
368 #[cfg(not(feature = "state_machine"))]
369 return Err(Error::query_execution(
370 "Advanced SELECT parsing requires state_machine feature",
371 ));
372
373 #[cfg(feature = "state_machine")]
374 {
375 let optimized_plan: Arc<OptimizedQueryPlan> =
380 if let Some(entry) = self.select_plan_cache.get(cql) {
381 self.record_cache_hit();
382 Arc::clone(&entry.value().1)
383 } else {
384 let select_statement = select_parser::parse_select(cql).inspect_err(|e| {
385 self.inc_error_queries();
386 crate::observability::record_error(e, "query");
387 })?;
388 let plan = Arc::new(crate::observability::record_result(
389 "query",
390 self.select_optimizer.optimize(select_statement).await,
391 )?);
392 self.cache_select_plan(cql, Arc::clone(&plan));
393 plan
394 };
395 let mut result = crate::observability::record_result(
396 "query",
397 self.select_executor
398 .execute((*optimized_plan).clone())
399 .await,
400 )?;
401 self.update_execution_stats(&mut result, start_time);
402 Ok(result)
403 }
404 }
405
406 #[cfg(feature = "state_machine")]
410 fn cache_select_plan(&self, cql: &str, plan: Arc<OptimizedQueryPlan>) {
411 let cache_size = self.config.query.query_cache_size.unwrap_or(0);
412 if cache_size == 0 {
413 return;
414 }
415 if self.select_plan_cache.len() >= cache_size {
416 let oldest_key = self
417 .select_plan_cache
418 .iter()
419 .min_by_key(|entry| entry.value().0)
420 .map(|entry| entry.key().clone());
421 if let Some(key) = oldest_key {
422 self.select_plan_cache.remove(&key);
423 }
424 }
425 self.select_plan_cache
426 .insert(cql.to_string(), (Instant::now(), plan));
427 }
428
429 pub async fn execute_with_params(&self, cql: &str, params: &[Value]) -> Result<QueryResult> {
457 let is_select = cql.trim().to_uppercase().starts_with("SELECT");
458
459 if !is_select {
460 if params.is_empty() {
464 return self.execute(cql).await;
465 }
466 self.inc_total_queries();
467 self.inc_error_queries();
468 return Err(Error::query_execution(
469 "Parameterized execution currently supports SELECT statements only",
470 ));
471 }
472
473 #[cfg(not(feature = "state_machine"))]
474 {
475 let _ = params;
476 self.inc_total_queries();
477 self.inc_error_queries();
478 return Err(Error::query_execution(
479 "Parameterized SELECT execution requires the state_machine feature",
480 ));
481 }
482
483 #[cfg(feature = "state_machine")]
484 {
485 let statement = select_parser::parse_select(cql).inspect_err(|_| {
488 self.inc_total_queries();
489 self.inc_error_queries();
490 })?;
491 let marker_count = statement.bind_marker_count();
492
493 if marker_count == 0 {
498 if params.is_empty() {
499 return self.execute(cql).await;
500 }
501 self.inc_total_queries();
502 self.inc_error_queries();
503 return Err(Error::query_execution(format!(
504 "Parameter count mismatch: query has 0 bind marker(s), got {} parameter(s)",
505 params.len()
506 )));
507 }
508
509 let start_time = Instant::now();
510 self.inc_total_queries();
511
512 let mut statement = statement;
517 statement
518 .bind_parameters(params)
519 .inspect_err(|_| self.inc_error_queries())?;
520
521 let optimized_plan = self.select_optimizer.optimize(statement).await?;
522 let mut result = self.select_executor.execute(optimized_plan).await?;
523 self.update_execution_stats(&mut result, start_time);
524 Ok(result)
525 }
526 }
527
528 pub async fn prepare(&self, cql: &str) -> Result<Arc<PreparedQuery>> {
530 if let Some(cached) = self.prepared_cache.get(cql) {
531 return Ok(cached.clone());
532 }
533
534 let parsed_query = self.parser.parse(cql)?;
535 let plan = self.planner.plan(&parsed_query).await?;
536
537 #[cfg(feature = "state_machine")]
543 let prepared = if cql.trim().to_uppercase().starts_with("SELECT") {
544 let statement = select_parser::parse_select(cql)?;
545 let marker_count = statement.bind_marker_count();
546 Arc::new(PreparedQuery::new_select(
547 parsed_query,
548 plan,
549 Arc::new(self.executor.clone()),
550 statement,
551 marker_count,
552 self.select_optimizer.clone(),
553 self.select_executor.clone(),
554 ))
555 } else {
556 Arc::new(PreparedQuery::new(
557 parsed_query,
558 plan,
559 Arc::new(self.executor.clone()),
560 ))
561 };
562
563 #[cfg(not(feature = "state_machine"))]
564 let prepared = Arc::new(PreparedQuery::new(
565 parsed_query,
566 plan,
567 Arc::new(self.executor.clone()),
568 ));
569
570 self.prepared_cache
571 .insert(cql.to_string(), prepared.clone());
572
573 Ok(prepared)
574 }
575
576 pub async fn execute_prepared(
578 &self,
579 prepared: &PreparedQuery,
580 params: &[Value],
581 ) -> Result<QueryResult> {
582 let start_time = Instant::now();
583 self.inc_total_queries();
584
585 let mut result = prepared.execute(params).await?;
586 self.update_execution_stats(&mut result, start_time);
587 Ok(result)
588 }
589
590 pub fn stats(&self) -> QueryStats {
592 self.stats.snapshot()
593 }
594
595 pub fn clear_caches(&self) {
597 self.prepared_cache.clear();
598 self.plan_cache.clear();
599 #[cfg(feature = "state_machine")]
600 self.select_plan_cache.clear();
601 }
602
603 pub fn clear_prepared_cache(&self) {
605 self.prepared_cache.clear();
606 }
607
608 pub fn clear_plan_cache(&self) {
610 self.plan_cache.clear();
611 #[cfg(feature = "state_machine")]
612 self.select_plan_cache.clear();
613 }
614
615 pub fn cache_stats(&self) -> CacheStats {
617 CacheStats {
618 prepared_cache_size: self.prepared_cache.len(),
619 plan_cache_size: self.plan_cache.len(),
620 prepared_cache_hits: self.prepared_cache.len() as u64,
621 plan_cache_hits: self.plan_cache.len() as u64,
622 }
623 }
624
625 pub async fn explain(&self, cql: &str) -> Result<ExplainResult> {
627 let parsed_query = self.parser.parse(cql)?;
629
630 let plan = self.planner.plan(&parsed_query).await?;
632
633 Ok(ExplainResult {
634 query_type: format!("{:?}", parsed_query.query_type),
635 plan_type: format!("{:?}", plan.plan_type),
636 estimated_cost: plan.estimated_cost,
637 estimated_rows: plan.estimated_rows,
638 selected_indexes: plan
639 .selected_indexes
640 .iter()
641 .map(|idx| format!("{} ({:?})", idx.index_name, idx.index_type))
642 .collect(),
643 execution_steps: plan
644 .steps
645 .iter()
646 .map(|step| {
647 format!(
648 "{:?}: {} (cost: {:.2})",
649 step.step_type,
650 step.columns.join(", "),
651 step.cost
652 )
653 })
654 .collect(),
655 parallelization_info: plan
656 .steps
657 .iter()
658 .filter(|step| step.parallelization.can_parallelize)
659 .map(|step| {
660 format!(
661 "Threads: {}, Partition: {:?}",
662 step.parallelization.suggested_threads, step.parallelization.partition_key
663 )
664 })
665 .collect(),
666 })
667 }
668
669 pub async fn analyze(&self, cql: &str) -> Result<AnalyzeResult> {
671 let start_time = Instant::now();
672
673 let mut execution_times = Vec::new();
675 let mut results = Vec::new();
676
677 for _ in 0..self.config.query.analyze_iterations.unwrap_or(5) {
678 let iter_start = Instant::now();
679 let result = self.execute(cql).await?;
680 execution_times.push(iter_start.elapsed());
681 results.push(result);
682 }
683
684 let total_time = start_time.elapsed();
685 let avg_time =
686 execution_times.iter().sum::<std::time::Duration>() / execution_times.len() as u32;
687 let no_times = || Error::query_execution("No execution times recorded for analysis");
688 let min_time = execution_times.iter().min().ok_or_else(no_times)?;
689 let max_time = execution_times.iter().max().ok_or_else(no_times)?;
690
691 let variance = execution_times
693 .iter()
694 .map(|time| {
695 let diff = time.as_nanos() as f64 - avg_time.as_nanos() as f64;
696 diff * diff
697 })
698 .sum::<f64>()
699 / execution_times.len() as f64;
700 let std_dev = variance.sqrt();
701
702 Ok(AnalyzeResult {
703 iterations: execution_times.len(),
704 total_time_ms: total_time.as_millis() as u64,
705 avg_time_ms: avg_time.as_millis() as u64,
706 min_time_ms: min_time.as_millis() as u64,
707 max_time_ms: max_time.as_millis() as u64,
708 std_dev_ms: (std_dev / 1_000_000.0) as u64, avg_rows_returned: results.iter().map(|r| r.rows.len()).sum::<usize>() / results.len(),
710 cache_hit_ratio: self.stats().cache_hit_ratio,
711 })
712 }
713
714 fn cache_query_plan(
716 &self,
717 cql: &str,
718 parsed_query: super::ParsedQuery,
719 plan: super::planner::QueryPlan,
720 ) {
721 let cache_size = self.config.query.query_cache_size.unwrap_or(0);
722 if cache_size == 0 {
723 return;
724 }
725
726 if self.plan_cache.len() >= cache_size {
727 let oldest_key = self
728 .plan_cache
729 .iter()
730 .min_by_key(|entry| entry.cached_at)
731 .map(|entry| entry.key().clone());
732 if let Some(key) = oldest_key {
733 self.plan_cache.remove(&key);
734 }
735 }
736
737 self.plan_cache.insert(
738 cql.to_string(),
739 QueryCacheEntry {
740 parsed_query,
741 plan,
742 cached_at: Instant::now(),
743 hit_count: AtomicU64::new(0),
744 },
745 );
746 }
747
748 pub async fn has_schema_for_table(&self, table: &str) -> bool {
750 self.schema_manager.get_table_schema(table).await.is_ok()
751 }
752
753 pub async fn schema_status(&self, table: &str) -> SchemaStatus {
755 match self.schema_manager.get_table_schema(table).await {
756 Ok(schema) => SchemaStatus::Available {
757 keyspace: schema.keyspace.clone(),
758 table: schema.table.clone(),
759 },
760 Err(Error::Schema(msg)) if msg.contains("not found") => {
761 SchemaStatus::Missing {
762 table: table.to_string(),
763 reason: msg,
764 }
765 }
766 Err(e) => SchemaStatus::ExtractionFailed {
767 table: table.to_string(),
768 cause: e.to_string(),
769 suggestion: "Verify SSTable files are valid Cassandra 5.0 format and Statistics.db contains SerializationHeader".to_string(),
770 },
771 }
772 }
773
774 fn update_execution_stats(&self, result: &mut QueryResult, start_time: Instant) {
786 use crate::observability::{self as obs, catalog, AttrValue};
787
788 let execution_time = start_time.elapsed();
789 result.execution_time_ms = if execution_time.is_zero() {
791 0
792 } else {
793 std::cmp::max(1, execution_time.as_millis() as u64)
794 };
795
796 let access_path_label: &'static str = result
802 .metadata
803 .access_path
804 .as_ref()
805 .map(|p| p.label())
806 .unwrap_or("unknown");
807 let plan_type_label: &'static str = result
808 .metadata
809 .plan_info
810 .as_ref()
811 .map(|p| Self::plan_type_label(&p.plan_type))
812 .unwrap_or("unknown");
813
814 obs::record_histogram(
816 catalog::QUERY_DURATION,
817 execution_time.as_secs_f64(),
818 &[
819 (catalog::attr::SUBSYSTEM, AttrValue::StaticStr("query")),
820 (
821 catalog::attr::ACCESS_PATH,
822 AttrValue::StaticStr(access_path_label),
823 ),
824 (
825 catalog::attr::PLAN_TYPE,
826 AttrValue::StaticStr(plan_type_label),
827 ),
828 ],
829 );
830 obs::add_counter(
831 catalog::QUERY_ROWS,
832 result.rows.len() as u64,
833 &[
834 (
835 catalog::attr::ACCESS_PATH,
836 AttrValue::StaticStr(access_path_label),
837 ),
838 (
839 catalog::attr::PLAN_TYPE,
840 AttrValue::StaticStr(plan_type_label),
841 ),
842 ],
843 );
844
845 let span = tracing::Span::current();
848 span.record(catalog::attr::PLAN_TYPE, plan_type_label);
849 span.record(catalog::attr::ACCESS_PATH, access_path_label);
850 span.record("cqlite.query.rows", result.rows.len());
851
852 let new_time_us = execution_time.as_micros() as u64;
856 self.stats
857 .record_execution(new_time_us, result.rows_affected);
858 }
859
860 fn plan_type_label(plan_type: &str) -> &'static str {
865 match plan_type {
866 "TableScan" => "table_scan",
867 "IndexScan" => "index_scan",
868 "PointLookup" => "point_lookup",
869 "RangeScan" => "range_scan",
870 "Join" => "join",
871 "Aggregation" => "aggregation",
872 "Subquery" => "subquery",
873 _ => "other",
874 }
875 }
876}
877
878#[derive(Debug, Clone)]
880pub struct CacheStats {
881 pub prepared_cache_size: usize,
883 pub plan_cache_size: usize,
885 pub prepared_cache_hits: u64,
887 pub plan_cache_hits: u64,
889}
890
891#[derive(Debug, Clone)]
893pub struct ExplainResult {
894 pub query_type: String,
896 pub plan_type: String,
898 pub estimated_cost: f64,
900 pub estimated_rows: u64,
902 pub selected_indexes: Vec<String>,
904 pub execution_steps: Vec<String>,
906 pub parallelization_info: Vec<String>,
908}
909
910#[derive(Debug, Clone)]
912pub struct AnalyzeResult {
913 pub iterations: usize,
915 pub total_time_ms: u64,
917 pub avg_time_ms: u64,
919 pub min_time_ms: u64,
921 pub max_time_ms: u64,
923 pub std_dev_ms: u64,
925 pub avg_rows_returned: usize,
927 pub cache_hit_ratio: f64,
929}
930
931#[cfg(all(test, feature = "state_machine"))]
932mod tests {
933 use super::*;
934 use crate::Config;
935 use std::sync::Arc;
936 use tempfile::TempDir;
937
938 #[tokio::test]
939 async fn test_query_engine_creation() {
940 let temp_dir = TempDir::new().unwrap();
941 let config = Config::default();
942 let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
943
944 let storage = Arc::new(
945 crate::storage::StorageEngine::open(
946 temp_dir.path(),
947 &config,
948 platform,
949 #[cfg(feature = "state_machine")]
950 None,
951 )
952 .await
953 .unwrap(),
954 );
955 let schema = Arc::new(
956 crate::schema::SchemaManager::new(temp_dir.path())
957 .await
958 .unwrap(),
959 );
960 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
961
962 let query_engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
963
964 assert_eq!(query_engine.stats().total_queries, 0);
965 assert_eq!(query_engine.cache_stats().prepared_cache_size, 0);
966 assert_eq!(query_engine.cache_stats().plan_cache_size, 0);
967 }
968
969 #[tokio::test]
970 #[ignore = "Hangs >60s; needs investigation - gated for M1"]
971 async fn test_query_caching() {
972 let temp_dir = TempDir::new().unwrap();
973 let mut config = Config::test_config();
974 config.query.query_cache_size = Some(10);
975
976 let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
977 let storage = Arc::new(
978 crate::storage::StorageEngine::open(
979 temp_dir.path(),
980 &config,
981 platform,
982 #[cfg(feature = "state_machine")]
983 None,
984 )
985 .await
986 .unwrap(),
987 );
988 let schema = Arc::new(
989 crate::schema::SchemaManager::new(temp_dir.path())
990 .await
991 .unwrap(),
992 );
993 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
994
995 let query_engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
996
997 let cql = "SELECT * FROM users WHERE id = 1";
999 let _ = query_engine.execute(cql).await;
1000 let _ = query_engine.execute(cql).await;
1001
1002 assert_eq!(query_engine.cache_stats().plan_cache_size, 1);
1004
1005 let stats = query_engine.stats();
1007 assert!(stats.cache_hit_ratio > 0.0);
1008 }
1009
1010 #[tokio::test]
1011 #[cfg(feature = "state_machine")]
1012 async fn test_prepared_statements() {
1013 let temp_dir = TempDir::new().unwrap();
1014 let config = Config::default();
1015 let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
1016
1017 let storage = Arc::new(
1018 crate::storage::StorageEngine::open(
1019 temp_dir.path(),
1020 &config,
1021 platform,
1022 #[cfg(feature = "state_machine")]
1023 None,
1024 )
1025 .await
1026 .unwrap(),
1027 );
1028 let schema = Arc::new(
1029 crate::schema::SchemaManager::new(temp_dir.path())
1030 .await
1031 .unwrap(),
1032 );
1033 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
1034
1035 let query_engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
1036
1037 let cql = "SELECT * FROM users WHERE id = ?";
1039 let prepared = query_engine.prepare(cql).await.unwrap();
1040
1041 let params = vec![Value::Integer(1)];
1043 let result = query_engine
1044 .execute_prepared(&prepared, ¶ms)
1045 .await
1046 .unwrap();
1047
1048 assert!(result.execution_time_ms > 0);
1050
1051 assert_eq!(query_engine.cache_stats().prepared_cache_size, 1);
1053 }
1054
1055 #[tokio::test]
1056 async fn test_query_explain() {
1057 let temp_dir = TempDir::new().unwrap();
1058 let config = Config::default();
1059 let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
1060
1061 let storage = Arc::new(
1062 crate::storage::StorageEngine::open(
1063 temp_dir.path(),
1064 &config,
1065 platform,
1066 #[cfg(feature = "state_machine")]
1067 None,
1068 )
1069 .await
1070 .unwrap(),
1071 );
1072 let schema = Arc::new(
1073 crate::schema::SchemaManager::new(temp_dir.path())
1074 .await
1075 .unwrap(),
1076 );
1077 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
1078
1079 let query_engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
1080
1081 let cql = "SELECT * FROM users WHERE id = 1";
1083 let explain_result = query_engine.explain(cql).await.unwrap();
1084
1085 assert_eq!(explain_result.query_type, "Select");
1086 assert!(explain_result.estimated_cost > 0.0);
1087 assert!(!explain_result.selected_indexes.is_empty());
1088 assert!(!explain_result.execution_steps.is_empty());
1089 }
1090
1091 #[tokio::test]
1092 #[cfg(feature = "state_machine")]
1093 async fn test_cache_eviction() {
1094 let temp_dir = TempDir::new().unwrap();
1095 let mut config = Config::default();
1096 config.query.query_cache_size = Some(2); let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
1099 let storage = Arc::new(
1100 crate::storage::StorageEngine::open(
1101 temp_dir.path(),
1102 &config,
1103 platform,
1104 #[cfg(feature = "state_machine")]
1105 None,
1106 )
1107 .await
1108 .unwrap(),
1109 );
1110 let schema = Arc::new(
1111 crate::schema::SchemaManager::new(temp_dir.path())
1112 .await
1113 .unwrap(),
1114 );
1115 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
1116
1117 let query_engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
1118
1119 let _ = query_engine
1122 .execute("SELECT * FROM users WHERE id = 1")
1123 .await;
1124 let _ = query_engine
1125 .execute("SELECT * FROM users WHERE id = 2")
1126 .await;
1127 let _ = query_engine
1128 .execute("SELECT * FROM users WHERE id = 3")
1129 .await;
1130
1131 assert_eq!(query_engine.select_plan_cache.len(), 2);
1133 }
1134
1135 #[tokio::test]
1136 #[cfg(feature = "state_machine")]
1137 async fn test_schema_validation_api() {
1138 let temp_dir = TempDir::new().unwrap();
1139 let config = Config::default();
1140 let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
1141
1142 let storage = Arc::new(
1143 crate::storage::StorageEngine::open(
1144 temp_dir.path(),
1145 &config,
1146 platform,
1147 #[cfg(feature = "state_machine")]
1148 None,
1149 )
1150 .await
1151 .unwrap(),
1152 );
1153 let schema = Arc::new(
1154 crate::schema::SchemaManager::new(temp_dir.path())
1155 .await
1156 .unwrap(),
1157 );
1158 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
1159
1160 let query_engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
1161
1162 let has_schema = query_engine.has_schema_for_table("nonexistent_table").await;
1164 assert!(!has_schema, "Should return false for non-existent table");
1165
1166 let status = query_engine.schema_status("nonexistent_table").await;
1168 match status {
1169 SchemaStatus::Missing { .. } | SchemaStatus::ExtractionFailed { .. } => {
1170 }
1172 SchemaStatus::Available { .. } => {
1173 panic!("Should not be Available for non-existent table");
1174 }
1175 }
1176 }
1177
1178 #[tokio::test]
1189 #[cfg(feature = "state_machine")]
1190 async fn adhoc_select_reuses_cached_optimized_plan() {
1191 use crate::query::select_optimizer::OPTIMIZE_INVOCATIONS;
1192
1193 let temp_dir = TempDir::new().unwrap();
1194 let mut config = Config::default();
1195 config.query.query_cache_size = Some(16);
1197 let platform = Arc::new(crate::platform::Platform::new(&config).await.unwrap());
1198 let storage = Arc::new(
1199 crate::storage::StorageEngine::open(
1200 temp_dir.path(),
1201 &config,
1202 platform,
1203 #[cfg(feature = "state_machine")]
1204 None,
1205 )
1206 .await
1207 .unwrap(),
1208 );
1209 let schema = Arc::new(
1210 crate::schema::SchemaManager::new(temp_dir.path())
1211 .await
1212 .unwrap(),
1213 );
1214 let memory = Arc::new(crate::memory::MemoryManager::new(&config).unwrap());
1215 let engine = QueryEngine::new(storage, schema, memory, &config).unwrap();
1216
1217 let cql = "SELECT * FROM t WHERE age > 5";
1220
1221 OPTIMIZE_INVOCATIONS.with(|c| c.set(0));
1222 for _ in 0..75 {
1223 let _ = engine.execute(cql).await;
1224 }
1225 let after_same = OPTIMIZE_INVOCATIONS.with(|c| c.get());
1226 assert!(
1227 after_same <= 1,
1228 "issue #1587: an ad-hoc SELECT must optimize at most once across 75 \
1229 identical executes (modern plan cache), got {after_same}"
1230 );
1231
1232 engine.clear_plan_cache();
1234 let _ = engine.execute(cql).await;
1235 let after_clear = OPTIMIZE_INVOCATIONS.with(|c| c.get());
1236 assert_eq!(
1237 after_clear,
1238 after_same + 1,
1239 "clearing the plan cache must re-optimize once (no stale reuse)"
1240 );
1241 }
1242}
1243
1244#[cfg(test)]
1250#[path = "engine_lock_hygiene_tests.rs"]
1251mod engine_lock_hygiene_tests;
1252
1253#[cfg(test)]
1254#[cfg(all(feature = "experimental", feature = "state_machine"))]
1255mod plan_cache_tests {
1256 use super::*;
1257 use crate::{
1258 memory::MemoryManager, platform::Platform, schema::SchemaManager, storage::StorageEngine,
1259 Config,
1260 };
1261 use std::sync::Arc;
1262 use tempfile::TempDir;
1263
1264 async fn setup_query_engine(config: &Config) -> (QueryEngine, TempDir) {
1265 let temp_dir = TempDir::new().unwrap();
1266 let platform = Arc::new(Platform::new(config).await.unwrap());
1267 let storage = Arc::new(
1268 StorageEngine::open(
1269 temp_dir.path(),
1270 config,
1271 platform,
1272 #[cfg(feature = "state_machine")]
1273 None,
1274 )
1275 .await
1276 .unwrap(),
1277 );
1278 let schema = Arc::new(SchemaManager::new(temp_dir.path()).await.unwrap());
1279 let memory = Arc::new(MemoryManager::new(config).unwrap());
1280
1281 let engine = QueryEngine::new(storage, schema, memory, config).unwrap();
1282 (engine, temp_dir)
1283 }
1284
1285 fn select_plan_cache_len(engine: &QueryEngine) -> usize {
1290 engine.select_plan_cache.len()
1291 }
1292
1293 async fn create_sample_table(engine: &QueryEngine) {
1298 engine
1299 .execute(
1300 "CREATE TABLE plan_cache_test (
1301 id INTEGER PRIMARY KEY,
1302 value TEXT
1303 )",
1304 )
1305 .await
1306 .unwrap();
1307 }
1308
1309 #[tokio::test]
1310 async fn test_plan_cache_disabled() {
1311 let mut config = Config::default();
1312 config.query.query_cache_size = Some(0);
1313
1314 let (engine, _temp_dir) = setup_query_engine(&config).await;
1315 create_sample_table(&engine).await;
1316
1317 engine
1318 .execute("SELECT * FROM plan_cache_test WHERE id = 1")
1319 .await
1320 .unwrap();
1321
1322 assert_eq!(select_plan_cache_len(&engine), 0);
1323 }
1324
1325 #[tokio::test]
1326 async fn test_plan_cache_reuse_point_lookup() {
1327 let mut config = Config::default();
1328 config.query.query_cache_size = Some(4);
1329
1330 let (engine, _temp_dir) = setup_query_engine(&config).await;
1331 create_sample_table(&engine).await;
1332
1333 engine.clear_plan_cache();
1334
1335 engine
1336 .execute("SELECT * FROM plan_cache_test WHERE id = 1")
1337 .await
1338 .unwrap();
1339 engine
1340 .execute("SELECT * FROM plan_cache_test WHERE id = 1")
1341 .await
1342 .unwrap();
1343
1344 assert_eq!(select_plan_cache_len(&engine), 1);
1345 assert!(engine.stats().cache_hit_ratio > 0.0);
1346 }
1347
1348 #[tokio::test]
1349 async fn test_plan_cache_eviction_limit() {
1350 let mut config = Config::default();
1351 config.query.query_cache_size = Some(2);
1352
1353 let (engine, _temp_dir) = setup_query_engine(&config).await;
1354 create_sample_table(&engine).await;
1355
1356 engine.clear_plan_cache();
1357
1358 for id in 1..=3 {
1359 engine
1360 .execute(&format!("SELECT * FROM plan_cache_test WHERE id = {}", id))
1361 .await
1362 .unwrap();
1363 }
1364
1365 assert_eq!(select_plan_cache_len(&engine), 2);
1366 }
1367}