Skip to content

Commit e92a91d

Browse files
Merge pull request #216 from code0-tech/#207-merge-execution-fix
fix: handled duration for nodes
2 parents bf9cdaa + f668d80 commit e92a91d

7 files changed

Lines changed: 325 additions & 55 deletions

File tree

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

Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -265,8 +265,12 @@ impl ExecutionEngine {
265265
#[cfg(test)]
266266
mod tests {
267267
use super::*;
268+
use crate::handler::argument::Argument;
269+
use crate::handler::registry::{FunctionRegistration, FunctionStore, ThunkRunner};
270+
use crate::runtime::execution::value_store::ValueStore;
268271
use crate::types::exit_reason::ExitReason;
269272
use std::cell::RefCell;
273+
use std::time::Duration;
270274
use tucana::shared::{
271275
InputType, ListValue, NodeParameter, NodeValue, ReferenceValue, Struct, SubFlow,
272276
SubFlowSetting, Value, node_execution_result, node_value, reference_value,
@@ -405,6 +409,15 @@ mod tests {
405409
}
406410
}
407411

412+
fn sleep_handler(
413+
_args: &[Argument],
414+
_ctx: &mut ValueStore,
415+
_run: &mut ThunkRunner<'_>,
416+
) -> Signal {
417+
std::thread::sleep(Duration::from_micros(2_000));
418+
Signal::Success(null_value())
419+
}
420+
408421
fn input_type_ref_param(
409422
database_id: i64,
410423
runtime_parameter_id: &str,
@@ -925,6 +938,103 @@ mod tests {
925938
));
926939
}
927940

941+
#[test]
942+
fn node_execution_result_tracks_actual_node_duration() {
943+
let mut handlers = FunctionStore::new();
944+
handlers.populate(&[FunctionRegistration::eager("test::sleep", sleep_handler, 0)]);
945+
let engine = ExecutionEngine { handlers };
946+
let sleep_node = node(1, "test::sleep", vec![], None);
947+
948+
let report = engine.execute_graph_report(1, vec![sleep_node], None, None, None, false);
949+
950+
assert_eq!(report.exit_reason, ExitReason::Success);
951+
assert_eq!(report.node_execution_results.len(), 1);
952+
953+
let node_result = &report.node_execution_results[0];
954+
assert_eq!(node_result.node_id, 1);
955+
assert!(node_result.started_at >= 1_000_000_000_000_000);
956+
assert!(node_result.finished_at > node_result.started_at);
957+
assert!(node_result.finished_at - node_result.started_at >= 1_000);
958+
}
959+
960+
#[test]
961+
fn execution_report_keeps_every_for_each_callback_execution() {
962+
let engine = ExecutionEngine::new();
963+
let for_each_node = node(
964+
1,
965+
"std::list::for_each",
966+
vec![
967+
literal_param(
968+
100,
969+
"list",
970+
list_value(vec![int_value(1), int_value(2), int_value(3)]),
971+
),
972+
thunk_param(101, "consumer", 2),
973+
],
974+
None,
975+
);
976+
let callback_node = node(
977+
2,
978+
"std::number::add",
979+
vec![
980+
input_type_ref_param(200, "first", 1, 1, 0),
981+
literal_param(201, "second", int_value(2)),
982+
],
983+
None,
984+
);
985+
986+
let report = engine.execute_graph_report(
987+
1,
988+
vec![for_each_node, callback_node],
989+
None,
990+
None,
991+
None,
992+
false,
993+
);
994+
995+
assert_eq!(report.exit_reason, ExitReason::Success);
996+
assert_eq!(report.node_execution_results.len(), 4);
997+
998+
let callback_results: Vec<_> = report
999+
.node_execution_results
1000+
.iter()
1001+
.filter(|result| result.node_id == 2)
1002+
.collect();
1003+
assert_eq!(callback_results.len(), 3);
1004+
1005+
let callback_values: Vec<_> = callback_results
1006+
.iter()
1007+
.map(|result| match result.result.as_ref() {
1008+
Some(node_execution_result::Result::Success(value)) => value.clone(),
1009+
other => panic!("expected callback success result, got {:?}", other),
1010+
})
1011+
.collect();
1012+
1013+
assert_eq!(
1014+
callback_values,
1015+
vec![int_value(3), int_value(4), int_value(5)]
1016+
);
1017+
let callback_parameters: Vec<_> = callback_results
1018+
.iter()
1019+
.map(|result| {
1020+
result
1021+
.parameter_results
1022+
.iter()
1023+
.map(|parameter| parameter.value.clone())
1024+
.collect::<Vec<_>>()
1025+
})
1026+
.collect();
1027+
assert_eq!(
1028+
callback_parameters,
1029+
vec![
1030+
vec![Some(int_value(1)), Some(int_value(2))],
1031+
vec![Some(int_value(2)), Some(int_value(2))],
1032+
vec![Some(int_value(3)), Some(int_value(2))],
1033+
]
1034+
);
1035+
assert_eq!(report.node_execution_results[3].node_id, 1);
1036+
}
1037+
9281038
#[test]
9291039
fn emitter_emits_start_and_finish_for_successful_execution() {
9301040
let engine = ExecutionEngine::new();

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

Lines changed: 53 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ use crate::runtime::execution::trace::{
2626
use crate::runtime::execution::tracer::{ExecutionTracer, Tracer};
2727
use crate::runtime::execution::value_store::{ValueStore, ValueStoreResult};
2828
use crate::runtime::remote::{RemoteExecution, RemoteRuntime};
29+
use crate::time::now_unix_micros;
2930
use crate::types::errors::runtime_error::RuntimeError;
3031
use crate::types::signal::Signal;
3132

@@ -66,6 +67,8 @@ struct NodeResult {
6667
signal: Signal,
6768
frame_id: Option<u64>,
6869
parameter_results: Vec<NodeParameterNodeExecutionResult>,
70+
started_at: i64,
71+
finished_at: i64,
6972
}
7073

7174
struct ExecutedNode {
@@ -125,10 +128,12 @@ impl<'a> EngineExecutor<'a> {
125128
emitter.emit(self.execution_id, EmitType::OngoingExec, value.clone());
126129
}
127130

128-
value_store.insert_success_with_parameters(
131+
value_store.insert_success_with_timing(
129132
node_id,
130133
value.clone(),
131134
result.parameter_results,
135+
result.started_at,
136+
result.finished_at,
132137
);
133138
match next_idx {
134139
Some(next) => current_idx = next,
@@ -243,25 +248,38 @@ impl<'a> EngineExecutor<'a> {
243248
let frame_id = self.trace_enter(node, value_store);
244249
let result = match &node.execution_target {
245250
NodeExecutionTarget::Local => {
251+
let started_at = now_unix_micros();
246252
let executed = self.execute_local_node(node, value_store, frame_id);
253+
let finished_at = now_unix_micros();
247254
let parameter_results = executed.parameter_results;
248255
let signal = self.commit_result(
249256
node.id,
250257
executed.signal,
251258
parameter_results.clone(),
259+
started_at,
260+
finished_at,
252261
value_store,
253262
);
254263
NodeResult {
255264
signal,
256265
frame_id,
257266
parameter_results,
267+
started_at,
268+
finished_at,
269+
}
270+
}
271+
NodeExecutionTarget::Remote { service } => {
272+
let started_at = now_unix_micros();
273+
let signal = self.execute_remote_node(node, service, value_store, frame_id);
274+
let finished_at = now_unix_micros();
275+
NodeResult {
276+
signal,
277+
frame_id,
278+
parameter_results: Vec::new(),
279+
started_at,
280+
finished_at,
258281
}
259282
}
260-
NodeExecutionTarget::Remote { service } => NodeResult {
261-
signal: self.execute_remote_node(node, service, value_store, frame_id),
262-
frame_id,
263-
parameter_results: Vec::new(),
264-
},
265283
};
266284
self.trace_exit(frame_id, &result.signal, value_store);
267285

@@ -331,6 +349,7 @@ impl<'a> EngineExecutor<'a> {
331349
value_store: &mut ValueStore,
332350
frame_id: Option<u64>,
333351
) -> Signal {
352+
let started_at = now_unix_micros();
334353
let remote_runtime = match self.remote {
335354
Some(remote) => remote,
336355
None => {
@@ -342,6 +361,8 @@ impl<'a> EngineExecutor<'a> {
342361
"Remote runtime not configured",
343362
)),
344363
Vec::new(),
364+
started_at,
365+
now_unix_micros(),
345366
value_store,
346367
);
347368
}
@@ -350,7 +371,14 @@ impl<'a> EngineExecutor<'a> {
350371
let mut args = match self.build_args(node, value_store, frame_id) {
351372
Ok(args) => args,
352373
Err(err) => {
353-
return self.commit_result(node.id, Signal::Failure(err), Vec::new(), value_store);
374+
return self.commit_result(
375+
node.id,
376+
Signal::Failure(err),
377+
Vec::new(),
378+
started_at,
379+
now_unix_micros(),
380+
value_store,
381+
);
354382
}
355383
};
356384

@@ -361,6 +389,8 @@ impl<'a> EngineExecutor<'a> {
361389
node.id,
362390
signal,
363391
parameter_results_from_args(&args),
392+
started_at,
393+
now_unix_micros(),
364394
value_store,
365395
);
366396
}
@@ -374,6 +404,8 @@ impl<'a> EngineExecutor<'a> {
374404
node.id,
375405
Signal::Failure(err),
376406
parameter_results,
407+
started_at,
408+
now_unix_micros(),
377409
value_store,
378410
);
379411
}
@@ -390,6 +422,8 @@ impl<'a> EngineExecutor<'a> {
390422
node.id,
391423
Signal::Failure(err),
392424
parameter_results,
425+
started_at,
426+
now_unix_micros(),
393427
value_store,
394428
),
395429
}
@@ -678,19 +712,29 @@ impl<'a> EngineExecutor<'a> {
678712
node_id: i64,
679713
signal: Signal,
680714
parameter_results: Vec<NodeParameterNodeExecutionResult>,
715+
started_at: i64,
716+
finished_at: i64,
681717
value_store: &mut ValueStore,
682718
) -> Signal {
683719
match signal {
684720
Signal::Success(value) => {
685-
value_store.insert_success_with_parameters(
721+
value_store.insert_success_with_timing(
686722
node_id,
687723
value.clone(),
688724
parameter_results,
725+
started_at,
726+
finished_at,
689727
);
690728
Signal::Success(value)
691729
}
692730
Signal::Failure(err) => {
693-
value_store.insert_error_with_parameters(node_id, err.clone(), parameter_results);
731+
value_store.insert_error_with_timing(
732+
node_id,
733+
err.clone(),
734+
parameter_results,
735+
started_at,
736+
finished_at,
737+
);
694738
Signal::Failure(err)
695739
}
696740
// Control signals are transient and should not be cached as node outputs.

0 commit comments

Comments
 (0)