use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
use yoagent::mcp::transport::McpTransport;
use yoagent::mcp::types::JsonRpcRequest;
use yoagent::mcp::HttpTransport;
async fn serve_pieces(pieces: Vec<Vec<u8>>, terminate: bool) -> String {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = format!("http://{}", listener.local_addr().unwrap());
tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let mut buf = [0u8; 4096];
let _ = socket.read(&mut buf).await;
let head = "HTTP/1.1 200 OK\r\n\
Content-Type: text/event-stream\r\n\
Transfer-Encoding: chunked\r\n\
\r\n";
let _ = socket.write_all(head.as_bytes()).await;
let _ = socket.flush().await;
for piece in pieces {
let _ = socket
.write_all(format!("{:x}\r\n", piece.len()).as_bytes())
.await;
let _ = socket.write_all(&piece).await;
let _ = socket.write_all(b"\r\n").await;
let _ = socket.flush().await;
}
if terminate {
let _ = socket.write_all(b"0\r\n\r\n").await;
let _ = socket.flush().await;
}
tokio::time::sleep(Duration::from_secs(300)).await;
});
addr
}
fn split_at(bytes: &[u8], at: usize) -> Vec<Vec<u8>> {
vec![bytes[..at].to_vec(), bytes[at..].to_vec()]
}
#[tokio::test]
async fn event_boundary_split_across_chunks_is_still_found() {
let request = JsonRpcRequest::new("tools/list", None);
let frames = format!(
"event: message\ndata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"tools\":[]}}}}\n\n",
request.id
);
let at = frames.len() - 1;
assert_eq!(
&frames.as_bytes()[at - 1..at + 1],
b"\n\n",
"split the boundary"
);
let addr = serve_pieces(split_at(frames.as_bytes(), at), false).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(10), transport.send(request))
.await
.expect("a boundary spanning a chunk seam must still be found")
.expect("the response must parse");
assert!(response.result.unwrap()["tools"].is_array());
}
#[tokio::test]
async fn multibyte_character_split_across_chunks_is_not_corrupted() {
let request = JsonRpcRequest::new("tools/call", None);
let text = "café 日本語 🎉";
let frames = format!(
"event: message\ndata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"text\":\"{}\"}}}}\n\n",
request.id, text
);
let at = frames.find('日').unwrap() + 1;
assert!(
!frames.is_char_boundary(at),
"split inside a multi-byte char"
);
let addr = serve_pieces(split_at(frames.as_bytes(), at), false).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(10), transport.send(request))
.await
.expect("must not hang")
.expect("the response must parse");
assert_eq!(
response.result.unwrap()["text"].as_str().unwrap(),
text,
"a character split across chunks must survive intact"
);
}
#[tokio::test]
async fn crlf_boundary_split_mid_sequence_is_normalized() {
let request = JsonRpcRequest::new("tools/list", None);
let frames = format!(
"event: message\r\ndata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"ok\":true}}}}\r\n\r\n",
request.id
);
let at = frames.len() - 1;
assert_eq!(&frames.as_bytes()[at - 1..at + 1], b"\r\n");
let addr = serve_pieces(split_at(frames.as_bytes(), at), false).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(10), transport.send(request))
.await
.expect("CRLF split across a seam must still frame")
.expect("the response must parse");
assert_eq!(response.result.unwrap()["ok"], true);
}
#[tokio::test]
async fn bare_cr_framing_is_supported() {
let request = JsonRpcRequest::new("tools/list", None);
let frames = format!(
"event: message\rdata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"ok\":true}}}}\r\r",
request.id
);
let addr = serve_pieces(vec![frames.into_bytes()], true).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(10), transport.send(request))
.await
.expect("must not hang")
.expect("bare-CR framing must parse");
assert_eq!(response.result.unwrap()["ok"], true);
}
#[tokio::test]
async fn unterminated_final_event_after_a_notification_is_parsed() {
let request = JsonRpcRequest::new("tools/call", None);
let frames = format!(
concat!(
"event: message\ndata: {{\"jsonrpc\":\"2.0\",\"method\":\"notifications/progress\"}}\n\n",
"event: message\ndata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"ok\":true}}}}\n"
),
request.id
);
let addr = serve_pieces(vec![frames.into_bytes()], true).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(10), transport.send(request))
.await
.expect("must not hang")
.expect("an unterminated final event must be parsed at EOF");
assert_eq!(response.result.unwrap()["ok"], true);
}
#[tokio::test]
async fn truncation_after_the_answer_still_succeeds() {
let request = JsonRpcRequest::new("tools/list", None);
let frames = format!(
"event: message\ndata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"ok\":true}}}}\n\n",
request.id
);
let pieces = vec![frames.into_bytes(), b"garbage-not-a-chunk".to_vec()];
let addr = serve_pieces(pieces, false).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(10), transport.send(request))
.await
.expect("must not hang")
.expect("a body that dies after the answer must not fail the call");
assert_eq!(response.result.unwrap()["ok"], true);
}
#[tokio::test]
async fn byte_at_a_time_delivery_assembles_correctly() {
let request = JsonRpcRequest::new("tools/call", None);
let text = "日本語";
let frames = format!(
"event: message\ndata: {{\"jsonrpc\":\"2.0\",\"id\":{},\"result\":{{\"text\":\"{}\"}}}}\n\n",
request.id, text
);
let pieces: Vec<Vec<u8>> = frames.as_bytes().iter().map(|b| vec![*b]).collect();
let addr = serve_pieces(pieces, false).await;
let transport = HttpTransport::new(&addr).unwrap();
let response = tokio::time::timeout(Duration::from_secs(20), transport.send(request))
.await
.expect("must not hang")
.expect("the response must parse");
assert_eq!(response.result.unwrap()["text"].as_str().unwrap(), text);
}