Skip to content

Commit 32443b9

Browse files
committed
feat: adjusted emitter api
1 parent 54bc8d9 commit 32443b9

3 files changed

Lines changed: 27 additions & 15 deletions

File tree

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

Lines changed: 14 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ use crate::runtime::remote::RemoteRuntime;
1616
use crate::types::exit_reason::ExitReason;
1717
use crate::types::signal::Signal;
1818
use compiler::compile_flow;
19-
pub use emitter::{EmitType, RespondEmitter};
19+
pub use emitter::{EmitType, ExecutionId, RespondEmitter};
2020

2121
fn null_value() -> Value {
2222
Value {
@@ -71,8 +71,10 @@ impl ExecutionEngine {
7171
respond_emitter: Option<&dyn RespondEmitter>,
7272
with_trace: bool,
7373
) -> (Signal, ExitReason) {
74+
let execution_id = ExecutionId::new_v4();
75+
7476
if let Some(emitter) = respond_emitter {
75-
emitter.emit(EmitType::StartingExec, null_value());
77+
emitter.emit(execution_id, EmitType::StartingExec, null_value());
7678
}
7779

7880
let mut value_store = match flow_input {
@@ -85,7 +87,7 @@ impl ExecutionEngine {
8587
Err(err) => {
8688
let runtime_error = err.as_runtime_error();
8789
if let Some(emitter) = respond_emitter {
88-
emitter.emit(EmitType::FailedExec, runtime_error.as_value());
90+
emitter.emit(execution_id, EmitType::FailedExec, runtime_error.as_value());
8991
}
9092
let signal = Signal::Failure(runtime_error);
9193
return (signal, ExitReason::Failure);
@@ -97,6 +99,7 @@ impl ExecutionEngine {
9799
&self.handlers,
98100
&mut value_store,
99101
remote,
102+
execution_id,
100103
respond_emitter,
101104
with_trace,
102105
);
@@ -108,11 +111,13 @@ impl ExecutionEngine {
108111
}
109112
if let Some(emitter) = respond_emitter {
110113
match &signal {
111-
Signal::Failure(err) => emitter.emit(EmitType::FailedExec, err.as_value()),
114+
Signal::Failure(err) => {
115+
emitter.emit(execution_id, EmitType::FailedExec, err.as_value())
116+
}
112117
Signal::Success(value) | Signal::Return(value) | Signal::Respond(value) => {
113-
emitter.emit(EmitType::FinishedExec, value.clone())
118+
emitter.emit(execution_id, EmitType::FinishedExec, value.clone())
114119
}
115-
Signal::Stop => emitter.emit(EmitType::FinishedExec, null_value()),
120+
Signal::Stop => emitter.emit(execution_id, EmitType::FinishedExec, null_value()),
116121
}
117122
}
118123
let exit_reason = signal.exit_reason();
@@ -386,7 +391,7 @@ mod tests {
386391
fn emitter_emits_start_and_finish_for_successful_execution() {
387392
let engine = ExecutionEngine::new();
388393
let events = RefCell::new(Vec::<EmitType>::new());
389-
let emitter = |emit_type: EmitType, _value: Value| {
394+
let emitter = |_execution_id, emit_type: EmitType, _value: Value| {
390395
events.borrow_mut().push(emit_type);
391396
};
392397

@@ -413,7 +418,7 @@ mod tests {
413418
fn emitter_emits_ongoing_for_intermediate_respond() {
414419
let engine = ExecutionEngine::new();
415420
let events = RefCell::new(Vec::<EmitType>::new());
416-
let emitter = |emit_type: EmitType, _value: Value| {
421+
let emitter = |_execution_id, emit_type: EmitType, _value: Value| {
417422
events.borrow_mut().push(emit_type);
418423
};
419424

@@ -466,7 +471,7 @@ mod tests {
466471
fn emitter_emits_failed_for_runtime_failure() {
467472
let engine = ExecutionEngine::new();
468473
let events = RefCell::new(Vec::<EmitType>::new());
469-
let emitter = |emit_type: EmitType, _value: Value| {
474+
let emitter = |_execution_id, emit_type: EmitType, _value: Value| {
470475
events.borrow_mut().push(emit_type);
471476
};
472477

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

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,10 @@
11
//! Respond emitter abstraction used by the engine.
22
33
use tucana::shared::Value;
4+
use uuid::Uuid;
5+
6+
/// Unique identifier for one top-level flow execution.
7+
pub type ExecutionId = Uuid;
48

59
/// Execution lifecycle event emitted by the runtime engine.
610
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -17,14 +21,14 @@ pub enum EmitType {
1721

1822
/// Callback interface for streaming execution lifecycle events.
1923
pub trait RespondEmitter {
20-
fn emit(&self, emit_type: EmitType, value: Value);
24+
fn emit(&self, execution_id: ExecutionId, emit_type: EmitType, value: Value);
2125
}
2226

2327
impl<F> RespondEmitter for F
2428
where
25-
F: Fn(EmitType, Value) + ?Sized,
29+
F: Fn(ExecutionId, EmitType, Value) + ?Sized,
2630
{
27-
fn emit(&self, emit_type: EmitType, value: Value) {
28-
self(emit_type, value);
31+
fn emit(&self, execution_id: ExecutionId, emit_type: EmitType, value: Value) {
32+
self(execution_id, emit_type, value);
2933
}
3034
}

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

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ use uuid::Uuid;
1212

1313
use crate::handler::argument::{Argument, ParameterNode};
1414
use crate::handler::registry::{FunctionStore, HandlerFunctionEntry};
15-
use crate::runtime::engine::emitter::{EmitType, RespondEmitter};
15+
use crate::runtime::engine::emitter::{EmitType, ExecutionId, RespondEmitter};
1616
use crate::runtime::engine::model::{CompiledArg, CompiledFlow, CompiledNode, NodeExecutionTarget};
1717
use crate::runtime::execution::trace::{
1818
ArgKind, ArgTrace, EdgeKind, Outcome, ReferenceKind, TraceRun,
@@ -28,6 +28,7 @@ pub fn execute_compiled(
2828
handlers: &FunctionStore,
2929
value_store: &mut ValueStore,
3030
remote: Option<&dyn RemoteRuntime>,
31+
execution_id: ExecutionId,
3132
respond_emitter: Option<&dyn RespondEmitter>,
3233
with_trace: bool,
3334
) -> (Signal, Option<TraceRun>) {
@@ -37,6 +38,7 @@ pub fn execute_compiled(
3738
flow,
3839
handlers,
3940
remote,
41+
execution_id,
4042
respond_emitter,
4143
tracer: tracer.as_ref(),
4244
};
@@ -63,6 +65,7 @@ struct EngineExecutor<'a> {
6365
flow: &'a CompiledFlow,
6466
handlers: &'a FunctionStore,
6567
remote: Option<&'a dyn RemoteRuntime>,
68+
execution_id: ExecutionId,
6669
respond_emitter: Option<&'a dyn RespondEmitter>,
6770
tracer: Option<&'a RefCell<Tracer>>,
6871
}
@@ -107,7 +110,7 @@ impl<'a> EngineExecutor<'a> {
107110
Signal::Respond(value) => {
108111
// `Respond` is an observable side effect; execution may still continue.
109112
if let Some(emitter) = self.respond_emitter {
110-
emitter.emit(EmitType::OngoingExec, value.clone());
113+
emitter.emit(self.execution_id, EmitType::OngoingExec, value.clone());
111114
}
112115

113116
value_store.insert_success(node_id, value.clone());

0 commit comments

Comments
 (0)