Skip to main content

navi_core/tool/builtin/workflow/
backends.rs

1//! Agent backends for workflow workers.
2//!
3//! - [`WorkerProbeBackend`]: exercises real tool registration + [`SecurityPolicy`]
4//!   (no live model). Used for permission/concurrency integration tests and as
5//!   the safe default when no parent executor is available.
6//! - [`SubagentBridgeBackend`]: production path — each `agent()` runs through
7//!   the real `subagent` tool / turn infrastructure with a filtered allowlist.
8
9use std::sync::atomic::{AtomicUsize, Ordering};
10use std::sync::{Arc, Weak};
11
12use async_trait::async_trait;
13use serde_json::json;
14
15use super::AgentBackend;
16use super::policy::EffectiveAgentPolicy;
17use super::types::{AgentBackendResult, AgentRequest, NESTED_WORKFLOW_TOOLS};
18use crate::security::{SecurityDecision, SecurityPolicy};
19use crate::tool::{ToolExecutor, ToolInvocation, ToolResult};
20
21/// Write-oriented tool names used when probing policy.
22const WRITE_TOOLS: &[&str] = &[
23    "write_file",
24    "write",
25    "edit",
26    "multiedit",
27    "apply_patch",
28    "code_edit",
29];
30const COMMAND_TOOLS: &[&str] = &["bash", "sandbox"];
31
32/// Builds a filtered worker executor and probes tool/path access via the real
33/// [`SecurityPolicy`] and tool registry path (no model call).
34pub struct WorkerProbeBackend {
35    policy: SecurityPolicy,
36    pub delay_ms: u64,
37    pub in_flight: Option<Arc<AtomicUsize>>,
38    pub peak_in_flight: Option<Arc<AtomicUsize>>,
39}
40
41impl WorkerProbeBackend {
42    pub fn new(policy: SecurityPolicy) -> Self {
43        Self {
44            policy,
45            delay_ms: 0,
46            in_flight: None,
47            peak_in_flight: None,
48        }
49    }
50
51    pub fn with_delay(mut self, delay_ms: u64) -> Self {
52        self.delay_ms = delay_ms;
53        self
54    }
55
56    pub fn with_inflight(mut self, in_flight: Arc<AtomicUsize>, peak: Arc<AtomicUsize>) -> Self {
57        self.in_flight = Some(in_flight);
58        self.peak_in_flight = Some(peak);
59        self
60    }
61}
62
63#[async_trait]
64impl AgentBackend for WorkerProbeBackend {
65    async fn run_agent(&self, request: AgentRequest) -> AgentBackendResult {
66        if let Some(ref inflight) = self.in_flight {
67            let n = inflight.fetch_add(1, Ordering::SeqCst) + 1;
68            if let Some(ref peak) = self.peak_in_flight {
69                peak.fetch_max(n, Ordering::SeqCst);
70            }
71        }
72
73        if self.delay_ms > 0 {
74            let delay = std::time::Duration::from_millis(self.delay_ms);
75            tokio::select! {
76                _ = tokio::time::sleep(delay) => {}
77                _ = request.cancel_token.notified() => {
78                    if let Some(ref inflight) = self.in_flight {
79                        inflight.fetch_sub(1, Ordering::SeqCst);
80                    }
81                    return AgentBackendResult {
82                        ok: false,
83                        output: json!({"error": "cancelled"}),
84                        error: Some("cancelled".into()),
85                    };
86                }
87            }
88        }
89
90        if request.cancel_token.is_requested() {
91            if let Some(ref inflight) = self.in_flight {
92                inflight.fetch_sub(1, Ordering::SeqCst);
93            }
94            return AgentBackendResult {
95                ok: false,
96                output: json!({"error": "cancelled"}),
97                error: Some("cancelled".into()),
98            };
99        }
100
101        let probe = probe_worker_capabilities_inner(&self.policy, &request.effective);
102
103        if let Some(ref inflight) = self.in_flight {
104            inflight.fetch_sub(1, Ordering::SeqCst);
105        }
106
107        AgentBackendResult {
108            ok: true,
109            output: json!({
110                "ok": true,
111                "backend": "worker_probe",
112                "prompt": request.prompt,
113                "label": request.label,
114                "agent_index": request.agent_index,
115                "profile": request.effective.profile,
116                "tools": request.effective.tools,
117                "create_files": request.effective.create_files,
118                "create_dirs": request.effective.create_dirs,
119                "write_allow": request.effective.write_allow,
120                "path_allow": request.effective.path_allow,
121                "path_deny": request.effective.path_deny,
122                "can_write_file": probe.can_write_file,
123                "can_edit": probe.can_edit,
124                "can_bash": probe.can_bash,
125                "can_subagent": probe.can_subagent,
126                "can_workflow": probe.can_workflow,
127                "write_path_allowed": probe.write_path_allowed,
128                "write_path_denied": probe.write_path_denied,
129                "create_new_file_allowed": probe.create_new_file_allowed,
130                "registered_tools": probe.registered_tools,
131                "policy_denials": probe.policy_denials,
132            }),
133            error: None,
134        }
135    }
136}
137
138#[derive(Debug, Default)]
139struct ProbeResult {
140    can_write_file: bool,
141    can_edit: bool,
142    can_bash: bool,
143    can_subagent: bool,
144    can_workflow: bool,
145    write_path_allowed: Vec<String>,
146    write_path_denied: Vec<String>,
147    create_new_file_allowed: bool,
148    registered_tools: Vec<String>,
149    policy_denials: Vec<String>,
150}
151
152/// Public probe summary for unit tests (fields only; no mock echo).
153#[derive(Debug, Clone)]
154pub struct WorkerProbeSummary {
155    pub can_bash: bool,
156    pub can_write_file: bool,
157    pub can_edit: bool,
158    pub can_subagent: bool,
159    pub can_workflow: bool,
160    pub registered_tools: Vec<String>,
161}
162
163/// Probe real tool allowlist + SecurityPolicy for a worker's effective policy.
164pub fn probe_worker_capabilities(
165    base_policy: &SecurityPolicy,
166    effective: &EffectiveAgentPolicy,
167) -> WorkerProbeSummary {
168    let inner = probe_worker_capabilities_inner(base_policy, effective);
169    WorkerProbeSummary {
170        can_bash: inner.can_bash,
171        can_write_file: inner.can_write_file,
172        can_edit: inner.can_edit,
173        can_subagent: inner.can_subagent,
174        can_workflow: inner.can_workflow,
175        registered_tools: inner.registered_tools,
176    }
177}
178
179fn scoped_policy(base: &SecurityPolicy, effective: &EffectiveAgentPolicy) -> SecurityPolicy {
180    base.clone()
181        .with_write_scope(crate::security::WritePathScope {
182            write_allow: effective.write_allow.clone(),
183            path_deny: effective.path_deny.clone(),
184            create_files: effective.create_files,
185            create_dirs: effective.create_dirs,
186        })
187}
188
189fn probe_worker_capabilities_inner(
190    base_policy: &SecurityPolicy,
191    effective: &EffectiveAgentPolicy,
192) -> ProbeResult {
193    let mut out = ProbeResult::default();
194    let project = base_policy.project_root().to_path_buf();
195
196    // Nested orchestration must never appear after intersection + strip.
197    out.can_subagent = effective.tools.iter().any(|t| t == "subagent");
198    out.can_workflow = effective.tools.iter().any(|t| t == "workflow");
199
200    // Worker executor with WritePathScope (same gate as production SubagentBridge).
201    let policy = scoped_policy(base_policy, effective);
202    let mut exec = ToolExecutor::empty(policy.clone());
203    register_filtered_tools(&mut exec, &project, effective);
204    out.registered_tools = exec.tool_names();
205    out.registered_tools.sort();
206
207    out.can_write_file = exec
208        .tool_names()
209        .iter()
210        .any(|t| t == "write_file" || t == "write");
211    out.can_edit = exec
212        .tool_names()
213        .iter()
214        .any(|t| t == "edit" || t == "multiedit");
215    out.can_bash = exec.tool_names().iter().any(|t| t == "bash");
216    out.can_subagent = out.registered_tools.iter().any(|t| t == "subagent");
217    out.can_workflow = out.registered_tools.iter().any(|t| t == "workflow");
218
219    // Real ToolExecutor::validate path for representative writes.
220    let probe_paths: Vec<String> = {
221        let mut c = effective.write_allow.clone();
222        if c.is_empty() {
223            c.push("src/a.rs".into());
224        }
225        c.push("__outside_write_allow__.rs".into());
226        for d in &effective.path_deny {
227            let clean = d
228                .trim_end_matches('/')
229                .trim_end_matches('*')
230                .trim_end_matches('/');
231            if !clean.is_empty() {
232                c.push(clean.to_string());
233            }
234        }
235        // Non-existent path under write_allow for create_files probe.
236        if let Some(first) = effective.write_allow.first() {
237            c.push(format!("__new_create_probe__/{first}"));
238        } else {
239            c.push("__new_create_probe__/file.rs".into());
240        }
241        c.sort();
242        c.dedup();
243        c
244    };
245
246    for path in &probe_paths {
247        let inv = ToolInvocation {
248            id: format!("probe-write-{path}"),
249            tool_name: "write_file".into(),
250            input: json!({"path": path, "content": "x"}),
251        };
252        match exec.validate(&inv) {
253            SecurityDecision::Deny(reason) => {
254                out.policy_denials
255                    .push(format!("write_file {path}: {reason}"));
256                out.write_path_denied.push(path.clone());
257            }
258            SecurityDecision::Allow | SecurityDecision::NeedsApproval(_) => {
259                // Only count as allowed if tool is registered AND validate ok.
260                if out.can_write_file {
261                    out.write_path_allowed.push(path.clone());
262                } else {
263                    out.write_path_denied.push(path.clone());
264                    out.policy_denials
265                        .push(format!("write_file {path}: tool not registered"));
266                }
267            }
268        }
269    }
270
271    // create_files: writing a write_allow path that does not exist yet must Deny
272    // when create_files=false (real SecurityPolicy WritePathScope).
273    if let Some(wa) = effective.write_allow.first() {
274        let abs = project.join(wa);
275        // Use a unique non-existent path that still matches write_allow when
276        // write_allow is a single file — probe that exact path if missing.
277        let probe_path = if abs.exists() {
278            // Existing path: also probe a sibling under same allow prefix if possible.
279            format!("__wf_create_probe__/{wa}")
280        } else {
281            wa.clone()
282        };
283        let inv = ToolInvocation {
284            id: "probe-create".into(),
285            tool_name: "write_file".into(),
286            input: json!({"path": probe_path, "content": "new"}),
287        };
288        match exec.validate(&inv) {
289            SecurityDecision::Deny(reason) => {
290                out.create_new_file_allowed = false;
291                out.policy_denials.push(format!("create_new: {reason}"));
292            }
293            SecurityDecision::Allow | SecurityDecision::NeedsApproval(_) => {
294                // Only true if tool registered, write_allow non-empty, create_files true,
295                // and path is in write_allow (validate already checked scope).
296                out.create_new_file_allowed = out.can_write_file && effective.create_files;
297            }
298        }
299    } else {
300        out.create_new_file_allowed = false;
301    }
302
303    // Empty write_allow ⇒ no writes even if tools listed.
304    if effective.write_allow.is_empty() {
305        out.can_write_file = false;
306        out.can_edit = false;
307        out.create_new_file_allowed = false;
308    }
309
310    out
311}
312
313fn register_filtered_tools(
314    exec: &mut ToolExecutor,
315    project: &std::path::Path,
316    effective: &EffectiveAgentPolicy,
317) {
318    use super::super::{
319        bash::BashTool, edit_tool::EditTool, read_tool::ReadTool, search_tool::SearchTool,
320        write_tool::WriteTool,
321    };
322
323    // Never register orchestration tools.
324    let allowed: Vec<&str> = effective
325        .tools
326        .iter()
327        .map(|s| s.as_str())
328        .filter(|t| !NESTED_WORKFLOW_TOOLS.contains(t))
329        .collect();
330
331    let has = |name: &str| allowed.contains(&name);
332
333    if has("read_file") || has("read") || has("view_file") {
334        exec.register_tool(Arc::new(ReadTool::new(project.to_path_buf())));
335    }
336    if has("search") || has("grep") || has("fs_browser") || has("list_dir") || has("glob") {
337        exec.register_tool(Arc::new(SearchTool::new(project.to_path_buf())));
338    }
339
340    // Writes only when write_allow is non-empty (empty ⇒ no writes even for implementer).
341    // create_files=false still registers tools; WritePathScope denies creates.
342    let writes_ok = !effective.write_allow.is_empty();
343    if writes_ok && (has("write_file") || has("write")) {
344        exec.register_tool(Arc::new(WriteTool::write_file(project.to_path_buf())));
345    }
346    if writes_ok && (has("edit") || has("multiedit")) {
347        exec.register_tool(Arc::new(EditTool::new(project.to_path_buf())));
348    }
349    if has("bash") {
350        exec.register_tool(Arc::new(BashTool::new(project.to_path_buf())));
351    }
352}
353
354/// Production backend: each worker is a real nested `subagent` turn with a
355/// tool allowlist derived from the effective workflow policy.
356pub struct SubagentBridgeBackend {
357    tool_executor: Weak<ToolExecutor>,
358}
359
360impl SubagentBridgeBackend {
361    pub fn new(tool_executor: Weak<ToolExecutor>) -> Self {
362        Self { tool_executor }
363    }
364}
365
366#[async_trait]
367impl AgentBackend for SubagentBridgeBackend {
368    async fn run_agent(&self, request: AgentRequest) -> AgentBackendResult {
369        let Some(executor) = self.tool_executor.upgrade() else {
370            return AgentBackendResult {
371                ok: false,
372                output: json!({"error": "tool executor unavailable"}),
373                error: Some("tool executor dropped".into()),
374            };
375        };
376
377        if request.cancel_token.is_requested() {
378            return AgentBackendResult {
379                ok: false,
380                output: json!({"error": "cancelled"}),
381                error: Some("cancelled".into()),
382            };
383        }
384
385        // Embed path policy in the prompt (guidance) AND pass write scope fields
386        // so SubagentTool forks a SecurityPolicy with WritePathScope (hard gate).
387        let tools_for_note: Vec<String> = {
388            let mut t = request.effective.tools.clone();
389            t.retain(|n| !NESTED_WORKFLOW_TOOLS.contains(&n.as_str()));
390            if request.effective.write_allow.is_empty() {
391                t.retain(|n| {
392                    !WRITE_TOOLS.contains(&n.as_str()) && !COMMAND_TOOLS.contains(&n.as_str())
393                });
394            }
395            t
396        };
397        let path_note = format!(
398            "\n\n[workflow worker policy]\n\
399             profile={}\n\
400             tools={:?}\n\
401             write_allow={:?}\n\
402             path_deny={:?}\n\
403             create_files={}\n\
404             create_dirs={}\n\
405             You MUST NOT call subagent or workflow. \
406             Writes are only allowed on write_allow paths (empty ⇒ no writes).",
407            request.effective.profile,
408            tools_for_note,
409            request.effective.write_allow,
410            request.effective.path_deny,
411            request.effective.create_files,
412            request.effective.create_dirs,
413        );
414
415        let prompt = format!("{}{path_note}", request.prompt);
416        let input = build_subagent_bridge_input(
417            &prompt,
418            request.label.as_deref(),
419            &request.effective,
420            request.model.as_deref(),
421            request.max_tokens,
422        );
423        let inv = ToolInvocation {
424            id: format!("wf-agent-{}", request.agent_index),
425            tool_name: "subagent".into(),
426            input,
427        };
428
429        let result: ToolResult = executor
430            .invoke_with_full_context(
431                inv,
432                crate::tool::ToolInvocationContext {
433                    cancel_token: Some(request.cancel_token.clone()),
434                    ..Default::default()
435                },
436                true, // workflow already approved at parent tool level
437            )
438            .await;
439
440        if request.cancel_token.is_requested() {
441            return AgentBackendResult {
442                ok: false,
443                output: json!({"error": "cancelled"}),
444                error: Some("cancelled".into()),
445            };
446        }
447
448        let err_msg = if result.ok {
449            None
450        } else {
451            Some(
452                result
453                    .output
454                    .get("error")
455                    .and_then(|e| e.as_str())
456                    .unwrap_or("subagent failed")
457                    .to_string(),
458            )
459        };
460        let mut output = result.output;
461        if let Some(obj) = output.as_object_mut() {
462            obj.insert("backend".into(), json!("subagent_bridge"));
463            obj.insert("agent_index".into(), json!(request.agent_index));
464            obj.insert("profile".into(), json!(request.effective.profile));
465            obj.insert("tools".into(), json!(request.effective.tools));
466            obj.insert("write_allow".into(), json!(request.effective.write_allow));
467            obj.insert("create_files".into(), json!(request.effective.create_files));
468        }
469
470        AgentBackendResult {
471            ok: result.ok,
472            output,
473            error: err_msg,
474        }
475    }
476}
477
478/// Build the JSON tool input the production bridge sends to `subagent`.
479/// Extracted for unit tests (schema + null-description regressions).
480pub(crate) fn build_subagent_bridge_input(
481    prompt: &str,
482    label: Option<&str>,
483    effective: &EffectiveAgentPolicy,
484    model: Option<&str>,
485    max_tokens: Option<usize>,
486) -> serde_json::Value {
487    let mut tools = effective.tools.clone();
488    tools.retain(|t| !NESTED_WORKFLOW_TOOLS.contains(&t.as_str()));
489    if effective.write_allow.is_empty() {
490        tools
491            .retain(|t| !WRITE_TOOLS.contains(&t.as_str()) && !COMMAND_TOOLS.contains(&t.as_str()));
492    }
493    let approval = if effective.write_allow.is_empty() {
494        "read_only"
495    } else if effective.approval == "escalate" {
496        "escalate"
497    } else {
498        effective.approval.as_str()
499    };
500    let mut options = json!({
501        "agent_profile": effective.profile,
502        "tools": tools,
503        "approval": approval,
504        "write_allow": effective.write_allow,
505        "path_deny": effective.path_deny,
506        "create_files": effective.create_files,
507        "create_dirs": effective.create_dirs,
508    });
509    if let Some(model) = model {
510        options
511            .as_object_mut()
512            .expect("options object")
513            .insert("model".into(), json!(model));
514    }
515    if let Some(max_tokens) = max_tokens {
516        options
517            .as_object_mut()
518            .expect("options object")
519            .insert("max_tokens".into(), json!(max_tokens));
520    }
521    let mut input = json!({
522        "prompt": prompt,
523        "options": options,
524    });
525    if let Some(label) = label.map(str::trim).filter(|s| !s.is_empty()) {
526        input
527            .as_object_mut()
528            .expect("input object")
529            .insert("description".into(), json!(label));
530    }
531    input
532}
533
534#[cfg(test)]
535mod tests {
536    use super::*;
537    use crate::config::{HarnessConfig, NaviConfig};
538    use crate::model::{ModelProvider, ModelRequest, ModelStream};
539    use crate::prompt::PromptCache;
540    use crate::runtime_components::RuntimeComponents;
541    use crate::tool::Tool;
542    use crate::tool::builtin::SubagentTool;
543    use crate::tool::builtin::workflow::policy::default_run_policy;
544    use std::sync::{Arc, RwLock};
545
546    struct NoopProvider;
547    impl ModelProvider for NoopProvider {
548        fn stream(&self, _req: ModelRequest) -> ModelStream {
549            Box::pin(futures_util::stream::empty())
550        }
551    }
552
553    /// Schema from a real SubagentTool — never a hand-rolled duplicate that can drift.
554    fn registered_subagent_schema() -> serde_json::Value {
555        let tool = SubagentTool::new(
556            std::sync::Weak::new(),
557            Arc::new(RwLock::new(Arc::new(NoopProvider) as Arc<dyn ModelProvider>)),
558            std::path::PathBuf::from("/tmp"),
559            std::path::PathBuf::from("/tmp"),
560            Arc::new(RwLock::new("test".into())),
561            HarnessConfig::default(),
562            Arc::new(RwLock::new(NaviConfig::default())),
563            Arc::new(PromptCache::new()),
564            RuntimeComponents::default(),
565        );
566        tool.definition().input_schema
567    }
568
569    #[test]
570    fn bridge_input_omits_null_description_when_label_missing() {
571        let mut run = default_run_policy();
572        run.create_files = true;
573        run.write_allow = vec!["scratch/a.txt".into()];
574        run.tools = vec![
575            "read_file".into(),
576            "write_file".into(),
577            "edit".into(),
578            "search".into(),
579        ];
580        let effective = crate::tool::builtin::workflow::policy::intersect_agent_policy(
581            &run,
582            &crate::tool::builtin::workflow::policy::AgentPolicyOpts {
583                profile: Some("implementer".into()),
584                ..Default::default()
585            },
586        );
587        assert!(effective.create_files);
588        let input = build_subagent_bridge_input("do work", None, &effective, None, None);
589        assert!(
590            input.get("description").is_none(),
591            "missing label must not serialize description:null, got {input}"
592        );
593        assert_eq!(input["options"]["create_files"], true);
594        assert_eq!(input["options"]["write_allow"], json!(["scratch/a.txt"]));
595        // Validate against the live SubagentTool schema (not a hand-rolled twin).
596        let schema = registered_subagent_schema();
597        let validator = jsonschema::validator_for(&schema).unwrap();
598        let errors: Vec<_> = validator
599            .iter_errors(&input)
600            .map(|e| e.to_string())
601            .collect();
602        assert!(
603            errors.is_empty(),
604            "bridge input invalid vs registered SubagentTool schema: {errors:?} input={input}"
605        );
606    }
607
608    #[test]
609    fn bridge_input_includes_non_empty_label() {
610        let run = default_run_policy();
611        let effective = crate::tool::builtin::workflow::policy::intersect_agent_policy(
612            &run,
613            &Default::default(),
614        );
615        let input = build_subagent_bridge_input("p", Some("  collect  "), &effective, None, None);
616        assert_eq!(input["description"], "collect");
617        let schema = registered_subagent_schema();
618        let validator = jsonschema::validator_for(&schema).unwrap();
619        let errors: Vec<_> = validator
620            .iter_errors(&input)
621            .map(|e| e.to_string())
622            .collect();
623        assert!(
624            errors.is_empty(),
625            "labeled bridge input invalid vs SubagentTool schema: {errors:?} input={input}"
626        );
627    }
628}