use dataflow_rs::{Engine, Message, Workflow};
use serde_json::{Value, json};
use std::time::Instant;
const CONFIGS: &[(usize, usize)] = &[(5, 150_000), (25, 60_000), (100, 15_000)];
fn subtree_workflow(k: usize) -> Workflow {
let mappings: Vec<String> = (0..k)
.map(|i| {
format!(
r#"{{ "path": "data.DOC.f{i}",
"logic": {{ "cat": [ {{ "var": "data.input.party.name" }}, " :: ",
{{ "var": "data.input.groupHeader.msgId" }}, " :: field-{i:04}" ] }} }}"#
)
})
.collect();
let json = format!(
r#"{{
"id": "wf_subtree",
"name": "Subtree writes",
"priority": 0,
"condition": true,
"tasks": [
{{
"id": "parse",
"name": "Load payload",
"function": {{ "name": "parse_json", "input": {{ "source": "payload", "target": "input" }} }}
}},
{{
"id": "assemble",
"name": "Assemble DOC",
"function": {{ "name": "map", "input": {{ "mappings": [ {} ] }} }}
}}
]
}}"#,
mappings.join(",")
);
Workflow::from_json(&json).unwrap()
}
fn build_payload() -> Value {
json!({
"party": { "name": "Acme Global Payments GmbH", "country": "DE" },
"groupHeader": { "msgId": "MSG-2026-000123", "creationDateTime": "2026-07-17T10:00:00Z" },
"amount": { "value": 2500, "currency": "EUR" }
})
}
async fn time_config(k: usize, n: usize, payload: &Value) -> (f64, f64) {
let engine = Engine::builder()
.with_workflow(subtree_workflow(k))
.build()
.unwrap();
{
let mut m = Message::from_value(payload);
engine.process_message(&mut m).await.unwrap();
assert_eq!(m.audit_trail().len(), 2, "parse + map must both run");
let doc = &m.context["data"]["DOC"];
for i in 0..k {
let field = doc.get(format!("f{i}"));
assert!(
field.map(|v| v.as_str().is_some()).unwrap_or(false),
"data.DOC.f{i} missing or not a string â mappings did not execute"
);
}
}
for _ in 0..n / 10 {
let mut m = Message::from_value(payload);
engine.process_message(&mut m).await.unwrap();
}
let start = Instant::now();
for _ in 0..n {
let mut m = Message::from_value(payload);
engine.process_message(&mut m).await.unwrap();
}
let el = start.elapsed();
let ns_msg = el.as_nanos() as f64 / n as f64;
let ns_mapping = ns_msg / k as f64;
println!(
" k={k:<4} {n:>7} msgs in {:>6.3}s => {ns_msg:>9.1} ns/msg => {ns_mapping:>7.1} ns/mapping",
el.as_secs_f64()
);
(ns_msg, ns_mapping)
}
#[tokio::main(flavor = "current_thread")]
async fn main() {
let payload = build_payload();
println!(
"Same-subtree write scaling: parse + k mappings into data.DOC.f*, capture_changes=on\n"
);
let mut per_mapping = Vec::new();
for &(k, n) in CONFIGS {
let (_, ns_mapping) = time_config(k, n, &payload).await;
per_mapping.push((k, ns_mapping));
}
let (k0, m0) = per_mapping[0];
let (kn, mn) = per_mapping[per_mapping.len() - 1];
println!(
"\nScaling check: ns/mapping at k={kn} vs k={k0} = {:.2}x \
(1.0x = linear write cost; the pre-write-through engine grows with k)",
mn / m0
);
}