#![allow(clippy::unwrap_used, clippy::pedantic, missing_docs)]
use std::{
collections::HashMap,
future::Future,
pin::Pin,
sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
},
};
use axum::body::Body;
use http::{Request, StatusCode, header};
use http_body_util::BodyExt as _;
use polyc_a2a::{
AppState, InMemoryTaskStore, TurnOutcome, TurnRequest, TurnRunner,
UnconfiguredApprovalResponder, card::signed_card, router,
};
use polyc_crypto::Signer;
use polyc_runtime::admission::AdmissionGate;
use serde_json::{Value, json};
use tower::ServiceExt as _;
struct StubRunner(TurnOutcome);
impl TurnRunner for StubRunner {
fn run_turn<'a>(
&'a self,
_req: TurnRequest,
) -> Pin<Box<dyn Future<Output = TurnOutcome> + Send + 'a>> {
let outcome = self.0.clone();
Box::pin(async move { outcome })
}
}
struct RecordingRunner {
outcomes: std::sync::Mutex<std::collections::VecDeque<TurnOutcome>>,
seen: Arc<std::sync::Mutex<Vec<TurnRequest>>>,
}
impl TurnRunner for RecordingRunner {
fn run_turn<'a>(
&'a self,
req: TurnRequest,
) -> Pin<Box<dyn Future<Output = TurnOutcome> + Send + 'a>> {
self.seen.lock().unwrap().push(req);
let outcome = self
.outcomes
.lock()
.unwrap()
.pop_front()
.unwrap_or_else(|| TurnOutcome::Completed {
text: "done".to_owned(),
});
Box::pin(async move { outcome })
}
}
struct AcceptingApprovals;
impl polyc_a2a::ApprovalResponder for AcceptingApprovals {
fn respond<'a>(
&'a self,
_turn_id: &'a str,
_request_id: &'a str,
_approved: bool,
_reason: &'a str,
_conversation_id: &'a str,
_resolve_token: &'a str,
) -> Pin<Box<dyn Future<Output = Result<bool, String>> + Send + 'a>> {
Box::pin(async { Ok(true) })
}
}
const TOKEN: &str = "test-bearer-token";
struct AuthorityRunner {
outcome: TurnOutcome,
admitted: Mutex<HashMap<polyc_rpc_client::IngressIdentity, String>>,
admissions: Arc<AtomicUsize>,
}
impl TurnRunner for AuthorityRunner {
fn receive_ingress<'a>(
&'a self,
req: TurnRequest,
) -> Pin<
Box<
dyn Future<
Output = Result<
polyc_a2a::task::IngressReceipt,
polyc_a2a::task::IngressReceiptError,
>,
> + Send
+ 'a,
>,
> {
let mut admitted = self.admitted.lock().unwrap();
let result = match admitted.get(&req.source_identity) {
Some(text) if text != &req.text => Err(polyc_a2a::task::IngressReceiptError {
message: "source event was already received with different content".to_owned(),
retryable: false,
content_conflict: true,
}),
Some(_) => Ok(polyc_a2a::task::IngressReceipt {
dispatch_id: "d1-dispatch".to_owned(),
}),
None => {
self.admissions.fetch_add(1, Ordering::SeqCst);
admitted.insert(req.source_identity, req.text);
Ok(polyc_a2a::task::IngressReceipt {
dispatch_id: "d1-dispatch".to_owned(),
})
}
};
drop(admitted);
Box::pin(async move { result })
}
fn run_turn<'a>(
&'a self,
_req: TurnRequest,
) -> Pin<Box<dyn Future<Output = TurnOutcome> + Send + 'a>> {
let outcome = self.outcome.clone();
Box::pin(async move { outcome })
}
}
fn app_with_authority() -> (axum::Router, Arc<AtomicUsize>) {
let admissions = Arc::new(AtomicUsize::new(0));
let signer = Signer::from_seed(7);
let card = signed_card(
&polyc_a2a::card::CardConfig {
name: "Polychrome".to_owned(),
description: "test".to_owned(),
url: "https://agent.example/".to_owned(),
version: "0.1.3".to_owned(),
},
&signer,
);
let app = router(AppState {
card: Arc::new(card),
runner: Arc::new(AuthorityRunner {
outcome: TurnOutcome::Completed {
text: "it is sunny".to_owned(),
},
admitted: Mutex::new(HashMap::new()),
admissions: admissions.clone(),
}),
approvals: Arc::new(UnconfiguredApprovalResponder),
store: Arc::new(InMemoryTaskStore::new()),
turn_limit: AdmissionGate::new(64),
peers: polyc_a2a::PeerAuthenticator::single("test-peer", TOKEN).unwrap(),
});
(app, admissions)
}
fn app_with(outcome: TurnOutcome) -> axum::Router {
let signer = Signer::from_seed(7);
let card = signed_card(
&polyc_a2a::card::CardConfig {
name: "Polychrome".to_owned(),
description: "test".to_owned(),
url: "https://agent.example/".to_owned(),
version: "0.1.3".to_owned(),
},
&signer,
);
router(AppState {
card: Arc::new(card),
runner: Arc::new(StubRunner(outcome)),
approvals: Arc::new(UnconfiguredApprovalResponder),
store: Arc::new(InMemoryTaskStore::new()),
turn_limit: AdmissionGate::new(64),
peers: polyc_a2a::PeerAuthenticator::single("test-peer", TOKEN).unwrap(),
})
}
async fn rpc(app: axum::Router, request: &Value) -> Value {
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/")
.header(header::CONTENT_TYPE, "application/json")
.header(header::AUTHORIZATION, format!("Bearer {TOKEN}"))
.body(Body::from(serde_json::to_vec(request).unwrap()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
let body = response.into_body().collect().await.unwrap().to_bytes();
serde_json::from_slice(&body).unwrap()
}
fn send_message(text: &str) -> Value {
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "SendMessage",
"params": {
"message": {
"role": "ROLE_USER",
"messageId": "m1",
"contextId": "ctx-1",
"parts": [{ "text": text }]
}
}
})
}
#[tokio::test]
async fn send_message_returns_a_wrapped_completed_task() {
let app = app_with(TurnOutcome::Completed {
text: "it is sunny".to_owned(),
});
let response = rpc(app, &send_message("what is the weather?")).await;
assert_eq!(response["jsonrpc"], "2.0");
assert_eq!(response["id"], 1);
assert!(response.get("error").is_none(), "{response}");
let task = &response["result"]["task"];
assert!(
task.is_object(),
"result must be wrapped under `task`: {response}"
);
assert!(task.get("kind").is_none(), "v1.0 tasks carry no `kind`");
assert_ne!(task["contextId"], "ctx-1");
assert!(
task["contextId"].as_str().is_some_and(|id| !id.is_empty()),
"the authenticated peer receives a stable peer-scoped context id"
);
assert_eq!(task["status"]["state"], "TASK_STATE_COMPLETED");
assert_eq!(task["status"]["message"]["parts"][0]["text"], "it is sunny");
assert_eq!(task["artifacts"][0]["parts"][0]["text"], "it is sunny");
assert_eq!(task["status"]["message"]["role"], "ROLE_AGENT");
}
#[tokio::test]
async fn approval_pause_maps_to_input_required() {
let app = app_with(TurnOutcome::InputRequired {
turn_id: "00000000-0000-0000-0000-000000000001".to_owned(),
request_id: "call-9".to_owned(),
tool_name: "wire_transfer".to_owned(),
prompt: "Approval required to run `wire_transfer`".to_owned(),
resolve_token: "tok-9".to_owned(),
});
let response = rpc(app, &send_message("send the money")).await;
let task = &response["result"]["task"];
assert_eq!(task["status"]["state"], "TASK_STATE_INPUT_REQUIRED");
assert!(
task["status"]["message"]["parts"][0]["text"]
.as_str()
.unwrap()
.contains("wire_transfer")
);
}
#[tokio::test]
async fn a_decision_reply_redrives_with_no_utterance() {
let seen = Arc::new(std::sync::Mutex::new(Vec::new()));
let signer = Signer::from_seed(7);
let card = signed_card(
&polyc_a2a::card::CardConfig {
name: "Polychrome".to_owned(),
description: "test".to_owned(),
url: "https://agent.example/".to_owned(),
version: "0.1.3".to_owned(),
},
&signer,
);
let app = router(AppState {
card: Arc::new(card),
runner: Arc::new(RecordingRunner {
outcomes: std::sync::Mutex::new(
[
TurnOutcome::InputRequired {
turn_id: "00000000-0000-0000-0000-000000000001".to_owned(),
request_id: "call-9".to_owned(),
tool_name: "wire_transfer".to_owned(),
prompt: "Approval required to run `wire_transfer`".to_owned(),
resolve_token: "tok-9".to_owned(),
},
TurnOutcome::Completed {
text: "sent".to_owned(),
},
]
.into_iter()
.collect(),
),
seen: Arc::clone(&seen),
}),
approvals: Arc::new(AcceptingApprovals),
store: Arc::new(InMemoryTaskStore::new()),
turn_limit: AdmissionGate::new(64),
peers: polyc_a2a::PeerAuthenticator::single("test-peer", TOKEN).unwrap(),
});
let paused = rpc(app.clone(), &send_message("send the money")).await;
let task = &paused["result"]["task"];
assert_eq!(task["status"]["state"], "TASK_STATE_INPUT_REQUIRED");
let context_id = task["contextId"].as_str().unwrap().to_owned();
let task_id = task["id"].as_str().unwrap().to_owned();
let decision = json!({
"jsonrpc": "2.0",
"id": 2,
"method": "SendMessage",
"params": {
"message": {
"role": "ROLE_USER",
"messageId": "m2",
"contextId": context_id,
"taskId": task_id,
"parts": [{ "text": "approve" }]
}
}
});
let resumed = rpc(app, &decision).await;
assert!(resumed.get("error").is_none(), "{resumed}");
let redrive_text = {
let seen = seen.lock().unwrap();
assert_eq!(seen.len(), 2, "the decision must redrive the turn");
seen[1].text.clone()
};
assert!(
polyc_rpc_client::brought_no_utterance(&[polyc_rpc_client::user_message(&redrive_text)]),
"the redrive carried an utterance ({redrive_text:?}), so the control plane \
would decide the room afresh instead of inheriting what the paused turn proved"
);
}
#[tokio::test]
async fn v0x_slash_method_is_rejected() {
let app = app_with(TurnOutcome::Completed {
text: String::new(),
});
let request = json!({ "jsonrpc": "2.0", "id": 5, "method": "message/send", "params": {} });
let response = rpc(app, &request).await;
assert_eq!(response["id"], 5);
assert_eq!(response["error"]["code"], -32601);
}
#[tokio::test]
async fn redelivered_message_id_admits_the_source_once() {
let (app, admissions) = app_with_authority();
let first = rpc(app.clone(), &send_message("what is the weather?")).await;
assert!(first.get("error").is_none(), "{first}");
let retry = rpc(app, &send_message("what is the weather?")).await;
assert!(retry.get("error").is_none(), "{retry}");
assert_eq!(
first["result"]["task"]["id"], retry["result"]["task"]["id"],
"the redelivery must return the first task, not mint a second"
);
assert_eq!(
admissions.load(Ordering::SeqCst),
1,
"one message id must admit one durable source, not two"
);
}
#[tokio::test]
async fn reused_message_id_with_changed_content_is_not_acknowledged() {
let (app, admissions) = app_with_authority();
let first = rpc(app.clone(), &send_message("what is the weather?")).await;
assert!(first.get("error").is_none(), "{first}");
let conflict = rpc(app, &send_message("delete the mailbox")).await;
assert!(
conflict.get("result").is_none(),
"changed content must not return a task: {conflict}"
);
assert_eq!(conflict["error"]["code"], -32602);
assert!(
conflict["error"]["message"]
.as_str()
.unwrap()
.contains("different content"),
"{conflict}"
);
assert_eq!(
admissions.load(Ordering::SeqCst),
1,
"the refused conflict must not admit a second source"
);
}