Skip to content

Commit 0eecd7d

Browse files
Merge pull request #225 from code0-tech/#224-empty-execution-results
empty execution results
2 parents 7710540 + 806fdb7 commit 0eecd7d

10 files changed

Lines changed: 537 additions & 16 deletions

File tree

crates/taurus-core/src/handler/argument.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ use tucana::shared::SubFlowSetting;
1313
#[derive(Clone)]
1414
pub struct FunctionThunk {
1515
pub identifier: String,
16+
pub result_id: Option<i64>,
1617
pub parameter_index: i64,
1718
pub settings: Vec<SubFlowSetting>,
1819
}
@@ -21,6 +22,7 @@ impl fmt::Debug for FunctionThunk {
2122
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2223
f.debug_struct("FunctionThunk")
2324
.field("identifier", &self.identifier)
25+
.field("result_id", &self.result_id)
2426
.field("parameter_index", &self.parameter_index)
2527
.field("settings_len", &self.settings.len())
2628
.finish()

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

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1176,6 +1176,94 @@ mod tests {
11761176
assert_node_result_id(&report.node_execution_results[3], 1);
11771177
}
11781178

1179+
#[test]
1180+
fn execution_report_keeps_every_for_each_function_identifier_callback_execution() {
1181+
let engine = ExecutionEngine::new();
1182+
let mut respond_node = node(
1183+
1,
1184+
"rest::control::respond",
1185+
vec![
1186+
literal_param(1, "http_status_code", int_value(200)),
1187+
literal_param(2, "headers", empty_struct_value()),
1188+
literal_param(3, "payload", string_value("20")),
1189+
],
1190+
None,
1191+
);
1192+
respond_node.definition_source = Some("draco-draco-cron".to_string());
1193+
let mut for_each_node = node(
1194+
2,
1195+
"std::list::for_each",
1196+
vec![
1197+
literal_param(
1198+
4,
1199+
"list",
1200+
list_value(vec![int_value(1), int_value(2), int_value(3)]),
1201+
),
1202+
function_thunk_param(
1203+
5,
1204+
"consumer",
1205+
"std::boolean::from_number",
1206+
vec![subflow_setting("value", Some(null_value()), false, false)],
1207+
),
1208+
],
1209+
Some(1),
1210+
);
1211+
for_each_node.definition_source = Some("draco-draco-cron".to_string());
1212+
1213+
let report = engine.execute_graph_report(
1214+
2,
1215+
vec![respond_node, for_each_node],
1216+
None,
1217+
None,
1218+
None,
1219+
false,
1220+
);
1221+
1222+
assert_eq!(report.exit_reason, ExitReason::Success);
1223+
assert_eq!(expect_success(report.signal), {
1224+
let mut fields = std::collections::HashMap::new();
1225+
fields.insert("http_status_code".to_string(), int_value(200));
1226+
fields.insert("headers".to_string(), empty_struct_value());
1227+
fields.insert("payload".to_string(), string_value("20"));
1228+
Value {
1229+
kind: Some(Kind::StructValue(Struct { fields })),
1230+
}
1231+
});
1232+
assert_eq!(report.node_execution_results.len(), 5);
1233+
1234+
let function_results: Vec<_> = report
1235+
.node_execution_results
1236+
.iter()
1237+
.filter(|result| result.id == Some(node_execution_result::Id::FunctionId(5)))
1238+
.collect();
1239+
assert_eq!(function_results.len(), 3);
1240+
1241+
for (index, result) in function_results.iter().enumerate() {
1242+
assert_eq!(result.parameter_results.len(), 1);
1243+
assert_eq!(
1244+
result.parameter_results[0].value,
1245+
Some(int_value(index as i64 + 1))
1246+
);
1247+
match result.result.as_ref() {
1248+
Some(node_execution_result::Result::Success(value)) => {
1249+
assert_eq!(
1250+
value,
1251+
&Value {
1252+
kind: Some(Kind::BoolValue(true)),
1253+
}
1254+
);
1255+
}
1256+
other => panic!("expected function success result, got {:?}", other),
1257+
}
1258+
}
1259+
1260+
assert_function_result_id(&report.node_execution_results[0], 5);
1261+
assert_function_result_id(&report.node_execution_results[1], 5);
1262+
assert_function_result_id(&report.node_execution_results[2], 5);
1263+
assert_node_result_id(&report.node_execution_results[3], 2);
1264+
assert_node_result_id(&report.node_execution_results[4], 1);
1265+
}
1266+
11791267
#[test]
11801268
fn emitter_emits_start_and_finish_for_successful_execution() {
11811269
let engine = ExecutionEngine::new();

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,7 @@ pub fn compile_flow(
162162
Some(sub_flow::ExecutionReference::FunctionIdentifier(identifier)) => {
163163
CompiledArg::Deferred(CompiledThunk::Function {
164164
identifier: identifier.clone(),
165+
result_id: identifier.parse().ok().or(Some(parameter.database_id)),
165166
parameter_index: parameter_index as i64,
166167
settings: sub_flow.settings.clone(),
167168
})

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

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -734,7 +734,7 @@ impl<'a> EngineExecutor<'a> {
734734
}
735735

736736
let mut fields = HashMap::new();
737-
for (parameter, value) in node.parameters.iter().zip(values.into_iter()) {
737+
for (parameter, value) in node.parameters.iter().zip(values) {
738738
fields.insert(parameter.runtime_parameter_id.clone(), value);
739739
}
740740

@@ -959,18 +959,22 @@ fn compiled_thunk_to_argument(thunk: &CompiledThunk) -> Thunk {
959959
CompiledThunk::Node(node_id) => Thunk::Node(*node_id),
960960
CompiledThunk::Function {
961961
identifier,
962+
result_id,
962963
parameter_index,
963964
settings,
964965
} => Thunk::Function(FunctionThunk {
965966
identifier: identifier.clone(),
967+
result_id: *result_id,
966968
parameter_index: *parameter_index,
967969
settings: settings.clone(),
968970
}),
969971
}
970972
}
971973

972974
fn parse_function_result_id(function: &FunctionThunk) -> Option<i64> {
973-
function.identifier.parse::<i64>().ok()
975+
function
976+
.result_id
977+
.or_else(|| function.identifier.parse::<i64>().ok())
974978
}
975979

976980
fn resolve_function_setting(

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ pub enum CompiledThunk {
2626
Node(i64),
2727
Function {
2828
identifier: String,
29+
result_id: Option<i64>,
2930
parameter_index: i64,
3031
settings: Vec<SubFlowSetting>,
3132
},

crates/taurus-core/src/runtime/functions/http.rs

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -297,11 +297,9 @@ fn decode_response_payload(response: http::Response<Body>) -> Result<Value, Stri
297297
.as_deref()
298298
.map(content_type_is_json)
299299
.unwrap_or(false)
300-
{
301-
if let Ok(json) = serde_json::from_str::<JsonValue>(&text) {
300+
&& let Ok(json) = serde_json::from_str::<JsonValue>(&text) {
302301
return Ok(from_json_value(json));
303302
}
304-
}
305303

306304
return Ok(text.to_value());
307305
}

crates/taurus-manual/src/main.rs

Lines changed: 81 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,12 @@ use taurus_provider::providers::remote::nats_remote_runtime::NATSRemoteRuntime;
1414
use tucana::shared::ExecutionFlow;
1515
use tucana::shared::NodeExecutionResult;
1616
use tucana::shared::ValidationFlow;
17+
use tucana::shared::Value;
1718
use tucana::shared::helper::value::from_json_value;
1819
use tucana::shared::helper::value::to_json_value;
1920
use tucana::shared::node_execution_result::Id as NodeExecutionResultId;
21+
use tucana::shared::node_execution_result::Result as NodeExecutionResultResult;
22+
use tucana::shared::value::Kind;
2023

2124
#[derive(Clone, Deserialize)]
2225
pub struct Input {
@@ -161,10 +164,7 @@ async fn main() {
161164
}
162165

163166
let flow_input = match case.inputs.get(index as usize) {
164-
Some(inp) => match inp.input.clone() {
165-
Some(json_input) => Some(from_json_value(json_input)),
166-
None => None,
167-
},
167+
Some(inp) => inp.input.clone().map(from_json_value),
168168
None => None,
169169
};
170170

@@ -182,7 +182,7 @@ async fn main() {
182182
);
183183
let duration_us = start.elapsed().as_micros();
184184
let finished_at = now_unix_micros();
185-
print_timing_debug(
185+
print_manual_execution_debug(
186186
started_at,
187187
finished_at,
188188
duration_us,
@@ -223,7 +223,7 @@ async fn main() {
223223
);
224224
let duration_us = start.elapsed().as_micros();
225225
let finished_at = now_unix_micros();
226-
print_timing_debug(
226+
print_manual_execution_debug(
227227
started_at,
228228
finished_at,
229229
duration_us,
@@ -301,6 +301,18 @@ async fn queue_execution(
301301
println!("{}", execution_id);
302302
}
303303

304+
fn print_manual_execution_debug(
305+
started_at: i64,
306+
finished_at: i64,
307+
duration_us: u128,
308+
node_results: &[NodeExecutionResult],
309+
) {
310+
let mut normalized_results = node_results.to_vec();
311+
normalize_node_execution_results(&mut normalized_results);
312+
print_timing_debug(started_at, finished_at, duration_us, &normalized_results);
313+
print_execution_result_debug(&normalized_results);
314+
}
315+
304316
fn print_timing_debug(
305317
started_at: i64,
306318
finished_at: i64,
@@ -342,6 +354,69 @@ fn print_timing_debug(
342354
}
343355
}
344356

357+
fn print_execution_result_debug(node_results: &[NodeExecutionResult]) {
358+
eprintln!("[manual execution result] {:#?}", node_results);
359+
}
360+
361+
fn normalize_node_execution_results(node_results: &mut [NodeExecutionResult]) {
362+
for result in node_results {
363+
normalize_node_execution_result(result);
364+
}
365+
}
366+
367+
fn normalize_node_execution_result(result: &mut NodeExecutionResult) {
368+
for parameter_result in &mut result.parameter_results {
369+
match &mut parameter_result.value {
370+
Some(value) => normalize_value(value),
371+
None => {
372+
parameter_result.value = Some(null_value());
373+
}
374+
}
375+
}
376+
377+
match &mut result.result {
378+
Some(NodeExecutionResultResult::Success(value)) => normalize_value(value),
379+
Some(NodeExecutionResultResult::Error(error)) => {
380+
if let Some(details) = &mut error.details {
381+
for value in details.fields.values_mut() {
382+
normalize_value(value);
383+
}
384+
}
385+
}
386+
None => {
387+
result.result = Some(NodeExecutionResultResult::Success(null_value()));
388+
}
389+
}
390+
}
391+
392+
fn normalize_value(value: &mut Value) {
393+
match &mut value.kind {
394+
Some(Kind::StructValue(struct_value)) => {
395+
for field in struct_value.fields.values_mut() {
396+
normalize_value(field);
397+
}
398+
}
399+
Some(Kind::ListValue(list_value)) => {
400+
for item in &mut list_value.values {
401+
normalize_value(item);
402+
}
403+
}
404+
Some(Kind::NumberValue(number)) if number.number.is_none() => {
405+
value.kind = Some(Kind::NullValue(0));
406+
}
407+
Some(_) => {}
408+
None => {
409+
value.kind = Some(Kind::NullValue(0));
410+
}
411+
}
412+
}
413+
414+
fn null_value() -> Value {
415+
Value {
416+
kind: Some(Kind::NullValue(0)),
417+
}
418+
}
419+
345420
fn execution_result_id_label(result: &NodeExecutionResult) -> String {
346421
match result.id {
347422
Some(NodeExecutionResultId::NodeId(id)) => format!("node_id={}", id),

crates/taurus/src/app/mod.rs

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -50,11 +50,10 @@ pub async fn run() {
5050
wait_for_shutdown(&mut worker_task, &mut health_task).await;
5151
if let Some(handle) = runtime_status_heartbeat_task.take() {
5252
handle.abort();
53-
if let Err(err) = handle.await {
54-
if !err.is_cancelled() {
53+
if let Err(err) = handle.await
54+
&& !err.is_cancelled() {
5555
log::warn!("Runtime status heartbeat task ended unexpectedly: {}", err);
5656
}
57-
}
5857
}
5958
update_stopped_status(runtime_status_service.as_ref()).await;
6059

0 commit comments

Comments
 (0)