roder_dynamic_workflows/execution/
task.rs1use async_trait::async_trait;
2use roder_api::extension::TaskExecutorId;
3use roder_api::tasks::{TaskExecutionContext, TaskExecutionResult, TaskExecutor, TaskSpec};
4
5use super::{WorkflowRunRequest, WorkflowRunner};
6
7pub const WORKFLOW_TASK_EXECUTOR_ID: &str = "dynamic-workflow";
8
9#[derive(Clone)]
10pub struct WorkflowTaskExecutor {
11 runner: WorkflowRunner,
12}
13
14impl WorkflowTaskExecutor {
15 pub fn new(runner: WorkflowRunner) -> Self {
16 Self { runner }
17 }
18}
19
20#[async_trait]
21impl TaskExecutor for WorkflowTaskExecutor {
22 fn id(&self) -> TaskExecutorId {
23 WORKFLOW_TASK_EXECUTOR_ID.to_string()
24 }
25
26 fn spec(&self) -> TaskSpec {
27 TaskSpec {
28 kind: "dynamic_workflow".to_string(),
29 description: "Run a dynamic workflow in the background.".to_string(),
30 input_schema: serde_json::json!({
31 "type": "object",
32 "required": ["runId", "script"],
33 "properties": {
34 "runId": { "type": "string" },
35 "threadId": { "type": ["string", "null"] },
36 "turnId": { "type": ["string", "null"] },
37 "script": { "type": "object" },
38 "arguments": { "type": "object" },
39 "startPaused": { "type": "boolean" }
40 }
41 }),
42 default_timeout_seconds: None,
43 metadata: serde_json::json!({ "kind": "dynamic_workflow" }),
44 }
45 }
46
47 async fn execute(
48 &self,
49 ctx: TaskExecutionContext,
50 input: serde_json::Value,
51 ) -> anyhow::Result<TaskExecutionResult> {
52 let mut request: WorkflowRunRequest = serde_json::from_value(input)?;
53 request.thread_id = request.thread_id.or(ctx.thread_id);
54 request.turn_id = request.turn_id.or(ctx.turn_id);
55 let handle = self.runner.start(request).await?;
56 let snapshot = handle.wait().await?;
57 Ok(TaskExecutionResult::success(serde_json::json!({
58 "run": snapshot.run,
59 "report": snapshot.report,
60 "reusedAgentResults": snapshot.reused_agent_results
61 })))
62 }
63}