#![cfg(all(
feature = "streamable-http",
feature = "v1-compat",
not(target_arch = "wasm32")
))]
use std::net::{Ipv4Addr, SocketAddr};
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use serde_json::{json, Value};
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpStream;
use tokio::sync::Mutex;
use tokio::task::JoinHandle;
use pmcp::server::streamable_http_server::{StreamableHttpServer, StreamableHttpServerConfig};
use pmcp::types::elicitation::ElicitRequestParams;
use pmcp::types::protocol::LATEST_PROTOCOL_VERSION;
use pmcp::types::sampling::{CreateMessageParams, SamplingMessage, SamplingMessageContent};
use pmcp::types::Role;
use pmcp::{RequestHandlerExtra, Server, ToolHandler};
const BOUND: Duration = Duration::from_secs(5);
const HOLD: Duration = Duration::from_secs(10);
const MUTEX_BOUND: Duration = Duration::from_secs(1);
const QUIET: Duration = Duration::from_millis(600);
const PEER_ABSENT: &str = "peer absent on this transport";
const HOST_MODEL: &str = "host-model";
const PROGRESS_TOKEN: &str = "progress-token-1";
const PROGRESS_TOTAL: f64 = 2.0;
const PROGRESS_STEPS: usize = 2;
const FLOOD_REPORTS: u32 = 200;
const FLOOD_FRAME_CEILING: usize = 20;
const PROGRESS_METHOD: &str = "notifications/progress";
struct SamplerTool;
#[async_trait]
impl ToolHandler for SamplerTool {
async fn handle(&self, _args: Value, extra: RequestHandlerExtra) -> pmcp::Result<Value> {
let peer = peer_or_absent(&extra)?;
let params = CreateMessageParams::new(vec![SamplingMessage::new(
Role::User,
SamplingMessageContent::Text {
text: "summarize".to_string(),
meta: None,
},
)]);
let result = peer.sample(params).await?;
Ok(json!(format!("sampled:{}", result.model)))
}
}
struct RootsTool;
#[async_trait]
impl ToolHandler for RootsTool {
async fn handle(&self, _args: Value, extra: RequestHandlerExtra) -> pmcp::Result<Value> {
let peer = peer_or_absent(&extra)?;
let roots = peer.list_roots().await?;
Ok(json!(format!("roots:{}", roots.roots.len())))
}
}
struct ElicitTool;
#[async_trait]
impl ToolHandler for ElicitTool {
async fn handle(&self, _args: Value, extra: RequestHandlerExtra) -> pmcp::Result<Value> {
let peer = peer_or_absent(&extra)?;
let answer = peer
.elicit(ElicitRequestParams::Form {
message: "which environment?".to_string(),
requested_schema: json!({
"type": "object",
"properties": { "env": { "type": "string" } },
"required": ["env"],
}),
})
.await?;
let env = answer
.content
.as_ref()
.and_then(|c| c.get("env"))
.and_then(Value::as_str)
.unwrap_or("<none>")
.to_string();
Ok(json!(format!("elicit:{:?}:{env}", answer.action)))
}
}
struct FastTool;
#[async_trait]
impl ToolHandler for FastTool {
async fn handle(&self, _args: Value, _extra: RequestHandlerExtra) -> pmcp::Result<Value> {
Ok(json!("fast-done"))
}
}
struct ProgressTool;
#[async_trait]
impl ToolHandler for ProgressTool {
async fn handle(&self, _args: Value, extra: RequestHandlerExtra) -> pmcp::Result<Value> {
extra
.report_progress(1.0, Some(PROGRESS_TOTAL), Some("half".to_string()))
.await?;
extra
.report_progress(
PROGRESS_TOTAL,
Some(PROGRESS_TOTAL),
Some("done".to_string()),
)
.await?;
Ok(json!("progress-done"))
}
}
struct ProgressFloodTool;
#[async_trait]
impl ToolHandler for ProgressFloodTool {
async fn handle(&self, _args: Value, extra: RequestHandlerExtra) -> pmcp::Result<Value> {
for step in 1..=FLOOD_REPORTS {
extra.report_progress(f64::from(step), None, None).await?;
}
Ok(json!(format!("flood:{FLOOD_REPORTS}")))
}
}
#[derive(Default)]
struct HoldGate {
entered: tokio::sync::Notify,
release: tokio::sync::Notify,
}
struct MutexHoldingTool(Arc<HoldGate>);
#[async_trait]
impl ToolHandler for MutexHoldingTool {
async fn handle(&self, _args: Value, _extra: RequestHandlerExtra) -> pmcp::Result<Value> {
self.0.entered.notify_one();
let _ = tokio::time::timeout(HOLD, self.0.release.notified()).await;
Ok(json!("held"))
}
}
fn peer_or_absent(
extra: &RequestHandlerExtra,
) -> pmcp::Result<Arc<dyn pmcp::shared::peer::PeerHandle>> {
extra.peer().cloned().ok_or_else(|| {
pmcp::Error::protocol(pmcp::ErrorCode::INTERNAL_ERROR, PEER_ABSENT.to_string())
})
}
fn build_server(gate: Arc<HoldGate>) -> Server {
Server::builder()
.name("http-peer-roundtrip")
.version("1.0.0")
.tool("sampler", SamplerTool)
.tool("roots", RootsTool)
.tool("elicit", ElicitTool)
.tool("fast", FastTool)
.tool("hold", MutexHoldingTool(gate))
.tool("progress", ProgressTool)
.tool("flood", ProgressFloodTool)
.build()
.expect("server builds")
}
fn build_dual_era_server() -> Server {
Server::builder()
.name("http-peer-roundtrip-v2")
.version("1.0.0")
.tool("progress", ProgressTool)
.with_supported_protocol_versions([
pmcp::types::protocol::ProtocolVersion(LATEST_PROTOCOL_VERSION.to_string()),
pmcp::types::protocol::ProtocolVersion(
pmcp::types::protocol::PROTOCOL_VERSION_2026_07_28.to_string(),
),
])
.build()
.expect("server builds")
}
async fn spawn() -> (SocketAddr, JoinHandle<()>, Arc<HoldGate>) {
let gate = Arc::new(HoldGate::default());
let (bound, handle) = spawn_server(build_server(gate.clone())).await;
(bound, handle, gate)
}
async fn spawn_server(server: Server) -> (SocketAddr, JoinHandle<()>) {
let addr = SocketAddr::new(Ipv4Addr::LOCALHOST.into(), 0);
let server = Arc::new(Mutex::new(server));
StreamableHttpServer::with_config(addr, server, StreamableHttpServerConfig::default())
.start()
.await
.expect("server starts on an ephemeral port")
}
struct Conn {
reader: BufReader<TcpStream>,
status: u16,
headers: Vec<(String, String)>,
buffer: String,
chunked: bool,
remaining: usize,
finished: bool,
}
impl Conn {
async fn open(addr: SocketAddr, verb: &str, extra: &[(String, String)], body: &str) -> Self {
let stream = TcpStream::connect(addr).await.expect("connects");
let accept = if verb == "GET" {
"text/event-stream"
} else {
"application/json, text/event-stream"
};
let mut request = format!(
"{verb} / HTTP/1.1\r\nHost: {addr}\r\nAccept: {accept}\r\n\
Content-Type: application/json\r\nContent-Length: {}\r\n",
body.len()
);
for (name, value) in extra {
request.push_str(&format!("{name}: {value}\r\n"));
}
request.push_str("\r\n");
request.push_str(body);
let mut reader = BufReader::new(stream);
reader
.get_mut()
.write_all(request.as_bytes())
.await
.expect("request written");
let mut status_line = String::new();
reader
.read_line(&mut status_line)
.await
.expect("status line");
let status = status_line
.split_whitespace()
.nth(1)
.and_then(|s| s.parse().ok())
.unwrap_or(0);
let mut headers = Vec::new();
loop {
let mut line = String::new();
reader.read_line(&mut line).await.expect("header line");
let line = line.trim_end();
if line.is_empty() {
break;
}
if let Some((name, value)) = line.split_once(':') {
headers.push((name.trim().to_ascii_lowercase(), value.trim().to_string()));
}
}
let chunked = headers
.iter()
.any(|(n, v)| n == "transfer-encoding" && v.contains("chunked"));
let remaining = headers
.iter()
.find(|(n, _)| n == "content-length")
.and_then(|(_, v)| v.parse().ok())
.unwrap_or(0);
Self {
reader,
status,
headers,
buffer: String::new(),
chunked,
remaining,
finished: false,
}
}
fn header(&self, name: &str) -> Option<&str> {
self.headers
.iter()
.find(|(n, _)| n == name)
.map(|(_, v)| v.as_str())
}
async fn pull(&mut self) -> bool {
if self.finished {
return false;
}
if !self.chunked {
let mut payload = vec![0u8; self.remaining];
let ok = self.remaining > 0 && self.reader.read_exact(&mut payload).await.is_ok();
self.finished = true;
if !ok {
return false;
}
self.buffer.push_str(&String::from_utf8_lossy(&payload));
return true;
}
let mut size_line = String::new();
if self.reader.read_line(&mut size_line).await.unwrap_or(0) == 0 {
self.finished = true;
return false;
}
let size_token = size_line.trim().split(';').next().unwrap_or("").to_string();
let Ok(size) = usize::from_str_radix(&size_token, 16) else {
self.finished = true;
return false;
};
if size == 0 {
self.finished = true;
return false;
}
let mut payload = vec![0u8; size];
if self.reader.read_exact(&mut payload).await.is_err() {
self.finished = true;
return false;
}
let mut crlf = [0u8; 2];
let _ = self.reader.read_exact(&mut crlf).await;
self.buffer.push_str(&String::from_utf8_lossy(&payload));
true
}
fn take_block(&mut self) -> Option<String> {
let end = self.buffer.find("\n\n")?;
let block = self.buffer[..end].to_string();
self.buffer.drain(..end + 2);
Some(block)
}
async fn next_data(&mut self) -> Option<String> {
loop {
if let Some(block) = self.take_block() {
let mut data = String::new();
for line in block.lines() {
if let Some(rest) = line.strip_prefix("data:") {
data.push_str(rest.trim_start());
}
}
if !data.is_empty() {
return Some(data);
}
continue;
}
if !self.pull().await {
if !self.buffer.trim().is_empty() {
let rest = std::mem::take(&mut self.buffer);
return Some(rest.trim().to_string());
}
return None;
}
}
}
async fn frame(&mut self, what: &str) -> Value {
let data = tokio::time::timeout(BOUND, self.next_data())
.await
.unwrap_or_else(|_| panic!("{what} must not hang"))
.unwrap_or_else(|| panic!("{what}: the stream ended with no frame"));
serde_json::from_str(&data).expect("every frame on this stream is JSON")
}
async fn silent_for(&mut self, window: Duration, whose: &str) {
if let Ok(Some(data)) = tokio::time::timeout(window, self.next_data()).await {
panic!("{whose} must receive nothing, but a frame was delivered to it: {data}");
}
}
async fn body(&mut self) -> String {
while self.pull().await {}
std::mem::take(&mut self.buffer)
}
}
fn envelope(method: &str, id: &Value, params: &Value) -> String {
json!({ "jsonrpc": "2.0", "id": id, "method": method, "params": params }).to_string()
}
fn init_body() -> String {
envelope(
"initialize",
&json!(1),
&json!({
"protocolVersion": LATEST_PROTOCOL_VERSION,
"capabilities": { "sampling": {}, "roots": {}, "elicitation": {} },
"clientInfo": { "name": "http-peer-fence", "version": "1.0.0" }
}),
)
}
fn call_body(id: i64, tool: &str) -> String {
envelope(
"tools/call",
&json!(id),
&json!({ "name": tool, "arguments": {} }),
)
}
fn call_body_with_progress(id: i64, tool: &str) -> String {
envelope(
"tools/call",
&json!(id),
&json!({
"name": tool,
"arguments": {},
"_meta": { "progressToken": PROGRESS_TOKEN },
}),
)
}
fn v2_call_body_with_progress(id: i64, tool: &str) -> String {
envelope(
"tools/call",
&json!(id),
&json!({
"name": tool,
"arguments": {},
"_meta": {
"progressToken": PROGRESS_TOKEN,
pmcp::testing::META_PROTOCOL_VERSION:
pmcp::types::protocol::PROTOCOL_VERSION_2026_07_28,
pmcp::testing::META_CLIENT_INFO:
{ "name": "http-peer-fence", "version": "1.0.0" },
pmcp::testing::META_CLIENT_CAPABILITIES:
{ "elicitation": {}, "sampling": {}, "roots": {} },
},
}),
)
}
fn session_header(session: &str) -> Vec<(String, String)> {
vec![("Mcp-Session-Id".to_string(), session.to_string())]
}
async fn post(
addr: SocketAddr,
extra: &[(String, String)],
body: &str,
what: &str,
) -> (u16, String) {
let mut conn = tokio::time::timeout(BOUND, Conn::open(addr, "POST", extra, body))
.await
.unwrap_or_else(|_| panic!("{what} must not hang"));
let status = conn.status;
let text = tokio::time::timeout(BOUND, conn.body())
.await
.unwrap_or_else(|_| panic!("{what} body must not hang"));
(status, text)
}
async fn open_session(addr: SocketAddr) -> String {
let mut conn = tokio::time::timeout(BOUND, Conn::open(addr, "POST", &[], &init_body()))
.await
.expect("initialize must not hang");
assert_eq!(conn.status, 200, "initialize is answered inline");
let session = conn
.header("mcp-session-id")
.expect("a stateful server mints a session id")
.to_string();
let _ = tokio::time::timeout(BOUND, conn.body())
.await
.expect("initialize body must not hang");
session
}
async fn open_stream(addr: SocketAddr, session: &str) -> Conn {
let conn = tokio::time::timeout(BOUND, Conn::open(addr, "GET", &session_header(session), ""))
.await
.expect("opening the SSE stream must not hang");
assert_eq!(conn.status, 200, "the SSE stream opens");
conn
}
fn spawn_call(addr: SocketAddr, session: &str, id: i64, tool: &str) -> JoinHandle<(u16, String)> {
let headers = session_header(session);
let body = call_body(id, tool);
tokio::spawn(async move { post(addr, &headers, &body, "a queued tools/call").await })
}
fn require_server_request(frame: &Value, method: &str) -> Value {
if frame.get("method").and_then(Value::as_str) == Some(method) {
return frame
.get("id")
.cloned()
.expect("a server-to-client request carries an id");
}
panic!(
"expected a server-to-client {method} request on this session's stream, got {frame} \
— the HTTP transport carries no peer ({PEER_ABSENT}); plan 11 injects it"
);
}
async fn answer(addr: SocketAddr, session: &str, id: &Value, result: Value) {
let body = json!({ "jsonrpc": "2.0", "id": id, "result": result }).to_string();
let (status, _) = post(
addr,
&session_header(session),
&body,
"the client's answer POST",
)
.await;
assert_eq!(status, 202, "an inbound response is accepted");
}
fn sampling_answer() -> Value {
json!({
"role": "assistant",
"content": { "type": "text", "text": "done" },
"model": HOST_MODEL,
})
}
fn reply_text(frame: &Value) -> String {
frame.to_string()
}
async fn collect_progress_until_reply(stream: &mut Conn, id: i64, what: &str) -> Vec<Value> {
let mut progress = Vec::new();
loop {
let frame = stream.frame(what).await;
if frame.get("method").and_then(Value::as_str) == Some(PROGRESS_METHOD) {
println!("observed {PROGRESS_METHOD} frame: {frame}");
progress.push(frame);
continue;
}
if frame.get("id").and_then(Value::as_i64) == Some(id) {
println!("observed the tools/call reply for id {id}: {frame}");
return progress;
}
}
}
fn progress_token_of(frame: &Value) -> Option<&str> {
frame
.get("params")
.and_then(|p| p.get("progressToken"))
.and_then(Value::as_str)
}
fn v2_call_headers(name: &str) -> Vec<(String, String)> {
vec![
("MCP-Method".to_string(), "tools/call".to_string()),
("Mcp-Name".to_string(), name.to_string()),
(
"MCP-Protocol-Version".to_string(),
pmcp::types::protocol::PROTOCOL_VERSION_2026_07_28.to_string(),
),
]
}
async fn teardown(handle: JoinHandle<()>, sockets: impl Send) {
drop(sockets);
handle.abort();
let _ = handle.await;
}
#[tokio::test]
async fn response_post_is_answered_while_a_tool_handler_holds_the_server_mutex() {
let (addr, handle, gate) = spawn().await;
let session = open_session(addr).await;
let parked = spawn_call(addr, &session, 2, "hold");
tokio::time::timeout(BOUND, gate.entered.notified())
.await
.expect("the parked handler must start while holding the server mutex");
let body = json!({ "jsonrpc": "2.0", "id": "dispatch-not-live", "result": {} }).to_string();
let started = std::time::Instant::now();
let posted = tokio::time::timeout(
MUTEX_BOUND,
post(
addr,
&session_header(&session),
&body,
"the inbound response POST",
),
)
.await;
let elapsed = started.elapsed();
gate.release.notify_one();
let (status, _) = posted.unwrap_or_else(|_| {
panic!(
"an inbound response POST must not hang behind a parked tool handler: \
still unanswered after {MUTEX_BOUND:?} while a handler holds Mutex<Server>. \
The inbound response is not classified before \
run_v2_header_gate / extract_and_validate_auth."
)
});
assert_eq!(status, 202, "an inbound response is accepted");
assert!(
elapsed < MUTEX_BOUND,
"the inbound response POST waited {elapsed:?} on the server mutex"
);
let _ = tokio::time::timeout(BOUND, parked).await;
teardown(handle, ()).await;
}
#[tokio::test]
async fn http_peer_sample_completes_over_a_v1_session() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let call = spawn_call(addr, &session, 2, "sampler");
let frame = stream.frame("the server-to-client sampling request").await;
let id = require_server_request(&frame, "sampling/createMessage");
answer(addr, &session, &id, sampling_answer()).await;
let reply = stream.frame("the tools/call reply").await;
assert!(
reply_text(&reply).contains(&format!("sampled:{HOST_MODEL}")),
"the tool must observe the host completion model: {reply}"
);
call.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn http_peer_list_roots_completes_over_a_v1_session() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let call = spawn_call(addr, &session, 2, "roots");
let frame = stream.frame("the server-to-client roots request").await;
let id = require_server_request(&frame, "roots/list");
answer(
addr,
&session,
&id,
json!({ "roots": [ { "uri": "file:///a" }, { "uri": "file:///b" } ] }),
)
.await;
let reply = stream.frame("the tools/call reply").await;
assert!(
reply_text(&reply).contains("roots:2"),
"the tool must observe the two host roots: {reply}"
);
call.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn http_peer_elicit_completes_over_a_v1_session() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let call = spawn_call(addr, &session, 2, "elicit");
let frame = stream
.frame("the server-to-client elicitation request")
.await;
let id = require_server_request(&frame, "elicitation/create");
answer(
addr,
&session,
&id,
json!({ "action": "accept", "content": { "env": "staging" } }),
)
.await;
let reply = stream.frame("the tools/call reply").await;
assert!(
reply_text(&reply).contains("staging"),
"the tool must observe the accepted form: {reply}"
);
call.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn a_second_request_is_processed_while_the_first_handler_parks() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let first = spawn_call(addr, &session, 2, "sampler");
let second = spawn_call(addr, &session, 3, "fast");
let mut sampler_reply: Option<Value> = None;
let mut fast_reply: Option<Value> = None;
let mut answered = false;
while sampler_reply.is_none() || fast_reply.is_none() {
let frame = stream.frame("a frame in the saturation driver").await;
if frame.get("method").and_then(Value::as_str) == Some("sampling/createMessage") {
let id = require_server_request(&frame, "sampling/createMessage");
answer(addr, &session, &id, sampling_answer()).await;
answered = true;
continue;
}
match frame.get("id").and_then(Value::as_i64) {
Some(2) => sampler_reply = Some(frame),
Some(3) => fast_reply = Some(frame),
_ => {},
}
}
let fast = fast_reply.expect("the second call is answered");
assert!(
reply_text(&fast).contains("fast-done"),
"the second request must be processed while the first parks: {fast}"
);
let sampled = sampler_reply.expect("the first call is answered");
assert!(
answered,
"the first handler must have reached its peer call — it never issued a \
sampling request ({PEER_ABSENT}); its reply was {sampled}"
);
assert!(
reply_text(&sampled).contains(&format!("sampled:{HOST_MODEL}")),
"the parked call must complete once answered: {sampled}"
);
first.abort();
second.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn a_server_to_client_request_reaches_only_the_issuing_session() {
let (addr, handle, _gate) = spawn().await;
let session_a = open_session(addr).await;
let session_b = open_session(addr).await;
assert_ne!(session_a, session_b, "two distinct sessions");
let mut stream_a = open_stream(addr, &session_a).await;
let mut stream_b = open_stream(addr, &session_b).await;
let call = spawn_call(addr, &session_a, 2, "sampler");
stream_b
.silent_for(QUIET, "the bystander session's stream")
.await;
let frame = stream_a
.frame("the issuing session's sampling request")
.await;
let id = require_server_request(&frame, "sampling/createMessage");
answer(addr, &session_a, &id, sampling_answer()).await;
let reply = stream_a
.frame("the issuing session's tools/call reply")
.await;
assert!(
reply_text(&reply).contains(&format!("sampled:{HOST_MODEL}")),
"the issuing session's call completes: {reply}"
);
stream_b
.silent_for(QUIET, "the bystander session's stream")
.await;
call.abort();
teardown(handle, (stream_a, stream_b)).await;
}
#[tokio::test]
async fn progress_frames_reach_the_issuing_session_before_the_tool_result() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let headers = session_header(&session);
let body = call_body_with_progress(2, "progress");
let call = tokio::spawn(async move { post(addr, &headers, &body, "the progress call").await });
let frames =
collect_progress_until_reply(&mut stream, 2, "a frame on the progress stream").await;
assert_eq!(
frames.len(),
PROGRESS_STEPS,
"the handler's `extra.report_progress(..)` calls must reach the client: observed \
{} frame(s) before the result, expected {PROGRESS_STEPS}",
frames.len()
);
for frame in &frames {
assert_eq!(
progress_token_of(frame),
Some(PROGRESS_TOKEN),
"every progress frame must echo the token the client sent: {frame}"
);
}
call.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn progress_frames_never_reach_an_unrelated_session() {
let (addr, handle, _gate) = spawn().await;
let session_a = open_session(addr).await;
let session_b = open_session(addr).await;
assert_ne!(session_a, session_b, "two distinct sessions");
let mut stream_a = open_stream(addr, &session_a).await;
let mut stream_b = open_stream(addr, &session_b).await;
let headers = session_header(&session_a);
let body = call_body_with_progress(2, "progress");
let call = tokio::spawn(async move { post(addr, &headers, &body, "the progress call").await });
let frames =
collect_progress_until_reply(&mut stream_a, 2, "a frame on the issuer's stream").await;
assert_eq!(
frames.len(),
PROGRESS_STEPS,
"the ISSUING session must receive every frame"
);
stream_b
.silent_for(QUIET, "the bystander session's stream")
.await;
call.abort();
teardown(handle, (stream_a, stream_b)).await;
}
#[tokio::test]
async fn no_progress_frame_is_emitted_without_a_progress_token() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let call = spawn_call(addr, &session, 2, "progress");
let frame = stream.frame("the tools/call reply").await;
assert_ne!(
frame.get("method").and_then(Value::as_str),
Some(PROGRESS_METHOD),
"with no progressToken the FIRST frame must be the result, not a progress \
notification: {frame}"
);
assert!(
reply_text(&frame).contains("progress-done"),
"the tool still succeeds — `report_progress` is a silent no-op with no reporter: {frame}"
);
call.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn a_tight_progress_loop_is_bounded_by_the_reporter_rate_limit() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let mut stream = open_stream(addr, &session).await;
let headers = session_header(&session);
let body = call_body_with_progress(2, "flood");
let call = tokio::spawn(async move { post(addr, &headers, &body, "the flood call").await });
let frames = collect_progress_until_reply(&mut stream, 2, "a frame on the flood stream").await;
assert!(
!frames.is_empty(),
"the flood must emit at least one frame — zero would mean the reporter is absent \
again, not that the bound works"
);
assert!(
frames.len() <= FLOOD_FRAME_CEILING,
"a tight loop of {FLOOD_REPORTS} `report_progress` calls produced {} frames, above the \
stated ceiling of {FLOOD_FRAME_CEILING}: the rate limit is the ONLY bound on an \
unbounded session sender (T-118.1-11-03)",
frames.len()
);
call.abort();
teardown(handle, stream).await;
}
#[tokio::test]
async fn progress_emission_still_succeeds_when_the_session_has_no_live_stream() {
let (addr, handle, _gate) = spawn().await;
let session = open_session(addr).await;
let (status, body) = post(
addr,
&session_header(&session),
&call_body_with_progress(2, "progress"),
"the progress call with no SSE stream",
)
.await;
assert_eq!(status, 200, "the call is answered inline: {body}");
assert!(
body.contains("progress-done"),
"the tool completes even though every progress frame was dropped: {body}"
);
teardown(handle, ()).await;
}
#[tokio::test]
async fn a_v2_call_with_a_progress_token_emits_progress_on_the_response_body() {
let (addr, handle) = spawn_server(build_dual_era_server()).await;
let (status, body) = post(
addr,
&v2_call_headers("progress"),
&v2_call_body_with_progress(2, "progress"),
"the v2 progress call",
)
.await;
assert_eq!(status, 200, "a v2 tools/call succeeds: {body}");
assert!(
body.contains("progress-done"),
"the tool runs to completion on v2: {body}"
);
let frames = body.matches(PROGRESS_METHOD).count();
assert_eq!(
frames, PROGRESS_STEPS,
"v2 delivers progress on the POST RESPONSE BODY as multi-frame SSE (plan 12 / D-16); \
expected {PROGRESS_STEPS} frames, saw {frames}: {body}"
);
assert!(
body.rfind("\"result\"") > body.rfind(PROGRESS_METHOD),
"and the result frame comes LAST, after every progress frame: {body}"
);
teardown(handle, ()).await;
}