Skip to content

Commit 970a4d6

Browse files
committed
fix: correct async impl for nats
1 parent 308d0df commit 970a4d6

10 files changed

Lines changed: 580 additions & 77 deletions

File tree

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

Lines changed: 281 additions & 23 deletions
Large diffs are not rendered by default.

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

Lines changed: 30 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,10 @@ pub enum CompileError {
3535
node_id: i64,
3636
parameter_index: usize,
3737
},
38+
EmptyRemoteService {
39+
node_id: i64,
40+
definition_source: String,
41+
},
3842
}
3943

4044
impl CompileError {
@@ -88,6 +92,17 @@ impl CompileError {
8892
node_id, parameter_index
8993
),
9094
),
95+
CompileError::EmptyRemoteService {
96+
node_id,
97+
definition_source,
98+
} => RuntimeError::new(
99+
"T-CORE-000106",
100+
"FlowCompileError",
101+
format!(
102+
"Node {} definition_source '{}' does not contain a remote service name",
103+
node_id, definition_source
104+
),
105+
),
91106
}
92107
}
93108
}
@@ -134,7 +149,7 @@ pub fn compile_flow(
134149
None => None,
135150
};
136151

137-
let execution_target = execution_target_for(&node);
152+
let execution_target = execution_target_for(node_id, &node)?;
138153

139154
let mut parameters = Vec::with_capacity(node.parameters.len());
140155
for (parameter_index, parameter) in node.parameters.iter().enumerate() {
@@ -198,12 +213,21 @@ pub fn compile_flow(
198213
})
199214
}
200215

201-
fn execution_target_for(node: &NodeFunction) -> NodeExecutionTarget {
216+
fn execution_target_for(
217+
node_id: i64,
218+
node: &NodeFunction,
219+
) -> Result<NodeExecutionTarget, CompileError> {
202220
match node.definition_source.as_deref() {
203-
None | Some("") | Some("taurus") => NodeExecutionTarget::Local,
204-
Some(source) if source.starts_with("draco") => NodeExecutionTarget::Local,
205-
Some(service) => NodeExecutionTarget::Remote {
206-
service: service.to_string(),
221+
None | Some("") | Some("taurus") => Ok(NodeExecutionTarget::Local),
222+
Some(source) if source.starts_with("draco") => Ok(NodeExecutionTarget::Local),
223+
Some(service) => match service.strip_prefix("action.").unwrap_or(service) {
224+
"" => Err(CompileError::EmptyRemoteService {
225+
node_id,
226+
definition_source: service.to_string(),
227+
}),
228+
service => Ok(NodeExecutionTarget::Remote {
229+
service: service.to_string(),
230+
}),
207231
},
208232
}
209233
}

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,13 +22,13 @@ pub enum EmitType {
2222
}
2323

2424
/// Callback interface for streaming execution lifecycle events.
25-
pub trait RespondEmitter {
25+
pub trait RespondEmitter: Send + Sync {
2626
fn emit(&self, execution_id: ExecutionId, emit_type: EmitType, value: Value);
2727
}
2828

2929
impl<F> RespondEmitter for F
3030
where
31-
F: Fn(ExecutionId, EmitType, Value) + ?Sized,
31+
F: Fn(ExecutionId, EmitType, Value) + Send + Sync + ?Sized,
3232
{
3333
fn emit(&self, execution_id: ExecutionId, emit_type: EmitType, value: Value) {
3434
self(execution_id, emit_type, value);

0 commit comments

Comments
 (0)