nodedb_lite/query/
engine.rs1use std::sync::{Arc, Mutex};
7
8use datafusion::execution::context::SessionContext;
9use datafusion::prelude::*;
10
11use nodedb_types::result::QueryResult;
12use nodedb_types::value::Value;
13
14use crate::engine::columnar::ColumnarEngine;
15use crate::engine::crdt::CrdtEngine;
16use crate::engine::htap::HtapBridge;
17use crate::engine::strict::StrictEngine;
18use crate::error::LiteError;
19use crate::storage::engine::StorageEngine;
20
21use super::arrow_convert::arrow_value_at;
22use super::columnar_provider::ColumnarTableProvider;
23use super::strict_provider::StrictTableProvider;
24use super::table_provider::LiteTableProvider;
25
26pub struct LiteQueryEngine<S: StorageEngine> {
32 pub(in crate::query) ctx: SessionContext,
33 pub(in crate::query) crdt: Arc<Mutex<CrdtEngine>>,
34 pub(in crate::query) strict: Arc<Mutex<StrictEngine<S>>>,
35 pub(in crate::query) columnar: Arc<Mutex<ColumnarEngine<S>>>,
36 pub(in crate::query) htap: Arc<Mutex<HtapBridge>>,
37 pub(in crate::query) storage: Arc<S>,
38}
39
40impl<S: StorageEngine> LiteQueryEngine<S> {
41 pub fn new(
43 crdt: Arc<Mutex<CrdtEngine>>,
44 strict: Arc<Mutex<StrictEngine<S>>>,
45 columnar: Arc<Mutex<ColumnarEngine<S>>>,
46 htap: Arc<Mutex<HtapBridge>>,
47 storage: Arc<S>,
48 ) -> Self {
49 let config = SessionConfig::new()
50 .with_information_schema(false)
51 .with_default_catalog_and_schema("nodedb", "public");
52
53 let ctx = SessionContext::new_with_config(config);
54 super::spatial_udf::register_spatial_udfs(&ctx);
55 nodedb_query::ts_udfs::register_timeseries_udfs(&ctx);
56 Self {
57 ctx,
58 crdt,
59 strict,
60 columnar,
61 htap,
62 storage,
63 }
64 }
65
66 pub fn register_collection(&self, name: &str) {
71 let provider = LiteTableProvider::new(name.to_string(), Arc::clone(&self.crdt));
72 let _ = self.ctx.register_table(name, Arc::new(provider));
74 }
75
76 pub fn register_strict_collection(&self, name: &str) {
78 let strict = match self.strict.lock() {
79 Ok(s) => s,
80 Err(p) => p.into_inner(),
81 };
82 if let Some(schema) = strict.schema(name) {
83 let provider =
84 StrictTableProvider::new(name.to_string(), schema, Arc::clone(&self.storage));
85 let _ = self.ctx.register_table(name, Arc::new(provider));
86 }
87 }
88
89 pub fn register_all_collections(&self) {
91 let crdt = match self.crdt.lock() {
93 Ok(c) => c,
94 Err(p) => p.into_inner(),
95 };
96 let crdt_collections = crdt.collection_names();
97 drop(crdt);
98
99 for name in &crdt_collections {
100 if name.starts_with("__") {
101 continue;
102 }
103 self.register_collection(name);
104 }
105
106 let strict = match self.strict.lock() {
108 Ok(s) => s,
109 Err(p) => p.into_inner(),
110 };
111 let strict_names: Vec<String> = strict
112 .collection_names()
113 .iter()
114 .map(|s| s.to_string())
115 .collect();
116 drop(strict);
117
118 for name in &strict_names {
119 self.register_strict_collection(name);
120 }
121
122 let columnar = match self.columnar.lock() {
124 Ok(c) => c,
125 Err(p) => p.into_inner(),
126 };
127 let columnar_names: Vec<String> = columnar
128 .collection_names()
129 .iter()
130 .map(|s| s.to_string())
131 .collect();
132 drop(columnar);
133
134 for name in &columnar_names {
135 self.register_columnar_collection(name);
136 }
137 }
138
139 fn apply_htap_routing(&self, sql: &str) {
146 use crate::engine::htap::routing::is_analytical_query;
147
148 if !is_analytical_query(sql) {
149 return;
150 }
151
152 let htap = match self.htap.lock() {
153 Ok(h) => h,
154 Err(p) => p.into_inner(),
155 };
156
157 if htap.is_empty() {
158 return;
159 }
160
161 for target_name in htap.all_targets() {
164 if let Some(view) = htap.view_by_target(target_name) {
165 let source = view.source.clone();
166 let target = view.target.clone();
167 drop(htap); self.register_columnar_collection_as(&target, &source);
173 return; }
175 }
176 }
177
178 fn register_columnar_collection_as(&self, collection: &str, table_name: &str) {
183 let columnar = match self.columnar.lock() {
184 Ok(c) => c,
185 Err(p) => p.into_inner(),
186 };
187 let Some(schema) = columnar.schema(collection) else {
188 return;
189 };
190 let schema = schema.clone();
191 drop(columnar);
192
193 let provider = ColumnarTableProvider::new(
194 collection.to_string(),
195 &schema,
196 Arc::clone(&self.storage),
197 Vec::new(),
198 Vec::new(),
199 );
200 let _ = self.ctx.register_table(table_name, Arc::new(provider));
201 }
202
203 pub fn register_columnar_collection(&self, name: &str) {
205 let columnar = match self.columnar.lock() {
206 Ok(c) => c,
207 Err(p) => p.into_inner(),
208 };
209 let Some(schema) = columnar.schema(name) else {
210 return;
211 };
212 let schema = schema.clone();
213
214 drop(columnar);
216
217 let provider = ColumnarTableProvider::new(
218 name.to_string(),
219 &schema,
220 Arc::clone(&self.storage),
221 Vec::new(),
222 Vec::new(),
223 );
224 let _ = self.ctx.register_table(name, Arc::new(provider));
225 }
226
227 pub async fn execute_sql(&self, sql: &str) -> Result<QueryResult, LiteError> {
232 if let Some(result) = self.try_handle_ddl(sql).await {
234 return result;
235 }
236
237 self.register_all_collections();
239
240 self.apply_htap_routing(sql);
243
244 let df = self
245 .ctx
246 .sql(sql)
247 .await
248 .map_err(|e| LiteError::Query(format!("SQL parse/plan: {e}")))?;
249
250 let batches = df
251 .collect()
252 .await
253 .map_err(|e| LiteError::Query(format!("SQL execute: {e}")))?;
254
255 let mut columns: Vec<String> = Vec::new();
257 let mut rows: Vec<Vec<Value>> = Vec::new();
258
259 for batch in &batches {
260 if columns.is_empty() {
261 columns = batch
262 .schema()
263 .fields()
264 .iter()
265 .map(|f| f.name().clone())
266 .collect();
267 }
268
269 let num_rows = batch.num_rows();
270 for row_idx in 0..num_rows {
271 let mut row = Vec::with_capacity(columns.len());
272 for col_idx in 0..batch.num_columns() {
273 let col = batch.column(col_idx);
274 let value = arrow_value_at(col, row_idx)?;
275 row.push(value);
276 }
277 rows.push(row);
278 }
279 }
280
281 Ok(QueryResult {
282 columns,
283 rows,
284 rows_affected: 0,
285 })
286 }
287}