Stream shell tool output

This commit is contained in:
2026-06-22 12:54:54 -05:00
parent 2ae4eda0a1
commit 169281d942
10 changed files with 330 additions and 19 deletions
Generated
+1
View File
@@ -1847,6 +1847,7 @@ dependencies = [
"libc",
"mio",
"pin-project-lite",
"signal-hook-registry",
"socket2",
"tokio-macros",
"windows-sys 0.61.2",
+1 -1
View File
@@ -35,7 +35,7 @@ reqwest = { version = "0.12", default-features = false, features = ["json", "rus
serde = { version = "1", features = ["derive"] }
serde_json = "1"
thiserror = "1"
tokio = { version = "1", features = ["macros", "rt-multi-thread", "sync", "time"] }
tokio = { version = "1", features = ["macros", "rt-multi-thread", "sync", "time", "process", "io-util"] }
tui-textarea = "0.6"
unicode-width = "0.1"
+2 -2
View File
@@ -98,6 +98,6 @@ Reasoning is hidden by default unless `show_reasoning` is enabled; press `Ctrl-S
Read-only mode allows `ls`, `read`, and `grep` within the launch cwd/`--cwd` and the bundled docs directory at `~/.cass/docs`.
Full-access mode additionally allows `write` and `edit`. Mutating tools use atomic writes where practical: Cass writes to a temporary file first, then renames it into place after validation/write success. `write` and `edit` are always blocked under `~/.cass/docs`.
Full-access mode additionally allows `write`, `edit`, and `shell`. Mutating tools use atomic writes where practical: Cass writes to a temporary file first, then renames it into place after validation/write success. `write` and `edit` are always blocked under `~/.cass/docs`. The `shell` tool runs commands via `sh -c` in the launch working directory with a configurable timeout (default 30 seconds) and streams stdout/stderr into the transcript while the command is running.
There is no shell/bash tool in v1.
The `shell` tool is available in full-access mode and runs shell commands in the launch working directory.
+55 -9
View File
@@ -4,7 +4,7 @@ use crate::conversation::{now_ts, Conversation, Record, StoredToolCall};
use crate::prompt;
use crate::providers::openai_compatible::{OpenAiCompatibleProvider, OpenAiCompatibleSettings};
use crate::providers::types::ModelMessage;
use crate::tools::{self, ToolContext};
use crate::tools::{self, ToolContext, ToolRuntimeEvent};
use anyhow::Result;
use serde_json::Value;
use std::path::PathBuf;
@@ -19,6 +19,12 @@ pub enum AgentEvent {
name: String,
arguments: Value,
},
ToolOutputChunk {
id: String,
name: String,
stream: String,
content: String,
},
ToolResult {
id: String,
name: String,
@@ -86,6 +92,7 @@ pub async fn run_turn(
read_roots: vec![settings.cwd.clone(), docs_dir.clone()],
blocked_write_roots: vec![docs_dir.clone()],
model_result_limit: settings.config.model_tool_result_limit,
runtime_tx: None,
};
let mut retrying_empty_final = false;
@@ -154,21 +161,42 @@ pub async fn run_turn(
break;
}
for call in tool_calls {
let call_id = call.id.clone();
let call_name = call.name.clone();
let call_arguments = call.arguments.clone();
let _ = tx.send(AgentEvent::ToolCallStarted {
id: call.id.clone(),
name: call.name.clone(),
arguments: call.arguments.clone(),
id: call_id.clone(),
name: call_name.clone(),
arguments: call_arguments.clone(),
});
let output = tools::execute(&call.name, call.arguments.clone(), &tool_ctx).await;
let (runtime_tx, mut runtime_rx) = mpsc::unbounded_channel::<ToolRuntimeEvent>();
let mut call_tool_ctx = tool_ctx.clone();
call_tool_ctx.runtime_tx = Some(runtime_tx);
let output = {
let execute = tools::execute(&call_name, call_arguments, &call_tool_ctx);
tokio::pin!(execute);
let output = loop {
tokio::select! {
output = &mut execute => break output,
Some(event) = runtime_rx.recv() => {
forward_tool_runtime_event(&tx, &call_id, &call_name, event);
}
}
};
while let Ok(event) = runtime_rx.try_recv() {
forward_tool_runtime_event(&tx, &call_id, &call_name, event);
}
output
};
let _ = tx.send(AgentEvent::ToolResult {
id: call.id.clone(),
name: call.name.clone(),
id: call_id.clone(),
name: call_name.clone(),
ok: output.ok,
content: output.content.clone(),
});
conversation.append(Record::Tool {
tool_call_id: call.id,
name: call.name,
tool_call_id: call_id,
name: call_name,
ok: output.ok,
content: output.content,
ts: now_ts(),
@@ -179,6 +207,24 @@ pub async fn run_turn(
Ok(conversation)
}
fn forward_tool_runtime_event(
tx: &mpsc::UnboundedSender<AgentEvent>,
id: &str,
name: &str,
event: ToolRuntimeEvent,
) {
match event {
ToolRuntimeEvent::OutputChunk { stream, content } => {
let _ = tx.send(AgentEvent::ToolOutputChunk {
id: id.to_string(),
name: name.to_string(),
stream,
content,
});
}
}
}
fn append_visible_assistant(
conversation: &mut Conversation,
tx: &mpsc::UnboundedSender<AgentEvent>,
+72
View File
@@ -9,6 +9,7 @@ use crate::ui::render::{self, TranscriptBlock, TranscriptKind};
use crate::ui::terminal;
use anyhow::{Context, Result};
use crossterm::event::{Event, KeyCode, KeyEventKind, KeyModifiers, MouseEventKind};
use std::collections::HashMap;
use std::fs;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
@@ -102,6 +103,7 @@ async fn run_tui(
let mut handle: Option<JoinHandle<Result<Conversation>>> = None;
let mut active_assistant: Option<usize> = None;
let mut active_reasoning: Option<usize> = None;
let mut active_tools: HashMap<String, usize> = HashMap::new();
let mut stick_to_bottom = true;
let mut chat_id = conversation.id.clone();
let mut autofill_selected = 0usize;
@@ -115,6 +117,7 @@ async fn run_tui(
transcript: &mut transcript,
active_assistant: &mut active_assistant,
active_reasoning: &mut active_reasoning,
active_tools: &mut active_tools,
status: &mut status,
stick_to_bottom,
show_full_tools,
@@ -134,6 +137,7 @@ async fn run_tui(
transcript: &mut transcript,
active_assistant: &mut active_assistant,
active_reasoning: &mut active_reasoning,
active_tools: &mut active_tools,
status: &mut status,
stick_to_bottom,
show_full_tools,
@@ -159,6 +163,7 @@ async fn run_tui(
}
active_assistant = None;
active_reasoning = None;
active_tools.clear();
status = "idle".into();
}
@@ -435,6 +440,7 @@ async fn run_tui(
autofill_selected = 0;
active_assistant = None;
active_reasoning = None;
active_tools.clear();
status = format!("new chat {chat_id}");
stick_to_bottom = true;
scroll = bottom_scroll(
@@ -484,6 +490,7 @@ async fn run_tui(
autofill_selected = 0;
active_assistant = None;
active_reasoning = None;
active_tools.clear();
status = format!("resumed chat {chat_id}");
stick_to_bottom = true;
scroll = bottom_scroll(
@@ -531,6 +538,7 @@ async fn run_tui(
});
active_assistant = None;
active_reasoning = None;
active_tools.clear();
status = "running".into();
let settings = AgentSettings {
config: config.clone(),
@@ -642,6 +650,7 @@ struct AgentEventContext<'a> {
transcript: &'a mut Vec<TranscriptBlock>,
active_assistant: &'a mut Option<usize>,
active_reasoning: &'a mut Option<usize>,
active_tools: &'a mut HashMap<String, usize>,
status: &'a mut String,
stick_to_bottom: bool,
show_full_tools: bool,
@@ -714,6 +723,19 @@ fn apply_agent_event(event: AgentEvent, ctx: &mut AgentEventContext<'_>) -> Resu
content: serde_json::to_string_pretty(&arguments)
.unwrap_or_else(|_| arguments.to_string()),
});
ctx.active_tools.insert(id, ctx.transcript.len() - 1);
update_bottom_scroll(ctx)?;
}
AgentEvent::ToolOutputChunk {
id,
name,
stream,
content,
} => {
*ctx.active_assistant = None;
*ctx.active_reasoning = None;
let idx = active_tool_block(ctx, &id, &name);
append_tool_output_chunk(&mut ctx.transcript[idx].content, &stream, &content);
update_bottom_scroll(ctx)?;
}
AgentEvent::ToolResult {
@@ -722,6 +744,24 @@ fn apply_agent_event(event: AgentEvent, ctx: &mut AgentEventContext<'_>) -> Resu
ok,
content,
} => {
if name == "shell" {
if let Some(idx) = ctx.active_tools.remove(&id) {
ctx.transcript[idx].kind = if ok {
TranscriptKind::Tool
} else {
TranscriptKind::Error
};
ctx.transcript[idx].title = format!(
"{name} {} ({})",
if ok { "✓" } else { "✗" },
short_call_id(&id)
);
ctx.transcript[idx].content = content;
update_bottom_scroll(ctx)?;
return Ok(());
}
}
ctx.active_tools.remove(&id);
ctx.transcript.push(TranscriptBlock {
kind: if ok {
TranscriptKind::Tool
@@ -748,12 +788,44 @@ fn apply_agent_event(event: AgentEvent, ctx: &mut AgentEventContext<'_>) -> Resu
AgentEvent::TurnFinished => {
*ctx.active_assistant = None;
*ctx.active_reasoning = None;
ctx.active_tools.clear();
*ctx.status = "turn finished".into();
}
}
Ok(())
}
fn active_tool_block(ctx: &mut AgentEventContext<'_>, id: &str, name: &str) -> usize {
if let Some(idx) = ctx.active_tools.get(id).copied() {
return idx;
}
ctx.transcript.push(TranscriptBlock {
kind: TranscriptKind::Tool,
title: format!("{name} … ({})", short_call_id(id)),
content: String::new(),
});
let idx = ctx.transcript.len() - 1;
ctx.active_tools.insert(id.to_string(), idx);
idx
}
fn append_tool_output_chunk(existing: &mut String, stream: &str, chunk: &str) {
if !existing.contains("streamed output:\n") {
if !existing.trim().is_empty() {
existing.push_str("\n\n");
}
existing.push_str("streamed output:\n");
}
if !existing.ends_with('\n') {
existing.push('\n');
}
existing.push_str(&format!("[{stream}] "));
existing.push_str(chunk);
if !chunk.ends_with('\n') {
existing.push('\n');
}
}
fn update_bottom_scroll(ctx: &mut AgentEventContext<'_>) -> Result<()> {
if ctx.stick_to_bottom {
*ctx.scroll = bottom_scroll(
+1 -1
View File
@@ -44,7 +44,7 @@ pub fn build_effective_system_prompt(
));
match mode {
AccessMode::ReadOnly => prompt.push_str("In read-only mode, you may inspect files with ls, read, and grep only inside the launch working directory or bundled Cass docs directory. Do not request write or edit. If a task requires modification, explain what needs full-access mode.\n\n"),
AccessMode::FullAccess => prompt.push_str("In full-access mode, you may request ls, read, grep, write, and edit when needed. Cass does not restrict read paths to the launch directory, but normal operating-system permissions still apply. write and edit are still blocked under the bundled Cass docs directory.\n\n"),
AccessMode::FullAccess => prompt.push_str("In full-access mode, you may request ls, read, grep, write, edit, and shell when needed. The shell tool runs commands in the launch working directory. Cass does not restrict read paths to the launch directory, but normal operating-system permissions still apply. write and edit are still blocked under the bundled Cass docs directory.\n\n"),
}
prompt.push_str("6. Response behavior\n\nAssistant output is streamed to the user. Keep user-facing text direct and useful. Tool calls and results are visible to the user, so avoid claiming work happened until the relevant tool result confirms it. After using tools or completing requested work, always end the turn with a concise final user-facing response. Do not finish a turn with only tool calls.\n");
prompt
+28
View File
@@ -4,6 +4,7 @@ pub mod ls;
pub mod path;
pub mod read;
pub mod schema;
pub mod shell;
pub mod write;
use crate::access::AccessMode;
@@ -11,6 +12,7 @@ use anyhow::Result;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::path::PathBuf;
use tokio::sync::mpsc;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ToolSpec {
@@ -26,6 +28,12 @@ pub struct ToolContext {
pub read_roots: Vec<PathBuf>,
pub blocked_write_roots: Vec<PathBuf>,
pub model_result_limit: usize,
pub runtime_tx: Option<mpsc::UnboundedSender<ToolRuntimeEvent>>,
}
#[derive(Debug, Clone)]
pub enum ToolRuntimeEvent {
OutputChunk { stream: String, content: String },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -43,6 +51,7 @@ pub fn available_tool_names(mode: AccessMode) -> Vec<String> {
"grep".into(),
"write".into(),
"edit".into(),
"shell".into(),
],
}
}
@@ -52,11 +61,26 @@ pub fn specs(mode: AccessMode) -> Vec<ToolSpec> {
if mode.can_write() {
specs.push(write::spec());
specs.push(edit::spec());
specs.push(shell::spec());
}
specs
}
pub async fn execute(name: &str, args: Value, ctx: &ToolContext) -> ToolOutput {
// Async tools
if name == "shell" {
let result = if ctx.mode.can_write() {
shell::run(args, ctx).await
} else {
Err(anyhow::anyhow!(
"tool `{name}` is unavailable in {} mode",
ctx.mode
))
};
return result_to_output(result, ctx);
}
// Sync tools
let result: Result<String> = match name {
"ls" => ls::run(args, ctx),
"read" => read::run(args, ctx),
@@ -69,6 +93,10 @@ pub async fn execute(name: &str, args: Value, ctx: &ToolContext) -> ToolOutput {
)),
_ => Err(anyhow::anyhow!("unknown tool `{name}`")),
};
result_to_output(result, ctx)
}
fn result_to_output(result: Result<String>, ctx: &ToolContext) -> ToolOutput {
match result {
Ok(content) => ToolOutput {
ok: true,
+133
View File
@@ -0,0 +1,133 @@
use super::{schema, ToolContext, ToolRuntimeEvent, ToolSpec};
use anyhow::{bail, Context, Result};
use serde::Deserialize;
use serde_json::{json, Value};
use std::process::Stdio;
use std::time::Duration;
use tokio::io::{AsyncRead, AsyncReadExt};
#[derive(Debug, Deserialize)]
struct Args {
command: String,
#[serde(default)]
timeout: Option<u64>,
}
pub fn spec() -> ToolSpec {
ToolSpec {
name: "shell".into(),
description: "Run a shell command in the launch cwd. Requires full-access mode. Streams stdout/stderr while running, then returns stdout, stderr, and exit code. Use timeout (seconds) to limit runtime."
.into(),
parameters: schema::object(
json!({
"command": {"type": "string", "description": "Shell command to execute"},
"timeout": {"type": "integer", "description": "Optional timeout in seconds (default 30)"}
}),
&["command"],
),
}
}
pub async fn run(args: Value, ctx: &ToolContext) -> Result<String> {
let args: Args = serde_json::from_value(args)?;
if !ctx.mode.can_write() {
bail!("shell tool requires full-access mode");
}
let mut cmd = tokio::process::Command::new("sh");
cmd.arg("-c").arg(&args.command);
cmd.current_dir(&ctx.cwd);
cmd.stdout(Stdio::piped());
cmd.stderr(Stdio::piped());
let mut child = cmd.spawn().context("spawning shell command")?;
let stdout = child.stdout.take().context("capturing command stdout")?;
let stderr = child.stderr.take().context("capturing command stderr")?;
let stdout_tx = ctx.runtime_tx.clone();
let stderr_tx = ctx.runtime_tx.clone();
let stdout_task = tokio::spawn(read_stream("stdout", stdout, stdout_tx));
let stderr_task = tokio::spawn(read_stream("stderr", stderr, stderr_tx));
let timeout = Duration::from_secs(args.timeout.unwrap_or(30));
let mut timed_out = false;
let status = match tokio::time::timeout(timeout, child.wait()).await {
Ok(status) => status.context("waiting for shell command")?,
Err(_) => {
timed_out = true;
let _ = child.kill().await;
child
.wait()
.await
.context("waiting for timed-out shell command to exit")?
}
};
let stdout = stdout_task
.await
.context("joining stdout reader")?
.context("reading command stdout")?;
let stderr = stderr_task
.await
.context("joining stderr reader")?
.context("reading command stderr")?;
let result = format_result(&stdout, &stderr, status.code().unwrap_or(-1));
if timed_out {
bail!("command timed out after {}s\n{}", timeout.as_secs(), result);
}
Ok(result)
}
async fn read_stream<R>(
stream: &'static str,
mut reader: R,
tx: Option<tokio::sync::mpsc::UnboundedSender<ToolRuntimeEvent>>,
) -> Result<Vec<u8>>
where
R: AsyncRead + Unpin,
{
let mut collected = Vec::new();
let mut buf = [0_u8; 4096];
loop {
let n = reader.read(&mut buf).await?;
if n == 0 {
break;
}
let chunk = &buf[..n];
collected.extend_from_slice(chunk);
if let Some(tx) = &tx {
let _ = tx.send(ToolRuntimeEvent::OutputChunk {
stream: stream.to_string(),
content: String::from_utf8_lossy(chunk).to_string(),
});
}
}
Ok(collected)
}
fn format_result(stdout: &[u8], stderr: &[u8], code: i32) -> String {
let stdout = String::from_utf8_lossy(stdout);
let stderr = String::from_utf8_lossy(stderr);
let mut result = String::new();
if !stdout.is_empty() {
result.push_str("stdout:\n");
result.push_str(&stdout);
if !stdout.ends_with('\n') {
result.push('\n');
}
}
if !stderr.is_empty() {
result.push_str("stderr:\n");
result.push_str(&stderr);
if !stderr.ends_with('\n') {
result.push('\n');
}
}
if result.is_empty() {
result.push_str("(no output)\n");
}
result.push_str(&format!("exit code: {}\n", code));
result
}
+8 -6
View File
@@ -330,7 +330,8 @@ fn heading_for(block: &TranscriptBlock) -> String {
}
fn display_content(block: &TranscriptBlock, show_full_tools: bool, show_reasoning: bool) -> String {
if matches!(block.kind, TranscriptKind::Tool) && !show_full_tools {
if matches!(block.kind, TranscriptKind::Tool) && !show_full_tools && !is_live_tool_output(block)
{
String::new()
} else if matches!(block.kind, TranscriptKind::Reasoning) && !show_reasoning {
String::new()
@@ -339,6 +340,10 @@ fn display_content(block: &TranscriptBlock, show_full_tools: bool, show_reasonin
}
}
fn is_live_tool_output(block: &TranscriptBlock) -> bool {
block.title.contains('…') && block.content.contains("streamed output:\n")
}
fn footer_text(state: &RenderState<'_>) -> String {
let busy = if state.busy { "running" } else { "idle" };
let mode = match state.mode {
@@ -445,8 +450,7 @@ fn input_height(input: &str, available_width: u16) -> u16 {
// available_width - 2. Count word-wrapped rows so that a single long
// line grows the input area instead of overflowing horizontally.
let prefix_width = 2usize;
let content_width = available_width
.saturating_sub(prefix_width as u16) as usize;
let content_width = available_width.saturating_sub(prefix_width as u16) as usize;
let content_width = content_width.max(1);
let mut rows = 0u16;
@@ -456,9 +460,7 @@ fn input_height(input: &str, available_width: u16) -> u16 {
input.lines().collect()
};
for line in &lines {
rows = rows.saturating_add(
ratatui_wrapped_row_count(line, content_width) as u16,
);
rows = rows.saturating_add(ratatui_wrapped_row_count(line, content_width) as u16);
}
if input.ends_with('\n') {
rows = rows.saturating_add(1);
+29
View File
@@ -10,6 +10,7 @@ fn ctx(root: &std::path::Path, mode: AccessMode) -> ToolContext {
read_roots: vec![root.to_path_buf()],
blocked_write_roots: Vec::new(),
model_result_limit: 100_000,
runtime_tx: None,
}
}
@@ -20,6 +21,7 @@ fn ctx_with_docs(root: &std::path::Path, docs: &std::path::Path, mode: AccessMod
read_roots: vec![root.to_path_buf(), docs.to_path_buf()],
blocked_write_roots: vec![docs.to_path_buf()],
model_result_limit: 100_000,
runtime_tx: None,
}
}
@@ -172,6 +174,33 @@ async fn full_access_blocks_write_and_edit_under_docs_root() {
);
}
#[tokio::test]
async fn shell_streams_output_chunks() {
let dir = tempdir().unwrap();
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let mut context = ctx(dir.path(), AccessMode::FullAccess);
context.runtime_tx = Some(tx);
let out = tools::execute(
"shell",
json!({"command":"printf hello; printf err >&2"}),
&context,
)
.await;
assert!(out.ok);
assert!(out.content.contains("stdout:\nhello"));
assert!(out.content.contains("stderr:\nerr"));
let mut streamed = String::new();
while let Ok(event) = rx.try_recv() {
let cassady::tools::ToolRuntimeEvent::OutputChunk { stream, content } = event;
streamed.push_str(&format!("{stream}:{content}"));
}
assert!(streamed.contains("stdout:hello"));
assert!(streamed.contains("stderr:err"));
}
#[cfg(unix)]
#[tokio::test]
async fn full_access_blocks_writes_through_symlinked_docs_dir() {