diff --git a/crates/tinyagents-harness/src/artifacts/README.md b/crates/tinyagents-harness/src/artifacts/README.md index f7c086a6c..7380fea35 100644 --- a/crates/tinyagents-harness/src/artifacts/README.md +++ b/crates/tinyagents-harness/src/artifacts/README.md @@ -66,6 +66,8 @@ already had. | `types.rs` | `ArtifactKind`, `OffloadedArtifact`, `OffloadError`, constants. | | `paths.rs` | Fail-closed path resolution (`resolve_artifact_path` and helpers). | | `policy.rs` | `ArtifactPathPolicy`, `ArtifactRedactor` traits and their defaults. | +| `tool_results.rs` | Per-tool-result persistence: `ToolResultArtifactStore`, `apply_per_result_persistence`, `spill_aggregate_tool_results`, paged artifact reads. Host supplies redactor, read/wrapper tool names, read limit. | +| `tool_results_test.rs` | Tests for the above, incl. a byte-stable envelope fixture. | | `ops.rs` | `ArtifactOffload` writer, `offload_oversized_result`, pointer/handoff plumbing. | | `test.rs` | Unit tests: happy path, fallback path, fail-closed hardening. | diff --git a/crates/tinyagents-harness/src/artifacts/mod.rs b/crates/tinyagents-harness/src/artifacts/mod.rs index ea0a868db..c6f2a5771 100644 --- a/crates/tinyagents-harness/src/artifacts/mod.rs +++ b/crates/tinyagents-harness/src/artifacts/mod.rs @@ -62,6 +62,7 @@ mod ops; mod paths; pub mod policy; +pub mod tool_results; mod types; pub use ops::{ diff --git a/crates/tinyagents-harness/src/artifacts/tool_results.rs b/crates/tinyagents-harness/src/artifacts/tool_results.rs new file mode 100644 index 000000000..14dbcd3da --- /dev/null +++ b/crates/tinyagents-harness/src/artifacts/tool_results.rs @@ -0,0 +1,774 @@ +//! Persist oversized tool outputs as action-workspace artifacts. +//! +//! Tool results enter the model context before the provider has seen them, so +//! this is the last cheap point to replace large raw output with a bounded +//! preview. The full, redacted body is written under `action_dir` so the +//! host's file-reading tool can inspect it later without exposing internal +//! host state. +//! +//! Distinct from [`offload_oversized_result`](super::offload_oversized_result), +//! which offloads a *worker's final result* under `outputs/`: this handles each +//! *tool result* mid-turn under `artifacts/tool-results///`. +//! +//! What the host supplies, because the crate cannot decide it: the +//! [`ArtifactRedactor`] every body passes through before it is stored, the name +//! of its file-reading tool, the wrapper tool a read may arrive inside, and the +//! largest body its reader will open. + +use std::path::{Path, PathBuf}; +use std::sync::Arc; + +use serde_json::Value; + +use super::policy::{ArtifactRedactor, Redacted}; + +const ARTIFACT_ROOT: &str = "artifacts/tool-results"; + +/// A read of a persisted artifact, recognised from a tool call's arguments. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ArtifactRead { + pub path: String, + pub offset: usize, +} + +/// The artifact a tool call reads, if any. +/// +/// `read_tool` and `wrapper_tool` are the host's vocabulary, passed in because +/// tool names are not the crate's to decide (OpenHuman passes `file_read` and +/// `use_skill`). +/// +/// Only a `read_tool` call counts, because only its result *is* the stored body: +/// `file_write`, `glob`, `list` or `apply_patch` can name an artifact path too, +/// and their output must still take the normal ladder. The read may arrive +/// wrapped in `wrapper_tool`, reported under that name +/// (`use_skill {"skill":"files","tool":"file_read","args":{"path":…}}`), so +/// `wrapper_tool` — and only it, the one tool whose result *is* the +/// wrapped tool's result — is followed into the tool it runs (#6284). Any +/// other tool that happens to carry `tool`/`args` fields is not a wrapper. +pub fn artifact_read_target( + tool_name: &str, + args: &Value, + read_tool: &str, + wrapper_tool: &str, +) -> Option { + if tool_name == wrapper_tool { + let inner_tool = args.get("tool").and_then(Value::as_str)?; + return artifact_read_target(inner_tool, args.get("args")?, read_tool, wrapper_tool); + } + if tool_name != read_tool { + return None; + } + let path = args.get("path").and_then(Value::as_str)?; + // A path component match: `artifacts/tool-results-backup/…` shares the + // prefix but is not the artifact directory. + let under_root = path + .trim_start_matches("./") + .strip_prefix(ARTIFACT_ROOT) + .is_some_and(|rest| rest.starts_with('/')); + // An absent or null offset starts at 0. A present one that is not a + // non-negative integer that fits `usize` is not a read `file_read` serves + // (it rejects it), so it is not an artifact read either; never reinterpret + // it as 0. + let offset = match args.get("offset") { + None | Some(Value::Null) => 0, + Some(value) => usize::try_from(value.as_u64()?).ok()?, + }; + under_root.then(|| ArtifactRead { + path: path.to_string(), + offset, + }) +} + +/// Bound one page of an artifact read to `budget_bytes`, naming the exact +/// `offset` the next read continues from. Never persists: re-persisting a read +/// of an artifact creates a new artifact whose preview is the same bounded +/// head, and the model can loop between previews without ever reaching the +/// body. +pub fn page_artifact_read( + content: String, + read: &ArtifactRead, + budget_bytes: usize, + read_tool: &str, +) -> String { + if budget_bytes == 0 { + return content; + } + // A page is useless without its continuation, so a budget too small to + // carry one is raised to the floor a persisted envelope already takes for + // the same reason (`MIN_ENVELOPE_ALLOWANCE_BYTES`). Every page therefore + // fits `max(budget_bytes, MIN_ENVELOPE_ALLOWANCE_BYTES)` and advances. + let budget_bytes = budget_bytes.max(MIN_ENVELOPE_ALLOWANCE_BYTES); + if content.len() <= budget_bytes { + return content; + } + let start = read.offset; + let Some(total) = start.checked_add(content.len()) else { + // Only reachable with an offset no real read carries (`file_read` + // rejects offsets past its at-most-10-MiB file). Bound the result but + // advertise no continuation, since none could advance. + let cut = floor_char_boundary(&content, budget_bytes); + return content[..cut].to_string(); + }; + let with_path = |next: usize| { + format!( + "\n\n[artifact page: bytes {start}..{next} of {total}. Continue with {read_tool} {{\"path\":\"{}\",\"offset\":{next}}}]", + read.path + ) + }; + // Without the path (the caller already has it). At most ~100 bytes, so it + // always leaves body room under the floor. + let without_path = |next: usize| { + format!( + "\n\n[artifact page: bytes {start}..{next} of {total}. Continue with {read_tool} at \"offset\":{next}]" + ) + }; + // Sized from the trailer this page will actually carry, not a fixed + // reservation. `next` never has more digits than `total`, so a trailer + // rendered with `total` is its longest form. + let use_path = with_path(total).len() + 4 <= budget_bytes; + let longest = if use_path { + with_path(total).len() + } else { + without_path(total).len() + }; + let cut = floor_char_boundary(&content, budget_bytes - longest); + let trailer = if use_path { + with_path(start + cut) + } else { + without_path(start + cut) + }; + format!("{}{trailer}", &content[..cut]) +} +const AGGREGATE_PREVIEW_BUDGET_BYTES: usize = 512; +/// #4469 item 6: floor for how tightly a persisted `[tool_result_preview]` +/// envelope may be bounded during aggregate spill. `allowed_len` can saturate to +/// `0` (or a handful of bytes) once earlier-spilled results have already consumed +/// the aggregate budget; bounding the envelope to that would return `""` — or a +/// header cut mid-line — discarding the `artifact_path` pointer the model needs +/// to `file_read` the full output. This floor keeps the envelope header (through +/// the `artifact_path` / `read_with` lines) intact even when the raw budget math +/// says zero; `apply_tool_result_budget` retains the head, so the pointer always +/// survives. Slightly overshooting the aggregate budget here is the correct +/// trade — a valid pointer is worth a few hundred bytes. +const MIN_ENVELOPE_ALLOWANCE_BYTES: usize = 512; +const TRAILER_RESERVED: usize = 256; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +struct BudgetOutcome { + original_bytes: usize, + final_bytes: usize, + truncated: bool, +} + +impl BudgetOutcome { + fn unchanged(len: usize) -> Self { + Self { + original_bytes: len, + final_bytes: len, + truncated: false, + } + } +} + +fn apply_tool_result_budget(content: String, budget_bytes: usize) -> (String, BudgetOutcome) { + let original_bytes = content.len(); + if budget_bytes == 0 || original_bytes <= budget_bytes { + return (content, BudgetOutcome::unchanged(original_bytes)); + } + + let head_capacity = budget_bytes.saturating_sub(TRAILER_RESERVED).max(1); + let mut cut = floor_char_boundary(&content, head_capacity); + if cut == 0 { + cut = content + .char_indices() + .next() + .map(|(_, c)| c.len_utf8()) + .unwrap_or(0); + } + + let dropped_bytes = original_bytes.saturating_sub(cut); + let mut out = String::with_capacity(cut + TRAILER_RESERVED); + out.push_str(&content[..cut]); + // Say what was lost, and say what NOT to do about it (#6408). + // + // The previous wording ended "re-run with a narrower query to see the + // rest", which is an instruction to repeat the call. For a listing tool + // with no narrowing argument in reach — `GITHUB_LIST_PULL_REQUESTS` is the + // reported case — the model has nothing to narrow, so it re-issues the + // identical call, gets the identical truncation, and repeats until the + // successful-repeat tracker halts the run: 8-14 calls per question, the + // token and quota burn, and an "Incomplete" the user cannot explain. + // + // Truncation here is deterministic: the same call returns the same bytes + // and the same cut. Saying so is what makes the retry stop, and giving the + // totals lets the model judge whether the head it kept is enough to answer + // from. Keep the `truncated by tool_result_budget` phrase — it is the + // grep handle several suites and the runbooks match on. + out.push_str(&format!( + "\n\n[… {dropped_bytes} of {original_bytes} bytes truncated by tool_result_budget. \ + Repeating this call returns the same truncation — narrow the request \ + (filter, paginate, or request a smaller range) or answer from the {cut} bytes above …]" + )); + + let final_bytes = out.len(); + ( + out, + BudgetOutcome { + original_bytes, + final_bytes, + truncated: true, + }, + ) +} + +/// Writes oversized tool results under `/artifacts/tool-results/`. +/// +/// The host supplies what the crate cannot decide: the redactor every body +/// passes through before it is stored (an artifact on disk is exactly as +/// readable as the result it replaces), the name of its file-reading tool +/// (quoted in the envelope so the model knows how to open the file), and the +/// largest body that tool will open. +#[derive(Debug, Clone)] +pub struct ToolResultArtifactStore { + action_dir: PathBuf, + session_key: String, + redactor: Arc, + read_tool: String, + max_readable_bytes: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PersistedToolResult { + pub output: String, + pub path: String, + pub original_bytes: usize, + pub stored_bytes: usize, + /// Whether the redactor rewrote anything before the body was stored. + pub redacted: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ToolResultArtifactOutcome { + pub original_bytes: usize, + pub final_bytes: usize, + pub persisted: bool, + pub artifact_path: Option, +} + +impl ToolResultArtifactOutcome { + pub fn unchanged(len: usize) -> Self { + Self { + original_bytes: len, + final_bytes: len, + persisted: false, + artifact_path: None, + } + } +} + +impl ToolResultArtifactStore { + /// `read_tool` is the host's file-reading tool and `max_readable_bytes` the + /// largest body it will open; a body whose redacted form exceeds it is + /// never stored, because the model could not read it back. + pub fn new( + action_dir: PathBuf, + session_key: impl Into, + redactor: Arc, + read_tool: impl Into, + max_readable_bytes: u64, + ) -> Self { + Self { + action_dir, + session_key: sanitize_component(&session_key.into()), + redactor, + read_tool: read_tool.into(), + max_readable_bytes, + } + } + + /// The root artifacts are written under. A caller choosing the wrong root + /// produces a pointer the model cannot dereference, and that is only + /// assertable from outside (#6483). + pub fn root(&self) -> &Path { + &self.action_dir + } + + /// The host's file-reading tool, as quoted in every envelope. + pub fn read_tool(&self) -> &str { + &self.read_tool + } + + /// Delete artifact directories for sessions other than this one that have + /// not been touched within `max_age`. + /// + /// Artifacts are written under the user's action workspace and nothing else + /// removes them, so without a bound the directory grows for the life of the + /// install — a connector returning 17-65 KB across 8-14 calls per question + /// (#6408) writes a file per call. Trading a token-burn bug for a + /// disk-growth bug is not a fix. + /// + /// Pruning on session START rather than on session end is deliberate. There + /// is no single point where a session host ends — cached agents are evicted + /// on fingerprint mismatch or poisoning, and a crash ends a session with no + /// hook at all — so an end-of-session sweep would miss exactly the runs most + /// likely to have left artifacts behind. An age sweep at start is + /// self-healing instead: whatever the last run did, the next one tidies it. + /// + /// The current session is never pruned, no matter its age, so a long-lived + /// session cannot delete artifacts the model may still be reading back. + /// + /// Best-effort by design: a failure here must not fail a turn. The caller + /// logs and carries on, because the worst case is disk left uncollected, + /// which the next session retries. + pub fn prune_stale_sessions(&self, max_age: std::time::Duration) -> std::io::Result { + let root = self.action_dir.join(ARTIFACT_ROOT); + let entries = match std::fs::read_dir(&root) { + Ok(entries) => entries, + // No artifact root yet is the common case on a first run. + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(0), + Err(error) => return Err(error), + }; + + let now = std::time::SystemTime::now(); + let mut removed = 0; + for entry in entries.flatten() { + let name = entry.file_name(); + // Never prune the directory this store is actively writing into. + if name.to_str() == Some(self.session_key.as_str()) { + continue; + } + if !entry.file_type().is_ok_and(|kind| kind.is_dir()) { + continue; + } + let stale = newest_modified(&entry.path()) + .ok() + .and_then(|modified| now.duration_since(modified).ok()) + .is_some_and(|age| age > max_age); + if stale && std::fs::remove_dir_all(entry.path()).is_ok() { + removed += 1; + } + } + Ok(removed) + } + + pub fn path_for_read_tool(&self, tool_name: &str, call_id: Option<&str>) -> String { + let call = call_id + .map(sanitize_component) + .filter(|value| !value.is_empty()) + .unwrap_or_else(random_call_id); + format!( + "{ARTIFACT_ROOT}/{}/{}/{}.txt", + self.session_key, + sanitize_component(tool_name), + call + ) + } + + /// Store `content` and return an envelope previewing it. When `content` + /// would be unreadable once sanitized (see [`readable_body`]) and a + /// `fallback` is given, the fallback is stored instead; when neither fits, + /// this errors and the caller truncates inline rather than writing an + /// artifact nobody can read. + async fn persist( + &self, + tool_name: &str, + call_id: Option<&str>, + content: &str, + fallback: Option<&str>, + preview_budget_bytes: usize, + reason: &str, + ) -> anyhow::Result { + let (content, sanitized) = readable_body( + content, + fallback, + self.max_readable_bytes, + self.redactor.as_ref(), + &self.read_tool, + )?; + let read_tool = self.read_tool.as_str(); + let relative_path = self.path_for_read_tool(tool_name, call_id); + let absolute_path = self.action_dir.join(&relative_path); + assert_within_action_dir(&self.action_dir, &absolute_path)?; + if let Some(parent) = absolute_path.parent() { + tokio::fs::create_dir_all(parent).await?; + let canonical_root = tokio::fs::canonicalize(&self.action_dir).await?; + let canonical_parent = tokio::fs::canonicalize(parent).await?; + if !canonical_parent.starts_with(&canonical_root) { + anyhow::bail!( + "tool-result artifact parent escaped action_dir: {}", + parent.display() + ); + } + } + tokio::fs::write(&absolute_path, sanitized.text.as_bytes()).await?; + + let (preview, preview_outcome) = + apply_tool_result_budget(sanitized.text.clone(), preview_budget_bytes); + let redaction_note = if sanitized.changed { + " Credential/PII redaction was applied before storage and preview exposure." + } else { + "" + }; + let truncation_note = if preview_outcome.truncated { + format!( + " Preview is bounded; {} stored bytes are available via {read_tool}.", + sanitized.text.len() + ) + } else { + String::new() + }; + + let envelope = format!( + "[tool_result_preview]\n\ + tool: {tool_name}\n\ + reason: {reason}\n\ + original_bytes: {}\n\ + stored_bytes: {}\n\ + artifact_path: {relative_path}\n\ + read_with: {read_tool} {{\"path\":\"{relative_path}\"}} (a long read returns one page and names the \"offset\" to continue from)\n\ + notes: Full scrubbed output was persisted under the action workspace.{redaction_note}{truncation_note}\n\n\ + [preview]\n{preview}", + content.len(), + sanitized.text.len(), + ); + + Ok(PersistedToolResult { + output: envelope, + path: relative_path, + original_bytes: content.len(), + stored_bytes: sanitized.text.len(), + redacted: sanitized.changed, + }) + } +} + +/// The body to store: `primary` if its sanitized form fits `limit`, else +/// `fallback` if *its* sanitized form does, else an error. Both checks are on +/// the sanitized size, the bytes actually written, because redaction can grow a +/// body (`+15551234567` becomes `[REDACTED_PII_PHONE]`): a raw body under the +/// limit can still produce an artifact the read tool refuses to open. +fn readable_body<'a>( + primary: &'a str, + fallback: Option<&'a str>, + limit: u64, + redactor: &dyn ArtifactRedactor, + read_tool: &str, +) -> anyhow::Result<(&'a str, Redacted)> { + let sanitized = redactor.redact(primary); + if sanitized.text.len() as u64 <= limit { + return Ok((primary, sanitized)); + } + if let Some(fallback) = fallback { + let sanitized_fallback = redactor.redact(fallback); + if sanitized_fallback.text.len() as u64 <= limit { + return Ok((fallback, sanitized_fallback)); + } + } + anyhow::bail!( + "tool result would not be readable once stored: {} sanitized bytes exceed the {limit}-byte {read_tool} limit", + sanitized.text.len() + ) +} + +/// Persist an over-budget result and return its envelope. +/// +/// `full_output` is the tool's output before any earlier stage rewrote it +/// (summarizer, TokenJuice). When given, *that* is what gets stored, so the +/// artifact holds what the tool returned rather than a compacted copy of it; +/// `content` still decides whether the budget was exceeded and is what the +/// model would otherwise have seen. +pub async fn apply_per_result_persistence( + content: String, + full_output: Option, + store: Option<&ToolResultArtifactStore>, + tool_name: &str, + call_id: Option<&str>, + budget_bytes: usize, +) -> (String, ToolResultArtifactOutcome) { + let original_bytes = content.len(); + if budget_bytes == 0 || original_bytes <= budget_bytes { + return ( + content, + ToolResultArtifactOutcome::unchanged(original_bytes), + ); + } + + if let Some(store) = store { + match store + .persist( + tool_name, + call_id, + full_output.as_deref().unwrap_or(&content), + full_output.as_ref().map(|_| content.as_str()), + budget_bytes, + "per-result budget exceeded", + ) + .await + { + Ok(persisted) => { + let (output, final_bytes) = bound_text_to_budget( + persisted.output, + budget_bytes.max(MIN_ENVELOPE_ALLOWANCE_BYTES), + ); + if final_bytes >= original_bytes { + // #4469 item 9: this branch does NOT fall back to inline + // truncation — the envelope is returned regardless, because it + // carries the `artifact_path` pointer to the full stored output + // (worth keeping even when the preview text nets no byte saving + // vs. the raw result). Log it as an observation only. + tracing::debug!( + "[agent][tool-result-artifacts] persisted envelope not smaller than raw result tool={} original_bytes={} final_bytes={} budget_bytes={} -- keeping envelope for its artifact_path pointer", + tool_name, + original_bytes, + final_bytes, + budget_bytes + ); + } + tracing::info!( + "[agent][tool-result-artifacts] persisted oversized tool result tool={} original_bytes={} stored_bytes={} path={} redacted={}", + tool_name, + persisted.original_bytes, + persisted.stored_bytes, + persisted.path, + persisted.redacted + ); + return ( + output, + ToolResultArtifactOutcome { + // The size of what was stored, which `full_output` can + // make larger than `content`; the artifact index and its + // contents list read this number. + original_bytes: persisted.original_bytes, + final_bytes, + persisted: true, + artifact_path: Some(persisted.path), + }, + ); + } + Err(err) => { + tracing::warn!( + "[agent][tool-result-artifacts] persist failed tool={} original_bytes={} err={} — falling back to inline truncation", + tool_name, + original_bytes, + err + ); + } + } + } + + // Reached two ways, and only one of them was audible: a persist that FAILED + // warns just above, but a run with no artifact store configured falls + // through to here silently — the oversized tail is discarded with nothing + // recording that it happened. Say so, so "where did the rest of my search + // result go" is answerable from the logs rather than by reading this + // function. + if store.is_none() { + tracing::info!( + "[agent][tool-result-artifacts] no artifact store configured; truncating oversized tool result inline tool={} original_bytes={} budget_bytes={} — the tail is discarded, not recoverable", + tool_name, + original_bytes, + budget_bytes + ); + } + let (output, BudgetOutcome { final_bytes, .. }) = + apply_tool_result_budget(content, budget_bytes); + ( + output, + ToolResultArtifactOutcome { + original_bytes, + final_bytes, + persisted: false, + artifact_path: None, + }, + ) +} + +pub async fn spill_aggregate_tool_results( + results: &mut [tinytools_agent::dialect::ToolOutcome], + store: Option<&ToolResultArtifactStore>, + budget_bytes: usize, +) { + if budget_bytes == 0 { + return; + } + let Some(store) = store else { + return; + }; + + let mut total: usize = results.iter().map(|result| result.output.len()).sum(); + if total <= budget_bytes { + return; + } + + let mut indexes: Vec = (0..results.len()).collect(); + indexes.sort_by_key(|idx| std::cmp::Reverse(results[*idx].output.len())); + + for idx in indexes { + if total <= budget_bytes { + break; + } + let original = results[idx].output.clone(); + let original_len = original.len(); + let allowed_len = budget_bytes.saturating_sub(total.saturating_sub(original_len)); + let persisted_output = if looks_like_preview_envelope(&original) { + Ok(PersistedToolResult { + output: original.clone(), + path: "".to_string(), + original_bytes: original_len, + stored_bytes: original_len, + redacted: false, + }) + } else { + store + .persist( + &results[idx].name, + results[idx].tool_call_id.as_deref(), + &original, + None, + allowed_len.min(AGGREGATE_PREVIEW_BUDGET_BYTES), + "aggregate tool-result budget exceeded", + ) + .await + }; + match persisted_output { + Ok(persisted) => { + // #4469 item 6: never bound the preview envelope below the minimum + // that preserves its `[tool_result_preview]` header + artifact + // pointer — `allowed_len` can be 0 here, which would blank the + // result and strip the `artifact_path` the model reads to recover + // the full output. + let envelope_allowance = allowed_len.max(MIN_ENVELOPE_ALLOWANCE_BYTES); + let (output, final_bytes) = + bound_text_to_budget(persisted.output, envelope_allowance); + total = total + .saturating_sub(original_len) + .saturating_add(final_bytes); + tracing::info!( + "[agent][tool-result-artifacts] aggregate spill tool={} original_bytes={} final_bytes={} total_bytes={} path={}", + results[idx].name, + original_len, + final_bytes, + total, + persisted.path + ); + results[idx].output = output; + } + Err(err) => { + tracing::warn!( + "[agent][tool-result-artifacts] aggregate spill failed tool={} bytes={} err={} -- falling back to inline budget trim", + results[idx].name, + original_len, + err + ); + let (output, final_bytes) = bound_text_to_budget(original, allowed_len); + total = total + .saturating_sub(original_len) + .saturating_add(final_bytes); + results[idx].output = output; + } + } + } +} + +fn looks_like_preview_envelope(value: &str) -> bool { + value.starts_with("[tool_result_preview]\n") +} + +fn bound_text_to_budget(content: String, budget_bytes: usize) -> (String, usize) { + if budget_bytes == 0 { + return (String::new(), 0); + } + let (mut output, BudgetOutcome { final_bytes, .. }) = + apply_tool_result_budget(content, budget_bytes); + if final_bytes <= budget_bytes { + return (output, final_bytes); + } + let cut = floor_char_boundary(&output, budget_bytes); + output.truncate(cut); + let final_bytes = output.len(); + (output, final_bytes) +} + +/// Round a byte index DOWN to the nearest UTF-8 character boundary. +fn floor_char_boundary(s: &str, index: usize) -> usize { + if index >= s.len() { + return s.len(); + } + let mut end = index; + while end > 0 && !s.is_char_boundary(end) { + end -= 1; + } + end +} + +/// 32 lowercase hex characters, unique per call within and across processes for +/// naming an artifact whose tool call carried no id. +fn random_call_id() -> String { + use std::collections::hash_map::RandomState; + use std::hash::{BuildHasher, Hasher}; + use std::sync::atomic::{AtomicU64, Ordering}; + + static COUNTER: AtomicU64 = AtomicU64::new(0); + let nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_nanos() as u64) + .unwrap_or(0); + let salt = COUNTER.fetch_add(1, Ordering::Relaxed); + let half = |extra: u64| { + let mut hasher = RandomState::new().build_hasher(); + hasher.write_u64(nanos); + hasher.write_u64(salt); + hasher.write_u64(extra); + hasher.finish() + }; + format!("{:016x}{:016x}", half(1), half(2)) +} + +fn sanitize_component(value: &str) -> String { + let mut out = String::with_capacity(value.len().min(80)); + for ch in value.chars().take(80) { + if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' { + out.push(ch); + } else { + out.push('_'); + } + } + if out.is_empty() { + "unknown".to_string() + } else { + out + } +} + +fn assert_within_action_dir(action_dir: &Path, path: &Path) -> anyhow::Result<()> { + if path.starts_with(action_dir) { + return Ok(()); + } + anyhow::bail!( + "tool-result artifact path escaped action_dir: {}", + path.display() + ); +} + +/// Return the newest modification time in a session tree without following +/// symlinks. Nested tool writes must keep their containing session alive. +fn newest_modified(path: &Path) -> std::io::Result { + let metadata = std::fs::symlink_metadata(path)?; + let mut newest = metadata.modified()?; + if metadata.is_dir() { + for entry in std::fs::read_dir(path)? { + let entry = entry?; + if entry.file_type()?.is_symlink() { + continue; + } + if let Ok(modified) = newest_modified(&entry.path()) { + newest = newest.max(modified); + } + } + } + Ok(newest) +} + +#[cfg(test)] +#[path = "tool_results_test.rs"] +mod test; diff --git a/crates/tinyagents-harness/src/artifacts/tool_results_test.rs b/crates/tinyagents-harness/src/artifacts/tool_results_test.rs new file mode 100644 index 000000000..60f9643cd --- /dev/null +++ b/crates/tinyagents-harness/src/artifacts/tool_results_test.rs @@ -0,0 +1,622 @@ +use super::*; +use crate::artifacts::policy::{ArtifactRedactor, Redacted}; +use serde_json::json; +use std::path::Path; +use std::sync::Arc; +use std::time::Duration; + +/// Stand-in for a host's credential/PII pass: rewrites a GitHub token and +/// phone numbers, and (like the real one) can grow a body. +#[derive(Debug)] +struct TestRedactor; + +impl ArtifactRedactor for TestRedactor { + fn redact(&self, content: &str) -> Redacted { + let mut out = content.replace("ghp_abcdefghijklmnopqrstuvwxyz123456", "[REDACTED_SECRET]"); + out = out.replace("+1555", "[REDACTED_PII_PHONE]"); + if out == content { + Redacted::unchanged(out) + } else { + Redacted::rewritten(out) + } + } +} + +fn store(dir: &Path, session: &str) -> ToolResultArtifactStore { + ToolResultArtifactStore::new( + dir.to_path_buf(), + session, + Arc::new(TestRedactor), + "file_read", + 10 * 1024 * 1024, + ) +} + +fn read_target(tool: &str, args: &serde_json::Value) -> Option { + artifact_read_target(tool, args, "file_read", "use_skill") +} + +#[tokio::test] +async fn threshold_persists_redacted_preview_and_file() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session/one"); + let raw = format!( + "{} {}", + "x".repeat(4096), + "ghp_abcdefghijklmnopqrstuvwxyz123456" + ); + + let (out, outcome) = + apply_per_result_persistence(raw, None, Some(&store), "shell", Some("call-1"), 1024).await; + + assert!(outcome.persisted); + assert!(out.contains("artifact_path: artifacts/tool-results/session_one/shell/call-1.txt")); + assert!(out.contains( + "read_with: file_read {\"path\":\"artifacts/tool-results/session_one/shell/call-1.txt\"}" + )); + assert!(out.contains("original_bytes:")); + assert!(out.contains("[preview]")); + assert!(out.contains("Credential/PII redaction was applied")); + assert!(!out.contains("ghp_abcdefghijklmnopqrstuvwxyz123456")); + + let stored = std::fs::read_to_string( + tmp.path() + .join("artifacts/tool-results/session_one/shell/call-1.txt"), + ) + .unwrap(); + assert!(stored.contains("xxxx")); + assert!(!stored.contains("ghp_abcdefghijklmnopqrstuvwxyz123456")); +} + +/// A call with no id still gets a unique, well-formed file name. +#[tokio::test] +async fn a_call_without_an_id_gets_a_random_hex_name() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "s"); + let a = store.path_for_read_tool("shell", None); + let b = store.path_for_read_tool("shell", None); + assert_ne!(a, b); + let stem = a + .strip_prefix("artifacts/tool-results/s/shell/") + .and_then(|rest| rest.strip_suffix(".txt")) + .expect("path shape"); + assert_eq!(stem.len(), 32); + assert!(stem.chars().all(|c| c.is_ascii_hexdigit())); +} + +#[tokio::test] +async fn fallback_truncates_when_store_missing() { + let raw = "z".repeat(4096); + let (out, outcome) = + apply_per_result_persistence(raw, None, None, "shell", Some("call"), 512).await; + assert!(!outcome.persisted); + assert!(out.contains("truncated by tool_result_budget")); + assert!(out.len() < 4096); +} + +/// The truncation trailer must not tell the model to re-run the call (#6408). +/// +/// The old wording ended "re-run with a narrower query to see the rest". A +/// listing tool with no narrowing argument in reach leaves the model only the +/// identical call, which returns the identical truncation, until the +/// successful-repeat tracker halts the run. Assert the retry instruction is +/// gone and the deterministic-truncation statement that replaces it is +/// present, so a revert to the old trailer fails here. +#[tokio::test] +async fn truncation_trailer_does_not_instruct_a_retry() { + let raw = "z".repeat(4096); + let (out, outcome) = + apply_per_result_persistence(raw, None, None, "GITHUB_LIST_PULL_REQUESTS", None, 512).await; + + assert!( + !outcome.persisted, + "fixture must truncate inline, not persist" + ); + assert!( + !out.contains("re-run"), + "trailer must not instruct a re-run; got: {out}" + ); + assert!( + out.contains("Repeating this call returns the same truncation"), + "trailer must say the truncation is deterministic; got: {out}" + ); + // Both totals, so the model can judge whether the retained head suffices. + assert!( + out.contains("of 4096 bytes truncated by tool_result_budget"), + "trailer must report dropped-of-original bytes; got: {out}" + ); +} + +/// The sweep bounds growth without touching the live session (#6408). +/// +/// Nothing else deletes these files, so an unbounded artifact directory would +/// trade a token-burn bug for a disk-growth one. Asserts both halves: a stale +/// OTHER session is removed, and the CURRENT session survives regardless of +/// age — a sweep that collected the running session's own artifacts would +/// delete the bodies the model is about to read back. +#[tokio::test] +async fn prune_removes_stale_sessions_but_never_the_current_one() { + let tmp = tempfile::tempdir().unwrap(); + let root = tmp.path().join("artifacts/tool-results"); + let current = root.join("live-session"); + let stale = root.join("old-session"); + std::fs::create_dir_all(¤t).unwrap(); + std::fs::create_dir_all(&stale).unwrap(); + std::fs::write(current.join("a.txt"), "current").unwrap(); + std::fs::write(stale.join("b.txt"), "stale").unwrap(); + + let store = store(tmp.path(), "live-session"); + + // Nothing is old enough yet: a generous window must collect nothing. + assert_eq!( + store + .prune_stale_sessions(Duration::from_secs(3600)) + .unwrap(), + 0 + ); + assert!(stale.exists(), "nothing is stale within the window"); + + // A zero window makes every directory stale, so only the current-session + // exemption can save `live-session`. + let removed = store.prune_stale_sessions(Duration::from_secs(0)).unwrap(); + assert_eq!( + removed, 1, + "exactly the one other session should be collected" + ); + assert!(!stale.exists(), "stale session dir must be removed"); + assert!( + current.join("a.txt").exists(), + "the current session's artifacts must survive its own sweep" + ); +} + +/// A missing artifact root is the normal first-run case, not an error. +#[tokio::test] +async fn prune_is_a_noop_when_no_artifacts_exist() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session"); + assert_eq!( + store.prune_stale_sessions(Duration::from_secs(0)).unwrap(), + 0 + ); +} + +#[tokio::test] +async fn persisted_preview_is_bounded_for_small_budget() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session"); + let raw = "x".repeat(800); + + let (out, outcome) = + apply_per_result_persistence(raw, None, Some(&store), "shell", Some("call"), 320).await; + + assert!(outcome.persisted); + assert!(outcome.final_bytes <= MIN_ENVELOPE_ALLOWANCE_BYTES); + assert_eq!(out.len(), outcome.final_bytes); + assert!(out.contains("[tool_result_preview]")); + assert!( + tmp.path() + .join("artifacts/tool-results/session/shell/call.txt") + .exists() + ); +} + +#[tokio::test] +async fn tiny_per_result_budget_keeps_the_recovery_pointer() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session"); + let (out, outcome) = apply_per_result_persistence( + "x".repeat(800), + None, + Some(&store), + "shell", + Some("call"), + 8, + ) + .await; + assert!(outcome.persisted); + assert!(out.contains("artifact_path: artifacts/tool-results/session/shell/call.txt")); + assert!(out.contains("read_with: file_read")); +} + +#[tokio::test] +async fn nested_artifact_activity_is_included_in_session_freshness() { + let tmp = tempfile::tempdir().unwrap(); + let root = tmp.path().join("artifacts/tool-results"); + let old_session = root.join("old-session"); + let tool_dir = old_session.join("shell"); + std::fs::create_dir_all(&tool_dir).unwrap(); + std::fs::write(tool_dir.join("recent.txt"), "recent").unwrap(); + let file_time = std::fs::metadata(tool_dir.join("recent.txt")) + .unwrap() + .modified() + .unwrap(); + assert_eq!(newest_modified(&old_session).unwrap(), file_time); +} + +#[cfg(unix)] +#[tokio::test] +async fn persist_rejects_symlinked_parent_outside_action_dir() { + use std::os::unix::fs::symlink; + let tmp = tempfile::tempdir().unwrap(); + let action = tmp.path().join("action"); + let outside = tmp.path().join("outside"); + std::fs::create_dir_all(action.join("artifacts/tool-results/session")).unwrap(); + std::fs::create_dir_all(&outside).unwrap(); + symlink( + &outside, + action.join("artifacts/tool-results/session/shell"), + ) + .unwrap(); + let store = store(&action, "session"); + let result = store + .persist("shell", Some("call"), "x", None, 0, "test") + .await; + assert!(result.is_err()); + assert!(!outside.join("call.txt").exists()); +} + +#[tokio::test] +async fn aggregate_spills_largest_until_under_budget() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session"); + let mut results = vec![ + tinytools_agent::dialect::ToolOutcome { + name: "small".into(), + output: "a".repeat(100), + success: true, + tool_call_id: Some("small".into()), + trusted_verbatim: false, + }, + tinytools_agent::dialect::ToolOutcome { + name: "largest".into(), + output: "b".repeat(2000), + success: true, + tool_call_id: Some("largest".into()), + trusted_verbatim: false, + }, + tinytools_agent::dialect::ToolOutcome { + name: "medium".into(), + output: "c".repeat(900), + success: true, + tool_call_id: Some("medium".into()), + trusted_verbatim: false, + }, + ]; + + spill_aggregate_tool_results(&mut results, Some(&store), 1800).await; + + assert!(results[1].output.starts_with("[tool_result_preview]\n")); + let total: usize = results.iter().map(|result| result.output.len()).sum(); + assert!(total <= 1800, "total={total}"); + assert!(!results[0].output.starts_with("[tool_result_preview]\n")); + assert!( + tmp.path() + .join("artifacts/tool-results/session/largest/largest.txt") + .exists() + ); +} + +#[tokio::test] +async fn aggregate_forces_budget_when_envelope_has_no_savings() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session"); + let mut results = vec![ + tinytools_agent::dialect::ToolOutcome { + name: "one".into(), + output: "a".repeat(350), + success: true, + tool_call_id: Some("one".into()), + trusted_verbatim: false, + }, + tinytools_agent::dialect::ToolOutcome { + name: "two".into(), + output: "b".repeat(350), + success: true, + tool_call_id: Some("two".into()), + trusted_verbatim: false, + }, + tinytools_agent::dialect::ToolOutcome { + name: "three".into(), + output: "c".repeat(350), + success: true, + tool_call_id: Some("three".into()), + trusted_verbatim: false, + }, + ]; + + spill_aggregate_tool_results(&mut results, Some(&store), 500).await; + + let total: usize = results.iter().map(|result| result.output.len()).sum(); + // #4469 item 6: the aggregate spill now floors each persisted envelope at + // MIN_ENVELOPE_ALLOWANCE_BYTES so the `[tool_result_preview]` header + + // `artifact_path` pointer always survives (previously an exhausted budget + // could blank a result to ""). That is a documented trade — the total may + // slightly overshoot the raw aggregate budget — so the invariant is now: + // (a) no envelope is blanked, and (b) the total stays bounded by the + // per-result floor rather than the raw budget. + assert!( + results.iter().all(|result| !result.output.is_empty()), + "no persisted envelope may be blanked — the artifact pointer must survive" + ); + assert!( + total <= results.len() * MIN_ENVELOPE_ALLOWANCE_BYTES, + "total={total} exceeds the per-result envelope floor bound" + ); + assert!( + tmp.path() + .join("artifacts/tool-results/session/one/one.txt") + .exists() + ); +} + +#[test] +fn artifact_read_target_finds_a_path_nested_in_a_wrapper_call() { + let args = json!({ + "skill": "files", + "tool": "file_read", + "args": {"path": "artifacts/tool-results/s/use_skill/c.txt", "offset": 42} + }); + assert_eq!( + read_target("use_skill", &args), + Some(ArtifactRead { + path: "artifacts/tool-results/s/use_skill/c.txt".to_string(), + offset: 42, + }) + ); + assert_eq!( + read_target("file_read", &json!({"path": "src/main.rs"})), + None + ); +} + +#[tokio::test] +async fn persisted_outcome_reports_the_size_of_the_stored_body() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "session"); + let rewritten = "s".repeat(3_000); + let full = "r".repeat(8_000); + + let (out, outcome) = apply_per_result_persistence( + rewritten, + Some(full), + Some(&store), + "shell", + Some("call"), + 1_000, + ) + .await; + + assert!(outcome.persisted); + assert_eq!( + outcome.original_bytes, 8_000, + "the outcome feeds the artifact index, so it must report the stored body's size" + ); + assert!(out.contains("original_bytes: 8000"), "{out}"); +} + +#[test] +fn artifact_read_target_matches_only_file_read_under_the_artifact_directory() { + let path = "artifacts/tool-results/s/shell/c.txt"; + assert!(read_target("file_read", &json!({"path": path})).is_some()); + assert_eq!( + read_target("file_write", &json!({"path": path, "content": "x"})), + None, + "a write to an artifact path is not a read of its content" + ); + assert_eq!( + read_target( + "use_skill", + &json!({"skill": "files", "tool": "glob", "args": {"path": path}}) + ), + None, + "a wrapped non-read tool is not a read of its content" + ); +} + +#[test] +fn artifact_read_target_matches_the_artifact_directory_as_a_path_component() { + assert!( + read_target( + "file_read", + &json!({"path": "./artifacts/tool-results/s/c.txt"}) + ) + .is_some() + ); + assert_eq!( + read_target( + "file_read", + &json!({"path": "artifacts/tool-results-backup/report.txt"}) + ), + None, + "a sibling directory sharing the prefix is not the artifact directory" + ); +} + +#[test] +fn an_artifact_page_stays_within_the_budget_with_a_long_path() { + let read = ArtifactRead { + path: format!( + "artifacts/tool-results/{}/use_skill/call.txt", + "s".repeat(240) + ), + offset: 1_234_567, + }; + + let page = page_artifact_read("y".repeat(5_000), &read, 1_000, "file_read"); + assert!( + page.len() <= 1_000, + "a page must fit the result budget, got {} bytes", + page.len() + ); + let body = page + .find("\n\n[artifact page") + .expect("continuation marker"); + assert!( + page.contains(&format!("\"offset\":{}", 1_234_567 + body)), + "the marker must name the offset right after this page: {page}" + ); +} + +#[test] +fn a_page_never_exceeds_the_floored_budget_and_always_advances() { + // A path too long for its full trailer to leave body room under the floor. + let read = ArtifactRead { + path: format!("artifacts/tool-results/{}/c.txt", "p".repeat(600)), + offset: 7, + }; + let content = "z".repeat(5_000); + + for budget in [2, 100, MIN_ENVELOPE_ALLOWANCE_BYTES, 700] { + let page = page_artifact_read(content.clone(), &read, budget, "file_read"); + let limit = budget.max(MIN_ENVELOPE_ALLOWANCE_BYTES); + assert!( + page.len() <= limit, + "budget {budget}: a page must fit max(budget, floor) = {limit}, got {} bytes", + page.len() + ); + let body = page + .find("\n\n[artifact page") + .expect("continuation marker"); + assert!( + body > 0, + "budget {budget}: every page must advance past its offset" + ); + assert!( + page.contains(&format!("\"offset\":{}", 7 + body)), + "budget {budget}: the marker must name the offset right after this page: {page}" + ); + } +} + +#[test] +fn a_body_redaction_grows_past_the_read_limit_falls_back_to_the_processed_copy() { + let raw = "call +15551234567 or +15557654321"; + // The limit is the raw size: the raw body fits, its redacted form does not. + let (chosen, stored) = readable_body( + raw, + Some("processed copy"), + raw.len() as u64, + &TestRedactor, + "file_read", + ) + .expect("the processed copy fits"); + assert_eq!( + chosen, "processed copy", + "a body whose sanitized form exceeds the read limit must not be stored, got {:?}", + stored.text + ); + + let (kept, _) = readable_body( + raw, + Some("processed copy"), + 10_000, + &TestRedactor, + "file_read", + ) + .expect("fits"); + assert_eq!( + kept, raw, + "a body that stays within the limit is stored as returned" + ); +} + +#[test] +fn a_body_is_refused_when_neither_candidate_fits_the_read_limit() { + let raw = "call +15551234567 or +15557654321"; + let fallback = "fallback +15550001111 +15550002222"; + let limit = raw.len().min(fallback.len()) as u64; + assert!( + readable_body(raw, Some(fallback), limit, &TestRedactor, "file_read").is_err(), + "when neither the raw body nor the fallback fits once sanitized, nothing may be stored" + ); + assert!( + readable_body(raw, None, limit, &TestRedactor, "file_read").is_err(), + "a body with no fallback that does not fit must not be stored either" + ); +} + +#[test] +fn artifact_read_target_follows_only_use_skill_into_a_wrapped_tool() { + let nested = json!({"tool": "file_read", "args": {"path": "artifacts/tool-results/s/c.txt"}}); + assert!( + read_target("use_skill", &nested).is_some(), + "use_skill forwards file_read's result, so its wrapped read counts" + ); + for outer in ["glob", "file_write", "shell"] { + assert_eq!( + read_target(outer, &nested), + None, + "{outer} carrying tool/args fields is not a wrapper; its result is its own" + ); + } +} + +#[test] +fn a_page_near_the_maximum_offset_neither_overflows_nor_advertises_a_stuck_continuation() { + let read = ArtifactRead { + path: "artifacts/tool-results/s/shell/c.txt".to_string(), + offset: usize::MAX - 10, + }; + // The offset comes straight from the model's arguments, so the page + // arithmetic must not overflow on an absurd one. + let page = std::panic::catch_unwind(|| { + page_artifact_read("q".repeat(5_000), &read, 1_000, "file_read") + }); + assert!( + page.is_ok(), + "an offset near usize::MAX must not overflow the page arithmetic" + ); + let page = page.unwrap(); + assert!( + page.len() <= 1_000, + "the result is still bounded, got {} bytes", + page.len() + ); + assert!( + !page.contains("Continue with"), + "no continuation may be advertised when the next offset cannot advance" + ); +} + +#[test] +fn artifact_read_target_rejects_an_explicit_invalid_offset() { + let path = "artifacts/tool-results/s/shell/c.txt"; + for bad in [json!(-1), json!(1.5), json!("12")] { + assert_eq!( + read_target("file_read", &json!({"path": path, "offset": bad.clone()})), + None, + "offset {bad} is not a read file_read serves, so it must not become an artifact read at 0" + ); + } + for absent in [json!({"path": path}), json!({"path": path, "offset": null})] { + assert_eq!( + read_target("file_read", &absent).map(|read| read.offset), + Some(0), + "an absent or null offset reads from the start" + ); + } +} + +/// Literal fixture: the envelope the model reads is a wire format hosts and +/// runbooks match on, so pin it byte for byte. +#[tokio::test] +async fn envelope_text_is_byte_stable() { + let tmp = tempfile::tempdir().unwrap(); + let store = store(tmp.path(), "s"); + let persisted = store + .persist("shell", Some("c"), "hello world", None, 1024, "r") + .await + .unwrap(); + let expected = "[tool_result_preview]\n\ +tool: shell\n\ +reason: r\n\ +original_bytes: 11\n\ +stored_bytes: 11\n\ +artifact_path: artifacts/tool-results/s/shell/c.txt\n\ +read_with: file_read {\"path\":\"artifacts/tool-results/s/shell/c.txt\"} (a long read returns one page and names the \"offset\" to continue from)\n\ +notes: Full scrubbed output was persisted under the action workspace.\n\n\ +[preview]\nhello world"; + assert_eq!(persisted.output, expected); + assert!(!persisted.redacted); +} diff --git a/crates/tinyagents-harness/src/tool/mod.rs b/crates/tinyagents-harness/src/tool/mod.rs index 5d95762e8..2c36bc0c4 100644 --- a/crates/tinyagents-harness/src/tool/mod.rs +++ b/crates/tinyagents-harness/src/tool/mod.rs @@ -7,6 +7,7 @@ pub mod deferred; pub mod discover; pub mod effects; +pub mod packs; mod prompt; mod schema; mod schema_compact; diff --git a/crates/tinyagents-harness/src/tool/packs/catalog.rs b/crates/tinyagents-harness/src/tool/packs/catalog.rs new file mode 100644 index 000000000..ae5911fbc --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/catalog.rs @@ -0,0 +1,103 @@ +//! The pack table a host hands to the generic `use_skill` machinery. +//! +//! The table itself is host data — which tools are packed, who owns them, what +//! the playbook says — and stays with the host. [`PackCatalog`] is the lookup +//! surface over it, so the disclosure tool, the listing renderer and the spec +//! scoper never name a global. + +use super::types::ToolPack; + +/// A host's compiled-in pack table plus the one piece of host vocabulary the +/// rendered text needs: the marker that makes a "no such skill / tool" result +/// classify as not-found in the host's status taxonomy. +#[derive(Debug, Clone, Copy)] +pub struct PackCatalog { + packs: &'static [ToolPack], + not_found_marker: &'static str, +} + +impl PackCatalog { + /// A catalog over `packs`. `not_found_marker` prefixes every "unknown skill" + /// or "no such tool" result so the host's error classifier can recognise it. + pub const fn new(packs: &'static [ToolPack], not_found_marker: &'static str) -> Self { + Self { + packs, + not_found_marker, + } + } + + /// Every pack, in table order. + pub const fn packs(&self) -> &'static [ToolPack] { + self.packs + } + + /// The not-found marker this catalog's host supplied. + pub const fn not_found_marker(&self) -> &'static str { + self.not_found_marker + } + + /// The pack with this id. + pub fn pack(&self, id: &str) -> Option<&'static ToolPack> { + self.packs.iter().find(|p| p.id == id) + } + + /// The pack owning `tool`, if any. + pub fn pack_for_tool(&self, tool: &str) -> Option<&'static ToolPack> { + self.packs.iter().find(|p| p.owns(tool)) + } + + /// Every packed tool name across all packs. + pub fn all_packed_tool_names(&self) -> Vec<&'static str> { + self.packs + .iter() + .flat_map(|p| p.tools.iter().copied()) + .collect() + } + + /// Every packed tool name that applies to `agent_id`. + /// + /// A pack is skipped entirely for agents listed as its owners — see + /// [`ToolPack::owners`]. + pub fn packed_tool_names_for_agent(&self, agent_id: &str) -> Vec<&'static str> { + self.packs + .iter() + .filter(|p| !p.is_owner(agent_id)) + .flat_map(|p| p.tools.iter().copied()) + .collect() + } + + /// The always-on index: one line per pack, rendered into `use_skill`'s own + /// description so the model can pick a pack without a round trip. + pub fn pack_index_markdown(&self) -> String { + self.pack_index_markdown_filtered(&|_| true) + } + + /// The pack index, limited to packs this session can call at least one tool in. + /// + /// A pack with nothing callable is not an answer to "which skills can I load", + /// and advertising it costs a round trip: the model loads it, learns it cannot + /// use it, and comes back. The capability does not disappear — a pack's owners + /// reach the model through their own delegation tools, whose descriptions are + /// already on the wire. Keeping the pack listed here would duplicate that + /// routing on every single turn. + pub fn pack_index_markdown_filtered(&self, is_callable: &dyn Fn(&str) -> bool) -> String { + let mut out = String::new(); + for p in self.packs { + if !p.tools.iter().any(|t| is_callable(t)) { + continue; + } + out.push_str(&format!("- `{}` — {}\n", p.id, p.summary)); + } + out + } + + /// Pack ids with at least one tool this session can call — the `skill` enum + /// `use_skill` should actually offer. + pub fn callable_pack_ids(&self, is_callable: &dyn Fn(&str) -> bool) -> Vec<&'static str> { + self.packs + .iter() + .filter(|p| p.tools.iter().any(|t| is_callable(t))) + .map(|p| p.id) + .collect() + } +} diff --git a/crates/tinyagents-harness/src/tool/packs/handle.rs b/crates/tinyagents-harness/src/tool/packs/handle.rs new file mode 100644 index 000000000..ea018cd34 --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/handle.rs @@ -0,0 +1,143 @@ +//! The late-bound view of the tool registries a pack tool dispatches into. + +use std::sync::{Arc, RwLock, Weak}; + +use tinytools::Tool; + +use super::catalog::PackCatalog; + +/// An `Arc`-shared, owned view of the tool registry a pack tool lives in. +pub(super) type ToolVec = Arc>>; + +/// A non-owning view of the tool registry, kept to break the binding cycle. +type ToolRegistryRef = Weak>>; + +/// A late-bound, non-owning view of the tool registries a pack tool reads. +/// +/// Late-bound because the pack tool is *inside* the registry it reads: the +/// vector cannot be built until it exists, and it cannot see it until it is +/// built. Non-owning because an `Arc` back into that same vector would be a +/// cycle that never drops. +/// +/// **Two registries, not one, and that is load-bearing.** An agent's tools live +/// in two `Arc`s: the durable registry, and the `synthesized_tools` set that +/// `collect_orchestrator_tools` rebuilds whenever the Composio connection set +/// changes (they were split in #6145 so a reconcile cannot block on a reader). +/// Every `delegate_*` tool is in the second one. Binding only the first is why +/// `use_skill` could not reach a single packed delegate — `do_crypto`, +/// `run_skill`, `build_workflow` and four more were withheld from the wire and +/// then unreachable through the route that was supposed to replace them, which +/// is strictly worse than not packing them at all. +#[derive(Clone, Default)] +pub struct PackRegistryHandle { + inner: Arc>, +} + +/// The two registries, each independently rebindable. +/// +/// They are separate slots rather than one vector because they are replaced on +/// different schedules: the durable registry is rebuilt when the agent is, the +/// synthesised one on every delegation refresh. +#[derive(Default)] +struct Slots { + durable: Option, + synthesized: Option, +} + +impl PackRegistryHandle { + /// Point this handle at the durable registry it lives in, replacing any + /// previous binding. + /// + /// Rebinding has to actually take effect. This was a `OnceLock` whose + /// second write was dropped, which silently contradicted + /// the host's `bind_pack_registry`'s own instruction to "re-bind after any + /// later rebuild of this `Arc`": once an agent replaced its tool vector the + /// handle still pointed at the old allocation, the `Weak` failed to + /// upgrade, and every `use_skill` call reported the registry as unavailable + /// for the rest of the session. Last write wins. + pub fn bind(&self, registry: ToolRegistryRef) { + self.with_slots(|slots| slots.durable = Some(registry)); + } + + /// Point this handle at the synthesised delegate set. + /// + /// Call it again after **every** `refresh_delegation_tools`, which replaces + /// that `Arc` wholesale — a stale `Weak` stops upgrading as soon as the last + /// reader of the old allocation goes, and the packed delegates silently + /// become unreachable. + pub fn bind_synthesized(&self, registry: ToolRegistryRef) { + self.with_slots(|slots| slots.synthesized = Some(registry)); + } + + fn with_slots(&self, edit: impl FnOnce(&mut Slots)) { + match self.inner.write() { + Ok(mut slots) => edit(&mut slots), + // The lock is only ever held for a pointer read or write, so a + // poisoned lock means a panic elsewhere. Recover rather than + // propagate: a stale binding degrades to "skill unavailable", + // which is the failure this rebinding exists to prevent. + Err(poisoned) => edit(&mut poisoned.into_inner()), + } + } + + /// Every live registry, durable first. + /// + /// Order matters on a name collision: `drop_synthesized_name_collisions` + /// gives the durable tool the name, so resolving durable-first is what + /// makes this agree with what the harness would actually execute. + pub(super) fn registries(&self) -> Vec { + let slots = match self.inner.read() { + Ok(slots) => slots, + Err(poisoned) => poisoned.into_inner(), + }; + [slots.durable.as_ref(), slots.synthesized.as_ref()] + .into_iter() + .flatten() + .filter_map(Weak::upgrade) + .collect() + } + + /// Resolve a packed tool by name, enforcing that it belongs to `skill`. + /// + /// The pack check is not decoration: without it `use_skill` would dispatch + /// into any packed tool regardless of the skill named, and the model could + /// reach a crypto write through a workflow skill. + pub(super) fn resolve( + &self, + catalog: &PackCatalog, + skill: &str, + tool: &str, + ) -> Option<(ToolVec, usize)> { + catalog.pack(skill).filter(|p| p.owns(tool))?; + self.find(tool) + } + + /// Resolves `tool` in `skill`'s pack, returning the exact registry `Arc` + /// it lives in — not a clone of the tool itself — so a caller can re-wrap + /// it in the same `CanonicalSharedToolAdapter` seam the harness uses at + /// registration for typed-dispatch selection. + /// + /// Public because a host's typed dispatch needs the same resolution + /// [`UseSkillTool`](super::UseSkillTool) performs, so a packed delegation + /// reached through `use_skill` can be re-dispatched through a live-parent + /// seam instead of falling back to plain `Tool::execute_with_context`. + pub fn resolve_registry_for( + &self, + catalog: &PackCatalog, + skill: &str, + tool: &str, + ) -> Option>>> { + self.resolve(catalog, skill, tool) + .map(|(tools, _idx)| tools) + } + + /// Locate `tool` in whichever registry holds it. + pub(super) fn find(&self, tool: &str) -> Option<(ToolVec, usize)> { + for tools in self.registries() { + if let Some(idx) = tools.iter().position(|t| t.name() == tool) { + return Some((tools, idx)); + } + } + None + } +} diff --git a/crates/tinyagents-harness/src/tool/packs/mod.rs b/crates/tinyagents-harness/src/tool/packs/mod.rs new file mode 100644 index 000000000..6c382fd1f --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/mod.rs @@ -0,0 +1,41 @@ +//! On-demand tool disclosure ("tool packs"). +//! +//! A tool's JSON schema is charged on every provider call of every turn, +//! whether or not the tool is used, so most of an orchestrator's fixed cost is +//! idle schema. A *pack* keeps its tools constructed and executable but +//! unadvertised; the agent sees one small tool instead. [`UseSkillTool`] +//! renders a pack's schemas into the conversation when called with a `skill` +//! alone, and executes one of them when also given a `tool`, forwarding +//! permission level, external-effect classification and timeout policy to the +//! real tool so nothing is laundered through the proxy. +//! +//! This is the mechanism only. What the host keeps: +//! +//! * the **pack table** (which tools are packed, who owns them, the playbooks), +//! handed in as a [`PackCatalog`]; +//! * the **posture** (which groups are withheld, advertised or off for a given +//! embedder, and which packs an agent may reach); +//! * the **binding** of a [`PackRegistryHandle`] to the registries a session +//! builds. +//! +//! Not to be confused with [`super::discover`], which is model-driven +//! `tool_search` over a BM25 catalog of deferred schemas: that one lets the model +//! *search* for a tool, this one lets a host *group* tools under named skills +//! with a playbook and an owner list. + +mod catalog; +mod handle; +mod render; +mod tool; +mod types; + +pub use catalog::PackCatalog; +pub use handle::PackRegistryHandle; +pub use render::{ + NoSuchPackTool, named_tool, render_pack_filtered, route_sentence, scope_use_skill_spec, +}; +pub use tool::{USE_SKILL, UseSkillTool}; +pub use types::ToolPack; + +#[cfg(test)] +mod test; diff --git a/crates/tinyagents-harness/src/tool/packs/render.rs b/crates/tinyagents-harness/src/tool/packs/render.rs new file mode 100644 index 000000000..8af3c3c43 --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/render.rs @@ -0,0 +1,247 @@ +//! Rendering for the `use_skill` disclosure half: pack listings, not-found +//! results, routing sentences and per-session spec scoping. + +use serde_json::Value; +use tinytools::ToolSpec; + +use super::catalog::PackCatalog; +use super::handle::PackRegistryHandle; + +/// Render a pack's listing, showing only the tools `is_callable` admits. +/// +/// **The filter is the whole point.** `load_skill` used to render every tool in +/// the pack this build compiled, and `use_skill` then refused any of them the +/// session's allowlist denies (`tinyagents::middleware::channel_permission_block`). +/// A non-owner was handed a menu it could not order from: the orchestrator +/// loaded `workflows`, read `propose_workflow` off the listing, called it, and +/// was told it "is not allowed in the current session". The denial named no +/// alternative, so the model retried — one live chat turn died on the +/// repeated-failure breaker after six identical denials. +/// +/// `registry.rs` used to claim non-owners "reach them through `use_skill`". +/// That was never true: the gate (`d5a09ea81`, 2026-08-21) predates the comment +/// asserting it (`a8f0a002b`, 2026-08-23). The listing is the side that was +/// wrong, so the listing is the side that changed. +/// +/// `route` is the sentence to append when the session can call nothing in the +/// pack — see [`route_sentence`]. Empty means "say nothing extra". +pub fn render_pack_filtered( + catalog: &PackCatalog, + skill: &str, + handle: &PackRegistryHandle, + is_callable: &dyn Fn(&str) -> bool, + route: &str, +) -> Result { + let Some(pack) = catalog.pack(skill) else { + // Scoped too: offering a hallucinating model a pack it cannot use is the + // same wrong turn the advertised index used to take, one error later. + return Err(format!( + "{} Unknown skill `{skill}`. Available:\n{}", + catalog.not_found_marker(), + catalog.pack_index_markdown_filtered(is_callable) + )); + }; + if handle.registries().is_empty() { + return Err( + "The skill registry is not available in this session; the tools in this skill \ + cannot be loaded." + .to_string(), + ); + } + + let mut out = format!("# Skill `{}`\n\n{}\n\n", pack.id, pack.summary); + if !pack.guide.trim().is_empty() { + out.push_str(pack.guide.trim()); + out.push_str("\n\n"); + } + out.push_str(&format!( + "Call these with `use_skill {{ \"skill\": \"{}\", \"tool\": \"\", \"args\": {{ … }} }}`. \ + `args` is the tool's own argument object, exactly as documented below.\n\n", + pack.id + )); + + let mut found = 0usize; + for name in pack.tools { + // A pack may name a tool this build compiled out (feature gate) or that + // this agent never had. Rendering the ones that exist beats failing the + // whole load. + let Some((tools, idx)) = handle.find(name) else { + continue; + }; + // Listing a tool the gate will refuse is worse than omitting it: a + // model cannot tell a policy denial from a transient failure, so it + // retries the same call instead of routing around it. + if !is_callable(name) { + continue; + } + let tool = &tools[idx]; + found += 1; + out.push_str(&format!( + "## `{}`\n\n{}\n\n", + tool.name(), + tool.description() + )); + // Minified, matching what the provider receives for a natively + // advertised tool. Pretty-printing costs roughly a third more tokens + // for indentation and newlines the model gains nothing from, and this + // text is charged to the context window exactly like a native schema. + out.push_str("```json\n"); + out.push_str( + &serde_json::to_string(&tool.parameters_schema()).unwrap_or_else(|_| "{}".to_string()), + ); + out.push_str("\n```\n\n"); + } + + if found == 0 { + let mut message = format!( + "Skill `{}` has no tools available in this session.", + pack.id + ); + if !route.is_empty() { + message.push(' '); + message.push_str(route); + } + return Err(message); + } + Ok(out) +} + +/// A `use_skill` call naming a tool its skill does not contain (#6302). +/// +/// Typed so the gate and its tests read one source: an invented name +/// (`install_skill` in `skills`) must read as "no such tool" plus what the +/// session can call instead, never as a permission denial. The not-found +/// marker makes the failure classify as `NotFound`. +pub struct NoSuchPackTool<'a> { + /// The not-found marker of the host's [`PackCatalog`]. + pub not_found_marker: &'a str, + pub skill: &'a str, + pub tool: &'a str, + /// Tools in the skill this session can call, in pack order. + pub callable: Vec<&'a str>, + /// The hand-off sentence to use when `callable` is empty. + pub route: String, +} + +impl NoSuchPackTool<'_> { + pub fn render(&self) -> String { + let mut out = format!( + "{} There is no tool `{}` in skill `{}`.", + self.not_found_marker, self.tool, self.skill + ); + if !self.callable.is_empty() { + let names = self + .callable + .iter() + .map(|name| format!("`{name}`")) + .collect::>() + .join(", "); + out.push_str(&format!(" The tools in it you can call: {names}.")); + } else if !self.route.is_empty() { + out.push(' '); + out.push_str(&self.route); + } else { + out.push_str(" Nothing in it is available in this session."); + } + out + } +} + +/// The "go here instead" sentence shared by the `use_skill` listing and the +/// `use_skill` denial, so a model never sees two different stories. +/// +/// `callable_delegates` are delegation tool names the caller has already +/// confirmed this session can invoke — naming the *tool* rather than the agent +/// is the difference between guidance and an instruction, and a model left to +/// guess the call retries. When none can be reached the owning agents are named +/// instead: strictly worse, but still better than a bare denial. +pub fn route_sentence(callable_delegates: &[String], owners: &[&str]) -> String { + if !callable_delegates.is_empty() { + let names = callable_delegates + .iter() + .map(|t| format!("`{t}`")) + .collect::>() + .join(" or "); + return format!( + "Call {names} instead — that agent owns these tools and runs them directly. \ + Do not retry this skill." + ); + } + if owners.is_empty() { + return String::new(); + } + format!( + "These tools belong to {}; hand the task to one of them rather than calling directly.", + owners + .iter() + .map(|o| format!("`{o}`")) + .collect::>() + .join(" or ") + ) +} + +pub(super) fn render_pack( + catalog: &PackCatalog, + skill: &str, + handle: &PackRegistryHandle, +) -> Result { + render_pack_filtered(catalog, skill, handle, &|_| true, "") +} + +/// Rewrite `use_skill`'s advertised spec to match what this session can do. +/// +/// The description is built once in [`super::UseSkillTool::new`], before any session +/// exists, so every agent was told all ten packs were loadable — including ones +/// it can call nothing in. Post-#(routing fix) that costs one wasted round trip +/// instead of a dead turn; it should cost zero. +/// +/// Both halves are rewritten, and the schema is the stronger one: narrowing the +/// `skill` enum makes an unusable pack *unrepresentable* rather than merely +/// discouraged in prose, and a shorter enum is fewer tokens, not more. +/// +/// Returns `false` when this session can call nothing in any pack — the caller +/// should then drop `use_skill` from the wire entirely, because an empty index +/// and an empty enum are not a tool. +pub fn scope_use_skill_spec( + catalog: &PackCatalog, + spec: &mut ToolSpec, + is_callable: &dyn Fn(&str) -> bool, +) -> bool { + let ids = catalog.callable_pack_ids(is_callable); + if ids.is_empty() { + return false; + } + if let Some(index) = spec.description.find("\n\nSkills:\n") { + spec.description.truncate(index); + spec.description.push_str("\n\nSkills:\n"); + spec.description + .push_str(&catalog.pack_index_markdown_filtered(is_callable)); + } + if let Some(enum_slot) = spec + .parameters + .pointer_mut("/properties/skill/enum") + .filter(|v| v.is_array()) + { + *enum_slot = Value::Array( + ids.iter() + .map(|id| Value::String((*id).to_string())) + .collect(), + ); + } + true +} + +/// The tool named in `args`, if the caller named one at all. +/// +/// An absent (or empty) `tool` is not a malformed call: it is the disclosure +/// half of this tool, and the distinction decides both which branch +/// `UseSkillTool::execute_with_context` takes and what permission level the +/// call is gated at. Public because the policy middleware has to draw the same +/// line — it intercepts the disclosure half to scope the listing to the session +/// and lets the execution half through to its gate — and two spellings of "did +/// the caller name a tool" would be two chances to disagree. +pub fn named_tool(args: &Value) -> Option<&str> { + args.get("tool") + .and_then(Value::as_str) + .filter(|name| !name.is_empty()) +} diff --git a/crates/tinyagents-harness/src/tool/packs/test.rs b/crates/tinyagents-harness/src/tool/packs/test.rs new file mode 100644 index 000000000..5aff90aa3 --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/test.rs @@ -0,0 +1,395 @@ +//! Mechanism tests over a two-pack test catalog. The product table's own +//! invariants (membership, owners, guides) are tested by the host that owns it. + +use std::sync::Arc; + +use async_trait::async_trait; +use serde_json::{Value, json}; +use tinytools::{PermissionLevel, Tool, ToolResult, ToolSpec, ToolTimeout}; + +use super::*; + +const MARKER: &str = "[not_found]"; + +static PACKS: &[ToolPack] = &[ + ToolPack { + id: "alpha", + summary: "Alpha things.", + tools: &["a_one", "a_two"], + owners: &["alpha_agent"], + guide: "Use alpha carefully.", + }, + ToolPack { + id: "beta", + summary: "Beta things.", + tools: &["b_one"], + owners: &[], + guide: "", + }, +]; + +const CATALOG: PackCatalog = PackCatalog::new(PACKS, MARKER); + +struct FakeTool { + name: &'static str, + level: PermissionLevel, + external: bool, + timeout: ToolTimeout, +} + +impl FakeTool { + fn plain(name: &'static str) -> Self { + Self { + name, + level: PermissionLevel::ReadOnly, + external: false, + timeout: ToolTimeout::Inherit, + } + } +} + +#[async_trait] +impl Tool for FakeTool { + fn name(&self) -> &str { + self.name + } + fn description(&self) -> &str { + "fake" + } + fn parameters_schema(&self) -> Value { + json!({"type": "object", "properties": {"marker": {"type": "string"}}}) + } + async fn execute(&self, args: Value) -> anyhow::Result { + Ok(ToolResult::success(format!("{}:{}", self.name, args))) + } + fn permission_level(&self) -> PermissionLevel { + self.level + } + fn external_effect_with_args(&self, _args: &Value) -> bool { + self.external + } + fn timeout_policy(&self, _args: &Value) -> ToolTimeout { + self.timeout + } +} + +/// A durable registry holding `tools` plus a bound `use_skill`. +fn bound(tools: Vec) -> Arc>> { + let handle = PackRegistryHandle::default(); + let mut all: Vec> = tools.into_iter().map(|t| Box::new(t) as _).collect(); + all.push(Box::new(UseSkillTool::new(handle.clone(), CATALOG))); + let all = Arc::new(all); + handle.bind(Arc::downgrade(&all)); + all +} + +fn handle_of(tools: &[Box]) -> &PackRegistryHandle { + find(tools, USE_SKILL) + .host_extension() + .and_then(|any| any.downcast_ref::()) + .expect("use_skill carries its handle as a host extension") +} + +fn find<'a>(tools: &'a [Box], name: &str) -> &'a dyn Tool { + tools + .iter() + .find(|t| t.name() == name) + .map(AsRef::as_ref) + .unwrap_or_else(|| panic!("{name} missing")) +} + +#[test] +fn the_use_skill_declaration_is_byte_stable() { + let tools = bound(vec![]); + let tool = find(&tools, USE_SKILL); + assert_eq!(tool.name(), "use_skill"); + assert_eq!( + tool.description(), + "Reach a skill's tools. Their names, descriptions and argument schemas are NOT in \ + your context until you ask for them: call this with `skill` alone to see them, then \ + again with `skill` + `tool` + `args` to run one.\n\nSkills:\n\ + - `alpha` — Alpha things.\n- `beta` — Beta things.\n" + ); + assert_eq!( + tool.parameters_schema(), + json!({ + "type": "object", + "properties": { + "skill": { "type": "string", "enum": ["alpha", "beta"], "description": "Skill to read or run a tool from." }, + "tool": { "type": "string", "description": "Tool to run. Omit to list the skill's tools and their arguments instead." }, + "args": { + "type": "object", + "description": "The tool's own arguments, as documented in the listing.", + "additionalProperties": true + } + }, + "required": ["skill"] + }) + ); +} + +#[tokio::test] +async fn a_skill_alone_renders_guide_and_the_schema_of_a_present_tool() { + let tools = bound(vec![FakeTool::plain("a_one")]); + let result = find(&tools, USE_SKILL) + .execute(json!({"skill": "alpha"})) + .await + .unwrap(); + assert!(!result.is_error); + let rendered = result.text(); + assert!(rendered.starts_with("# Skill `alpha`\n\nAlpha things.\n\nUse alpha carefully.\n\n")); + assert!(rendered.contains("## `a_one`\n\nfake\n\n```json\n")); + assert!(rendered.contains(r#"{"properties":{"marker":{"type":"string"}},"type":"object"}"#)); + // A pack tool this session lacks is skipped, not fatal. + assert!(!rendered.contains("a_two")); +} + +#[tokio::test] +async fn an_unknown_skill_is_a_marked_error_listing_the_alternatives() { + let tools = bound(vec![FakeTool::plain("a_one")]); + let result = find(&tools, USE_SKILL) + .execute(json!({"skill": "nope"})) + .await + .unwrap(); + assert!(result.is_error); + let message = result.text(); + assert!(message.starts_with(&format!("{MARKER} Unknown skill `nope`. Available:\n"))); + assert!(message.contains("- `alpha` — Alpha things.")); +} + +#[tokio::test] +async fn dispatch_forwards_args_to_the_packed_tool() { + let tools = bound(vec![FakeTool::plain("a_one")]); + let result = find(&tools, USE_SKILL) + .execute(json!({"skill": "alpha", "tool": "a_one", "args": {"marker": "x"}})) + .await + .unwrap(); + assert!(!result.is_error); + assert_eq!(result.text(), r#"a_one:{"marker":"x"}"#); +} + +#[tokio::test] +async fn a_tool_from_another_skill_is_refused() { + // Cross-skill dispatch would make `skill` decoration and let a harmless + // skill reach a dangerous tool. + let tools = bound(vec![FakeTool { + level: PermissionLevel::Dangerous, + ..FakeTool::plain("b_one") + }]); + let result = find(&tools, USE_SKILL) + .execute(json!({"skill": "alpha", "tool": "b_one", "args": {}})) + .await + .unwrap(); + assert!(result.is_error, "cross-skill dispatch was admitted"); + assert!(result.text().starts_with(&format!( + "{MARKER} No tool `b_one` in skill `alpha`. Call `use_skill {{ \"skill\": \"alpha\" }}` to see what it contains." + ))); +} + +#[test] +fn permission_level_is_the_inner_tools_not_the_proxys() { + let tools = bound(vec![FakeTool { + level: PermissionLevel::Dangerous, + ..FakeTool::plain("a_one") + }]); + let use_skill = find(&tools, USE_SKILL); + assert_eq!( + use_skill.permission_level_with_args(&json!({"skill": "alpha", "tool": "a_one"})), + PermissionLevel::Dangerous + ); + // Naming no tool (or an empty one) only reads a schema. + for args in [ + json!({"skill": "alpha"}), + json!({"skill": "alpha", "tool": ""}), + ] { + assert_eq!( + use_skill.permission_level_with_args(&args), + PermissionLevel::ReadOnly + ); + } + // Unresolvable: the ceiling, never a permissive default. + assert_eq!( + use_skill.permission_level_with_args(&json!({"skill": "alpha", "tool": "ghost"})), + PermissionLevel::Dangerous + ); +} + +#[test] +fn external_effect_and_timeout_are_forwarded() { + let tools = bound(vec![FakeTool { + external: true, + timeout: ToolTimeout::Unbounded, + ..FakeTool::plain("a_one") + }]); + let use_skill = find(&tools, USE_SKILL); + let call = json!({"skill": "alpha", "tool": "a_one"}); + assert!(use_skill.external_effect_with_args(&call)); + assert_eq!(use_skill.timeout_policy(&call), ToolTimeout::Unbounded); + let ghost = json!({"skill": "alpha", "tool": "ghost"}); + assert!(!use_skill.external_effect_with_args(&ghost)); + assert_eq!(use_skill.timeout_policy(&ghost), ToolTimeout::Inherit); +} + +#[test] +fn an_unbound_handle_degrades_closed() { + let tool = UseSkillTool::new(PackRegistryHandle::default(), CATALOG); + assert_eq!(tool.permission_level(), PermissionLevel::Dangerous); +} + +#[tokio::test] +async fn an_unbound_handle_reports_the_registry_unavailable() { + let tool = UseSkillTool::new(PackRegistryHandle::default(), CATALOG); + let result = tool.execute(json!({"skill": "alpha"})).await.unwrap(); + assert!(result.is_error); + assert!(result.text().contains("registry is not available")); +} + +/// The synthesised set is a second registry; a packed tool living only there +/// must be reachable, and rebinding it must repoint the handle. +#[tokio::test] +async fn a_tool_in_the_synthesised_set_is_reachable_and_rebinding_repoints() { + let durable = bound(vec![]); + let first: Arc>> = Arc::new(vec![Box::new(FakeTool::plain("a_one"))]); + handle_of(&durable).bind_synthesized(Arc::downgrade(&first)); + let use_skill = find(&durable, USE_SKILL); + + let ran = use_skill + .execute(json!({"skill": "alpha", "tool": "a_one", "args": {}})) + .await + .unwrap(); + assert!(!ran.is_error, "{}", ran.text()); + + let second: Arc>> = Arc::new(vec![Box::new(FakeTool::plain("a_two"))]); + handle_of(&durable).bind_synthesized(Arc::downgrade(&second)); + drop(first); + let ran = use_skill + .execute(json!({"skill": "alpha", "tool": "a_two", "args": {}})) + .await + .unwrap(); + assert!(!ran.is_error, "handle did not follow the rebind"); + let stale = use_skill + .execute(json!({"skill": "alpha", "tool": "a_one", "args": {}})) + .await + .unwrap(); + assert!(stale.is_error, "a dropped tool stayed reachable"); +} + +#[test] +fn resolve_registry_for_enforces_pack_membership() { + let durable = bound(vec![FakeTool::plain("a_one")]); + let handle = handle_of(&durable); + assert!( + handle + .resolve_registry_for(&CATALOG, "alpha", "a_one") + .is_some() + ); + assert!( + handle + .resolve_registry_for(&CATALOG, "beta", "a_one") + .is_none() + ); + assert!( + handle + .resolve_registry_for(&CATALOG, "alpha", "ghost") + .is_none() + ); +} + +#[tokio::test] +async fn the_filtered_listing_omits_tools_the_session_cannot_call() { + let tools = bound(vec![FakeTool::plain("a_one"), FakeTool::plain("a_two")]); + let handle = handle_of(&tools); + let listing = render_pack_filtered(&CATALOG, "alpha", handle, &|name| name == "a_one", "") + .expect("one callable tool"); + assert!(listing.contains("## `a_one`")); + assert!(!listing.contains("## `a_two`")); + + let denied = render_pack_filtered(&CATALOG, "alpha", handle, &|_| false, "Hand off instead.") + .unwrap_err(); + assert_eq!( + denied, + "Skill `alpha` has no tools available in this session. Hand off instead." + ); + let bare = render_pack_filtered(&CATALOG, "alpha", handle, &|_| false, "").unwrap_err(); + assert_eq!( + bare, + "Skill `alpha` has no tools available in this session." + ); +} + +#[test] +fn no_such_pack_tool_renders_each_shape() { + let base = |callable: Vec<&'static str>, route: &str| NoSuchPackTool { + not_found_marker: MARKER, + skill: "alpha", + tool: "x", + callable, + route: route.to_string(), + }; + assert_eq!( + base(vec!["a_one", "a_two"], "ignored").render(), + "[not_found] There is no tool `x` in skill `alpha`. The tools in it you can call: `a_one`, `a_two`." + ); + assert_eq!( + base(vec![], "Go elsewhere.").render(), + "[not_found] There is no tool `x` in skill `alpha`. Go elsewhere." + ); + assert_eq!( + base(vec![], "").render(), + "[not_found] There is no tool `x` in skill `alpha`. Nothing in it is available in this session." + ); +} + +#[test] +fn route_sentence_prefers_a_callable_tool_and_falls_back_to_owners() { + assert_eq!( + route_sentence(&["a".to_string(), "b".to_string()], &["o"]), + "Call `a` or `b` instead — that agent owns these tools and runs them directly. Do not retry this skill." + ); + assert_eq!( + route_sentence(&[], &["x", "y"]), + "These tools belong to `x` or `y`; hand the task to one of them rather than calling directly." + ); + assert!(route_sentence(&[], &[]).is_empty()); +} + +#[test] +fn scoping_narrows_the_enum_and_the_index_or_drops_the_tool() { + let tools = bound(vec![]); + let mut spec = ToolSpec { + name: USE_SKILL.into(), + description: find(&tools, USE_SKILL).description().to_string(), + parameters: find(&tools, USE_SKILL).parameters_schema(), + }; + assert!(scope_use_skill_spec(&CATALOG, &mut spec, &|name| name == "b_one")); + assert!( + spec.description + .ends_with("Skills:\n- `beta` — Beta things.\n") + ); + assert_eq!( + spec.parameters.pointer("/properties/skill/enum"), + Some(&json!(["beta"])) + ); + assert!(!scope_use_skill_spec(&CATALOG, &mut spec, &|_| false)); +} + +#[test] +fn catalog_lookups_honour_owners() { + assert_eq!(CATALOG.pack_for_tool("b_one").map(|p| p.id), Some("beta")); + assert!(CATALOG.pack_for_tool("ghost").is_none()); + assert_eq!(CATALOG.all_packed_tool_names(), ["a_one", "a_two", "b_one"]); + // An owner keeps its belt, including a thread-renamed `_` id. + assert_eq!( + CATALOG.packed_tool_names_for_agent("alpha_agent"), + ["b_one"] + ); + assert_eq!( + CATALOG.packed_tool_names_for_agent("alpha_agent_t1"), + ["b_one"] + ); + assert_eq!( + CATALOG.packed_tool_names_for_agent("alpha_agentx"), + ["a_one", "a_two", "b_one"] + ); + assert_eq!(CATALOG.callable_pack_ids(&|n| n == "a_two"), ["alpha"]); +} diff --git a/crates/tinyagents-harness/src/tool/packs/tool.rs b/crates/tinyagents-harness/src/tool/packs/tool.rs new file mode 100644 index 000000000..61c23e738 --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/tool.rs @@ -0,0 +1,206 @@ +//! The always-on tool that stands in for every packed tool. + +use async_trait::async_trait; +use serde_json::{Value, json}; +use tinytools::{PermissionLevel, Tool, ToolCallOptions, ToolResult, ToolRunContext}; + +use super::catalog::PackCatalog; +use super::handle::{PackRegistryHandle, ToolVec}; +use super::render::{named_tool, render_pack}; + +/// Name of the disclosure-and-dispatch tool. +pub const USE_SKILL: &str = "use_skill"; + +/// The one always-on tool that stands in for every packed tool. +/// +/// It is both halves of the pack seam. Called with a `skill` alone it renders +/// that pack's tool schemas into the conversation; called with a `skill` and a +/// `tool` it executes that tool. These were two tools — `load_skill` and +/// `use_skill` — until the pack index that each carried in its own description +/// made the pair spend 3.3 kB of every single turn saying one list twice, and +/// made the first call of any packed tool a mandatory two-call round trip. +/// A model that learned the retired name gets the harness's "unknown tool" +/// recovery result (#4249) and retries against the schema it can actually see, +/// so no alias is carried for it. +/// +/// **Permission forwarding is load-bearing.** The harness gates a call on the +/// tool's `permission_level_with_args`, so a proxy reporting its own level would +/// launder every packed tool's risk down to this one's — a crypto write would +/// be admitted on a channel that refuses crypto writes. Both accessors resolve +/// the inner tool and defer to it; the arg-less one has nothing to resolve +/// from, so it reports the highest level any packed tool needs rather than +/// guessing low. A call that names no `tool` reads a schema and nothing else, +/// so that one branch is genuinely `ReadOnly`. +pub struct UseSkillTool { + handle: PackRegistryHandle, + catalog: PackCatalog, + description: String, +} + +impl UseSkillTool { + pub fn new(handle: PackRegistryHandle, catalog: PackCatalog) -> Self { + let description = format!( + "Reach a skill's tools. Their names, descriptions and argument schemas are NOT in \ + your context until you ask for them: call this with `skill` alone to see them, then \ + again with `skill` + `tool` + `args` to run one.\n\nSkills:\n{}", + catalog.pack_index_markdown() + ); + Self { + handle, + catalog, + description, + } + } + + fn skill_enum(&self) -> Vec<&'static str> { + self.catalog.packs().iter().map(|p| p.id).collect() + } + + fn resolve(&self, args: &Value) -> Option<(ToolVec, usize)> { + let skill = args.get("skill").and_then(Value::as_str)?; + let tool = named_tool(args)?; + self.handle.resolve(&self.catalog, skill, tool) + } +} + +#[async_trait] +impl Tool for UseSkillTool { + fn name(&self) -> &str { + USE_SKILL + } + + fn description(&self) -> &str { + &self.description + } + + fn parameters_schema(&self) -> Value { + json!({ + "type": "object", + "properties": { + "skill": { "type": "string", "enum": self.skill_enum(), "description": "Skill to read or run a tool from." }, + "tool": { "type": "string", "description": "Tool to run. Omit to list the skill's tools and their arguments instead." }, + "args": { + "type": "object", + "description": "The tool's own arguments, as documented in the listing.", + "additionalProperties": true + } + }, + "required": ["skill"] + }) + } + + async fn execute(&self, args: Value) -> anyhow::Result { + self.execute_with_context(args, ToolCallOptions::default(), None) + .await + } + + async fn execute_with_options( + &self, + args: Value, + options: ToolCallOptions, + ) -> anyhow::Result { + self.execute_with_context(args, options, None).await + } + + async fn execute_with_context( + &self, + args: Value, + options: ToolCallOptions, + context: Option<&dyn ToolRunContext>, + ) -> anyhow::Result { + let Some(skill) = args.get("skill").and_then(Value::as_str) else { + return Ok(ToolResult::error(format!( + "`skill` is required.\n\nSkills:\n{}", + self.catalog.pack_index_markdown() + ))); + }; + + // Disclosure half: no tool named, so render the pack's schemas. + let Some(name) = named_tool(&args) else { + return Ok(match render_pack(&self.catalog, skill, &self.handle) { + Ok(text) => ToolResult::success(text), + Err(message) => ToolResult::error(message), + }); + }; + + let Some((tools, idx)) = self.handle.resolve(&self.catalog, skill, name) else { + return Ok(ToolResult::error(format!( + "{} No tool `{name}` in skill `{skill}`. Call `use_skill {{ \"skill\": \"{skill}\" }}` \ + to see what it contains.\n\nSkills:\n{}", + self.catalog.not_found_marker(), + self.catalog.pack_index_markdown() + ))); + }; + let inner_args = args.get("args").cloned().unwrap_or_else(|| json!({})); + tracing::debug!(tool = tools[idx].name(), "[toolpacks] use_skill dispatch"); + tools[idx] + .execute_with_context(inner_args, options, context) + .await + } + + fn supports_markdown(&self) -> bool { + // The inner result is forwarded verbatim, markdown rendering included, + // so advertise the capability rather than suppressing a real saving. + true + } + + fn external_effect_with_args(&self, args: &Value) -> bool { + match self.resolve(args) { + Some((tools, idx)) => { + let inner_args = args.get("args").cloned().unwrap_or_else(|| json!({})); + tools[idx].external_effect_with_args(&inner_args) + } + None => false, + } + } + + fn timeout_policy(&self, args: &Value) -> tinytools::ToolTimeout { + match self.resolve(args) { + Some((tools, idx)) => { + let inner_args = args.get("args").cloned().unwrap_or_else(|| json!({})); + tools[idx].timeout_policy(&inner_args) + } + None => tinytools::ToolTimeout::Inherit, + } + } + + fn permission_level(&self) -> PermissionLevel { + let packed = self.catalog.all_packed_tool_names(); + self.handle + .registries() + .iter() + .flat_map(|tools| tools.iter()) + .filter(|t| packed.contains(&t.name())) + .map(|t| t.permission_level()) + .max() + // Unbound or empty: report the ceiling, never a permissive default. + .unwrap_or(PermissionLevel::Dangerous) + } + + fn permission_level_with_args(&self, args: &Value) -> PermissionLevel { + // Naming no tool renders a schema and does nothing else. Reporting the + // packed ceiling here would put an approval prompt in front of reading + // a tool list, which is the round trip this tool exists to remove. + if named_tool(args).is_none() { + return PermissionLevel::ReadOnly; + } + match self.resolve(args) { + Some((tools, idx)) => { + let inner = args.get("args").cloned().unwrap_or_else(|| json!({})); + tools[idx].permission_level_with_args(&inner) + } + // Unresolvable: the call will fail anyway, but report the ceiling so + // a malformed call can never be admitted on a channel that would + // have refused the real tool. + None => self.permission_level(), + } + } + + /// The registry handle rides on the vocabulary's erased host extension: + /// `PackRegistryHandle` is a harness concept, and `tinytools` has no + /// business naming it. A host reads it back by downcasting to + /// [`PackRegistryHandle`]. + fn host_extension(&self) -> Option<&(dyn std::any::Any + Send + Sync)> { + Some(&self.handle) + } +} diff --git a/crates/tinyagents-harness/src/tool/packs/types.rs b/crates/tinyagents-harness/src/tool/packs/types.rs new file mode 100644 index 000000000..82a606059 --- /dev/null +++ b/crates/tinyagents-harness/src/tool/packs/types.rs @@ -0,0 +1,62 @@ +//! Tool-pack types: the unit of on-demand tool disclosure. + +/// A named bundle of tools that is **not** advertised to the model by default. +/// +/// The pack's tools stay fully constructed and executable; what changes is that +/// their JSON schemas never reach the provider until the agent asks for them. +/// That trade is the whole point: an orchestrator carrying ~77 tool schemas +/// spends far more of its fixed per-turn budget on schemas than on its own +/// instructions, and most of those tools go untouched in most conversations. +#[derive(Debug)] +pub struct ToolPack { + /// Stable id the agent names in `use_skill`. + pub id: &'static str, + /// One line, rendered in the always-on pack index. This is the only text + /// about the pack the model sees before loading it, so it has to carry + /// enough intent for the model to know when to reach for it. + pub summary: &'static str, + /// Tool names this pack owns. A name listed here is removed from the + /// agent's advertised surface and reachable only through `use_skill`. + pub tools: &'static [&'static str], + /// Agent ids for which this pack is **not** applied. + /// + /// Withholding is a bet that the tools are idle in most turns. That bet is + /// wrong for the specialist a family was delegated to: `workflow_builder` + /// exists precisely to run the flow authoring tools, so packing them would + /// put a `use_skill` round trip in front of the first call of every one of + /// its turns and buy nothing — its whole belt is the pack. + /// + /// The earlier packs did not need this because they held only synthesised + /// `delegate_*` tools, which exist on the orchestrator alone. Packs over + /// raw tools do, and an owner list is the narrowest way to say so. + pub owners: &'static [&'static str], + /// The skill's playbook, printed when `use_skill` loads the pack, between + /// the summary and the tool schemas. Empty for a pack that is only a + /// schema bundle. + /// + /// This is what replaced most single-belt specialists: a sub-agent whose + /// whole value was a prompt over a handful of tools is a ~500-token guide + /// here plus `Deferred` tools the orchestrator reaches itself, instead of + /// a separate context, model call and hand-off envelope. Lookup detail + /// belongs in the guide; a rule that must bind before the model thinks to + /// load the skill (confirm before moving money) stays in the orchestrator + /// prompt. Kept under ~550 tokens by `toolpacks_tests.rs`. + pub guide: &'static str, +} + +impl ToolPack { + pub fn owns(&self, tool: &str) -> bool { + self.tools.contains(&tool) + } + + /// Whether `agent_id` owns this pack, including a web-chat thread name + /// formed as `_` after the session is built. + pub fn is_owner(&self, agent_id: &str) -> bool { + self.owners.iter().any(|owner| { + agent_id == *owner + || agent_id + .strip_prefix(*owner) + .is_some_and(|suffix| suffix.starts_with('_')) + }) + } +} diff --git a/vendor/tinyinference b/vendor/tinyinference index 98d669455..04a6c429e 160000 --- a/vendor/tinyinference +++ b/vendor/tinyinference @@ -1 +1 @@ -Subproject commit 98d6694559e52b90ceaeda713a17c1b380a153a4 +Subproject commit 04a6c429ee86dcbdb3fcb29cc7368c6f5cb1b4e2