Skip to main content

laminar_sql/datafusion/
mod.rs

1//! DataFusion integration for SQL processing.
2
3/// Marker UDFs for the `ai_*` SQL functions (rewritten by the AI operator).
4pub mod ai_udf;
5mod bridge;
6mod channel_source;
7/// Lambda higher-order functions for arrays and maps (F-SCHEMA-015 Tier 3)
8pub mod complex_type_lambda;
9/// Array, Struct, and Map scalar UDFs (F-SCHEMA-015)
10pub mod complex_type_udf;
11mod exec;
12/// End-to-end streaming SQL execution
13pub mod execute;
14/// Format bridge UDFs for inline format conversion
15pub mod format_bridge_udf;
16/// LaminarDB streaming JSON extension UDFs (F-SCHEMA-013)
17pub mod json_extensions;
18/// SQL/JSON path query compiler and scalar UDFs
19pub mod json_path;
20/// JSON table-valued functions (array/object expansion)
21pub mod json_tvf;
22/// JSONB binary format types for JSON UDF evaluation
23pub mod json_types;
24/// PostgreSQL-compatible JSON aggregate UDAFs
25pub mod json_udaf;
26/// PostgreSQL-compatible JSON scalar UDFs
27pub mod json_udf;
28/// Live source provider for streaming execution with plan caching
29pub mod live_source;
30/// Lookup join plan node for DataFusion.
31pub mod lookup_join;
32/// Physical execution plan and extension planner for lookup joins.
33pub mod lookup_join_exec;
34/// Processing-time UDF for `PROCTIME()` support
35pub mod proctime_udf;
36mod source;
37mod table_provider;
38/// Dynamic watermark filter for scan-level late-data pruning
39/// Watermark UDF for current watermark access
40pub mod watermark_udf;
41/// Window function UDFs (TUMBLE, HOP, SESSION, CUMULATE)
42pub mod window_udf;
43
44pub use ai_udf::{ai_function_markers, AiFunctionMarker};
45pub use bridge::{BridgeSendError, BridgeSender, BridgeStream, BridgeTrySendError, StreamBridge};
46pub use channel_source::ChannelStreamSource;
47pub use complex_type_lambda::{
48    register_lambda_functions, ArrayFilter, ArrayReduce, ArrayTransform, MapFilter,
49    MapTransformValues,
50};
51pub use complex_type_udf::{
52    register_complex_type_functions, MapContainsKey, MapFromArrays, MapKeys, MapValues, StructDrop,
53    StructExtract, StructMerge, StructRename, StructSet,
54};
55pub use exec::StreamingScanExec;
56pub use execute::{execute_streaming_sql, DdlResult, QueryResult, StreamingSqlResult};
57pub use format_bridge_udf::{FromJsonUdf, ParseEpochUdf, ParseTimestampUdf, ToJsonUdf};
58pub use json_extensions::{
59    register_json_extensions, JsonInferSchema, JsonToColumns, JsonbDeepMerge, JsonbExcept,
60    JsonbFlatten, JsonbMerge, JsonbPick, JsonbRenameKeys, JsonbStripNulls, JsonbUnflatten,
61};
62pub use json_path::{CompiledJsonPath, JsonPathStep, JsonbPathExistsUdf, JsonbPathMatchUdf};
63pub use json_tvf::{
64    register_json_table_functions, JsonbArrayElementsTextTvf, JsonbArrayElementsTvf,
65    JsonbEachTextTvf, JsonbEachTvf, JsonbObjectKeysTvf,
66};
67pub use json_udaf::{JsonAgg, JsonObjectAgg};
68pub use json_udf::{
69    JsonBuildArray, JsonBuildObject, JsonTypeof, JsonbContainedBy, JsonbContains, JsonbExists,
70    JsonbExistsAll, JsonbExistsAny, JsonbGet, JsonbGetIdx, JsonbGetPath, JsonbGetPathText,
71    JsonbGetText, JsonbGetTextIdx, ToJsonb,
72};
73pub use live_source::{LiveSourceHandle, LiveSourceProvider};
74pub use lookup_join_exec::{
75    LookupJoinExec, LookupJoinExtensionPlanner, LookupSnapshot, LookupTableRegistry,
76    PartialLookupJoinExec, PartialLookupState, RegisteredLookup,
77};
78pub use proctime_udf::ProcTimeUdf;
79pub use source::{SortColumn, StreamSource, StreamSourceRef};
80pub use table_provider::StreamingTableProvider;
81pub use watermark_udf::WatermarkUdf;
82pub use window_udf::{
83    CumulateWindowEnd, CumulateWindowStart, HopWindowEnd, HopWindowStart, SessionWindowStart,
84    TumbleWindowEnd, TumbleWindowStart,
85};
86
87use std::sync::atomic::AtomicI64;
88use std::sync::Arc;
89
90use datafusion::execution::SessionStateBuilder;
91use datafusion::prelude::*;
92use datafusion_expr::ScalarUDF;
93
94use crate::planner::streaming_optimizer::{StreamingPhysicalValidator, StreamingValidatorMode};
95
96/// Returns a base `SessionConfig` with identifier normalization disabled.
97///
98/// DataFusion's default behaviour lowercases all unquoted SQL identifiers
99/// (per the SQL standard). LaminarDB disables this so that mixed-case
100/// column names from external sources (Kafka, CDC, WebSocket) can be
101/// referenced without double-quoting.
102#[must_use]
103pub fn base_session_config() -> SessionConfig {
104    let mut config = SessionConfig::new();
105    config.options_mut().sql_parser.enable_ident_normalization = false;
106    // Single partition for streaming micro-batch execution. Multi-partition
107    // plans contain stateful operators (RepartitionExec) that cannot be
108    // reused across cycles, causing panics on cached physical plans.
109    config = config.with_target_partitions(1);
110    config
111}
112
113/// Creates a `DataFusion` session context with identifier normalization
114/// disabled.
115///
116/// Suitable for ad-hoc / non-streaming queries (filters, lookups).
117/// For streaming workloads prefer [`create_streaming_context`].
118/// This standalone factory uses upstream unbounded memory and disk-manager defaults.
119/// It does not share a `LaminarDB` budget; use `SessionContext::new_with_config_rt`
120/// with a caller-owned runtime when reservation limits are required.
121#[must_use]
122pub fn create_session_context() -> SessionContext {
123    SessionContext::new_with_config(base_session_config())
124}
125
126/// Creates a `DataFusion` session context configured for streaming queries.
127///
128/// The context is configured with:
129/// - Batch size of 8192 (balanced for streaming throughput)
130/// - Single partition (streaming sources are typically not partitioned)
131/// - Identifier normalization disabled (mixed-case columns work unquoted)
132/// - All streaming UDFs registered (TUMBLE, HOP, SESSION, WATERMARK)
133/// - `StreamingPhysicalValidator` in `Reject` mode (blocks unsafe plans)
134///
135/// This standalone factory retains upstream memory/disk defaults and does not share
136/// a `LaminarDB` reservation budget.
137///
138/// The watermark UDF is initialized with no watermark set (returns NULL).
139/// Use [`register_streaming_functions_with_watermark`] to provide a live
140/// watermark source.
141///
142/// # Example
143///
144/// ```rust,ignore
145/// let ctx = create_streaming_context();
146/// ctx.register_table("events", provider)?;
147/// let df = ctx.sql("SELECT * FROM events").await?;
148/// ```
149#[must_use]
150pub fn create_streaming_context() -> SessionContext {
151    create_streaming_context_with_validator(StreamingValidatorMode::Reject)
152}
153
154/// Creates a streaming context with a configurable validator mode.
155///
156/// Same as [`create_streaming_context`] but allows choosing how the
157/// [`StreamingPhysicalValidator`] handles plan violations.
158///
159/// Use [`StreamingValidatorMode::Off`] to get the previous behaviour
160/// (no plan-time validation).
161#[must_use]
162pub fn create_streaming_context_with_validator(mode: StreamingValidatorMode) -> SessionContext {
163    let config = base_session_config().with_batch_size(8192);
164
165    let ctx = if matches!(mode, StreamingValidatorMode::Off) {
166        SessionContext::new_with_config(config)
167    } else {
168        // Build a default state to get the standard optimizer rules, then
169        // prepend our streaming validator so it fires before DataFusion's
170        // built-in SanityCheckPlan (which produces generic error messages).
171        let default_state = SessionStateBuilder::new()
172            .with_config(config.clone())
173            .with_default_features()
174            .build();
175        let mut rules: Vec<
176            Arc<dyn datafusion::physical_optimizer::PhysicalOptimizerRule + Send + Sync>,
177        > = vec![Arc::new(StreamingPhysicalValidator::new(mode))];
178        rules.extend(default_state.physical_optimizers().iter().cloned());
179
180        let state = SessionStateBuilder::new()
181            .with_config(config)
182            .with_default_features()
183            .with_physical_optimizer_rules(rules)
184            .build();
185        SessionContext::new_with_state(state)
186    };
187
188    register_streaming_functions(&ctx);
189    ctx
190}
191
192/// Window-time, JSON, complex-type, lambda, and `proctime()` UDFs —
193/// every streaming UDF except `watermark()`. Pulled out of the public
194/// `register_streaming_functions*` entry points so they share a single
195/// list and stay in sync.
196fn register_non_watermark_udfs(ctx: &SessionContext) {
197    ctx.register_udf(ScalarUDF::new_from_impl(TumbleWindowStart::new()));
198    ctx.register_udf(ScalarUDF::new_from_impl(TumbleWindowEnd::new()));
199    ctx.register_udf(ScalarUDF::new_from_impl(HopWindowStart::new()));
200    ctx.register_udf(ScalarUDF::new_from_impl(HopWindowEnd::new()));
201    ctx.register_udf(ScalarUDF::new_from_impl(SessionWindowStart::new()));
202    ctx.register_udf(ScalarUDF::new_from_impl(CumulateWindowStart::new()));
203    ctx.register_udf(ScalarUDF::new_from_impl(CumulateWindowEnd::new()));
204    ctx.register_udf(ScalarUDF::new_from_impl(ProcTimeUdf::new()));
205    for marker in ai_function_markers() {
206        ctx.register_udf(marker);
207    }
208    register_json_functions(ctx);
209    register_json_extensions(ctx);
210    register_complex_type_functions(ctx);
211    register_lambda_functions(ctx);
212}
213
214/// Registers `LaminarDB` streaming UDFs with a session context. The
215/// `watermark()` UDF is registered in unset mode (always returns NULL);
216/// use [`register_streaming_functions_with_watermark`] to provide a
217/// live watermark source from Ring 0.
218pub fn register_streaming_functions(ctx: &SessionContext) {
219    register_non_watermark_udfs(ctx);
220    ctx.register_udf(ScalarUDF::new_from_impl(WatermarkUdf::unset()));
221}
222
223/// Registers streaming UDFs with a live watermark source — same as
224/// [`register_streaming_functions`] but `watermark()` reads
225/// `watermark_ms` (in milliseconds since epoch; values < 0 mean "no
226/// watermark", returning NULL).
227pub fn register_streaming_functions_with_watermark(
228    ctx: &SessionContext,
229    watermark_ms: Arc<AtomicI64>,
230) {
231    register_non_watermark_udfs(ctx);
232    ctx.register_udf(ScalarUDF::new_from_impl(WatermarkUdf::new(watermark_ms)));
233}
234
235/// Registers all PostgreSQL-compatible JSON UDFs and UDAFs
236/// with the given `SessionContext`.
237pub fn register_json_functions(ctx: &SessionContext) {
238    // Extraction operators
239    ctx.register_udf(ScalarUDF::new_from_impl(JsonbGet::new()));
240    ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetIdx::new()));
241    ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetText::new()));
242    ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetTextIdx::new()));
243    ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetPath::new()));
244    ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetPathText::new()));
245
246    // Existence operators
247    ctx.register_udf(ScalarUDF::new_from_impl(JsonbExists::new()));
248    ctx.register_udf(ScalarUDF::new_from_impl(JsonbExistsAny::new()));
249    ctx.register_udf(ScalarUDF::new_from_impl(JsonbExistsAll::new()));
250
251    // Containment operators
252    ctx.register_udf(ScalarUDF::new_from_impl(JsonbContains::new()));
253    ctx.register_udf(ScalarUDF::new_from_impl(JsonbContainedBy::new()));
254
255    // Interrogation / construction
256    ctx.register_udf(ScalarUDF::new_from_impl(JsonTypeof::new()));
257    ctx.register_udf(ScalarUDF::new_from_impl(JsonBuildObject::new()));
258    ctx.register_udf(ScalarUDF::new_from_impl(JsonBuildArray::new()));
259    ctx.register_udf(ScalarUDF::new_from_impl(ToJsonb::new()));
260
261    // Aggregates
262    ctx.register_udaf(datafusion_expr::AggregateUDF::new_from_impl(JsonAgg::new()));
263    ctx.register_udaf(datafusion_expr::AggregateUDF::new_from_impl(
264        JsonObjectAgg::new(),
265    ));
266
267    // Format bridge functions
268    ctx.register_udf(ScalarUDF::new_from_impl(ParseEpochUdf::new()));
269    ctx.register_udf(ScalarUDF::new_from_impl(ParseTimestampUdf::new()));
270    ctx.register_udf(ScalarUDF::new_from_impl(ToJsonUdf::new()));
271    ctx.register_udf(ScalarUDF::new_from_impl(FromJsonUdf::new()));
272
273    // JSON path query functions (scalar)
274    ctx.register_udf(ScalarUDF::new_from_impl(JsonbPathExistsUdf::new()));
275    ctx.register_udf(ScalarUDF::new_from_impl(JsonbPathMatchUdf::new()));
276
277    // JSON table-valued functions
278    register_json_table_functions(ctx);
279}
280
281#[cfg(test)]
282mod tests;