Skip to content

Commit 7a33d9b

Browse files
fix(cluster-engine): properly handle upgrade failures
1 parent 8f4b825 commit 7a33d9b

3 files changed

Lines changed: 112 additions & 1 deletion

File tree

packages/socket.io-cluster-engine/docs/sequence_diagrams.md

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
* [HTTP long-polling (read)](#http-long-polling-read)
66
* [HTTP long-polling (write)](#http-long-polling-write)
77
* [WebSocket upgrade](#websocket-upgrade)
8+
* [WebSocket upgrade failure](#websocket-upgrade-failure)
89
* [Message types](#message-types)
910
<!-- TOC -->
1011

@@ -124,6 +125,48 @@ sequenceDiagram
124125
end
125126
```
126127

128+
### WebSocket upgrade failure
129+
130+
```mermaid
131+
sequenceDiagram
132+
autonumber
133+
134+
participant C as Client
135+
participant W1 as Worker A
136+
participant B as Cluster bus
137+
participant W2 as Worker B
138+
139+
Note over W1: Client is connected with HTTP long-polling on Worker A
140+
141+
C->>W2: HTTP Upgrade WebSocket<br/>sid=abc
142+
Note over W2: Session ID unknown locally
143+
144+
W2->>B: ACQUIRE_LOCK<br/>sid=abc, transport=websocket, type=read
145+
B->>W1: ACQUIRE_LOCK
146+
Note over W1: Lockable if current transport is polling<br/>and not already upgrading/upgraded
147+
148+
W1->>B: ACQUIRE_LOCK_RESPONSE<br/>success=true
149+
B->>W2: ACQUIRE_LOCK_RESPONSE
150+
151+
Note over W2: Accept WebSocket upgrade and start upgrade probe
152+
153+
C->>W2: ping "probe"
154+
W2->>C: pong "probe"
155+
156+
Note over W2: Probe failure (timeout or invalid packet)
157+
158+
W2->>B: UPGRADE<br/>sid=abc, success=false
159+
B->>W1: UPGRADE
160+
161+
alt connection was still delayed
162+
Note over W1: Emit "connection" event
163+
else connection was already emitted
164+
Note over W1: Stay on HTTP long-polling
165+
end
166+
167+
Note over W2: Close WebSocket transport
168+
```
169+
127170
## Message types
128171

129172
| Message type | Direction | Purpose |

packages/socket.io-cluster-engine/lib/engine.ts

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ const kDelayedTimer = Symbol("delayedTimer");
1818
const kBuffer = Symbol("buffer");
1919
const kPacketListener = Symbol("packetListener");
2020
const kNoopTimer = Symbol("noopTimer");
21+
const kUpgradeTimer = Symbol("upgradeTimer");
2122
const kSenderId = Symbol("senderId");
2223

2324
type Brand<K, T> = K & { __brand: T };
@@ -227,11 +228,25 @@ export abstract class ClusterEngine extends Server {
227228
case "websocket":
228229
case "webtransport": {
229230
client.upgrading = true;
231+
230232
client[kNoopTimer] = setTimeout(() => {
231233
debug("writing a noop packet to polling for fast upgrade");
232234
// @ts-expect-error sendPacket() is private
233235
client.sendPacket("noop");
234236
}, this._opts.noopUpgradeInterval);
237+
238+
client[kUpgradeTimer] = setTimeout(() => {
239+
if (client.upgrading) {
240+
debug("upgrade did not complete, resetting upgrade state");
241+
client.upgrading = false;
242+
clearTimeout(client[kNoopTimer]);
243+
if (client[kDelayed]) {
244+
this._doConnect(client);
245+
}
246+
}
247+
}, this.opts.upgradeTimeout);
248+
249+
break;
235250
}
236251
}
237252
break;
@@ -288,6 +303,7 @@ export abstract class ClusterEngine extends Server {
288303
}
289304

290305
clearTimeout(client[kNoopTimer]);
306+
clearTimeout(client[kUpgradeTimer]);
291307
client.upgrading = false;
292308

293309
if (message.data.success) {
@@ -323,6 +339,8 @@ export abstract class ClusterEngine extends Server {
323339
},
324340
});
325341
}
342+
} else if (client[kDelayed]) {
343+
this._doConnect(client);
326344
}
327345
break;
328346
}
@@ -578,6 +596,16 @@ export abstract class ClusterEngine extends Server {
578596
() => this._onUpgradeSuccess(sid, transport, req, senderId),
579597
() => {
580598
debug("upgrade failure");
599+
this.publishMessage({
600+
requestId: ++this._requestCount as RequestId,
601+
senderId: this._nodeId,
602+
recipientId: senderId,
603+
type: MessageType.UPGRADE,
604+
data: {
605+
sid,
606+
success: false,
607+
},
608+
});
581609
},
582610
);
583611
}

packages/socket.io-cluster-engine/test/in-memory.test.ts

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,9 @@ describe("in-memory", () => {
4747
const port1 = (httpServer1.address() as AddressInfo).port;
4848

4949
httpServer2 = createServer();
50-
engine2 = new InMemoryEngine(eventBus);
50+
engine2 = new InMemoryEngine(eventBus, {
51+
upgradeTimeout: 100,
52+
});
5153
engine2.attach(httpServer2);
5254
httpServer2.listen(0);
5355
const port2 = (httpServer2.address() as AddressInfo).port;
@@ -430,4 +432,42 @@ describe("in-memory", () => {
430432

431433
await promise;
432434
});
435+
436+
it("should resume after upgrade failure", async () => {
437+
const sid = await handshake(ports[2]);
438+
439+
const socket = new WebSocket(
440+
`ws://localhost:${ports[1]}/engine.io/?EIO=4&transport=websocket&sid=${sid}`,
441+
);
442+
443+
await new Promise<void>((resolve, reject) => {
444+
engine3.on("connection", (socket) => {
445+
if (socket.upgrading) {
446+
return reject("should not be upgrading");
447+
}
448+
socket.on("message", (data: string) => {
449+
assert.equal(data, "ping");
450+
socket.send("pong");
451+
});
452+
resolve();
453+
});
454+
455+
// don't send the probe pong, so the upgrade should time out on engine 1
456+
socket.onopen = () => {};
457+
});
458+
459+
{
460+
const res = await fetch(url(ports[0], sid), {
461+
method: "POST",
462+
body: "4ping",
463+
});
464+
assert.equal(res.status, 200);
465+
}
466+
467+
{
468+
const res = await fetch(url(ports[1], sid));
469+
assert.equal(res.status, 200);
470+
assert.equal(await res.text(), "2\x1e4pong");
471+
}
472+
});
433473
});

0 commit comments

Comments
 (0)