Skip to main content

a3s_code_core/
dynamic_workflow.rs

1//! A3S Flow-backed dynamic workflow runtime.
2//!
3//! `DynamicWorkflowRuntime` lets hosts run a sandboxed PTC script as an A3S
4//! Flow runtime. Flow owns durable replay and step lifecycle; A3S Code's
5//! existing `program` tool remains the sandbox and tool-call boundary.
6
7use crate::tools::{
8    registry_tool_invoker, Tool, ToolContext, ToolInvoker, ToolOutput, ToolRegistry, ToolResult,
9};
10use crate::{
11    agent::AgentEvent,
12    flow_graph::FlowGraphObserver,
13    planning::{Complexity, ExecutionPlan, Task, TaskStatus},
14};
15use a3s_flow::{
16    FanoutFlowEventObserver, FlowEngine, FlowEvent, FlowEventEnvelope, FlowEventObserver,
17    FlowEventStore, FlowRuntime, InMemoryEventStore, LocalFileEventStore, RuntimeCommand,
18    StepInvocation, StepStatus, WorkflowInvocation, WorkflowRunSnapshot, WorkflowRunStatus,
19    WorkflowSpec,
20};
21use anyhow::{Context, Result};
22use async_trait::async_trait;
23use chrono::Utc;
24use serde::{Deserialize, Serialize};
25use serde_json::{json, Map, Value};
26use std::collections::BTreeSet;
27use std::path::{Path, PathBuf};
28use std::sync::Arc;
29use std::time::Duration;
30use tokio::sync::{broadcast, Mutex};
31
32const DYNAMIC_WORKFLOW_TOOL: &str = "dynamic_workflow";
33const GENERATE_OBJECT_TOOL: &str = "generate_object";
34const PROGRAM_TOOL: &str = "program";
35const PARALLEL_TASK_TOOL: &str = "parallel_task";
36const MAX_INLINE_RETRY_RESUMES: usize = 8;
37const MAX_INLINE_RETRY_DELAY: Duration = Duration::from_secs(5);
38
39/// Project-relative directory used for durable dynamic workflow history.
40pub const DYNAMIC_WORKFLOW_STORE_RELATIVE_PATH: &str = ".a3s/workflow";
41
42/// Resolve the durable dynamic workflow history directory for a local workspace.
43pub fn dynamic_workflow_store_path(workspace_root: impl AsRef<Path>) -> PathBuf {
44    workspace_root
45        .as_ref()
46        .join(DYNAMIC_WORKFLOW_STORE_RELATIVE_PATH)
47}
48
49/// Limits forwarded to the underlying PTC `program` tool.
50#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
51#[serde(rename_all = "camelCase")]
52pub struct DynamicWorkflowScriptLimits {
53    #[serde(skip_serializing_if = "Option::is_none")]
54    pub timeout_ms: Option<u64>,
55    #[serde(skip_serializing_if = "Option::is_none")]
56    pub max_tool_calls: Option<usize>,
57    #[serde(skip_serializing_if = "Option::is_none")]
58    pub max_output_bytes: Option<usize>,
59}
60
61/// Runs A3S Flow workflow and step invocations through a sandboxed PTC script.
62#[derive(Clone)]
63pub struct DynamicWorkflowRuntime {
64    invoker: Arc<dyn ToolInvoker>,
65    context: ToolContext,
66    source: Arc<str>,
67    allowed_tools: Vec<String>,
68    limits: DynamicWorkflowScriptLimits,
69}
70
71impl DynamicWorkflowRuntime {
72    pub fn new(
73        registry: Arc<ToolRegistry>,
74        context: ToolContext,
75        source: impl Into<String>,
76    ) -> Self {
77        let allowed_tools = default_allowed_tools(&registry);
78        // Session/agent callers install the governed gateway in ToolContext.
79        // The raw registry adapter is retained only for explicit low-level
80        // callers that construct this public runtime outside an AgentSession.
81        let invoker = context
82            .tool_invoker()
83            .unwrap_or_else(|| registry_tool_invoker(registry));
84        Self {
85            invoker,
86            context,
87            source: Arc::from(source.into()),
88            allowed_tools,
89            limits: DynamicWorkflowScriptLimits::default(),
90        }
91    }
92
93    pub fn with_allowed_tools(mut self, allowed_tools: impl IntoIterator<Item = String>) -> Self {
94        self.allowed_tools = sanitize_allowed_tools(allowed_tools);
95        self
96    }
97
98    pub fn with_limits(mut self, limits: DynamicWorkflowScriptLimits) -> Self {
99        self.limits = limits;
100        self
101    }
102
103    async fn run_script(
104        &self,
105        payload: Value,
106        context: &ToolContext,
107    ) -> a3s_flow::Result<ToolResult> {
108        let mut args = json!({
109            "type": "script",
110            "language": "javascript",
111            "source": self.source.as_ref(),
112            "inputs": payload,
113            "allowed_tools": self.allowed_tools,
114        });
115        if let Some(object) = args.as_object_mut() {
116            if let Ok(Value::Object(limits)) = serde_json::to_value(&self.limits) {
117                if !limits.is_empty() {
118                    object.insert("limits".to_string(), Value::Object(limits));
119                }
120            }
121        }
122
123        let result = self
124            .invoker
125            .invoke(context.nested_tool_invocation(PROGRAM_TOOL, args), context)
126            .await;
127        if result.exit_code != 0 {
128            return Err(a3s_flow::FlowError::Runtime(result.output));
129        }
130        Ok(result)
131    }
132
133    async fn context_for_step(&self, step_name: &str) -> a3s_flow::Result<ToolContext> {
134        if step_name != GENERATE_OBJECT_TOOL {
135            return Ok(self.context.clone());
136        }
137        let Some(admission) = self.context.model_generation_admission() else {
138            return Ok(self.context.clone());
139        };
140        let permit = admission
141            .acquire(&self.context.cancellation_token())
142            .await
143            .map_err(|error| {
144                a3s_flow::FlowError::Runtime(format!(
145                    "model-generation admission failed before workflow step: {error}"
146                ))
147            })?;
148        self.context
149            .clone()
150            .with_model_generation_permit(admission, Arc::new(permit))
151            .map_err(|error| {
152                a3s_flow::FlowError::Runtime(format!(
153                    "bind model-generation admission to workflow step: {error}"
154                ))
155            })
156    }
157
158    async fn run_tool_step(&self, tool_name: &str, args: Value) -> a3s_flow::Result<Value> {
159        let result = self
160            .invoker
161            .invoke(
162                self.context
163                    .nested_tool_invocation(tool_name.to_string(), args),
164                &self.context,
165            )
166            .await;
167        if result.exit_code != 0 {
168            return Err(a3s_flow::FlowError::Runtime(result.output));
169        }
170        Ok(json!({
171            "tool": result.name,
172            "output": result.output,
173            "exit_code": result.exit_code,
174            "metadata": result.metadata,
175        }))
176    }
177}
178
179#[async_trait]
180impl FlowRuntime for DynamicWorkflowRuntime {
181    async fn run_workflow(
182        &self,
183        invocation: WorkflowInvocation,
184    ) -> a3s_flow::Result<RuntimeCommand> {
185        let payload = invocation_payload("workflow", &invocation.run_id, &invocation.history)
186            .with("input", invocation.input);
187        let result = self.run_script(payload.into_value(), &self.context).await?;
188        serde_json::from_value(script_result(&result)?).map_err(a3s_flow::FlowError::from)
189    }
190
191    async fn run_step(&self, invocation: StepInvocation) -> a3s_flow::Result<Value> {
192        if invocation.step_name == PARALLEL_TASK_TOOL {
193            return self
194                .run_tool_step(PARALLEL_TASK_TOOL, invocation.input)
195                .await;
196        }
197
198        let context = self.context_for_step(&invocation.step_name).await?;
199        let payload = invocation_payload("step", &invocation.run_id, &invocation.history)
200            .with("step_id", invocation.step_id)
201            .with("step_name", invocation.step_name)
202            .with("input", invocation.input);
203        let result = self.run_script(payload.into_value(), &context).await?;
204        script_result(&result)
205    }
206}
207
208struct WorkflowProgressState {
209    tasks: Vec<Task>,
210}
211
212impl WorkflowProgressState {
213    fn new() -> Self {
214        Self { tasks: Vec::new() }
215    }
216
217    fn upsert_step(
218        &mut self,
219        step_id: &str,
220        step_name: &str,
221        input: Option<&Value>,
222        status: TaskStatus,
223    ) {
224        let content = workflow_step_description(step_id, step_name, input);
225        if let Some(task) = self.tasks.iter_mut().find(|task| task.id == step_id) {
226            task.content = content;
227            task.status = status;
228            task.tool = Some(step_name.to_string());
229        } else {
230            self.tasks
231                .push(Task::new(step_id.to_string(), content).with_tool(step_name));
232            if let Some(task) = self.tasks.last_mut() {
233                task.status = status;
234            }
235        }
236    }
237
238    fn mark_status(&mut self, step_id: &str, status: TaskStatus) {
239        if let Some(task) = self.tasks.iter_mut().find(|task| task.id == step_id) {
240            task.status = status;
241        }
242    }
243
244    fn step_position(&self, step_id: &str) -> (usize, usize) {
245        let total = self.tasks.len().max(1);
246        let number = self
247            .tasks
248            .iter()
249            .position(|task| task.id == step_id)
250            .map(|idx| idx + 1)
251            .unwrap_or(total);
252        (number, total)
253    }
254
255    fn step_description(&self, step_id: &str) -> String {
256        self.tasks
257            .iter()
258            .find(|task| task.id == step_id)
259            .map(|task| task.content.clone())
260            .unwrap_or_else(|| step_id.to_string())
261    }
262}
263
264struct AgentEventFlowObserver {
265    tx: broadcast::Sender<AgentEvent>,
266    session_id: String,
267    state: Mutex<WorkflowProgressState>,
268}
269
270impl AgentEventFlowObserver {
271    fn new(tx: broadcast::Sender<AgentEvent>, session_id: String) -> Self {
272        Self {
273            tx,
274            session_id,
275            state: Mutex::new(WorkflowProgressState::new()),
276        }
277    }
278
279    fn emit_task_update(&self, tasks: &[Task]) {
280        let _ = self.tx.send(AgentEvent::TaskUpdated {
281            session_id: self.session_id.clone(),
282            tasks: tasks.to_vec(),
283        });
284    }
285}
286
287#[async_trait]
288impl FlowEventObserver for AgentEventFlowObserver {
289    async fn observe(&self, envelope: FlowEventEnvelope) {
290        match envelope.event {
291            FlowEvent::RunStarted => {
292                let _ = self.tx.send(AgentEvent::PlanningStart {
293                    prompt: "dynamic_workflow".to_string(),
294                });
295            }
296            FlowEvent::StepCreated {
297                step_id,
298                step_name,
299                input,
300                ..
301            } => {
302                let mut state = self.state.lock().await;
303                state.upsert_step(&step_id, &step_name, Some(&input), TaskStatus::Pending);
304                self.emit_task_update(&state.tasks);
305                let mut plan = ExecutionPlan::new("dynamic workflow", Complexity::Medium);
306                for task in state.tasks.iter().cloned() {
307                    plan.add_step(task);
308                }
309                let _ = self.tx.send(AgentEvent::PlanningEnd {
310                    estimated_steps: plan.steps.len(),
311                    plan,
312                });
313            }
314            FlowEvent::StepStarted { step_id, .. } => {
315                let mut state = self.state.lock().await;
316                state.mark_status(&step_id, TaskStatus::InProgress);
317                self.emit_task_update(&state.tasks);
318                let (step_number, total_steps) = state.step_position(&step_id);
319                let _ = self.tx.send(AgentEvent::StepStart {
320                    description: state.step_description(&step_id),
321                    step_id,
322                    step_number,
323                    total_steps,
324                });
325            }
326            FlowEvent::StepCompleted { step_id, .. } => {
327                let mut state = self.state.lock().await;
328                state.mark_status(&step_id, TaskStatus::Completed);
329                self.emit_task_update(&state.tasks);
330                let (step_number, total_steps) = state.step_position(&step_id);
331                let _ = self.tx.send(AgentEvent::StepEnd {
332                    step_id,
333                    status: TaskStatus::Completed,
334                    step_number,
335                    total_steps,
336                });
337            }
338            FlowEvent::StepRetrying { step_id, .. } => {
339                let mut state = self.state.lock().await;
340                state.mark_status(&step_id, TaskStatus::InProgress);
341                self.emit_task_update(&state.tasks);
342            }
343            FlowEvent::StepFailed { step_id, .. } => {
344                let mut state = self.state.lock().await;
345                state.mark_status(&step_id, TaskStatus::Failed);
346                self.emit_task_update(&state.tasks);
347                let (step_number, total_steps) = state.step_position(&step_id);
348                let _ = self.tx.send(AgentEvent::StepEnd {
349                    step_id,
350                    status: TaskStatus::Failed,
351                    step_number,
352                    total_steps,
353                });
354            }
355            FlowEvent::RunFailed { .. } => {
356                let mut state = self.state.lock().await;
357                for task in &mut state.tasks {
358                    if task.status.is_active() {
359                        task.status = TaskStatus::Failed;
360                    }
361                }
362                self.emit_task_update(&state.tasks);
363            }
364            FlowEvent::RunCancelled { .. } => {
365                let mut state = self.state.lock().await;
366                for task in &mut state.tasks {
367                    if task.status.is_active() {
368                        task.status = TaskStatus::Cancelled;
369                    }
370                }
371                self.emit_task_update(&state.tasks);
372            }
373            _ => {}
374        }
375    }
376}
377
378fn workflow_step_description(step_id: &str, step_name: &str, input: Option<&Value>) -> String {
379    if step_name == PARALLEL_TASK_TOOL {
380        let count = input
381            .and_then(|value| value.get("tasks"))
382            .and_then(Value::as_array)
383            .map(Vec::len)
384            .unwrap_or(0);
385        if count > 0 {
386            return format!("Fan out {count} parallel subagent task(s)");
387        }
388    }
389
390    input
391        .and_then(|value| value.get("description").or_else(|| value.get("title")))
392        .and_then(Value::as_str)
393        .map(ToString::to_string)
394        .unwrap_or_else(|| {
395            if step_name == step_id {
396                step_id.to_string()
397            } else {
398                format!("{step_name}: {step_id}")
399            }
400        })
401}
402
403/// Model-visible tool that executes a dynamic workflow through A3S Flow.
404pub struct DynamicWorkflowTool {
405    registry: Arc<ToolRegistry>,
406    graph_observer: Option<FlowGraphObserver>,
407}
408
409impl DynamicWorkflowTool {
410    pub fn new(registry: Arc<ToolRegistry>) -> Self {
411        Self {
412            registry,
413            graph_observer: None,
414        }
415    }
416
417    /// Project committed Flow events into an optional reactive state graph.
418    /// A3S Flow remains the workflow execution source of truth.
419    pub fn with_graph_observer(mut self, observer: FlowGraphObserver) -> Self {
420        self.graph_observer = Some(observer);
421        self
422    }
423}
424
425#[async_trait]
426impl Tool for DynamicWorkflowTool {
427    fn name(&self) -> &str {
428        DYNAMIC_WORKFLOW_TOOL
429    }
430
431    fn description(&self) -> &str {
432        "Run a local dynamic workflow with A3S Flow. The workflow source is a sandboxed JavaScript PTC script that may call allowed ctx tools; A3S Flow records workflow and step history."
433    }
434
435    fn parameters(&self) -> Value {
436        json!({
437            "type": "object",
438            "additionalProperties": false,
439            "properties": {
440                "source": {
441                    "type": "string",
442                    "description": "JavaScript PTC source defining async function run(ctx, inputs). For inputs.kind='workflow', return a Flow command: {type:'complete', output}, {type:'fail', error}, {type:'schedule_step', step_id, step_name, input, retry?}, or {type:'schedule_steps', steps:[...]}. For inputs.kind='step', return the step JSON output. A scheduled step with step_name='parallel_task' bypasses QuickJS and calls the host parallel_task tool directly with input as its arguments."
443                },
444                "input": {
445                    "type": "object",
446                    "description": "Initial workflow input."
447                },
448                "run_id": {
449                    "type": "string",
450                    "description": "Optional durable run id. Reusing it with the same source and input is idempotent."
451                },
452                "allowed_tools": {
453                    "type": "array",
454                    "description": "Tool names the workflow script may call through ctx. Defaults to all registered tools except program, dynamic_workflow, and parallel_task. Login-registered tools such as runtime are allowed when present.",
455                    "items": { "type": "string" }
456                },
457                "limits": {
458                    "type": "object",
459                    "additionalProperties": false,
460                    "properties": {
461                        "timeoutMs": { "type": "integer", "minimum": 1 },
462                        "maxToolCalls": { "type": "integer", "minimum": 1 },
463                        "maxOutputBytes": { "type": "integer", "minimum": 1 }
464                    }
465                }
466            },
467            "required": ["source"]
468        })
469    }
470
471    async fn execute(&self, args: &Value, ctx: &ToolContext) -> Result<ToolOutput> {
472        let Some(source) = args.get("source").and_then(Value::as_str) else {
473            return Ok(ToolOutput::error("dynamic_workflow requires source"));
474        };
475        let input = args.get("input").cloned().unwrap_or_else(|| json!({}));
476        let allowed_tools = args
477            .get("allowed_tools")
478            .and_then(Value::as_array)
479            .map(|items| {
480                items
481                    .iter()
482                    .filter_map(Value::as_str)
483                    .map(ToString::to_string)
484                    .collect::<Vec<_>>()
485            })
486            .unwrap_or_else(|| default_allowed_tools(&self.registry));
487        let limits = args
488            .get("limits")
489            .cloned()
490            .and_then(|value| serde_json::from_value(value).ok())
491            .unwrap_or_default();
492
493        let runtime = Arc::new(
494            DynamicWorkflowRuntime::new(Arc::clone(&self.registry), ctx.clone(), source)
495                .with_allowed_tools(allowed_tools)
496                .with_limits(limits),
497        );
498        let requested_run_id = args.get("run_id").and_then(Value::as_str);
499        let store = match flow_store_for_context(ctx, requested_run_id).await {
500            Ok(store) => store,
501            Err(error) => return Ok(ToolOutput::error(error.to_string())),
502        };
503        let mut observers: Vec<Arc<dyn FlowEventObserver>> = Vec::new();
504        if let Some(tx) = ctx.agent_event_tx.clone() {
505            observers.push(Arc::new(AgentEventFlowObserver::new(
506                tx,
507                ctx.session_id.clone().unwrap_or_default(),
508            )));
509        }
510        if let Some(observer) = &self.graph_observer {
511            observers.push(Arc::new(observer.clone()));
512        }
513        let engine = if observers.is_empty() {
514            FlowEngine::new(store, runtime)
515        } else {
516            FlowEngine::builder(runtime)
517                .with_store(store)
518                .with_observer(Arc::new(FanoutFlowEventObserver::from_observers(observers)))
519                .build()
520        };
521        let source_hash = source_hash(source);
522        let spec = WorkflowSpec::rust_embedded(
523            "a3s-code.dynamic-workflow",
524            source_hash.as_str(),
525            "ptc",
526            "run",
527        );
528
529        let run_id = match requested_run_id {
530            Some(run_id) => match engine.start_with_id(run_id, spec, input).await {
531                Ok(run_id) => run_id,
532                Err(err) => return Ok(ToolOutput::error(err.to_string())),
533            },
534            None => match engine.start(spec, input).await {
535                Ok(run_id) => run_id,
536                Err(err) => return Ok(ToolOutput::error(err.to_string())),
537            },
538        };
539
540        let snapshot = match drive_inline_retries(&engine, &run_id, ctx).await {
541            Ok(snapshot) => snapshot,
542            Err(err) => return Ok(ToolOutput::error(err.to_string())),
543        };
544        let history = match engine.history(&run_id).await {
545            Ok(history) => history,
546            Err(err) => return Ok(ToolOutput::error(err.to_string())),
547        };
548
549        let output = match &snapshot.output {
550            Some(output) => {
551                serde_json::to_string_pretty(output).unwrap_or_else(|_| output.to_string())
552            }
553            None => snapshot
554                .error
555                .clone()
556                .unwrap_or_else(|| format!("workflow status: {:?}", snapshot.status)),
557        };
558
559        let status = snapshot.status;
560        let metadata = json!({
561            "dynamic_workflow": {
562                "run_id": run_id,
563                "status": format!("{:?}", snapshot.status),
564                "last_sequence": snapshot.last_sequence,
565                "source_hash": source_hash,
566                "snapshot": snapshot,
567                "history": history,
568            }
569        });
570        let output = match status {
571            WorkflowRunStatus::Completed => ToolOutput::success(output),
572            WorkflowRunStatus::Failed | WorkflowRunStatus::Cancelled => ToolOutput::error(output),
573            _ => ToolOutput::error(format!(
574                "dynamic_workflow ended without a terminal result: {status:?}; {output}"
575            )),
576        };
577
578        Ok(output.with_metadata(metadata))
579    }
580}
581
582/// Drive short, persisted step retries inside the originating tool call.
583///
584/// A3S Flow deliberately suspends at a delayed retry boundary. Interactive
585/// waits and hooks must remain suspended for an external host, but a bounded
586/// retry delay is ordinary fault recovery: returning it as a terminal tool
587/// error forces every caller to reimplement the scheduler and previously made
588/// DeepResearch abandon its event-sourced run. Retry attempts and their delay
589/// remain authoritative in the Flow journal; this helper only waits for the
590/// due time and asks the engine to replay the same run.
591async fn drive_inline_retries(
592    engine: &FlowEngine,
593    run_id: &str,
594    ctx: &ToolContext,
595) -> Result<WorkflowRunSnapshot> {
596    for _ in 0..MAX_INLINE_RETRY_RESUMES {
597        let snapshot = engine.snapshot(run_id).await?;
598        if snapshot.status.is_terminal() {
599            return Ok(snapshot);
600        }
601        let Some(retry_after) = snapshot
602            .steps
603            .values()
604            .filter(|step| step.status == StepStatus::Pending)
605            .filter_map(|step| step.retry_after)
606            .min()
607        else {
608            return Ok(snapshot);
609        };
610        let delay = retry_after
611            .signed_duration_since(Utc::now())
612            .to_std()
613            .unwrap_or_default();
614        if delay > MAX_INLINE_RETRY_DELAY {
615            return Ok(snapshot);
616        }
617        let cancellation = ctx.cancellation_token();
618        tokio::select! {
619            biased;
620            _ = cancellation.cancelled() => {
621                anyhow::bail!("dynamic_workflow cancelled while waiting for a scheduled retry");
622            }
623            _ = tokio::time::sleep(delay) => {}
624        }
625        engine.drive(run_id).await?;
626    }
627    engine.snapshot(run_id).await.map_err(Into::into)
628}
629
630pub fn register_dynamic_workflow(registry: &Arc<ToolRegistry>) {
631    registry.register(Arc::new(DynamicWorkflowTool::new(Arc::clone(registry))));
632}
633
634async fn flow_store_for_context(
635    ctx: &ToolContext,
636    requested_run_id: Option<&str>,
637) -> Result<Arc<dyn FlowEventStore>> {
638    match ctx.workspace_services.local_root() {
639        Some(root) => {
640            let store = dynamic_workflow_store_path(root);
641            validate_dynamic_workflow_directory(&root.join(".a3s"), ".a3s").await?;
642            validate_dynamic_workflow_directory(&store, ".a3s/workflow").await?;
643            if let Some(run_id) = requested_run_id.filter(|run_id| safe_workflow_run_id(run_id)) {
644                validate_dynamic_workflow_log(&store.join(format!("{run_id}.jsonl"))).await?;
645            }
646            Ok(Arc::new(LocalFileEventStore::new(store)))
647        }
648        None => Ok(Arc::new(InMemoryEventStore::new())),
649    }
650}
651
652async fn validate_dynamic_workflow_directory(path: &Path, label: &str) -> Result<()> {
653    match tokio::fs::symlink_metadata(path).await {
654        Ok(metadata) if metadata.file_type().is_symlink() => {
655            anyhow::bail!("refusing to use symlinked dynamic workflow directory {label}")
656        }
657        Ok(metadata) if !metadata.is_dir() => {
658            anyhow::bail!("dynamic workflow path {label} exists but is not a directory")
659        }
660        Ok(_) => Ok(()),
661        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
662        Err(error) => Err(error).with_context(|| format!("inspect dynamic workflow path {label}")),
663    }
664}
665
666async fn validate_dynamic_workflow_log(path: &Path) -> Result<()> {
667    match tokio::fs::symlink_metadata(path).await {
668        Ok(metadata) if metadata.file_type().is_symlink() => anyhow::bail!(
669            "refusing to read or append symlinked dynamic workflow history {}",
670            path.display()
671        ),
672        Ok(metadata) if !metadata.is_file() => anyhow::bail!(
673            "dynamic workflow history path {} exists but is not a file",
674            path.display()
675        ),
676        Ok(_) => Ok(()),
677        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
678        Err(error) => Err(error)
679            .with_context(|| format!("inspect dynamic workflow history {}", path.display())),
680    }
681}
682
683fn safe_workflow_run_id(run_id: &str) -> bool {
684    !run_id.is_empty()
685        && run_id
686            .chars()
687            .all(|ch| ch.is_ascii_alphanumeric() || ch == '-' || ch == '_')
688}
689
690struct PayloadBuilder {
691    value: Map<String, Value>,
692}
693
694impl PayloadBuilder {
695    fn with(mut self, key: &str, value: impl Serialize) -> Self {
696        self.value.insert(
697            key.to_string(),
698            serde_json::to_value(value).unwrap_or(Value::Null),
699        );
700        self
701    }
702
703    fn into_value(self) -> Value {
704        Value::Object(self.value)
705    }
706}
707
708fn invocation_payload(kind: &str, run_id: &str, history: &[FlowEventEnvelope]) -> PayloadBuilder {
709    let mut value = Map::new();
710    value.insert("kind".to_string(), json!(kind));
711    value.insert("run_id".to_string(), json!(run_id));
712    value.insert("history".to_string(), json!(history));
713    value.insert("step_outputs".to_string(), completed_step_outputs(history));
714    value.insert("step_failures".to_string(), failed_step_outputs(history));
715    PayloadBuilder { value }
716}
717
718fn completed_step_outputs(history: &[FlowEventEnvelope]) -> Value {
719    let mut outputs = Map::new();
720    for envelope in history {
721        if let FlowEvent::StepCompleted { step_id, output } = &envelope.event {
722            outputs.insert(step_id.clone(), output.clone());
723        }
724    }
725    Value::Object(outputs)
726}
727
728fn failed_step_outputs(history: &[FlowEventEnvelope]) -> Value {
729    let mut outputs = Map::new();
730    for envelope in history {
731        if let FlowEvent::StepFailed {
732            step_id,
733            attempt,
734            error,
735        } = &envelope.event
736        {
737            outputs.insert(
738                step_id.clone(),
739                json!({
740                    "attempt": attempt,
741                    "error": error,
742                }),
743            );
744        }
745    }
746    Value::Object(outputs)
747}
748
749fn script_result(result: &ToolResult) -> a3s_flow::Result<Value> {
750    result
751        .metadata
752        .as_ref()
753        .and_then(|metadata| metadata.get("script_result"))
754        .cloned()
755        .ok_or_else(|| {
756            a3s_flow::FlowError::Runtime(
757                "PTC program result did not include script_result metadata".to_string(),
758            )
759        })
760}
761
762fn default_allowed_tools(registry: &ToolRegistry) -> Vec<String> {
763    sanitize_allowed_tools(registry.list())
764}
765
766fn sanitize_allowed_tools(items: impl IntoIterator<Item = String>) -> Vec<String> {
767    let mut tools = items.into_iter().collect::<BTreeSet<_>>();
768    tools.remove(PROGRAM_TOOL);
769    tools.remove(DYNAMIC_WORKFLOW_TOOL);
770    tools.remove(PARALLEL_TASK_TOOL);
771    tools.into_iter().collect()
772}
773
774fn source_hash(source: &str) -> String {
775    sha256::digest(source.as_bytes())
776}
777
778#[cfg(test)]
779#[path = "dynamic_workflow/tests.rs"]
780mod tests;