1use 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
21const 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
32pub 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#[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
163pub 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 out.can_subagent = effective.tools.iter().any(|t| t == "subagent");
198 out.can_workflow = effective.tools.iter().any(|t| t == "workflow");
199
200 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 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 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 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 if let Some(wa) = effective.write_allow.first() {
274 let abs = project.join(wa);
275 let probe_path = if abs.exists() {
278 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 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 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 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 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
354pub 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 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, )
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
478pub(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 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 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}