use crate::mcp::era::{
classify_message, correlation_id, fold_envelope, id_is_acceptable, observe_client_capabilities,
observe_header, observe_request_metadata, observe_result, resolve_era, CapabilityObservation,
EnvelopeObservation, McpEraContext, MessageKind, ParsedMcpEvent, RequestMetadata,
};
use crate::mcp::era::{CorrelationId, DuplicateAwareSink, EraResolution, SeenMembers, UniqueValue};
use crate::mcp::json_depth;
use crate::mcp::types::*;
use anyhow::{bail, Context, Result};
use serde::Deserialize;
pub fn parse_mcp_transcript(text: &str, format: McpInputFormat) -> Result<Vec<McpEvent>> {
Ok(parse_mcp_transcript_detailed(text, format)?
.into_iter()
.map(|parsed| parsed.event)
.collect())
}
pub(crate) fn parse_mcp_transcript_detailed(
text: &str,
format: McpInputFormat,
) -> Result<Vec<ParsedMcpEvent>> {
parse_mcp_transcript_detailed_internal(text, format, None)
}
pub(crate) fn parse_mcp_transcript_detailed_with_depth(
text: &str,
format: McpInputFormat,
max_json_depth: usize,
) -> Result<Vec<ParsedMcpEvent>> {
parse_mcp_transcript_detailed_internal(text, format, Some(max_json_depth))
}
fn parse_mcp_transcript_detailed_internal(
text: &str,
format: McpInputFormat,
max_embedded_json_depth: Option<usize>,
) -> Result<Vec<ParsedMcpEvent>> {
let (events, envelopes) = parse_events_with_envelopes(text, format, max_embedded_json_depth)?;
let framed = is_framed(format);
let parsed: Vec<ParsedMcpEvent> = events
.into_iter()
.map(|event| {
let (request_metadata, result_observation, capability_observation, is_error_response) =
match payload_raw(&event.payload) {
Some(raw) => match classify_message(raw) {
Ok(MessageKind::Request { .. }) => (
Some(observe_request_metadata(raw)),
None,
observe_client_capabilities(raw),
false,
),
Ok(MessageKind::Notification { .. }) => (None, None, None, false),
Ok(MessageKind::Response) => {
(None, observe_result(raw), None, raw.get("error").is_some())
}
Err(_) => (None, None, None, false),
},
None => (None, None, None, false),
};
let envelope = if framed {
envelopes
.get(event.source_line.saturating_sub(1) as usize)
.cloned()
.unwrap_or(EnvelopeObservation::Malformed)
} else {
EnvelopeObservation::NotApplicable
};
let era = resolve_era(
&envelope,
request_metadata
.as_ref()
.unwrap_or(&RequestMetadata::Absent),
);
let correlation = payload_raw(&event.payload).and_then(correlation_id);
ParsedMcpEvent {
event,
context: McpEraContext {
envelope,
era,
correlation,
request_metadata,
result_observation,
capability_observation,
},
is_error_response,
}
})
.collect();
correlate_calls(parsed)
}
struct CallSignals {
era: EraResolution,
capability: Option<CapabilityObservation>,
}
fn correlate_calls(mut parsed: Vec<ParsedMcpEvent>) -> Result<Vec<ParsedMcpEvent>> {
let mut outstanding: std::collections::HashMap<CorrelationId, CallSignals> =
std::collections::HashMap::new();
for p in &mut parsed {
let Some(raw) = payload_raw(&p.event.payload) else {
continue;
};
let Some(id) = correlation_id(raw) else {
continue;
};
let kind = classify_message(raw).ok();
match kind {
Some(MessageKind::Request { .. }) => {
let signals = CallSignals {
era: p.context.era.clone(),
capability: p.context.capability_observation.clone(),
};
if outstanding.insert(id.clone(), signals).is_some() {
bail!(
"two outstanding JSON-RPC requests share an id at source line {}",
p.event.source_line
);
}
}
Some(MessageKind::Response) => {
if let Some(signals) = outstanding.remove(&id) {
p.context.era = signals.era;
p.context.capability_observation = signals.capability;
}
}
Some(MessageKind::Notification { .. }) | None => {}
}
}
Ok(parsed)
}
fn observe_transport_context(ctx: Option<&serde_json::Value>) -> Option<EnvelopeObservation> {
let ctx = ctx?;
let Some(map) = ctx.as_object() else {
return Some(EnvelopeObservation::Malformed);
};
observe_header(map.get("headers"))
}
fn is_framed(format: McpInputFormat) -> bool {
matches!(
format,
McpInputFormat::StreamableHttp | McpInputFormat::HttpSse
)
}
fn payload_raw(payload: &McpPayload) -> Option<&serde_json::Value> {
match payload {
McpPayload::SessionStart { raw }
| McpPayload::ToolsListRequest { raw }
| McpPayload::ToolsListResponse { raw, .. }
| McpPayload::ToolCallRequest { raw, .. }
| McpPayload::ToolCallResponse { raw, .. }
| McpPayload::SessionEnd { raw, .. }
| McpPayload::Other { raw, .. } => Some(raw),
}
}
fn parse_events_with_envelopes(
text: &str,
format: McpInputFormat,
max_embedded_json_depth: Option<usize>,
) -> Result<(Vec<McpEvent>, Vec<EnvelopeObservation>)> {
let (events, envelopes) = match format {
McpInputFormat::JsonRpc => (parse_jsonrpc_jsonl(text)?, Vec::new()),
McpInputFormat::Inspector => (parse_inspector_best_effort(text)?, Vec::new()),
McpInputFormat::StreamableHttp => parse_transport_transcript_detailed(
text,
"streamable-http",
"streamable-http transcript",
false,
max_embedded_json_depth,
)?,
McpInputFormat::HttpSse => parse_transport_transcript_detailed(
text,
"http-sse",
"http-sse transcript",
true,
max_embedded_json_depth,
)?,
};
Ok((events, envelopes))
}
fn parse_jsonrpc_jsonl(text: &str) -> Result<Vec<McpEvent>> {
let mut out = Vec::new();
for (lineno, line) in text.lines().enumerate() {
let line = line.trim();
if line.is_empty() {
continue;
}
let UniqueValue(v) = serde_json::from_str::<UniqueValue>(line)
.with_context(|| format!("invalid JSON on line {}", lineno + 1))?;
let event = parse_jsonrpc_message(
v,
(lineno + 1) as u64,
None,
McpAuthorizationDiscovery::default(),
)?;
out.push(event);
}
Ok(out)
}
fn parse_inspector_best_effort(text: &str) -> Result<Vec<McpEvent>> {
let UniqueValue(v) =
serde_json::from_str::<UniqueValue>(text).context("invalid inspector JSON")?;
let arr = v
.get("events")
.cloned()
.or_else(|| v.as_array().cloned().map(serde_json::Value::Array))
.and_then(|x| x.as_array().cloned())
.unwrap_or_default();
let mut out = Vec::new();
for (idx, item) in arr.into_iter().enumerate() {
let event = parse_jsonrpc_message(
item,
(idx + 1) as u64,
None,
McpAuthorizationDiscovery::default(),
)?;
out.push(event);
}
Ok(out)
}
fn parse_transport_transcript_detailed(
text: &str,
expected_transport: &str,
source_label: &str,
allow_endpoint_event: bool,
max_embedded_json_depth: Option<usize>,
) -> Result<(Vec<McpEvent>, Vec<EnvelopeObservation>)> {
let transcript: TransportTranscript =
serde_json::from_str(text).with_context(|| format!("invalid {}", source_label))?;
let actual_transport = transcript.transport.as_deref().unwrap_or("missing");
if actual_transport != expected_transport {
bail!(
"{} transport must be {:?}, found {:?}",
source_label,
expected_transport,
actual_transport
);
}
let mut transcript_slots = Vec::new();
if let Some(o) = observe_transport_context(transcript.transport_context.as_ref().map(|u| &u.0))
{
transcript_slots.push(o);
}
if let Some(o) = observe_header(transcript.headers.as_ref().map(|u| &u.0)) {
transcript_slots.push(o);
}
let mut envelopes = Vec::new();
let mut out = Vec::new();
for (idx, entry) in transcript.entries.into_iter().enumerate() {
let mut slots = transcript_slots.clone();
if let Some(o) = observe_transport_context(entry.transport_context.as_ref().map(|u| &u.0)) {
slots.push(o);
}
if let Some(o) = observe_header(entry.headers.as_ref().map(|u| &u.0)) {
slots.push(o);
}
envelopes.push(fold_envelope(slots, true));
let source_line = (idx + 1) as u64;
let present = usize::from(entry.request.is_some())
+ usize::from(entry.response.is_some())
+ usize::from(entry.sse.is_some());
if present != 1 {
bail!(
"{} entry {} must contain exactly one of request, response, or sse",
source_label,
source_line
);
}
if let Some(UniqueValue(request)) = entry.request {
out.push(parse_jsonrpc_message(
request,
source_line,
entry.timestamp_ms,
McpAuthorizationDiscovery::default(),
)?);
continue;
}
let auth_discovery = parse_transport_auth_discovery(&entry);
if let Some(UniqueValue(response)) = entry.response {
out.push(parse_jsonrpc_message(
response,
source_line,
entry.timestamp_ms,
auth_discovery,
)?);
continue;
}
if let Some(sse) = entry.sse {
if let Some(jsonrpc) =
extract_jsonrpc_from_sse(&sse, allow_endpoint_event, max_embedded_json_depth)?
{
out.push(parse_jsonrpc_message(
jsonrpc,
source_line,
entry.timestamp_ms,
McpAuthorizationDiscovery::default(),
)?);
}
}
}
Ok((out, envelopes))
}
fn parse_jsonrpc_message(
v: serde_json::Value,
source_line: u64,
timestamp_ms_override: Option<u64>,
auth_discovery: McpAuthorizationDiscovery,
) -> Result<McpEvent> {
if !v.is_object() {
bail!(
"MCP event at source line {} must be a JSON object",
source_line
);
}
if v.get("jsonrpc").and_then(serde_json::Value::as_str) != Some("2.0") {
bail!(
"MCP event at source line {} must carry JSON-RPC version 2.0",
source_line
);
}
let ts_ms = timestamp_ms_override.or_else(|| extract_ts_ms(&v));
let kind = classify_message(&v)
.map_err(|e| anyhow::anyhow!("MCP event at source line {}: {}", source_line, e))?;
let raw_id = v.get("id");
let id_str = match kind {
MessageKind::Notification { .. } => None,
MessageKind::Request { .. } => Some(require_acceptable_id(raw_id, source_line)?),
MessageKind::Response => {
let is_error = v.get("error").is_some();
match raw_id {
None | Some(serde_json::Value::Null) if is_error => None,
_ => Some(require_acceptable_id(raw_id, source_line)?),
}
}
};
let payload =
if let MessageKind::Request { method } | MessageKind::Notification { method } = kind {
match method {
"tools/list" => McpPayload::ToolsListRequest { raw: v.clone() },
"tools/call" => {
let params = v.get("params").cloned().unwrap_or(serde_json::Value::Null);
let name = params
.get("name")
.and_then(|x| x.as_str())
.unwrap_or("unknown_tool")
.to_string();
let arguments = params
.get("arguments")
.cloned()
.unwrap_or(serde_json::Value::Null);
McpPayload::ToolCallRequest {
name,
arguments,
raw: v.clone(),
}
}
_ => McpPayload::Other { raw: v.clone() },
}
} else {
if v.get("result").is_some() {
if looks_like_tools_list_result(&v) {
let tools = parse_tools_list_result(&v)?;
McpPayload::ToolsListResponse {
tools,
raw: v.clone(),
}
} else {
McpPayload::ToolCallResponse {
result: v.get("result").cloned().unwrap_or(serde_json::Value::Null),
is_error: false,
raw: v.clone(),
}
}
} else if v.get("error").is_some() {
McpPayload::ToolCallResponse {
result: v.get("error").cloned().unwrap_or(serde_json::Value::Null),
is_error: true,
raw: v.clone(),
}
} else {
McpPayload::Other { raw: v.clone() }
}
};
Ok(McpEvent {
source_line,
timestamp_ms: ts_ms,
jsonrpc_id: id_str,
auth_discovery,
payload,
})
}
fn parse_transport_auth_discovery(entry: &TransportTranscriptEntry) -> McpAuthorizationDiscovery {
let Some(status) = extract_http_status(entry) else {
return McpAuthorizationDiscovery::default();
};
if status != 401 {
return McpAuthorizationDiscovery::default();
}
let header_value = entry
.transport_context
.as_ref()
.and_then(|value| find_header_case_insensitive(&value.0, "www-authenticate"))
.or_else(|| {
entry
.headers
.as_ref()
.and_then(|value| find_header_case_insensitive(&value.0, "www-authenticate"))
});
let Some(www_authenticate) = header_value else {
return McpAuthorizationDiscovery::default();
};
let resource_metadata_visible = auth_param_visible(&www_authenticate, "resource_metadata");
let scope_challenge_visible = auth_param_visible(&www_authenticate, "scope");
if !resource_metadata_visible && !scope_challenge_visible {
return McpAuthorizationDiscovery::default();
}
McpAuthorizationDiscovery {
visible: true,
source_kind: McpAuthorizationDiscoverySourceKind::WwwAuthenticate,
resource_metadata_visible,
authorization_servers_visible: false,
scope_challenge_visible,
}
}
fn extract_http_status(entry: &TransportTranscriptEntry) -> Option<u16> {
entry
.transport_context
.as_ref()
.and_then(|v| extract_http_status_from_value(&v.0))
.or_else(|| {
entry
.headers
.as_ref()
.and_then(|v| extract_http_status_from_value(&v.0))
})
}
fn extract_http_status_from_value(value: &serde_json::Value) -> Option<u16> {
match value {
serde_json::Value::Object(map) => {
for key in ["status", "status_code", "http_status"] {
if let Some(status) = map.get(key).and_then(json_value_to_u16) {
return Some(status);
}
}
map.get("response").and_then(extract_http_status_from_value)
}
_ => None,
}
}
fn json_value_to_u16(value: &serde_json::Value) -> Option<u16> {
match value {
serde_json::Value::Number(n) => n.as_u64().and_then(|n| u16::try_from(n).ok()),
serde_json::Value::String(s) => s.parse::<u16>().ok(),
_ => None,
}
}
fn find_header_case_insensitive(value: &serde_json::Value, header_name: &str) -> Option<String> {
match value {
serde_json::Value::Object(map) => {
if let Some(headers) = map.get("headers") {
if let Some(found) = find_header_case_insensitive(headers, header_name) {
return Some(found);
}
}
if let Some(response) = map.get("response") {
if let Some(found) = find_header_case_insensitive(response, header_name) {
return Some(found);
}
}
map.iter().find_map(|(key, value)| {
if key.eq_ignore_ascii_case(header_name) {
value.as_str().map(ToString::to_string)
} else {
None
}
})
}
_ => None,
}
}
fn auth_param_visible(header_value: &str, param_name: &str) -> bool {
let lower = header_value.to_ascii_lowercase();
let needle = format!("{param_name}=");
lower
.match_indices(&needle)
.any(|(idx, _)| idx == 0 || matches!(lower.as_bytes()[idx - 1], b' ' | b',' | b'\t'))
}
fn require_acceptable_id(raw_id: Option<&serde_json::Value>, source_line: u64) -> Result<String> {
reject_unusable_id_shape(raw_id, source_line)?;
match raw_id.filter(|v| id_is_acceptable(v)) {
Some(serde_json::Value::String(id)) => Ok(id.clone()),
Some(serde_json::Value::Number(id)) => Ok(id.to_string()),
_ => bail!(
"JSON-RPC id on source line {} must be a string or a number",
source_line
),
}
}
fn reject_unusable_id_shape(raw_id: Option<&serde_json::Value>, source_line: u64) -> Result<()> {
match raw_id {
None
| Some(serde_json::Value::Null)
| Some(serde_json::Value::String(_))
| Some(serde_json::Value::Number(_)) => Ok(()),
Some(serde_json::Value::Bool(_)) => {
bail!(
"JSON-RPC id on source line {} must not be a boolean",
source_line
)
}
Some(serde_json::Value::Array(_)) => {
bail!(
"JSON-RPC id on source line {} must not be an array",
source_line
)
}
Some(serde_json::Value::Object(_)) => {
bail!(
"JSON-RPC id on source line {} must not be an object",
source_line
)
}
}
}
fn extract_jsonrpc_from_sse(
sse: &TransportSseEnvelope,
allow_endpoint_event: bool,
max_embedded_json_depth: Option<usize>,
) -> Result<Option<serde_json::Value>> {
let event_name = sse.event.as_deref().unwrap_or("message");
if event_name == "endpoint" && allow_endpoint_event {
return Ok(None);
}
if event_name != "message" {
return Ok(None);
}
extract_jsonrpc_like_value(&sse.data.0, max_embedded_json_depth)
}
#[derive(Debug, thiserror::Error)]
#[error("embedded SSE JSON exceeds its depth limit")]
pub(crate) struct EmbeddedJsonDepthExceeded;
fn extract_jsonrpc_like_value(
value: &serde_json::Value,
max_embedded_json_depth: Option<usize>,
) -> Result<Option<serde_json::Value>> {
match value {
serde_json::Value::Object(map)
if map.contains_key("method")
|| map.contains_key("result")
|| map.contains_key("error")
|| map.contains_key("jsonrpc") =>
{
Ok(Some(value.clone()))
}
serde_json::Value::String(text) => {
if max_embedded_json_depth
.is_some_and(|limit| json_depth::exceeds_limit(text.as_bytes(), limit, false))
{
return Err(EmbeddedJsonDepthExceeded.into());
}
match serde_json::from_str::<UniqueValue>(text) {
Ok(UniqueValue(parsed)) => {
extract_jsonrpc_like_value(&parsed, max_embedded_json_depth)
}
Err(e) if e.classify() == serde_json::error::Category::Data => {
Err(anyhow::Error::new(e).context("invalid SSE data payload"))
}
Err(_) => Ok(None),
}
}
_ => Ok(None),
}
}
fn extract_ts_ms(v: &serde_json::Value) -> Option<u64> {
if let Some(t) = v.get("timestamp_ms").and_then(|t| t.as_u64()) {
return Some(t);
}
if let Some(t) = v.get("timestamp").and_then(|t| t.as_u64()) {
return Some(t); }
None
}
fn looks_like_tools_list_result(v: &serde_json::Value) -> bool {
v.get("result")
.and_then(|r| r.get("tools"))
.and_then(|t| t.as_array())
.is_some()
}
fn parse_tools_list_result(v: &serde_json::Value) -> Result<Vec<McpToolDef>> {
let tools = v
.get("result")
.and_then(|r| r.get("tools"))
.and_then(|t| t.as_array())
.cloned()
.unwrap_or_default();
let mut out = Vec::new();
for tool in tools {
let name = tool
.get("name")
.and_then(|x| x.as_str())
.unwrap_or("unknown")
.to_string();
let description = tool
.get("description")
.and_then(|x| x.as_str())
.map(|s| s.to_string());
let input_schema = tool
.get("inputSchema")
.cloned()
.or_else(|| tool.get("input_schema").cloned());
out.push(McpToolDef {
name,
description,
input_schema,
tool_identity: None,
});
}
Ok(out)
}
#[derive(Debug, Default)]
struct TransportTranscript {
transport: Option<String>,
#[allow(dead_code)]
transport_context: Option<UniqueValue>,
#[allow(dead_code)]
headers: Option<UniqueValue>,
entries: Vec<TransportTranscriptEntry>,
}
impl<'de> Deserialize<'de> for TransportTranscript {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
struct V;
impl<'de> serde::de::Visitor<'de> for V {
type Value = TransportTranscript;
fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("a transport transcript with unique members")
}
fn visit_map<A: serde::de::MapAccess<'de>>(
self,
mut map: A,
) -> Result<TransportTranscript, A::Error> {
let mut out = TransportTranscript::default();
let mut seen = SeenMembers::default();
while let Some(key) = map.next_key::<String>()? {
seen.insert::<A::Error>(&key)?;
match key.as_str() {
"transport" => out.transport = map.next_value()?,
"transport_context" => out.transport_context = Some(map.next_value()?),
"headers" => out.headers = Some(map.next_value()?),
"entries" => out.entries = map.next_value()?,
_ => {
map.next_value::<DuplicateAwareSink>()?;
}
}
}
Ok(out)
}
}
d.deserialize_map(V)
}
}
#[derive(Debug, Default)]
struct TransportTranscriptEntry {
timestamp_ms: Option<u64>,
#[allow(dead_code)]
transport_context: Option<UniqueValue>,
#[allow(dead_code)]
headers: Option<UniqueValue>,
request: Option<UniqueValue>,
response: Option<UniqueValue>,
sse: Option<TransportSseEnvelope>,
}
impl<'de> Deserialize<'de> for TransportTranscriptEntry {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
struct V;
impl<'de> serde::de::Visitor<'de> for V {
type Value = TransportTranscriptEntry;
fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("a transport transcript entry with unique members")
}
fn visit_map<A: serde::de::MapAccess<'de>>(
self,
mut map: A,
) -> Result<TransportTranscriptEntry, A::Error> {
let mut out = TransportTranscriptEntry::default();
let mut seen = SeenMembers::default();
while let Some(key) = map.next_key::<String>()? {
seen.insert::<A::Error>(&key)?;
match key.as_str() {
"timestamp_ms" => out.timestamp_ms = map.next_value()?,
"transport_context" => out.transport_context = Some(map.next_value()?),
"headers" => out.headers = Some(map.next_value()?),
"request" => out.request = Some(map.next_value()?),
"response" => out.response = Some(map.next_value()?),
"sse" => match map.next_value::<Option<TransportSseEnvelope>>()? {
Some(envelope) => out.sse = Some(envelope),
None => {
return Err(serde::de::Error::custom(
"sse slot is present and null",
))
}
},
_ => {
map.next_value::<DuplicateAwareSink>()?;
}
}
}
Ok(out)
}
}
d.deserialize_map(V)
}
}
#[derive(Debug, Default)]
struct TransportSseEnvelope {
event: Option<String>,
#[allow(dead_code)]
id: Option<String>,
data: UniqueValue,
}
impl<'de> Deserialize<'de> for TransportSseEnvelope {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
struct V;
impl<'de> serde::de::Visitor<'de> for V {
type Value = TransportSseEnvelope;
fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("an SSE envelope with unique members")
}
fn visit_map<A: serde::de::MapAccess<'de>>(
self,
mut map: A,
) -> Result<TransportSseEnvelope, A::Error> {
let mut out = TransportSseEnvelope::default();
let mut seen = SeenMembers::default();
let mut saw_data = false;
while let Some(key) = map.next_key::<String>()? {
seen.insert::<A::Error>(&key)?;
match key.as_str() {
"event" => out.event = map.next_value()?,
"id" => out.id = map.next_value()?,
"data" => {
out.data = map.next_value()?;
saw_data = true;
}
_ => {
map.next_value::<DuplicateAwareSink>()?;
}
}
}
if !saw_data {
return Err(serde::de::Error::missing_field("data"));
}
Ok(out)
}
}
d.deserialize_map(V)
}
}
#[cfg(test)]
mod era_wiring_tests;