1#![forbid(unsafe_code)]
2
3use std::collections::HashMap;
30use std::sync::RwLock;
31
32use crate::{SqlError, SqlResult};
33
34#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct LatenessAnnotation {
39 pub column: String,
40 pub lateness_ms: u64,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
47pub enum IncrementalViewStatement {
48 Create {
49 name: String,
50 body_sql: String,
51 is_materialized: bool,
52 lateness: Vec<LatenessAnnotation>,
53 },
54 DeclareRecursive {
55 name: String,
56 body_sql: String,
57 },
58 Refresh {
59 name: String,
60 },
61 Drop {
62 name: String,
63 },
64}
65
66#[derive(Debug, Clone)]
70pub struct IncrementalViewEntry {
71 pub body_sql: String,
72 pub is_materialized: bool,
73 pub is_recursive: bool,
74 pub lateness: Vec<LatenessAnnotation>,
75}
76
77#[derive(Debug, Default)]
83pub struct IncrementalViewRegistry {
84 views: RwLock<HashMap<String, IncrementalViewEntry>>,
85}
86
87impl IncrementalViewRegistry {
88 pub fn new() -> Self {
89 Self::default()
90 }
91
92 pub fn register(&self, name: impl Into<String>, entry: IncrementalViewEntry) -> SqlResult<()> {
93 let mut views = self.views.write().map_err(|_| SqlError::DataFusion {
94 message: "incremental view registry lock poisoned".into(),
95 })?;
96 views.insert(name.into(), entry);
97 Ok(())
98 }
99
100 pub fn remove(&self, name: &str) -> SqlResult<bool> {
101 let mut views = self.views.write().map_err(|_| SqlError::DataFusion {
102 message: "incremental view registry lock poisoned".into(),
103 })?;
104 Ok(views.remove(name).is_some())
105 }
106
107 pub fn get(&self, name: &str) -> SqlResult<Option<IncrementalViewEntry>> {
108 let views = self.views.read().map_err(|_| SqlError::DataFusion {
109 message: "incremental view registry lock poisoned".into(),
110 })?;
111 Ok(views.get(name).cloned())
112 }
113
114 pub fn contains(&self, name: &str) -> bool {
115 self.views
116 .read()
117 .map(|v| v.contains_key(name))
118 .unwrap_or(false)
119 }
120
121 pub fn view_names(&self) -> SqlResult<Vec<String>> {
122 let views = self.views.read().map_err(|_| SqlError::DataFusion {
123 message: "incremental view registry lock poisoned".into(),
124 })?;
125 Ok(views.keys().cloned().collect())
126 }
127}
128
129pub fn parse_incremental_view_statement(sql: &str) -> SqlResult<Option<IncrementalViewStatement>> {
135 let trimmed = sql.trim().trim_end_matches(';');
136 let upper = trimmed.to_uppercase();
137
138 let is_materialized = upper.starts_with("CREATE MATERIALIZED INCREMENTAL VIEW ");
141 if is_materialized || upper.starts_with("CREATE INCREMENTAL VIEW ") {
142 let prefix = if is_materialized {
143 "CREATE MATERIALIZED INCREMENTAL VIEW "
144 } else {
145 "CREATE INCREMENTAL VIEW "
146 };
147 let rest = trimmed
148 .get(prefix.len()..)
149 .ok_or_else(|| SqlError::Unsupported {
150 feature: "CREATE INCREMENTAL VIEW".into(),
151 })?;
152 let (name, body_with_lateness) = split_name_and_body(rest)?;
153 let (body_sql, lateness) = split_body_and_lateness(&body_with_lateness)?;
154 return Ok(Some(IncrementalViewStatement::Create {
155 name,
156 body_sql,
157 is_materialized,
158 lateness,
159 }));
160 }
161
162 let mv_prefix = if upper.starts_with("CREATE OR REPLACE MATERIALIZED VIEW ") {
169 Some("CREATE OR REPLACE MATERIALIZED VIEW ")
170 } else if upper.starts_with("CREATE MATERIALIZED VIEW ") {
171 Some("CREATE MATERIALIZED VIEW ")
172 } else {
173 None
174 };
175 if let Some(prefix) = mv_prefix {
176 let rest = trimmed
177 .get(prefix.len()..)
178 .ok_or_else(|| SqlError::Unsupported {
179 feature: "CREATE MATERIALIZED VIEW".into(),
180 })?;
181 let (name, body_with_lateness) = split_name_and_body(rest)?;
182 let (body_sql, lateness) = split_body_and_lateness(&body_with_lateness)?;
183 return Ok(Some(IncrementalViewStatement::Create {
184 name,
185 body_sql,
186 is_materialized: true,
187 lateness,
188 }));
189 }
190
191 if upper.starts_with("DECLARE RECURSIVE VIEW ") {
193 let rest = trimmed
194 .get("DECLARE RECURSIVE VIEW ".len()..)
195 .ok_or_else(|| SqlError::Unsupported {
196 feature: "DECLARE RECURSIVE VIEW".into(),
197 })?;
198 let (name, body_sql) = split_name_and_body(rest)?;
199 let (body_sql, _lateness) = split_body_and_lateness(&body_sql)?;
200 return Ok(Some(IncrementalViewStatement::DeclareRecursive {
201 name,
202 body_sql,
203 }));
204 }
205
206 if upper.starts_with("REFRESH INCREMENTAL VIEW ") {
208 let name = trimmed
209 .get("REFRESH INCREMENTAL VIEW ".len()..)
210 .ok_or_else(|| SqlError::Unsupported {
211 feature: "REFRESH INCREMENTAL VIEW".into(),
212 })?
213 .trim()
214 .to_string();
215 if name.is_empty() {
216 return Err(SqlError::EmptyTableName);
217 }
218 return Ok(Some(IncrementalViewStatement::Refresh { name }));
219 }
220
221 if upper.starts_with("DROP INCREMENTAL VIEW ") {
223 let name = trimmed
224 .get("DROP INCREMENTAL VIEW ".len()..)
225 .ok_or_else(|| SqlError::Unsupported {
226 feature: "DROP INCREMENTAL VIEW".into(),
227 })?
228 .trim()
229 .to_string();
230 if name.is_empty() {
231 return Err(SqlError::EmptyTableName);
232 }
233 return Ok(Some(IncrementalViewStatement::Drop { name }));
234 }
235
236 if upper.starts_with("REFRESH MATERIALIZED VIEW ") {
238 let name = trimmed
239 .get("REFRESH MATERIALIZED VIEW ".len()..)
240 .ok_or_else(|| SqlError::Unsupported {
241 feature: "REFRESH MATERIALIZED VIEW".into(),
242 })?
243 .trim()
244 .to_string();
245 if name.is_empty() {
246 return Err(SqlError::EmptyTableName);
247 }
248 return Ok(Some(IncrementalViewStatement::Refresh { name }));
249 }
250
251 if upper.starts_with("DROP MATERIALIZED VIEW ") {
253 let name = trimmed
254 .get("DROP MATERIALIZED VIEW ".len()..)
255 .ok_or_else(|| SqlError::Unsupported {
256 feature: "DROP MATERIALIZED VIEW".into(),
257 })?
258 .trim()
259 .to_string();
260 if name.is_empty() {
261 return Err(SqlError::EmptyTableName);
262 }
263 return Ok(Some(IncrementalViewStatement::Drop { name }));
264 }
265
266 Ok(None)
267}
268
269pub enum IncrementalViewResult {
276 Created(String),
278 Dropped(String),
280 Refresh(String),
282 Recursive(String),
284}
285
286pub fn execute_incremental_view_ddl(
287 registry: &IncrementalViewRegistry,
288 sql: &str,
289) -> SqlResult<Option<IncrementalViewResult>> {
290 let Some(stmt) = parse_incremental_view_statement(sql)? else {
291 return Ok(None);
292 };
293
294 match stmt {
295 IncrementalViewStatement::Create {
296 ref name,
297 ref body_sql,
298 is_materialized,
299 ref lateness,
300 } => {
301 registry.register(
302 name.clone(),
303 IncrementalViewEntry {
304 body_sql: body_sql.clone(),
305 is_materialized,
306 is_recursive: false,
307 lateness: lateness.clone(),
308 },
309 )?;
310 Ok(Some(IncrementalViewResult::Created(name.clone())))
311 }
312
313 IncrementalViewStatement::DeclareRecursive {
314 ref name,
315 ref body_sql,
316 } => {
317 registry.register(
318 name.clone(),
319 IncrementalViewEntry {
320 body_sql: body_sql.clone(),
321 is_materialized: false,
322 is_recursive: true,
323 lateness: vec![],
324 },
325 )?;
326 Ok(Some(IncrementalViewResult::Recursive(name.clone())))
327 }
328
329 IncrementalViewStatement::Refresh { ref name } => {
330 if !registry.contains(name) {
331 return Err(SqlError::Unsupported {
332 feature: format!("REFRESH INCREMENTAL VIEW: view '{name}' is not registered"),
333 });
334 }
335 Ok(Some(IncrementalViewResult::Refresh(name.clone())))
336 }
337
338 IncrementalViewStatement::Drop { ref name } => {
339 registry.remove(name)?;
340 Ok(Some(IncrementalViewResult::Dropped(name.clone())))
341 }
342 }
343}
344
345fn split_name_and_body(rest: &str) -> SqlResult<(String, String)> {
349 let upper = rest.to_uppercase();
350 let as_pos = upper.find(" AS ").ok_or_else(|| SqlError::Unsupported {
351 feature: "CREATE INCREMENTAL VIEW / DECLARE RECURSIVE VIEW requires AS <query>".into(),
352 })?;
353 let name = rest[..as_pos].trim().to_string();
354 let body = rest[as_pos + 4..].trim().to_string();
355 if name.is_empty() {
356 return Err(SqlError::EmptyTableName);
357 }
358 if body.is_empty() {
359 return Err(SqlError::EmptyQuery);
360 }
361 Ok((name, body))
362}
363
364fn split_body_and_lateness(
371 body_with_lateness: &str,
372) -> SqlResult<(String, Vec<LatenessAnnotation>)> {
373 let upper = body_with_lateness.to_uppercase();
374
375 let Some(lat_pos) = find_lateness_clause_start(&upper) else {
378 return Ok((body_with_lateness.trim().to_string(), vec![]));
379 };
380
381 let body_sql = body_with_lateness[..lat_pos].trim().to_string();
382 let lateness_str = &body_with_lateness[lat_pos..];
383 let lateness = parse_lateness_clauses(lateness_str)?;
384 Ok((body_sql, lateness))
385}
386
387fn find_lateness_clause_start(upper: &str) -> Option<usize> {
389 let bytes = upper.as_bytes();
391 let keyword = b"LATENESS";
392 let mut depth = 0usize;
393 let mut i = 0usize;
394 while i + keyword.len() <= bytes.len() {
395 let Some(&b) = bytes.get(i) else {
396 break;
397 };
398 match b {
399 b'(' => {
400 depth += 1;
401 i += 1;
402 }
403 b')' => {
404 depth = depth.saturating_sub(1);
405 i += 1;
406 }
407 _ if depth == 0 && bytes.get(i..).is_some_and(|s| s.starts_with(keyword)) => {
408 let before_ok =
409 i == 0 || bytes.get(i - 1).is_some_and(|b| !b.is_ascii_alphanumeric());
410 let after = i + keyword.len();
411 let after_ok = bytes.get(after).is_none_or(|b| !b.is_ascii_alphanumeric());
412 if before_ok && after_ok {
413 return Some(i);
414 }
415 i += 1;
416 }
417 _ => {
418 i += 1;
419 }
420 }
421 }
422 None
423}
424
425fn parse_lateness_clauses(lateness_str: &str) -> SqlResult<Vec<LatenessAnnotation>> {
427 let upper = lateness_str.to_uppercase();
429 let mut result = Vec::new();
430 let mut remaining = lateness_str.trim();
431
432 loop {
433 let upper_rem = remaining.to_uppercase();
434 let stripped = if upper_rem.starts_with("LATENESS ") {
435 &remaining["LATENESS ".len()..]
436 } else if upper_rem.starts_with(", LATENESS ") {
437 &remaining[", LATENESS ".len()..]
438 } else {
439 break;
440 };
441
442 let tokens: Vec<&str> = stripped.splitn(5, char::is_whitespace).collect();
444 if tokens.len() < 4 {
445 break;
446 }
447 let col = tokens
448 .first()
449 .copied()
450 .unwrap_or("")
451 .trim_matches(',')
452 .to_string();
453 let interval_str = tokens.get(2).copied().unwrap_or("").trim_matches('\'');
455 let unit_str = tokens.get(3).copied().unwrap_or("").trim_matches(',');
456 let n: u64 = interval_str.parse().map_err(|_| SqlError::Unsupported {
457 feature: format!("LATENESS INTERVAL value '{interval_str}' is not a valid integer"),
458 })?;
459 let ms = match unit_str.to_uppercase().as_str() {
460 "SECOND" | "SECONDS" => n * 1000,
461 "MINUTE" | "MINUTES" => n * 60_000,
462 "HOUR" | "HOURS" => n * 3_600_000,
463 "DAY" | "DAYS" => n * 86_400_000,
464 "MILLISECOND" | "MILLISECONDS" | "MS" => n,
465 _ => {
466 return Err(SqlError::Unsupported {
467 feature: format!(
468 "LATENESS interval unit '{unit_str}' is not supported \
469 (expected SECOND, MINUTE, HOUR, DAY, or MILLISECOND)"
470 ),
471 });
472 }
473 };
474
475 result.push(LatenessAnnotation {
476 column: col,
477 lateness_ms: ms,
478 });
479
480 let consumed_upper: String = upper_rem
482 .chars()
483 .take("LATENESS ".len() + stripped.len() - stripped.trim_start().len())
484 .collect();
485 let _ = consumed_upper; let next = remaining[1..].to_uppercase().find("LATENESS");
488 match next {
489 Some(pos) => {
490 remaining = &remaining[1 + pos..];
491 }
492 None => break,
493 }
494 }
495
496 let _ = upper; Ok(result)
498}
499
500#[cfg(test)]
503mod tests {
504 use super::*;
505
506 #[test]
507 fn parse_create_incremental_view() {
508 let sql = "CREATE INCREMENTAL VIEW revenue AS SELECT SUM(amount) FROM orders";
509 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
510 assert!(matches!(
511 stmt,
512 IncrementalViewStatement::Create { ref name, is_materialized: false, .. }
513 if name == "revenue"
514 ));
515 }
516
517 #[test]
518 fn parse_create_materialized_incremental_view() {
519 let sql = "CREATE MATERIALIZED INCREMENTAL VIEW snap AS SELECT * FROM t";
520 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
521 assert!(matches!(
522 stmt,
523 IncrementalViewStatement::Create {
524 is_materialized: true,
525 ..
526 }
527 ));
528 }
529
530 #[test]
531 fn parse_create_materialized_view_maps_to_ivm() {
532 let sql = "CREATE MATERIALIZED VIEW revenue AS SELECT SUM(amount) AS t FROM orders";
535 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
536 assert!(matches!(
537 stmt,
538 IncrementalViewStatement::Create { ref name, body_sql: ref body, is_materialized: true, .. }
539 if name == "revenue" && body.starts_with("SELECT SUM(amount)")
540 ));
541 }
542
543 #[test]
544 fn parse_create_or_replace_materialized_view() {
545 let sql = "CREATE OR REPLACE MATERIALIZED VIEW mv AS SELECT * FROM t";
546 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
547 assert!(matches!(
548 stmt,
549 IncrementalViewStatement::Create { ref name, is_materialized: true, .. } if name == "mv"
550 ));
551 }
552
553 #[test]
554 fn parse_create_materialized_view_with_lateness() {
555 let sql = "CREATE MATERIALIZED VIEW ev AS SELECT * FROM s \
556 LATENESS event_ts INTERVAL '5' MINUTE";
557 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
558 match stmt {
559 IncrementalViewStatement::Create {
560 is_materialized,
561 lateness,
562 body_sql,
563 ..
564 } => {
565 assert!(is_materialized);
566 assert_eq!(lateness.len(), 1, "LATENESS annotation is parsed");
567 assert!(!body_sql.to_uppercase().contains("LATENESS"));
568 }
569 other => panic!("expected Create, got {other:?}"),
570 }
571 }
572
573 #[test]
574 fn parse_refresh_and_drop_materialized_view() {
575 let refresh = parse_incremental_view_statement("REFRESH MATERIALIZED VIEW revenue")
576 .unwrap()
577 .unwrap();
578 assert!(matches!(
579 refresh,
580 IncrementalViewStatement::Refresh { ref name } if name == "revenue"
581 ));
582 let drop = parse_incremental_view_statement("DROP MATERIALIZED VIEW revenue;")
583 .unwrap()
584 .unwrap();
585 assert!(matches!(
586 drop,
587 IncrementalViewStatement::Drop { ref name } if name == "revenue"
588 ));
589 }
590
591 #[test]
592 fn materialized_view_does_not_shadow_materialized_incremental_view() {
593 let sql = "CREATE MATERIALIZED INCREMENTAL VIEW snap AS SELECT * FROM t";
596 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
597 assert!(matches!(
598 stmt,
599 IncrementalViewStatement::Create { ref name, is_materialized: true, .. } if name == "snap"
600 ));
601 }
602
603 #[test]
604 fn parse_declare_recursive_view() {
605 let sql = "DECLARE RECURSIVE VIEW reach AS SELECT dst FROM edges WHERE src = 0";
606 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
607 assert!(matches!(
608 stmt,
609 IncrementalViewStatement::DeclareRecursive { ref name, .. } if name == "reach"
610 ));
611 }
612
613 #[test]
614 fn parse_refresh_incremental_view() {
615 let sql = "REFRESH INCREMENTAL VIEW revenue";
616 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
617 assert!(matches!(
618 stmt,
619 IncrementalViewStatement::Refresh { ref name } if name == "revenue"
620 ));
621 }
622
623 #[test]
624 fn parse_drop_incremental_view() {
625 let sql = "DROP INCREMENTAL VIEW revenue;";
626 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
627 assert!(matches!(
628 stmt,
629 IncrementalViewStatement::Drop { ref name } if name == "revenue"
630 ));
631 }
632
633 #[test]
634 fn non_incremental_sql_returns_none() {
635 let sql = "SELECT 1";
636 assert!(parse_incremental_view_statement(sql).unwrap().is_none());
637 }
638
639 #[test]
640 fn parse_create_with_lateness() {
641 let sql =
642 "CREATE INCREMENTAL VIEW ev AS SELECT * FROM s LATENESS event_ts INTERVAL '5' MINUTE";
643 let stmt = parse_incremental_view_statement(sql).unwrap().unwrap();
644 if let IncrementalViewStatement::Create { lateness, .. } = stmt {
645 assert_eq!(lateness.len(), 1);
646 assert_eq!(lateness[0].column, "event_ts");
647 assert_eq!(lateness[0].lateness_ms, 5 * 60_000);
648 } else {
649 panic!("expected Create");
650 }
651 }
652
653 #[test]
654 fn registry_register_and_get() {
655 let reg = IncrementalViewRegistry::new();
656 reg.register(
657 "v1",
658 IncrementalViewEntry {
659 body_sql: "SELECT 1".into(),
660 is_materialized: false,
661 is_recursive: false,
662 lateness: vec![],
663 },
664 )
665 .unwrap();
666 assert!(reg.contains("v1"));
667 let entry = reg.get("v1").unwrap().unwrap();
668 assert_eq!(entry.body_sql, "SELECT 1");
669 }
670
671 #[test]
672 fn execute_ddl_create_and_drop() {
673 let reg = IncrementalViewRegistry::new();
674 let result =
675 execute_incremental_view_ddl(®, "CREATE INCREMENTAL VIEW v AS SELECT 1").unwrap();
676 assert!(matches!(result, Some(IncrementalViewResult::Created(_))));
677 assert!(reg.contains("v"));
678
679 execute_incremental_view_ddl(®, "DROP INCREMENTAL VIEW v").unwrap();
680 assert!(!reg.contains("v"));
681 }
682
683 #[test]
684 fn execute_ddl_refresh_returns_refresh_variant() {
685 let reg = IncrementalViewRegistry::new();
686 execute_incremental_view_ddl(®, "CREATE INCREMENTAL VIEW v AS SELECT 1").unwrap();
687 let result = execute_incremental_view_ddl(®, "REFRESH INCREMENTAL VIEW v").unwrap();
688 assert!(matches!(result, Some(IncrementalViewResult::Refresh(_))));
689 }
690
691 #[test]
692 fn execute_ddl_refresh_missing_returns_error() {
693 let reg = IncrementalViewRegistry::new();
694 let err = execute_incremental_view_ddl(®, "REFRESH INCREMENTAL VIEW nonexistent");
695 assert!(err.is_err());
696 }
697}