Skip to content
Draft
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
fix(scaling): make the flush deferral actually install, and test it
`require('engine.io/build/socket')` cannot resolve, in two independent ways:
engine.io is a transitive dependency of socket.io and therefore invisible from
`src` under pnpm (MODULE_NOT_FOUND), and `build/socket` is not in engine.io's
`exports` map, so the deep path fails with ERR_PACKAGE_PATH_NOT_EXPORTED even
where the package is resolvable. The failure was downgraded to a warning, so
`engineFlushDefer: true` left the stock synchronous flush in place and the feature
was inert — which means the N=3 numbers in the PR description cannot have measured
this patch.

Resolve engine.io through its public entry point relative to socket.io instead,
which also guarantees we patch the copy socket.io actually loaded rather than a
second one a hoisting layout might provide.

Also sync the copied body to the installed engine.io 6.6.9: upstream attaches the
payload on a truthiness check, not `!== undefined`.

New spec covers coalescing, packet order, the closing/closed early return, send
callbacks, the callback-in-options-slot form, upstream's payload semantics, and
pins the SHA of engine.io's own sendPacket so an upgrade that touches it fails
here and forces a re-vet of the copy. It restores the prototype afterwards —
otherwise every later spec in the suite runs with the feature silently enabled.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
  • Loading branch information
SamTV12345 and claude committed Jul 27, 2026
commit f5b91431fd16049122f156e4016a819e58777e03
24 changes: 18 additions & 6 deletions src/node/utils/EngineFlushDeferral.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,10 +50,19 @@ export const installEngineFlushDeferral = (): void => {

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
SocketProto = require('engine.io/build/socket').Socket.prototype;
const engineIo = require(require.resolve('engine.io', {paths: [require.resolve('socket.io')]}));
SocketProto = engineIo.Socket.prototype;
} catch (err: any) {
logger.warn(`Unable to install engine.io flush deferral (module not found): ${err && err.message || err}`);
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') {
Expand All @@ -68,9 +77,10 @@ export const installEngineFlushDeferral = (): void => {

// Re-implementing sendPacket inline rather than wrapping the original
// so the single closing `this.flush()` becomes a microtask-coalesced
// schedule. The body is intentionally a near-verbatim copy of the
// engine.io 6.6.5 implementation so future engine.io upgrades that
// change packet-shape semantics still need re-vetting.
// 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) {
Expand All @@ -82,7 +92,9 @@ export const installEngineFlushDeferral = (): void => {
options = options || {};
options.compress = options.compress !== false;
const packet: any = {type, options};
if (data !== undefined) packet.data = data;
// 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);
Expand Down
143 changes: 143 additions & 0 deletions src/tests/backend/specs/engineFlushDeferral.ts
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}`);
});
});
Loading