1169 lines
50 KiB
Diff
1169 lines
50 KiB
Diff
diff --git a/dist/defaultWorker.js b/dist/defaultWorker.js
|
|
index 306d691077f81e96bf977ed13aa1a962e13e5d0c..a4e2a21c12fd2ad2888aec00bd7c66406bc51104 100644
|
|
--- a/dist/defaultWorker.js
|
|
+++ b/dist/defaultWorker.js
|
|
@@ -366,6 +366,8 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
__name(this, "WebSocketShard");
|
|
}
|
|
connection = null;
|
|
+ connectionEpoch = 0;
|
|
+ destroyPromise = null;
|
|
useIdentifyCompress = false;
|
|
inflate = null;
|
|
textDecoder = new import_node_util.TextDecoder();
|
|
@@ -415,6 +417,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
if (this.#status !== 0 /* Idle */) {
|
|
throw new Error("Tried to connect a shard that wasn't idle");
|
|
}
|
|
+ const epoch = this.connectionEpoch;
|
|
const { version, encoding, compression } = this.strategy.options;
|
|
const params = new import_node_url.URLSearchParams({ v: version, encoding });
|
|
if (compression) {
|
|
@@ -433,6 +436,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
const url = `${session?.resumeURL ?? this.strategy.options.gatewayInformation.url}?${params.toString()}`;
|
|
this.debug([`Connecting to ${url}`]);
|
|
const connection = new WebSocketConstructor(url, [], {
|
|
@@ -440,21 +444,25 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
});
|
|
connection.binaryType = "arraybuffer";
|
|
connection.onmessage = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onMessage(event.data, event.data instanceof ArrayBuffer);
|
|
};
|
|
connection.onerror = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.onError(event.error);
|
|
};
|
|
connection.onclose = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onClose(event.code);
|
|
};
|
|
connection.onopen = () => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.sendRateLimitState = getInitialSendRateLimitState();
|
|
};
|
|
this.connection = connection;
|
|
this.#status = 1 /* Connecting */;
|
|
const { ok } = await this.waitForEvent("hello" /* Hello */, this.strategy.options.helloTimeout);
|
|
- if (!ok) {
|
|
+ if (!ok || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
if (session?.shardCount === this.strategy.options.shardCount) {
|
|
@@ -464,6 +472,26 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
}
|
|
async destroy(options = {}) {
|
|
+ // Invalidate pending setup and event waits before cancellation can resume them.
|
|
+ const epoch = ++this.connectionEpoch;
|
|
+ if (this.destroyPromise) {
|
|
+ await this.destroyPromise;
|
|
+ return;
|
|
+ }
|
|
+ if (this.#status === 0 /* Idle */) return;
|
|
+ const cleanup = this.destroyConnection(options);
|
|
+ this.destroyPromise = cleanup;
|
|
+ try {
|
|
+ await cleanup;
|
|
+ } finally {
|
|
+ if (this.destroyPromise === cleanup) this.destroyPromise = null;
|
|
+ }
|
|
+ if (options.recover !== void 0 && epoch === this.connectionEpoch) {
|
|
+ await (0, import_promises2.setTimeout)(500);
|
|
+ if (epoch === this.connectionEpoch) return this.internalConnect();
|
|
+ }
|
|
+ }
|
|
+ async destroyConnection(options = {}) {
|
|
if (this.#status === 0 /* Idle */) {
|
|
this.debug(["Tried to destroy a shard that was idle"]);
|
|
return;
|
|
@@ -494,36 +522,27 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
if (options.recover !== 1 /* Resume */) {
|
|
await this.strategy.updateSessionInfo(this.id, null);
|
|
}
|
|
- if (this.connection) {
|
|
- this.connection.onmessage = null;
|
|
- this.connection.onclose = null;
|
|
- const shouldClose = this.connection.readyState === import_ws.WebSocket.OPEN;
|
|
- this.debug([
|
|
- "Connection status during destroy",
|
|
- `Needs closing: ${shouldClose}`,
|
|
- `Ready state: ${this.connection.readyState}`
|
|
- ]);
|
|
- if (shouldClose) {
|
|
- let outerResolve;
|
|
- const promise = new Promise((resolve2) => {
|
|
- outerResolve = resolve2;
|
|
- });
|
|
- this.connection.onclose = outerResolve;
|
|
- this.connection.close(options.code, options.reason);
|
|
- await promise;
|
|
+ const connection = this.connection;
|
|
+ if (connection) {
|
|
+ connection.onmessage = null;
|
|
+ connection.onopen = null;
|
|
+ connection.onclose = null;
|
|
+ // A connecting ws can still emit its handshake error. Keep its epoch-bound
|
|
+ // error handler until terminate/close has actually completed.
|
|
+ if (connection.readyState !== 3 /* CLOSED */) {
|
|
+ const closed = new Promise((resolve) => { connection.onclose = resolve; });
|
|
+ if (connection.readyState === 0 /* CONNECTING */) connection.terminate();
|
|
+ else if (connection.readyState === 1 /* OPEN */) connection.close(options.code, options.reason);
|
|
+ await closed;
|
|
this.emit("closed" /* Closed */, { code: options.code });
|
|
}
|
|
- this.connection.onerror = null;
|
|
- } else {
|
|
- this.debug(["Destroying a shard that has no connection; please open an issue on GitHub"]);
|
|
+ connection.onerror = null;
|
|
+ if (this.connection === connection) this.connection = null;
|
|
}
|
|
this.#status = 0 /* Idle */;
|
|
- if (options.recover !== void 0) {
|
|
- await (0, import_promises2.setTimeout)(500);
|
|
- return this.internalConnect();
|
|
- }
|
|
}
|
|
async waitForEvent(event, timeoutDuration) {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug([`Waiting for event ${event} ${timeoutDuration ? `for ${timeoutDuration}ms` : "indefinitely"}`]);
|
|
const timeoutController = new AbortController();
|
|
const timeout = timeoutDuration ? (0, import_node_timers.setTimeout)(() => timeoutController.abort(), timeoutDuration).unref() : null;
|
|
@@ -536,6 +555,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
]);
|
|
return { ok: !closed };
|
|
} catch {
|
|
+ if (epoch !== this.connectionEpoch) return { ok: false };
|
|
void this.destroy({
|
|
code: 1e3 /* Normal */,
|
|
reason: "Something timed out or went wrong while waiting for an event",
|
|
@@ -546,7 +566,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
if (timeout) {
|
|
(0, import_node_timers.clearTimeout)(timeout);
|
|
}
|
|
- this.timeoutAbortControllers.delete(event);
|
|
+ if (this.timeoutAbortControllers.get(event) === timeoutController) this.timeoutAbortControllers.delete(event);
|
|
if (!closeController.signal.aborted) {
|
|
closeController.abort();
|
|
}
|
|
@@ -593,6 +613,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
this.connection.send(JSON.stringify(payload));
|
|
}
|
|
async identify() {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug(["Waiting for identify throttle"]);
|
|
const controller = new AbortController();
|
|
const closeHandler = /* @__PURE__ */ __name(() => {
|
|
@@ -602,7 +623,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
try {
|
|
await this.strategy.waitForIdentify(this.id, controller.signal);
|
|
} catch {
|
|
- if (controller.signal.aborted) {
|
|
+ if (controller.signal.aborted || epoch !== this.connectionEpoch) {
|
|
this.debug(["Was waiting for an identify, but the shard closed in the meantime"]);
|
|
return;
|
|
}
|
|
@@ -618,6 +639,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
} finally {
|
|
this.off("closed" /* Closed */, closeHandler);
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([
|
|
"Identifying",
|
|
`shard id: ${this.id.toString()}`,
|
|
@@ -642,6 +664,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
op: import_v102.GatewayOpcodes.Identify,
|
|
d
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.waitForEvent("ready" /* Ready */, this.strategy.options.readyTimeout);
|
|
}
|
|
async resume(session) {
|
|
@@ -663,14 +686,17 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
});
|
|
}
|
|
async heartbeat(requested = false) {
|
|
+ const epoch = this.connectionEpoch;
|
|
if (!this.isAck && !requested) {
|
|
return this.destroy({ reason: "Zombie connection", recover: 1 /* Resume */ });
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.send({
|
|
op: import_v102.GatewayOpcodes.Heartbeat,
|
|
d: session?.sequence ?? null
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.lastHeartbeatAt = Date.now();
|
|
this.isAck = false;
|
|
}
|
|
@@ -722,8 +748,9 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
return null;
|
|
}
|
|
async onMessage(data, isBinary) {
|
|
+ const epoch = this.connectionEpoch;
|
|
const payload = await this.unpackMessage(data, isBinary);
|
|
- if (!payload) {
|
|
+ if (!payload || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
switch (payload.op) {
|
|
@@ -733,7 +760,6 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
switch (payload.t) {
|
|
case import_v102.GatewayDispatchEvents.Ready: {
|
|
- this.#status = 3 /* Ready */;
|
|
const session2 = {
|
|
sequence: payload.s,
|
|
sessionId: payload.d.session_id,
|
|
@@ -742,6 +768,8 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
resumeURL: payload.d.resume_gateway_url
|
|
};
|
|
await this.strategy.updateSessionInfo(this.id, session2);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
+ this.#status = 3 /* Ready */;
|
|
this.emit("ready" /* Ready */, { data: payload.d });
|
|
break;
|
|
}
|
|
@@ -756,9 +784,11 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (session) {
|
|
if (payload.s > session.sequence) {
|
|
await this.strategy.updateSessionInfo(this.id, { ...session, sequence: payload.s });
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
}
|
|
} else {
|
|
this.debug([
|
|
@@ -782,6 +812,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
case import_v102.GatewayOpcodes.InvalidSession: {
|
|
this.debug([`Invalid session; will attempt to resume: ${payload.d.toString()}`]);
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (payload.d && session) {
|
|
await this.resume(session);
|
|
} else {
|
|
@@ -797,17 +828,19 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
const jitter = Math.random();
|
|
const firstWait = Math.floor(payload.d.heartbeat_interval * jitter);
|
|
this.debug([`Preparing first heartbeat of the connection with a jitter of ${jitter}; waiting ${firstWait}ms`]);
|
|
+ const controller = new AbortController();
|
|
try {
|
|
- const controller = new AbortController();
|
|
this.initialHeartbeatTimeoutController = controller;
|
|
await (0, import_promises2.setTimeout)(firstWait, void 0, { signal: controller.signal });
|
|
} catch {
|
|
this.debug(["Cancelled initial heartbeat due to #destroy being called"]);
|
|
return;
|
|
} finally {
|
|
- this.initialHeartbeatTimeoutController = null;
|
|
+ if (this.initialHeartbeatTimeoutController === controller) this.initialHeartbeatTimeoutController = null;
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.heartbeat();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([`First heartbeat sent, starting to beat every ${payload.d.heartbeat_interval}ms`]);
|
|
this.heartbeatInterval = (0, import_node_timers.setInterval)(() => void this.heartbeat(), payload.d.heartbeat_interval);
|
|
break;
|
|
diff --git a/dist/defaultWorker.mjs b/dist/defaultWorker.mjs
|
|
index 1302a02537062b489b0be156eccb60a960f015c5..b8f1aff44c28f5deede48fca0287a6aa860e2e8a 100644
|
|
--- a/dist/defaultWorker.mjs
|
|
+++ b/dist/defaultWorker.mjs
|
|
@@ -348,6 +348,8 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
__name(this, "WebSocketShard");
|
|
}
|
|
connection = null;
|
|
+ connectionEpoch = 0;
|
|
+ destroyPromise = null;
|
|
useIdentifyCompress = false;
|
|
inflate = null;
|
|
textDecoder = new TextDecoder();
|
|
@@ -397,6 +399,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
if (this.#status !== 0 /* Idle */) {
|
|
throw new Error("Tried to connect a shard that wasn't idle");
|
|
}
|
|
+ const epoch = this.connectionEpoch;
|
|
const { version, encoding, compression } = this.strategy.options;
|
|
const params = new URLSearchParams({ v: version, encoding });
|
|
if (compression) {
|
|
@@ -415,6 +418,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
const url = `${session?.resumeURL ?? this.strategy.options.gatewayInformation.url}?${params.toString()}`;
|
|
this.debug([`Connecting to ${url}`]);
|
|
const connection = new WebSocketConstructor(url, [], {
|
|
@@ -422,21 +426,25 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
});
|
|
connection.binaryType = "arraybuffer";
|
|
connection.onmessage = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onMessage(event.data, event.data instanceof ArrayBuffer);
|
|
};
|
|
connection.onerror = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.onError(event.error);
|
|
};
|
|
connection.onclose = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onClose(event.code);
|
|
};
|
|
connection.onopen = () => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.sendRateLimitState = getInitialSendRateLimitState();
|
|
};
|
|
this.connection = connection;
|
|
this.#status = 1 /* Connecting */;
|
|
const { ok } = await this.waitForEvent("hello" /* Hello */, this.strategy.options.helloTimeout);
|
|
- if (!ok) {
|
|
+ if (!ok || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
if (session?.shardCount === this.strategy.options.shardCount) {
|
|
@@ -446,6 +454,26 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
}
|
|
async destroy(options = {}) {
|
|
+ // Invalidate pending setup and event waits before cancellation can resume them.
|
|
+ const epoch = ++this.connectionEpoch;
|
|
+ if (this.destroyPromise) {
|
|
+ await this.destroyPromise;
|
|
+ return;
|
|
+ }
|
|
+ if (this.#status === 0 /* Idle */) return;
|
|
+ const cleanup = this.destroyConnection(options);
|
|
+ this.destroyPromise = cleanup;
|
|
+ try {
|
|
+ await cleanup;
|
|
+ } finally {
|
|
+ if (this.destroyPromise === cleanup) this.destroyPromise = null;
|
|
+ }
|
|
+ if (options.recover !== void 0 && epoch === this.connectionEpoch) {
|
|
+ await sleep2(500);
|
|
+ if (epoch === this.connectionEpoch) return this.internalConnect();
|
|
+ }
|
|
+ }
|
|
+ async destroyConnection(options = {}) {
|
|
if (this.#status === 0 /* Idle */) {
|
|
this.debug(["Tried to destroy a shard that was idle"]);
|
|
return;
|
|
@@ -476,36 +504,27 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
if (options.recover !== 1 /* Resume */) {
|
|
await this.strategy.updateSessionInfo(this.id, null);
|
|
}
|
|
- if (this.connection) {
|
|
- this.connection.onmessage = null;
|
|
- this.connection.onclose = null;
|
|
- const shouldClose = this.connection.readyState === WebSocket.OPEN;
|
|
- this.debug([
|
|
- "Connection status during destroy",
|
|
- `Needs closing: ${shouldClose}`,
|
|
- `Ready state: ${this.connection.readyState}`
|
|
- ]);
|
|
- if (shouldClose) {
|
|
- let outerResolve;
|
|
- const promise = new Promise((resolve2) => {
|
|
- outerResolve = resolve2;
|
|
- });
|
|
- this.connection.onclose = outerResolve;
|
|
- this.connection.close(options.code, options.reason);
|
|
- await promise;
|
|
+ const connection = this.connection;
|
|
+ if (connection) {
|
|
+ connection.onmessage = null;
|
|
+ connection.onopen = null;
|
|
+ connection.onclose = null;
|
|
+ // A connecting ws can still emit its handshake error. Keep its epoch-bound
|
|
+ // error handler until terminate/close has actually completed.
|
|
+ if (connection.readyState !== 3 /* CLOSED */) {
|
|
+ const closed = new Promise((resolve) => { connection.onclose = resolve; });
|
|
+ if (connection.readyState === 0 /* CONNECTING */) connection.terminate();
|
|
+ else if (connection.readyState === 1 /* OPEN */) connection.close(options.code, options.reason);
|
|
+ await closed;
|
|
this.emit("closed" /* Closed */, { code: options.code });
|
|
}
|
|
- this.connection.onerror = null;
|
|
- } else {
|
|
- this.debug(["Destroying a shard that has no connection; please open an issue on GitHub"]);
|
|
+ connection.onerror = null;
|
|
+ if (this.connection === connection) this.connection = null;
|
|
}
|
|
this.#status = 0 /* Idle */;
|
|
- if (options.recover !== void 0) {
|
|
- await sleep2(500);
|
|
- return this.internalConnect();
|
|
- }
|
|
}
|
|
async waitForEvent(event, timeoutDuration) {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug([`Waiting for event ${event} ${timeoutDuration ? `for ${timeoutDuration}ms` : "indefinitely"}`]);
|
|
const timeoutController = new AbortController();
|
|
const timeout = timeoutDuration ? setTimeout(() => timeoutController.abort(), timeoutDuration).unref() : null;
|
|
@@ -518,6 +537,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
]);
|
|
return { ok: !closed };
|
|
} catch {
|
|
+ if (epoch !== this.connectionEpoch) return { ok: false };
|
|
void this.destroy({
|
|
code: 1e3 /* Normal */,
|
|
reason: "Something timed out or went wrong while waiting for an event",
|
|
@@ -528,7 +548,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
if (timeout) {
|
|
clearTimeout(timeout);
|
|
}
|
|
- this.timeoutAbortControllers.delete(event);
|
|
+ if (this.timeoutAbortControllers.get(event) === timeoutController) this.timeoutAbortControllers.delete(event);
|
|
if (!closeController.signal.aborted) {
|
|
closeController.abort();
|
|
}
|
|
@@ -575,6 +595,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
this.connection.send(JSON.stringify(payload));
|
|
}
|
|
async identify() {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug(["Waiting for identify throttle"]);
|
|
const controller = new AbortController();
|
|
const closeHandler = /* @__PURE__ */ __name(() => {
|
|
@@ -584,7 +605,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
try {
|
|
await this.strategy.waitForIdentify(this.id, controller.signal);
|
|
} catch {
|
|
- if (controller.signal.aborted) {
|
|
+ if (controller.signal.aborted || epoch !== this.connectionEpoch) {
|
|
this.debug(["Was waiting for an identify, but the shard closed in the meantime"]);
|
|
return;
|
|
}
|
|
@@ -600,6 +621,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
} finally {
|
|
this.off("closed" /* Closed */, closeHandler);
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([
|
|
"Identifying",
|
|
`shard id: ${this.id.toString()}`,
|
|
@@ -624,6 +646,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
op: GatewayOpcodes2.Identify,
|
|
d
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.waitForEvent("ready" /* Ready */, this.strategy.options.readyTimeout);
|
|
}
|
|
async resume(session) {
|
|
@@ -645,14 +668,17 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
});
|
|
}
|
|
async heartbeat(requested = false) {
|
|
+ const epoch = this.connectionEpoch;
|
|
if (!this.isAck && !requested) {
|
|
return this.destroy({ reason: "Zombie connection", recover: 1 /* Resume */ });
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.send({
|
|
op: GatewayOpcodes2.Heartbeat,
|
|
d: session?.sequence ?? null
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.lastHeartbeatAt = Date.now();
|
|
this.isAck = false;
|
|
}
|
|
@@ -704,8 +730,9 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
return null;
|
|
}
|
|
async onMessage(data, isBinary) {
|
|
+ const epoch = this.connectionEpoch;
|
|
const payload = await this.unpackMessage(data, isBinary);
|
|
- if (!payload) {
|
|
+ if (!payload || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
switch (payload.op) {
|
|
@@ -715,7 +742,6 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
switch (payload.t) {
|
|
case GatewayDispatchEvents.Ready: {
|
|
- this.#status = 3 /* Ready */;
|
|
const session2 = {
|
|
sequence: payload.s,
|
|
sessionId: payload.d.session_id,
|
|
@@ -724,6 +750,8 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
resumeURL: payload.d.resume_gateway_url
|
|
};
|
|
await this.strategy.updateSessionInfo(this.id, session2);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
+ this.#status = 3 /* Ready */;
|
|
this.emit("ready" /* Ready */, { data: payload.d });
|
|
break;
|
|
}
|
|
@@ -738,9 +766,11 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (session) {
|
|
if (payload.s > session.sequence) {
|
|
await this.strategy.updateSessionInfo(this.id, { ...session, sequence: payload.s });
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
}
|
|
} else {
|
|
this.debug([
|
|
@@ -764,6 +794,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
case GatewayOpcodes2.InvalidSession: {
|
|
this.debug([`Invalid session; will attempt to resume: ${payload.d.toString()}`]);
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (payload.d && session) {
|
|
await this.resume(session);
|
|
} else {
|
|
@@ -779,17 +810,19 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
const jitter = Math.random();
|
|
const firstWait = Math.floor(payload.d.heartbeat_interval * jitter);
|
|
this.debug([`Preparing first heartbeat of the connection with a jitter of ${jitter}; waiting ${firstWait}ms`]);
|
|
+ const controller = new AbortController();
|
|
try {
|
|
- const controller = new AbortController();
|
|
this.initialHeartbeatTimeoutController = controller;
|
|
await sleep2(firstWait, void 0, { signal: controller.signal });
|
|
} catch {
|
|
this.debug(["Cancelled initial heartbeat due to #destroy being called"]);
|
|
return;
|
|
} finally {
|
|
- this.initialHeartbeatTimeoutController = null;
|
|
+ if (this.initialHeartbeatTimeoutController === controller) this.initialHeartbeatTimeoutController = null;
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.heartbeat();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([`First heartbeat sent, starting to beat every ${payload.d.heartbeat_interval}ms`]);
|
|
this.heartbeatInterval = setInterval(() => void this.heartbeat(), payload.d.heartbeat_interval);
|
|
break;
|
|
diff --git a/dist/index.js b/dist/index.js
|
|
index 106b2ec450a248cd5bd7c36c038215145438eade..80a09fa2a8c953187e26d22e2e4a87bdb0893133 100644
|
|
--- a/dist/index.js
|
|
+++ b/dist/index.js
|
|
@@ -605,6 +605,8 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
__name(this, "WebSocketShard");
|
|
}
|
|
connection = null;
|
|
+ connectionEpoch = 0;
|
|
+ destroyPromise = null;
|
|
useIdentifyCompress = false;
|
|
inflate = null;
|
|
textDecoder = new import_node_util.TextDecoder();
|
|
@@ -654,6 +656,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
if (this.#status !== 0 /* Idle */) {
|
|
throw new Error("Tried to connect a shard that wasn't idle");
|
|
}
|
|
+ const epoch = this.connectionEpoch;
|
|
const { version: version2, encoding, compression } = this.strategy.options;
|
|
const params = new import_node_url.URLSearchParams({ v: version2, encoding });
|
|
if (compression) {
|
|
@@ -672,6 +675,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
const url = `${session?.resumeURL ?? this.strategy.options.gatewayInformation.url}?${params.toString()}`;
|
|
this.debug([`Connecting to ${url}`]);
|
|
const connection = new WebSocketConstructor(url, [], {
|
|
@@ -679,21 +683,25 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
});
|
|
connection.binaryType = "arraybuffer";
|
|
connection.onmessage = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onMessage(event.data, event.data instanceof ArrayBuffer);
|
|
};
|
|
connection.onerror = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.onError(event.error);
|
|
};
|
|
connection.onclose = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onClose(event.code);
|
|
};
|
|
connection.onopen = () => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.sendRateLimitState = getInitialSendRateLimitState();
|
|
};
|
|
this.connection = connection;
|
|
this.#status = 1 /* Connecting */;
|
|
const { ok } = await this.waitForEvent("hello" /* Hello */, this.strategy.options.helloTimeout);
|
|
- if (!ok) {
|
|
+ if (!ok || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
if (session?.shardCount === this.strategy.options.shardCount) {
|
|
@@ -703,6 +711,26 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
}
|
|
async destroy(options = {}) {
|
|
+ // Invalidate pending setup and event waits before cancellation can resume them.
|
|
+ const epoch = ++this.connectionEpoch;
|
|
+ if (this.destroyPromise) {
|
|
+ await this.destroyPromise;
|
|
+ return;
|
|
+ }
|
|
+ if (this.#status === 0 /* Idle */) return;
|
|
+ const cleanup = this.destroyConnection(options);
|
|
+ this.destroyPromise = cleanup;
|
|
+ try {
|
|
+ await cleanup;
|
|
+ } finally {
|
|
+ if (this.destroyPromise === cleanup) this.destroyPromise = null;
|
|
+ }
|
|
+ if (options.recover !== void 0 && epoch === this.connectionEpoch) {
|
|
+ await (0, import_promises2.setTimeout)(500);
|
|
+ if (epoch === this.connectionEpoch) return this.internalConnect();
|
|
+ }
|
|
+ }
|
|
+ async destroyConnection(options = {}) {
|
|
if (this.#status === 0 /* Idle */) {
|
|
this.debug(["Tried to destroy a shard that was idle"]);
|
|
return;
|
|
@@ -733,36 +761,27 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
if (options.recover !== 1 /* Resume */) {
|
|
await this.strategy.updateSessionInfo(this.id, null);
|
|
}
|
|
- if (this.connection) {
|
|
- this.connection.onmessage = null;
|
|
- this.connection.onclose = null;
|
|
- const shouldClose = this.connection.readyState === import_ws.WebSocket.OPEN;
|
|
- this.debug([
|
|
- "Connection status during destroy",
|
|
- `Needs closing: ${shouldClose}`,
|
|
- `Ready state: ${this.connection.readyState}`
|
|
- ]);
|
|
- if (shouldClose) {
|
|
- let outerResolve;
|
|
- const promise = new Promise((resolve2) => {
|
|
- outerResolve = resolve2;
|
|
- });
|
|
- this.connection.onclose = outerResolve;
|
|
- this.connection.close(options.code, options.reason);
|
|
- await promise;
|
|
+ const connection = this.connection;
|
|
+ if (connection) {
|
|
+ connection.onmessage = null;
|
|
+ connection.onopen = null;
|
|
+ connection.onclose = null;
|
|
+ // A connecting ws can still emit its handshake error. Keep its epoch-bound
|
|
+ // error handler until terminate/close has actually completed.
|
|
+ if (connection.readyState !== 3 /* CLOSED */) {
|
|
+ const closed = new Promise((resolve) => { connection.onclose = resolve; });
|
|
+ if (connection.readyState === 0 /* CONNECTING */) connection.terminate();
|
|
+ else if (connection.readyState === 1 /* OPEN */) connection.close(options.code, options.reason);
|
|
+ await closed;
|
|
this.emit("closed" /* Closed */, { code: options.code });
|
|
}
|
|
- this.connection.onerror = null;
|
|
- } else {
|
|
- this.debug(["Destroying a shard that has no connection; please open an issue on GitHub"]);
|
|
+ connection.onerror = null;
|
|
+ if (this.connection === connection) this.connection = null;
|
|
}
|
|
this.#status = 0 /* Idle */;
|
|
- if (options.recover !== void 0) {
|
|
- await (0, import_promises2.setTimeout)(500);
|
|
- return this.internalConnect();
|
|
- }
|
|
}
|
|
async waitForEvent(event, timeoutDuration) {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug([`Waiting for event ${event} ${timeoutDuration ? `for ${timeoutDuration}ms` : "indefinitely"}`]);
|
|
const timeoutController = new AbortController();
|
|
const timeout = timeoutDuration ? (0, import_node_timers.setTimeout)(() => timeoutController.abort(), timeoutDuration).unref() : null;
|
|
@@ -775,6 +794,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
]);
|
|
return { ok: !closed };
|
|
} catch {
|
|
+ if (epoch !== this.connectionEpoch) return { ok: false };
|
|
void this.destroy({
|
|
code: 1e3 /* Normal */,
|
|
reason: "Something timed out or went wrong while waiting for an event",
|
|
@@ -785,7 +805,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
if (timeout) {
|
|
(0, import_node_timers.clearTimeout)(timeout);
|
|
}
|
|
- this.timeoutAbortControllers.delete(event);
|
|
+ if (this.timeoutAbortControllers.get(event) === timeoutController) this.timeoutAbortControllers.delete(event);
|
|
if (!closeController.signal.aborted) {
|
|
closeController.abort();
|
|
}
|
|
@@ -832,6 +852,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
this.connection.send(JSON.stringify(payload));
|
|
}
|
|
async identify() {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug(["Waiting for identify throttle"]);
|
|
const controller = new AbortController();
|
|
const closeHandler = /* @__PURE__ */ __name(() => {
|
|
@@ -841,7 +862,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
try {
|
|
await this.strategy.waitForIdentify(this.id, controller.signal);
|
|
} catch {
|
|
- if (controller.signal.aborted) {
|
|
+ if (controller.signal.aborted || epoch !== this.connectionEpoch) {
|
|
this.debug(["Was waiting for an identify, but the shard closed in the meantime"]);
|
|
return;
|
|
}
|
|
@@ -857,6 +878,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
} finally {
|
|
this.off("closed" /* Closed */, closeHandler);
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([
|
|
"Identifying",
|
|
`shard id: ${this.id.toString()}`,
|
|
@@ -881,6 +903,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
op: import_v102.GatewayOpcodes.Identify,
|
|
d
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.waitForEvent("ready" /* Ready */, this.strategy.options.readyTimeout);
|
|
}
|
|
async resume(session) {
|
|
@@ -902,14 +925,17 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
});
|
|
}
|
|
async heartbeat(requested = false) {
|
|
+ const epoch = this.connectionEpoch;
|
|
if (!this.isAck && !requested) {
|
|
return this.destroy({ reason: "Zombie connection", recover: 1 /* Resume */ });
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.send({
|
|
op: import_v102.GatewayOpcodes.Heartbeat,
|
|
d: session?.sequence ?? null
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.lastHeartbeatAt = Date.now();
|
|
this.isAck = false;
|
|
}
|
|
@@ -961,8 +987,9 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
return null;
|
|
}
|
|
async onMessage(data, isBinary) {
|
|
+ const epoch = this.connectionEpoch;
|
|
const payload = await this.unpackMessage(data, isBinary);
|
|
- if (!payload) {
|
|
+ if (!payload || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
switch (payload.op) {
|
|
@@ -972,7 +999,6 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
switch (payload.t) {
|
|
case import_v102.GatewayDispatchEvents.Ready: {
|
|
- this.#status = 3 /* Ready */;
|
|
const session2 = {
|
|
sequence: payload.s,
|
|
sessionId: payload.d.session_id,
|
|
@@ -981,6 +1007,8 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
resumeURL: payload.d.resume_gateway_url
|
|
};
|
|
await this.strategy.updateSessionInfo(this.id, session2);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
+ this.#status = 3 /* Ready */;
|
|
this.emit("ready" /* Ready */, { data: payload.d });
|
|
break;
|
|
}
|
|
@@ -995,9 +1023,11 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (session) {
|
|
if (payload.s > session.sequence) {
|
|
await this.strategy.updateSessionInfo(this.id, { ...session, sequence: payload.s });
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
}
|
|
} else {
|
|
this.debug([
|
|
@@ -1021,6 +1051,7 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
case import_v102.GatewayOpcodes.InvalidSession: {
|
|
this.debug([`Invalid session; will attempt to resume: ${payload.d.toString()}`]);
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (payload.d && session) {
|
|
await this.resume(session);
|
|
} else {
|
|
@@ -1036,17 +1067,19 @@ var WebSocketShard = class extends import_async_event_emitter.AsyncEventEmitter
|
|
const jitter = Math.random();
|
|
const firstWait = Math.floor(payload.d.heartbeat_interval * jitter);
|
|
this.debug([`Preparing first heartbeat of the connection with a jitter of ${jitter}; waiting ${firstWait}ms`]);
|
|
+ const controller = new AbortController();
|
|
try {
|
|
- const controller = new AbortController();
|
|
this.initialHeartbeatTimeoutController = controller;
|
|
await (0, import_promises2.setTimeout)(firstWait, void 0, { signal: controller.signal });
|
|
} catch {
|
|
this.debug(["Cancelled initial heartbeat due to #destroy being called"]);
|
|
return;
|
|
} finally {
|
|
- this.initialHeartbeatTimeoutController = null;
|
|
+ if (this.initialHeartbeatTimeoutController === controller) this.initialHeartbeatTimeoutController = null;
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.heartbeat();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([`First heartbeat sent, starting to beat every ${payload.d.heartbeat_interval}ms`]);
|
|
this.heartbeatInterval = (0, import_node_timers.setInterval)(() => void this.heartbeat(), payload.d.heartbeat_interval);
|
|
break;
|
|
@@ -1388,6 +1421,8 @@ var WebSocketManager = class extends import_async_event_emitter2.AsyncEventEmitt
|
|
this.options = { ...DefaultWebSocketManagerOptions, ...options };
|
|
this.strategy = this.options.buildStrategy(this);
|
|
}
|
|
+ connectionEpoch = 0;
|
|
+ destroyed = false;
|
|
/**
|
|
* Fetches the gateway information from Discord - or returns it from cache if available
|
|
*
|
|
@@ -1450,10 +1485,19 @@ var WebSocketManager = class extends import_async_event_emitter2.AsyncEventEmitt
|
|
return shardIds;
|
|
}
|
|
async connect() {
|
|
+ if (this.destroyed) return;
|
|
+ const epoch = this.connectionEpoch;
|
|
const shardCount = await this.getShardCount();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.updateShardCount(shardCount);
|
|
+ if (epoch !== this.connectionEpoch) {
|
|
+ await this.strategy.destroy();
|
|
+ return;
|
|
+ }
|
|
const shardIds = await this.getShardIds();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
const data = await this.fetchGatewayInformation();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (data.session_start_limit.remaining < shardIds.length) {
|
|
throw new Error(
|
|
`Not enough sessions remaining to spawn ${shardIds.length} shards; only ${data.session_start_limit.remaining} remaining; resets at ${new Date(Date.now() + data.session_start_limit.reset_after).toISOString()}`
|
|
@@ -1462,6 +1506,8 @@ var WebSocketManager = class extends import_async_event_emitter2.AsyncEventEmitt
|
|
await this.strategy.connect();
|
|
}
|
|
destroy(options) {
|
|
+ this.destroyed = true;
|
|
+ this.connectionEpoch += 1;
|
|
return this.strategy.destroy(options);
|
|
}
|
|
send(shardId, payload) {
|
|
diff --git a/dist/index.mjs b/dist/index.mjs
|
|
index 0ebc06d477ac54303c7542261cc5bcc9a9537519..bfb1f0880bf43f70e5c9176b62a1c23ca21596bc 100644
|
|
--- a/dist/index.mjs
|
|
+++ b/dist/index.mjs
|
|
@@ -566,6 +566,8 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
__name(this, "WebSocketShard");
|
|
}
|
|
connection = null;
|
|
+ connectionEpoch = 0;
|
|
+ destroyPromise = null;
|
|
useIdentifyCompress = false;
|
|
inflate = null;
|
|
textDecoder = new TextDecoder();
|
|
@@ -615,6 +617,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
if (this.#status !== 0 /* Idle */) {
|
|
throw new Error("Tried to connect a shard that wasn't idle");
|
|
}
|
|
+ const epoch = this.connectionEpoch;
|
|
const { version: version2, encoding, compression } = this.strategy.options;
|
|
const params = new URLSearchParams({ v: version2, encoding });
|
|
if (compression) {
|
|
@@ -633,6 +636,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
const url = `${session?.resumeURL ?? this.strategy.options.gatewayInformation.url}?${params.toString()}`;
|
|
this.debug([`Connecting to ${url}`]);
|
|
const connection = new WebSocketConstructor(url, [], {
|
|
@@ -640,21 +644,25 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
});
|
|
connection.binaryType = "arraybuffer";
|
|
connection.onmessage = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onMessage(event.data, event.data instanceof ArrayBuffer);
|
|
};
|
|
connection.onerror = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.onError(event.error);
|
|
};
|
|
connection.onclose = (event) => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
void this.onClose(event.code);
|
|
};
|
|
connection.onopen = () => {
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.sendRateLimitState = getInitialSendRateLimitState();
|
|
};
|
|
this.connection = connection;
|
|
this.#status = 1 /* Connecting */;
|
|
const { ok } = await this.waitForEvent("hello" /* Hello */, this.strategy.options.helloTimeout);
|
|
- if (!ok) {
|
|
+ if (!ok || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
if (session?.shardCount === this.strategy.options.shardCount) {
|
|
@@ -664,6 +672,26 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
}
|
|
async destroy(options = {}) {
|
|
+ // Invalidate pending setup and event waits before cancellation can resume them.
|
|
+ const epoch = ++this.connectionEpoch;
|
|
+ if (this.destroyPromise) {
|
|
+ await this.destroyPromise;
|
|
+ return;
|
|
+ }
|
|
+ if (this.#status === 0 /* Idle */) return;
|
|
+ const cleanup = this.destroyConnection(options);
|
|
+ this.destroyPromise = cleanup;
|
|
+ try {
|
|
+ await cleanup;
|
|
+ } finally {
|
|
+ if (this.destroyPromise === cleanup) this.destroyPromise = null;
|
|
+ }
|
|
+ if (options.recover !== void 0 && epoch === this.connectionEpoch) {
|
|
+ await sleep2(500);
|
|
+ if (epoch === this.connectionEpoch) return this.internalConnect();
|
|
+ }
|
|
+ }
|
|
+ async destroyConnection(options = {}) {
|
|
if (this.#status === 0 /* Idle */) {
|
|
this.debug(["Tried to destroy a shard that was idle"]);
|
|
return;
|
|
@@ -694,36 +722,27 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
if (options.recover !== 1 /* Resume */) {
|
|
await this.strategy.updateSessionInfo(this.id, null);
|
|
}
|
|
- if (this.connection) {
|
|
- this.connection.onmessage = null;
|
|
- this.connection.onclose = null;
|
|
- const shouldClose = this.connection.readyState === WebSocket.OPEN;
|
|
- this.debug([
|
|
- "Connection status during destroy",
|
|
- `Needs closing: ${shouldClose}`,
|
|
- `Ready state: ${this.connection.readyState}`
|
|
- ]);
|
|
- if (shouldClose) {
|
|
- let outerResolve;
|
|
- const promise = new Promise((resolve2) => {
|
|
- outerResolve = resolve2;
|
|
- });
|
|
- this.connection.onclose = outerResolve;
|
|
- this.connection.close(options.code, options.reason);
|
|
- await promise;
|
|
+ const connection = this.connection;
|
|
+ if (connection) {
|
|
+ connection.onmessage = null;
|
|
+ connection.onopen = null;
|
|
+ connection.onclose = null;
|
|
+ // A connecting ws can still emit its handshake error. Keep its epoch-bound
|
|
+ // error handler until terminate/close has actually completed.
|
|
+ if (connection.readyState !== 3 /* CLOSED */) {
|
|
+ const closed = new Promise((resolve) => { connection.onclose = resolve; });
|
|
+ if (connection.readyState === 0 /* CONNECTING */) connection.terminate();
|
|
+ else if (connection.readyState === 1 /* OPEN */) connection.close(options.code, options.reason);
|
|
+ await closed;
|
|
this.emit("closed" /* Closed */, { code: options.code });
|
|
}
|
|
- this.connection.onerror = null;
|
|
- } else {
|
|
- this.debug(["Destroying a shard that has no connection; please open an issue on GitHub"]);
|
|
+ connection.onerror = null;
|
|
+ if (this.connection === connection) this.connection = null;
|
|
}
|
|
this.#status = 0 /* Idle */;
|
|
- if (options.recover !== void 0) {
|
|
- await sleep2(500);
|
|
- return this.internalConnect();
|
|
- }
|
|
}
|
|
async waitForEvent(event, timeoutDuration) {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug([`Waiting for event ${event} ${timeoutDuration ? `for ${timeoutDuration}ms` : "indefinitely"}`]);
|
|
const timeoutController = new AbortController();
|
|
const timeout = timeoutDuration ? setTimeout(() => timeoutController.abort(), timeoutDuration).unref() : null;
|
|
@@ -736,6 +755,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
]);
|
|
return { ok: !closed };
|
|
} catch {
|
|
+ if (epoch !== this.connectionEpoch) return { ok: false };
|
|
void this.destroy({
|
|
code: 1e3 /* Normal */,
|
|
reason: "Something timed out or went wrong while waiting for an event",
|
|
@@ -746,7 +766,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
if (timeout) {
|
|
clearTimeout(timeout);
|
|
}
|
|
- this.timeoutAbortControllers.delete(event);
|
|
+ if (this.timeoutAbortControllers.get(event) === timeoutController) this.timeoutAbortControllers.delete(event);
|
|
if (!closeController.signal.aborted) {
|
|
closeController.abort();
|
|
}
|
|
@@ -793,6 +813,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
this.connection.send(JSON.stringify(payload));
|
|
}
|
|
async identify() {
|
|
+ const epoch = this.connectionEpoch;
|
|
this.debug(["Waiting for identify throttle"]);
|
|
const controller = new AbortController();
|
|
const closeHandler = /* @__PURE__ */ __name(() => {
|
|
@@ -802,7 +823,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
try {
|
|
await this.strategy.waitForIdentify(this.id, controller.signal);
|
|
} catch {
|
|
- if (controller.signal.aborted) {
|
|
+ if (controller.signal.aborted || epoch !== this.connectionEpoch) {
|
|
this.debug(["Was waiting for an identify, but the shard closed in the meantime"]);
|
|
return;
|
|
}
|
|
@@ -818,6 +839,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
} finally {
|
|
this.off("closed" /* Closed */, closeHandler);
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([
|
|
"Identifying",
|
|
`shard id: ${this.id.toString()}`,
|
|
@@ -842,6 +864,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
op: GatewayOpcodes2.Identify,
|
|
d
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.waitForEvent("ready" /* Ready */, this.strategy.options.readyTimeout);
|
|
}
|
|
async resume(session) {
|
|
@@ -863,14 +886,17 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
});
|
|
}
|
|
async heartbeat(requested = false) {
|
|
+ const epoch = this.connectionEpoch;
|
|
if (!this.isAck && !requested) {
|
|
return this.destroy({ reason: "Zombie connection", recover: 1 /* Resume */ });
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.send({
|
|
op: GatewayOpcodes2.Heartbeat,
|
|
d: session?.sequence ?? null
|
|
});
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.lastHeartbeatAt = Date.now();
|
|
this.isAck = false;
|
|
}
|
|
@@ -922,8 +948,9 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
return null;
|
|
}
|
|
async onMessage(data, isBinary) {
|
|
+ const epoch = this.connectionEpoch;
|
|
const payload = await this.unpackMessage(data, isBinary);
|
|
- if (!payload) {
|
|
+ if (!payload || epoch !== this.connectionEpoch) {
|
|
return;
|
|
}
|
|
switch (payload.op) {
|
|
@@ -933,7 +960,6 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
switch (payload.t) {
|
|
case GatewayDispatchEvents.Ready: {
|
|
- this.#status = 3 /* Ready */;
|
|
const session2 = {
|
|
sequence: payload.s,
|
|
sessionId: payload.d.session_id,
|
|
@@ -942,6 +968,8 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
resumeURL: payload.d.resume_gateway_url
|
|
};
|
|
await this.strategy.updateSessionInfo(this.id, session2);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
+ this.#status = 3 /* Ready */;
|
|
this.emit("ready" /* Ready */, { data: payload.d });
|
|
break;
|
|
}
|
|
@@ -956,9 +984,11 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
}
|
|
}
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (session) {
|
|
if (payload.s > session.sequence) {
|
|
await this.strategy.updateSessionInfo(this.id, { ...session, sequence: payload.s });
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
}
|
|
} else {
|
|
this.debug([
|
|
@@ -982,6 +1012,7 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
case GatewayOpcodes2.InvalidSession: {
|
|
this.debug([`Invalid session; will attempt to resume: ${payload.d.toString()}`]);
|
|
const session = await this.strategy.retrieveSessionInfo(this.id);
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (payload.d && session) {
|
|
await this.resume(session);
|
|
} else {
|
|
@@ -997,17 +1028,19 @@ var WebSocketShard = class extends AsyncEventEmitter {
|
|
const jitter = Math.random();
|
|
const firstWait = Math.floor(payload.d.heartbeat_interval * jitter);
|
|
this.debug([`Preparing first heartbeat of the connection with a jitter of ${jitter}; waiting ${firstWait}ms`]);
|
|
+ const controller = new AbortController();
|
|
try {
|
|
- const controller = new AbortController();
|
|
this.initialHeartbeatTimeoutController = controller;
|
|
await sleep2(firstWait, void 0, { signal: controller.signal });
|
|
} catch {
|
|
this.debug(["Cancelled initial heartbeat due to #destroy being called"]);
|
|
return;
|
|
} finally {
|
|
- this.initialHeartbeatTimeoutController = null;
|
|
+ if (this.initialHeartbeatTimeoutController === controller) this.initialHeartbeatTimeoutController = null;
|
|
}
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.heartbeat();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
this.debug([`First heartbeat sent, starting to beat every ${payload.d.heartbeat_interval}ms`]);
|
|
this.heartbeatInterval = setInterval(() => void this.heartbeat(), payload.d.heartbeat_interval);
|
|
break;
|
|
@@ -1351,6 +1384,8 @@ var WebSocketManager = class extends AsyncEventEmitter2 {
|
|
this.options = { ...DefaultWebSocketManagerOptions, ...options };
|
|
this.strategy = this.options.buildStrategy(this);
|
|
}
|
|
+ connectionEpoch = 0;
|
|
+ destroyed = false;
|
|
/**
|
|
* Fetches the gateway information from Discord - or returns it from cache if available
|
|
*
|
|
@@ -1413,10 +1448,19 @@ var WebSocketManager = class extends AsyncEventEmitter2 {
|
|
return shardIds;
|
|
}
|
|
async connect() {
|
|
+ if (this.destroyed) return;
|
|
+ const epoch = this.connectionEpoch;
|
|
const shardCount = await this.getShardCount();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
await this.updateShardCount(shardCount);
|
|
+ if (epoch !== this.connectionEpoch) {
|
|
+ await this.strategy.destroy();
|
|
+ return;
|
|
+ }
|
|
const shardIds = await this.getShardIds();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
const data = await this.fetchGatewayInformation();
|
|
+ if (epoch !== this.connectionEpoch) return;
|
|
if (data.session_start_limit.remaining < shardIds.length) {
|
|
throw new Error(
|
|
`Not enough sessions remaining to spawn ${shardIds.length} shards; only ${data.session_start_limit.remaining} remaining; resets at ${new Date(Date.now() + data.session_start_limit.reset_after).toISOString()}`
|
|
@@ -1425,6 +1469,8 @@ var WebSocketManager = class extends AsyncEventEmitter2 {
|
|
await this.strategy.connect();
|
|
}
|
|
destroy(options) {
|
|
+ this.destroyed = true;
|
|
+ this.connectionEpoch += 1;
|
|
return this.strategy.destroy(options);
|
|
}
|
|
send(shardId, payload) {
|