use super::types::McpEvent;
pub(crate) const RESULT_TYPE_SINCE: &str = "2026-07-28";
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum EnvelopeObservation {
NotApplicable,
Absent,
Present(String),
Malformed,
}
#[derive(Debug, Clone, Default)]
pub(crate) struct UniqueValue(pub(crate) serde_json::Value);
impl<'de> serde::Deserialize<'de> for UniqueValue {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
struct Visitor;
impl<'de> serde::de::Visitor<'de> for Visitor {
type Value = UniqueValue;
fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("any JSON value with unique object members")
}
fn visit_unit<E>(self) -> Result<UniqueValue, E> {
Ok(UniqueValue(serde_json::Value::Null))
}
fn visit_none<E>(self) -> Result<UniqueValue, E> {
Ok(UniqueValue(serde_json::Value::Null))
}
fn visit_some<D: serde::Deserializer<'de>>(
self,
d: D,
) -> Result<UniqueValue, D::Error> {
<UniqueValue as serde::Deserialize>::deserialize(d)
}
fn visit_bool<E>(self, v: bool) -> Result<UniqueValue, E> {
Ok(UniqueValue(v.into()))
}
fn visit_i64<E>(self, v: i64) -> Result<UniqueValue, E> {
Ok(UniqueValue(v.into()))
}
fn visit_u64<E>(self, v: u64) -> Result<UniqueValue, E> {
Ok(UniqueValue(v.into()))
}
fn visit_f64<E>(self, v: f64) -> Result<UniqueValue, E> {
Ok(UniqueValue(
serde_json::Number::from_f64(v).map_or(serde_json::Value::Null, Into::into),
))
}
fn visit_str<E>(self, v: &str) -> Result<UniqueValue, E> {
Ok(UniqueValue(v.into()))
}
fn visit_seq<A: serde::de::SeqAccess<'de>>(
self,
mut seq: A,
) -> Result<UniqueValue, A::Error> {
let mut out = Vec::new();
while let Some(UniqueValue(v)) = seq.next_element()? {
out.push(v);
}
Ok(UniqueValue(serde_json::Value::Array(out)))
}
fn visit_map<A: serde::de::MapAccess<'de>>(
self,
mut map: A,
) -> Result<UniqueValue, A::Error> {
let mut out = serde_json::Map::new();
while let Some(key) = map.next_key::<String>()? {
if out.contains_key(&key) {
return Err(serde::de::Error::custom(
"JSON object contains a duplicate member",
));
}
let UniqueValue(value) = map.next_value()?;
out.insert(key, value);
}
Ok(UniqueValue(serde_json::Value::Object(out)))
}
}
deserializer.deserialize_any(Visitor)
}
}
#[derive(Default)]
pub(crate) struct SeenMembers(std::collections::HashSet<String>);
impl SeenMembers {
pub(crate) fn insert<E: serde::de::Error>(&mut self, key: &str) -> Result<(), E> {
if !self.0.insert(key.to_string()) {
return Err(serde::de::Error::custom(
"JSON object contains a duplicate member",
));
}
Ok(())
}
}
pub(crate) struct DuplicateAwareSink;
impl<'de> serde::Deserialize<'de> for DuplicateAwareSink {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
struct V;
impl<'de> serde::de::Visitor<'de> for V {
type Value = DuplicateAwareSink;
fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("any JSON value with unique object members")
}
fn visit_unit<E>(self) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_none<E>(self) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_some<D: serde::Deserializer<'de>>(
self,
d: D,
) -> Result<DuplicateAwareSink, D::Error> {
<DuplicateAwareSink as serde::Deserialize>::deserialize(d)
}
fn visit_bool<E>(self, _: bool) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_i64<E>(self, _: i64) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_u64<E>(self, _: u64) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_f64<E>(self, _: f64) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_str<E>(self, _: &str) -> Result<DuplicateAwareSink, E> {
Ok(DuplicateAwareSink)
}
fn visit_seq<A: serde::de::SeqAccess<'de>>(
self,
mut seq: A,
) -> Result<DuplicateAwareSink, A::Error> {
while seq.next_element::<DuplicateAwareSink>()?.is_some() {}
Ok(DuplicateAwareSink)
}
fn visit_map<A: serde::de::MapAccess<'de>>(
self,
mut map: A,
) -> Result<DuplicateAwareSink, A::Error> {
let mut seen = SeenMembers::default();
while let Some(key) = map.next_key::<String>()? {
seen.insert::<A::Error>(&key)?;
map.next_value::<DuplicateAwareSink>()?;
}
Ok(DuplicateAwareSink)
}
}
deserializer.deserialize_any(V)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub(crate) enum CorrelationId {
Str(String),
Num(String),
}
pub(crate) fn id_is_acceptable(v: &serde_json::Value) -> bool {
matches!(
v,
serde_json::Value::String(_) | serde_json::Value::Number(_)
)
}
pub(crate) fn correlation_id(v: &serde_json::Value) -> Option<CorrelationId> {
v.get("id").and_then(correlation_key)
}
pub(crate) fn correlation_key(v: &serde_json::Value) -> Option<CorrelationId> {
match v {
serde_json::Value::String(s) => Some(CorrelationId::Str(s.clone())),
serde_json::Value::Number(n) if n.is_i64() || n.is_u64() => {
Some(CorrelationId::Num(n.to_string()))
}
_ => None,
}
}
pub(crate) const REQUEST_ONLY_METHODS: &[&str] = &[
"completion/complete",
"elicitation/create",
"initialize",
"logging/setLevel",
"ping",
"prompts/get",
"prompts/list",
"resources/list",
"resources/read",
"resources/subscribe",
"resources/templates/list",
"resources/unsubscribe",
"roots/list",
"sampling/createMessage",
"server/discover",
"subscriptions/listen",
"tasks/cancel",
"tasks/get",
"tasks/list",
"tasks/result",
"tools/call",
"tools/list",
];
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MessageKind<'a> {
Request { method: &'a str },
Notification { method: &'a str },
Response,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub(crate) enum MessageShapeError {
#[error("JSON-RPC method must be a string")]
NonStringMethod,
#[error("a JSON-RPC response must carry exactly one of result or error")]
NotExactlyOneResponseBody,
#[error("a JSON-RPC error response must carry an integer code and string message")]
MalformedErrorBody,
}
pub(crate) fn classify_message(
v: &serde_json::Value,
) -> Result<MessageKind<'_>, MessageShapeError> {
let Some(method) = v.get("method") else {
let has_result = v.get("result").is_some();
let has_error = v.get("error").is_some();
if has_result == has_error {
return Err(MessageShapeError::NotExactlyOneResponseBody);
}
if has_error {
let Some(error) = v.get("error").and_then(serde_json::Value::as_object) else {
return Err(MessageShapeError::MalformedErrorBody);
};
let integer_code = error
.get("code")
.and_then(serde_json::Value::as_i64)
.is_some()
|| error
.get("code")
.and_then(serde_json::Value::as_u64)
.is_some();
let string_message = error
.get("message")
.and_then(serde_json::Value::as_str)
.is_some();
if !integer_code || !string_message {
return Err(MessageShapeError::MalformedErrorBody);
}
}
return Ok(MessageKind::Response);
};
let Some(method) = method.as_str() else {
return Err(MessageShapeError::NonStringMethod);
};
Ok(
if v.get("id").is_some() || REQUEST_ONLY_METHODS.contains(&method) {
MessageKind::Request { method }
} else {
MessageKind::Notification { method }
},
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum RequestMetadata {
Absent,
Present(String),
Malformed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum UnknownReason {
NoSignal,
UnsupportedVersion(String),
MalformedSignal,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum EraResolution {
Known(String),
Unknown(UnknownReason),
Conflicting {
header: String,
body: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ResultObservation {
Missing,
Complete,
InputRequired,
Unrecognized,
CompleteWithContinuation,
InputRequiredWithMalformedContinuation,
InputRequiredWithoutContinuation,
Malformed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum CapabilityObservation {
CoreOnly,
ExtensionNotUnderstood,
Absent,
Malformed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum IncompleteReason {
EraUnknown(UnknownReason),
UnrecognizedResultType,
RecognitionUndeterminable,
ContradictoryResult,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum InvalidReason {
EraConflicting {
header: String,
body: String,
},
MissingResultType,
MissingRequestMetadata,
MalformedRequestMetadata,
MalformedEraSignal,
MalformedResultType,
MalformedCapabilities,
MissingCapabilities,
MalformedContinuation,
UncontinuableInputRequired,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum RequestAssessment {
Valid,
Incomplete(IncompleteReason),
Invalid(InvalidReason),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ResultConclusion {
Terminal,
NonTerminal,
Incomplete(IncompleteReason),
Invalid(InvalidReason),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct McpEraContext {
pub(crate) envelope: EnvelopeObservation,
pub(crate) era: EraResolution,
pub(crate) request_metadata: Option<RequestMetadata>,
pub(crate) result_observation: Option<ResultObservation>,
pub(crate) correlation: Option<CorrelationId>,
pub(crate) capability_observation: Option<CapabilityObservation>,
}
#[derive(Debug, Clone)]
pub(crate) struct ParsedMcpEvent {
pub(crate) event: McpEvent,
pub(crate) context: McpEraContext,
pub(crate) is_error_response: bool,
}
pub(crate) const SUPPORTED_VERSIONS: &[&str] = &[
"2024-11-05",
"2025-03-26",
"2025-06-18",
"2025-11-25",
RESULT_TYPE_SINCE,
];
pub(crate) fn is_version_shaped(v: &str) -> bool {
let b = v.as_bytes();
if b.len() != 10 || b[4] != b'-' || b[7] != b'-' {
return false;
}
if !b
.iter()
.enumerate()
.all(|(i, c)| matches!(i, 4 | 7) || c.is_ascii_digit())
{
return false;
}
let num = |from: usize, to: usize| v[from..to].parse::<u32>().unwrap_or(0);
let (year, month, day) = (num(0, 4), num(5, 7), num(8, 10));
if !matches!(month, 1..=12) || day == 0 {
return false;
}
let leap = year % 4 == 0 && (year % 100 != 0 || year % 400 == 0);
let days = match month {
1 | 3 | 5 | 7 | 8 | 10 | 12 => 31,
4 | 6 | 9 | 11 => 30,
_ if leap => 29,
_ => 28,
};
day <= days
}
fn classify(version: &str) -> EraResolution {
if SUPPORTED_VERSIONS.contains(&version) {
EraResolution::Known(version.to_string())
} else {
EraResolution::Unknown(UnknownReason::UnsupportedVersion(version.to_string()))
}
}
pub(crate) fn resolve_era(
envelope: &EnvelopeObservation,
metadata: &RequestMetadata,
) -> EraResolution {
let malformed = EraResolution::Unknown(UnknownReason::MalformedSignal);
let header = match envelope {
EnvelopeObservation::Present(v) => Some(v.as_str()),
EnvelopeObservation::Absent | EnvelopeObservation::NotApplicable => None,
EnvelopeObservation::Malformed => return malformed,
};
let body = match metadata {
RequestMetadata::Present(v) => Some(v.as_str()),
RequestMetadata::Absent => None,
RequestMetadata::Malformed => return malformed,
};
if [header, body]
.into_iter()
.flatten()
.any(|v| !is_version_shaped(v))
{
return malformed;
}
match (header, body) {
(Some(h), Some(b)) if h != b => EraResolution::Conflicting {
header: h.to_string(),
body: b.to_string(),
},
(Some(v), _) | (None, Some(v)) => classify(v),
(None, None) => EraResolution::Unknown(UnknownReason::NoSignal),
}
}
fn requires_result_type(version: &str) -> bool {
version >= RESULT_TYPE_SINCE
}
pub(crate) fn conclude(
era: &EraResolution,
observed: &ResultObservation,
capability: Option<&CapabilityObservation>,
) -> ResultConclusion {
if matches!(observed, ResultObservation::Malformed) {
return ResultConclusion::Invalid(InvalidReason::MalformedResultType);
}
let version = match era {
EraResolution::Known(v) => v,
EraResolution::Unknown(UnknownReason::MalformedSignal) => {
return ResultConclusion::Invalid(InvalidReason::MalformedEraSignal)
}
EraResolution::Unknown(reason) => {
return ResultConclusion::Incomplete(IncompleteReason::EraUnknown(reason.clone()))
}
EraResolution::Conflicting { header, body } => {
return ResultConclusion::Invalid(InvalidReason::EraConflicting {
header: header.clone(),
body: body.clone(),
})
}
};
let modern = requires_result_type(version);
if modern && matches!(capability, Some(CapabilityObservation::Malformed)) {
return ResultConclusion::Invalid(InvalidReason::MalformedCapabilities);
}
match observed {
ResultObservation::Complete => ResultConclusion::Terminal,
ResultObservation::InputRequired => ResultConclusion::NonTerminal,
ResultObservation::CompleteWithContinuation => {
ResultConclusion::Incomplete(IncompleteReason::ContradictoryResult)
}
ResultObservation::InputRequiredWithoutContinuation => {
ResultConclusion::Invalid(InvalidReason::UncontinuableInputRequired)
}
ResultObservation::InputRequiredWithMalformedContinuation => {
ResultConclusion::Invalid(InvalidReason::MalformedContinuation)
}
ResultObservation::Unrecognized if !modern => {
ResultConclusion::Incomplete(IncompleteReason::UnrecognizedResultType)
}
ResultObservation::Unrecognized => match capability {
Some(CapabilityObservation::CoreOnly) => {
ResultConclusion::Incomplete(IncompleteReason::UnrecognizedResultType)
}
Some(CapabilityObservation::Absent | CapabilityObservation::ExtensionNotUnderstood) => {
ResultConclusion::Incomplete(IncompleteReason::RecognitionUndeterminable)
}
None => ResultConclusion::Incomplete(IncompleteReason::RecognitionUndeterminable),
Some(CapabilityObservation::Malformed) => unreachable!("handled above"),
},
ResultObservation::Malformed => {
ResultConclusion::Invalid(InvalidReason::MalformedResultType)
}
ResultObservation::Missing => {
if modern {
ResultConclusion::Invalid(InvalidReason::MissingResultType)
} else {
ResultConclusion::Terminal
}
}
}
}
pub(crate) fn conclude_request(
era: &EraResolution,
metadata: &RequestMetadata,
capability: Option<&CapabilityObservation>,
) -> RequestAssessment {
if matches!(metadata, RequestMetadata::Malformed) {
return RequestAssessment::Invalid(InvalidReason::MalformedRequestMetadata);
}
match era {
EraResolution::Conflicting { header, body } => {
RequestAssessment::Invalid(InvalidReason::EraConflicting {
header: header.clone(),
body: body.clone(),
})
}
EraResolution::Unknown(UnknownReason::MalformedSignal) => {
RequestAssessment::Invalid(InvalidReason::MalformedEraSignal)
}
EraResolution::Unknown(reason) => {
RequestAssessment::Incomplete(IncompleteReason::EraUnknown(reason.clone()))
}
EraResolution::Known(version) if !requires_result_type(version) => RequestAssessment::Valid,
EraResolution::Known(_) => match metadata {
RequestMetadata::Present(_) => match capability {
Some(CapabilityObservation::CoreOnly)
| Some(CapabilityObservation::ExtensionNotUnderstood) => RequestAssessment::Valid,
Some(CapabilityObservation::Malformed) => {
RequestAssessment::Invalid(InvalidReason::MalformedCapabilities)
}
Some(CapabilityObservation::Absent) | None => {
RequestAssessment::Invalid(InvalidReason::MissingCapabilities)
}
},
RequestMetadata::Absent => {
RequestAssessment::Invalid(InvalidReason::MissingRequestMetadata)
}
RequestMetadata::Malformed => unreachable!("handled above"),
},
}
}
pub(crate) const PROTOCOL_VERSION_META_KEY: &str = "io.modelcontextprotocol/protocolVersion";
pub(crate) const CLIENT_CAPABILITIES_META_KEY: &str = "io.modelcontextprotocol/clientCapabilities";
const CAPABILITY_ELICITATION_KEY: &str = "elicitation";
const CAPABILITY_ROOTS_KEY: &str = "roots";
const CAPABILITY_SAMPLING_KEY: &str = "sampling";
const CAPABILITY_EXPERIMENTAL_KEY: &str = "experimental";
const CAPABILITY_EXTENSIONS_KEY: &str = "extensions";
pub(crate) const PROTOCOL_VERSION_HEADER: &str = "MCP-Protocol-Version";
pub(crate) fn observe_header(headers: Option<&serde_json::Value>) -> Option<EnvelopeObservation> {
let headers = headers?;
let Some(map) = headers.as_object() else {
return Some(EnvelopeObservation::Malformed);
};
let mut found: Option<&str> = None;
let mut any = false;
for (_, value) in map
.iter()
.filter(|(k, _)| k.eq_ignore_ascii_case(PROTOCOL_VERSION_HEADER))
{
any = true;
match value.as_str() {
Some(v) if is_version_shaped(v) && found.is_none_or(|prev| prev == v) => {
found = Some(v)
}
_ => return Some(EnvelopeObservation::Malformed),
}
}
if !any {
return None;
}
Some(match found {
Some(v) => EnvelopeObservation::Present(v.to_string()),
None => EnvelopeObservation::Malformed,
})
}
pub(crate) fn fold_envelope(
observations: Vec<EnvelopeObservation>,
framed: bool,
) -> EnvelopeObservation {
if observations.is_empty() {
return if framed {
EnvelopeObservation::Absent
} else {
EnvelopeObservation::NotApplicable
};
}
let mut seen: Option<&str> = None;
for o in &observations {
#[expect(
clippy::wildcard_enum_match_arm,
reason = "anything that is not a single consistent Present observation is malformed, including Absent and NotApplicable; folding them to malformed is the fail-closed direction and a new variant inherits it"
)]
match o {
EnvelopeObservation::Present(v) => match seen {
Some(prev) if prev != v => return EnvelopeObservation::Malformed,
_ => seen = Some(v),
},
_ => return EnvelopeObservation::Malformed,
}
}
seen.map(|v| EnvelopeObservation::Present(v.to_string()))
.unwrap_or(EnvelopeObservation::Malformed)
}
pub(crate) fn observe_request_metadata(raw: &serde_json::Value) -> RequestMetadata {
let Some(params) = raw.get("params") else {
return RequestMetadata::Absent;
};
if !params.is_object() {
return RequestMetadata::Malformed;
}
let Some(meta) = params.get("_meta") else {
return RequestMetadata::Absent;
};
if !meta.is_object() {
return RequestMetadata::Malformed;
}
match meta.get(PROTOCOL_VERSION_META_KEY) {
None => RequestMetadata::Absent,
Some(v) => match v.as_str() {
Some(s) if is_version_shaped(s) => RequestMetadata::Present(s.to_string()),
_ => RequestMetadata::Malformed,
},
}
}
pub(crate) fn observe_client_capabilities(
raw: &serde_json::Value,
) -> Option<CapabilityObservation> {
let params = raw.get("params")?;
if !params.is_object() {
return Some(CapabilityObservation::Malformed);
}
let meta = params.get("_meta")?;
if !meta.is_object() {
return Some(CapabilityObservation::Malformed);
}
Some(match meta.get(CLIENT_CAPABILITIES_META_KEY) {
None => CapabilityObservation::Absent,
Some(caps) => {
let Some(caps) = caps.as_object() else {
return Some(CapabilityObservation::Malformed);
};
let mut beyond_core = false;
for (name, value) in caps {
match name.as_str() {
CAPABILITY_ELICITATION_KEY | CAPABILITY_ROOTS_KEY | CAPABILITY_SAMPLING_KEY => {
if !value.is_object() {
return Some(CapabilityObservation::Malformed);
}
}
CAPABILITY_EXPERIMENTAL_KEY | CAPABILITY_EXTENSIONS_KEY => {
let Some(map) = value.as_object() else {
return Some(CapabilityObservation::Malformed);
};
beyond_core |= !map.is_empty();
}
_ => beyond_core = true,
}
}
if beyond_core {
CapabilityObservation::ExtensionNotUnderstood
} else {
CapabilityObservation::CoreOnly
}
}
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ContinuationShape {
Absent,
Present,
Malformed,
}
fn continuation_shape(result: &serde_json::Value) -> ContinuationShape {
let mut present = false;
match result.get(CONTINUATION_REQUEST_STATE_KEY) {
None => {}
Some(v) if v.is_null() => {}
Some(v) if v.is_string() => present = true,
Some(_) => return ContinuationShape::Malformed,
}
match result.get(CONTINUATION_INPUT_REQUESTS_KEY) {
None => {}
Some(v) if v.is_null() => {}
Some(v) if v.is_object() => present = true,
Some(_) => return ContinuationShape::Malformed,
}
if present {
ContinuationShape::Present
} else {
ContinuationShape::Absent
}
}
const CONTINUATION_INPUT_REQUESTS_KEY: &str = "inputRequests";
const CONTINUATION_REQUEST_STATE_KEY: &str = "requestState";
pub(crate) fn observe_result(raw: &serde_json::Value) -> Option<ResultObservation> {
let result = raw.get("result")?;
if !result.is_object() {
return Some(ResultObservation::Malformed);
}
Some(match result.get("resultType") {
None => ResultObservation::Missing,
Some(v) => match v.as_str() {
Some("complete") => match continuation_shape(result) {
ContinuationShape::Absent => ResultObservation::Complete,
ContinuationShape::Present | ContinuationShape::Malformed => {
ResultObservation::CompleteWithContinuation
}
},
Some("input_required") => match continuation_shape(result) {
ContinuationShape::Present => ResultObservation::InputRequired,
ContinuationShape::Absent => ResultObservation::InputRequiredWithoutContinuation,
ContinuationShape::Malformed => {
ResultObservation::InputRequiredWithMalformedContinuation
}
},
Some(_) => ResultObservation::Unrecognized,
None => ResultObservation::Malformed,
},
})
}
#[cfg(test)]
mod tests;