Skip to main content

sim_lib_topology/
verb.rs

1//! Built-in topology verb execution.
2
3use sim_kernel::{Cx, Error, Expr, Result, Symbol};
4
5use crate::{
6    CompiledGraph, Graph,
7    adapter::{call_target_expr, resolve_target},
8    run::{
9        BudgetLedger, TopologyCells, TopologyNonlinearState, TopologyPacket, WorkItem,
10        predicate_accepts,
11    },
12};
13
14/// Result of executing one core topology node.
15#[derive(Clone, Debug, PartialEq, Eq)]
16pub enum VerbAction {
17    /// Emit a packet from a named output port.
18    Emit(TopologyPacket),
19    /// Complete the graph with one public output expression.
20    Complete {
21        /// Index of the completing node.
22        node_index: usize,
23        /// The public output expression.
24        expr: Expr,
25    },
26}
27
28/// Runs one built-in topology verb.
29pub fn run_core_node(
30    cx: &mut Cx,
31    graph: &Graph,
32    plan: &CompiledGraph,
33    budget: &mut BudgetLedger,
34    cells: &mut TopologyCells,
35    nonlinear: &mut TopologyNonlinearState,
36    item: &WorkItem,
37) -> Result<Vec<VerbAction>> {
38    let node = &graph.nodes[plan.nodes[item.node_index].source_index];
39    match node.verb.name.as_ref() {
40        "in" => emit(item.node_index, "out", item.expr.clone()),
41        "out" => Ok(vec![VerbAction::Complete {
42            node_index: item.node_index,
43            expr: item.expr.clone(),
44        }]),
45        "wire" => emit(item.node_index, "out", item.expr.clone()),
46        "tee" => emit(item.node_index, "out", item.expr.clone()),
47        "branch" => {
48            let accepted = match &item.expr {
49                Expr::Bool(value) => *value,
50                _ => {
51                    let predicate = node_option(node, "when").ok_or_else(|| {
52                        Error::Eval(format!(
53                            "topology run: branch node {} requires a when predicate for non-bool input",
54                            node.id.as_symbol()
55                        ))
56                    })?;
57                    predicate_accepts(cx, predicate, &item.expr)?
58                }
59            };
60            let desired = if accepted { "true" } else { "false" };
61            let port = if has_route(plan, item.node_index, desired) {
62                desired
63            } else if has_route(plan, item.node_index, "else") {
64                "else"
65            } else {
66                desired
67            };
68            emit(item.node_index, port, item.expr.clone())
69        }
70        "cell" => run_cell_node(cx, cells, node, item),
71        "merge" => run_merge_node(plan, budget, nonlinear, node, item),
72        "race" => run_race_node(cx, nonlinear, node, item),
73        "quorum" => run_quorum_node(cx, nonlinear, node, item),
74        "reduce" => run_reduce_node(cx, plan, nonlinear, node, item),
75        "patch" => run_patch_node(cx, node, item),
76        "call" => {
77            let target = node.target.as_ref().ok_or_else(|| {
78                Error::Eval(format!(
79                    "topology run: call node {} has no target",
80                    node.id.as_symbol()
81                ))
82            })?;
83            let target = resolve_target(cx, target)?;
84            if target.object().as_eval_fabric().is_some() {
85                budget.record_child_run()?;
86            }
87            let output = call_target_expr(cx, target, item.expr.clone())?;
88            emit(item.node_index, "out", output)
89        }
90        other => Err(Error::Eval(format!(
91            "topology run: unsupported core verb {other}"
92        ))),
93    }
94}
95
96fn run_merge_node(
97    plan: &CompiledGraph,
98    budget: &BudgetLedger,
99    nonlinear: &mut TopologyNonlinearState,
100    node: &crate::Node,
101    item: &WorkItem,
102) -> Result<Vec<VerbAction>> {
103    let mode = option_symbol(node, "mode")?.unwrap_or_else(|| Symbol::new("all"));
104    match mode.name.as_ref() {
105        "any" => {
106            if nonlinear.mark_merge_any_complete(item.node_index) {
107                emit(item.node_index, "out", item.expr.clone())
108            } else {
109                Ok(Vec::new())
110            }
111        }
112        "latest" => {
113            let output =
114                nonlinear.update_latest(item.node_index, item.port.clone(), item.expr.clone());
115            emit(item.node_index, "out", output)
116        }
117        "all" => {
118            let buffered =
119                nonlinear.push_merge(item.node_index, item.port.clone(), item.expr.clone());
120            budget.check_merge_buffer(buffered)?;
121            let required = incoming_count(plan, item.node_index);
122            if buffered >= required {
123                emit(
124                    item.node_index,
125                    "out",
126                    Expr::List(nonlinear.drain_merge(item.node_index, required)),
127                )
128            } else {
129                Ok(Vec::new())
130            }
131        }
132        "count" => {
133            let required = option_u32(node, "count")?.unwrap_or(1).max(1) as usize;
134            let buffered =
135                nonlinear.push_merge(item.node_index, item.port.clone(), item.expr.clone());
136            budget.check_merge_buffer(buffered)?;
137            if buffered >= required {
138                emit(
139                    item.node_index,
140                    "out",
141                    Expr::List(nonlinear.drain_merge(item.node_index, required)),
142                )
143            } else {
144                Ok(Vec::new())
145            }
146        }
147        other => Err(Error::Eval(format!(
148            "topology run: unsupported merge mode {other}"
149        ))),
150    }
151}
152
153fn run_race_node(
154    cx: &mut Cx,
155    nonlinear: &mut TopologyNonlinearState,
156    node: &crate::Node,
157    item: &WorkItem,
158) -> Result<Vec<VerbAction>> {
159    let accepted = match node_option(node, "accept") {
160        Some(predicate) => predicate_accepts(cx, predicate, &item.expr)?,
161        None => true,
162    };
163    if accepted && nonlinear.mark_race_complete(item.node_index) {
164        emit(item.node_index, "out", item.expr.clone())
165    } else {
166        Ok(Vec::new())
167    }
168}
169
170fn run_quorum_node(
171    cx: &mut Cx,
172    nonlinear: &mut TopologyNonlinearState,
173    node: &crate::Node,
174    item: &WorkItem,
175) -> Result<Vec<VerbAction>> {
176    let required = option_u32(node, "n")?.unwrap_or(2).max(1);
177    let key = match option_target(node, "key")? {
178        Some(target) => {
179            let target = resolve_target(cx, &target)?;
180            call_target_expr(cx, target, item.expr.clone())?
181        }
182        None => item.expr.clone(),
183    };
184    let (value, count) = nonlinear.record_quorum(item.node_index, key, item.expr.clone());
185    if count >= required && nonlinear.mark_quorum_complete(item.node_index) {
186        emit(item.node_index, "out", value)
187    } else {
188        Ok(Vec::new())
189    }
190}
191
192fn run_reduce_node(
193    cx: &mut Cx,
194    plan: &CompiledGraph,
195    nonlinear: &mut TopologyNonlinearState,
196    node: &crate::Node,
197    item: &WorkItem,
198) -> Result<Vec<VerbAction>> {
199    let initial = node_option(node, "initial").cloned().unwrap_or(Expr::Nil);
200    let current = nonlinear.reduce_current(item.node_index, initial);
201    let next = match option_target(node, "target")? {
202        Some(target) => {
203            let target = resolve_target(cx, &target)?;
204            call_target_expr(cx, target, Expr::List(vec![current, item.expr.clone()]))?
205        }
206        None => match current {
207            Expr::Nil => item.expr.clone(),
208            other => Expr::List(vec![other, item.expr.clone()]),
209        },
210    };
211    let count = nonlinear.record_reduce(item.node_index, next.clone());
212    let emit_partials = option_bool(node, "emit_partials")?.unwrap_or(false);
213    if emit_partials || count >= incoming_count(plan, item.node_index) {
214        if !emit_partials {
215            nonlinear.reset_reduce(item.node_index);
216        }
217        emit(item.node_index, "out", next)
218    } else {
219        Ok(Vec::new())
220    }
221}
222
223fn run_cell_node(
224    cx: &mut Cx,
225    cells: &mut TopologyCells,
226    node: &crate::Node,
227    item: &WorkItem,
228) -> Result<Vec<VerbAction>> {
229    let name = option_symbol(node, "name")?.ok_or_else(|| {
230        Error::Eval(format!(
231            "topology run: cell node {} requires name",
232            node.id.as_symbol()
233        ))
234    })?;
235    let op = option_symbol(node, "op")?.unwrap_or_else(|| Symbol::new("read"));
236    let cell_value = match op.name.as_ref() {
237        "read" => cells.read(&name)?,
238        "write" => cells.write(cx, &name, item.expr.clone())?,
239        "append" => cells.append(cx, &name, item.expr.clone())?,
240        "merge" => cells.merge(cx, &name, item.expr.clone())?,
241        "clear" => cells.clear(cx, &name)?,
242        other => {
243            return Err(Error::Eval(format!(
244                "topology run: unsupported cell op {other}"
245            )));
246        }
247    };
248    let default_emit = if op.name.as_ref() == "read" {
249        Symbol::new("cell")
250    } else {
251        Symbol::new("input")
252    };
253    let emit_mode = option_symbol(node, "emit")?.unwrap_or(default_emit);
254    let output = match emit_mode.name.as_ref() {
255        "input" => item.expr.clone(),
256        "cell" => cell_value,
257        "both" => Expr::Map(vec![
258            (Expr::Symbol(Symbol::new("input")), item.expr.clone()),
259            (Expr::Symbol(Symbol::new("cell")), cell_value),
260        ]),
261        other => {
262            return Err(Error::Eval(format!(
263                "topology run: unsupported cell emit mode {other}"
264            )));
265        }
266    };
267    emit(item.node_index, "out", output)
268}
269
270fn run_patch_node(cx: &mut Cx, node: &crate::Node, item: &WorkItem) -> Result<Vec<VerbAction>> {
271    let mode = option_symbol(node, "mode")?.unwrap_or_else(|| Symbol::new("produce"));
272    match mode.name.as_ref() {
273        "produce" => {
274            let proposal = match option_target(node, "target")? {
275                Some(target) => {
276                    let target = resolve_target(cx, &target)?;
277                    call_target_expr(cx, target, item.expr.clone())?
278                }
279                None => node_option(node, "patch")
280                    .cloned()
281                    .unwrap_or_else(|| item.expr.clone()),
282            };
283            crate::TopologyPatch::from_expr(&proposal)?;
284            emit(item.node_index, "out", proposal)
285        }
286        "apply" => Err(Error::Eval(
287            "topology run: patch apply must use topology/patch".to_owned(),
288        )),
289        other => Err(Error::Eval(format!(
290            "topology run: unsupported patch mode {other}"
291        ))),
292    }
293}
294
295fn emit(node_index: usize, port: &str, expr: Expr) -> Result<Vec<VerbAction>> {
296    Ok(vec![VerbAction::Emit(TopologyPacket {
297        node_index,
298        port: Symbol::new(port),
299        expr,
300    })])
301}
302
303fn has_route(plan: &CompiledGraph, node_index: usize, port: &str) -> bool {
304    plan.outgoing_edges[node_index]
305        .iter()
306        .any(|edge_index| plan.edges[*edge_index].from.port.name.as_ref() == port)
307}
308
309fn incoming_count(plan: &CompiledGraph, node_index: usize) -> usize {
310    plan.incoming_edges[node_index].len().max(1)
311}
312
313fn node_option<'a>(node: &'a crate::Node, key: &str) -> Option<&'a Expr> {
314    node.options
315        .iter()
316        .find(|(name, _)| name.namespace.is_none() && name.name.as_ref() == key)
317        .map(|(_, value)| value)
318}
319
320fn option_symbol(node: &crate::Node, key: &str) -> Result<Option<Symbol>> {
321    let Some(value) = node_option(node, key) else {
322        return Ok(None);
323    };
324    match value {
325        Expr::Symbol(symbol) => Ok(Some(symbol.clone())),
326        Expr::String(text) => Ok(Some(Symbol::new(text.clone()))),
327        other => Err(Error::Eval(format!(
328            "topology run: node option {key} expects symbol or string, got {other:?}"
329        ))),
330    }
331}
332
333fn option_target(node: &crate::Node, key: &str) -> Result<Option<Expr>> {
334    match node_option(node, key) {
335        Some(value) => Ok(Some(value.clone())),
336        None if key == "target" => Ok(node.target.clone()),
337        None => Ok(None),
338    }
339}
340
341fn option_bool(node: &crate::Node, key: &str) -> Result<Option<bool>> {
342    let Some(value) = node_option(node, key) else {
343        return Ok(None);
344    };
345    match value {
346        Expr::Bool(value) => Ok(Some(*value)),
347        Expr::Symbol(symbol) if symbol.namespace.is_none() && symbol.name.as_ref() == "true" => {
348            Ok(Some(true))
349        }
350        Expr::Symbol(symbol) if symbol.namespace.is_none() && symbol.name.as_ref() == "false" => {
351            Ok(Some(false))
352        }
353        other => Err(Error::Eval(format!(
354            "topology run: node option {key} expects bool, got {other:?}"
355        ))),
356    }
357}
358
359fn option_u32(node: &crate::Node, key: &str) -> Result<Option<u32>> {
360    let Some(value) = node_option(node, key) else {
361        return Ok(None);
362    };
363    match value {
364        Expr::Number(number) => number.canonical.parse::<u32>().map(Some).map_err(|_| {
365            Error::Eval(format!(
366                "topology run: node option {key} expects u32, got {}",
367                number.canonical
368            ))
369        }),
370        Expr::String(text) => text.parse::<u32>().map(Some).map_err(|_| {
371            Error::Eval(format!(
372                "topology run: node option {key} expects u32, got {text}"
373            ))
374        }),
375        other => Err(Error::Eval(format!(
376            "topology run: node option {key} expects u32, got {other:?}"
377        ))),
378    }
379}