const FLOW_DSL_SRC: &str = include_str!("flow_dsl.lua");
const BP_DSL_SRC: &str = include_str!("bp_dsl.lua");
const EMPTY_OBJECT_MARKER_KEY: &str = "__mse_empty_object__";
fn replace_empty_object_markers(value: &mut serde_json::Value) {
match value {
serde_json::Value::Object(map) => {
let is_marker = map.len() == 1
&& map.get(EMPTY_OBJECT_MARKER_KEY) == Some(&serde_json::Value::Bool(true));
if is_marker {
*value = serde_json::Value::Object(serde_json::Map::new());
return;
}
for v in map.values_mut() {
replace_empty_object_markers(v);
}
}
serde_json::Value::Array(arr) => {
for v in arr.iter_mut() {
replace_empty_object_markers(v);
}
}
_ => {}
}
}
pub fn preload(lua: &mlua::Lua) -> mlua::Result<()> {
let package: mlua::Table = lua.globals().get("package")?;
let preload: mlua::Table = package.get("preload")?;
preload.set(
"flow_dsl",
lua.create_function(|lua, ()| {
lua.load(FLOW_DSL_SRC)
.set_name("flow_dsl.lua")
.eval::<mlua::Value>()
})?,
)?;
preload.set(
"bp_dsl",
lua.create_function(|lua, ()| {
lua.load(BP_DSL_SRC)
.set_name("bp_dsl.lua")
.eval::<mlua::Value>()
})?,
)?;
Ok(())
}
pub fn build_bp_from_script(script: &str) -> anyhow::Result<serde_json::Value> {
use mlua::LuaSerdeExt;
let lua = mlua::Lua::new();
preload(&lua).map_err(|e| anyhow::anyhow!("dsl preload failed: {e}"))?;
let result: mlua::Value = lua
.load(script)
.set_name("<bp-script>")
.eval()
.map_err(|e| anyhow::anyhow!("bp-script eval failed: {e}"))?;
let options = mlua::serde::de::Options::new().encode_empty_tables_as_array(true);
let mut value: serde_json::Value = lua
.from_value_with(result, options)
.map_err(|e| anyhow::anyhow!("lua value -> json conversion failed: {e}"))?;
replace_empty_object_markers(&mut value);
Ok(value)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn preload_exposes_flow_dsl_and_bp_dsl() {
let lua = mlua::Lua::new();
preload(&lua).expect("preload must succeed");
let ok: bool = lua
.load(
r#"
local F = require("flow_dsl")
local B = require("bp_dsl")
return F ~= nil and B ~= nil
"#,
)
.eval()
.expect("require must succeed for both modules");
assert!(ok, "flow_dsl / bp_dsl must both resolve via require()");
}
#[test]
fn build_bp_from_script_returns_json_value() {
let out = build_bp_from_script(
r#"
local F = require("flow_dsl")
return { id = "t", flow = F.assign{ at = F.p("$.x"), value = F.lit(1) } }
"#,
)
.expect("script must build");
assert_eq!(out["id"], serde_json::json!("t"));
assert_eq!(out["flow"]["kind"], serde_json::json!("assign"));
assert_eq!(
out["flow"]["at"],
serde_json::json!({"op": "path", "at": "$.x"})
);
}
#[test]
fn build_bp_from_script_surfaces_lua_errors() {
let err = build_bp_from_script("error(\"boom\")").expect_err("must propagate the error");
assert!(err.to_string().contains("boom"));
}
#[test]
fn f_obj_marker_becomes_a_genuine_empty_json_object() {
let out = build_bp_from_script(
r#"
local F = require("flow_dsl")
return { spec = F.obj(), other = {} }
"#,
)
.expect("script must build");
assert_eq!(out["spec"], serde_json::json!({}));
assert!(
out["spec"].is_object(),
"F.obj() must become an object, not an array"
);
assert_eq!(out["other"], serde_json::json!([]));
}
#[test]
fn bp_dsl_skip_on_compiles_to_branch_in_verdict_skip_on_list() {
let out = build_bp_from_script(
r#"
local F = require("flow_dsl")
local B = require("bp_dsl")
return B.pipeline({
B.stage "gate" { agent = "mock-gate" },
B.stage "worker" {
agent = "mock-worker",
input = B.from "gate",
skip_on = { "SKIP", "NOT_APPLICABLE" },
},
halted_at = "$.halted_at",
})
"#,
)
.expect("skip_on pipeline must build");
assert_eq!(out["kind"], serde_json::json!("seq"));
let top_children = out["children"].as_array().expect("top seq children");
assert_eq!(top_children.len(), 2);
assert_eq!(top_children[0]["kind"], serde_json::json!("step"));
assert_eq!(top_children[0]["ref"], serde_json::json!("mock-gate"));
let rest = &top_children[1];
assert_eq!(rest["kind"], serde_json::json!("seq"));
let rest_children = rest["children"].as_array().expect("rest seq children");
let worker_guarded = &rest_children[0];
assert_eq!(worker_guarded["kind"], serde_json::json!("branch"));
let cond = &worker_guarded["cond"];
assert_eq!(cond["op"], serde_json::json!("in"));
assert_eq!(
cond["needle"],
serde_json::json!({"op": "path", "at": "$.gate.parts[\"verdict\"]"})
);
assert_eq!(cond["haystack"]["op"], serde_json::json!("lit"));
assert_eq!(
cond["haystack"]["value"],
serde_json::json!(["SKIP", "NOT_APPLICABLE"])
);
assert_eq!(
worker_guarded["then"],
serde_json::json!({"kind": "seq", "children": []})
);
let body = &worker_guarded["else"];
assert_eq!(body["kind"], serde_json::json!("seq"));
assert_eq!(body["children"][0]["ref"], serde_json::json!("mock-worker"));
}
#[test]
fn bp_dsl_skip_on_coexists_with_halt_on() {
let out = build_bp_from_script(
r#"
local F = require("flow_dsl")
local B = require("bp_dsl")
return B.pipeline({
B.stage "planner" { agent = "mock-planner" },
B.stage "worker" {
agent = "mock-worker",
input = B.from "planner",
skip_on = { "SKIP" },
halt_on = { "BLOCKED" },
},
B.stage "publisher" { agent = "mock-publisher" },
halted_at = "$.halted_at",
})
"#,
)
.expect("skip_on + halt_on pipeline must build");
let rest = &out["children"][1];
let worker_seq = rest;
assert_eq!(worker_seq["kind"], serde_json::json!("seq"));
let worker_children = worker_seq["children"]
.as_array()
.expect("worker seq children");
assert_eq!(
worker_children.len(),
2,
"skip guard + halt_on gate (with publisher threaded into gate else)"
);
let skip_branch = &worker_children[0];
assert_eq!(skip_branch["kind"], serde_json::json!("branch"));
assert_eq!(skip_branch["cond"]["op"], serde_json::json!("in"));
let halt_gate = &worker_children[1];
assert_eq!(halt_gate["kind"], serde_json::json!("branch"));
assert_eq!(halt_gate["cond"]["op"], serde_json::json!("eq"));
assert_eq!(
halt_gate["cond"]["lhs"],
serde_json::json!({"op": "path", "at": "$.worker.parts[\"verdict\"]"})
);
let gate_else = &halt_gate["else"];
assert_eq!(gate_else["kind"], serde_json::json!("seq"));
let contains_publisher = gate_else["children"]
.as_array()
.map(|arr| {
arr.iter()
.any(|c| c["ref"] == serde_json::json!("mock-publisher"))
})
.unwrap_or(false);
assert!(
contains_publisher,
"halt_on gate else must thread the publisher stage through: {gate_else}"
);
}
#[test]
fn bp_dsl_skip_on_empty_list_is_noop() {
let with_empty = build_bp_from_script(
r#"
local B = require("bp_dsl")
return B.pipeline({
B.stage "worker" { agent = "mock-worker", skip_on = {} },
halted_at = "$.halted_at",
})
"#,
)
.expect("skip_on={} pipeline must build");
let children = with_empty["children"].as_array().expect("seq children");
assert_eq!(
children.len(),
2,
"no skip guard emitted for empty skip_on: {with_empty}"
);
assert_eq!(children[0]["kind"], serde_json::json!("step"));
assert_eq!(children[0]["ref"], serde_json::json!("mock-worker"));
let baseline = build_bp_from_script(
r#"
local B = require("bp_dsl")
return B.pipeline({
B.stage "worker" { agent = "mock-worker" },
halted_at = "$.halted_at",
})
"#,
)
.expect("baseline pipeline must build");
assert_eq!(with_empty, baseline, "skip_on = {{}} must be a no-op");
}
#[test]
fn empty_object_marker_replacement_does_not_misfire_on_ordinary_data() {
let out = build_bp_from_script(
r#"
return {
a = { __mse_empty_object__ = false },
b = { __mse_empty_object__ = true, extra = 1 },
}
"#,
)
.expect("script must build");
assert_eq!(out["a"], serde_json::json!({"__mse_empty_object__": false}));
assert_eq!(
out["b"],
serde_json::json!({"__mse_empty_object__": true, "extra": 1})
);
}
}