Skip to content

Commit d41d3c8

Browse files
committed
fix(transport): downgrade idle timeout log from error to debug
Idle keep-alive timeout is normal zombie-session cleanup, not a transport failure. Route it through a dedicated WorkerQuitReason::IdleTimeout variant. Log it at debug level instead of treating it as a fatal error. Remove the unused LocalSessionWorkerError::KeepAliveTimeout variant. Closes #817
1 parent 63583b1 commit d41d3c8

3 files changed

Lines changed: 155 additions & 5 deletions

File tree

crates/rmcp/src/transport/streamable_http_server/session/local.rs

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -928,8 +928,6 @@ pub enum LocalSessionWorkerError {
928928
FailToSendInitializeRequest(SessionError),
929929
#[error("fail to handle message: {0}")]
930930
FailToHandleMessage(SessionError),
931-
#[error("keep alive timeout after {}ms", _0.as_millis())]
932-
KeepAliveTimeout(Duration),
933931
#[error("Transport closed")]
934932
TransportClosed,
935933
#[error("Tokio join error {0}")]
@@ -1008,7 +1006,7 @@ impl Worker for LocalSessionWorker {
10081006
return Err(WorkerQuitReason::Cancelled)
10091007
}
10101008
_ = keep_alive_timeout => {
1011-
return Err(WorkerQuitReason::fatal(LocalSessionWorkerError::KeepAliveTimeout(keep_alive), "poll next session event"))
1009+
return Err(WorkerQuitReason::IdleTimeout(keep_alive))
10121010
}
10131011
};
10141012
match event {

crates/rmcp/src/transport/worker.rs

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
use std::borrow::Cow;
1+
use std::{borrow::Cow, time::Duration};
22

33
use tokio_util::sync::CancellationToken;
44
use tracing::{Instrument, Level};
@@ -22,6 +22,8 @@ pub enum WorkerQuitReason<E> {
2222
TransportClosed,
2323
#[error("Handler terminated")]
2424
HandlerTerminated,
25+
#[error("Worker idle timeout ({}ms)", _0.as_millis())]
26+
IdleTimeout(Duration),
2527
}
2628

2729
impl<E: std::error::Error + Send + 'static> WorkerQuitReason<E> {
@@ -122,7 +124,8 @@ impl<W: Worker> WorkerTransport<W> {
122124
.inspect_err(|e| match e {
123125
WorkerQuitReason::Cancelled
124126
| WorkerQuitReason::TransportClosed
125-
| WorkerQuitReason::HandlerTerminated => {
127+
| WorkerQuitReason::HandlerTerminated
128+
| WorkerQuitReason::IdleTimeout(_) => {
126129
tracing::debug!("worker quit with reason: {:?}", e);
127130
}
128131
WorkerQuitReason::Join(e) => {
Lines changed: 149 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,149 @@
1+
#![cfg(all(
2+
feature = "transport-streamable-http-server",
3+
feature = "transport-streamable-http-client-reqwest",
4+
not(feature = "local")
5+
))]
6+
7+
use std::{
8+
sync::{Arc, Mutex},
9+
time::Duration,
10+
};
11+
12+
use rmcp::transport::streamable_http_server::{
13+
StreamableHttpServerConfig, StreamableHttpService, session::local::LocalSessionManager,
14+
};
15+
use tokio_util::sync::CancellationToken;
16+
use tracing_subscriber::layer::SubscriberExt;
17+
18+
mod common;
19+
use common::calculator::Calculator;
20+
21+
// Issue #817: keep-alive timeout emits tracing::error! for normal idle reaping.
22+
23+
struct CapturedEvent {
24+
level: tracing::Level,
25+
message: String,
26+
}
27+
28+
struct CapturingLayer {
29+
events: Arc<Mutex<Vec<CapturedEvent>>>,
30+
}
31+
32+
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for CapturingLayer {
33+
fn on_event(
34+
&self,
35+
event: &tracing::Event<'_>,
36+
_ctx: tracing_subscriber::layer::Context<'_, S>,
37+
) {
38+
let mut visitor = MessageVisitor(String::new());
39+
event.record(&mut visitor);
40+
self.events.lock().unwrap().push(CapturedEvent {
41+
level: *event.metadata().level(),
42+
message: visitor.0,
43+
});
44+
}
45+
}
46+
47+
struct MessageVisitor(String);
48+
49+
impl tracing::field::Visit for MessageVisitor {
50+
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
51+
if field.name() == "message" {
52+
self.0 = format!("{:?}", value);
53+
}
54+
}
55+
}
56+
57+
#[tokio::test(flavor = "current_thread")]
58+
async fn test_keep_alive_timeout_does_not_emit_error_log() {
59+
let events = Arc::new(Mutex::new(Vec::<CapturedEvent>::new()));
60+
61+
let subscriber = tracing_subscriber::registry().with(CapturingLayer {
62+
events: events.clone(),
63+
});
64+
65+
let _guard = tracing::subscriber::set_default(subscriber);
66+
67+
let ct = CancellationToken::new();
68+
let mut session_manager = LocalSessionManager::default();
69+
session_manager.session_config.keep_alive = Some(Duration::from_millis(200));
70+
let session_manager = Arc::new(session_manager);
71+
72+
let service = StreamableHttpService::new(
73+
|| Ok(Calculator::new()),
74+
session_manager.clone(),
75+
StreamableHttpServerConfig::default()
76+
.with_sse_keep_alive(None)
77+
.with_cancellation_token(ct.child_token()),
78+
);
79+
80+
let router = axum::Router::new().nest_service("/mcp", service);
81+
let tcp_listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
82+
let addr = tcp_listener.local_addr().unwrap();
83+
84+
tokio::spawn({
85+
let ct = ct.clone();
86+
async move {
87+
let _ = axum::serve(tcp_listener, router)
88+
.with_graceful_shutdown(async move { ct.cancelled_owned().await })
89+
.await;
90+
}
91+
});
92+
93+
let client = reqwest::Client::new();
94+
95+
// Initialize session
96+
let response = client
97+
.post(format!("http://{addr}/mcp"))
98+
.header("Content-Type", "application/json")
99+
.header("Accept", "application/json, text/event-stream")
100+
.body(r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"test","version":"1.0"}}}"#)
101+
.send()
102+
.await
103+
.unwrap();
104+
assert_eq!(response.status(), 200);
105+
let session_id = response.headers()["mcp-session-id"]
106+
.to_str()
107+
.unwrap()
108+
.to_string();
109+
110+
// Complete handshake
111+
client
112+
.post(format!("http://{addr}/mcp"))
113+
.header("Content-Type", "application/json")
114+
.header("Accept", "application/json, text/event-stream")
115+
.header("mcp-session-id", &session_id)
116+
.header("Mcp-Protocol-Version", "2025-06-18")
117+
.body(r#"{"jsonrpc":"2.0","method":"notifications/initialized"}"#)
118+
.send()
119+
.await
120+
.unwrap();
121+
122+
// Wait for keep_alive timeout (200ms) plus margin
123+
tokio::time::sleep(Duration::from_millis(400)).await;
124+
125+
let captured = events.lock().unwrap();
126+
127+
let error_events: Vec<_> = captured
128+
.iter()
129+
.filter(|e| e.level == tracing::Level::ERROR)
130+
.filter(|e| e.message.contains("keep alive timeout") || e.message.contains("IdleTimeout"))
131+
.collect();
132+
assert!(
133+
error_events.is_empty(),
134+
"keep-alive timeout should not produce ERROR logs, found {}: {:?}",
135+
error_events.len(),
136+
error_events.iter().map(|e| &e.message).collect::<Vec<_>>()
137+
);
138+
139+
let debug_events: Vec<_> = captured
140+
.iter()
141+
.filter(|e| e.level == tracing::Level::DEBUG && e.message.contains("IdleTimeout"))
142+
.collect();
143+
assert!(
144+
!debug_events.is_empty(),
145+
"expected a DEBUG log with IdleTimeout, but found none"
146+
);
147+
148+
ct.cancel();
149+
}

0 commit comments

Comments
 (0)