use super::*;
#[derive(Clone)]
pub(super) enum TimeSource {
Column(String),
Ingest,
None,
}
#[derive(Clone, Copy)]
pub(super) enum EmitShape {
Window,
Rank,
Ring,
}
#[derive(Clone, Debug)]
pub(super) struct WindowSpec {
pub(super) time_column: String,
pub(super) ingest_time: bool,
pub(super) window_ms: i64,
pub(super) slide_ms: Option<i64>,
pub(super) allowed_lateness_ms: i64,
pub(super) idle_timeout_ms: i64,
pub(super) key_by: bool,
pub(super) group_by: Vec<String>,
pub(super) aggs: Vec<fv_streams_ops::Agg>,
pub(super) trigger: fv_streams_ops::Trigger,
}
#[derive(Clone, Debug)]
pub(super) struct SessionSpec {
pub(super) time_column: String,
pub(super) ingest_time: bool,
pub(super) gap_ms: i64,
pub(super) allowed_lateness_ms: i64,
pub(super) idle_timeout_ms: i64,
pub(super) key_by: bool,
pub(super) group_by: Vec<String>,
pub(super) aggs: Vec<fv_streams_ops::Agg>,
}
#[derive(Clone, Debug)]
pub(super) struct TopNSpec {
pub(super) group_by: Vec<String>,
pub(super) order_by: String,
pub(super) descending: bool,
pub(super) n: usize,
pub(super) tie_by: Option<String>,
pub(super) key_by: bool,
}
#[derive(Clone, Debug)]
pub(super) struct LastNSpec {
pub(super) group_by: Vec<String>,
pub(super) n: usize,
pub(super) aggs: Vec<fv_streams_ops::Agg>,
pub(super) key_by: bool,
}
#[derive(Clone)]
pub(super) struct JoinSpec {
pub(super) join_key: String,
pub(super) time_column: String,
pub(super) window_ms: i64,
pub(super) allowed_lateness_ms: i64,
pub(super) idle_timeout_ms: i64,
pub(super) key_by: bool,
}
#[derive(Clone)]
pub(super) struct LookupJoinSpec {
pub(super) join_key: String,
pub(super) distribution: Distribution,
}
pub(super) enum StageOp {
None,
Compute(J),
Windowed(WindowSpec),
Session(SessionSpec),
Join(JoinSpec),
LookupJoin(LookupJoinSpec),
TopN(TopNSpec),
LastN(LastNSpec),
}
pub(super) struct StagePlan {
pub(super) pre: Vec<fv_plan::inline::Step>,
pub(super) op: StageOp,
pub(super) post: Vec<fv_plan::inline::Step>,
}
pub(super) fn idle_timeout_ms(step: &J) -> i64 {
step["idleTimeoutMs"]
.as_i64()
.unwrap_or_else(|| env("STREAM_IDLE_TIMEOUT_MS", "0").parse().unwrap_or(0))
.max(0)
}
pub(super) fn parse_join_spec(step: &J) -> Result<JoinSpec, String> {
Ok(JoinSpec {
join_key: step["joinKey"]
.as_str()
.ok_or("streamJoin: `joinKey` is required")?
.to_string(),
time_column: step["timeColumn"]
.as_str()
.ok_or("streamJoin: `timeColumn` (event-time column, epoch-ms) is required")?
.to_string(),
window_ms: step["windowMs"]
.as_i64()
.filter(|&m| m >= 0)
.ok_or("streamJoin: `windowMs` (>= 0) is required")?,
allowed_lateness_ms: step["allowedLatenessMs"].as_i64().unwrap_or(0).max(0),
idle_timeout_ms: idle_timeout_ms(step),
key_by: step["keyBy"].as_bool().unwrap_or(false),
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum Distribution {
Auto,
Broadcast,
CoPartition,
}
impl Distribution {
pub(super) fn resolve(self, estimated_bytes: Option<u64>, broadcast_max: u64) -> (bool, String) {
match self {
Distribution::Broadcast => (true, "broadcast".into()),
Distribution::CoPartition => (false, "co-partition".into()),
Distribution::Auto => match estimated_bytes {
Some(b) if b > broadcast_max => (
false,
format!(
"co-partition (auto: table ≈ {} MiB > {} MiB)",
b >> 20,
broadcast_max >> 20
),
),
Some(b) => (true, format!("broadcast (auto: table ≈ {} KiB)", b >> 10)),
None => (true, "broadcast (auto: table size unknown)".into()),
},
}
}
}
pub(super) fn parse_lookup_join_spec(step: &J) -> Result<LookupJoinSpec, String> {
let join_key = step["joinKey"]
.as_str()
.filter(|s| !s.is_empty())
.ok_or("lookupJoin: `joinKey` is required (present on both the stream and the table)")?
.to_string();
let distribution = match step["distribution"].as_str() {
Some("auto") => Distribution::Auto,
Some("broadcast") => Distribution::Broadcast,
Some("coPartition") | Some("co-partition") => Distribution::CoPartition,
Some(other) => {
return Err(format!(
"lookupJoin: `distribution` must be \"auto\", \"broadcast\" or \"coPartition\"; got {other:?}"
))
}
None if step["coPartition"].as_bool().unwrap_or(false) => Distribution::CoPartition,
None => Distribution::Auto,
};
Ok(LookupJoinSpec { join_key, distribution })
}
pub(super) fn parse_window_spec(step: &J) -> Result<WindowSpec, String> {
let (time_column, ingest_time) = parse_time_source(step, "windowedAggregate")?;
let window_ms = step["windowMs"]
.as_i64()
.filter(|&m| m > 0)
.ok_or("windowedAggregate: `windowMs` (> 0) is required")?;
let slide_ms = match step["slideMs"].as_i64() {
None => None,
Some(s) if s > 0 && s <= window_ms => Some(s),
Some(s) => {
return Err(format!(
"windowedAggregate: `slideMs` must be in (0, windowMs]; got {s}"
))
}
};
let allowed_lateness_ms = step["allowedLatenessMs"].as_i64().unwrap_or(0).max(0);
let group_by: Vec<String> = step["groupBy"]
.as_array()
.map(|a| a.iter().filter_map(|v| v.as_str().map(String::from)).collect())
.unwrap_or_default();
let aggs = parse_aggs(step, "windowedAggregate")?;
let key_by = step["keyBy"].as_bool().unwrap_or(false);
let trigger = parse_trigger(step)?;
Ok(WindowSpec {
time_column,
ingest_time,
window_ms,
slide_ms,
allowed_lateness_ms,
idle_timeout_ms: idle_timeout_ms(step),
key_by,
group_by,
aggs,
trigger,
})
}
fn parse_trigger(step: &J) -> Result<fv_streams_ops::Trigger, String> {
use fv_streams_ops::Trigger;
let emit = &step["emit"];
if emit.is_null() {
return Ok(Trigger::OnWatermark);
}
let rows = emit["everyRows"].as_i64();
let ms = emit["everyMs"].as_i64();
match (rows, ms) {
(Some(_), Some(_)) => Err("windowedAggregate `emit`: set at most one of `everyRows` / `everyMs`".into()),
(Some(n), None) if n > 0 => Ok(Trigger::EveryRows(n as u64)),
(Some(n), None) => Err(format!("windowedAggregate `emit.everyRows` must be > 0; got {n}")),
(None, Some(t)) if t > 0 => Ok(Trigger::EveryMs(t)),
(None, Some(t)) => Err(format!("windowedAggregate `emit.everyMs` must be > 0; got {t}")),
(None, None) => Ok(Trigger::OnWatermark), }
}
pub(super) fn parse_session_spec(step: &J) -> Result<SessionSpec, String> {
let (time_column, ingest_time) = parse_time_source(step, "sessionAggregate")?;
let gap_ms = step["gapMs"]
.as_i64()
.filter(|&g| g > 0)
.ok_or("sessionAggregate: `gapMs` (> 0, inactivity gap that closes a session) is required")?;
let allowed_lateness_ms = step["allowedLatenessMs"].as_i64().unwrap_or(0).max(0);
let group_by: Vec<String> = step["groupBy"]
.as_array()
.map(|a| a.iter().filter_map(|v| v.as_str().map(String::from)).collect())
.unwrap_or_default();
let aggs = parse_aggs(step, "sessionAggregate")?;
let key_by = step["keyBy"].as_bool().unwrap_or(false);
Ok(SessionSpec {
time_column,
ingest_time,
gap_ms,
allowed_lateness_ms,
idle_timeout_ms: idle_timeout_ms(step),
key_by,
group_by,
aggs,
})
}
pub(super) fn parse_time_source(step: &J, op_name: &str) -> Result<(String, bool), String> {
let ingest = step["ingestTime"].as_bool().unwrap_or(false);
match (ingest, step["timeColumn"].as_str()) {
(true, Some(_)) => Err(format!(
"{op_name}: `ingestTime: true` and `timeColumn` are mutually exclusive — processing-time windows have no event-time column"
)),
(true, None) => Ok((String::new(), true)),
(false, Some(t)) => Ok((t.to_string(), false)),
(false, None) => Err(format!(
"{op_name}: `timeColumn` (event-time column, epoch-ms) is required (or `ingestTime: true` for processing-time windows)"
)),
}
}
pub(super) fn parse_topn_spec(step: &J) -> Result<TopNSpec, String> {
let order_by = step["orderBy"]
.as_str()
.ok_or("topN: `orderBy` (the numeric ranking column) is required")?
.to_string();
let n = step["n"]
.as_i64()
.filter(|&n| n > 0)
.ok_or("topN: `n` (> 0) is required")? as usize;
let descending = match step["direction"].as_str() {
None | Some("desc") => true,
Some("asc") => false,
Some(other) => {
return Err(format!(
"topN: `direction` must be \"desc\" or \"asc\"; got \"{other}\""
))
}
};
let group_by: Vec<String> = step["groupBy"]
.as_array()
.map(|a| a.iter().filter_map(|v| v.as_str().map(String::from)).collect())
.unwrap_or_default();
let tie_by = step["tieBy"].as_str().map(String::from);
let key_by = step["keyBy"].as_bool().unwrap_or(false);
Ok(TopNSpec {
group_by,
order_by,
descending,
n,
tie_by,
key_by,
})
}
pub(super) fn parse_lastn_spec(step: &J) -> Result<LastNSpec, String> {
let n = step["n"]
.as_i64()
.filter(|&n| n > 0)
.ok_or("lastN: `n` (> 0, ring size per key) is required")? as usize;
let group_by: Vec<String> = step["groupBy"]
.as_array()
.map(|a| a.iter().filter_map(|v| v.as_str().map(String::from)).collect())
.unwrap_or_default();
let aggs = parse_aggs(step, "lastN")?;
let key_by = step["keyBy"].as_bool().unwrap_or(false);
Ok(LastNSpec {
group_by,
n,
aggs,
key_by,
})
}
pub(super) fn parse_aggs(step: &J, op_name: &str) -> Result<Vec<fv_streams_ops::Agg>, String> {
let aggs = step["aggs"]
.as_array()
.ok_or_else(|| format!("{op_name}: `aggs` array is required"))?
.iter()
.map(|a| {
let op = a["op"].as_str().ok_or_else(|| format!("{op_name} agg: `op` is required"))?;
let alias = a["alias"].as_str().ok_or_else(|| format!("{op_name} agg: `alias` is required"))?.to_string();
let op = match op {
"count" => fv_streams_ops::AggOp::Count,
"sum" => fv_streams_ops::AggOp::Sum,
"avg" => fv_streams_ops::AggOp::Avg,
"min" => fv_streams_ops::AggOp::Min,
"max" => fv_streams_ops::AggOp::Max,
"countDistinct" => fv_streams_ops::AggOp::CountDistinct,
"varPop" => fv_streams_ops::AggOp::VarPop,
"varSamp" => fv_streams_ops::AggOp::VarSamp,
"stddevPop" => fv_streams_ops::AggOp::StddevPop,
"stddevSamp" => fv_streams_ops::AggOp::StddevSamp,
"boolAnd" => fv_streams_ops::AggOp::BoolAnd,
"boolOr" => fv_streams_ops::AggOp::BoolOr,
"bitAnd" => fv_streams_ops::AggOp::BitAnd,
"bitOr" => fv_streams_ops::AggOp::BitOr,
"bitXor" => fv_streams_ops::AggOp::BitXor,
other => return Err(format!("{op_name} agg: unknown op `{other}` (count/sum/avg/min/max/countDistinct/varPop/varSamp/stddevPop/stddevSamp/boolAnd/boolOr/bitAnd/bitOr/bitXor)")),
};
let column = a["column"].as_str().unwrap_or("").to_string();
if op != fv_streams_ops::AggOp::Count && column.is_empty() {
return Err(format!("{op_name} agg `{alias}`: `column` is required for {op:?}"));
}
Ok(fv_streams_ops::Agg { op, column, alias })
})
.collect::<Result<Vec<_>, String>>()?;
if aggs.is_empty() {
return Err(format!("{op_name}: at least one agg is required"));
}
Ok(aggs)
}
pub(super) fn classify_steps(raw: &[J]) -> Result<StagePlan, String> {
let is_heavy = |s: &J| {
matches!(
s["op"].as_str(),
Some(
"windowedAggregate"
| "sessionAggregate"
| "streamJoin"
| "lookupJoin"
| "topN"
| "lastN"
| "wasm"
| "container"
)
)
};
let heavy: Vec<usize> = raw
.iter()
.enumerate()
.filter(|(_, s)| is_heavy(s))
.map(|(i, _)| i)
.collect();
if heavy.len() > 1 {
return Err(format!(
"a stage may contain at most ONE stateful/compute step (found {}) — split the pipeline into multiple stages (chained topologies)",
heavy.len()
));
}
let parse_inline = |slice: &[J], place: &str| -> Result<Vec<fv_plan::inline::Step>, String> {
let steps = slice
.iter()
.map(fv_plan::inline::parse)
.collect::<Result<Vec<_>, _>>()
.map_err(|e| format!("unsupported {place} step for streaming (stateless ops only — select/rename/drop/filter/applyExpression): {e}"))?;
fv_plan::inline::validate_steps(&steps).map_err(|e| format!("{place} step expression invalid: {e}"))?;
Ok(steps)
};
match heavy.first() {
None => Ok(StagePlan {
pre: parse_inline(raw, "inline")?,
op: StageOp::None,
post: Vec::new(),
}),
Some(&i) => {
let pre = parse_inline(&raw[..i], "pre")?;
let post = parse_inline(&raw[i + 1..], "post")?;
let s = &raw[i];
let op = match s["op"].as_str() {
Some("windowedAggregate") => StageOp::Windowed(parse_window_spec(s)?),
Some("sessionAggregate") => StageOp::Session(parse_session_spec(s)?),
Some("streamJoin") => {
if !pre.is_empty() {
return Err(
"a streamJoin stage cannot have PRE steps (which of the two inputs would they apply to?) — transform each input in its own preceding stage".into(),
);
}
StageOp::Join(parse_join_spec(s)?)
}
Some("lookupJoin") => {
if !pre.is_empty() {
return Err(
"a lookupJoin stage cannot have PRE steps (which of the two inputs would they apply to?) — transform the stream and the table each in its own preceding stage".into(),
);
}
StageOp::LookupJoin(parse_lookup_join_spec(s)?)
}
Some("topN") => StageOp::TopN(parse_topn_spec(s)?),
Some("lastN") => StageOp::LastN(parse_lastn_spec(s)?),
Some("wasm") | Some("container") => StageOp::Compute(s.clone()),
_ => unreachable!("is_heavy gate"),
};
Ok(StagePlan { pre, op, post })
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn classify_inline_steps() {
let raw = vec![json!({"op": "applyExpression", "column": "d", "expression": "amount * 2"})];
let p = classify_steps(&raw).unwrap();
assert!(matches!(p.op, StageOp::None));
assert_eq!(p.pre.len(), 1);
}
#[test]
fn classify_single_compute_step() {
for op in ["wasm", "container"] {
let raw = vec![json!({"op": op, "ref": "spikeTotal"})];
match classify_steps(&raw).unwrap().op {
StageOp::Compute(s) => assert_eq!(s["ref"], "spikeTotal"),
_ => panic!("expected compute for op={op}"),
}
}
}
#[test]
fn classify_mixed_inline_and_compute_splits_pre_and_op() {
let raw = vec![
json!({"op": "filter", "expression": "amount > 0"}),
json!({"op": "wasm", "ref": "x"}),
];
let p = classify_steps(&raw).unwrap();
assert_eq!(p.pre.len(), 1);
assert!(matches!(p.op, StageOp::Compute(_)));
assert!(p.post.is_empty());
}
#[test]
fn classify_rejects_a_bad_inline_step() {
assert!(classify_steps(&[json!({"op": "bogusOp"})]).is_err());
}
#[test]
fn classify_windowed_aggregate_step() {
let raw = vec![json!({
"op": "windowedAggregate", "timeColumn": "ts", "windowMs": 1000, "allowedLatenessMs": 200,
"groupBy": ["k"],
"aggs": [{"op": "count", "alias": "n"}, {"op": "sum", "column": "amount", "alias": "total"}]
})];
match classify_steps(&raw).unwrap().op {
StageOp::Windowed(w) => {
assert_eq!(w.time_column, "ts");
assert_eq!((w.window_ms, w.allowed_lateness_ms), (1000, 200));
assert_eq!(w.group_by, vec!["k".to_string()]);
assert_eq!(w.aggs.len(), 2);
assert_eq!(w.aggs[0].op, fv_streams_ops::AggOp::Count);
assert_eq!(w.aggs[1].op, fv_streams_ops::AggOp::Sum);
assert_eq!(w.trigger, fv_streams_ops::Trigger::OnWatermark, "no `emit` ⇒ plain");
}
_ => panic!("expected windowed"),
}
}
#[test]
fn windowed_aggregate_parses_the_emit_trigger() {
use fv_streams_ops::Trigger;
let base = |emit: J| {
json!({
"op": "windowedAggregate", "timeColumn": "ts", "windowMs": 86_400_000,
"keyBy": true, "groupBy": ["auction"], "aggs": [{"op": "count", "alias": "n"}],
"emit": emit,
})
};
let trig = |emit: J| parse_window_spec(&base(emit)).unwrap().trigger;
assert_eq!(
parse_window_spec(
&json!({"op":"windowedAggregate","timeColumn":"ts","windowMs":1000,"aggs":[{"op":"count","alias":"n"}]})
)
.unwrap()
.trigger,
Trigger::OnWatermark
);
assert_eq!(trig(json!({})), Trigger::OnWatermark);
assert_eq!(trig(json!({"everyRows": 1000})), Trigger::EveryRows(1000));
assert_eq!(trig(json!({"everyMs": 500})), Trigger::EveryMs(500));
assert!(parse_window_spec(&base(json!({"everyRows": 100, "everyMs": 100}))).is_err());
assert!(parse_window_spec(&base(json!({"everyRows": 0}))).is_err());
assert!(parse_window_spec(&base(json!({"everyMs": -1}))).is_err());
}
#[test]
fn window_spec_rejects_bad_config() {
assert!(parse_window_spec(&json!({"windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]})).is_err());
assert!(parse_window_spec(&json!({"timeColumn": "ts", "aggs": [{"op":"count","alias":"n"}]})).is_err());
assert!(
parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 0, "aggs": [{"op":"count","alias":"n"}]}))
.is_err()
);
assert!(parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "aggs": []})).is_err());
assert!(parse_window_spec(
&json!({"timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"median","alias":"m"}]})
)
.is_err());
assert!(
parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"sum","alias":"s"}]}))
.is_err()
);
}
#[test]
fn window_spec_reads_slide_for_sliding_windows() {
let base = |slide: J| json!({"timeColumn": "ts", "windowMs": 1000, "slideMs": slide, "aggs": [{"op":"count","alias":"n"}]});
assert_eq!(parse_window_spec(&base(json!(500))).unwrap().slide_ms, Some(500));
assert_eq!(
parse_window_spec(&json!({"timeColumn":"ts","windowMs":1000,"aggs":[{"op":"count","alias":"n"}]}))
.unwrap()
.slide_ms,
None
); assert!(
parse_window_spec(&base(json!(1500))).is_err(),
"slide > window rejected"
);
assert!(parse_window_spec(&base(json!(0))).is_err(), "slide 0 rejected");
}
#[test]
fn classify_session_aggregate_step() {
let raw = vec![json!({
"op": "sessionAggregate", "timeColumn": "ts", "gapMs": 5000, "allowedLatenessMs": 100,
"groupBy": ["bidder"],
"aggs": [{"op": "count", "alias": "bids"}, {"op": "max", "column": "price", "alias": "top"}]
})];
match classify_steps(&raw).unwrap().op {
StageOp::Session(s) => {
assert_eq!(s.time_column, "ts");
assert_eq!((s.gap_ms, s.allowed_lateness_ms), (5000, 100));
assert_eq!(s.group_by, vec!["bidder".to_string()]);
assert_eq!(s.aggs.len(), 2);
assert_eq!(s.aggs[0].op, fv_streams_ops::AggOp::Count);
assert_eq!(s.aggs[1].op, fv_streams_ops::AggOp::Max);
}
_ => panic!("expected session"),
}
}
#[test]
fn session_spec_rejects_bad_config() {
assert!(parse_session_spec(&json!({"gapMs": 5000, "aggs": [{"op":"count","alias":"n"}]})).is_err());
assert!(parse_session_spec(&json!({"timeColumn": "ts", "aggs": [{"op":"count","alias":"n"}]})).is_err());
assert!(
parse_session_spec(&json!({"timeColumn": "ts", "gapMs": 0, "aggs": [{"op":"count","alias":"n"}]})).is_err()
);
assert!(parse_session_spec(&json!({"timeColumn": "ts", "gapMs": 5000})).is_err());
assert!(
parse_session_spec(&json!({"timeColumn": "ts", "gapMs": 5000, "aggs": [{"op":"sum","alias":"s"}]}))
.is_err()
);
}
#[test]
fn classify_session_with_pre_step() {
let raw = vec![
json!({"op": "filter", "expression": "amount > 0"}),
json!({"op": "sessionAggregate", "timeColumn": "ts", "gapMs": 1000, "aggs": [{"op":"count","alias":"n"}]}),
];
let p = classify_steps(&raw).unwrap();
assert_eq!(p.pre.len(), 1);
assert!(matches!(p.op, StageOp::Session(_)));
}
#[test]
fn key_by_is_parsed() {
let w = parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "keyBy": true, "groupBy": ["k"], "aggs": [{"op":"count","alias":"n"}]})).unwrap();
assert!(w.key_by);
let s = parse_session_spec(&json!({"timeColumn": "ts", "gapMs": 1000, "keyBy": true, "groupBy": ["k"], "aggs": [{"op":"count","alias":"n"}]})).unwrap();
assert!(s.key_by);
let j =
parse_join_spec(&json!({"joinKey": "pid", "timeColumn": "ts", "windowMs": 1000, "keyBy": true})).unwrap();
assert!(j.key_by);
let w0 =
parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]}))
.unwrap();
assert!(!w0.key_by);
let j0 = parse_join_spec(&json!({"joinKey": "pid", "timeColumn": "ts", "windowMs": 1000})).unwrap();
assert!(!j0.key_by);
}
#[test]
fn idle_timeout_is_parsed_for_stateful_steps() {
let w = parse_window_spec(
&json!({"timeColumn": "ts", "windowMs": 1000, "idleTimeoutMs": 5000, "aggs": [{"op":"count","alias":"n"}]}),
)
.unwrap();
assert_eq!(w.idle_timeout_ms, 5000);
let s = parse_session_spec(
&json!({"timeColumn": "ts", "gapMs": 1000, "idleTimeoutMs": 7000, "aggs": [{"op":"count","alias":"n"}]}),
)
.unwrap();
assert_eq!(s.idle_timeout_ms, 7000);
let j =
parse_join_spec(&json!({"joinKey": "pid", "timeColumn": "ts", "windowMs": 1000, "idleTimeoutMs": 9000}))
.unwrap();
assert_eq!(j.idle_timeout_ms, 9000);
let w0 =
parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]}))
.unwrap();
assert_eq!(w0.idle_timeout_ms, 0);
}
#[test]
fn classify_lookup_join_step() {
match classify_steps(&[json!({"op": "lookupJoin", "joinKey": "key"})])
.unwrap()
.op
{
StageOp::LookupJoin(l) => {
assert_eq!(l.join_key, "key");
assert_eq!(l.distribution, Distribution::Auto, "default = auto");
}
_ => panic!("expected lookupJoin"),
}
let dist = |s: J| match classify_steps(&[s]).unwrap().op {
StageOp::LookupJoin(l) => l.distribution,
_ => panic!("expected lookupJoin"),
};
assert_eq!(
dist(json!({"op": "lookupJoin", "joinKey": "k", "distribution": "broadcast"})),
Distribution::Broadcast
);
assert_eq!(
dist(json!({"op": "lookupJoin", "joinKey": "k", "distribution": "coPartition"})),
Distribution::CoPartition
);
assert_eq!(
dist(json!({"op": "lookupJoin", "joinKey": "k", "coPartition": true})),
Distribution::CoPartition
);
assert_eq!(
dist(json!({"op": "lookupJoin", "joinKey": "k", "distribution": "auto"})),
Distribution::Auto
);
let max = 64 << 20;
assert!(Distribution::Auto.resolve(Some(1 << 20), max).0);
assert!(Distribution::Auto.resolve(None, max).0);
let (b, how) = Distribution::Auto.resolve(Some(200 << 20), max);
assert!(!b && how.contains("200 MiB > 64 MiB"), "{how}");
assert!(Distribution::Broadcast.resolve(Some(200 << 20), max).0, "forced");
assert!(!Distribution::CoPartition.resolve(Some(1), max).0, "forced");
assert!(parse_lookup_join_spec(&json!({})).is_err());
assert!(parse_lookup_join_spec(&json!({"joinKey": "k", "distribution": "sideways"})).is_err());
assert!(classify_steps(&[
json!({"op": "filter", "expression": "price > 0"}),
json!({"op": "lookupJoin", "joinKey": "key"}),
])
.is_err());
}
#[test]
fn classify_stream_join_step() {
let raw = vec![
json!({"op": "streamJoin", "joinKey": "pid", "timeColumn": "ts", "windowMs": 5000, "allowedLatenessMs": 100}),
];
match classify_steps(&raw).unwrap().op {
StageOp::Join(j) => {
assert_eq!((j.join_key.as_str(), j.time_column.as_str()), ("pid", "ts"));
assert_eq!((j.window_ms, j.allowed_lateness_ms), (5000, 100));
}
_ => panic!("expected join"),
}
assert!(
parse_join_spec(&json!({"timeColumn": "ts", "windowMs": 5000})).is_err(),
"joinKey required"
);
assert!(
parse_join_spec(&json!({"joinKey": "pid", "windowMs": 5000})).is_err(),
"timeColumn required"
);
assert!(
parse_join_spec(&json!({"joinKey": "pid", "timeColumn": "ts"})).is_err(),
"windowMs required"
);
}
#[test]
fn classify_mixed_pre_window_post() {
let raw = vec![
json!({"op": "filter", "expression": "amount > 0"}),
json!({"op": "windowedAggregate", "timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]}),
json!({"op": "filter", "expression": "n > 5"}),
];
let p = classify_steps(&raw).unwrap();
assert_eq!((p.pre.len(), p.post.len()), (1, 1));
assert!(matches!(p.op, StageOp::Windowed(_)));
}
#[test]
fn classify_rejects_two_heavy_steps_and_join_pre_steps() {
let raw = vec![
json!({"op": "windowedAggregate", "timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]}),
json!({"op": "windowedAggregate", "timeColumn": "windowStart", "windowMs": 5000, "aggs": [{"op":"sum","column":"n","alias":"t"}]}),
];
assert!(classify_steps(&raw)
.err()
.expect("two heavy steps must fail")
.contains("multiple stages"));
let raw = vec![
json!({"op": "filter", "expression": "amount > 0"}),
json!({"op": "streamJoin", "joinKey": "pid", "timeColumn": "ts", "windowMs": 5000}),
];
assert!(classify_steps(&raw)
.err()
.expect("join pre steps must fail")
.contains("preceding stage"));
}
#[test]
fn parse_topn_spec_defaults_and_validation() {
let t = parse_topn_spec(
&json!({"orderBy": "price", "n": 10, "groupBy": ["auction"], "tieBy": "ts", "keyBy": true}),
)
.unwrap();
assert_eq!((t.n, t.descending, t.key_by), (10, true, true)); assert_eq!(t.tie_by.as_deref(), Some("ts"));
let asc = parse_topn_spec(&json!({"orderBy": "ts", "n": 1, "direction": "asc"})).unwrap();
assert!(!asc.descending);
assert!(asc.group_by.is_empty()); assert!(parse_topn_spec(&json!({"n": 1})).unwrap_err().contains("orderBy"));
assert!(parse_topn_spec(&json!({"orderBy": "p", "n": 0}))
.unwrap_err()
.contains("`n`"));
assert!(
parse_topn_spec(&json!({"orderBy": "p", "n": 1, "direction": "sideways"}))
.unwrap_err()
.contains("direction")
);
}
#[test]
fn parse_lastn_spec_requires_n_and_aggs() {
let l = parse_lastn_spec(
&json!({"n": 10, "groupBy": ["seller"], "aggs": [{"op": "avg", "column": "price", "alias": "avgPrice"}]}),
)
.unwrap();
assert_eq!((l.n, l.group_by.len(), l.aggs.len()), (10, 1, 1));
assert!(
parse_lastn_spec(&json!({"groupBy": ["s"], "aggs": [{"op":"count","alias":"n"}]}))
.unwrap_err()
.contains("`n`")
);
assert!(parse_lastn_spec(&json!({"n": 5})).unwrap_err().contains("aggs"));
}
#[test]
fn parse_func2_aggregate_ops() {
for op in [
"varPop",
"varSamp",
"stddevPop",
"stddevSamp",
"boolAnd",
"boolOr",
"bitAnd",
"bitOr",
"bitXor",
] {
let w = parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "groupBy": [],
"aggs": [{"op": op, "column": "v", "alias": "r"}]}))
.unwrap();
assert_eq!(w.aggs.len(), 1, "op {op} parses");
}
let err = parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "groupBy": [],
"aggs": [{"op": "varPop", "alias": "r"}]}))
.unwrap_err();
assert!(err.contains("column"), "varPop needs a column: {err}");
let err = parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "groupBy": [],
"aggs": [{"op": "bogus", "alias": "r"}]}))
.unwrap_err();
assert!(err.contains("varPop") && err.contains("bitXor"));
}
#[test]
fn parse_count_distinct_agg() {
let w = parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "groupBy": [],
"aggs": [{"op": "countDistinct", "column": "bidder", "alias": "uniq"}]}))
.unwrap();
assert_eq!(w.aggs[0].op, fv_streams_ops::AggOp::CountDistinct);
let err = parse_window_spec(&json!({"timeColumn": "ts", "windowMs": 1000, "groupBy": [],
"aggs": [{"op": "countDistinct", "alias": "uniq"}]}))
.unwrap_err();
assert!(err.contains("column"));
}
#[test]
fn ingest_time_is_exclusive_with_time_column() {
let w = parse_window_spec(
&json!({"ingestTime": true, "windowMs": 1000, "groupBy": ["k"], "aggs": [{"op":"count","alias":"n"}]}),
)
.unwrap();
assert!(w.ingest_time && w.time_column.is_empty());
let both = parse_window_spec(
&json!({"ingestTime": true, "timeColumn": "ts", "windowMs": 1000, "groupBy": [], "aggs": [{"op":"count","alias":"n"}]}),
);
assert!(both.unwrap_err().contains("mutually exclusive"));
let neither =
parse_window_spec(&json!({"windowMs": 1000, "groupBy": [], "aggs": [{"op":"count","alias":"n"}]}));
assert!(neither.unwrap_err().contains("timeColumn"));
let sess = parse_session_spec(
&json!({"ingestTime": true, "gapMs": 1000, "groupBy": ["k"], "aggs": [{"op":"count","alias":"n"}]}),
)
.unwrap();
assert!(sess.ingest_time);
}
#[test]
fn classify_topn_and_lastn_are_heavy_steps() {
let raw = vec![
json!({"op": "filter", "expression": "price > 0"}),
json!({"op": "topN", "orderBy": "price", "n": 1, "groupBy": ["auction"]}),
];
let p = classify_steps(&raw).unwrap();
assert_eq!(p.pre.len(), 1);
assert!(matches!(p.op, StageOp::TopN(_)));
let raw = vec![
json!({"op": "lastN", "n": 10, "groupBy": ["seller"], "aggs": [{"op":"avg","column":"price","alias":"a"}]}),
];
assert!(matches!(classify_steps(&raw).unwrap().op, StageOp::LastN(_)));
let raw = vec![
json!({"op": "topN", "orderBy": "price", "n": 1}),
json!({"op": "lastN", "n": 10, "aggs": [{"op":"count","alias":"n"}]}),
];
assert!(classify_steps(&raw)
.err()
.expect("two heavy steps")
.contains("multiple stages"));
}
}