From 9f5ecbbc3f36482049c681e82c3aaf33c70730b4 Mon Sep 17 00:00:00 2001 From: xingsy97 <87063252+xingsy97@users.noreply.github.com> Date: Tue, 22 Sep 2026 18:26:06 +0800 Subject: [PATCH 1/2] fix(sio): isolate recovery offsets across child namespaces --- packages/socket.io/lib/parent-namespace.ts | 3 +- .../test/connection-state-recovery.ts | 66 ++++++++++++++++++- packages/socket.io/test/namespaces.ts | 49 ++++++++++++++ 3 files changed, 116 insertions(+), 2 deletions(-) diff --git a/packages/socket.io/lib/parent-namespace.ts b/packages/socket.io/lib/parent-namespace.ts index d3468bf46f..5e2263b903 100644 --- a/packages/socket.io/lib/parent-namespace.ts +++ b/packages/socket.io/lib/parent-namespace.ts @@ -115,7 +115,8 @@ export class ParentNamespace< class ParentBroadcastAdapter extends Adapter { broadcast(packet: any, opts: BroadcastOptions) { this.nsp.children.forEach((nsp) => { - nsp.adapter.broadcast(packet, opts); + // Each child adapter may append its own recovery offset to the packet. + nsp.adapter.broadcast({ ...packet, data: [...packet.data] }, opts); }); } } diff --git a/packages/socket.io/test/connection-state-recovery.ts b/packages/socket.io/test/connection-state-recovery.ts index c0dcbf0703..e313d7596d 100644 --- a/packages/socket.io/test/connection-state-recovery.ts +++ b/packages/socket.io/test/connection-state-recovery.ts @@ -1,6 +1,12 @@ import { Server, Socket } from ".."; import expect from "expect.js"; -import { waitFor, eioHandshake, eioPush, eioPoll } from "./support/util"; +import { + waitFor, + eioHandshake, + eioPush, + eioPoll, + createClient, +} from "./support/util"; import { createServer, Server as HttpServer } from "http"; import { Adapter } from "socket.io-adapter"; @@ -77,6 +83,64 @@ describe("connection state recovery", () => { io.close(); }); + it("should recover twice after replaying a parent namespace room broadcast", async () => { + const io = new Server(0, { connectionStateRecovery: {} }); + const parent = io.of(/^\/dynamic-\d+$/); + parent.on("connection", (socket) => socket.join("some-room")); + + const clients = ["/dynamic-101", "/dynamic-102"].map((nsp) => + createClient(io, nsp, { forceNew: true, reconnection: false }), + ); + + try { + await Promise.all(clients.map((client) => waitFor(client, "connect"))); + const ids = clients.map((client) => client.id); + + const initialPackets = clients.map((client) => waitFor(client, "seed")); + parent.emit("seed"); + await Promise.all(initialPackets); + + const missedPackets = clients.map((client) => waitFor(client, "hello")); + + for (let attempt = 0; attempt < 2; attempt++) { + await Promise.all( + clients.map((client) => { + const socket = io.of(client.nsp).sockets.get(client.id); + const disconnected = Promise.all([ + waitFor(client, "disconnect"), + waitFor(socket, "disconnect"), + ]); + socket.conn.close(); + return disconnected; + }), + ); + + if (attempt === 0) { + parent.to("some-room").emit("hello", "world"); + } + // The second recovery must use the replayed offset, without a new event. + await Promise.all( + clients.map((client) => { + const connected = waitFor(client, "connect"); + client.connect(); + return connected; + }), + ); + + for (let i = 0; i < clients.length; i++) { + expect(clients[i].recovered).to.be(true); + expect(clients[i].id).to.be(ids[i]); + } + if (attempt === 0) { + expect(await Promise.all(missedPackets)).to.eql(["world", "world"]); + } + } + } finally { + clients.forEach((client) => client.disconnect()); + await io.close(); + } + }); + it("should restore rooms and data attributes", async () => { const httpServer = createServer().listen(0); const io = new Server(httpServer, { diff --git a/packages/socket.io/test/namespaces.ts b/packages/socket.io/test/namespaces.ts index a8b3f08fcf..98c32f2d34 100644 --- a/packages/socket.io/test/namespaces.ts +++ b/packages/socket.io/test/namespaces.ts @@ -7,6 +7,7 @@ import { successFn, createPartialDone, assert, + waitFor, } from "./support/util"; describe("namespaces", () => { @@ -556,6 +557,54 @@ describe("namespaces", () => { }); }); + for (const recovery of [false, true]) { + for (const binary of [false, true]) { + it(`should isolate room broadcasts across child namespaces (recovery: ${recovery}, binary: ${binary})`, async () => { + const io = new Server( + 0, + recovery ? { connectionStateRecovery: {} } : {}, + ); + const parent = io.of(/^\/dynamic-\d+$/); + parent.on("connection", (socket) => socket.join("some-room")); + + const clients = ["/dynamic-101", "/dynamic-102"].map((nsp) => + createClient(io, nsp, { forceNew: true }), + ); + const payload = binary ? Buffer.from("hello") : { hello: "world" }; + + try { + await Promise.all( + clients.map((client) => waitFor(client, "connect")), + ); + + const received = clients.map( + (client) => + new Promise((resolve) => { + client.once("hello", (...args) => resolve(args)); + }), + ); + + parent.to("some-room").emit("hello", payload, 42); + + const messages = await Promise.all(received); + for (const args of messages) { + expect(args.length).to.be(recovery ? 3 : 2); + expect(args.slice(0, 2)).to.eql([payload, 42]); + if (recovery) { + expect(args[2]).to.be.a("string"); + } + } + if (recovery) { + expect(messages[0][2]).to.not.be(messages[1][2]); + } + } finally { + clients.forEach((client) => client.disconnect()); + await io.close(); + } + }); + } + } + it("should allow connections to dynamic namespaces with a function", (done) => { const io = new Server(0); const socket = createClient(io, "/dynamic-101"); From a8d606ca94a513dd2a02e614f34674b9f9ea465e Mon Sep 17 00:00:00 2001 From: xingsy97 <87063252+xingsy97@users.noreply.github.com> Date: Mon, 28 Sep 2026 16:11:31 +0800 Subject: [PATCH 2/2] resolve comment --- .../socket.io-adapter/lib/cluster-adapter.ts | 22 +++-- .../lib/in-memory-adapter.ts | 5 +- .../socket.io-adapter/test/cluster-adapter.ts | 80 ++++++++++++++++++- packages/socket.io-adapter/test/index.ts | 59 +++++++++++++- packages/socket.io/lib/parent-namespace.ts | 3 +- 5 files changed, 152 insertions(+), 17 deletions(-) diff --git a/packages/socket.io-adapter/lib/cluster-adapter.ts b/packages/socket.io-adapter/lib/cluster-adapter.ts index 170bf51b05..bb0ac6e2a4 100644 --- a/packages/socket.io-adapter/lib/cluster-adapter.ts +++ b/packages/socket.io-adapter/lib/cluster-adapter.ts @@ -251,10 +251,12 @@ export abstract class ClusterAdapter extends Adapter { }, ); } else { - const packet = message.data.packet; const opts = decodeOptions(message.data.opts); - - this.addOffsetIfNecessary(packet, opts, offset); + const packet = this.addOffsetIfNecessary( + message.data.packet, + opts, + offset, + ); super.broadcast(packet, opts); } @@ -436,7 +438,7 @@ export abstract class ClusterAdapter extends Adapter { opts: encodeOptions(opts), }, }); - this.addOffsetIfNecessary(packet, opts, offset); + packet = this.addOffsetIfNecessary(packet, opts, offset); } catch (e) { debug("[%s] error while broadcasting message: %s", this.uid, e.message); } @@ -446,8 +448,8 @@ export abstract class ClusterAdapter extends Adapter { } /** - * Adds an offset at the end of the data array in order to allow the client to receive any missed packets when it - * reconnects after a temporary disconnection. + * Returns a copy of the packet with an offset at the end of the data array, if necessary, in order to allow the client + * to receive any missed packets when it reconnects after a temporary disconnection. * * @param packet * @param opts @@ -460,7 +462,7 @@ export abstract class ClusterAdapter extends Adapter { offset: Offset, ) { if (!this.nsp.server.opts.connectionStateRecovery) { - return; + return packet; } const isEventPacket = packet.type === 2; // packets with acknowledgement are not stored because the acknowledgement function cannot be serialized and @@ -469,8 +471,12 @@ export abstract class ClusterAdapter extends Adapter { const notVolatile = opts.flags?.volatile === undefined; if (isEventPacket && withoutAcknowledgement && notVolatile) { - packet.data.push(offset); + return { + ...packet, + data: [...packet.data, offset], + }; } + return packet; } override broadcastWithAck( diff --git a/packages/socket.io-adapter/lib/in-memory-adapter.ts b/packages/socket.io-adapter/lib/in-memory-adapter.ts index cf178170e0..cf9eee7d47 100644 --- a/packages/socket.io-adapter/lib/in-memory-adapter.ts +++ b/packages/socket.io-adapter/lib/in-memory-adapter.ts @@ -486,7 +486,10 @@ export class SessionAwareAdapter extends Adapter { const id = yeast(); // the offset is stored at the end of the data array, so the client knows the ID of the last packet it has // processed (and the format is backward-compatible) - packet.data.push(id); + packet = { + ...packet, + data: [...packet.data, id], + }; this.packets.push({ id, opts, diff --git a/packages/socket.io-adapter/test/cluster-adapter.ts b/packages/socket.io-adapter/test/cluster-adapter.ts index bfd043c48d..4159c11165 100644 --- a/packages/socket.io-adapter/test/cluster-adapter.ts +++ b/packages/socket.io-adapter/test/cluster-adapter.ts @@ -22,8 +22,8 @@ class EventEmitterAdapter extends ClusterAdapterWithHeartbeat { readonly eventBus, ) { super(nsp, {}); - this.eventBus.on("message", (message) => { - this.onMessage(message as ClusterMessage); + this.eventBus.on("message", (message, offset) => { + this.onMessage(message as ClusterMessage, offset); }); } @@ -31,8 +31,9 @@ class EventEmitterAdapter extends ClusterAdapterWithHeartbeat { if (this.shouldFailPublish) { return Promise.reject(new Error("publish failed")); } - this.eventBus.emit("message", message); - return Promise.resolve(String(++this.offset)); + const offset = String(++this.offset); + this.eventBus.emit("message", message, offset); + return Promise.resolve(offset); } protected doPublishResponse( @@ -44,6 +45,77 @@ class EventEmitterAdapter extends ClusterAdapterWithHeartbeat { } } +describe("cluster adapter connection state recovery", () => { + for (const recovery of [false, true]) { + it(`should only copy packets when recovery is enabled (recovery: ${recovery})`, async () => { + const payload = { hello: "world" }; + const binary = Buffer.from("hello"); + const packet = { + nsp: "/", + type: 2, + data: ["hello", payload, binary], + }; + const encodedPackets: (typeof packet)[] = []; + const eventBus = new EventEmitter(); + const adapters = Array.from( + { length: 2 }, + () => + new EventEmitterAdapter( + { + name: "/", + server: { + encoder: { + encode(packet) { + encodedPackets.push(packet); + return []; + }, + }, + opts: { + connectionStateRecovery: recovery ? {} : undefined, + }, + }, + }, + eventBus, + ), + ); + + try { + for (let i = 0; i < 2; i++) { + await adapters[0].broadcast(packet, { + rooms: new Set(), + except: new Set(), + }); + } + + expect(packet.data).to.eql(["hello", payload, binary]); + expect(encodedPackets).to.have.length(4); + for (const encoded of encodedPackets) { + expect(encoded.data[1]).to.be(payload); + expect(encoded.data[2]).to.be(binary); + if (recovery) { + expect(encoded).to.not.be(packet); + expect(encoded.data).to.not.be(packet.data); + expect(encoded.data).to.have.length(4); + expect(encoded.data[3]).to.be.a("string"); + } else { + expect(encoded).to.be(packet); + expect(encoded.data).to.be(packet.data); + } + } + if (recovery) { + expect(encodedPackets[0].data[3]).to.be(encodedPackets[1].data[3]); + expect(encodedPackets[2].data[3]).to.be(encodedPackets[3].data[3]); + expect(encodedPackets[0].data[3]).to.not.be( + encodedPackets[2].data[3], + ); + } + } finally { + adapters.forEach((adapter) => adapter.close()); + } + }); + } +}); + describe("cluster adapter", () => { let servers: Server[], serverSockets: ServerSocket[], diff --git a/packages/socket.io-adapter/test/index.ts b/packages/socket.io-adapter/test/index.ts index 285ae66f02..87fc064d24 100644 --- a/packages/socket.io-adapter/test/index.ts +++ b/packages/socket.io-adapter/test/index.ts @@ -319,10 +319,12 @@ describe("socket.io-adapter", () => { describe("connection state recovery", () => { it("should persist and restore session", async () => { + let offset: string; const adapter = new SessionAwareAdapter({ server: { encoder: { encode(packet) { + offset = packet.data[1]; return packet; }, }, @@ -355,7 +357,7 @@ describe("socket.io-adapter", () => { }, ); - const offset = packetData[1]; + expect(packetData).to.eql(["hello"]); const session = await adapter.restoreSession("def", offset); expect(session).to.not.be(null); @@ -365,10 +367,14 @@ describe("socket.io-adapter", () => { }); it("should restore missed packets", async () => { + let offset: string; const adapter = new SessionAwareAdapter({ server: { encoder: { encode(packet) { + if (offset === undefined) { + offset = packet.data[1]; + } return packet; }, }, @@ -489,7 +495,7 @@ describe("socket.io-adapter", () => { }, ); - const offset = packetData[1]; + expect(packetData).to.eql(["hello"]); const session = await adapter.restoreSession("def", offset); expect(session).to.not.be(null); @@ -503,6 +509,55 @@ describe("socket.io-adapter", () => { expect(session.missedPackets[2][0]).to.eql("no except"); }); + it("should isolate recovery offsets across namespaces", () => { + const payload = { hello: "world" }; + const binary = Buffer.from("hello"); + const packet = { + nsp: "/", + type: 2, + data: ["hello", payload, binary], + }; + const encodedPackets: (typeof packet)[] = []; + + for (const name of ["/first", "/second"]) { + const adapter = new SessionAwareAdapter({ + name, + server: { + encoder: { + encode(packet) { + encodedPackets.push(packet); + return []; + }, + }, + opts: { + connectionStateRecovery: { + maxDisconnectionDuration: 5000, + }, + }, + }, + }); + adapter.broadcast(packet, { + rooms: new Set(), + except: new Set(), + }); + } + + expect(packet.data).to.eql(["hello", payload, binary]); + expect(encodedPackets).to.have.length(2); + for (const encoded of encodedPackets) { + expect(encoded.data[1]).to.be(payload); + expect(encoded.data[2]).to.be(binary); + expect(encoded).to.not.be(packet); + expect(encoded.data).to.not.be(packet.data); + expect(encoded.data).to.have.length(4); + expect(encoded.data[3]).to.be.a("string"); + } + expect(packet.nsp).to.be("/"); + expect(encodedPackets[0].nsp).to.be("/first"); + expect(encodedPackets[1].nsp).to.be("/second"); + expect(encodedPackets[0].data[3]).to.not.be(encodedPackets[1].data[3]); + }); + it("should fail to restore an unknown session", async () => { const adapter = new SessionAwareAdapter({ server: { diff --git a/packages/socket.io/lib/parent-namespace.ts b/packages/socket.io/lib/parent-namespace.ts index 5e2263b903..d3468bf46f 100644 --- a/packages/socket.io/lib/parent-namespace.ts +++ b/packages/socket.io/lib/parent-namespace.ts @@ -115,8 +115,7 @@ export class ParentNamespace< class ParentBroadcastAdapter extends Adapter { broadcast(packet: any, opts: BroadcastOptions) { this.nsp.children.forEach((nsp) => { - // Each child adapter may append its own recovery offset to the packet. - nsp.adapter.broadcast({ ...packet, data: [...packet.data] }, opts); + nsp.adapter.broadcast(packet, opts); }); } }