|
13 | 13 | EventStore, |
14 | 14 | StreamableHTTPServerTransport, |
15 | 15 | StreamId, |
| 16 | + _request_stream_key, |
16 | 17 | ) |
17 | 18 | from mcp.shared.message import SessionMessage |
18 | 19 |
|
@@ -56,11 +57,14 @@ async def test_router_unconsumed_request_stream_does_not_block_siblings() -> Non |
56 | 57 | async with transport.connect() as (_read_stream, write_stream): |
57 | 58 | # Model two concurrent POSTs at the point _handle_post_request has |
58 | 59 | # registered the per-request stream but A's sse_writer has not yet |
59 | | - # reached its first receive(). |
60 | | - streams["A"] = anyio.create_memory_object_stream[EventMessage](REQUEST_STREAM_BUFFER_SIZE) |
61 | | - streams["B"] = anyio.create_memory_object_stream[EventMessage](REQUEST_STREAM_BUFFER_SIZE) |
62 | | - a_send, a_recv = streams["A"] |
63 | | - b_reader = streams["B"][1] |
| 60 | + # reached its first receive(). Routing keys must match `_request_stream_key` |
| 61 | + # (type-preserving prefixes), same as production registration. |
| 62 | + key_a = _request_stream_key("A") |
| 63 | + key_b = _request_stream_key("B") |
| 64 | + streams[key_a] = anyio.create_memory_object_stream[EventMessage](REQUEST_STREAM_BUFFER_SIZE) |
| 65 | + streams[key_b] = anyio.create_memory_object_stream[EventMessage](REQUEST_STREAM_BUFFER_SIZE) |
| 66 | + a_send, a_recv = streams[key_a] |
| 67 | + b_reader = streams[key_b][1] |
64 | 68 | b_received = anyio.Event() |
65 | 69 |
|
66 | 70 | async def consume_b() -> None: |
|
0 commit comments