From 61acc5781f778a54e9ca8eb437a43a22351208ad Mon Sep 17 00:00:00 2001 From: Fleny Date: Fri, 8 Nov 2024 15:16:10 +0100 Subject: [PATCH] fix(gateway): Fix shard reconnecting in a loop (#3907) * Attempt to fix shard reconnecting in a loop * Add a bit of logging on READY & RESUMED * Fix integration tests * Revert turbo.json change * Pass down the logger to LeakyBucket on HELLO * Update some logs --------- Co-authored-by: Awesome Stickz --- packages/gateway/src/Shard.ts | 124 ++++++++++-------- packages/gateway/src/manager.ts | 20 +-- .../tests/integration/connection.spec.ts | 6 +- 3 files changed, 77 insertions(+), 73 deletions(-) diff --git a/packages/gateway/src/Shard.ts b/packages/gateway/src/Shard.ts index 147233a50..f874bb709 100644 --- a/packages/gateway/src/Shard.ts +++ b/packages/gateway/src/Shard.ts @@ -227,7 +227,7 @@ export class DiscordenoShard { // A new identify has been requested even though there is already a connection open. // Therefore we need to close the old connection and heartbeating before creating a new one. if (this.isOpen()) { - this.logger.debug(`CLOSING EXISTING SHARD: #${this.id}`) + this.logger.debug(`[Shard] Identifying open Shard #${this.id}, closing the connection`) this.close(ShardSocketCloseCodes.ReIdentifying, 'Re-identifying closure of old connection.') } @@ -250,7 +250,7 @@ export class DiscordenoShard { properties: this.gatewayConfig.properties, intents: this.gatewayConfig.intents, shard: [this.id, this.gatewayConfig.totalShards], - presence: await this.makePresence?.(), + presence: await this.makePresence(), }, }, true, @@ -263,8 +263,7 @@ export class DiscordenoShard { this.shardIsReady() resolve() }) - // When identifying too fast, - // Discord sends an invalid session payload. + // When identifying too fast, Discord sends an invalid session payload. // This can safely be ignored though and the shard starts a new identify action. this.resolves.set('INVALID_SESSION', () => { this.resolves.delete('READY') @@ -280,27 +279,29 @@ export class DiscordenoShard { /** Attempt to resume the previous shards session with the gateway. */ async resume(): Promise { - this.logger.debug(`[Gateway] Resuming Shard #${this.id}`) + this.logger.debug(`[Shard] Resuming Shard #${this.id}`) + // It has been requested to resume the Shards session. // It's possible that the shard is still connected with Discord's gateway therefore we need to forcefully close it. if (this.isOpen()) { - this.logger.debug(`[Gateway] Resuming Shard #${this.id} in isOpen`) + this.logger.debug(`[Shard] Resuming open Shard #${this.id}, closing the connection`) this.close(ShardSocketCloseCodes.ResumeClosingOldConnection, 'Reconnecting the shard, closing old connection.') } // Shard has never identified, so we cannot resume. if (!this.sessionId) { - this.logger.debug(`[Shard] Trying to resume a shard #${this.id} that was NOT first identified. (No session id found)`) + this.logger.debug(`[Shard] Trying to resume Shard #${this.id} without the session id. Identifying the shard instead.`) - return await this.identify() + await this.identify() + return } this.state = ShardState.Resuming - this.logger.debug(`[Gateway] Resuming Shard #${this.id}, before connecting`) // Before we can resume, we need to create a new connection with Discord's gateway. await this.connect() - this.logger.debug(`[Gateway] Resuming Shard #${this.id}, after connecting. ${this.sessionId} | ${this.previousSequenceNumber}`) + + this.logger.debug(`[Shard] Resuming Shard #${this.id} connected. Session id: ${this.sessionId} | Sequence: ${this.previousSequenceNumber}`) this.send( { @@ -313,12 +314,11 @@ export class DiscordenoShard { }, true, ) - this.logger.debug(`[Shard] Resuming Shard #${this.id} after send resume`) return await new Promise((resolve) => { this.resolves.set('RESUMED', () => resolve()) - // If it is attempted to resume with an invalid session id, - // Discord sends an invalid session payload + + // If it is attempted to resume with an invalid session id, Discord sends an invalid session payload // Not erroring here since it is easy that this happens, also it would be not catchable this.resolves.set('INVALID_SESSION', () => { this.resolves.delete('RESUMED') @@ -327,8 +327,9 @@ export class DiscordenoShard { }) } - /** Send a message to Discord. - * @param {boolean} [highPriority=false] - Whether this message should be send asap. + /** + * Send a message to Discord. + * @param highPriority - Whether this message should be send asap. */ async send(message: ShardSocketRequest, highPriority: boolean = false): Promise { // Before acquiring a token from the bucket, check whether the shard is currently offline or not. @@ -351,11 +352,12 @@ export class DiscordenoShard { /** Handle a gateway connection error */ handleError(error: Event): void { - this.logger.error(`[Shard] There was an error connecting shard ${this.id}.`, error) + this.logger.error(`[Shard] There was an error connecting Shard #${this.id}.`, error) } /** Handle a gateway connection close. */ async handleClose(close: CloseEvent): Promise { + this.socket = undefined this.stopHeartbeating() // Clear the zlib/zstd data @@ -364,9 +366,7 @@ export class DiscordenoShard { this.inflateBuffer = null this.decompressionPromisesQueue = [] - this.logger.debug( - `[Shard] Gateway connection closed with code ${close.code} (${close.reason || ''}).`, - ) + this.logger.debug(`[Shard] Shard #${this.id} closed with code ${close.code}${close.reason ? `, and reason: ${close.reason}` : ''}.`) switch (close.code) { case ShardSocketCloseCodes.TestingFinished: { @@ -384,21 +384,8 @@ export class DiscordenoShard { this.state = ShardState.Disconnected this.events.disconnected?.(this) - // gateway.debug("GW CLOSED_RECONNECT", { shardId, payload: event }); return } - // Gateway connection closes which require a new identify. - case GatewayCloseEventCodes.UnknownOpcode: - case GatewayCloseEventCodes.NotAuthenticated: - case GatewayCloseEventCodes.InvalidSeq: - case GatewayCloseEventCodes.RateLimited: - case GatewayCloseEventCodes.SessionTimedOut: { - this.logger.debug('[Shard] Gateway connection closing requiring re-identify.') - this.state = ShardState.Identifying - this.events.disconnected?.(this) - - return await this.identify() - } // When these codes are received something went really wrong. // On those we cannot start a reconnect attempt. case GatewayCloseEventCodes.AuthenticationFailed: @@ -412,16 +399,36 @@ export class DiscordenoShard { throw new Error(close.reason || 'Discord gave no reason! GG! You broke Discord!') } - // Gateway connection closes on which a resume is allowed. - case GatewayCloseEventCodes.UnknownError: - case GatewayCloseEventCodes.DecodeError: - case GatewayCloseEventCodes.AlreadyAuthenticated: - default: { - this.logger.info(`[Shard] Closed shard #${this.id} with code ${close.code}. Attempting to resume...`) - this.state = ShardState.Resuming + // Gateway connection closes which require a new identify. + case GatewayCloseEventCodes.NotAuthenticated: + case GatewayCloseEventCodes.InvalidSeq: + case GatewayCloseEventCodes.SessionTimedOut: { + this.logger.debug(`[Shard] Shard #${this.id} closed requiring re-identify.`) + this.state = ShardState.Identifying this.events.disconnected?.(this) - return await this.resume() + await this.identify() + return + } + // Gateway connection closes on which a resume is allowed. + case GatewayCloseEventCodes.UnknownError: + case GatewayCloseEventCodes.UnknownOpcode: + case GatewayCloseEventCodes.DecodeError: + case GatewayCloseEventCodes.RateLimited: + case GatewayCloseEventCodes.AlreadyAuthenticated: + default: { + this.logger.info(`[Shard] Shard #${this.id} closed with code ${close.code}. Attempting to resume...`) + // We don't want to get into an infinite loop where we resume forever, so if we were already resuming we identify instead + this.state = this.state === ShardState.Resuming ? ShardState.Identifying : ShardState.Resuming + this.events.disconnected?.(this) + + if (this.state === ShardState.Resuming) { + await this.resume() + } else { + await this.identify() + } + + return } } } @@ -507,7 +514,6 @@ export class DiscordenoShard { switch (packet.op) { case GatewayOpcodes.Heartbeat: { - // TODO: can this actually happen if (!this.isOpen()) return this.heart.lastBeat = Date.now() @@ -525,7 +531,7 @@ export class DiscordenoShard { } case GatewayOpcodes.Hello: { const interval = (packet.d as DiscordHello).heartbeat_interval - this.logger.debug(`[Gateway] Hello on Shard #${this.id}`) + this.logger.debug(`[Shard] Shard #${this.id} received Hello`) this.startHeartbeating(interval) if (this.state !== ShardState.Resuming) { @@ -537,6 +543,7 @@ export class DiscordenoShard { max: this.calculateSafeRequests(), refillInterval: 60000, refillAmount: this.calculateSafeRequests(), + logger: this.logger, }) // Queue should not be lost on a re-identify. @@ -558,8 +565,6 @@ export class DiscordenoShard { break } case GatewayOpcodes.Reconnect: { - // gateway.debug("GW RECONNECT", { shardId }); - this.events.requestedReconnect?.(this) await this.resume() @@ -568,7 +573,7 @@ export class DiscordenoShard { } case GatewayOpcodes.InvalidSession: { const resumable = packet.d as boolean - this.logger.debug(`[Shard] Received Invalid Session for Shard #${this.id} with resumeable as ${resumable.toString()}`) + this.logger.debug(`[Shard] Received Invalid Session for Shard #${this.id} with resumable as ${resumable}`) this.events.invalidSession?.(this, resumable) @@ -598,6 +603,8 @@ export class DiscordenoShard { this.state = ShardState.Connected this.events.resumed?.(this) + this.logger.debug(`[Shard] Shard #${this.id} received RESUMED`) + // Continue the requests which have been queued since the shard went offline. this.offlineSendQueue.forEach((resolve) => resolve()) // Setting the length to 0 will delete the elements in it @@ -607,14 +614,16 @@ export class DiscordenoShard { this.resolves.delete('RESUMED') break case 'READY': { - // Important for future resumes. const payload = packet.d as DiscordReady + // Important for future resumes. this.resumeGatewayUrl = payload.resume_gateway_url - this.sessionId = payload.session_id + this.state = ShardState.Connected + this.logger.debug(`[Shard] Shard #${this.id} received READY`) + // Continue the requests which have been queued since the shard went offline. // Important when this is a re-identify this.offlineSendQueue.forEach((resolve) => resolve()) @@ -652,7 +661,10 @@ export class DiscordenoShard { return } - /** This function communicates with the management process, in order to know whether its free to identify. When this function resolves, this means that the shard is allowed to send an identify payload to discord. */ + /** + * This function communicates with the management process, in order to know whether its free to identify. + * When this function resolves, this means that the shard is allowed to send an identify payload to discord. + */ async requestIdentify(): Promise {} /** This function communicates with the management process, in order to tell it can identify the next shard. */ @@ -660,11 +672,10 @@ export class DiscordenoShard { /** Start sending heartbeat payloads to Discord in the provided interval. */ startHeartbeating(interval: number): void { - this.logger.debug(`[Shard] Start heartbeating on shard #${this.id}`) + this.logger.debug(`[Shard] Start heartbeating on Shard #${this.id}`) // If old heartbeast exist like after resume, clear the old ones. - if (this.heart.intervalId) clearInterval(this.heart.intervalId) - if (this.heart.timeoutId) clearTimeout(this.heart.timeoutId) + this.stopHeartbeating() this.heart.interval = interval @@ -682,11 +693,11 @@ export class DiscordenoShard { const jitter = Math.ceil(this.heart.interval * (Math.random() || 0.5)) this.heart.timeoutId = setTimeout(() => { - this.logger.debug(`[Shard] Beginning heartbeating process for shard #${this.id}`) + this.logger.debug(`[Shard] Beginning heartbeating process for Shard #${this.id}`) if (!this.isOpen()) return - this.logger.debug(`[Shard] Heartbeating on #${this.id}. Previous sequence number: ${this.previousSequenceNumber}`) + this.logger.debug(`[Shard] Heartbeating on Shard #${this.id}. Previous sequence number: ${this.previousSequenceNumber}`) // Using a direct socket.send call here because heartbeat requests are reserved by us. this.socket?.send( @@ -711,13 +722,14 @@ export class DiscordenoShard { // The Shard needs to start a re-identify action accordingly. // Reference: https://discord.com/developers/docs/topics/gateway#heartbeating-example-gateway-heartbeat-ack if (!this.heart.acknowledged) { - this.logger.debug(`[Shard] Heartbeat not acknowledged for shard #${this.id}. Assuming zombied connection.`) + this.logger.debug(`[Shard] Heartbeat not acknowledged for Shard #${this.id}. Assuming zombied connection.`) this.close(ShardSocketCloseCodes.ZombiedConnection, 'Zombied connection, did not receive an heartbeat ACK in time.') - return await this.identify() + await this.resume() + return } - this.logger.debug(`[Shard] Heartbeating on #${this.id}. Previous sequence number: ${this.previousSequenceNumber}`) + this.logger.debug(`[Shard] Heartbeating on Shard #${this.id}. Previous sequence number: ${this.previousSequenceNumber}`) // Using a direct socket.send call here because heartbeat requests are reserved by us. this.socket?.send( diff --git a/packages/gateway/src/manager.ts b/packages/gateway/src/manager.ts index a24e5880b..14634f051 100644 --- a/packages/gateway/src/manager.ts +++ b/packages/gateway/src/manager.ts @@ -196,7 +196,7 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate return await new Promise((resolve) => { // Mark that we are making an identify request so another is not made. bucket.identifyRequests.push(resolve) - gateway.logger.debug(`[Gateway] identifying shard #(${shardId}).`) + gateway.logger.debug(`[Gateway] Identifying Shard #${shardId}.`) // This will trigger identify and when READY is received it will resolve the above request. shard?.identify().then(async () => { // Tell the manager that this shard is online @@ -371,12 +371,12 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate await shard.send(payload) }, async tellWorkerToIdentify(workerId, shardId, bucketId) { - gateway.logger.debug(`[Gateway] tell worker to identify (${workerId}, ${shardId}, ${bucketId})`) + gateway.logger.debug(`[Gateway] Tell worker to identify (${workerId}, ${shardId}, ${bucketId})`) await gateway.identify(shardId) }, async identify(shardId: number) { let shard = this.shards.get(shardId) - gateway.logger.debug(`[Gateway] identifying ${shard ? 'existing' : 'new'} shard (${shardId})`) + gateway.logger.debug(`[Gateway] Identifying ${shard ? 'existing' : 'new'} shard (${shardId})`) if (!shard) { shard = new Shard({ @@ -420,7 +420,7 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate return await new Promise((resolve) => { // Mark that we are making an identify request so another is not made. bucket.identifyRequests.push(resolve) - gateway.logger.debug(`[Gateway] identifying shard #(${shardId}).`) + gateway.logger.debug(`[Gateway] Identifying Shard #${shardId}.`) // This will trigger identify and when READY is received it will resolve the above request. shard?.identify() }) @@ -428,21 +428,15 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate async kill(shardId: number) { const shard = this.shards.get(shardId) if (!shard) { - return gateway.logger.debug(`[Gateway] kill shard but not found (${shardId})`) + return gateway.logger.debug(`[Gateway] A kill for Shard #${shardId} was requested, but the shard could not be found`) } - gateway.logger.debug(`[Gateway] kill shard (${shardId})`) + gateway.logger.debug(`[Gateway] Killing Shard #${shardId}`) this.shards.delete(shardId) await shard.shutdown() }, async requestIdentify(_shardId: number) { - gateway.logger.debug(`[Gateway] requesting identify`) - // const bucket = gateway.buckets.get(shardId % gateway.connection.sessionStartLimit.maxConcurrency) - // if (!bucket) return - - // return await new Promise((resolve) => { - // bucket.identifyRequests.push(resolve) - // }) + gateway.logger.debug(`[Gateway] Requesting identify`) }, // Helpers methods below this diff --git a/packages/gateway/tests/integration/connection.spec.ts b/packages/gateway/tests/integration/connection.spec.ts index 359ce2752..781577624 100644 --- a/packages/gateway/tests/integration/connection.spec.ts +++ b/packages/gateway/tests/integration/connection.spec.ts @@ -154,8 +154,7 @@ describe('gateway', () => { await connected uwsOptions.closing = true - - await gateway.shutdown(ShardSocketCloseCodes.Shutdown, 'User requested bot stop') + await gateway.shutdown(ShardSocketCloseCodes.Shutdown, 'User requested bot stop', true) uWS.us_listen_socket_close(uwsToken) }) @@ -191,8 +190,7 @@ describe('gateway', () => { clearTimeout(timeout) uwsOptions.closing = true - - await gateway.shutdown(ShardSocketCloseCodes.Shutdown, 'User requested bot stop') + await gateway.shutdown(ShardSocketCloseCodes.Shutdown, 'User requested bot stop', true) uWS.us_listen_socket_close(uwsToken) })