fix: use workers in big bot template with resharding

This commit is contained in:
Skillz4Killz
2022-03-30 17:38:21 +00:00
committed by GitHub
parent 4c9522b83f
commit 5a7a8c2c90
2 changed files with 388 additions and 51 deletions
+134 -51
View File
@@ -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<number, Worker>();
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);
+254
View File
@@ -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<GatewayManager>) {
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<string>) {
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<GatewayManager>;
}
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,
}),
);
}
};