1use 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#[derive(Clone, Debug, PartialEq, Eq)]
16pub enum VerbAction {
17 Emit(TopologyPacket),
19 Complete {
21 node_index: usize,
23 expr: Expr,
25 },
26}
27
28pub 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}