Skip to main content

roder_dynamic_workflows/execution/
executor.rs

1use std::sync::Arc;
2
3use async_trait::async_trait;
4use roder_api::subagents::{SubagentDispatcher, SubagentResult};
5use roder_api::trace::SubagentTraceSink;
6
7use super::{WorkflowAgentExecutionContext, WorkflowAgentExecutionRequest};
8
9#[async_trait]
10pub trait WorkflowAgentExecutor: Send + Sync + 'static {
11    async fn execute_agent(
12        &self,
13        context: WorkflowAgentExecutionContext,
14        request: WorkflowAgentExecutionRequest,
15    ) -> anyhow::Result<SubagentResult>;
16}
17
18pub struct SubagentDispatcherWorkflowExecutor {
19    dispatcher: Arc<dyn SubagentDispatcher>,
20    trace_sink: Option<Arc<dyn SubagentTraceSink>>,
21}
22
23impl SubagentDispatcherWorkflowExecutor {
24    pub fn new(dispatcher: Arc<dyn SubagentDispatcher>) -> Self {
25        Self {
26            dispatcher,
27            trace_sink: None,
28        }
29    }
30
31    pub fn with_trace_sink(mut self, trace_sink: Arc<dyn SubagentTraceSink>) -> Self {
32        self.trace_sink = Some(trace_sink);
33        self
34    }
35}
36
37#[async_trait]
38impl WorkflowAgentExecutor for SubagentDispatcherWorkflowExecutor {
39    async fn execute_agent(
40        &self,
41        context: WorkflowAgentExecutionContext,
42        request: WorkflowAgentExecutionRequest,
43    ) -> anyhow::Result<SubagentResult> {
44        let parent_thread = context
45            .thread_id
46            .clone()
47            .unwrap_or_else(|| context.run_id.clone());
48        let parent_turn = context
49            .turn_id
50            .clone()
51            .unwrap_or_else(|| context.agent_id.clone());
52        self.dispatcher
53            .dispatch_traced(
54                parent_thread,
55                parent_turn,
56                request.subagent_request,
57                self.trace_sink.clone(),
58            )
59            .await
60    }
61}