1use 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
39pub const DYNAMIC_WORKFLOW_STORE_RELATIVE_PATH: &str = ".a3s/workflow";
41
42pub 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#[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#[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(®istry);
78 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
403pub 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 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
582async 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;