Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 97 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions dozer-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -64,4 +64,5 @@ mongodb = ["dozer-ingestion/mongodb"]
onnx = ["dozer-sql/onnx"]
tokio-console = ["dozer-tracing/tokio-console"]
javascript = ["dozer-ingestion/javascript", "dozer-sql/javascript"]
wasm = ["dozer-sql/wasm"]
datafusion = ["dozer-ingestion/datafusion"]
3 changes: 3 additions & 0 deletions dozer-sql/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,11 @@ tokio = { version = "1", features = ["rt", "macros"] }

[dev-dependencies]
proptest = "1.3.1"
tempfile = "3.10.1"
wat = "=1.0.85"

[features]
python = ["dozer-sql-expression/python"]
onnx = ["dozer-sql-expression/onnx"]
javascript = ["dozer-sql-expression/javascript"]
wasm = ["dozer-sql-expression/wasm"]
8 changes: 8 additions & 0 deletions dozer-sql/expression/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,23 @@ jsonpath = { path = "../jsonpath" }
bincode = { workspace = true }
tokio = "1.34.0"
async-recursion = "1.0.5"
wasmi = { version = "0.31.2", optional = true }

dozer-deno = { path = "../../dozer-deno", optional = true }
deno_core = { workspace = true, optional = true }

[dev-dependencies]
proptest = "1.2.0"
wat = "=1.0.85"
tempfile = "3.10.1"

[features]
bigdecimal = ["dep:bigdecimal", "sqlparser/bigdecimal"]
python = ["dozer-types/python-auto-initialize"]
onnx = ["dep:ort", "dep:ndarray", "dep:half"]
javascript = ["dep:dozer-deno", "dep:deno_core"]
wasm = ["dep:wasmi"]

[[example]]
name = "wasm_udf"
required-features = ["wasm"]
91 changes: 91 additions & 0 deletions dozer-sql/expression/examples/wasm_udf.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@
//! Evaluate the AssemblyScript example through Dozer's SQL parser and expression planner.
use std::{error::Error, sync::Arc};

use dozer_sql_expression::{
builder::ExpressionBuilder,
sqlparser::{
ast::{SelectItem, SetExpr, Statement},
dialect::DozerDialect,
parser::Parser,
},
};
use dozer_types::{
models::config::Config,
serde_yaml,
types::{Field, FieldDefinition, FieldType, Record, Schema, SourceDefinition},
};

fn main() -> Result<(), Box<dyn Error>> {
let config_path = std::env::args()
.nth(1)
.ok_or("usage: wasm_udf <config.yaml>")?;
let config: Config = serde_yaml::from_str(&std::fs::read_to_string(config_path)?)?;
let sql = config.sql.as_deref().ok_or("config must include sql")?;
let statements = Parser::parse_sql(&DozerDialect {}, sql)?;
let Statement::Query(query) = &statements[0] else {
return Err("expected SELECT".into());
};
let SetExpr::Select(select) = &*query.body else {
return Err("expected SELECT body".into());
};
let schema = Schema::default()
.field(
FieldDefinition::new(
"value".into(),
FieldType::Int,
true,
SourceDefinition::Dynamic,
),
false,
)
.field(
FieldDefinition::new(
"amount".into(),
FieldType::Float,
true,
SourceDefinition::Dynamic,
),
false,
)
.to_owned();
let runtime = Arc::new(tokio::runtime::Builder::new_current_thread().build()?);
let mut builder = ExpressionBuilder::new(schema.fields.len(), runtime.clone());
println!("SQL: {}", sql.trim());
let mut expressions = Vec::new();
for item in &select.projection {
let SelectItem::UnnamedExpr(expr) = item else {
return Err("example expects expressions without aliases".into());
};
let expression = runtime.block_on(builder.build(false, expr, &schema, &config.udfs))?;
expression.get_type(&schema)?;
expressions.push(expression);
}
for input in [Field::Int(41), Field::Null, Field::Int(-1)] {
let amount = match &input {
Field::Int(value) => Field::Float((*value as f64).into()),
_ => Field::Null,
};
println!(
"\nInput value: {}, amount: {}",
display(&input),
display(&amount)
);
let record = Record::new(vec![input, amount]);
for expression in &mut expressions {
let label = expression.to_string(&schema);
match expression.evaluate(&record, &schema) {
Ok(value) => println!(" {label} = {}", display(&value)),
Err(error) => println!(" {label}: {error}"),
}
}
}
Ok(())
}

