use std::sync::Arc;
use async_trait::async_trait;
use nexo_extensions::runtime::admin_router::AdminRouter;
use serde_json::Value;
use super::dispatcher::{AdminRpcDispatcher, AdminRpcResult};
#[async_trait]
pub trait AdminOutboundWriter: Send + Sync + std::fmt::Debug {
async fn send(&self, line: String);
}
#[derive(Clone)]
pub struct DispatcherAdminRouter {
dispatcher: Arc<AdminRpcDispatcher>,
writer: Arc<dyn AdminOutboundWriter>,
}
impl std::fmt::Debug for DispatcherAdminRouter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DispatcherAdminRouter")
.field("dispatcher", &self.dispatcher)
.finish_non_exhaustive()
}
}
impl DispatcherAdminRouter {
pub fn new(dispatcher: Arc<AdminRpcDispatcher>, writer: Arc<dyn AdminOutboundWriter>) -> Self {
Self { dispatcher, writer }
}
}
#[async_trait]
impl AdminRouter for DispatcherAdminRouter {
async fn handle_frame(&self, extension_id: &str, line: String) {
let frame: Value = match serde_json::from_str(&line) {
Ok(v) => v,
Err(e) => {
tracing::warn!(
ext = %extension_id,
error = %e,
"admin frame is not valid json; dropping",
);
return;
}
};
let id = frame.get("id").cloned().unwrap_or(Value::Null);
let method = match frame.get("method").and_then(Value::as_str) {
Some(m) => m.to_string(),
None => {
tracing::debug!(
ext = %extension_id,
"admin frame has no method (response shape); dropping",
);
return;
}
};
let params = frame.get("params").cloned().unwrap_or(Value::Null);
let result = self
.dispatcher
.dispatch(extension_id, &method, params)
.await;
let response = frame_response(id, result);
let line = match serde_json::to_string(&response) {
Ok(s) => s,
Err(e) => {
tracing::warn!(
ext = %extension_id,
method = %method,
error = %e,
"failed to serialize admin response; dropping",
);
return;
}
};
self.writer.send(line).await;
}
}
fn frame_response(id: Value, result: AdminRpcResult) -> Value {
match (result.result, result.error) {
(Some(value), _) => serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"result": value,
}),
(None, Some(err)) => {
let mut error = serde_json::json!({
"code": err.code(),
"message": format!("{err}"),
});
if let Some(data) = err.data() {
error["data"] = data;
}
serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"error": error,
})
}
(None, None) => serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"result": Value::Null,
}),
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::{HashMap, HashSet};
use std::sync::Mutex;
use crate::agent::admin_rpc::CapabilitySet;
#[derive(Debug, Default)]
struct CapturingWriter {
sent: Mutex<Vec<String>>,
}
#[async_trait]
impl AdminOutboundWriter for CapturingWriter {
async fn send(&self, line: String) {
self.sent.lock().unwrap().push(line);
}
}
fn build_dispatcher_granting(microapp_id: &str, caps: &[&str]) -> AdminRpcDispatcher {
let mut grants = HashMap::new();
grants.insert(
microapp_id.to_string(),
caps.iter().map(|s| s.to_string()).collect::<HashSet<_>>(),
);
AdminRpcDispatcher::new().with_capabilities(CapabilitySet::from_grants(grants))
}
#[tokio::test]
async fn handle_frame_dispatches_and_writes_response() {
let dispatcher = Arc::new(build_dispatcher_granting("agent-creator", &["_echo"]));
let writer = Arc::new(CapturingWriter::default());
let router = DispatcherAdminRouter::new(dispatcher, writer.clone());
let frame = serde_json::json!({
"jsonrpc": "2.0",
"id": "app:abc",
"method": "nexo/admin/echo",
"params": { "x": 7 },
})
.to_string();
router.handle_frame("agent-creator", frame).await;
let sent = writer.sent.lock().unwrap();
assert_eq!(sent.len(), 1);
let response: Value = serde_json::from_str(&sent[0]).unwrap();
assert_eq!(response["jsonrpc"], "2.0");
assert_eq!(response["id"], "app:abc");
assert_eq!(response["result"]["echoed"]["x"], 7);
assert_eq!(response["result"]["microapp_id"], "agent-creator");
}
#[tokio::test]
async fn handle_frame_emits_error_envelope_on_capability_denial() {
let dispatcher = Arc::new(AdminRpcDispatcher::new());
let writer = Arc::new(CapturingWriter::default());
let router = DispatcherAdminRouter::new(dispatcher, writer.clone());
let frame = serde_json::json!({
"jsonrpc": "2.0",
"id": "app:def",
"method": "nexo/admin/echo",
"params": Value::Null,
})
.to_string();
router.handle_frame("agent-creator", frame).await;
let sent = writer.sent.lock().unwrap();
assert_eq!(sent.len(), 1);
let response: Value = serde_json::from_str(&sent[0]).unwrap();
assert_eq!(response["jsonrpc"], "2.0");
assert_eq!(response["id"], "app:def");
assert_eq!(response["error"]["code"], -32004);
assert_eq!(response["error"]["data"]["capability"], "_echo");
}
#[tokio::test]
async fn handle_frame_drops_invalid_json_silently() {
let dispatcher = Arc::new(AdminRpcDispatcher::new());
let writer = Arc::new(CapturingWriter::default());
let router = DispatcherAdminRouter::new(dispatcher, writer.clone());
router
.handle_frame("agent-creator", "{not valid json".to_string())
.await;
assert!(
writer.sent.lock().unwrap().is_empty(),
"no response should be sent for unparseable input"
);
}
#[tokio::test]
async fn handle_frame_drops_response_shaped_frames() {
let dispatcher = Arc::new(AdminRpcDispatcher::new());
let writer = Arc::new(CapturingWriter::default());
let router = DispatcherAdminRouter::new(dispatcher, writer.clone());
let frame = serde_json::json!({
"jsonrpc": "2.0",
"id": "app:ghi",
"result": { "ok": true },
})
.to_string();
router.handle_frame("agent-creator", frame).await;
assert!(writer.sent.lock().unwrap().is_empty());
}
}