hive-runtime: shared runtime crate with claude and acp backends
A `Runtime` trait (run / compact / archive) with two backends: - claude: a pass-through to hive_claude's InfiniteSession and SessionStore, so a claude turn is the same spawn, session handling and errors as before. - acp: a generic Agent Client Protocol client. It spawns the command, args and env from RuntimeSpec (HIVE_RUNTIME / HIVE_ACP_COMMAND / HIVE_ACP_ARGS / HIVE_ACP_ENV), refuses an agent whose mcpCapabilities.http is not true, passes the claude --mcp-config servers as ACP mcpServers, keeps one session id in a file (session/load after a restart, session/new otherwise), and maps session/update into claude stream-json events plus usage_update into Telemetry. Permission requests are answered by a caller-supplied policy on the ACP tool kind. compact returns Unsupported for now. The crate depends on no hyperhive binary crate, so the subagent daemon can move onto it without pulling in hive-agent. Refs #4391
This commit is contained in:
parent
d98d407bf1
commit
b0e26e7e44
10 changed files with 1583 additions and 0 deletions
316
hive-runtime/src/acp/rpc.rs
Normal file
316
hive-runtime/src/acp/rpc.rs
Normal file
|
|
@ -0,0 +1,316 @@
|
|||
//! JSON-RPC 2.0 over an ACP agent's stdio: newline-delimited messages, the
|
||||
//! agent's requests answered here, its notifications queued for the turn.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::path::Path;
|
||||
use std::process::Stdio;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::{Arc, Mutex, PoisonError};
|
||||
|
||||
use serde_json::{Value, json};
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
||||
use tokio::process::{Child, ChildStdin, Command};
|
||||
use tokio::sync::{mpsc, oneshot};
|
||||
|
||||
use super::{AcpError, PermissionPolicy};
|
||||
use crate::spec::AcpCommand;
|
||||
|
||||
/// Something the agent sent that is not a response to one of our requests.
|
||||
pub(super) enum Incoming {
|
||||
/// The `params` of a `session/update` notification.
|
||||
Update(Value),
|
||||
Stdout(String),
|
||||
Stderr(String),
|
||||
/// The agent closed its stdout: it has exited or is about to.
|
||||
Closed,
|
||||
}
|
||||
|
||||
type Reply = Result<Value, AcpError>;
|
||||
type Pending = Arc<Mutex<HashMap<u64, (&'static str, oneshot::Sender<Reply>)>>>;
|
||||
|
||||
pub(super) struct Connection {
|
||||
child: Child,
|
||||
stdin: Arc<tokio::sync::Mutex<ChildStdin>>,
|
||||
next_id: AtomicU64,
|
||||
pending: Pending,
|
||||
pub(super) incoming: mpsc::UnboundedReceiver<Incoming>,
|
||||
}
|
||||
|
||||
impl Connection {
|
||||
/// Spawn the agent in `cwd` and start reading its output.
|
||||
pub(super) fn spawn(
|
||||
command: &AcpCommand,
|
||||
cwd: &Path,
|
||||
permit: PermissionPolicy,
|
||||
) -> Result<Self, AcpError> {
|
||||
let mut child = Command::new(&command.command)
|
||||
.args(&command.args)
|
||||
.envs(&command.env)
|
||||
.current_dir(cwd)
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.kill_on_drop(true)
|
||||
.spawn()
|
||||
.map_err(|source| AcpError::Spawn {
|
||||
program: command.command.clone(),
|
||||
source,
|
||||
})?;
|
||||
let (Some(stdin), Some(stdout), Some(stderr)) =
|
||||
(child.stdin.take(), child.stdout.take(), child.stderr.take())
|
||||
else {
|
||||
unreachable!("all three stdio handles were requested as pipes");
|
||||
};
|
||||
let stdin = Arc::new(tokio::sync::Mutex::new(stdin));
|
||||
let pending: Pending = Arc::default();
|
||||
let (tx, incoming) = mpsc::unbounded_channel();
|
||||
|
||||
let stderr_tx = tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut lines = BufReader::new(stderr).lines();
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
if stderr_tx.send(Incoming::Stderr(line)).is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let reader = Reader {
|
||||
stdin: stdin.clone(),
|
||||
pending: pending.clone(),
|
||||
tx,
|
||||
permit,
|
||||
};
|
||||
tokio::spawn(async move {
|
||||
let mut lines = BufReader::new(stdout).lines();
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
reader.dispatch(&line).await;
|
||||
}
|
||||
// Fail every request still waiting: no reply can come now.
|
||||
reader
|
||||
.pending
|
||||
.lock()
|
||||
.unwrap_or_else(PoisonError::into_inner)
|
||||
.clear();
|
||||
let _ = reader.tx.send(Incoming::Closed);
|
||||
});
|
||||
|
||||
Ok(Self {
|
||||
child,
|
||||
stdin,
|
||||
next_id: AtomicU64::new(1),
|
||||
pending,
|
||||
incoming,
|
||||
})
|
||||
}
|
||||
|
||||
/// Whether the agent process has exited.
|
||||
pub(super) fn exited(&mut self) -> bool {
|
||||
!matches!(self.child.try_wait(), Ok(None))
|
||||
}
|
||||
|
||||
/// Send a request and return the receiver its reply arrives on.
|
||||
pub(super) async fn send(
|
||||
&self,
|
||||
method: &'static str,
|
||||
params: Value,
|
||||
) -> Result<oneshot::Receiver<Reply>, AcpError> {
|
||||
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
|
||||
let (tx, rx) = oneshot::channel();
|
||||
self.pending
|
||||
.lock()
|
||||
.unwrap_or_else(PoisonError::into_inner)
|
||||
.insert(id, (method, tx));
|
||||
let message = json!({ "jsonrpc": "2.0", "id": id, "method": method, "params": params });
|
||||
if let Err(e) = write(&self.stdin, &message).await {
|
||||
self.pending
|
||||
.lock()
|
||||
.unwrap_or_else(PoisonError::into_inner)
|
||||
.remove(&id);
|
||||
return Err(e);
|
||||
}
|
||||
Ok(rx)
|
||||
}
|
||||
|
||||
/// Send a request and wait for its reply, queueing anything else the agent
|
||||
/// sends meanwhile on [`Self::incoming`].
|
||||
pub(super) async fn request(
|
||||
&self,
|
||||
method: &'static str,
|
||||
params: Value,
|
||||
) -> Result<Value, AcpError> {
|
||||
let rx = self.send(method, params).await?;
|
||||
rx.await.unwrap_or(Err(AcpError::Closed))
|
||||
}
|
||||
}
|
||||
|
||||
/// The stdout side: routes responses to their waiting request, answers the
|
||||
/// agent's own requests, and queues its notifications.
|
||||
struct Reader {
|
||||
stdin: Arc<tokio::sync::Mutex<ChildStdin>>,
|
||||
pending: Pending,
|
||||
tx: mpsc::UnboundedSender<Incoming>,
|
||||
permit: PermissionPolicy,
|
||||
}
|
||||
|
||||
impl Reader {
|
||||
async fn dispatch(&self, line: &str) {
|
||||
let Ok(message) = serde_json::from_str::<Value>(line) else {
|
||||
if !line.trim().is_empty() {
|
||||
let _ = self.tx.send(Incoming::Stdout(line.to_owned()));
|
||||
}
|
||||
return;
|
||||
};
|
||||
let method = message.get("method").and_then(Value::as_str);
|
||||
let id = message.get("id");
|
||||
match (method, id) {
|
||||
(None, Some(id)) => self.resolve(id, &message),
|
||||
(Some(method), Some(id)) => {
|
||||
let reply = match method {
|
||||
"session/request_permission" => json!({
|
||||
"jsonrpc": "2.0", "id": id,
|
||||
"result": permission_outcome(&message["params"], &self.permit),
|
||||
}),
|
||||
_ => json!({
|
||||
"jsonrpc": "2.0", "id": id,
|
||||
"error": { "code": -32601, "message": format!("method not found: {method}") },
|
||||
}),
|
||||
};
|
||||
if let Err(e) = write(&self.stdin, &reply).await {
|
||||
tracing::warn!(error = %e, method, "failed to answer ACP agent request");
|
||||
}
|
||||
}
|
||||
(Some("session/update"), None) => {
|
||||
let params = message.get("params").cloned().unwrap_or(Value::Null);
|
||||
let _ = self.tx.send(Incoming::Update(params));
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
fn resolve(&self, id: &Value, message: &Value) {
|
||||
let Some(id) = id.as_u64() else { return };
|
||||
let Some((method, tx)) = self
|
||||
.pending
|
||||
.lock()
|
||||
.unwrap_or_else(PoisonError::into_inner)
|
||||
.remove(&id)
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let reply = match message.get("error") {
|
||||
Some(error) => Err(AcpError::Rpc {
|
||||
method,
|
||||
code: error.get("code").and_then(Value::as_i64).unwrap_or(0),
|
||||
message: error
|
||||
.get("message")
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or("")
|
||||
.to_owned(),
|
||||
data: error.get("data").map(Value::to_string).unwrap_or_default(),
|
||||
}),
|
||||
None => Ok(message.get("result").cloned().unwrap_or(Value::Null)),
|
||||
};
|
||||
let _ = tx.send(reply);
|
||||
}
|
||||
}
|
||||
|
||||
async fn write(stdin: &tokio::sync::Mutex<ChildStdin>, message: &Value) -> Result<(), AcpError> {
|
||||
let mut line = message.to_string();
|
||||
line.push('\n');
|
||||
let mut stdin = stdin.lock().await;
|
||||
stdin.write_all(line.as_bytes()).await?;
|
||||
stdin.flush().await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// The `session/request_permission` result: the agent's first option of the
|
||||
/// matching kind, allow if `permit` accepts the tool call's kind and reject
|
||||
/// otherwise. With no such option the request is answered `cancelled`.
|
||||
pub(super) fn permission_outcome(params: &Value, permit: &PermissionPolicy) -> Value {
|
||||
let kind = params
|
||||
.get("toolCall")
|
||||
.and_then(|t| t.get("kind"))
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or("other");
|
||||
let wanted = if permit(kind) {
|
||||
["allow_once", "allow_always"]
|
||||
} else {
|
||||
["reject_once", "reject_always"]
|
||||
};
|
||||
let options = params
|
||||
.get("options")
|
||||
.and_then(Value::as_array)
|
||||
.map_or(&[][..], Vec::as_slice);
|
||||
let chosen = wanted.iter().find_map(|want| {
|
||||
options
|
||||
.iter()
|
||||
.find(|o| o.get("kind").and_then(Value::as_str) == Some(want))
|
||||
.and_then(|o| o.get("optionId"))
|
||||
});
|
||||
match chosen {
|
||||
Some(option) => json!({ "outcome": { "outcome": "selected", "optionId": option } }),
|
||||
None => json!({ "outcome": { "outcome": "cancelled" } }),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::permission_outcome;
|
||||
use crate::acp::PermissionPolicy;
|
||||
use serde_json::json;
|
||||
use std::sync::Arc;
|
||||
|
||||
fn request(kind: &str) -> serde_json::Value {
|
||||
json!({
|
||||
"sessionId": "s",
|
||||
"toolCall": { "toolCallId": "t", "kind": kind },
|
||||
"options": [
|
||||
{ "optionId": "once", "kind": "allow_once", "name": "Allow once" },
|
||||
{ "optionId": "always", "kind": "allow_always", "name": "Always allow" },
|
||||
{ "optionId": "reject", "kind": "reject_once", "name": "Reject" },
|
||||
],
|
||||
})
|
||||
}
|
||||
|
||||
fn no_execute() -> PermissionPolicy {
|
||||
Arc::new(|kind: &str| kind != "execute")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_permitted_kind_is_allowed_once() {
|
||||
assert_eq!(
|
||||
permission_outcome(&request("edit"), &no_execute()),
|
||||
json!({ "outcome": { "outcome": "selected", "optionId": "once" } })
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_refused_kind_is_rejected() {
|
||||
assert_eq!(
|
||||
permission_outcome(&request("execute"), &no_execute()),
|
||||
json!({ "outcome": { "outcome": "selected", "optionId": "reject" } })
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_missing_kind_is_judged_as_other() {
|
||||
let policy: PermissionPolicy = Arc::new(|kind: &str| kind == "other");
|
||||
let mut req = request("x");
|
||||
req["toolCall"].as_object_mut().unwrap().remove("kind");
|
||||
assert_eq!(
|
||||
permission_outcome(&req, &policy)["outcome"]["optionId"],
|
||||
"once"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn no_matching_option_cancels() {
|
||||
let req = json!({ "toolCall": { "kind": "execute" },
|
||||
"options": [{ "optionId": "once", "kind": "allow_once" }] });
|
||||
assert_eq!(
|
||||
permission_outcome(&req, &no_execute()),
|
||||
json!({ "outcome": { "outcome": "cancelled" } })
|
||||
);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue