fix: correctly handle timeouts and keepalive - #41
Conversation
|
I think this should be sufficient @honungsburk! we need to update docs tho |
|
Caution Review failedThe pull request is closed. Walkthrough
Sequence Diagram(s)sequenceDiagram
autonumber
actor Client
participant Plug as StreamableHTTP.Plug
participant Transport as StreamableHTTP (GenServer)
participant TaskSup as Task.Supervisor
participant Server as MCP Server
note right of Plug: Build RequestParams (message, transport,\nsession_id, session_header, timeout, context)
Client->>Plug: HTTP POST /mcp (JSON)
Plug->>Transport: handle_message(RequestParams)
Transport->>TaskSup: async_nolink(fn -> Server.handle_request(...) end)
Transport-->>Transport: start per-request timer (RequestParams.timeout)
alt success before timeout
TaskSup-->>Transport: {:ok, result}
Transport-->>Plug: JSON-RPC result
Plug-->>Client: 200 JSON
else timeout
Transport-->>Plug: JSON-RPC error (timeout, session_id)
Plug-->>Client: 200 JSON error
end
sequenceDiagram
autonumber
actor Client
participant Plug as StreamableHTTP.Plug
participant Transport as StreamableHTTP (GenServer)
participant SSE as SSE.Handler
participant Stream as SSE.Streaming
note right of Plug: Build RequestParams and register SSE handler
Client->>Plug: HTTP GET/POST wants SSE
Plug->>Transport: register_sse_handler(session_id) with RequestParams
Transport-->>Plug: {:ok, SSE pid} or existing handler
Plug->>Stream: start_sse_streaming(conn, RequestParams)
note over Transport: keepalive_enabled true when handlers > 0
loop every keepalive_interval
Transport-->>SSE: :sse_keepalive
SSE->>Stream: keep_alive() -> ": keepalive" chunk
end
Plug->>Transport: handle_message_for_sse(RequestParams)
Transport->>SSE: stream events/chunks for result
SSE-->>Client: event-stream data
sequenceDiagram
autonumber
participant Transport as StreamableHTTP (GenServer)
participant TaskSup as Task.Supervisor
participant Server as MCP Server
note right of Transport: Use RequestParams.timeout and task_timeout handling
Transport->>TaskSup: async_nolink(request)
Transport-->>Transport: set timer (RequestParams.timeout)
alt task finishes before timer
TaskSup-->>Transport: {:ok, result}
Transport-->>Caller: JSON-RPC success/error
else timer fires first
Transport-->>Transport: :task_timeout
Transport-->>Caller: JSON-RPC timeout error with session_id
Transport->>TaskSup: ignore/cancel late result
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes P0: Verify RequestParams @enforce_keys and propagation of timeout/request_timeout into all call sites; ensure no missing Keyword.fetch! uses that crash runtime. Review with confidence, clarity & light humor 😎. Pre-merge checks and finishing touches❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
📜 Recent review detailsConfiguration used: CodeRabbit UI Review profile: ASSERTIVE Plan: Pro 📒 Files selected for processing (1)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 11
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (4)
lib/anubis/server/transport/streamable_http/plug.ex (1)
303-319: Consider using Task.Supervisor for background requests (P2)The
start_background_request/1function usesTask.start/1to create an unsupervised, unlinked task. According to the PR description, one of the objectives is to "switch Task.Supervisor usage to async_nolink to avoid cascading failures."However, this code doesn't use the task supervisor at all—it creates a completely unsupervised task. If this task crashes, there's no supervision or monitoring, and the SSE client might wait indefinitely for a response that never comes.
Consider using the task supervisor that's already available in the transport state:
- defp start_background_request(params) do + defp start_background_request(params, task_supervisor) do self_pid = self() - Task.start(fn -> + Task.Supervisor.start_child(task_supervisor, fn -> case StreamableHTTP.handle_message(params) do {:ok, response} when is_binary(response) -> send(self_pid, {:sse_message, response}) {:error, reason} -> Logging.transport_event( "sse_background_request_error", %{reason: reason}, level: :error ) end end) endYou'll need to pass the task_supervisor through from the opts map. Note: If you want to prevent cascading failures, use
Task.Supervisor.async_nolink/2instead ofstart_child/2, but ensure you handle the task's result or don't care about it.lib/anubis/server/transport/streamable_http.ex (3)
252-281: Task supervision change LGTM (async_nolink + per‑task timeout). P2Good isolation; no linking to transport.
Also store the task in task_info so we can cancel it on timeout (see timeout handler comment). Example:
- task_info = %{ + task_info = %{ type: :handle_message, session_id: session_id, - from: from, - task_timeout: task_timeout_ref + from: from, + task: task, + task_timeout: task_timeout_ref }
285-311: SSE aware task path LGTM; align with task cancel refactor. P2Mirror storing the task for cancellation on timeout.
- task_info = %{ + task_info = %{ type: :handle_message_for_sse, session_id: session_id, from: from, has_sse_handler: sse_handler?, - task_timeout: task_timeout_ref + task_timeout: task_timeout_ref, + task: task }
433-445: Fix timer cancel key; current check is dead code. P2You store :task_timeout, but check :timeout_ref. Timer isn’t canceled.
Apply:
- if Map.has_key?(task_info, :timeout_ref) do - Process.cancel_timer(task_info.task_timeout) - end + Process.cancel_timer(task_info.task_timeout)This is safe even if it already fired.
📜 Review details
Configuration used: CodeRabbit UI
Review profile: ASSERTIVE
Plan: Pro
⛔ Files ignored due to path filters (1)
priv/dev/upcase/mix.lockis excluded by!**/*.lock
📒 Files selected for processing (9)
.gitignore(1 hunks)lib/anubis/server/transport/streamable_http.ex(16 hunks)lib/anubis/server/transport/streamable_http/plug.ex(9 hunks)lib/anubis/server/transport/streamable_http/request_params.ex(1 hunks)lib/anubis/sse/streaming.ex(2 hunks)priv/dev/upcase/.formatter.exs(1 hunks)priv/dev/upcase/lib/upcase/router.ex(1 hunks)priv/dev/upcase/lib/upcase/server.ex(3 hunks)test/anubis/server/transport/streamable_http_test.exs(2 hunks)
🧰 Additional context used
🧬 Code graph analysis (4)
priv/dev/upcase/lib/upcase/server.ex (2)
priv/dev/echo-elixir/lib/echo_mcp/server.ex (2)
handle_tool_call(60-64)handle_tool_call(66-78)test/support/stub_server.ex (2)
handle_tool_call(98-103)handle_tool_call(105-107)
lib/anubis/server/transport/streamable_http/plug.ex (5)
lib/anubis/server/transport/streamable_http.ex (4)
register_sse_handler(139-141)handle_message(159-162)handle_message_for_sse(172-175)get_sse_handler(184-186)lib/anubis/server/transport/stdio.ex (1)
process_message(250-274)lib/anubis/server/transport/streamable_http/request_params.ex (1)
new(24-33)lib/anubis/server/transport/sse/plug.ex (1)
send_jsonrpc_error(322-329)lib/anubis/sse/streaming.ex (1)
start(32-41)
test/anubis/server/transport/streamable_http_test.exs (1)
lib/anubis/server/transport/streamable_http.ex (1)
handle_message_for_sse(172-175)
lib/anubis/server/transport/streamable_http.ex (3)
test/support/stub_transport.ex (1)
send_message(83-85)lib/anubis/server/transport/sse.ex (6)
register_sse_handler(166-168)handle_message(187-189)get_sse_handler(198-200)forward_request_to_server(341-361)handle_info(396-407)handle_info(410-412)lib/anubis/mcp/error.ex (6)
protocol(95-102)protocol(104-111)protocol(113-120)protocol(122-129)protocol(131-138)to_json_rpc(241-252)
🪛 GitHub Check: static-analysis (1.18, 26)
lib/anubis/sse/streaming.ex
[warning] 82-82: call
The function call loop will not succeed.
🪛 GitHub Check: static-analysis (1.18, 27)
lib/anubis/sse/streaming.ex
[warning] 82-82: call
The function call loop will not succeed.
🪛 GitHub Check: static-analysis (1.18, 28)
lib/anubis/sse/streaming.ex
[warning] 82-82: call
The function call loop will not succeed.
🔇 Additional comments (13)
lib/anubis/sse/streaming.ex (1)
81-82: Static analysis warning is a false positive ✓The static analysis tool is flagging the
loop/4call, but this is a false positive. The function is tail-recursive and will succeed as long askeep_alive/1returns aPlug.Conn.t()(which it does, assuming my suggestion above is implemented).priv/dev/upcase/lib/upcase/router.ex (1)
10-12: LGTM! Timeout configuration is appropriate for dev testing ✓The 10-second
request_timeoutaligns well with the "timeout" tool in the server (which can sleep up to that duration). This configuration properly exercises the new timeout handling mechanism introduced in this PR.lib/anubis/server/transport/streamable_http/request_params.ex (1)
1-34: LGTM! Clean encapsulation of request parameters ✓The RequestParams struct provides excellent encapsulation for the various parameters that flow through the request pipeline. The use of
@enforce_keysfor critical fields and the clear separation between required/optional fields makes the API safer and more maintainable.This addresses the PR objective of improving "encapsulation and API maintainability" effectively.
lib/anubis/server/transport/streamable_http/plug.ex (3)
79-87: LGTM! Clean consolidation of timeout configuration ✓The refactoring from multiple timeout options (
:timeout,:call_timeout) to a single:request_timeoutsimplifies configuration and aligns well with the RequestParams pattern. This directly addresses the PR objective of providing "per-request timeouts" and makes the API more consistent.
133-143: RequestParams construction looks solid ✓The construction of RequestParams with all required fields is correct and complete. The code properly extracts session_id, builds context, and assembles everything into a single parameter object for downstream processing.
170-179: Good separation of request vs notification handling ✓The refactored
process_message/2cleanly separates requests (which might need SSE) from notifications (which are fire-and-forget). The pattern matching on the params struct and delegation to specialized handlers maintains good modularity.lib/anubis/server/transport/streamable_http.ex (7)
54-57: Aliases LGTM. P3Imports for Error, ID, and RequestParams are correct and used appropriately.
138-141: Registration API change LGTM. P3Spec and call consistent with new API.
158-162: Per‑request timeout passthrough LGTM. P2+1s buffer is fine to let the internal timer reply before the outer call times out.
Please ensure params.timeout is in ms everywhere it’s set/passed to avoid unit drift.
170-175: SSE path: per‑request timeout LGTM. P2Consistent with handle_message/1. Good.
183-186: get_sse_handler no explicit timeout. P3Fine as a quick lookup, but consider adding a small explicit timeout for consistency.
[ suggest_optional_refactor ]- def get_sse_handler(transport, session_id) do - GenServer.call(transport, {:get_sse_handler, session_id}) - end + def get_sse_handler(transport, session_id) do + GenServer.call(transport, {:get_sse_handler, session_id}, 1000) + end
348-366: Server call wrapper LGTM. P2Clean mapping to {:request, msg, session_id, context}; good error logging.
482-492: Keepalive loop LGTM. P2Broadcasting :sse_keepalive and rescheduling is fine; combined with the “start on first register” fix, this will meet the objective.
There was a problem hiding this comment.
Actionable comments posted: 1
📜 Review details
Configuration used: CodeRabbit UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (1)
lib/anubis/sse/streaming.ex(2 hunks)
🔇 Additional comments (1)
lib/anubis/sse/streaming.ex (1)
81-90: ✅ LGTM - Keepalive error handling correctly implemented (P0 concern resolved) 😎The past review concern about handling
Plug.Conn.chunk/2errors has been properly addressed. The current implementation pattern-matches on the result fromkeep_alive/1and handles both success and error cases explicitly:
- Success: continues loop + increments counter
- Error: logs with session context + exits gracefully
This caller-side error handling is actually cleaner than the suggested approach since it keeps logging concerns in the loop where
session_idis available. Solid work! 🎯
- Fix pattern mismatch in determine_session_id/3: change [message] to message This was preventing session ID generation on initialize requests (P0 issue) - Remove call_timeout parameter to avoid conflict with PR zoedsoupe#41 Addresses PR zoedsoupe#48 review comments from zoedsoupe and CodeRabbit
Problem
The HTTP transport layer had issues with:
Closes #8
Closes #10
Solution
RequestParamsstruct to encapsulate request context and timeoutasync_nolinkto prevent cascading failurescall_timeoutandrequest_timeoutoptions in favor of per-request timeoutsRationale
RequestParamsprovides clean encapsulation of request data and makes the API more maintainableasync_nolinkprevents task failures from taking down the supervisorSummary by CodeRabbit
New Features
Refactor
Chores
Tests