From f3dab47085f3ef133a8128634dab9e4cbec12723 Mon Sep 17 00:00:00 2001 From: Skillz Date: Wed, 19 Aug 2020 20:56:06 -0400 Subject: [PATCH] settle them promises --- src/module/requestManager.ts | 107 ++++++++++++++++++++++------------- 1 file changed, 69 insertions(+), 38 deletions(-) diff --git a/src/module/requestManager.ts b/src/module/requestManager.ts index da35c8531..862e85d76 100644 --- a/src/module/requestManager.ts +++ b/src/module/requestManager.ts @@ -4,8 +4,9 @@ import { delay } from "https://deno.land/std@0.61.0/async/delay.ts"; import { Errors } from "../types/errors.ts"; import { HttpResponseCode } from "../types/discord.ts"; import { logRed } from "../utils/logger.ts"; +import { baseEndpoints } from "../constants/discord.ts"; -const queue: QueuedRequest[] = []; +const pathQueues: { [key: string]: QueuedRequest[] } = {}; const ratelimitedPaths = new Map(); let globallyRateLimited = false; let queueInProcess = false; @@ -40,49 +41,79 @@ async function processRateLimitedPaths() { processRateLimitedPaths(); } -async function processQueue() { - if (queue.length && !globallyRateLimited) { - const request = queue.shift(); - if (!request) return; +function addToQueue(request: QueuedRequest) { + const route = request.url.substring(baseEndpoints.BASE_URL.length + 1); + const parts = route.split("/"); + // Remove the major param + parts.shift(); + const [id] = parts; - const rateLimitedURLResetIn = await checkRatelimits(request.url); + if (pathQueues[id]) { + pathQueues[id].push(request); + } else { + pathQueues[id] = [request]; + } +} - if (request.bucketID) { - const rateLimitResetIn = await checkRatelimits(request.bucketID); - if (rateLimitResetIn) { - // This request is still rate limited readd to queue - queue.push(request); - } else if (rateLimitedURLResetIn) { - // This URL is rate limited readd to queue - queue.push(request); - } else { - // This request is not rate limited so it should be run - const result = await request.callback(); - if (result && result.rateLimited) { - queue.push( - { ...request, bucketID: result.bucketID || request.bucketID }, - ); - } - } - } else { - if (rateLimitedURLResetIn) { - // This URL is rate limited readd to queue - queue.push(request); - } else { - // This request has no bucket id so it should be processed - const result = await request.callback(); - if (request && result && result.rateLimited) { - queue.push( - { ...request, bucketID: result.bucketID || request.bucketID }, - ); - } - } +async function cleanupQueues() { + Object.entries(pathQueues).map(([key, value]) => { + if (!value.length) { + // Remove it entirely + delete pathQueues[key]; } + }); +} + +async function processQueue() { + if ( + (Object.keys(pathQueues).length) && !globallyRateLimited + ) { + await Promise.allSettled( + Object.values(pathQueues).map(async (pathQueue) => { + const request = pathQueue.shift(); + if (!request) return; + + const rateLimitedURLResetIn = await checkRatelimits(request.url); + + if (request.bucketID) { + const rateLimitResetIn = await checkRatelimits(request.bucketID); + if (rateLimitResetIn) { + // This request is still rate limited readd to queue + addToQueue(request); + } else if (rateLimitedURLResetIn) { + // This URL is rate limited readd to queue + addToQueue(request); + } else { + // This request is not rate limited so it should be run + const result = await request.callback(); + if (result && result.rateLimited) { + addToQueue( + { ...request, bucketID: result.bucketID || request.bucketID }, + ); + } + } + } else { + if (rateLimitedURLResetIn) { + // This URL is rate limited readd to queue + addToQueue(request); + } else { + // This request has no bucket id so it should be processed + const result = await request.callback(); + if (request && result && result.rateLimited) { + addToQueue( + { ...request, bucketID: result.bucketID || request.bucketID }, + ); + } + } + } + }), + ); } - if (queue.length) { + if (Object.keys(pathQueues).length) { await delay(1000); processQueue(); + cleanupQueues() } else queueInProcess = false; } @@ -231,7 +262,7 @@ async function runMethod( } }; - queue.push({ + addToQueue({ callback, bucketID, url,