Skip to main content

roder_dynamic_workflows/execution/
task.rs

1use 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}