Skip to content

Commit 65a33b7

Browse files
committed
feat(gateway): resolve bindings from X-Mcpmux-Machine-Id header
Parse the per-device machine header on each MCP request and use it as the top-priority Tier 1 signal when present, so tunneled callers are not mistaken for the gateway host. Add integration tests for header outranking and deny. Signed-off-by: Joe Sangiorgio <jsangio1@gmail.com>
1 parent 1e7c51f commit 65a33b7

11 files changed

Lines changed: 229 additions & 61 deletions

File tree

crates/mcpmux-gateway/src/consumers/mcp_notifier.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -936,7 +936,7 @@ impl MCPNotifier {
936936
}
937937
match self
938938
.feature_set_resolver
939-
.resolve(Some(&session_id), Some(&client_id))
939+
.resolve(Some(&session_id), Some(&client_id), None)
940940
.await
941941
{
942942
Ok(resolved) if resolved.space_id == Some(space_id) => {

crates/mcpmux-gateway/src/mcp/context.rs

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,11 +5,15 @@ use http;
55
use rmcp::{model::Extensions, service::RequestContext, RoleServer};
66
use uuid::Uuid;
77

8-
/// OAuth claims extracted from JWT token
8+
/// OAuth claims extracted from JWT token, plus optional per-device machine identity.
99
#[derive(Debug, Clone)]
1010
pub struct OAuthContext {
1111
pub client_id: String,
1212
pub space_id: Uuid,
13+
/// Physical device identity from the client's `X-Mcpmux-Machine-Id` header.
14+
/// Distinct from the gateway's `local_machine_id` when callers reach a shared
15+
/// tunneled gateway from multiple machines.
16+
pub request_machine_id: Option<Uuid>,
1317
}
1418

1519
/// Extract OAuth context from extensions
@@ -46,9 +50,16 @@ pub fn extract_oauth_context(extensions: &Extensions) -> Result<OAuthContext> {
4650
let space_id =
4751
Uuid::parse_str(space_id_str).map_err(|e| anyhow!("Failed to parse space_id: {}", e))?;
4852

53+
let request_machine_id = parts
54+
.headers
55+
.get("x-mcpmux-machine-id")
56+
.and_then(|v| v.to_str().ok())
57+
.and_then(|s| Uuid::parse_str(s).ok());
58+
4959
Ok(OAuthContext {
5060
client_id,
5161
space_id,
62+
request_machine_id,
5263
})
5364
}
5465

crates/mcpmux-gateway/src/mcp/handler.rs

Lines changed: 45 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -101,9 +101,12 @@ impl McpMuxGatewayHandler {
101101
client_id: &str,
102102
session_id: Option<&str>,
103103
root_for_prompt: Option<&str>,
104+
request_machine_id: Option<uuid::Uuid>,
104105
) {
105106
let resolver = &services.feature_set_resolver;
106-
match resolver.resolve(session_id, Some(client_id)).await {
107+
match resolver
108+
.resolve(session_id, Some(client_id), request_machine_id)
109+
.await {
107110
Ok(resolved) => {
108111
info!(
109112
%client_id,
@@ -192,11 +195,12 @@ impl McpMuxGatewayHandler {
192195
&self,
193196
session_id: Option<&str>,
194197
client_id: &str,
198+
request_machine_id: Option<uuid::Uuid>,
195199
) -> Result<(uuid::Uuid, Vec<String>), McpError> {
196200
let resolved = self
197201
.services
198202
.authorization_service
199-
.resolve(session_id, Some(client_id))
203+
.resolve(session_id, Some(client_id), request_machine_id)
200204
.await
201205
.map_err(|e| McpError::internal_error(format!("Failed to resolve: {e}"), None))?;
202206
let space_id = resolved.space_id.ok_or_else(|| {
@@ -226,6 +230,7 @@ impl McpMuxGatewayHandler {
226230
peer: &rmcp::service::Peer<RoleServer>,
227231
session_id: Option<&str>,
228232
client_id: &str,
233+
request_machine_id: Option<uuid::Uuid>,
229234
) {
230235
let Some(sid) = session_id else { return };
231236
// Fast path: already have a definitive answer (Some(roots),
@@ -318,6 +323,7 @@ impl McpMuxGatewayHandler {
318323
&client_id,
319324
Some(&session_id),
320325
root_for_prompt.as_deref(),
326+
request_machine_id,
321327
)
322328
.await;
323329
});
@@ -532,6 +538,7 @@ impl ServerHandler for McpMuxGatewayHandler {
532538
let notifier = self.notification_bridge.clone();
533539
let client_id_str = oauth_ctx.client_id.clone();
534540
let session_id_for_task = session_id.clone();
541+
let request_machine_id = oauth_ctx.request_machine_id;
535542
tokio::spawn(async move {
536543
// Retry list_roots() on transport errors with bounded
537544
// backoff. Without roots a roots-capable session is
@@ -620,6 +627,7 @@ impl ServerHandler for McpMuxGatewayHandler {
620627
&client_id_str,
621628
Some(&session_id_for_task),
622629
root_for_prompt.as_deref(),
630+
request_machine_id,
623631
)
624632
.await;
625633
});
@@ -632,6 +640,7 @@ impl ServerHandler for McpMuxGatewayHandler {
632640
&oauth_ctx.client_id,
633641
Some(&session_id),
634642
None,
643+
oauth_ctx.request_machine_id,
635644
)
636645
.await;
637646
}
@@ -671,6 +680,7 @@ impl ServerHandler for McpMuxGatewayHandler {
671680
let notifier = self.notification_bridge.clone();
672681
let client_id_str = oauth_ctx.client_id.clone();
673682
let session_id_for_task = session_id.clone();
683+
let request_machine_id = oauth_ctx.request_machine_id;
674684
tokio::spawn(async move {
675685
match peer.list_roots().await {
676686
Ok(result) => {
@@ -702,6 +712,7 @@ impl ServerHandler for McpMuxGatewayHandler {
702712
&client_id_str,
703713
Some(&session_id_for_task),
704714
root_for_prompt.as_deref(),
715+
request_machine_id,
705716
)
706717
.await;
707718
}
@@ -734,13 +745,18 @@ impl ServerHandler for McpMuxGatewayHandler {
734745
&context.peer,
735746
session_id_owned.as_deref(),
736747
&oauth_ctx.client_id,
748+
oauth_ctx.request_machine_id,
737749
)
738750
.await;
739751
// Resolve routing once: the resolver returns the authoritative
740752
// (Space, FS) for this session — this may differ from oauth_ctx
741753
// when a WorkspaceBinding redirects to another space.
742754
let (space_id, feature_set_ids) = self
743-
.resolve_routing(session_id_owned.as_deref(), &oauth_ctx.client_id)
755+
.resolve_routing(
756+
session_id_owned.as_deref(),
757+
&oauth_ctx.client_id,
758+
oauth_ctx.request_machine_id,
759+
)
744760
.await?;
745761

746762
// Get advertised (surfaced) tools only — full invokable set is reachable
@@ -810,14 +826,14 @@ impl ServerHandler for McpMuxGatewayHandler {
810826
// connection). Without the probe it resolves to empty FS ids and
811827
// fails "not allowed by the current grants" — breaking the
812828
// list==call invariant the list handlers already uphold.
813-
self.ensure_roots_probed(&context.peer, session_id, &oauth_ctx.client_id)
829+
self.ensure_roots_probed(&context.peer, session_id, &oauth_ctx.client_id, oauth_ctx.request_machine_id)
814830
.await;
815831

816832
// Resolve routing once — the binding's target space is authoritative
817833
// (may differ from oauth_ctx.space_id). Needed both to gate the
818834
// per-Space meta tools below and to route a normal tool call.
819835
let (space_id, feature_set_ids) = self
820-
.resolve_routing(session_id, &oauth_ctx.client_id)
836+
.resolve_routing(session_id, &oauth_ctx.client_id, oauth_ctx.request_machine_id)
821837
.await?;
822838

823839
// Intercept meta tools (mcpmux_*) BEFORE feature-set filtering, gated
@@ -1032,10 +1048,15 @@ impl ServerHandler for McpMuxGatewayHandler {
10321048
&context.peer,
10331049
session_id_owned.as_deref(),
10341050
&oauth_ctx.client_id,
1051+
oauth_ctx.request_machine_id,
10351052
)
10361053
.await;
10371054
let (space_id, feature_set_ids) = self
1038-
.resolve_routing(session_id_owned.as_deref(), &oauth_ctx.client_id)
1055+
.resolve_routing(
1056+
session_id_owned.as_deref(),
1057+
&oauth_ctx.client_id,
1058+
oauth_ctx.request_machine_id,
1059+
)
10391060
.await?;
10401061

10411062
// Get advertised (surfaced) prompts only — full fetchable set is reachable
@@ -1086,10 +1107,15 @@ impl ServerHandler for McpMuxGatewayHandler {
10861107
&context.peer,
10871108
session_id_owned.as_deref(),
10881109
&oauth_ctx.client_id,
1110+
oauth_ctx.request_machine_id,
10891111
)
10901112
.await;
10911113
let (space_id, feature_set_ids) = self
1092-
.resolve_routing(session_id_owned.as_deref(), &oauth_ctx.client_id)
1114+
.resolve_routing(
1115+
session_id_owned.as_deref(),
1116+
&oauth_ctx.client_id,
1117+
oauth_ctx.request_machine_id,
1118+
)
10931119
.await?;
10941120

10951121
let (server_id, prompt_name) = self
@@ -1180,10 +1206,15 @@ impl ServerHandler for McpMuxGatewayHandler {
11801206
&context.peer,
11811207
session_id_owned.as_deref(),
11821208
&oauth_ctx.client_id,
1209+
oauth_ctx.request_machine_id,
11831210
)
11841211
.await;
11851212
let (space_id, feature_set_ids) = self
1186-
.resolve_routing(session_id_owned.as_deref(), &oauth_ctx.client_id)
1213+
.resolve_routing(
1214+
session_id_owned.as_deref(),
1215+
&oauth_ctx.client_id,
1216+
oauth_ctx.request_machine_id,
1217+
)
11871218
.await?;
11881219

11891220
// Get advertised (surfaced) resources only — full readable set is reachable
@@ -1232,10 +1263,15 @@ impl ServerHandler for McpMuxGatewayHandler {
12321263
&context.peer,
12331264
session_id_owned.as_deref(),
12341265
&oauth_ctx.client_id,
1266+
oauth_ctx.request_machine_id,
12351267
)
12361268
.await;
12371269
let (space_id, feature_set_ids) = self
1238-
.resolve_routing(session_id_owned.as_deref(), &oauth_ctx.client_id)
1270+
.resolve_routing(
1271+
session_id_owned.as_deref(),
1272+
&oauth_ctx.client_id,
1273+
oauth_ctx.request_machine_id,
1274+
)
12391275
.await?;
12401276

12411277
let authorized_resources = self

crates/mcpmux-gateway/src/services/authorization.rs

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,10 @@ impl AuthorizationService {
3333
_space_id: &Uuid,
3434
session_id: Option<&str>,
3535
) -> Result<Vec<String>> {
36-
let resolved = self.resolver.resolve(session_id, Some(client_id)).await?;
36+
let resolved = self
37+
.resolver
38+
.resolve(session_id, Some(client_id), None)
39+
.await?;
3740
Ok(resolved.feature_set_ids)
3841
}
3942

@@ -44,8 +47,11 @@ impl AuthorizationService {
4447
&self,
4548
session_id: Option<&str>,
4649
client_id: Option<&str>,
50+
request_machine_id: Option<Uuid>,
4751
) -> Result<ResolvedFeatureSet> {
48-
self.resolver.resolve(session_id, client_id).await
52+
self.resolver
53+
.resolve(session_id, client_id, request_machine_id)
54+
.await
4955
}
5056

5157
/// Does this session/client resolve to any FeatureSet?

crates/mcpmux-gateway/src/services/feature_set_resolver.rs

Lines changed: 25 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -30,8 +30,11 @@
3030
//! return (default_space, [], Unbound)
3131
//! ```
3232
//!
33-
//! Signal 3 (machine) scopes binding lookup: client machine → local machine
34-
//! → global. A path bound only on another machine does not match here.
33+
//! Signal 3 (machine) scopes binding lookup. When the client sends
34+
//! `X-Mcpmux-Machine-Id`, only that machine and global bindings are
35+
//! considered — gateway `local_machine_id` and `inbound_clients.machine_id`
36+
//! are skipped so a tunneled caller is not mistaken for the gateway host.
37+
//! Without the header: client machine → local machine → global.
3538
//!
3639
//! ## Deny by default
3740
//!
@@ -225,13 +228,28 @@ impl FeatureSetResolverService {
225228
*self.local_machine_id.write().await = id;
226229
}
227230

228-
/// Tier 1 exact binding lookup: client machine, then local machine, then global.
231+
/// Tier 1 exact binding lookup: request machine header, then client machine,
232+
/// then local machine, then global.
229233
async fn find_binding_for_roots(
230234
&self,
231235
roots: &[String],
232236
client_id: Option<&str>,
237+
request_machine_id: Option<Uuid>,
233238
) -> Result<Option<mcpmux_core::WorkspaceBinding>> {
234239
for root in roots {
240+
if let Some(header_machine) = request_machine_id {
241+
if let Some(binding) = self
242+
.binding_repo
243+
.find_exact_for_machine(&header_machine, root, client_id)
244+
.await?
245+
{
246+
return Ok(Some(binding));
247+
}
248+
if let Some(binding) = self.binding_repo.find_exact_global(root).await? {
249+
return Ok(Some(binding));
250+
}
251+
continue;
252+
}
235253
if let Some(cid) = client_id {
236254
if let Some(client_machine) = self.client_repo.get_machine_id(cid).await? {
237255
if let Some(binding) = self
@@ -324,10 +342,13 @@ impl FeatureSetResolverService {
324342
/// the caller is stateless — e.g. desktop UI HTTP path).
325343
/// `client_id`: the OAuth client identity. Used only for the Tier-2
326344
/// `client_grants` lookup; ignored for binding-based routing.
345+
/// `request_machine_id`: optional per-device identity from
346+
/// `X-Mcpmux-Machine-Id`; highest-priority machine signal for Tier 1.
327347
pub async fn resolve(
328348
&self,
329349
session_id: Option<&str>,
330350
client_id: Option<&str>,
351+
request_machine_id: Option<Uuid>,
331352
) -> Result<ResolvedFeatureSet> {
332353
let default_space_id = match self.space_repo.get_default().await? {
333354
Some(s) => s.id,
@@ -376,7 +397,7 @@ impl FeatureSetResolverService {
376397
if has_roots {
377398
let reported_roots = roots.expect("has_roots implies Some");
378399
if let Some(binding) = self
379-
.find_binding_for_roots(&reported_roots, client_id)
400+
.find_binding_for_roots(&reported_roots, client_id, request_machine_id)
380401
.await?
381402
{
382403
debug!(

crates/mcpmux-gateway/src/services/meta_tools/meta_tool_common.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ pub(crate) async fn caller_space_id(call: &MetaToolCall<'_>) -> Result<Uuid, Met
5858
let resolved = call
5959
.ctx
6060
.resolver
61-
.resolve(call.session_id, Some(call.client_id))
61+
.resolve(call.session_id, Some(call.client_id), None)
6262
.await?;
6363
if let Some(space_id) = resolved.space_id {
6464
return Ok(space_id);
@@ -76,7 +76,7 @@ pub(crate) async fn caller_resolution(
7676
) -> Result<ResolvedFeatureSet, MetaToolError> {
7777
call.ctx
7878
.resolver
79-
.resolve(call.session_id, Some(call.client_id))
79+
.resolve(call.session_id, Some(call.client_id), None)
8080
.await
8181
.map_err(|e| MetaToolError::Internal(e.to_string()))
8282
}

crates/mcpmux-gateway/src/services/meta_tools/set_workspace_root.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ impl MetaTool for SetWorkspaceRootTool {
7777
let resolved = call
7878
.ctx
7979
.resolver
80-
.resolve(Some(session_id), Some(call.client_id))
80+
.resolve(Some(session_id), Some(call.client_id), None)
8181
.await?;
8282

8383
let space_id = resolved

tests/rust/tests/integration/effective_features.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -152,7 +152,7 @@ impl Ctx {
152152
/// The tools a session actually sees — the exact two-step the request
153153
/// handler runs: resolve the mapping, then pull its FS-filtered tools.
154154
async fn effective_tools(&self, session_id: &str) -> Vec<String> {
155-
let resolved = self.resolver.resolve(Some(session_id), None).await.unwrap();
155+
let resolved = self.resolver.resolve(Some(session_id), None, None).await.unwrap();
156156
let tools = self
157157
.feature_service
158158
.get_tools_for_grants(&self.space_id_str, &resolved.feature_set_ids)
@@ -263,7 +263,7 @@ async fn empty_mapping_yields_zero_effective_tools() {
263263
ctx.session_roots.set("sess", [root]);
264264
ctx.session_roots.set_roots_capable("sess", true);
265265

266-
let resolved = ctx.resolver.resolve(Some("sess"), None).await.unwrap();
266+
let resolved = ctx.resolver.resolve(Some("sess"), None, None).await.unwrap();
267267
assert_eq!(resolved.source, ResolutionSource::WorkspaceBinding);
268268
assert!(resolved.feature_set_ids.is_empty());
269269
assert!(ctx.effective_tools("sess").await.is_empty());
@@ -296,7 +296,7 @@ async fn unbound_session_returns_no_tools() {
296296
ctx.session_roots.set("sess", [root]);
297297
ctx.session_roots.set_roots_capable("sess", true);
298298

299-
let resolved = ctx.resolver.resolve(Some("sess"), None).await.unwrap();
299+
let resolved = ctx.resolver.resolve(Some("sess"), None, None).await.unwrap();
300300
assert_eq!(resolved.source, ResolutionSource::Unbound);
301301
assert!(resolved.feature_set_ids.is_empty());
302302
assert!(ctx.effective_tools("sess").await.is_empty());
@@ -332,7 +332,7 @@ async fn empty_starter_grants_nothing_to_unbound_session() {
332332
ctx.session_roots.set("sess", [root]);
333333
ctx.session_roots.set_roots_capable("sess", true);
334334

335-
let resolved = ctx.resolver.resolve(Some("sess"), None).await.unwrap();
335+
let resolved = ctx.resolver.resolve(Some("sess"), None, None).await.unwrap();
336336
assert_eq!(resolved.source, ResolutionSource::Unbound);
337337
assert!(resolved.feature_set_ids.is_empty());
338338
assert!(ctx.effective_tools("sess").await.is_empty());

0 commit comments

Comments
 (0)