diff --git a/events/Twitch/Connected.js b/events/Twitch/Connected.js index 998d415..093680e 100644 --- a/events/Twitch/Connected.js +++ b/events/Twitch/Connected.js @@ -1,8 +1,129 @@ // Event triggered by client connected to server + +const ReconnectingWebSocket = require("reconnecting-websocket"); +const { mClient } = require("../.."); +const { getIDByName } = require("../../functions"); +const { WebSocket, EventEmitter } = require("ws"); +const { default: axios } = require("axios"); +const events = new EventEmitter() +const WebsocketEvents = { + CONNECTED: "connected", + DISCONNECTED: "disconnected" +} module.exports = { name: 'Twitch/Connected', once: false, - execute(address, port) { + async execute(address, port) { console.log(`Connected to ${address}:${port}`) + const endpoint = `wss://eventsub.wss.twitch.tv/ws` + + const db = mClient.db('clients') + const coll = db.collection('credentials') + const credentials = await coll.findOne({ service: 'twitch' }) + + const headers = { + 'Client-Id': credentials.client_id, + 'Authorization': `Bearer ${credentials.token.access_token}`, + 'Content-Type': 'application/json', + }; + + var connection + var transport = {} + var subscribedEvents = {} + + const options = { + debug: true + } + function onConnect(data) { + transport.session_id = data.payload.session.id, + transport.conncted_at = data.payload.session.connected_at, + events.emit(WebsocketEvents.CONNECTED) + } + function onDisconnect(reason){ + console.log(`Disconnecting due to reason`, reason) + } + function connect() { + connection = new ReconnectingWebSocket(endpoint, [], { + WebSocket: WebSocket, + maxRetries: Infinity + }) + + return new Promise((resolve, reject) => { + connection.onclose = ({ reason }) => onDisconnect(reason) + connection.onmessage = ({ data }) => onMessage(data) + }) + } + function onMessage(data){ + const parsed = JSON.parse(data) + + if(options.debug){ + console.log(`[Debug] EvenSub Data:`, parsed) + } + + if (parsed.metadata.message_type === 'session_welcome'){ + return onConnect(parsed) + } + return onEventMessage(parsed) + } + function onEventMessage(data){ + const messageType = data.metadata.message_type + + switch (messageType) { + case 'notification': + const subscriptionType = data.metadata.subscription_type + return subscribedEvents[subscriptionType]?.(data.payload?.event) + default: + return subscribedEvents[messageType]?.(data.payload) + } + } + async function subscribe(type, condition, listener, version = 1){ + const sessionId = transport.session_id + const subscriptionPayload = { + type: type, + version: version, + condition: condition, + transport: { + method: 'websocket', + session_id: sessionId, + }, + }; + const res = await axios.post('https://api.twitch.tv/helix/eventsub/subscriptions', + subscriptionPayload, { + headers: headers + }) + console.log(`Subscribed to ${type}`) + subscribedEvents[type] = listener + return true + } + connect(options) + if (options.debug){ + events.on(WebsocketEvents.CONNECTED, () => { + console.log("Connected to EventSub"); + }) + + events.on(WebsocketEvents.DISCONNECTED, () => { + console.log("Disconnected from EventSub"); + }) + } + + const conditions = [{ + broadcaster_user_id: String(await getIDByName("desq_blocki")) + }] + + subscribe( + "stream.online", + conditions[0], + stream => { + console.log(`${stream.broadcaster_user_login} went offline`) + } + ) + + subscribe( + "stream.offline", + conditions[0], + stream => { + console.log(`${stream.broadcaster_user_login} went online`) + } + ) } } \ No newline at end of file diff --git a/events/Twitch/Disconnected.js b/events/Twitch/Disconnected.js new file mode 100644 index 0000000..b72e31d --- /dev/null +++ b/events/Twitch/Disconnected.js @@ -0,0 +1,49 @@ +// Event triggered by client connected to server + +const { RWSConnect, RWSDisconnect } = require("../../functions") + +module.exports = { + name: 'Twitch/Disconnected', + once: false, + async execute(reason) { + console.log('Twitch disconnected:', reason); + try { + RWSDisconnect() + } catch (error) { + console.error + } + + if (reason.includes('Login authentication failed')) { + try { + await refreshAccessToken(); + console.log('Reinitializing Twitch client after token refresh...'); + + try { + await tClient.disconnect(); + } catch (err) { + console.warn('Twitch client already disconnected:', err.message); + } + + const db = mClient.db("clients"); + const credentialCollection = db.collection('credentials'); + const tCreds = await credentialCollection.findOne({ service: 'twitch' }); + + // Recreate the client + tClient = new tmi.Client({ + options: { debug: false }, + connection: { reconnect: true, secure: true }, + identity: { + username: tCreds.username, + password: `oauth:${tCreds.token.access_token}`, + }, + channels: tClient.channels, + }); + + await tClient.connect(); + console.log('Twitch client reconnected successfully.'); + } catch (error) { + console.error('Failed to refresh token and reconnect:', error); + } + } + } +} diff --git a/events/Twitch/Message.js b/events/Twitch/Message.js index b2bacea..d759019 100644 --- a/events/Twitch/Message.js +++ b/events/Twitch/Message.js @@ -1,11 +1,11 @@ -const { mClient, dClient } = require('../..') +const { mClient, dClient, tClient } = require('../..') const { updateChatMode, getIDByName } = require('../../functions') require('dotenv').configDotenv module.exports = { name: 'Twitch/Message', once: false, - async execute(channel, userstate, message, self, tClient) { + async execute(channel, userstate, message, self) { const knownBots = new Set(['streamlabs', 'nightbot', 'moobot', 'soundalerts', 'streamelements', 'remasuri_bot', 'commanderroot', 'x__hel_bot__x']) async function emitShoutoutInfo(channel, username) { // Function sends information to the Shoutout Handler if they haven't been shouted out before, default target: VIPs diff --git a/events/Twitch/Shoutout.js b/events/Twitch/Shoutout.js index 90e9bec..37addf8 100644 --- a/events/Twitch/Shoutout.js +++ b/events/Twitch/Shoutout.js @@ -1,4 +1,5 @@ -const { getChannelInformation, getIDByName } = require("../../functions"); +const { tClient } = require("../.."); +const { getChannelInformation } = require("../../functions"); // Event triggered by Custom Shoutout emitted in Messages.js const soBuffer = new Map(); @@ -7,7 +8,7 @@ let isShoutoutInProgress = false; // Track shoutout status to avoid overlapping module.exports = { name: 'Twitch/Shoutout', once: false, - async execute(reason, channel, username, tClient) { + async execute(reason, channel, username) { // Function to perform shoutout (Nightbot or custom) async function doShoutouts(channel, user) { diff --git a/functions.js b/functions.js index 970f97d..d03c4ea 100644 --- a/functions.js +++ b/functions.js @@ -1,5 +1,8 @@ -const { default: axios } = require("axios"); +const axios = require("axios"); const { mClient } = require("."); +const ReconnectingWebSocket = require("reconnecting-websocket"); +const { WebSocket } = require("ws"); +const twitchEventSubWsUrl = 'wss://eventsub.wss.twitch.tv/ws'; // Delay function for pauses function delay(ms) { return new Promise(resolve => setTimeout(resolve, ms)); } @@ -161,7 +164,7 @@ async function getChannelInformation(streamer) { return await makeHelixRequest('GET', 'channels', { broadcaster_id }); } -async function getStreams(user_login){ +async function getStreams(user_login) { return await makeHelixRequest('GET', 'streams', { user_login }); } diff --git a/handlers/events.js b/handlers/events.js index 83f957b..bbaa44b 100644 --- a/handlers/events.js +++ b/handlers/events.js @@ -28,53 +28,21 @@ for (const folder of eventFolders) { // Handle Twitch events tClient.on('connected', (address, port) => { - dClient.emit('Twitch/Connected', address, port, tClient); + dClient.emit('Twitch/Connected', address, port); }); tClient.on('message', async (channel, userstate, message, self) => { - dClient.emit('Twitch/Message', channel, userstate, message, self, tClient); + dClient.emit('Twitch/Message', channel, userstate, message, self); }); tClient.on('raided', (channel, username, viewers) => { - dClient.emit('Twitch/Raided', channel, username, viewers, tClient); + dClient.emit('Twitch/Raided', channel, username, viewers); }); tClient.on('subscribers', (channel, enabled) => { - dClient.emit('Twitch/Subscribers', channel, enabled, tClient); + dClient.emit('Twitch/Subscribers', channel, enabled); }); tClient.on('disconnected', async (reason) => { - console.log('Twitch disconnected:', reason); - if (reason.includes('Login authentication failed')) { - try { - await refreshAccessToken(); - console.log('Reinitializing Twitch client after token refresh...'); - - try { - await tClient.disconnect(); - } catch (err) { - console.warn('Twitch client already disconnected:', err.message); - } - - const db = mClient.db("clients"); - const credentialCollection = db.collection('credentials'); - const tCreds = await credentialCollection.findOne({ service: 'twitch' }); - - // Recreate the client - tClient = new tmi.Client({ - options: { debug: false }, - connection: { reconnect: true, secure: true }, - identity: { - username: tCreds.username, - password: `oauth:${tCreds.token.access_token}`, - }, - channels: tClient.channels, - }); - - await tClient.connect(); - console.log('Twitch client reconnected successfully.'); - } catch (error) { - console.error('Failed to refresh token and reconnect:', error); - } - } + dClient.emit('Twitch/Disconnected', reason) }); \ No newline at end of file