Skip to main content

scv_tools/
lib.rs

1//! SCV's bounded, workspace-aware built-in tools.
2
3mod 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    /// Default `bash` timeout when a call does not choose one.
54    pub command_timeout: Duration,
55    /// Default native-agent timeout when a call does not choose one.
56    pub agent_timeout: Duration,
57    /// The longest timeout any single call may request.
58    pub max_timeout: Duration,
59    pub output_limit_bytes: usize,
60    pub max_read_bytes: usize,
61    pub max_write_bytes: usize,
62    /// Agent tools are offered only below this delegation depth.
63    pub max_delegation_depth: u32,
64    /// How many delegated conversations a session remembers, and for how long.
65    pub conversations: ConversationLimits,
66    /// Records delegated runs for listing and cleanup; `None` runs them untracked.
67    pub delegation: Option<DelegationContext>,
68    /// Background jobs an agent call may start at once (`background: true`);
69    /// 0 turns background calls and `agent_wait` / `agent_status` off.
70    pub max_background: usize,
71    /// The session's background job store, when the server reports finished
72    /// jobs; otherwise the registry makes its own.
73    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/// The registry and parent session that delegated runs are recorded under.
98#[derive(Debug, Clone)]
99pub struct DelegationContext {
100    pub registry: Arc<DelegationRegistry>,
101    pub session: String,
102    /// Delegation depth the session's client declared (0 for a direct
103    /// client). Runs count from the larger of this and the process's own.
104    pub depth: u32,
105}
106
107impl DelegationContext {
108    /// The depth delegated runs of this session start from.
109    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    /// Arguments placed immediately before the prompt, for CLIs that take the
119    /// prompt as a flag value.
120    pub prompt_args: Vec<String>,
121    /// The CLI's own full-autonomy arguments, placed after `args`, when the
122    /// user configured `permissions = "full"`; the approval summary says so.
123    pub full_permission_args: Option<Vec<String>>,
124    /// Arguments appended for a per-call model; `{model}` is substituted.
125    /// Empty means the adapter does not offer model selection.
126    pub model_args: Vec<String>,
127    /// Arguments appended for a per-call effort; `{effort}` is substituted.
128    /// Empty means the adapter does not offer effort selection.
129    pub effort_args: Vec<String>,
130    /// Describes the `model` argument for the calling model.
131    pub model_hint: String,
132    /// Environment for the nested process. SCV supplies an instance-private home.
133    pub environment: Vec<(OsString, OsString)>,
134    /// Per-user install directories searched when `command` is not on `PATH`.
135    pub search_dirs: Vec<PathBuf>,
136    /// What the CLI prints, and so how its reply is read.
137    pub output: OutputFormat,
138    /// How a conversation with the CLI is continued, if it can be.
139    pub resume: Resume,
140    /// SCV's private home for this agent, for files SCV hands the CLI.
141    pub home: Option<PathBuf>,
142    /// How SCV talks to the agent.
143    pub transport: Transport,
144    /// The agent's ACP server, when `[agents.<name>] transport` allows it and
145    /// the adapter table has one.
146    pub acp: Option<AcpAgentLaunch>,
147}
148
149/// An agent's Agent Client Protocol server, resolved from its adapter-table
150/// entry and `[agents.<name>] transport`.
151#[derive(Debug, Clone)]
152pub struct AcpAgentLaunch {
153    pub command: String,
154    /// Arguments with the `permissions = "full"` switches already applied.
155    pub args: Vec<String>,
156    /// The ACP session mode that grants full permissions, selected in every
157    /// new session when `permissions = "full"`.
158    pub full_mode: Option<String>,
159    /// Extra environment for the ACP server, such as permission settings the
160    /// server reads only from its environment.
161    pub environment: Vec<(OsString, OsString)>,
162    /// `transport = "acp"`: never fall back to one CLI process per turn, so
163    /// the agent is not offered while its ACP server is missing.
164    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    // A delegated SCV at the depth limit may not delegate further.
194    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    // One job store per session, shared by its agent tools and dropped with
204    // it, which cancels the jobs still running.
205    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    // One store per session, shared by its agent tools and dropped with it.
222    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            // An agent that is not installed is not offered to the model.
234            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                // `transport = "acp"` without its server: not offered.
279                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        // An agent that is not installed is not offered to the model.
294        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/// A process tool's default timeout and the ceiling a call may raise it to.
662#[derive(Debug, Clone, Copy)]
663pub(crate) struct Timeouts {
664    pub(crate) default: Duration,
665    pub(crate) max: Duration,
666}
667
668impl Timeouts {
669    /// The call's timeout: its own request up to the ceiling, else the
670    /// default. A request above the ceiling is refused, never clamped, so the
671    /// caller learns the limit instead of being cut off early.
672    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
786/// Effort levels accepted by the built-in adapters' CLIs.
787const AGENT_EFFORTS: [&str; 5] = ["low", "medium", "high", "xhigh", "max"];
788
789impl NativeAgentTool {
790    /// The fixed arguments, validated model and effort selections, and the
791    /// prompt arguments; the prompt is appended separately as the final argument.
792    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        // The prompt follows the flags as a positional argument, so it must
815        // not be readable as one.
816        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    /// Arguments that name per-run files or IDs, placed after the fixed and
879    /// format arguments. Returns the Codex last-message file to read and
880    /// remove afterwards.
881    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
900/// `template` with `{session}` replaced by `session`.
901fn 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
914/// Read and delete a file the CLI wrote for SCV, bounded to `limit` bytes.
915fn 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
941/// Models often send an optional string they mean to leave unset as `""`, so
942/// a blank value selects the default rather than failing the call.
943fn 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
950/// Longest `cwd` argument accepted, in bytes.
951const 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
962/// Resolve a requested agent directory against the workspace. Resolution
963/// follows symlinks, so a link pointing outside the workspace is refused
964/// rather than trusted by name.
965pub(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
983/// Model names are passed as one argument, so only reject values that could
984/// read as a flag, name an `@file` argument, or carry unexpected characters.
985pub(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
1242/// Point a failed agent run that reads like a missing sign-in at the host
1243/// command that fixes it, since the agent's own advice (`/login`) cannot be
1244/// followed from a remote chat.
1245pub(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
1277/// Give a native agent command its adapter environment: remove every
1278/// inherited credential, endpoint, and state-location variable any adapter
1279/// declares, then set `environment` (such as the relocated config home).
1280pub 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
1357/// A delegated run's stdout reader, stderr tail, and how it ended.
1358struct AgentRun {
1359    stream: AgentStream,
1360    exit: RunExit,
1361    exit_code: Option<i32>,
1362    stderr_tail: String,
1363}
1364
1365/// Run a native agent: stdout is parsed as it arrives rather than buffered,
1366/// stderr keeps only its tail, and the run is recorded in the delegation
1367/// registry while it lasts.
1368async 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    // Bookkeeping must not fail the delegation: an unrecorded run is still
1378    // tagged, so a later sweep can find what it leaves behind.
1379    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
1463/// Wait for a spawned process group until it exits, times out, or is
1464/// cancelled, always finishing the whole group and draining its output.
1465async 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    // Always finish the process group, even if its original leader already exited.
1542    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    // Negative PID addresses the process group created at spawn.
1554    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
1572/// Where a child's output goes as it is read.
1573pub(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    /// `bash -l` sources the host's login profile before it runs a command,
1720    /// and CI images can spend seconds there under parallel test load. Waits
1721    /// that include shell startup use this ceiling; they end as soon as their
1722    /// condition holds.
1723    const SHELL_STARTUP: Duration = Duration::from_secs(30);
1724
1725    /// Polls `probe` every 10 ms until it yields a value or `limit` passes.
1726    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        // Time from the shell's last write, which excludes its startup.
1934        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    /// Fake agents run through `bash` so no test ever executes a file that a
1993    /// concurrently forked test process may still hold open for writing
1994    /// (which fails spawning with ETXTBSY).
1995    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        // DeepSeek Harness 0.1.7-rc.1's startup error without a key.
2189        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    /// A fake agent CLI in `format`, run through `bash script`, optionally
2558    /// recorded in `delegation`.
2559    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    /// Like [`structured_agent`], continuing conversations as `resume` says,
2582    /// in `conversations` (shared by one session's tools).
2583    #[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        // No event log reaches the parent.
2678        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        // The run's record is gone once it ends.
2696        assert!(context_home.registry.list(true).is_empty());
2697    }
2698
2699    /// A fake Codex that records each call's arguments, reports thread
2700    /// `th-1`, and answers with the prompt it was given. With `slow_start`,
2701    /// a first (non-resume) turn hangs after reporting its thread.
2702    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        // The script path is the only fixed argument, so `$@` starts after it:
2769        // `resume` comes right after the fixed arguments, and the CLI's thread
2770        // ID sits just before the prompt.
2771        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        // A conversation stays in its cwd.
2776        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        // Another session's tools do not know this session's handles.
2785        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        // The CLI's own ID is never accepted in place of a handle.
2814        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            // The detached sleep leaves the agent's process group and session.
2992            "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        // A client that is itself delegated (`session.start.delegation_depth`)
3066        // counts too, even though this process is not delegated.
3067        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}