Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions crates/rmcp/src/handler/server.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -89,6 +89,9 @@ impl<H: ServerHandler> Service<RoleServer> for H {
ClientNotification::RootsListChangedNotification(_notification) => {
self.on_roots_list_changed(context).await
}
ClientNotification::CustomClientNotification(notification) => {
self.on_custom_notification(notification, context).await
}
};
Ok(())
}
Expand DownExpand Up@@ -224,6 +227,14 @@ pub trait ServerHandler: Sized + Send + Sync + 'static {
) -> impl Future<Output = ()> + Send + '_ {
std::future::ready(())
}
fn on_custom_notification(
&self,
notification: CustomClientNotification,
context: NotificationContext<RoleServer>,
) -> impl Future<Output = ()> + Send + '_ {
let _ = (notification, context);
std::future::ready(())
}

fn get_info(&self) -> ServerInfo {
ServerInfo::default()
Expand Down
69 changes: 68 additions & 1 deletion crates/rmcp/src/model.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -627,6 +627,40 @@ const_string!(CancelledNotificationMethod = "notifications/cancelled");
pub type CancelledNotification =
Notification<CancelledNotificationMethod, CancelledNotificationParam>;

/// A catch-all notification the client can use to send custom messages to a server.
///
/// This preserves the raw `method` name and `params` payload so handlers can
/// deserialize them into domain-specific types.
#[derive(Debug, Clone)]
#[cfg_attr(feature = "schemars", derive(schemars::JsonSchema))]
pub struct CustomClientNotification {
pub method: String,
pub params: Option<Value>,
/// extensions will carry anything possible in the context, including [`Meta`]
///
/// this is similar with the Extensions in `http` crate
#[cfg_attr(feature = "schemars", schemars(skip))]
pub extensions: Extensions,
}

impl CustomClientNotification {
pub fn new(method: impl Into<String>, params: Option<Value>) -> Self {
Self {
method: method.into(),
params,
extensions: Extensions::default(),
}
}

/// Deserialize `params` into a strongly-typed structure.
pub fn params_as<T: DeserializeOwned>(&self) -> Result<Option<T>, serde_json::Error> {
self.params
.as_ref()
.map(|params| serde_json::from_value(params.clone()))
.transpose()
}
}

const_string!(InitializeResultMethod = "initialize");
/// # Initialization
/// This request is sent from the client to the server when it first connects, asking it to begin initialization.
Expand DownExpand Up@@ -1748,7 +1782,8 @@ ts_union!(
| CancelledNotification
| ProgressNotification
| InitializedNotification
| RootsListChangedNotification;
| RootsListChangedNotification
| CustomClientNotification;
);

ts_union!(
Expand DownExpand Up@@ -1857,6 +1892,38 @@ mod tests {
assert_eq!(json, raw);
}

#[test]
fn test_custom_client_notification_roundtrip() {
let raw = json!( {
"jsonrpc": JsonRpcVersion2_0,
"method": "notifications/custom",
"params": {"foo": "bar"},
});

let message: ClientJsonRpcMessage =
serde_json::from_value(raw.clone()).expect("invalid notification");
match &message {
ClientJsonRpcMessage::Notification(JsonRpcNotification {
notification: ClientNotification::CustomClientNotification(notification),
..
}) => {
assert_eq!(notification.method, "notifications/custom");
assert_eq!(
notification
.params
.as_ref()
.and_then(|p| p.get("foo"))
.expect("foo present"),
"bar"
);
}
_ => panic!("Expected custom client notification"),
}

let json = serde_json::to_value(message).expect("valid json");
assert_eq!(json, raw);
}

#[test]
fn test_request_conversion() {
let raw = json!( {
Expand Down
25 changes: 23 additions & 2 deletions crates/rmcp/src/model/meta.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;

use super::{
ClientNotification, ClientRequest, Extensions, JsonObject, JsonRpcMessage, NumberOrString,
ProgressToken, ServerNotification, ServerRequest,
ClientNotification, ClientRequest, CustomClientNotification, Extensions, JsonObject,
JsonRpcMessage, NumberOrString, ProgressToken, ServerNotification, ServerRequest,
};

pub trait GetMeta {
Expand All@@ -18,6 +18,26 @@ pub trait GetExtensions {
fn extensions_mut(&mut self) -> &mut Extensions;
}

impl GetExtensions for CustomClientNotification {
fn extensions(&self) -> &Extensions {
&self.extensions
}
fn extensions_mut(&mut self) -> &mut Extensions {
&mut self.extensions
}
}

impl GetMeta for CustomClientNotification {
fn get_meta_mut(&mut self) -> &mut Meta {
self.extensions_mut().get_or_insert_default()
}
fn get_meta(&self) -> &Meta {
self.extensions()
.get::<Meta>()
.unwrap_or(Meta::static_empty())
}
}

macro_rules! variant_extension {
(
$Enum: ident {
Expand DownExpand Up@@ -84,6 +104,7 @@ variant_extension! {
ProgressNotification
InitializedNotification
RootsListChangedNotification
CustomClientNotification
}
}

Expand Down
57 changes: 55 additions & 2 deletions crates/rmcp/src/model/serde_impl.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,8 +3,8 @@ use std::borrow::Cow;
use serde::{Deserialize, Serialize};

use super::{
Extensions, Meta, Notification, NotificationNoParam, Request, RequestNoParam,
RequestOptionalParam,
CustomClientNotification, Extensions, Meta, Notification, NotificationNoParam, Request,
RequestNoParam, RequestOptionalParam,
};
#[derive(Serialize, Deserialize)]
struct WithMeta<'a, P> {
Expand DownExpand Up@@ -249,6 +249,59 @@ where
}
}

impl Serialize for CustomClientNotification {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
let extensions = &self.extensions;
let _meta = extensions.get::<Meta>().map(Cow::Borrowed);
let params = self.params.as_ref();

let params = if _meta.is_some() || params.is_some() {
Some(WithMeta {
_meta,
_rest: &self.params,
})
} else {
None
};

ProxyOptionalParam::serialize(
&ProxyOptionalParam {
method: &self.method,
params,
},
serializer,
)
}
}

impl<'de> Deserialize<'de> for CustomClientNotification {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let body =
ProxyOptionalParam::<'_, _, Option<serde_json::Value>>::deserialize(deserializer)?;
let mut params = None;
let mut _meta = None;
if let Some(body_params) = body.params {
params = body_params._rest;
_meta = body_params._meta.map(|m| m.into_owned());
}
let mut extensions = Extensions::new();
if let Some(meta) = _meta {
extensions.insert(meta);
}
Ok(CustomClientNotification {
extensions,
method: body.method,
params,
})
}
}

#[cfg(test)]
mod test {
use serde_json::json;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -396,6 +396,19 @@
"content"
]
},
"CustomClientNotification": {
"description": "A catch-all notification the client can use to send custom messages to a server.\n\nThis preserves the raw `method` name and `params` payload so handlers can\ndeserialize them into domain-specific types.",
"type": "object",
"properties": {
"method": {
"type": "string"
},
"params": true
},
"required": [
"method"
]
},
"ElicitationAction": {
"description": "Represents the possible actions a user can take in response to an elicitation request.\n\nWhen a server requests user input through elicitation, the user can:\n- Accept: Provide the requested information and continue\n- Decline: Refuse to provide the information but continue the operation\n- Cancel: Stop the entire operation",
"oneOf": [
Expand DownExpand Up@@ -636,6 +649,9 @@
},
{
"$ref": "#/definitions/NotificationNoParam2"
},
{
"$ref": "#/definitions/CustomClientNotification"
}
],
"required": [
Expand Down
76 changes: 74 additions & 2 deletions crates/rmcp/tests/test_notification.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,10 +3,12 @@ use std::sync::Arc;
use rmcp::{
ClientHandler, ServerHandler, ServiceExt,
model::{
ResourceUpdatedNotificationParam, ServerCapabilities, ServerInfo, SubscribeRequestParam,
ClientNotification, CustomClientNotification, ResourceUpdatedNotificationParam,
ServerCapabilities, ServerInfo, SubscribeRequestParam,
},
};
use tokio::sync::Notify;
use serde_json::json;
use tokio::sync::{Mutex, Notify};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

pub struct Server {}
Expand DownExpand Up@@ -93,3 +95,73 @@ async fn test_server_notification() -> anyhow::Result<()> {
client.cancel().await?;
Ok(())
}

type CustomNotificationPayload = (String, Option<serde_json::Value>);

struct CustomServer {
receive_signal: Arc<Notify>,
payload: Arc<Mutex<Option<CustomNotificationPayload>>>,
}

impl ServerHandler for CustomServer {
async fn on_custom_notification(
&self,
notification: CustomClientNotification,
_context: rmcp::service::NotificationContext<rmcp::RoleServer>,
) {
let CustomClientNotification { method, params, .. } = notification;
let mut payload = self.payload.lock().await;
*payload = Some((method, params));
self.receive_signal.notify_one();
}
}

#[tokio::test]
async fn test_custom_client_notification_reaches_server() -> anyhow::Result<()> {
let _ = tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "debug".to_string().into()),
)
.with(tracing_subscriber::fmt::layer())
.try_init();

let (server_transport, client_transport) = tokio::io::duplex(4096);
let receive_signal = Arc::new(Notify::new());
let payload = Arc::new(Mutex::new(None));

{
let receive_signal = receive_signal.clone();
let payload = payload.clone();
tokio::spawn(async move {
let server = CustomServer {
receive_signal,
payload,
}
.serve(server_transport)
.await?;
server.waiting().await?;
anyhow::Ok(())
});
}

let client = ().serve(client_transport).await?;

client
.send_notification(ClientNotification::CustomClientNotification(
CustomClientNotification::new(
"notifications/custom-test",
Some(json!({ "foo": "bar" })),
),
))
.await?;

tokio::time::timeout(std::time::Duration::from_secs(5), receive_signal.notified()).await?;

let (method, params) = payload.lock().await.clone().expect("payload set");
assert_eq!("notifications/custom-test", method);
assert_eq!(Some(json!({ "foo": "bar" })), params);

client.cancel().await?;
Ok(())
}