#![cfg(all(
feature = "streamable-http",
feature = "http-client",
not(target_arch = "wasm32")
))]
mod common;
use std::collections::HashMap;
#[cfg(feature = "v1-compat")]
use std::net::SocketAddr;
#[cfg(feature = "v1-compat")]
use std::time::Duration;
#[cfg(feature = "v1-compat")]
use common::example_process::{
spawn_example, target_dir, wait_until_listening, wait_until_released,
};
#[cfg(feature = "v1-compat")]
use common::v2::{header, post, spawn_default_config, teardown, v1_body, v2_body, v2_headers_for};
use pmcp::server::builder::ServerCoreBuilder;
use pmcp::server::core::ProtocolHandler;
#[cfg(feature = "v1-compat")]
use pmcp::shared::http_constants::MCP_SESSION_ID;
use pmcp::types::completable::{
CompletionItem, CompletionProviderTrait, CompletionRequest, CompletionResponse,
StaticCompletionProvider,
};
use pmcp::types::jsonrpc::ResponsePayload;
#[cfg(feature = "v1-compat")]
use pmcp::types::protocol::LATEST_PROTOCOL_VERSION;
use pmcp::types::protocol::{
CompleteRequest, CompletionArgument, CompletionReference, InitializeRequest,
};
use pmcp::types::{ClientCapabilities, ClientRequest, JSONRPCResponse, Request, RequestId};
use pmcp::Implementation;
#[cfg(feature = "v1-compat")]
use pmcp::Server;
use proptest::prelude::*;
use proptest::test_runner::Config as ProptestConfig;
#[cfg(feature = "v1-compat")]
use serde_json::json;
use serde_json::Value;
const SUITE_PROMPT: &str = "test_prompt_with_arguments";
const SUITE_ARG_NAME: &str = "arg1";
const SUITE_ARG_VALUE: &str = "test";
const METHOD: &str = "completion/complete";
#[cfg(feature = "v1-compat")]
fn suite_params() -> Value {
params_with_partial(SUITE_ARG_VALUE)
}
#[cfg(feature = "v1-compat")]
fn params_with_partial(partial: &str) -> Value {
json!({
"ref": { "type": "ref/prompt", "name": SUITE_PROMPT },
"argument": { "name": SUITE_ARG_NAME, "value": partial },
})
}
fn request_with_partial(partial: &str) -> Request {
Request::Client(Box::new(ClientRequest::Complete(CompleteRequest {
r#ref: CompletionReference::Prompt {
name: SUITE_PROMPT.to_string(),
},
argument: CompletionArgument {
name: SUITE_ARG_NAME.to_string(),
value: partial.to_string(),
},
})))
}
const PROVIDER_VALUES: [&str; 3] = ["alpha-118-1-04", "beta-118-1-04", "gamma-118-1-04"];
fn registered_provider() -> StaticCompletionProvider {
StaticCompletionProvider::from_strings(
PROVIDER_VALUES.iter().map(|v| (*v).to_string()).collect(),
)
}
const EMPTY_PARTIAL: &str = "";
#[cfg(feature = "v1-compat")]
fn server_without_completions() -> Server {
Server::builder()
.name("completion-complete-http")
.version("1.0.0")
.build()
.expect("server builds")
}
#[cfg(feature = "v1-compat")]
fn server_with_completions() -> Server {
Server::builder()
.name("completion-complete-http-provider")
.version("1.0.0")
.completions(registered_provider())
.build()
.expect("server builds")
}
#[cfg(feature = "v1-compat")]
async fn v1_open_session(addr: SocketAddr) -> String {
let params = json!({
"protocolVersion": LATEST_PROTOCOL_VERSION,
"capabilities": {},
"clientInfo": { "name": "completion-complete", "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 {METHOD}: HTTP {} {}",
response.status, response.raw
);
response.mcp_session_id.unwrap_or_else(|| {
panic!(
"the server minted no Mcp-Session-Id on initialize, so the v1 leg cannot \
proceed. Response was: {}",
response.raw
)
})
}
#[cfg(feature = "v1-compat")]
async fn http_complete(addr: SocketAddr) -> common::v2::Resp {
http_complete_with(addr, suite_params()).await
}
#[cfg(feature = "v1-compat")]
async fn http_complete_with(addr: SocketAddr, params: Value) -> common::v2::Resp {
let session = v1_open_session(addr).await;
post(
addr,
&[header(MCP_SESSION_ID, &session)],
&v1_body(METHOD, json!(1), params),
)
.await
}
fn values_of(result: &Value, context: &str) -> Vec<String> {
let values = result["completion"]["values"]
.as_array()
.unwrap_or_else(|| {
panic!("{context}: result.completion.values must be an array. Was: {result}")
});
values
.iter()
.map(|value| {
value.as_str().map_or_else(
|| panic!("{context}: values must be string[] (schema.ts:2649). Was: {result}"),
str::to_string,
)
})
.collect()
}
#[cfg(feature = "v1-compat")]
#[tokio::test]
async fn http_server_answers_the_spec_completion_shape_with_no_handler_registered() {
let (addr, handle) = spawn_default_config(server_without_completions()).await;
let response = http_complete(addr).await;
assert_eq!(
response.status, 200,
"{METHOD} must be a SUCCESS even with no completion handler registered — the \
suite's own comment sanctions an empty array. Got HTTP {}: {}",
response.status, response.raw
);
assert!(
response.body["error"].is_null(),
"{METHOD} must not answer a JSON-RPC error with no handler registered: {}",
response.raw
);
let values = response.body["result"]["completion"]["values"].clone();
assert!(
values.is_array(),
"G-4: result.completion.values MUST be an array \
(CompleteResult, schema.ts:2644-2663). The suite fails this exact assertion \
with `Missing completion field`. Raw response was: {}",
response.raw
);
assert_eq!(
values.as_array().map(Vec::len),
Some(0),
"with no handler registered the array must be EMPTY, not populated: {}",
response.raw
);
teardown(handle, ()).await;
}
fn core_init_request() -> Request {
Request::Client(Box::new(ClientRequest::Initialize(InitializeRequest::new(
Implementation::new("completion-complete", "0.0.0"),
ClientCapabilities::default(),
))))
}
fn payload_of(response: &JSONRPCResponse) -> String {
serde_json::to_string(&response.payload)
.unwrap_or_else(|error| format!("<payload did not serialize: {error}>"))
}
async fn core_complete(core: &pmcp::server::core::ServerCore) -> JSONRPCResponse {
core_complete_with(core, SUITE_ARG_VALUE).await
}
async fn core_complete_with(
core: &pmcp::server::core::ServerCore,
partial: &str,
) -> JSONRPCResponse {
let init = core
.handle_request(RequestId::from(1i64), core_init_request(), None)
.await;
assert!(
matches!(init.payload, ResponsePayload::Result(_)),
"the ServerCore handshake must succeed before {METHOD}: {}",
payload_of(&init)
);
core.handle_request(RequestId::from(2i64), request_with_partial(partial), None)
.await
}
fn result_of(response: &JSONRPCResponse, context: &str) -> Value {
match &response.payload {
ResponsePayload::Result(value) => value.clone(),
ResponsePayload::Error(_) => panic!(
"{context}: {METHOD} answered a JSON-RPC error rather than the spec shape. \
Payload was: {}",
payload_of(response)
),
}
}
#[tokio::test]
async fn server_core_answers_the_spec_completion_shape_with_no_handler_registered() {
let core = ServerCoreBuilder::new()
.name("completion-complete-core")
.version("1.0.0")
.build()
.expect("core builds");
let response = core_complete(&core).await;
let result = result_of(
&response,
"G-4 (the ServerCore half): `ServerCore`'s `_ =>` catch-all used to answer -32601 \
here, which is a DIFFERENT defect from the high-level `Server`'s empty result object",
);
assert_eq!(
values_of(&result, "ServerCore, no handler"),
Vec::<String>::new(),
"with no handler registered the array must be EMPTY, not populated: {result}"
);
}
#[cfg(feature = "v1-compat")]
#[tokio::test]
async fn http_server_returns_the_values_of_a_provider_registered_through_server_builder() {
let (addr, handle) = spawn_default_config(server_with_completions()).await;
let response = http_complete_with(addr, params_with_partial(EMPTY_PARTIAL)).await;
assert_eq!(
response.status, 200,
"{METHOD} with a registered provider must succeed: HTTP {} {}",
response.status, response.raw
);
assert_eq!(
values_of(
&response.body["result"],
"ServerBuilder-registered provider"
),
PROVIDER_VALUES.to_vec(),
"a provider registered through `ServerBuilder::completions` must reach the \
high-level `Server` dispatcher. An empty array here means the slot exists but \
is wired to nothing. Raw response was: {}",
response.raw
);
teardown(handle, ()).await;
}
#[tokio::test]
async fn server_core_returns_the_values_of_a_provider_registered_through_core_builder() {
let core = ServerCoreBuilder::new()
.name("completion-complete-core-provider")
.version("1.0.0")
.completions(registered_provider())
.build()
.expect("core builds");
let response = core_complete_with(&core, EMPTY_PARTIAL).await;
let result = result_of(&response, "ServerCoreBuilder-registered provider");
assert_eq!(
values_of(&result, "ServerCoreBuilder-registered provider"),
PROVIDER_VALUES.to_vec(),
"a provider registered through `ServerCoreBuilder::completions` must reach the \
`ServerCore` dispatcher. An empty array here means the slot exists but is wired \
to nothing. Result was: {result}"
);
}
struct CountingProvider {
count: usize,
}
#[async_trait::async_trait]
impl CompletionProviderTrait for CountingProvider {
async fn complete(&self, _request: CompletionRequest) -> pmcp::Result<CompletionResponse> {
Ok(CompletionResponse {
completions: (0..self.count)
.map(|index| CompletionItem {
value: format!("candidate-{index}"),
label: None,
description: None,
icon: None,
metadata: HashMap::new(),
})
.collect(),
has_more: false,
continuation_token: None,
})
}
}
const SPEC_MAX_VALUES: usize = 100;
proptest! {
#![proptest_config(ProptestConfig::with_cases(48))]
#[test]
fn emitted_values_never_exceed_the_spec_bound(count in 0usize..250) {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime builds");
runtime.block_on(async move {
let core = ServerCoreBuilder::new()
.name("completion-complete-property")
.version("1.0.0")
.completions(CountingProvider { count })
.build()
.expect("core builds");
let response = core_complete_with(&core, EMPTY_PARTIAL).await;
let result = result_of(&response, "property: any provider output length");
let values = values_of(&result, "property: any provider output length");
prop_assert!(
values.len() <= SPEC_MAX_VALUES,
"a provider returning {count} values must be truncated to at most \
{SPEC_MAX_VALUES} (schema.ts:2649 @maxItems 100). Got {}: {result}",
values.len()
);
prop_assert_eq!(
values.len(),
count.min(SPEC_MAX_VALUES),
"truncation must keep as many values as the bound allows, not fewer: {}",
result
);
let truncated = count > SPEC_MAX_VALUES;
prop_assert_eq!(
result["completion"]["hasMore"].as_bool(),
Some(truncated),
"hasMore must be true EXACTLY when elements were dropped, so a truncated \
list is never readable as exhaustive: {}",
result
);
prop_assert_eq!(
result["completion"]["total"].as_u64(),
Some(count as u64),
"this provider reports has_more = false, so it returned everything it had \
and `total` is the TRUE total including the dropped elements: {}",
result
);
Ok(())
})?;
}
}
#[cfg(feature = "v1-compat")]
const EXAMPLE_REL_PATH: &str = "debug/examples/s54_v2_dual_conformance";
#[cfg(feature = "v1-compat")]
const BIND_ADDR: &str = "127.0.0.1:8153";
#[cfg(feature = "v1-compat")]
const ARTIFACT_REL_PATH: &str = "118.1-04-example-response.json";
#[cfg(feature = "v1-compat")]
const READY_TIMEOUT: Duration = Duration::from_secs(30);
#[cfg(feature = "v1-compat")]
const RELEASE_TIMEOUT: Duration = Duration::from_secs(10);
#[cfg(feature = "v1-compat")]
fn assert_example_answer(era: &str, response: &common::v2::Resp) {
assert_eq!(
response.status, 200,
"{era}: the example must serve {METHOD}, got HTTP {}: {}",
response.status, response.raw
);
let values = values_of(&response.body["result"], era);
assert!(
!values.is_empty(),
"{era}: the example registers a completion provider whose values are prefixed \
`test`, and the suite's partial is `test`, so at least one candidate must come \
back. An empty array means the registered provider was not reached. Raw \
response was: {}",
response.raw
);
for value in &values {
assert!(
value.starts_with(SUITE_ARG_VALUE),
"{era}: `{value}` does not start with the requested partial `{SUITE_ARG_VALUE}`, \
so the provider was not consulted with the request's argument value: {}",
response.raw
);
}
}
#[cfg(feature = "v1-compat")]
#[tokio::test]
async fn the_dual_conformance_example_serves_completion_complete_on_both_eras() {
let (addr, mut guard) = spawn_example(EXAMPLE_REL_PATH, BIND_ADDR);
wait_until_listening(addr, &mut guard, READY_TIMEOUT).await;
let params = suite_params();
let session = v1_open_session(addr).await;
let v1 = post(
addr,
&[header(MCP_SESSION_ID, &session)],
&v1_body(METHOD, json!(1), params.clone()),
)
.await;
let v2 = post(
addr,
&v2_headers_for(METHOD, ¶ms),
&v2_body(METHOD, json!(2), params.clone()),
)
.await;
let artifact = json!({
"note": format!(
"Live `{METHOD}` on {SUITE_PROMPT} (argument {SUITE_ARG_NAME} = \
\"{SUITE_ARG_VALUE}\"), served by target/{EXAMPLE_REL_PATH} bound to \
{BIND_ADDR}. Phase 118.1-04, G-4 / CONF-05."
),
"request": params,
"v1": { "raw": v1.raw, "body": v1.body },
"v2": { "raw": v2.raw, "body": v2.body },
});
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_example_answer("v1 (2025-11-25)", &v1);
assert_example_answer("v2 (2026-07-28)", &v2);
assert!(
v2.raw.contains("resultType"),
"the v2 leg carries no v2 result-envelope key, so it was not served as v2: {}",
v2.raw
);
assert!(
!v1.raw.contains("resultType"),
"the v1 leg carries a v2 result-envelope key, so the per-request era gate was \
bypassed: {}",
v1.raw
);
drop(guard);
wait_until_released(addr, RELEASE_TIMEOUT).await;
}