use std::path::Path;
use serde_json::{Value, json};
use crate::routes::verdict::ActionVerdict;
use crate::routes::{ACTION_TIMEOUT, MEMORY_SERVICE};
pub(crate) async fn call_tool(
socket: &Path,
tool: &str,
arguments: Value,
id: &str,
) -> Result<Value, ActionVerdict> {
let request = json!({
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": { "name": tool, "arguments": arguments },
});
let sent =
trusty_common::uds::send_framed_request::<_, trusty_common::uds::server::RpcResponse>(
socket,
&request,
ACTION_TIMEOUT,
)
.await;
let response = match sent {
Ok(r) => r,
Err(e) => {
return Err(ActionVerdict::Unreachable {
id: id.to_string(),
reason: format!("{MEMORY_SERVICE} did not answer {tool}: {e}"),
});
}
};
if let Some(error) = response.error {
return Err(ActionVerdict::Refused {
id: id.to_string(),
reason: format!(
"{MEMORY_SERVICE} refused {tool} (code {}): {}",
error.code, error.message
),
detail: json!({ "code": error.code }),
});
}
Ok(tool_payload(&response.result.unwrap_or(Value::Null)).unwrap_or(Value::Null))
}
fn tool_payload(result: &Value) -> Option<Value> {
let text = result
.get("content")?
.as_array()?
.first()?
.get("text")?
.as_str()?;
serde_json::from_str::<Value>(text).ok()
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
fn stub_memory_daemon(dir: &Path, reply: impl Into<String>) -> PathBuf {
let socket = dir.join("sockets").join("memory.sock");
let reply = reply.into();
let listener = trusty_common::uds::bind_hardened(&socket).expect("bind");
tokio::spawn(async move {
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
let Ok((mut conn, _)) = listener.accept().await else {
return;
};
let mut sink = Vec::new();
let _ = conn.read_to_end(&mut sink).await;
let _ = conn.write_all(reply.as_bytes()).await;
let _ = conn.write_all(b"\n").await;
let _ = conn.flush().await;
});
socket
}
#[tokio::test(flavor = "multi_thread")]
async fn call_tool_unwraps_the_tool_payload() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let reply = json!({
"jsonrpc": "2.0",
"id": 1,
"result": { "content": [{ "type": "text", "text": r#"{"palace":"scratch","orphans_removed":4}"# }] },
})
.to_string();
let socket = stub_memory_daemon(tmp.path(), reply);
let payload = call_tool(&socket, "palace_compact", json!({}), "scratch")
.await
.expect("the exchange succeeds");
assert_eq!(payload["palace"], json!("scratch"));
assert_eq!(payload["orphans_removed"], json!(4));
}
#[tokio::test(flavor = "multi_thread")]
async fn call_tool_reports_a_jsonrpc_error_as_a_refusal() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(
tmp.path(),
r#"{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"unknown palace 'scratch'"}}"#,
);
let verdict = call_tool(&socket, "palace_compact", json!({}), "scratch")
.await
.expect_err("a JSON-RPC error is not a success");
assert!(
matches!(&verdict, ActionVerdict::Refused { reason, .. } if reason.contains("unknown palace")),
"the refusal must carry the daemon's words: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn call_tool_reports_a_dead_socket_as_unreachable() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let verdict = call_tool(
&tmp.path().join("absent.sock"),
"palace_compact",
json!({}),
"scratch",
)
.await
.expect_err("a dead socket is not a success");
assert!(
matches!(verdict, ActionVerdict::Unreachable { .. }),
"a dead socket must read as unreachable: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn call_tool_returns_null_for_an_answer_with_no_payload() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(tmp.path(), r#"{"jsonrpc":"2.0","id":1,"result":{}}"#);
let payload = call_tool(&socket, "palace_compact", json!({}), "scratch")
.await
.expect("a payload-less answer is still an answer");
assert_eq!(payload, Value::Null);
}
}