Watch
0
0
Fork
You've already forked hyperhive
0
hyperhive/hive-runtime/src/acp/stream.rs
atlas f9c6a56ab9 hive-runtime: report an empty ACP end_turn as a turn error (#4819)
opencode answers `end_turn` even when its provider rejected the request
with a non-retryable error (401, 400), and forwards nothing over ACP, so
the turn looked like an empty success. A turn that ends with `end_turn`,
no event, no `usage_update` and no `usage` in the prompt response now
fails with `AcpError::EmptyEndTurn`.

A turn that really produced nothing and reported no usage is reported
the same way; that false positive is accepted.

The test agent's plain `end_turn` replies now carry a response `usage`,
so its ordinary turns stay successes; a response `usage` feeds only cost
telemetry, not the compaction watermark. A new `blank` mode keeps the
empty reply for the error case.
2026-09-30 09:11:36 +02:00

554 lines
20 KiB
Rust

//! Pure translation between ACP and the claude-shaped values the rest of the
//! crate's callers read: `session/update` → `stream-json` events, a turn's
//! usage → [`Telemetry`], and a claude `--mcp-config` → ACP `mcpServers`.
use std::collections::HashMap;
use hive_claude::{Telemetry, TokenUsage};
use serde_json::{Map, Value, json};
/// Turns one prompt's `session/update` notifications into claude
/// `stream-json` events.
///
/// Text and thought chunks are buffered and emitted as one block when the
/// update kind changes or the turn ends, so a consumer sees a paragraph where
/// claude would have sent one, not a row per token. A tool call is emitted as
/// a `tool_use` once its input is known, and its `tool_result` when it
/// completes or fails.
pub(super) struct StreamMapper {
servers: Vec<String>,
text: String,
thought: String,
tools: HashMap<String, ToolCall>,
usage: Option<(u64, u64)>,
emitted: bool,
}
struct ToolCall {
name: String,
input: Value,
announced: bool,
}
impl StreamMapper {
/// `servers` are the MCP server names handed to the agent, used to give
/// their tools claude's `mcp__<server>__<tool>` names.
pub(super) fn new(servers: Vec<String>) -> Self {
Self {
servers,
text: String::new(),
thought: String::new(),
tools: HashMap::new(),
usage: None,
emitted: false,
}
}
/// Fold one `update` (the `update` field of a `session/update`), returning
/// the events it completes.
pub(super) fn push(&mut self, update: &Value) -> Vec<Value> {
let kind = update.get("sessionUpdate").and_then(Value::as_str);
let mut out = Vec::new();
match kind {
Some("agent_message_chunk") => {
self.flush_thought(&mut out);
self.text.push_str(chunk_text(update));
}
Some("agent_thought_chunk") => {
self.flush_text(&mut out);
self.thought.push_str(chunk_text(update));
}
Some("tool_call" | "tool_call_update") => {
self.flush_text(&mut out);
self.flush_thought(&mut out);
self.tool_update(update, &mut out);
}
Some("usage_update") => {
let field = |k: &str| update.get(k).and_then(Value::as_u64);
if let (Some(used), Some(size)) = (field("used"), field("size")) {
self.usage = Some((used, size));
}
}
_ => {}
}
self.emitted |= !out.is_empty();
out
}
/// Emit whatever is still buffered at the end of the turn.
pub(super) fn finish(&mut self) -> Vec<Value> {
let mut out = Vec::new();
self.flush_text(&mut out);
self.flush_thought(&mut out);
let mut pending: Vec<_> = self.tools.drain().filter(|(_, t)| !t.announced).collect();
pending.sort_by(|a, b| a.0.cmp(&b.0));
for (id, tool) in pending {
out.push(tool_use(&id, &tool.name, &tool.input));
}
self.emitted |= !out.is_empty();
out
}
/// Whether the turn has emitted any event.
pub(super) fn has_content(&self) -> bool {
self.emitted
}
/// Whether the turn has had a `usage_update`.
pub(super) fn has_usage_update(&self) -> bool {
self.usage.is_some()
}
/// The turn's telemetry: context from the last `usage_update`, cost from
/// the `session/prompt` response's `usage` when the agent sends one.
pub(super) fn telemetry(&self, response: &Value, model: Option<&str>) -> Telemetry {
let mut telemetry = Telemetry::default();
if let Some((used, size)) = self.usage {
telemetry.context.input_tokens = used;
telemetry.context_window = Some(size);
}
if let Some(usage) = response.get("usage") {
let field = |k: &str| usage.get(k).and_then(Value::as_u64).unwrap_or(0);
let mut cost = TokenUsage::default();
cost.input_tokens = field("inputTokens");
cost.output_tokens = field("outputTokens");
cost.cache_read_input_tokens = field("cachedReadTokens");
cost.cache_creation_input_tokens = field("cachedWriteTokens");
telemetry.cost = cost;
}
telemetry.model = model.map(str::to_owned);
telemetry
}
fn tool_update(&mut self, update: &Value, out: &mut Vec<Value>) {
let Some(id) = update.get("toolCallId").and_then(Value::as_str) else {
return;
};
let tool = self.tools.entry(id.to_owned()).or_insert_with(|| ToolCall {
// The first title is the one kept: later updates may retitle the
// call with a description of what it did.
name: canonical_tool_name(
update
.get("title")
.and_then(Value::as_str)
.unwrap_or("tool"),
&self.servers,
),
input: json!({}),
announced: false,
});
if !tool.announced
&& let Some(input) = update.get("rawInput").filter(|v| !is_empty(v))
{
tool.input = input.clone();
}
let status = update.get("status").and_then(Value::as_str);
let done = matches!(status, Some("completed" | "failed"));
if !tool.announced && (done || !is_empty(&tool.input)) {
tool.announced = true;
out.push(tool_use(id, &tool.name, &tool.input));
}
if done {
self.tools.remove(id);
out.push(tool_result(
id,
&tool_output(update),
status == Some("failed"),
));
}
}
fn flush_text(&mut self, out: &mut Vec<Value>) {
if !self.text.is_empty() {
let text = std::mem::take(&mut self.text);
out.push(assistant(json!({ "type": "text", "text": text })));
}
}
fn flush_thought(&mut self, out: &mut Vec<Value>) {
if !self.thought.is_empty() {
let thinking = std::mem::take(&mut self.thought);
out.push(assistant(
json!({ "type": "thinking", "thinking": thinking }),
));
}
}
}
/// The model an agent reports for a session: the current value of its
/// `model`-category config option, else the `models.currentModelId` some
/// agents send instead.
pub(super) fn session_model(response: &Value) -> Option<String> {
let from_options = response
.get("configOptions")
.and_then(Value::as_array)
.and_then(|options| {
options
.iter()
.find(|o| o.get("category").and_then(Value::as_str) == Some("model"))
})
.and_then(|o| o.get("currentValue"))
.and_then(Value::as_str);
from_options
.or_else(|| {
response
.get("models")
.and_then(|m| m.get("currentModelId"))
.and_then(Value::as_str)
})
.map(str::to_owned)
}
/// The command names an `available_commands_update` advertises, or `None`
/// for any other `update`. Each replaces the session's previous list.
pub(super) fn advertised_commands(update: &Value) -> Option<Vec<String>> {
if update.get("sessionUpdate").and_then(Value::as_str) != Some("available_commands_update") {
return None;
}
let commands = update.get("availableCommands").and_then(Value::as_array);
Some(
commands
.into_iter()
.flatten()
.filter_map(|c| c.get("name").and_then(Value::as_str))
.map(str::to_owned)
.collect(),
)
}
/// Convert a claude `--mcp-config` document into ACP's `mcpServers` list.
/// Returns the list and the server names, in the same order.
pub(super) fn mcp_servers(config: &Value) -> (Vec<Value>, Vec<String>) {
let Some(servers) = config.get("mcpServers").and_then(Value::as_object) else {
return (Vec::new(), Vec::new());
};
let mut list = Vec::new();
let mut names = Vec::new();
for (name, entry) in servers {
let kind = entry.get("type").and_then(Value::as_str);
let server = match kind {
Some(kind @ ("http" | "sse")) => json!({
"type": kind,
"name": name,
"url": entry.get("url").cloned().unwrap_or(Value::Null),
"headers": name_values(entry.get("headers")),
}),
_ => json!({
"name": name,
"command": entry.get("command").cloned().unwrap_or(Value::Null),
"args": entry.get("args").cloned().unwrap_or_else(|| json!([])),
"env": name_values(entry.get("env")),
}),
};
list.push(server);
names.push(name.clone());
}
(list, names)
}
/// The claude name for a tool an ACP agent reports as `<server>_<tool>` or
/// `<server>__<tool>`: `mcp__<server>__<tool>`, the form tool-call counters
/// and the terminal renderer key on. Anything else is returned unchanged.
pub(super) fn canonical_tool_name(name: &str, servers: &[String]) -> String {
if name.starts_with("mcp__") {
return name.to_owned();
}
split_mcp_name(name, servers).map_or_else(
|| name.to_owned(),
|(server, tool)| format!("mcp__{server}__{tool}"),
)
}
/// Which of `servers` a tool named `<server>_<tool>`, `<server>__<tool>` or
/// `mcp__<server>__<tool>` belongs to, if any.
pub(super) fn mcp_server_of<'a>(name: &str, servers: &'a [String]) -> Option<&'a str> {
split_mcp_name(name, servers).map(|(server, _)| server)
}
/// `(server, tool)` for an MCP tool name, matching the longest server name so
/// one that prefixes another's never claims the other's tools.
fn split_mcp_name<'a, 'n>(name: &'n str, servers: &'a [String]) -> Option<(&'a str, &'n str)> {
let bare = name.strip_prefix("mcp__").unwrap_or(name);
servers
.iter()
.filter_map(|server| {
let rest = bare.strip_prefix(server.as_str())?.strip_prefix('_')?;
let tool = rest.strip_prefix('_').unwrap_or(rest);
(!tool.is_empty()).then_some((server.as_str(), tool))
})
.max_by_key(|(server, _)| server.len())
}
fn chunk_text(update: &Value) -> &str {
update
.get("content")
.filter(|c| c.get("type").and_then(Value::as_str) == Some("text"))
.and_then(|c| c.get("text"))
.and_then(Value::as_str)
.unwrap_or("")
}
/// A tool call's output: its text content blocks, else its raw output.
fn tool_output(update: &Value) -> String {
let text: Vec<&str> = update
.get("content")
.and_then(Value::as_array)
.into_iter()
.flatten()
.filter_map(|c| c.get("content").and_then(|c| c.get("text")))
.filter_map(Value::as_str)
.collect();
if !text.is_empty() {
return text.join("\n");
}
match update.get("rawOutput") {
None | Some(Value::Null) => String::new(),
Some(Value::String(s)) => s.clone(),
Some(other) => other.to_string(),
}
}
fn is_empty(v: &Value) -> bool {
match v {
Value::Null => true,
Value::Object(m) => m.is_empty(),
_ => false,
}
}
/// `{"K": "V"}` → `[{"name": "K", "value": "V"}]`, ACP's form for env and
/// headers.
fn name_values(map: Option<&Value>) -> Value {
let pairs = map
.and_then(Value::as_object)
.map(Map::iter)
.into_iter()
.flatten()
.map(|(k, v)| json!({ "name": k, "value": v.as_str().unwrap_or_default() }))
.collect();
Value::Array(pairs)
}
fn assistant(block: Value) -> Value {
let mut event = json!({ "type": "assistant", "message": { "role": "assistant" } });
event["message"]["content"] = Value::Array(vec![block]);
event
}
fn tool_use(id: &str, name: &str, input: &Value) -> Value {
assistant(json!({ "type": "tool_use", "id": id, "name": name, "input": input }))
}
fn tool_result(id: &str, output: &str, is_error: bool) -> Value {
json!({
"type": "user",
"message": { "role": "user", "content": [{
"type": "tool_result", "tool_use_id": id, "content": output, "is_error": is_error,
}]},
})
}
#[cfg(test)]
mod tests {
use super::{StreamMapper, canonical_tool_name, mcp_servers, session_model};
use serde_json::{Value, json};
fn mapper() -> StreamMapper {
StreamMapper::new(vec!["hyperhive".into(), "bash".into()])
}
fn feed(m: &mut StreamMapper, updates: &[Value]) -> Vec<Value> {
let mut out: Vec<Value> = updates.iter().flat_map(|u| m.push(u)).collect();
out.extend(m.finish());
out
}
fn text_chunk(kind: &str, text: &str) -> Value {
json!({ "sessionUpdate": kind, "content": { "type": "text", "text": text } })
}
#[test]
fn message_chunks_become_one_assistant_text_block() {
let out = feed(
&mut mapper(),
&[
text_chunk("agent_message_chunk", "Hel"),
text_chunk("agent_message_chunk", "lo"),
],
);
assert_eq!(
out,
vec![
json!({ "type": "assistant", "message": { "role": "assistant",
"content": [{ "type": "text", "text": "Hello" }] } })
]
);
}
#[test]
fn a_thought_then_text_flushes_in_order() {
let out = feed(
&mut mapper(),
&[
text_chunk("agent_thought_chunk", "hmm"),
text_chunk("agent_message_chunk", "done"),
],
);
let blocks: Vec<&Value> = out.iter().map(|e| &e["message"]["content"][0]).collect();
assert_eq!(blocks[0], &json!({ "type": "thinking", "thinking": "hmm" }));
assert_eq!(blocks[1], &json!({ "type": "text", "text": "done" }));
}
/// The sequence an agent sends for one MCP tool call: announced with no
/// input, updated with the input while running, then completed with its
/// output and a new, descriptive title.
#[test]
fn a_tool_call_becomes_tool_use_then_tool_result() {
let out = feed(
&mut mapper(),
&[
text_chunk("agent_message_chunk", "sending"),
json!({ "sessionUpdate": "tool_call", "toolCallId": "c1",
"title": "hyperhive_send", "kind": "other", "status": "pending",
"rawInput": {} }),
json!({ "sessionUpdate": "tool_call_update", "toolCallId": "c1",
"status": "in_progress", "rawInput": { "to": "atlas", "body": "hi" } }),
json!({ "sessionUpdate": "tool_call_update", "toolCallId": "c1",
"status": "completed", "title": "Sent a message",
"content": [{ "type": "content", "content": { "type": "text", "text": "ok" } }] }),
],
);
assert_eq!(out.len(), 3, "{out:#?}");
assert_eq!(out[0]["message"]["content"][0]["text"], "sending");
assert_eq!(
out[1],
json!({ "type": "assistant", "message": { "role": "assistant", "content": [{
"type": "tool_use", "id": "c1", "name": "mcp__hyperhive__send",
"input": { "to": "atlas", "body": "hi" } }] } })
);
assert_eq!(
out[2],
json!({ "type": "user", "message": { "role": "user", "content": [{
"type": "tool_result", "tool_use_id": "c1", "content": "ok", "is_error": false }] } })
);
}
#[test]
fn a_failed_tool_call_is_an_error_result_carrying_raw_output() {
let out = feed(
&mut mapper(),
&[
json!({ "sessionUpdate": "tool_call_update", "toolCallId": "c2",
"title": "read", "status": "failed",
"rawOutput": { "error": "no such file" } }),
],
);
assert_eq!(out[0]["message"]["content"][0]["name"], "read");
assert_eq!(out[1]["message"]["content"][0]["is_error"], true);
assert_eq!(
out[1]["message"]["content"][0]["content"],
r#"{"error":"no such file"}"#
);
}
#[test]
fn a_tool_call_still_running_at_turn_end_is_still_counted() {
let out = feed(
&mut mapper(),
&[json!({ "sessionUpdate": "tool_call", "toolCallId": "c3",
"title": "bash_run", "status": "pending" })],
);
assert_eq!(out.len(), 1);
assert_eq!(out[0]["message"]["content"][0]["name"], "mcp__bash__run");
}
#[test]
fn usage_update_and_prompt_usage_become_telemetry() {
let mut m = mapper();
feed(
&mut m,
&[
json!({ "sessionUpdate": "usage_update", "used": 1000, "size": 262_144 }),
json!({ "sessionUpdate": "usage_update", "used": 5000, "size": 262_144,
"cost": { "amount": 0.0, "currency": "USD" } }),
],
);
let response = json!({ "stopReason": "end_turn", "usage": {
"inputTokens": 4000, "outputTokens": 300, "cachedReadTokens": 1000 } });
let t = m.telemetry(&response, Some("provider/model"));
assert_eq!(t.context.context_tokens(), 5000);
assert_eq!(t.context_window, Some(262_144));
assert_eq!(t.cost.input_tokens, 4000);
assert_eq!(t.cost.output_tokens, 300);
assert_eq!(t.cost.cache_read_input_tokens, 1000);
assert_eq!(t.model.as_deref(), Some("provider/model"));
}
#[test]
fn updates_with_nothing_to_show_emit_nothing() {
let out = feed(
&mut mapper(),
&[
json!({ "sessionUpdate": "available_commands_update", "availableCommands": [] }),
json!({ "sessionUpdate": "plan", "entries": [] }),
json!({ "sessionUpdate": "something_new" }),
],
);
assert!(out.is_empty(), "{out:#?}");
}
#[test]
fn tool_names_are_canonicalised_only_for_known_servers() {
let servers = vec!["hyperhive".to_owned(), "bash".to_owned()];
assert_eq!(
canonical_tool_name("hyperhive_get_loose_ends", &servers),
"mcp__hyperhive__get_loose_ends"
);
assert_eq!(
canonical_tool_name("hyperhive__send", &servers),
"mcp__hyperhive__send"
);
assert_eq!(
canonical_tool_name("mcp__bash__run", &servers),
"mcp__bash__run"
);
assert_eq!(canonical_tool_name("edit", &servers), "edit");
assert_eq!(canonical_tool_name("bashful", &servers), "bashful");
assert_eq!(canonical_tool_name("bash_", &servers), "bash_");
}
#[test]
fn claude_mcp_config_becomes_acp_mcp_servers() {
let config = json!({ "mcpServers": {
"hyperhive": { "type": "http", "url": "http://127.0.0.1:8790/mcp" },
"scraper": { "command": "/bin/scrape", "args": ["--x"],
"env": { "HYPERHIVE_STATE_DIR": "/state" } },
}});
let (servers, names) = mcp_servers(&config);
assert_eq!(names, vec!["hyperhive", "scraper"]);
assert_eq!(
servers,
vec![
json!({ "type": "http", "name": "hyperhive",
"url": "http://127.0.0.1:8790/mcp", "headers": [] }),
json!({ "name": "scraper", "command": "/bin/scrape", "args": ["--x"],
"env": [{ "name": "HYPERHIVE_STATE_DIR", "value": "/state" }] }),
]
);
}
#[test]
fn session_model_prefers_the_model_config_option() {
let resp = json!({
"configOptions": [
{ "id": "mode", "category": "mode", "currentValue": "build" },
{ "id": "model", "category": "model", "currentValue": "p/m-1" },
],
"models": { "currentModelId": "p/m-2" },
});
assert_eq!(session_model(&resp).as_deref(), Some("p/m-1"));
let resp = json!({ "models": { "currentModelId": "p/m-2" } });
assert_eq!(session_model(&resp).as_deref(), Some("p/m-2"));
assert_eq!(session_model(&json!({})), None);
}
}