-
-
Notifications
You must be signed in to change notification settings - Fork 3k
feat(scaling): engine.io socket flush deferral — modest tail-latency reduction at mid-load #7774
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Draft
JohnMcLear
wants to merge
4
commits into
develop
Choose a base branch
from
feat/engine-io-flush-deferral
base: develop
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Draft
Changes from all commits
Commits
Show all changes
4 commits
Select commit
Hold shift + click to select a range
073954d
feat(scaling): engine.io socket flush deferral (#7756 / #7767)
JohnMcLear f0661a9
fix(engine-flush-defer): address Qodo review
JohnMcLear df6b84c
Merge remote-tracking branch 'origin/develop' into pr7774
SamTV12345 f5b9143
fix(scaling): make the flush deferral actually install, and test it
SamTV12345 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
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
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
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
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,111 @@ | ||
| // Engine.io socket flush deferral — #7756 / #7767 deeper investigation | ||
| // after the simple WS transport-level packing prototype (#7772) showed | ||
| // that the writeBuffer almost never accumulates because flush() drains | ||
| // immediately on `transport.writable === true`. | ||
| // | ||
| // engine.io's Socket.sendPacket(...) ends with: | ||
| // | ||
| // this.writeBuffer.push(packet); | ||
| // if (callback) this.packetsFn.push(callback); | ||
| // this.flush(); // <-- synchronous | ||
| // | ||
| // flush() reads writeBuffer and hands it to transport.send. For | ||
| // WebSocket, transport.writable is true again within microseconds of | ||
| // each write, so each sendPacket() call drains a buffer of size 1. The | ||
| // transport.send([packets]) function then iterates packets and writes | ||
| // one WS frame per packet — which is what the polling transport's | ||
| // natural encodePayload batching avoids. | ||
| // | ||
| // This patch coalesces synchronous-task sendPacket calls onto a single | ||
| // microtask-scheduled flush. Inside the same JS task, multiple | ||
| // sendPacket() calls accumulate in writeBuffer; the queued microtask | ||
| // then calls flush() once with the whole batch. The transport's | ||
| // send([batch]) sees N > 1 packets and the WS payload-encoding fast | ||
| // path (also added by lever 8) coalesces them into one frame. | ||
| // | ||
| // Microtask deferral adds zero meaningful wall-clock latency: | ||
| // microtasks drain before the next macrotask, so any consumer waiting | ||
| // on the next setImmediate / setTimeout / I/O callback still sees the | ||
| // flush completed. | ||
| // | ||
| // Forward-compatible. Existing clients receive identical wire bytes | ||
| // because the engine.io packet encoding is unchanged; the difference | ||
| // is only how many engine.io packets share one transport-level send | ||
| // call. The WS transport's send([packets]) path is then where lever 8 | ||
| // (or this patch's accompanying engine-packing branch) decides | ||
| // whether to ship them as N frames or one payload-encoded frame. | ||
| // | ||
| // Gated by settings.engineFlushDefer. Default off; production unaffected. | ||
|
|
||
| import log4js from 'log4js'; | ||
|
|
||
| const logger = log4js.getLogger('engine-flush-defer'); | ||
|
|
||
| let installed = false; | ||
|
|
||
| const SCHEDULED = Symbol('engineFlushScheduled'); | ||
|
|
||
| export const installEngineFlushDeferral = (): void => { | ||
| if (installed) return; | ||
|
|
||
| let SocketProto: {sendPacket: (...a: unknown[]) => unknown}; | ||
| try { | ||
| // Resolve engine.io through socket.io, and only via its public entry point: | ||
| // - engine.io is a transitive dependency of socket.io, so under pnpm's | ||
| // layout it is not resolvable from `src` at all (MODULE_NOT_FOUND), and | ||
| // - `engine.io/build/socket` is not in engine.io's `exports` map, so a deep | ||
| // require fails with ERR_PACKAGE_PATH_NOT_EXPORTED even when it is. | ||
| // Resolving relative to socket.io also guarantees we patch the very copy | ||
| // socket.io loaded, not a second one that a hoisting layout might provide. | ||
| // eslint-disable-next-line @typescript-eslint/no-require-imports | ||
| const engineIo = require(require.resolve('engine.io', {paths: [require.resolve('socket.io')]})); | ||
| SocketProto = engineIo.Socket.prototype; | ||
| } catch (err: any) { | ||
| logger.warn('engineFlushDefer is enabled but the engine.io Socket class could not be ' + | ||
| `resolved, so the patch is NOT active: ${err && err.message || err}`); | ||
| return; // Leave `installed` false so a later boot path can retry. | ||
| } | ||
| if (typeof SocketProto.sendPacket !== 'function') { | ||
| logger.warn('engine.io Socket shape unexpected; skipping flush deferral patch'); | ||
| return; // Leave `installed` false so a later boot path can retry. | ||
| } | ||
| // Only after both require and shape check succeed do we record that the | ||
| // patch is installed. Setting the flag before validation (the original | ||
| // code) would have permanently disabled retries after a transient | ||
| // require failure in test/CI environments where socket.io may load late. | ||
| installed = true; | ||
|
|
||
| // Re-implementing sendPacket inline rather than wrapping the original | ||
| // so the single closing `this.flush()` becomes a microtask-coalesced | ||
| // schedule. The body is a verbatim copy of engine.io 6.6.9's | ||
| // implementation apart from that last line — engineFlushDeferral.ts pins the | ||
| // upstream source with a fingerprint, so an engine.io upgrade that touches | ||
| // sendPacket fails the test suite and forces a re-vet of this copy. | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any | ||
| SocketProto.sendPacket = function (this: any, type: any, data: any, options: any, callback: any) { | ||
| if ('function' === typeof options) { | ||
| callback = options; | ||
| options = {}; | ||
| } | ||
| if ('closing' === this.readyState || 'closed' === this.readyState) return; | ||
|
|
||
| options = options || {}; | ||
| options.compress = options.compress !== false; | ||
| const packet: any = {type, options}; | ||
| // Upstream uses a truthiness check, not `!== undefined`: a falsy payload is | ||
| // deliberately not attached (encodePacket renders `type` alone for it). | ||
| if (data) packet.data = data; | ||
| this.emit('packetCreate', packet); | ||
| this.writeBuffer.push(packet); | ||
| if ('function' === typeof callback) this.packetsFn.push(callback); | ||
|
|
||
| if (this[SCHEDULED]) return; | ||
| this[SCHEDULED] = true; | ||
| queueMicrotask(() => { | ||
| this[SCHEDULED] = false; | ||
| this.flush(); | ||
| }); | ||
| }; | ||
|
|
||
| logger.info('engine.io socket flush deferral enabled (#7756 / #7767)'); | ||
| }; |
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
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
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,143 @@ | ||
| 'use strict'; | ||
|
|
||
| /** | ||
| * Tests for the opt-in engine.io flush deferral (settings.engineFlushDefer). | ||
| * | ||
| * Two jobs: | ||
| * 1. Prove the patch actually installs and coalesces. The first cut resolved the | ||
| * Socket class via `require('engine.io/build/socket')`, which cannot work: | ||
| * engine.io is a transitive dependency of socket.io (not resolvable from | ||
| * `src` under pnpm) and `build/socket` is not in its `exports` map. The | ||
| * failure was caught and logged, so the feature was silently inert. | ||
| * 2. Pin engine.io's own sendPacket source. The patch re-implements that method, | ||
| * so an engine.io upgrade that changes it must fail here and force a re-vet | ||
| * rather than silently diverging. | ||
| */ | ||
|
|
||
| const assert = require('assert').strict; | ||
| const crypto = require('crypto'); | ||
| const {installEngineFlushDeferral} = require('../../../node/utils/EngineFlushDeferral'); | ||
|
|
||
| // sha256 of engine.io 6.6.9's Socket.prototype.sendPacket source. Update this | ||
| // ONLY together with a re-read of the upstream method against the copy in | ||
| // EngineFlushDeferral.ts. | ||
| const PINNED_SEND_PACKET_SHA = | ||
| 'a94abb8dd747d4f55bc7611143b185d384463afd5ee9f006a5cde5a1851c0a05'; | ||
|
|
||
| const engineIoSocketProto = () => { | ||
| const path = require.resolve('engine.io', {paths: [require.resolve('socket.io')]}); | ||
| return require(path).Socket.prototype; | ||
| }; | ||
|
|
||
| // Minimal stand-in for an engine.io Socket: the patch only touches these. | ||
| const fakeSocket = () => ({ | ||
| readyState: 'open', | ||
| writeBuffer: [] as any[], | ||
| packetsFn: [] as any[], | ||
| flushed: 0, | ||
| flushedSizes: [] as number[], | ||
| emit() {}, | ||
| flush() { | ||
| this.flushed++; | ||
| this.flushedSizes.push(this.writeBuffer.length); | ||
| this.writeBuffer = []; | ||
| }, | ||
| }); | ||
|
|
||
| describe(__filename, function () { | ||
| let sendPacket: any; | ||
| let upstreamSendPacket: any; | ||
|
|
||
| before(function () { | ||
| // Capture upstream BEFORE patching, so the drift check below compares against | ||
| // the real engine.io implementation. | ||
| upstreamSendPacket = engineIoSocketProto().sendPacket; | ||
| installEngineFlushDeferral(); | ||
| sendPacket = engineIoSocketProto().sendPacket; | ||
| }); | ||
|
|
||
| after(function () { | ||
| // The patch mutates a shared prototype, so leaving it in place would silently | ||
| // run every later spec in the suite with the feature enabled. | ||
| engineIoSocketProto().sendPacket = upstreamSendPacket; | ||
| }); | ||
|
|
||
| it('installs onto the engine.io copy socket.io actually loaded', function () { | ||
| // The bug this guards: a resolution failure downgraded to a warning, leaving | ||
| // the stock synchronous implementation in place. | ||
| assert.match(sendPacket.toString(), /queueMicrotask/, | ||
| 'engineFlushDefer did not patch engine.io Socket.prototype.sendPacket'); | ||
| }); | ||
|
|
||
| it('coalesces packets sent in one task into a single flush', async function () { | ||
| const s = fakeSocket(); | ||
| sendPacket.call(s, 'message', 'a'); | ||
| sendPacket.call(s, 'message', 'b'); | ||
| sendPacket.call(s, 'message', 'c'); | ||
| assert.equal(s.flushed, 0, 'flush must not happen synchronously'); | ||
| assert.equal(s.writeBuffer.length, 3, 'packets accumulate in writeBuffer'); | ||
|
|
||
| await Promise.resolve(); | ||
| assert.equal(s.flushed, 1, 'exactly one flush per task'); | ||
| assert.deepEqual(s.flushedSizes, [3], 'the flush carries the whole batch'); | ||
| }); | ||
|
|
||
| it('preserves packet order and payloads', async function () { | ||
| const s = fakeSocket(); | ||
| const seen: any[] = []; | ||
| s.flush = function () { seen.push(...this.writeBuffer.map((p: any) => p.data)); }; | ||
| for (const d of ['1', '2', '3']) sendPacket.call(s, 'message', d); | ||
| await Promise.resolve(); | ||
| assert.deepEqual(seen, ['1', '2', '3']); | ||
| }); | ||
|
|
||
| it('schedules a fresh flush for the next task', async function () { | ||
| const s = fakeSocket(); | ||
| sendPacket.call(s, 'message', 'a'); | ||
| await Promise.resolve(); | ||
| sendPacket.call(s, 'message', 'b'); | ||
| await Promise.resolve(); | ||
| assert.equal(s.flushed, 2); | ||
| assert.deepEqual(s.flushedSizes, [1, 1]); | ||
| }); | ||
|
|
||
| it('drops packets once the socket is closing or closed', async function () { | ||
| for (const readyState of ['closing', 'closed']) { | ||
| const s = fakeSocket(); | ||
| s.readyState = readyState; | ||
| sendPacket.call(s, 'message', 'a'); | ||
| await Promise.resolve(); | ||
| assert.equal(s.writeBuffer.length, 0, `${readyState}: nothing queued`); | ||
| assert.equal(s.flushed, 0, `${readyState}: nothing flushed`); | ||
| } | ||
| }); | ||
|
|
||
| it('keeps send callbacks, and treats a callback in the options slot as one', async function () { | ||
| const s = fakeSocket(); | ||
| const cb = () => {}; | ||
| sendPacket.call(s, 'message', 'a', cb); | ||
| assert.deepEqual(s.packetsFn, [cb], 'callback passed in the options position'); | ||
| assert.equal(s.writeBuffer[0].options.compress, true, 'compression defaults to on'); | ||
| await Promise.resolve(); | ||
| }); | ||
|
|
||
| it('attaches only truthy payloads, like upstream', async function () { | ||
| const s = fakeSocket(); | ||
| sendPacket.call(s, 'pong'); | ||
| sendPacket.call(s, 'message', ''); | ||
| assert.equal('data' in s.writeBuffer[0], false, 'no data for a bare pong'); | ||
| assert.equal('data' in s.writeBuffer[1], false, 'no data for an empty payload'); | ||
| await Promise.resolve(); | ||
| }); | ||
|
|
||
| it('engine.io sendPacket is unchanged since the copy was vetted', function () { | ||
| assert.doesNotMatch(upstreamSendPacket.toString(), /queueMicrotask/, | ||
| 'the patch was already installed before this spec ran, so upstream could not ' + | ||
| 'be captured — run this spec in its own process to check for drift'); | ||
| const sha = crypto.createHash('sha256').update(upstreamSendPacket.toString()).digest('hex'); | ||
| assert.equal(sha, PINNED_SEND_PACKET_SHA, | ||
| 'engine.io Socket.prototype.sendPacket changed upstream — re-vet the copy in ' + | ||
| 'EngineFlushDeferral.ts against the new implementation, then update ' + | ||
| `PINNED_SEND_PACKET_SHA to ${sha}`); | ||
| }); | ||
| }); |
Oops, something went wrong.
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.
Uh oh!
There was an error while loading. Please reload this page.