Repository navigation
Cross-worker WebSocket topics: pub/sub with MQTT filters, interest filter, publish rate limit (#2, #120) - #118
Merged
Merged
Conversation
) A worker is a thread with its own PHP context, so a room could not be a PHP array of connections: a broadcast only ever reached peers of the sending worker, and examples/ws-chat-server.php had to run setWorkers(1). Rooms now live in the server. $room = $server->room('chat'); // same room in every worker $room->join($ws); $room->broadcast($line, $ws); Share-nothing, the model uWebSockets uses for its topic tree: each worker owns its member table, and cross-worker work is handed to the owner through its mailbox (thread_mailbox.h, #81). A ws_session_t pointer never leaves its thread — only a room pointer, a ws id and a refcounted body do — so UAF is impossible by construction rather than by discipline. A room is just an identity (name + atomic refcount). It holds no member list and no member count, so join/leave are purely thread-local: not one atomic, not one lock. The refcount is what makes the pointer safe to hand across threads, and it lets a room die once its last member and last in-flight command are gone — a name->id dictionary would never release one, so dynamic names ("room:{uuid}") would leak for the life of the process. Broadcast allocates the body once for the whole fan-out and posts a command to every other worker, including ones holding no members: they find an empty table and do nothing. Delivery uses the non-suspending sink, so a backed-up peer drops the message rather than stalling the room. count() has no global tally to read, so it is a scatter/gather: each worker adds its local count into the shared query and the last to answer wakes the caller — the same shape as Socket.IO's fetchSockets() or Phoenix's per-node counters. It is a snapshot, not a live number, and a worker that misses the 1s bound is left out rather than hanging the call. Test 038: 6 clients across 4 workers — one sends, the other five receive it whichever worker owns them, the sender is excluded, and count() returns 6.
Contributor
CoverageTotal lines: 82.32% → 82.74% (+0.42 pp)
|
Review of the initial rooms implementation turned up five memory-safety bugs, one of them reproducible on the shipped example. - The hub was freed when the server stopped, while a WebSocketRoom the script still held kept pointing at it — the room's destructor then took a mutex in freed memory (valgrind: invalid read, freed by ws_hub_free). The hub is refcounted now; the server and every room object hold a reference. - count() handed a trigger event to the answering workers and disposed it on timeout, so a late answer fired a freed event. An event's ref_count is not atomic, so it cannot be shared across threads at all: answers now travel home through the asker's own mailbox and the query settles on one thread, with no cross-thread event. That also drops a uv_async handle per count() call. - broadcast()/count() read a worker's mailbox pointer without the hub lock, racing a detach that frees it. Claiming, retiring and posting to a slot now all happen under the lock, and a slot carries a generation so a reply cannot land in a reused slot. - Releasing the last reference to a room could free it twice: a lookup could revive it from zero while the first dropper waited for the mutex. Replaced with refcount_dec_and_lock. - getName() took a non-atomic refcount on a zend_string shared by every worker; it returns a copy now. - A detaching worker discarded its mailbox without draining it, leaking every queued command — once per worker, on every hot reload. Also: a full worker mailbox dropped room traffic silently (counted now, see ws_hub_dropped); ws_session_try_send duplicated the tail of ws_do_send, so the queue+flush mechanics moved to ws_session_queue_and_flush and the PHP layer kept only the policy; count() no longer swallows cancellation; broadcast scans only claimed slots, mints ws_id with an atomic, and skips the payload copy on a single-worker server; the broadcast() stub promised a room-wide delivery count it never returned. Tests: rooms had one test, which SIGKILLed the process, so a clean shutdown was never exercised. Adds single-worker + room-outlives-server, leave/auto-leave on disconnect, and room isolation + broadcastBinary + getName; the WebSocket client helpers move to a shared include.
Eight comments from the hardening pass either repeated a fact already stated at the field, in ws_hub.h, or on the very next line. What survives is only what the code cannot show: thread affinity, the deadlock and revive hazards, and why a cancellation must not come back as a number.
…tree fix (#2) Replaces the room API with MQTT-style topics. A room object was the wrong primitive: it belongs to the server but has to be carried into a worker by closure, and a wildcard filter cannot be interned as an object at all — `user/42/#` is a predicate, not a room. Nobody models it that way (Socket.IO keys an adapter by string, uWS has ws.subscribe/app.publish), so membership now lives in the connection and a topic is a string at the call site: $ws->subscribe('chat'); // or 'user/42/#', 'chat/+/typing' $ws->publish('chat', $msg); // publish topics must be concrete $ws->subscriberCount('chat'); Each worker keeps its own topic tree (ws_topic_tree.c) over the sessions it owns; a publish travels to the other workers as a STRING through the existing mailbox, and each matches it locally. A ws_session_t pointer never leaves its thread, so there is no shared registry to lock and no lifetime to get wrong. Interest filter: a publish used to wake EVERY worker, copy the payload into each mailbox and let most of them find no subscriber. Each worker now summarises its subscriptions in a counting Bloom filter and a publisher skips the workers that cannot match — what NATS propagates between nodes. A wildcard cannot be keyed by name, so the filter holds each subscription's leading literal prefix and the publisher probes every level-prefix of its topic: that can only ever produce a wasted wake-up, never a lost message. Measured on 8 workers: 44800 useless cross-worker posts -> 0, wildcard subscribers still reached. Fixes a SIGSEGV on every stop() with a subscriber still connected: start() detached from the hub — freeing the topic tree — before the scope drain destroyed the sessions, and a session unsubscribes itself as it is destroyed, so the teardown walked freed nodes. The other topic tests fclose() first, which is why they hid it; 043 now pins it. HttpServer::getRuntimeStats() reports ws_topic_posted / _skipped / _dropped; the drop counter existed but was never wired to anything. Also fixes the test client: it read the handshake with fread(4096) and silently discarded any frames that arrived in the same segment, and decoded only frames shorter than 126 bytes.
Conflicts were all additive — telemetry (#5) and the WS handler-exception fix (#119) landed in the same files as the topics work: - http_server_class.c: keep both the hub teardown and the stats slab release. - php_websocket.c: keep both the topic methods and getRemotePort(). - CHANGELOG: both sets of entries under [Unreleased]. - The arginfo headers are generated, so they were regenerated from the stubs rather than hand-merged. main added src/core/stats_registry.c, so configure has to be regenerated from config.m4 (phpize) — a stale Makefile builds a .so without http_stats_field_count and every #5 test dies on the undefined symbol. Full suite on the merge result: 340 passed, 0 failed.
Three things, all found by asking why the earlier review's "minor" notes were true at all. setWsMaxSubscriptions() caps the distinct topic filters one connection may hold. Default 0 — no limit, which is what every self-hosted broker ships (EMQX max_subscriptions, NATS max_subs): only the application knows how many topics it needs, and a number invented here would break a legitimate handler at runtime. The knob is for the application that pipes client input into subscribe(). Over the cap the FILTER is refused and the connection stays up — EMQX answers SUBACK 0x97, NATS -ERR 'Maximum Subscriptions Exceeded'; neither drops the client. A re-subscribe to a filter already held is idempotent and spends no quota. Filters may be 128 levels deep, up from 32 — the ceiling EMQX uses, and the only one of these limits that is on by default anywhere. Topics were unusable past 64 workers: setWorkers() allows 1024 but the hub had 64 slots, so a worker beyond that got no topic tree — subscribe() threw and publish() quietly did nothing. The slot table matches setWorkers() now, and a worker that still cannot attach fails the start instead of serving half a feature. The hub is reached through the connection's SERVER now, not through a thread-local: ws_session_t carries it, snapshotted at init. A thread-global was the wrong shape — it is one per thread, and two HttpServers can share a thread. Not fixed here, and it is the bigger one: http_server_class.c keeps the accepting server in a __thread global (current_server), so the second server to start() on a thread takes over EVERY listener's connections. Reproduced with plain HTTP, no WebSocket involved: two servers on two ports, both answer with the second one's handler. It predates all of this (it is in main and in the first commit), the accept callback has to learn its server from its own listener, and that is a central enough path to deserve its own issue. Until then, cross-server topic isolation cannot be observed even though the hub is now per-server — which is why this commit does not test it. fuzz_ws_frame had not linked on this branch since before topics landed (http_connection_destroy_if_idle_deferred had no stub); it links again.
…ure (#2) Topics were only ever exercised over plaintext HTTP/1, which left the three delivery paths that could plausibly differ untested. - 046 permessage-deflate. Each session negotiated its OWN extensions, so one publish() has to serve a compressed peer and a plain one side by side: the first gets an RSV1 frame it can inflate, the second the raw text. - 047 HTTP/2. A session is bound to a STREAM, not a connection, which is why a publish is addressed by ws_id per session. Two WS streams multiplexed on one TCP connection: excludeSelf must skip the publisher's stream and still reach its sibling on the same connection — keyed by connection it would silently swallow both. - 048 wss. Delivery goes out through whatever transport the session was built on, TLS included; 018 only echoed. - 049 backpressure. Delivery runs on the reactor with no coroutine to park, so a subscriber whose socket is backed up must DROP the message, not queue it (one dead reader would grow the worker without bound) and publish() must not suspend (one slow peer would stall delivery to everyone else). 200 x 2KB at a peer that never reads: some served, the rest dropped, all accounted for, loop completes. Also drops a clause from the hub-attach error that stopped being true when the thread's attachment list became keyed by hub: a second WebSocket server on the thread now attaches fine, so the only way to fail is running out of slots. NOTE for the runner: php-src/build-debug has no zlib, so its SKIPIF silently skips every permessage-deflate test. Run the suite with /usr/local/bin/php (same build id) — the WebSocket suite is then 49/49 with no skips.
…nst a foreign client (#2) publish() is the one WebSocket call an unprivileged peer can turn into work on EVERY worker in the process: send()/trySend() only ever touch its own socket, and the outbound FIFO is already capped per session. Unmetered, a client looping on a message the handler relays fills every worker's inbox — and once an inbox is full the drops hit OTHER topics' traffic too, so the damage is process-wide rather than confined to the abuser. setWsPublishRateLimit($perSecond, $burst = 0) puts a token bucket on it, per connection, off by default — the same default EMQX ships for messages_rate, because only the application knows its own traffic. Over the rate publish() throws WebSocketBackpressureException and the connection stays up: the sender is TOLD. That was the gap the issue actually named — an over-rate message otherwise disappears into a full mailbox where neither the sender nor anyone else can see it. Credit is carried in nanoseconds of allowance rather than whole tokens: the coarse clock ticks every 4-10ms, so integer tokens would round away to nothing at any rate finer than that. Per-topic fair queueing (the issue's other suggestion) is not needed once no single connection can fill the mailbox — the eviction it was meant to prevent took a flood to cause. The drop counter it asked to surface is already in getRuntimeStats() as ws_topic_dropped. e2e/topics drives all of this with the Python `websockets` library — an implementation nobody here wrote, negotiating permessage-deflate by default, so the publishes it receives are compressed frames decoded by code that has never seen ours. Every topic test until now spoke to a hand-rolled client in this repo, which cannot catch a framing habit we got wrong on both sides. 9/9: cross-worker delivery over 4 workers, '+' spanning exactly one level, '#' at any depth, an unsubscribed topic reaching nobody, subscriberCount tallying every worker, and a 4 KiB payload arriving intact. Closes #120.
…ue (#2) A quality pass over the topic code before review. No behaviour change; suite still 350/0, e2e 9/9, fuzz_ws_frame links. Names (CODING_STANDARDS §3 — nothing single-letter crosses a function boundary): - ws_hub.c had one type under three names (`local`, `me`, `mine`) and, worse, the same two names meaning something else in ws_hub_count: `local` was the match COUNT and `me` was the coroutine. Now `local` is always the attachment, the count is `local_matches`, the coroutine is `coroutine`. - ws_topic_tree.c: `v` (a visit context threaded through five functions) is `visit`; `lv`/`ln` are `level`/`level_len`; `n`/`out` are `count`/`kept`; `deliver` reads as the predicate `should_deliver`. const and types: - `except_id` is computed once, so it is const. - ws_topic_unsubscribe only compares its filter — `const zend_string *`. - `walking` sits beside two uint32_t counters; it is one now. - The two setters cast to zend_ulong before comparing against UINT32_MAX, like their neighbours; `count > UINT32_MAX` is a tautology on a 32-bit build. - zend_throw_exception_ex with no format argument is just zend_throw_exception. Comments that had quietly stopped being true — the ones worth catching: - ws_hub.h claimed the hub "publishes to everyone and the workers filter". The interest filter is exactly what stopped that being so. - The stats block still said "there is no rate limit yet" (#120 added one). - Four separate places argued that the hub is per-server "because two HttpServers can share a thread" — which CODING_STANDARDS §1.1 rules out as a deliberate invariant. The real reason is §1.2: per-server state does not live in a thread-global. Said once, correctly. - A comment pointed at `node_remove`, a function that does not exist (ws_node_detach). - ws_interest_bump carried the docs of the function it replaced (talking about a NULL return it does not have). - Dropped the ones that restate the line below them.
… refcount (#2) The only conflict was CHANGELOG — both sides added bullets under [Unreleased], kept both. It also still carried a claim I had already removed from the code: that the hub is per-server "because two HttpServers can share a thread", which CODING_STANDARDS §1.1 rules out. The real reason is §1.2. ws_hub_addref() was declared, defined and called from nowhere — a leftover from the room API, where a WebSocketRoom held a reference. With rooms gone the hub has exactly one owner (the server that created it; clones borrow the pointer and never free it, and http_server_free releases it under ws_hub_owner), so the refcount was permanently 1. The field and its atomics go with it. Full suite on the merge result: 350 passed, 0 failed.
Second quality pass — the items the first one left on the table. No behaviour change; 350/0. - ws_hub_count was 85 lines, over the ~80 CODING_STANDARDS §13a.3 calls a smell. The fan-out — build the interest probe, take `admin`, post one COUNT per interested worker, count what went out — is a phase of its own, so it is ws_hub_ask_others() now. 58 lines and 30, and the critical section is no longer buried in the middle of a function that also allocates the query and parks the coroutine. - The interest filter was updated by the same three-line expression twice, once wrapped so awkwardly it took three lines on its own. One helper, one place that knows a filter's prefix is what the hub wants. - ws_node_free(node, bool self) is two functions: ws_node_free_contents() for the root, which is embedded in the tree and must not be freed, and ws_node_free() for everything else. `ws_node_free(node, true)` told a reader nothing.
macOS CI caught it: the test asserted that a matching publish woke another worker, which is only true if the kernel spread the connections. With SO_REUSEPORT it does; on the shared-fd pool path all eight peers can land on the publisher's own worker, delivery is purely local, and ws_topic_posted stays 0. ws_topic_posted is a property of where the connections went, not of the interest filter. What the filter actually promises is still pinned: an unsubscribed topic wakes nobody (posted 0, skipped > 0 — true whatever the distribution, since the other workers are attached either way), and a matching one reaches every subscriber.
, #2) Merges main (issue #117 — pool-mode stop() — landed) and fixes three bugs an adversarial review of the refactors turned up. Full suite on the merge: 355 passed, 0 failed. The rate limit could be walked straight past. ws_do_publish gated on `w->session != NULL && !allowed(...)`, and the session is built lazily by the first WS I/O — so a handler whose first act is publish() has none, and the whole token bucket was skipped. Reproduced: 50 publishes accepted against a limit of 10/s. It is now 5 through, 45 refused, and 050 pins it. The same bucket was per-STREAM on HTTP/2. It lived in ws_session_t, and H2 puts many sessions on one TCP — a client multiplied its allowance by the number of streams it opened. Both are the same fix: the bucket moves to http_connection_t, which is the granularity #120 is actually about (one client, one connection). The arena zeroes a slot on hand-out, so the new fields need no explicit reset. And w->conn is never cleared on teardown — only w->session and w->closed are — so publish()/subscriberCount() on a $ws that outlived its handler (the spawned-writer pattern of 025) followed a pointer into a connection slot the arena had already given to someone else. ws_live_conn_of() refuses to follow it once `closed`. With #117 fixed upstream, the topic tests stop() cleanly instead of SIGKILLing themselves. That was not cosmetic: a SIGKILLed process writes no .gcda, so the coverage of everything 038/042 exercise — the whole cross-worker path, which is unreachable with one worker — was being thrown away. New, all four in pool mode, because the single-worker versions could not have caught what they cover: - 051 both limits reach the worker CLONES. A field missed in one of the three config copy paths would leave a limit silently off in the pool while 045/050 stayed green. - 052 permessage-deflate across workers. The payload crosses raw and the RECEIVING worker deflates it per session — compress once at the source and share it by refcount, and the far side decodes garbage. - 053 backpressure across workers: delivery to a stalled peer happens on the receiving worker's reactor. It drops its own traffic and stalls nobody else. - 054 the counting Bloom gives interest BACK on unsubscribe. A leaked decrement breaks nothing visible — delivery stays correct — it just wakes every worker for nobody, forever. No delivery assertion can see that; this one asserts the workers go back to being skipped.
#2) Covers the untrusted-input surface of the cross-worker topic feature — topic names and subscribe filters arrive verbatim from the network. - fuzz/fuzz_ws_topic.c: a libFuzzer target for ws_topic_tree.c (the MQTT filter parser, the interest-prefix builder, and the wildcard matcher + node lifecycle). This TU was previously weak-stubbed out of the frame fuzzer, so the parser/matcher had no fuzz coverage. 365k runs clean under ASAN+UBSan. - tests/phpt/websocket/055-topics-matching-oracle.phpt: a differential test of the C matcher against a PHP MQTT reference oracle over a full filter x topic matrix, driven through subscriberCount() (same tree walk as publish, but synchronous, so no timing flakes). Also releases a one-time zend_string leak in the shared fuzz harness that a leak-clean run of the new target surfaced (harness_common.c) — it had been leaking on every embedded fuzzer's init.
test(websocket): fuzz the topic engine + differential-test the matcher (#2)
This was referenced Jul 14, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #2. Closes #120.
What
A publish now reaches subscribers on every worker of the process. Until now it could only reach peers of the sending worker — a worker is a thread with its own PHP context — so a chat had to run
setWorkers(1).Why there is no room object
This branch started with
HttpServer::room()returning aWebSocketRoom. That was the wrong primitive: a room belongs to the server but has to be carried into a worker by closure, and a wildcard filter cannot be interned as an object at all —user/42/#is a predicate, not a room. Nobody models it that way either: Socket.IO keys an adapter by string, uWS hasws.subscribe()/app.publish(), Swoole has no primitive. So membership lives in the connection and a topic is a string at the call site — nothing to obtain, hold, or pass into a handler.Design
Each worker keeps its own topic tree over the sessions it owns. A publish travels to the other workers as a string through the existing mailbox, and each matches it locally. A
ws_session_tpointer never leaves its thread, so there is no shared registry to lock and no lifetime to get wrong.Interest filter. A publish used to wake every worker — copy the payload, post a command, fire an eventfd — and most of them then found no subscriber. Each worker now summarises its subscriptions in a counting Bloom filter and a publisher skips the workers that cannot match; this is the "interest" NATS propagates between nodes. A wildcard cannot be keyed by name, so the filter holds each subscription's leading literal prefix and the publisher probes every level-prefix of its topic: that can only ever cost a wasted wake-up, never lose a message. Measured on 8 workers: 44 800 useless cross-worker posts → 0, wildcard subscribers still reached.
It degrades honestly: an unbounded topic space (
order/{uuid}/status) saturates the filter and the hub is back to waking everyone — never to losing a message.Fixed along the way
stop()with a subscriber still connected. The worker let go of its topic tree before the scope drain destroyed the sessions — and a session unsubscribes itself as it is destroyed, so the teardown walked freed nodes. The other topic testsfclose()first, which is why they hid it;043pins it now.setWorkers()allows 1024 but the hub had 64 slots, so workers beyond that got no tree:subscribe()threw andpublish()quietly did nothing. A worker that cannot attach now fails the start instead of serving half a feature.Limits (#120)
publish()is the one WebSocket call an unprivileged peer can turn into work on every worker —send()/trySend()only ever touch its own socket.setWsPublishRateLimit($perSecond, $burst = 0)puts a per-connection token bucket on it, off by default (as EMQX shipsmessages_rate). Over the ratepublish()throwsWebSocketBackpressureExceptionand the connection stays up: the sender is told, instead of the message vanishing into a full mailbox where nobody can see it.setWsMaxSubscriptions()caps distinct filters per connection, also 0 = unlimited by default — every self-hosted broker defaults that way (EMQXmax_subscriptions, NATSmax_subs), because only the application knows how many topics it needs. Filter depth is capped always: 128 levels, the ceiling EMQX uses.getRuntimeStats()reportsws_topic_posted/ws_topic_skipped/ws_topic_dropped.Tests
WebSocket suite 53/53, no skips; full suite 350 passed, 0 failed.
Topics are covered on every transport, not just plaintext HTTP/1:
publish()serves a compressed peer and a plain one side by side.excludeSelfmust skip the publisher's stream and still reach its sibling on the same connection.websocketslibrary, an implementation nobody here wrote, negotiating permessage-deflate by default. Every other topic test speaks to a hand-rolled client in this repo, which cannot catch a framing habit we got wrong on both sides. 9/9.Notes for the reviewer
038and042— the cross-worker tests — must exit withSIGKILL, because HttpServer::stop() unsupported in pool mode (workers > 1) — cross-thread shutdown not wired up #117 (HttpServer::stop()unsupported in pool mode) is open. A SIGKILLed process never writes.gcda, so everything they exercise is invisible to gcov:ws_hub_post_locked,ws_hub_drain,ws_hub_answer_count,ws_query_settle,ws_interest_matches, the payload/cmd allocators. With one worker the fan-out loop skips its own slot and none of it runs. I tried converting both tostop(): they pass alone and hang the suite — reverted. This cannot be both tested and counted until HttpServer::stop() unsupported in pool mode (workers > 1) — cross-thread shutdown not wired up #117 lands.HttpServers on one thread are out of scope by design (CODING_STANDARDS §1.1), and the accepting server is a__threadglobal (current_server), so the second tostart()takes over every listener's connections. Pre-existing, reproducible with plain HTTP and no WebSocket; not touched here.format-code.shwas not run. The changed files were read against §3/§4/§6/§13a by hand.