1mod acp_agent;
4pub mod adapters;
5mod agent_output;
6mod agent_progress;
7pub mod background;
8pub mod conversation;
9pub mod delegation;
10mod live;
11mod scv_agent;
12pub mod web;
13
14use std::{
15 collections::HashMap,
16 ffi::OsString,
17 io::{Read as _, Write as _},
18 os::unix::process::CommandExt as _,
19 path::{Component, Path, PathBuf},
20 sync::{
21 Arc,
22 atomic::{AtomicU64, Ordering},
23 },
24 time::Duration,
25};
26
27use async_trait::async_trait;
28use cap_std::{
29 ambient_authority,
30 fs::{Dir, OpenOptions},
31};
32use scv_core::{Tool, ToolContext, ToolError, ToolOutput, ToolRegistry, ToolRisk, ToolSpec};
33
34use crate::{
35 adapters::{OutputFormat, Resume, Transport},
36 agent_output::{AgentStream, RunExit, STDERR_TAIL_BYTES, TailBuffer},
37 conversation::{ConversationLimits, ConversationStore},
38 delegation::{DelegationGuard, DelegationRegistry},
39};
40use serde::Deserialize;
41use serde_json::{Value, json};
42use sha2::{Digest, Sha256};
43use tokio::{
44 io::AsyncReadExt,
45 process::Command,
46 sync::Mutex,
47 task::JoinHandle,
48 time::{Instant, sleep, sleep_until, timeout, timeout_at},
49};
50
51#[derive(Debug, Clone)]
52pub struct ToolsConfig {
53 pub command_timeout: Duration,
55 pub agent_timeout: Duration,
57 pub max_timeout: Duration,
59 pub output_limit_bytes: usize,
60 pub max_read_bytes: usize,
61 pub max_write_bytes: usize,
62 pub max_delegation_depth: u32,
64 pub conversations: ConversationLimits,
66 pub delegation: Option<DelegationContext>,
68 pub max_background: usize,
71 pub background: Option<Arc<background::BackgroundJobs>>,
74}
75
76impl Default for ToolsConfig {
77 fn default() -> Self {
78 Self {
79 command_timeout: Duration::from_secs(600),
80 agent_timeout: Duration::from_secs(3600),
81 max_timeout: Duration::from_secs(14400),
82 output_limit_bytes: 64 * 1024,
83 max_read_bytes: 256 * 1024,
84 max_write_bytes: 1024 * 1024,
85 max_delegation_depth: 2,
86 conversations: ConversationLimits {
87 max: 8,
88 idle: Duration::from_secs(86400),
89 },
90 delegation: None,
91 max_background: 2,
92 background: None,
93 }
94 }
95}
96
97#[derive(Debug, Clone)]
99pub struct DelegationContext {
100 pub registry: Arc<DelegationRegistry>,
101 pub session: String,
102 pub depth: u32,
105}
106
107impl DelegationContext {
108 pub fn owner_depth(&self) -> u32 {
110 self.registry.depth().max(self.depth)
111 }
112}
113
114#[derive(Debug, Clone)]
115pub struct AgentAdapterConfig {
116 pub command: String,
117 pub args: Vec<String>,
118 pub prompt_args: Vec<String>,
121 pub full_permission_args: Option<Vec<String>>,
124 pub model_args: Vec<String>,
127 pub effort_args: Vec<String>,
130 pub model_hint: String,
132 pub environment: Vec<(OsString, OsString)>,
134 pub search_dirs: Vec<PathBuf>,
136 pub output: OutputFormat,
138 pub resume: Resume,
140 pub home: Option<PathBuf>,
142 pub transport: Transport,
144 pub acp: Option<AcpAgentLaunch>,
147}
148
149#[derive(Debug, Clone)]
152pub struct AcpAgentLaunch {
153 pub command: String,
154 pub args: Vec<String>,
156 pub full_mode: Option<String>,
159 pub environment: Vec<(OsString, OsString)>,
162 pub required: bool,
165}
166
167pub type SkillMap = HashMap<String, PathBuf>;
168
169pub fn builtin_registry(
170 config: ToolsConfig,
171 skills: SkillMap,
172 skill_roots: Vec<PathBuf>,
173 max_skill_bytes: usize,
174 adapters: HashMap<String, AgentAdapterConfig>,
175) -> Result<ToolRegistry, ToolError> {
176 let mut registry = ToolRegistry::default();
177 registry.register(Arc::new(ReadTool {
178 max_bytes: config.max_read_bytes,
179 }))?;
180 registry.register(Arc::new(ReadSkillTool {
181 skills,
182 roots: skill_roots,
183 max_bytes: max_skill_bytes,
184 }))?;
185 registry.register(Arc::new(WriteTool {
186 max_bytes: config.max_write_bytes,
187 }))?;
188 registry.register(Arc::new(BashTool {
189 timeout: config.command_timeout,
190 max_timeout: config.max_timeout,
191 output_limit: config.output_limit_bytes,
192 }))?;
193 let depth = config
195 .delegation
196 .as_ref()
197 .map_or_else(delegation::current_depth, DelegationContext::owner_depth);
198 let adapters = if depth < config.max_delegation_depth {
199 adapters
200 } else {
201 HashMap::new()
202 };
203 let jobs = (config.max_background > 0).then(|| {
206 config.background.clone().unwrap_or_else(|| {
207 Arc::new(background::BackgroundJobs::new(config.max_background, None))
208 })
209 });
210 let mut agents = 0;
211 let mut register_agent = |registry: &mut ToolRegistry, tool: Arc<dyn Tool>| {
212 agents += 1;
213 match &jobs {
214 Some(jobs) => registry.register(Arc::new(background::BackgroundCapable {
215 inner: tool,
216 jobs: Arc::clone(jobs),
217 })),
218 None => registry.register(tool),
219 }
220 };
221 let conversations = Arc::new(ConversationStore::new(
223 config.conversations,
224 config
225 .delegation
226 .as_ref()
227 .map(|context| context.registry.conversation_dir()),
228 ));
229 for (name, adapter) in adapters {
230 if adapter.transport == Transport::ScvProtocol {
231 let resolved =
232 adapters::resolve_agent_executable(&adapter.command, &adapter.search_dirs);
233 if resolved.is_some() {
235 register_agent(
236 &mut registry,
237 Arc::new(scv_agent::ScvAgentTool {
238 name,
239 command: adapter.command,
240 resolved,
241 args: adapter.args,
242 environment: adapter.environment,
243 timeouts: Timeouts {
244 default: config.agent_timeout,
245 max: config.max_timeout,
246 },
247 output_limit: config.output_limit_bytes,
248 delegation: config.delegation.clone(),
249 conversations: Arc::clone(&conversations),
250 }),
251 )?;
252 }
253 continue;
254 }
255 if let Some(launch) = adapter.acp.clone() {
256 let resolved =
257 adapters::resolve_agent_executable(&launch.command, &adapter.search_dirs);
258 if resolved.is_some() {
259 register_agent(
260 &mut registry,
261 Arc::new(acp_agent::AcpAgentTool::new(
262 name,
263 &adapter,
264 launch,
265 resolved,
266 Timeouts {
267 default: config.agent_timeout,
268 max: config.max_timeout,
269 },
270 config.output_limit_bytes,
271 config.delegation.clone(),
272 Arc::clone(&conversations),
273 )),
274 )?;
275 continue;
276 }
277 if launch.required {
278 continue;
280 }
281 }
282 let tool = NativeAgentTool::new(
283 name,
284 adapter,
285 Timeouts {
286 default: config.agent_timeout,
287 max: config.max_timeout,
288 },
289 config.output_limit_bytes,
290 config.delegation.clone(),
291 Arc::clone(&conversations),
292 );
293 if tool.resolved.is_some() {
295 register_agent(&mut registry, Arc::new(tool))?;
296 }
297 }
298 if let Some(jobs) = jobs.filter(|_| agents > 0) {
299 registry.register(Arc::new(background::WaitTool {
300 jobs: Arc::clone(&jobs),
301 timeouts: Timeouts {
302 default: config.agent_timeout,
303 max: config.max_timeout,
304 },
305 }))?;
306 registry.register(Arc::new(background::StatusTool { jobs }))?;
307 }
308 Ok(registry)
309}
310
311struct ReadTool {
312 max_bytes: usize,
313}
314
315#[derive(Deserialize)]
316#[serde(deny_unknown_fields)]
317struct ReadArgs {
318 path: String,
319 #[serde(default)]
320 offset: usize,
321 limit: Option<usize>,
322}
323
324#[async_trait]
325impl Tool for ReadTool {
326 fn spec(&self) -> ToolSpec {
327 ToolSpec {
328 name: "read".into(),
329 description: "Read a bounded UTF-8 file inside the workspace".into(),
330 parameters: json!({
331 "type":"object",
332 "properties":{
333 "path":{"type":"string"},
334 "offset":{"type":"integer","minimum":0},
335 "limit":{"type":"integer","minimum":1}
336 },
337 "required":["path"],
338 "additionalProperties":false
339 }),
340 }
341 }
342
343 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
344 let args: ReadArgs = parse_args(arguments)?;
345 validate_read_args(&args)?;
346 Ok(if is_secret_like(Path::new(&args.path)) {
347 ToolRisk::Filesystem
348 } else {
349 ToolRisk::ReadOnly
350 })
351 }
352
353 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
354 let args: ReadArgs = parse_args(arguments)?;
355 validate_read_args(&args)?;
356 Ok(format!("Read {}", args.path))
357 }
358
359 async fn execute(
360 &self,
361 arguments: Value,
362 context: ToolContext,
363 ) -> Result<ToolOutput, ToolError> {
364 let args: ReadArgs = parse_args(&arguments)?;
365 validate_read_args(&args)?;
366 let requested = args.limit.unwrap_or(self.max_bytes).min(self.max_bytes);
367 let offset = u64::try_from(args.offset).unwrap_or(u64::MAX);
368 let workspace = context.workspace.clone();
369 let display_path = args.path.clone();
370 let relative = PathBuf::from(&args.path);
371 validate_relative(&relative)?;
372 let read = tokio::task::spawn_blocking(move || {
373 let root = open_workspace(&workspace)?;
374 let mut file = root
375 .open(&relative)
376 .map_err(|error| map_cap_error("read", &display_path, error))?;
377 let total_bytes = file
378 .metadata()
379 .map_err(|error| ToolError(format!("stat {display_path}: {error}")))?
380 .len();
381 let start = offset.min(total_bytes);
382 std::io::Seek::seek(&mut file, std::io::SeekFrom::Start(start))
383 .map_err(|error| ToolError(format!("seek {display_path}: {error}")))?;
384 let mut bytes = Vec::with_capacity(requested.min(8192));
385 std::io::Read::take(&mut file, u64::try_from(requested).unwrap_or(u64::MAX))
386 .read_to_end(&mut bytes)
387 .map_err(|error| ToolError(format!("read {display_path}: {error}")))?;
388 Ok::<_, ToolError>((bytes, total_bytes, start))
389 });
390 let (bytes, total_bytes, start) = tokio::select! {
391 result = read => result.map_err(|error| ToolError(format!("read task failed: {error}")))??,
392 _ = context.cancellation.cancelled() => return Err(ToolError("read cancelled".into())),
393 };
394 let content = std::str::from_utf8(&bytes)
395 .map_err(|_| ToolError(format!("selected range of {} is not UTF-8", args.path)))?;
396 let end = start.saturating_add(u64::try_from(bytes.len()).unwrap_or(u64::MAX));
397 let truncated = start > 0 || end < total_bytes;
398 Ok(ToolOutput {
399 content: json!({
400 "path": args.path,
401 "content": content,
402 "total_bytes": total_bytes,
403 "offset": start,
404 "truncated": truncated
405 })
406 .to_string(),
407 is_error: false,
408 truncated,
409 })
410 }
411}
412
413struct ReadSkillTool {
414 skills: SkillMap,
415 roots: Vec<PathBuf>,
416 max_bytes: usize,
417}
418
419#[derive(Deserialize)]
420#[serde(deny_unknown_fields)]
421struct ReadSkillArgs {
422 name: String,
423}
424
425#[async_trait]
426impl Tool for ReadSkillTool {
427 fn spec(&self) -> ToolSpec {
428 ToolSpec {
429 name: "read_skill".into(),
430 description: "Load a discovered SCV skill by name".into(),
431 parameters: json!({
432 "type":"object",
433 "properties":{"name":{"type":"string"}},
434 "required":["name"],
435 "additionalProperties":false
436 }),
437 }
438 }
439
440 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
441 let _: ReadSkillArgs = parse_args(arguments)?;
442 Ok(ToolRisk::ReadOnly)
443 }
444
445 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
446 let args: ReadSkillArgs = parse_args(arguments)?;
447 Ok(format!("Load skill {}", args.name))
448 }
449
450 async fn execute(
451 &self,
452 arguments: Value,
453 context: ToolContext,
454 ) -> Result<ToolOutput, ToolError> {
455 let args: ReadSkillArgs = parse_args(&arguments)?;
456 let configured = self
457 .skills
458 .get(&args.name)
459 .ok_or_else(|| ToolError(format!("unknown skill: {}", args.name)))?;
460 let path = std::fs::canonicalize(configured)
461 .map_err(|error| ToolError(format!("load skill {}: {error}", args.name)))?;
462 if !self.roots.iter().any(|root| path.starts_with(root)) {
463 return Err(ToolError("skill path escaped its configured root".into()));
464 }
465 let max_bytes = self.max_bytes;
466 let skill_name = args.name.clone();
467 let bytes = tokio::select! {
468 result = tokio::task::spawn_blocking(move || {
469 let mut file = std::fs::File::open(&path)
470 .map_err(|error| ToolError(format!("load skill {skill_name}: {error}")))?;
471 let mut bytes = Vec::with_capacity(max_bytes.min(8192));
472 std::io::Read::take(
473 &mut file,
474 u64::try_from(max_bytes).unwrap_or(u64::MAX).saturating_add(1),
475 )
476 .read_to_end(&mut bytes)
477 .map_err(|error| ToolError(format!("load skill {skill_name}: {error}")))?;
478 Ok::<_, ToolError>(bytes)
479 }) => result.map_err(|error| ToolError(format!("skill read task failed: {error}")))??,
480 _ = context.cancellation.cancelled() => return Err(ToolError("skill read cancelled".into())),
481 };
482 let end = bytes.len().min(self.max_bytes);
483 let content = std::str::from_utf8(&bytes[..end])
484 .map_err(|_| ToolError("skill is not UTF-8".into()))?;
485 Ok(ToolOutput {
486 content: content.to_owned(),
487 is_error: false,
488 truncated: end < bytes.len(),
489 })
490 }
491}
492
493struct WriteTool {
494 max_bytes: usize,
495}
496
497#[derive(Deserialize)]
498#[serde(deny_unknown_fields)]
499struct WriteArgs {
500 path: String,
501 content: String,
502 mode: WriteMode,
503 expected_sha256: Option<String>,
504}
505
506#[derive(Deserialize)]
507#[serde(rename_all = "snake_case")]
508enum WriteMode {
509 Create,
510 Replace,
511}
512
513#[async_trait]
514impl Tool for WriteTool {
515 fn spec(&self) -> ToolSpec {
516 ToolSpec {
517 name: "write".into(),
518 description: "Atomically create or replace a UTF-8 file inside the workspace".into(),
519 parameters: json!({
520 "type":"object",
521 "properties":{
522 "path":{"type":"string"},
523 "content":{"type":"string"},
524 "mode":{"type":"string","enum":["create","replace"]},
525 "expected_sha256":{"type":"string"}
526 },
527 "required":["path","content","mode"],
528 "additionalProperties":false
529 }),
530 }
531 }
532
533 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
534 let _: WriteArgs = parse_args(arguments)?;
535 Ok(ToolRisk::Filesystem)
536 }
537
538 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
539 let args: WriteArgs = parse_args(arguments)?;
540 let mode = match args.mode {
541 WriteMode::Create => "Create",
542 WriteMode::Replace => "Replace",
543 };
544 Ok(format!(
545 "{mode} {} ({} bytes)",
546 args.path,
547 args.content.len()
548 ))
549 }
550
551 async fn execute(
552 &self,
553 arguments: Value,
554 context: ToolContext,
555 ) -> Result<ToolOutput, ToolError> {
556 let args: WriteArgs = parse_args(&arguments)?;
557 if args.content.len() > self.max_bytes {
558 return Err(ToolError(format!(
559 "write exceeds {} byte limit",
560 self.max_bytes
561 )));
562 }
563 let workspace = context.workspace.clone();
564 let cancellation = context.cancellation.clone();
565 tokio::task::spawn_blocking(move || {
566 if cancellation.is_cancelled() {
567 return Err(ToolError("write cancelled".into()));
568 }
569 let path = PathBuf::from(&args.path);
570 validate_relative(&path)?;
571 let root = open_workspace(&workspace)?;
572 let exists = match root.symlink_metadata(&path) {
573 Ok(_) => true,
574 Err(error) if error.kind() == std::io::ErrorKind::NotFound => false,
575 Err(error) => return Err(map_cap_error("inspect", &args.path, error)),
576 };
577 match args.mode {
578 WriteMode::Create if exists => {
579 return Err(ToolError(format!("{} already exists", args.path)));
580 }
581 WriteMode::Replace if !exists => {
582 return Err(ToolError(format!("{} does not exist", args.path)));
583 }
584 _ => {}
585 }
586 if let Some(expected) = args.expected_sha256 {
587 let mut current_file = root
588 .open(&path)
589 .map_err(|error| map_cap_error("hash", &args.path, error))?;
590 let mut current = Vec::new();
591 current_file
592 .read_to_end(&mut current)
593 .map_err(|error| ToolError(format!("hash {}: {error}", args.path)))?;
594 let actual = format!("{:x}", Sha256::digest(current));
595 if actual != expected.to_ascii_lowercase() {
596 return Err(ToolError(format!(
597 "{} changed: expected sha256 {}, found {}",
598 args.path, expected, actual
599 )));
600 }
601 }
602 let parent = path.parent().unwrap_or_else(|| Path::new("."));
603 root.create_dir_all(parent)
604 .map_err(|error| map_cap_error("create directory for", &args.path, error))?;
605 let temporary_path = unique_temporary_path(parent);
606 let mut options = OpenOptions::new();
607 options.write(true).create_new(true);
608 let mut temporary = root
609 .open_with(&temporary_path, &options)
610 .map_err(|error| map_cap_error("create temporary file for", &args.path, error))?;
611 let write_result = (|| {
612 temporary
613 .write_all(args.content.as_bytes())
614 .and_then(|_| temporary.sync_all())
615 .map_err(|error| ToolError(format!("write {}: {error}", args.path)))?;
616 if cancellation.is_cancelled() {
617 return Err(ToolError("write cancelled".into()));
618 }
619 match args.mode {
620 WriteMode::Create => root
621 .hard_link(&temporary_path, &root, &path)
622 .map_err(|error| map_cap_error("create", &args.path, error)),
623 WriteMode::Replace => root
624 .rename(&temporary_path, &root, &path)
625 .map_err(|error| map_cap_error("replace", &args.path, error)),
626 }
627 })();
628 if matches!(args.mode, WriteMode::Create) || write_result.is_err() {
629 let _ = root.remove_file(&temporary_path);
630 }
631 write_result?;
632 Ok(ToolOutput::success(
633 json!({
634 "path":args.path,
635 "bytes":args.content.len(),
636 "sha256":format!("{:x}", Sha256::digest(args.content.as_bytes()))
637 })
638 .to_string(),
639 ))
640 })
641 .await
642 .map_err(|error| ToolError(format!("write task failed: {error}")))?
643 }
644}
645
646struct BashTool {
647 timeout: Duration,
648 max_timeout: Duration,
649 output_limit: usize,
650}
651
652impl BashTool {
653 fn timeouts(&self) -> Timeouts {
654 Timeouts {
655 default: self.timeout,
656 max: self.max_timeout,
657 }
658 }
659}
660
661#[derive(Debug, Clone, Copy)]
663pub(crate) struct Timeouts {
664 pub(crate) default: Duration,
665 pub(crate) max: Duration,
666}
667
668impl Timeouts {
669 pub(crate) fn resolve(self, requested: Option<u64>) -> Result<Duration, ToolError> {
673 match requested {
674 None => Ok(self.default.min(self.max)),
675 Some(0) => Err(ToolError("timeout_seconds must be positive".into())),
676 Some(seconds) if seconds > self.max.as_secs() => Err(ToolError(format!(
677 "timeout_seconds {seconds} exceeds the configured maximum of {} seconds \
678 (tools.max_timeout_seconds)",
679 self.max.as_secs()
680 ))),
681 Some(seconds) => Ok(Duration::from_secs(seconds)),
682 }
683 }
684}
685
686pub(crate) fn timeout_schema(timeouts: Timeouts) -> Value {
687 json!({
688 "type":"integer",
689 "minimum":1,
690 "maximum":timeouts.max.as_secs(),
691 "description":format!(
692 "Seconds before the process is killed. Defaults to {}; at most {}. \
693 Raise it for long work such as builds, releases, or landing a change.",
694 timeouts.default.min(timeouts.max).as_secs(),
695 timeouts.max.as_secs()
696 )
697 })
698}
699
700#[derive(Deserialize)]
701#[serde(deny_unknown_fields)]
702struct BashArgs {
703 command: String,
704 timeout_seconds: Option<u64>,
705}
706
707#[async_trait]
708impl Tool for BashTool {
709 fn spec(&self) -> ToolSpec {
710 ToolSpec {
711 name: "bash".into(),
712 description: "Run a Bash command in the workspace (not sandboxed)".into(),
713 parameters: json!({
714 "type":"object",
715 "properties":{
716 "command":{"type":"string"},
717 "timeout_seconds":timeout_schema(self.timeouts())
718 },
719 "required":["command"],
720 "additionalProperties":false
721 }),
722 }
723 }
724
725 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
726 let args: BashArgs = parse_args(arguments)?;
727 validate_process_args(&args.command)?;
728 self.timeouts().resolve(args.timeout_seconds)?;
729 Ok(ToolRisk::Process)
730 }
731
732 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
733 let args: BashArgs = parse_args(arguments)?;
734 validate_process_args(&args.command)?;
735 self.timeouts().resolve(args.timeout_seconds)?;
736 Ok(format!(
737 "Run with /bin/bash -lc: {}",
738 bounded(&args.command, 2000)
739 ))
740 }
741
742 async fn execute(
743 &self,
744 arguments: Value,
745 context: ToolContext,
746 ) -> Result<ToolOutput, ToolError> {
747 let args: BashArgs = parse_args(&arguments)?;
748 validate_process_args(&args.command)?;
749 let requested = self.timeouts().resolve(args.timeout_seconds)?;
750 execute_process(
751 ProcessSpec {
752 executable: OsString::from("/bin/bash"),
753 args: vec![OsString::from("-lc"), OsString::from(args.command)],
754 cwd: context.workspace,
755 environment: Vec::new(),
756 sanitize_scv_environment: false,
757 timeout: requested,
758 output_limit: self.output_limit,
759 },
760 context.cancellation,
761 )
762 .await
763 }
764}
765
766struct NativeAgentTool {
767 name: String,
768 command: String,
769 resolved: Option<PathBuf>,
770 args: Vec<String>,
771 prompt_args: Vec<String>,
772 full_permission_args: Option<Vec<String>>,
773 model_args: Vec<String>,
774 effort_args: Vec<String>,
775 model_hint: String,
776 environment: Vec<(OsString, OsString)>,
777 timeouts: Timeouts,
778 output_limit: usize,
779 output: OutputFormat,
780 resume: Resume,
781 home: Option<PathBuf>,
782 delegation: Option<DelegationContext>,
783 conversations: Arc<ConversationStore>,
784}
785
786const AGENT_EFFORTS: [&str; 5] = ["low", "medium", "high", "xhigh", "max"];
788
789impl NativeAgentTool {
790 fn command_args(&self, args: &AgentArgs) -> Result<Vec<String>, ToolError> {
793 validate_process_args(&args.prompt)?;
794 self.timeouts.resolve(args.timeout_seconds)?;
795 if let Some(cwd) = &args.cwd {
796 validate_agent_cwd(cwd)?;
797 }
798 if let Some(session) = &args.session {
799 if !self.resume.is_supported() {
800 return Err(ToolError(format!(
801 "{} cannot continue a conversation; omit session to start a new one",
802 self.name
803 )));
804 }
805 if !conversation::is_handle(session) {
806 return Err(ToolError(format!(
807 "session {:?} is not a conversation handle; pass the `session` value an \
808 earlier {} call returned, or omit it to start a new conversation",
809 bounded(session, 80),
810 self.name
811 )));
812 }
813 }
814 if args.prompt.starts_with('-') {
817 return Err(ToolError("agent prompt must not start with '-'".into()));
818 }
819 let mut command = self.args.clone();
820 command.extend(self.full_permission_args.iter().flatten().cloned());
821 command.extend(self.output.args().iter().map(|arg| (*arg).to_owned()));
822 for (field, value, template, placeholder) in [
823 ("model", &args.model, &self.model_args, "{model}"),
824 ("effort", &args.effort, &self.effort_args, "{effort}"),
825 ] {
826 let Some(value) = value else {
827 continue;
828 };
829 if template.is_empty() {
830 return Err(ToolError(format!(
831 "{} does not support selecting a {field}",
832 self.name
833 )));
834 }
835 let valid = if field == "model" {
836 valid_model_name(value)
837 } else {
838 AGENT_EFFORTS.contains(&value.as_str())
839 };
840 if !valid {
841 return Err(ToolError(format!("invalid {field} {value:?}")));
842 }
843 command.extend(template.iter().map(|part| part.replace(placeholder, value)));
844 }
845 command.extend(self.prompt_args.iter().cloned());
846 Ok(command)
847 }
848 fn new(
849 name: String,
850 config: AgentAdapterConfig,
851 timeouts: Timeouts,
852 output_limit: usize,
853 delegation: Option<DelegationContext>,
854 conversations: Arc<ConversationStore>,
855 ) -> Self {
856 let resolved = adapters::resolve_agent_executable(&config.command, &config.search_dirs);
857 Self {
858 name,
859 command: config.command,
860 resolved,
861 args: config.args,
862 prompt_args: config.prompt_args,
863 full_permission_args: config.full_permission_args,
864 model_args: config.model_args,
865 effort_args: config.effort_args,
866 model_hint: config.model_hint,
867 environment: config.environment,
868 timeouts,
869 output_limit,
870 output: config.output,
871 resume: config.resume,
872 home: config.home,
873 delegation,
874 conversations,
875 }
876 }
877
878 fn run_args(&self, id: &str) -> (Vec<OsString>, Option<PathBuf>) {
882 match self.output {
883 OutputFormat::CodexJsonl => {
884 let Some(dir) = self.home.as_ref().map(|home| home.join("tmp")) else {
885 return (Vec::new(), None);
886 };
887 if private_dir(&dir).is_err() {
888 return (Vec::new(), None);
889 }
890 let file = dir.join(format!("scv-{id}.last-message"));
891 (vec!["-o".into(), file.clone().into()], Some(file))
892 }
893 OutputFormat::Text | OutputFormat::ClaudeStreamJson | OutputFormat::PiJson => {
894 (Vec::new(), None)
895 }
896 }
897 }
898}
899
900fn session_args(template: &[&str], session: &str) -> Vec<OsString> {
902 template
903 .iter()
904 .map(|part| OsString::from(part.replace("{session}", session)))
905 .collect()
906}
907
908fn private_dir(dir: &Path) -> std::io::Result<()> {
909 use std::os::unix::fs::PermissionsExt as _;
910 std::fs::create_dir_all(dir)?;
911 std::fs::set_permissions(dir, std::fs::Permissions::from_mode(0o700))
912}
913
914fn take_file(path: &Path, limit: usize) -> Option<String> {
916 let file = std::fs::File::open(path).ok();
917 let _ = std::fs::remove_file(path);
918 let mut bytes = Vec::new();
919 std::io::Read::take(file?, u64::try_from(limit).unwrap_or(u64::MAX))
920 .read_to_end(&mut bytes)
921 .ok()?;
922 let text = String::from_utf8_lossy(&bytes).trim().to_owned();
923 (!text.is_empty()).then_some(text)
924}
925
926#[derive(Deserialize)]
927#[serde(deny_unknown_fields)]
928pub(crate) struct AgentArgs {
929 pub(crate) prompt: String,
930 pub(crate) timeout_seconds: Option<u64>,
931 #[serde(default, deserialize_with = "blank_as_none")]
932 pub(crate) session: Option<String>,
933 #[serde(default, deserialize_with = "blank_as_none")]
934 pub(crate) cwd: Option<String>,
935 #[serde(default, deserialize_with = "blank_as_none")]
936 pub(crate) model: Option<String>,
937 #[serde(default, deserialize_with = "blank_as_none")]
938 pub(crate) effort: Option<String>,
939}
940
941fn blank_as_none<'de, D: serde::Deserializer<'de>>(
944 deserializer: D,
945) -> Result<Option<String>, D::Error> {
946 let value = Option::<String>::deserialize(deserializer)?;
947 Ok(value.filter(|value| !value.trim().is_empty()))
948}
949
950const MAX_AGENT_CWD_BYTES: usize = 4096;
952
953pub(crate) fn validate_agent_cwd(cwd: &str) -> Result<(), ToolError> {
954 if cwd.trim().is_empty() || cwd.len() > MAX_AGENT_CWD_BYTES || cwd.contains('\0') {
955 return Err(ToolError(format!(
956 "cwd must be a non-empty directory path of at most {MAX_AGENT_CWD_BYTES} bytes"
957 )));
958 }
959 Ok(())
960}
961
962pub(crate) fn resolve_agent_cwd(workspace: &Path, cwd: Option<&str>) -> Result<PathBuf, ToolError> {
966 let root = std::fs::canonicalize(workspace)
967 .map_err(|error| ToolError(format!("resolve workspace: {error}")))?;
968 let Some(cwd) = cwd else {
969 return Ok(root);
970 };
971 validate_agent_cwd(cwd)?;
972 let resolved = std::fs::canonicalize(root.join(cwd))
973 .map_err(|error| ToolError(format!("cwd {cwd:?}: {error}")))?;
974 if !resolved.starts_with(&root) {
975 return Err(ToolError(format!("cwd {cwd:?} is outside the workspace")));
976 }
977 if !resolved.is_dir() {
978 return Err(ToolError(format!("cwd {cwd:?} is not a directory")));
979 }
980 Ok(resolved)
981}
982
983pub(crate) fn valid_model_name(value: &str) -> bool {
986 !value.is_empty()
987 && value.len() <= 128
988 && !value.starts_with(['-', '@'])
989 && value
990 .chars()
991 .all(|c| c.is_ascii_alphanumeric() || "._:/@[]-".contains(c))
992}
993
994#[async_trait]
995impl Tool for NativeAgentTool {
996 fn spec(&self) -> ToolSpec {
997 let mut properties = json!({
998 "prompt":{"type":"string"},
999 "cwd":{
1000 "type":"string",
1001 "description":"Directory inside the workspace to run in, such as a project directory (\"scv\"). \
1002 The agent loads that directory's AGENTS.md or CLAUDE.md and its project skills. \
1003 Defaults to the workspace root."
1004 },
1005 "timeout_seconds":timeout_schema(self.timeouts)
1006 });
1007 if self.resume.is_supported() {
1008 properties["session"] = json!({
1009 "type":"string",
1010 "description":"The `session` handle an earlier call to this tool returned, such as \"codex-1\". \
1011 Pass it to continue that conversation: the agent keeps its context, in the same cwd. \
1012 Omit it to start a new conversation for unrelated work."
1013 });
1014 }
1015 if !self.model_args.is_empty() {
1016 properties["model"] = json!({
1017 "type":"string",
1018 "description":format!(
1019 "{} Set only when the user asks for a specific model; \
1020 omit to use the agent's configured default.",
1021 self.model_hint
1022 )
1023 });
1024 }
1025 if !self.effort_args.is_empty() {
1026 properties["effort"] = json!({
1027 "type":"string",
1028 "enum":AGENT_EFFORTS,
1029 "description":"Reasoning effort. Set only when the user asks for one; \
1030 omit to use the agent's configured default."
1031 });
1032 }
1033 ToolSpec {
1034 name: self.name.clone(),
1035 description: format!(
1036 "Launch the configured {} CLI as a nested coding agent (not sandboxed). \
1037 Delegate substantial work here rather than doing it step by step with \
1038 bash: research and web lookups, multi-file coding, and running tools, \
1039 builds, and tests. Set cwd to the project the work is in so the agent \
1040 follows that project's instructions and skills.{}",
1041 self.name,
1042 if self.resume.is_supported() {
1043 " Each result carries a `session` handle: pass it back to follow up on the \
1044 same work (answers, fixes, next steps) instead of repeating the context."
1045 } else {
1046 " Each call starts a fresh conversation."
1047 }
1048 ),
1049 parameters: json!({
1050 "type":"object",
1051 "properties":properties,
1052 "required":["prompt"],
1053 "additionalProperties":false
1054 }),
1055 }
1056 }
1057
1058 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
1059 let args: AgentArgs = parse_args(arguments)?;
1060 self.command_args(&args)?;
1061 Ok(ToolRisk::Delegate)
1062 }
1063
1064 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
1065 let args: AgentArgs = parse_args(arguments)?;
1066 let command_args = self.command_args(&args)?;
1067 let executable = self.resolved.as_ref().map_or_else(
1068 || self.command.as_str().into(),
1069 |path| path.display().to_string(),
1070 );
1071 let directory = args.cwd.as_deref().map_or_else(
1072 || "the workspace root".to_owned(),
1073 |cwd| format!("{:?} (inside the workspace)", bounded(cwd, 200)),
1074 );
1075 let conversation = args.session.as_deref().map_or_else(
1076 || " in a new conversation".to_owned(),
1077 |session| format!(", continuing conversation {session},"),
1078 );
1079 let timeout = self.timeouts.resolve(args.timeout_seconds)?;
1080 let permissions = if self.full_permission_args.is_some() {
1081 " FULL PERMISSIONS (permissions = \"full\"): the agent's own approval prompts \
1082 and sandbox are off, so it edits files, runs commands, and uses the network \
1083 without asking."
1084 } else {
1085 ""
1086 };
1087 Ok(format!(
1088 "Launch {executable} with args {command_args:?} and prompt {:?}{conversation} in {directory} for up to {} seconds. The nested agent has your user permissions.{permissions}",
1089 bounded(&args.prompt, 2000),
1090 timeout.as_secs()
1091 ))
1092 }
1093
1094 async fn execute(
1095 &self,
1096 arguments: Value,
1097 context: ToolContext,
1098 ) -> Result<ToolOutput, ToolError> {
1099 let args: AgentArgs = parse_args(&arguments)?;
1100 let command_args = self.command_args(&args)?;
1101 let cwd = resolve_agent_cwd(&context.workspace, args.cwd.as_deref())?;
1102 let executable = self.resolved.as_ref().ok_or_else(|| {
1103 ToolError(format!(
1104 "{} executable {:?} was not found on PATH or in the user's install directories",
1105 self.name, self.command
1106 ))
1107 })?;
1108 let agent = self.name.trim_start_matches("agent_");
1109 let turn = if self.resume.is_supported() {
1110 Some(self.conversations.begin(
1111 agent,
1112 args.session.as_deref(),
1113 &cwd,
1114 self.resume.assigns_id(),
1115 )?)
1116 } else {
1117 None
1118 };
1119 let pending = self.delegation.as_ref().map(|delegation| {
1120 delegation.registry.begin_at(
1121 delegation.owner_depth(),
1122 agent,
1123 &delegation.session,
1124 &cwd,
1125 turn.as_ref().map(|turn| (turn.handle.as_str(), turn.turn)),
1126 )
1127 });
1128 let run_id = pending.as_ref().map_or_else(
1129 || uuid::Uuid::new_v4().simple().to_string(),
1130 |pending| pending.handle.clone(),
1131 );
1132 let fixed_len = self.args.len()
1133 + self.full_permission_args.as_ref().map_or(0, Vec::len)
1134 + self.output.args().len();
1135 let (fixed, rest) = command_args.split_at(fixed_len);
1136 let (selections, prompt_args) = rest.split_at(rest.len() - self.prompt_args.len());
1137 let (base, modes) = fixed.split_at(self.args.len());
1138 let continuing = args.session.is_some();
1139 let (run_args, last_message) = self.run_args(&run_id);
1140 let resume = match (self.resume, &turn) {
1141 (
1142 Resume::Supported {
1143 start,
1144 subcommand,
1145 options,
1146 positional,
1147 },
1148 Some(turn),
1149 ) => {
1150 let vendor = turn.vendor.as_deref().unwrap_or_default();
1151 if continuing {
1152 (
1153 subcommand.iter().map(OsString::from).collect(),
1154 session_args(options, vendor),
1155 session_args(positional, vendor),
1156 )
1157 } else if turn.vendor.is_some() {
1158 (Vec::new(), session_args(start, vendor), Vec::new())
1159 } else {
1160 Default::default()
1161 }
1162 }
1163 _ => Default::default(),
1164 };
1165 let (subcommand, session_options, positional): (
1166 Vec<OsString>,
1167 Vec<OsString>,
1168 Vec<OsString>,
1169 ) = resume;
1170 let mut process_args: Vec<OsString> = base.iter().map(OsString::from).collect();
1171 process_args.extend(subcommand);
1172 process_args.extend(modes.iter().map(OsString::from));
1173 process_args.extend(run_args);
1174 process_args.extend(session_options);
1175 process_args.extend(selections.iter().map(OsString::from));
1176 process_args.extend(positional);
1177 process_args.extend(prompt_args.iter().map(OsString::from));
1178 process_args.push(OsString::from(args.prompt));
1179 let mut environment = self.environment.clone();
1180 match &pending {
1181 Some(pending) => environment.extend(pending.environment.iter().cloned()),
1182 None => environment.push((
1183 delegation::DEPTH_VARIABLE.into(),
1184 (delegation::current_depth() + 1).to_string().into(),
1185 )),
1186 }
1187 let requested = self.timeouts.resolve(args.timeout_seconds)?;
1188 let registration = self
1189 .delegation
1190 .as_ref()
1191 .map(|delegation| Arc::clone(&delegation.registry))
1192 .zip(pending);
1193 let run = execute_agent_process(
1194 ProcessSpec {
1195 executable: executable.as_os_str().to_owned(),
1196 args: process_args,
1197 cwd,
1198 environment,
1199 sanitize_scv_environment: true,
1200 timeout: requested,
1201 output_limit: self.output_limit,
1202 },
1203 AgentStream::new(self.output, self.output_limit).with_progress(context.progress),
1204 registration,
1205 context.cancellation,
1206 )
1207 .await;
1208 let fallback = last_message
1209 .as_deref()
1210 .and_then(|path| take_file(path, self.output_limit));
1211 let run = run?;
1212 let result = run.stream.finish(run.exit, fallback);
1213 let conversation = turn.and_then(|turn| {
1214 let number = turn.turn;
1215 turn.finish(
1216 result.session.clone(),
1217 result.status == agent_output::RunStatus::Completed,
1218 )
1219 .map(|handle| (handle, number))
1220 });
1221 let (content, truncated) = result.to_json(
1222 agent,
1223 conversation
1224 .as_ref()
1225 .map(|(handle, turn)| (handle.as_str(), *turn)),
1226 run.exit_code,
1227 &run.stderr_tail,
1228 self.output_limit,
1229 );
1230 let mut output = ToolOutput {
1231 content,
1232 is_error: result.status != agent_output::RunStatus::Completed,
1233 truncated,
1234 };
1235 if output.is_error {
1236 add_sign_in_hint(&mut output, agent);
1237 }
1238 Ok(output)
1239 }
1240}
1241
1242pub(crate) fn add_sign_in_hint(output: &mut ToolOutput, agent: &str) {
1246 let lower = output.content.to_ascii_lowercase();
1247 let unauthenticated = [
1248 "not logged in",
1249 "not signed in",
1250 "not authenticated",
1251 "login",
1252 "log in",
1253 "unauthorized",
1254 "authentication",
1255 "missing_credential",
1256 "no api key",
1257 "auth_required",
1258 ]
1259 .iter()
1260 .any(|needle| lower.contains(needle));
1261 if !unauthenticated {
1262 return;
1263 }
1264 if let Ok(Value::Object(mut content)) = serde_json::from_str::<Value>(&output.content) {
1265 content.insert(
1266 "hint".into(),
1267 format!(
1268 "The {agent} CLI appears to be signed out of SCV's private agent home. \
1269 The host owner can sign it in with: scv agents login {agent}"
1270 )
1271 .into(),
1272 );
1273 output.content = Value::Object(content).to_string();
1274 }
1275}
1276
1277pub fn apply_agent_environment(
1281 command: &mut std::process::Command,
1282 environment: &[(OsString, OsString)],
1283) {
1284 apply_agent_environment_from(
1285 command,
1286 std::env::vars_os().map(|(variable, _)| variable),
1287 environment,
1288 );
1289}
1290
1291fn apply_agent_environment_from(
1292 command: &mut std::process::Command,
1293 inherited: impl IntoIterator<Item = OsString>,
1294 environment: &[(OsString, OsString)],
1295) {
1296 for variable in inherited {
1297 if adapters::is_removed_agent_variable(&variable) {
1298 command.env_remove(variable);
1299 }
1300 }
1301 command.envs(environment.iter().map(|(key, value)| (key, value)));
1302}
1303
1304struct ProcessSpec {
1305 executable: OsString,
1306 args: Vec<OsString>,
1307 cwd: PathBuf,
1308 environment: Vec<(OsString, OsString)>,
1309 sanitize_scv_environment: bool,
1310 timeout: Duration,
1311 output_limit: usize,
1312}
1313
1314async fn execute_process(
1315 spec: ProcessSpec,
1316 cancellation: tokio_util::sync::CancellationToken,
1317) -> Result<ToolOutput, ToolError> {
1318 let deadline = Instant::now() + spec.timeout;
1319 let mut child = spawn_process(&spec)?;
1320 let pid = child_pid(&child)?;
1321 let output = Arc::new(Mutex::new(BoundedOutput::new(spec.output_limit)));
1322 let stdout_task = child
1323 .stdout
1324 .take()
1325 .map(|stdout| tokio::spawn(drain_output(stdout, Arc::clone(&output))));
1326 let stderr_task = child
1327 .stderr
1328 .take()
1329 .map(|stderr| tokio::spawn(drain_output(stderr, Arc::clone(&output))));
1330 let finished = supervise(
1331 &mut child,
1332 pid,
1333 deadline,
1334 cancellation,
1335 stdout_task,
1336 stderr_task,
1337 )
1338 .await;
1339 delegation::untrack_spawned(pid as u32);
1340 let finished = finished?;
1341 let collected = output.lock().await;
1342 let text = String::from_utf8_lossy(&collected.bytes).into_owned();
1343 let content = json!({
1344 "exit_code": finished.status.code(),
1345 "timed_out": finished.timed_out,
1346 "output": text,
1347 "truncated": collected.truncated
1348 })
1349 .to_string();
1350 Ok(ToolOutput {
1351 content,
1352 is_error: finished.timed_out || !finished.status.success(),
1353 truncated: collected.truncated,
1354 })
1355}
1356
1357struct AgentRun {
1359 stream: AgentStream,
1360 exit: RunExit,
1361 exit_code: Option<i32>,
1362 stderr_tail: String,
1363}
1364
1365async fn execute_agent_process(
1369 spec: ProcessSpec,
1370 stream: AgentStream,
1371 registration: Option<(Arc<DelegationRegistry>, delegation::PendingDelegation)>,
1372 cancellation: tokio_util::sync::CancellationToken,
1373) -> Result<AgentRun, ToolError> {
1374 let deadline = Instant::now() + spec.timeout;
1375 let mut child = spawn_process(&spec)?;
1376 let pid = child_pid(&child)?;
1377 let guard: Option<DelegationGuard> =
1380 registration.and_then(|(registry, pending)| registry.register(pending, pid as u32).ok());
1381 let stdout = Arc::new(Mutex::new(stream));
1382 let stderr = Arc::new(Mutex::new(TailBuffer::new(STDERR_TAIL_BYTES)));
1383 let stdout_task = child
1384 .stdout
1385 .take()
1386 .map(|reader| tokio::spawn(drain_output(reader, Arc::clone(&stdout))));
1387 let stderr_task = child
1388 .stderr
1389 .take()
1390 .map(|reader| tokio::spawn(drain_output(reader, Arc::clone(&stderr))));
1391 let finished = supervise(
1392 &mut child,
1393 pid,
1394 deadline,
1395 cancellation,
1396 stdout_task,
1397 stderr_task,
1398 )
1399 .await;
1400 delegation::untrack_spawned(pid as u32);
1401 let killed = guard.as_ref().is_some_and(DelegationGuard::was_killed);
1402 if let Some(guard) = guard {
1403 guard.finish().await;
1404 }
1405 let finished = finished?;
1406 let exit = if finished.timed_out {
1407 RunExit::TimedOut
1408 } else if killed {
1409 RunExit::Killed
1410 } else {
1411 RunExit::Exited {
1412 success: finished.status.success(),
1413 }
1414 };
1415 let stderr_tail = stderr.lock().await.text();
1416 let stream = Arc::try_unwrap(stdout)
1417 .map_err(|_| ToolError("agent output reader is still running".into()))?
1418 .into_inner();
1419 Ok(AgentRun {
1420 stream,
1421 exit,
1422 exit_code: finished.status.code(),
1423 stderr_tail,
1424 })
1425}
1426
1427fn spawn_process(spec: &ProcessSpec) -> Result<tokio::process::Child, ToolError> {
1428 let mut command = Command::new(&spec.executable);
1429 if spec.sanitize_scv_environment {
1430 apply_agent_environment(command.as_std_mut(), &spec.environment);
1431 } else {
1432 command.envs(spec.environment.iter().map(|(key, value)| (key, value)));
1433 }
1434 command
1435 .args(&spec.args)
1436 .current_dir(&spec.cwd)
1437 .stdin(std::process::Stdio::null())
1438 .stdout(std::process::Stdio::piped())
1439 .stderr(std::process::Stdio::piped())
1440 .kill_on_drop(true);
1441 command.as_std_mut().process_group(0);
1442 let child = command
1443 .spawn()
1444 .map_err(|error| ToolError(format!("launch {:?}: {error}", spec.executable)))?;
1445 if let Some(pid) = child.id() {
1446 delegation::track_spawned(pid);
1447 }
1448 Ok(child)
1449}
1450
1451fn child_pid(child: &tokio::process::Child) -> Result<i32, ToolError> {
1452 child
1453 .id()
1454 .and_then(|pid| i32::try_from(pid).ok())
1455 .ok_or_else(|| ToolError("child process has no pid".into()))
1456}
1457
1458struct Finished {
1459 status: std::process::ExitStatus,
1460 timed_out: bool,
1461}
1462
1463async fn supervise(
1466 child: &mut tokio::process::Child,
1467 pid: i32,
1468 deadline: Instant,
1469 cancellation: tokio_util::sync::CancellationToken,
1470 stdout_task: Option<JoinHandle<()>>,
1471 stderr_task: Option<JoinHandle<()>>,
1472) -> Result<Finished, ToolError> {
1473 enum Completion {
1474 Exited(std::process::ExitStatus),
1475 TimedOut,
1476 Cancelled,
1477 }
1478 let completion = tokio::select! {
1479 status = child.wait() => Completion::Exited(status.map_err(|error| ToolError(format!("wait for child: {error}")))?),
1480 _ = cancellation.cancelled() => {
1481 Completion::Cancelled
1482 },
1483 _ = sleep_until(deadline) => Completion::TimedOut,
1484 };
1485
1486 let (status, timed_out, drain_deadline) = match completion {
1487 Completion::Exited(status) => {
1488 let cleanup_deadline = deadline.min(Instant::now() + Duration::from_secs(2));
1489 let status = terminate_group(pid, child, Some(status), cleanup_deadline, true).await?;
1490 (
1491 status,
1492 false,
1493 deadline.min(Instant::now() + Duration::from_millis(250)),
1494 )
1495 }
1496 Completion::TimedOut => {
1497 let status = terminate_group(pid, child, None, Instant::now(), false).await?;
1498 (status, true, Instant::now() + Duration::from_millis(250))
1499 }
1500 Completion::Cancelled => {
1501 let cleanup_deadline = Instant::now() + Duration::from_secs(2);
1502 let _ = terminate_group(pid, child, None, cleanup_deadline, true).await;
1503 finish_drain(stdout_task, Instant::now() + Duration::from_millis(250)).await;
1504 finish_drain(stderr_task, Instant::now() + Duration::from_millis(250)).await;
1505 return Err(ToolError("process cancelled".into()));
1506 }
1507 };
1508 finish_drain(stdout_task, drain_deadline).await;
1509 finish_drain(stderr_task, drain_deadline).await;
1510 Ok(Finished { status, timed_out })
1511}
1512
1513async fn terminate_group(
1514 pid: i32,
1515 child: &mut tokio::process::Child,
1516 mut status: Option<std::process::ExitStatus>,
1517 deadline: Instant,
1518 graceful: bool,
1519) -> Result<std::process::ExitStatus, ToolError> {
1520 signal_group(
1521 pid,
1522 if graceful {
1523 libc::SIGTERM
1524 } else {
1525 libc::SIGKILL
1526 },
1527 );
1528 while Instant::now() < deadline {
1529 if status.is_none() {
1530 status = child
1531 .try_wait()
1532 .map_err(|error| ToolError(format!("wait for child: {error}")))?;
1533 }
1534 if !process_group_exists(pid)
1535 && let Some(status) = status
1536 {
1537 return Ok(status);
1538 }
1539 sleep(Duration::from_millis(20)).await;
1540 }
1541 signal_group(pid, libc::SIGKILL);
1543 if let Some(status) = status {
1544 return Ok(status);
1545 }
1546 timeout(Duration::from_secs(1), child.wait())
1547 .await
1548 .map_err(|_| ToolError("child did not exit after process-group kill".into()))?
1549 .map_err(|error| ToolError(format!("wait after KILL: {error}")))
1550}
1551
1552fn signal_group(pid: i32, signal: i32) {
1553 unsafe {
1555 libc::kill(-pid, signal);
1556 }
1557}
1558
1559fn process_group_exists(pid: i32) -> bool {
1560 let result = unsafe { libc::kill(-pid, 0) };
1561 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
1562}
1563
1564async fn finish_drain(task: Option<JoinHandle<()>>, deadline: Instant) {
1565 let Some(mut task) = task else { return };
1566 if timeout_at(deadline, &mut task).await.is_err() {
1567 task.abort();
1568 let _ = task.await;
1569 }
1570}
1571
1572pub(crate) trait OutputSink: Send + 'static {
1574 fn push(&mut self, bytes: &[u8]);
1575}
1576
1577impl OutputSink for BoundedOutput {
1578 fn push(&mut self, bytes: &[u8]) {
1579 BoundedOutput::push(self, bytes);
1580 }
1581}
1582
1583impl OutputSink for AgentStream {
1584 fn push(&mut self, bytes: &[u8]) {
1585 AgentStream::push(self, bytes);
1586 }
1587}
1588
1589impl OutputSink for TailBuffer {
1590 fn push(&mut self, bytes: &[u8]) {
1591 TailBuffer::push(self, bytes);
1592 }
1593}
1594
1595pub(crate) async fn drain_output<R, S>(mut reader: R, output: Arc<Mutex<S>>)
1596where
1597 R: tokio::io::AsyncRead + Unpin,
1598 S: OutputSink,
1599{
1600 let mut chunk = [0u8; 8192];
1601 loop {
1602 match reader.read(&mut chunk).await {
1603 Ok(0) | Err(_) => break,
1604 Ok(read) => output.lock().await.push(&chunk[..read]),
1605 }
1606 }
1607}
1608
1609struct BoundedOutput {
1610 bytes: Vec<u8>,
1611 limit: usize,
1612 truncated: bool,
1613}
1614
1615impl BoundedOutput {
1616 fn new(limit: usize) -> Self {
1617 Self {
1618 bytes: Vec::with_capacity(limit.min(8192)),
1619 limit,
1620 truncated: false,
1621 }
1622 }
1623
1624 fn push(&mut self, bytes: &[u8]) {
1625 let remaining = self.limit.saturating_sub(self.bytes.len());
1626 self.bytes
1627 .extend_from_slice(&bytes[..bytes.len().min(remaining)]);
1628 self.truncated |= bytes.len() > remaining;
1629 }
1630}
1631
1632pub(crate) fn parse_args<T: for<'de> Deserialize<'de>>(value: &Value) -> Result<T, ToolError> {
1633 serde_json::from_value(value.clone())
1634 .map_err(|error| ToolError(format!("invalid arguments: {error}")))
1635}
1636
1637fn validate_read_args(args: &ReadArgs) -> Result<(), ToolError> {
1638 if args.limit == Some(0) {
1639 return Err(ToolError("read limit must be positive".into()));
1640 }
1641 Ok(())
1642}
1643
1644pub(crate) fn validate_process_args(value: &str) -> Result<(), ToolError> {
1645 if value.trim().is_empty() {
1646 return Err(ToolError("command or prompt must be non-empty".into()));
1647 }
1648 Ok(())
1649}
1650
1651fn validate_relative(path: &Path) -> Result<(), ToolError> {
1652 if path.as_os_str().is_empty() || path.is_absolute() {
1653 return Err(ToolError("path must be non-empty and relative".into()));
1654 }
1655 for component in path.components() {
1656 if matches!(
1657 component,
1658 Component::ParentDir | Component::RootDir | Component::Prefix(_)
1659 ) {
1660 return Err(ToolError(
1661 "parent traversal and absolute paths are not allowed".into(),
1662 ));
1663 }
1664 }
1665 Ok(())
1666}
1667
1668static TEMPORARY_COUNTER: AtomicU64 = AtomicU64::new(0);
1669
1670fn open_workspace(workspace: &Path) -> Result<Dir, ToolError> {
1671 Dir::open_ambient_dir(workspace, ambient_authority())
1672 .map_err(|error| ToolError(format!("open workspace capability: {error}")))
1673}
1674
1675fn unique_temporary_path(parent: &Path) -> PathBuf {
1676 let id = TEMPORARY_COUNTER.fetch_add(1, Ordering::Relaxed);
1677 parent.join(format!(".scv-write-{}-{id}.tmp", std::process::id()))
1678}
1679
1680fn map_cap_error(action: &str, path: &str, error: std::io::Error) -> ToolError {
1681 ToolError(format!(
1682 "{action} {path}: {error}; path must remain within workspace"
1683 ))
1684}
1685
1686fn is_secret_like(path: &Path) -> bool {
1687 path.components().any(|component| {
1688 let value = component.as_os_str().to_string_lossy().to_ascii_lowercase();
1689 value == ".env"
1690 || value.starts_with(".env.")
1691 || value.contains("credential")
1692 || value.contains("private_key")
1693 || value.ends_with(".pem")
1694 || value.ends_with(".key")
1695 })
1696}
1697
1698pub(crate) fn bounded(value: &str, max_chars: usize) -> String {
1699 let mut output: String = value.chars().take(max_chars).collect();
1700 if value.chars().count() > max_chars {
1701 output.push('โฆ');
1702 }
1703 output
1704}
1705
1706#[cfg(test)]
1707mod tests {
1708 use std::os::unix::fs::symlink;
1709
1710 use super::*;
1711
1712 fn test_conversations() -> Arc<ConversationStore> {
1713 Arc::new(ConversationStore::new(
1714 ToolsConfig::default().conversations,
1715 None,
1716 ))
1717 }
1718
1719 const SHELL_STARTUP: Duration = Duration::from_secs(30);
1724
1725 async fn wait_for<T>(limit: Duration, mut probe: impl FnMut() -> Option<T>) -> Option<T> {
1727 let deadline = std::time::Instant::now() + limit;
1728 loop {
1729 if let Some(value) = probe() {
1730 return Some(value);
1731 }
1732 if std::time::Instant::now() >= deadline {
1733 return None;
1734 }
1735 tokio::time::sleep(Duration::from_millis(10)).await;
1736 }
1737 }
1738
1739 fn is_gone(pid: i32) -> Option<()> {
1740 (unsafe { libc::kill(pid, 0) } != 0).then_some(())
1741 }
1742
1743 #[test]
1744 fn rejects_parent_traversal() {
1745 assert!(validate_relative(Path::new("../secret")).is_err());
1746 assert!(validate_relative(Path::new("/etc/passwd")).is_err());
1747 }
1748
1749 #[test]
1750 fn detects_secret_like_paths() {
1751 assert!(is_secret_like(Path::new(".env")));
1752 assert!(is_secret_like(Path::new("keys/id.pem")));
1753 assert!(!is_secret_like(Path::new("src/main.rs")));
1754 }
1755
1756 #[tokio::test]
1757 async fn read_is_contained_and_bounded() {
1758 let directory = tempfile::tempdir().unwrap();
1759 std::fs::write(directory.path().join("hello.txt"), "abcdef").unwrap();
1760 let tool = ReadTool { max_bytes: 3 };
1761 let output = tool
1762 .execute(
1763 json!({"path":"hello.txt"}),
1764 ToolContext::new(
1765 directory.path().canonicalize().unwrap(),
1766 tokio_util::sync::CancellationToken::new(),
1767 ),
1768 )
1769 .await
1770 .unwrap();
1771 assert!(output.truncated);
1772 assert!(output.content.contains("abc"));
1773 }
1774
1775 #[tokio::test]
1776 async fn read_rejects_symlink_escape() {
1777 let workspace = tempfile::tempdir().unwrap();
1778 let outside = tempfile::tempdir().unwrap();
1779 std::fs::write(outside.path().join("secret"), "nope").unwrap();
1780 symlink(outside.path(), workspace.path().join("escape")).unwrap();
1781 let tool = ReadTool { max_bytes: 100 };
1782 let result = tool
1783 .execute(
1784 json!({"path":"escape/secret"}),
1785 ToolContext::new(
1786 workspace.path().canonicalize().unwrap(),
1787 tokio_util::sync::CancellationToken::new(),
1788 ),
1789 )
1790 .await;
1791 assert!(result.unwrap_err().to_string().contains("workspace"));
1792 }
1793
1794 #[tokio::test]
1795 async fn write_is_atomic_and_checks_hash() {
1796 let workspace = tempfile::tempdir().unwrap();
1797 let root = workspace.path().canonicalize().unwrap();
1798 let tool = WriteTool { max_bytes: 100 };
1799 tool.execute(
1800 json!({"path":"file.txt","content":"first","mode":"create"}),
1801 ToolContext::new(root.clone(), tokio_util::sync::CancellationToken::new()),
1802 )
1803 .await
1804 .unwrap();
1805 let hash = format!("{:x}", Sha256::digest(b"first"));
1806 tool.execute(
1807 json!({"path":"file.txt","content":"second","mode":"replace","expected_sha256":hash}),
1808 ToolContext::new(root.clone(), tokio_util::sync::CancellationToken::new()),
1809 )
1810 .await
1811 .unwrap();
1812 assert_eq!(
1813 std::fs::read_to_string(root.join("file.txt")).unwrap(),
1814 "second"
1815 );
1816 let result = tool
1817 .execute(
1818 json!({"path":"file.txt","content":"third","mode":"replace","expected_sha256":"deadbeef"}),
1819 ToolContext::new(root, tokio_util::sync::CancellationToken::new()),
1820 )
1821 .await;
1822 assert!(result.unwrap_err().to_string().contains("changed"));
1823 }
1824
1825 #[tokio::test]
1826 async fn write_rejects_symlink_escape() {
1827 let workspace = tempfile::tempdir().unwrap();
1828 let outside = tempfile::tempdir().unwrap();
1829 symlink(outside.path(), workspace.path().join("escape")).unwrap();
1830 let tool = WriteTool { max_bytes: 100 };
1831 let result = tool
1832 .execute(
1833 json!({"path":"escape/file.txt","content":"nope","mode":"create"}),
1834 ToolContext::new(
1835 workspace.path().canonicalize().unwrap(),
1836 tokio_util::sync::CancellationToken::new(),
1837 ),
1838 )
1839 .await;
1840 assert!(result.unwrap_err().to_string().contains("workspace"));
1841 assert!(!outside.path().join("file.txt").exists());
1842 }
1843
1844 #[tokio::test]
1845 async fn bash_timeout_terminates_the_process() {
1846 let workspace = tempfile::tempdir().unwrap();
1847 let tool = BashTool {
1848 timeout: Duration::from_millis(50),
1849 max_timeout: Duration::from_millis(50),
1850 output_limit: 100,
1851 };
1852 let started = std::time::Instant::now();
1853 let output = tool
1854 .execute(
1855 json!({"command":"sleep 5"}),
1856 ToolContext::new(
1857 workspace.path().canonicalize().unwrap(),
1858 tokio_util::sync::CancellationToken::new(),
1859 ),
1860 )
1861 .await
1862 .unwrap();
1863 assert!(output.is_error);
1864 assert!(started.elapsed() < Duration::from_secs(3));
1865 }
1866
1867 #[tokio::test]
1868 async fn bash_output_is_bounded_and_reports_truncation() {
1869 let workspace = tempfile::tempdir().unwrap();
1870 let tool = BashTool {
1871 timeout: SHELL_STARTUP,
1872 max_timeout: SHELL_STARTUP,
1873 output_limit: 8,
1874 };
1875 let output = tool
1876 .execute(
1877 json!({"command":"printf 12345678901234567890"}),
1878 ToolContext::new(
1879 workspace.path().canonicalize().unwrap(),
1880 tokio_util::sync::CancellationToken::new(),
1881 ),
1882 )
1883 .await
1884 .unwrap();
1885 assert!(output.truncated);
1886 assert!(output.content.contains("12345678"));
1887 assert!(!output.content.contains("123456789"));
1888 }
1889
1890 #[tokio::test]
1891 async fn bash_cancellation_terminates_the_process_group() {
1892 let workspace = tempfile::tempdir().unwrap();
1893 let tool = BashTool {
1894 timeout: Duration::from_secs(30),
1895 max_timeout: Duration::from_secs(30),
1896 output_limit: 100,
1897 };
1898 let cancellation = tokio_util::sync::CancellationToken::new();
1899 let cancel = cancellation.clone();
1900 let started = std::time::Instant::now();
1901 let execution = tokio::spawn(async move {
1902 tool.execute(
1903 json!({"command":"sleep 30"}),
1904 ToolContext::new(workspace.path().canonicalize().unwrap(), cancellation),
1905 )
1906 .await
1907 });
1908 tokio::time::sleep(Duration::from_millis(50)).await;
1909 cancel.cancel();
1910 let error = execution.await.unwrap().unwrap_err();
1911 assert!(error.to_string().contains("cancelled"));
1912 assert!(started.elapsed() < Duration::from_secs(3));
1913 }
1914
1915 #[tokio::test]
1916 async fn background_descendant_cannot_hold_output_pipes_open() {
1917 let workspace = tempfile::tempdir().unwrap();
1918 let root = workspace.path().canonicalize().unwrap();
1919 let tool = BashTool {
1920 timeout: SHELL_STARTUP,
1921 max_timeout: SHELL_STARTUP,
1922 output_limit: 100,
1923 };
1924 let output = tool
1925 .execute(
1926 json!({"command":"sleep 60 & echo $! > background.pid; exit 0"}),
1927 ToolContext::new(root.clone(), tokio_util::sync::CancellationToken::new()),
1928 )
1929 .await
1930 .unwrap();
1931 let returned = std::time::SystemTime::now();
1932 assert!(!output.is_error);
1933 let exited = std::fs::metadata(root.join("background.pid"))
1935 .unwrap()
1936 .modified()
1937 .unwrap();
1938 assert!(returned.duration_since(exited).unwrap_or_default() < Duration::from_secs(3));
1939 let pid: i32 = std::fs::read_to_string(root.join("background.pid"))
1940 .unwrap()
1941 .trim()
1942 .parse()
1943 .unwrap();
1944 assert!(
1945 wait_for(Duration::from_secs(5), || is_gone(pid))
1946 .await
1947 .is_some(),
1948 "background descendant {pid} survived tool completion"
1949 );
1950 }
1951
1952 #[tokio::test]
1953 async fn cancellation_kills_a_term_ignoring_descendant() {
1954 let workspace = tempfile::tempdir().unwrap();
1955 let root = workspace.path().canonicalize().unwrap();
1956 let tool = BashTool {
1957 timeout: Duration::from_secs(30),
1958 max_timeout: Duration::from_secs(30),
1959 output_limit: 100,
1960 };
1961 let cancellation = tokio_util::sync::CancellationToken::new();
1962 let cancel = cancellation.clone();
1963 let command_root = root.clone();
1964 let execution = tokio::spawn(async move {
1965 tool.execute(
1966 json!({"command":"trap '' TERM; (trap '' TERM; sleep 30) & echo $! > stubborn.pid; wait"}),
1967 ToolContext::new(command_root, cancellation),
1968 )
1969 .await
1970 });
1971 let pid_path = root.join("stubborn.pid");
1972 let pid = wait_for(SHELL_STARTUP, || {
1973 std::fs::read_to_string(&pid_path)
1974 .ok()
1975 .and_then(|value| value.trim().parse::<i32>().ok())
1976 })
1977 .await
1978 .expect("command did not report its descendant pid");
1979 let started = std::time::Instant::now();
1980 cancel.cancel();
1981 let error = execution.await.unwrap().unwrap_err();
1982 assert!(error.to_string().contains("cancelled"));
1983 assert!(started.elapsed() < Duration::from_secs(3));
1984 assert!(
1985 wait_for(Duration::from_secs(5), || is_gone(pid))
1986 .await
1987 .is_some(),
1988 "TERM-ignoring descendant {pid} survived cancellation"
1989 );
1990 }
1991
1992 fn fake_agent(
1996 workspace: &Path,
1997 name: &str,
1998 script: &str,
1999 args: &[&str],
2000 environment: Vec<(OsString, OsString)>,
2001 ) -> NativeAgentTool {
2002 fake_agent_with_prompt_args(workspace, name, script, args, &[], environment)
2003 }
2004
2005 fn fake_agent_with_prompt_args(
2006 workspace: &Path,
2007 name: &str,
2008 script: &str,
2009 args: &[&str],
2010 prompt_args: &[&str],
2011 environment: Vec<(OsString, OsString)>,
2012 ) -> NativeAgentTool {
2013 let script_path = workspace.join("fake-agent.sh");
2014 std::fs::write(&script_path, script).unwrap();
2015 let mut fixed = vec![script_path.display().to_string()];
2016 fixed.extend(args.iter().map(|arg| arg.to_string()));
2017 NativeAgentTool::new(
2018 name.into(),
2019 AgentAdapterConfig {
2020 command: "bash".into(),
2021 args: fixed,
2022 prompt_args: prompt_args.iter().map(|arg| arg.to_string()).collect(),
2023 full_permission_args: None,
2024 model_args: vec!["--model".into(), "{model}".into()],
2025 effort_args: vec!["--effort".into(), "{effort}".into()],
2026 model_hint: adapters::adapter(name.trim_start_matches("agent_"))
2027 .map_or(
2028 "Model ID in the form this agent's CLI accepts.",
2029 |adapter| adapter.model_hint,
2030 )
2031 .into(),
2032 environment,
2033 search_dirs: Vec::new(),
2034 output: OutputFormat::Text,
2035 resume: Resume::Unsupported,
2036 home: None,
2037 transport: Transport::Process,
2038 acp: None,
2039 },
2040 Timeouts {
2041 default: Duration::from_secs(2),
2042 max: Duration::from_secs(5),
2043 },
2044 1024,
2045 None,
2046 test_conversations(),
2047 )
2048 }
2049
2050 fn context(workspace: &Path) -> ToolContext {
2051 ToolContext::new(
2052 workspace.canonicalize().unwrap(),
2053 tokio_util::sync::CancellationToken::new(),
2054 )
2055 }
2056
2057 #[tokio::test]
2058 async fn native_agent_preserves_argument_boundaries() {
2059 let workspace = tempfile::tempdir().unwrap();
2060 let tool = fake_agent(
2061 workspace.path(),
2062 "agent_fake",
2063 "pwd\nprintf '%s\\n' \"$@\"\n",
2064 &["--fixed"],
2065 Vec::new(),
2066 );
2067 let output = tool
2068 .execute(
2069 json!({"prompt":"hello; echo unsafe"}),
2070 context(workspace.path()),
2071 )
2072 .await
2073 .unwrap();
2074 assert!(output.content.contains("--fixed"));
2075 assert!(output.content.contains("hello; echo unsafe"));
2076 assert!(
2077 output
2078 .content
2079 .contains(&workspace.path().display().to_string())
2080 );
2081 }
2082
2083 #[tokio::test]
2084 async fn native_agent_maps_model_and_effort_to_adapter_flags() {
2085 let workspace = tempfile::tempdir().unwrap();
2086 let tool = fake_agent(
2087 workspace.path(),
2088 "agent_claude",
2089 "printf '%s\\n' \"$@\"\n",
2090 &["-p"],
2091 Vec::new(),
2092 );
2093 let properties = &tool.spec().parameters["properties"];
2094 assert_eq!(properties["effort"]["enum"], json!(AGENT_EFFORTS));
2095 assert_eq!(properties["model"]["type"], "string");
2096 let arguments = json!({"prompt":"hi","model":"sonnet","effort":"medium"});
2097 assert!(
2098 tool.approval_summary(&arguments)
2099 .unwrap()
2100 .contains(r#""--model", "sonnet", "--effort", "medium""#)
2101 );
2102 let output = tool
2103 .execute(arguments, context(workspace.path()))
2104 .await
2105 .unwrap();
2106 let output: Value = serde_json::from_str(&output.content).unwrap();
2107 assert_eq!(output["reply"], "-p\n--model\nsonnet\n--effort\nmedium\nhi");
2108 for invalid in [
2109 json!({"prompt":"hi","model":"--dangerously-skip-permissions"}),
2110 json!({"prompt":"hi","model":"sonnet medium"}),
2111 json!({"prompt":"hi","effort":"extreme"}),
2112 json!({"prompt":"hi","model":"@/etc/passwd"}),
2113 json!({"prompt":"--resume"}),
2114 ] {
2115 assert!(tool.risk(&invalid).is_err());
2116 }
2117 let fixed_only = NativeAgentTool::new(
2118 "agent_pi".into(),
2119 AgentAdapterConfig {
2120 command: "pi".into(),
2121 args: vec!["-p".into()],
2122 prompt_args: Vec::new(),
2123 full_permission_args: None,
2124 model_args: Vec::new(),
2125 effort_args: Vec::new(),
2126 model_hint: String::new(),
2127 environment: Vec::new(),
2128 search_dirs: Vec::new(),
2129 output: OutputFormat::Text,
2130 resume: Resume::Unsupported,
2131 home: None,
2132 transport: Transport::Process,
2133 acp: None,
2134 },
2135 Timeouts {
2136 default: Duration::from_secs(2),
2137 max: Duration::from_secs(2),
2138 },
2139 1024,
2140 None,
2141 test_conversations(),
2142 );
2143 assert!(
2144 fixed_only.spec().parameters["properties"]
2145 .get("model")
2146 .is_none()
2147 );
2148 let error = fixed_only
2149 .risk(&json!({"prompt":"hi","model":"sonnet"}))
2150 .unwrap_err();
2151 assert!(
2152 error
2153 .to_string()
2154 .contains("does not support selecting a model")
2155 );
2156 }
2157
2158 #[test]
2159 fn native_agent_model_hints_name_the_adapter_family_and_default() {
2160 let workspace = tempfile::tempdir().unwrap();
2161 let description = |name: &str, field: &str| {
2162 fake_agent(workspace.path(), name, "", &[], Vec::new())
2163 .spec()
2164 .parameters["properties"][field]["description"]
2165 .as_str()
2166 .unwrap()
2167 .to_owned()
2168 };
2169 let claude = description("agent_claude", "model");
2170 let codex = description("agent_codex", "model");
2171 let other = description("agent_other", "model");
2172 assert!(claude.contains("sonnet or opus"));
2173 for text in [&codex, &other] {
2174 assert!(!text.contains("sonnet"), "{text}");
2175 }
2176 assert!(codex.contains("not a Claude alias"));
2177 for text in [claude, codex, other, description("agent_codex", "effort")] {
2178 assert!(
2179 text.contains("omit to use the agent's configured default"),
2180 "{text}"
2181 );
2182 }
2183 }
2184
2185 #[tokio::test]
2186 async fn signed_out_dsh_failure_names_the_host_login_command() {
2187 let workspace = tempfile::tempdir().unwrap();
2188 let tool = fake_agent(
2190 workspace.path(),
2191 "agent_dsh",
2192 "echo 'dsh: MISSING_CREDENTIAL: llm-deepseek: no API key for provider route \"deepseek-official\"' >&2\nexit 1\n",
2193 &[],
2194 Vec::new(),
2195 );
2196 let output = tool
2197 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2198 .await
2199 .unwrap();
2200 assert!(output.is_error);
2201 let content: Value = serde_json::from_str(&output.content).unwrap();
2202 assert!(
2203 content["hint"]
2204 .as_str()
2205 .unwrap()
2206 .ends_with("scv agents login dsh"),
2207 "{content}"
2208 );
2209 }
2210
2211 #[tokio::test]
2212 async fn signed_out_agent_failure_names_the_host_login_command() {
2213 let workspace = tempfile::tempdir().unwrap();
2214 let tool = fake_agent(
2215 workspace.path(),
2216 "agent_claude",
2217 "echo 'Not logged in ยท Please run /login'\nexit 1\n",
2218 &[],
2219 Vec::new(),
2220 );
2221 let output = tool
2222 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2223 .await
2224 .unwrap();
2225 assert!(output.is_error);
2226 let content: Value = serde_json::from_str(&output.content).unwrap();
2227 assert!(
2228 content["hint"]
2229 .as_str()
2230 .unwrap()
2231 .ends_with("scv agents login claude")
2232 );
2233 let other = fake_agent(
2234 workspace.path(),
2235 "agent_claude",
2236 "echo 'disk full'\nexit 1\n",
2237 &[],
2238 Vec::new(),
2239 );
2240 let output = other
2241 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2242 .await
2243 .unwrap();
2244 assert!(output.is_error);
2245 assert!(!output.content.contains("hint"));
2246 }
2247
2248 #[tokio::test]
2249 async fn native_agent_uses_instance_private_environment() {
2250 let workspace = tempfile::tempdir().unwrap();
2251 let home = workspace.path().join("private-home");
2252 let tool = fake_agent(
2253 workspace.path(),
2254 "agent_codex",
2255 "printf 'HOME=%s\\nSCV_HOME=%s\\nCODEX_HOME=%s\\nSCV_CONFIG=%s\\nOPENAI_API_KEY=%s\\nCODEX_API_KEY=%s\\n' \"$HOME\" \"$SCV_HOME\" \"$CODEX_HOME\" \"${SCV_CONFIG-unset}\" \"${OPENAI_API_KEY-unset}\" \"${CODEX_API_KEY-unset}\"\n",
2256 &[],
2257 vec![
2258 ("HOME".into(), home.clone().into()),
2259 ("SCV_HOME".into(), home.clone().into()),
2260 ("CODEX_HOME".into(), home.join("codex").into()),
2261 ],
2262 );
2263 let output = tool
2264 .execute(
2265 json!({"prompt":"print environment"}),
2266 context(workspace.path()),
2267 )
2268 .await
2269 .unwrap();
2270 assert!(output.content.contains(&format!("HOME={}", home.display())));
2271 assert!(
2272 output
2273 .content
2274 .contains(&format!("CODEX_HOME={}/codex", home.display()))
2275 );
2276 assert!(output.content.contains("SCV_CONFIG=unset"));
2277 assert!(output.content.contains("OPENAI_API_KEY=unset"));
2278 assert!(output.content.contains("CODEX_API_KEY=unset"));
2279 }
2280
2281 #[tokio::test]
2282 async fn native_agent_places_prompt_flags_just_before_the_prompt() {
2283 let workspace = tempfile::tempdir().unwrap();
2284 let tool = fake_agent_with_prompt_args(
2285 workspace.path(),
2286 "agent_grok",
2287 "printf '%s\\n' \"$@\"\n",
2288 &[],
2289 &["-p"],
2290 Vec::new(),
2291 );
2292 let arguments = json!({"prompt":"hi","model":"grok-4","effort":"high"});
2293 assert!(
2294 tool.approval_summary(&arguments)
2295 .unwrap()
2296 .contains(r#""--model", "grok-4", "--effort", "high", "-p""#)
2297 );
2298 let output = tool
2299 .execute(arguments, context(workspace.path()))
2300 .await
2301 .unwrap();
2302 let output: Value = serde_json::from_str(&output.content).unwrap();
2303 assert_eq!(output["reply"], "--model\ngrok-4\n--effort\nhigh\n-p\nhi");
2304 }
2305
2306 #[tokio::test]
2307 async fn full_permissions_follow_the_fixed_arguments_and_are_announced() {
2308 let workspace = tempfile::tempdir().unwrap();
2309 let mut tool = fake_agent(
2310 workspace.path(),
2311 "agent_claude",
2312 "printf '%s\\n' \"$@\"\n",
2313 &["-p"],
2314 Vec::new(),
2315 );
2316 let arguments = json!({"prompt":"hi","model":"opus"});
2317 assert!(!tool.approval_summary(&arguments).unwrap().contains("FULL"));
2318 tool.full_permission_args =
2319 Some(vec!["--permission-mode".into(), "bypassPermissions".into()]);
2320 let summary = tool.approval_summary(&arguments).unwrap();
2321 assert!(summary.contains("FULL PERMISSIONS"), "{summary}");
2322 assert!(
2323 summary
2324 .contains(r#""-p", "--permission-mode", "bypassPermissions", "--model", "opus""#)
2325 );
2326 let output = tool
2327 .execute(arguments, context(workspace.path()))
2328 .await
2329 .unwrap();
2330 let output: Value = serde_json::from_str(&output.content).unwrap();
2331 assert_eq!(
2332 output["reply"],
2333 "-p\n--permission-mode\nbypassPermissions\n--model\nopus\nhi"
2334 );
2335 }
2336
2337 #[test]
2338 fn agent_environment_drops_inherited_credentials_but_keeps_its_own_home() {
2339 let mut command = std::process::Command::new("true");
2340 apply_agent_environment_from(
2341 &mut command,
2342 [
2343 "GROK_HOME",
2344 "XAI_API_KEY",
2345 "PI_CODING_AGENT_DIR",
2346 "DEEPSEEK_API_KEY",
2347 "ANTHROPIC_API_KEY",
2348 "OPENROUTER_API_KEY",
2349 "PATH",
2350 ]
2351 .map(OsString::from),
2352 &[("GROK_HOME".into(), "/private/.grok".into())],
2353 );
2354 let envs: HashMap<_, _> = command
2355 .get_envs()
2356 .map(|(key, value)| (key.to_owned(), value.map(ToOwned::to_owned)))
2357 .collect();
2358 assert_eq!(
2359 envs[&OsString::from("GROK_HOME")],
2360 Some(OsString::from("/private/.grok"))
2361 );
2362 for removed in [
2363 "XAI_API_KEY",
2364 "PI_CODING_AGENT_DIR",
2365 "DEEPSEEK_API_KEY",
2366 "ANTHROPIC_API_KEY",
2367 "OPENROUTER_API_KEY",
2368 ] {
2369 assert_eq!(envs[&OsString::from(removed)], None, "{removed}");
2370 }
2371 assert!(!envs.contains_key(&OsString::from("PATH")));
2372 }
2373
2374 #[test]
2375 fn uninstalled_agents_are_not_offered() {
2376 let adapter = |command: &str| AgentAdapterConfig {
2377 command: command.into(),
2378 args: Vec::new(),
2379 prompt_args: Vec::new(),
2380 full_permission_args: None,
2381 model_args: Vec::new(),
2382 effort_args: Vec::new(),
2383 model_hint: String::new(),
2384 environment: Vec::new(),
2385 search_dirs: Vec::new(),
2386 output: OutputFormat::Text,
2387 resume: Resume::Unsupported,
2388 home: None,
2389 transport: Transport::Process,
2390 acp: None,
2391 };
2392 let registry = builtin_registry(
2393 ToolsConfig::default(),
2394 SkillMap::new(),
2395 Vec::new(),
2396 1024,
2397 HashMap::from([
2398 ("agent_present".to_owned(), adapter("bash")),
2399 (
2400 "agent_missing".to_owned(),
2401 adapter("scv-test-agent-that-is-not-installed"),
2402 ),
2403 ]),
2404 )
2405 .unwrap();
2406 assert!(registry.get("agent_present").is_some());
2407 assert!(registry.get("agent_missing").is_none());
2408 }
2409
2410 #[tokio::test]
2411 async fn native_agent_runs_in_a_contained_directory() {
2412 let workspace = tempfile::tempdir().unwrap();
2413 let outside = tempfile::tempdir().unwrap();
2414 let root = workspace.path().canonicalize().unwrap();
2415 std::fs::create_dir(root.join("project")).unwrap();
2416 std::fs::write(root.join("notes.txt"), "not a directory").unwrap();
2417 symlink(outside.path(), root.join("escape")).unwrap();
2418 symlink(root.join("project"), root.join("inner-link")).unwrap();
2419 let tool = fake_agent(&root, "agent_codex", "pwd\n", &[], Vec::new());
2420 let run = |arguments: Value| tool.execute(arguments, context(&root));
2421
2422 for arguments in [
2423 json!({"prompt":"hi"}),
2424 json!({"prompt":"hi","cwd":""}),
2425 json!({"prompt":"hi","cwd":" ","model":"","effort":" "}),
2426 ] {
2427 let output = run(arguments.clone()).await.unwrap();
2428 let output: Value = serde_json::from_str(&output.content).unwrap();
2429 assert_eq!(output["reply"], root.display().to_string(), "{arguments}");
2430 }
2431 for cwd in [
2432 "project".to_owned(),
2433 "project/".to_owned(),
2434 "inner-link".to_owned(),
2435 root.join("project").display().to_string(),
2436 ] {
2437 let output = run(json!({"prompt":"hi","cwd":cwd})).await.unwrap();
2438 let output: Value = serde_json::from_str(&output.content).unwrap();
2439 assert_eq!(
2440 output["reply"],
2441 root.join("project").display().to_string(),
2442 "{cwd}"
2443 );
2444 }
2445 for (cwd, error) in [
2446 ("..", "outside the workspace"),
2447 ("escape", "outside the workspace"),
2448 ("/", "outside the workspace"),
2449 ("notes.txt", "not a directory"),
2450 ("missing", "No such file"),
2451 ] {
2452 let result = run(json!({"prompt":"hi","cwd":cwd})).await;
2453 assert!(
2454 result.as_ref().unwrap_err().to_string().contains(error),
2455 "{cwd}: {result:?}"
2456 );
2457 }
2458 assert!(tool.risk(&json!({"prompt":"hi","cwd":"a\0b"})).is_err());
2459 assert!(
2460 tool.risk(&json!({"prompt":"hi","cwd":"x".repeat(MAX_AGENT_CWD_BYTES + 1)}))
2461 .is_err()
2462 );
2463 let summary = tool
2464 .approval_summary(&json!({"prompt":"hi","cwd":"project","timeout_seconds":4}))
2465 .unwrap();
2466 assert!(summary.contains(r#"in "project" (inside the workspace) for up to 4 seconds"#));
2467 assert!(
2468 tool.approval_summary(&json!({"prompt":"hi"}))
2469 .unwrap()
2470 .contains("in the workspace root for up to 2 seconds")
2471 );
2472 let description = tool.spec().parameters["properties"]["cwd"]["description"]
2473 .as_str()
2474 .unwrap()
2475 .to_owned();
2476 assert!(description.contains("AGENTS.md"));
2477 }
2478
2479 #[tokio::test]
2480 async fn per_call_timeouts_may_rise_to_the_ceiling_but_not_past_it() {
2481 let timeouts = Timeouts {
2482 default: Duration::from_secs(120),
2483 max: Duration::from_secs(1800),
2484 };
2485 assert_eq!(timeouts.resolve(None).unwrap(), Duration::from_secs(120));
2486 assert_eq!(timeouts.resolve(Some(30)).unwrap(), Duration::from_secs(30));
2487 assert_eq!(
2488 timeouts.resolve(Some(1800)).unwrap(),
2489 Duration::from_secs(1800)
2490 );
2491 assert!(timeouts.resolve(Some(0)).is_err());
2492 assert!(
2493 timeouts
2494 .resolve(Some(1801))
2495 .unwrap_err()
2496 .to_string()
2497 .contains("maximum of 1800 seconds (tools.max_timeout_seconds)")
2498 );
2499
2500 let workspace = tempfile::tempdir().unwrap();
2501 let agent = fake_agent(
2502 workspace.path(),
2503 "agent_codex",
2504 "echo ran\n",
2505 &[],
2506 Vec::new(),
2507 );
2508 let schema = &agent.spec().parameters["properties"]["timeout_seconds"];
2509 assert_eq!(schema["maximum"], 5);
2510 assert!(
2511 schema["description"]
2512 .as_str()
2513 .unwrap()
2514 .contains("Defaults to 2; at most 5")
2515 );
2516 assert!(
2517 agent
2518 .risk(&json!({"prompt":"hi","timeout_seconds":5}))
2519 .is_ok()
2520 );
2521 assert!(
2522 agent
2523 .risk(&json!({"prompt":"hi","timeout_seconds":6}))
2524 .is_err()
2525 );
2526 assert!(
2527 agent
2528 .execute(
2529 json!({"prompt":"hi","timeout_seconds":6}),
2530 context(workspace.path())
2531 )
2532 .await
2533 .is_err()
2534 );
2535
2536 let bash = BashTool {
2537 timeout: Duration::from_secs(1),
2538 max_timeout: Duration::from_secs(3),
2539 output_limit: 100,
2540 };
2541 assert_eq!(
2542 bash.spec().parameters["properties"]["timeout_seconds"]["maximum"],
2543 3
2544 );
2545 assert!(
2546 bash.risk(&json!({"command":"true","timeout_seconds":3}))
2547 .is_ok()
2548 );
2549 assert!(
2550 bash.risk(&json!({"command":"true","timeout_seconds":4}))
2551 .unwrap_err()
2552 .to_string()
2553 .contains("tools.max_timeout_seconds")
2554 );
2555 }
2556
2557 fn structured_agent(
2560 workspace: &Path,
2561 name: &str,
2562 format: OutputFormat,
2563 script: &str,
2564 home: Option<PathBuf>,
2565 delegation: Option<DelegationContext>,
2566 timeout: Duration,
2567 ) -> NativeAgentTool {
2568 conversing_agent(
2569 workspace,
2570 name,
2571 format,
2572 Resume::Unsupported,
2573 script,
2574 home,
2575 delegation,
2576 timeout,
2577 test_conversations(),
2578 )
2579 }
2580
2581 #[allow(clippy::too_many_arguments)]
2584 fn conversing_agent(
2585 workspace: &Path,
2586 name: &str,
2587 format: OutputFormat,
2588 resume: Resume,
2589 script: &str,
2590 home: Option<PathBuf>,
2591 delegation: Option<DelegationContext>,
2592 timeout: Duration,
2593 conversations: Arc<ConversationStore>,
2594 ) -> NativeAgentTool {
2595 let script_path = workspace.join(format!("fake-{name}.sh"));
2596 std::fs::write(&script_path, script).unwrap();
2597 NativeAgentTool::new(
2598 name.into(),
2599 AgentAdapterConfig {
2600 command: "bash".into(),
2601 args: vec![script_path.display().to_string()],
2602 prompt_args: Vec::new(),
2603 full_permission_args: None,
2604 model_args: Vec::new(),
2605 effort_args: Vec::new(),
2606 model_hint: String::new(),
2607 environment: Vec::new(),
2608 search_dirs: Vec::new(),
2609 output: format,
2610 resume,
2611 home,
2612 transport: Transport::Process,
2613 acp: None,
2614 },
2615 Timeouts {
2616 default: timeout,
2617 max: Duration::from_secs(30),
2618 },
2619 64 * 1024,
2620 delegation,
2621 conversations,
2622 )
2623 }
2624
2625 fn delegation_context(home: &Path) -> DelegationContext {
2626 DelegationContext {
2627 registry: Arc::new(DelegationRegistry::new(home)),
2628 session: "session-1".into(),
2629 depth: 0,
2630 }
2631 }
2632
2633 #[tokio::test]
2634 async fn claude_stream_json_becomes_a_structured_result() {
2635 let workspace = tempfile::tempdir().unwrap();
2636 let args_file = workspace.path().join("args.txt");
2637 let script = format!(
2638 r#"printf '%s\n' "$@" > {args}
2639printf '%s\n' "$SCV_PARENT" "$SCV_DELEGATION_DEPTH" >> {args}
2640echo '{{"type":"system","subtype":"init","session_id":"x","unknown":[1,2]}}'
2641echo '{{"type":"assistant","message":{{"content":[{{"type":"text","text":"thinking"}}]}}}}'
2642echo 'stray diagnostic' >&2
2643echo '{{"type":"result","subtype":"success","is_error":false,"result":"all done","usage":{{"input_tokens":12,"output_tokens":3}}}}'
2644"#,
2645 args = args_file.display()
2646 );
2647 let home = tempfile::tempdir().unwrap();
2648 let context_home = delegation_context(home.path());
2649 let tool = conversing_agent(
2650 workspace.path(),
2651 "agent_claude",
2652 OutputFormat::ClaudeStreamJson,
2653 adapters::adapter("claude").unwrap().resume,
2654 &script,
2655 None,
2656 Some(context_home.clone()),
2657 Duration::from_secs(10),
2658 test_conversations(),
2659 );
2660 let output = tool
2661 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2662 .await
2663 .unwrap();
2664 assert!(!output.is_error, "{}", output.content);
2665 let value: Value = serde_json::from_str(&output.content).unwrap();
2666 assert_eq!(value["agent"], "claude");
2667 assert_eq!(
2668 (value["session"].as_str(), value["turn"].as_u64()),
2669 (Some("claude-1"), Some(1))
2670 );
2671 assert_eq!(value["status"], "completed");
2672 assert_eq!(value["reply"], "all done");
2673 assert_eq!(value["usage"]["input_tokens"], 12);
2674 assert_eq!(value["exit_code"], 0);
2675 assert_eq!(value["stderr_tail"], "stray diagnostic");
2676 assert_eq!(value["truncated"], false);
2677 assert!(!output.content.contains("thinking"));
2679 let recorded = std::fs::read_to_string(&args_file).unwrap();
2680 let lines: Vec<&str> = recorded.lines().collect();
2681 assert_eq!(
2682 &lines[..4],
2683 [
2684 "--output-format",
2685 "stream-json",
2686 "--verbose",
2687 "--session-id"
2688 ]
2689 );
2690 assert!(uuid::Uuid::parse_str(lines[4]).is_ok());
2691 assert_eq!(lines[5], "hi");
2692 let chain = lines[6];
2693 assert!(chain.contains("/session-1/claude-"), "{chain}");
2694 assert_eq!(lines[7], "1");
2695 assert!(context_home.registry.list(true).is_empty());
2697 }
2698
2699 fn fake_codex(workspace: &Path, slow_start: bool) -> String {
2703 let log = workspace.join("calls.txt");
2704 format!(
2705 r#"printf '%s\n' "$@" '--' >> {log}
2706case " $* " in *" resume "*) ;; *) echo '{{"type":"thread.started","thread_id":"th-1"}}'; {hang} ;; esac
2707for last; do :; done
2708echo "{{\"type\":\"item.completed\",\"item\":{{\"type\":\"agent_message\",\"text\":\"echo: $last\"}}}}"
2709echo '{{"type":"turn.completed","usage":{{"input_tokens":1,"output_tokens":1}}}}'
2710"#,
2711 log = log.display(),
2712 hang = if slow_start { "sleep 30" } else { ":" }
2713 )
2714 }
2715
2716 fn calls(workspace: &Path) -> Vec<Vec<String>> {
2717 std::fs::read_to_string(workspace.join("calls.txt"))
2718 .unwrap()
2719 .split("--\n")
2720 .filter(|call| !call.is_empty())
2721 .map(|call| call.lines().map(str::to_owned).collect())
2722 .collect()
2723 }
2724
2725 #[tokio::test]
2726 async fn conversations_continue_the_cli_session_in_the_same_cwd() {
2727 let workspace = tempfile::tempdir().unwrap();
2728 std::fs::create_dir(workspace.path().join("sub")).unwrap();
2729 let codex_resume = adapters::adapter("codex").unwrap().resume;
2730 let store = test_conversations();
2731 let tool = conversing_agent(
2732 workspace.path(),
2733 "agent_codex",
2734 OutputFormat::CodexJsonl,
2735 codex_resume,
2736 &fake_codex(workspace.path(), false),
2737 None,
2738 None,
2739 Duration::from_secs(10),
2740 Arc::clone(&store),
2741 );
2742 let first = tool
2743 .execute(
2744 json!({"prompt":"remember heron"}),
2745 context(workspace.path()),
2746 )
2747 .await
2748 .unwrap();
2749 let value: Value = serde_json::from_str(&first.content).unwrap();
2750 assert_eq!(value["status"], "completed", "{value}");
2751 assert_eq!(
2752 (value["session"].as_str(), value["turn"].as_u64()),
2753 (Some("codex-1"), Some(1))
2754 );
2755 let second = tool
2756 .execute(
2757 json!({"prompt":"what word?","session":"codex-1"}),
2758 context(workspace.path()),
2759 )
2760 .await
2761 .unwrap();
2762 let value: Value = serde_json::from_str(&second.content).unwrap();
2763 assert_eq!(value["reply"], "echo: what word?");
2764 assert_eq!(
2765 (value["session"].as_str(), value["turn"].as_u64()),
2766 (Some("codex-1"), Some(2))
2767 );
2768 let recorded = calls(workspace.path());
2772 assert_eq!(recorded[0], ["--json", "remember heron"]);
2773 assert_eq!(recorded[1], ["resume", "--json", "th-1", "what word?"]);
2774
2775 let moved = tool
2777 .execute(
2778 json!({"prompt":"x","session":"codex-1","cwd":"sub"}),
2779 context(workspace.path()),
2780 )
2781 .await
2782 .unwrap_err();
2783 assert!(moved.0.contains("runs in"), "{}", moved.0);
2784 let other_session = conversing_agent(
2786 workspace.path(),
2787 "agent_codex",
2788 OutputFormat::CodexJsonl,
2789 codex_resume,
2790 &fake_codex(workspace.path(), false),
2791 None,
2792 None,
2793 Duration::from_secs(10),
2794 test_conversations(),
2795 );
2796 let unknown = other_session
2797 .execute(
2798 json!({"prompt":"x","session":"codex-1"}),
2799 context(workspace.path()),
2800 )
2801 .await
2802 .unwrap_err();
2803 assert!(
2804 unknown.0.contains("unknown in this session"),
2805 "{}",
2806 unknown.0
2807 );
2808 assert_eq!(
2809 calls(workspace.path()).len(),
2810 2,
2811 "rejected turns never launch the CLI"
2812 );
2813 let vendor = json!({"prompt":"x","session":"01a0cd5a-7195-7b31-a503-e235d5da7b45"});
2815 assert!(
2816 tool.risk(&vendor)
2817 .unwrap_err()
2818 .0
2819 .contains("not a conversation handle")
2820 );
2821 assert!(
2822 tool.spec().parameters["properties"]
2823 .get("session")
2824 .is_some()
2825 );
2826 }
2827
2828 #[tokio::test]
2829 async fn a_timed_out_turn_stays_resumable_and_unsupported_agents_refuse_sessions() {
2830 let workspace = tempfile::tempdir().unwrap();
2831 let tool = conversing_agent(
2832 workspace.path(),
2833 "agent_codex",
2834 OutputFormat::CodexJsonl,
2835 adapters::adapter("codex").unwrap().resume,
2836 &fake_codex(workspace.path(), true),
2837 None,
2838 None,
2839 Duration::from_secs(1),
2840 test_conversations(),
2841 );
2842 let first = tool
2843 .execute(json!({"prompt":"start"}), context(workspace.path()))
2844 .await
2845 .unwrap();
2846 let value: Value = serde_json::from_str(&first.content).unwrap();
2847 assert_eq!(value["status"], "timeout", "{value}");
2848 assert_eq!(value["session"], "codex-1");
2849 let resumed = tool
2850 .execute(
2851 json!({"prompt":"continue where you left off","session":"codex-1"}),
2852 context(workspace.path()),
2853 )
2854 .await
2855 .unwrap();
2856 let value: Value = serde_json::from_str(&resumed.content).unwrap();
2857 assert_eq!(value["status"], "completed", "{value}");
2858 assert_eq!(value["turn"], 2);
2859
2860 let plain = structured_agent(
2861 workspace.path(),
2862 "agent_grok",
2863 OutputFormat::Text,
2864 "echo hi\n",
2865 None,
2866 None,
2867 Duration::from_secs(5),
2868 );
2869 let refused = plain
2870 .risk(&json!({"prompt":"x","session":"grok-1"}))
2871 .unwrap_err();
2872 assert!(
2873 refused.0.contains("cannot continue a conversation"),
2874 "{}",
2875 refused.0
2876 );
2877 assert!(
2878 plain.spec().parameters["properties"]
2879 .get("session")
2880 .is_none()
2881 );
2882 let output = plain
2883 .execute(json!({"prompt":"x"}), context(workspace.path()))
2884 .await
2885 .unwrap();
2886 assert!(
2887 !output.content.contains("\"session\""),
2888 "{}",
2889 output.content
2890 );
2891 }
2892
2893 #[tokio::test]
2894 async fn codex_json_reads_the_last_message_file_and_removes_it() {
2895 let workspace = tempfile::tempdir().unwrap();
2896 let home = tempfile::tempdir().unwrap();
2897 let script = r#"while [ "$#" -gt 0 ]; do
2898 if [ "$1" = "-o" ]; then printf 'final from file\n' > "$2"; echo "$2" > last-path.txt; fi
2899 shift
2900done
2901echo '{"type":"thread.started","thread_id":"t"}'
2902echo '{"type":"turn.completed","usage":{"input_tokens":5,"output_tokens":1}}'
2903"#;
2904 let tool = structured_agent(
2905 workspace.path(),
2906 "agent_codex",
2907 OutputFormat::CodexJsonl,
2908 script,
2909 Some(home.path().to_owned()),
2910 None,
2911 Duration::from_secs(10),
2912 );
2913 let output = tool
2914 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2915 .await
2916 .unwrap();
2917 let value: Value = serde_json::from_str(&output.content).unwrap();
2918 assert_eq!(value["status"], "completed", "{value}");
2919 assert_eq!(value["reply"], "final from file");
2920 let path = std::fs::read_to_string(workspace.path().join("last-path.txt")).unwrap();
2921 let path = PathBuf::from(path.trim());
2922 assert!(path.starts_with(home.path().join("tmp")));
2923 assert!(!path.exists(), "the last-message file is removed");
2924 use std::os::unix::fs::PermissionsExt as _;
2925 let mode = std::fs::metadata(home.path().join("tmp"))
2926 .unwrap()
2927 .permissions()
2928 .mode();
2929 assert_eq!(mode & 0o777, 0o700);
2930 }
2931
2932 #[tokio::test]
2933 async fn pi_json_and_signed_out_claude_results() {
2934 let workspace = tempfile::tempdir().unwrap();
2935 let pi = structured_agent(
2936 workspace.path(),
2937 "agent_pi",
2938 OutputFormat::PiJson,
2939 r#"echo '{"type":"session","id":"p"}'
2940echo '{"type":"message_end","message":{"role":"assistant","content":[{"type":"text","text":"pi ok"}],"usage":{"input":7,"output":2}}}'
2941"#,
2942 None,
2943 None,
2944 Duration::from_secs(10),
2945 );
2946 let output = pi
2947 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2948 .await
2949 .unwrap();
2950 let value: Value = serde_json::from_str(&output.content).unwrap();
2951 assert_eq!(value["reply"], "pi ok");
2952 assert_eq!(value["usage"]["output_tokens"], 2);
2953
2954 let claude = structured_agent(
2955 workspace.path(),
2956 "agent_claude",
2957 OutputFormat::ClaudeStreamJson,
2958 r#"echo '{"type":"result","subtype":"success","is_error":true,"result":"Not logged in ยท Please run /login"}'
2959exit 1
2960"#,
2961 None,
2962 None,
2963 Duration::from_secs(10),
2964 );
2965 let output = claude
2966 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2967 .await
2968 .unwrap();
2969 assert!(output.is_error);
2970 let value: Value = serde_json::from_str(&output.content).unwrap();
2971 assert_eq!(value["status"], "failed");
2972 assert_eq!(value["exit_code"], 1);
2973 assert!(
2974 value["hint"]
2975 .as_str()
2976 .unwrap()
2977 .contains("scv agents login claude")
2978 );
2979 }
2980
2981 #[cfg(target_os = "linux")]
2982 #[tokio::test]
2983 async fn a_timed_out_run_and_its_detached_descendants_are_stopped() {
2984 let workspace = tempfile::tempdir().unwrap();
2985 let home = tempfile::tempdir().unwrap();
2986 let delegation = delegation_context(home.path());
2987 let tool = structured_agent(
2988 workspace.path(),
2989 "agent_codex",
2990 OutputFormat::CodexJsonl,
2991 "setsid sleep 60 &\necho \"$SCV_PARENT\" > chain.txt\nexec sleep 60\n",
2993 None,
2994 Some(delegation.clone()),
2995 Duration::from_secs(1),
2996 );
2997 let output = tool
2998 .execute(json!({"prompt":"hi"}), context(workspace.path()))
2999 .await
3000 .unwrap();
3001 let value: Value = serde_json::from_str(&output.content).unwrap();
3002 assert_eq!(value["status"], "timeout");
3003 let chain = std::fs::read_to_string(workspace.path().join("chain.txt")).unwrap();
3004 let handle = chain.trim().rsplit('/').next().unwrap().to_owned();
3005 let tagged = || {
3006 std::fs::read_dir("/proc")
3007 .unwrap()
3008 .filter_map(Result::ok)
3009 .filter(|entry| {
3010 std::fs::read(entry.path().join("environ")).is_ok_and(|environ| {
3011 environ
3012 .split(|byte| *byte == 0)
3013 .any(|entry| entry == format!("SCV_PARENT={}", chain.trim()).as_bytes())
3014 })
3015 })
3016 .count()
3017 };
3018 let mut remaining = tagged();
3019 for _ in 0..100 {
3020 if remaining == 0 {
3021 break;
3022 }
3023 tokio::time::sleep(Duration::from_millis(50)).await;
3024 remaining = tagged();
3025 }
3026 assert_eq!(remaining, 0, "tagged processes of {handle} survived");
3027 assert!(delegation.registry.list(true).is_empty());
3028 }
3029
3030 #[test]
3031 fn agents_are_not_offered_at_the_delegation_depth_limit() {
3032 let adapter = AgentAdapterConfig {
3033 command: "bash".into(),
3034 args: Vec::new(),
3035 prompt_args: Vec::new(),
3036 full_permission_args: None,
3037 model_args: Vec::new(),
3038 effort_args: Vec::new(),
3039 model_hint: String::new(),
3040 environment: Vec::new(),
3041 search_dirs: Vec::new(),
3042 output: OutputFormat::Text,
3043 resume: Resume::Unsupported,
3044 home: None,
3045 transport: Transport::Process,
3046 acp: None,
3047 };
3048 let home = tempfile::tempdir().unwrap();
3049 for (max_depth, offered) in [(0, false), (1, true)] {
3050 let registry = builtin_registry(
3051 ToolsConfig {
3052 max_delegation_depth: max_depth,
3053 delegation: Some(delegation_context(home.path())),
3054 ..ToolsConfig::default()
3055 },
3056 SkillMap::new(),
3057 Vec::new(),
3058 1024,
3059 HashMap::from([("agent_claude".to_owned(), adapter.clone())]),
3060 )
3061 .unwrap();
3062 assert_eq!(registry.get("agent_claude").is_some(), offered);
3063 assert!(registry.get("bash").is_some());
3064 }
3065 for (declared, max_depth, offered) in [(1, 1, false), (1, 2, true), (5, 2, false)] {
3068 let registry = builtin_registry(
3069 ToolsConfig {
3070 max_delegation_depth: max_depth,
3071 delegation: Some(DelegationContext {
3072 depth: declared,
3073 ..delegation_context(home.path())
3074 }),
3075 ..ToolsConfig::default()
3076 },
3077 SkillMap::new(),
3078 Vec::new(),
3079 1024,
3080 HashMap::from([("agent_claude".to_owned(), adapter.clone())]),
3081 )
3082 .unwrap();
3083 assert_eq!(
3084 registry.get("agent_claude").is_some(),
3085 offered,
3086 "declared {declared}, limit {max_depth}"
3087 );
3088 }
3089 }
3090}