Bump alkcall 0.2 -> 0.3 (published 2026-09-04; the remediation of alkcall reviews 004-005, the upstream resolution of alkhttp review 003 WS-24/WS-25). Two source-break classes, both mechanical: - take_call_connection now returns Option<Arc<CallConnection>> (was bare value) — drop the double-wrap in WssSession::connect. - register/register_on/register_openable take &OperationRegistry (was &mut) — drop 51 stale `mut` bindings (clippy --fix), which is the interior-mutability change that makes the per-session fork surface possible. No behavior change in this crate yet; the 0.3 mechanisms (fork, serve_single_stream, op/register, install_bootstrap_discovery, collision policy) are the composition surface for alkhttp review 003 Unit 2 (WS-20..22, WS-26 wiring) — not wired here yet. Verification: cargo test 453 passed (default) / 575 passed (all-features), clippy both sides -D warnings clean, fmt clean, cargo doc --no-deps clean.
687 lines
25 KiB
Rust
687 lines
25 KiB
Rust
//! Full-surface integration suite (phase 4): one in-process `HttpAdapter`
|
|
//! serving the gateway endpoints, the `/openapi.json` projection, and the
|
|
//! WS channels session over real I/O; `from_openapi` imported against a
|
|
//! local HTTP echo server; `to_openapi`/`to_mcp` projections consumed
|
|
//! back. Exercises the composition the assembly layer performs.
|
|
|
|
#![cfg(feature = "mcp")]
|
|
|
|
use std::collections::HashMap;
|
|
use std::sync::Arc;
|
|
|
|
use alkcall::core::auth::{AuthContext, Identity, IdentityProvider};
|
|
use alkcall::core::types::{Capabilities, Connection};
|
|
use alkcall::protocol::wire::{EventEnvelope, ResponseEnvelope};
|
|
use alkcall::registry::discovery::{
|
|
services_list_handler, services_list_spec, services_schema_handler, services_schema_spec,
|
|
};
|
|
use alkcall::registry::registration::{
|
|
make_handler, make_streaming_handler, HandlerKind, HandlerRegistration, OperationProvenance,
|
|
OperationRegistry,
|
|
};
|
|
use alkcall::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility};
|
|
use alkhttp::adapters::FromOpenAPI;
|
|
use alkhttp::client::HttpClientConfig;
|
|
use alkhttp::server::HttpAdapter;
|
|
use alkhttp::websocket::{frame_channel0_chunk, ChunkAssembler, FrameAssembler, WsClient};
|
|
|
|
fn identity(id: &str, scopes: &[&str]) -> Identity {
|
|
Identity {
|
|
id: id.to_string(),
|
|
scopes: scopes.iter().map(|s| s.to_string()).collect(),
|
|
resources: HashMap::new(),
|
|
}
|
|
}
|
|
|
|
struct StaticTokens {
|
|
tokens: std::sync::Mutex<HashMap<String, Identity>>,
|
|
}
|
|
|
|
impl IdentityProvider for StaticTokens {
|
|
fn resolve_from_fingerprint(&self, _: &str) -> Option<Identity> {
|
|
None
|
|
}
|
|
fn resolve_from_token(&self, token: &alkcall::core::auth::AuthToken) -> Option<Identity> {
|
|
let s = String::from_utf8_lossy(&token.raw).to_string();
|
|
self.tokens
|
|
.lock()
|
|
.unwrap_or_else(|e| e.into_inner())
|
|
.get(&s)
|
|
.cloned()
|
|
}
|
|
}
|
|
|
|
fn provider_with(tokens: Vec<(&str, Identity)>) -> Arc<dyn IdentityProvider> {
|
|
let map: HashMap<String, Identity> = tokens
|
|
.into_iter()
|
|
.map(|(t, i)| (t.to_string(), i))
|
|
.collect();
|
|
Arc::new(StaticTokens {
|
|
tokens: std::sync::Mutex::new(map),
|
|
})
|
|
}
|
|
|
|
/// The local operation registry the `HttpAdapter` serves: an echo op (open
|
|
/// and echo-restricted variants), a streaming sub op, and the discovery
|
|
/// ops the adapters need.
|
|
fn local_registry() -> Arc<OperationRegistry> {
|
|
let inner = OperationRegistry::new();
|
|
inner
|
|
.register(HandlerRegistration::new(
|
|
OperationSpec::new(
|
|
"echo/run",
|
|
OperationType::Query,
|
|
Visibility::External,
|
|
serde_json::json!({}),
|
|
serde_json::json!({}),
|
|
vec![],
|
|
AccessControl::default(),
|
|
None,
|
|
),
|
|
HandlerKind::Once(make_handler(|input, ctx| async move {
|
|
ResponseEnvelope::ok(ctx.request_id, input)
|
|
})),
|
|
OperationProvenance::Local,
|
|
None,
|
|
None,
|
|
Capabilities::new(),
|
|
))
|
|
.unwrap();
|
|
inner
|
|
.register(HandlerRegistration::new(
|
|
OperationSpec::new(
|
|
"events/tick",
|
|
OperationType::Sub,
|
|
Visibility::External,
|
|
serde_json::json!({}),
|
|
serde_json::json!({}),
|
|
vec![],
|
|
AccessControl::default(),
|
|
None,
|
|
),
|
|
HandlerKind::Stream(make_streaming_handler(|input, ctx| {
|
|
futures::stream::iter(vec![
|
|
ResponseEnvelope::ok(
|
|
ctx.request_id.clone(),
|
|
serde_json::json!({ "n": 1, "input": input }),
|
|
),
|
|
ResponseEnvelope::ok(
|
|
ctx.request_id.clone(),
|
|
serde_json::json!({ "n": 2, "input": input }),
|
|
),
|
|
])
|
|
})),
|
|
OperationProvenance::Local,
|
|
None,
|
|
None,
|
|
Capabilities::new(),
|
|
))
|
|
.unwrap();
|
|
let inner = Arc::new(inner);
|
|
|
|
let registry = OperationRegistry::new();
|
|
registry
|
|
.register(HandlerRegistration::new(
|
|
services_list_spec(),
|
|
HandlerKind::Once(services_list_handler(Arc::clone(&inner))),
|
|
OperationProvenance::Local,
|
|
None,
|
|
None,
|
|
Capabilities::new(),
|
|
))
|
|
.unwrap();
|
|
registry
|
|
.register(HandlerRegistration::new(
|
|
services_schema_spec(),
|
|
HandlerKind::Once(services_schema_handler(Arc::clone(&inner))),
|
|
OperationProvenance::Local,
|
|
None,
|
|
None,
|
|
Capabilities::new(),
|
|
))
|
|
.unwrap();
|
|
for spec in inner.list_operations() {
|
|
let name = spec.name.clone();
|
|
let reg = inner.registration(&name).unwrap();
|
|
registry
|
|
.register(HandlerRegistration::new(
|
|
reg.spec.clone(),
|
|
reg.handler.clone(),
|
|
reg.provenance,
|
|
reg.composition_authority.clone(),
|
|
reg.scoped_env.clone(),
|
|
reg.capabilities.clone(),
|
|
))
|
|
.unwrap();
|
|
}
|
|
Arc::new(registry)
|
|
}
|
|
|
|
/// Serve the full adapter surface over a real TCP listener. Each accepted
|
|
/// TCP connection is wrapped as an alkcall `Connection` (single-stream,
|
|
/// `http/1.1` ALPN) and handed to the adapter's `ProtocolHandler::handle`
|
|
/// — the same path a production endpoint drives. Returns the base URL.
|
|
async fn spawn_full_server(
|
|
registry: Arc<OperationRegistry>,
|
|
provider: Arc<dyn IdentityProvider>,
|
|
) -> String {
|
|
let adapter = std::sync::Arc::new(HttpAdapter::new(Arc::clone(&provider), registry));
|
|
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let addr = listener.local_addr().unwrap();
|
|
tokio::spawn(async move {
|
|
loop {
|
|
let Ok((sock, _)) = listener.accept().await else {
|
|
break;
|
|
};
|
|
let conn = Connection::from_bidi(sock, b"http/1.1".to_vec(), None);
|
|
// The bearer middleware resolves the identity per request from
|
|
// the Authorization header; the transport-level AuthContext
|
|
// carries no identity (TLS-identity binding is the endpoint's
|
|
// job, not this test's).
|
|
let auth = AuthContext::anonymous(b"http/1.1");
|
|
let a = std::sync::Arc::clone(&adapter);
|
|
tokio::spawn(async move {
|
|
let _ =
|
|
alkcall::core::types::ProtocolHandler::handle(a.as_ref(), conn, &auth).await;
|
|
});
|
|
}
|
|
});
|
|
format!("http://{addr}")
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn full_surface_gateway_over_http() {
|
|
// A minimal registry with the discovery ops; the gateway endpoints
|
|
// must serve search/schema/call/subscribe against it over HTTP.
|
|
let registry = local_registry();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &["user"]))]);
|
|
let base = spawn_full_server(Arc::clone(®istry), Arc::clone(&provider)).await;
|
|
let client = reqwest::Client::new();
|
|
|
|
// /healthz — no auth.
|
|
let resp = client.get(format!("{base}/healthz")).send().await.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
|
|
// /search — ACL-filtered discovery.
|
|
let resp = client
|
|
.get(format!("{base}/search"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
let names: Vec<&str> = body["output"]["operations"]
|
|
.as_array()
|
|
.unwrap()
|
|
.iter()
|
|
.filter_map(|o| o["name"].as_str())
|
|
.collect();
|
|
assert!(names.contains(&"echo/run"), "got {names:?}");
|
|
assert!(names.contains(&"events/tick"), "got {names:?}");
|
|
|
|
// /call — request/response round trip.
|
|
let resp = client
|
|
.post(format!("{base}/call"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.json(&serde_json::json!({ "operation": "echo/run", "input": { "v": 42 } }))
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(body["output"]["v"], 42);
|
|
|
|
// /subscribe — SSE stream of the Sub op.
|
|
let resp = client
|
|
.post(format!("{base}/subscribe"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.json(&serde_json::json!({ "operation": "events/tick", "input": {} }))
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
let text = resp.text().await.unwrap();
|
|
assert!(text.contains("\"n\":1"), "first chunk in SSE: {text}");
|
|
assert!(text.contains("\"n\":2"), "second chunk in SSE: {text}");
|
|
|
|
// /schema — the full spec (GET with a name query param).
|
|
let resp = client
|
|
.get(format!("{base}/schema?name=echo%2Frun"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(body["output"]["name"], "echo/run");
|
|
|
|
// /openapi.json — the 6-endpoint projection including /publish.
|
|
let resp = client
|
|
.get(format!("{base}/openapi.json"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(body["info"]["version"], "1.4.0");
|
|
assert!(body["paths"].get("/publish").is_some());
|
|
assert!(body["paths"].get("/call").is_some());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn full_surface_ws_call_round_trip() {
|
|
let registry = local_registry();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &["user"]))]);
|
|
let base = spawn_full_server(registry, provider).await;
|
|
// WS endpoint rides the same TCP listener.
|
|
let ws_base = base.replacen("http://", "ws://", 1);
|
|
|
|
let mut ws = WsClient::connect_authorized(&format!("{ws_base}/alk/channels"), "tok-1")
|
|
.await
|
|
.unwrap();
|
|
let frame = EventEnvelope::requested(
|
|
"ws-full-1",
|
|
serde_json::json!({ "operationId": "echo/run", "input": { "v": 7 } }),
|
|
);
|
|
ws.send_binary(frame_channel0_chunk(&frame)).await;
|
|
|
|
let mut chunks = ChunkAssembler::new();
|
|
let mut frames = FrameAssembler::new();
|
|
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
|
|
loop {
|
|
assert!(tokio::time::Instant::now() < deadline, "timed out");
|
|
if let Some(env) = frames.next_frame() {
|
|
assert_eq!(env.r#type, "call.responded");
|
|
assert_eq!(env.id, "ws-full-1");
|
|
assert_eq!(env.payload["output"]["v"], 7);
|
|
break;
|
|
}
|
|
let bin = ws.next_binary(std::time::Duration::from_millis(500)).await;
|
|
match bin {
|
|
Some(bytes) => {
|
|
chunks.push(&bytes);
|
|
while let Some((channel_id, payload)) = chunks.next_chunk() {
|
|
assert_eq!(channel_id, 0);
|
|
frames.push(&payload);
|
|
}
|
|
}
|
|
None => panic!("ws closed unexpectedly"),
|
|
}
|
|
}
|
|
ws.close().await;
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn from_openapi_import_then_gateway_call() {
|
|
// A local HTTP service the adapter imports; the gateway dispatch then
|
|
// reaches it through the imported forwarding handler.
|
|
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let addr = listener.local_addr().unwrap();
|
|
let upstream = format!("http://{addr}");
|
|
tokio::spawn(async move {
|
|
let app = axum::Router::new().route(
|
|
"/widgets",
|
|
axum::routing::get(|| async {
|
|
axum::Json(serde_json::json!({ "widgets": ["a", "b"] }))
|
|
}),
|
|
);
|
|
axum::serve(listener, app).await.unwrap();
|
|
});
|
|
|
|
let doc = r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"},"paths":{"/widgets":{"get":{"operationId":"listWidgets","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}}}"#;
|
|
let spec = alkhttp::adapters::OpenAPISpec::from_json(doc).unwrap();
|
|
let config = alkhttp::adapters::HttpServiceConfig {
|
|
namespace: "upstream".to_string(),
|
|
base_url: upstream.clone(),
|
|
auth: None,
|
|
default_headers: HashMap::new(),
|
|
};
|
|
let http_client =
|
|
Arc::new(alkhttp::client::SharedHttpClient::new(HttpClientConfig::default()).unwrap());
|
|
let adapter = FromOpenAPI::new(spec, config, http_client);
|
|
let bundles = alkcall::client::OperationAdapter::import(&adapter)
|
|
.await
|
|
.expect("import succeeds");
|
|
assert_eq!(bundles.len(), 1);
|
|
assert_eq!(bundles[0].spec.name, "upstream/listWidgets");
|
|
|
|
// The imported bundles are Internal (ADR-015 — composition material);
|
|
// a wire call to them is NOT_FOUND (ADR-015 §2). The assembly layer
|
|
// composes them under an External facade. Verify Internal-not-callable
|
|
// through the gateway, then compose the External facade and call that.
|
|
let registry = OperationRegistry::new();
|
|
for b in bundles {
|
|
registry.register(b).unwrap();
|
|
}
|
|
registry
|
|
.register(HandlerRegistration::new(
|
|
OperationSpec::new(
|
|
"widgets/list",
|
|
OperationType::Query,
|
|
Visibility::External,
|
|
serde_json::json!({}),
|
|
serde_json::json!({}),
|
|
vec![],
|
|
AccessControl::default(),
|
|
None,
|
|
),
|
|
HandlerKind::Once(make_handler(|_input, ctx| {
|
|
// The facade composes the imported leaf (env.invoke —
|
|
// composition-only per ADR-015). Its scoped_env declares
|
|
// the imported leaf as the reachable set.
|
|
async move {
|
|
let response = ctx
|
|
.env
|
|
.invoke("upstream", "listWidgets", serde_json::json!({}), &ctx)
|
|
.await;
|
|
ResponseEnvelope {
|
|
request_id: ctx.request_id,
|
|
result: response.result,
|
|
}
|
|
}
|
|
})),
|
|
OperationProvenance::Local,
|
|
None,
|
|
Some(alkcall::registry::context::ScopedPeerEnv::new([
|
|
"upstream/listWidgets",
|
|
])),
|
|
Capabilities::new(),
|
|
))
|
|
.unwrap();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &[]))]);
|
|
let base = spawn_full_server(Arc::new(registry), provider).await;
|
|
|
|
let client = reqwest::Client::new();
|
|
|
|
// Internal op from the wire → NOT_FOUND (does not leak existence).
|
|
let resp = client
|
|
.post(format!("{base}/call"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.json(&serde_json::json!({ "operation": "upstream/listWidgets", "input": {} }))
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(
|
|
resp.status(),
|
|
404,
|
|
"Internal imported op is invisible from the wire"
|
|
);
|
|
|
|
// The External facade composes it: external HTTP API → from_openapi
|
|
// forwarding handler → upstream HTTP API.
|
|
let resp = client
|
|
.post(format!("{base}/call"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.json(&serde_json::json!({ "operation": "widgets/list", "input": {} }))
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200, "facade composes the imported op");
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(body["output"]["widgets"], serde_json::json!(["a", "b"]));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn to_openapi_and_to_mcp_projections_over_served_registry() {
|
|
use rmcp::model::{CallToolRequestParams, ClientCapabilities, ClientInfo, Implementation};
|
|
use rmcp::service::RoleClient;
|
|
use rmcp::transport::streamable_http_client::{
|
|
StreamableHttpClientTransport, StreamableHttpClientTransportConfig,
|
|
};
|
|
use rmcp::{Peer, ServiceExt};
|
|
|
|
let registry = local_registry();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &["user"]))]);
|
|
|
|
// /openapi.json served from the same registry: the to_openapi
|
|
// projection sees the local ops through services/list.
|
|
let base = spawn_full_server(Arc::clone(®istry), Arc::clone(&provider)).await;
|
|
let client = reqwest::Client::new();
|
|
let resp = client
|
|
.get(format!("{base}/openapi.json"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 200);
|
|
let doc: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(doc["info"]["title"], "alk gateway");
|
|
assert_eq!(doc["paths"].as_object().unwrap().len(), 6);
|
|
// Validates against openapiv3 (ADR-042 contract).
|
|
let text = serde_json::to_string(&doc).unwrap();
|
|
let _: openapiv3::OpenAPI = serde_json::from_str(&text).unwrap();
|
|
|
|
// /mcp served by the same adapter: an MCP client connects, lists the
|
|
// 4 gateway tools, calls search — the tool-gateway pattern (ADR-041).
|
|
let url = format!("{base}/mcp");
|
|
let transport = StreamableHttpClientTransport::from_config(
|
|
StreamableHttpClientTransportConfig::with_uri(url),
|
|
);
|
|
let client_info = ClientInfo::new(
|
|
ClientCapabilities::default(),
|
|
Implementation::new("integration-test", "0.1.0"),
|
|
);
|
|
let running = client_info.serve(transport).await.expect("initialize");
|
|
let peer: Peer<RoleClient> = running.peer().clone();
|
|
|
|
let tools = peer
|
|
.list_tools(Default::default())
|
|
.await
|
|
.expect("tools/list");
|
|
let names: Vec<String> = tools.tools.iter().map(|t| t.name.to_string()).collect();
|
|
assert_eq!(names.len(), 4);
|
|
assert!(names.contains(&"search".to_string()));
|
|
assert!(names.contains(&"schema".to_string()));
|
|
assert!(names.contains(&"call".to_string()));
|
|
assert!(names.contains(&"batch".to_string()));
|
|
|
|
let mut args = serde_json::Map::new();
|
|
args.insert("query".to_string(), serde_json::Value::Null);
|
|
let params = CallToolRequestParams::new("search".to_string()).with_arguments(args);
|
|
let result = peer.call_tool(params).await.expect("search call");
|
|
assert_eq!(result.is_error, Some(false));
|
|
let structured = result.structured_content.expect("structured present");
|
|
let ops = structured
|
|
.get("operations")
|
|
.and_then(serde_json::Value::as_array)
|
|
.expect("operations array");
|
|
let names: Vec<&str> = ops
|
|
.iter()
|
|
.filter_map(|o| o.get("name").and_then(|v| v.as_str()))
|
|
.collect();
|
|
assert!(names.contains(&"echo/run"), "got {names:?}");
|
|
assert!(
|
|
!names.contains(&"events/tick"),
|
|
"Sub ops excluded from search"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn gateway_error_fidelity_end_to_end() {
|
|
// Unknown op through /call → 404 NOT_FOUND; internal op → 404.
|
|
let registry = local_registry();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &["user"]))]);
|
|
let base = spawn_full_server(registry, provider).await;
|
|
let client = reqwest::Client::new();
|
|
|
|
let resp = client
|
|
.post(format!("{base}/call"))
|
|
.header("Authorization", "Bearer tok-1")
|
|
.json(&serde_json::json!({ "operation": "no/such", "input": {} }))
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(resp.status(), 404);
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(body["code"], "NOT_FOUND");
|
|
}
|
|
|
|
use http::header::AUTHORIZATION;
|
|
|
|
#[tokio::test]
|
|
async fn gateway_endpoints_exist_with_bearer_enforcement() {
|
|
let registry = local_registry();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &["user"]))]);
|
|
let base = spawn_full_server(registry, provider).await;
|
|
let client = reqwest::Client::new();
|
|
|
|
// GET endpoints: /search, /schema (schema needs a known name).
|
|
let resp = client
|
|
.get(format!("{base}/search"))
|
|
.header(AUTHORIZATION, "Bearer tok-1")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_ne!(
|
|
resp.status(),
|
|
404,
|
|
"/search must exist on the gateway surface"
|
|
);
|
|
let resp = client
|
|
.get(format!("{base}/schema?name=echo/run"))
|
|
.header(AUTHORIZATION, "Bearer tok-1")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_ne!(
|
|
resp.status(),
|
|
404,
|
|
"/schema must exist on the gateway surface"
|
|
);
|
|
for path in ["/call", "/batch", "/subscribe", "/publish"] {
|
|
let resp = client
|
|
.post(format!("{base}{path}"))
|
|
.header(AUTHORIZATION, "Bearer tok-1")
|
|
.header("Content-Type", "application/json")
|
|
.body("{}")
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_ne!(
|
|
resp.status(),
|
|
404,
|
|
"{path} must exist on the gateway surface"
|
|
);
|
|
}
|
|
|
|
// Unauthenticated call to an op with no restrictions: allowed
|
|
// (AccessControl::default() passes for any identity, including none).
|
|
let resp = client
|
|
.post(format!("{base}/call"))
|
|
.header("Content-Type", "application/json")
|
|
.body(r#"{"operation": "no/such", "input": {}}"#)
|
|
.send()
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(
|
|
resp.status(),
|
|
404,
|
|
"unknown op is NOT_FOUND regardless of auth"
|
|
);
|
|
}
|
|
|
|
/// COV-11b: the rmcp-entered `ServerHandler::call_tool` routing shell is
|
|
/// what production runs; the dispatch-level tests below it bypass it.
|
|
/// One real rmcp-protocol `peer.call_tool` round-trip per gateway tool
|
|
/// through the served `/mcp` mount.
|
|
#[tokio::test]
|
|
async fn to_mcp_call_tool_production_dispatch_round_trips_all_tools() {
|
|
use rmcp::model::{CallToolRequestParams, ClientCapabilities, ClientInfo, Implementation};
|
|
use rmcp::service::RoleClient;
|
|
use rmcp::transport::streamable_http_client::{
|
|
StreamableHttpClientTransport, StreamableHttpClientTransportConfig,
|
|
};
|
|
use rmcp::{Peer, ServiceExt};
|
|
|
|
let registry = local_registry();
|
|
let provider = provider_with(vec![("tok-1", identity("alice", &["user"]))]);
|
|
let base = spawn_full_server(registry, provider).await;
|
|
|
|
let mut default_headers = reqwest::header::HeaderMap::new();
|
|
default_headers.insert(
|
|
reqwest::header::AUTHORIZATION,
|
|
reqwest::header::HeaderValue::from_static("Bearer tok-1"),
|
|
);
|
|
let http = reqwest::Client::builder()
|
|
.default_headers(default_headers)
|
|
.build()
|
|
.unwrap();
|
|
let url = format!("{base}/mcp");
|
|
let transport = StreamableHttpClientTransport::with_client(
|
|
http,
|
|
StreamableHttpClientTransportConfig::with_uri(url),
|
|
);
|
|
let client_info = ClientInfo::new(
|
|
ClientCapabilities::default(),
|
|
Implementation::new("integration-test", "0.1.0"),
|
|
);
|
|
let running = client_info.serve(transport).await.expect("initialize");
|
|
let peer: Peer<RoleClient> = running.peer().clone();
|
|
|
|
let schema_params = CallToolRequestParams::new("schema".to_string()).with_arguments(
|
|
serde_json::json!({ "name": "echo/run" })
|
|
.as_object()
|
|
.unwrap()
|
|
.clone(),
|
|
);
|
|
let schema = peer.call_tool(schema_params).await.expect("schema call");
|
|
assert_eq!(schema.is_error, Some(false));
|
|
let structured = schema.structured_content.expect("structured present");
|
|
assert_eq!(structured["name"], "echo/run");
|
|
assert!(structured.get("input_schema").is_some());
|
|
|
|
let call_params = CallToolRequestParams::new("call".to_string()).with_arguments(
|
|
serde_json::json!({ "operation": "echo/run", "input": { "v": 7 } })
|
|
.as_object()
|
|
.unwrap()
|
|
.clone(),
|
|
);
|
|
let call_result = peer.call_tool(call_params).await.expect("call call");
|
|
assert_eq!(call_result.is_error, Some(false));
|
|
assert_eq!(
|
|
call_result.structured_content,
|
|
Some(serde_json::json!({ "v": 7 }))
|
|
);
|
|
|
|
let batch_params = CallToolRequestParams::new("batch".to_string()).with_arguments(
|
|
serde_json::json!({ "calls": [
|
|
{ "operation": "echo/run", "input": { "n": 1 } },
|
|
{ "operation": "echo/run", "input": { "n": 2 } }
|
|
] })
|
|
.as_object()
|
|
.unwrap()
|
|
.clone(),
|
|
);
|
|
let batch = peer.call_tool(batch_params).await.expect("batch call");
|
|
assert_eq!(batch.is_error, Some(false));
|
|
let results = batch
|
|
.structured_content
|
|
.and_then(|v| v.get("results").cloned())
|
|
.expect("results array");
|
|
assert_eq!(
|
|
results,
|
|
serde_json::json!([
|
|
{ "isError": false, "output": { "n": 1 } },
|
|
{ "isError": false, "output": { "n": 2 } }
|
|
])
|
|
);
|
|
|
|
let unknown = peer
|
|
.call_tool(CallToolRequestParams::new("bogus".to_string()))
|
|
.await
|
|
.expect("unknown tool call resolves, not errors");
|
|
assert_eq!(unknown.is_error, Some(true));
|
|
let err = unknown
|
|
.structured_content
|
|
.expect("structured error present");
|
|
assert_eq!(err["code"], "NOT_FOUND");
|
|
assert!(err["message"]
|
|
.as_str()
|
|
.unwrap_or_default()
|
|
.contains("unknown gateway tool"));
|
|
|
|
let _ = running.cancel().await;
|
|
}
|