fn display(value: &Field) -> String {
if matches!(value, Field::Null) {
"NULL".into()
} else {
value.to_string()
}
}
29 changes: 29 additions & 0 deletions dozer-sql/expression/src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -573,6 +573,35 @@ impl ExpressionBuilder {
Err(Error::JavaScriptNotEnabled)
}
}
UdfType::Wasm(config) => {
#[cfg(feature = "wasm")]
{
let mut args = Vec::with_capacity(sql_function.args.len());
for argument in &sql_function.args {
args.push(
self.parse_sql_function_arg(
parse_aggregations,
argument,
schema,
udfs,
)
.await?,
);
}
let udf = crate::wasm::Udf::new(function_name, config, args)?;
// Selection factories do not request expression types separately.
// Aggregate projections must wait for the post-aggregation schema.
if !parse_aggregations {
udf.get_type(schema)?;
}
Ok(Expression::WasmUdf(Box::new(udf)))
}
#[cfg(not(feature = "wasm"))]
{
let _ = config;
Err(Error::WasmNotEnabled)
}
}
};
}

Expand Down
7 changes: 7 additions & 0 deletions dozer-sql/expression/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,13 @@ pub enum Error {
#[error("Javascript is not enabled")]
JavaScriptNotEnabled,

#[error("WASM UDF support is not enabled; build with --features wasm")]
WasmNotEnabled,

#[cfg(feature = "wasm")]
#[error("WASM UDF error: {0}")]
Wasm(String),

#[cfg(feature = "javascript")]
#[error("JavaScript UDF error: {0}")]
JavaScript(#[from] crate::javascript::Error),
Expand Down
8 changes: 8 additions & 0 deletions dozer-sql/expression/src/execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,8 @@ pub enum Expression {
},
#[cfg(feature = "javascript")]
JavaScriptUdf(crate::javascript::Udf),
#[cfg(feature = "wasm")]
WasmUdf(Box<crate::wasm::Udf>),
}

impl Expression {
Expand Down Expand Up @@ -285,6 +287,8 @@ impl Expression {
}
#[cfg(feature = "javascript")]
Expression::JavaScriptUdf(udf) => udf.to_string(schema),
#[cfg(feature = "wasm")]
Expression::WasmUdf(udf) => udf.to_string(schema),
Expression::IsNull { arg } => arg.to_string(schema) + " IS NULL ",
Expression::IsNotNull { arg } => arg.to_string(schema) + " IS NOT NULL ",
}
Expand Down Expand Up @@ -378,6 +382,8 @@ impl Expression {
Expression::IsNotNull { arg } => evaluate_is_not_null(schema, arg, record),
#[cfg(feature = "javascript")]
Expression::JavaScriptUdf(udf) => udf.evaluate(record, schema),
#[cfg(feature = "wasm")]
Expression::WasmUdf(udf) => udf.evaluate(record, schema),
}
}

Expand Down Expand Up @@ -487,6 +493,8 @@ impl Expression {
)),
#[cfg(feature = "javascript")]
Expression::JavaScriptUdf(udf) => Ok(udf.get_type()),
#[cfg(feature = "wasm")]
Expression::WasmUdf(udf) => udf.get_type(schema),
Expression::IsNull { arg: _ } => Ok(ExpressionType::new(
FieldType::Boolean,
false,
Expand Down
Loading