use std::collections::HashMap;
use std::fs::OpenOptions;
use std::io::{BufWriter, Write};
use std::path::Path;
use std::time::Instant;
use agent_client_protocol::schema::SuccessorMessage;
use agent_client_protocol::schema::v1::{
MessageMcpNotification, MessageMcpRequest, MessageMcpResponse, Notification as RpcNotification,
Request as RpcRequest, RequestId,
};
use agent_client_protocol::{
DynConnectTo, JsonRpcMessage, RawJsonRpcMessage, RawJsonRpcParams,
RawJsonRpcResponse as RpcResponse, Role, UntypedMessage,
};
use rustc_hash::FxHashMap;
use serde::{Deserialize, Serialize};
use crate::ComponentIndex;
use crate::snoop::SnooperComponent;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
#[non_exhaustive]
pub enum TraceEvent {
Request(RequestEvent),
Response(ResponseEvent),
Notification(NotificationEvent),
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum Protocol {
Acp,
Mcp,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[non_exhaustive]
pub struct RequestEvent {
pub ts: f64,
pub protocol: Protocol,
pub from: String,
pub to: String,
pub id: serde_json::Value,
pub method: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub session: Option<String>,
pub params: serde_json::Value,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[non_exhaustive]
pub struct ResponseEvent {
pub ts: f64,
pub from: String,
pub to: String,
pub id: serde_json::Value,
pub is_error: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error_domain: Option<Protocol>,
pub payload: serde_json::Value,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[non_exhaustive]
pub struct NotificationEvent {
pub ts: f64,
pub protocol: Protocol,
pub from: String,
pub to: String,
pub method: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub session: Option<String>,
pub params: serde_json::Value,
}
pub trait WriteEvent: Send + 'static {
fn write_event(&mut self, event: &TraceEvent) -> std::io::Result<()>;
}
pub(crate) struct EventWriter<W> {
writer: W,
}
impl<W: Write> EventWriter<W> {
pub fn new(writer: W) -> Self {
Self { writer }
}
}
impl<W: Write + Send + 'static> WriteEvent for EventWriter<W> {
fn write_event(&mut self, event: &TraceEvent) -> std::io::Result<()> {
serde_json::to_writer(&mut self.writer, event).map_err(std::io::Error::other)?;
self.writer.write_all(b"\n")?;
self.writer.flush()
}
}
impl WriteEvent for futures::channel::mpsc::UnboundedSender<TraceEvent> {
fn write_event(&mut self, event: &TraceEvent) -> std::io::Result<()> {
self.unbounded_send(event.clone())
.map_err(|e| std::io::Error::new(std::io::ErrorKind::BrokenPipe, e))
}
}
pub struct TraceWriter {
dest: Box<dyn WriteEvent>,
start_time: Instant,
request_details: FxHashMap<serde_json::Value, RequestDetails>,
}
impl std::fmt::Debug for TraceWriter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TraceWriter")
.field("start_time", &self.start_time)
.finish_non_exhaustive()
}
}
struct RequestDetails {
protocol: Protocol,
request_from: ComponentIndex,
request_to: ComponentIndex,
}
impl TraceWriter {
pub fn new<D: WriteEvent>(dest: D) -> Self {
Self {
dest: Box::new(dest),
start_time: Instant::now(),
request_details: HashMap::default(),
}
}
pub fn from_path(path: impl AsRef<Path>) -> std::io::Result<Self> {
let file = OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(path.as_ref())?;
Ok(Self::new(EventWriter::new(BufWriter::new(file))))
}
fn elapsed(&self) -> f64 {
self.start_time.elapsed().as_secs_f64()
}
fn write_event(&mut self, event: &TraceEvent) {
drop(self.dest.write_event(event));
}
#[expect(clippy::too_many_arguments)]
fn request(
&mut self,
protocol: Protocol,
from: ComponentIndex,
to: ComponentIndex,
id: serde_json::Value,
method: String,
session: Option<String>,
mut params: serde_json::Value,
) {
redact_http_credentials(&mut params);
self.request_details.insert(
id.clone(),
RequestDetails {
protocol,
request_from: from,
request_to: to,
},
);
self.write_event(&TraceEvent::Request(RequestEvent {
ts: self.elapsed(),
protocol,
from: format!("{from:?}"),
to: format!("{to:?}"),
id,
method,
session,
params,
}));
}
fn response(
&mut self,
from: ComponentIndex,
to: ComponentIndex,
id: serde_json::Value,
error_domain: Option<Protocol>,
mut payload: serde_json::Value,
) {
redact_http_credentials(&mut payload);
self.write_event(&TraceEvent::Response(ResponseEvent {
ts: self.elapsed(),
from: format!("{from:?}"),
to: format!("{to:?}"),
id,
is_error: error_domain.is_some(),
error_domain,
payload,
}));
}
fn notification(
&mut self,
protocol: Protocol,
from: ComponentIndex,
to: ComponentIndex,
method: impl Into<String>,
session: Option<String>,
mut params: serde_json::Value,
) {
redact_http_credentials(&mut params);
self.write_event(&TraceEvent::Notification(NotificationEvent {
ts: self.elapsed(),
protocol,
from: format!("{from:?}"),
to: format!("{to:?}"),
method: method.into(),
session,
params,
}));
}
fn trace_message(&mut self, traced_message: TracedMessage) {
let TracedMessage {
component_index,
successor_index,
incoming,
message,
} = traced_message;
match message {
RawJsonRpcMessage::Request(req) => {
let MessageInfo {
successor,
id,
protocol,
method,
params,
} = MessageInfo::from_request(req);
self.trace_request_or_notification(
incoming,
component_index,
successor_index,
successor,
id,
protocol,
method,
params,
);
}
RawJsonRpcMessage::Notification(notification) => {
let MessageInfo {
successor,
id,
protocol,
method,
params,
} = MessageInfo::from_notification(notification);
self.trace_request_or_notification(
incoming,
component_index,
successor_index,
successor,
id,
protocol,
method,
params,
);
}
RawJsonRpcMessage::Response(resp) => {
let (id, is_error, payload) = match resp {
RpcResponse::Result { id, result } => (id, false, result),
RpcResponse::Error { id, error } => {
(id, true, serde_json::to_value(error).unwrap_or_default())
}
};
let id = id_to_json(&id);
if let Some(RequestDetails {
protocol,
request_from,
request_to,
}) = self.request_details.remove(&id)
{
let (error_domain, payload) = response_outcome(protocol, is_error, payload);
self.response(request_to, request_from, id, error_domain, payload);
}
}
}
}
#[expect(clippy::too_many_arguments)]
fn trace_request_or_notification(
&mut self,
incoming: Incoming,
component_index: ComponentIndex,
successor_index: ComponentIndex,
successor: Successor,
id: Option<RequestId>,
protocol: Protocol,
method: String,
params: serde_json::Value,
) {
let (from, to) = match (successor, incoming, component_index, successor_index) {
(Successor(false), Incoming(true), ComponentIndex::Proxy(proxy_index), _) => (
ComponentIndex::predecessor_of(proxy_index),
ComponentIndex::Proxy(proxy_index),
),
(Successor(true), Incoming(true), component_index, successor_index) => {
(successor_index, component_index)
}
(Successor(true), Incoming(false), component_index, ComponentIndex::Agent) => {
(component_index, ComponentIndex::Agent)
}
_ => return,
};
match id {
Some(id) => {
self.request(protocol, from, to, id_to_json(&id), method, None, params);
}
None => {
self.notification(protocol, from, to, method, None, params);
}
}
}
pub(crate) fn spawn(
mut self: TraceWriter,
) -> (
TraceHandle,
impl std::future::Future<Output = Result<(), agent_client_protocol::Error>>,
) {
use futures::StreamExt;
let (tx, mut rx) = futures::channel::mpsc::unbounded();
let future = async move {
while let Some(event) = rx.next().await {
self.trace_message(event);
}
Ok(())
};
(TraceHandle { tx }, future)
}
}
#[derive(Clone, Debug)]
pub(crate) struct TraceHandle {
tx: futures::channel::mpsc::UnboundedSender<TracedMessage>,
}
impl TraceHandle {
fn trace_message(
&self,
component_index: ComponentIndex,
successor_index: ComponentIndex,
incoming: Incoming,
message: &RawJsonRpcMessage,
) -> Result<(), agent_client_protocol::Error> {
self.tx
.unbounded_send(TracedMessage {
component_index,
successor_index,
incoming,
message: message.clone(),
})
.map_err(agent_client_protocol::util::internal_error)
}
pub fn bridge_component<R: Role>(
&self,
proxy_index: ComponentIndex,
successor_index: ComponentIndex,
proxy: impl agent_client_protocol::ConnectTo<R>,
) -> DynConnectTo<R> {
DynConnectTo::new(SnooperComponent::new(
proxy,
{
let trace_handle = self.clone();
move |msg| {
trace_handle.trace_message(proxy_index, successor_index, Incoming(true), msg)
}
},
{
let trace_handle = self.clone();
move |msg| {
trace_handle.trace_message(proxy_index, successor_index, Incoming(false), msg)
}
},
))
}
}
fn id_to_json(id: &RequestId) -> serde_json::Value {
serde_json::to_value(id).expect("RequestId serializes infallibly")
}
fn params_from_transport(params: Option<RawJsonRpcParams>) -> serde_json::Value {
params.map_or(serde_json::Value::Null, RawJsonRpcParams::into_value)
}
fn response_outcome(
protocol: Protocol,
outer_error: bool,
payload: serde_json::Value,
) -> (Option<Protocol>, serde_json::Value) {
if outer_error {
return (Some(Protocol::Acp), payload);
}
if protocol == Protocol::Mcp {
match serde_json::from_value::<MessageMcpResponse>(payload.clone()) {
Ok(MessageMcpResponse::Result { result, .. }) => return (None, result),
Ok(MessageMcpResponse::Error { error, .. }) => {
return (
Some(Protocol::Mcp),
serde_json::to_value(error).expect("MCP errors contain only JSON values"),
);
}
_ => {}
}
}
(None, payload)
}
fn redact_http_credentials(value: &mut serde_json::Value) {
fn is_credential(name: &str) -> bool {
[
"authorization",
"proxy-authorization",
"cookie",
"set-cookie",
"x-api-key",
]
.iter()
.any(|candidate| name.eq_ignore_ascii_case(candidate))
}
match value {
serde_json::Value::Object(object) => {
if let Some(serde_json::Value::String(url)) = object.get_mut("url")
&& let Some(redacted) = redact_http_url(url)
{
*url = redacted;
}
match object.get_mut("headers") {
Some(serde_json::Value::Array(headers)) => {
for header in headers {
if header
.get("name")
.and_then(serde_json::Value::as_str)
.is_some_and(is_credential)
&& let Some(value) = header.get_mut("value")
{
*value = serde_json::Value::String("[REDACTED]".to_owned());
}
}
}
Some(serde_json::Value::Object(headers)) => {
for (name, value) in headers {
if is_credential(name) {
*value = serde_json::Value::String("[REDACTED]".to_owned());
}
}
}
_ => {}
}
for value in object.values_mut() {
redact_http_credentials(value);
}
}
serde_json::Value::Array(values) => {
for value in values {
redact_http_credentials(value);
}
}
_ => {}
}
}
fn is_secret_query_key(name: &str) -> bool {
let normalized: String = name
.chars()
.filter(|c| !matches!(c, '-' | '_'))
.map(|c| c.to_ascii_lowercase())
.collect();
matches!(
normalized.as_str(),
"token"
| "accesstoken"
| "refreshtoken"
| "idtoken"
| "apikey"
| "key"
| "secret"
| "clientsecret"
| "password"
| "passwd"
| "pwd"
| "auth"
| "authorization"
| "bearer"
| "signature"
| "sig"
| "credential"
| "credentials"
)
}
fn redact_http_url(value: &str) -> Option<String> {
let mut url = match url::Url::parse(value) {
Ok(url) if matches!(url.scheme(), "http" | "https") => url,
Ok(_) => return None,
Err(_) => {
let scheme: String = value
.trim_start_matches(|c: char| c <= '\u{20}')
.split(':')
.next()
.unwrap_or_default()
.chars()
.filter(|c| !matches!(c, '\t' | '\n' | '\r'))
.collect();
return (scheme.eq_ignore_ascii_case("http") || scheme.eq_ignore_ascii_case("https"))
.then(|| "[REDACTED]".to_owned());
}
};
let mut changed = false;
if !url.username().is_empty() || url.password().is_some() {
url.set_password(None).expect("HTTP URL has an authority");
url.set_username("").expect("HTTP URL has an authority");
changed = true;
}
if let Some(query) = url.query() {
let redacted = query
.split('&')
.map(|pair| {
let is_secret = url::form_urlencoded::parse(pair.as_bytes())
.next()
.is_some_and(|(key, _)| is_secret_query_key(&key));
if is_secret {
changed = true;
let key = pair.split('=').next().unwrap_or_default();
format!("{key}=[REDACTED]")
} else {
pair.to_owned()
}
})
.collect::<Vec<_>>()
.join("&");
url.set_query(Some(&redacted));
}
changed.then(|| url.into())
}
#[derive(Debug)]
struct TracedMessage {
component_index: ComponentIndex,
successor_index: ComponentIndex,
incoming: Incoming,
message: RawJsonRpcMessage,
}
#[derive(Debug)]
struct MessageInfo {
successor: Successor,
id: Option<RequestId>,
protocol: Protocol,
method: String,
params: serde_json::Value,
}
#[derive(Copy, Clone, Debug)]
struct Successor(bool);
#[derive(Copy, Clone, Debug)]
struct Incoming(bool);
impl MessageInfo {
fn from_request(req: RpcRequest<RawJsonRpcParams>) -> Self {
let untyped =
UntypedMessage::parse_message(&req.method, ¶ms_from_transport(req.params))
.expect("untyped message is infallible");
Self::from_untyped_request(Successor(false), Some(req.id), Protocol::Acp, untyped)
}
fn from_notification(notification: RpcNotification<RawJsonRpcParams>) -> Self {
let untyped = UntypedMessage::parse_message(
¬ification.method,
¶ms_from_transport(notification.params),
)
.expect("untyped message is infallible");
Self::from_untyped_notification(Successor(false), Protocol::Acp, untyped)
}
fn from_untyped_request(
successor: Successor,
id: Option<RequestId>,
protocol: Protocol,
untyped: UntypedMessage,
) -> Self {
if let Ok(m) = SuccessorMessage::parse_message(&untyped.method, &untyped.params) {
return Self::from_untyped_request(Successor(true), id, protocol, m.message);
}
if let Ok(m) = MessageMcpRequest::parse_message(&untyped.method, &untyped.params) {
let params = m
.params
.map_or(serde_json::Value::Null, serde_json::Value::Object);
return Self::from_untyped_request(
successor,
id,
Protocol::Mcp,
UntypedMessage {
method: m.method,
params,
},
);
}
Self::new(successor, id, protocol, untyped)
}
fn from_untyped_notification(
successor: Successor,
protocol: Protocol,
untyped: UntypedMessage,
) -> Self {
if let Ok(m) = SuccessorMessage::parse_message(&untyped.method, &untyped.params) {
return Self::from_untyped_notification(Successor(true), protocol, m.message);
}
if let Ok(m) = MessageMcpNotification::parse_message(&untyped.method, &untyped.params) {
let params = m
.params
.map_or(serde_json::Value::Null, serde_json::Value::Object);
return Self::from_untyped_notification(
successor,
Protocol::Mcp,
UntypedMessage {
method: m.method,
params,
},
);
}
Self::new(successor, None, protocol, untyped)
}
fn new(
successor: Successor,
id: Option<RequestId>,
protocol: Protocol,
untyped: UntypedMessage,
) -> Self {
Self {
successor,
id,
protocol,
method: untyped.method,
params: untyped.params,
}
}
}
#[cfg(test)]
mod tests {
use agent_client_protocol::RawJsonRpcMessage;
use serde_json::json;
use super::{MessageInfo, Protocol, ResponseEvent, redact_http_credentials, response_outcome};
#[test]
fn http_url_redaction_handles_encoded_keys_duplicates_and_userinfo() {
let mut value = json!({
"url":"HTTPS://user:p%40ss@example.test:8443/mcp?mode=a%20b&%61ccess_TOKEN=one&api-key=two&token&token=three&&count=2#section"
});
redact_http_credentials(&mut value);
assert_eq!(
value["url"],
"https://example.test:8443/mcp?mode=a%20b&%61ccess_TOKEN=[REDACTED]&api-key=[REDACTED]&token=[REDACTED]&token=[REDACTED]&&count=2#section"
);
let once = value.clone();
redact_http_credentials(&mut value);
assert_eq!(value, once, "redaction must be idempotent");
for key in [
"TOKEN",
"access_token",
"refresh-token",
"idToken",
"api_key",
"key",
"secret",
"client_secret",
"password",
"passwd",
"pwd",
"auth",
"authorization",
"bearer",
"signature",
"sig",
"credential",
"credentials",
] {
let mut value = json!({"url":format!("http://example.test/mcp?{key}=private&mode=ok")});
redact_http_credentials(&mut value);
assert_eq!(
value["url"],
format!("http://example.test/mcp?{key}=[REDACTED]&mode=ok")
);
}
for url in [
"https://user@example.test/mcp",
"https://:private@example.test/mcp",
"https://u%40ser:p%40ss@example.test/mcp",
] {
let mut value = json!({"url":url});
redact_http_credentials(&mut value);
assert_eq!(value["url"], "https://example.test/mcp");
}
for scheme in [
"https",
"http\t",
"ht\ntps",
"h\rttp",
"\u{0}https",
"\u{1f}http",
] {
let mut invalid = json!({
"url":format!("{scheme}://user:private@[bad-host]/?token=private")
});
redact_http_credentials(&mut invalid);
assert_eq!(invalid["url"], "[REDACTED]");
}
}
#[test]
fn credential_free_urls_and_explicit_payloads_are_preserved() {
let original = json!({
"url":"HTTPS://EXAMPLE.test:443/mcp?mode=a+b&mode=a%20b&&count=2#section",
"nonHttp":{"url":"file:///tmp/data?token=visible"},
"prompt":"an intentionally recorded prompt",
"image":{"data":"intentionally recorded image"},
"file":{"content":"intentionally recorded file"},
"customSecret":"not a recognized credential field"
});
let mut trace_copy = original.clone();
redact_http_credentials(&mut trace_copy);
assert_eq!(trace_copy, original);
}
#[tokio::test]
async fn recording_redacts_all_event_kinds_without_changing_wire_messages() {
use super::{ComponentIndex, TraceEvent, TraceWriter};
use agent_client_protocol::{Channel, ConnectTo, TransportFrame, UntypedRole};
use futures::StreamExt as _;
tokio::task::LocalSet::new().run_until(async {
let payload = json!({
"mcpServers":[{
"type":"http", "url":"https://user:private@example.test/mcp?token=private&mode=ok",
"headers":[{"name":"Authorization","value":"Bearer private"},
{"name":"visible","value":"ok"}]
}],
"prompt":"recorded prompt", "image":{"data":"recorded image"},
"file":{"content":"recorded file"}
});
let (events_tx, mut events_rx) = futures::channel::mpsc::unbounded();
let (handle, recording) = TraceWriter::new(events_tx).spawn();
let recording = tokio::task::spawn_local(recording);
let (client, mut client_peer) = Channel::duplex();
let (base, mut base_peer) = Channel::duplex();
let bridge = handle.bridge_component::<UntypedRole>(
ComponentIndex::Proxy(0), ComponentIndex::Agent, base);
let bridge = tokio::task::spawn_local(bridge.connect_to(client));
drop(handle);
let request = RawJsonRpcMessage::request("test/request".into(), payload.clone(), 1.into()).unwrap();
let notification = RawJsonRpcMessage::notification("test/notification".into(), payload.clone()).unwrap();
for message in [request, notification] {
let frame = TransportFrame::Single(message);
let original_wire = frame.to_json().unwrap();
client_peer.tx.unbounded_send(frame).unwrap();
let forwarded = base_peer.rx.next().await.unwrap();
assert_eq!(forwarded.to_json().unwrap(), original_wire);
}
let response = TransportFrame::Single(RawJsonRpcMessage::response(1.into(), Ok(payload.clone())));
let original_wire = response.to_json().unwrap();
base_peer.tx.unbounded_send(response).unwrap();
assert_eq!(client_peer.rx.next().await.unwrap().to_json().unwrap(), original_wire);
drop(client_peer.tx);
drop(base_peer.tx);
bridge.await.unwrap().unwrap();
recording.await.unwrap().unwrap();
let mut expected = payload.clone();
expected["mcpServers"][0]["url"] =
json!("https://example.test/mcp?token=[REDACTED]&mode=ok");
expected["mcpServers"][0]["headers"][0]["value"] = json!("[REDACTED]");
let mut count = 0;
while let Some(event) = events_rx.next().await {
let recorded = match event {
TraceEvent::Request(event) => event.params,
TraceEvent::Notification(event) => event.params,
TraceEvent::Response(event) => event.payload,
};
assert_eq!(recorded, expected);
count += 1;
}
assert_eq!(count, 3);
assert_eq!(payload["mcpServers"][0]["headers"][0]["value"], "Bearer private");
}).await;
}
#[test]
fn mcp_and_binding_errors_keep_their_domains() {
let inner = json!({"code":-32000,"message":"peer","data":null,"extension":true});
assert_eq!(
response_outcome(Protocol::Mcp, false, json!({"error":inner})),
(Some(Protocol::Mcp), inner)
);
let outer = json!({"code":-33002,"message":"binding failure"});
assert_eq!(
response_outcome(Protocol::Mcp, true, outer.clone()),
(Some(Protocol::Acp), outer)
);
for result in [
json!(null),
json!({"resultType":"input_required","requestState":"opaque"}),
] {
assert_eq!(
response_outcome(Protocol::Mcp, false, json!({"result":result})),
(None, result)
);
}
}
#[test]
fn old_trace_response_without_domain_remains_readable() {
let event: ResponseEvent = serde_json::from_value(json!({
"ts":0,"from":"client","to":"agent","id":1,"is_error":true,
"payload":{"code":-32602,"message":"old trace"}
}))
.unwrap();
assert!(event.is_error);
assert_eq!(event.error_domain, None);
}
#[test]
fn declaration_credentials_are_redacted_without_changing_transport_payload() {
let original = json!({"mcpServers":[
{"headers":[{"name":"Authorization","value":"Bearer private"},{"name":"visible","value":"ok"}]},
{"headers":{"COOKIE":"private","visible":"ok"}}
]});
let mut trace_copy = original.clone();
redact_http_credentials(&mut trace_copy);
assert_eq!(
trace_copy["mcpServers"][0]["headers"][0]["value"],
"[REDACTED]"
);
assert_eq!(
trace_copy["mcpServers"][1]["headers"]["COOKIE"],
"[REDACTED]"
);
assert_eq!(trace_copy["mcpServers"][0]["headers"][1]["value"], "ok");
assert_eq!(
original["mcpServers"][0]["headers"][0]["value"],
"Bearer private"
);
}
#[test]
fn malformed_mcp_notification_preserves_the_observed_envelope() {
let params = json!({
"serverId": "server-1",
"requestId": "request-1",
"method": "notifications/progress",
"params": ["invalid named params"]
});
let RawJsonRpcMessage::Notification(notification) =
RawJsonRpcMessage::notification("mcp/message".into(), params.clone())
.expect("notification is valid JSON-RPC")
else {
unreachable!("notification constructor returned a different message kind")
};
let info = MessageInfo::from_notification(notification);
assert_eq!(info.protocol, Protocol::Acp);
assert_eq!(info.method, "mcp/message");
assert_eq!(info.params, params);
}
#[test]
fn valid_mcp_notification_is_traced_as_inner_mcp() {
let params = json!({"progressToken":"token", "progress":1});
let RawJsonRpcMessage::Notification(notification) = RawJsonRpcMessage::notification(
"mcp/message".into(),
json!({
"serverId":"server-1","requestId":"request-1",
"method":"notifications/progress","params":params
}),
)
.unwrap() else {
unreachable!("notification constructor")
};
let info = MessageInfo::from_notification(notification);
assert_eq!(info.protocol, Protocol::Mcp);
assert_eq!(info.method, "notifications/progress");
assert_eq!(info.params, params);
}
}