Repository navigation
Fixes #7313: Close stalled websocket sessions and bound the admin send queue - #7343
Conversation
| } | ||
| inFlightMessage = message; | ||
| } | ||
| final ScheduledFuture<?> future = SEND_WATCHDOG.schedule( |
There was a problem hiding this comment.
timeoutFuture is never assigned anywhere — sendNext() keeps the ScheduledFuture in a local variable and passes it to onSendResult, so the field stays null forever and both cancel blocks (forceClose():533-536 and close():552-555) are dead code.
It is not a correctness problem today: onSendTimeout() re-checks closed first, so a task that outlives its session returns immediately. But it does mean every session closed while a message is in flight keeps its watchdog task (and a strong reference to this queue and the Session) alive for up to sendTimeoutMillis (30s by default), on a scheduler that is never drained.
Either assign it where the send is scheduled, i.e. inside the synchronized block that sets inFlightMessage:
inFlightMessage = message;
// then, after scheduling:
timeoutFuture = future;or drop the field and the two cancel blocks and say explicitly that closed is the only guard. As written it reads like the timeout is cancelled on close when it is not.
Aias00
left a comment
There was a problem hiding this comment.
Approving — this is the right pair of fixes for a half-open connection, and the two failure modes are genuinely different: a failed send gives you a SendResult callback, but a stalled one gives you nothing at all, which is why master could keep a dead session in SESSION_SET forever and keep appending to an unbounded queue behind it. Closing so the gateway reconnects and re-pulls everything (MYSELF) is the correct remedy, and it is consistent with what #7344 does on the client side.
I checked the things that usually go wrong here:
- No
ConcurrentModificationExceptionon the publish paths.forceClose→removeSessionIndexesmutatesSESSION_SETandNAMESPACE_SESSION_MAP, but both broadcast loops already iterate a snapshot (new ArrayList<>(SESSION_SET)at:293,new ArrayList<>(sessions)at:333). Good. - No map leak.
removeSessionIndexes:360-362removes fromSESSION_SET, drops theSessionSendQueueviaremoveSessionSendQueue, and unwinds the namespace index, so a forced close cleans up the same way@OnClose→clearSessiondoes. - The queue is not left "stuck sending".
forceClosesetsclosed = trueandsending = falseand clears the backlog, so the nextsend()returns without queueing and the nextsendNext()bails out. This is the piece master was missing — a failed async send leftsendingtrue in the old code path in some orders.
One thing to fix (inline, non-blocking but real)
timeoutFuture is never assigned: sendNext() keeps the ScheduledFuture in a local and hands it to onSendResult, so the field is always null and both cancel blocks in forceClose/close never run. Not a correctness problem — onSendTimeout re-checks closed and returns — but every session closed mid-send pins its queue and Session in the scheduler for up to 30s. Either assign it or delete the field and rely on closed explicitly.
Worth considering
session.close()is called while the caller holds the queue monitor.send():451andonSendTimeout():521both callforceCloseinsidesynchronized (this), andforceCloseonly releases that monitor after its own block, sosession.close()at:542runs under it. It is safe today because containers dispatch@OnCloseasynchronously, but a container that invoked the close handler inline would takeclearSession→removeSessionSendQueue→close()onto the same monitor. Cheap to avoid: haveforceClosereturn a "close me" signal and do the actualsession.close()outside the lock.- One scheduled task per message.
sendNextschedules and immediately (in the normal case) cancels a watchdog for every single message pushed. That is an allocation plus a delay-queue insert/remove on the config-push path, and it is bounded only because there is at most one in-flight message per session. A single periodic sweep over "in-flight since" timestamps would scale better if admin pushes become chatty. - Both limits are hard-coded (256 queued / 30s). A slow gateway receiving a large full-sync burst will hit the queue cap, get closed, reconnect, and receive the same large burst again. That is the intended "don't drop silently" behaviour, but it can flap; making the two values configurable would let operators trade memory against resync churn per deployment.
- The watchdog is a static single-thread executor with no shutdown hook (daemon, so it does not block exit) — fine for admin, just noting there is no lifecycle tie-in.
The tests are good: the stalled-callback case actually stubs a container that never invokes the handler, which is exactly the scenario, and the overflow test verifies the second message never reaches sendText. resetSendGuards() in setUp keeps the static overrides from leaking between tests.
…send-queue-guard # Conflicts: # shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
Aias00
left a comment
There was a problem hiding this comment.
Re-reviewed: bounding the send queue and adding a watchdog that closes sessions whose async send callback never fires is the right way to deal with half-open connections (#7313) — the gateway then reconnects and does a full resync instead of staying stalled forever. Daemon watchdog thread and configurable timeout/queue are sensible. The red e2e / e2e-logging-rocketmq jobs are the known flakes, unrelated to this change — approving.
Aias00
left a comment
There was a problem hiding this comment.
Re-approved after the latest commits. The WebsocketCollector change is materially the one I reviewed — the send watchdog that closes sessions whose async send callback never fires (so a half-open connection triggers a reconnect and full resync rather than stalling), plus the bounded send queue and configurable timeout. The extra additions in this push are the e2e test files rather than new collector logic. CI is green.
Note that this branch, like several other open PRs, now also carries identical DividePluginTest / LoggingRuleSyncTest changes, so whichever merges first will force a rebase on the others.
Fixes #7313
Background
WebsocketCollector.SessionSendQueueadvances only when theAsyncRemote.sendTextcallbackruns. On a half-open connection the callback never fires:
sendingstays true forever,later messages pile up in an unbounded
ArrayDeque, and nothing closes the session orrequests reconciliation. A failed
SendResultor a synchronous send exception drops themessage but leaves the session open, so the gateway silently misses configuration updates.
Changes
Schedule a per-send watchdog timeout (default 30s): if the callback is not observed, the
session is treated as broken and closed;
A failed
SendResultor a synchronous send exception now closes the session instead ofdiscarding the message and continuing;
Bound the per-session queue (default 256 messages); overflow closes the session rather
than silently dropping messages;
Closing removes the session from all indexes and cancels pending timeouts; the gateway
reconnects and performs a full
MYSELFsynchronization, which makes every failure pathrecoverable;
Healthy sessions keep strict message ordering; test hooks allow tuning the timeout and
queue limit.
Assign the active watchdog future under the send-queue monitor and clear it after completion. Closing a session now actually cancels the in-flight timeout, as requested in review; a regression test checks cancellation rather than only waiting for the closed guard.
Merge current master while preserving InitialSync request/sequence framing. Independently apply only the RocketMQ test readiness fix: await exact created rule IDs before the single log-producing request, without relaxing log assertions or consumption timeouts.
Verification
./mvnw.cmd -pl shenyu-admin -am test -Dtest=WebsocketCollectorTest -Dsurefire.failIfNoSpecifiedTests=false -Dmaven.javadoc.skip=true: BUILD SUCCESS; 37 tests, zero failures/errors/skips; Checkstyle passed.git diff --checkpassed.verifywithLoggingRuleSyncTestpassed (2 tests, zero failures/errors/skips), including Checkstyle and RAT.03769725ais pushed; checks are pending and not yet verified green. Earlier failures were missing build artifacts and a rule-synchronization race, not failing collector assertions.