Skip to main content

laminar_sql/translator/
streaming_ddl.rs

1//! SQL DDL to streaming API translation.
2
3#[allow(clippy::disallowed_types)] // cold path: SQL translation
4use std::collections::HashMap;
5use std::str::FromStr;
6use std::sync::Arc;
7use std::time::Duration;
8
9use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
10use sqlparser::ast::{ColumnDef, DataType as SqlDataType};
11
12use laminar_core::streaming::config::{
13    BackpressureStrategy, WaitStrategy, DEFAULT_BUFFER_SIZE, MAX_BUFFER_SIZE, MIN_BUFFER_SIZE,
14};
15
16use crate::parser::ParseError;
17use crate::parser::{CreateSinkStatement, CreateSourceStatement, SinkFrom, WatermarkDef};
18
19/// Watermark specification for a source.
20#[derive(Debug, Clone)]
21pub struct WatermarkSpec {
22    /// Column name for event time.
23    pub column: String,
24    /// Bounded out-of-orderness duration.
25    pub max_out_of_orderness: Duration,
26    /// Whether this is a processing-time watermark (`PROCTIME()`).
27    ///
28    /// When `true`, the runtime should use `ProcessingTimeGenerator`
29    /// instead of `BoundedOutOfOrdernessGenerator`.
30    pub is_processing_time: bool,
31}
32
33/// Configuration options for a streaming source.
34#[derive(Debug, Clone)]
35pub struct SourceConfigOptions {
36    /// Buffer size for the channel.
37    pub buffer_size: usize,
38    /// Backpressure strategy.
39    pub backpressure: BackpressureStrategy,
40    /// Wait strategy for consumers.
41    pub wait_strategy: WaitStrategy,
42    /// Whether to track statistics.
43    pub track_stats: bool,
44}
45
46impl Default for SourceConfigOptions {
47    fn default() -> Self {
48        Self {
49            buffer_size: DEFAULT_BUFFER_SIZE,
50            backpressure: BackpressureStrategy::Block,
51            wait_strategy: WaitStrategy::SpinYield,
52            track_stats: false,
53        }
54    }
55}
56
57/// Column definition for a streaming source.
58#[derive(Debug, Clone)]
59pub struct ColumnDefinition {
60    /// Column name.
61    pub name: String,
62    /// Arrow data type.
63    pub data_type: DataType,
64    /// Whether the column is nullable.
65    pub nullable: bool,
66}
67
68/// A validated streaming source definition.
69///
70/// This is the output of translating a `CreateSourceStatement` to a typed
71/// configuration that can be used to create runtime sources.
72#[derive(Debug, Clone)]
73pub struct SourceDefinition {
74    /// Source name.
75    pub name: String,
76    /// Column definitions.
77    pub columns: Vec<ColumnDefinition>,
78    /// Arrow schema.
79    pub schema: SchemaRef,
80    /// Watermark specification, if defined.
81    pub watermark: Option<WatermarkSpec>,
82    /// Configuration options.
83    pub config: SourceConfigOptions,
84}
85
86impl TryFrom<CreateSourceStatement> for SourceDefinition {
87    type Error = ParseError;
88
89    fn try_from(stmt: CreateSourceStatement) -> Result<Self, Self::Error> {
90        translate_create_source(stmt)
91    }
92}
93
94/// A validated streaming sink definition.
95#[derive(Debug, Clone)]
96pub struct SinkDefinition {
97    /// Sink name.
98    pub name: String,
99    /// Input source or query.
100    pub input: String,
101    /// Configuration options.
102    pub config: SourceConfigOptions,
103}
104
105impl TryFrom<CreateSinkStatement> for SinkDefinition {
106    type Error = ParseError;
107
108    fn try_from(stmt: CreateSinkStatement) -> Result<Self, Self::Error> {
109        translate_create_sink(stmt)
110    }
111}
112
113/// Translates a CREATE SOURCE statement to a typed SourceDefinition.
114///
115/// # Errors
116///
117/// Returns `ParseError::ValidationError` if:
118/// - The `channel` option is specified (not user-configurable)
119/// - An invalid option value is provided
120/// - Column types cannot be converted to Arrow types
121pub fn translate_create_source(
122    stmt: CreateSourceStatement,
123) -> Result<SourceDefinition, ParseError> {
124    let columns = convert_columns(&stmt.columns)?;
125    translate_create_source_with_columns(stmt, columns)
126}
127
128/// Translate a CREATE SOURCE statement using an already-resolved column
129/// list. Used by the DDL layer after `discover_schema` so `WATERMARK FOR`
130/// validates against the discovered columns rather than the SQL text.
131///
132/// # Errors
133///
134/// Returns `ParseError` from option validation or watermark parsing.
135pub fn translate_create_source_with_columns(
136    stmt: CreateSourceStatement,
137    columns: Vec<ColumnDefinition>,
138) -> Result<SourceDefinition, ParseError> {
139    validate_source_options(&stmt.with_options)?;
140    let config = parse_source_options(&stmt.with_options)?;
141
142    let fields: Vec<Field> = columns
143        .iter()
144        .map(|col| Field::new(&col.name, col.data_type.clone(), col.nullable))
145        .collect();
146    let schema = Arc::new(Schema::new(fields));
147
148    let watermark = if let Some(wm) = stmt.watermark {
149        Some(parse_watermark(&wm, &columns)?)
150    } else {
151        None
152    };
153
154    Ok(SourceDefinition {
155        name: stmt.name.to_string(),
156        columns,
157        schema,
158        watermark,
159        config,
160    })
161}
162
163/// Translates a CREATE SINK statement to a typed SinkDefinition.
164///
165/// # Errors
166///
167/// Returns `ParseError::ValidationError` if:
168/// - The `channel` option is specified (not user-configurable)
169/// - An invalid option value is provided
170pub fn translate_create_sink(stmt: CreateSinkStatement) -> Result<SinkDefinition, ParseError> {
171    // Validate options first
172    validate_source_options(&stmt.with_options)?;
173
174    // Parse configuration options
175    let config = parse_source_options(&stmt.with_options)?;
176
177    // Get input name
178    let input = match stmt.from {
179        SinkFrom::Table(name) => name.to_string(),
180        SinkFrom::Query(_) => {
181            // For now, we don't support inline queries - need to create a view first
182            return Err(ParseError::ValidationError(
183                "inline queries not yet supported in CREATE SINK - use a view".to_string(),
184            ));
185        }
186    };
187
188    Ok(SinkDefinition {
189        name: stmt.name.to_string(),
190        input,
191        config,
192    })
193}
194
195/// Validates that source options don't include disallowed keys.
196fn validate_source_options(options: &HashMap<String, String>) -> Result<(), ParseError> {
197    // Reject 'channel' option - channel type is auto-derived
198    if options.contains_key("channel") {
199        return Err(ParseError::ValidationError(
200            "the 'channel' option is not user-configurable - channel type is automatically derived from usage patterns".to_string(),
201        ));
202    }
203
204    // Reject 'type' option for same reason
205    if options.contains_key("type") {
206        return Err(ParseError::ValidationError(
207            "the 'type' option is not user-configurable for in-memory streaming sources"
208                .to_string(),
209        ));
210    }
211
212    Ok(())
213}
214
215/// Parses source options from WITH clause.
216fn parse_source_options(
217    options: &HashMap<String, String>,
218) -> Result<SourceConfigOptions, ParseError> {
219    let mut config = SourceConfigOptions::default();
220
221    for (key, value) in options {
222        match key.to_lowercase().as_str() {
223            "buffer_size" | "buffersize" => {
224                config.buffer_size = parse_buffer_size(value)?;
225            }
226            "backpressure" => {
227                config.backpressure =
228                    BackpressureStrategy::from_str(value).map_err(ParseError::ValidationError)?;
229            }
230            "wait_strategy" | "waitstrategy" => {
231                config.wait_strategy =
232                    WaitStrategy::from_str(value).map_err(ParseError::ValidationError)?;
233            }
234            "track_stats" | "trackstats" | "stats" => {
235                config.track_stats = parse_bool(value)?;
236            }
237            // Ignore connector-specific and unknown options.
238            // Connector-specific: handled by connector implementations.
239            // Unknown: allow forward compatibility with new options.
240            _ => {}
241        }
242    }
243
244    Ok(config)
245}
246
247/// Parses buffer_size option.
248fn parse_buffer_size(value: &str) -> Result<usize, ParseError> {
249    let size: usize = value.parse().map_err(|_| {
250        ParseError::ValidationError(format!(
251            "invalid buffer_size: '{}' - must be a number",
252            value
253        ))
254    })?;
255
256    if size < MIN_BUFFER_SIZE {
257        return Err(ParseError::ValidationError(format!(
258            "buffer_size {} is too small - minimum is {}",
259            size, MIN_BUFFER_SIZE
260        )));
261    }
262
263    if size > MAX_BUFFER_SIZE {
264        return Err(ParseError::ValidationError(format!(
265            "buffer_size {} is too large - maximum is {}",
266            size, MAX_BUFFER_SIZE
267        )));
268    }
269
270    Ok(size)
271}
272
273/// Parses a boolean option.
274fn parse_bool(value: &str) -> Result<bool, ParseError> {
275    match value.to_lowercase().as_str() {
276        "true" | "yes" | "on" | "1" => Ok(true),
277        "false" | "no" | "off" | "0" => Ok(false),
278        _ => Err(ParseError::ValidationError(format!(
279            "invalid boolean value: '{}' - expected true/false",
280            value
281        ))),
282    }
283}
284
285/// Converts SQL column definitions to Arrow types.
286fn convert_columns(columns: &[ColumnDef]) -> Result<Vec<ColumnDefinition>, ParseError> {
287    columns.iter().map(convert_column).collect()
288}
289
290/// Converts a single SQL column definition to Arrow type.
291fn convert_column(col: &ColumnDef) -> Result<ColumnDefinition, ParseError> {
292    let data_type = sql_type_to_arrow(&col.data_type)?;
293
294    // Check for NOT NULL constraint
295    let nullable = !col
296        .options
297        .iter()
298        .any(|opt| matches!(opt.option, sqlparser::ast::ColumnOption::NotNull));
299
300    Ok(ColumnDefinition {
301        name: col.name.value.clone(),
302        data_type,
303        nullable,
304    })
305}
306
307/// Converts SQL data type to Arrow data type.
308///
309/// # Errors
310///
311/// Returns `ParseError::ValidationError` for unsupported SQL data types.
312pub fn sql_type_to_arrow(sql_type: &SqlDataType) -> Result<DataType, ParseError> {
313    match sql_type {
314        // Integer types
315        SqlDataType::TinyInt(_) => Ok(DataType::Int8),
316        SqlDataType::SmallInt(_) => Ok(DataType::Int16),
317        SqlDataType::Int(_) | SqlDataType::Integer(_) => Ok(DataType::Int32),
318        SqlDataType::BigInt(_) => Ok(DataType::Int64),
319
320        // Unsigned integer types - wrapped in Unsigned variant
321        // Note: sqlparser wraps unsigned types differently in different versions
322
323        // Floating point types
324        SqlDataType::Float(_) | SqlDataType::Real => Ok(DataType::Float32),
325        SqlDataType::Double(_) | SqlDataType::DoublePrecision => Ok(DataType::Float64),
326
327        // Decimal types
328        SqlDataType::Decimal(info) | SqlDataType::Numeric(info) => {
329            #[allow(clippy::cast_possible_truncation)] // Precision/scale are typically small values
330            let (precision, scale) = match info {
331                sqlparser::ast::ExactNumberInfo::PrecisionAndScale(p, s) => (*p as u8, *s as i8),
332                sqlparser::ast::ExactNumberInfo::Precision(p) => (*p as u8, 0),
333                sqlparser::ast::ExactNumberInfo::None => (38, 9), // Default precision/scale
334            };
335            Ok(DataType::Decimal128(precision, scale))
336        }
337
338        // String types (including JSON/UUID stored as strings)
339        SqlDataType::Char(_)
340        | SqlDataType::Character(_)
341        | SqlDataType::Varchar(_)
342        | SqlDataType::CharacterVarying(_)
343        | SqlDataType::Text
344        | SqlDataType::String(_)
345        | SqlDataType::JSON
346        | SqlDataType::JSONB
347        | SqlDataType::Uuid => Ok(DataType::Utf8),
348
349        // Binary types
350        SqlDataType::Binary(_)
351        | SqlDataType::Varbinary(_)
352        | SqlDataType::Blob(_)
353        | SqlDataType::Bytea => Ok(DataType::Binary),
354
355        // Boolean type
356        SqlDataType::Boolean | SqlDataType::Bool => Ok(DataType::Boolean),
357
358        // Date/time types
359        SqlDataType::Date => Ok(DataType::Date32),
360        SqlDataType::Time(_, _) => Ok(DataType::Time64(TimeUnit::Microsecond)),
361        SqlDataType::Timestamp(_, _) => Ok(DataType::Timestamp(TimeUnit::Microsecond, None)),
362
363        // Interval type
364        SqlDataType::Interval { .. } => Ok(DataType::Interval(
365            arrow::datatypes::IntervalUnit::MonthDayNano,
366        )),
367
368        // Array type: ARRAY<T>, T[], Array(T)
369        SqlDataType::Array(elem_def) => {
370            let item_type = match elem_def {
371                sqlparser::ast::ArrayElemTypeDef::AngleBracket(t)
372                | sqlparser::ast::ArrayElemTypeDef::SquareBracket(t, _)
373                | sqlparser::ast::ArrayElemTypeDef::Parenthesis(t) => sql_type_to_arrow(t)?,
374                sqlparser::ast::ArrayElemTypeDef::None => {
375                    return Err(ParseError::ValidationError(
376                        "ARRAY type requires element type, e.g. ARRAY<INT>".into(),
377                    ));
378                }
379            };
380            Ok(DataType::List(Arc::new(Field::new(
381                "item", item_type, true,
382            ))))
383        }
384
385        // Complex types (MAP, STRUCT, nested records) — use auto-discovery instead.
386        _ => Err(ParseError::ValidationError(format!(
387            "unsupported data type in hand-declared column: {sql_type:?} \
388             — use auto-discovery with an Avro source for complex types"
389        ))),
390    }
391}
392
393/// Checks if an expression is a `PROCTIME()` function call.
394fn is_proctime_call(expr: &sqlparser::ast::Expr) -> bool {
395    if let sqlparser::ast::Expr::Function(func) = expr {
396        if let Some(name) = func.name.0.last() {
397            return name.to_string().eq_ignore_ascii_case("proctime");
398        }
399    }
400    false
401}
402
403/// Parses watermark definition.
404fn parse_watermark(
405    wm: &WatermarkDef,
406    columns: &[ColumnDefinition],
407) -> Result<WatermarkSpec, ParseError> {
408    let column_name = wm.column.value.clone();
409
410    // Verify column exists and is a timestamp type
411    let col = columns
412        .iter()
413        .find(|c| c.name == column_name)
414        .ok_or_else(|| {
415            ParseError::ValidationError(format!(
416                "watermark column '{}' not found in column list",
417                column_name
418            ))
419        })?;
420
421    if !matches!(col.data_type, DataType::Timestamp(_, _)) {
422        return Err(ParseError::ValidationError(format!(
423            "watermark column '{}' must be a TIMESTAMP, found {:?}",
424            column_name, col.data_type
425        )));
426    }
427
428    // Check for PROCTIME() watermark expression
429    if let Some(expr) = &wm.expression {
430        if is_proctime_call(expr) {
431            return Ok(WatermarkSpec {
432                column: column_name,
433                max_out_of_orderness: Duration::ZERO,
434                is_processing_time: true,
435            });
436        }
437    }
438
439    // Parse the watermark expression to extract out-of-orderness.
440    // When expression is None (WATERMARK FOR col without AS), use zero delay.
441    let max_out_of_orderness = match &wm.expression {
442        Some(expr) => parse_watermark_expression(expr),
443        None => Duration::ZERO,
444    };
445
446    Ok(WatermarkSpec {
447        column: column_name,
448        max_out_of_orderness,
449        is_processing_time: false,
450    })
451}
452
453/// Parses watermark expression to extract the bounded out-of-orderness.
454fn parse_watermark_expression(expr: &sqlparser::ast::Expr) -> Duration {
455    use sqlparser::ast::Expr;
456
457    match expr {
458        Expr::BinaryOp { op, right, .. } => match op {
459            sqlparser::ast::BinaryOperator::Minus => parse_interval_expr(right),
460            _ => Duration::ZERO,
461        },
462        // If just the column name, assume zero lateness
463        Expr::Identifier(_) => Duration::ZERO,
464        // Default to 1 second for complex expressions
465        _ => Duration::from_secs(1),
466    }
467}
468
469/// Parses an interval expression to a Duration.
470fn parse_interval_expr(expr: &sqlparser::ast::Expr) -> Duration {
471    use sqlparser::ast::Expr;
472
473    let Expr::Interval(interval) = expr else {
474        return Duration::from_secs(1);
475    };
476
477    // Extract value and unit from interval
478    let value_str = match interval.value.as_ref() {
479        Expr::Value(v) => {
480            // v is ValueWithSpan, access the inner value
481            match &v.value {
482                sqlparser::ast::Value::SingleQuotedString(s) => s.clone(),
483                sqlparser::ast::Value::Number(n, _) => n.clone(),
484                _ => return Duration::from_secs(1),
485            }
486        }
487        _ => return Duration::from_secs(1),
488    };
489
490    let value: u64 = value_str.parse().unwrap_or(1);
491
492    // Determine unit
493    let unit = interval
494        .leading_field
495        .as_ref()
496        .map_or("second", |u| match u {
497            sqlparser::ast::DateTimeField::Microsecond => "microsecond",
498            sqlparser::ast::DateTimeField::Millisecond => "millisecond",
499            sqlparser::ast::DateTimeField::Minute => "minute",
500            sqlparser::ast::DateTimeField::Hour => "hour",
501            sqlparser::ast::DateTimeField::Day => "day",
502            _ => "second",
503        });
504
505    match unit {
506        "microsecond" | "microseconds" => Duration::from_micros(value),
507        "millisecond" | "milliseconds" => Duration::from_millis(value),
508        "minute" | "minutes" => Duration::from_secs(value * 60),
509        "hour" | "hours" => Duration::from_secs(value * 3600),
510        "day" | "days" => Duration::from_secs(value * 86400),
511        _ => Duration::from_secs(value),
512    }
513}
514
515#[cfg(test)]
516mod tests {
517    use super::*;
518    use crate::parser::{parse_streaming_sql, StreamingStatement};
519
520    fn parse_and_translate(sql: &str) -> Result<SourceDefinition, ParseError> {
521        let statements = parse_streaming_sql(sql)?;
522        let stmt = statements
523            .into_iter()
524            .next()
525            .ok_or_else(|| ParseError::StreamingError("No statement found".to_string()))?;
526        match stmt {
527            StreamingStatement::CreateSource(source) => translate_create_source(*source),
528            _ => Err(ParseError::StreamingError(
529                "Expected CREATE SOURCE".to_string(),
530            )),
531        }
532    }
533
534    #[test]
535    fn test_basic_source() {
536        let def =
537            parse_and_translate("CREATE SOURCE events (id BIGINT NOT NULL, name VARCHAR)").unwrap();
538
539        assert_eq!(def.name, "events");
540        assert_eq!(def.columns.len(), 2);
541        assert_eq!(def.columns[0].name, "id");
542        assert_eq!(def.columns[0].data_type, DataType::Int64);
543        assert!(!def.columns[0].nullable);
544        assert_eq!(def.columns[1].name, "name");
545        assert!(def.columns[1].nullable);
546    }
547
548    #[test]
549    fn test_source_with_options() {
550        let def = parse_and_translate(
551            "CREATE SOURCE events (id BIGINT) WITH (
552                'buffer_size' = '4096',
553                'backpressure' = 'reject'
554            )",
555        )
556        .unwrap();
557
558        assert_eq!(def.config.buffer_size, 4096);
559        assert_eq!(def.config.backpressure, BackpressureStrategy::Reject);
560    }
561
562    #[test]
563    fn test_source_with_watermark() {
564        let def = parse_and_translate(
565            "CREATE SOURCE events (
566                id BIGINT,
567                ts TIMESTAMP,
568                WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
569            )",
570        )
571        .unwrap();
572
573        assert!(def.watermark.is_some());
574        let wm = def.watermark.unwrap();
575        assert_eq!(wm.column, "ts");
576        assert_eq!(wm.max_out_of_orderness, Duration::from_secs(5));
577    }
578
579    #[test]
580    fn test_reject_channel_option() {
581        let result =
582            parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('channel' = 'mpsc')");
583
584        assert!(result.is_err());
585        let err = result.unwrap_err();
586        assert!(err.to_string().contains("channel"));
587    }
588
589    #[test]
590    fn test_reject_type_option() {
591        let result = parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('type' = 'spsc')");
592
593        assert!(result.is_err());
594    }
595
596    #[test]
597    fn test_buffer_size_bounds() {
598        // Too small
599        let result =
600            parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('buffer_size' = '1')");
601        assert!(result.is_err());
602
603        // Too large
604        let result = parse_and_translate(
605            "CREATE SOURCE events (id BIGINT) WITH ('buffer_size' = '999999999')",
606        );
607        assert!(result.is_err());
608
609        // Valid
610        let result =
611            parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('buffer_size' = '1024')");
612        assert!(result.is_ok());
613    }
614
615    #[test]
616    fn test_backpressure_strategies() {
617        assert_eq!(
618            BackpressureStrategy::from_str("block").unwrap(),
619            BackpressureStrategy::Block
620        );
621        assert_eq!(
622            BackpressureStrategy::from_str("drop_oldest").unwrap(),
623            BackpressureStrategy::DropOldest
624        );
625        assert_eq!(
626            BackpressureStrategy::from_str("reject").unwrap(),
627            BackpressureStrategy::Reject
628        );
629        assert!(BackpressureStrategy::from_str("invalid").is_err());
630    }
631
632    #[test]
633    fn test_wait_strategies() {
634        assert_eq!(WaitStrategy::from_str("spin").unwrap(), WaitStrategy::Spin);
635        assert_eq!(
636            WaitStrategy::from_str("spin_yield").unwrap(),
637            WaitStrategy::SpinYield
638        );
639        assert_eq!(WaitStrategy::from_str("park").unwrap(), WaitStrategy::Park);
640        assert!(WaitStrategy::from_str("invalid").is_err());
641    }
642
643    #[test]
644    fn test_sql_type_conversions() {
645        let def = parse_and_translate(
646            "CREATE SOURCE types (
647                a TINYINT,
648                b SMALLINT,
649                c INT,
650                d BIGINT,
651                e FLOAT,
652                f DOUBLE,
653                g DECIMAL(10,2),
654                h VARCHAR(255),
655                i TEXT,
656                j BOOLEAN,
657                k TIMESTAMP,
658                l DATE
659            )",
660        )
661        .unwrap();
662
663        assert_eq!(def.columns.len(), 12);
664        assert_eq!(def.columns[0].data_type, DataType::Int8);
665        assert_eq!(def.columns[1].data_type, DataType::Int16);
666        assert_eq!(def.columns[2].data_type, DataType::Int32);
667        assert_eq!(def.columns[3].data_type, DataType::Int64);
668        assert_eq!(def.columns[4].data_type, DataType::Float32);
669        assert_eq!(def.columns[5].data_type, DataType::Float64);
670        assert_eq!(def.columns[6].data_type, DataType::Decimal128(10, 2));
671        assert_eq!(def.columns[7].data_type, DataType::Utf8);
672        assert_eq!(def.columns[8].data_type, DataType::Utf8);
673        assert_eq!(def.columns[9].data_type, DataType::Boolean);
674        assert!(matches!(
675            def.columns[10].data_type,
676            DataType::Timestamp(_, _)
677        ));
678        assert_eq!(def.columns[11].data_type, DataType::Date32);
679    }
680
681    #[test]
682    fn test_schema_generation() {
683        let def = parse_and_translate(
684            "CREATE SOURCE events (id BIGINT NOT NULL, name VARCHAR NOT NULL, value DOUBLE)",
685        )
686        .unwrap();
687
688        let schema = def.schema;
689        assert_eq!(schema.fields().len(), 3);
690        assert_eq!(schema.field(0).name(), "id");
691        assert!(!schema.field(0).is_nullable());
692        assert_eq!(schema.field(1).name(), "name");
693        assert!(!schema.field(1).is_nullable());
694        assert_eq!(schema.field(2).name(), "value");
695        assert!(schema.field(2).is_nullable());
696    }
697
698    #[test]
699    fn test_watermark_column_not_found() {
700        let result = parse_and_translate(
701            "CREATE SOURCE events (
702                id BIGINT,
703                WATERMARK FOR nonexistent AS nonexistent - INTERVAL '1' SECOND
704            )",
705        );
706
707        assert!(result.is_err());
708        assert!(result.unwrap_err().to_string().contains("not found"));
709    }
710
711    #[test]
712    fn test_watermark_wrong_type() {
713        let result = parse_and_translate(
714            "CREATE SOURCE events (
715                name VARCHAR,
716                WATERMARK FOR name AS name - INTERVAL '1' SECOND
717            )",
718        );
719
720        assert!(result.is_err());
721        assert!(result
722            .unwrap_err()
723            .to_string()
724            .contains("must be a TIMESTAMP"));
725    }
726
727    #[test]
728    fn test_watermark_milliseconds() {
729        let def = parse_and_translate(
730            "CREATE SOURCE events (
731                ts TIMESTAMP,
732                WATERMARK FOR ts AS ts - INTERVAL '100' MILLISECOND
733            )",
734        )
735        .unwrap();
736
737        let wm = def.watermark.unwrap();
738        assert_eq!(wm.max_out_of_orderness, Duration::from_millis(100));
739    }
740
741    #[test]
742    fn test_watermark_minutes() {
743        let def = parse_and_translate(
744            "CREATE SOURCE events (
745                ts TIMESTAMP,
746                WATERMARK FOR ts AS ts - INTERVAL '5' MINUTE
747            )",
748        )
749        .unwrap();
750
751        let wm = def.watermark.unwrap();
752        assert_eq!(wm.max_out_of_orderness, Duration::from_secs(300));
753    }
754
755    #[test]
756    fn test_track_stats_option() {
757        let def =
758            parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('track_stats' = 'true')")
759                .unwrap();
760
761        assert!(def.config.track_stats);
762    }
763
764    #[test]
765    fn test_wait_strategy_option() {
766        let def =
767            parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('wait_strategy' = 'park')")
768                .unwrap();
769
770        assert_eq!(def.config.wait_strategy, WaitStrategy::Park);
771    }
772
773    #[test]
774    fn test_default_config() {
775        let def = parse_and_translate("CREATE SOURCE events (id BIGINT)").unwrap();
776
777        assert_eq!(def.config.buffer_size, DEFAULT_BUFFER_SIZE);
778        assert_eq!(def.config.backpressure, BackpressureStrategy::Block);
779        assert_eq!(def.config.wait_strategy, WaitStrategy::SpinYield);
780        assert!(!def.config.track_stats);
781    }
782
783    #[test]
784    fn test_external_connector_options_ignored() {
785        // External connector options should be accepted but not affect config
786        let def = parse_and_translate(
787            "CREATE SOURCE events (id BIGINT) WITH (
788                'connector' = 'kafka',
789                'topic' = 'events',
790                'bootstrap.servers' = 'localhost:9092',
791                'buffer_size' = '8192'
792            )",
793        )
794        .unwrap();
795
796        // Only buffer_size should affect config
797        assert_eq!(def.config.buffer_size, 8192);
798    }
799
800    #[test]
801    fn test_source_watermark_no_expression() {
802        let def = parse_and_translate(
803            "CREATE SOURCE events (
804                ts TIMESTAMP,
805                WATERMARK FOR ts
806            )",
807        )
808        .unwrap();
809
810        assert!(def.watermark.is_some());
811        let wm = def.watermark.unwrap();
812        assert_eq!(wm.column, "ts");
813        assert_eq!(wm.max_out_of_orderness, Duration::ZERO);
814    }
815
816    #[test]
817    fn test_source_watermark_bigint_column_rejected() {
818        let result = parse_and_translate(
819            "CREATE SOURCE events (
820                ts BIGINT,
821                WATERMARK FOR ts
822            )",
823        );
824
825        let error = result.unwrap_err().to_string();
826        assert!(error.contains("must be a TIMESTAMP"), "{error}");
827    }
828
829    #[test]
830    fn test_watermark_proctime() {
831        let def = parse_and_translate(
832            "CREATE SOURCE events (
833                ts TIMESTAMP,
834                WATERMARK FOR ts AS PROCTIME()
835            )",
836        )
837        .unwrap();
838
839        assert!(def.watermark.is_some());
840        let wm = def.watermark.unwrap();
841        assert_eq!(wm.column, "ts");
842        assert!(wm.is_processing_time);
843        assert_eq!(wm.max_out_of_orderness, Duration::ZERO);
844    }
845
846    #[test]
847    fn test_watermark_event_time_not_proctime() {
848        let def = parse_and_translate(
849            "CREATE SOURCE events (
850                ts TIMESTAMP,
851                WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
852            )",
853        )
854        .unwrap();
855
856        let wm = def.watermark.unwrap();
857        assert!(!wm.is_processing_time);
858    }
859
860    #[test]
861    fn array_of_int_parses() {
862        let def = parse_and_translate("CREATE SOURCE events (tags ARRAY<INT>)").unwrap();
863        match &def.columns[0].data_type {
864            DataType::List(field) => assert_eq!(field.data_type(), &DataType::Int32),
865            other => panic!("expected DataType::List, got {other:?}"),
866        }
867    }
868
869    #[test]
870    fn decimal_with_precision_parses() {
871        let def = parse_and_translate("CREATE SOURCE events (amount DECIMAL(10, 2))").unwrap();
872        assert_eq!(def.columns[0].data_type, DataType::Decimal128(10, 2));
873    }
874
875    /// Hand-declared MAP columns point users at auto-discovery.
876    #[test]
877    fn hand_declared_map_column_errors_actionably() {
878        let err =
879            parse_and_translate("CREATE SOURCE events (data MAP(VARCHAR, VARCHAR))").unwrap_err();
880        let msg = err.to_string();
881        assert!(
882            msg.contains("use auto-discovery") || msg.contains("unsupported"),
883            "expected actionable error for hand-declared MAP, got: {msg}"
884        );
885    }
886}