Skip to content

Commit 2f8d3b7

Browse files
authored
Fix/issue 817 idle timeout log level (#824)
* 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 * fix(session): tolerate dead worker in close_session Swallow SessionServiceTerminated in close_session when the worker has already exited. This prevents a spurious ERROR log during the post-exit cleanup path in spawn_session_worker. * fix(transport): address PR review feedback - deprecate KeepAliveTimeout - harden tests
1 parent 014fb2e commit 2f8d3b7

3 files changed

Lines changed: 269 additions & 6 deletions

File tree

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

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -66,9 +66,16 @@ impl SessionManager for LocalSessionManager {
6666
Ok(response)
6767
}
6868
async fn close_session(&self, id: &SessionId) -> Result<(), Self::Error> {
69-
let mut sessions = self.sessions.write().await;
70-
if let Some(handle) = sessions.remove(id) {
71-
handle.close().await?;
69+
let handle = {
70+
let mut sessions = self.sessions.write().await;
71+
sessions.remove(id)
72+
};
73+
if let Some(handle) = handle {
74+
match handle.close().await {
75+
// Worker already exited — nothing left to clean up.
76+
Ok(()) | Err(SessionError::SessionServiceTerminated) => {}
77+
Err(e) => return Err(e.into()),
78+
}
7279
}
7380
Ok(())
7481
}
@@ -928,6 +935,7 @@ pub enum LocalSessionWorkerError {
928935
FailToSendInitializeRequest(SessionError),
929936
#[error("fail to handle message: {0}")]
930937
FailToHandleMessage(SessionError),
938+
#[deprecated(note = "idle timeout now surfaces as WorkerQuitReason::IdleTimeout")]
931939
#[error("keep alive timeout after {}ms", _0.as_millis())]
932940
KeepAliveTimeout(Duration),
933941
#[error("init timeout after {}ms", _0.as_millis())]
@@ -1021,7 +1029,7 @@ impl Worker for LocalSessionWorker {
10211029
return Err(WorkerQuitReason::Cancelled)
10221030
}
10231031
_ = keep_alive_timeout => {
1024-
return Err(WorkerQuitReason::fatal(LocalSessionWorkerError::KeepAliveTimeout(keep_alive), "poll next session event"))
1032+
return Err(WorkerQuitReason::IdleTimeout(keep_alive))
10251033
}
10261034
};
10271035
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 after {}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: 252 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,252 @@
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,
14+
session::{SessionManager, local::LocalSessionManager},
15+
};
16+
use tokio_util::sync::CancellationToken;
17+
use tracing_subscriber::layer::SubscriberExt;
18+
19+
mod common;
20+
use common::calculator::Calculator;
21+
22+
struct CapturedEvent {
23+
level: tracing::Level,
24+
target: String,
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+
target: event.metadata().target().to_string(),
43+
message: visitor.0,
44+
});
45+
}
46+
}
47+
48+
struct MessageVisitor(String);
49+
50+
impl tracing::field::Visit for MessageVisitor {
51+
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
52+
if field.name() == "message" {
53+
self.0 = format!("{:?}", value);
54+
}
55+
}
56+
}
57+
58+
#[tokio::test(flavor = "current_thread")]
59+
async fn test_keep_alive_timeout_does_not_emit_error_log() {
60+
let events = Arc::new(Mutex::new(Vec::<CapturedEvent>::new()));
61+
62+
let subscriber = tracing_subscriber::registry().with(CapturingLayer {
63+
events: events.clone(),
64+
});
65+
66+
let _guard = tracing::subscriber::set_default(subscriber);
67+
68+
let ct = CancellationToken::new();
69+
let mut session_manager = LocalSessionManager::default();
70+
session_manager.session_config.keep_alive = Some(Duration::from_millis(200));
71+
let session_manager = Arc::new(session_manager);
72+
73+
let service = StreamableHttpService::new(
74+
|| Ok(Calculator::new()),
75+
session_manager.clone(),
76+
StreamableHttpServerConfig::default()
77+
.with_sse_keep_alive(None)
78+
.with_cancellation_token(ct.child_token()),
79+
);
80+
81+
let router = axum::Router::new().nest_service("/mcp", service);
82+
let tcp_listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
83+
let addr = tcp_listener.local_addr().unwrap();
84+
85+
tokio::spawn({
86+
let ct = ct.clone();
87+
async move {
88+
let _ = axum::serve(tcp_listener, router)
89+
.with_graceful_shutdown(async move { ct.cancelled_owned().await })
90+
.await;
91+
}
92+
});
93+
94+
let client = reqwest::Client::new();
95+
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+
client
111+
.post(format!("http://{addr}/mcp"))
112+
.header("Content-Type", "application/json")
113+
.header("Accept", "application/json, text/event-stream")
114+
.header("mcp-session-id", &session_id)
115+
.header("Mcp-Protocol-Version", "2025-06-18")
116+
.body(r#"{"jsonrpc":"2.0","method":"notifications/initialized"}"#)
117+
.send()
118+
.await
119+
.unwrap();
120+
121+
tokio::time::sleep(Duration::from_millis(400)).await;
122+
123+
// Wait until close_session() has completed so all logs are captured.
124+
let session_id_parsed: Arc<str> = Arc::from(session_id.as_str());
125+
for _ in 0..20 {
126+
if !session_manager
127+
.has_session(&session_id_parsed)
128+
.await
129+
.unwrap()
130+
{
131+
break;
132+
}
133+
tokio::time::sleep(Duration::from_millis(50)).await;
134+
}
135+
assert!(
136+
!session_manager
137+
.has_session(&session_id_parsed)
138+
.await
139+
.unwrap(),
140+
"session should have been removed after idle reap"
141+
);
142+
143+
let captured = events.lock().unwrap();
144+
145+
let error_events: Vec<_> = captured
146+
.iter()
147+
.filter(|e| e.level == tracing::Level::ERROR && e.target.starts_with("rmcp"))
148+
.collect();
149+
assert!(
150+
error_events.is_empty(),
151+
"idle reap should not produce any ERROR logs, found {}: {:?}",
152+
error_events.len(),
153+
error_events.iter().map(|e| &e.message).collect::<Vec<_>>()
154+
);
155+
156+
let debug_events: Vec<_> = captured
157+
.iter()
158+
.filter(|e| {
159+
e.level == tracing::Level::DEBUG
160+
&& e.target.starts_with("rmcp")
161+
&& e.message.contains("IdleTimeout")
162+
})
163+
.collect();
164+
assert!(
165+
!debug_events.is_empty(),
166+
"expected a DEBUG log with IdleTimeout, but found none"
167+
);
168+
169+
ct.cancel();
170+
}
171+
172+
#[tokio::test(flavor = "current_thread")]
173+
async fn test_explicit_close_on_live_session_succeeds() {
174+
let ct = CancellationToken::new();
175+
let mut session_manager = LocalSessionManager::default();
176+
session_manager.session_config.keep_alive = Some(Duration::from_secs(60));
177+
let session_manager = Arc::new(session_manager);
178+
179+
let service = StreamableHttpService::new(
180+
|| Ok(Calculator::new()),
181+
session_manager.clone(),
182+
StreamableHttpServerConfig::default()
183+
.with_sse_keep_alive(None)
184+
.with_cancellation_token(ct.child_token()),
185+
);
186+
187+
let router = axum::Router::new().nest_service("/mcp", service);
188+
let tcp_listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
189+
let addr = tcp_listener.local_addr().unwrap();
190+
191+
tokio::spawn({
192+
let ct = ct.clone();
193+
async move {
194+
let _ = axum::serve(tcp_listener, router)
195+
.with_graceful_shutdown(async move { ct.cancelled_owned().await })
196+
.await;
197+
}
198+
});
199+
200+
let client = reqwest::Client::new();
201+
202+
let response = client
203+
.post(format!("http://{addr}/mcp"))
204+
.header("Content-Type", "application/json")
205+
.header("Accept", "application/json, text/event-stream")
206+
.body(r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"test","version":"1.0"}}}"#)
207+
.send()
208+
.await
209+
.unwrap();
210+
assert_eq!(response.status(), 200);
211+
let session_id = response.headers()["mcp-session-id"]
212+
.to_str()
213+
.unwrap()
214+
.to_string();
215+
216+
client
217+
.post(format!("http://{addr}/mcp"))
218+
.header("Content-Type", "application/json")
219+
.header("Accept", "application/json, text/event-stream")
220+
.header("mcp-session-id", &session_id)
221+
.header("Mcp-Protocol-Version", "2025-06-18")
222+
.body(r#"{"jsonrpc":"2.0","method":"notifications/initialized"}"#)
223+
.send()
224+
.await
225+
.unwrap();
226+
227+
let session_id_parsed: Arc<str> = Arc::from(session_id.as_str());
228+
229+
assert!(
230+
session_manager
231+
.has_session(&session_id_parsed)
232+
.await
233+
.unwrap(),
234+
"session should exist before explicit close"
235+
);
236+
237+
let result = session_manager.close_session(&session_id_parsed).await;
238+
assert!(
239+
result.is_ok(),
240+
"close_session on a live worker should succeed: {result:?}"
241+
);
242+
243+
assert!(
244+
!session_manager
245+
.has_session(&session_id_parsed)
246+
.await
247+
.unwrap(),
248+
"session should not exist after explicit close"
249+
);
250+
251+
ct.cancel();
252+
}

0 commit comments

Comments
 (0)