use std::net::{Ipv4Addr, SocketAddr};
use std::sync::Arc;
use async_trait::async_trait;
use serde_json::{json, Value};
use tokio::sync::Mutex;
use pmcp::server::streamable_http_server::{StreamableHttpServer, StreamableHttpServerConfig};
use pmcp::types::capabilities::ServerCapabilities;
use pmcp::types::protocol::{protocol_era, Era, LATEST_PROTOCOL_VERSION};
use pmcp::{RequestHandlerExtra, Server, ToolHandler};
use pmcp_agent::invoker::{
ClientToolInvoker, ConnectorClient, ConnectorClientFactory, UrlConnectorClientFactory,
};
use pmcp_agent::seams::{ToolCall, ToolInvoker};
const DEFAULT_ADDR: &str = "127.0.0.1:8147";
const V2_TOOL: &str = "weather";
const CITY_KEY: &str = "city";
const CITY: &str = "Berlin";
const V1_TOOL: &str = "echo";
const POLL_CAP_SECS: u64 = 15;
const START_PAIRED_SERVER: &str = "cargo run --example s47_v2_stateless_mrtr --features full";
struct EchoTool;
#[async_trait]
impl ToolHandler for EchoTool {
async fn handle(&self, args: Value, _extra: RequestHandlerExtra) -> pmcp::Result<Value> {
Ok(json!({ "echoed": args }))
}
}
#[tokio::main]
async fn main() -> std::result::Result<(), Box<dyn std::error::Error>> {
let addr = std::env::args()
.nth(1)
.unwrap_or_else(|| DEFAULT_ADDR.to_string());
let endpoint = format!("http://{addr}/");
println!();
println!("=============================================================");
println!(" pmcp-agent CONNECTOR -> {endpoint}");
println!("=============================================================");
demo_v2_connection(&endpoint).await?;
demo_v1_fallback().await?;
demo_unreachable_propagates().await?;
println!();
println!("=============================================================");
println!(" All three demonstrations behaved as documented.");
println!("=============================================================");
Ok(())
}
async fn demo_v2_connection(endpoint: &str) -> std::result::Result<(), Box<dyn std::error::Error>> {
println!();
println!("[1] v2 (2026-07-28) connection — the paired s47 server");
println!("-------------------------------------------------------------");
let factory = UrlConnectorClientFactory::new();
let connector = factory.client_for(endpoint).await.map_err(|error| {
format!(
"could not connect to {endpoint}: {error}\n \
Is the paired server running? Start it with:\n {START_PAIRED_SERVER}"
)
})?;
let negotiated = report_era(connector.as_ref(), Era::V2)?;
println!(" negotiated : {negotiated} (classified as the v2 era)");
println!(" handshake : none — server/discover was the first request");
let invoker = ClientToolInvoker::new(Arc::clone(&connector), POLL_CAP_SECS);
let tools = invoker.list_tools().await;
if tools.is_empty() {
return Err(
format!("the v2 server advertised no tools; expected at least {V2_TOOL}").into(),
);
}
println!(" tools/list : {}", tool_names(&tools));
let outcome = invoker
.invoke(ToolCall {
id: "demo-1".to_string(),
name: V2_TOOL.to_string(),
arguments: json!({ CITY_KEY: CITY }),
connector: None,
})
.await;
if outcome.is_error {
return Err(format!(
"calling {V2_TOOL} over v2 failed: {}",
outcome.error.unwrap_or_default()
)
.into());
}
println!(" tools/call : {}", outcome.content);
Ok(())
}
async fn demo_v1_fallback() -> std::result::Result<(), Box<dyn std::error::Error>> {
println!();
println!("[2] v1 (2025-11-25) fallback — an in-process v1-only server");
println!("-------------------------------------------------------------");
let server = Server::builder()
.name("s53-v1-only")
.version("1.0.0")
.capabilities(ServerCapabilities::tools_only())
.tool(V1_TOOL, EchoTool)
.build()?;
let http = StreamableHttpServer::with_config(
SocketAddr::new(Ipv4Addr::LOCALHOST.into(), 0),
Arc::new(Mutex::new(server)),
StreamableHttpServerConfig::default(),
);
let (bound, handle) = http.start().await?;
let endpoint = format!("http://{bound}/");
println!(" v1-only server: {endpoint}");
let outcome = run_v1_fallback(&endpoint).await;
handle.abort();
let _ = handle.await;
outcome
}
async fn run_v1_fallback(endpoint: &str) -> std::result::Result<(), Box<dyn std::error::Error>> {
let factory = UrlConnectorClientFactory::new();
let connector = factory
.client_for(endpoint)
.await
.map_err(|error| format!("the v1 fallback did not connect to {endpoint}: {error}"))?;
let negotiated = report_era(connector.as_ref(), Era::V1)?;
println!(" negotiated : {negotiated} (classified as the v1 era)");
if negotiated != LATEST_PROTOCOL_VERSION {
return Err(format!(
"the fallback must report the version the server ECHOED in its \
initialize result ({LATEST_PROTOCOL_VERSION}), not {negotiated}"
)
.into());
}
println!(" fallback rule : the endpoint ANSWERED, so v2 rejection => try v1");
let invoker = ClientToolInvoker::new(Arc::clone(&connector), POLL_CAP_SECS);
let outcome = invoker
.invoke(ToolCall {
id: "demo-2".to_string(),
name: V1_TOOL.to_string(),
arguments: json!({ "message": "hello from the fallback" }),
connector: None,
})
.await;
if outcome.is_error {
return Err(format!(
"calling {V1_TOOL} over the v1 fallback failed: {}",
outcome.error.unwrap_or_default()
)
.into());
}
println!(" tools/call : {}", outcome.content);
Ok(())
}
async fn demo_unreachable_propagates() -> std::result::Result<(), Box<dyn std::error::Error>> {
println!();
println!("[3] Unreachable host — the error PROPAGATES, no silent downgrade");
println!("-------------------------------------------------------------");
let endpoint = closed_loopback_endpoint()?;
println!(" closed port : {endpoint}");
let factory = UrlConnectorClientFactory::new();
match factory.client_for(&endpoint).await {
Ok(connector) => Err(format!(
"a host that never answered must NOT yield a connector; got one \
reporting era {}",
connector.negotiated_protocol_version().unwrap_or("<none>")
)
.into()),
Err(error) => {
println!(" factory says : {error}");
println!(" no v1 attempt was made — nothing answered, so there was");
println!(" no protocol signal to fall back on.");
Ok(())
},
}
}
fn report_era(
connector: &dyn ConnectorClient,
want: Era,
) -> std::result::Result<String, Box<dyn std::error::Error>> {
let Some(negotiated) = connector.negotiated_protocol_version() else {
return Err("the connector reported no negotiated protocol version".into());
};
let era = protocol_era(negotiated);
if era != want {
return Err(
format!("expected the {want:?} era, but {negotiated} classifies as {era:?}").into(),
);
}
Ok(negotiated.to_string())
}
fn tool_names(tools: &[pmcp::types::ToolInfo]) -> String {
tools
.iter()
.map(|tool| tool.name.clone())
.collect::<Vec<_>>()
.join(", ")
}
fn closed_loopback_endpoint() -> std::result::Result<String, Box<dyn std::error::Error>> {
let listener = std::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))?;
let addr = listener.local_addr()?;
drop(listener);
Ok(format!("http://{addr}/"))
}