1mod 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#[async_trait]
47pub trait AgentBackend: Send + Sync {
48 async fn run_agent(&self, request: AgentRequest) -> AgentBackendResult;
49}
50
51#[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
126pub 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
149pub struct WorkflowTool {
151 policy: SecurityPolicy,
152 config: WorkflowConfig,
153 backend: Arc<dyn AgentBackend>,
154}
155
156impl WorkflowTool {
157 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 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 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
857pub fn workflow_tool_description() -> &'static str {
859 TOOL_DESCRIPTION
860}