diff --git a/src/backend/common/vendor/discord/DiscordWSClient.ts b/src/backend/common/vendor/discord/DiscordWSClient.ts index 03be7e2b..7c9635c0 100644 --- a/src/backend/common/vendor/discord/DiscordWSClient.ts +++ b/src/backend/common/vendor/discord/DiscordWSClient.ts @@ -19,7 +19,7 @@ import { isSuperAgentResponseError } from "../../errors/ErrorUtils.js"; import { urlToMusicService } from "../ListenbrainzApiClient.js"; import { urlContainsKnownMediaDomain } from "../../../utils/RequestUtils.js"; import { CoverArtApiClient } from "../musicbrainz/CoverArtApiClient.js"; -import { formatWebsocketClose, isCloseEvent, isErrorEvent } from "../../../utils/NetworkUtils.js"; +import { formatWebsocketClose, isCloseEvent, isErrorEvent, wsReadyStateToStr } from "../../../utils/NetworkUtils.js"; const ARTWORK_PLACEHOLDER = 'https://raw.githubusercontent.com/FoxxMD/multi-scrobbler/master/assets/default-artwork.png'; const MS_ART = 'https://raw.githubusercontent.com/FoxxMD/multi-scrobbler/master/assets/icon.png'; @@ -40,7 +40,7 @@ export class DiscordWSClient extends AbstractApiClient { declare config: DiscordStrongData; - heartbeatInterval: NodeJS.Timeout + heartbeatInterval?: NodeJS.Timeout acknowledged: boolean = true; // https://docs.discord.com/developers/events/gateway#ready-event @@ -63,7 +63,10 @@ export class DiscordWSClient extends AbstractApiClient { lastActiveStatus?: PresenceUpdateStatus = PresenceUpdateStatus.Offline; lastActivities: GatewayActivity[] = []; - activityTimeout: NodeJS.Timeout; + activityTimeout?: NodeJS.Timeout; + clearLastActivitiesTimeout?: NodeJS.Timeout; + + get friendlySocketState() { return `Socket state: ${wsReadyStateToStr(this.client.readyState)}`} emitter: EventEmitter; @@ -81,24 +84,6 @@ export class DiscordWSClient extends AbstractApiClient { } initClient = async () => { - // let baseUrl: string; - // if (this.resume_gateway_url !== undefined) { - // baseUrl = this.resume_gateway_url; - // } else { - // if (this.initialGatewayUrl === undefined) { - // try { - // await this.fetchGatewayUrl(); - // baseUrl = this.initialGatewayUrl; - // } catch (e) { - // throw new Error('Could not get initial gateway url', { cause: e }); - // } - // } else { - // baseUrl = this.initialGatewayUrl; - // } - // } - - // const gatewayUrl = `${baseUrl}?encoding=json&v=10`; - // this.logger.debug(`Using Gateway URL ${gatewayUrl}`); const url = () => { let baseUrl: string; @@ -219,10 +204,10 @@ export class DiscordWSClient extends AbstractApiClient { protected authenticate = () => { if (this.canResume && this.session_id !== undefined) { // using resume - this.handleResume(); + this.sendResume(); } else { // initial identify - this.handleIdentify(); + this.sendIdentify(); } } @@ -295,7 +280,7 @@ export class DiscordWSClient extends AbstractApiClient { } } - handleIdentify() { + sendIdentify() { const data: GatewayIdentify = { op: GatewayOpcodes.Identify, d: { @@ -321,7 +306,7 @@ export class DiscordWSClient extends AbstractApiClient { // jitter // https://docs.discord.com/developers/events/gateway#heartbeat-interval const sleepTime = randomInt(data.heartbeat_interval - 1); - this.logger.debug(`Heartbeat Interval: ${data.heartbeat_interval}ms (${Math.floor(data.heartbeat_interval / 1000)}s), waiting ${Math.floor(sleepTime / 1000)}s before sending first heartbeat.`); + this.logger.debug(`Heartbeat Interval: ${data.heartbeat_interval}ms (${Math.floor(data.heartbeat_interval / 1000)}s), waiting ${Math.floor(sleepTime / 1000)}s before sending first heartbeat. (Seq ${this.sequence})`); await sleep(sleepTime); if (this.client.readyState !== this.client.OPEN) { this.logger.warn(`Not continuing with heartbeat because connection is not open`); @@ -336,7 +321,7 @@ export class DiscordWSClient extends AbstractApiClient { return; } if (!this.acknowledged) { - // zombied! + this.logger.warn('Did not recieve heartbeat acknowledgment! May be a zombie so trying to reconnect.'); return this.handleReconnect().then(() => null).catch((e) => this.logger.error(e)); } this.sendHeartbeat(); @@ -356,7 +341,7 @@ export class DiscordWSClient extends AbstractApiClient { this.logger.warn(`Cannot send heartbeat because connection is not open`); return; } else if(isDebugMode()) { - this.logger.debug('Sending heartbeat'); + this.logger.debug(`Sending heartbeat (Seq ${this.sequence})`); } this.client.send(JSON.stringify(heartbeatRequest)); } @@ -373,7 +358,18 @@ export class DiscordWSClient extends AbstractApiClient { this.ready = true; this.authOK = true; this.reconnecting = false; - this.logger.verbose(`Gateway Connection READY for ${this.user.username}`); + this.logger.verbose(`Gateway Connection READY for ${this.user.username} | Session ${this.session_id} (Seq ${this.sequence})`); + this.cancelClearLastActivities(); + this.emitter.emit('ready', {ready: true}); + } + + handleResume() { + this.logger.verbose(`Recieved Resumed | Session ${this.session_id} (Seq ${this.sequence})`); + this.canResume = true; + this.ready = true; + this.authOK = true; + this.reconnecting = false; + this.cancelClearLastActivities(); this.emitter.emit('ready', {ready: true}); } @@ -384,29 +380,36 @@ export class DiscordWSClient extends AbstractApiClient { } async handleReconnect() { - + this.logger.verbose('Starting reconnect attempt'); // on a manual close ws-isosocket does not retry // so we need to do it manually this.client.close(); - const result = await Promise.race([ - pEvent(this.client, 'close'), - sleep(3000), - ]); - if(result === undefined) { - throw new Error('Waited too long for client to close'); + + // closing manually also does not trigger a 'close' event so we need to check readyState + const now = dayjs(); + let elapsed = 0; + if(this.client.readyState !== this.client.CLOSED) { + this.logger.debug(`${this.friendlySocketState}, waiting for it to be closed...`); + while(this.client.readyState !== this.client.CLOSED && elapsed < 10000) { + elapsed = Math.abs(dayjs().diff(now, 'ms')); + this.logger.debug(`Elapsed ${elapsed}ms | ${this.friendlySocketState}`); + await sleep(1000); + } + } + + if(this.client.readyState !== this.client.CLOSED && this.client.readyState !== this.client.CONNECTING) { + throw new Error('Waited too long for socket to close'); } + this.logger.debug('Socket closed, reconnecting'); + this.cleanupConnectionSync(); try { await this.tryAuthenticate(); } catch (e) { throw new Error('Could not manually reconnect', {cause: e}); } - //await this.cleanupConnection(); - // maybe don't do this if we've failed N times - //this.initClient(); - //this.connect(); } - handleResume() { + sendResume() { const data: GatewayResumeData = { token: this.config.token, session_id: this.session_id, @@ -417,40 +420,42 @@ export class DiscordWSClient extends AbstractApiClient { this.logger.warn(`Cannot send resume because connection is not open`); return; } else if(isDebugMode()) { - this.logger.debug('Sending resume'); + this.logger.verbose(`Sending resume | Session ${this.session_id} | Seq ${this.sequence}`); } this.client.send(JSON.stringify({ op: GatewayOpcodes.Resume, d: data })); } - async cleanupConnection() { - clearInterval(this.heartbeatInterval); - this.heartbeatInterval = undefined; - clearTimeout(this.activityTimeout); - this.activityTimeout = undefined; - - if (this.client.CLOSED !== this.client.readyState) { - this.client.close(); - // wait for close or just give it a few seconds - const result = await Promise.race([ - pEvent(this.client, 'close'), - sleep(3000), - ]); - } else if(this.client.CLOSING === this.client.readyState) { - this.logger.debug('Giving the client time to close...'); - await sleep(3000); - } - this.ready = false; - if (!this.canResume) { - this.logger.debug('Cannot resume session, clearing session data for clean reconnect'); - this.session_id = undefined; - this.sequence = undefined; - this.resume_gateway_url = undefined; - this.user = undefined; + clearLastActivities(now: boolean = false) { + this.cancelClearLastActivities(); + if(!now) { + // clear activities if takes longer than 5 seconds to get READY + // + // allows us to reasonably assume activities have not changed during a clean reconnect that did not take very long + // so that we don't need to wait for another SESSION_REPLACE and can immediately update now playing + // + // but if user has to restart later, or something else goes wrong, we are sure we have reset activities when next READY + this.logger.debug('Delaying last activities clear for 5 seconds...'); + this.clearLastActivitiesTimeout = setTimeout(() => { + this.logger.debug('Clearing last activities'); + this.lastActiveStatus = PresenceUpdateStatus.Offline; + this.lastActivities = []; + this.clearLastActivitiesTimeout = undefined; + }, 5000); + } else { + this.logger.debug('Clearing last activities'); this.lastActiveStatus = PresenceUpdateStatus.Offline; this.lastActivities = []; } } + cancelClearLastActivities() { + if(this.clearLastActivitiesTimeout !== undefined) { + this.logger.debug('Cancelling clear activities timeout'); + clearTimeout(this.clearLastActivitiesTimeout); + this.clearLastActivitiesTimeout = undefined; + } + } + cleanupConnectionSync() { if(this.heartbeatInterval !== undefined) { clearInterval(this.heartbeatInterval); @@ -462,20 +467,18 @@ export class DiscordWSClient extends AbstractApiClient { } this.ready = false; if (!this.canResume) { - this.logger.debug('Cannot resume session, clearing session data for clean reconnect'); + this.logger.verbose('Cannot resume session, clearing session data for clean reconnect'); this.session_id = undefined; this.sequence = undefined; this.resume_gateway_url = undefined; this.user = undefined; - this.lastActiveStatus = PresenceUpdateStatus.Offline; - this.lastActivities = []; + this.clearLastActivities(); } } handleUserSessionUpdates = (data: UserSession[]) => { - this.logger.debug('Recieved updated user sessions'); if (data.filter(x => x.session_id !== this.session_id && x.session_id !== 'all').length === 0) { - this.logger.debug('No other user sessions exist, marking our session presence as inactive'); + this.logger.debug(`Recieved updated user sessions (Seq ${this.sequence}) => No other user sessions exist, marking our session presence as inactive`); this.lastActiveStatus = PresenceUpdateStatus.Offline; this.lastActivities = []; return; @@ -491,7 +494,7 @@ export class DiscordWSClient extends AbstractApiClient { }; return sessionId; }); - this.logger.debug(sessionSummaries.join('\n')); + this.logger.debug(`Recieved updated user sessions (Seq ${this.sequence})\n${sessionSummaries.join('\n')}`); const last = this.lastActiveStatus; @@ -518,25 +521,18 @@ export class DiscordWSClient extends AbstractApiClient { const ourSession = data.find(x => x.session_id === this.session_id); if(ourSession !== undefined && ourSession.activities.length > 0) { this.logger.debug(`Clearing our session presence, MS presence no longer allowed because ${reason}`); - this.clearActivity(); + this.sendClearActivity(); } } } async handleMessage(message: _DataPayload | _NonDispatchPayload) { + const { op, s } = message; + if (s !== null && s !== undefined) { + this.sequence = s; + } + const seqHint = `Seq ${this.sequence}`; try { - const { op, s } = message; - if (s !== null && s !== undefined) { - this.sequence = s; - } - if (isDebugMode()) { - const friendlyOp = opcodeToFriendly(op); - let handleHint = `Got opcode ${op}${friendlyOp !== op ? ` (${friendlyOp})` : ''}`; - if (op === GatewayOpcodes.Dispatch) { - handleHint += ` w/ Dispatch Event ${message.t}`; - } - this.logger.debug(handleHint); - } switch (op) { case GatewayOpcodes.Hello: @@ -544,13 +540,13 @@ export class DiscordWSClient extends AbstractApiClient { break; case GatewayOpcodes.HeartbeatAck: if (isDebugMode()) { - this.logger.debug("Heartbeat acknowledged"); + this.logger.debug(`Received Heartbeat acknowledgement (${seqHint})`); } this.acknowledged = true; break; case GatewayOpcodes.Heartbeat: if (isDebugMode()) { - this.logger.debug("Received Heartbeat"); + this.logger.debug(`Received Heartbeat request (${seqHint})`); } this.sendHeartbeat(); break; @@ -569,36 +565,40 @@ export class DiscordWSClient extends AbstractApiClient { // @ts-expect-error this.handleUserSessionUpdates(message.d as UserSession[]); break; + case GatewayDispatchEvents.Resumed: + this.handleResume(); + break; + default: + if (isDebugMode()) { + this.logger.debug(`Recieved Dispatch Event ${message.t} (${seqHint})`); + } }; break; case GatewayOpcodes.InvalidSession: - this.logger.debug('Recieved invalid session opcode'); + this.logger.verbose(`Recieved invalid session opcode (${seqHint})`); this.handleInvalidSession(message.d as GatewayInvalidSessionData); break; case GatewayOpcodes.Reconnect: - this.logger.debug('Recieved reconnect opcode'); + this.logger.verbose(`Recieved reconnect opcode (${seqHint})`); await this.handleReconnect(); break; - case GatewayOpcodes.Resume: - this.logger.debug({ data: message.d }, 'Recieved Resumed session'); - this.canResume = true; - this.ready = true; - this.authOK = true; - this.reconnecting = false; - this.emitter.emit('ready', {ready: true}); - break; case GatewayOpcodes.Identify: - this.logger.debug({ data: message.d }, 'Recieved Identifiy opcode'); + this.logger.debug({ data: message.d }, `Recieved Identify opcode (${seqHint})`); break; case GatewayOpcodes.PresenceUpdate: - this.logger.debug({ data: message.d }, 'Recieved Presence Update opcode'); + this.logger.debug({ data: message.d }, `Recieved Presence Update opcode (${seqHint})`); break; default: - this.logger.debug(`Recieved unhandled opcode: ${op}`); + this.logger.debug(`Recieved opcode: ${op} (${seqHint})`); break; } } catch (error) { - throw new Error('Error handling gateway message', { cause: error }); + const friendlyOp = opcodeToFriendly(op); + let handleHint = `opcode ${op}${friendlyOp !== op ? ` (${friendlyOp})` : ''}`; + if (op === GatewayOpcodes.Dispatch) { + handleHint += ` w/ Dispatch Event ${message.t}`; + } + throw new Error(`Error handling gateway message for ${handleHint} (${seqHint})`, { cause: error }); } } @@ -666,7 +666,7 @@ export class DiscordWSClient extends AbstractApiClient { sendActivity = async (data: SourceData | undefined) => { if(data === undefined) { - this.clearActivity(); + this.sendClearActivity(); return; } const [sendOk, reasons] = this.checkOkToSend(); @@ -702,11 +702,11 @@ export class DiscordWSClient extends AbstractApiClient { clearTimeout(this.activityTimeout); } this.activityTimeout = setTimeout(() => { - this.clearActivity(); + this.sendClearActivity(); }, Math.abs(clearTime.diff(dayjs(), 'ms'))); } - clearActivity = () => { + sendClearActivity = () => { if (this.activityTimeout !== undefined) { clearTimeout(this.activityTimeout); this.activityTimeout = undefined; diff --git a/src/backend/utils/NetworkUtils.ts b/src/backend/utils/NetworkUtils.ts index 1e871be9..3e50bceb 100644 --- a/src/backend/utils/NetworkUtils.ts +++ b/src/backend/utils/NetworkUtils.ts @@ -250,4 +250,19 @@ export const isErrorEvent = (e: Event): e is ErrorEvent => { } export const isRetryEvent = (e: Event): e is RetryEvent => { return e.type === 'retry'; +} + +export const wsReadyStateToStr = (state: number): string => { + switch(state) { + case 0: + return 'connecting'; + case 1: + return 'open'; + case 2: + return 'closing'; + case 3: + return 'closed'; + default: + return state.toString(); + } } \ No newline at end of file