fix WS member& bot status stuff for basic sharding

This commit is contained in:
Skillz
2020-07-16 20:41:25 -04:00
parent e4f683b4cf
commit 107ef59fd2
5 changed files with 163 additions and 4 deletions
+130
View File
@@ -17,6 +17,8 @@ import {
import { DiscordHeartbeatPayload } from "../types/discord.ts";
import { logRed } from "../utils/logger.ts";
import { handleDiscordPayload } from "./shardingManager.ts";
import { FetchMembersOptions } from "../types/guild.ts";
import { BotStatusRequest } from "../utils/utils.ts";
const basicShards = new Map<number, BasicShard>();
@@ -29,6 +31,16 @@ export interface BasicShard {
needToResume: boolean;
}
const RequestMembersQueue: RequestMemberQueuedRequest[] = [];
let processQueue = false;
interface RequestMemberQueuedRequest {
guildID: string;
shardID: number;
nonce: string;
options?: FetchMembersOptions;
}
export async function createBasicShard(
data: DiscordBotGatewayData,
identifyPayload: IdentifyPayload,
@@ -163,6 +175,8 @@ async function heartbeat(
interval: number,
) {
await delay(interval);
if (shard.socket.isClosed) return;
shard.socket.send(
JSON.stringify(
{ op: GatewayOpcode.Heartbeat, d: shard.previousSequenceNumber },
@@ -201,3 +215,119 @@ async function resumeConnection(
await delay(1000 * 15);
if (shard.needToResume) resumeConnection(botGatewayData, payload);
}
export function requestGuildMembers(
guildID: string,
shardID: number,
nonce: string,
options?: FetchMembersOptions,
queuedRequest = false,
) {
const shard = basicShards.get(shardID);
// This request was not from this queue so we add it to queue first
if (!queuedRequest) {
RequestMembersQueue.push({
guildID,
shardID,
nonce,
options,
});
if (!processQueue) {
processQueue = true;
processGatewayQueue();
}
return;
}
// If its closed add back to queue to redo on resume
if (shard?.socket.isClosed) {
requestGuildMembers(guildID, shardID, nonce, options);
return;
}
shard?.socket.send(JSON.stringify({
op: GatewayOpcode.RequestGuildMembers,
d: {
guild_id: guildID,
query: options?.query || "",
limit: options?.query || 0,
presences: options?.presences || false,
user_ids: options?.userIDs,
nonce,
},
}));
}
async function processGatewayQueue() {
if (!RequestMembersQueue.length) {
processQueue = false;
return;
}
basicShards.forEach((shard) => {
const index = RequestMembersQueue.findIndex((q) => q.shardID === shard.id);
// 2 events per second is the rate limit.
const request = RequestMembersQueue[index];
if (request) {
requestGuildMembers(
request.guildID,
request.shardID,
request.nonce,
request.options,
true,
);
// Remove item from queue
RequestMembersQueue.splice(index, 1);
const secondIndex = RequestMembersQueue.findIndex((q) =>
q.shardID === shard.id
);
const secondRequest = RequestMembersQueue[secondIndex];
if (secondRequest) {
requestGuildMembers(
request.guildID,
request.shardID,
secondRequest.nonce,
secondRequest.options,
true,
);
// Remove item from queue
RequestMembersQueue.splice(secondIndex, 1);
}
}
});
await delay(1500);
eventHandlers.debug?.(
{
type: "requestMembersProcessing",
data: {
remaining: RequestMembersQueue.length,
first: RequestMembersQueue[0],
},
},
);
processGatewayQueue();
}
export function botGatewayStatusRequest(payload: BotStatusRequest) {
basicShards.forEach((shard) => {
shard.socket.send(JSON.stringify({
op: GatewayOpcode.StatusUpdate,
d: {
since: null,
game: payload.game.name
? {
name: payload.game.name,
type: payload.game.type,
}
: null,
status: payload.status,
afk: false,
},
}));
});
}
+1 -1
View File
@@ -49,7 +49,7 @@ export const createClient = async (data: ClientOptions) => {
(bits, next) => (bits |= next),
0,
);
identifyPayload.shard = [0, botGatewayData.shards]
identifyPayload.shard = [0, botGatewayData.shards];
spawnShards(botGatewayData, identifyPayload);
};
+23 -3
View File
@@ -55,9 +55,15 @@ import {
} from "../types/message.ts";
import { createMessage } from "../structures/message.ts";
import { GuildUpdateChange } from "../types/options.ts";
import { createBasicShard } from "./basicShard.ts";
import {
createBasicShard,
requestGuildMembers,
botGatewayStatusRequest,
} from "./basicShard.ts";
import { BotStatusRequest } from "../utils/utils.ts";
let shardCounter = 0;
let basicSharding = false;
export interface FetchAllMembersRequest {
resolve: Function;
@@ -114,6 +120,7 @@ export const spawnShards = async (
createNextShard = false;
if (data.shards >= 25) createShardWorker();
else {
basicSharding = true;
createBasicShard(data, payload, false, id - 1);
}
spawnShards(data, payload, id + 1);
@@ -136,9 +143,11 @@ export async function handleDiscordPayload(
return eventHandlers.heartbeat?.();
case GatewayOpcode.Dispatch:
if (data.t === "READY") {
setBotID((data.d as ReadyPayload).user.id);
const payload = data.d as ReadyPayload;
setBotID(payload.user.id);
// Triggered on each shard
eventHandlers.ready?.();
eventHandlers.shardReady?.(shardID);
if (payload.shard && shardID === payload.shard[1] - 1) eventHandlers.ready?.()
// Wait 5 seconds to spawn next shard
await delay(5000);
createNextShard = true;
@@ -643,6 +652,10 @@ export async function requestAllMembers(
const nonce = Math.random().toString();
fetchAllMembersProcessingRequests.set(nonce, resolve);
if (basicSharding) {
return requestGuildMembers(guild.id, guild.shardID, nonce, options);
}
shards[guild.shardID].postMessage({
type: "FETCH_MEMBERS",
guildID: guild.id,
@@ -652,6 +665,13 @@ export async function requestAllMembers(
}
export function sendGatewayCommand(type: "EDIT_BOTS_STATUS", payload: object) {
if (basicSharding) {
if (type === "EDIT_BOTS_STATUS") {
botGatewayStatusRequest(payload as BotStatusRequest);
}
return;
}
shards.forEach((shard) => {
shard.postMessage({
type,
+1
View File
@@ -125,6 +125,7 @@ export interface EventHandlers {
roleUpdate?: (guild: Guild, role: Role, cachedRole: Role) => unknown;
roleGained?: (guild: Guild, member: Member, roleID: string) => unknown;
roleLost?: (guild: Guild, member: Member, roleID: string) => unknown;
shardReady?: (shardID: number) => unknown;
typingStart?: (data: TypingStartPayload) => unknown;
voiceChannelJoin?: (member: Member, channelID: string) => unknown;
voiceChannelLeave?: (member: Member, channelID: string) => unknown;
+8
View File
@@ -7,6 +7,14 @@ export const sleep = (timeout: number) => {
return new Promise((resolve) => setTimeout(resolve, timeout));
};
export interface BotStatusRequest {
status: StatusType;
game: {
name?: string;
type: ActivityType;
};
}
export function editBotsStatus(
status: StatusType,
name?: string,