Skip to content

Fixes #7313: Close stalled websocket sessions and bound the admin send queue - #7343

Merged
Aias00 merged 6 commits into
apache:masterfrom
BobSong-dev:fix/7313-websocket-send-queue-guard
Oct 1, 2026
Merged

Aias00 merged 6 commits into
apache:masterfrom
BobSong-dev:fix/7313-websocket-send-queue-guard

Conversation

@BobSong-dev

@BobSong-dev BobSong-dev commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #7313

Background

WebsocketCollector.SessionSendQueue advances only when the AsyncRemote.sendText callback

runs. On a half-open connection the callback never fires: sending stays true forever,

later messages pile up in an unbounded ArrayDeque, and nothing closes the session or

requests reconciliation. A failed SendResult or a synchronous send exception drops the

message 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 SendResult or a synchronous send exception now closes the session instead of

    discarding 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 MYSELF synchronization, which makes every failure path

    recoverable;

  • 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

  • Local: ./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.
  • Local: Admin Apache RAT and git diff --check passed.
  • Local: the six-module RocketMQ e2e reactor verify with LoggingRuleSyncTest passed (2 tests, zero failures/errors/skips), including Checkstyle and RAT.
  • GitHub CI: latest head 03769725a is pushed; checks are pending and not yet verified green. Earlier failures were missing build artifacts and a rule-synchronization race, not failing collector assertions.
  • Not run locally: full-project tests or Docker e2e.

}
inFlightMessage = message;
}
final ScheduledFuture<?> future = SEND_WATCHDOG.schedule(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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
Aias00 previously approved these changes Sep 29, 2026

@Aias00 Aias00 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 ConcurrentModificationException on the publish paths. forceClose → removeSessionIndexes mutates SESSION_SET and NAMESPACE_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-362 removes from SESSION_SET, drops the SessionSendQueue via removeSessionSendQueue, and unwinds the namespace index, so a forced close cleans up the same way @OnClose → clearSession does.
  • The queue is not left "stuck sending". forceClose sets closed = true and sending = false and clears the backlog, so the next send() returns without queueing and the next sendNext() bails out. This is the piece master was missing — a failed async send left sending true 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

  1. session.close() is called while the caller holds the queue monitor. send():451 and onSendTimeout():521 both call forceClose inside synchronized (this), and forceClose only releases that monitor after its own block, so session.close() at :542 runs under it. It is safe today because containers dispatch @OnClose asynchronously, but a container that invoked the close handler inline would take clearSession → removeSessionSendQueue → close() onto the same monitor. Cheap to avoid: have forceClose return a "close me" signal and do the actual session.close() outside the lock.
  2. One scheduled task per message. sendNext schedules 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.
  3. 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.
  4. 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
Aias00 previously approved these changes Sep 30, 2026

@Aias00 Aias00 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 Aias00 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@Aias00
Aias00 merged commit 40e6879 into apache:master Oct 1, 2026
39 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[BUG] WebSocket SessionSendQueue can stall indefinitely and lose configuration updates

2 participants