Skip to content

Commit e901e9d

Browse files
committed
feat: added subflow handling to execution results
1 parent 573f84e commit e901e9d

7 files changed

Lines changed: 358 additions & 66 deletions

File tree

crates/taurus-core/src/runtime/engine.rs

Lines changed: 88 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -411,6 +411,20 @@ mod tests {
411411
}
412412
}
413413

414+
fn assert_node_result_id(result: &NodeExecutionResult, expected_id: i64) {
415+
assert_eq!(
416+
result.id,
417+
Some(node_execution_result::Id::NodeId(expected_id))
418+
);
419+
}
420+
421+
fn assert_function_result_id(result: &NodeExecutionResult, expected_id: i64) {
422+
assert_eq!(
423+
result.id,
424+
Some(node_execution_result::Id::FunctionId(expected_id))
425+
);
426+
}
427+
414428
fn sleep_handler(
415429
_args: &[Argument],
416430
_ctx: &mut ValueStore,
@@ -420,6 +434,21 @@ mod tests {
420434
Signal::Success(null_value())
421435
}
422436

437+
fn echo_first_arg_handler(
438+
args: &[Argument],
439+
_ctx: &mut ValueStore,
440+
_run: &mut ThunkRunner<'_>,
441+
) -> Signal {
442+
match args.first() {
443+
Some(Argument::Eval(value)) => Signal::Success(value.clone()),
444+
_ => Signal::Failure(crate::types::errors::runtime_error::RuntimeError::new(
445+
"T-TEST-000001",
446+
"MissingEchoArgument",
447+
"expected first eager argument",
448+
)),
449+
}
450+
}
451+
423452
#[derive(Clone)]
424453
struct StubRemoteRuntime {
425454
result: NodeExecutionResult,
@@ -848,6 +877,57 @@ mod tests {
848877
assert_eq!(expect_success(signal), int_value(42));
849878
}
850879

880+
#[test]
881+
fn execution_report_includes_function_identifier_subflow_results() {
882+
let mut handlers = FunctionStore::default();
883+
handlers.populate(&[FunctionRegistration::eager("42", echo_first_arg_handler, 1)]);
884+
let engine = ExecutionEngine { handlers };
885+
886+
let add_node = node(
887+
1,
888+
"std::number::add",
889+
vec![
890+
function_thunk_param(
891+
100,
892+
"lhs",
893+
"42",
894+
vec![subflow_setting("value", Some(int_value(20)), false, true)],
895+
),
896+
literal_param(101, "rhs", int_value(2)),
897+
],
898+
None,
899+
);
900+
901+
let report = engine.execute_graph_report(1, vec![add_node], None, None, None, false);
902+
903+
assert_eq!(report.exit_reason, ExitReason::Success);
904+
assert_eq!(expect_success(report.signal), int_value(22));
905+
assert_eq!(report.node_execution_results.len(), 2);
906+
907+
let function_result = &report.node_execution_results[0];
908+
assert_function_result_id(function_result, 42);
909+
assert_eq!(function_result.parameter_results.len(), 1);
910+
assert_eq!(
911+
function_result.parameter_results[0].value,
912+
Some(int_value(20))
913+
);
914+
match function_result.result.as_ref() {
915+
Some(node_execution_result::Result::Success(value)) => {
916+
assert_eq!(value, &int_value(20));
917+
}
918+
other => panic!("expected function success result, got {:?}", other),
919+
}
920+
921+
let node_result = &report.node_execution_results[1];
922+
assert_node_result_id(node_result, 1);
923+
match node_result.result.as_ref() {
924+
Some(node_execution_result::Result::Success(value)) => {
925+
assert_eq!(value, &int_value(22));
926+
}
927+
other => panic!("expected node success result, got {:?}", other),
928+
}
929+
}
930+
851931
#[test]
852932
fn execution_report_includes_literal_node_parameter_results() {
853933
let engine = ExecutionEngine::new();
@@ -867,7 +947,7 @@ mod tests {
867947
assert_eq!(report.node_execution_results.len(), 1);
868948

869949
let node_result = &report.node_execution_results[0];
870-
assert_eq!(node_result.node_id, 1);
950+
assert_node_result_id(node_result, 1);
871951
assert_eq!(node_result.parameter_results.len(), 2);
872952
assert_eq!(node_result.parameter_results[0].value, Some(int_value(1)));
873953
assert_eq!(node_result.parameter_results[1].value, Some(int_value(2)));
@@ -906,7 +986,7 @@ mod tests {
906986
assert_eq!(report.node_execution_results.len(), 2);
907987

908988
let node_result = &report.node_execution_results[1];
909-
assert_eq!(node_result.node_id, 2);
989+
assert_node_result_id(node_result, 2);
910990
assert_eq!(node_result.parameter_results.len(), 2);
911991
assert_eq!(node_result.parameter_results[0].value, Some(int_value(7)));
912992
assert_eq!(node_result.parameter_results[1].value, Some(int_value(5)));
@@ -939,7 +1019,7 @@ mod tests {
9391019
assert_eq!(report.node_execution_results.len(), 1);
9401020

9411021
let node_result = &report.node_execution_results[0];
942-
assert_eq!(node_result.node_id, 1);
1022+
assert_node_result_id(node_result, 1);
9431023
assert_eq!(node_result.parameter_results.len(), 3);
9441024
assert_eq!(node_result.parameter_results[0].value, Some(int_value(200)));
9451025
assert_eq!(
@@ -961,10 +1041,10 @@ mod tests {
9611041
let engine = ExecutionEngine::new();
9621042
let remote = StubRemoteRuntime {
9631043
result: NodeExecutionResult {
964-
node_id: 99,
9651044
started_at: 1,
9661045
finished_at: 2,
9671046
parameter_results: Vec::new(),
1047+
id: Some(node_execution_result::Id::NodeId(99)),
9681048
result: None,
9691049
},
9701050
};
@@ -987,7 +1067,7 @@ mod tests {
9871067
assert_eq!(report.node_execution_results.len(), 1);
9881068

9891069
let node_result = &report.node_execution_results[0];
990-
assert_eq!(node_result.node_id, 1);
1070+
assert_node_result_id(node_result, 1);
9911071
assert_eq!(node_result.parameter_results.len(), 1);
9921072
assert_eq!(node_result.parameter_results[0].value, Some(int_value(20)));
9931073
match node_result.result.as_ref() {
@@ -1012,7 +1092,7 @@ mod tests {
10121092
assert_eq!(report.node_execution_results.len(), 1);
10131093

10141094
let node_result = &report.node_execution_results[0];
1015-
assert_eq!(node_result.node_id, 1);
1095+
assert_node_result_id(node_result, 1);
10161096
assert!(node_result.started_at >= 1_000_000_000_000_000);
10171097
assert!(node_result.finished_at > node_result.started_at);
10181098
assert!(node_result.finished_at - node_result.started_at >= 1_000);
@@ -1059,7 +1139,7 @@ mod tests {
10591139
let callback_results: Vec<_> = report
10601140
.node_execution_results
10611141
.iter()
1062-
.filter(|result| result.node_id == 2)
1142+
.filter(|result| result.id == Some(node_execution_result::Id::NodeId(2)))
10631143
.collect();
10641144
assert_eq!(callback_results.len(), 3);
10651145

@@ -1093,7 +1173,7 @@ mod tests {
10931173
vec![Some(int_value(3)), Some(int_value(2))],
10941174
]
10951175
);
1096-
assert_eq!(report.node_execution_results[3].node_id, 1);
1176+
assert_node_result_id(&report.node_execution_results[3], 1);
10971177
}
10981178

10991179
#[test]

crates/taurus-core/src/runtime/engine/executor.rs

Lines changed: 77 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -183,15 +183,27 @@ impl<'a> EngineExecutor<'a> {
183183
function: &FunctionThunk,
184184
value_store: &mut ValueStore,
185185
) -> ExecutionResult {
186+
let started_at = now_unix_micros();
187+
let function_result_id = parse_function_result_id(function);
186188
let entry = match self.handlers.get(function.identifier.as_str()).copied() {
187189
Some(entry) => entry,
188190
None => {
191+
let error = RuntimeError::new(
192+
"T-CORE-000002",
193+
"FunctionNotFound",
194+
format!("Function {} not found", function.identifier),
195+
);
196+
if let Some(function_id) = function_result_id {
197+
value_store.insert_function_error_with_timing(
198+
function_id,
199+
error.clone(),
200+
Vec::new(),
201+
started_at,
202+
now_unix_micros(),
203+
);
204+
}
189205
return ExecutionResult {
190-
signal: Signal::Failure(RuntimeError::new(
191-
"T-CORE-000002",
192-
"FunctionNotFound",
193-
format!("Function {} not found", function.identifier),
194-
)),
206+
signal: Signal::Failure(error),
195207
root_frame: None,
196208
};
197209
}
@@ -208,12 +220,24 @@ impl<'a> EngineExecutor<'a> {
208220
Err(err) => {
209221
let signal = Signal::Failure(err);
210222
self.trace_exit(frame_id, &signal, value_store);
223+
if let Some(function_id) = function_result_id {
224+
let parameter_results = Vec::new();
225+
self.commit_function_result(
226+
function_id,
227+
signal.clone(),
228+
parameter_results,
229+
started_at,
230+
now_unix_micros(),
231+
value_store,
232+
);
233+
}
211234
return ExecutionResult {
212235
signal,
213236
root_frame: frame_id,
214237
};
215238
}
216239
};
240+
let parameter_results = parameter_results_from_args(&args);
217241

218242
let signal =
219243
if let Some(signal) = self.force_eager_args(&entry, &mut args, value_store, frame_id) {
@@ -233,6 +257,16 @@ impl<'a> EngineExecutor<'a> {
233257
};
234258

235259
self.trace_exit(frame_id, &signal, value_store);
260+
if let Some(function_id) = function_result_id {
261+
self.commit_function_result(
262+
function_id,
263+
signal.clone(),
264+
parameter_results,
265+
started_at,
266+
now_unix_micros(),
267+
value_store,
268+
);
269+
}
236270

237271
ExecutionResult {
238272
signal,
@@ -747,6 +781,40 @@ impl<'a> EngineExecutor<'a> {
747781
}
748782
}
749783

784+
fn commit_function_result(
785+
&self,
786+
function_id: i64,
787+
signal: Signal,
788+
parameter_results: Vec<NodeParameterNodeExecutionResult>,
789+
started_at: i64,
790+
finished_at: i64,
791+
value_store: &mut ValueStore,
792+
) -> Signal {
793+
match signal {
794+
Signal::Success(value) => {
795+
value_store.insert_function_success_with_timing(
796+
function_id,
797+
value.clone(),
798+
parameter_results,
799+
started_at,
800+
finished_at,
801+
);
802+
Signal::Success(value)
803+
}
804+
Signal::Failure(err) => {
805+
value_store.insert_function_error_with_timing(
806+
function_id,
807+
err.clone(),
808+
parameter_results,
809+
started_at,
810+
finished_at,
811+
);
812+
Signal::Failure(err)
813+
}
814+
other => other,
815+
}
816+
}
817+
750818
fn commit_remote_result(
751819
&self,
752820
node_id: i64,
@@ -901,6 +969,10 @@ fn compiled_thunk_to_argument(thunk: &CompiledThunk) -> Thunk {
901969
}
902970
}
903971

972+
fn parse_function_result_id(function: &FunctionThunk) -> Option<i64> {
973+
function.identifier.parse::<i64>().ok()
974+
}
975+
904976
fn resolve_function_setting(
905977
function: &FunctionThunk,
906978
setting: &SubFlowSetting,

crates/taurus-core/src/runtime/execution/value_store.rs

Lines changed: 38 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
33
use std::collections::HashMap;
44

5-
use tucana::shared::node_execution_result::Result as TucanaNodeResult;
5+
use tucana::shared::node_execution_result::{Id as TucanaNodeResultId, Result as TucanaNodeResult};
66
use tucana::shared::{
77
InputType, NodeExecutionResult, NodeParameterNodeExecutionResult, ReferenceValue, Value,
88
value::Kind,
@@ -157,10 +157,10 @@ impl ValueStore {
157157
self.insert_node_result(
158158
id,
159159
NodeExecutionResult {
160-
node_id: id,
161160
started_at,
162161
finished_at,
163162
parameter_results,
163+
id: Some(TucanaNodeResultId::NodeId(id)),
164164
result: Some(TucanaNodeResult::Success(value)),
165165
},
166166
);
@@ -177,21 +177,55 @@ impl ValueStore {
177177
self.insert_node_result(
178178
id,
179179
NodeExecutionResult {
180-
node_id: id,
181180
started_at,
182181
finished_at,
183182
parameter_results,
183+
id: Some(TucanaNodeResultId::NodeId(id)),
184184
result: Some(TucanaNodeResult::Error(runtime_error.as_tucana_error())),
185185
},
186186
);
187187
}
188188

189189
pub fn insert_node_result(&mut self, id: i64, mut result: NodeExecutionResult) {
190-
result.node_id = id;
190+
result.id = Some(TucanaNodeResultId::NodeId(id));
191191
self.latest_results.insert(id, result.clone());
192192
self.result_history.push(result);
193193
}
194194

195+
pub fn insert_function_success_with_timing(
196+
&mut self,
197+
id: i64,
198+
value: Value,
199+
parameter_results: Vec<NodeParameterNodeExecutionResult>,
200+
started_at: i64,
201+
finished_at: i64,
202+
) {
203+
self.result_history.push(NodeExecutionResult {
204+
started_at,
205+
finished_at,
206+
parameter_results,
207+
id: Some(TucanaNodeResultId::FunctionId(id)),
208+
result: Some(TucanaNodeResult::Success(value)),
209+
});
210+
}
211+
212+
pub fn insert_function_error_with_timing(
213+
&mut self,
214+
id: i64,
215+
runtime_error: RuntimeError,
216+
parameter_results: Vec<NodeParameterNodeExecutionResult>,
217+
started_at: i64,
218+
finished_at: i64,
219+
) {
220+
self.result_history.push(NodeExecutionResult {
221+
started_at,
222+
finished_at,
223+
parameter_results,
224+
id: Some(TucanaNodeResultId::FunctionId(id)),
225+
result: Some(TucanaNodeResult::Error(runtime_error.as_tucana_error())),
226+
});
227+
}
228+
195229
pub fn node_execution_results(&self) -> Vec<NodeExecutionResult> {
196230
self.result_history.clone()
197231
}

0 commit comments

Comments
 (0)