#[cfg(feature = "ir")]
use crate::expr::evaluate as evaluate_expression;
use crate::expr::ExprFailure;
#[cfg(feature = "ir")]
use crate::plan::{run_plan, ExecutionPlanSpec, OpSpec, RelationKind};
use crate::plan::{ExecOutcome, PlanFailure};
use crate::value::Value;
#[cfg(feature = "ir")]
use serde_json::Value as J;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BehaviorFailureCode {
UnknownComponent,
UnknownNodeKind,
MapOverNotArray,
MapIntoElementNotObject,
MapBatchResultMismatch,
FanoutOverNotArray,
FanoutBatchResultMismatch,
UnknownEntry,
}
impl BehaviorFailureCode {
pub fn as_str(self) -> &'static str {
match self {
BehaviorFailureCode::UnknownComponent => "UNKNOWN_COMPONENT",
BehaviorFailureCode::UnknownNodeKind => "UNKNOWN_NODE_KIND",
BehaviorFailureCode::MapOverNotArray => "MAP_OVER_NOT_ARRAY",
BehaviorFailureCode::MapIntoElementNotObject => "MAP_INTO_ELEMENT_NOT_OBJECT",
BehaviorFailureCode::MapBatchResultMismatch => "MAP_BATCH_RESULT_MISMATCH",
BehaviorFailureCode::FanoutOverNotArray => "FANOUT_OVER_NOT_ARRAY",
BehaviorFailureCode::FanoutBatchResultMismatch => "FANOUT_BATCH_RESULT_MISMATCH",
BehaviorFailureCode::UnknownEntry => "UNKNOWN_ENTRY",
}
}
}
#[derive(Debug, Clone)]
pub struct BehaviorError {
code: String,
pub message: String,
pub detail: Option<Box<crate::plan::ErrorDetail>>,
}
impl BehaviorError {
pub fn new(code: impl Into<String>, message: impl Into<String>) -> Self {
BehaviorError {
code: code.into(),
message: message.into(),
detail: None,
}
}
pub fn code(&self) -> &str {
&self.code
}
pub fn detail(&self) -> Option<&crate::plan::ErrorDetail> {
self.detail.as_deref()
}
}
impl std::fmt::Display for BehaviorError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}: {}", self.code, self.message)
}
}
impl std::error::Error for BehaviorError {}
impl From<ExprFailure> for BehaviorError {
fn from(e: ExprFailure) -> Self {
BehaviorError {
code: e.code.as_str().to_string(),
message: e.message,
detail: e.detail,
}
}
}
impl From<PlanFailure> for BehaviorError {
fn from(e: PlanFailure) -> Self {
BehaviorError {
code: e.code.as_str().to_string(),
message: e.message,
detail: e.detail,
}
}
}
#[cfg(feature = "ir")]
fn bfail<T>(code: BehaviorFailureCode, message: impl Into<String>) -> Result<T, BehaviorError> {
Err(BehaviorError {
code: code.as_str().to_string(),
message: message.into(),
detail: None,
})
}
#[cfg(feature = "ir")]
fn wire_type_name(v: &Value) -> &'static str {
v.type_name()
}
#[cfg(feature = "ir")]
fn raw_value_of(v: &Value) -> String {
match v {
Value::Str(s) => s.clone(),
Value::Int(i) => i.to_string(),
Value::Float(f) => crate::canonical::py_float_repr(*f).unwrap_or_else(|_| f.to_string()),
Value::Bool(b) => b.to_string(),
Value::Null => "null".to_string(),
_ => crate::canonical::canonical_json(v).unwrap_or_else(|_| v.type_name().to_string()),
}
}
#[cfg(feature = "ir")]
fn portable_type_notation(t: &J) -> String {
if let Some(s) = t.as_str() {
return s.to_string();
}
for key in ["opt", "arr", "map"] {
if let Some(inner) = t.get(key) {
return format!("{key}({})", portable_type_notation(inner));
}
}
if let Some(J::Object(o)) = t.get("obj") {
let parts: Vec<String> = o
.iter()
.map(|(k, v)| format!("{k}:{}", portable_type_notation(v)))
.collect();
return format!("obj{{{}}}", parts.join(","));
}
"?".to_string()
}
#[cfg(feature = "ir")]
fn result_type_of(node: &J, out_type: &J) -> J {
if node.get("map").is_some() {
serde_json::json!({ "arr": out_type })
} else {
out_type.clone()
}
}
#[cfg(feature = "ir")]
fn conform_fail(
code: crate::expr::ExprFailureCode,
message: String,
kind: crate::plan::ErrorKind,
node_id: &str,
field: &str,
expected: &J,
actual: Option<&Value>,
) -> ExprFailure {
let mut detail = crate::plan::ErrorDetail {
kind: Some(kind),
model: Some(node_id.to_string()),
field: Some(field.to_string()),
expected_type: Some(portable_type_notation(expected)),
..Default::default()
};
if let Some(a) = actual {
detail.actual_wire_type = Some(wire_type_name(a).to_string());
detail.raw_value = Some(raw_value_of(a));
}
ExprFailure {
code,
message,
detail: Some(Box::new(detail)),
}
}
#[cfg(feature = "ir")]
fn obj_set(pairs: &mut Vec<(String, Value)>, key: &str, val: Value) {
if let Some(slot) = pairs.iter_mut().find(|(k, _)| k == key) {
slot.1 = val;
} else {
pairs.push((key.to_string(), val));
}
}
#[cfg(feature = "ir")]
fn assert_conforms_to_out_type(
node_id: &str,
v: &Value,
t: &J,
field: &str,
) -> Result<Value, ExprFailure> {
use crate::expr::ExprFailureCode::{MissingProp, TypeMismatch};
use crate::plan::ErrorKind;
if let Some(scalar) = t.as_str() {
if scalar == "float" {
if let Value::Int(i) = v {
return Ok(Value::Float(*i as f64)); }
}
let ok = match scalar {
"string" => matches!(v, Value::Str(_)),
"int" => matches!(v, Value::Int(_)),
"float" => matches!(v, Value::Float(_)),
"bool" => matches!(v, Value::Bool(_)),
"null" => matches!(v, Value::Null),
_ => true,
};
if !ok {
return Err(conform_fail(
TypeMismatch,
format!(
"node '{node_id}': {field}: expected {scalar}, got {}",
wire_type_name(v)
),
ErrorKind::TypeMismatch,
node_id,
field,
t,
Some(v),
));
}
return Ok(v.clone());
}
if let Some(inner) = t.get("opt") {
if matches!(v, Value::Null) {
return Ok(Value::Null);
}
return assert_conforms_to_out_type(node_id, v, inner, field);
}
if let Some(inner) = t.get("arr") {
match v {
Value::Arr(items) => {
let mut out = Vec::with_capacity(items.len());
for (i, el) in items.iter().enumerate() {
out.push(assert_conforms_to_out_type(
node_id,
el,
inner,
&format!("{field}[{i}]"),
)?);
}
return Ok(Value::Arr(out));
}
_ => {
return Err(conform_fail(
TypeMismatch,
format!(
"node '{node_id}': {field}: expected arr, got {}",
wire_type_name(v)
),
ErrorKind::TypeMismatch,
node_id,
field,
t,
Some(v),
))
}
}
}
if let Some(inner) = t.get("map") {
match v {
Value::Obj(pairs) => {
let mut out = Vec::with_capacity(pairs.len());
for (k, mv) in pairs {
out.push((
k.clone(),
assert_conforms_to_out_type(node_id, mv, inner, &format!("{field}.{k}"))?,
));
}
return Ok(Value::Obj(out));
}
_ => {
return Err(conform_fail(
TypeMismatch,
format!(
"node '{node_id}': {field}: expected map, got {}",
wire_type_name(v)
),
ErrorKind::TypeMismatch,
node_id,
field,
t,
Some(v),
))
}
}
}
if let Some(J::Object(fields)) = t.get("obj") {
let mut out: Vec<(String, Value)> = match v {
Value::Obj(pairs) => pairs.clone(),
_ => {
return Err(conform_fail(
TypeMismatch,
format!(
"node '{node_id}': {field}: expected obj, got {}",
wire_type_name(v)
),
ErrorKind::TypeMismatch,
node_id,
field,
t,
Some(v),
));
}
};
for (k, ft) in fields {
match v.obj_get(k) {
None => {
if ft.get("opt").is_some() {
obj_set(&mut out, k, Value::Null); } else {
return Err(conform_fail(
MissingProp,
format!("node '{node_id}': {field}: missing property .{k}"),
ErrorKind::MissingField,
node_id,
&format!("{field}.{k}"),
ft,
None,
));
}
}
Some(fv) => {
let nv = assert_conforms_to_out_type(node_id, fv, ft, &format!("{field}.{k}"))?;
obj_set(&mut out, k, nv);
}
}
}
return Ok(Value::Obj(out));
}
Ok(v.clone())
}
pub trait ComponentExec {
fn exec(
&mut self,
component: &str,
ports: &[(String, Value)],
bound: Option<&Value>,
) -> Option<ExecOutcome>;
fn exec_ctx(
&mut self,
node_id: &str,
component: &str,
ports: &[(String, Value)],
bound: Option<&Value>,
) -> Option<ExecOutcome> {
let _ = node_id;
self.exec(component, ports, bound)
}
}
impl ComponentExec for &mut dyn ComponentExec {
fn exec(
&mut self,
component: &str,
ports: &[(String, Value)],
bound: Option<&Value>,
) -> Option<ExecOutcome> {
(**self).exec(component, ports, bound)
}
fn exec_ctx(
&mut self,
node_id: &str,
component: &str,
ports: &[(String, Value)],
bound: Option<&Value>,
) -> Option<ExecOutcome> {
(**self).exec_ctx(node_id, component, ports, bound)
}
}
#[cfg(feature = "ir")]
fn scope_get<'a>(scope: &'a [(String, Value)], key: &str) -> Option<&'a Value> {
scope.iter().find(|(k, _)| k == key).map(|(_, v)| v)
}
#[cfg(feature = "ir")]
fn fanout_dedup_drop(
aligned_bodies: Vec<Value>,
dedupe_key: &str,
drop: &str,
implicit_source: Option<&str>,
) -> Vec<Value> {
use std::collections::HashSet;
let mut items: Vec<Value> = Vec::with_capacity(aligned_bodies.len());
let mut seen: HashSet<String> = HashSet::new();
for body in aligned_bodies.into_iter() {
let fields = match &body {
Value::Obj(f) => Some(f),
_ => None,
};
let key_val = fields.and_then(|f| f.iter().find(|(k, _)| k == dedupe_key).map(|(_, v)| v));
let has_key = matches!(key_val, Some(v) if !matches!(v, Value::Null));
if !has_key {
if drop == "dangling" {
continue;
}
items.push(body); continue;
}
let seen_key = match key_val.unwrap() {
Value::Str(s) => format!("s:{s}"),
other => format!("j:{other:?}"),
};
if seen.contains(&seen_key) {
continue;
}
seen.insert(seen_key);
match (implicit_source, fields) {
(Some(src), Some(f)) if f.iter().any(|(k, _)| k == src) => {
let stripped: Vec<(String, Value)> =
f.iter().filter(|(k, _)| k != src).cloned().collect();
items.push(Value::Obj(stripped));
}
_ => items.push(body),
}
}
items
}
#[cfg(feature = "ir")]
fn node_kind(n: &J) -> Result<&'static str, BehaviorError> {
if n.get("fanout").is_some() {
Ok("fanout")
} else if n.get("map").is_some() {
Ok("map")
} else if n.get("cond").is_some() {
Ok("cond")
} else if n.get("component").is_some() {
Ok("componentRef")
} else {
bfail(
BehaviorFailureCode::UnknownNodeKind,
format!(
"body node '{}' is not componentRef/map/cond/fanout",
n.get("id").and_then(|v| v.as_str()).unwrap_or("?")
),
)
}
}
#[cfg(feature = "ir")]
fn node_sub(n: &J) -> &J {
if let Some(f) = n.get("fanout") {
f
} else if let Some(m) = n.get("map") {
m
} else if let Some(c) = n.get("cond") {
c
} else {
n
}
}
#[cfg(feature = "ir")]
fn node_parent(n: &J) -> Option<&str> {
node_sub(n).get("parent").and_then(|v| v.as_str())
}
#[cfg(feature = "ir")]
fn node_bind_field(n: &J) -> Option<String> {
if n.get("map").is_some() || n.get("cond").is_some() || n.get("fanout").is_some() {
return None;
}
n.get("bindField")
.and_then(|v| v.as_str())
.map(str::to_string)
}
#[cfg(feature = "ir")]
fn node_relation_kind(n: &J) -> Option<RelationKind> {
if n.get("cond").is_some() {
return None;
}
let sub = if n.get("map").is_some() || n.get("fanout").is_some() {
node_sub(n)
} else {
n
};
match sub.get("relationKind").and_then(|v| v.as_str()) {
Some("connection") => Some(RelationKind::Connection),
Some(_) => Some(RelationKind::Single),
None => None,
}
}
#[cfg(feature = "ir")]
fn node_relation_kind_str(n: &J) -> Option<&str> {
if n.get("cond").is_some() {
return None;
}
let sub = if n.get("map").is_some() || n.get("fanout").is_some() {
node_sub(n)
} else {
n
};
sub.get("relationKind").and_then(|v| v.as_str())
}
#[cfg(feature = "ir")]
fn node_policy(n: &J) -> Option<String> {
if n.get("cond").is_some() {
return None;
}
let sub = if n.get("map").is_some() || n.get("fanout").is_some() {
node_sub(n)
} else {
n
};
sub.get("policy")
.and_then(|v| v.as_str())
.map(str::to_string)
}
#[cfg(feature = "ir")]
fn eval_ports(ports: &J, scope: &[(String, Value)]) -> Result<Vec<(String, Value)>, BehaviorError> {
let obj = ports.as_object().ok_or_else(|| BehaviorError {
code: BehaviorFailureCode::UnknownNodeKind.as_str().to_string(),
message: "ports must be an object".into(),
detail: None,
})?;
let mut out = Vec::with_capacity(obj.len());
for (k, v) in obj {
out.push((k.clone(), evaluate_expression(v, scope)?));
}
Ok(out)
}
#[cfg(feature = "ir")]
fn parse_plan(plan: Option<&J>) -> Option<ExecutionPlanSpec> {
let p = plan?;
if p.is_null() {
return None;
}
let groups = p
.get("groups")?
.as_array()?
.iter()
.map(|g| {
g.as_array()
.map(|row| {
row.iter()
.filter_map(|i| i.as_u64().map(|x| x as usize))
.collect()
})
.unwrap_or_default()
})
.collect();
let concurrency = p.get("concurrency").and_then(|c| c.as_i64()).unwrap_or(1);
Some(ExecutionPlanSpec {
groups,
concurrency,
})
}
#[cfg(feature = "ir")]
pub fn run_behavior(
ir: &J,
handlers: &mut dyn ComponentExec,
input: &[(String, Value)],
entry: Option<&str>,
) -> Result<Value, BehaviorError> {
let components = ir
.get("components")
.and_then(|c| c.as_array())
.ok_or_else(|| BehaviorError {
code: BehaviorFailureCode::UnknownEntry.as_str().to_string(),
message: "IR.components must be an array".into(),
detail: None,
})?;
let comp = match entry {
Some(name) => components
.iter()
.find(|c| c.get("name").and_then(|n| n.as_str()) == Some(name)),
None => components.first(),
};
let comp = match comp {
Some(c) => c,
None => {
return bfail(
BehaviorFailureCode::UnknownEntry,
format!("component '{}' not found in IR", entry.unwrap_or("<first>")),
)
}
};
let body: Vec<J> = comp
.get("body")
.and_then(|b| b.as_array())
.cloned()
.unwrap_or_default();
let index_of = |id: &str| {
body.iter()
.position(|n| n.get("id").and_then(|v| v.as_str()) == Some(id))
};
let mut ops: Vec<OpSpec> = Vec::with_capacity(body.len());
for n in &body {
let parent = node_parent(n).and_then(index_of);
ops.push(OpSpec {
id: n
.get("id")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string(),
parent,
bind_field: node_bind_field(n),
relation_kind: node_relation_kind(n),
policy: node_policy(n),
});
}
run_behavior_inner(
&body,
&ops,
parse_plan(comp.get("plan")),
comp.get("output").cloned().unwrap_or(J::Null),
&bind_input_scope(comp, input),
handlers,
)
}
#[cfg(feature = "ir")]
fn bind_input_scope(comp: &J, input: &[(String, Value)]) -> Vec<(String, Value)> {
let mut scope: Vec<(String, Value)> = input.to_vec();
if let Some(ports) = comp.get("inputPorts").and_then(|p| p.as_object()) {
for (name, schema) in ports {
let optional = schema.get("required").and_then(|r| r.as_bool()) == Some(false);
if optional && !scope.iter().any(|(k, _)| k == name) {
scope.push((name.clone(), Value::Null));
}
}
}
scope
}
#[cfg(feature = "ir")]
fn run_behavior_inner(
body: &[J],
ops: &[OpSpec],
plan: Option<ExecutionPlanSpec>,
output: J,
input: &[(String, Value)],
handlers: &mut dyn ComponentExec,
) -> Result<Value, BehaviorError> {
let index_of = |id: &str| {
body.iter()
.position(|n| n.get("id").and_then(|v| v.as_str()) == Some(id))
};
let mut results: Vec<(String, Value)> = Vec::new();
let mut pending_err: Option<BehaviorError> = None;
let run = {
let results_cell = &mut results;
let err_cell = &mut pending_err;
let exec = |op: &OpSpec, _bound: Option<&Value>| -> ExecOutcome {
if err_cell.is_some() {
return ExecOutcome::error("aborted");
}
let idx = match index_of(&op.id) {
Some(i) => i,
None => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::UnknownNodeKind.as_str().to_string(),
message: format!("no body node for op '{}'", op.id),
detail: None,
});
return ExecOutcome::error("aborted");
}
};
let node = &body[idx];
let base_scope = |extra: Option<(&str, &Value)>| -> Vec<(String, Value)> {
let mut s: Vec<(String, Value)> = input.to_vec();
for (k, v) in results_cell.iter() {
s.push((k.clone(), v.clone()));
}
if let Some((k, v)) = extra {
s.push((k.to_string(), v.clone()));
}
s
};
let outcome = match node_kind(node) {
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
Ok("cond") => {
let c = node.get("cond").unwrap();
let cond_expr = serde_json::json!({
"cond": [c.get("if"), c.get("then"), c.get("else")]
});
match evaluate_expression(&cond_expr, &base_scope(None)) {
Ok(v) => ExecOutcome::Ok(v),
Err(e) => {
*err_cell = Some(e.into());
return ExecOutcome::error("aborted");
}
}
}
Ok("map") => {
let m = node.get("map").unwrap();
let over = match evaluate_expression(
m.get("over").unwrap_or(&J::Null),
&base_scope(None),
) {
Ok(v) => v,
Err(e) => {
*err_cell = Some(e.into());
return ExecOutcome::error("aborted");
}
};
let arr = match &over {
Value::Arr(a) => a.clone(),
_ => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::MapOverNotArray.as_str().to_string(),
message: format!("map '{}': 'over' is not an array", op.id),
detail: None,
});
return ExecOutcome::error("aborted");
}
};
let component = m.get("component").and_then(|v| v.as_str()).unwrap_or("");
let as_name = m.get("as").and_then(|v| v.as_str()).unwrap_or("$");
let ports_j = m
.get("ports")
.cloned()
.unwrap_or(J::Object(Default::default()));
let when = m.get("when");
let batched = m.get("batched").and_then(|v| v.as_bool()).unwrap_or(false);
let element_policy = {
let raw = m
.get("elementPolicy")
.and_then(|v| v.as_str())
.unwrap_or("error");
match crate::plan::ElementPolicyKind::parse(raw) {
None => {
*err_cell = Some(BehaviorError {
code: "UNKNOWN_ELEMENT_POLICY".to_string(),
message: format!(
"map '{}': unknown element policy '{raw}' (fail-closed)",
op.id
),
detail: None,
});
return ExecOutcome::error("aborted");
}
Some(crate::plan::ElementPolicyKind::Skip) if batched => {
*err_cell = Some(BehaviorError {
code: "ELEMENT_POLICY_NOT_APPLICABLE".to_string(),
message: format!("map '{}': elementPolicy 'skip' needs a per-element Failure, but a batched map takes ONE outcome for the whole batch (fail-closed)", op.id),
detail: None,
});
return ExecOutcome::error("aborted");
}
Some(p) => p,
}
};
let into = m.get("into").and_then(|v| v.as_str());
let keep = |scope: &[(String, Value)]| -> Result<bool, BehaviorError> {
match when {
None => Ok(true),
Some(w) => {
let cond_expr = serde_json::json!({ "cond": [w, true, false] });
Ok(matches!(
evaluate_expression(&cond_expr, scope)?,
Value::Bool(true)
))
}
}
};
let mut kept_idx: Vec<usize> = Vec::new(); let collected: Vec<Value>;
if batched {
let mut items: Vec<Value> = Vec::new();
for (i, el) in arr.iter().enumerate() {
let scope = base_scope(Some((as_name, el)));
match keep(&scope) {
Ok(false) => continue,
Ok(true) => {}
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
}
match eval_ports(&ports_j, &scope) {
Ok(p) => items.push(Value::Obj(p)),
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
}
kept_idx.push(i);
}
if items.is_empty() {
collected = Vec::new(); } else {
let want = items.len();
let batch_ports = vec![("items".to_string(), Value::Arr(items))];
match handlers.exec_ctx(&op.id, component, &batch_ports, None) {
None => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::UnknownComponent
.as_str()
.to_string(),
message: format!(
"component '{component}' has no handler (fail-closed)"
),
detail: None,
});
return ExecOutcome::error("aborted");
}
Some(ExecOutcome::Error(e, d)) => return ExecOutcome::Error(e, d),
Some(ExecOutcome::Ok(Value::Arr(r))) if r.len() == want => {
collected = r;
}
Some(ExecOutcome::Ok(_)) => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::MapBatchResultMismatch
.as_str()
.to_string(),
message: format!(
"map '{}': batched handler must return a list aligned to items (want {want})",
op.id
),
detail: None,
});
return ExecOutcome::error("aborted");
}
}
}
} else {
let mut out: Vec<Value> = Vec::with_capacity(arr.len());
for (i, el) in arr.iter().enumerate() {
let scope = base_scope(Some((as_name, el)));
match keep(&scope) {
Ok(false) => continue,
Ok(true) => {}
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
}
let ports = match eval_ports(&ports_j, &scope) {
Ok(p) => p,
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
};
match handlers.exec_ctx(&op.id, component, &ports, Some(el)) {
None => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::UnknownComponent
.as_str()
.to_string(),
message: format!(
"component '{component}' has no handler (fail-closed)"
),
detail: None,
});
return ExecOutcome::error("aborted");
}
Some(ExecOutcome::Error(e, d)) => {
if element_policy == crate::plan::ElementPolicyKind::Skip {
continue;
}
return ExecOutcome::Error(e, d);
}
Some(ExecOutcome::Ok(v)) => out.push(v),
}
kept_idx.push(i);
}
collected = out;
}
match into {
None => ExecOutcome::Ok(Value::Arr(collected)),
Some(key) => {
let mut augmented: Vec<Value> = Vec::with_capacity(arr.len());
let mut k = 0usize;
for (i, el) in arr.iter().enumerate() {
if k < kept_idx.len() && kept_idx[k] == i {
let fields = match el {
Value::Obj(f) => f.clone(),
_ => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::MapIntoElementNotObject
.as_str()
.to_string(),
message: format!(
"map '{}': 'into' requires object elements (element {i} is not an object)",
op.id
),
detail: None,
});
return ExecOutcome::error("aborted");
}
};
let mut fields = fields;
let val = collected[k].clone();
match fields.iter_mut().find(|(fk, _)| fk == key) {
Some(slot) => slot.1 = val,
None => fields.push((key.to_string(), val)),
}
augmented.push(Value::Obj(fields));
k += 1;
} else {
augmented.push(el.clone());
}
}
ExecOutcome::Ok(Value::Arr(augmented))
}
}
}
Ok("fanout") => {
let f = node.get("fanout").unwrap();
let over = match evaluate_expression(
f.get("over").unwrap_or(&J::Null),
&base_scope(None),
) {
Ok(v) => v,
Err(e) => {
*err_cell = Some(e.into());
return ExecOutcome::error("aborted");
}
};
let arr = match &over {
Value::Arr(a) => a.clone(),
_ => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::FanoutOverNotArray.as_str().to_string(),
message: format!("fanout '{}': 'over' is not an array", op.id),
detail: None,
});
return ExecOutcome::error("aborted");
}
};
let component = f.get("component").and_then(|v| v.as_str()).unwrap_or("");
let as_name = f.get("as").and_then(|v| v.as_str()).unwrap_or("$");
let ports_j = f
.get("ports")
.cloned()
.unwrap_or(J::Object(Default::default()));
let dedupe_key = f.get("dedupeKey").and_then(|v| v.as_str()).unwrap_or("");
let drop = f.get("drop").and_then(|v| v.as_str()).unwrap_or("dangling");
let implicit_source = f.get("implicitSource").and_then(|v| v.as_str());
let mut items: Vec<Value> = Vec::with_capacity(arr.len());
for el in arr.iter() {
let scope = base_scope(Some((as_name, el)));
match eval_ports(&ports_j, &scope) {
Ok(p) => items.push(Value::Obj(p)),
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
}
}
if items.is_empty() {
ExecOutcome::Ok(Value::Obj(vec![
("items".into(), Value::Arr(vec![])),
("cursor".into(), Value::Null),
]))
} else {
let want = items.len();
let batch_ports = vec![("items".to_string(), Value::Arr(items))];
match handlers.exec_ctx(&op.id, component, &batch_ports, None) {
None => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::UnknownComponent
.as_str()
.to_string(),
message: format!(
"component '{component}' has no handler (fail-closed)"
),
detail: None,
});
ExecOutcome::error("aborted")
}
Some(ExecOutcome::Error(e, d)) => ExecOutcome::Error(e, d),
Some(ExecOutcome::Ok(Value::Arr(r))) if r.len() == want => {
let deduped =
fanout_dedup_drop(r, dedupe_key, drop, implicit_source);
ExecOutcome::Ok(Value::Obj(vec![
("items".into(), Value::Arr(deduped)),
("cursor".into(), Value::Null),
]))
}
Some(ExecOutcome::Ok(_)) => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::FanoutBatchResultMismatch
.as_str()
.to_string(),
message: format!(
"fanout '{}': batched handler must return a list aligned to the deduped id list (want {want})",
op.id
),
detail: None,
});
ExecOutcome::error("aborted")
}
}
}
}
Ok(_) => {
let component = node.get("component").and_then(|v| v.as_str()).unwrap_or("");
let ports_j = node
.get("ports")
.cloned()
.unwrap_or(J::Object(Default::default()));
let ports = match eval_ports(&ports_j, &base_scope(None)) {
Ok(p) => p,
Err(e) => {
*err_cell = Some(e);
return ExecOutcome::error("aborted");
}
};
match handlers.exec_ctx(&op.id, component, &ports, None) {
None => {
*err_cell = Some(BehaviorError {
code: BehaviorFailureCode::UnknownComponent.as_str().to_string(),
message: format!(
"component '{component}' has no handler (fail-closed)"
),
detail: None,
});
return ExecOutcome::error("aborted");
}
Some(o) => o,
}
}
};
if let ExecOutcome::Ok(v) = &outcome {
let mut value = v.clone();
if let Some(ot) = node.get("outType") {
let want = result_type_of(node, ot);
match assert_conforms_to_out_type(&op.id, &value, &want, "result") {
Ok(normalized) => value = normalized,
Err(e) => {
*err_cell = Some(e.into());
return ExecOutcome::error("aborted");
}
}
}
results_cell.push((op.id.clone(), value));
}
outcome
};
run_plan(plan.as_ref(), ops, exec)
};
if let Some(e) = pending_err {
return Err(e);
}
let run = run?;
for (i, id_val) in run.skipped.iter().enumerate() {
let _ = i;
if let Some(idx) = index_of(id_val) {
let rk = node_relation_kind_str(&body[idx]);
let unproduced = if rk == Some("connection") {
Value::Obj(vec![
("items".into(), Value::Arr(vec![])),
("cursor".into(), Value::Null),
])
} else {
Value::Null
};
if scope_get(&results, id_val).is_none() {
results.push((id_val.clone(), unproduced));
}
}
}
let scope: Vec<(String, Value)> = {
let mut s: Vec<(String, Value)> = input.to_vec();
s.extend(results.iter().cloned());
s
};
Ok(evaluate_expression(&output, &scope)?)
}