use std::collections::BTreeMap;
use nojson::{Json, RawJsonValue};
use crate::metrics::Counter;
use crate::sansio::deepseek::{ChatMessage, ToolCall, ToolDef};
pub const ARGUMENTS_MAX_BYTES: usize = 64 * 1024;
pub const TURN_TOOL_CALL_LIMIT: usize = 20;
pub const DEFAULT_LIST_MAX_ENTRIES: usize = 200;
pub const DEFAULT_SEARCH_MAX_RESULTS: usize = 50;
pub const READ_MAX_BYTES: usize = 1024 * 1024;
pub const PATCH_MAX_EDITS: usize = 20;
pub const PATCH_MAX_FILE_BYTES: usize = READ_MAX_BYTES;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReadOnlyTool {
List {
path: String,
recursive: bool,
max_entries: usize,
include_hidden: bool,
},
Read {
path: String,
line_range: Option<(usize, usize)>,
},
Search {
pattern: String,
path_prefix: Option<String>,
case_sensitive: bool,
max_results: usize,
},
}
impl ReadOnlyTool {
pub fn definitions() -> Vec<ToolDef> {
vec![
ToolDef {
name: "list".to_string(),
description: "List files and directories under a workspace-relative path. \
Returns a JSON array of {path, kind, size} entries; the \
result includes truncated:true when max_entries is hit. \
A path outside the workspace requires one-shot human \
approval before it can be listed."
.to_string(),
parameters_json: LIST_PARAMS_SCHEMA.to_string(),
},
ToolDef {
name: "read".to_string(),
description: "Read a UTF-8 text file at a workspace-relative path. \
Optionally restrict to a 1-indexed inclusive [start, end] \
line range. Content is truncated to the first 1 MiB. \
A path outside the workspace requires one-shot human \
approval before it can be read."
.to_string(),
parameters_json: READ_PARAMS_SCHEMA.to_string(),
},
ToolDef {
name: "search".to_string(),
description: "Literal substring search across text files under an \
optional workspace-relative prefix. Returns \
{path, line, snippet} hits; binary or non-UTF-8 files \
are silently skipped. A path_prefix outside the \
workspace requires one-shot human approval."
.to_string(),
parameters_json: SEARCH_PARAMS_SCHEMA.to_string(),
},
]
}
pub fn parse(function_name: &str, arguments_json: &str) -> Result<Self, ToolExecutionError> {
match function_name {
"list" => parse_list(arguments_json),
"read" => parse_read(arguments_json),
"search" => parse_search(arguments_json),
_ => Err(ToolExecutionError::UnknownTool),
}
}
}
const LIST_PARAMS_SCHEMA: &str = r#"{
"type":"object",
"properties":{
"path":{"type":"string","description":"Workspace-relative directory path (e.g. \".\" or \"src\")."},
"recursive":{"type":"boolean","default":false},
"max_entries":{"type":"integer","default":200,"minimum":1},
"include_hidden":{"type":"boolean","default":false}
},
"required":["path"]
}"#;
const READ_PARAMS_SCHEMA: &str = r#"{
"type":"object",
"properties":{
"path":{"type":"string","description":"Workspace-relative file path."},
"line_range":{"type":"array","items":{"type":"integer","minimum":1},"minItems":2,"maxItems":2,"description":"1-indexed inclusive [start, end] range."}
},
"required":["path"]
}"#;
const SEARCH_PARAMS_SCHEMA: &str = r#"{
"type":"object",
"properties":{
"pattern":{"type":"string","description":"Literal substring to match (no regex)."},
"path_prefix":{"type":"string","description":"Restrict to a workspace-relative subtree."},
"case_sensitive":{"type":"boolean","default":false},
"max_results":{"type":"integer","default":50,"minimum":1}
},
"required":["pattern"]
}"#;
fn parse_list(arguments_json: &str) -> Result<ReadOnlyTool, ToolExecutionError> {
let json = nojson::RawJson::parse(arguments_json).map_err(map_parse_err)?;
let root = json.value();
let path = required_string(root, "path")?;
let recursive = optional_bool(root, "recursive")?.unwrap_or(false);
let max_entries = optional_usize(root, "max_entries")?.unwrap_or(DEFAULT_LIST_MAX_ENTRIES);
let include_hidden = optional_bool(root, "include_hidden")?.unwrap_or(false);
Ok(ReadOnlyTool::List {
path,
recursive,
max_entries,
include_hidden,
})
}
fn parse_read(arguments_json: &str) -> Result<ReadOnlyTool, ToolExecutionError> {
let json = nojson::RawJson::parse(arguments_json).map_err(map_parse_err)?;
let root = json.value();
let path = required_string(root, "path")?;
let line_range = match root
.to_member("line_range")
.map_err(map_parse_err)?
.optional()
{
None => None,
Some(value) => {
let mut iter = value.to_array().map_err(map_parse_err)?;
let start = next_usize(&mut iter, "line_range")?;
let end = next_usize(&mut iter, "line_range")?;
if iter.next().is_some() {
return Err(ToolExecutionError::ArgumentsParseFailed(
"line_range must have exactly 2 elements".to_string(),
));
}
Some((start, end))
}
};
Ok(ReadOnlyTool::Read { path, line_range })
}
fn parse_search(arguments_json: &str) -> Result<ReadOnlyTool, ToolExecutionError> {
let json = nojson::RawJson::parse(arguments_json).map_err(map_parse_err)?;
let root = json.value();
let pattern = required_string(root, "pattern")?;
let path_prefix = optional_string(root, "path_prefix")?;
let case_sensitive = optional_bool(root, "case_sensitive")?.unwrap_or(false);
let max_results = optional_usize(root, "max_results")?.unwrap_or(DEFAULT_SEARCH_MAX_RESULTS);
Ok(ReadOnlyTool::Search {
pattern,
path_prefix,
case_sensitive,
max_results,
})
}
fn map_parse_err(err: nojson::JsonParseError) -> ToolExecutionError {
ToolExecutionError::ArgumentsParseFailed(err.to_string())
}
fn next_usize<'text, 'raw, I>(iter: &mut I, field: &str) -> Result<usize, ToolExecutionError>
where
I: Iterator<Item = RawJsonValue<'text, 'raw>>,
{
iter.next()
.ok_or_else(|| {
ToolExecutionError::ArgumentsParseFailed(format!(
"{field} must have exactly 2 elements"
))
})?
.try_into()
.map_err(map_parse_err)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PatchTool {
Add { path: String, content: String },
Update {
path: String,
before: String,
after: String,
},
}
impl PatchTool {
pub fn path(&self) -> &str {
match self {
Self::Add { path, .. } | Self::Update { path, .. } => path,
}
}
pub fn kind_label(&self) -> &'static str {
match self {
Self::Add { .. } => "add",
Self::Update { .. } => "update",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PatchInvocation {
pub edits: Vec<PatchTool>,
}
impl PatchInvocation {
pub fn definition() -> ToolDef {
ToolDef {
name: "patch".to_string(),
description: "Apply a batch of file edits to the workspace. \
Each edit is either an add (create a new file) or \
an update (replace a unique substring). All edits \
in one call must target distinct paths. Patches \
that only update git-tracked files are applied \
immediately; any add, non-tracked edit, or target \
outside the workspace requires user approval before \
it touches the filesystem. A path outside the \
workspace is given either as an absolute path or as \
a relative path that escapes with `..`."
.to_string(),
parameters_json: PATCH_PARAMS_SCHEMA.to_string(),
}
}
pub fn parse(arguments_json: &str) -> Result<Self, ToolExecutionError> {
let json = nojson::RawJson::parse(arguments_json).map_err(map_parse_err)?;
let root = json.value();
let edits_value = root
.to_member("edits")
.map_err(map_parse_err)?
.required()
.map_err(map_parse_err)?;
let mut edits: Vec<PatchTool> = Vec::new();
let mut seen_paths: std::collections::HashSet<String> = std::collections::HashSet::new();
for item in edits_value.to_array().map_err(map_parse_err)? {
let kind = required_string(item, "kind")?;
let path = required_string(item, "path")?;
if !seen_paths.insert(path.clone()) {
return Err(ToolExecutionError::Patch(
PatchError::MultipleEditsSamePath { path },
));
}
let tool = match kind.as_str() {
"add" => {
let content = required_string(item, "content")?;
if content.len() > PATCH_MAX_FILE_BYTES {
return Err(ToolExecutionError::Patch(PatchError::FileTooLarge { path }));
}
PatchTool::Add { path, content }
}
"update" => {
let before = required_string(item, "before")?;
let after = required_string(item, "after")?;
if after.len() > PATCH_MAX_FILE_BYTES {
return Err(ToolExecutionError::Patch(PatchError::FileTooLarge { path }));
}
PatchTool::Update {
path,
before,
after,
}
}
other => {
return Err(ToolExecutionError::ArgumentsParseFailed(format!(
"unknown edit kind: {other}"
)));
}
};
edits.push(tool);
}
if edits.is_empty() {
return Err(ToolExecutionError::ArgumentsParseFailed(
"edits must not be empty".to_string(),
));
}
if edits.len() > PATCH_MAX_EDITS {
return Err(ToolExecutionError::Patch(PatchError::TooManyEdits {
count: edits.len() as u64,
}));
}
Ok(Self { edits })
}
}
const PATCH_PARAMS_SCHEMA: &str = r#"{
"type":"object",
"properties":{
"edits":{"type":"array","minItems":1,"items":{
"type":"object",
"properties":{
"kind":{"type":"string","enum":["add","update"]},
"path":{"type":"string","description":"Workspace-relative target path. Must be unique across edits in one call."},
"content":{"type":"string","description":"Full contents for add."},
"before":{"type":"string","description":"For update: byte-exact substring to replace. Must match exactly once."},
"after":{"type":"string","description":"For update: replacement text."}
},
"required":["kind","path"]
}}
},
"required":["edits"]
}"#;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CommandInvocation {
pub argv: Vec<String>,
}
impl CommandInvocation {
pub fn definition() -> ToolDef {
ToolDef {
name: "command".to_string(),
description: "Run a command in the workspace by executing argv[0] with argv[1..] \
directly (no shell). Every call requires user approval unless a matching \
argv_prefix rule pre-approves it. Output is capped at 256 KiB per stream; if \
a stream is truncated the result sets `truncated: true`. Non-zero exit status \
is returned as a normal result (not an error). Runtime is capped (default 180 \
seconds); a command killed on timeout reports `termination_reason: \"timeout\"`."
.to_string(),
parameters_json: COMMAND_PARAMS_SCHEMA.to_string(),
}
}
pub fn parse(arguments_json: &str) -> Result<Self, ToolExecutionError> {
let json = nojson::RawJson::parse(arguments_json).map_err(map_parse_err)?;
let root = json.value();
let argv = required_string_array(root, "argv")?;
if argv.is_empty() {
return Err(ToolExecutionError::Command(CommandError::EmptyArgv));
}
Ok(Self { argv })
}
}
const COMMAND_PARAMS_SCHEMA: &str = r#"{
"type":"object",
"properties":{
"argv":{"type":"array","items":{"type":"string"},"minItems":1,"description":"Command and arguments to exec directly (no shell interpretation). Use each program's own flags for pipe / redirect / glob equivalents (for example --max-count instead of piping to head). For a shell pipe or chain, invoke it explicitly as [\"bash\", \"-c\", \"...\"]; that will still require user approval unless a matching argv_prefix rule pre-approves it."}
},
"required":["argv"]
}"#;
fn required_string(root: RawJsonValue<'_, '_>, name: &str) -> Result<String, ToolExecutionError> {
let value = root
.to_member(name)
.map_err(map_parse_err)?
.required()
.map_err(map_parse_err)?;
value.try_into().map_err(map_parse_err)
}
fn required_string_array(
root: RawJsonValue<'_, '_>,
name: &str,
) -> Result<Vec<String>, ToolExecutionError> {
let value = root
.to_member(name)
.map_err(map_parse_err)?
.required()
.map_err(map_parse_err)?;
let array = value.to_array().map_err(map_parse_err)?;
let mut out = Vec::new();
for item in array {
let s: String = item.try_into().map_err(map_parse_err)?;
out.push(s);
}
Ok(out)
}
fn optional_string(
root: RawJsonValue<'_, '_>,
name: &str,
) -> Result<Option<String>, ToolExecutionError> {
match root.to_member(name).map_err(map_parse_err)?.optional() {
None => Ok(None),
Some(value) => value.try_into().map_err(map_parse_err),
}
}
fn optional_bool(
root: RawJsonValue<'_, '_>,
name: &str,
) -> Result<Option<bool>, ToolExecutionError> {
match root.to_member(name).map_err(map_parse_err)?.optional() {
None => Ok(None),
Some(value) => Ok(Some(value.try_into().map_err(map_parse_err)?)),
}
}
fn optional_usize(
root: RawJsonValue<'_, '_>,
name: &str,
) -> Result<Option<usize>, ToolExecutionError> {
match root.to_member(name).map_err(map_parse_err)?.optional() {
None => Ok(None),
Some(value) => Ok(Some(value.try_into().map_err(map_parse_err)?)),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ToolOutcome {
Ok(String),
Err(ToolExecutionError),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ToolExecutionError {
OutsideWorkspace,
NotUtf8,
Binary,
IoError(String),
ArgumentsParseFailed(String),
ArgumentsTooLarge,
UnknownTool,
TurnToolCallLimitExceeded,
Patch(PatchError),
Command(CommandError),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PatchError {
Rejected,
Conflict { path: String },
NoMatch { path: String },
AmbiguousMatch { path: String, match_count: u64 },
AddOnExistingFile { path: String },
ParentDirMissing { path: String },
UpdateOnMissingFile { path: String },
ExcludedPath { path: String, reason: String },
IgnoredParent { path: String },
TooManyEdits { count: u64 },
MultipleEditsSamePath { path: String },
FileTooLarge { path: String },
IoError { path: String, message: String },
CrossDeviceRename { path: String },
}
impl PatchError {
pub fn to_code_and_message(&self) -> (&'static str, String) {
match self {
Self::Rejected => (
"patch_rejected",
"user rejected the patch preview".to_string(),
),
Self::Conflict { path } => (
"patch_conflict",
format!("target file changed between preview and apply: {path}"),
),
Self::NoMatch { path } => (
"patch_no_match",
format!("`before` did not match any content in {path}"),
),
Self::AmbiguousMatch { path, match_count } => (
"patch_ambiguous_match",
format!("`before` matched {match_count} places in {path}; expected exactly 1"),
),
Self::AddOnExistingFile { path } => (
"patch_add_on_existing_file",
format!("cannot add: file already exists at {path}"),
),
Self::ParentDirMissing { path } => (
"patch_parent_dir_missing",
format!("parent directory does not exist for {path}"),
),
Self::UpdateOnMissingFile { path } => (
"patch_update_on_missing_file",
format!("cannot update: file does not exist at {path}"),
),
Self::ExcludedPath { path, reason } => (
"patch_excluded_path",
format!("path is runtime-critical ({reason}): {path}"),
),
Self::IgnoredParent { path } => (
"patch_ignored_parent",
format!("cannot add into git-ignored directory: {path}"),
),
Self::TooManyEdits { count } => (
"patch_too_many_edits",
format!("edits count {count} exceeded the {PATCH_MAX_EDITS} limit"),
),
Self::MultipleEditsSamePath { path } => (
"patch_multiple_edits_same_path",
format!("more than one edit targets {path} within the same patch call"),
),
Self::FileTooLarge { path } => (
"patch_file_too_large",
format!(
"target or new content for {path} exceeded the {PATCH_MAX_FILE_BYTES} byte limit"
),
),
Self::IoError { path, message } => (
"patch_io_error",
format!("filesystem I/O failed for {path}: {message}"),
),
Self::CrossDeviceRename { path } => (
"patch_cross_device_rename",
format!("cross-device rename not supported for {path}"),
),
}
}
pub fn to_hint(&self) -> Option<&'static str> {
match self {
Self::MultipleEditsSamePath { .. } => Some(
"each path may appear in only one edit per call; split the edits into \
separate patch calls (one per target path) or merge this path's changes \
into a single update with one before/after pair",
),
Self::NoMatch { .. } => Some(
"`before` is not present verbatim in the file; read the file first, then \
copy an exact existing substring into `before`",
),
Self::AmbiguousMatch { .. } => Some(
"`before` matched more than one place; extend `before` with surrounding \
lines so it matches exactly once",
),
Self::AddOnExistingFile { .. } => Some(
"path already exists; use update with a before/after pair instead of add, \
or choose a different path",
),
Self::UpdateOnMissingFile { .. } => {
Some("path does not exist; use add to create it, or fix the path")
}
Self::ParentDirMissing { .. } => Some(
"the parent directory does not exist; create it first (e.g. mkdir -p), \
then retry",
),
Self::TooManyEdits { .. } => {
Some("too many edits in one call; split across multiple patch calls")
}
Self::FileTooLarge { .. } => {
Some("content is too large; reduce it or split across multiple patch calls")
}
Self::Conflict { .. } => {
Some("target changed between preview and apply; re-read the file and retry")
}
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CommandError {
Rejected,
SpawnFailed { message: String },
EmptyArgv,
}
impl CommandError {
pub fn to_code_and_message(&self) -> (&'static str, String) {
match self {
Self::Rejected => (
"command_rejected",
"user rejected the command preview".to_string(),
),
Self::SpawnFailed { message } => (
"command_spawn_failed",
format!("failed to spawn command: {message}"),
),
Self::EmptyArgv => ("command_empty_argv", "argv must not be empty".to_string()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PreviewContent {
pub path: String,
pub content: Option<Vec<u8>>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct PatchPreview {
pub target_paths: Vec<String>,
pub added_lines: u64,
pub removed_lines: u64,
pub edit_count: u64,
pub auto_approve: bool,
pub not_revertible: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ApprovalState {
NotRequired,
Pending,
Approved,
Rejected,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CommandPreview {
pub argv: Vec<String>,
pub working_directory: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct CommandOutputTail {
pub stdout_tail: String,
pub stderr_tail: String,
pub stdout_bytes_total: u64,
pub stderr_bytes_total: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommandOutputStream {
Stdout,
Stderr,
}
const COMMAND_TAIL_CHARS: usize = 4 * 1024;
impl ToolExecutionError {
pub fn message(&self) -> String {
let (_code, message) = self.code_and_message();
message
}
pub fn to_json_string(&self) -> String {
let (code, message) = self.code_and_message();
let hint: Option<&'static str> = match self {
Self::Patch(err) => err.to_hint(),
_ => None,
};
Json(ToolErrorJson {
code,
message: &message,
hint,
})
.to_string()
}
fn code_and_message(&self) -> (&'static str, String) {
match self {
Self::OutsideWorkspace => (
"outside_workspace",
"path escapes the workspace root".to_string(),
),
Self::NotUtf8 => ("not_utf8", "file is not valid UTF-8".to_string()),
Self::Binary => (
"binary",
"file contains binary data and cannot be read as text".to_string(),
),
Self::IoError(msg) => ("io_error", msg.clone()),
Self::ArgumentsParseFailed(msg) => ("arguments_parse_failed", msg.clone()),
Self::ArgumentsTooLarge => (
"arguments_too_large",
format!(
"tool call arguments exceeded the {} byte limit",
ARGUMENTS_MAX_BYTES
),
),
Self::UnknownTool => (
"unknown_tool",
"function_name does not match a known tool".to_string(),
),
Self::TurnToolCallLimitExceeded => (
"turn_tool_call_limit_exceeded",
format!(
"this user turn exceeded the {} tool call limit",
TURN_TOOL_CALL_LIMIT
),
),
Self::Patch(err) => err.to_code_and_message(),
Self::Command(err) => err.to_code_and_message(),
}
}
}
struct ToolErrorJson<'a> {
code: &'a str,
message: &'a str,
hint: Option<&'a str>,
}
impl nojson::DisplayJson for ToolErrorJson<'_> {
fn fmt(&self, f: &mut nojson::JsonFormatter<'_, '_>) -> std::fmt::Result {
f.object(|f| {
f.member("error", self.code)?;
f.member("message", self.message)?;
f.member("hint", self.hint)?;
Ok(())
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct RequestId(u64);
impl RequestId {
pub const fn new(raw: u64) -> Self {
Self(raw)
}
pub const fn as_u64(self) -> u64 {
self.0
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum Status {
#[default]
Idle,
AwaitingModel,
Streaming,
ToolRunning,
AwaitingApproval,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct PendingResponse {
pub content: String,
pub finish_reason: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Event {
UserMessage(String),
Cancel,
ContentDelta { request: RequestId, text: String },
ToolCallDelta {
request: RequestId,
index: u64,
id: Option<String>,
function_name: Option<String>,
arguments_fragment: Option<String>,
},
Finish {
request: RequestId,
reason: Option<String>,
},
ToolResult {
request: RequestId,
call_id: String,
outcome: ToolOutcome,
},
PatchPreviewReady {
request: RequestId,
call_id: String,
preview_content: Vec<PreviewContent>,
preview: PatchPreview,
},
ApproveToolCall { call_id: String },
RejectToolCall { call_id: String },
CommandOutputChunk {
request: RequestId,
call_id: String,
stream: CommandOutputStream,
bytes: Vec<u8>,
},
TransportError { request: RequestId, message: String },
Timeout { request: RequestId },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Action {
StartRequest {
id: RequestId,
messages: Vec<ChatMessage>,
},
CancelRequest { id: RequestId },
ExecuteTool {
request: RequestId,
call_id: String,
invocation: ReadOnlyTool,
},
PreviewPatch {
request: RequestId,
call_id: String,
invocation: PatchInvocation,
},
ApplyPatch {
request: RequestId,
call_id: String,
invocation: PatchInvocation,
preview_content: Vec<PreviewContent>,
},
ExecuteCommand {
request: RequestId,
call_id: String,
invocation: CommandInvocation,
},
CancelToolExecution { request: RequestId },
ReportError { message: String },
Redraw,
}
#[derive(Debug, Clone, Default)]
pub struct AgentCore {
conversation: Vec<ChatMessage>,
pending: Option<Pending>,
status: Status,
next_id: u64,
tool_calls_this_turn: usize,
workspace_display: String,
metrics: AgentMetrics,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Pending {
id: RequestId,
phase: PendingPhase,
response: PendingResponse,
tool_call_slots: BTreeMap<u64, ToolCallSlot>,
tool_results: Vec<PendingToolResult>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PendingPhase {
Streaming,
ToolRunning,
AwaitingApproval,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ToolCallSlot {
id: Option<String>,
function_name: Option<String>,
arguments: String,
over_limit: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PendingToolResult {
call_id: String,
function_name: String,
arguments_json: String,
approval: ApprovalState,
patch_preview: Option<PatchPreview>,
preview_content: Vec<PreviewContent>,
command_preview: Option<CommandPreview>,
command_output_tail: Option<CommandOutputTail>,
outcome: Option<ToolOutcome>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ActiveToolCall {
pub call_id: String,
pub function_name: String,
pub arguments_json: String,
pub outcome: Option<ToolOutcome>,
pub is_streaming: bool,
pub approval: ApprovalState,
pub patch_preview: Option<PatchPreview>,
pub preview_content: Vec<PreviewContent>,
pub command_preview: Option<CommandPreview>,
pub command_output_tail: Option<CommandOutputTail>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct AgentMetrics {
pub user_messages_accepted: Counter,
pub user_messages_rejected_while_active: Counter,
pub cancels_applied: Counter,
pub cancels_ignored_when_idle: Counter,
pub content_deltas_appended: Counter,
pub content_deltas_dropped_as_stale: Counter,
pub finishes_committed: Counter,
pub finishes_dropped_as_stale: Counter,
pub transport_errors_recorded: Counter,
pub transport_errors_dropped_as_stale: Counter,
pub timeouts_applied: Counter,
pub timeouts_dropped_as_stale: Counter,
pub tool_call_deltas_appended: Counter,
pub tool_call_deltas_dropped_as_stale: Counter,
pub tool_call_arguments_fragments_dropped_over_limit: Counter,
pub tool_results_committed: Counter,
pub tool_results_dropped_as_stale: Counter,
pub tool_calls_executed: Counter,
pub tool_calls_rejected_by_turn_limit: Counter,
pub tool_calls_rejected_by_arguments_limit: Counter,
pub patch_calls_previewed: Counter,
pub patch_previews_committed: Counter,
pub patch_previews_dropped_as_stale: Counter,
pub tool_call_approvals_committed: Counter,
pub tool_call_approvals_dropped_as_stale: Counter,
pub tool_call_rejections_committed: Counter,
pub tool_call_rejections_dropped_as_stale: Counter,
pub command_calls_dispatched: Counter,
pub command_executions_started: Counter,
pub command_output_chunks_appended: Counter,
pub command_output_chunks_dropped_as_stale: Counter,
}
impl AgentCore {
pub fn new() -> Self {
Self::default()
}
pub fn set_workspace_display(&mut self, display: String) {
self.workspace_display = display;
}
pub fn conversation(&self) -> &[ChatMessage] {
&self.conversation
}
pub fn pending_response(&self) -> Option<&PendingResponse> {
self.pending.as_ref().map(|p| &p.response)
}
pub fn active_request(&self) -> Option<RequestId> {
self.pending.as_ref().map(|p| p.id)
}
pub fn pending_approval_call_id(&self) -> Option<String> {
self.pending
.as_ref()?
.tool_results
.iter()
.find(|r| r.approval == ApprovalState::Pending)
.map(|r| r.call_id.clone())
}
pub fn status(&self) -> Status {
self.status
}
pub fn active_tool_calls(&self) -> Vec<ActiveToolCall> {
let Some(pending) = self.pending.as_ref() else {
return Vec::new();
};
match pending.phase {
PendingPhase::Streaming => pending
.tool_call_slots
.iter()
.map(|(index, slot)| ActiveToolCall {
call_id: slot
.id
.clone()
.unwrap_or_else(|| format!("__pending_{index}")),
function_name: slot.function_name.clone().unwrap_or_default(),
arguments_json: slot.arguments.clone(),
outcome: None,
is_streaming: true,
approval: ApprovalState::NotRequired,
patch_preview: None,
preview_content: Vec::new(),
command_preview: None,
command_output_tail: None,
})
.collect(),
PendingPhase::ToolRunning | PendingPhase::AwaitingApproval => pending
.tool_results
.iter()
.map(|r| ActiveToolCall {
call_id: r.call_id.clone(),
function_name: r.function_name.clone(),
arguments_json: r.arguments_json.clone(),
outcome: r.outcome.clone(),
is_streaming: false,
approval: r.approval,
patch_preview: r.patch_preview.clone(),
preview_content: r.preview_content.clone(),
command_preview: r.command_preview.clone(),
command_output_tail: r.command_output_tail.clone(),
})
.collect(),
}
}
pub fn metrics(&self) -> &AgentMetrics {
&self.metrics
}
pub fn handle_event(&mut self, event: Event) -> Vec<Action> {
match event {
Event::UserMessage(text) => self.on_user_message(text),
Event::Cancel => self.on_cancel(),
Event::ContentDelta { request, text } => self.on_content_delta(request, text),
Event::ToolCallDelta {
request,
index,
id,
function_name,
arguments_fragment,
} => self.on_tool_call_delta(request, index, id, function_name, arguments_fragment),
Event::Finish { request, reason } => self.on_finish(request, reason),
Event::ToolResult {
request,
call_id,
outcome,
} => self.on_tool_result(request, call_id, outcome),
Event::PatchPreviewReady {
request,
call_id,
preview_content,
preview,
} => self.on_patch_preview_ready(request, call_id, preview_content, preview),
Event::ApproveToolCall { call_id } => self.on_approve_tool_call(call_id),
Event::RejectToolCall { call_id } => self.on_reject_tool_call(call_id),
Event::CommandOutputChunk {
request,
call_id,
stream,
bytes,
} => self.on_command_output_chunk(request, call_id, stream, bytes),
Event::TransportError { request, message } => self.on_transport_error(request, message),
Event::Timeout { request } => self.on_timeout(request),
}
}
fn on_user_message(&mut self, text: String) -> Vec<Action> {
if self.pending.is_some() {
self.metrics.user_messages_rejected_while_active.inc();
return Vec::new();
}
self.conversation.push(ChatMessage::User(text));
let id = self.mint_id();
self.pending = Some(Pending {
id,
phase: PendingPhase::Streaming,
response: PendingResponse::default(),
tool_call_slots: BTreeMap::new(),
tool_results: Vec::new(),
});
self.status = Status::AwaitingModel;
self.tool_calls_this_turn = 0;
self.metrics.user_messages_accepted.inc();
vec![
Action::StartRequest {
id,
messages: self.conversation.clone(),
},
Action::Redraw,
]
}
fn on_cancel(&mut self) -> Vec<Action> {
let Some(pending) = self.pending.take() else {
self.metrics.cancels_ignored_when_idle.inc();
return Vec::new();
};
let cancel_action = match pending.phase {
PendingPhase::Streaming => Action::CancelRequest { id: pending.id },
PendingPhase::ToolRunning | PendingPhase::AwaitingApproval => {
Action::CancelToolExecution {
request: pending.id,
}
}
};
self.status = Status::Idle;
self.tool_calls_this_turn = 0;
self.metrics.cancels_applied.inc();
vec![cancel_action, Action::Redraw]
}
fn on_content_delta(&mut self, request: RequestId, text: String) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.content_deltas_dropped_as_stale.inc();
return Vec::new();
};
if pending.id != request || pending.phase != PendingPhase::Streaming {
self.metrics.content_deltas_dropped_as_stale.inc();
return Vec::new();
}
pending.response.content.push_str(&text);
self.status = Status::Streaming;
self.metrics.content_deltas_appended.inc();
vec![Action::Redraw]
}
fn on_tool_call_delta(
&mut self,
request: RequestId,
index: u64,
id: Option<String>,
function_name: Option<String>,
arguments_fragment: Option<String>,
) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.tool_call_deltas_dropped_as_stale.inc();
return Vec::new();
};
if pending.id != request || pending.phase != PendingPhase::Streaming {
self.metrics.tool_call_deltas_dropped_as_stale.inc();
return Vec::new();
}
let slot = pending
.tool_call_slots
.entry(index)
.or_insert_with(|| ToolCallSlot {
id: None,
function_name: None,
arguments: String::new(),
over_limit: false,
});
if slot.id.is_none()
&& let Some(new_id) = id
{
slot.id = Some(new_id);
}
if slot.function_name.is_none()
&& let Some(new_name) = function_name
{
slot.function_name = Some(new_name);
}
if let Some(fragment) = arguments_fragment {
if slot.over_limit {
self.metrics
.tool_call_arguments_fragments_dropped_over_limit
.inc();
} else if slot.arguments.len().saturating_add(fragment.len()) > ARGUMENTS_MAX_BYTES {
slot.over_limit = true;
self.metrics
.tool_call_arguments_fragments_dropped_over_limit
.inc();
} else {
slot.arguments.push_str(&fragment);
}
}
self.status = Status::Streaming;
self.metrics.tool_call_deltas_appended.inc();
vec![Action::Redraw]
}
fn on_finish(&mut self, request: RequestId, reason: Option<String>) -> Vec<Action> {
let Some(pending_ref) = self.pending.as_ref() else {
self.metrics.finishes_dropped_as_stale.inc();
return Vec::new();
};
if pending_ref.id != request || pending_ref.phase != PendingPhase::Streaming {
self.metrics.finishes_dropped_as_stale.inc();
return Vec::new();
}
let mut pending = self.pending.take().expect("checked above");
pending.response.finish_reason = reason.clone();
let content = std::mem::take(&mut pending.response.content);
let is_tool_calls_reason = reason.as_deref() == Some("tool_calls");
let has_slots = !pending.tool_call_slots.is_empty();
if !is_tool_calls_reason || !has_slots {
self.conversation.push(ChatMessage::Assistant {
content,
tool_calls: Vec::new(),
});
self.status = Status::Idle;
self.tool_calls_this_turn = 0;
self.metrics.finishes_committed.inc();
return vec![Action::Redraw];
}
let request_id = pending.id;
let (tool_calls, over_limit_ids) =
finalize_tool_call_slots(std::mem::take(&mut pending.tool_call_slots));
self.conversation.push(ChatMessage::Assistant {
content,
tool_calls: tool_calls.clone(),
});
self.metrics.finishes_committed.inc();
let mut actions = Vec::new();
let mut pending_results: Vec<PendingToolResult> = Vec::with_capacity(tool_calls.len());
for call in tool_calls.into_iter() {
if over_limit_ids.contains(&call.id) {
pending_results.push(synthetic_err_result(
call,
ToolExecutionError::ArgumentsTooLarge,
));
self.metrics.tool_calls_rejected_by_arguments_limit.inc();
continue;
}
if self.tool_calls_this_turn >= TURN_TOOL_CALL_LIMIT {
pending_results.push(synthetic_err_result(
call,
ToolExecutionError::TurnToolCallLimitExceeded,
));
self.metrics.tool_calls_rejected_by_turn_limit.inc();
continue;
}
if call.function_name == "patch" {
match PatchInvocation::parse(&call.arguments_json) {
Ok(invocation) => {
actions.push(Action::PreviewPatch {
request: request_id,
call_id: call.id.clone(),
invocation,
});
pending_results.push(patch_pending_result(call));
self.tool_calls_this_turn += 1;
self.metrics.patch_calls_previewed.inc();
}
Err(err) => {
pending_results.push(synthetic_err_result(call, err));
}
}
} else if call.function_name == "command" {
match CommandInvocation::parse(&call.arguments_json) {
Ok(invocation) => {
let preview = CommandPreview {
argv: invocation.argv.clone(),
working_directory: self.workspace_display.clone(),
};
pending_results.push(command_pending_result(call, preview));
self.tool_calls_this_turn += 1;
self.metrics.command_calls_dispatched.inc();
}
Err(err) => {
pending_results.push(synthetic_err_result(call, err));
}
}
} else {
match ReadOnlyTool::parse(&call.function_name, &call.arguments_json) {
Ok(invocation) => {
actions.push(Action::ExecuteTool {
request: request_id,
call_id: call.id.clone(),
invocation,
});
pending_results.push(read_only_pending_result(call));
self.tool_calls_this_turn += 1;
self.metrics.tool_calls_executed.inc();
}
Err(err) => {
pending_results.push(synthetic_err_result(call, err));
}
}
}
}
let all_resolved = pending_results.iter().all(|r| r.outcome.is_some());
if all_resolved {
actions.extend(self.advance_to_next_request(pending_results));
return actions;
}
pending.phase = PendingPhase::ToolRunning;
pending.tool_results = pending_results;
self.pending = Some(pending);
self.status = Status::ToolRunning;
self.recompute_phase_and_status();
actions.push(Action::Redraw);
actions
}
fn on_tool_result(
&mut self,
request: RequestId,
call_id: String,
outcome: ToolOutcome,
) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.tool_results_dropped_as_stale.inc();
return Vec::new();
};
if pending.id != request || matches!(pending.phase, PendingPhase::Streaming) {
self.metrics.tool_results_dropped_as_stale.inc();
return Vec::new();
}
let Some(entry) = pending
.tool_results
.iter_mut()
.find(|r| r.call_id == call_id && r.outcome.is_none())
else {
self.metrics.tool_results_dropped_as_stale.inc();
return Vec::new();
};
entry.outcome = Some(outcome);
self.metrics.tool_results_committed.inc();
self.recompute_phase_and_status();
self.maybe_advance_to_next_request()
}
fn on_patch_preview_ready(
&mut self,
request: RequestId,
call_id: String,
preview_content: Vec<PreviewContent>,
preview: PatchPreview,
) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.patch_previews_dropped_as_stale.inc();
return Vec::new();
};
if pending.id != request || matches!(pending.phase, PendingPhase::Streaming) {
self.metrics.patch_previews_dropped_as_stale.inc();
return Vec::new();
}
let Some(entry) = pending.tool_results.iter_mut().find(|r| {
r.call_id == call_id
&& r.approval == ApprovalState::NotRequired
&& r.outcome.is_none()
&& r.patch_preview.is_none()
}) else {
self.metrics.patch_previews_dropped_as_stale.inc();
return Vec::new();
};
entry.patch_preview = Some(preview);
entry.preview_content = preview_content;
entry.approval = ApprovalState::Pending;
self.metrics.patch_previews_committed.inc();
self.recompute_phase_and_status();
vec![Action::Redraw]
}
fn on_approve_tool_call(&mut self, call_id: String) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.tool_call_approvals_dropped_as_stale.inc();
return Vec::new();
};
let request_id = pending.id;
let Some(entry) = pending
.tool_results
.iter_mut()
.find(|r| r.call_id == call_id && r.approval == ApprovalState::Pending)
else {
self.metrics.tool_call_approvals_dropped_as_stale.inc();
return Vec::new();
};
entry.approval = ApprovalState::Approved;
let action = match entry.function_name.as_str() {
"patch" => {
let preview_content = entry.preview_content.clone();
match PatchInvocation::parse(&entry.arguments_json) {
Ok(invocation) => Action::ApplyPatch {
request: request_id,
call_id: call_id.clone(),
invocation,
preview_content,
},
Err(err) => {
entry.outcome = Some(ToolOutcome::Err(err));
self.metrics.tool_call_approvals_committed.inc();
self.recompute_phase_and_status();
return self.maybe_advance_to_next_request();
}
}
}
"command" => match CommandInvocation::parse(&entry.arguments_json) {
Ok(invocation) => {
self.metrics.command_executions_started.inc();
Action::ExecuteCommand {
request: request_id,
call_id: call_id.clone(),
invocation,
}
}
Err(err) => {
entry.outcome = Some(ToolOutcome::Err(err));
self.metrics.tool_call_approvals_committed.inc();
self.recompute_phase_and_status();
return self.maybe_advance_to_next_request();
}
},
other => {
entry.outcome = Some(ToolOutcome::Err(ToolExecutionError::ArgumentsParseFailed(
format!("no approval flow for tool {other}"),
)));
self.metrics.tool_call_approvals_committed.inc();
self.recompute_phase_and_status();
return self.maybe_advance_to_next_request();
}
};
self.metrics.tool_call_approvals_committed.inc();
self.recompute_phase_and_status();
vec![action, Action::Redraw]
}
fn on_reject_tool_call(&mut self, call_id: String) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.tool_call_rejections_dropped_as_stale.inc();
return Vec::new();
};
let Some(entry) = pending
.tool_results
.iter_mut()
.find(|r| r.call_id == call_id && r.approval == ApprovalState::Pending)
else {
self.metrics.tool_call_rejections_dropped_as_stale.inc();
return Vec::new();
};
entry.approval = ApprovalState::Rejected;
entry.outcome = Some(match entry.function_name.as_str() {
"command" => ToolOutcome::Err(ToolExecutionError::Command(CommandError::Rejected)),
_ => ToolOutcome::Err(ToolExecutionError::Patch(PatchError::Rejected)),
});
self.metrics.tool_call_rejections_committed.inc();
self.recompute_phase_and_status();
self.maybe_advance_to_next_request()
}
fn on_command_output_chunk(
&mut self,
request: RequestId,
call_id: String,
stream: CommandOutputStream,
bytes: Vec<u8>,
) -> Vec<Action> {
let Some(pending) = self.pending.as_mut() else {
self.metrics.command_output_chunks_dropped_as_stale.inc();
return Vec::new();
};
if pending.id != request || matches!(pending.phase, PendingPhase::Streaming) {
self.metrics.command_output_chunks_dropped_as_stale.inc();
return Vec::new();
}
let Some(entry) = pending.tool_results.iter_mut().find(|r| {
r.call_id == call_id
&& r.approval == ApprovalState::Approved
&& r.outcome.is_none()
&& r.function_name == "command"
}) else {
self.metrics.command_output_chunks_dropped_as_stale.inc();
return Vec::new();
};
let text = String::from_utf8_lossy(&bytes);
let byte_len = bytes.len() as u64;
let tail = entry
.command_output_tail
.get_or_insert_with(CommandOutputTail::default);
match stream {
CommandOutputStream::Stdout => {
append_bounded(&mut tail.stdout_tail, &text, COMMAND_TAIL_CHARS);
tail.stdout_bytes_total = tail.stdout_bytes_total.saturating_add(byte_len);
}
CommandOutputStream::Stderr => {
append_bounded(&mut tail.stderr_tail, &text, COMMAND_TAIL_CHARS);
tail.stderr_bytes_total = tail.stderr_bytes_total.saturating_add(byte_len);
}
}
self.metrics.command_output_chunks_appended.inc();
vec![Action::Redraw]
}
fn recompute_phase_and_status(&mut self) {
let Some(pending) = self.pending.as_mut() else {
return;
};
let has_pending_approval = pending
.tool_results
.iter()
.any(|r| r.approval == ApprovalState::Pending && r.outcome.is_none());
match (pending.phase, has_pending_approval) {
(PendingPhase::Streaming, _) => {}
(_, true) => {
pending.phase = PendingPhase::AwaitingApproval;
self.status = Status::AwaitingApproval;
}
(_, false) => {
pending.phase = PendingPhase::ToolRunning;
self.status = Status::ToolRunning;
}
}
}
fn maybe_advance_to_next_request(&mut self) -> Vec<Action> {
let Some(pending) = self.pending.as_ref() else {
return Vec::new();
};
if pending.tool_results.iter().all(|r| r.outcome.is_some()) {
let pending = self.pending.take().expect("checked above");
self.advance_to_next_request(pending.tool_results)
} else {
vec![Action::Redraw]
}
}
fn on_transport_error(&mut self, request: RequestId, message: String) -> Vec<Action> {
let Some(pending_ref) = self.pending.as_ref() else {
self.metrics.transport_errors_dropped_as_stale.inc();
return Vec::new();
};
if pending_ref.id != request || pending_ref.phase != PendingPhase::Streaming {
self.metrics.transport_errors_dropped_as_stale.inc();
return Vec::new();
}
self.pending = None;
self.status = Status::Idle;
self.tool_calls_this_turn = 0;
self.metrics.transport_errors_recorded.inc();
vec![Action::ReportError { message }, Action::Redraw]
}
fn on_timeout(&mut self, request: RequestId) -> Vec<Action> {
let Some(pending_ref) = self.pending.as_ref() else {
self.metrics.timeouts_dropped_as_stale.inc();
return Vec::new();
};
if pending_ref.id != request || pending_ref.phase != PendingPhase::Streaming {
self.metrics.timeouts_dropped_as_stale.inc();
return Vec::new();
}
let id = pending_ref.id;
self.pending = None;
self.status = Status::Idle;
self.tool_calls_this_turn = 0;
self.metrics.timeouts_applied.inc();
vec![
Action::CancelRequest { id },
Action::ReportError {
message: "request timed out".to_string(),
},
Action::Redraw,
]
}
fn mint_id(&mut self) -> RequestId {
let id = RequestId(self.next_id);
self.next_id = self.next_id.wrapping_add(1);
id
}
fn advance_to_next_request(&mut self, results: Vec<PendingToolResult>) -> Vec<Action> {
for result in results {
let outcome = result
.outcome
.expect("advance_to_next_request called with unresolved tool result");
let content = match outcome {
ToolOutcome::Ok(s) => s,
ToolOutcome::Err(err) => err.to_json_string(),
};
self.conversation.push(ChatMessage::Tool {
tool_call_id: result.call_id,
content,
});
}
let id = self.mint_id();
self.pending = Some(Pending {
id,
phase: PendingPhase::Streaming,
response: PendingResponse::default(),
tool_call_slots: BTreeMap::new(),
tool_results: Vec::new(),
});
self.status = Status::AwaitingModel;
vec![
Action::StartRequest {
id,
messages: self.conversation.clone(),
},
Action::Redraw,
]
}
}
fn synthetic_err_result(call: ToolCall, err: ToolExecutionError) -> PendingToolResult {
PendingToolResult {
call_id: call.id,
function_name: call.function_name,
arguments_json: call.arguments_json,
approval: ApprovalState::NotRequired,
patch_preview: None,
preview_content: Vec::new(),
command_preview: None,
command_output_tail: None,
outcome: Some(ToolOutcome::Err(err)),
}
}
fn read_only_pending_result(call: ToolCall) -> PendingToolResult {
PendingToolResult {
call_id: call.id,
function_name: call.function_name,
arguments_json: call.arguments_json,
approval: ApprovalState::NotRequired,
patch_preview: None,
preview_content: Vec::new(),
command_preview: None,
command_output_tail: None,
outcome: None,
}
}
fn patch_pending_result(call: ToolCall) -> PendingToolResult {
PendingToolResult {
call_id: call.id,
function_name: call.function_name,
arguments_json: call.arguments_json,
approval: ApprovalState::NotRequired,
patch_preview: None,
preview_content: Vec::new(),
command_preview: None,
command_output_tail: None,
outcome: None,
}
}
fn append_bounded(tail: &mut String, text: &str, max_chars: usize) {
tail.push_str(text);
let count = tail.chars().count();
if count > max_chars {
let drop = count - max_chars;
let split = tail
.char_indices()
.nth(drop)
.map(|(idx, _)| idx)
.unwrap_or(0);
tail.drain(..split);
}
}
fn command_pending_result(call: ToolCall, preview: CommandPreview) -> PendingToolResult {
PendingToolResult {
call_id: call.id,
function_name: call.function_name,
arguments_json: call.arguments_json,
approval: ApprovalState::Pending,
patch_preview: None,
preview_content: Vec::new(),
command_preview: Some(preview),
command_output_tail: None,
outcome: None,
}
}
fn finalize_tool_call_slots(
slots: BTreeMap<u64, ToolCallSlot>,
) -> (Vec<ToolCall>, std::collections::HashSet<String>) {
let mut calls = Vec::with_capacity(slots.len());
let mut over_limit = std::collections::HashSet::new();
for (index, slot) in slots.into_iter() {
let id = slot.id.unwrap_or_else(|| format!("__missing_id__{index}"));
let function_name = slot.function_name.unwrap_or_default();
if slot.over_limit {
over_limit.insert(id.clone());
}
calls.push(ToolCall {
id,
function_name,
arguments_json: slot.arguments,
});
}
(calls, over_limit)
}
#[cfg(test)]
mod tests {
use super::*;
fn user(core: &mut AgentCore, text: &str) -> Vec<Action> {
core.handle_event(Event::UserMessage(text.to_string()))
}
fn last_start_id(actions: &[Action]) -> RequestId {
for action in actions {
if let Action::StartRequest { id, .. } = action {
return *id;
}
}
panic!("no StartRequest in {actions:?}");
}
#[test]
fn user_message_from_idle_starts_request_and_appends_user_turn() {
let mut core = AgentCore::new();
let actions = user(&mut core, "hello");
assert_eq!(core.status(), Status::AwaitingModel);
assert_eq!(core.conversation().len(), 1);
assert!(matches!(core.conversation()[0], ChatMessage::User(_)));
assert!(matches!(actions[0], Action::StartRequest { .. }));
assert!(actions.contains(&Action::Redraw));
assert!(core.active_request().is_some());
}
#[test]
fn user_message_while_active_is_ignored() {
let mut core = AgentCore::new();
let _ = user(&mut core, "first");
let actions = user(&mut core, "second");
assert!(actions.is_empty());
assert_eq!(core.conversation().len(), 1);
}
#[test]
fn content_delta_accumulates_and_flips_status_to_streaming() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = core.handle_event(Event::ContentDelta {
request: id,
text: "he".to_string(),
});
assert_eq!(actions, vec![Action::Redraw]);
assert_eq!(core.status(), Status::Streaming);
let _ = core.handle_event(Event::ContentDelta {
request: id,
text: "llo".to_string(),
});
let pending = core.pending_response().expect("pending");
assert_eq!(pending.content, "hello");
}
#[test]
fn finish_commits_assistant_turn_and_returns_to_idle() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::ContentDelta {
request: id,
text: "hello".to_string(),
});
let actions = core.handle_event(Event::Finish {
request: id,
reason: Some("stop".to_string()),
});
assert_eq!(actions, vec![Action::Redraw]);
assert_eq!(core.status(), Status::Idle);
assert!(core.active_request().is_none());
assert!(core.pending_response().is_none());
assert_eq!(core.conversation().len(), 2);
match &core.conversation()[1] {
ChatMessage::Assistant { content, .. } => assert_eq!(content, "hello"),
other => panic!("expected assistant, got {other:?}"),
}
}
#[test]
fn cancel_drops_pending_response_and_asks_transport_to_cancel() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::ContentDelta {
request: id,
text: "partial".to_string(),
});
let actions = core.handle_event(Event::Cancel);
assert!(actions.contains(&Action::CancelRequest { id }));
assert!(actions.contains(&Action::Redraw));
assert_eq!(core.status(), Status::Idle);
assert!(core.pending_response().is_none());
assert_eq!(core.conversation().len(), 1);
}
#[test]
fn cancel_from_idle_is_no_op() {
let mut core = AgentCore::new();
assert!(core.handle_event(Event::Cancel).is_empty());
}
#[test]
fn transport_error_ends_request_and_reports_message() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = core.handle_event(Event::TransportError {
request: id,
message: "boom".to_string(),
});
assert!(actions.contains(&Action::ReportError {
message: "boom".to_string(),
}));
assert!(actions.contains(&Action::Redraw));
assert_eq!(core.status(), Status::Idle);
assert!(core.pending_response().is_none());
assert_eq!(core.conversation().len(), 1);
}
#[test]
fn timeout_cancels_transport_and_reports_error() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = core.handle_event(Event::Timeout { request: id });
assert!(actions.contains(&Action::CancelRequest { id }));
assert!(
actions
.iter()
.any(|a| matches!(a, Action::ReportError { .. })),
"expected ReportError, got {actions:?}",
);
assert!(actions.contains(&Action::Redraw));
assert!(core.pending_response().is_none());
}
#[test]
fn stale_deltas_are_dropped_without_state_change() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::Finish {
request: id,
reason: Some("stop".to_string()),
});
let actions = core.handle_event(Event::ContentDelta {
request: id,
text: "late".to_string(),
});
assert!(actions.is_empty());
assert_eq!(core.conversation().len(), 2);
}
#[test]
fn events_tagged_with_unknown_id_are_dropped() {
let mut core = AgentCore::new();
let _ = user(&mut core, "hi");
let actions = core.handle_event(Event::ContentDelta {
request: RequestId::new(u64::MAX),
text: "x".to_string(),
});
assert!(actions.is_empty());
assert!(core.pending_response().expect("pending").content.is_empty());
}
#[test]
fn each_start_request_gets_a_unique_id() {
let mut core = AgentCore::new();
let id1 = last_start_id(&user(&mut core, "first"));
let _ = core.handle_event(Event::Finish {
request: id1,
reason: None,
});
let id2 = last_start_id(&user(&mut core, "second"));
assert_ne!(id1, id2);
}
#[test]
fn late_finish_after_cancel_is_ignored() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::Cancel);
let actions = core.handle_event(Event::Finish {
request: id,
reason: Some("stop".to_string()),
});
assert!(actions.is_empty());
assert_eq!(core.conversation().len(), 1);
}
#[test]
fn start_request_carries_full_conversation_snapshot() {
let mut core = AgentCore::new();
let id1 = last_start_id(&user(&mut core, "first"));
let _ = core.handle_event(Event::ContentDelta {
request: id1,
text: "one".to_string(),
});
let _ = core.handle_event(Event::Finish {
request: id1,
reason: None,
});
let actions = user(&mut core, "second");
let messages = actions.iter().find_map(|a| match a {
Action::StartRequest { messages, .. } => Some(messages.clone()),
_ => None,
});
let messages = messages.expect("StartRequest present");
assert_eq!(messages.len(), 3);
assert!(matches!(messages[0], ChatMessage::User(_)));
assert!(matches!(messages[1], ChatMessage::Assistant { .. }));
match &messages[2] {
ChatMessage::User(content) => assert_eq!(content, "second"),
other => panic!("expected user, got {other:?}"),
}
}
#[test]
fn metrics_start_at_zero() {
let core = AgentCore::new();
assert_eq!(*core.metrics(), AgentMetrics::default());
}
#[test]
fn user_message_accepted_and_rejected_counters() {
let mut core = AgentCore::new();
let _ = user(&mut core, "first");
assert_eq!(core.metrics().user_messages_accepted.get(), 1);
assert_eq!(core.metrics().user_messages_rejected_while_active.get(), 0);
let _ = user(&mut core, "second");
assert_eq!(core.metrics().user_messages_accepted.get(), 1);
assert_eq!(core.metrics().user_messages_rejected_while_active.get(), 1);
}
#[test]
fn cancel_applied_and_ignored_counters() {
let mut core = AgentCore::new();
let _ = core.handle_event(Event::Cancel);
assert_eq!(core.metrics().cancels_applied.get(), 0);
assert_eq!(core.metrics().cancels_ignored_when_idle.get(), 1);
let _ = user(&mut core, "hi");
let _ = core.handle_event(Event::Cancel);
assert_eq!(core.metrics().cancels_applied.get(), 1);
assert_eq!(core.metrics().cancels_ignored_when_idle.get(), 1);
}
#[test]
fn content_delta_appended_and_dropped_counters() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::ContentDelta {
request: id,
text: "a".to_string(),
});
let _ = core.handle_event(Event::ContentDelta {
request: RequestId::new(u64::MAX),
text: "b".to_string(),
});
assert_eq!(core.metrics().content_deltas_appended.get(), 1);
assert_eq!(core.metrics().content_deltas_dropped_as_stale.get(), 1);
}
#[test]
fn finish_committed_and_dropped_counters() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::Finish {
request: id,
reason: None,
});
let _ = core.handle_event(Event::Finish {
request: id,
reason: None,
});
assert_eq!(core.metrics().finishes_committed.get(), 1);
assert_eq!(core.metrics().finishes_dropped_as_stale.get(), 1);
}
#[test]
fn transport_error_recorded_and_dropped_counters() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::TransportError {
request: id,
message: "boom".to_string(),
});
let _ = core.handle_event(Event::TransportError {
request: id,
message: "late".to_string(),
});
assert_eq!(core.metrics().transport_errors_recorded.get(), 1);
assert_eq!(core.metrics().transport_errors_dropped_as_stale.get(), 1);
}
#[test]
fn timeout_applied_and_dropped_counters() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(Event::Timeout { request: id });
let _ = core.handle_event(Event::Timeout { request: id });
assert_eq!(core.metrics().timeouts_applied.get(), 1);
assert_eq!(core.metrics().timeouts_dropped_as_stale.get(), 1);
}
#[test]
fn other_counters_do_not_move_on_a_single_event() {
let mut core = AgentCore::new();
let _ = user(&mut core, "hi");
let m = core.metrics();
assert_eq!(m.user_messages_accepted.get(), 1);
assert_eq!(m.user_messages_rejected_while_active.get(), 0);
assert_eq!(m.cancels_applied.get(), 0);
assert_eq!(m.cancels_ignored_when_idle.get(), 0);
assert_eq!(m.content_deltas_appended.get(), 0);
assert_eq!(m.content_deltas_dropped_as_stale.get(), 0);
assert_eq!(m.finishes_committed.get(), 0);
assert_eq!(m.finishes_dropped_as_stale.get(), 0);
assert_eq!(m.transport_errors_recorded.get(), 0);
assert_eq!(m.transport_errors_dropped_as_stale.get(), 0);
assert_eq!(m.timeouts_applied.get(), 0);
assert_eq!(m.timeouts_dropped_as_stale.get(), 0);
assert_eq!(m.tool_call_deltas_appended.get(), 0);
assert_eq!(m.tool_call_deltas_dropped_as_stale.get(), 0);
assert_eq!(m.tool_call_arguments_fragments_dropped_over_limit.get(), 0);
assert_eq!(m.tool_results_committed.get(), 0);
assert_eq!(m.tool_results_dropped_as_stale.get(), 0);
assert_eq!(m.tool_calls_executed.get(), 0);
assert_eq!(m.tool_calls_rejected_by_turn_limit.get(), 0);
assert_eq!(m.tool_calls_rejected_by_arguments_limit.get(), 0);
assert_eq!(m.patch_calls_previewed.get(), 0);
assert_eq!(m.patch_previews_committed.get(), 0);
assert_eq!(m.patch_previews_dropped_as_stale.get(), 0);
assert_eq!(m.tool_call_approvals_committed.get(), 0);
assert_eq!(m.tool_call_approvals_dropped_as_stale.get(), 0);
assert_eq!(m.tool_call_rejections_committed.get(), 0);
assert_eq!(m.tool_call_rejections_dropped_as_stale.get(), 0);
assert_eq!(m.command_calls_dispatched.get(), 0);
assert_eq!(m.command_executions_started.get(), 0);
assert_eq!(m.command_output_chunks_appended.get(), 0);
assert_eq!(m.command_output_chunks_dropped_as_stale.get(), 0);
}
fn tool_call_delta(
request: RequestId,
index: u64,
id: Option<&str>,
function_name: Option<&str>,
arguments_fragment: Option<&str>,
) -> Event {
Event::ToolCallDelta {
request,
index,
id: id.map(str::to_string),
function_name: function_name.map(str::to_string),
arguments_fragment: arguments_fragment.map(str::to_string),
}
}
fn drive_single_tool_call(
core: &mut AgentCore,
request: RequestId,
call_id: &str,
function_name: &str,
arguments_json: &str,
) -> Vec<Action> {
let _ = core.handle_event(tool_call_delta(
request,
0,
Some(call_id),
Some(function_name),
Some(arguments_json),
));
core.handle_event(Event::Finish {
request,
reason: Some("tool_calls".to_string()),
})
}
#[test]
fn tool_call_delta_merges_fragments_and_emits_execute_tool_on_finish() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(tool_call_delta(
id,
0,
Some("call_1"),
Some("read"),
Some(r#"{"path":""#),
));
let _ = core.handle_event(tool_call_delta(id, 0, None, None, Some(r#"src/main.rs"}"#)));
assert_eq!(core.metrics().tool_call_deltas_appended.get(), 2);
let actions = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
assert_eq!(core.status(), Status::ToolRunning);
assert_eq!(core.metrics().tool_calls_executed.get(), 1);
assert!(actions.iter().any(|a| matches!(
a,
Action::ExecuteTool {
call_id, invocation: ReadOnlyTool::Read { path, .. }, ..
} if call_id == "call_1" && path == "src/main.rs"
)));
match &core.conversation()[1] {
ChatMessage::Assistant { tool_calls, .. } => {
assert_eq!(tool_calls.len(), 1);
assert_eq!(tool_calls[0].function_name, "read");
}
other => panic!("expected assistant with tool_calls, got {other:?}"),
}
}
#[test]
fn multiple_parallel_tool_calls_are_ordered_by_index() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(tool_call_delta(
id,
1,
Some("b"),
Some("list"),
Some(r#"{"path":"src"}"#),
));
let _ = core.handle_event(tool_call_delta(
id,
0,
Some("a"),
Some("read"),
Some(r#"{"path":"README.md"}"#),
));
let actions = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
let execute_ids: Vec<String> = actions
.iter()
.filter_map(|a| match a {
Action::ExecuteTool { call_id, .. } => Some(call_id.clone()),
_ => None,
})
.collect();
assert_eq!(execute_ids, vec!["a".to_string(), "b".to_string()]);
}
#[test]
fn tool_result_completes_slot_and_advances_to_next_request() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_tool_call(&mut core, id, "call_1", "list", r#"{"path":"."}"#);
assert_eq!(core.status(), Status::ToolRunning);
let actions = core.handle_event(Event::ToolResult {
request: id,
call_id: "call_1".to_string(),
outcome: ToolOutcome::Ok(r#"[{"path":"a"}]"#.to_string()),
});
assert_eq!(core.metrics().tool_results_committed.get(), 1);
let start = actions
.iter()
.find_map(|a| match a {
Action::StartRequest { id, messages } => Some((*id, messages.clone())),
_ => None,
})
.expect("StartRequest emitted");
assert_ne!(start.0, id);
assert_eq!(start.1.len(), 3);
assert!(matches!(start.1[2], ChatMessage::Tool { .. }));
assert_eq!(core.status(), Status::AwaitingModel);
}
#[test]
fn tool_result_before_all_arrive_stays_in_tool_running() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(tool_call_delta(
id,
0,
Some("a"),
Some("read"),
Some(r#"{"path":"a"}"#),
));
let _ = core.handle_event(tool_call_delta(
id,
1,
Some("b"),
Some("read"),
Some(r#"{"path":"b"}"#),
));
let _ = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
let actions = core.handle_event(Event::ToolResult {
request: id,
call_id: "a".to_string(),
outcome: ToolOutcome::Ok("ok".to_string()),
});
assert_eq!(core.status(), Status::ToolRunning);
assert_eq!(actions, vec![Action::Redraw]);
}
#[test]
fn tool_result_with_unknown_call_id_is_dropped_as_stale() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_tool_call(&mut core, id, "call_1", "list", r#"{"path":"."}"#);
let actions = core.handle_event(Event::ToolResult {
request: id,
call_id: "does_not_exist".to_string(),
outcome: ToolOutcome::Ok("ok".to_string()),
});
assert!(actions.is_empty());
assert_eq!(core.metrics().tool_results_dropped_as_stale.get(), 1);
}
#[test]
fn arguments_over_limit_yields_synthetic_arguments_too_large_err() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(tool_call_delta(
id,
0,
Some("call_big"),
Some("read"),
Some(&"x".repeat(ARGUMENTS_MAX_BYTES + 1)),
));
assert_eq!(
core.metrics()
.tool_call_arguments_fragments_dropped_over_limit
.get(),
1
);
let actions = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
assert_eq!(
core.metrics().tool_calls_rejected_by_arguments_limit.get(),
1
);
assert_eq!(core.metrics().tool_calls_executed.get(), 0);
let messages = actions
.iter()
.find_map(|a| match a {
Action::StartRequest { messages, .. } => Some(messages.clone()),
_ => None,
})
.expect("StartRequest emitted");
let tool_msg = messages
.iter()
.find_map(|m| match m {
ChatMessage::Tool { content, .. } => Some(content),
_ => None,
})
.expect("Tool message present");
assert!(tool_msg.contains("arguments_too_large"));
}
#[test]
fn turn_tool_call_limit_forces_synthetic_err_for_extras() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
for i in 0..=TURN_TOOL_CALL_LIMIT as u64 {
let call_id = format!("c{i}");
let _ = core.handle_event(tool_call_delta(
id,
i,
Some(&call_id),
Some("list"),
Some(r#"{"path":"."}"#),
));
}
let _ = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
assert_eq!(
core.metrics().tool_calls_executed.get(),
TURN_TOOL_CALL_LIMIT as u64
);
assert_eq!(core.metrics().tool_calls_rejected_by_turn_limit.get(), 1);
}
#[test]
fn unknown_function_name_yields_synthetic_unknown_tool_err() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = drive_single_tool_call(&mut core, id, "call_1", "bogus", r#"{}"#);
let messages = actions
.iter()
.find_map(|a| match a {
Action::StartRequest { messages, .. } => Some(messages.clone()),
_ => None,
})
.expect("StartRequest emitted");
let tool_content = messages
.iter()
.find_map(|m| match m {
ChatMessage::Tool { content, .. } => Some(content.clone()),
_ => None,
})
.expect("Tool message present");
assert!(tool_content.contains("unknown_tool"));
assert_eq!(core.metrics().tool_calls_executed.get(), 0);
}
#[test]
fn cancel_in_tool_running_emits_cancel_tool_execution() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_tool_call(&mut core, id, "call_1", "list", r#"{"path":"."}"#);
assert_eq!(core.status(), Status::ToolRunning);
let actions = core.handle_event(Event::Cancel);
assert!(actions.contains(&Action::CancelToolExecution { request: id }));
assert_eq!(core.status(), Status::Idle);
assert!(core.pending_response().is_none());
}
#[test]
fn transport_error_dropped_in_tool_running_phase() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_tool_call(&mut core, id, "call_1", "list", r#"{"path":"."}"#);
let actions = core.handle_event(Event::TransportError {
request: id,
message: "should be ignored".to_string(),
});
assert!(actions.is_empty());
assert_eq!(core.metrics().transport_errors_dropped_as_stale.get(), 1);
assert_eq!(core.status(), Status::ToolRunning);
}
#[test]
fn tool_result_dropped_in_streaming_phase() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = core.handle_event(Event::ToolResult {
request: id,
call_id: "call_1".to_string(),
outcome: ToolOutcome::Ok("x".to_string()),
});
assert!(actions.is_empty());
assert_eq!(core.metrics().tool_results_dropped_as_stale.get(), 1);
}
#[test]
fn tool_calls_this_turn_resets_between_user_turns() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_tool_call(&mut core, id, "call_1", "list", r#"{"path":"."}"#);
let next_id = {
let actions = core.handle_event(Event::ToolResult {
request: id,
call_id: "call_1".to_string(),
outcome: ToolOutcome::Ok("ok".to_string()),
});
last_start_id(&actions)
};
let _ = core.handle_event(Event::Finish {
request: next_id,
reason: Some("stop".to_string()),
});
assert_eq!(core.status(), Status::Idle);
let id2 = last_start_id(&user(&mut core, "second"));
let _ = drive_single_tool_call(&mut core, id2, "c2", "list", r#"{"path":"."}"#);
assert_eq!(core.metrics().tool_calls_executed.get(), 2);
assert_eq!(core.metrics().tool_calls_rejected_by_turn_limit.get(), 0);
}
#[test]
fn finish_with_tool_calls_reason_but_no_slots_falls_through_to_idle() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
assert_eq!(actions, vec![Action::Redraw]);
assert_eq!(core.status(), Status::Idle);
assert_eq!(core.conversation().len(), 2);
}
#[test]
fn readonly_tool_definitions_expose_the_three_tools() {
let defs = ReadOnlyTool::definitions();
let names: Vec<&str> = defs.iter().map(|d| d.name.as_str()).collect();
assert_eq!(names, vec!["list", "read", "search"]);
}
#[test]
fn readonly_tool_parse_list_defaults() {
let tool = ReadOnlyTool::parse("list", r#"{"path":"src"}"#).expect("parses");
assert_eq!(
tool,
ReadOnlyTool::List {
path: "src".to_string(),
recursive: false,
max_entries: DEFAULT_LIST_MAX_ENTRIES,
include_hidden: false,
}
);
}
#[test]
fn readonly_tool_parse_read_with_line_range() {
let tool = ReadOnlyTool::parse("read", r#"{"path":"src/main.rs","line_range":[10,20]}"#)
.expect("parses");
assert_eq!(
tool,
ReadOnlyTool::Read {
path: "src/main.rs".to_string(),
line_range: Some((10, 20)),
}
);
}
#[test]
fn readonly_tool_parse_search_with_defaults() {
let tool = ReadOnlyTool::parse("search", r#"{"pattern":"TODO"}"#).expect("parses");
assert_eq!(
tool,
ReadOnlyTool::Search {
pattern: "TODO".to_string(),
path_prefix: None,
case_sensitive: false,
max_results: DEFAULT_SEARCH_MAX_RESULTS,
}
);
}
#[test]
fn readonly_tool_parse_unknown_function_name() {
let err = ReadOnlyTool::parse("foo", "{}").expect_err("unknown");
assert_eq!(err, ToolExecutionError::UnknownTool);
}
#[test]
fn readonly_tool_parse_invalid_json_is_reported() {
let err = ReadOnlyTool::parse("list", "not-json").expect_err("invalid");
assert!(matches!(err, ToolExecutionError::ArgumentsParseFailed(_)));
}
#[test]
fn tool_execution_error_to_json_includes_code_and_message() {
let s = ToolExecutionError::OutsideWorkspace.to_json_string();
assert!(s.contains(r#""error":"outside_workspace""#));
assert!(s.contains(r#""message""#));
}
#[test]
fn tool_execution_error_message_is_human_readable_and_matches_json() {
let err = ToolExecutionError::OutsideWorkspace;
assert_eq!(err.message(), "path escapes the workspace root");
let json = err.to_json_string();
assert!(json.contains(err.message().as_str()), "got {json}");
let parse = ToolExecutionError::ArgumentsParseFailed("bad shape".to_string());
assert_eq!(parse.message(), "bad shape");
assert!(!parse.message().contains("ArgumentsParseFailed"));
}
#[test]
fn patch_definition_advertises_the_patch_function_name() {
let def = PatchInvocation::definition();
assert_eq!(def.name, "patch");
assert!(def.description.contains("approval"));
assert!(def.parameters_json.contains("edits"));
}
#[test]
fn patch_parse_add_and_update_edits() {
let inv = PatchInvocation::parse(
r#"{"edits":[
{"kind":"add","path":"a.txt","content":"hello"},
{"kind":"update","path":"b.txt","before":"foo","after":"bar"}
]}"#,
)
.expect("parses");
assert_eq!(inv.edits.len(), 2);
assert_eq!(
inv.edits[0],
PatchTool::Add {
path: "a.txt".to_string(),
content: "hello".to_string(),
}
);
assert_eq!(
inv.edits[1],
PatchTool::Update {
path: "b.txt".to_string(),
before: "foo".to_string(),
after: "bar".to_string(),
}
);
}
#[test]
fn patch_parse_rejects_same_path_twice() {
let err = PatchInvocation::parse(
r#"{"edits":[
{"kind":"update","path":"dup","before":"a","after":"b"},
{"kind":"update","path":"dup","before":"c","after":"d"}
]}"#,
)
.expect_err("same-path rejected");
assert_eq!(
err,
ToolExecutionError::Patch(PatchError::MultipleEditsSamePath {
path: "dup".to_string(),
})
);
}
#[test]
fn patch_parse_rejects_empty_edits() {
let err = PatchInvocation::parse(r#"{"edits":[]}"#).expect_err("empty rejected");
assert!(matches!(err, ToolExecutionError::ArgumentsParseFailed(_)));
}
#[test]
fn patch_parse_rejects_too_many_edits() {
let mut edits = String::from("[");
for i in 0..(PATCH_MAX_EDITS + 1) {
if i > 0 {
edits.push(',');
}
edits.push_str(&format!(r#"{{"kind":"add","path":"f{i}","content":""}}"#));
}
edits.push(']');
let json = format!(r#"{{"edits":{edits}}}"#);
let err = PatchInvocation::parse(&json).expect_err("too many rejected");
assert!(matches!(
err,
ToolExecutionError::Patch(PatchError::TooManyEdits { .. })
));
}
#[test]
fn patch_parse_rejects_unknown_kind() {
let err = PatchInvocation::parse(r#"{"edits":[{"kind":"delete","path":"x"}]}"#)
.expect_err("unknown kind rejected");
assert!(matches!(err, ToolExecutionError::ArgumentsParseFailed(_)));
}
#[test]
fn patch_parse_rejects_add_over_size_limit() {
let big = "x".repeat(PATCH_MAX_FILE_BYTES + 1);
let json = format!(
r#"{{"edits":[{{"kind":"add","path":"big","content":{}}}]}}"#,
nojson::Json(&big),
);
let err = PatchInvocation::parse(&json).expect_err("too big rejected");
assert!(matches!(
err,
ToolExecutionError::Patch(PatchError::FileTooLarge { .. })
));
}
#[test]
fn patch_error_json_encodes_code_and_message() {
let e = ToolExecutionError::Patch(PatchError::NoMatch {
path: "a.txt".to_string(),
});
let s = e.to_json_string();
assert!(s.contains(r#""error":"patch_no_match""#), "got {s}");
assert!(s.contains("a.txt"));
}
#[test]
fn patch_error_json_includes_hint_for_same_path() {
let e = ToolExecutionError::Patch(PatchError::MultipleEditsSamePath {
path: "dup".to_string(),
});
let s = e.to_json_string();
assert!(
s.contains(r#""error":"patch_multiple_edits_same_path""#),
"got {s}"
);
assert!(s.contains(r#""hint":"#), "expected a hint member in {s}");
assert!(s.contains("separate patch calls"), "got {s}");
}
#[test]
fn patch_error_json_emits_null_hint_when_not_actionable() {
let e = ToolExecutionError::Patch(PatchError::Rejected);
let s = e.to_json_string();
assert!(s.contains(r#""error":"patch_rejected""#), "got {s}");
assert!(
s.contains(r#""hint":null"#),
"Rejected should emit null hint: {s}"
);
}
#[test]
fn command_definition_advertises_the_command_function_name() {
let def = CommandInvocation::definition();
assert_eq!(def.name, "command");
assert!(def.description.contains("approval"));
assert!(def.parameters_json.contains("argv"));
assert!(!def.parameters_json.contains("timeout_seconds"));
}
#[test]
fn command_parse_extracts_argv() {
let inv = CommandInvocation::parse(r#"{"argv":["ls","-la"]}"#).expect("parses");
assert_eq!(inv.argv, vec!["ls".to_string(), "-la".to_string()]);
}
#[test]
fn command_parse_ignores_stray_timeout_seconds() {
let inv = CommandInvocation::parse(r#"{"argv":["cargo","test"],"timeout_seconds":120}"#)
.expect("parses");
assert_eq!(inv.argv, vec!["cargo".to_string(), "test".to_string()]);
}
#[test]
fn command_parse_rejects_empty_argv() {
let err = CommandInvocation::parse(r#"{"argv":[]}"#).expect_err("empty rejected");
assert_eq!(err, ToolExecutionError::Command(CommandError::EmptyArgv));
}
#[test]
fn command_error_json_encodes_code_and_message() {
let e = ToolExecutionError::Command(CommandError::SpawnFailed {
message: "no such file".to_string(),
});
let s = e.to_json_string();
assert!(s.contains(r#""error":"command_spawn_failed""#), "got {s}");
assert!(s.contains("no such file"));
}
fn drive_single_patch_call(
core: &mut AgentCore,
request: RequestId,
call_id: &str,
arguments_json: &str,
) -> Vec<Action> {
let _ = core.handle_event(tool_call_delta(
request,
0,
Some(call_id),
Some("patch"),
Some(arguments_json),
));
core.handle_event(Event::Finish {
request,
reason: Some("tool_calls".to_string()),
})
}
fn valid_add_patch_json() -> &'static str {
r#"{"edits":[{"kind":"add","path":"new.txt","content":"hi"}]}"#
}
#[test]
fn patch_finish_emits_preview_patch_and_stays_in_tool_running() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let actions = drive_single_patch_call(&mut core, id, "p1", valid_add_patch_json());
assert_eq!(core.status(), Status::ToolRunning);
assert!(actions.iter().any(|a| matches!(
a,
Action::PreviewPatch { call_id, .. } if call_id == "p1"
)));
assert_eq!(core.metrics().patch_calls_previewed.get(), 1);
assert_eq!(core.pending_approval_call_id(), None);
}
#[test]
fn patch_preview_ready_transitions_to_awaiting_approval_and_publishes_call_id() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_patch_call(&mut core, id, "p1", valid_add_patch_json());
let snapshot = PreviewContent {
path: "new.txt".to_string(),
content: None,
};
let preview = PatchPreview {
target_paths: vec!["new.txt".to_string()],
added_lines: 1,
removed_lines: 0,
edit_count: 1,
auto_approve: false,
not_revertible: None,
};
let actions = core.handle_event(Event::PatchPreviewReady {
request: id,
call_id: "p1".to_string(),
preview_content: vec![snapshot],
preview,
});
assert_eq!(actions, vec![Action::Redraw]);
assert_eq!(core.status(), Status::AwaitingApproval);
assert_eq!(core.pending_approval_call_id(), Some("p1".to_string()));
assert_eq!(core.metrics().patch_previews_committed.get(), 1);
let active = core.active_tool_calls();
assert_eq!(active.len(), 1);
assert_eq!(active[0].approval, ApprovalState::Pending);
assert!(active[0].patch_preview.is_some());
}
#[test]
fn patch_approve_emits_apply_patch_with_stored_snapshots() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_patch_call(&mut core, id, "p1", valid_add_patch_json());
let snapshot = PreviewContent {
path: "new.txt".to_string(),
content: Some(b"before".to_vec()),
};
let _ = core.handle_event(Event::PatchPreviewReady {
request: id,
call_id: "p1".to_string(),
preview_content: vec![snapshot.clone()],
preview: PatchPreview::default(),
});
let actions = core.handle_event(Event::ApproveToolCall {
call_id: "p1".to_string(),
});
let apply = actions
.iter()
.find_map(|a| match a {
Action::ApplyPatch {
call_id,
preview_content,
..
} => Some((call_id.clone(), preview_content.clone())),
_ => None,
})
.expect("ApplyPatch emitted");
assert_eq!(apply.0, "p1");
assert_eq!(apply.1, vec![snapshot]);
assert_eq!(core.metrics().tool_call_approvals_committed.get(), 1);
assert_eq!(core.status(), Status::ToolRunning);
assert!(core.pending_approval_call_id().is_none());
}
#[test]
fn patch_reject_synthesizes_err_and_advances_when_last() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_patch_call(&mut core, id, "p1", valid_add_patch_json());
let _ = core.handle_event(Event::PatchPreviewReady {
request: id,
call_id: "p1".to_string(),
preview_content: Vec::new(),
preview: PatchPreview::default(),
});
let actions = core.handle_event(Event::RejectToolCall {
call_id: "p1".to_string(),
});
assert!(
actions
.iter()
.any(|a| matches!(a, Action::StartRequest { .. }))
);
assert_eq!(core.metrics().tool_call_rejections_committed.get(), 1);
let last = core.conversation().last().expect("has tool message");
match last {
ChatMessage::Tool { content, .. } => {
assert!(content.contains("patch_rejected"), "content={content}");
}
other => panic!("expected Tool, got {other:?}"),
}
}
#[test]
fn patch_approve_dropped_if_no_pending_call() {
let mut core = AgentCore::new();
let _ = user(&mut core, "hi");
let actions = core.handle_event(Event::ApproveToolCall {
call_id: "nope".to_string(),
});
assert!(actions.is_empty());
assert_eq!(core.metrics().tool_call_approvals_dropped_as_stale.get(), 1);
}
#[test]
fn patch_preview_ready_dropped_if_call_id_unknown() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_patch_call(&mut core, id, "p1", valid_add_patch_json());
let actions = core.handle_event(Event::PatchPreviewReady {
request: id,
call_id: "does_not_exist".to_string(),
preview_content: Vec::new(),
preview: PatchPreview::default(),
});
assert!(actions.is_empty());
assert_eq!(core.metrics().patch_previews_dropped_as_stale.get(), 1);
}
#[test]
fn read_only_result_lands_in_awaiting_approval_phase() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = core.handle_event(tool_call_delta(
id,
0,
Some("r1"),
Some("list"),
Some(r#"{"path":"."}"#),
));
let _ = core.handle_event(tool_call_delta(
id,
1,
Some("p1"),
Some("patch"),
Some(valid_add_patch_json()),
));
let _ = core.handle_event(Event::Finish {
request: id,
reason: Some("tool_calls".to_string()),
});
let _ = core.handle_event(Event::PatchPreviewReady {
request: id,
call_id: "p1".to_string(),
preview_content: Vec::new(),
preview: PatchPreview::default(),
});
assert_eq!(core.status(), Status::AwaitingApproval);
let actions = core.handle_event(Event::ToolResult {
request: id,
call_id: "r1".to_string(),
outcome: ToolOutcome::Ok("ok".to_string()),
});
assert_eq!(core.metrics().tool_results_committed.get(), 1);
assert_eq!(actions, vec![Action::Redraw]);
assert_eq!(core.status(), Status::AwaitingApproval);
}
fn drive_single_command_call(
core: &mut AgentCore,
request: RequestId,
call_id: &str,
arguments_json: &str,
) -> Vec<Action> {
let _ = core.handle_event(tool_call_delta(
request,
0,
Some(call_id),
Some("command"),
Some(arguments_json),
));
core.handle_event(Event::Finish {
request,
reason: Some("tool_calls".to_string()),
})
}
fn valid_command_json() -> &'static str {
r#"{"argv":["echo","hi"]}"#
}
#[test]
fn command_finish_populates_command_preview_and_marks_pending() {
let mut core = AgentCore::new();
core.set_workspace_display("/tmp/wksp".to_string());
let id = last_start_id(&user(&mut core, "hi"));
let actions = drive_single_command_call(&mut core, id, "c1", valid_command_json());
assert!(
!actions
.iter()
.any(|a| matches!(a, Action::ExecuteCommand { .. }))
);
assert_eq!(core.status(), Status::AwaitingApproval);
assert_eq!(core.metrics().command_calls_dispatched.get(), 1);
let active = core.active_tool_calls();
assert_eq!(active.len(), 1);
assert_eq!(active[0].approval, ApprovalState::Pending);
let preview = active[0]
.command_preview
.as_ref()
.expect("preview populated");
assert_eq!(preview.argv, vec!["echo".to_string(), "hi".to_string()]);
assert_eq!(preview.working_directory, "/tmp/wksp");
assert_eq!(core.pending_approval_call_id(), Some("c1".to_string()));
}
#[test]
fn command_approve_emits_execute_command_action() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_command_call(&mut core, id, "c1", valid_command_json());
let actions = core.handle_event(Event::ApproveToolCall {
call_id: "c1".to_string(),
});
assert_eq!(core.metrics().tool_call_approvals_committed.get(), 1);
assert_eq!(core.metrics().command_executions_started.get(), 1);
let matched = actions.iter().any(|a| {
matches!(
a,
Action::ExecuteCommand { call_id, invocation, .. }
if call_id == "c1"
&& invocation.argv == vec!["echo".to_string(), "hi".to_string()]
)
});
assert!(matched, "expected ExecuteCommand, got {actions:?}");
assert_eq!(core.status(), Status::ToolRunning);
}
#[test]
fn command_reject_synthesizes_command_rejected_err() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_command_call(&mut core, id, "c1", valid_command_json());
let actions = core.handle_event(Event::RejectToolCall {
call_id: "c1".to_string(),
});
assert_eq!(core.metrics().tool_call_rejections_committed.get(), 1);
let start = actions
.iter()
.find_map(|a| match a {
Action::StartRequest { messages, .. } => Some(messages.clone()),
_ => None,
})
.expect("StartRequest emitted");
let tool_msg = start
.iter()
.rev()
.find_map(|m| match m {
ChatMessage::Tool { content, .. } => Some(content.clone()),
_ => None,
})
.expect("Tool message present");
assert!(tool_msg.contains("command_rejected"), "content={tool_msg}");
}
#[test]
fn command_output_chunk_appends_to_tail_when_approved() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_command_call(&mut core, id, "c1", valid_command_json());
let _ = core.handle_event(Event::ApproveToolCall {
call_id: "c1".to_string(),
});
let _ = core.handle_event(Event::CommandOutputChunk {
request: id,
call_id: "c1".to_string(),
stream: CommandOutputStream::Stdout,
bytes: b"line 1\n".to_vec(),
});
let _ = core.handle_event(Event::CommandOutputChunk {
request: id,
call_id: "c1".to_string(),
stream: CommandOutputStream::Stderr,
bytes: b"warn\n".to_vec(),
});
assert_eq!(core.metrics().command_output_chunks_appended.get(), 2);
let active = core.active_tool_calls();
let tail = active[0]
.command_output_tail
.as_ref()
.expect("tail populated");
assert_eq!(tail.stdout_tail, "line 1\n");
assert_eq!(tail.stderr_tail, "warn\n");
assert_eq!(tail.stdout_bytes_total, 7);
assert_eq!(tail.stderr_bytes_total, 5);
}
#[test]
fn command_output_chunk_dropped_if_not_approved() {
let mut core = AgentCore::new();
let id = last_start_id(&user(&mut core, "hi"));
let _ = drive_single_command_call(&mut core, id, "c1", valid_command_json());
let actions = core.handle_event(Event::CommandOutputChunk {
request: id,
call_id: "c1".to_string(),
stream: CommandOutputStream::Stdout,
bytes: b"early".to_vec(),
});
assert!(actions.is_empty());
assert_eq!(
core.metrics().command_output_chunks_dropped_as_stale.get(),
1
);
}
}