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