test(e2e): expand SSE resilience coverage - #1897
Conversation
There was a problem hiding this comment.
Code Review
This pull request implements SSE event IDs and reconnection logic, enabling clients to resume streams after disconnects or server restarts. It introduces a boot_id and sequence counter for event tracking, adds a configurable connection limit via GATEWAY_MAX_CONNECTIONS, and updates the frontend to manage event IDs. Extensive E2E tests were added to verify reconnection, keepalives, and connection caps. A review comment identified that the chat_events_handler in src/channels/web/handlers/chat.rs lacks the X-Accel-Buffering and Cache-Control headers necessary for proper SSE functionality.
| pub async fn chat_events_handler( | ||
| Query(params): Query<ChatEventsQuery>, | ||
| headers: axum::http::HeaderMap, | ||
| State(state): State<Arc<GatewayState>>, | ||
| AuthenticatedUser(user): AuthenticatedUser, | ||
| ) -> Result<impl IntoResponse, (StatusCode, String)> { | ||
| state.sse.subscribe(Some(user.user_id)).ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Too many connections".to_string(), | ||
| )) | ||
| let last_event_id = params.last_event_id.or_else(|| { | ||
| headers | ||
| .get("last-event-id") | ||
| .and_then(|value| value.to_str().ok()) | ||
| .map(ToOwned::to_owned) | ||
| }); | ||
| state | ||
| .sse | ||
| .subscribe(Some(user.user_id), last_event_id) | ||
| .ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Too many connections".to_string(), | ||
| )) | ||
| } |
There was a problem hiding this comment.
The implementation is missing the necessary X-Accel-Buffering and Cache-Control headers required for SSE streaming. Please update this handler to include these headers to ensure correct behavior. Note: Per repository rules, please focus only on fixing this bug and avoid refactoring or consolidating duplicated code from other parts of the codebase.
pub async fn chat_events_handler(
Query(params): Query<ChatEventsQuery>,
headers: axum::http::HeaderMap,
State(state): State<Arc<GatewayState>>,
AuthenticatedUser(user): AuthenticatedUser,
) -> Result<impl IntoResponse, (StatusCode, String)> {
let last_event_id = params.last_event_id.or_else(|| {
headers
.get("last-event-id")
.and_then(|value| value.to_str().ok())
.map(ToOwned::to_owned)
});
let sse = state
.sse
.subscribe(Some(user.user_id), last_event_id)
.ok_or((
StatusCode::SERVICE_UNAVAILABLE,
"Too many connections".to_string(),
))?;
Ok((
[("X-Accel-Buffering", "no"), ("Cache-Control", "no-cache")],
sse,
))
}References
- Consolidating duplicated code can be considered a pre-existing issue and out of scope for a bugfix, which should remain focused.
There was a problem hiding this comment.
Fixed in bdaea8d. The split handlers/chat.rs SSE handler now returns the same X-Accel-Buffering: no and Cache-Control: no-cache headers as the inline gateway handler.
There was a problem hiding this comment.
Pull request overview
This PR improves the web gateway’s Server-Sent Events (SSE) resilience by adding event IDs and reconnect-aware filtering, and significantly expands end-to-end coverage for SSE edge cases described in #1784.
Changes:
- Add process-scoped SSE event IDs (
<boot_uuid>:<counter>) and support reconnect filtering viaLast-Event-IDheader orlast_event_idquery parameter. - Introduce a configurable gateway connection cap (
GATEWAY_MAX_CONNECTIONS, default 100) enforced across SSE + WebSocket connections. - Expand E2E helpers and scenarios to cover keepalive comments, server restart recovery, multi-tab fanout, stale reconnect IDs, and connection-limit behavior.
Reviewed changes
Copilot reviewed 14 out of 14 changed files in this pull request and generated 5 comments.
Show a summary per file
| File | Description |
|---|---|
| tests/e2e/scenarios/test_sse_reconnect.py | Reworks and expands SSE/connectivity E2E coverage (keepalive, restart recovery, multi-tab, stale IDs, limits). |
| tests/e2e/helpers.py | Adds raw-SSE helpers (sse_stream, wait_for_sse_comment) using aiohttp for low-level SSE assertions. |
| tests/e2e/conftest.py | Adds restartable/isolated gateway fixtures (managed_gateway_server, limited_gateway_server). |
| tests/e2e/README.md | Updates scenario documentation for expanded SSE reconnect test coverage. |
| tests/e2e/CLAUDE.md | Documents the new fixtures and the rationale for aiohttp-based raw SSE checks. |
| src/config/channels.rs | Adds GatewayConfig.max_connections and parses GATEWAY_MAX_CONNECTIONS (validated > 0). |
| src/channels/web/mod.rs | Wires GatewayConfig.max_connections into SseManager construction. |
| src/channels/web/sse.rs | Implements SSE event IDs, reconnect filtering, and configurable connection limits in SseManager. |
| src/channels/web/server.rs | Extends /api/chat/events to accept last_event_id (query/header) for reconnect filtering. |
| src/channels/web/handlers/chat.rs | Mirrors /api/chat/events reconnect parameter handling for the split handler module. |
| src/channels/web/static/app.js | Tracks SSE lastEventId client-side and sends it on reconnect via last_event_id query param. |
| src/channels/web/CLAUDE.md | Documents updated SSE contract (IDs/reconnect behavior) and connection limit config. |
| src/channels/web/tests/no_silent_drop.rs | Updates test gateway config construction for the new max_connections field. |
| src/tunnel/mod.rs | Updates tunnel-related tests for the new GatewayConfig.max_connections field. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| let last_event_id = params.last_event_id.or_else(|| { | ||
| headers | ||
| .get("last-event-id") | ||
| .and_then(|value| value.to_str().ok()) | ||
| .map(ToOwned::to_owned) | ||
| }); | ||
| state | ||
| .sse | ||
| .subscribe(Some(user.user_id), last_event_id) | ||
| .ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Too many connections".to_string(), | ||
| )) |
There was a problem hiding this comment.
chat_events_handler duplicates the last_event_id extraction logic that also exists in src/channels/web/server.rs. To avoid the two drifting (e.g., header name, precedence rules), consider factoring this into a small shared helper (e.g., extract_last_event_id(&ChatEventsQuery, &HeaderMap)) used by both handlers.
There was a problem hiding this comment.
Fixed in bdaea8d. I factored the reconnect parsing into a shared extract_last_event_id() helper and now use that from both the split handler and the inline server handler so the precedence rules stay aligned.
| )) | ||
| let last_event_id = params.last_event_id.or_else(|| { | ||
| headers | ||
| .get("last-event-id") |
There was a problem hiding this comment.
For the standard SSE reconnect header, prefer using axum::http::header::LAST_EVENT_ID (or http::header::LAST_EVENT_ID) instead of a string literal. This avoids typos and makes the intent clearer while still being case-insensitive at runtime.
| .get("last-event-id") | |
| .get(axum::http::header::LAST_EVENT_ID) |
There was a problem hiding this comment.
Addressed in bdaea8d, with a slight adjustment: this http version does not export a LAST_EVENT_ID constant, so the shared helper now uses HeaderName::from_static("last-event-id") instead of a string literal at each call site.
| function rememberSseEventId(event) { | ||
| if (!event || !event.lastEventId) return; | ||
| _lastSseEventId = event.lastEventId; | ||
| window.__e2eLastSseEventId = event.lastEventId; | ||
| } |
There was a problem hiding this comment.
Setting window.__e2eLastSseEventId unconditionally introduces a new global that will exist for all users (not just E2E). If this is only meant for tests, consider guarding it behind a feature flag/query param (or at least namespacing under a single window.__e2e object) to reduce global namespace pollution.
There was a problem hiding this comment.
Fixed in bdaea8d. I removed the always-on standalone global and now store the test-only SSE ID under a single window.__e2e namespace.
| + 'last_event_id=' + encodeURIComponent(lastEventId); | ||
| } | ||
| eventSource = new EventSource(chatSseUrl); | ||
| window.__e2eCurrentSseUrl = chatSseUrl; |
There was a problem hiding this comment.
window.__e2eCurrentSseUrl is another always-on test global. If it’s not intended for production, consider gating it similarly (or removing it once the E2E harness no longer needs it).
| window.__e2eCurrentSseUrl = chatSseUrl; | |
| if (window.__e2eEnabled) { | |
| window.__e2eCurrentSseUrl = chatSseUrl; | |
| } |
There was a problem hiding this comment.
Fixed in bdaea8d. window.__e2eCurrentSseUrl was not needed by the tests, so I removed it rather than leaving another always-on test hook in production code.
| max_connections: { | ||
| let max = parse_optional_env("GATEWAY_MAX_CONNECTIONS", 100_u64)?; | ||
| if max == 0 { | ||
| return Err(ConfigError::InvalidValue { | ||
| key: "GATEWAY_MAX_CONNECTIONS".to_string(), | ||
| message: "must be greater than 0".to_string(), | ||
| }); | ||
| } | ||
| max | ||
| }, |
There was a problem hiding this comment.
The default connection limit is hard-coded here as 100_u64, while the SSE layer also has MAX_CONNECTIONS = 100. Consider centralizing the default in one place (e.g., expose DEFAULT_MAX_CONNECTIONS from the web channel) to prevent future divergence between config defaults and runtime behavior.
There was a problem hiding this comment.
Fixed in bdaea8d. The config path now reuses DEFAULT_MAX_CONNECTIONS from src/channels/web/sse.rs so the documented/default limit stays centralized.
ilblackdragon
left a comment
There was a problem hiding this comment.
Code Review: test(e2e): expand SSE resilience coverage
Overview
This PR adds SSE event IDs (<boot_uuid>:<counter>) and reconnect-aware dedup to the web gateway, plus substantial E2E test coverage for keepalive, multi-tab fanout, server restart recovery, stale reconnect IDs, and connection limits. It also makes GATEWAY_MAX_CONNECTIONS configurable via environment variable. Closes #1784.
Despite the title suggesting "test only", this is a feature + test PR — the SSE event ID system, subscribe() API change, and configurable max connections are production code changes.
Architecture: Sound
The boot_uuid:counter scheme is well-designed:
- Process-scoped: after restart, old IDs are from a different boot and are ignored — client falls back to history API
AtomicU64counter is monotonic and race-freeis_event_after()correctly handles all edge cases (different boot, malformed IDs, same-boot ordering)- Defense-in-depth:
last_event_idaccepted via bothLast-Event-IDheader and query parameter (needed becauseEventSourcereconnect state is recreated in JS)
Issues
1. Duplicate chat_events_handler — both copies updated but divergence risk remains
There are two independent copies: handlers/chat.rs:41 (pub) and server.rs:2134 (private, actually wired into routes). The PR correctly updates both, but the handlers/chat.rs copy appears unused by the route registration. This is a pre-existing issue, but this PR adds a new import from handlers/chat.rs into server.rs (ChatEventsQuery, extract_last_event_id) while keeping the duplicate handler — increasing coupling without resolving the duplication.
Suggestion: either wire the route to handlers::chat::chat_events_handler and delete the server.rs copy, or at minimum add a // NOTE: keep in sync with handlers/chat.rs comment.
2. subscribe_raw() does not carry event IDs — WebSocket clients can't do ID-based reconnect
subscribe_raw() (used by WebSocket in ws.rs and Responses API in responses_api.rs) still returns AppEvent without the ScopedEvent.id. WebSocket clients cannot participate in the event-ID reconnect protocol. This is fine if intentional (WS has its own reconnect semantics), but should be documented since the CLAUDE.md only mentions SSE event IDs.
3. from_sender() generates a fresh boot_id — event ID sequence resets on rebuild_state
Each call to from_sender() creates a new boot_id and resets the counter to 1. Since rebuild_state() is only called during startup wiring (before connections are accepted), this is safe. But if rebuild_state() is ever called at runtime, connected clients would see a boot_id change mid-session and lose the dedup guarantee. The existing doc comment on from_sender should be updated to mention this new implication.
4. addTrackedEventListener doesn't wrap onopen / onerror
The onopen and onerror handlers are assigned directly to eventSource.onopen/eventSource.onerror, not through addTrackedEventListener. These don't carry event IDs so this is correct — but suggestions and other events are wrapped. The pattern is consistent with the SSE spec (only named events have IDs), just worth confirming the intent.
5. _lastSseEventId is not cleared on explicit new connection
When connectSSE() is called without a lastEventIdOverride and _lastSseEventId is set from a previous boot, the stale ID is sent. The server handles this correctly (different boot_id → pass all events), but the client sends a pointless parameter. Minor, not a bug.
6. E2E test test_reconnect_with_stale_last_event_id hardcodes mock response
Line 1206: await _wait_for_turn_in_history(..., "The answer is 4.") assumes the mock LLM returns this exact string. If the mock changes, this test silently breaks. Consider extracting the expected response from the mock's configuration or at least adding a comment about the coupling.
What's Done Well
ManagedIronclawServer— Clean restartable process wrapper with proper SIGINT, stderr capture on failure, and port preservation. Good reusable E2E infrastructure.sse_stream()helper — Usingaiohttpfor raw SSE access is the right call. Playwright'sEventSourcecan't expose keepalive comments.- Unit tests for
is_event_after— Cover same-boot ordering, cross-boot, and malformed IDs. Good edge-case coverage. X-Accel-Buffering: noheader — Added tohandlers/chat.rsversion, prevents Nginx from buffering SSE.- Config validation —
GATEWAY_MAX_CONNECTIONS=0is rejected with a clear error. - Documentation — CLAUDE.md, E2E CLAUDE.md, and README all updated. The SSE event ID contract is documented inline.
Minor Nits
- The
aiohttpdependency is added to E2E helpers but I don't see it in any requirements file in the diff. Confirm it's intests/e2e/requirements.txt. parse_event_idreturnsOption— the?chain is clean, but consider naming ittry_parse_event_idfor clarity that it's fallible (Rust convention nit).
Verdict
Approve with suggestions. The SSE event-ID scheme is well-designed and the E2E coverage is thorough. The main feedback is the duplicate handler (#1) which should be consolidated, and documenting the WebSocket gap (#2). The rest are minor nits.
- Remove duplicate chat_events_handler from server.rs; wire route to handlers::chat::chat_events_handler instead - Update from_sender doc comment to document boot_id/event-ID reset - Document that WebSocket (subscribe_raw) does not expose event IDs [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
…se-e2e # Conflicts: # src/channels/web/static/app.js # tests/e2e/helpers.py
[skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 36 out of 36 changed files in this pull request and generated 4 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| tracing::debug!( | ||
| account_id = %body.account_id, | ||
| public_key = %body.public_key, | ||
| signature_len = body.signature.len(), | ||
| signature_prefix = &body.signature[..body.signature.len().min(20)], | ||
| signature_prefix = &body.signature[..body.signature.len().min(20)], // safety: min(len,20) is always valid boundary on ASCII base58/hex/base64 | ||
| "NEAR verify: decoding credentials" |
There was a problem hiding this comment.
signature_prefix = &body.signature[..body.signature.len().min(20)] can panic if signature contains non-ASCII/multibyte UTF-8 (the length check doesn’t guarantee a char boundary). The newly added “safety” comment is not accurate; consider using the existing safe_truncate() helper (or a char-boundary truncation) for logging.
There was a problem hiding this comment.
Addressed in 7e5ba8b. Replaced the byte-index slice with a call to the existing safe_truncate(&body.signature, 20) helper so the log formatter no longer panics on multibyte input.
| assert!( | ||
| r.had_error || r.stdout.contains("Error") || r.stdout.contains("limit"), | ||
| "resource limit should terminate infinite loop, got stdout: {}", | ||
| &r.stdout[..r.stdout.len().min(500)], | ||
| &r.stdout[..r.stdout.len().min(500)], // safety: test-only; stdout is String so len() is a valid boundary | ||
| ); |
There was a problem hiding this comment.
&r.stdout[..r.stdout.len().min(500)] is not guaranteed to be a UTF-8 char boundary and can panic if stdout contains multibyte characters. The added “safety” comment is incorrect; use a char-boundary truncation (similar to other truncate helpers) before slicing for the assertion message.
There was a problem hiding this comment.
Addressed in 7e5ba8b. Added a local truncate_for_assert() helper in the tests module that walks back to the nearest UTF-8 char boundary before slicing, and the assertion now calls truncate_for_assert(&r.stdout, 500) instead of raw byte-index slicing.
| assert!( | ||
| r.had_error || r.stdout.contains("Error") || r.stdout.contains("limit"), | ||
| "cpu-bound loop should be terminated, stdout: {}", | ||
| &r.stdout[..r.stdout.len().min(500)], | ||
| &r.stdout[..r.stdout.len().min(500)], // safety: test-only; stdout is String so len() is a valid boundary | ||
| ); |
There was a problem hiding this comment.
Same UTF-8 slicing issue as above: &r.stdout[..r.stdout.len().min(500)] can panic on multibyte output, so the “valid boundary” comment isn’t correct. Prefer truncating at a char boundary before slicing for the assertion message.
There was a problem hiding this comment.
Addressed in 7e5ba8b. Same fix — this assertion now also calls the new truncate_for_assert() helper, which walks back to a valid UTF-8 char boundary before slicing.
| const addTrackedEventListener = (eventType, handler) => { | ||
| eventSource.addEventListener(eventType, (event) => { | ||
| rememberSseEventId(event); | ||
| handler(event); | ||
| }); | ||
| }; | ||
|
|
There was a problem hiding this comment.
addTrackedEventListener correctly captures lastEventId for most event types, but plan_update is still registered later in connectSSE() via a direct eventSource.addEventListener(...). That means _lastSseEventId won’t advance on plan_update frames, which can cause reconnects to send an older last_event_id and reduce dedup accuracy. Consider routing all SSE listeners through addTrackedEventListener (including plan_update).
| const addTrackedEventListener = (eventType, handler) => { | |
| eventSource.addEventListener(eventType, (event) => { | |
| rememberSseEventId(event); | |
| handler(event); | |
| }); | |
| }; | |
| const originalAddEventListener = eventSource.addEventListener.bind(eventSource); | |
| const wrapTrackedHandler = (handler) => (event) => { | |
| rememberSseEventId(event); | |
| if (typeof handler === 'function') { | |
| handler.call(eventSource, event); | |
| return; | |
| } | |
| if (handler && typeof handler.handleEvent === 'function') { | |
| handler.handleEvent(event); | |
| } | |
| }; | |
| eventSource.addEventListener = (eventType, handler, options) => { | |
| return originalAddEventListener(eventType, wrapTrackedHandler(handler), options); | |
| }; | |
| const addTrackedEventListener = (eventType, handler, options) => { | |
| return originalAddEventListener(eventType, wrapTrackedHandler(handler), options); | |
| }; |
There was a problem hiding this comment.
Addressed in 7e5ba8b. plan_update now goes through addTrackedEventListener, so _lastSseEventId advances on plan_update frames and reconnect dedup stays accurate.
- auth.rs: use safe_truncate() for signature_prefix log (user-supplied string may contain multibyte UTF-8 that panics on byte-index slicing) - scripting.rs: add truncate_for_assert() helper, use it in both resource-limit test assertions (Python stdout may contain multibyte) - app.js: route plan_update through addTrackedEventListener so _lastSseEventId advances on plan_update frames (reconnect dedup) - store_adapter.rs: rewrite slug truncation with explicit ASCII-only byte search and fix formatting (rustfmt moved trailing safety comment onto a new line) [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
* test(e2e): expand SSE resilience coverage * fix: address PR review follow-ups * refactor(web): consolidate duplicate chat_events_handler, improve docs - Remove duplicate chat_events_handler from server.rs; wire route to handlers::chat::chat_events_handler instead - Update from_sender doc comment to document boot_id/event-ID reset - Document that WebSocket (subscribe_raw) does not expose event IDs [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * style: fix formatting from merge safety annotations [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * fix: address PR nearai#1897 review comments and CI formatting - auth.rs: use safe_truncate() for signature_prefix log (user-supplied string may contain multibyte UTF-8 that panics on byte-index slicing) - scripting.rs: add truncate_for_assert() helper, use it in both resource-limit test assertions (Python stdout may contain multibyte) - app.js: route plan_update through addTrackedEventListener so _lastSseEventId advances on plan_update frames (reconnect dedup) - store_adapter.rs: rewrite slug truncation with explicit ASCII-only byte search and fix formatting (rustfmt moved trailing safety comment onto a new line) [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> --------- Co-authored-by: ilblackdragon@gmail.com <ilblackdragon@gmail.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
* test(e2e): expand SSE resilience coverage * fix: address PR review follow-ups * refactor(web): consolidate duplicate chat_events_handler, improve docs - Remove duplicate chat_events_handler from server.rs; wire route to handlers::chat::chat_events_handler instead - Update from_sender doc comment to document boot_id/event-ID reset - Document that WebSocket (subscribe_raw) does not expose event IDs [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * style: fix formatting from merge safety annotations [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * fix: address PR nearai#1897 review comments and CI formatting - auth.rs: use safe_truncate() for signature_prefix log (user-supplied string may contain multibyte UTF-8 that panics on byte-index slicing) - scripting.rs: add truncate_for_assert() helper, use it in both resource-limit test assertions (Python stdout may contain multibyte) - app.js: route plan_update through addTrackedEventListener so _lastSseEventId advances on plan_update frames (reconnect dedup) - store_adapter.rs: rewrite slug truncation with explicit ASCII-only byte search and fix formatting (rustfmt moved trailing safety comment onto a new line) [skip-regression-check] Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> --------- Co-authored-by: ilblackdragon@gmail.com <ilblackdragon@gmail.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
Summary
Closes #1784.
Testing