Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
22 changes: 14 additions & 8 deletions packages/socket.io-adapter/lib/cluster-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
}
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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(
Expand Down
5 changes: 4 additions & 1 deletion packages/socket.io-adapter/lib/in-memory-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
80 changes: 76 additions & 4 deletions packages/socket.io-adapter/test/cluster-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,17 +22,18 @@ 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);
});
}

protected doPublish(message: ClusterMessage): Promise<string> {
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(
Expand All @@ -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[],
Expand Down
59 changes: 57 additions & 2 deletions packages/socket.io-adapter/test/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
},
},
Expand Down Expand Up @@ -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);
Expand All @@ -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;
},
},
Expand Down Expand Up @@ -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);
Expand All @@ -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: {
Expand Down
66 changes: 65 additions & 1 deletion packages/socket.io/test/connection-state-recovery.ts
Original file line number Diff line number Diff line change
@@ -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";

Expand Down Expand Up @@ -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, {
Expand Down
Loading