use reifydb_core::sort::SortDirection;
use reifydb_value::{
error::{AstErrorKind, Error, TypeError},
fragment::Fragment,
value::{duration::Duration, identity::IdentityKind},
};
use crate::{
Result,
ast::{
ast::{
AstBindingProtocolKind, AstColumnProperty, AstColumnPropertyEntry, AstColumnPropertyKind,
AstColumnToCreate, AstCreate, AstCreateColumnProperty, AstCreateDeferredView,
AstCreateDictionary, AstCreateEvent, AstCreateHandler, AstCreateMigration, AstCreateNamespace,
AstCreatePrimaryKey, AstCreateProcedure, AstCreateQueue, AstCreateRelationship,
AstCreateRemoteNamespace, AstCreateRingBuffer, AstCreateSeries, AstCreateSubscription,
AstCreateSumType, AstCreateTable, AstCreateTag, AstCreateTest, AstCreateTransactionalView,
AstHydrationConfig, AstIndexColumn, AstJoinPick, AstJoinRetention, AstPersistent,
AstPolicyTargetType, AstPrimaryKey, AstProcedureParam, AstQueueDeduplicate, AstQueueDispatch,
AstQueueFifo, AstQueueRetention, AstQueueRetry, AstRelationshipCardinality,
AstRelationshipJunction, AstRowSettings, AstStatement, AstTimeDeclaration,
AstTimestampPrecision, AstTtl, AstType, AstVariant, AstViewStorageKind, AstViewWithClause,
},
identifier::{
MaybeQualifiedDeferredViewIdentifier, MaybeQualifiedDictionaryIdentifier,
MaybeQualifiedNamespaceIdentifier, MaybeQualifiedProcedureIdentifier,
MaybeQualifiedQueueIdentifier, MaybeQualifiedRingBufferIdentifier,
MaybeQualifiedSeriesIdentifier, MaybeQualifiedSumTypeIdentifier, MaybeQualifiedTableIdentifier,
MaybeQualifiedTestIdentifier, MaybeQualifiedTransactionalViewIdentifier,
},
parse::{Parser, Precedence},
},
bump::{BumpBox, BumpFragment},
duration::{DurationBound, FOREVER, compile_duration},
error::{OperationKind, RqlError},
token::{
keyword::{
Keyword,
Keyword::{
Create, Deferred, Dictionary, Exists, For, If, Namespace, Remote, Replace, Ringbuffer,
Series, Subscription, Table, Tag, Test, Transactional, View,
},
},
operator::{
Operator,
Operator::{Colon, Not, Or},
},
separator::{Separator, Separator::Comma},
token::{Literal, Token, TokenKind},
},
};
const QUEUE_OPTION_KEYS: &str = "'fifo', 'deduplicate', 'retention', 'retry', or 'time'";
const TIME_VALUES: &str = "'none', 'event(<column>)', or 'processing'";
const QUEUE_FIFO_KEYS: &str = "'partitions' or 'ordered_by'";
const QUEUE_DEDUPLICATE_KEYS: &str = "'by' or 'ttl'";
const QUEUE_RETENTION_KEYS: &str = "'done'";
const QUEUE_RETRY_KEYS: &str = "'attempts' or 'backoff'";
const ROW_CONFIG_KEYS: &str = "'ttl', 'persistent', or 'on'";
const JOIN_WITH_KEYS: &str = "'retention', 'snapshot', 'latest', or 'earliest'";
fn unexpected_queue_option(token: &Token<'_>, expected: &str) -> Error {
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: expected.to_string(),
},
message: format!("expected {}, found `{}`", expected, token.fragment.text()),
fragment: token.fragment.to_owned(),
})
}
fn missing_queue_dispatch(token: &Token<'_>) -> Error {
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'fifo'".to_string(),
},
message: "CREATE QUEUE requires a dispatch block, for example WITH { fifo: {} }".to_string(),
fragment: token.fragment.to_owned(),
})
}
fn duplicate_queue_dispatch(token: &Token<'_>) -> Error {
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "exactly one dispatch block".to_string(),
},
message: format!(
"a queue declares exactly one dispatch block, found a second `{}`",
token.fragment.text()
),
fragment: token.fragment.to_owned(),
})
}
impl<'bump> Parser<'bump> {
pub(crate) fn parse_create(&mut self) -> Result<AstCreate<'bump>> {
let token = self.consume_keyword(Create)?;
let or_replace = if (self.consume_if(TokenKind::Operator(Or))?).is_some() {
self.consume_keyword(Replace)?;
true
} else {
false
};
if or_replace {
let fragment = self.current()?.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "FLOW after CREATE OR REPLACE".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"FLOW after CREATE OR REPLACE",
fragment.text()
),
fragment,
}));
}
if (self.consume_if(TokenKind::Keyword(Remote))?).is_some() {
self.consume_keyword(Namespace)?;
return self.parse_remote_namespace(token);
}
if (self.consume_if(TokenKind::Keyword(Namespace))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::Namespace);
}
return self.parse_namespace(token);
}
if (self.consume_if(TokenKind::Keyword(View))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::View);
}
return self.parse_deferred_view(token);
}
if (self.consume_if(TokenKind::Keyword(Deferred))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Ringbuffer))?).is_some() {
self.consume_keyword(View)?;
return self.parse_deferred_view_with_storage(token, ViewStorageKindHint::RingBuffer);
}
if (self.consume_if(TokenKind::Keyword(Series))?).is_some() {
self.consume_keyword(View)?;
return self.parse_deferred_view_with_storage(token, ViewStorageKindHint::Series);
}
if (self.consume_if(TokenKind::Keyword(View))?).is_some() {
return self.parse_deferred_view(token);
}
unimplemented!()
}
if (self.consume_if(TokenKind::Keyword(Transactional))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Ringbuffer))?).is_some() {
self.consume_keyword(View)?;
return self
.parse_transactional_view_with_storage(token, ViewStorageKindHint::RingBuffer);
}
if (self.consume_if(TokenKind::Keyword(Series))?).is_some() {
self.consume_keyword(View)?;
return self.parse_transactional_view_with_storage(token, ViewStorageKindHint::Series);
}
if (self.consume_if(TokenKind::Keyword(View))?).is_some() {
return self.parse_transactional_view(token);
}
unimplemented!()
}
if (self.consume_if(TokenKind::Keyword(Table))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::Table);
}
return self.parse_table(token);
}
if (self.consume_if(TokenKind::Keyword(Ringbuffer))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::RingBuffer);
}
return self.parse_ringbuffer(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Queue))?).is_some() {
return self.parse_queue(token);
}
if (self.consume_if(TokenKind::Keyword(Dictionary))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::Dictionary);
}
return self.parse_dictionary(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Enum))?).is_some() {
return self.parse_enum(token);
}
if (self.consume_if(TokenKind::Keyword(Series))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::Series);
}
return self.parse_series(token);
}
if (self.consume_if(TokenKind::Keyword(Subscription))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::Subscription);
}
return self.parse_subscription(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Primary))?).is_some() {
self.consume_keyword(Keyword::Key)?;
return self.parse_create_primary_key(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Column))?).is_some() {
self.consume_keyword(Keyword::Property)?;
return self.parse_create_column_property(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Procedure))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Policy))?).is_some() {
return self.parse_create_policy(token, AstPolicyTargetType::Procedure);
}
return self.parse_procedure(token);
}
if (self.consume_if(TokenKind::Keyword(Test))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Procedure))?).is_some() {
return self.parse_test_procedure(token);
}
return self.parse_test(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Event))?).is_some() {
return self.parse_event(token);
}
if (self.consume_if(TokenKind::Keyword(Tag))?).is_some() {
return self.parse_tag(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Handler))?).is_some() {
return self.parse_handler(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Authentication))?).is_some() {
return self.parse_create_authentication(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::User))?).is_some() {
if (self.consume_if(TokenKind::Keyword(Keyword::Attribute))?).is_some() {
return self.parse_create_identity_attribute(token);
}
return self.parse_create_identity(token, IdentityKind::User);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Service))?).is_some() {
return self.parse_create_identity(token, IdentityKind::Service);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Role))?).is_some() {
return self.parse_create_role(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Session))?).is_some() {
self.consume_keyword(Keyword::Policy)?;
return self.parse_create_policy(token, AstPolicyTargetType::Session);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Feature))?).is_some() {
self.consume_keyword(Keyword::Policy)?;
return self.parse_create_policy(token, AstPolicyTargetType::Feature);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Function))?).is_some() {
self.consume_keyword(Keyword::Policy)?;
return self.parse_create_policy(token, AstPolicyTargetType::Function);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Migration))?).is_some() {
return self.parse_migration(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Source))?).is_some() {
return self.parse_source(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Sink))?).is_some() {
return self.parse_sink(token);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Http))?).is_some() {
self.consume_keyword(Keyword::Binding)?;
return self.parse_create_binding(token, AstBindingProtocolKind::Http);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Grpc))?).is_some() {
self.consume_keyword(Keyword::Binding)?;
return self.parse_create_binding(token, AstBindingProtocolKind::Grpc);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Ws))?).is_some() {
self.consume_keyword(Keyword::Binding)?;
return self.parse_create_binding(token, AstBindingProtocolKind::Ws);
}
if (self.consume_if(TokenKind::Keyword(Keyword::Relationship))?).is_some() {
return self.parse_create_relationship(token);
}
if self.peek_is_index_creation()? {
return self.parse_create_index(token);
}
unimplemented!();
}
fn parse_procedure(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let proc_ident = MaybeQualifiedProcedureIdentifier::new(name).with_namespace(namespace);
let params = if self.current()?.is_operator(Operator::OpenCurly) {
self.parse_procedure_params()?
} else {
Vec::new()
};
self.consume_operator(Operator::As)?;
self.consume_operator(Operator::OpenCurly)?;
let body_start_pos = self.position;
let mut body = Vec::new();
loop {
self.skip_new_line()?;
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
body.push(node);
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
self.consume_if(TokenKind::Separator(Separator::Semicolon))?;
}
let body_end_pos = self.position;
let body_source = if body_start_pos < body_end_pos {
let start = self.tokens[body_start_pos].fragment.offset();
let end = self.tokens[body_end_pos - 1].fragment.source_end();
self.source[start..end].trim().to_string()
} else {
String::new()
};
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Procedure(AstCreateProcedure {
token,
name: proc_ident,
params,
body,
body_source,
is_test: false,
}))
}
fn parse_test_procedure(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let proc_ident = MaybeQualifiedProcedureIdentifier::new(name).with_namespace(namespace);
let params = if self.current()?.is_operator(Operator::OpenCurly) {
self.parse_procedure_params()?
} else {
Vec::new()
};
self.consume_operator(Operator::As)?;
self.consume_operator(Operator::OpenCurly)?;
let body_start_pos = self.position;
let mut body = Vec::new();
loop {
self.skip_new_line()?;
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
body.push(node);
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
self.consume_if(TokenKind::Separator(Separator::Semicolon))?;
}
let body_end_pos = self.position;
let body_source = if body_start_pos < body_end_pos {
let start = self.tokens[body_start_pos].fragment.offset();
let end = self.tokens[body_end_pos - 1].fragment.source_end();
self.source[start..end].trim().to_string()
} else {
String::new()
};
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Procedure(AstCreateProcedure {
token,
name: proc_ident,
params,
body,
body_source,
is_test: true,
}))
}
fn parse_test(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let test_ident = MaybeQualifiedTestIdentifier::new(name).with_namespace(namespace);
let cases = if !self.is_eof() && self.current()?.kind == TokenKind::Operator(Operator::OpenBracket) {
let params_start_pos = self.position;
let mut depth = 0u32;
loop {
let tok = &self.tokens[self.position];
match tok.kind {
TokenKind::Operator(Operator::OpenBracket) => depth += 1,
TokenKind::Operator(Operator::CloseBracket) => {
depth -= 1;
if depth == 0 {
self.position += 1;
break;
}
}
_ => {}
}
self.position += 1;
}
let params_end_pos = self.position;
let start = self.tokens[params_start_pos].fragment.offset();
let end = self.tokens[params_end_pos - 1].fragment.offset()
+ self.tokens[params_end_pos - 1].fragment.text().len();
Some(self.source[start..end].to_string())
} else {
None
};
self.consume_operator(Operator::OpenCurly)?;
let body_start_pos = self.position;
let mut body = Vec::new();
loop {
self.skip_new_line()?;
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
body.push(node);
if !self.is_eof() && self.current()?.is_operator(Operator::Pipe) {
self.advance()?;
continue;
}
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
self.consume_if(TokenKind::Separator(Separator::Semicolon))?;
}
let body_end_pos = self.position;
let body_source = if body_start_pos < body_end_pos {
let start = self.tokens[body_start_pos].fragment.offset();
let end = self.tokens[body_end_pos - 1].fragment.source_end();
self.source[start..end].trim().to_string()
} else {
String::new()
};
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Test(AstCreateTest {
token,
name: test_ident,
cases,
body,
body_source,
}))
}
fn parse_procedure_params(&mut self) -> Result<Vec<AstProcedureParam<'bump>>> {
let mut params = Vec::new();
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let name = self.parse_identifier_with_hyphens()?.into_fragment();
self.consume_operator(Colon)?;
let param_type = self.parse_type()?;
params.push(AstProcedureParam {
name,
param_type,
});
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(params)
}
fn parse_namespace(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut if_not_exists = if (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
true
} else {
false
};
let segments = self.parse_double_colon_separated_identifiers()?;
if !if_not_exists && (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
if_not_exists = true;
}
let namespace = MaybeQualifiedNamespaceIdentifier::new(
segments.into_iter().map(|s| s.into_fragment()).collect(),
);
Ok(AstCreate::Namespace(AstCreateNamespace {
token,
namespace,
if_not_exists,
}))
}
fn parse_remote_namespace(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut if_not_exists = if (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
true
} else {
false
};
let segments = self.parse_double_colon_separated_identifiers()?;
if !if_not_exists && (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
if_not_exists = true;
}
let namespace = MaybeQualifiedNamespaceIdentifier::new(
segments.into_iter().map(|s| s.into_fragment()).collect(),
);
self.consume_keyword(Keyword::With)?;
self.consume_operator(Operator::OpenCurly)?;
let mut grpc = None;
let mut remote_token = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"grpc" => {
let value = self.consume_literal(Literal::Text)?;
grpc = Some(value.fragment);
}
"token" => {
let value = self.consume_literal(Literal::Text)?;
remote_token = Some(value.fragment);
}
_other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'grpc' or 'token'".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'grpc' or 'token'",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
let grpc = grpc.ok_or_else(|| {
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'grpc' key in WITH block".to_string(),
},
message: "CREATE REMOTE NAMESPACE requires 'grpc' in WITH block".to_string(),
fragment: token.fragment.to_owned(),
})
})?;
Ok(AstCreate::RemoteNamespace(AstCreateRemoteNamespace {
token,
namespace,
if_not_exists,
grpc,
token_value: remote_token,
}))
}
fn parse_series(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let series = MaybeQualifiedSeriesIdentifier::new(name).with_namespace(namespace);
let mut tag = None;
let mut key_field = None;
let mut precision = None;
let mut settings = None;
let mut partition_by: Vec<String> = Vec::new();
let mut time_declaration = AstTimeDeclaration::default();
if self.consume_if(TokenKind::Keyword(Keyword::With))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let with_key = {
let current = self.current()?;
match current.kind {
TokenKind::Identifier => self.consume_identifier()?,
TokenKind::Keyword(Keyword::Tag) => {
let token = self.advance()?;
Token {
kind: TokenKind::Identifier,
..token
}
}
_ => {
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'key', 'tag', 'precision', or 'ttl'"
.to_string(),
},
message: format!(
"expected 'key', 'tag', 'precision', or 'row', found `{}`",
current.fragment.text()
),
fragment: current.fragment.to_owned(),
}));
}
}
};
self.consume_operator(Operator::Colon)?;
match with_key.fragment.text() {
"key" => {
let key_token = self.consume_identifier()?;
key_field = Some(key_token.fragment);
}
"tag" => {
let mut tag_segments =
self.parse_double_colon_separated_identifiers()?;
let tag_name = tag_segments.pop().unwrap().into_fragment();
let tag_namespace: Vec<_> =
tag_segments.into_iter().map(|s| s.into_fragment()).collect();
tag = Some(MaybeQualifiedSumTypeIdentifier::new(tag_name)
.with_namespace(tag_namespace));
}
"precision" => {
let prec_token = self.consume_identifier()?;
precision = Some(match prec_token.fragment.text() {
"second" => AstTimestampPrecision::Second,
"millisecond" => AstTimestampPrecision::Millisecond,
"microsecond" => AstTimestampPrecision::Microsecond,
"nanosecond" => AstTimestampPrecision::Nanosecond,
_ => {
let fragment = prec_token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'second', 'millisecond', 'microsecond', or 'nanosecond'"
.to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'second', 'millisecond', 'microsecond', or 'nanosecond'",
fragment.text()
),
fragment,
}));
}
});
}
"row" => {
settings = Some(self.parse_row_config()?);
}
"partition" => {
partition_by = self.parse_partition_config()?;
}
"time" => {
time_declaration = self.parse_time_declaration()?;
}
_other => {
let fragment = with_key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'key', 'tag', 'precision', 'partition', 'row', or 'time'"
.to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'key', 'tag', 'precision', 'partition', 'row', or 'time'",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
}
let key_fragment = match key_field {
Some(k) => Some(k),
None => {
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "WITH block containing 'key' field".to_string(),
},
message: "CREATE SERIES requires a WITH block with a 'key' field specifying the ordering column".to_string(),
fragment: token.fragment.to_owned(),
}));
}
};
Ok(AstCreate::Series(AstCreateSeries {
token,
series,
columns,
tag,
key: key_fragment,
precision,
partition_by,
settings,
time_declaration,
}))
}
fn parse_subscription(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let columns = if self.current()?.is_operator(Operator::As) {
Vec::new()
} else if self.current()?.is_operator(Operator::OpenCurly) {
self.parse_columns()?
} else if self.current()?.is_keyword(Keyword::With) {
Vec::new()
} else {
let fragment = self.current()?.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'{', 'WITH', or 'AS'".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'{', 'WITH', or 'AS'",
fragment.text()
),
fragment,
}));
};
let mut hydration = AstHydrationConfig::default();
let mut throttle: Option<Duration> = None;
let mut linger: Option<Duration> = None;
if self.consume_if(TokenKind::Keyword(Keyword::With))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"hydration" => {
hydration = self.parse_hydration_with_value()?;
}
"throttle" => {
throttle = Some(self.parse_throttle_duration()?);
}
"linger" => {
linger = Some(self.parse_linger_duration()?);
}
_ => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'hydration', 'throttle', or 'linger'"
.to_string(),
},
message: format!(
"expected 'hydration', 'throttle', or 'linger', found `{}`",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
}
let as_clause = if self.consume_if(TokenKind::Operator(Operator::As))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
let mut query_nodes = Vec::new();
let mut has_pipes = false;
loop {
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
query_nodes.push(node);
if !self.is_eof() && self.current()?.is_operator(Operator::Pipe) {
self.advance()?;
has_pipes = true;
} else {
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
}
}
self.consume_operator(Operator::CloseCurly)?;
Some(AstStatement {
nodes: query_nodes,
has_pipes,
is_output: false,
rql: "",
})
} else {
None
};
if columns.is_empty() && as_clause.is_none() {
let fragment = self
.current()
.ok()
.map(|t| t.fragment.to_owned())
.unwrap_or_else(|| Fragment::internal("end of input"));
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "AS clause (shape-less CREATE SUBSCRIPTION requires AS clause)"
.to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"AS clause (shape-less CREATE SUBSCRIPTION requires AS clause)",
fragment.text()
),
fragment,
}));
}
Ok(AstCreate::Subscription(AstCreateSubscription {
token,
columns,
as_clause,
hydration,
throttle,
linger,
}))
}
fn parse_deferred_view(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let view = MaybeQualifiedDeferredViewIdentifier::new(name).with_namespace(namespace);
let clause = if !self.is_eof() && self.current()?.is_keyword(Keyword::With) {
self.advance()?;
self.parse_view_with_clause()?
} else {
AstViewWithClause::default()
};
let AstViewWithClause {
settings,
partition_by,
} = clause;
let as_clause = if self.consume_if(TokenKind::Operator(Operator::As))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
let mut query_nodes = Vec::new();
let mut has_pipes = false;
loop {
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
query_nodes.push(node);
if !self.is_eof() && self.current()?.is_operator(Operator::Pipe) {
self.advance()?;
has_pipes = true;
} else {
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
}
}
self.consume_operator(Operator::CloseCurly)?;
Some(AstStatement {
nodes: query_nodes,
has_pipes,
is_output: false,
rql: "",
})
} else {
None
};
Ok(AstCreate::DeferredView(AstCreateDeferredView {
token,
view,
columns,
as_clause,
storage_kind: AstViewStorageKind::Table {
partition_by,
},
settings,
}))
}
fn parse_deferred_view_with_storage(
&mut self,
token: Token<'bump>,
hint: ViewStorageKindHint,
) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let view = MaybeQualifiedDeferredViewIdentifier::new(name).with_namespace(namespace);
let (storage_kind, settings) = self.parse_view_storage_with_clause(hint)?;
let as_clause = self.parse_view_as_clause()?;
Ok(AstCreate::DeferredView(AstCreateDeferredView {
token,
view,
columns,
as_clause,
storage_kind,
settings,
}))
}
fn parse_transactional_view(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let view = MaybeQualifiedTransactionalViewIdentifier::new(name).with_namespace(namespace);
let AstViewWithClause {
settings,
partition_by,
} = if !self.is_eof() && self.current()?.is_keyword(Keyword::With) {
self.advance()?;
self.parse_view_with_clause()?
} else {
AstViewWithClause::default()
};
let as_clause = if self.consume_if(TokenKind::Operator(Operator::As))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
let mut query_nodes = Vec::new();
let mut has_pipes = false;
loop {
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
query_nodes.push(node);
if !self.is_eof() && self.current()?.is_operator(Operator::Pipe) {
self.advance()?;
has_pipes = true;
} else {
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
}
}
self.consume_operator(Operator::CloseCurly)?;
Some(AstStatement {
nodes: query_nodes,
has_pipes,
is_output: false,
rql: "",
})
} else {
None
};
Ok(AstCreate::TransactionalView(AstCreateTransactionalView {
token,
view,
columns,
as_clause,
storage_kind: AstViewStorageKind::Table {
partition_by,
},
settings,
}))
}
fn parse_transactional_view_with_storage(
&mut self,
token: Token<'bump>,
hint: ViewStorageKindHint,
) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let view = MaybeQualifiedTransactionalViewIdentifier::new(name).with_namespace(namespace);
let (storage_kind, settings) = self.parse_view_storage_with_clause(hint)?;
let as_clause = self.parse_view_as_clause()?;
Ok(AstCreate::TransactionalView(AstCreateTransactionalView {
token,
view,
columns,
as_clause,
storage_kind,
settings,
}))
}
fn parse_table(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let if_not_exists = if (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
true
} else {
false
};
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let table = MaybeQualifiedTableIdentifier::new(name).with_namespace(namespace);
let mut settings = None;
let mut partition_by: Vec<String> = Vec::new();
let mut time_declaration = AstTimeDeclaration::default();
if self.consume_if(TokenKind::Keyword(Keyword::With))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"row" => {
settings = Some(self.parse_row_config()?);
}
"partition" => {
partition_by = self.parse_partition_config()?;
}
"time" => {
time_declaration = self.parse_time_declaration()?;
}
_other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'row', 'partition', or 'time'".to_string(),
},
message: format!(
"expected 'row', 'partition', or 'time', found `{}`",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
}
Ok(AstCreate::Table(AstCreateTable {
token,
table,
if_not_exists,
columns,
partition_by,
settings,
time_declaration,
}))
}
fn parse_ringbuffer(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
self.consume_keyword(Keyword::With)?;
self.consume_operator(Operator::OpenCurly)?;
let mut capacity: Option<u64> = None;
let mut partition_by: Vec<String> = Vec::new();
let mut settings = None;
let mut time_declaration = AstTimeDeclaration::default();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = {
let current = self.current()?;
match current.kind {
TokenKind::Identifier => self.consume_identifier()?,
TokenKind::Keyword(Keyword::Tag) => {
let token = self.advance()?;
Token {
kind: TokenKind::Identifier,
..token
}
}
_ => {
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'capacity', 'partition', or 'row'"
.to_string(),
},
message: format!(
"expected 'capacity', 'partition', or 'row', found `{}`",
current.fragment.text()
),
fragment: current.fragment.to_owned(),
}));
}
}
};
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"capacity" => {
let capacity_token = self.consume(TokenKind::Literal(Literal::Number))?;
capacity =
Some(capacity_token.fragment.text().parse::<u64>().map_err(|_| {
let fragment = capacity_token.fragment.to_owned();
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "valid capacity number".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"valid capacity number",
fragment.text()
),
fragment,
})
})?);
}
"partition" => {
partition_by = self.parse_partition_config()?;
}
"row" => {
settings = Some(self.parse_row_config()?);
}
"time" => {
time_declaration = self.parse_time_declaration()?;
}
_other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'capacity', 'partition', 'row', or 'time'"
.to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'capacity', 'partition', 'row', or 'time'",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
let capacity = capacity.ok_or_else(|| {
let fragment = self
.current()
.ok()
.map(|t| t.fragment.to_owned())
.unwrap_or_else(|| Fragment::internal("end of input"));
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'capacity' is required for RINGBUFFER".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'capacity' is required for RINGBUFFER",
fragment.text()
),
fragment,
})
})?;
let ringbuffer = MaybeQualifiedRingBufferIdentifier::new(name).with_namespace(namespace);
Ok(AstCreate::RingBuffer(AstCreateRingBuffer {
token,
ringbuffer,
columns,
capacity,
partition_by,
settings,
time_declaration,
}))
}
fn parse_queue(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let columns = self.parse_columns()?;
let mut dispatch = None;
let mut deduplicate = None;
let mut retention = None;
let mut retry = None;
let mut time_declaration = AstTimeDeclaration::default();
if (self.consume_if(TokenKind::Keyword(Keyword::With))?).is_some() {
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_queue_option_key(QUEUE_OPTION_KEYS)?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"fifo" => {
if dispatch.is_some() {
return Err(duplicate_queue_dispatch(&key));
}
dispatch = Some(AstQueueDispatch::Fifo(self.parse_queue_fifo(key)?));
}
"deduplicate" => {
deduplicate = Some(self.parse_queue_deduplicate(key)?);
}
"retention" => {
retention = Some(self.parse_queue_retention()?);
}
"retry" => {
retry = Some(self.parse_queue_retry()?);
}
"time" => {
time_declaration = self.parse_time_declaration()?;
}
_other => return Err(unexpected_queue_option(&key, QUEUE_OPTION_KEYS)),
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
}
let queue = MaybeQualifiedQueueIdentifier::new(name).with_namespace(namespace);
let Some(dispatch) = dispatch else {
return Err(missing_queue_dispatch(&token));
};
Ok(AstCreate::Queue(AstCreateQueue {
token,
queue,
columns,
dispatch,
deduplicate,
retention,
retry,
time_declaration,
}))
}
fn parse_queue_fifo(&mut self, token: Token<'bump>) -> Result<AstQueueFifo<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut partitions = None;
let mut ordered_by = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_queue_option_key(QUEUE_FIFO_KEYS)?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"partitions" => {
partitions = Some(self.consume(TokenKind::Literal(Literal::Number))?);
}
"ordered_by" => {
ordered_by = Some(self.consume_identifier()?);
}
_other => return Err(unexpected_queue_option(&key, QUEUE_FIFO_KEYS)),
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstQueueFifo {
token,
partitions,
ordered_by,
})
}
fn parse_queue_deduplicate(&mut self, token: Token<'bump>) -> Result<AstQueueDeduplicate<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut by = Vec::new();
let mut ttl = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_queue_option_key(QUEUE_DEDUPLICATE_KEYS)?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"by" => {
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
by.push(self.consume_identifier()?);
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
break;
}
}
self.skip_new_line()?;
self.consume_operator(Operator::CloseCurly)?;
}
"ttl" => {
ttl = Some(self.consume_duration_or_forever()?);
}
_other => return Err(unexpected_queue_option(&key, QUEUE_DEDUPLICATE_KEYS)),
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstQueueDeduplicate {
token,
by,
ttl,
})
}
fn parse_queue_retention(&mut self) -> Result<AstQueueRetention<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut done = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_queue_option_key(QUEUE_RETENTION_KEYS)?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"done" => {
done = Some(self.consume_duration("'done'")?);
}
_other => return Err(unexpected_queue_option(&key, QUEUE_RETENTION_KEYS)),
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstQueueRetention {
done,
})
}
fn parse_queue_retry(&mut self) -> Result<AstQueueRetry<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut attempts = None;
let mut backoff = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_queue_option_key(QUEUE_RETRY_KEYS)?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"attempts" => {
attempts = Some(self.consume(TokenKind::Literal(Literal::Number))?);
}
"backoff" => {
backoff = Some(self.consume_duration("'backoff'")?);
}
_other => return Err(unexpected_queue_option(&key, QUEUE_RETRY_KEYS)),
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstQueueRetry {
attempts,
backoff,
})
}
fn consume_queue_option_key(&mut self, expected: &str) -> Result<Token<'bump>> {
let current = self.current()?;
match current.kind {
TokenKind::Identifier => self.consume_identifier(),
_ => Err(unexpected_queue_option(¤t, expected)),
}
}
fn parse_time_declaration(&mut self) -> Result<AstTimeDeclaration<'bump>> {
if let Some(token) = self.consume_if(TokenKind::Literal(Literal::None))? {
return Ok(AstTimeDeclaration::None(token));
}
let keyword = self.consume_identifier()?;
match keyword.fragment.text().to_ascii_lowercase().as_str() {
"processing" => Ok(AstTimeDeclaration::Processing(keyword)),
"event" => {
self.consume_operator(Operator::OpenParen)?;
let column = self.consume_identifier()?;
self.consume_operator(Operator::CloseParen)?;
Ok(AstTimeDeclaration::Event {
keyword,
column,
})
}
_ => {
let fragment = keyword.fragment.to_owned();
Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: TIME_VALUES.to_string(),
},
message: format!("expected {}, found `{}`", TIME_VALUES, fragment.text()),
fragment,
}))
}
}
}
fn parse_partition_config(&mut self) -> Result<Vec<String>> {
let mut partition_by: Vec<String> = Vec::new();
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let inner_key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match inner_key.fragment.text() {
"by" => {
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let col = self.consume_identifier()?;
partition_by.push(col.fragment.text().to_string());
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
}
_other => {
let fragment = inner_key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'by'".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"'by'",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(partition_by)
}
fn parse_primary_keyinition(&mut self) -> Result<AstPrimaryKey<'bump>> {
let mut columns = Vec::new();
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let column = self.parse_column_identifier()?;
let sort_direction = if self.current()?.is_operator(Operator::Colon) {
self.consume_operator(Operator::Colon)?;
if self.current()?.is_keyword(Keyword::Asc) {
self.consume_keyword(Keyword::Asc)?;
SortDirection::Asc
} else if self.current()?.is_keyword(Keyword::Desc) {
self.consume_keyword(Keyword::Desc)?;
SortDirection::Desc
} else {
SortDirection::Desc
}
} else {
SortDirection::Desc
};
columns.push(AstIndexColumn {
column,
order: Some(sort_direction),
});
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Separator::Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
if columns.is_empty() {
let fragment = self.current()?.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "at least one column in primary key".to_string(),
},
message: format!(
"Unexpected token: expected {}, got {}",
"at least one column in primary key",
fragment.text()
),
fragment,
}));
}
Ok(AstPrimaryKey {
columns,
})
}
fn parse_create_primary_key(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
self.consume_keyword(Keyword::On)?;
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let table = MaybeQualifiedTableIdentifier::new(name).with_namespace(namespace);
let pk_def = self.parse_primary_keyinition()?;
Ok(AstCreate::PrimaryKey(AstCreatePrimaryKey {
token,
table,
columns: pk_def.columns,
}))
}
fn parse_create_column_property(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
self.consume_keyword(Keyword::On)?;
let column = self.parse_column_identifier()?;
self.consume_operator(Operator::OpenCurly)?;
let mut properties = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let kind_token = self.consume_identifier()?;
let kind = match kind_token.fragment.text() {
"saturation" => AstColumnPropertyKind::Saturation,
"default" => AstColumnPropertyKind::Default,
_ => {
let fragment = kind_token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::InvalidPolicy,
message: format!("Invalid property token: {}", fragment.text()),
fragment,
}));
}
};
self.consume_operator(Operator::Colon)?;
let value = BumpBox::new_in(self.parse_node(Precedence::None)?, self.bump());
properties.push(AstColumnPropertyEntry {
kind,
value,
});
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::ColumnProperty(AstCreateColumnProperty {
token,
column,
properties,
}))
}
fn parse_dictionary(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let if_not_exists = if (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
true
} else {
false
};
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let dictionary = if namespace.is_empty() {
MaybeQualifiedDictionaryIdentifier::new(name)
} else {
MaybeQualifiedDictionaryIdentifier::new(name).with_namespace(namespace)
};
self.consume_keyword(For)?;
let value_type = self.parse_type()?;
self.consume_operator(Operator::As)?;
let id_type = self.parse_type()?;
Ok(AstCreate::Dictionary(AstCreateDictionary {
token,
if_not_exists,
dictionary,
value_type,
id_type,
}))
}
fn parse_enum(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let if_not_exists = if (self.consume_if(TokenKind::Keyword(If))?).is_some() {
self.consume_operator(Not)?;
self.consume_keyword(Exists)?;
true
} else {
false
};
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name_frag = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let sumtype_ident = if namespace.is_empty() {
MaybeQualifiedSumTypeIdentifier::new(name_frag)
} else {
MaybeQualifiedSumTypeIdentifier::new(name_frag).with_namespace(namespace)
};
self.consume_operator(Operator::OpenCurly)?;
let mut variants = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let variant_name = self.parse_identifier_with_hyphens()?.into_fragment();
let columns = if !self.is_eof() && self.current()?.is_operator(Operator::OpenCurly) {
self.parse_columns()?
} else {
Vec::new()
};
variants.push(AstVariant {
name: variant_name,
columns,
});
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
self.skip_new_line()?;
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Enum(AstCreateSumType {
token,
if_not_exists,
name: sumtype_ident,
variants,
}))
}
pub(crate) fn parse_type(&mut self) -> Result<AstType<'bump>> {
let ty_token = self.consume_identifier()?;
if ty_token.fragment.text().eq_ignore_ascii_case("option") {
self.consume_operator(Operator::OpenParen)?;
let inner = self.parse_type()?;
self.consume_operator(Operator::CloseParen)?;
return Ok(AstType::Optional(Box::new(inner)));
}
if !self.is_eof() && self.current()?.is_operator(Operator::DoubleColon) {
self.consume_operator(Operator::DoubleColon)?;
let name_token = self.consume_identifier()?;
return Ok(AstType::Qualified {
namespace: ty_token.fragment,
name: name_token.fragment,
});
}
if !self.is_eof() && self.current()?.is_operator(Operator::OpenParen) {
self.consume_operator(Operator::OpenParen)?;
let mut params = Vec::new();
params.push(self.parse_literal_number()?);
while self.consume_if(TokenKind::Separator(Comma))?.is_some() {
params.push(self.parse_literal_number()?);
}
self.consume_operator(Operator::CloseParen)?;
Ok(AstType::Constrained {
name: ty_token.fragment,
params,
})
} else {
Ok(AstType::Unconstrained(ty_token.fragment))
}
}
fn parse_columns(&mut self) -> Result<Vec<AstColumnToCreate<'bump>>> {
let mut result = Vec::new();
self.consume_operator(Operator::OpenCurly)?;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
result.push(self.parse_column()?);
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
break;
};
}
self.consume_operator(Operator::CloseCurly)?;
Ok(result)
}
pub(crate) fn parse_column(&mut self) -> Result<AstColumnToCreate<'bump>> {
let name_identifier = self.parse_identifier_with_hyphens()?;
self.consume_operator(Colon)?;
let ty_token = self.consume_identifier()?;
let name = name_identifier.into_fragment();
let ty = if ty_token.fragment.text().eq_ignore_ascii_case("option") {
self.consume_operator(Operator::OpenParen)?;
let inner = self.parse_type()?;
self.consume_operator(Operator::CloseParen)?;
AstType::Optional(Box::new(inner))
} else if !self.is_eof() && self.current()?.is_operator(Operator::DoubleColon) {
self.consume_operator(Operator::DoubleColon)?;
let name_token = self.consume_identifier()?;
AstType::Qualified {
namespace: ty_token.fragment,
name: name_token.fragment,
}
} else if !self.is_eof() && self.current()?.is_operator(Operator::OpenParen) {
self.consume_operator(Operator::OpenParen)?;
let mut params = Vec::new();
params.push(self.parse_literal_number()?);
while self.consume_if(TokenKind::Separator(Comma))?.is_some() {
params.push(self.parse_literal_number()?);
}
self.consume_operator(Operator::CloseParen)?;
AstType::Constrained {
name: ty_token.fragment,
params,
}
} else {
AstType::Unconstrained(ty_token.fragment)
};
let properties = if !self.is_eof() && self.current()?.is_keyword(Keyword::With) {
self.parse_column_properties()?
} else {
vec![]
};
Ok(AstColumnToCreate {
name,
ty,
properties,
})
}
fn parse_column_properties(&mut self) -> Result<Vec<AstColumnProperty<'bump>>> {
self.consume_keyword(Keyword::With)?;
self.consume_operator(Operator::OpenCurly)?;
let mut properties = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key_token = {
let current = self.current()?;
match current.kind {
TokenKind::Identifier => self.consume_identifier()?,
TokenKind::Keyword(Keyword::Dictionary) => {
let token = self.advance()?;
Token {
kind: TokenKind::Identifier,
..token
}
}
_ => {
let fragment = current.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::InvalidColumnProperty,
message: format!(
"Invalid column property: {}",
fragment.text()
),
fragment,
}));
}
}
};
let key = key_token.fragment.text();
let property = match key {
"auto_increment" => AstColumnProperty::AutoIncrement,
"dictionary" => {
self.consume_operator(Colon)?;
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> =
segments.into_iter().map(|s| s.into_fragment()).collect();
let dict_ident = if namespace.is_empty() {
MaybeQualifiedDictionaryIdentifier::new(name)
} else {
MaybeQualifiedDictionaryIdentifier::new(name).with_namespace(namespace)
};
AstColumnProperty::Dictionary(dict_ident)
}
"saturation" => {
self.consume_operator(Colon)?;
let value = BumpBox::new_in(self.parse_node(Precedence::None)?, self.bump());
AstColumnProperty::Saturation(value)
}
"default" => {
self.consume_operator(Colon)?;
let value = BumpBox::new_in(self.parse_node(Precedence::None)?, self.bump());
AstColumnProperty::Default(value)
}
_ => {
let fragment = key_token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::InvalidColumnProperty,
message: format!("Invalid column property: {}", fragment.text()),
fragment,
}));
}
};
properties.push(property);
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(properties)
}
pub(crate) fn parse_event(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name_frag = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let sumtype_ident = if namespace.is_empty() {
MaybeQualifiedSumTypeIdentifier::new(name_frag)
} else {
MaybeQualifiedSumTypeIdentifier::new(name_frag).with_namespace(namespace)
};
self.consume_operator(Operator::OpenCurly)?;
let mut variants = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let variant_name = self.parse_identifier_with_hyphens()?.into_fragment();
let columns = if !self.is_eof() && self.current()?.is_operator(Operator::OpenCurly) {
self.parse_columns()?
} else {
Vec::new()
};
variants.push(AstVariant {
name: variant_name,
columns,
});
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
self.skip_new_line()?;
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Event(AstCreateEvent {
token,
name: sumtype_ident,
variants,
}))
}
pub(crate) fn parse_tag(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name_frag = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let sumtype_ident = if namespace.is_empty() {
MaybeQualifiedSumTypeIdentifier::new(name_frag)
} else {
MaybeQualifiedSumTypeIdentifier::new(name_frag).with_namespace(namespace)
};
self.consume_operator(Operator::OpenCurly)?;
let mut variants = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let variant_name = self.parse_identifier_with_hyphens()?.into_fragment();
let columns = if !self.is_eof() && self.current()?.is_operator(Operator::OpenCurly) {
self.parse_columns()?
} else {
Vec::new()
};
variants.push(AstVariant {
name: variant_name,
columns,
});
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
self.skip_new_line()?;
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Tag(AstCreateTag {
token,
name: sumtype_ident,
variants,
}))
}
pub(crate) fn parse_handler(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let name_frag = segments.pop().unwrap().into_fragment();
let namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let handler_name = MaybeQualifiedTableIdentifier::new(name_frag).with_namespace(namespace);
self.consume_keyword(Keyword::On)?;
let mut event_segments = self.parse_double_colon_separated_identifiers()?;
let on_variant = event_segments.pop().unwrap().into_fragment();
let event_name_frag = event_segments.pop().unwrap().into_fragment();
let event_namespace: Vec<_> = event_segments.into_iter().map(|s| s.into_fragment()).collect();
let on_event = if event_namespace.is_empty() {
MaybeQualifiedSumTypeIdentifier::new(event_name_frag)
} else {
MaybeQualifiedSumTypeIdentifier::new(event_name_frag).with_namespace(event_namespace)
};
self.consume_operator(Operator::OpenCurly)?;
let body_start_pos = self.position;
let mut body = Vec::new();
loop {
self.skip_new_line()?;
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
body.push(node);
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
self.consume_if(TokenKind::Separator(Separator::Semicolon))?;
}
let body_end_pos = self.position;
let body_source = if body_start_pos < body_end_pos {
let start = self.tokens[body_start_pos].fragment.offset();
let end = self.tokens[body_end_pos - 1].fragment.source_end();
self.source[start..end].trim().to_string()
} else {
String::new()
};
self.consume_operator(Operator::CloseCurly)?;
Ok(AstCreate::Handler(AstCreateHandler {
token,
name: handler_name,
on_event,
on_variant,
body,
body_source,
}))
}
fn parse_migration(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let name = match &self.current()?.kind {
TokenKind::Literal(Literal::Text) => {
let text = self.current()?.fragment.text().to_string();
self.advance()?;
text
}
_ => {
let fragment = self.current()?.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "migration name as string literal".to_string(),
},
message: format!(
"Expected migration name as string literal, got {}",
fragment.text()
),
fragment,
}));
}
};
self.consume_operator(Operator::OpenCurly)?;
let body_start_pos = self.position;
let mut depth = 1u32;
while depth > 0 {
if self.is_eof() {
let fragment = self.tokens[body_start_pos].fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "closing '}'".to_string(),
},
message: "Unexpected end of input while parsing migration body".to_string(),
fragment,
}));
}
match self.current()?.kind {
TokenKind::Operator(Operator::OpenCurly) => {
depth += 1;
self.advance()?;
}
TokenKind::Operator(Operator::CloseCurly) => {
depth -= 1;
if depth > 0 {
self.advance()?;
}
}
_ => {
self.advance()?;
}
}
}
let body_end_pos = self.position;
let body_source = if body_start_pos < body_end_pos {
let start = self.tokens[body_start_pos].fragment.offset();
let end = self.tokens[body_end_pos - 1].fragment.source_end();
self.source[start..end].trim().to_string()
} else {
String::new()
};
self.consume_operator(Operator::CloseCurly)?;
let rollback_body_source = if (self.consume_if(TokenKind::Keyword(Keyword::Rollback))?).is_some() {
self.consume_operator(Operator::OpenCurly)?;
let rb_start_pos = self.position;
let mut depth = 1u32;
while depth > 0 {
if self.is_eof() {
let fragment = self.tokens[rb_start_pos].fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "closing '}'".to_string(),
},
message: "Unexpected end of input while parsing rollback body"
.to_string(),
fragment,
}));
}
match self.current()?.kind {
TokenKind::Operator(Operator::OpenCurly) => {
depth += 1;
self.advance()?;
}
TokenKind::Operator(Operator::CloseCurly) => {
depth -= 1;
if depth > 0 {
self.advance()?;
}
}
_ => {
self.advance()?;
}
}
}
let rb_end_pos = self.position;
let rb_source = if rb_start_pos < rb_end_pos {
let start = self.tokens[rb_start_pos].fragment.offset();
let end = self.tokens[rb_end_pos - 1].fragment.offset()
+ self.tokens[rb_end_pos - 1].fragment.text().len();
self.source[start..end].trim().to_string()
} else {
String::new()
};
self.consume_operator(Operator::CloseCurly)?;
Some(rb_source)
} else {
None
};
Ok(AstCreate::Migration(AstCreateMigration {
token,
name,
body_source,
rollback_body_source,
}))
}
fn parse_view_as_clause(&mut self) -> Result<Option<AstStatement<'bump>>> {
if self.consume_if(TokenKind::Operator(Operator::As))?.is_some() {
self.consume_operator(Operator::OpenCurly)?;
let mut query_nodes = Vec::new();
let mut has_pipes = false;
loop {
if self.is_eof() || self.current()?.kind == TokenKind::Operator(Operator::CloseCurly) {
break;
}
let node = self.parse_node(Precedence::None)?;
query_nodes.push(node);
if !self.is_eof() && self.current()?.is_operator(Operator::Pipe) {
self.advance()?;
has_pipes = true;
} else {
self.consume_if(TokenKind::Separator(Separator::NewLine))?;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(Some(AstStatement {
nodes: query_nodes,
has_pipes,
is_output: false,
rql: "",
}))
} else {
Ok(None)
}
}
fn parse_view_storage_with_clause(
&mut self,
hint: ViewStorageKindHint,
) -> Result<(AstViewStorageKind, Option<AstRowSettings<'bump>>)> {
self.consume_keyword(Keyword::With)?;
self.consume_operator(Operator::OpenCurly)?;
let mut settings = None;
match hint {
ViewStorageKindHint::RingBuffer => {
let mut capacity: Option<u64> = None;
let mut partition_by: Vec<String> = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"capacity" => {
let token =
self.consume(TokenKind::Literal(Literal::Number))?;
capacity =
Some(token.fragment.text().parse::<u64>().map_err(
|_| {
let fragment =
token.fragment.to_owned();
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "valid capacity number".to_string(),
},
message: format!("expected valid capacity number, got {}", fragment.text()),
fragment,
})
},
)?);
}
"partition" => {
partition_by = self.parse_partition_config()?;
}
"row" => {
settings = Some(self.parse_row_config()?);
}
other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'capacity', 'partition', or 'row'"
.to_string(),
},
message: format!(
"unexpected key '{}' in WITH clause",
other
),
fragment,
}));
}
}
self.consume_if(TokenKind::Separator(Comma))?;
}
self.consume_operator(Operator::CloseCurly)?;
let capacity = capacity.ok_or_else(|| {
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "capacity".to_string(),
},
message: "ringbuffer view requires 'capacity' in WITH clause"
.to_string(),
fragment: Fragment::internal(""),
})
})?;
Ok((
AstViewStorageKind::RingBuffer {
capacity,
partition_by,
},
settings,
))
}
ViewStorageKindHint::Series => {
let mut key_column: Option<String> = None;
let mut precision: Option<AstTimestampPrecision> = None;
let mut partition_by: Vec<String> = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"key" => {
let token = self.consume_identifier()?;
key_column = Some(token.fragment.text().to_string());
}
"precision" => {
let token = self.consume_identifier()?;
precision = Some(match token.fragment.text() {
"second" => AstTimestampPrecision::Second,
"millisecond" => AstTimestampPrecision::Millisecond,
"microsecond" => AstTimestampPrecision::Microsecond,
"nanosecond" => AstTimestampPrecision::Nanosecond,
_ => {
let fragment = token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "second, millisecond, microsecond, or nanosecond".to_string(),
},
message: format!("unexpected precision '{}'", fragment.text()),
fragment,
}));
}
});
}
"partition" => {
partition_by = self.parse_partition_config()?;
}
"row" => {
settings = Some(self.parse_row_config()?);
}
other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'key', 'precision', 'partition', or 'row'"
.to_string(),
},
message: format!(
"unexpected key '{}' in WITH clause",
other
),
fragment,
}));
}
}
self.consume_if(TokenKind::Separator(Comma))?;
}
self.consume_operator(Operator::CloseCurly)?;
let key_column = key_column.unwrap_or_default();
Ok((
AstViewStorageKind::Series {
key_column,
precision,
partition_by,
},
settings,
))
}
}
}
fn parse_view_with_clause(&mut self) -> Result<AstViewWithClause<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut clause = AstViewWithClause::default();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"partition" => {
clause.partition_by = self.parse_partition_config()?;
}
"row" => {
clause.settings = Some(self.parse_row_config()?);
}
other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'partition' or 'row'".to_string(),
},
message: format!("unexpected key '{}' in WITH clause", other),
fragment,
}));
}
}
self.consume_if(TokenKind::Separator(Comma))?;
}
self.consume_operator(Operator::CloseCurly)?;
Ok(clause)
}
fn parse_throttle_duration(&mut self) -> Result<Duration> {
let token = self.consume_duration("'throttle'")?;
compile_duration(&token, DurationBound::AllowZero, "'throttle'")
}
fn parse_linger_duration(&mut self) -> Result<Duration> {
let token = self.consume_duration("'linger'")?;
compile_duration(&token, DurationBound::AllowZero, "'linger'")
}
fn parse_hydration_with_value(&mut self) -> Result<AstHydrationConfig> {
self.consume_operator(Operator::OpenCurly)?;
let mut enabled: Option<bool> = None;
let mut max_rows: Option<u64> = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"enabled" => {
let current = self.current()?;
match current.kind {
TokenKind::Literal(Literal::True) => {
self.advance()?;
enabled = Some(true);
}
TokenKind::Literal(Literal::False) => {
self.advance()?;
enabled = Some(false);
}
_ => {
let fragment = current.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "boolean literal".to_string(),
},
message: format!(
"expected boolean literal for hydration.enabled, found `{}`",
fragment.text()
),
fragment,
}));
}
}
}
"max_rows" => {
let token = self.consume(TokenKind::Literal(Literal::Number))?;
let text = token.fragment.text();
if text.starts_with('-') {
let fragment = token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "non-negative integer".to_string(),
},
message: format!(
"hydration.max_rows must be a positive integer, found `{}`",
fragment.text()
),
fragment,
}));
}
let parsed = text.parse::<u64>().map_err(|_| {
let fragment = token.fragment.to_owned();
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "valid u64 integer".to_string(),
},
message: format!(
"hydration.max_rows must be a valid u64, found `{}`",
fragment.text()
),
fragment,
})
})?;
if parsed == 0 {
let fragment = token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "positive integer".to_string(),
},
message:
"hydration.max_rows must be greater than zero (use enabled: false to disable)"
.to_string(),
fragment,
}));
}
max_rows = Some(parsed);
}
_other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'enabled' or 'max_rows'".to_string(),
},
message: format!(
"expected 'enabled' or 'max_rows' in hydration config, found `{}`",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
Ok(AstHydrationConfig {
enabled: enabled.unwrap_or(true),
max_rows,
})
}
fn parse_row_config(&mut self) -> Result<AstRowSettings<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut duration: Option<Token<'bump>> = None;
let mut anchor: Option<Token<'bump>> = None;
let mut persistent: Option<AstPersistent<'bump>> = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"ttl" => {
duration = Some(self.consume_duration(ROW_CONFIG_KEYS)?);
}
"on" => {
anchor = Some(self.consume_identifier()?);
}
"persistent" => {
let token = self.consume_boolean("persistent")?;
persistent = Some(AstPersistent {
value: token.kind == TokenKind::Literal(Literal::True),
token,
});
}
_other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: ROW_CONFIG_KEYS.to_string(),
},
message: format!(
"expected {} in row config, found `{}`",
ROW_CONFIG_KEYS,
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
if duration.is_none() && persistent.is_none() {
let fragment = self
.current()
.ok()
.map(|t| t.fragment.to_owned())
.unwrap_or_else(|| Fragment::internal("end of input"));
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'ttl' is required in row config".to_string(),
},
message: "'ttl' is required in row config".to_string(),
fragment,
}));
}
if duration.is_none()
&& let Some(orphan) = anchor
{
let fragment = orphan.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "a 'ttl' alongside it".to_string(),
},
message: "'on' qualifies a 'ttl'; add 'ttl' to the row config".to_string(),
fragment,
}));
}
if let Some(p) = &persistent
&& !p.value && duration.is_none()
{
let fragment = p.token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "a 'ttl' alongside 'persistent: false'".to_string(),
},
message: "a non-persistent object requires a row ttl; add 'ttl' to the row config"
.to_string(),
fragment,
}));
}
Ok(AstRowSettings {
ttl: duration.map(|duration| AstTtl {
duration,
anchor,
}),
persistent,
})
}
fn consume_duration(&mut self, context: &str) -> Result<Token<'bump>> {
let current = self.current()?;
if current.kind == TokenKind::Literal(Literal::Duration) {
return self.advance();
}
let fragment = current.fragment.to_owned();
Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "duration literal".to_string(),
},
message: format!(
"expected a bare duration literal such as `2h` for {}, found `{}`",
context,
fragment.text()
),
fragment,
}))
}
fn consume_duration_or_forever(&mut self) -> Result<Token<'bump>> {
let current = self.current()?;
if current.kind == TokenKind::Identifier && current.fragment.text() == FOREVER {
return self.advance();
}
self.consume_duration("'ttl', or the keyword `forever`,")
}
fn consume_boolean(&mut self, key: &str) -> Result<Token<'bump>> {
let current = self.current()?;
match current.kind {
TokenKind::Literal(Literal::True) | TokenKind::Literal(Literal::False) => self.advance(),
_ => {
let fragment = current.fragment.to_owned();
Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'true' or 'false'".to_string(),
},
message: format!(
"expected boolean literal for '{}', found `{}`",
key,
fragment.text()
),
fragment,
}))
}
}
}
fn parse_join_retention(&mut self) -> Result<AstJoinRetention<'bump>> {
self.consume_operator(Operator::OpenCurly)?;
let mut left: Option<Token<'bump>> = None;
let mut right: Option<Token<'bump>> = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"left" => {
if left.is_some() {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "single 'left' entry".to_string(),
},
message: "'left' specified more than once in join retention"
.to_string(),
fragment,
}));
}
left = Some(self.consume_duration("the 'left' side")?);
}
"right" => {
if right.is_some() {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "single 'right' entry".to_string(),
},
message: "'right' specified more than once in join retention"
.to_string(),
fragment,
}));
}
right = Some(self.consume_duration("the 'right' side")?);
}
other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'left' or 'right'".to_string(),
},
message: format!(
"unexpected key '{}' in join retention; expected 'left' or 'right'",
other
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
if left.is_none() && right.is_none() {
let fragment = self
.current()
.ok()
.map(|t| t.fragment.to_owned())
.unwrap_or_else(|| Fragment::internal("end of input"));
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "at least one of 'left' or 'right'".to_string(),
},
message: "join retention must specify at least one side ('left' or 'right')"
.to_string(),
fragment,
}));
}
Ok(AstJoinRetention {
left,
right,
})
}
pub(crate) fn reject_with_clause(&mut self, kind: OperationKind) -> Result<()> {
if self.is_eof() || !self.current()?.is_keyword(Keyword::With) {
return Ok(());
}
Err(RqlError::OperatorNoWithClause {
kind,
fragment: self.current()?.fragment.to_owned(),
}
.into())
}
pub(crate) fn parse_with_clause_for_join(
&mut self,
) -> Result<(Option<AstJoinRetention<'bump>>, bool, Option<AstJoinPick<'bump>>)> {
if self.is_eof() || !self.current()?.is_keyword(Keyword::With) {
return Ok((None, false, None));
}
self.advance()?;
self.consume_operator(Operator::OpenCurly)?;
let mut retention: Option<AstJoinRetention<'bump>> = None;
let mut snapshot: Option<bool> = None;
let mut pick: Option<AstJoinPick<'bump>> = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key.fragment.text() {
"retention" => {
retention = Some(self.parse_join_retention()?);
}
"snapshot" => {
snapshot = Some(self.parse_join_bool("snapshot")?);
}
"latest" | "earliest" => {
let default_direction = match key.fragment.text() {
"latest" => SortDirection::Desc,
_ => SortDirection::Asc,
};
if pick.is_some() {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "one of 'latest' or 'earliest'".to_string(),
},
message:
"'latest' and 'earliest' cannot both be set on one join"
.to_string(),
fragment,
}));
}
pick = self.parse_join_pick(key, default_direction)?;
}
other => {
let fragment = key.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: JOIN_WITH_KEYS.to_string(),
},
message: format!("unexpected key '{}' in join WITH clause", other),
fragment,
}));
}
}
self.consume_if(TokenKind::Separator(Comma))?;
}
self.consume_operator(Operator::CloseCurly)?;
Ok((retention, snapshot.unwrap_or(false), pick))
}
fn parse_join_bool(&mut self, key: &str) -> Result<bool> {
let value = self.advance()?;
match value.kind {
TokenKind::Literal(Literal::True) => Ok(true),
TokenKind::Literal(Literal::False) => Ok(false),
_ => {
let fragment = value.fragment.to_owned();
Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "boolean literal 'true' or 'false'".to_string(),
},
message: format!(
"expected boolean literal for '{}', got '{}'",
key,
value.fragment.text()
),
fragment,
}))
}
}
}
fn parse_join_pick(
&mut self,
key: Token<'bump>,
default_direction: SortDirection,
) -> Result<Option<AstJoinPick<'bump>>> {
if !self.current()?.is_operator(Operator::OpenCurly) {
return Ok(match self.parse_join_bool(key.fragment.text())? {
true => Some(AstJoinPick {
token: key,
default_direction,
columns: Vec::new(),
directions: Vec::new(),
}),
false => None,
});
}
self.advance()?;
let mut columns = Vec::new();
let mut directions = Vec::new();
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
columns.push(self.parse_column_identifier()?);
self.skip_new_line()?;
if self.current()?.is_operator(Operator::Colon) {
self.advance()?;
let token = self.advance()?;
directions.push(Some(if token.is_keyword(Keyword::Asc) {
SortDirection::Asc
} else if token.is_keyword(Keyword::Desc) {
SortDirection::Desc
} else {
let fragment = token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'asc' or 'desc'".to_string(),
},
message: format!(
"expected 'asc' or 'desc' after the column, got '{}'",
token.fragment.text()
),
fragment,
}));
}));
} else {
directions.push(None);
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_none() {
break;
}
}
self.skip_new_line()?;
self.consume_operator(Operator::CloseCurly)?;
Ok(Some(AstJoinPick {
token: key,
default_direction,
columns,
directions,
}))
}
fn parse_create_relationship(&mut self, token: Token<'bump>) -> Result<AstCreate<'bump>> {
let name_token = self.consume_identifier()?;
let name = name_token.fragment;
self.consume_keyword(Keyword::On)?;
let (source, source_column) = self.parse_relationship_table_with_column(1)?;
let junction = if self.consume_if(TokenKind::Keyword(Keyword::Through))?.is_some() {
let (jtable, jcols) = self.parse_relationship_table_with_columns(2)?;
let mut iter = jcols.into_iter();
let first = iter.next().unwrap();
let second = iter.next().unwrap();
Some(AstRelationshipJunction {
table: jtable,
source_column: first,
target_column: second,
})
} else {
None
};
self.consume_keyword(Keyword::References)?;
let (target, target_column) = self.parse_relationship_table_with_column(1)?;
self.consume_keyword(Keyword::With)?;
self.consume_operator(Operator::OpenCurly)?;
let mut cardinality: Option<AstRelationshipCardinality> = None;
loop {
self.skip_new_line()?;
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
let key_token = self.consume_identifier()?;
self.consume_operator(Operator::Colon)?;
match key_token.fragment.text() {
"cardinality" => {
let value = self.consume_literal(Literal::Text)?;
let text = value.fragment.text();
let parsed = AstRelationshipCardinality::parse(text);
match parsed {
Some(c) => cardinality = Some(c),
None => {
let fragment = value.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'1:1', 'N:1', '1:N', or 'N:M'"
.to_string(),
},
message: format!(
"unknown relationship cardinality `{}`: expected '1:1', 'N:1', '1:N', or 'N:M'",
text
),
fragment,
}));
}
}
}
_other => {
let fragment = key_token.fragment.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "'cardinality'".to_string(),
},
message: format!(
"unexpected key `{}`: relationship WITH block accepts only 'cardinality'",
fragment.text()
),
fragment,
}));
}
}
self.skip_new_line()?;
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
if self.current()?.is_operator(Operator::CloseCurly) {
break;
}
}
self.consume_operator(Operator::CloseCurly)?;
let cardinality = cardinality.ok_or_else(|| {
Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "WITH { cardinality: '<value>' }".to_string(),
},
message: "CREATE RELATIONSHIP requires 'cardinality' in WITH block".to_string(),
fragment: token.fragment.to_owned(),
})
})?;
if cardinality.requires_junction() && junction.is_none() {
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "THROUGH <table>(<src>, <tgt>)".to_string(),
},
message: "N:M relationship requires a THROUGH junction table".to_string(),
fragment: name.to_owned(),
}));
}
if !cardinality.requires_junction() && junction.is_some() {
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: "no THROUGH clause".to_string(),
},
message: "THROUGH junction is only allowed for N:M cardinality".to_string(),
fragment: name.to_owned(),
}));
}
Ok(AstCreate::Relationship(AstCreateRelationship {
token,
name,
source,
source_column,
target,
target_column,
junction,
cardinality,
}))
}
fn parse_relationship_table_with_column(
&mut self,
_count: usize,
) -> Result<(MaybeQualifiedTableIdentifier<'bump>, BumpFragment<'bump>)> {
let (table, mut cols) = self.parse_relationship_table_with_columns(1)?;
Ok((table, cols.pop().unwrap()))
}
fn parse_relationship_table_with_columns(
&mut self,
count: usize,
) -> Result<(MaybeQualifiedTableIdentifier<'bump>, Vec<BumpFragment<'bump>>)> {
let mut segments = self.parse_double_colon_separated_identifiers()?;
let table_name = segments.pop().unwrap().into_fragment();
let table_namespace: Vec<_> = segments.into_iter().map(|s| s.into_fragment()).collect();
let table = MaybeQualifiedTableIdentifier::new(table_name).with_namespace(table_namespace);
self.consume_operator(Operator::OpenParen)?;
let mut cols = Vec::with_capacity(count);
loop {
let col_token = self.consume_identifier()?;
cols.push(col_token.fragment);
if self.consume_if(TokenKind::Separator(Comma))?.is_some() {
continue;
}
break;
}
self.consume_operator(Operator::CloseParen)?;
if cols.len() != count {
let fragment = table.name.to_owned();
return Err(Error::from(TypeError::Ast {
kind: AstErrorKind::UnexpectedToken {
expected: format!("{} column name(s)", count),
},
message: format!(
"expected {} column name(s) in parentheses, found {}",
count,
cols.len()
),
fragment,
}));
}
Ok((table, cols))
}
}
enum ViewStorageKindHint {
RingBuffer,
Series,
}
#[cfg(test)]
pub mod tests {
use bumpalo::Bump;
use reifydb_value::value::duration::Duration;
use crate::{
ast::{
ast::{
Ast, AstColumnProperty, AstCreate, AstCreateDeferredView, AstCreateDictionary,
AstCreateNamespace, AstCreateQueue, AstCreateRingBuffer, AstCreateSeries,
AstCreateSubscription, AstCreateSumType, AstCreateTable, AstCreateTransactionalView,
AstHydrationConfig, AstQueueDispatch, AstType,
},
parse::Parser,
},
token::tokenize,
};
#[test]
fn test_create_namespace() {
let bump = Bump::new();
let source = "CREATE NAMESPACE REIFYDB";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "REIFYDB");
assert!(!if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_with_hyphen() {
let bump = Bump::new();
let source = "CREATE NAMESPACE my-namespace";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my-namespace");
assert!(!if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_if_not_exists() {
let bump = Bump::new();
let source = "CREATE NAMESPACE IF NOT EXISTS my_namespace";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my_namespace");
assert!(if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_if_not_exists_with_hyphen() {
let bump = Bump::new();
let source = "CREATE NAMESPACE IF NOT EXISTS my-test-namespace";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my-test-namespace");
assert!(if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_if_not_exists_with_backtick() {
let bump = Bump::new();
let source = "CREATE NAMESPACE IF NOT EXISTS `my-namespace`";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my-namespace");
assert!(if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_name_if_not_exists() {
let bump = Bump::new();
let source = "CREATE NAMESPACE my_namespace IF NOT EXISTS";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my_namespace");
assert!(if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_name_if_not_exists_with_hyphen() {
let bump = Bump::new();
let source = "CREATE NAMESPACE my-test-namespace IF NOT EXISTS";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my-test-namespace");
assert!(if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_name_if_not_exists_with_backtick() {
let bump = Bump::new();
let source = "CREATE NAMESPACE `my-namespace` IF NOT EXISTS";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
if_not_exists,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my-namespace");
assert!(if_not_exists);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_table_with_hyphen() {
let bump = Bump::new();
let source = "CREATE TABLE my-shape::my-table { id: Int4 }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Table(AstCreateTable {
table,
..
}) => {
assert_eq!(table.namespace[0].text(), "my-shape");
assert_eq!(table.name.text(), "my-table");
}
_ => unreachable!(),
}
}
#[test]
fn test_create_ringbuffer_with_hyphen() {
let bump = Bump::new();
let source = "CREATE RINGBUFFER my-ns::my-buffer { id: Int4 } WITH { capacity: 100 }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::RingBuffer(AstCreateRingBuffer {
ringbuffer,
capacity,
..
}) => {
assert_eq!(ringbuffer.namespace[0].text(), "my-ns");
assert_eq!(ringbuffer.name.text(), "my-buffer");
assert_eq!(*capacity, 100);
}
_ => unreachable!(),
}
}
#[test]
fn test_create_dictionary_with_hyphen() {
let bump = Bump::new();
let source = "CREATE DICTIONARY my-dict FOR Text AS Int4";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Dictionary(AstCreateDictionary {
dictionary,
..
}) => {
assert_eq!(dictionary.name.text(), "my-dict");
}
_ => unreachable!(),
}
}
#[test]
fn test_create_table_with_hyphenated_columns() {
let bump = Bump::new();
let source = "CREATE TABLE test::user-data { user-id: Int4, user-name: Text }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Table(AstCreateTable {
table,
columns,
..
}) => {
assert_eq!(table.name.text(), "user-data");
assert_eq!(columns.len(), 2);
assert_eq!(columns[0].name.text(), "user-id");
assert_eq!(columns[1].name.text(), "user-name");
}
_ => unreachable!(),
}
}
#[test]
fn test_create_namespace_with_backtick() {
let bump = Bump::new();
let source = "CREATE NAMESPACE `my-namespace`";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Namespace(AstCreateNamespace {
namespace,
..
}) => {
assert_eq!(namespace.segments[0].text(), "my-namespace");
}
_ => unreachable!(),
}
}
#[test]
fn test_create_series() {
let bump = Bump::new();
let source = r#"
create series test::metrics{ts: datetime, value: Int2} WITH { key: ts }
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Series(AstCreateSeries {
series,
columns,
key,
..
}) => {
assert_eq!(series.namespace[0].text(), "test");
assert_eq!(series.name.text(), "metrics");
assert_eq!(columns.len(), 2);
assert_eq!(columns[0].name.text(), "ts");
match &columns[0].ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "datetime")
}
_ => panic!("Expected simple type"),
}
assert_eq!(columns[1].name.text(), "value");
match &columns[1].ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "Int2")
}
_ => panic!("Expected simple type"),
}
assert!(columns[1].properties.is_empty());
assert!(key.is_some());
assert_eq!(key.as_ref().unwrap().text(), "ts");
}
_ => unreachable!(),
}
}
#[test]
fn test_create_table() {
let bump = Bump::new();
let source = r#"
create table test::users{id: int2, name: text, is_premium: bool}
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Table(AstCreateTable {
table,
columns,
..
}) => {
assert_eq!(table.namespace[0].text(), "test");
assert_eq!(table.name.text(), "users");
assert_eq!(columns.len(), 3);
{
let col = &columns[0];
assert_eq!(col.name.text(), "id");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "int2")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
{
let col = &columns[1];
assert_eq!(col.name.text(), "name");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "text")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
{
let col = &columns[2];
assert_eq!(col.name.text(), "is_premium");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "bool")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
}
_ => unreachable!(),
}
}
#[test]
fn test_create_table_with_auto_increment() {
let bump = Bump::new();
let source = r#"
create table test::users { id: int4 with { auto_increment }, name: utf8 }
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Table(AstCreateTable {
table,
columns,
..
}) => {
assert_eq!(table.namespace[0].text(), "test");
assert_eq!(table.name.text(), "users");
assert_eq!(columns.len(), 2);
{
let col = &columns[0];
assert_eq!(col.name.text(), "id");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "int4")
}
_ => panic!("Expected simple type"),
}
assert_eq!(col.properties.len(), 1);
assert!(matches!(col.properties[0], AstColumnProperty::AutoIncrement));
}
{
let col = &columns[1];
assert_eq!(col.name.text(), "name");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "utf8")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
}
_ => unreachable!(),
}
}
#[test]
fn test_create_deferred_view() {
let bump = Bump::new();
let source = r#"
create deferred view test::views{field: int2}
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::DeferredView(AstCreateDeferredView {
view,
columns,
..
}) => {
assert_eq!(view.namespace[0].text(), "test");
assert_eq!(view.name.text(), "views");
assert_eq!(columns.len(), 1);
let col = &columns[0];
assert_eq!(col.name.text(), "field");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "int2")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
_ => unreachable!(),
}
}
#[test]
fn test_bare_create_view_defaults_to_deferred() {
let bump = Bump::new();
let source = r#"
create view test::views{field: int2}
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::DeferredView(AstCreateDeferredView {
view,
..
}) => {
assert_eq!(view.namespace[0].text(), "test");
assert_eq!(view.name.text(), "views");
}
other => panic!("bare create view must parse as a deferred view, got {other:?}"),
}
}
#[test]
fn test_create_transactional_view() {
let bump = Bump::new();
let source = r#"
create transactional view test::myview{id: int4, name: utf8}
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::TransactionalView(AstCreateTransactionalView {
view,
columns,
..
}) => {
assert_eq!(view.namespace[0].text(), "test");
assert_eq!(view.name.text(), "myview");
assert_eq!(columns.len(), 2);
{
let col = &columns[0];
assert_eq!(col.name.text(), "id");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "int4")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
{
let col = &columns[1];
assert_eq!(col.name.text(), "name");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "utf8")
}
_ => panic!("Expected simple type"),
}
assert!(col.properties.is_empty());
}
}
_ => unreachable!(),
}
}
#[test]
fn test_create_ringbuffer() {
let bump = Bump::new();
let source = r#"
create ringbuffer test::events { id: int4, data: utf8 } with { capacity: 10 }
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::RingBuffer(AstCreateRingBuffer {
ringbuffer,
columns,
capacity,
..
}) => {
assert_eq!(ringbuffer.namespace[0].text(), "test");
assert_eq!(ringbuffer.name.text(), "events");
assert_eq!(*capacity, 10);
assert_eq!(columns.len(), 2);
{
let col = &columns[0];
assert_eq!(col.name.text(), "id");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "int4")
}
_ => panic!("Expected simple type"),
}
}
{
let col = &columns[1];
assert_eq!(col.name.text(), "data");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "utf8")
}
_ => panic!("Expected simple type"),
}
}
}
_ => unreachable!("Expected ring buffer create"),
}
}
#[test]
fn test_create_transactional_view_with_query() {
let bump = Bump::new();
let source = r#"
create transactional view test::myview{id: int4, name: utf8} as {
from test::users
where age > 18
}
"#;
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::TransactionalView(AstCreateTransactionalView {
view,
columns,
as_clause,
..
}) => {
assert_eq!(view.namespace[0].text(), "test");
assert_eq!(view.name.text(), "myview");
assert_eq!(columns.len(), 2);
assert!(as_clause.is_some());
if let Some(as_statement) = as_clause {
assert!(as_statement.len() > 0);
}
}
_ => unreachable!(),
}
}
#[test]
fn test_create_dictionary_basic() {
let bump = Bump::new();
let source = "CREATE DICTIONARY token_mints FOR Utf8 AS Uint2";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Dictionary(dict) => {
assert!(dict.dictionary.namespace.is_empty());
assert_eq!(dict.dictionary.name.text(), "token_mints");
match &dict.value_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Utf8"),
_ => panic!("Expected unconstrained type"),
}
match &dict.id_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Uint2"),
_ => panic!("Expected unconstrained type"),
}
}
_ => unreachable!("Expected Dictionary create"),
}
}
#[test]
fn test_create_dictionary_qualified() {
let bump = Bump::new();
let source = "CREATE DICTIONARY analytics::token_mints FOR Utf8 AS Uint4";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Dictionary(dict) => {
assert_eq!(dict.dictionary.namespace[0].text(), "analytics");
assert_eq!(dict.dictionary.name.text(), "token_mints");
match &dict.value_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Utf8"),
_ => panic!("Expected unconstrained type"),
}
match &dict.id_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Uint4"),
_ => panic!("Expected unconstrained type"),
}
}
_ => unreachable!("Expected Dictionary create"),
}
}
#[test]
fn test_create_dictionary_blob_value() {
let bump = Bump::new();
let source = "CREATE DICTIONARY hashes FOR Blob AS Uint8";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Dictionary(dict) => {
assert_eq!(dict.dictionary.name.text(), "hashes");
match &dict.value_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Blob"),
_ => panic!("Expected unconstrained type"),
}
match &dict.id_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Uint8"),
_ => panic!("Expected unconstrained type"),
}
}
_ => unreachable!("Expected Dictionary create"),
}
}
#[test]
fn test_create_dictionary_if_not_exists() {
let bump = Bump::new();
let source = "CREATE DICTIONARY IF NOT EXISTS token_mints FOR Utf8 AS Uint4";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Dictionary(dict) => {
assert!(dict.if_not_exists);
assert!(dict.dictionary.namespace.is_empty());
assert_eq!(dict.dictionary.name.text(), "token_mints");
match &dict.value_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Utf8"),
_ => panic!("Expected unconstrained type"),
}
match &dict.id_type {
AstType::Unconstrained(ty) => assert_eq!(ty.text(), "Uint4"),
_ => panic!("Expected unconstrained type"),
}
}
_ => unreachable!("Expected Dictionary create"),
}
}
#[test]
fn test_create_enum_basic() {
let bump = Bump::new();
let source = "CREATE ENUM Status { Active, Inactive, Pending }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Enum(AstCreateSumType {
if_not_exists,
name,
variants,
..
}) => {
assert!(!if_not_exists);
assert!(name.namespace.is_empty());
assert_eq!(name.name.text(), "Status");
assert_eq!(variants.len(), 3);
assert_eq!(variants[0].name.text(), "Active");
assert_eq!(variants[1].name.text(), "Inactive");
assert_eq!(variants[2].name.text(), "Pending");
assert!(variants[0].columns.is_empty());
assert!(variants[1].columns.is_empty());
assert!(variants[2].columns.is_empty());
}
_ => unreachable!("Expected Enum create"),
}
}
#[test]
fn test_create_enum_with_fields() {
let bump = Bump::new();
let source =
"CREATE ENUM Object { Circle { radius: Float8 }, Rectangle { width: Float8, height: Float8 } }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Enum(AstCreateSumType {
name,
variants,
..
}) => {
assert_eq!(name.name.text(), "Object");
assert_eq!(variants.len(), 2);
assert_eq!(variants[0].name.text(), "Circle");
assert_eq!(variants[0].columns.len(), 1);
assert_eq!(variants[0].columns[0].name.text(), "radius");
assert_eq!(variants[1].name.text(), "Rectangle");
assert_eq!(variants[1].columns.len(), 2);
assert_eq!(variants[1].columns[0].name.text(), "width");
assert_eq!(variants[1].columns[1].name.text(), "height");
}
_ => unreachable!("Expected Enum create"),
}
}
#[test]
fn test_create_enum_qualified_name() {
let bump = Bump::new();
let source = "CREATE ENUM analytics::Status { Active, Inactive }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Enum(AstCreateSumType {
name,
variants,
..
}) => {
assert_eq!(name.namespace[0].text(), "analytics");
assert_eq!(name.name.text(), "Status");
assert_eq!(variants.len(), 2);
}
_ => unreachable!("Expected Enum create"),
}
}
#[test]
fn test_create_enum_if_not_exists() {
let bump = Bump::new();
let source = "CREATE ENUM IF NOT EXISTS Status { Active }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Enum(AstCreateSumType {
if_not_exists,
name,
variants,
..
}) => {
assert!(if_not_exists);
assert!(name.namespace.is_empty());
assert_eq!(name.name.text(), "Status");
assert_eq!(variants.len(), 1);
assert_eq!(variants[0].name.text(), "Active");
}
_ => unreachable!("Expected Enum create"),
}
}
#[test]
fn test_create_subscription_basic() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION { id: Int4, name: Utf8 }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
..
}) => {
assert_eq!(columns.len(), 2);
{
let col = &columns[0];
assert_eq!(col.name.text(), "id");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "Int4")
}
_ => panic!("Expected simple type"),
}
}
{
let col = &columns[1];
assert_eq!(col.name.text(), "name");
match &col.ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "Utf8")
}
_ => panic!("Expected simple type"),
}
}
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_single_column() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION { value: Float8 }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
..
}) => {
assert_eq!(columns.len(), 1);
assert_eq!(columns[0].name.text(), "value");
match &columns[0].ty {
AstType::Unconstrained(ident) => {
assert_eq!(ident.text(), "Float8")
}
_ => panic!("Expected simple type"),
}
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_with_simple_query() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION { id: Int4, name: Utf8 } AS { from test::products }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
as_clause,
..
}) => {
assert_eq!(columns.len(), 2);
assert_eq!(columns[0].name.text(), "id");
assert_eq!(columns[1].name.text(), "name");
assert!(as_clause.is_some(), "AS clause should be present");
let as_clause = as_clause.as_ref().unwrap();
assert_eq!(as_clause.nodes.len(), 1, "Should have one FROM node");
match &as_clause.nodes[0] {
Ast::From(_) => {}
_ => panic!("Expected FROM node in AS clause"),
}
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_with_piped_query() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION { id: Int4, price: Float8 } AS { from test::products | filter {price > 50} | filter {stock > 0} }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
as_clause,
..
}) => {
assert_eq!(columns.len(), 2);
assert!(as_clause.is_some(), "AS clause should be present");
let as_clause = as_clause.as_ref().unwrap();
assert!(as_clause.nodes.len() >= 1, "Should have at least FROM node");
match &as_clause.nodes[0] {
Ast::From(_) => {}
_ => panic!("Expected FROM node as first node in AS clause"),
}
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_without_as_clause() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION { value: Float8 }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
as_clause,
..
}) => {
assert!(as_clause.is_none(), "AS clause should not be present");
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_shapeless() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION AS { FROM demo::events }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
as_clause,
..
}) => {
assert_eq!(columns.len(), 0, "Object-less should have no columns");
assert!(as_clause.is_some(), "AS clause should be present");
let as_clause = as_clause.as_ref().unwrap();
assert_eq!(as_clause.nodes.len(), 1, "Should have one FROM node");
match &as_clause.nodes[0] {
Ast::From(_) => {}
_ => panic!("Expected FROM node in AS clause"),
}
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_shapeless_with_filter() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION AS { FROM demo::events | FILTER {id > 1 and id < 3} }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
as_clause,
..
}) => {
assert_eq!(columns.len(), 0, "Object-less should have no columns");
assert!(as_clause.is_some(), "AS clause should be present");
let as_clause = as_clause.as_ref().unwrap();
assert!(as_clause.nodes.len() >= 1, "Should have at least FROM node");
assert!(as_clause.has_pipes, "Should have pipes");
match &as_clause.nodes[0] {
Ast::From(_) => {}
_ => panic!("Expected FROM node as first node"),
}
}
_ => unreachable!("Expected Subscription create"),
}
}
#[test]
fn test_create_subscription_shapeless_missing_as_fails() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let result = parser.parse();
assert!(result.is_err(), "Object-less subscription without AS should fail");
}
#[test]
fn test_create_subscription_backward_compat_with_columns() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION { id: Int4 } AS { FROM demo::events }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Subscription(AstCreateSubscription {
columns,
as_clause,
..
}) => {
assert_eq!(columns.len(), 1, "Should have one column");
assert_eq!(columns[0].name.text(), "id");
assert!(as_clause.is_some(), "AS clause should be present");
}
_ => unreachable!("Expected Subscription create"),
}
}
fn parse_subscription_hydration(source: &str) -> AstHydrationConfig {
let bump = Bump::new();
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
let r = result.pop().unwrap();
match r.first_unchecked().as_create() {
AstCreate::Subscription(s) => s.hydration.clone(),
_ => unreachable!("expected subscription"),
}
}
#[test]
fn test_subscription_hydration_disabled() {
let cfg = parse_subscription_hydration(
"CREATE SUBSCRIPTION WITH { hydration: { enabled: false } } AS { FROM demo::events }",
);
assert!(!cfg.enabled);
assert_eq!(cfg.max_rows, None);
}
#[test]
fn test_subscription_hydration_with_max_rows() {
let cfg = parse_subscription_hydration(
"CREATE SUBSCRIPTION WITH { hydration: { enabled: true, max_rows: 1000 } } AS { FROM demo::events }",
);
assert!(cfg.enabled);
assert_eq!(cfg.max_rows, Some(1000));
}
#[test]
fn test_subscription_hydration_max_rows_only_defaults_enabled() {
let cfg = parse_subscription_hydration(
"CREATE SUBSCRIPTION WITH { hydration: { max_rows: 250 } } AS { FROM demo::events }",
);
assert!(cfg.enabled);
assert_eq!(cfg.max_rows, Some(250));
}
#[test]
fn test_subscription_hydration_empty_struct() {
let cfg = parse_subscription_hydration(
"CREATE SUBSCRIPTION WITH { hydration: { } } AS { FROM demo::events }",
);
assert!(cfg.enabled);
assert_eq!(cfg.max_rows, None);
}
#[test]
fn test_subscription_with_clause_omitted_defaults_to_enabled() {
let cfg = parse_subscription_hydration("CREATE SUBSCRIPTION AS { FROM demo::events }");
assert!(cfg.enabled);
assert_eq!(cfg.max_rows, None);
}
fn parse_subscription_linger(source: &str) -> Option<Duration> {
let bump = Bump::new();
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
let r = result.pop().unwrap();
match r.first_unchecked().as_create() {
AstCreate::Subscription(s) => s.linger,
_ => unreachable!("expected subscription"),
}
}
#[test]
fn test_subscription_linger_parsed() {
let linger = parse_subscription_linger(
"CREATE SUBSCRIPTION WITH { linger: 250ms } AS { FROM demo::events }",
);
assert_eq!(
linger,
Some(Duration::from_milliseconds(250).unwrap()),
"linger is a sibling WITH key parsed as a duration"
);
}
#[test]
fn test_subscription_linger_absent_is_none() {
let linger = parse_subscription_linger("CREATE SUBSCRIPTION AS { FROM demo::events }");
assert_eq!(
linger, None,
"an omitted linger is None, not a default - the default is applied downstream at clamp time"
);
}
#[test]
fn test_subscription_throttle_and_linger_coexist() {
let bump = Bump::new();
let source = "CREATE SUBSCRIPTION WITH { throttle: 1s, linger: 5ms } AS { FROM demo::events }";
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
let r = result.pop().unwrap();
match r.first_unchecked().as_create() {
AstCreate::Subscription(s) => {
assert_eq!(
s.throttle,
Some(Duration::from_seconds(1).unwrap()),
"throttle and linger are orthogonal WITH keys"
);
assert_eq!(s.linger, Some(Duration::from_milliseconds(5).unwrap()));
}
_ => unreachable!("expected subscription"),
}
}
fn parse_subscription_must_fail(source: &str) -> String {
let bump = Bump::new();
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let err = parser.parse().expect_err("expected parse error");
err.to_string()
}
#[test]
fn test_subscription_with_unknown_top_level_key_rejected() {
let msg = parse_subscription_must_fail(
"CREATE SUBSCRIPTION WITH { snapshot: { enabled: false } } AS { FROM demo::events }",
);
assert!(msg.contains("hydration"), "error should reference 'hydration', got: {}", msg);
}
#[test]
fn test_subscription_with_unknown_sub_key_rejected() {
let msg = parse_subscription_must_fail(
"CREATE SUBSCRIPTION WITH { hydration: { foo: 1 } } AS { FROM demo::events }",
);
assert!(msg.contains("'enabled' or 'max_rows'"), "error should reference legal sub-keys, got: {}", msg);
}
#[test]
fn test_subscription_with_non_struct_hydration_value_rejected() {
let msg = parse_subscription_must_fail(
"CREATE SUBSCRIPTION WITH { hydration: false } AS { FROM demo::events }",
);
assert!(!msg.is_empty(), "expected a parse error message");
}
#[test]
fn test_subscription_with_max_rows_zero_rejected() {
let msg = parse_subscription_must_fail(
"CREATE SUBSCRIPTION WITH { hydration: { max_rows: 0 } } AS { FROM demo::events }",
);
assert!(
msg.contains("greater than zero") || msg.contains("positive"),
"error should explain zero rejection, got: {}",
msg
);
}
#[test]
fn test_subscription_with_non_bool_enabled_rejected() {
let msg = parse_subscription_must_fail(
"CREATE SUBSCRIPTION WITH { hydration: { enabled: 1 } } AS { FROM demo::events }",
);
assert!(msg.contains("boolean"), "error should reference boolean expectation, got: {}", msg);
}
fn parse_queue_ast(source: &str) -> String {
let bump = Bump::new();
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let mut result = parser.parse().unwrap();
assert_eq!(result.len(), 1);
let result = result.pop().unwrap();
let create = result.first_unchecked().as_create();
match create {
AstCreate::Queue(AstCreateQueue {
queue,
columns,
dispatch,
retention,
retry,
..
}) => {
let AstQueueDispatch::Fifo(fifo) = dispatch;
format!(
"ns={:?} name={} columns={:?} partitions={:?} ordered_by={:?} retention={:?} retry={:?}",
queue.namespace.iter().map(|n| n.text().to_string()).collect::<Vec<_>>(),
queue.name.text(),
columns.iter().map(|c| c.name.text().to_string()).collect::<Vec<_>>(),
fifo.partitions.map(|t| t.fragment.text().to_string()),
fifo.ordered_by.map(|t| t.fragment.text().to_string()),
retention.as_ref().map(|r| r.done.map(|t| t.fragment.text().to_string())),
retry.as_ref().map(|r| (
r.attempts.map(|t| t.fragment.text().to_string()),
r.backoff.map(|t| t.fragment.text().to_string())
)),
)
}
_ => unreachable!("expected a CREATE QUEUE ast"),
}
}
fn parse_queue_must_fail(source: &str) -> String {
let bump = Bump::new();
let tokens = tokenize(&bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(&bump, source, tokens);
let err = parser.parse().expect_err("expected CREATE QUEUE to fail parsing");
format!("{:?}", err)
}
#[test]
fn test_create_queue_with_all_options() {
let rendered = parse_queue_ast(
r#"CREATE QUEUE ns::jobs { order_id: uuid7, kind: utf8 } WITH {
fifo: { partitions: 32, ordered_by: order_id },
retention: { done: 7d },
retry: { attempts: 5, backoff: 10s }
}"#,
);
assert_eq!(
rendered,
r#"ns=["ns"] name=jobs columns=["order_id", "kind"] partitions=Some("32") ordered_by=Some("order_id") retention=Some(Some("7d")) retry=Some((Some("5"), Some("10s")))"#
);
}
#[test]
fn test_create_queue_with_only_the_dispatch_block() {
let rendered = parse_queue_ast("CREATE QUEUE ns::jobs { order_id: uuid7 } WITH { fifo: {} }");
assert_eq!(
rendered,
r#"ns=["ns"] name=jobs columns=["order_id"] partitions=None ordered_by=None retention=None retry=None"#
);
}
#[test]
fn test_create_queue_without_a_dispatch_block_is_rejected() {
let bare = parse_queue_must_fail("CREATE QUEUE ns::jobs { id: int4 }");
assert!(bare.contains("dispatch block"), "got: {}", bare);
let other_options =
parse_queue_must_fail(r#"CREATE QUEUE ns::jobs { id: int4 } WITH { retention: { done: 7d } }"#);
assert!(other_options.contains("dispatch block"), "got: {}", other_options);
}
#[test]
fn test_create_queue_with_two_dispatch_blocks_is_rejected() {
let err = parse_queue_must_fail("CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, fifo: {} }");
assert!(err.contains("exactly one dispatch block"), "got: {}", err);
}
#[test]
fn test_create_queue_rejects_dispatch_options_at_the_top_level() {
let err = parse_queue_must_fail("CREATE QUEUE ns::jobs { id: int4 } WITH { partitions: 8 }");
assert!(err.contains("partitions"), "got: {}", err);
}
#[test]
fn test_create_queue_with_partitions_only() {
let rendered = parse_queue_ast("CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: { partitions: 1 } }");
assert!(rendered.contains(r#"partitions=Some("1")"#), "got: {}", rendered);
assert!(rendered.contains("ordered_by=None retention=None retry=None"), "got: {}", rendered);
}
#[test]
fn test_create_queue_with_ordered_by_only() {
let rendered = parse_queue_ast("CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: { ordered_by: id } }");
assert!(rendered.contains(r#"ordered_by=Some("id")"#), "got: {}", rendered);
assert!(rendered.contains("partitions=None"), "got: {}", rendered);
}
#[test]
fn test_create_queue_with_retention_only() {
let rendered = parse_queue_ast(
r#"CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, retention: { done: 1h } }"#,
);
assert!(rendered.contains(r#"retention=Some(Some("1h"))"#), "got: {}", rendered);
assert!(rendered.contains("retry=None"), "got: {}", rendered);
}
#[test]
fn test_create_queue_with_partial_retry() {
let rendered =
parse_queue_ast("CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, retry: { attempts: 2 } }");
assert!(rendered.contains(r#"retry=Some((Some("2"), None))"#), "got: {}", rendered);
}
#[test]
fn test_create_queue_unknown_option_rejected() {
let msg = parse_queue_must_fail("CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, partition: 4 }");
assert!(msg.contains("'fifo'"), "error should name the valid options, got: {}", msg);
}
#[test]
fn test_create_queue_unknown_retention_key_rejected() {
let msg = parse_queue_must_fail(
r#"CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, retention: { dead: 7d } }"#,
);
assert!(msg.contains("done"), "error should name the valid retention key, got: {}", msg);
}
#[test]
fn test_create_queue_unknown_retry_key_rejected() {
let msg = parse_queue_must_fail(
"CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, retry: { tries: 3 } }",
);
assert!(msg.contains("attempts"), "error should name the valid retry keys, got: {}", msg);
}
#[test]
fn test_create_queue_non_text_backoff_rejected() {
let msg = parse_queue_must_fail(
"CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: {}, retry: { backoff: 10 } }",
);
assert!(!msg.is_empty(), "expected a parse error message");
}
#[test]
fn test_create_queue_non_number_partitions_rejected() {
let msg = parse_queue_must_fail(
r#"CREATE QUEUE ns::jobs { id: int4 } WITH { fifo: { partitions: "32" } }"#,
);
assert!(!msg.is_empty(), "expected a parse error message");
}
}
#[cfg(test)]
mod time_declaration_tests {
use crate::{
Result,
ast::{
ast::{Ast, AstCreate, AstTimeDeclaration},
parse::Parser,
},
bump::Bump,
token::tokenize,
};
fn declared_in(bump: &Bump, source: &str) -> Result<String> {
let tokens = tokenize(bump, source).unwrap().into_iter().collect();
let mut parser = Parser::new(bump, source, tokens);
let mut result = parser.parse()?;
let statement = result.pop().unwrap();
let Ast::Create(create) = statement.first_unchecked() else {
panic!("expected a CREATE statement");
};
let decl = match &**create {
AstCreate::Table(t) => &t.time_declaration,
AstCreate::Series(s) => &s.time_declaration,
AstCreate::RingBuffer(r) => &r.time_declaration,
AstCreate::Queue(q) => &q.time_declaration,
other => panic!("statement carries no time declaration: {other:?}"),
};
Ok(match decl {
AstTimeDeclaration::Undeclared => "undeclared".to_string(),
AstTimeDeclaration::None(_) => "none".to_string(),
AstTimeDeclaration::Processing(_) => "processing".to_string(),
AstTimeDeclaration::Event {
column,
..
} => format!("event({})", column.fragment.text()),
})
}
fn declared(source: &str) -> String {
let bump = Bump::new();
declared_in(&bump, source).expect("statement must parse")
}
fn rejected(source: &str) {
let bump = Bump::new();
assert!(declared_in(&bump, source).is_err(), "must be rejected: {source}");
}
#[test]
fn the_event_keyword_carries_its_column_inside_the_parens() {
assert_eq!(
declared(r#"create table ns::trades { a: int4 } with { time: event(block_time) }"#),
"event(block_time)"
);
}
#[test]
fn the_processing_value_parses_through_the_same_key() {
assert_eq!(declared(r#"create table ns::audit { a: int4 } with { time: processing }"#), "processing");
}
#[test]
fn the_none_value_parses_through_the_same_key() {
assert_eq!(declared(r#"create table ns::token { a: int4 } with { time: none }"#), "none");
}
#[test]
fn a_populator_column_may_be_named_with_a_keyword() {
assert_eq!(declared(r#"create table ns::t { a: int4 } with { time: event(event) }"#), "event(event)");
}
#[test]
fn event_without_a_column_does_not_parse() {
rejected(r#"create table ns::t { a: int4 } with { time: event }"#);
rejected(r#"create table ns::t { a: int4 } with { time: event() }"#);
}
#[test]
fn a_column_may_not_be_attached_to_none_or_processing() {
rejected(r#"create table ns::t { a: int4 } with { time: processing(at) }"#);
rejected(r#"create table ns::t { a: int4 } with { time: none(at) }"#);
}
#[test]
fn the_retired_ts_key_is_no_longer_accepted() {
rejected(r#"create table ns::t { a: int4 } with { ts: block_time }"#);
rejected(r#"create table ns::t { a: int4 } with { time: event(at), ts: at }"#);
}
#[test]
fn an_unknown_time_value_is_rejected() {
rejected(r#"create table ns::t { a: int4 } with { time: wallclock }"#);
}
#[test]
fn a_source_object_may_omit_the_declaration_entirely() {
assert_eq!(declared(r#"create table ns::t { a: int4 }"#), "undeclared");
assert_eq!(
declared(r#"create table ns::t { a: int4 } with { partition: { by: { a } } }"#),
"undeclared",
"a WITH block that declares only other keys still leaves the time declaration empty"
);
}
#[test]
fn every_source_object_accepts_the_declaration() {
assert_eq!(
declared(r#"create table ns::t { a: int4 } with { time: event(at) }"#),
"event(at)",
"table"
);
assert_eq!(
declared(r#"create series ns::s { a: int4 } with { key: a, time: event(at) }"#),
"event(at)",
"series"
);
assert_eq!(
declared(r#"create ringbuffer ns::r { a: int4 } with { capacity: 10, time: event(at) }"#),
"event(at)",
"ringbuffer"
);
assert_eq!(
declared(r#"create queue ns::q { a: int4 } with { fifo: {}, time: event(at) }"#),
"event(at)",
"queue"
);
}
#[test]
fn a_view_rejects_a_time_declaration_outright() {
for source in [
r#"create deferred view ns::v { a: int4 } with { time: event(at) } as { from ns::t }"#,
r#"create deferred view ns::v { a: int4 } with { time: none } as { from ns::t }"#,
r#"create transactional view ns::v { a: int4 } with { time: processing } as { from ns::t }"#,
] {
rejected(source);
}
}
#[test]
fn a_storage_backed_view_rejects_a_time_declaration_too() {
for source in [
r#"create deferred ringbuffer view ns::v { a: int4 } with { capacity: 10, time: event(at) } as { from ns::t }"#,
r#"create transactional ringbuffer view ns::v { a: int4 } with { capacity: 10, time: none } as { from ns::t }"#,
r#"create deferred series view ns::v { a: int4 } with { key: a, time: event(at) } as { from ns::t }"#,
r#"create transactional series view ns::v { a: int4 } with { key: a, time: processing } as { from ns::t }"#,
] {
rejected(source);
}
}
#[test]
fn an_unknown_with_key_is_still_rejected_and_lists_the_time_key() {
let bump = Bump::new();
let err = declared_in(&bump, r#"create table ns::t { a: int4 } with { bogus: 1 }"#).unwrap_err();
let message = format!("{:?}", err.diagnostic());
assert!(message.contains("time"), "the diagnostic must offer `time`: {message}");
assert!(!message.contains("'ts'"), "the diagnostic must not still offer the retired `ts`: {message}");
}
}