Repository navigation
fix(streamable-http-server): release the session map lock before waiting on a session - #1323
TheSeydiCharyyev wants to merge 2 commits into
Conversation
|
The MSRV failure isn't from your change. It was fixed on |
…ing on a session LocalSessionManager kept the read guard on its session map while it waited on a session worker. initialize_session keeps it for the whole ServerHandler::initialize, so one slow initialize blocks create_session for every new client. tokio's RwLock is fair, so has_session calls for other sessions then wait behind that writer, and the whole server waits until the slow initialize returns. Clone the session handle out of the map and drop the guard before waiting, like close_session already does.
8f64805 to
1de7ee4
Compare
|
Rebased onto the latest |
| let mut stuck = vec![spawn_call(&manager, &stuck_id, StuckCall::Initialize)]; | ||
| if call != StuckCall::Initialize { | ||
| tokio::time::sleep(Duration::from_millis(50)).await; | ||
| stuck.push(spawn_call(&manager, &stuck_id, call)); | ||
| } | ||
| tokio::time::sleep(Duration::from_millis(100)).await; | ||
|
|
||
| // A new client connects while the call is waiting... | ||
| let create = tokio::spawn({ | ||
| let manager = manager.clone(); | ||
| async move { manager.create_session().await.map(|(id, _)| id) } | ||
| }); | ||
| tokio::time::sleep(Duration::from_millis(100)).await; |
There was a problem hiding this comment.
The test only reproduces the bug when the stuck call already holds the guard and the create_session writer is queued before has_session runs. Can we make the ordering deterministic instead of relying on timing?
There was a problem hiding this comment.
Done in 763b2c4. The test now uses #[tokio::test(start_paused = true)]. It first waits until the stuck worker has passed initialize on (it reads the request from the stuck transport). Each next step waits for settle(): with the clock paused, time moves only when all tasks are blocked. So on the old code the stuck call already holds the guard, and the create_session writer is queued before has_session runs. On the code before the fix all five cases fail in 30 of 30 runs; with the fix they pass in 100 of 100 runs.
| let handle = sessions | ||
| .get(id) | ||
| .ok_or(LocalSessionManagerError::SessionNotFound(id.clone()))?; | ||
| let handle = self.session_handle(id).await?; |
There was a problem hiding this comment.
The PR description says that after a session is closed, a call waiting for a reply gets SessionServiceTerminated. But while initialize is in progress, the worker is blocked in recv_from_handler() and doesn't read SessionEvent::Close. So, as far as I can tell, a pending initialize_session keeps waiting even after close_session returns.
There was a problem hiding this comment.
You are right. I checked it with a probe test: before this change, close_session blocked until initialize finished. With it, close_session returns at once and removes the session from the map, and initialize keeps waiting for the handler, as before. I fixed the PR description. Should close_session also stop a pending initialize? That needs the worker to read session events during initialization, so I kept it out of this PR.
| .await | ||
| .get(id) | ||
| .cloned() | ||
| .ok_or(LocalSessionManagerError::SessionNotFound(id.clone())) |
There was a problem hiding this comment.
Nit: ok_or creates the error, including id.clone(), on every call, even when it finds the session. Since all five callers are now in one place, this is a good time to make the clone lazy.
| .ok_or(LocalSessionManagerError::SessionNotFound(id.clone())) | |
| .ok_or_else(|| LocalSessionManagerError::SessionNotFound(id.clone())) |
…d clone the id lazily The test now runs with the tokio clock paused, so each step starts only after all spawned tasks are blocked, instead of after a fixed sleep. It also waits until the worker has passed initialize on before it adds the second call. session_handle now builds the SessionNotFound error only when the session is missing.
While one client's
initializeis running, no other client can connect. Once a new client tries, requests on all other sessions wait too. With a server whoseinitializetakes 3 s, a second client'sinitializetook 2.8 s and apingon an existing session took 2.6 s. With this change, both take about 3 ms.Motivation and Context
LocalSessionManagerkeeps all sessions in oneRwLock<HashMap<..>>. Five methods take the read guard and keep it while they wait on the session worker:initialize_sessionwaits for the initialize response, so it keeps the guard during the wholeServerHandler::initialize.create_stream,create_standalone_streamandresumewait for the worker to reply.accept_messagewaits when the session's event channel is full.create_sessionneeds the write lock, so a new client waits for that guard. tokio'sRwLockis fair: when a writer is waiting, new readers wait behind it. The tower service callshas_sessionfor every POST and GET that has a session id, so requests on other sessions wait as well. The wait lasts as long as the slowinitializeruns.Fix
A private helper,
session_handle, clones theLocalSessionHandleout of the map and drops the guard before the caller waits on the worker. The handle is only an id and a channel sender, so the clone is cheap. The five methods use the helper.close_sessionalready drops the lock before it awaits; now the other methods work the same way.How Has This Been Tested?
tests/test_streamable_http_session_isolation.rs. One session waits on its worker: nothing answers on its transport, the same setup astest_streamable_http_init_timeout.rs. Thencreate_sessionandhas_sessionfor another session must finish within 1 s. There is one case for each of the five methods. The test runs with the tokio clock paused: it waits until the worker has passedinitializeon, and each next step starts only after all spawned tasks are blocked, so the order does not depend on timing. Onmainall five cases fail withhas_session blocked: Err(Elapsed(()))(30 of 30 runs). With this change they pass (100 of 100 runs).initialize_sessionfails all five cases, because every case starts with a pending initialize.StreamableHttpServicewith aServerHandlerwhoseinitializesleeps 3 s, then a second client'sinitializeand apingon an existing session. The numbers are at the top.cargo test -p rmcpwith the CI feature set withoutlocal, for the 19 streamable HTTP server test files: all 153 tests pass. Also the nightly rustfmt check and both clippy commands from CI.Tested on Windows 11.
Breaking Changes
None. The public API does not change: the helper is private.
Types of changes
Checklist
Additional context
One small change in behavior:
close_sessionno longer waits for calls that hold the read guard. During initialization the worker does not read session events, so a pendinginitialize_sessionstill waits for the handler's response, as before, and the worker handlesCloseonly after that. I checked it: before this changeclose_sessionblocked untilinitializefinished; now it returns at once and removes the session from the map, whileinitializekeeps waiting. After initialization the worker handles events in order, so events queued beforeCloseare still handled.