From 5a7a8c2c9040f72a63b37bee994634ae3768f619 Mon Sep 17 00:00:00 2001 From: Skillz4Killz <23035000+Skillz4Killz@users.noreply.github.com> Date: Wed, 30 Mar 2022 17:38:21 +0000 Subject: [PATCH] fix: use workers in big bot template with resharding --- template/bigbot/src/gateway/mod.ts | 185 +++++++++++++------ template/bigbot/src/gateway/worker.ts | 254 ++++++++++++++++++++++++++ 2 files changed, 388 insertions(+), 51 deletions(-) create mode 100644 template/bigbot/src/gateway/worker.ts diff --git a/template/bigbot/src/gateway/mod.ts b/template/bigbot/src/gateway/mod.ts index f6f61bbaf..a96e6e546 100644 --- a/template/bigbot/src/gateway/mod.ts +++ b/template/bigbot/src/gateway/mod.ts @@ -1,13 +1,10 @@ +import { Collection, createGatewayManager, createRestManager, endpoints } from "../../deps.ts"; import { DISCORD_TOKEN, - EVENT_HANDLER_PORT, EVENT_HANDLER_SECRET_KEY, - EVENT_HANDLER_URL, - GATEWAY_INTENTS, REST_AUTHORIZATION_KEY, REST_PORT, } from "../../configs.ts"; -import { createGatewayManager, createRestManager, endpoints } from "../../deps.ts"; // CREATE A SIMPLE MANAGER FOR REST const rest = createRestManager({ @@ -16,57 +13,143 @@ const rest = createRestManager({ customUrl: `http://localhost:${REST_PORT}`, }); -// CALL THE REST PROCESS TO GET GATEWAY DATA -const result = await rest.runMethod(rest, "get", endpoints.GATEWAY_BOT()).then(( - res, -) => ({ - url: res.url, - shards: res.shards, - sessionStartLimit: { - total: res.session_start_limit.total, - remaining: res.session_start_limit.remaining, - resetAfter: res.session_start_limit.reset_after, - maxConcurrency: res.session_start_limit.max_concurrency, - }, -})); - const gateway = createGatewayManager({ - // FOR DEBUGGING - // debug: console.log, - // THE AUTHORIZATION WE WILL USE ON OUR EVENT HANDLER PROCESS secretKey: EVENT_HANDLER_SECRET_KEY, token: DISCORD_TOKEN, - intents: GATEWAY_INTENTS, - // LOAD DATA FROM DISCORDS RECOMMENDATIONS OR YOUR OWN CUSTOM ONES HERE - shardsRecommended: result.shards, - sessionStartLimitTotal: result.sessionStartLimit.total, - sessionStartLimitRemaining: result.sessionStartLimit.remaining, - sessionStartLimitResetAfter: result.sessionStartLimit.resetAfter, - maxConcurrency: result.sessionStartLimit.maxConcurrency, - maxShards: result.shards, - lastShardId: result.shards, - + intents: ["GuildMessages", "Guilds"], // THIS WILL BASICALLY BE YOUR HANDLER FOR YOUR EVENTS. - handleDiscordPayload: async function (_, data, shardId) { - // TODO: CHANGE FROM SENDING THROUGH HTTP TO USING A WS FOR FASTER PROCESSING! OR HTTP3 OR WHATEVER! - if (!data.t) return; - - await fetch(`${EVENT_HANDLER_URL}:${EVENT_HANDLER_PORT}`, { - headers: { - Authorization: gateway.secretKey, - }, - method: "POST", - body: JSON.stringify({ - shardId, - data, - }), - }) - // BELOW IS FOR DENO MEMORY LEAK - .then((res) => res.text()) - .catch(() => null); - }, + handleDiscordPayload: async function (_, data, shardId) {}, }); -// START THE GATEWAY -gateway.spawnShards(gateway); +const workers = new Collection(); + +async function startGateway() { + // CALL THE REST PROCESS TO GET GATEWAY DATA + const result = await rest.runMethod(rest, "get", endpoints.GATEWAY_BOT()).then((res) => ({ + url: res.url, + shards: res.shards, + sessionStartLimit: { + total: res.session_start_limit.total, + remaining: res.session_start_limit.remaining, + resetAfter: res.session_start_limit.reset_after, + maxConcurrency: res.session_start_limit.max_concurrency, + }, + })); + + // LOAD DATA FROM DISCORDS RECOMMENDATIONS OR YOUR OWN CUSTOM ONES HERE + gateway.shardsRecommended = result.shards; + gateway.sessionStartLimitTotal = result.sessionStartLimit.total; + gateway.sessionStartLimitRemaining = result.sessionStartLimit.remaining; + gateway.sessionStartLimitResetAfter = result.sessionStartLimit.resetAfter; + gateway.maxConcurrency = result.sessionStartLimit.maxConcurrency; + gateway.maxShards = result.shards + gateway.lastShardId = result.shards + + // PREPARE BUCKETS FOR IDENTIFYING + gateway.prepareBuckets(gateway, 0, result.shards); + + function startWorker(workerId: number, bucketId: number, firstShardId: number, lastShardId: number) { + const worker = workers.get(workerId); + if (!worker) return; + + // TRIGGER IDENTIFY IN WORKER + worker.postMessage( + JSON.stringify({ + type: "IDENTIFY", + shardId: firstShardId, + shardsRecommended: result.shards, + sessionStartLimitTotal: result.sessionStartLimit.total, + sessionStartLimitRemaining: result.sessionStartLimit.remaining, + sessionStartLimitResetAfter: result.sessionStartLimit.resetAfter, + maxConcurrency: result.sessionStartLimit.maxConcurrency, + maxShards: gateway.maxShards, + lastShardId: lastShardId, + workerId, + }), + ); + } + + gateway.buckets.forEach((bucket, bucketId) => { + for (let i = 0; i < bucket.workers.length; i++) { + const workerId = bucket.workers[i][0]; + const worker = new Worker(new URL("./worker.js", import.meta.url).href, { + name: `w-${workerId}-b${bucketId}`, + type: "module", + }); + workers.set(workerId, worker); + + if (bucket.workers[i + 1]) { + worker.onmessage = function (message) { + const data = JSON.parse(message.data); + if (data.type === "ALL_SHARDS_READY") { + const queue = bucket.workers[i + 1]; + if (queue) startWorker(queue[0], bucketId, queue[1], queue[queue.length - 1]); + } + + if (data.type === "RESHARDED") { + const nextWorker = workers.get(workerId + 1); + if (nextWorker) { + nextWorker.postMessage( + JSON.stringify({ + type: "RESHARD", + results: data.results, + }), + ); + } + } + }; + } else { + // THIS IS FINAL WORKER + worker.onmessage = function (message) { + const data = JSON.parse(message.data); + if (data.type === "RESHARDED") { + // THERE IS NO NEXT WORKER SO TELL ALL WORKERS TO CLOSE OLD GATEWAYS + workers.forEach((workerx) => { + workerx.postMessage( + JSON.stringify({ + type: "RESHARDED-CLOSEOLD", + }), + ); + }); + } + }; + } + } + + const queue = bucket.workers[0]; + startWorker(queue[0], bucketId, queue[1], queue[queue.length - 1]); + }); +} + +startGateway(); + + +setInterval(async () => { + console.log("GW DEBUG", "[Resharding] Checking if resharding is needed."); + + const results = await rest.runMethod(rest, "get", endpoints.GATEWAY_BOT()).then((res) => ({ + url: res.url, + shards: res.shards, + sessionStartLimit: { + total: res.session_start_limit.total, + remaining: res.session_start_limit.remaining, + resetAfter: res.session_start_limit.reset_after, + maxConcurrency: res.session_start_limit.max_concurrency, + }, + })); + const percentage = ((results.shards - gateway.maxShards) / gateway.maxShards) * 100; + // Less than necessary% being used so do nothing + if (percentage < gateway.reshardPercentage) return; + + // Don't have enough identify rate limits to reshard + if (results.sessionStartLimit.remaining < results.shards) return; + + workers.first()?.postMessage( + JSON.stringify({ + type: "RESHARD", + results, + }), + ); + // DAILY +}, 1000 * 60 * 60 * 24); diff --git a/template/bigbot/src/gateway/worker.ts b/template/bigbot/src/gateway/worker.ts new file mode 100644 index 000000000..fe3999da9 --- /dev/null +++ b/template/bigbot/src/gateway/worker.ts @@ -0,0 +1,254 @@ +import { + BOT_ID, + DISCORD_TOKEN, + EVENT_HANDLER_PORT, + EVENT_HANDLER_SECRET_KEY, + EVENT_HANDLER_URL, +} from "../../configs.ts"; +import { Collection, createGatewayManager, DiscordReady, GatewayManager, GetGatewayBot } from "../../deps.ts"; + +let gateway: GatewayManager; +// FOR RESHARDED +let gatewayPendingClosing: GatewayManager; +let workerId: number; + +function spawnGateway(shardId: number, options: Partial) { + console.log(`[Worker #${workerId}]`, "[Worker] Spawning the worker gateway.", shardId, options); + gateway = createGatewayManager({ + // LOAD DATA FROM DISCORDS RECOMMENDATIONS OR YOUR OWN CUSTOM ONES HERE + shardsRecommended: options.shardsRecommended, + sessionStartLimitTotal: options.sessionStartLimitTotal, + sessionStartLimitRemaining: options.sessionStartLimitRemaining, + sessionStartLimitResetAfter: options.sessionStartLimitResetAfter, + maxConcurrency: options.maxConcurrency, + maxShards: options.maxShards, + // SET STARTING SHARD ID + firstShardId: shardId, + // SET LAST SHARD ID + lastShardId: options.lastShardId ?? shardId, + // THE AUTHORIZATION WE WILL USE ON OUR EVENT HANDLER PROCESS + secretKey: EVENT_HANDLER_SECRET_KEY, + token: DISCORD_TOKEN, + intents: ["GuildMessages", "Guilds", "GuildMembers"], + handleDiscordPayload: async function (_, data, shardId) { + // TRIGGER RAW EVENT + if (!data.t) return; + + const id = (data.t && ["GUILD_CREATE", "GUILD_DELETE", "GUILD_UPDATE"].includes(data.t) + ? (data.d as any)?.id + : (data.d as any)?.guild_id) ?? "000000000000000000"; + + // IF FINAL SHARD BECAME READY TRIGGER NEXT WORKER + if (data.t === "READY") { + console.log(`[Worker #${workerId}]`, `[Worker] Shard #${shardId} online`); + + if (shardId === gateway.lastShardId) { + // @ts-ignore + postMessage( + JSON.stringify({ + type: "ALL_SHARDS_READY", + }), + ); + } + } + + // DONT SEND THESE EVENTS USELESS TO BOT + if (["GUILD_LOADED_DD"].includes(data.t)) return; + + await fetch(`${EVENT_HANDLER_URL}:${EVENT_HANDLER_PORT}`, { + headers: { + Authorization: gateway.secretKey, + "Content-Type": "application/json", + }, + method: "POST", + body: JSON.stringify({ + shardId, + data, + }), + }) + // BELOW IS FOR DENO MEMORY LEAK + .then((res) => res.text()) + .catch(() => null); + }, + }); + + // START THE GATEWAY + gateway.spawnShards(gateway, shardId); + + return gateway; +} + +interface IdentifyPayload { + type: "IDENTIFY"; + shardId: number; + shards: number; + sessionStartLimit: { + total: number; + remaining: number; + resetAfter: number; + maxConcurrency: number; + }; + shardsRecommended: number; + sessionStartLimitTotal: number; + sessionStartLimitRemaining: number; + sessionStartLimitResetAfter: number; + maxConcurrency: number; + maxShards: number; + lastShardId: number; + workerId: number; +} + +interface ReshardPayload { + type: "RESHARD"; + results: GetGatewayBot; +} + +interface FullyReshardedPayload { + type: "RESHARDED-CLOSEOLD"; +} + +// @ts-ignore this should not be erroring +self.onmessage = async function (message: MessageEvent) { + const data = JSON.parse(message.data) as IdentifyPayload | ReshardPayload | FullyReshardedPayload; + + if (data.type === "IDENTIFY") { + workerId = data.workerId; + + gateway = spawnGateway(data.shardId, { + shardsRecommended: data.shardsRecommended, + sessionStartLimitTotal: data.sessionStartLimitTotal, + sessionStartLimitRemaining: data.sessionStartLimitRemaining, + sessionStartLimitResetAfter: data.sessionStartLimitResetAfter, + maxConcurrency: data.maxConcurrency, + maxShards: data.maxShards, + lastShardId: data.lastShardId, + spawnShardDelay: 5000, + }); + } + + if (data.type === "RESHARDED-CLOSEOLD") { + console.log(`[Worker #${workerId}]`, "[Resharding] Closing old gateways."); + await gateway.resharding.closeOldShards(gatewayPendingClosing); + } + + if (data.type === "RESHARD") { + console.log(`[Worker #${workerId}]`, "[Worker] Resharding the worker."); + gateway.resharding.isPending = async function (gateway: GatewayManager) { + for (let i = gateway.firstShardId; i < gateway.lastShardId; i++) { + const shard = gateway.shards.get(i); + if (!shard?.ready) { + return true; + } + } + + return false; + }; + + async function processResharding(oldGateway: GatewayManager, results: GetGatewayBot) { + oldGateway.debug("GW DEBUG", "[Resharding] Starting the reshard process."); + + const gateway = createGatewayManager({ + ...oldGateway, + // RESET THE SETS AND COLLECTIONS + cache: { + guildIds: new Set(), + loadingGuildIds: new Set(), + editedMessages: new Collection(), + }, + shards: new Collection(), + loadingShards: new Collection(), + buckets: new Collection(), + utf8decoder: new TextDecoder(), + }); + + for (const [key, value] of Object.entries(oldGateway)) { + if (key === "handleDiscordPayload") { + gateway.handleDiscordPayload = async function (_, data, shardId) { + if (data.t === "READY") { + const payload = data.d as DiscordReady; + console.log(`[Worker - ${workerId}] Shard #${payload.shard?.[0]} online`); + if (shardId === gateway.lastShardId) { + // @ts-ignore + postMessage( + JSON.stringify({ + type: "RESHARDED", + results, + }), + ); + } + + await gateway.resharding.markNewGuildShardId( + payload.guilds.map((g) => BigInt(g.id)), + shardId, + ); + } + }; + continue; + } + + // DON"T OVERRIDE THESE + if (["cache", "shards", "loadingShards", "buckets", "utf8decoder"].includes(key)) continue; + + // USE ANY CUSTOMIZED OPTIONS FROM OLD GATEWAY + // @ts-ignore silly ts error + gateway[key] = oldGateway[key as keyof typeof oldGateway]; + } + + // Begin resharding + // If more than 100K servers, begin switching to 16x sharding + if (gateway.useOptimalLargeBotSharding) { + console.log(`[Worker - ${workerId}]`, "[Resharding] Using optimal large bot sharding solution."); + gateway.maxShards = gateway.calculateMaxShards(results.shards, results.sessionStartLimit.maxConcurrency); + } else { + gateway.maxShards = results.shards; + } + + // FOR MANUAL SHARD CONTROL, OVERRIDE THIS SHARD ID! + gateway.lastShardId = oldGateway.lastShardId === oldGateway.maxShards - 1 + ? gateway.maxShards - 1 + : oldGateway.lastShardId; + gateway.shardsRecommended = results.shards; + gateway.sessionStartLimitTotal = results.sessionStartLimit.total; + gateway.sessionStartLimitRemaining = results.sessionStartLimit.remaining; + gateway.sessionStartLimitResetAfter = results.sessionStartLimit.resetAfter; + gateway.maxConcurrency = results.sessionStartLimit.maxConcurrency; + + gateway.spawnShards(gateway, gateway.firstShardId); + + return new Promise((resolve) => { + // TIMER TO KEEP CHECKING WHEN ALL SHARDS HAVE RESHARDED + const timer = setInterval(async () => { + const pending = await gateway.resharding.isPending(gateway, oldGateway); + // STILL PENDING ON SOME SHARDS TO BE CREATED + if (pending) return; + + // ENABLE EVENTS ON NEW SHARDS AND IGNORE EVENTS ON OLD + const oldHandler = oldGateway.handleDiscordPayload; + gateway.handleDiscordPayload = oldHandler; + oldGateway.handleDiscordPayload = function (og, data, shardId) { + // ALLOW EXCEPTION FOR CHUNKING TO PREVENT REQUESTS FREEZING + if (data.t !== "GUILD_MEMBERS_CHUNK") return; + oldHandler(og, data, shardId); + }; + + // STOP TIMER + clearInterval(timer); + await gateway.resharding.editGuildShardIds(); + gatewayPendingClosing = oldGateway; + gateway.debug("GW DEBUG", "[Resharding] Complete."); + resolve(gateway); + }, 30000); + }) as Promise; + } + + gateway = await processResharding(gateway, data.results); + console.log(`[Worker - ${workerId}] Resharded the worker.`); + // @ts-ignore this should not be erroring + postMessage( + JSON.stringify({ + type: "RESHARDED", + results: data.results, + }), + ); + } +};