Skip to main content

nodedb_lite/query/
engine.rs

1//! Lite query engine: full SQL via DataFusion over Loro documents.
2//!
3//! Manages a DataFusion `SessionContext` with collections registered
4//! as table providers backed by the CRDT engine.
5
6use 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
26/// Lite-side query engine wrapping DataFusion.
27///
28/// Registered collections appear as tables. SQL queries execute
29/// entirely in-process against the Loro CRDT state, strict Binary Tuple store,
30/// or columnar compressed segments.
31pub 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    /// Create a new query engine.
42    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    /// Register a collection as a queryable table.
67    ///
68    /// Call this before executing SQL that references the collection.
69    /// For auto-registration, call `register_all_collections()`.
70    pub fn register_collection(&self, name: &str) {
71        let provider = LiteTableProvider::new(name.to_string(), Arc::clone(&self.crdt));
72        // Register directly via the session context.
73        let _ = self.ctx.register_table(name, Arc::new(provider));
74    }
75
76    /// Register a strict collection as a queryable table.
77    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    /// Register all existing collections as tables (both CRDT and strict).
90    pub fn register_all_collections(&self) {
91        // Register CRDT (schemaless) collections.
92        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        // Register strict document collections.
107        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        // Register columnar collections.
123        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    /// Apply HTAP routing: for analytical queries, re-register source strict
140    /// tables to point at their materialized columnar views.
141    ///
142    /// Only applies when the HtapBridge has registered views. The source table
143    /// name is re-registered to point at the columnar view, so DataFusion
144    /// transparently reads from the faster columnar format.
145    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 each materialized view, re-register the source table name to
162        // point at the columnar view's table provider.
163        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); // Release lock before registering.
168
169                // Re-register the SOURCE name to point at the COLUMNAR target.
170                // This makes `SELECT ... FROM customers GROUP BY ...` read from
171                // the columnar materialized view instead of the strict B-tree.
172                self.register_columnar_collection_as(&target, &source);
173                return; // Only one routing redirect per query.
174            }
175        }
176    }
177
178    /// Register a columnar collection under a different table name.
179    ///
180    /// Used by HTAP routing to make a source table name point at its
181    /// materialized columnar view.
182    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    /// Register a columnar collection as a queryable table.
204    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        // Collect segment IDs and delete bitmaps.
215        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    /// Execute a SQL query and return results.
228    ///
229    /// DDL statements (CREATE/DROP COLLECTION) are intercepted and handled
230    /// directly. All other statements are passed to DataFusion.
231    pub async fn execute_sql(&self, sql: &str) -> Result<QueryResult, LiteError> {
232        // Intercept DDL before DataFusion.
233        if let Some(result) = self.try_handle_ddl(sql).await {
234            return result;
235        }
236
237        // Auto-register collections mentioned in the query.
238        self.register_all_collections();
239
240        // HTAP routing: for analytical queries, re-register source tables to point
241        // at their materialized columnar views (if any exist and session allows it).
242        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        // Convert Arrow RecordBatches to QueryResult.
256        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}