#![cfg(all(
feature = "streamable-http",
feature = "http-client",
feature = "v1-compat",
not(target_arch = "wasm32")
))]
mod common;
use common::example_process::{
spawn_example, target_dir, wait_until_listening, wait_until_released,
};
use common::v2::{header, post, post_with_accept, v1_body, v2_body, v2_headers_for};
use pmcp::shared::http_constants::MCP_SESSION_ID;
use pmcp::types::protocol::LATEST_PROTOCOL_VERSION;
use serde_json::{json, Value};
use std::net::SocketAddr;
use std::time::Duration;
const EXAMPLE_REL_PATH: &str = "debug/examples/s54_v2_dual_conformance";
const BIND_ADDR: &str = "127.0.0.1:8159";
const LOGGING_TOOL: &str = "test_tool_with_logging";
const DIAGNOSTIC_TOOL: &str = "test_logging_tool";
const LOG_METHOD: &str = "notifications/message";
const LOG_LEVEL_META_KEY: &str = "io.modelcontextprotocol/logLevel";
const EXPECTED_LOG_RECORDS: usize = 3;
const ACCEPT_SSE: &str = "application/json, text/event-stream";
const READY_TIMEOUT: Duration = Duration::from_secs(30);
const RELEASE_TIMEOUT: Duration = Duration::from_secs(10);
const V1_COLLECT_WINDOW: Duration = Duration::from_secs(5);
const ARTIFACT_REL_PATH: &str = "118.2-09-log-records.json";
async fn v1_open_session(addr: SocketAddr) -> String {
let params = json!({
"protocolVersion": LATEST_PROTOCOL_VERSION,
"capabilities": {},
"clientInfo": { "name": "log-records-example-run", "version": "0.0.0" },
});
let response = post(addr, &[], &v1_body("initialize", json!(0), params)).await;
assert_eq!(
response.status, 200,
"the v1 handshake must succeed before the logging call: HTTP {} {}",
response.status, response.raw
);
response.mcp_session_id.unwrap_or_else(|| {
panic!(
"the example minted no Mcp-Session-Id on initialize, so the v1 leg cannot \
proceed. Response was: {}",
response.raw
)
})
}
async fn read_v1_session_stream(
addr: SocketAddr,
session: &str,
trigger: impl std::future::Future<Output = ()>,
) -> String {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let mut stream = tokio::net::TcpStream::connect(addr)
.await
.expect("the example accepts a GET connection");
let request = format!(
"GET / HTTP/1.1\r\nHost: {addr}\r\nAccept: text/event-stream\r\n\
Mcp-Session-Id: {session}\r\nConnection: keep-alive\r\n\r\n"
);
stream
.write_all(request.as_bytes())
.await
.expect("the GET request is written");
let mut observed = String::new();
let mut buffer = [0_u8; 8192];
let deadline = tokio::time::Instant::now() + V1_COLLECT_WINDOW;
while !observed.contains("\r\n\r\n") {
let read = tokio::time::timeout_at(deadline, stream.read(&mut buffer))
.await
.expect("the example answers the GET before the deadline")
.expect("the GET stream is readable");
assert!(read > 0, "the example closed the SSE stream immediately");
observed.push_str(&String::from_utf8_lossy(&buffer[..read]));
}
trigger.await;
while tokio::time::Instant::now() < deadline && !observed.contains("\"result\"") {
match tokio::time::timeout_at(deadline, stream.read(&mut buffer)).await {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(read)) => observed.push_str(&String::from_utf8_lossy(&buffer[..read])),
Ok(Err(error)) => panic!("the v1 SSE stream errored: {error}"),
}
}
observed
}
fn v2_call_body(tool: &str, id: Value, level: Option<&str>) -> String {
let params = json!({ "name": tool, "arguments": {} });
let base = v2_body("tools/call", id, params);
let Some(level) = level else {
return base;
};
let mut value: Value = serde_json::from_str(&base).expect("the harness body parses");
value["params"]["_meta"][LOG_LEVEL_META_KEY] = json!(level);
value.to_string()
}
#[tokio::test]
async fn the_dual_conformance_example_emits_log_records_on_the_wire() {
let (addr, mut guard) = spawn_example(EXAMPLE_REL_PATH, BIND_ADDR);
wait_until_listening(addr, &mut guard, READY_TIMEOUT).await;
let session = v1_open_session(addr).await;
let set_level = post(
addr,
&[header(MCP_SESSION_ID, &session)],
&v1_body("logging/setLevel", json!(1), json!({ "level": "debug" })),
)
.await;
let v1_frames = read_v1_session_stream(addr, &session, async {
let call = post(
addr,
&[header(MCP_SESSION_ID, &session)],
&v1_body(
"tools/call",
json!(2),
json!({ "name": LOGGING_TOOL, "arguments": {} }),
),
)
.await;
assert_eq!(
call.status, 202,
"the v1 logging call must be ACCEPTED into the open session stream: {}",
call.raw
);
})
.await;
let params = json!({ "name": DIAGNOSTIC_TOOL, "arguments": {} });
let unauthorized = post_with_accept(
addr,
ACCEPT_SSE,
&v2_headers_for("tools/call", ¶ms),
&v2_call_body(DIAGNOSTIC_TOOL, json!(3), None),
)
.await;
let authorized = post_with_accept(
addr,
ACCEPT_SSE,
&v2_headers_for("tools/call", ¶ms),
&v2_call_body(DIAGNOSTIC_TOOL, json!(4), Some("info")),
)
.await;
let v1_records = v1_frames.matches(LOG_METHOD).count();
let unauthorized_records = unauthorized.raw.matches(LOG_METHOD).count();
let authorized_records = authorized.raw.matches(LOG_METHOD).count();
let artifact = json!({
"note": format!(
"Live `tools/call` logging records served by target/{EXAMPLE_REL_PATH} \
bound to {BIND_ADDR}. Phase 118.2-09, CONF-10."
),
"v1_set_level": { "status": set_level.status, "raw": set_level.raw },
"v1_session_stream": { "records": v1_records, "raw": v1_frames },
"v2_unauthorized": { "records": unauthorized_records, "raw": unauthorized.raw },
"v2_authorized": { "records": authorized_records, "raw": authorized.raw },
});
let artifact_path = target_dir().join(ARTIFACT_REL_PATH);
std::fs::write(
&artifact_path,
serde_json::to_string_pretty(&artifact).expect("the artifact always serializes"),
)
.unwrap_or_else(|error| panic!("could not write {}: {error}", artifact_path.display()));
assert_eq!(
set_level.status, 200,
"v1 `logging/setLevel` must be served: {}",
set_level.raw
);
assert!(
set_level.raw.contains(r#""result":{}"#),
"v1 `logging/setLevel` must answer a literal empty object: {}",
set_level.raw
);
assert!(
v1_records >= EXPECTED_LOG_RECORDS,
"v1 (2025-11-25): `{LOGGING_TOOL}` must put at least {EXPECTED_LOG_RECORDS} \
`{LOG_METHOD}` frames on the session stream — that is the floor the pinned \
`tools-call-with-logging` scenario asserts (`a.length < 3` fails it). Saw \
{v1_records}. If this is 0, check whether the tool is emitting with \
`tracing::info!` again, which reaches the operator and never the client. \
Recorded at {}. Stream was:\n{v1_frames}",
artifact_path.display()
);
assert!(
authorized_records >= 1,
"v2 (2026-07-28): `{DIAGNOSTIC_TOOL}` called WITH \
`_meta[\"{LOG_LEVEL_META_KEY}\"]` must emit at least one `{LOG_METHOD}` frame \
on the POST response body. Saw {authorized_records}. Without this the \
zero-record assertion below proves nothing. Body was:\n{}",
authorized.raw
);
assert_eq!(
unauthorized_records, 0,
"v2 (2026-07-28): SEP-2575 — a request that did NOT set \
`_meta[\"{LOG_LEVEL_META_KEY}\"]` must receive NO `{LOG_METHOD}` frame at all \
(\"If absent, the server MUST NOT send any notifications/message\"). Saw \
{unauthorized_records}. Body was:\n{}",
unauthorized.raw
);
drop(guard);
wait_until_released(addr, RELEASE_TIMEOUT).await;
}