10 Commits
Author SHA1 Message Date
darman dbc676332a M5: slash commands (command/*.md)
Adds slash-command support: markdown command files under command/*.md, expansion in harness-app, and TUI input handling so typing /mycommand runs it.
2026-07-10 16:20:22 +02:00
darman d6af65f494 M5: skills (SKILL.md and skill tool)
Adds harness-core::config::markdown to parse SKILL.md frontmatter/body, a skill tool in harness-tools to invoke them, and system-prompt wiring so available skills are advertised to the model. Includes the rustfmt pass over the M5 files touched by this and the preceding LSP/rmcp work.
2026-07-10 16:20:22 +02:00
darman fe6d00ce6d M5: rmcp stdio client and tool adapters
Adds an rmcp-based stdio MCP client in harness-mcp plus adapters exposing a configured server's tools as namespaced harness-tools, with an echo-server fixture and integration test, and wiring through harness-app.
2026-07-10 16:20:22 +02:00
darman f71a347061 M5: LSP pool and diagnostics wiring
Adds the DiagnosticsSource seam trait in harness-core::lsp (kept out of harness-tools so tools never link the LSP crate directly), and implements it in harness-lsp as a pool that lazily spawns one language server per file extension (rust-analyzer, typescript-language-server, gopls, pyright-langserver) only when the binary is on PATH. Adds harness-tools::diagnostics and wires diagnostics into edit/write/bash/glob/grep output.
2026-07-10 16:20:22 +02:00
darman e32f61e4f8 M4: dedicated opencode (Zen) provider
Adds a dedicated provider entry for opencode's Zen proxy on top of the existing OpenAI-compatible codec, with its own config wiring in harness-app and harness-core::config::load.
2026-07-10 16:20:11 +02:00
darman 9ee5a1940d M4: jobs pane, subtask drill-in, reminders
Adds the jobs pane to the TUI with a new snapshot test, subtask drill-in navigation, and periodic reminders injected into the engine loop for long-running background jobs.
2026-07-10 16:20:11 +02:00
darman ac516a7e23 M4: subagent context-file reporting to the job board
Lets subagents report context files back to the job board (surfaced from the read tool) so the parent session's next-turn context includes what a background child has been looking at.
2026-07-10 16:20:11 +02:00
darman 8976b5dc09 M4: task tool, subagent spawner, board wiring
Adds the task tool and subagent spawner in harness-tools, and wires spawning/aliasing/reuse and depth limits through harness-app and the job board, so an orchestrator session can launch and track child sessions.
2026-07-10 16:20:11 +02:00
darman 07476d1318 M4: permission intersection and job board
Adds harness-core::engine::jobs, the JobBoard tracking foreground/background subagent jobs with persistence, plus permission-rule intersection so a spawned subagent's effective permissions are the parent's rules narrowed by its own.
2026-07-10 16:20:11 +02:00
darman f3b7666d3b M4: agent registry and bundled markdown agents
Adds harness-core::agent, an agent registry loading bundled markdown agent definitions (designer, explorer, fixer, librarian, oracle, orchestrator) with override support.
2026-07-10 16:20:11 +02:00
60 changed files with 5282 additions and 109 deletions
Generated
+214 -2
View File
@@ -35,6 +35,15 @@ version = "0.2.21"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923"
[[package]]
name = "android_system_properties"
version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "819e7219dbd41043ac279b19830f2efc897156490d7fd6ea916720117ee66311"
dependencies = [
"libc",
]
[[package]]
name = "anyhow"
version = "1.0.103"
@@ -92,6 +101,18 @@ version = "1.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0"
[[package]]
name = "autocfg"
version = "1.5.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53"
[[package]]
name = "base64"
version = "0.21.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9d297deb1925b89f2ccc13d7635fa0714f12c87adce1c75356b39ca9b7178567"
[[package]]
name = "base64"
version = "0.22.1"
@@ -175,6 +196,20 @@ dependencies = [
"rand_core 0.10.1",
]
[[package]]
name = "chrono"
version = "0.4.45"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1aa79e62e7697b8e29b513a68abacf485adcd1fe8284a4316c5ae868e6633327"
dependencies = [
"iana-time-zone",
"js-sys",
"num-traits",
"serde",
"wasm-bindgen",
"windows-link",
]
[[package]]
name = "compact_str"
version = "0.8.2"
@@ -217,6 +252,12 @@ dependencies = [
"windows-sys 0.61.2",
]
[[package]]
name = "core-foundation-sys"
version = "0.8.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b"
[[package]]
name = "cpufeatures"
version = "0.3.0"
@@ -669,12 +710,15 @@ dependencies = [
name = "harness-app"
version = "0.1.0"
dependencies = [
"async-trait",
"dirs",
"futures",
"harness-core",
"harness-lsp",
"harness-mcp",
"harness-providers",
"harness-tools",
"serde_json",
"tempfile",
"thiserror 2.0.18",
"tokio",
@@ -695,6 +739,7 @@ dependencies = [
"schemars",
"serde",
"serde_json",
"serde_yaml_ng",
"tempfile",
"thiserror 2.0.18",
"tokio",
@@ -707,14 +752,28 @@ dependencies = [
name = "harness-lsp"
version = "0.1.0"
dependencies = [
"async-trait",
"harness-core",
"serde_json",
"tempfile",
"tokio",
"tracing",
]
[[package]]
name = "harness-mcp"
version = "0.1.0"
dependencies = [
"async-trait",
"harness-core",
"rmcp",
"serde",
"serde_json",
"tempfile",
"thiserror 2.0.18",
"tokio",
"tokio-util",
"tracing",
]
[[package]]
@@ -803,6 +862,12 @@ dependencies = [
"foldhash",
]
[[package]]
name = "hashbrown"
version = "0.17.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a"
[[package]]
name = "hashlink"
version = "0.9.1"
@@ -899,7 +964,7 @@ version = "0.1.20"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0"
dependencies = [
"base64",
"base64 0.22.1",
"bytes",
"futures-channel",
"futures-util",
@@ -916,6 +981,30 @@ dependencies = [
"tracing",
]
[[package]]
name = "iana-time-zone"
version = "0.1.65"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e31bc9ad994ba00e440a8aa5c9ef0ec67d5cb5e5cb0cc7f8b744a35b389cc470"
dependencies = [
"android_system_properties",
"core-foundation-sys",
"iana-time-zone-haiku",
"js-sys",
"log",
"wasm-bindgen",
"windows-core",
]
[[package]]
name = "iana-time-zone-haiku"
version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f"
dependencies = [
"cc",
]
[[package]]
name = "icu_collections"
version = "2.2.0"
@@ -1041,6 +1130,16 @@ dependencies = [
"winapi-util",
]
[[package]]
name = "indexmap"
version = "2.14.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d466e9454f08e4a911e14806c24e16fba1b4c121d1ea474396f396069cf949d9"
dependencies = [
"equivalent",
"hashbrown 0.17.1",
]
[[package]]
name = "indoc"
version = "2.0.7"
@@ -1273,6 +1372,15 @@ version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "521739c6d2bac4aa25192232afe6841231376b2b26d4d9fae5ecf8ca5772e441"
[[package]]
name = "num-traits"
version = "0.2.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841"
dependencies = [
"autocfg",
]
[[package]]
name = "once_cell"
version = "1.21.4"
@@ -1580,7 +1688,7 @@ version = "0.12.28"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147"
dependencies = [
"base64",
"base64 0.22.1",
"bytes",
"futures-core",
"futures-util",
@@ -1629,6 +1737,38 @@ dependencies = [
"windows-sys 0.52.0",
]
[[package]]
name = "rmcp"
version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "33a0110d28bd076f39e14bfd5b0340216dd18effeb5d02b43215944cc3e5c751"
dependencies = [
"base64 0.21.7",
"chrono",
"futures",
"paste",
"pin-project-lite",
"rmcp-macros",
"schemars",
"serde",
"serde_json",
"thiserror 2.0.18",
"tokio",
"tokio-util",
"tracing",
]
[[package]]
name = "rmcp-macros"
version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a6e2b2fd7497540489fa2db285edd43b7ed14c49157157438664278da6e42a7a"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "rusqlite"
version = "0.32.1"
@@ -1827,6 +1967,19 @@ dependencies = [
"serde",
]
[[package]]
name = "serde_yaml_ng"
version = "0.10.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7b4db627b98b36d4203a7b458cf3573730f2bb591b28871d916dfa9efabfd41f"
dependencies = [
"indexmap",
"itoa",
"ryu",
"serde",
"unsafe-libyaml",
]
[[package]]
name = "sharded-slab"
version = "0.1.7"
@@ -2357,6 +2510,12 @@ version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1fc81956842c57dac11422a97c3b8195a1ff727f06e85c84ed2e8aa277c9a0fd"
[[package]]
name = "unsafe-libyaml"
version = "0.2.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "673aac59facbab8a9007c7f6108d11f63b603f7cabff99fabf650fea5c32b861"
[[package]]
name = "untrusted"
version = "0.9.0"
@@ -2561,12 +2720,65 @@ version = "0.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f"
[[package]]
name = "windows-core"
version = "0.62.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb"
dependencies = [
"windows-implement",
"windows-interface",
"windows-link",
"windows-result",
"windows-strings",
]
[[package]]
name = "windows-implement"
version = "0.60.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "windows-interface"
version = "0.59.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "windows-link"
version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5"
[[package]]
name = "windows-result"
version = "0.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7781fa89eaf60850ac3d2da7af8e5242a5ea78d1a11c49bf2910bb5a73853eb5"
dependencies = [
"windows-link",
]
[[package]]
name = "windows-strings"
version = "0.5.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7837d08f69c77cf6b07689544538e017c1bfcf57e34b4c0ff58e6c2cd3b37091"
dependencies = [
"windows-link",
]
[[package]]
name = "windows-sys"
version = "0.48.0"
+3
View File
@@ -12,12 +12,15 @@ harness-mcp = { workspace = true }
harness-lsp = { workspace = true }
tokio = { workspace = true }
tokio-util = { workspace = true }
async-trait = { workspace = true }
dirs = { workspace = true }
thiserror = { workspace = true }
tracing = { workspace = true }
[dev-dependencies]
tempfile = { workspace = true }
futures = { workspace = true }
serde_json = { workspace = true }
[lints]
workspace = true
File diff suppressed because it is too large Load Diff
+1
View File
@@ -11,6 +11,7 @@ futures = { workspace = true }
async-trait = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
serde_yaml_ng = { workspace = true }
schemars = { workspace = true }
rusqlite = { workspace = true }
globset = { workspace = true }
@@ -0,0 +1,12 @@
---
description: UI and interaction design specialist for user-facing surfaces
mode: subagent
temperature: 0.4
tools: { task: false }
---
You are Designer, a UI and interaction design specialist. You handle user-facing surfaces:
layout, component structure, styling, and interaction details.
- Understand the existing design language before proposing changes; stay consistent with it.
- When implementing, make focused edits and describe the visual/interaction effect.
- Call out accessibility and responsive concerns relevant to the change.
@@ -0,0 +1,16 @@
---
description: Read-only reconnaissance specialist for mapping code and finding relevant files
mode: subagent
temperature: 0.1
tools: { write: false, edit: false, task: false }
permission:
- { permission: "edit", pattern: "*", action: deny }
- { permission: "write", pattern: "*", action: deny }
---
You are Explorer, a read-only reconnaissance specialist. You locate the code, files, and
facts the orchestrator needs and report back concisely.
- Use read, glob, grep, and bash (read-only commands) to investigate.
- Never modify files. Return a focused summary with concrete `path:line` references, not
file dumps.
- State what you found and, briefly, what you could not find.
@@ -0,0 +1,13 @@
---
description: Implementation specialist that makes focused code changes and verifies them
mode: subagent
temperature: 0.2
tools: { task: false }
---
You are Fixer, an implementation specialist. You take a well-scoped change, implement it,
and verify it compiles/tests.
- Make the smallest change that satisfies the request; match the surrounding code's style.
- Use read/grep to understand context before editing; use bash to build and run tests.
- Report exactly what you changed (files and the essence of the diff) and the result of any
verification you ran.
@@ -0,0 +1,15 @@
---
description: Documentation and knowledge lookup specialist
mode: subagent
temperature: 0.1
tools: { write: false, edit: false, task: false }
permission:
- { permission: "edit", pattern: "*", action: deny }
- { permission: "write", pattern: "*", action: deny }
---
You are Librarian. You find and summarize documentation, comments, READMEs, config, and
other in-repo knowledge on request.
- Search docs and source for the relevant material with read, glob, and grep.
- Quote the authoritative source with its `path:line`; do not invent details.
- Return a concise, well-organized summary with pointers back to the sources.
@@ -0,0 +1,15 @@
---
description: Deep-reasoning analyst for architecture, debugging, and design trade-offs
mode: subagent
temperature: 0.3
tools: { write: false, edit: false, task: false }
permission:
- { permission: "edit", pattern: "*", action: deny }
- { permission: "write", pattern: "*", action: deny }
---
You are Oracle, a deep-reasoning analyst. You are consulted for hard questions:
root-causing bugs, weighing architectural trade-offs, and reviewing designs.
- Read whatever code and context you need, but do not modify anything.
- Reason carefully and explicitly; state assumptions and the evidence behind conclusions.
- Return a decisive recommendation with the reasoning that supports it.
@@ -0,0 +1,19 @@
---
description: Primary coordinator that plans work and delegates to specialists
mode: primary
temperature: 0.2
---
You are the orchestrator. You plan the work, delegate focused pieces to specialist
subagents via the `task` tool, and synthesize their results into a final answer.
Guidelines:
- Break the request into concrete, independently-verifiable pieces.
- Prefer delegating reconnaissance and analysis to subagents so your own context stays
focused; do the integration and final write-up yourself.
- Launch background tasks for long-running independent work, then continue planning.
Do not poll running jobs — wait for completion and reconcile terminal jobs before your
final response.
- Reuse a completed specialist session (by its job alias) when following up on the same
thread of work.
{{SUBAGENTS}}
+496
View File
@@ -0,0 +1,496 @@
//! Agent definitions and registry.
//!
//! All agent *behavior* lives in markdown + config — the engine only understands `mode`, tool
//! filters, permissions, model, and depth. Definitions are layered (bundled → global → project
//! → config patch); later layers win by name. See `docs/04-multiagent.md`.
use std::collections::HashMap;
use std::path::Path;
use serde::Deserialize;
use crate::config::AgentPatch;
use crate::permission::{Rule, Ruleset};
use crate::types::ModelRef;
/// Marker in a primary agent's prompt replaced at load time with the routing list of enabled
/// subagents (name + description). Keeps routing text in sync with the enabled agent set.
const SUBAGENTS_MARKER: &str = "{{SUBAGENTS}}";
const BUNDLED: &[(&str, &str)] = &[
(
"orchestrator",
include_str!("../../assets/agents/orchestrator.md"),
),
("explorer", include_str!("../../assets/agents/explorer.md")),
("oracle", include_str!("../../assets/agents/oracle.md")),
(
"librarian",
include_str!("../../assets/agents/librarian.md"),
),
("fixer", include_str!("../../assets/agents/fixer.md")),
("designer", include_str!("../../assets/agents/designer.md")),
];
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum AgentMode {
Primary,
#[default]
Subagent,
All,
}
impl AgentMode {
fn parse(s: &str) -> Option<Self> {
match s.trim().to_ascii_lowercase().as_str() {
"primary" => Some(Self::Primary),
"subagent" => Some(Self::Subagent),
"all" => Some(Self::All),
_ => None,
}
}
/// Whether this agent can be invoked as a subagent via the `task` tool.
pub fn is_subagent(self) -> bool {
matches!(self, Self::Subagent | Self::All)
}
/// Whether this agent can drive a top-level (primary) session.
pub fn is_primary(self) -> bool {
matches!(self, Self::Primary | Self::All)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AgentSource {
Bundled,
Global,
Project,
Config,
}
#[derive(Debug, Clone)]
pub struct AgentDef {
pub name: String,
pub description: String,
pub mode: AgentMode,
/// `None` = follow the session model.
pub model: Option<ModelRef>,
pub temperature: Option<f32>,
pub prompt: String,
pub permissions: Ruleset,
/// Tool enable/disable overrides (wildcard keys allowed); absent = inherit default.
pub tools: HashMap<String, bool>,
pub max_steps: Option<u32>,
pub source: AgentSource,
}
#[derive(Debug, thiserror::Error)]
pub enum AgentError {
#[error("agent {0}: {1}")]
Frontmatter(String, String),
}
/// YAML frontmatter shape (opencode-compatible).
#[derive(Debug, Default, Deserialize)]
#[serde(default)]
struct Frontmatter {
description: Option<String>,
mode: Option<String>,
model: Option<String>,
temperature: Option<f32>,
tools: Option<HashMap<String, bool>>,
permission: Option<Vec<Rule>>,
max_steps: Option<u32>,
disable: Option<bool>,
}
fn parse_model_ref(s: &str) -> Option<ModelRef> {
s.split_once('/')
.map(|(p, m)| ModelRef::new(p.trim(), m.trim()))
}
/// Splits a markdown agent file into (frontmatter, body). A file without a leading `---`
/// fence is treated as an all-body prompt with empty frontmatter.
fn split_frontmatter(content: &str) -> (&str, &str) {
let rest = match content
.strip_prefix("---\n")
.or_else(|| content.strip_prefix("---\r\n"))
{
Some(r) => r,
None => return ("", content),
};
// Find the closing fence line.
for delim in ["\n---\n", "\n---\r\n"] {
if let Some(idx) = rest.find(delim) {
let body_start = idx + delim.len();
return (&rest[..idx], &rest[body_start..]);
}
}
// Trailing fence with no body / no trailing newline.
if let Some(fm) = rest.strip_suffix("\n---").or(Some(rest)) {
if rest.ends_with("\n---") {
return (fm, "");
}
}
("", content)
}
/// Parses one markdown agent definition. Returns `Ok(None)` when the file marks itself
/// `disable: true`.
fn parse_agent(
name: &str,
source: AgentSource,
content: &str,
) -> Result<Option<AgentDef>, AgentError> {
let (fm_raw, body) = split_frontmatter(content);
let fm: Frontmatter = if fm_raw.trim().is_empty() {
Frontmatter::default()
} else {
serde_yaml_ng::from_str(fm_raw)
.map_err(|e| AgentError::Frontmatter(name.to_string(), e.to_string()))?
};
if fm.disable == Some(true) {
return Ok(None);
}
Ok(Some(AgentDef {
name: name.to_string(),
description: fm.description.unwrap_or_default(),
mode: fm
.mode
.as_deref()
.and_then(AgentMode::parse)
.unwrap_or_default(),
model: fm.model.as_deref().and_then(parse_model_ref),
temperature: fm.temperature,
prompt: body.trim_end().to_string(),
permissions: fm.permission.unwrap_or_default(),
tools: fm.tools.unwrap_or_default(),
max_steps: fm.max_steps,
source,
}))
}
/// Applies a config `AgentPatch` onto an existing definition (only set fields override).
fn apply_patch(def: &mut AgentDef, patch: &AgentPatch) {
if let Some(mode) = patch.mode.as_deref().and_then(AgentMode::parse) {
def.mode = mode;
}
if let Some(model) = patch.model.as_deref().and_then(parse_model_ref) {
def.model = Some(model);
}
if let Some(temp) = patch.temperature {
def.temperature = Some(temp);
}
if let Some(prompt) = &patch.prompt {
def.prompt = prompt.clone();
}
if let Some(tools) = &patch.tools {
def.tools.extend(tools.clone());
}
if let Some(permission) = &patch.permission {
def.permissions = permission.clone();
}
}
#[derive(Debug, Clone, Default)]
pub struct AgentRegistry {
agents: HashMap<String, AgentDef>,
}
impl AgentRegistry {
pub fn get(&self, name: &str) -> Option<&AgentDef> {
self.agents.get(name)
}
pub fn all(&self) -> Vec<&AgentDef> {
self.agents.values().collect()
}
pub fn len(&self) -> usize {
self.agents.len()
}
pub fn is_empty(&self) -> bool {
self.agents.is_empty()
}
/// Just the bundled agents — the default when no overrides are configured (tests, headless).
pub fn bundled() -> Self {
let mut reg = Self::default();
reg.load_markdown_layer(
BUNDLED.iter().map(|(n, c)| (n.to_string(), *c)),
AgentSource::Bundled,
);
reg.generate_routing();
reg
}
/// Full layered load: bundled → global dir → project dir → config patches.
pub fn load(
config_agents: &HashMap<String, AgentPatch>,
global_dir: Option<&Path>,
project_dir: Option<&Path>,
) -> Self {
let mut reg = Self::default();
reg.load_markdown_layer(
BUNDLED.iter().map(|(n, c)| (n.to_string(), *c)),
AgentSource::Bundled,
);
if let Some(dir) = global_dir {
reg.load_dir(dir, AgentSource::Global);
}
if let Some(dir) = project_dir {
reg.load_dir(dir, AgentSource::Project);
}
reg.apply_config(config_agents);
reg.generate_routing();
reg
}
fn load_markdown_layer<I, S>(&mut self, files: I, source: AgentSource)
where
I: IntoIterator<Item = (String, S)>,
S: AsRef<str>,
{
for (name, content) in files {
match parse_agent(&name, source, content.as_ref()) {
Ok(Some(def)) => {
self.agents.insert(name, def);
}
Ok(None) => {
self.agents.remove(&name); // disable: true removes an earlier layer
}
Err(e) => tracing::warn!(error = %e, "skipping malformed agent"),
}
}
}
fn load_dir(&mut self, dir: &Path, source: AgentSource) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
let mut files: Vec<(String, String)> = Vec::new();
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("md") {
continue;
}
let Some(name) = path.file_stem().and_then(|s| s.to_str()) else {
continue;
};
if let Ok(content) = std::fs::read_to_string(&path) {
files.push((name.to_string(), content));
}
}
files.sort();
self.load_markdown_layer(files, source);
}
fn apply_config(&mut self, config_agents: &HashMap<String, AgentPatch>) {
for (name, patch) in config_agents {
if patch.disable == Some(true) {
self.agents.remove(name);
continue;
}
if let Some(def) = self.agents.get_mut(name) {
apply_patch(def, patch);
} else if patch.model.is_some() || patch.prompt.is_some() {
// Unknown name with enough to stand on its own → custom agent.
let mut def = AgentDef {
name: name.clone(),
description: String::new(),
mode: AgentMode::default(),
model: None,
temperature: None,
prompt: String::new(),
permissions: Vec::new(),
tools: HashMap::new(),
max_steps: None,
source: AgentSource::Config,
};
apply_patch(&mut def, patch);
self.agents.insert(name.clone(), def);
}
}
}
/// Replaces `{{SUBAGENTS}}` in every primary agent's prompt with a generated routing list
/// of the enabled subagents (so disabling an agent removes it from routing text).
fn generate_routing(&mut self) {
let mut subagents: Vec<(String, String)> = self
.agents
.values()
.filter(|a| a.mode.is_subagent())
.map(|a| (a.name.clone(), a.description.clone()))
.collect();
subagents.sort();
let routing = if subagents.is_empty() {
"## Agents\n\nNo specialist subagents are available.".to_string()
} else {
let mut s = String::from("## Agents\n\nDelegate to these specialists via `task`:\n");
for (name, desc) in &subagents {
s.push_str(&format!("- **{name}** — {desc}\n"));
}
s
};
for agent in self.agents.values_mut() {
if agent.prompt.contains(SUBAGENTS_MARKER) {
agent.prompt = agent.prompt.replace(SUBAGENTS_MARKER, routing.trim_end());
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::permission::Action;
#[test]
fn bundled_loads_all_six_agents() {
let reg = AgentRegistry::bundled();
for name in [
"orchestrator",
"explorer",
"oracle",
"librarian",
"fixer",
"designer",
] {
assert!(reg.get(name).is_some(), "missing {name}");
}
assert_eq!(reg.len(), 6);
}
#[test]
fn parses_mode_model_temperature_tools_and_permissions() {
let md = "---\n\
description: test agent\n\
mode: subagent\n\
model: anthropic/claude-haiku-4-5\n\
temperature: 0.1\n\
tools: { write: false, bash: true }\n\
permission:\n\
\x20 - { permission: \"edit\", pattern: \"*\", action: deny }\n\
---\n\
You are a test agent.\n";
let def = parse_agent("tester", AgentSource::Bundled, md)
.unwrap()
.unwrap();
assert_eq!(def.description, "test agent");
assert_eq!(def.mode, AgentMode::Subagent);
assert_eq!(
def.model,
Some(ModelRef::new("anthropic", "claude-haiku-4-5"))
);
assert_eq!(def.temperature, Some(0.1));
assert_eq!(def.tools.get("write"), Some(&false));
assert_eq!(def.tools.get("bash"), Some(&true));
assert_eq!(def.permissions.len(), 1);
assert_eq!(def.permissions[0].action, Action::Deny);
assert_eq!(def.prompt, "You are a test agent.");
}
#[test]
fn body_without_frontmatter_is_all_prompt() {
let def = parse_agent("x", AgentSource::Global, "just a prompt")
.unwrap()
.unwrap();
assert_eq!(def.prompt, "just a prompt");
assert_eq!(def.mode, AgentMode::Subagent); // default
}
#[test]
fn disable_true_removes_agent() {
assert!(
parse_agent("x", AgentSource::Config, "---\ndisable: true\n---\nbody")
.unwrap()
.is_none()
);
}
#[test]
fn explorer_denies_writes_and_disables_edit_tool() {
let reg = AgentRegistry::bundled();
let explorer = reg.get("explorer").unwrap();
assert_eq!(explorer.mode, AgentMode::Subagent);
assert_eq!(explorer.tools.get("edit"), Some(&false));
assert!(explorer
.permissions
.iter()
.any(|r| r.permission == "write" && r.action == Action::Deny));
}
#[test]
fn orchestrator_routing_lists_subagents_and_drops_marker() {
let reg = AgentRegistry::bundled();
let prompt = &reg.get("orchestrator").unwrap().prompt;
assert!(!prompt.contains("{{SUBAGENTS}}"));
assert!(prompt.contains("explorer"));
assert!(prompt.contains("fixer"));
// The orchestrator itself is primary and must not list itself.
assert!(!prompt.contains("- **orchestrator**"));
}
#[test]
fn config_patch_overrides_only_set_fields() {
let mut patches = HashMap::new();
patches.insert(
"explorer".to_string(),
AgentPatch {
temperature: Some(0.9),
..Default::default()
},
);
let reg = AgentRegistry::load(&patches, None, None);
let explorer = reg.get("explorer").unwrap();
assert_eq!(explorer.temperature, Some(0.9)); // overridden
assert_eq!(explorer.mode, AgentMode::Subagent); // untouched
assert_eq!(explorer.tools.get("edit"), Some(&false)); // untouched
}
#[test]
fn config_disable_removes_and_unknown_with_model_creates() {
let mut patches = HashMap::new();
patches.insert(
"designer".to_string(),
AgentPatch {
disable: Some(true),
..Default::default()
},
);
patches.insert(
"custom".to_string(),
AgentPatch {
model: Some("openai/gpt-5".into()),
prompt: Some("custom prompt".into()),
..Default::default()
},
);
let reg = AgentRegistry::load(&patches, None, None);
assert!(reg.get("designer").is_none());
let custom = reg.get("custom").unwrap();
assert_eq!(custom.source, AgentSource::Config);
assert_eq!(custom.model, Some(ModelRef::new("openai", "gpt-5")));
assert_eq!(custom.prompt, "custom prompt");
}
#[test]
fn unknown_config_agent_without_model_or_prompt_is_ignored() {
let mut patches = HashMap::new();
patches.insert(
"ghost".to_string(),
AgentPatch {
temperature: Some(0.5),
..Default::default()
},
);
let reg = AgentRegistry::load(&patches, None, None);
assert!(reg.get("ghost").is_none());
}
}
+1
View File
@@ -29,6 +29,7 @@ const ENV_CONFIG_PATH: &str = "AI_HARNESS_CONFIG";
const PROVIDER_ENV_KEYS: &[(&str, &str)] = &[
("anthropic", "ANTHROPIC_API_KEY"),
("openai", "OPENAI_API_KEY"),
("opencode", "OPENCODE_API_KEY"),
];
fn read_jsonc(path: &Path) -> Result<Value, ConfigError> {
+296
View File
@@ -0,0 +1,296 @@
//! Markdown-defined commands and skills (M5). Both are YAML-frontmatter + body files loaded
//! from a global dir (`~/.config/ai-harness/<kind>/`) and a project dir (`<project>/.harness/
//! <kind>/`), project winning by name. See `docs/06-config.md` and `docs/09-integrations.md`.
//!
//! - Commands (`command/*.md`) are a pure input-layer concern: `/name args` expands the body
//! template (`$ARGUMENTS`, `$1..$9`) into the user message, optionally switching agent/model.
//! - Skills (`skill/<name>/SKILL.md`) advertise `name + description` in the system prompt; the
//! model pulls a skill's body on demand via the built-in `skill` tool.
use std::collections::HashMap;
use std::path::Path;
use serde::Deserialize;
/// A slash command: a named prompt template with an optional agent/model override.
#[derive(Debug, Clone, PartialEq)]
pub struct CommandDef {
/// Invocation name (the file stem); used as `/name`.
pub name: String,
pub description: String,
/// Run the expanded prompt under this agent instead of the session's default.
pub agent: Option<String>,
/// Run under this `provider/model` instead of the session's default.
pub model: Option<String>,
/// The body, with `$ARGUMENTS` / `$1..$9` placeholders.
pub template: String,
}
impl CommandDef {
/// Substitutes `$ARGUMENTS` (the whole argument string) and `$1..$9` (whitespace-split
/// positionals; missing ones become empty) into the template.
pub fn expand(&self, arguments: &str) -> String {
let positionals: Vec<&str> = arguments.split_whitespace().collect();
let mut out = self.template.replace("$ARGUMENTS", arguments);
for i in 1..=9 {
let value = positionals.get(i - 1).copied().unwrap_or("");
out = out.replace(&format!("${i}"), value);
}
out
}
}
/// A skill: advertised by `name + description`, body loaded on demand by the `skill` tool.
#[derive(Debug, Clone, PartialEq)]
pub struct SkillDef {
/// Skill name (the containing directory name); used as the `skill` tool argument.
pub name: String,
pub description: String,
pub body: String,
}
/// Frontmatter fields shared by commands (all optional).
#[derive(Debug, Default, Deserialize)]
#[serde(default)]
struct CommandFrontmatter {
description: Option<String>,
agent: Option<String>,
model: Option<String>,
}
#[derive(Debug, Default, Deserialize)]
#[serde(default)]
struct SkillFrontmatter {
description: Option<String>,
}
/// Splits `---\n…\n---\n` frontmatter from the body. A file with no leading fence is all body.
fn split_frontmatter(content: &str) -> (&str, &str) {
let rest = match content
.strip_prefix("---\n")
.or_else(|| content.strip_prefix("---\r\n"))
{
Some(r) => r,
None => return ("", content),
};
for delim in ["\n---\n", "\n---\r\n"] {
if let Some(idx) = rest.find(delim) {
return (&rest[..idx], &rest[idx + delim.len()..]);
}
}
if rest.ends_with("\n---") {
return (rest.trim_end_matches("\n---"), "");
}
("", content)
}
fn parse_command(name: &str, content: &str) -> CommandDef {
let (fm_raw, body) = split_frontmatter(content);
let fm: CommandFrontmatter = if fm_raw.trim().is_empty() {
CommandFrontmatter::default()
} else {
serde_yaml_ng::from_str(fm_raw).unwrap_or_default()
};
CommandDef {
name: name.to_string(),
description: fm.description.unwrap_or_default(),
agent: fm.agent,
model: fm.model,
template: body.trim().to_string(),
}
}
fn parse_skill(name: &str, content: &str) -> SkillDef {
let (fm_raw, body) = split_frontmatter(content);
let fm: SkillFrontmatter = if fm_raw.trim().is_empty() {
SkillFrontmatter::default()
} else {
serde_yaml_ng::from_str(fm_raw).unwrap_or_default()
};
SkillDef {
name: name.to_string(),
description: fm.description.unwrap_or_default(),
body: body.trim().to_string(),
}
}
/// Loads `command/*.md` from global then project dirs (project wins by name).
pub fn load_commands(
global_dir: Option<&Path>,
project_dir: Option<&Path>,
) -> HashMap<String, CommandDef> {
let mut commands = HashMap::new();
for dir in [global_dir, project_dir].into_iter().flatten() {
let Ok(entries) = std::fs::read_dir(dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("md") {
continue;
}
let Some(name) = path.file_stem().and_then(|s| s.to_str()) else {
continue;
};
if let Ok(content) = std::fs::read_to_string(&path) {
commands.insert(name.to_string(), parse_command(name, &content));
}
}
}
commands
}
/// Loads `skill/<name>/SKILL.md` from global then project dirs (project wins by name),
/// returned name-sorted for a stable system-prompt listing.
pub fn load_skills(global_dir: Option<&Path>, project_dir: Option<&Path>) -> Vec<SkillDef> {
let mut by_name: HashMap<String, SkillDef> = HashMap::new();
for dir in [global_dir, project_dir].into_iter().flatten() {
let Ok(entries) = std::fs::read_dir(dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if !path.is_dir() {
continue;
}
let Some(name) = path.file_name().and_then(|s| s.to_str()) else {
continue;
};
let skill_file = path.join("SKILL.md");
if let Ok(content) = std::fs::read_to_string(&skill_file) {
by_name.insert(name.to_string(), parse_skill(name, &content));
}
}
}
let mut skills: Vec<SkillDef> = by_name.into_values().collect();
skills.sort_by(|a, b| a.name.cmp(&b.name));
skills
}
/// The system-prompt section advertising available skills (name + description). `None` when
/// there are no skills, so no empty section is injected.
pub fn skills_prompt(skills: &[SkillDef]) -> Option<String> {
if skills.is_empty() {
return None;
}
let mut section = String::from(
"## Skills\n\nThese skills are available. Load a skill's full instructions on demand \
by calling the `skill` tool with its name before doing the related work:\n",
);
for skill in skills {
section.push_str(&format!("- **{}** — {}\n", skill.name, skill.description));
}
Some(section.trim_end().to_string())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn expand_substitutes_arguments_and_positionals() {
let cmd = CommandDef {
name: "greet".into(),
description: String::new(),
agent: None,
model: None,
template: "Say $1 to $2. All: $ARGUMENTS".into(),
};
assert_eq!(cmd.expand("hi there"), "Say hi to there. All: hi there");
// Missing positionals collapse to empty.
assert_eq!(cmd.expand("solo"), "Say solo to . All: solo");
}
#[test]
fn parse_command_reads_frontmatter_and_body() {
let md = "---\ndescription: review a PR\nagent: oracle\nmodel: openai/gpt-5\n---\nReview $ARGUMENTS please.\n";
let cmd = parse_command("review", md);
assert_eq!(cmd.description, "review a PR");
assert_eq!(cmd.agent.as_deref(), Some("oracle"));
assert_eq!(cmd.model.as_deref(), Some("openai/gpt-5"));
assert_eq!(cmd.template, "Review $ARGUMENTS please.");
}
#[test]
fn parse_command_without_frontmatter_is_all_template() {
let cmd = parse_command("x", "just do $1");
assert_eq!(cmd.template, "just do $1");
assert!(cmd.agent.is_none());
}
#[test]
fn parse_skill_reads_description_and_body() {
let md = "---\ndescription: format code\n---\nRun the formatter.\n";
let skill = parse_skill("formatter", md);
assert_eq!(skill.name, "formatter");
assert_eq!(skill.description, "format code");
assert_eq!(skill.body, "Run the formatter.");
}
#[test]
fn skills_prompt_lists_each_and_is_none_when_empty() {
assert!(skills_prompt(&[]).is_none());
let skills = vec![
SkillDef {
name: "a".into(),
description: "does a".into(),
body: "".into(),
},
SkillDef {
name: "b".into(),
description: "does b".into(),
body: "".into(),
},
];
let prompt = skills_prompt(&skills).unwrap();
assert!(prompt.contains("- **a** — does a"));
assert!(prompt.contains("- **b** — does b"));
}
#[test]
fn load_skills_reads_skill_dirs_and_project_wins() {
let dir = tempfile::tempdir().unwrap();
let global = dir.path().join("global");
let project = dir.path().join("project");
std::fs::create_dir_all(global.join("fmt")).unwrap();
std::fs::create_dir_all(project.join("fmt")).unwrap();
std::fs::create_dir_all(global.join("lint")).unwrap();
std::fs::write(
global.join("fmt/SKILL.md"),
"---\ndescription: global fmt\n---\nglobal body",
)
.unwrap();
std::fs::write(
project.join("fmt/SKILL.md"),
"---\ndescription: project fmt\n---\nproject body",
)
.unwrap();
std::fs::write(
global.join("lint/SKILL.md"),
"---\ndescription: lint\n---\nlint body",
)
.unwrap();
let skills = load_skills(Some(&global), Some(&project));
assert_eq!(skills.len(), 2);
// Sorted by name: fmt, lint.
assert_eq!(skills[0].name, "fmt");
assert_eq!(skills[0].description, "project fmt"); // project overrode global
assert_eq!(skills[0].body, "project body");
assert_eq!(skills[1].name, "lint");
}
#[test]
fn load_commands_project_overrides_global() {
let dir = tempfile::tempdir().unwrap();
let global = dir.path().join("g");
let project = dir.path().join("p");
std::fs::create_dir_all(&global).unwrap();
std::fs::create_dir_all(&project).unwrap();
std::fs::write(global.join("deploy.md"), "global deploy").unwrap();
std::fs::write(project.join("deploy.md"), "project deploy").unwrap();
let commands = load_commands(Some(&global), Some(&project));
assert_eq!(commands.len(), 1);
assert_eq!(commands["deploy"].template, "project deploy");
}
}
+2
View File
@@ -1,7 +1,9 @@
pub mod load;
pub mod markdown;
pub mod schema;
pub use load::{load, ConfigError};
pub use markdown::{load_commands, load_skills, skills_prompt, CommandDef, SkillDef};
pub use schema::{
AgentPatch, Config, LspServerConfig, McpServerConfig, OrchestrationConfig, ProviderConfig,
TuiConfig,
+579
View File
@@ -0,0 +1,579 @@
//! Background job board — tracks subagent tasks spawned via the `task` tool so the
//! orchestrator can see running work, reconcile terminal results, and reuse completed
//! child sessions by alias. Simplified native port of oh-my-opencode-slim's
//! `background-job-board.ts`. See `docs/04-multiagent.md`.
//!
//! The board is an in-memory `RwLock<HashMap>` mirrored to the `job` table so it survives a
//! resume. All mutations persist through the passed-in `Store`.
use std::collections::HashMap;
use std::sync::RwLock;
use serde::{Deserialize, Serialize};
use crate::event::{AppEvent, EventBus, JobRecordEvent};
use crate::store::{Store, StoreError};
use crate::types::SessionId;
/// A file a child session read, surfaced on the board so the orchestrator knows what a
/// completed specialist already looked at.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ContextFile {
pub path: String,
pub lines: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JobState {
Running,
Completed,
Error,
Cancelled,
}
impl JobState {
/// Terminal jobs are candidates for reconciliation; a completed one is reusable.
pub fn is_terminal(self) -> bool {
!matches!(self, JobState::Running)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct JobRecord {
pub task_id: String,
/// Human-friendly handle: first 3 chars of the agent name + a per-agent counter (`exp-1`).
pub alias: String,
pub parent_session: SessionId,
pub child_session: SessionId,
pub agent: String,
pub description: String,
pub objective: Option<String>,
pub state: JobState,
/// Whether the orchestrator has already seen this job's terminal result.
pub reconciled: bool,
pub result_summary: Option<String>,
pub context_files: Vec<ContextFile>,
pub launched_at: i64,
pub updated_at: i64,
pub last_used_at: i64,
}
/// Parameters for registering a newly launched task on the board.
pub struct LaunchSpec {
pub task_id: String,
pub parent_session: SessionId,
pub child_session: SessionId,
pub agent: String,
pub description: String,
pub objective: Option<String>,
}
/// Result summaries are truncated to this many chars on the board.
const SUMMARY_MAX: usize = 2000;
/// Context files shown per job in the prompt injection.
const CONTEXT_FILES_SHOWN: usize = 8;
/// A read must cover at least this many lines to be worth reporting to the board.
pub const MIN_REPORTED_LINES: u32 = 10;
pub struct JobBoard {
store: Store,
bus: EventBus,
jobs: RwLock<HashMap<String, JobRecord>>,
max_reusable_per_agent: u32,
}
impl JobBoard {
/// Builds a board and loads any persisted jobs for `parent_session`'s tree from the store.
pub async fn load(
store: Store,
bus: EventBus,
parent_session: &SessionId,
max_reusable_per_agent: u32,
) -> Result<Self, StoreError> {
let existing = store.jobs_for_parent(parent_session.clone()).await?;
let mut jobs = HashMap::new();
for job in existing {
jobs.insert(job.task_id.clone(), job);
}
Ok(Self {
store,
bus,
jobs: RwLock::new(jobs),
max_reusable_per_agent,
})
}
/// Assigns the next alias for `agent` under this board: `<3-char-prefix>-<n>`.
fn next_alias(&self, agent: &str) -> String {
let prefix: String = agent.chars().take(3).collect();
let prefix = if prefix.is_empty() {
"job".to_string()
} else {
prefix.to_ascii_lowercase()
};
let n = self
.jobs
.read()
.unwrap()
.values()
.filter(|j| j.agent == agent)
.count()
+ 1;
format!("{prefix}-{n}")
}
/// Registers a freshly launched task and returns its assigned alias.
pub async fn register_launch(&self, spec: LaunchSpec, now: i64) -> Result<String, StoreError> {
let alias = self.next_alias(&spec.agent);
let record = JobRecord {
task_id: spec.task_id,
alias: alias.clone(),
parent_session: spec.parent_session,
child_session: spec.child_session,
agent: spec.agent,
description: spec.description,
objective: spec.objective,
state: JobState::Running,
reconciled: false,
result_summary: None,
context_files: Vec::new(),
launched_at: now,
updated_at: now,
last_used_at: now,
};
self.upsert(record).await?;
Ok(alias)
}
/// Marks a job terminal with an optional result summary (truncated).
pub async fn finish(
&self,
task_id: &str,
state: JobState,
result_summary: Option<String>,
now: i64,
) -> Result<(), StoreError> {
let Some(mut record) = self.get(task_id) else {
return Ok(());
};
record.state = state;
record.result_summary = result_summary.map(|s| truncate_summary(&s));
record.updated_at = now;
record.last_used_at = now;
self.upsert(record).await?;
if state == JobState::Completed {
self.trim_reusable(now).await?;
}
Ok(())
}
/// Records a file a child session read (deduping by path, keeping the largest read).
pub async fn report_context_file(
&self,
task_id: &str,
path: String,
lines: u32,
now: i64,
) -> Result<(), StoreError> {
if lines < MIN_REPORTED_LINES {
return Ok(());
}
let Some(mut record) = self.get(task_id) else {
return Ok(());
};
match record.context_files.iter_mut().find(|f| f.path == path) {
Some(existing) => existing.lines = existing.lines.max(lines),
None => record.context_files.push(ContextFile { path, lines }),
}
record.updated_at = now;
self.upsert(record).await
}
/// Resolves an alias or task id to a job for the given parent (reuse lookup). Only
/// completed (reusable) jobs match.
pub fn resolve_reusable(&self, parent: &SessionId, alias_or_id: &str) -> Option<JobRecord> {
self.jobs
.read()
.unwrap()
.values()
.find(|j| {
&j.parent_session == parent
&& j.state == JobState::Completed
&& (j.alias == alias_or_id || j.task_id == alias_or_id)
})
.cloned()
}
/// Bumps `last_used_at` when a completed session is reused, keeping it fresh in the LRU.
pub async fn touch(&self, task_id: &str, now: i64) -> Result<(), StoreError> {
let Some(mut record) = self.get(task_id) else {
return Ok(());
};
record.last_used_at = now;
self.upsert(record).await
}
/// Marks all terminal jobs `reconciled` — the orchestrator has now seen them.
pub async fn reconcile_terminal(&self, now: i64) -> Result<(), StoreError> {
let to_update: Vec<JobRecord> = {
let jobs = self.jobs.read().unwrap();
jobs.values()
.filter(|j| j.state.is_terminal() && !j.reconciled)
.cloned()
.collect()
};
for mut record in to_update {
record.reconciled = true;
record.updated_at = now;
self.upsert(record).await?;
}
Ok(())
}
pub fn is_empty(&self) -> bool {
self.jobs.read().unwrap().is_empty()
}
pub fn snapshot(&self) -> Vec<JobRecord> {
self.jobs.read().unwrap().values().cloned().collect()
}
/// Renders the board as a synthetic prompt block, or `None` when there is nothing to show.
/// Mirrors slim's `formatForPrompt`.
pub fn format_for_prompt(&self) -> Option<String> {
let jobs = self.jobs.read().unwrap();
if jobs.is_empty() {
return None;
}
let mut active: Vec<&JobRecord> = jobs
.values()
.filter(|j| !j.state.is_terminal() || !j.reconciled)
.collect();
let mut reusable: Vec<&JobRecord> = jobs
.values()
.filter(|j| j.state == JobState::Completed && j.reconciled)
.collect();
active.sort_by(|a, b| a.alias.cmp(&b.alias));
reusable.sort_by(|a, b| a.alias.cmp(&b.alias));
if active.is_empty() && reusable.is_empty() {
return None;
}
let mut out = String::from(
"### Background Job Board\n\
Do not poll running jobs; wait for completion. Reconcile terminal jobs before your \
final response.\nCompleted sessions are reusable by alias for the same specialist.\n",
);
if !active.is_empty() {
out.push_str("\n#### Active / Unreconciled\n");
for job in &active {
let state = state_label(job.state);
out.push_str(&format!(
"- {} / {} / {} / {state}",
job.alias, job.child_session, job.agent
));
if let Some(obj) = &job.objective {
out.push_str(&format!(" — Objective: {obj}"));
}
out.push('\n');
if let Some(summary) = &job.result_summary {
out.push_str(&format!(" Result: {summary}\n"));
}
}
}
if !reusable.is_empty() {
out.push_str("\n#### Reusable Sessions\n");
for job in &reusable {
out.push_str(&format!(
"- {} / {} / {} / completed\n",
job.alias, job.child_session, job.agent
));
if let Some(obj) = &job.objective {
out.push_str(&format!(" Objective: {obj}\n"));
}
if !job.context_files.is_empty() {
let files: Vec<&str> = job
.context_files
.iter()
.take(CONTEXT_FILES_SHOWN)
.map(|f| f.path.as_str())
.collect();
out.push_str(&format!(" Context read: {}\n", files.join(", ")));
}
}
}
Some(out)
}
fn get(&self, task_id: &str) -> Option<JobRecord> {
self.jobs.read().unwrap().get(task_id).cloned()
}
async fn upsert(&self, record: JobRecord) -> Result<(), StoreError> {
self.store.upsert_job(record.clone()).await?;
self.jobs
.write()
.unwrap()
.insert(record.task_id.clone(), record.clone());
self.bus.publish(AppEvent::JobUpdated {
job: JobRecordEvent(serde_json::to_value(&record).unwrap_or(serde_json::Value::Null)),
});
Ok(())
}
/// Keeps at most `max_reusable_per_agent` completed jobs per agent (LRU by `last_used_at`);
/// older completed jobs are dropped from the board and the store.
async fn trim_reusable(&self, _now: i64) -> Result<(), StoreError> {
let to_remove: Vec<String> = {
let jobs = self.jobs.read().unwrap();
let mut by_agent: HashMap<&str, Vec<&JobRecord>> = HashMap::new();
for job in jobs.values().filter(|j| j.state == JobState::Completed) {
by_agent.entry(job.agent.as_str()).or_default().push(job);
}
let mut remove = Vec::new();
for group in by_agent.values_mut() {
if group.len() as u32 <= self.max_reusable_per_agent {
continue;
}
// Oldest last_used_at first; drop the excess from the front.
group.sort_by_key(|j| j.last_used_at);
let excess = group.len() - self.max_reusable_per_agent as usize;
for job in group.iter().take(excess) {
remove.push(job.task_id.clone());
}
}
remove
};
for task_id in to_remove {
self.store.delete_job(task_id.clone()).await?;
self.jobs.write().unwrap().remove(&task_id);
}
Ok(())
}
}
fn state_label(state: JobState) -> &'static str {
match state {
JobState::Running => "running",
JobState::Completed => "completed",
JobState::Error => "error",
JobState::Cancelled => "cancelled",
}
}
fn truncate_summary(s: &str) -> String {
if s.len() <= SUMMARY_MAX {
return s.to_string();
}
let mut end = SUMMARY_MAX;
while !s.is_char_boundary(end) {
end -= 1;
}
format!("{}", &s[..end])
}
#[cfg(test)]
mod tests {
use super::*;
async fn board(max_reusable: u32) -> (Store, JobBoard, SessionId) {
let store = Store::open_in_memory().unwrap();
let bus = EventBus::new();
let parent = SessionId::new();
let board = JobBoard::load(store.clone(), bus, &parent, max_reusable)
.await
.unwrap();
(store, board, parent)
}
fn spec(
task_id: &str,
parent: SessionId,
child: SessionId,
agent: &str,
objective: Option<&str>,
) -> LaunchSpec {
LaunchSpec {
task_id: task_id.into(),
parent_session: parent,
child_session: child,
agent: agent.into(),
description: "d".into(),
objective: objective.map(Into::into),
}
}
#[tokio::test]
async fn alias_increments_per_agent() {
let (_store, board, parent) = board(2).await;
let a1 = board
.register_launch(
spec("t1", parent.clone(), SessionId::new(), "explorer", None),
1,
)
.await
.unwrap();
let a2 = board
.register_launch(
spec("t2", parent.clone(), SessionId::new(), "explorer", None),
2,
)
.await
.unwrap();
let f1 = board
.register_launch(
spec("t3", parent.clone(), SessionId::new(), "fixer", None),
3,
)
.await
.unwrap();
assert_eq!(a1, "exp-1");
assert_eq!(a2, "exp-2");
assert_eq!(f1, "fix-1");
}
#[tokio::test]
async fn finish_makes_job_reusable_and_resolvable_by_alias() {
let (_store, board, parent) = board(2).await;
let child = SessionId::new();
let alias = board
.register_launch(
spec(
"t1",
parent.clone(),
child.clone(),
"explorer",
Some("map the auth flow"),
),
1,
)
.await
.unwrap();
// Running jobs are not reusable.
assert!(board.resolve_reusable(&parent, &alias).is_none());
board
.finish("t1", JobState::Completed, Some("done".into()), 2)
.await
.unwrap();
let resolved = board.resolve_reusable(&parent, &alias).unwrap();
assert_eq!(resolved.child_session, child);
assert_eq!(resolved.result_summary.as_deref(), Some("done"));
// Also resolvable by task id.
assert!(board.resolve_reusable(&parent, "t1").is_some());
}
#[tokio::test]
async fn context_files_dedupe_and_respect_min_lines() {
let (_store, board, parent) = board(2).await;
board
.register_launch(spec("t1", parent, SessionId::new(), "explorer", None), 1)
.await
.unwrap();
// Below threshold — ignored.
board
.report_context_file("t1", "small.rs".into(), 3, 2)
.await
.unwrap();
board
.report_context_file("t1", "a.rs".into(), 20, 2)
.await
.unwrap();
// Same file again with a larger read keeps the max.
board
.report_context_file("t1", "a.rs".into(), 50, 3)
.await
.unwrap();
let job = board.snapshot().into_iter().next().unwrap();
assert_eq!(job.context_files.len(), 1);
assert_eq!(job.context_files[0].path, "a.rs");
assert_eq!(job.context_files[0].lines, 50);
}
#[tokio::test]
async fn trim_reusable_keeps_lru_within_limit() {
let (_store, board, parent) = board(2).await;
for (i, ts) in [(1, 10), (2, 20), (3, 30)] {
board
.register_launch(
spec(
&format!("t{i}"),
parent.clone(),
SessionId::new(),
"explorer",
None,
),
ts,
)
.await
.unwrap();
board
.finish(&format!("t{i}"), JobState::Completed, None, ts)
.await
.unwrap();
}
// max_reusable = 2, so the oldest (t1, last_used 10) is dropped.
let ids: Vec<String> = board.snapshot().into_iter().map(|j| j.task_id).collect();
assert_eq!(ids.len(), 2);
assert!(!ids.contains(&"t1".to_string()));
assert!(ids.contains(&"t2".to_string()));
assert!(ids.contains(&"t3".to_string()));
}
#[tokio::test]
async fn reconcile_flips_terminal_jobs_and_moves_them_to_reusable_section() {
let (_store, board, parent) = board(2).await;
board
.register_launch(
spec("t1", parent, SessionId::new(), "explorer", Some("obj")),
1,
)
.await
.unwrap();
board
.finish("t1", JobState::Completed, Some("res".into()), 2)
.await
.unwrap();
// Before reconcile: appears under Active/Unreconciled.
let prompt = board.format_for_prompt().unwrap();
assert!(prompt.contains("Active / Unreconciled"));
board.reconcile_terminal(3).await.unwrap();
let prompt = board.format_for_prompt().unwrap();
assert!(prompt.contains("Reusable Sessions"));
assert!(prompt.contains("exp-1"));
}
#[tokio::test]
async fn board_reloads_persisted_jobs() {
let store = Store::open_in_memory().unwrap();
let bus = EventBus::new();
let parent = SessionId::new();
{
let board = JobBoard::load(store.clone(), bus.clone(), &parent, 2)
.await
.unwrap();
board
.register_launch(
spec("t1", parent.clone(), SessionId::new(), "explorer", None),
1,
)
.await
.unwrap();
}
// Fresh board over the same store sees the persisted job.
let board = JobBoard::load(store, bus, &parent, 2).await.unwrap();
assert!(!board.is_empty());
assert_eq!(board.snapshot().len(), 1);
}
#[tokio::test]
async fn empty_board_formats_to_none() {
let (_store, board, _parent) = board(2).await;
assert!(board.format_for_prompt().is_none());
}
}
+252
View File
@@ -22,6 +22,14 @@ pub struct RunConfig {
pub instructions: Vec<String>,
/// Pricing for `model`, from models.dev metadata. `None` leaves cost at 0.
pub cost: Option<crate::types::ModelCost>,
/// Whether to append the background job board to requests (primary/delegating agents).
pub inject_job_board: bool,
/// Optional user-provided reminder injected at the start of every turn (off by default).
pub reminder_turn_start: Option<String>,
/// Optional user-provided reminder injected on the turn after a file tool ran.
pub reminder_after_file_tool: Option<String>,
/// Pre-rendered "## Skills" system block advertising loadable skills. `None` = no skills.
pub skills_prompt: Option<String>,
}
/// Adds a step's usage/cost onto the persisted session and republishes it. Cost accounting is
@@ -160,6 +168,8 @@ pub async fn run_session(
) -> RunOutcome {
let mut doomloop = DoomLoopGuard::new();
let mut steps = 0u32;
// Whether the previous step ran a file tool, gating the `after_file_tool` reminder.
let mut prev_used_file_tool = false;
loop {
match should_continue(&ctx.store, &ctx.session_id).await {
@@ -198,9 +208,41 @@ pub async fn run_session(
wire_messages.extend(convert_message(message, &parts));
}
// Collect synthetic (non-persisted) blocks to append to the last user message this
// turn: the optional turn-start reminder, the job board, and — if the previous step
// ran a file tool — the optional after-file-tool reminder. docs/04-multiagent.md.
let mut synthetic: Vec<String> = Vec::new();
if let Some(reminder) = &run_config.reminder_turn_start {
synthetic.push(reminder.clone());
}
if run_config.inject_job_board {
if let Some(board) = &ctx.job_board {
if let Some(block) = board.format_for_prompt() {
synthetic.push(block);
}
}
}
if prev_used_file_tool {
if let Some(reminder) = &run_config.reminder_after_file_tool {
synthetic.push(reminder.clone());
}
}
if !synthetic.is_empty() {
if let Some(last_user) = wire_messages
.iter_mut()
.rev()
.find(|m| m.role == WireRole::User)
{
for text in synthetic {
last_user.content.push(WireContent::Text { text });
}
}
}
let system_blocks = system::assemble(
system::env_header(&ctx.cwd),
&run_config.agent_prompt,
run_config.skills_prompt.as_deref(),
&run_config.instructions,
);
let tools: Vec<ToolSchema> = ctx
@@ -260,6 +302,16 @@ pub async fn run_session(
}
Ok(outcome) => {
accumulate_session_usage(&ctx, &outcome.usage, outcome.cost, now_fn()).await;
prev_used_file_tool = outcome.used_file_tool;
// A completed step means the orchestrator has now seen any terminal jobs
// that were on the board this turn; mark them reconciled.
if run_config.inject_job_board {
if let Some(board) = &ctx.job_board {
if let Err(e) = board.reconcile_terminal(now_fn()).await {
tracing::warn!(error = %e, "failed to reconcile job board");
}
}
}
match outcome.result {
StepResult::Continue => continue,
StepResult::Stop => return RunOutcome::Stopped,
@@ -376,11 +428,16 @@ mod tests {
permissions,
static_rules: Vec::new(),
extra_rules: Arc::new(std::sync::Mutex::new(Vec::new())),
parent_rules: Vec::new(),
session_id,
cwd: cwd.clone(),
data_dir: cwd.join("tool-output"),
cancel: CancellationToken::new(),
now: 1,
spawner: None,
job_board: None,
context_reporter: None,
diagnostics: None,
}
}
@@ -462,6 +519,10 @@ mod tests {
output: 15.0,
..Default::default()
}),
inject_job_board: false,
reminder_turn_start: None,
reminder_after_file_tool: None,
skills_prompt: None,
};
let outcome = run_session(Arc::new(provider), ctx, &run_config, || 2).await;
@@ -569,6 +630,10 @@ mod tests {
max_steps: 10,
instructions: Vec::new(),
cost: None,
inject_job_board: false,
reminder_turn_start: None,
reminder_after_file_tool: None,
skills_prompt: None,
};
let outcome = run_session(Arc::new(provider), ctx, &run_config, || 2).await;
@@ -578,4 +643,191 @@ mod tests {
let messages = store.messages(session_id).await.unwrap();
assert_eq!(messages.len(), 1);
}
/// Records the last request it was asked to stream so tests can assert on prompt content.
struct CapturingProvider {
last: StdMutex<Option<LlmRequest>>,
}
#[async_trait]
impl Provider for CapturingProvider {
fn id(&self) -> &str {
"mock"
}
async fn list_models(&self) -> Result<Vec<ModelInfo>, ProviderError> {
Ok(vec![])
}
async fn stream(
&self,
req: LlmRequest,
_cancel: CancellationToken,
) -> Result<LlmEventStream, ProviderError> {
*self.last.lock().unwrap() = Some(req);
let events = vec![
Ok(LlmEvent::TextStart { id: "t".into() }),
Ok(LlmEvent::TextDelta {
id: "t".into(),
text: "ok".into(),
}),
Ok(LlmEvent::TextEnd { id: "t".into() }),
Ok(LlmEvent::Finish {
reason: FinishReason::Stop,
usage: usage(1, 1),
}),
];
Ok(Box::pin(futures::stream::iter(events)))
}
}
#[tokio::test]
async fn job_board_is_injected_into_the_last_user_message() {
use crate::engine::jobs::{JobBoard, LaunchSpec};
let store = Store::open_in_memory().unwrap();
let bus = EventBus::new();
let model = ModelRef::new("mock", "mock-model");
let session = Session::new_root("orchestrator", model.clone(), 1);
let session_id = session.id.clone();
store.upsert_session(session).await.unwrap();
let user_message = Message::new_user(session_id.clone(), 1);
store.upsert_message(user_message.clone()).await.unwrap();
store
.upsert_part(Part {
id: crate::types::PartId::new(),
message_id: user_message.id.clone(),
session_id: session_id.clone(),
idx: 0,
body: PartBody::Text {
text: "carry on".into(),
synthetic: false,
},
})
.await
.unwrap();
// A board with one running job for this session.
let board = std::sync::Arc::new(
JobBoard::load(store.clone(), bus.clone(), &session_id, 2)
.await
.unwrap(),
);
board
.register_launch(
LaunchSpec {
task_id: "t1".into(),
parent_session: session_id.clone(),
child_session: SessionId::new(),
agent: "explorer".into(),
description: "map auth".into(),
objective: Some("map the auth flow".into()),
},
1,
)
.await
.unwrap();
let cwd = tempfile::tempdir().unwrap();
let mut ctx = make_ctx(store, bus, session_id, cwd.path().to_path_buf()).await;
ctx.job_board = Some(board);
let run_config = RunConfig {
agent_name: "orchestrator".into(),
agent_prompt: "You orchestrate.".into(),
model,
temperature: None,
max_steps: 1,
instructions: Vec::new(),
cost: None,
inject_job_board: true,
reminder_turn_start: None,
reminder_after_file_tool: None,
skills_prompt: None,
};
let provider = std::sync::Arc::new(CapturingProvider {
last: StdMutex::new(None),
});
let outcome = run_session(provider.clone(), ctx, &run_config, || 2).await;
assert!(matches!(outcome, RunOutcome::Stopped));
let req = provider.last.lock().unwrap().clone().expect("a request");
let last_user = req
.messages
.iter()
.rev()
.find(|m| m.role == WireRole::User)
.expect("a user message");
let text: String = last_user
.content
.iter()
.filter_map(|c| match c {
WireContent::Text { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>()
.join("\n");
assert!(text.contains("Background Job Board"), "got: {text}");
assert!(text.contains("exp-1"), "got: {text}");
assert!(text.contains("map the auth flow"), "got: {text}");
}
#[tokio::test]
async fn turn_start_reminder_is_injected_into_the_request() {
let store = Store::open_in_memory().unwrap();
let bus = EventBus::new();
let model = ModelRef::new("mock", "mock-model");
let session = Session::new_root("orchestrator", model.clone(), 1);
let session_id = session.id.clone();
store.upsert_session(session).await.unwrap();
let user_message = Message::new_user(session_id.clone(), 1);
store.upsert_message(user_message.clone()).await.unwrap();
store
.upsert_part(Part {
id: crate::types::PartId::new(),
message_id: user_message.id.clone(),
session_id: session_id.clone(),
idx: 0,
body: PartBody::Text {
text: "do the thing".into(),
synthetic: false,
},
})
.await
.unwrap();
let cwd = tempfile::tempdir().unwrap();
let ctx = make_ctx(store, bus, session_id, cwd.path().to_path_buf()).await;
let run_config = RunConfig {
agent_name: "orchestrator".into(),
agent_prompt: "You orchestrate.".into(),
model,
temperature: None,
max_steps: 1,
instructions: Vec::new(),
cost: None,
inject_job_board: false,
reminder_turn_start: Some("REMEMBER: stay on task.".into()),
reminder_after_file_tool: None,
skills_prompt: None,
};
let provider = std::sync::Arc::new(CapturingProvider {
last: StdMutex::new(None),
});
let outcome = run_session(provider.clone(), ctx, &run_config, || 2).await;
assert!(matches!(outcome, RunOutcome::Stopped));
let req = provider.last.lock().unwrap().clone().expect("a request");
let last_user = req
.messages
.iter()
.rev()
.find(|m| m.role == WireRole::User)
.expect("a user message");
let has_reminder = last_user.content.iter().any(
|c| matches!(c, WireContent::Text { text } if text.contains("REMEMBER: stay on task.")),
);
assert!(has_reminder, "turn-start reminder should be injected");
}
}
+2
View File
@@ -1,4 +1,5 @@
pub mod doomloop;
pub mod jobs;
pub mod processor;
pub mod retry;
#[path = "loop.rs"]
@@ -6,5 +7,6 @@ pub mod session_loop;
pub mod system;
pub use doomloop::DoomLoopGuard;
pub use jobs::{ContextFile, JobBoard, JobRecord, JobState};
pub use processor::{process_step, StepContext, StepError, StepOutcome, StepResult};
pub use session_loop::{run_session, RunConfig};
+36 -2
View File
@@ -10,10 +10,14 @@ use crate::event::{AppEvent, EventBus};
use crate::llm::{FinishReason, LlmEvent, LlmEventStream, ProviderError};
use crate::permission::{PermissionService, Ruleset};
use crate::store::Store;
use crate::tool::{MetadataSink, PermissionHandle, Tool, ToolCtx, ToolError, ToolRegistry};
use crate::tool::{
ContextReporter, MetadataSink, PermissionHandle, SubagentSpawner, Tool, ToolCtx, ToolError,
ToolRegistry,
};
use crate::types::{Message, MessageId, Part, PartBody, PartId, SessionId, TokenUsage, ToolState};
use super::doomloop::DoomLoopGuard;
use super::jobs::JobBoard;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StepResult {
@@ -29,6 +33,9 @@ pub struct StepOutcome {
/// Dollar cost of this step's usage (0.0 when no pricing is available).
pub cost: f64,
pub aborted: bool,
/// Whether a file-mutating tool (`edit`/`write`) ran this step — drives the optional
/// `after_file_tool` reminder injection on the next turn.
pub used_file_tool: bool,
}
pub struct StepError {
@@ -47,6 +54,9 @@ pub struct StepContext {
pub permissions: Arc<PermissionService>,
pub static_rules: Ruleset,
pub extra_rules: Arc<Mutex<Ruleset>>,
/// Parent-effective ruleset for a subagent session; empty for a root session. Enables
/// permission intersection on this session's tool calls.
pub parent_rules: Ruleset,
pub session_id: SessionId,
pub cwd: PathBuf,
/// Session's `tool-output` spill directory (see `tool::truncate`).
@@ -54,6 +64,16 @@ pub struct StepContext {
pub cancel: CancellationToken,
/// Wall-clock for `created_at` stamps — passed in so tests stay deterministic.
pub now: i64,
/// Lets the `task` tool spawn subagents. `None` disables delegation (headless/tests).
pub spawner: Option<Arc<dyn SubagentSpawner>>,
/// This session's background job board (as a parent). Injected into requests when the
/// running agent can delegate; `None` disables the board.
pub job_board: Option<Arc<JobBoard>>,
/// Present in subagent sessions: reports read files to this session's job on the parent
/// board. `None` for root sessions (nothing to report to).
pub context_reporter: Option<Arc<dyn ContextReporter>>,
/// Language-server diagnostics source shared by the session's edit/write tool calls.
pub diagnostics: Option<Arc<dyn crate::lsp::DiagnosticsSource>>,
}
struct FlushTracker {
@@ -83,9 +103,13 @@ impl FlushTracker {
}
}
/// Tool names that mutate files — after one runs, the optional `after_file_tool` reminder fires.
const FILE_TOOLS: &[&str] = &["edit", "write"];
struct Run<'a> {
ctx: &'a StepContext,
assistant: Option<Message>,
used_file_tool: bool,
next_idx: u32,
active_text: Option<PartId>,
active_reasoning: Option<PartId>,
@@ -102,6 +126,7 @@ impl<'a> Run<'a> {
Self {
ctx,
assistant: None,
used_file_tool: false,
next_idx: 0,
active_text: None,
active_reasoning: None,
@@ -371,6 +396,9 @@ impl<'a> Run<'a> {
input: serde_json::Value,
doomloop: &mut DoomLoopGuard,
) -> Result<(), ProviderError> {
if FILE_TOOLS.contains(&name.as_str()) {
self.used_file_tool = true;
}
let part_id = self.pending_tools.remove(&call_id).unwrap_or_default();
let running = Part {
id: part_id.clone(),
@@ -483,7 +511,8 @@ impl<'a> Run<'a> {
self.ctx.static_rules.clone(),
self.ctx.extra_rules.clone(),
call_cancel.clone(),
);
)
.with_parent_rules(self.ctx.parent_rules.clone());
let tool_ctx = ToolCtx {
session_id: self.ctx.session_id.clone(),
message_id: self.message_id(),
@@ -493,6 +522,9 @@ impl<'a> Run<'a> {
cancel: call_cancel.clone(),
ask,
metadata: metadata_sink,
spawner: self.ctx.spawner.clone(),
context_reporter: self.ctx.context_reporter.clone(),
diagnostics: self.ctx.diagnostics.clone(),
};
let result = tokio::select! {
@@ -589,6 +621,7 @@ pub async fn process_step(
usage,
cost: step_cost,
aborted: true,
used_file_tool: run.used_file_tool,
});
}
};
@@ -661,5 +694,6 @@ pub async fn process_step(
usage,
cost: step_cost,
aborted: false,
used_file_tool: run.used_file_tool,
})
}
+23 -3
View File
@@ -25,10 +25,18 @@ pub fn env_header(cwd: &Path) -> String {
header
}
/// Ordered system prompt blocks: environment header, agent prompt, then project
/// instructions (e.g. AGENTS.md contents). Order matches `02-engine.md`.
pub fn assemble(env_header: String, agent_prompt: &str, instructions: &[String]) -> Vec<String> {
/// Ordered system prompt blocks: environment header, agent prompt, the optional skills
/// listing, then project instructions (e.g. AGENTS.md contents). Order matches `02-engine.md`.
pub fn assemble(
env_header: String,
agent_prompt: &str,
skills: Option<&str>,
instructions: &[String],
) -> Vec<String> {
let mut blocks = vec![env_header, agent_prompt.to_string()];
if let Some(skills) = skills {
blocks.push(skills.to_string());
}
blocks.extend(instructions.iter().cloned());
blocks
}
@@ -49,8 +57,20 @@ mod tests {
let blocks = assemble(
"ENV".to_string(),
"AGENT",
None,
&["AGENTS.md contents".to_string()],
);
assert_eq!(blocks, vec!["ENV", "AGENT", "AGENTS.md contents"]);
}
#[test]
fn assemble_inserts_skills_after_agent_before_instructions() {
let blocks = assemble(
"ENV".to_string(),
"AGENT",
Some("SKILLS"),
&["INSTR".to_string()],
);
assert_eq!(blocks, vec!["ENV", "AGENT", "SKILLS", "INSTR"]);
}
}
+2
View File
@@ -1,7 +1,9 @@
pub mod agent;
pub mod config;
pub mod engine;
pub mod event;
pub mod llm;
pub mod lsp;
pub mod permission;
pub mod store;
pub mod tool;
+46
View File
@@ -0,0 +1,46 @@
//! Seam trait for language-server diagnostics, implemented by `harness-lsp` and consumed by
//! the edit/write tools. Kept in core so `harness-tools` never links the LSP crate directly.
//! See `docs/09-integrations.md`.
use std::path::Path;
use std::time::Duration;
use async_trait::async_trait;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Severity {
Error,
Warning,
Info,
Hint,
}
/// One diagnostic reported by a language server. Line/character are 1-based for display.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Diagnostic {
pub line: u32,
pub character: u32,
pub severity: Severity,
pub message: String,
pub source: Option<String>,
}
impl Diagnostic {
/// `{file}: L{line}: {message}` — the one-line form appended to tool output.
pub fn display_line(&self, file: &str) -> String {
format!("{file}: L{}: {}", self.line, self.message)
}
}
/// A source of file diagnostics (an LSP pool). All methods are best-effort: failures are
/// swallowed (logged) so a broken language server never fails a tool call.
#[async_trait]
pub trait DiagnosticsSource: Send + Sync {
/// Ensure a server for `path`'s language is running and told about the file's current
/// contents (spawn-if-needed + didOpen/didChange).
async fn touch(&self, path: &Path);
/// Diagnostics for `path`, waiting up to `wait` for the server to (re)publish after a
/// change. Returns whatever is known on timeout.
async fn diagnostics(&self, path: &Path, wait: Duration) -> Vec<Diagnostic>;
}
+1 -1
View File
@@ -1,7 +1,7 @@
pub mod rule;
pub mod service;
pub use rule::{evaluate, Action, Rule, Ruleset};
pub use rule::{evaluate, evaluate_intersected, Action, Rule, Ruleset};
pub use service::{
spawn_auto_approve, AskDecision, AskError, AskInput, PermissionReply, PermissionService,
};
@@ -44,6 +44,37 @@ pub fn evaluate(stack: &[&Ruleset], permission: &str, pattern: &str) -> Action {
result
}
impl Action {
/// How restrictive this verdict is: `Deny` > `Ask` > `Allow`. Used to intersect a
/// parent and child verdict when a subagent runs (docs/04-multiagent.md).
fn restrictiveness(self) -> u8 {
match self {
Action::Allow => 0,
Action::Ask => 1,
Action::Deny => 2,
}
}
}
/// Evaluate `permission`/`pattern` against a parent-effective stack and a child stack
/// independently, returning the **more restrictive** of the two verdicts
/// (`deny > ask > allow`). A subagent's tool call must satisfy both the rules it inherits
/// from its spawning chain and its own agent ruleset.
pub fn evaluate_intersected(
parent_stack: &[&Ruleset],
child_stack: &[&Ruleset],
permission: &str,
pattern: &str,
) -> Action {
let parent = evaluate(parent_stack, permission, pattern);
let child = evaluate(child_stack, permission, pattern);
if child.restrictiveness() >= parent.restrictiveness() {
child
} else {
parent
}
}
#[cfg(test)]
mod tests {
use super::*;
@@ -112,4 +143,49 @@ mod tests {
fn empty_stack_defaults_to_ask() {
assert_eq!(evaluate(&[], "bash", "ls"), Action::Ask);
}
#[test]
fn intersected_takes_the_more_restrictive_verdict() {
let parent_allow: Ruleset = vec![rule("edit", "*", Action::Allow)];
let child_deny: Ruleset = vec![rule("edit", "*", Action::Deny)];
// Child denies what the parent would allow → deny wins.
assert_eq!(
evaluate_intersected(&[&parent_allow], &[&child_deny], "edit", "main.rs"),
Action::Deny
);
// Symmetric: parent denies what the child would allow → deny still wins.
assert_eq!(
evaluate_intersected(&[&child_deny], &[&parent_allow], "edit", "main.rs"),
Action::Deny
);
}
#[test]
fn intersected_ask_beats_allow_but_loses_to_deny() {
let allow: Ruleset = vec![rule("bash", "*", Action::Allow)];
let ask: Ruleset = vec![rule("bash", "*", Action::Ask)];
let deny: Ruleset = vec![rule("bash", "*", Action::Deny)];
assert_eq!(
evaluate_intersected(&[&allow], &[&ask], "bash", "ls"),
Action::Ask
);
assert_eq!(
evaluate_intersected(&[&ask], &[&deny], "bash", "ls"),
Action::Deny
);
}
#[test]
fn intersected_allows_only_when_both_allow() {
let allow: Ruleset = vec![rule("read", "*", Action::Allow)];
assert_eq!(
evaluate_intersected(&[&allow], &[&allow], "read", "src/a.rs"),
Action::Allow
);
// Empty child stack defaults to Ask, which is more restrictive than parent Allow.
assert_eq!(
evaluate_intersected(&[&allow], &[], "read", "src/a.rs"),
Action::Ask
);
}
}
+19 -1
View File
@@ -8,7 +8,7 @@ use ulid::Ulid;
use crate::event::{AppEvent, EventBus, PermissionRequest};
use crate::types::SessionId;
use super::rule::{evaluate, Action, Rule, Ruleset};
use super::rule::{evaluate, evaluate_intersected, Action, Rule, Ruleset};
pub struct AskInput {
pub permission: String,
@@ -73,6 +73,24 @@ impl PermissionService {
}
}
/// Like [`ask`](Self::ask), but for a subagent: the verdict is the more restrictive of
/// the `parent_stack` (rules inherited from the spawning chain) and `child_stack` (the
/// subagent's own rules). Used so a child can never widen what its parent forbids.
pub async fn ask_intersected(
&self,
session_id: &SessionId,
parent_stack: &[&Ruleset],
child_stack: &[&Ruleset],
input: AskInput,
cancel: &CancellationToken,
) -> Result<AskDecision, AskError> {
match evaluate_intersected(parent_stack, child_stack, &input.permission, &input.pattern) {
Action::Allow => Ok(AskDecision::Allowed),
Action::Deny => Err(AskError::Denied),
Action::Ask => self.ask_user(session_id, input, cancel).await,
}
}
/// Bypasses ruleset evaluation entirely — used by the doom-loop guard, which must ask
/// regardless of any `Allow` rule.
pub async fn force_ask(
+41
View File
@@ -1,6 +1,7 @@
use rusqlite::{params, Connection};
use tokio::sync::oneshot;
use crate::engine::jobs::JobRecord;
use crate::types::{Message, MessageId, Part, Session, SessionId};
use super::api::StoreError;
@@ -17,6 +18,9 @@ pub enum StoreCmd {
Sessions(Reply<Vec<Session>>),
Messages(SessionId, Reply<Vec<Message>>),
Parts(MessageId, Reply<Vec<Part>>),
UpsertJob(JobRecord, Reply<()>),
DeleteJob(String, Reply<()>),
JobsForParent(SessionId, Reply<Vec<JobRecord>>),
}
fn init_schema(conn: &Connection) -> rusqlite::Result<()> {
@@ -109,6 +113,34 @@ fn list_parts(conn: &Connection, message_id: &MessageId) -> Result<Vec<Part>, St
.collect()
}
fn upsert_job(conn: &Connection, job: &JobRecord) -> Result<(), StoreError> {
let data = serde_json::to_string(job)?;
conn.execute(
"INSERT INTO job (task_id, parent_session_id, data) VALUES (?1, ?2, ?3)
ON CONFLICT(task_id) DO UPDATE SET data = ?3",
params![job.task_id, job.parent_session.as_ref(), data],
)?;
Ok(())
}
fn delete_job(conn: &Connection, task_id: &str) -> Result<(), StoreError> {
conn.execute("DELETE FROM job WHERE task_id = ?1", params![task_id])?;
Ok(())
}
fn list_jobs_for_parent(
conn: &Connection,
parent: &SessionId,
) -> Result<Vec<JobRecord>, StoreError> {
let mut stmt = conn.prepare("SELECT data FROM job WHERE parent_session_id = ?1")?;
let rows = stmt
.query_map(params![parent.as_ref()], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
rows.iter()
.map(|data| serde_json::from_str(data).map_err(StoreError::from))
.collect()
}
/// Runs on a dedicated OS thread; the async facade in `api.rs` talks to it over `mpsc`.
pub fn run(conn: Connection, mut rx: tokio::sync::mpsc::Receiver<StoreCmd>) {
if let Err(e) = init_schema(&conn) {
@@ -138,6 +170,15 @@ pub fn run(conn: Connection, mut rx: tokio::sync::mpsc::Receiver<StoreCmd>) {
StoreCmd::Parts(message_id, reply) => {
let _ = reply.send(list_parts(&conn, &message_id));
}
StoreCmd::UpsertJob(job, reply) => {
let _ = reply.send(upsert_job(&conn, &job));
}
StoreCmd::DeleteJob(task_id, reply) => {
let _ = reply.send(delete_job(&conn, &task_id));
}
StoreCmd::JobsForParent(parent, reply) => {
let _ = reply.send(list_jobs_for_parent(&conn, &parent));
}
}
}
}
+14
View File
@@ -3,6 +3,7 @@ use std::path::Path;
use rusqlite::Connection;
use tokio::sync::{mpsc, oneshot};
use crate::engine::jobs::JobRecord;
use crate::types::{Message, MessageId, Part, Session, SessionId};
use super::actor::{self, StoreCmd};
@@ -94,6 +95,19 @@ impl Store {
pub async fn parts(&self, message_id: MessageId) -> Result<Vec<Part>, StoreError> {
self.call(|reply| StoreCmd::Parts(message_id, reply)).await
}
pub async fn upsert_job(&self, job: JobRecord) -> Result<(), StoreError> {
self.call(|reply| StoreCmd::UpsertJob(job, reply)).await
}
pub async fn delete_job(&self, task_id: String) -> Result<(), StoreError> {
self.call(|reply| StoreCmd::DeleteJob(task_id, reply)).await
}
pub async fn jobs_for_parent(&self, parent: SessionId) -> Result<Vec<JobRecord>, StoreError> {
self.call(|reply| StoreCmd::JobsForParent(parent, reply))
.await
}
}
#[cfg(test)]
+93 -15
View File
@@ -10,6 +10,58 @@ use tokio_util::sync::CancellationToken;
use crate::permission::{AskDecision, AskError, AskInput, PermissionService, Ruleset};
use crate::types::{MessageId, SessionId};
/// A request from the `task` tool to run a subagent. The spawner (owned by the composition
/// root) resolves the agent, enforces the depth limit, applies permission intersection, and
/// runs the child session foreground or background. See `docs/04-multiagent.md`.
pub struct SpawnRequest {
pub parent_session_id: SessionId,
pub parent_message_id: MessageId,
pub agent: String,
pub description: String,
pub prompt: String,
/// Alias or task id of a completed job to reuse (continue its child session).
pub reuse_task_id: Option<String>,
pub background: bool,
/// The tool call's cancellation token — used for foreground child runs. Background runs
/// are childed from the parent session's run token by the spawner instead.
pub cancel: CancellationToken,
}
#[derive(Debug)]
pub struct SpawnOutcome {
pub child_session_id: SessionId,
pub background: bool,
/// Board alias assigned to a background launch.
pub alias: Option<String>,
/// Final assistant text of a foreground run.
pub final_text: Option<String>,
}
#[derive(Debug, thiserror::Error)]
pub enum SpawnError {
#[error("unknown subagent {0:?}, or it is not usable as a subagent")]
InvalidAgent(String),
#[error("subagent depth limit reached — do this work yourself instead of delegating further")]
DepthExceeded,
#[error("cannot reuse {0:?}: no completed job with that alias for this session")]
ReuseNotFound(String),
#[error("{0}")]
Other(String),
}
#[async_trait]
pub trait SubagentSpawner: Send + Sync {
async fn spawn(&self, req: SpawnRequest) -> Result<SpawnOutcome, SpawnError>;
}
/// Lets a child session report the files it read to its job board entry, so a completed
/// specialist advertises what it already looked at (docs/04-multiagent.md). Present only in
/// subagent sessions; the spawner wires it to the right board + job.
#[async_trait]
pub trait ContextReporter: Send + Sync {
async fn report_file(&self, path: String, lines: u32);
}
#[derive(Debug, thiserror::Error)]
pub enum ToolError {
#[error("permission denied")]
@@ -60,6 +112,9 @@ pub struct PermissionHandle {
session_id: SessionId,
static_rules: Ruleset,
extra_rules: Arc<Mutex<Ruleset>>,
/// Parent-effective ruleset for a subagent session; empty for a root session. When
/// non-empty, verdicts are intersected so a child can only ever be *more* restricted.
parent_rules: Ruleset,
cancel: CancellationToken,
}
@@ -76,10 +131,17 @@ impl PermissionHandle {
session_id,
static_rules,
extra_rules,
parent_rules: Vec::new(),
cancel,
}
}
/// Sets the parent-effective ruleset so this handle intersects verdicts (subagent runs).
pub fn with_parent_rules(mut self, parent_rules: Ruleset) -> Self {
self.parent_rules = parent_rules;
self
}
pub async fn ask(
&self,
permission: impl Into<String>,
@@ -88,21 +150,29 @@ impl PermissionHandle {
metadata: serde_json::Value,
) -> Result<(), ToolError> {
let extra_snapshot = self.extra_rules.lock().unwrap().clone();
let stack: [&Ruleset; 2] = [&self.static_rules, &extra_snapshot];
let decision = self
.service
.ask(
&self.session_id,
&stack,
AskInput {
permission: permission.into(),
pattern: pattern.into(),
always_pattern: always_pattern.into(),
metadata,
},
&self.cancel,
)
.await?;
let child_stack: [&Ruleset; 2] = [&self.static_rules, &extra_snapshot];
let input = AskInput {
permission: permission.into(),
pattern: pattern.into(),
always_pattern: always_pattern.into(),
metadata,
};
let decision = if self.parent_rules.is_empty() {
self.service
.ask(&self.session_id, &child_stack, input, &self.cancel)
.await?
} else {
let parent_stack: [&Ruleset; 1] = [&self.parent_rules];
self.service
.ask_intersected(
&self.session_id,
&parent_stack,
&child_stack,
input,
&self.cancel,
)
.await?
};
if let AskDecision::AllowedAlways(rule) = decision {
self.extra_rules.lock().unwrap().push(rule);
}
@@ -120,6 +190,14 @@ pub struct ToolCtx {
pub cancel: CancellationToken,
pub ask: PermissionHandle,
pub metadata: MetadataSink,
/// Present when the engine can spawn subagents (the `task` tool's capability). `None`
/// in headless/test contexts with no orchestration wired in.
pub spawner: Option<Arc<dyn SubagentSpawner>>,
/// Present in subagent sessions: lets the read tool report files to the job board.
pub context_reporter: Option<Arc<dyn ContextReporter>>,
/// Language-server diagnostics source (edit/write surface errors after a change). `None`
/// disables LSP integration.
pub diagnostics: Option<Arc<dyn crate::lsp::DiagnosticsSource>>,
}
#[derive(Debug)]
+7
View File
@@ -6,6 +6,13 @@ license.workspace = true
[dependencies]
harness-core = { workspace = true }
tokio = { workspace = true }
async-trait = { workspace = true }
serde_json = { workspace = true }
tracing = { workspace = true }
[dev-dependencies]
tempfile = { workspace = true }
[lints]
workspace = true
+325
View File
@@ -0,0 +1,325 @@
//! Minimal JSON-RPC client over an LSP server's child stdio: `Content-Length` framing, a
//! request-id → oneshot map, and a diagnostics store keyed by document URI. Only the handful
//! of methods the edit/write flow needs are implemented (docs/09-integrations.md).
use std::collections::HashMap;
use std::path::Path;
use std::sync::atomic::{AtomicI64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use harness_core::lsp::{Diagnostic, Severity};
use serde_json::{json, Value};
use tokio::io::{AsyncReadExt, AsyncWriteExt, BufReader};
use tokio::process::{Child, Command};
use tokio::sync::{oneshot, Notify};
type Pending = Arc<Mutex<HashMap<i64, oneshot::Sender<Value>>>>;
type DiagStore = Arc<Mutex<HashMap<String, Vec<Diagnostic>>>>;
/// A running language server plus the state needed to talk to it.
pub struct LspClient {
outgoing: tokio::sync::mpsc::UnboundedSender<String>,
next_id: AtomicI64,
pending: Pending,
diagnostics: DiagStore,
diag_notify: Arc<Notify>,
// Kept alive so the child is killed on drop (`kill_on_drop`).
_child: Child,
}
impl LspClient {
/// Spawns `command args`, performs the `initialize`/`initialized` handshake rooted at
/// `root`, and returns a ready client.
pub async fn spawn(command: &str, args: &[String], root: &Path) -> std::io::Result<Self> {
let mut child = Command::new(command)
.args(args)
.current_dir(root)
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::null())
.kill_on_drop(true)
.spawn()?;
let stdin = child.stdin.take().expect("piped stdin");
let stdout = child.stdout.take().expect("piped stdout");
let (outgoing, mut out_rx) = tokio::sync::mpsc::unbounded_channel::<String>();
let pending: Pending = Arc::default();
let diagnostics: DiagStore = Arc::default();
let diag_notify = Arc::new(Notify::new());
// Writer task: frame and forward outgoing payloads.
tokio::spawn(async move {
let mut stdin = stdin;
while let Some(payload) = out_rx.recv().await {
let frame = format!("Content-Length: {}\r\n\r\n{}", payload.len(), payload);
if stdin.write_all(frame.as_bytes()).await.is_err() {
break;
}
let _ = stdin.flush().await;
}
});
// Reader task: parse frames, route responses, collect diagnostics, ack server requests.
{
let pending = pending.clone();
let diagnostics = diagnostics.clone();
let diag_notify = diag_notify.clone();
let outgoing_ack = outgoing.clone();
tokio::spawn(async move {
let mut reader = BufReader::new(stdout);
while let Some(msg) = read_message(&mut reader).await {
dispatch(msg, &pending, &diagnostics, &diag_notify, &outgoing_ack);
}
});
}
let client = Self {
outgoing,
next_id: AtomicI64::new(1),
pending,
diagnostics,
diag_notify,
_child: child,
};
client.initialize(root).await?;
Ok(client)
}
async fn initialize(&self, root: &Path) -> std::io::Result<()> {
let root_uri = path_to_uri(root);
let params = json!({
"processId": std::process::id(),
"rootUri": root_uri,
"workspaceFolders": [{ "uri": root_uri, "name": "root" }],
"capabilities": {
"textDocument": {
"publishDiagnostics": { "relatedInformation": false }
}
},
"clientInfo": { "name": "ai-harness" }
});
// A slow server (rust-analyzer indexing) can take a while to answer initialize.
let _ = self
.request("initialize", params, Duration::from_secs(30))
.await;
self.notify("initialized", json!({}));
Ok(())
}
fn send(&self, payload: Value) {
if let Ok(text) = serde_json::to_string(&payload) {
let _ = self.outgoing.send(text);
}
}
pub fn notify(&self, method: &str, params: Value) {
self.send(json!({ "jsonrpc": "2.0", "method": method, "params": params }));
}
async fn request(&self, method: &str, params: Value, timeout: Duration) -> Option<Value> {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
let (tx, rx) = oneshot::channel();
self.pending.lock().unwrap().insert(id, tx);
self.send(json!({ "jsonrpc": "2.0", "id": id, "method": method, "params": params }));
match tokio::time::timeout(timeout, rx).await {
Ok(Ok(value)) => Some(value),
_ => {
self.pending.lock().unwrap().remove(&id);
None
}
}
}
pub fn did_open(&self, path: &Path, language_id: &str, version: i32, text: &str) {
self.notify(
"textDocument/didOpen",
json!({
"textDocument": {
"uri": path_to_uri(path),
"languageId": language_id,
"version": version,
"text": text,
}
}),
);
}
pub fn did_change(&self, path: &Path, version: i32, text: &str) {
self.notify(
"textDocument/didChange",
json!({
"textDocument": { "uri": path_to_uri(path), "version": version },
"contentChanges": [{ "text": text }],
}),
);
}
/// Clears the stored diagnostics for `path` so the next `wait_diagnostics` observes a fresh
/// publish rather than a stale one.
pub fn clear(&self, path: &Path) {
self.diagnostics.lock().unwrap().remove(&path_to_uri(path));
}
/// Waits up to `wait` for the server to publish diagnostics for `path`, returning whatever
/// is stored on timeout (possibly empty).
pub async fn wait_diagnostics(&self, path: &Path, wait: Duration) -> Vec<Diagnostic> {
let uri = path_to_uri(path);
let deadline = tokio::time::Instant::now() + wait;
loop {
if let Some(diags) = self.diagnostics.lock().unwrap().get(&uri) {
return diags.clone();
}
let notified = self.diag_notify.notified();
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
break;
}
if tokio::time::timeout(remaining, notified).await.is_err() {
break;
}
}
self.diagnostics
.lock()
.unwrap()
.get(&uri)
.cloned()
.unwrap_or_default()
}
}
fn dispatch(
msg: Value,
pending: &Pending,
diagnostics: &DiagStore,
diag_notify: &Arc<Notify>,
outgoing: &tokio::sync::mpsc::UnboundedSender<String>,
) {
let method = msg.get("method").and_then(|m| m.as_str());
let id = msg.get("id");
match (method, id) {
// Server → client request: ack with a null result so the server can proceed
// (e.g. client/registerCapability, window/workDoneProgress/create).
(Some(_), Some(id)) => {
let reply = json!({ "jsonrpc": "2.0", "id": id, "result": Value::Null });
if let Ok(text) = serde_json::to_string(&reply) {
let _ = outgoing.send(text);
}
}
// Notification from the server.
(Some("textDocument/publishDiagnostics"), None) => {
if let Some(params) = msg.get("params") {
if let Some((uri, diags)) = parse_diagnostics(params) {
diagnostics.lock().unwrap().insert(uri, diags);
diag_notify.notify_waiters();
}
}
}
(Some(_), None) => {}
// Response to one of our requests.
(None, Some(id)) => {
if let Some(id) = id.as_i64() {
if let Some(tx) = pending.lock().unwrap().remove(&id) {
let result = msg.get("result").cloned().unwrap_or(Value::Null);
let _ = tx.send(result);
}
}
}
(None, None) => {}
}
}
fn parse_diagnostics(params: &Value) -> Option<(String, Vec<Diagnostic>)> {
let uri = params.get("uri")?.as_str()?.to_string();
let items = params.get("diagnostics")?.as_array()?;
let diags = items
.iter()
.filter_map(|d| {
let start = d.get("range")?.get("start")?;
Some(Diagnostic {
line: start.get("line")?.as_u64().unwrap_or(0) as u32 + 1,
character: start.get("character")?.as_u64().unwrap_or(0) as u32 + 1,
severity: match d.get("severity").and_then(|s| s.as_u64()) {
Some(1) => Severity::Error,
Some(2) => Severity::Warning,
Some(3) => Severity::Info,
_ => Severity::Hint,
},
message: d.get("message")?.as_str().unwrap_or("").to_string(),
source: d
.get("source")
.and_then(|s| s.as_str())
.map(|s| s.to_string()),
})
})
.collect();
Some((uri, diags))
}
/// Reads one `Content-Length`-framed JSON-RPC message, or `None` at EOF.
async fn read_message<R: tokio::io::AsyncBufRead + Unpin>(reader: &mut R) -> Option<Value> {
use tokio::io::AsyncBufReadExt;
let mut content_length: Option<usize> = None;
loop {
let mut line = String::new();
let n = reader.read_line(&mut line).await.ok()?;
if n == 0 {
return None; // EOF
}
let trimmed = line.trim_end();
if trimmed.is_empty() {
break; // end of headers
}
if let Some(value) = trimmed.strip_prefix("Content-Length:") {
content_length = value.trim().parse().ok();
}
}
let len = content_length?;
let mut buf = vec![0u8; len];
reader.read_exact(&mut buf).await.ok()?;
serde_json::from_slice(&buf).ok()
}
/// `file://` URI for an absolute path (best-effort; assumes UTF-8, no percent-encoding needed
/// for the local paths we handle).
fn path_to_uri(path: &Path) -> String {
let s = path.to_string_lossy();
if s.starts_with('/') {
format!("file://{s}")
} else {
format!("file:///{s}")
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_publish_diagnostics_into_one_based_positions() {
let params = json!({
"uri": "file:///tmp/a.rs",
"diagnostics": [{
"range": {"start": {"line": 4, "character": 8}, "end": {"line": 4, "character": 12}},
"severity": 1,
"message": "cannot find value `x`",
"source": "rustc"
}]
});
let (uri, diags) = parse_diagnostics(&params).unwrap();
assert_eq!(uri, "file:///tmp/a.rs");
assert_eq!(diags.len(), 1);
assert_eq!(diags[0].line, 5); // 0-based 4 → 1-based 5
assert_eq!(diags[0].character, 9);
assert_eq!(diags[0].severity, Severity::Error);
assert_eq!(diags[0].source.as_deref(), Some("rustc"));
}
#[test]
fn path_to_uri_prefixes_file_scheme() {
assert_eq!(path_to_uri(Path::new("/tmp/a.rs")), "file:///tmp/a.rs");
}
}
+241 -1
View File
@@ -1 +1,241 @@
// LSP client pool + diagnostics service land here in M5.
//! LSP diagnostics pool: maps a file extension to a language server, lazily spawns one server
//! per language, and answers the core `DiagnosticsSource` seam. Built-in servers are only
//! offered when their binary is on `PATH`; config can add or override them.
//! See `docs/09-integrations.md`.
mod client;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use async_trait::async_trait;
use harness_core::lsp::{Diagnostic, DiagnosticsSource};
use client::LspClient;
/// One language server: which binary to run and which extensions it handles.
#[derive(Debug, Clone)]
pub struct ServerConfig {
pub name: String,
pub command: String,
pub args: Vec<String>,
pub extensions: Vec<String>,
}
/// Built-in servers, tried only when their binary exists on `PATH`.
fn builtin_servers() -> Vec<ServerConfig> {
vec![
ServerConfig {
name: "rust".into(),
command: "rust-analyzer".into(),
args: vec![],
extensions: vec!["rs".into()],
},
ServerConfig {
name: "typescript".into(),
command: "typescript-language-server".into(),
args: vec!["--stdio".into()],
extensions: vec!["ts".into(), "tsx".into(), "js".into(), "jsx".into()],
},
ServerConfig {
name: "go".into(),
command: "gopls".into(),
args: vec![],
extensions: vec!["go".into()],
},
ServerConfig {
name: "python".into(),
command: "pyright-langserver".into(),
args: vec!["--stdio".into()],
extensions: vec!["py".into()],
},
]
}
/// LSP `languageId` for a file extension (falls back to the extension itself).
fn language_id_for_ext(ext: &str) -> &str {
match ext {
"rs" => "rust",
"ts" => "typescript",
"tsx" => "typescriptreact",
"js" => "javascript",
"jsx" => "javascriptreact",
"go" => "go",
"py" => "python",
other => other,
}
}
/// Returns true if `command` is an absolute existing file or resolvable on `PATH`.
fn binary_exists(command: &str) -> bool {
let p = Path::new(command);
if p.is_absolute() {
return p.is_file();
}
let Some(paths) = std::env::var_os("PATH") else {
return false;
};
std::env::split_paths(&paths).any(|dir| dir.join(command).is_file())
}
enum Slot {
Ready(Arc<LspClient>),
Failed,
}
pub struct LspPool {
root: PathBuf,
servers: Vec<ServerConfig>,
/// One slot per server name; absent until first spawn attempt.
clients: tokio::sync::Mutex<HashMap<String, Slot>>,
/// Open documents → last-sent version, so `didChange` bumps monotonically.
open: Mutex<HashMap<PathBuf, i32>>,
}
impl LspPool {
/// Builds a pool: available built-ins plus any `config` servers (which override built-ins
/// by name; an empty command disables a built-in).
pub fn new(root: PathBuf, config: Vec<ServerConfig>) -> Self {
let mut by_name: HashMap<String, ServerConfig> = HashMap::new();
for server in builtin_servers() {
if binary_exists(&server.command) {
by_name.insert(server.name.clone(), server);
}
}
for server in config {
if server.command.is_empty() {
by_name.remove(&server.name);
} else {
by_name.insert(server.name.clone(), server);
}
}
Self {
root,
servers: by_name.into_values().collect(),
clients: tokio::sync::Mutex::new(HashMap::new()),
open: Mutex::new(HashMap::new()),
}
}
fn server_for(&self, path: &Path) -> Option<&ServerConfig> {
let ext = path.extension().and_then(|e| e.to_str())?;
self.servers
.iter()
.find(|s| s.extensions.iter().any(|e| e == ext))
}
/// Gets or lazily spawns the client for `server`. Caches a failure so we don't respawn a
/// broken server on every edit.
async fn client_for(&self, server: &ServerConfig) -> Option<Arc<LspClient>> {
let mut clients = self.clients.lock().await;
match clients.get(&server.name) {
Some(Slot::Ready(c)) => return Some(c.clone()),
Some(Slot::Failed) => return None,
None => {}
}
match LspClient::spawn(&server.command, &server.args, &self.root).await {
Ok(client) => {
let client = Arc::new(client);
clients.insert(server.name.clone(), Slot::Ready(client.clone()));
Some(client)
}
Err(e) => {
tracing::warn!(server = %server.name, error = %e, "LSP server failed to start");
clients.insert(server.name.clone(), Slot::Failed);
None
}
}
}
}
#[async_trait]
impl DiagnosticsSource for LspPool {
async fn touch(&self, path: &Path) {
let Some(server) = self.server_for(path).cloned() else {
return;
};
let Ok(text) = tokio::fs::read_to_string(path).await else {
return;
};
let Some(client) = self.client_for(&server).await else {
return;
};
let ext = path
.extension()
.and_then(|e| e.to_str())
.unwrap_or_default();
let language_id = language_id_for_ext(ext);
let existing_version = {
let mut open = self.open.lock().unwrap();
match open.get_mut(path) {
Some(version) => {
*version += 1;
Some(*version)
}
None => {
open.insert(path.to_path_buf(), 1);
None
}
}
};
// Fresh diagnostics only reflect the current contents — drop any stale set first.
client.clear(path);
match existing_version {
None => client.did_open(path, language_id, 1, &text),
Some(version) => client.did_change(path, version, &text),
}
}
async fn diagnostics(&self, path: &Path, wait: Duration) -> Vec<Diagnostic> {
let Some(server) = self.server_for(path).cloned() else {
return Vec::new();
};
let Some(client) = self.client_for(&server).await else {
return Vec::new();
};
client.wait_diagnostics(path, wait).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn language_ids_map_common_extensions() {
assert_eq!(language_id_for_ext("rs"), "rust");
assert_eq!(language_id_for_ext("tsx"), "typescriptreact");
assert_eq!(language_id_for_ext("unknown"), "unknown");
}
#[test]
fn config_server_overrides_builtin_by_name() {
let pool = LspPool::new(
std::env::temp_dir(),
vec![ServerConfig {
name: "rust".into(),
command: "my-custom-ra".into(),
args: vec!["--flag".into()],
extensions: vec!["rs".into()],
}],
);
let server = pool.server_for(Path::new("/x/a.rs")).unwrap();
assert_eq!(server.command, "my-custom-ra");
}
#[test]
fn no_server_for_unknown_extension() {
let pool = LspPool::new(std::env::temp_dir(), vec![]);
assert!(pool.server_for(Path::new("/x/a.zzz")).is_none());
}
#[test]
fn binary_exists_finds_a_known_tool() {
assert!(binary_exists("sh"));
assert!(!binary_exists("definitely-not-a-real-binary-xyz"));
}
}
+38
View File
@@ -0,0 +1,38 @@
//! End-to-end LSP test against a real `rust-analyzer`. Ignored by default (requires the binary
//! on PATH and is slow — RA indexes the project). Run with:
//! cargo test -p harness-lsp --test rust_analyzer -- --ignored --nocapture
use std::time::Duration;
use harness_core::lsp::{DiagnosticsSource, Severity};
use harness_lsp::LspPool;
#[tokio::test]
#[ignore = "requires rust-analyzer on PATH; slow"]
async fn rust_analyzer_reports_a_type_error() {
// Minimal cargo project with a deliberate type error.
let dir = tempfile::tempdir().unwrap();
std::fs::write(
dir.path().join("Cargo.toml"),
"[package]\nname = \"probe\"\nversion = \"0.1.0\"\nedition = \"2021\"\n",
)
.unwrap();
std::fs::create_dir_all(dir.path().join("src")).unwrap();
let main_rs = dir.path().join("src/main.rs");
std::fs::write(
&main_rs,
"fn main() {\n let x: i32 = \"not an integer\";\n let _ = x;\n}\n",
)
.unwrap();
let pool = LspPool::new(dir.path().to_path_buf(), vec![]);
pool.touch(&main_rs).await;
// Generous wait: rust-analyzer must index the workspace before it reports anything.
let diags = pool.diagnostics(&main_rs, Duration::from_secs(60)).await;
println!("diagnostics: {diags:#?}");
assert!(
diags.iter().any(|d| d.severity == Severity::Error),
"expected at least one error diagnostic, got: {diags:?}"
);
}
+11
View File
@@ -6,6 +6,17 @@ license.workspace = true
[dependencies]
harness-core = { workspace = true }
rmcp = { workspace = true, features = ["client", "transport-child-process"] }
tokio = { workspace = true }
async-trait = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
tracing = { workspace = true }
thiserror = { workspace = true }
[dev-dependencies]
tokio-util = { workspace = true }
tempfile = { workspace = true }
[lints]
workspace = true
+264 -1
View File
@@ -1 +1,264 @@
// MCP stdio client → Tool adapters land here in M5.
//! MCP stdio client → `Tool` adapters (M5). For each configured server we spawn the child
//! over rmcp's `TokioChildProcess` transport, `initialize`, `list_tools`, and wrap every
//! remote tool as an [`McpTool`] named `{server}_{tool}`. Servers are gated behind the `mcp`
//! permission key. See `docs/09-integrations.md`.
//!
//! Out of scope for v1 (matching the doc): resources, prompts, sampling, and non-stdio
//! transports.
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use async_trait::async_trait;
use harness_core::tool::{Tool, ToolCtx, ToolError, ToolOutput};
use rmcp::model::{CallToolRequestParam, RawContent};
use rmcp::service::RunningService;
use rmcp::transport::TokioChildProcess;
use rmcp::{RoleClient, ServiceExt};
use tokio::process::Command;
/// Sanitized-name cap so a `{server}_{tool}` name stays a legal tool identifier.
const MAX_NAME_LEN: usize = 64;
/// One configured MCP server: the child command plus its environment.
#[derive(Debug, Clone)]
pub struct ServerConfig {
pub command: String,
pub args: Vec<String>,
pub env: HashMap<String, String>,
}
#[derive(Debug, thiserror::Error)]
enum ConnectError {
#[error("spawn/transport failed: {0}")]
Transport(#[from] std::io::Error),
#[error("MCP service error: {0}")]
Service(#[from] rmcp::service::ServiceError),
}
/// Connects every configured server and returns ready tool adapters. A server that fails to
/// start (or list its tools) is logged and skipped — the rest are unaffected, and the session
/// still runs with whatever connected. `named_servers` is `{server_name: config}`.
pub async fn connect_all(named_servers: HashMap<String, ServerConfig>) -> Vec<Arc<dyn Tool>> {
let mut tools: Vec<Arc<dyn Tool>> = Vec::new();
for (name, config) in named_servers {
match connect(&name, &config).await {
Ok(mut server_tools) => {
tracing::info!(server = %name, count = server_tools.len(), "MCP server connected");
tools.append(&mut server_tools);
}
Err(e) => {
tracing::warn!(server = %name, error = %e, "MCP server failed to start; skipping");
}
}
}
tools
}
async fn connect(name: &str, config: &ServerConfig) -> Result<Vec<Arc<dyn Tool>>, ConnectError> {
let mut command = Command::new(&config.command);
command.args(&config.args);
for (key, value) in &config.env {
command.env(key, value);
}
// rmcp sets stdin/stdout to piped and kill-on-drop; the child dies with the `RunningService`.
let transport = TokioChildProcess::new(&mut command)?;
let service = Arc::new(().serve(transport).await?);
let remote_tools = service.peer().list_all_tools().await?;
let adapters = remote_tools
.into_iter()
.map(|tool| {
let full_name = qualified_name(name, &tool.name);
let parameters = serde_json::Value::Object((*tool.input_schema).clone());
Arc::new(McpTool {
full_name,
remote_name: tool.name.to_string(),
description: tool.description.to_string(),
parameters,
service: service.clone(),
}) as Arc<dyn Tool>
})
.collect();
Ok(adapters)
}
/// `{server}_{tool}` sanitized to `[a-zA-Z0-9_-]` and capped at [`MAX_NAME_LEN`] chars.
fn qualified_name(server: &str, tool: &str) -> String {
let mut name: String = format!("{server}_{tool}")
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '_' || c == '-' {
c
} else {
'_'
}
})
.collect();
name.truncate(MAX_NAME_LEN);
name
}
/// A single remote MCP tool exposed to the engine as a `Tool`. Holds a shared handle to the
/// server's `RunningService` (kept alive for the whole session so the child stays up).
struct McpTool {
/// Engine-facing name: sanitized `{server}_{tool}`; also the permission pattern.
full_name: String,
/// The server's own tool name, sent back verbatim in `call_tool`.
remote_name: String,
description: String,
parameters: serde_json::Value,
service: Arc<RunningService<RoleClient, ()>>,
}
#[async_trait]
impl Tool for McpTool {
fn name(&self) -> &str {
&self.full_name
}
fn description(&self) -> &str {
&self.description
}
fn parameters(&self) -> serde_json::Value {
self.parameters.clone()
}
async fn execute(
&self,
input: serde_json::Value,
ctx: ToolCtx,
) -> Result<ToolOutput, ToolError> {
ctx.ask
.ask(
"mcp",
self.full_name.clone(),
self.full_name.clone(),
input.clone(),
)
.await?;
let arguments = match input {
serde_json::Value::Object(map) => Some(map),
serde_json::Value::Null => None,
other => {
return Err(ToolError::Invalid(format!(
"MCP tool arguments must be a JSON object, got {other}"
)))
}
};
let result = self
.service
.peer()
.call_tool(CallToolRequestParam {
name: self.remote_name.clone().into(),
arguments,
})
.await
.map_err(|e| ToolError::Other(e.to_string()))?;
let mut text = String::new();
let mut image_index = 0;
for content in &result.content {
match &content.raw {
RawContent::Text(t) => {
if !text.is_empty() {
text.push('\n');
}
text.push_str(&t.text);
}
RawContent::Image(image) => {
let note = save_image(&ctx.data_dir, &self.full_name, image_index, image).await;
if !text.is_empty() {
text.push('\n');
}
text.push_str(&note);
image_index += 1;
}
RawContent::Resource(resource) => {
let embedded = resource_text(resource);
if !embedded.is_empty() {
if !text.is_empty() {
text.push('\n');
}
text.push_str(&embedded);
}
}
}
}
// MCP surfaces tool-level failures as `is_error` with the message in `content`; map
// that to a tool error so the model sees it as a failed call rather than a result.
if result.is_error.unwrap_or(false) {
return Err(ToolError::Other(if text.is_empty() {
"MCP tool reported an error".to_string()
} else {
text
}));
}
Ok(ToolOutput::new(self.full_name.clone(), text))
}
}
/// Writes an image payload to the session data dir and returns a one-line note for the tool
/// output. Best-effort: a write failure still yields a note (without a path).
async fn save_image(
data_dir: &PathBuf,
tool_name: &str,
index: usize,
image: &rmcp::model::RawImageContent,
) -> String {
let ext = image.mime_type.rsplit('/').next().unwrap_or("bin");
let file_name = format!("{tool_name}-image-{index}.{ext}.b64");
let path = data_dir.join(&file_name);
let saved = tokio::fs::create_dir_all(data_dir).await.is_ok()
&& tokio::fs::write(&path, &image.data).await.is_ok();
if saved {
format!(
"[image: {} ({} base64 bytes) saved to {}]",
image.mime_type,
image.data.len(),
path.display()
)
} else {
format!(
"[image: {} ({} base64 bytes, not saved)]",
image.mime_type,
image.data.len()
)
}
}
/// Best-effort text extraction from an embedded resource (text resources only in v1).
fn resource_text(resource: &rmcp::model::RawEmbeddedResource) -> String {
match &resource.resource {
rmcp::model::ResourceContents::TextResourceContents { text, .. } => text.clone(),
_ => String::new(),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn qualified_name_prefixes_and_sanitizes() {
assert_eq!(qualified_name("fs", "read_file"), "fs_read_file");
assert_eq!(
qualified_name("my.server", "do/thing"),
"my_server_do_thing"
);
}
#[test]
fn qualified_name_caps_length() {
let long_tool = "t".repeat(100);
let name = qualified_name("srv", &long_tool);
assert_eq!(name.len(), MAX_NAME_LEN);
assert!(name.starts_with("srv_t"));
}
}
+153
View File
@@ -0,0 +1,153 @@
//! End-to-end MCP integration test: spawn a real stdio MCP server (a small Python fixture),
//! connect through the real rmcp client, and verify a discovered tool is callable and gated
//! behind the `mcp` permission key. This is the M5 milestone's ✅ for MCP.
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use harness_core::event::{AppEvent, EventBus};
use harness_core::permission::{PermissionReply, PermissionService};
use harness_core::tool::{MetadataSink, PermissionHandle, ToolCtx, ToolError};
use harness_core::types::{MessageId, SessionId};
use harness_mcp::{connect_all, ServerConfig};
use tokio_util::sync::CancellationToken;
/// A permission frontend that replies `Once` to every ask and records the `permission`/`pattern`
/// of each, so a test can assert the call was actually gated.
fn recording_auto_approve(
bus: EventBus,
service: Arc<PermissionService>,
) -> Arc<Mutex<Vec<(String, String)>>> {
let asks = Arc::new(Mutex::new(Vec::new()));
let asks_task = asks.clone();
// Subscribe before spawning: a subscription created inside the task could miss the ask
// (tokio broadcast only delivers to receivers that exist at publish time).
let mut rx = bus.subscribe();
tokio::spawn(async move {
while let Ok(event) = rx.recv().await {
if let AppEvent::PermissionAsked { request } = event {
asks_task
.lock()
.unwrap()
.push((request.permission.clone(), request.pattern.clone()));
service.reply(&request.id, PermissionReply::Once);
}
}
});
asks
}
struct Harness {
ctx_data_dir: std::path::PathBuf,
service: Arc<PermissionService>,
}
impl Harness {
fn ctx(&self) -> ToolCtx {
let (metadata, _rx) = MetadataSink::channel();
ToolCtx {
session_id: SessionId::new(),
message_id: MessageId::new(),
call_id: "call_1".into(),
data_dir: self.ctx_data_dir.clone(),
cwd: std::env::temp_dir(),
cancel: CancellationToken::new(),
ask: PermissionHandle::new(
self.service.clone(),
SessionId::new(),
Vec::new(),
Arc::new(Mutex::new(Vec::new())),
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
}
fn fixture_server() -> HashMap<String, ServerConfig> {
let script = concat!(env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/echo_server.py");
HashMap::from([(
"fix".to_string(),
ServerConfig {
command: "python3".to_string(),
args: vec![script.to_string()],
env: HashMap::new(),
},
)])
}
#[tokio::test]
async fn discovers_and_calls_a_real_mcp_tool_with_permission() {
let tools = connect_all(fixture_server()).await;
let names: Vec<_> = tools.iter().map(|t| t.name().to_string()).collect();
assert!(
names.contains(&"fix_echo".to_string()),
"expected fix_echo among {names:?}"
);
assert!(names.contains(&"fix_boom".to_string()));
let echo = tools.iter().find(|t| t.name() == "fix_echo").unwrap();
// Schema passes through untouched from the server.
assert_eq!(echo.parameters()["properties"]["text"]["type"], "string");
let bus = EventBus::new();
let service = Arc::new(PermissionService::new(bus.clone()));
let asks = recording_auto_approve(bus, service.clone());
let dir = tempfile::tempdir().unwrap();
let harness = Harness {
ctx_data_dir: dir.path().to_path_buf(),
service,
};
let out = echo
.execute(serde_json::json!({"text": "hi there"}), harness.ctx())
.await
.expect("echo call succeeds");
assert_eq!(out.output, "hi there");
// The call was gated on the `mcp` key with the qualified tool name as the pattern.
let recorded = asks.lock().unwrap().clone();
assert_eq!(recorded, vec![("mcp".to_string(), "fix_echo".to_string())]);
}
#[tokio::test]
async fn tool_error_result_maps_to_tool_error() {
let tools = connect_all(fixture_server()).await;
let boom = tools.iter().find(|t| t.name() == "fix_boom").unwrap();
let bus = EventBus::new();
let service = Arc::new(PermissionService::new(bus.clone()));
let _asks = recording_auto_approve(bus, service.clone());
let dir = tempfile::tempdir().unwrap();
let harness = Harness {
ctx_data_dir: dir.path().to_path_buf(),
service,
};
let err = boom
.execute(serde_json::json!({}), harness.ctx())
.await
.expect_err("boom reports an error result");
match err {
ToolError::Other(msg) => assert_eq!(msg, "kaboom"),
other => panic!("expected ToolError::Other, got {other:?}"),
}
}
#[tokio::test]
async fn a_failed_server_is_skipped_not_fatal() {
let servers = HashMap::from([(
"broken".to_string(),
ServerConfig {
command: "definitely-not-a-real-binary-xyz".to_string(),
args: vec![],
env: HashMap::new(),
},
)]);
// No panic, no tools — the missing server is logged and skipped.
let tools = connect_all(servers).await;
assert!(tools.is_empty());
}
+86
View File
@@ -0,0 +1,86 @@
#!/usr/bin/env python3
"""Minimal MCP stdio server fixture for harness-mcp integration tests.
Speaks newline-delimited JSON-RPC (the framing rmcp's child-process transport uses) and
implements just enough of the protocol to be discovered and called: `initialize`,
`notifications/initialized`, `tools/list`, and `tools/call`. Exposes one tool, `echo`,
which returns its `text` argument, plus `boom`, which returns an error result.
"""
import json
import sys
PROTOCOL_VERSION = "2024-11-05"
TOOLS = [
{
"name": "echo",
"description": "Returns the text it is given.",
"inputSchema": {
"type": "object",
"properties": {"text": {"type": "string"}},
"required": ["text"],
},
},
{
"name": "boom",
"description": "Always fails.",
"inputSchema": {"type": "object", "properties": {}},
},
]
def reply(msg_id, result):
sys.stdout.write(json.dumps({"jsonrpc": "2.0", "id": msg_id, "result": result}) + "\n")
sys.stdout.flush()
def main():
# readline() rather than `for line in sys.stdin`: the latter's read-ahead buffer blocks
# until it fills, which would stall the JSON-RPC handshake line-by-line.
while True:
line = sys.stdin.readline()
if line == "": # EOF: parent closed stdin
break
line = line.strip()
if not line:
continue
msg = json.loads(line)
method = msg.get("method")
msg_id = msg.get("id")
if method == "initialize":
reply(msg_id, {
"protocolVersion": PROTOCOL_VERSION,
"capabilities": {"tools": {}},
"serverInfo": {"name": "echo-fixture", "version": "0.1.0"},
})
elif method == "notifications/initialized":
pass # notification: no response
elif method == "tools/list":
reply(msg_id, {"tools": TOOLS})
elif method == "tools/call":
params = msg.get("params") or {}
name = params.get("name")
args = params.get("arguments") or {}
if name == "echo":
reply(msg_id, {
"content": [{"type": "text", "text": args.get("text", "")}],
"isError": False,
})
elif name == "boom":
reply(msg_id, {
"content": [{"type": "text", "text": "kaboom"}],
"isError": True,
})
else:
reply(msg_id, {
"content": [{"type": "text", "text": f"unknown tool {name}"}],
"isError": True,
})
elif msg_id is not None:
# Unknown request: empty result keeps the client happy.
reply(msg_id, {})
if __name__ == "__main__":
main()
+50 -3
View File
@@ -8,6 +8,8 @@ use tokio_util::sync::CancellationToken;
use crate::codec::{openai_chat, openai_responses};
const DEFAULT_BASE_URL: &str = "https://api.openai.com/v1";
/// OpenCode Zen's OpenAI-compatible gateway (chat-completions only).
const OPENCODE_BASE_URL: &str = "https://opencode.ai/zen/v1";
/// Which OpenAI wire format to speak for a given model.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -33,28 +35,57 @@ fn flavor_for(model: &str) -> ApiFlavor {
}
pub struct OpenAiProvider {
id: String,
api_key: String,
base_url: String,
/// When set, always speak chat-completions regardless of model name — required for
/// OpenAI-compatible gateways (e.g. OpenCode Zen) that don't implement `/responses`.
chat_only: bool,
client: reqwest::Client,
}
impl OpenAiProvider {
pub fn new(api_key: impl Into<String>) -> Self {
Self {
id: "openai".to_string(),
api_key: api_key.into(),
base_url: DEFAULT_BASE_URL.to_string(),
chat_only: false,
client: reqwest::Client::new(),
}
}
pub fn with_base_url(api_key: impl Into<String>, base_url: impl Into<String>) -> Self {
Self {
id: "openai".to_string(),
api_key: api_key.into(),
base_url: base_url.into(),
chat_only: false,
client: reqwest::Client::new(),
}
}
/// An OpenCode Zen provider: id `opencode`, chat-completions only, defaulting to Zen's
/// gateway. Pass `base_url = None` to use the default endpoint.
pub fn opencode(api_key: impl Into<String>, base_url: Option<String>) -> Self {
Self {
id: "opencode".to_string(),
api_key: api_key.into(),
base_url: base_url.unwrap_or_else(|| OPENCODE_BASE_URL.to_string()),
chat_only: true,
client: reqwest::Client::new(),
}
}
/// The wire format for `model`, honoring `chat_only`.
fn flavor(&self, model: &str) -> ApiFlavor {
if self.chat_only {
ApiFlavor::Chat
} else {
flavor_for(model)
}
}
fn classify_error(
status: reqwest::StatusCode,
body: String,
@@ -77,7 +108,7 @@ impl OpenAiProvider {
#[async_trait]
impl Provider for OpenAiProvider {
fn id(&self) -> &str {
"openai"
&self.id
}
async fn list_models(&self) -> Result<Vec<ModelInfo>, ProviderError> {
@@ -90,7 +121,7 @@ impl Provider for OpenAiProvider {
req: LlmRequest,
cancel: CancellationToken,
) -> Result<LlmEventStream, ProviderError> {
let (path, body) = match flavor_for(&req.model) {
let (path, body) = match self.flavor(&req.model) {
ApiFlavor::Responses => ("/responses", openai_responses::build_request(&req)),
ApiFlavor::Chat => ("/chat/completions", openai_chat::build_request(&req)),
};
@@ -120,7 +151,7 @@ impl Provider for OpenAiProvider {
return Err(Self::classify_error(status, body, retry_after));
}
Ok(match flavor_for(&req.model) {
Ok(match self.flavor(&req.model) {
ApiFlavor::Responses => openai_responses::decode(response.bytes_stream()),
ApiFlavor::Chat => openai_chat::decode(response.bytes_stream()),
})
@@ -151,6 +182,22 @@ mod tests {
assert_eq!(OpenAiProvider::new("k").id(), "openai");
}
#[test]
fn opencode_provider_is_chat_only_and_ids_as_opencode() {
let p = OpenAiProvider::opencode("k", None);
assert_eq!(p.id(), "opencode");
assert_eq!(p.base_url, OPENCODE_BASE_URL);
// Even a gpt-*/o* model must route to chat completions on a chat-only gateway.
assert_eq!(p.flavor("gpt-5.5"), ApiFlavor::Chat);
assert_eq!(p.flavor("claude-sonnet-5"), ApiFlavor::Chat);
}
#[test]
fn opencode_honors_custom_base_url() {
let p = OpenAiProvider::opencode("k", Some("https://example.test/v1".into()));
assert_eq!(p.base_url, "https://example.test/v1");
}
#[test]
fn classifies_context_overflow_from_400_body() {
assert!(matches!(
+3
View File
@@ -149,6 +149,9 @@ mod tests {
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
+73
View File
@@ -0,0 +1,73 @@
//! Shared LSP-diagnostics reporting for the edit/write tools. After a file is written we ask
//! the (optional) diagnostics source to re-analyze it and append any error-severity items to
//! the tool output so the model sees mistakes it just introduced. Best-effort: no source, a
//! slow server, or a timeout all just mean "no diagnostics" — never a tool failure.
use std::path::Path;
use std::time::Duration;
use harness_core::lsp::Severity;
use harness_core::tool::{ToolCtx, ToolOutput};
/// docs/09-integrations.md: wait up to 1.5s for the server to (re)publish after the change.
const DIAGNOSTICS_WAIT: Duration = Duration::from_millis(1500);
/// Touches `path` in the language server and appends error-severity diagnostics to `output`
/// (both as a human-readable block in the text and the full set in metadata under `diagnostics`).
pub async fn append_diagnostics(
ctx: &ToolCtx,
path: &Path,
display_name: &str,
output: &mut ToolOutput,
) {
let Some(source) = &ctx.diagnostics else {
return;
};
source.touch(path).await;
let diagnostics = source.diagnostics(path, DIAGNOSTICS_WAIT).await;
if diagnostics.is_empty() {
return;
}
let errors: Vec<_> = diagnostics
.iter()
.filter(|d| d.severity == Severity::Error)
.collect();
if !errors.is_empty() {
output
.output
.push_str("\n\nLSP errors detected in this file, please fix:");
for diag in &errors {
output.output.push('\n');
output.output.push_str(&diag.display_line(display_name));
}
}
// Full set (all severities) into metadata for the TUI.
let items: Vec<serde_json::Value> = diagnostics
.iter()
.map(|d| {
serde_json::json!({
"line": d.line,
"character": d.character,
"severity": severity_str(d.severity),
"message": d.message,
"source": d.source,
})
})
.collect();
if let serde_json::Value::Object(map) = &mut output.metadata {
map.insert("diagnostics".into(), serde_json::Value::Array(items));
} else {
output.metadata = serde_json::json!({ "diagnostics": items });
}
}
fn severity_str(severity: Severity) -> &'static str {
match severity {
Severity::Error => "error",
Severity::Warning => "warning",
Severity::Info => "info",
Severity::Hint => "hint",
}
}
+5 -1
View File
@@ -196,8 +196,9 @@ impl Tool for EditTool {
.map_err(|e| ToolError::Other(format!("{}: {e}", path.display())))?;
let (added, removed) = diff_stats(&content_old, &content_new);
let mut output = ToolOutput::new(pattern, "Edit applied successfully.".to_string());
let mut output = ToolOutput::new(pattern.clone(), "Edit applied successfully.".to_string());
output.metadata = serde_json::json!({"diff": diff, "added": added, "removed": removed});
crate::diagnostics::append_diagnostics(&ctx, &path, &pattern, &mut output).await;
Ok(output)
}
}
@@ -232,6 +233,9 @@ mod tests {
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
+3
View File
@@ -130,6 +130,9 @@ mod tests {
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
+3
View File
@@ -159,6 +159,9 @@ mod tests {
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
+20
View File
@@ -1,9 +1,12 @@
mod bash;
mod diagnostics;
mod edit;
mod glob;
mod grep;
mod paths;
mod read;
mod skill;
mod task;
mod write;
pub use bash::BashTool;
@@ -11,6 +14,8 @@ pub use edit::EditTool;
pub use glob::GlobTool;
pub use grep::GrepTool;
pub use read::ReadTool;
pub use skill::SkillTool;
pub use task::TaskTool;
pub use write::WriteTool;
use std::sync::Arc;
@@ -26,3 +31,18 @@ pub fn register_builtins(registry: &mut ToolRegistry) {
registry.register(Arc::new(GlobTool));
registry.register(Arc::new(GrepTool));
}
/// Registers the multiagent `task` tool (M4). Kept separate from [`register_builtins`] so
/// non-orchestrating contexts can omit delegation; the tool no-ops with an error if the
/// session has no spawner wired in.
pub fn register_task_tool(registry: &mut ToolRegistry) {
registry.register(Arc::new(TaskTool));
}
/// Registers the `skill` tool (M5) over a loaded skill set. No-op when there are no skills,
/// so the tool is only advertised when something can be loaded.
pub fn register_skill_tool(registry: &mut ToolRegistry, skills: &[harness_core::config::SkillDef]) {
if !skills.is_empty() {
registry.register(Arc::new(SkillTool::new(skills)));
}
}
+10
View File
@@ -102,6 +102,13 @@ impl Tool for ReadTool {
"(empty file or offset past end)".to_string(),
));
}
// In a subagent session, advertise this read on the job board.
if let Some(reporter) = &ctx.context_reporter {
let reported = paths::relative_pattern(&ctx.cwd, &path);
reporter.report_file(reported, numbered.len() as u32).await;
}
Ok(ToolOutput::new(
params.file_path.clone(),
numbered.join("\n"),
@@ -139,6 +146,9 @@ mod tests {
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
+154
View File
@@ -0,0 +1,154 @@
//! The `skill` tool (M5): the system prompt advertises each skill's name + description; when
//! the model decides a skill is relevant it calls this tool with the skill name to pull the
//! full instructions on demand. See `docs/09-integrations.md`.
use std::collections::HashMap;
use async_trait::async_trait;
use harness_core::config::SkillDef;
use harness_core::tool::{Tool, ToolCtx, ToolError, ToolOutput};
use schemars::JsonSchema;
use serde::Deserialize;
#[derive(Debug, Deserialize, JsonSchema)]
struct SkillParams {
/// The name of the skill to load, as advertised in the system prompt.
name: String,
}
/// Serves skill bodies by name. Built from the loaded skill set; if there are no skills the
/// caller simply doesn't register the tool.
pub struct SkillTool {
/// name → (description, body).
skills: HashMap<String, (String, String)>,
/// Sorted names, for a stable "unknown skill" hint.
names: Vec<String>,
description: String,
}
impl SkillTool {
pub fn new(skills: &[SkillDef]) -> Self {
let mut names: Vec<String> = skills.iter().map(|s| s.name.clone()).collect();
names.sort();
let map = skills
.iter()
.map(|s| (s.name.clone(), (s.description.clone(), s.body.clone())))
.collect();
let description = format!(
"Load the full instructions for a named skill before doing the related work. \
Available skills: {}.",
names.join(", ")
);
Self {
skills: map,
names,
description,
}
}
}
#[async_trait]
impl Tool for SkillTool {
fn name(&self) -> &str {
"skill"
}
fn description(&self) -> &str {
&self.description
}
fn parameters(&self) -> serde_json::Value {
serde_json::to_value(schemars::schema_for!(SkillParams)).unwrap()
}
async fn execute(
&self,
input: serde_json::Value,
_ctx: ToolCtx,
) -> Result<ToolOutput, ToolError> {
let params: SkillParams =
serde_json::from_value(input).map_err(|e| ToolError::Invalid(e.to_string()))?;
match self.skills.get(&params.name) {
Some((_description, body)) => Ok(ToolOutput::new(
format!("skill: {}", params.name),
body.clone(),
)),
None => Err(ToolError::Invalid(format!(
"unknown skill {:?}; available: {}",
params.name,
self.names.join(", ")
))),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use harness_core::event::EventBus;
use harness_core::permission::{spawn_auto_approve, PermissionService};
use harness_core::tool::{MetadataSink, PermissionHandle};
use harness_core::types::SessionId;
use std::sync::{Arc, Mutex};
use tokio_util::sync::CancellationToken;
fn ctx() -> ToolCtx {
let bus = EventBus::new();
let service = Arc::new(PermissionService::new(bus.clone()));
spawn_auto_approve(bus, service.clone());
let (metadata, _rx) = MetadataSink::channel();
ToolCtx {
session_id: SessionId::new(),
message_id: harness_core::types::MessageId::new(),
call_id: "c1".into(),
data_dir: std::env::temp_dir(),
cwd: std::env::temp_dir(),
cancel: CancellationToken::new(),
ask: PermissionHandle::new(
service,
SessionId::new(),
Vec::new(),
Arc::new(Mutex::new(Vec::new())),
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
fn sample() -> Vec<SkillDef> {
vec![SkillDef {
name: "formatter".into(),
description: "format code".into(),
body: "Run cargo fmt.".into(),
}]
}
#[tokio::test]
async fn returns_skill_body_by_name() {
let tool = SkillTool::new(&sample());
let out = tool
.execute(serde_json::json!({"name": "formatter"}), ctx())
.await
.unwrap();
assert_eq!(out.output, "Run cargo fmt.");
}
#[tokio::test]
async fn unknown_skill_is_an_input_error() {
let tool = SkillTool::new(&sample());
let err = tool
.execute(serde_json::json!({"name": "nope"}), ctx())
.await
.unwrap_err();
assert!(matches!(err, ToolError::Invalid(_)));
}
#[test]
fn description_lists_available_skills() {
let tool = SkillTool::new(&sample());
assert!(tool.description().contains("formatter"));
}
}
+140
View File
@@ -0,0 +1,140 @@
//! The `task` tool: delegate work to a specialist subagent, foreground or background.
//!
//! This tool is deliberately thin — it validates input, gates on a `task/<agent>` permission,
//! and hands off to the engine's `SubagentSpawner` (owned by the composition root), which
//! resolves the agent, enforces the depth limit, applies permission intersection, and runs
//! the child session. See `docs/04-multiagent.md`.
use async_trait::async_trait;
use serde::Deserialize;
use serde_json::json;
use harness_core::tool::{
invalid_input, SpawnError, SpawnRequest, Tool, ToolCtx, ToolError, ToolOutput,
};
#[derive(Debug, Deserialize)]
struct TaskInput {
/// Short human-facing label for the subtask (shown on the job board).
description: String,
/// The full instruction handed to the subagent.
prompt: String,
/// Which specialist to run (must be a subagent-capable agent).
subagent_type: String,
/// Alias or task id of a completed job to continue instead of starting fresh.
#[serde(default)]
task_id: Option<String>,
/// Run in the background and return immediately (tracked on the job board).
#[serde(default)]
background: bool,
}
pub struct TaskTool;
#[async_trait]
impl Tool for TaskTool {
fn name(&self) -> &str {
"task"
}
fn description(&self) -> &str {
"Delegate a self-contained unit of work to a specialist subagent. Set `background: \
true` to launch it without blocking (track progress on the job board); reuse a \
completed subagent by passing its `task_id`/alias to continue the same session."
}
fn parameters(&self) -> serde_json::Value {
json!({
"type": "object",
"properties": {
"description": { "type": "string", "description": "Short label for the subtask." },
"prompt": { "type": "string", "description": "Full instruction for the subagent." },
"subagent_type": { "type": "string", "description": "Specialist to run." },
"task_id": { "type": "string", "description": "Alias/id of a completed job to reuse." },
"background": { "type": "boolean", "description": "Launch without blocking." }
},
"required": ["description", "prompt", "subagent_type"]
})
}
async fn execute(
&self,
input: serde_json::Value,
ctx: ToolCtx,
) -> Result<ToolOutput, ToolError> {
let args: TaskInput = serde_json::from_value(input).map_err(|e| invalid_input(self, e))?;
let Some(spawner) = ctx.spawner.clone() else {
return Err(ToolError::Other(
"subagent delegation is not available in this session".into(),
));
};
// Gate on task/<agent>. `Always` grants blanket delegation to this specialist.
ctx.ask
.ask(
"task",
&args.subagent_type,
&args.subagent_type,
json!({
"agent": args.subagent_type,
"description": args.description,
"background": args.background,
}),
)
.await?;
let req = SpawnRequest {
parent_session_id: ctx.session_id.clone(),
parent_message_id: ctx.message_id.clone(),
agent: args.subagent_type.clone(),
description: args.description.clone(),
prompt: args.prompt,
reuse_task_id: args.task_id,
background: args.background,
cancel: ctx.cancel.clone(),
};
let outcome = spawner.spawn(req).await.map_err(map_spawn_error)?;
if outcome.background {
let alias = outcome.alias.unwrap_or_default();
Ok(ToolOutput {
title: format!("launched {} ({alias})", args.subagent_type),
output: format!(
"Launched background task {alias} ({}). Check the job board; do not poll — \
wait for completion.",
outcome.child_session_id
),
metadata: json!({
"child_session": outcome.child_session_id,
"agent": args.subagent_type,
"alias": alias,
"background": true,
}),
})
} else {
let text = outcome.final_text.unwrap_or_default();
Ok(ToolOutput {
title: format!("{} — {}", args.subagent_type, args.description),
output: text,
metadata: json!({
"child_session": outcome.child_session_id,
"agent": args.subagent_type,
"background": false,
}),
})
}
}
}
fn map_spawn_error(err: SpawnError) -> ToolError {
match err {
// Depth/agent problems are the model's to fix — surface as tool errors it can read
// and act on, not hard failures.
SpawnError::DepthExceeded => ToolError::Other(err.to_string()),
SpawnError::InvalidAgent(_) => ToolError::Invalid(err.to_string()),
SpawnError::ReuseNotFound(_) => ToolError::Invalid(err.to_string()),
SpawnError::Other(msg) => ToolError::Other(msg),
}
}
+4
View File
@@ -80,6 +80,7 @@ impl Tool for WriteTool {
format!("wrote {} bytes", params.content.len()),
);
output.metadata = serde_json::json!({"diff": diff, "added": added, "removed": removed});
crate::diagnostics::append_diagnostics(&ctx, &path, &params.file_path, &mut output).await;
Ok(output)
}
}
@@ -114,6 +115,9 @@ mod tests {
CancellationToken::new(),
),
metadata,
spawner: None,
context_reporter: None,
diagnostics: None,
}
}
+1 -1
View File
@@ -31,7 +31,7 @@ pub struct App {
impl App {
pub async fn new(cwd: PathBuf) -> anyhow::Result<Self> {
let engine = EngineHandle::init(cwd)?;
let engine = EngineHandle::init(cwd).await?;
let bus_rx = engine.bus().subscribe();
let config = engine.config();
+126 -2
View File
@@ -9,7 +9,16 @@ use crate::state::{AppState, ModalState};
#[derive(Debug)]
pub enum InputAction {
None,
Submit { text: String },
Submit {
text: String,
},
/// A user-defined slash command: `name` (no leading `/`), its `args`, and the `raw` input
/// to fall back to submitting verbatim if no such command is defined.
RunCommand {
name: String,
args: String,
raw: String,
},
Abort,
Quit,
LoadSessions,
@@ -23,6 +32,10 @@ pub enum InputAction {
SessionPickerUp,
SessionPickerDown,
SessionPickerSelect,
OpenJobs,
JobsUp,
JobsDown,
JobsDrillIn,
}
/// Translate a crossterm event into an action and/or mutate `state` directly.
@@ -49,10 +62,21 @@ fn handle_key(key: KeyEvent, state: &mut AppState) -> InputAction {
match &state.modal {
ModalState::Permission { .. } => handle_permission_key(key, state),
ModalState::SessionPicker { .. } => handle_session_picker_key(key, state),
ModalState::JobsPane { .. } => handle_jobs_pane_key(key, state),
ModalState::None => handle_normal_key(key, state),
}
}
fn handle_jobs_pane_key(key: KeyEvent, _state: &mut AppState) -> InputAction {
match key.code {
KeyCode::Up => InputAction::JobsUp,
KeyCode::Down => InputAction::JobsDown,
KeyCode::Enter => InputAction::JobsDrillIn,
KeyCode::Esc => InputAction::CloseModal,
_ => InputAction::None,
}
}
fn handle_permission_key(key: KeyEvent, _state: &mut AppState) -> InputAction {
match key.code {
KeyCode::Char('y') | KeyCode::Char('Y') => {
@@ -119,6 +143,7 @@ fn handle_normal_key(key: KeyEvent, state: &mut AppState) -> InputAction {
}
}
KeyCode::Char('s') if ctrl => InputAction::LoadSessions,
KeyCode::Char('j') if ctrl => InputAction::OpenJobs,
KeyCode::Up => {
let (row, _) = state.input.cursor();
if row == 0 {
@@ -158,8 +183,15 @@ fn parse_slash_command(text: &str) -> Option<InputAction> {
"/model" if !rest.is_empty() => Some(InputAction::SetModel(rest)),
"/agent" if !rest.is_empty() => Some(InputAction::SetAgent(rest)),
"/sessions" => Some(InputAction::LoadSessions),
"/jobs" => Some(InputAction::OpenJobs),
"/quit" => Some(InputAction::Quit),
_ => None,
// Any other `/word` is treated as a user-defined command, resolved against the engine's
// loaded commands when applied; if none matches, the raw text is submitted as-is.
other => Some(InputAction::RunCommand {
name: other.trim_start_matches('/').to_string(),
args: rest,
raw: trimmed.to_string(),
}),
}
}
@@ -173,6 +205,30 @@ pub async fn apply_action(action: InputAction, state: &mut AppState, engine: &En
}
}
}
InputAction::RunCommand { name, args, raw } => {
let Some(session_id) = state.session_id.clone() else {
return;
};
match engine.command(&name) {
Some(cmd) => {
let text = cmd.expand(&args);
// A command may switch model/agent for this one run only.
let model_ref = cmd.model.clone().unwrap_or_else(|| state.model_ref.clone());
if let Err(e) = engine
.prompt_with(session_id, text, &model_ref, cmd.agent.as_deref())
.await
{
tracing::error!(error = %e, "command prompt failed");
}
}
// Unknown command: submit the original text as an ordinary message.
None => {
if let Err(e) = engine.prompt(session_id, raw, &state.model_ref).await {
tracing::error!(error = %e, "prompt failed");
}
}
}
}
InputAction::Abort => {
if let Some(session_id) = state.session_id.clone() {
engine.abort(&session_id);
@@ -235,6 +291,45 @@ pub async fn apply_action(action: InputAction, state: &mut AppState, engine: &En
}
}
}
InputAction::OpenJobs => {
if let Some(session_id) = state.session_id.clone() {
match engine.jobs(session_id).await {
Ok(jobs) => {
state.set_jobs(jobs);
state.open_jobs_pane();
}
Err(e) => tracing::error!(error = %e, "failed to load jobs"),
}
}
}
InputAction::JobsUp => {
if let ModalState::JobsPane { selected } = &mut state.modal {
*selected = selected.saturating_sub(1);
}
state.dirty = true;
}
InputAction::JobsDown => {
let job_count = state.jobs.len();
if let ModalState::JobsPane { selected } = &mut state.modal {
if *selected + 1 < job_count {
*selected += 1;
}
}
state.dirty = true;
}
InputAction::JobsDrillIn => {
if let Some(child) = state.selected_job_child() {
match engine.get_session(child).await {
Ok(Some(session)) => {
if let Err(e) = modal::select_session(state, engine, session).await {
tracing::error!(error = %e, "failed to open subtask session");
}
}
Ok(None) => tracing::warn!("subtask session no longer exists"),
Err(e) => tracing::error!(error = %e, "failed to load subtask session"),
}
}
}
InputAction::None => {}
}
}
@@ -261,4 +356,33 @@ mod tests {
assert!(state.dirty, "a keystroke must request a redraw");
assert_eq!(state.input.lines().join("\n"), "x");
}
#[test]
fn builtin_slash_commands_still_parse() {
assert!(matches!(
parse_slash_command("/new"),
Some(InputAction::NewSession)
));
assert!(matches!(
parse_slash_command("/model openai/gpt-5"),
Some(InputAction::SetModel(m)) if m == "openai/gpt-5"
));
}
#[test]
fn unknown_slash_becomes_a_run_command() {
match parse_slash_command("/deploy prod now") {
Some(InputAction::RunCommand { name, args, raw }) => {
assert_eq!(name, "deploy");
assert_eq!(args, "prod now");
assert_eq!(raw, "/deploy prod now");
}
other => panic!("expected RunCommand, got {other:?}"),
}
}
#[test]
fn non_slash_text_is_not_a_command() {
assert!(parse_slash_command("hello world").is_none());
}
}
+1 -1
View File
@@ -50,7 +50,7 @@ async fn run_headless(args: &[String]) -> i32 {
}
};
let app = match harness_app::App::init(cwd) {
let app = match harness_app::App::init(cwd).await {
Ok(app) => app,
Err(e) => {
eprintln!("error: {e}");
+5 -1
View File
@@ -26,13 +26,17 @@ pub async fn select_session(
session: Session,
) -> Result<(), harness_app::AppError> {
let session_id = session.id.clone();
let messages = engine.session_messages(session_id).await?;
let messages = engine.session_messages(session_id.clone()).await?;
let mut all_parts = Vec::new();
for message in &messages {
let parts = engine.message_parts(message.id.clone()).await?;
all_parts.extend(parts);
}
state.set_session(session, messages, all_parts);
// Surface any subtasks this session spawned (drill-in and session-switch both land here).
if let Ok(jobs) = engine.jobs(session_id).await {
state.set_jobs(jobs);
}
state.modal = ModalState::None;
Ok(())
}
+137 -1
View File
@@ -33,6 +33,9 @@ pub fn render(frame: &mut Frame, state: &mut AppState) {
ModalState::SessionPicker { sessions, selected } => {
render_session_picker(frame, sessions, *selected, area);
}
ModalState::JobsPane { selected } => {
render_jobs_pane(frame, &state.jobs, *selected, area);
}
ModalState::None => {}
}
}
@@ -169,7 +172,7 @@ fn render_status(frame: &mut Frame, state: &AppState, area: Rect) {
let hints = if state.ctrl_c_pressed {
"Press Ctrl+C again to quit"
} else {
"Ctrl+S: sessions | Esc: abort | Ctrl+C: quit"
"Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C: quit"
};
let usage = format!(
"{} · ${:.4}",
@@ -242,6 +245,86 @@ fn render_session_picker(
frame.render_widget(paragraph, popup);
}
fn render_jobs_pane(
frame: &mut Frame,
jobs: &[harness_core::engine::JobRecord],
selected: usize,
area: Rect,
) {
let popup = centered_rect(70, 70, area);
frame.render_widget(Clear, popup);
let mut text: Vec<Line<'static>> =
vec![Line::from("Background jobs").style(Style::new().add_modifier(Modifier::BOLD))];
if jobs.is_empty() {
text.push(Line::default());
text.push(Line::from("No subtasks spawned yet.").style(Style::new().fg(Color::DarkGray)));
} else {
for (i, job) in jobs.iter().enumerate() {
let marker = if i == selected { "> " } else { " " };
let (icon, color) = job_state_style(job.state);
let header = format!(
"{marker}{icon} {} · {} · {}",
job.alias,
job.agent,
job_state_label(job.state),
);
let style = if i == selected {
Style::new().bg(Color::Blue).fg(Color::White)
} else {
Style::new().fg(color)
};
text.push(Line::from(Span::styled(header, style)));
if let Some(objective) = &job.objective {
text.push(
Line::from(format!(" {objective}")).style(Style::new().fg(Color::Gray)),
);
}
if !job.context_files.is_empty() {
let files: Vec<&str> = job
.context_files
.iter()
.take(8)
.map(|f| f.path.as_str())
.collect();
text.push(
Line::from(format!(" read: {}", files.join(", ")))
.style(Style::new().fg(Color::DarkGray)),
);
}
}
text.push(Line::default());
text.push(
Line::from("Enter: open subtask · Esc: close").style(Style::new().fg(Color::DarkGray)),
);
}
let block = Block::default().borders(Borders::ALL).title(" jobs ");
let paragraph = Paragraph::new(text).block(block).wrap(Wrap { trim: false });
frame.render_widget(paragraph, popup);
}
fn job_state_label(state: harness_core::engine::JobState) -> &'static str {
use harness_core::engine::JobState;
match state {
JobState::Running => "running",
JobState::Completed => "completed",
JobState::Error => "error",
JobState::Cancelled => "cancelled",
}
}
fn job_state_style(state: harness_core::engine::JobState) -> (&'static str, Color) {
use harness_core::engine::JobState;
match state {
JobState::Running => ("", Color::Yellow),
JobState::Completed => ("", Color::Green),
JobState::Error => ("", Color::Red),
JobState::Cancelled => ("", Color::DarkGray),
}
}
fn centered_rect(percent_x: u16, percent_y: u16, r: Rect) -> Rect {
let popup_layout = Layout::default()
.direction(Direction::Vertical)
@@ -487,6 +570,59 @@ mod tests {
insta::assert_snapshot!(buffer_to_string(terminal.backend()));
}
#[test]
fn snapshot_jobs_pane() {
use harness_core::engine::{ContextFile, JobRecord, JobState};
let backend = TestBackend::new(80, 24);
let mut terminal = Terminal::new(backend).unwrap();
let mut state = AppState::new(
"anthropic/claude-sonnet-4-5".to_string(),
"orchestrator".to_string(),
);
state.session_id = Some(SessionId("ses_test_001".to_string()));
state.jobs = vec![
JobRecord {
task_id: "t1".to_string(),
alias: "exp-1".to_string(),
parent_session: SessionId("ses_test_001".to_string()),
child_session: SessionId("ses_child_001".to_string()),
agent: "explorer".to_string(),
description: "map auth".to_string(),
objective: Some("map the auth flow".to_string()),
state: JobState::Completed,
reconciled: true,
result_summary: Some("done".to_string()),
context_files: vec![ContextFile {
path: "src/auth.rs".to_string(),
lines: 42,
}],
launched_at: 1,
updated_at: 2,
last_used_at: 2,
},
JobRecord {
task_id: "t2".to_string(),
alias: "fix-1".to_string(),
parent_session: SessionId("ses_test_001".to_string()),
child_session: SessionId("ses_child_002".to_string()),
agent: "fixer".to_string(),
description: "patch bug".to_string(),
objective: Some("fix the null deref".to_string()),
state: JobState::Running,
reconciled: false,
result_summary: None,
context_files: Vec::new(),
launched_at: 3,
updated_at: 3,
last_used_at: 3,
},
];
state.modal = ModalState::JobsPane { selected: 1 };
terminal.draw(|frame| render(frame, &mut state)).unwrap();
insta::assert_snapshot!(buffer_to_string(terminal.backend()));
}
#[test]
fn render_empty_state_produces_frame() {
let backend = TestBackend::new(80, 24);
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 300
assertion_line: 381
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -0,0 +1,29 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 621
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
┌ chat ────────────────────────────────────────────────────────────────────────┐
│ │
│ │
│ ┌ jobs ────────────────────────────────────────────────┐ │
│ │Background jobs │ │
│ │ ✓ exp-1 · explorer · completed │ │
│ │ map the auth flow │ │
│ │ read: src/auth.rs │ │
│ │> ⚙ fix-1 · fixer · running │ │
│ │ fix the null deref │ │
│ │ │ │
│ │Enter: open subtask · Esc: close │ │
│ │ │ │
│ │ │ │
│ │ │ │
│ │ │ │
│ │ │ │
│ │ │ │
└───────────└──────────────────────────────────────────────────────┘───────────┘
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 430
assertion_line: 511
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 487
assertion_line: 568
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 357
assertion_line: 438
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 407
assertion_line: 488
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 382
assertion_line: 463
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
@@ -1,6 +1,6 @@
---
source: crates/harness-tui/src/render.rs
assertion_line: 331
assertion_line: 412
expression: buffer_to_string(terminal.backend())
---
new session · orchestrator · anthropic/claude-sonnet-4-5
@@ -26,4 +26,4 @@ expression: buffer_to_string(terminal.backend())
┌ input ───────────────────────────────────────────────────────────────────────┐
│ │
└──────────────────────────────────────────────────────────────────────────────┘
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Esc: abort | Ctrl+C: quit
idle · 0 tok · $0.0000 · Ctrl+S: sessions | Ctrl+J: jobs | Esc: abort | Ctrl+C:
+50 -3
View File
@@ -1,4 +1,5 @@
use harness_app::EngineHandle;
use harness_core::engine::JobRecord;
use harness_core::event::{AppEvent, PermissionRequest, RunOutcome};
use harness_core::permission::PermissionReply;
use harness_core::types::{Message, Part, PartBody, PartId, Role, Session, SessionId, ToolState};
@@ -86,6 +87,10 @@ pub enum ModalState {
sessions: Vec<Session>,
selected: usize,
},
/// Background job board for the current session, with a cursor for drill-in.
JobsPane {
selected: usize,
},
}
/// All mutable UI state lives here.
@@ -109,6 +114,9 @@ pub struct AppState {
pub session_cost: f64,
/// Accumulated session tokens (input + output), for the status bar.
pub session_tokens: u64,
/// Background jobs spawned by the current session, newest activity last. Populated on
/// session load and kept live via `JobUpdated` events; surfaced in the jobs pane.
pub jobs: Vec<JobRecord>,
}
impl AppState {
@@ -131,6 +139,7 @@ impl AppState {
render_width: 0,
session_cost: 0.0,
session_tokens: 0,
jobs: Vec::new(),
}
}
@@ -156,6 +165,8 @@ impl AppState {
self.session_cost = session.cost;
self.session_tokens = session.usage.input + session.usage.output;
self.messages.clear();
// Jobs are reloaded for the newly-selected session by the caller.
self.jobs.clear();
let mut messages = messages;
messages.sort_by_key(|m| m.created_at);
@@ -258,14 +269,50 @@ impl AppState {
}
self.dirty = true;
}
AppEvent::JobUpdated { .. }
| AppEvent::AuthPrompt { .. }
| AppEvent::ServerNotice { .. } => {
AppEvent::JobUpdated { job } => {
if let Ok(record) = serde_json::from_value::<JobRecord>(job.0) {
if self.session_id.as_ref() == Some(&record.parent_session) {
self.upsert_job(record);
self.dirty = true;
}
}
}
AppEvent::AuthPrompt { .. } | AppEvent::ServerNotice { .. } => {
self.dirty = true;
}
}
}
/// Replaces this session's job list (e.g. after loading a session).
pub fn set_jobs(&mut self, jobs: Vec<JobRecord>) {
self.jobs = jobs;
self.dirty = true;
}
/// Inserts or replaces a job by `task_id`, preserving list order for stable rendering.
fn upsert_job(&mut self, record: JobRecord) {
if let Some(existing) = self.jobs.iter_mut().find(|j| j.task_id == record.task_id) {
*existing = record;
} else {
self.jobs.push(record);
}
}
/// Opens the jobs pane (cursor at the top).
pub fn open_jobs_pane(&mut self) {
self.modal = ModalState::JobsPane { selected: 0 };
self.dirty = true;
}
/// Child session of the currently-selected job in the jobs pane, for drill-in.
pub fn selected_job_child(&self) -> Option<SessionId> {
if let ModalState::JobsPane { selected } = &self.modal {
self.jobs.get(*selected).map(|j| j.child_session.clone())
} else {
None
}
}
pub fn scroll_up(&mut self, n: u16) {
self.scroll_offset = self.scroll_offset.saturating_sub(n);
self.dirty = true;