Skip to main content

navi_core/tool/builtin/workflow/
mod.rs

1//! Built-in `workflow` tool — sandboxed Lua 5.4 multi-agent orchestration.
2
3mod backends;
4mod journal;
5mod policy;
6mod runtime;
7mod types;
8
9#[cfg(test)]
10mod tests;
11
12pub use backends::{SubagentBridgeBackend, WorkerProbeBackend};
13pub use policy::{
14    AgentPolicyOpts, EffectiveAgentPolicy, MAX_AGENTS_CEILING, MAX_PARALLEL_CEILING, RunPolicy,
15    clamp_max_agents, clamp_max_parallel, default_run_policy, intersect_agent_policy,
16};
17pub use types::{
18    AGENT_RESULT_MAX_BYTES, AgentBackendResult, AgentRequest, DEFAULT_MAX_AGENTS,
19    DEFAULT_MAX_PARALLEL, DEFAULT_MAX_SCRIPT_BYTES, DEFAULT_RUN_TIMEOUT_MS, NESTED_WORKFLOW_TOOLS,
20    WorkflowErrorCode, WorkflowRunStatus, WorkflowStats,
21};
22
23use std::sync::Arc;
24use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
25use std::time::Instant;
26
27use anyhow::Result;
28use async_trait::async_trait;
29use serde_json::{Value, json};
30use tokio::sync::Semaphore;
31
32use self::journal::WorkflowJournal;
33use self::runtime::{LuaRunInput, run_lua_workflow};
34use self::types::*;
35use super::helpers;
36use crate::cancel::CancelToken;
37use crate::config::WorkflowConfig;
38use crate::security::{SecurityPolicy, redact_secrets};
39use crate::tool::{
40    Tool, ToolDefinition, ToolInvocation, ToolInvocationContext, ToolKind, ToolResult,
41};
42
43static RUN_COUNTER: AtomicU64 = AtomicU64::new(1);
44
45/// Pluggable worker backend.
46#[async_trait]
47pub trait AgentBackend: Send + Sync {
48    async fn run_agent(&self, request: AgentRequest) -> AgentBackendResult;
49}
50
51/// Mock backend for unit tests (no live model).
52#[derive(Default)]
53pub struct MockAgentBackend {
54    pub calls: std::sync::Mutex<Vec<AgentRequest>>,
55    pub delay_ms: u64,
56    pub in_flight: Option<Arc<AtomicUsize>>,
57    pub peak_in_flight: Option<Arc<AtomicUsize>>,
58}
59
60#[async_trait]
61impl AgentBackend for MockAgentBackend {
62    async fn run_agent(&self, request: AgentRequest) -> AgentBackendResult {
63        if let Some(ref inflight) = self.in_flight {
64            let n = inflight.fetch_add(1, Ordering::SeqCst) + 1;
65            if let Some(ref peak) = self.peak_in_flight {
66                peak.fetch_max(n, Ordering::SeqCst);
67            }
68        }
69        {
70            let mut guard = self.calls.lock().unwrap_or_else(|e| e.into_inner());
71            guard.push(request.clone());
72        }
73
74        if self.delay_ms > 0 {
75            let delay = std::time::Duration::from_millis(self.delay_ms);
76            tokio::select! {
77                _ = tokio::time::sleep(delay) => {}
78                _ = request.cancel_token.notified() => {
79                    if let Some(ref inflight) = self.in_flight {
80                        inflight.fetch_sub(1, Ordering::SeqCst);
81                    }
82                    return AgentBackendResult {
83                        ok: false,
84                        output: json!({"error": "cancelled"}),
85                        error: Some("cancelled".into()),
86                    };
87                }
88            }
89        }
90
91        if request.cancel_token.is_requested() {
92            if let Some(ref inflight) = self.in_flight {
93                inflight.fetch_sub(1, Ordering::SeqCst);
94            }
95            return AgentBackendResult {
96                ok: false,
97                output: json!({"error": "cancelled"}),
98                error: Some("cancelled".into()),
99            };
100        }
101
102        if let Some(ref inflight) = self.in_flight {
103            inflight.fetch_sub(1, Ordering::SeqCst);
104        }
105
106        AgentBackendResult {
107            ok: true,
108            output: json!({
109                "ok": true,
110                "prompt": request.prompt,
111                "label": request.label,
112                "agent_index": request.agent_index,
113                "profile": request.effective.profile,
114                "tools": request.effective.tools,
115                "create_files": request.effective.create_files,
116                "create_dirs": request.effective.create_dirs,
117                "write_allow": request.effective.write_allow,
118                "path_allow": request.effective.path_allow,
119                "path_deny": request.effective.path_deny,
120            }),
121            error: None,
122        }
123    }
124}
125
126/// Strips nested orchestration tools and delegates.
127pub struct PolicyAgentBackend {
128    pub inner: Arc<dyn AgentBackend>,
129}
130
131#[async_trait]
132impl AgentBackend for PolicyAgentBackend {
133    async fn run_agent(&self, request: AgentRequest) -> AgentBackendResult {
134        for banned in NESTED_WORKFLOW_TOOLS {
135            if request.effective.tools.iter().any(|t| t == *banned) {
136                return AgentBackendResult {
137                    ok: false,
138                    output: json!({"error": "policy_denied", "tool": banned}),
139                    error: Some(format!(
140                        "worker must not receive orchestration tool {banned}"
141                    )),
142                };
143            }
144        }
145        self.inner.run_agent(request).await
146    }
147}
148
149/// Built-in workflow tool.
150pub struct WorkflowTool {
151    policy: SecurityPolicy,
152    config: WorkflowConfig,
153    backend: Arc<dyn AgentBackend>,
154}
155
156impl WorkflowTool {
157    /// Default constructor: [`WorkerProbeBackend`] (real SecurityPolicy tool
158    /// filtering, no live model). Production runtimes should prefer
159    /// [`Self::with_subagent_bridge`] once a `ToolExecutor` weak handle exists.
160    pub fn new(policy: SecurityPolicy, config: WorkflowConfig) -> Self {
161        let backend = Arc::new(WorkerProbeBackend::new(policy.clone()));
162        Self {
163            policy,
164            config,
165            backend,
166        }
167    }
168
169    /// Production constructor: each `agent()` runs a real nested `subagent` turn.
170    pub fn with_subagent_bridge(
171        policy: SecurityPolicy,
172        config: WorkflowConfig,
173        tool_executor: std::sync::Weak<crate::tool::ToolExecutor>,
174    ) -> Self {
175        Self {
176            policy,
177            config,
178            backend: Arc::new(SubagentBridgeBackend::new(tool_executor)),
179        }
180    }
181
182    pub fn with_backend(
183        policy: SecurityPolicy,
184        config: WorkflowConfig,
185        backend: Arc<dyn AgentBackend>,
186    ) -> Self {
187        Self {
188            policy,
189            config,
190            backend,
191        }
192    }
193
194    pub fn with_mock(
195        policy: SecurityPolicy,
196        config: WorkflowConfig,
197        mock: MockAgentBackend,
198    ) -> Self {
199        Self {
200            policy,
201            config,
202            backend: Arc::new(mock),
203        }
204    }
205
206    /// Integration tests: SecurityPolicy-probing backend with optional delay.
207    pub fn with_probe(
208        policy: SecurityPolicy,
209        config: WorkflowConfig,
210        probe: WorkerProbeBackend,
211    ) -> Self {
212        Self {
213            policy,
214            config,
215            backend: Arc::new(probe),
216        }
217    }
218}
219
220pub(crate) struct AgentJob {
221    pub request: AgentRequest,
222    pub response: std::sync::mpsc::Sender<AgentBackendResult>,
223}
224
225pub(crate) struct WorkflowHostError {
226    pub code: WorkflowErrorCode,
227    pub message: String,
228    pub hint: Option<String>,
229}
230
231const TOOL_DESCRIPTION: &str = "\
232Run a multi-agent workflow authored as a sandboxed Lua 5.4 script. \
233The script only orchestrates workers; workers perform all filesystem/shell IO.
234
235Entrypoint (primary): define `function workflow() ... end` and return a value.
236
237Host builtins (only these — no require/io/os/debug/JSON.parse):
238  agent(prompt, opts?)  — spawn one worker; blocks until done; returns a table
239  parallel(thunks)      — array of zero-arg functions only; order-preserving results
240  pipeline(items, fn)   — map each ipairs item through fn (may parallelize)
241  phase(title)          — progress boundary
242  log(message)          — progress log
243  args                  — read-only tool args table
244  budget                — {total, spent, remaining}
245
246Default run policy is read-only (explorer): tools like read_file+search, \
247create_files/dirs=false, write_allow={}. Grant writes only via write_allow paths \
248(intersected with run policy). Empty write_allow ⇒ no writes even for implementer.
249
250Caps: default max_parallel=16, max_agents=1000 (clamped ceilings 64 / 5000).
251Workers never get nested `subagent` or `workflow` tools.
252
253Example:
254  function workflow()
255    phase(\"scan\")
256    local hits = pipeline(args.files or {}, function(f)
257      return agent(\"Audit \" .. f, {label = f})
258    end)
259    return { count = #hits, hits = hits }
260  end
261
262Do NOT use require, io, os, package, loadfile, or JSON.parse. \
263Agent results are already Lua tables. \
264Do NOT invent host APIs beyond agent/parallel/pipeline/phase/log/args/budget.";
265
266#[async_trait]
267impl Tool for WorkflowTool {
268    fn definition(&self) -> ToolDefinition {
269        helpers::definition(
270            "workflow",
271            TOOL_DESCRIPTION,
272            ToolKind::Command,
273            json!({
274                "type": "object",
275                "properties": {
276                    "script": {
277                        "type": "string",
278                        "description": "Non-empty Lua 5.4 source. Must define function workflow() or return a value from the chunk."
279                    },
280                    "args": {
281                        "type": "object",
282                        "description": "JSON object injected as read-only Lua global `args`."
283                    },
284                    "max_parallel": {
285                        "type": "integer",
286                        "description": "Max concurrent workers (default 16, ceiling 64)."
287                    },
288                    "max_agents": {
289                        "type": "integer",
290                        "description": "Max agents per run (default 1000, ceiling 5000)."
291                    },
292                    "timeout_ms": {
293                        "type": "integer",
294                        "description": "Wall-clock timeout for the entire run in milliseconds."
295                    },
296                    "name": {
297                        "type": "string",
298                        "description": "Optional label for UI / journal."
299                    },
300                    "policy": {
301                        "type": "object",
302                        "description": "Run-level default agent policy.",
303                        "properties": {
304                            "profile": { "type": "string" },
305                            "tools": { "type": "array", "items": { "type": "string" } },
306                            "path_allow": { "type": "array", "items": { "type": "string" } },
307                            "path_deny": { "type": "array", "items": { "type": "string" } },
308                            "create_files": { "type": "boolean" },
309                            "create_dirs": { "type": "boolean" },
310                            "write_allow": { "type": "array", "items": { "type": "string" } },
311                            "approval": { "type": "string" }
312                        }
313                    },
314                    "resume_from_run_id": {
315                        "type": "string",
316                        "description": "Resume is not implemented in v1."
317                    }
318                },
319                "required": ["script"],
320                "additionalProperties": false
321            }),
322        )
323    }
324
325    async fn invoke(&self, invocation: ToolInvocation) -> Result<ToolResult> {
326        self.invoke_with_context(invocation, ToolInvocationContext::default())
327            .await
328    }
329
330    async fn invoke_with_context(
331        &self,
332        invocation: ToolInvocation,
333        context: ToolInvocationContext,
334    ) -> Result<ToolResult> {
335        let started = Instant::now();
336        let invocation_id = invocation.id.clone();
337
338        if !self.config.enabled {
339            return Ok(fail(
340                &invocation_id,
341                None,
342                WorkflowRunStatus::Failed,
343                WorkflowErrorCode::PolicyDenied,
344                "workflow tool is disabled in config",
345                Some("Set [workflow] enabled = true."),
346                WorkflowStats::default(),
347                None,
348            ));
349        }
350
351        if helpers::optional_string(&invocation.input, "resume_from_run_id").is_some() {
352            return Ok(fail(
353                &invocation_id,
354                None,
355                WorkflowRunStatus::Failed,
356                WorkflowErrorCode::NotImplemented,
357                "resume_from_run_id is not implemented in v1",
358                None,
359                WorkflowStats::default(),
360                None,
361            ));
362        }
363
364        let script = match helpers::required_string(&invocation.input, "script") {
365            Ok(s) if !s.trim().is_empty() => s.to_string(),
366            Ok(_) => {
367                return Ok(fail(
368                    &invocation_id,
369                    None,
370                    WorkflowRunStatus::Failed,
371                    WorkflowErrorCode::InvalidHostCall,
372                    "script must be non-empty",
373                    None,
374                    WorkflowStats::default(),
375                    None,
376                ));
377            }
378            Err(err) => {
379                return Ok(fail(
380                    &invocation_id,
381                    None,
382                    WorkflowRunStatus::Failed,
383                    WorkflowErrorCode::InvalidHostCall,
384                    &format!("missing or invalid script: {err}"),
385                    Some("Provide a non-empty Lua script string."),
386                    WorkflowStats::default(),
387                    None,
388                ));
389            }
390        };
391
392        let max_script = if self.config.max_script_bytes == 0 {
393            DEFAULT_MAX_SCRIPT_BYTES
394        } else {
395            self.config.max_script_bytes
396        };
397        if script.len() > max_script {
398            return Ok(fail(
399                &invocation_id,
400                None,
401                WorkflowRunStatus::Failed,
402                WorkflowErrorCode::ScriptTooLarge,
403                &format!("script is {} bytes; max is {max_script}", script.len()),
404                Some("Shorten the Lua script."),
405                WorkflowStats::default(),
406                None,
407            ));
408        }
409
410        let max_parallel = clamp_max_parallel(
411            optional_usize(&invocation.input, "max_parallel")
412                .unwrap_or(self.config.max_parallel.max(1)),
413        );
414        let max_agents = clamp_max_agents(
415            optional_usize(&invocation.input, "max_agents")
416                .unwrap_or(self.config.max_agents.max(1)),
417        );
418        let timeout_ms = optional_u64(&invocation.input, "timeout_ms").unwrap_or(
419            if self.config.run_timeout_ms == 0 {
420                DEFAULT_RUN_TIMEOUT_MS
421            } else {
422                self.config.run_timeout_ms
423            },
424        );
425        let name = helpers::optional_string(&invocation.input, "name");
426        let args = invocation
427            .input
428            .get("args")
429            .cloned()
430            .unwrap_or_else(|| json!({}));
431        let run_policy = parse_run_policy(invocation.input.get("policy"));
432
433        let run_id = new_run_id();
434        let journal_dir = self.policy.data_dir().join("workflows").join(&run_id);
435        let mut journal = match WorkflowJournal::create(&journal_dir, &run_id, name.as_deref()) {
436            Ok(j) => j,
437            Err(err) => {
438                return Ok(fail(
439                    &invocation_id,
440                    Some(run_id),
441                    WorkflowRunStatus::Failed,
442                    WorkflowErrorCode::ScriptRuntimeError,
443                    &format!("failed to create journal: {err}"),
444                    None,
445                    WorkflowStats::default(),
446                    None,
447                ));
448            }
449        };
450        let _ = journal.write_meta_start(
451            &script,
452            &args,
453            max_parallel,
454            max_agents,
455            self.policy.project_root(),
456        );
457
458        let cancel_token = context
459            .cancel_token
460            .clone()
461            .unwrap_or_else(CancelToken::new);
462        let semaphore = Arc::new(Semaphore::new(max_parallel.max(1)));
463        let backend: Arc<dyn AgentBackend> = Arc::new(PolicyAgentBackend {
464            inner: self.backend.clone(),
465        });
466
467        let (job_tx, mut job_rx) = tokio::sync::mpsc::unbounded_channel::<AgentJob>();
468        let (lua_done_tx, lua_done_rx) =
469            tokio::sync::oneshot::channel::<Result<runtime::LuaRunOutcome, WorkflowHostError>>();
470
471        let journal_path = journal.journal_path().to_path_buf();
472        let stats = Arc::new(std::sync::Mutex::new(WorkflowStats::default()));
473        let in_flight = Arc::new(AtomicUsize::new(0));
474
475        let lua_input = LuaRunInput {
476            script,
477            args,
478            run_policy,
479            max_agents,
480            max_parallel,
481            job_tx,
482            cancel_token: cancel_token.clone(),
483        };
484
485        let _lua_thread = std::thread::Builder::new()
486            .name("navi-workflow-lua".into())
487            .spawn(move || {
488                let outcome = run_lua_workflow(lua_input);
489                let _ = lua_done_tx.send(outcome);
490            })
491            .map_err(|e| anyhow::anyhow!("spawn lua thread: {e}"))?;
492
493        let stats_j = stats.clone();
494        let cancel_j = cancel_token.clone();
495        let in_flight_j = in_flight.clone();
496        let job_loop_handle = tokio::spawn(async move {
497            let mut handles = Vec::new();
498            while let Some(job) = job_rx.recv().await {
499                if cancel_j.is_requested() {
500                    let _ = job.response.send(AgentBackendResult {
501                        ok: false,
502                        output: json!({"error": "cancelled"}),
503                        error: Some("cancelled".into()),
504                    });
505                    continue;
506                }
507                let permit = match semaphore.clone().acquire_owned().await {
508                    Ok(p) => p,
509                    Err(_) => {
510                        let _ = job.response.send(AgentBackendResult {
511                            ok: false,
512                            output: json!({"error": "semaphore closed"}),
513                            error: Some("semaphore closed".into()),
514                        });
515                        continue;
516                    }
517                };
518                let backend = backend.clone();
519                let stats_j = stats_j.clone();
520                let journal_path = journal_path.clone();
521                let cancel_j = cancel_j.clone();
522                let in_flight_j = in_flight_j.clone();
523                handles.push(tokio::spawn(async move {
524                    let n = in_flight_j.fetch_add(1, Ordering::SeqCst) + 1;
525                    {
526                        let mut s = stats_j.lock().unwrap_or_else(|e| e.into_inner());
527                        s.agents_started += 1;
528                        s.max_parallel_used = s.max_parallel_used.max(n);
529                    }
530                    let agent_index = job.request.agent_index;
531                    let label = job.request.label.clone();
532                    let prompt = job.request.prompt.clone();
533                    append_journal_line(
534                        &journal_path,
535                        &json!({
536                            "event": "agent_started",
537                            "agent_index": agent_index,
538                            "label": label,
539                            "prompt": redact_secrets(&prompt),
540                        }),
541                    );
542                    let mut req = job.request;
543                    req.cancel_token = cancel_j;
544                    let mut result = backend.run_agent(req).await;
545                    result.output = truncate_json(result.output, AGENT_RESULT_MAX_BYTES);
546                    {
547                        let mut s = stats_j.lock().unwrap_or_else(|e| e.into_inner());
548                        if result.ok {
549                            s.agents_completed += 1;
550                        } else {
551                            s.agents_failed += 1;
552                        }
553                    }
554                    append_journal_line(
555                        &journal_path,
556                        &json!({
557                            "event": "agent_completed",
558                            "agent_index": agent_index,
559                            "ok": result.ok,
560                        }),
561                    );
562                    in_flight_j.fetch_sub(1, Ordering::SeqCst);
563                    let _ = job.response.send(result);
564                    drop(permit);
565                }));
566            }
567            for h in handles {
568                let _ = h.await;
569            }
570        });
571
572        let timeout = std::time::Duration::from_millis(timeout_ms.max(1));
573        enum WaitKind {
574            Cancelled,
575            TimedOut,
576            Lua(
577                Result<
578                    Result<runtime::LuaRunOutcome, WorkflowHostError>,
579                    tokio::sync::oneshot::error::RecvError,
580                >,
581            ),
582        }
583        let wait = tokio::select! {
584            biased;
585            _ = cancel_token.notified() => WaitKind::Cancelled,
586            _ = tokio::time::sleep(timeout) => WaitKind::TimedOut,
587            outcome = lua_done_rx => WaitKind::Lua(outcome),
588        };
589        let finish = match wait {
590            WaitKind::Cancelled => {
591                cancel_token.cancel();
592                let _ =
593                    tokio::time::timeout(std::time::Duration::from_secs(5), job_loop_handle).await;
594                Finish::Cancelled
595            }
596            WaitKind::TimedOut => {
597                cancel_token.cancel();
598                let _ =
599                    tokio::time::timeout(std::time::Duration::from_secs(5), job_loop_handle).await;
600                Finish::TimedOut
601            }
602            WaitKind::Lua(outcome) => {
603                let _ =
604                    tokio::time::timeout(std::time::Duration::from_secs(30), job_loop_handle).await;
605                match outcome {
606                    Ok(Ok(o)) => Finish::Lua(o),
607                    Ok(Err(e)) => Finish::Err(e),
608                    Err(_) => Finish::Err(WorkflowHostError {
609                        code: WorkflowErrorCode::ScriptRuntimeError,
610                        message: "workflow Lua task dropped".into(),
611                        hint: None,
612                    }),
613                }
614            }
615        };
616
617        let mut final_stats = stats.lock().unwrap_or_else(|e| e.into_inner()).clone();
618        final_stats.elapsed_ms = started.elapsed().as_millis() as u64;
619        let journal_path_str = journal.journal_path().display().to_string();
620
621        let tool_result = match finish {
622            Finish::Cancelled => {
623                final_stats.phases = journal.take_phases();
624                let _ = journal.finalize(&run_id, WorkflowRunStatus::Cancelled, &final_stats, None);
625                fail(
626                    &invocation_id,
627                    Some(run_id),
628                    WorkflowRunStatus::Cancelled,
629                    WorkflowErrorCode::Cancelled,
630                    "workflow cancelled",
631                    None,
632                    final_stats,
633                    Some(journal_path_str),
634                )
635            }
636            Finish::TimedOut => {
637                final_stats.phases = journal.take_phases();
638                let _ = journal.finalize(&run_id, WorkflowRunStatus::TimedOut, &final_stats, None);
639                fail(
640                    &invocation_id,
641                    Some(run_id),
642                    WorkflowRunStatus::TimedOut,
643                    WorkflowErrorCode::Timeout,
644                    "workflow timed out",
645                    None,
646                    final_stats,
647                    Some(journal_path_str),
648                )
649            }
650            Finish::Err(e) => {
651                final_stats.phases = journal.take_phases();
652                let status = status_for_code(e.code);
653                let _ = journal.finalize(&run_id, status, &final_stats, Some(&e.message));
654                fail(
655                    &invocation_id,
656                    Some(run_id),
657                    status,
658                    e.code,
659                    &e.message,
660                    e.hint.as_deref(),
661                    final_stats,
662                    Some(journal_path_str),
663                )
664            }
665            Finish::Lua(outcome) => {
666                for p in &outcome.phases {
667                    journal.record_phase(p);
668                }
669                for line in &outcome.logs {
670                    journal.record_log(line);
671                }
672                final_stats.phases = outcome.phases.clone();
673                if outcome.agents_started > final_stats.agents_started {
674                    final_stats.agents_started = outcome.agents_started;
675                }
676                final_stats.elapsed_ms = started.elapsed().as_millis() as u64;
677
678                if let Some(err) = outcome.error {
679                    let status = status_for_code(err.code);
680                    let _ = journal.finalize(&run_id, status, &final_stats, Some(&err.message));
681                    fail(
682                        &invocation_id,
683                        Some(run_id),
684                        status,
685                        err.code,
686                        &err.message,
687                        err.hint.as_deref(),
688                        final_stats,
689                        Some(journal_path_str),
690                    )
691                } else {
692                    let _ =
693                        journal.finalize(&run_id, WorkflowRunStatus::Completed, &final_stats, None);
694                    let compact = truncate_json(outcome.result, 32 * 1024);
695                    ToolResult {
696                        invocation_id,
697                        ok: true,
698                        output: json!({
699                            "ok": true,
700                            "run_id": run_id,
701                            "status": WorkflowRunStatus::Completed,
702                            "result": compact,
703                            "stats": final_stats,
704                            "journal_path": journal_path_str,
705                            "error": null,
706                            "name": name,
707                        }),
708                    }
709                }
710            }
711        };
712
713        Ok(tool_result)
714    }
715}
716
717enum Finish {
718    Lua(runtime::LuaRunOutcome),
719    Err(WorkflowHostError),
720    Cancelled,
721    TimedOut,
722}
723
724fn status_for_code(code: WorkflowErrorCode) -> WorkflowRunStatus {
725    match code {
726        WorkflowErrorCode::Cancelled => WorkflowRunStatus::Cancelled,
727        WorkflowErrorCode::Timeout => WorkflowRunStatus::TimedOut,
728        _ => WorkflowRunStatus::Failed,
729    }
730}
731
732fn new_run_id() -> String {
733    let n = RUN_COUNTER.fetch_add(1, Ordering::SeqCst);
734    let millis = std::time::SystemTime::now()
735        .duration_since(std::time::UNIX_EPOCH)
736        .map(|d| d.as_millis())
737        .unwrap_or(0);
738    format!("wf_{millis}_{n}")
739}
740
741fn parse_run_policy(value: Option<&Value>) -> RunPolicy {
742    let mut policy = default_run_policy();
743    let Some(obj) = value.and_then(|v| v.as_object()) else {
744        return policy;
745    };
746    if let Some(p) = obj.get("profile").and_then(|v| v.as_str()) {
747        policy.profile = p.to_string();
748    }
749    if let Some(tools) = obj.get("tools").and_then(|v| v.as_array()) {
750        policy.tools = tools
751            .iter()
752            .filter_map(|v| v.as_str().map(|s| s.to_string()))
753            .collect();
754    }
755    if let Some(v) = obj.get("path_allow").and_then(|v| v.as_array()) {
756        policy.path_allow = v
757            .iter()
758            .filter_map(|x| x.as_str().map(|s| s.to_string()))
759            .collect();
760    }
761    if let Some(v) = obj.get("path_deny").and_then(|v| v.as_array()) {
762        policy.path_deny = v
763            .iter()
764            .filter_map(|x| x.as_str().map(|s| s.to_string()))
765            .collect();
766    }
767    if let Some(b) = obj.get("create_files").and_then(|v| v.as_bool()) {
768        policy.create_files = b;
769    }
770    if let Some(b) = obj.get("create_dirs").and_then(|v| v.as_bool()) {
771        policy.create_dirs = b;
772    }
773    if let Some(v) = obj.get("write_allow").and_then(|v| v.as_array()) {
774        policy.write_allow = v
775            .iter()
776            .filter_map(|x| x.as_str().map(|s| s.to_string()))
777            .collect();
778    }
779    if let Some(a) = obj.get("approval").and_then(|v| v.as_str()) {
780        policy.approval = a.to_string();
781    }
782    policy
783}
784
785fn fail(
786    invocation_id: &str,
787    run_id: Option<String>,
788    status: WorkflowRunStatus,
789    code: WorkflowErrorCode,
790    message: &str,
791    hint: Option<&str>,
792    stats: WorkflowStats,
793    journal_path: Option<String>,
794) -> ToolResult {
795    ToolResult {
796        invocation_id: invocation_id.to_string(),
797        ok: false,
798        output: json!({
799            "ok": false,
800            "run_id": run_id,
801            "status": status,
802            "result": null,
803            "stats": stats,
804            "journal_path": journal_path,
805            "error": {
806                "code": code,
807                "message": message,
808                "hint": hint,
809            },
810            "error_code": code,
811            "message": message,
812        }),
813    }
814}
815
816fn append_journal_line(path: &std::path::Path, value: &Value) {
817    use std::io::Write;
818    if let Ok(mut f) = std::fs::OpenOptions::new()
819        .create(true)
820        .append(true)
821        .open(path)
822    {
823        if let Ok(line) = serde_json::to_string(value) {
824            let _ = writeln!(f, "{line}");
825        }
826    }
827}
828
829fn truncate_json(value: Value, max_bytes: usize) -> Value {
830    let Ok(s) = serde_json::to_string(&value) else {
831        return value;
832    };
833    if s.len() <= max_bytes {
834        return value;
835    }
836    json!({
837        "truncated": true,
838        "original_bytes": s.len(),
839        "preview": redact_secrets(&s.chars().take(max_bytes.min(4096)).collect::<String>()),
840    })
841}
842
843fn optional_usize(input: &Value, key: &str) -> Option<usize> {
844    input.get(key).and_then(|v| {
845        v.as_u64()
846            .map(|n| n as usize)
847            .or_else(|| v.as_i64().map(|n| n.max(0) as usize))
848    })
849}
850
851fn optional_u64(input: &Value, key: &str) -> Option<u64> {
852    input
853        .get(key)
854        .and_then(|v| v.as_u64().or_else(|| v.as_i64().map(|n| n.max(0) as u64)))
855}
856
857/// Description text for snapshot tests (§12).
858pub fn workflow_tool_description() -> &'static str {
859    TOOL_DESCRIPTION
860}