1#[allow(clippy::disallowed_types)] use 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#[derive(Debug, Clone)]
21pub struct WatermarkSpec {
22 pub column: String,
24 pub max_out_of_orderness: Duration,
26 pub is_processing_time: bool,
31}
32
33#[derive(Debug, Clone)]
35pub struct SourceConfigOptions {
36 pub buffer_size: usize,
38 pub backpressure: BackpressureStrategy,
40 pub wait_strategy: WaitStrategy,
42 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#[derive(Debug, Clone)]
59pub struct ColumnDefinition {
60 pub name: String,
62 pub data_type: DataType,
64 pub nullable: bool,
66}
67
68#[derive(Debug, Clone)]
73pub struct SourceDefinition {
74 pub name: String,
76 pub columns: Vec<ColumnDefinition>,
78 pub schema: SchemaRef,
80 pub watermark: Option<WatermarkSpec>,
82 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#[derive(Debug, Clone)]
96pub struct SinkDefinition {
97 pub name: String,
99 pub input: String,
101 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
113pub 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
128pub 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
163pub fn translate_create_sink(stmt: CreateSinkStatement) -> Result<SinkDefinition, ParseError> {
171 validate_source_options(&stmt.with_options)?;
173
174 let config = parse_source_options(&stmt.with_options)?;
176
177 let input = match stmt.from {
179 SinkFrom::Table(name) => name.to_string(),
180 SinkFrom::Query(_) => {
181 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
195fn validate_source_options(options: &HashMap<String, String>) -> Result<(), ParseError> {
197 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 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
215fn 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 _ => {}
241 }
242 }
243
244 Ok(config)
245}
246
247fn 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
273fn 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
285fn convert_columns(columns: &[ColumnDef]) -> Result<Vec<ColumnDefinition>, ParseError> {
287 columns.iter().map(convert_column).collect()
288}
289
290fn convert_column(col: &ColumnDef) -> Result<ColumnDefinition, ParseError> {
292 let data_type = sql_type_to_arrow(&col.data_type)?;
293
294 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
307pub fn sql_type_to_arrow(sql_type: &SqlDataType) -> Result<DataType, ParseError> {
313 match sql_type {
314 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 SqlDataType::Float(_) | SqlDataType::Real => Ok(DataType::Float32),
325 SqlDataType::Double(_) | SqlDataType::DoublePrecision => Ok(DataType::Float64),
326
327 SqlDataType::Decimal(info) | SqlDataType::Numeric(info) => {
329 #[allow(clippy::cast_possible_truncation)] 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), };
335 Ok(DataType::Decimal128(precision, scale))
336 }
337
338 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 SqlDataType::Binary(_)
351 | SqlDataType::Varbinary(_)
352 | SqlDataType::Blob(_)
353 | SqlDataType::Bytea => Ok(DataType::Binary),
354
355 SqlDataType::Boolean | SqlDataType::Bool => Ok(DataType::Boolean),
357
358 SqlDataType::Date => Ok(DataType::Date32),
360 SqlDataType::Time(_, _) => Ok(DataType::Time64(TimeUnit::Microsecond)),
361 SqlDataType::Timestamp(_, _) => Ok(DataType::Timestamp(TimeUnit::Microsecond, None)),
362
363 SqlDataType::Interval { .. } => Ok(DataType::Interval(
365 arrow::datatypes::IntervalUnit::MonthDayNano,
366 )),
367
368 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 _ => 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
393fn 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
403fn parse_watermark(
405 wm: &WatermarkDef,
406 columns: &[ColumnDefinition],
407) -> Result<WatermarkSpec, ParseError> {
408 let column_name = wm.column.value.clone();
409
410 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 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 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
453fn 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 Expr::Identifier(_) => Duration::ZERO,
464 _ => Duration::from_secs(1),
466 }
467}
468
469fn 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 let value_str = match interval.value.as_ref() {
479 Expr::Value(v) => {
480 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 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 let result =
600 parse_and_translate("CREATE SOURCE events (id BIGINT) WITH ('buffer_size' = '1')");
601 assert!(result.is_err());
602
603 let result = parse_and_translate(
605 "CREATE SOURCE events (id BIGINT) WITH ('buffer_size' = '999999999')",
606 );
607 assert!(result.is_err());
608
609 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 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 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 #[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}