mirror of
https://github.com/discordeno/discordeno.git
synced 2026-09-17 08:47:22 +00:00
formatter: Use semicolons (#4686)
I prefer semicolors, they also help avoiding certain pitfalls in JavaScript/TypeScript, such as the following code sample: ```js const xyz = "test" (something.else as string) = "another" ``` This results in a TypeError: "test" is not a function, this is because js thinks we are trying to call the string "test" as a function. To fix this it requires a `;` somewhere before the `(`, such as `;(something ... ` which in my opinion is ugly and less clean overall.
This commit is contained in:
+309
-309
File diff suppressed because it is too large
Load Diff
@@ -1,3 +1,3 @@
|
||||
export * from './manager.js'
|
||||
export * from './Shard.js'
|
||||
export * from './types.js'
|
||||
export * from './manager.js';
|
||||
export * from './Shard.js';
|
||||
export * from './types.js';
|
||||
|
||||
+212
-208
@@ -1,4 +1,4 @@
|
||||
import { randomBytes } from 'node:crypto'
|
||||
import { randomBytes } from 'node:crypto';
|
||||
import {
|
||||
type AtLeastOne,
|
||||
type BigString,
|
||||
@@ -10,10 +10,10 @@ import {
|
||||
GatewayIntents,
|
||||
GatewayOpcodes,
|
||||
type RequestGuildMembers,
|
||||
} from '@discordeno/types'
|
||||
import { Collection, jsonSafeReplacer, LeakyBucket, logger } from '@discordeno/utils'
|
||||
import Shard from './Shard.js'
|
||||
import { type ShardEvents, ShardSocketCloseCodes, type ShardSocketRequest, type TransportCompression, type UpdateVoiceState } from './types.js'
|
||||
} from '@discordeno/types';
|
||||
import { Collection, jsonSafeReplacer, LeakyBucket, logger } from '@discordeno/utils';
|
||||
import Shard from './Shard.js';
|
||||
import { type ShardEvents, ShardSocketCloseCodes, type ShardSocketRequest, type TransportCompression, type UpdateVoiceState } from './types.js';
|
||||
|
||||
export function createGatewayManager(options: CreateGatewayManagerOptions): GatewayManager {
|
||||
const connectionOptions = options.connection ?? {
|
||||
@@ -25,7 +25,7 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
total: 1000,
|
||||
resetAfter: 1000 * 60 * 60 * 24,
|
||||
},
|
||||
}
|
||||
};
|
||||
|
||||
const gateway: GatewayManager = {
|
||||
events: options.events ?? {},
|
||||
@@ -66,87 +66,87 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
getSessionInfo: options.resharding?.getSessionInfo,
|
||||
updateGuildsShardId: options.resharding?.updateGuildsShardId,
|
||||
async checkIfReshardingIsNeeded() {
|
||||
gateway.logger.debug('[Resharding] Checking if resharding is needed.')
|
||||
gateway.logger.debug('[Resharding] Checking if resharding is needed.');
|
||||
|
||||
if (!gateway.resharding.enabled) {
|
||||
gateway.logger.debug('[Resharding] Resharding is disabled.')
|
||||
gateway.logger.debug('[Resharding] Resharding is disabled.');
|
||||
|
||||
return { needed: false }
|
||||
return { needed: false };
|
||||
}
|
||||
|
||||
if (!gateway.resharding.getSessionInfo) {
|
||||
throw new Error("[Resharding] Resharding is enabled but no 'resharding.getSessionInfo()' is not provided.")
|
||||
throw new Error("[Resharding] Resharding is enabled but no 'resharding.getSessionInfo()' is not provided.");
|
||||
}
|
||||
|
||||
gateway.logger.debug('[Resharding] Resharding is enabled.')
|
||||
gateway.logger.debug('[Resharding] Resharding is enabled.');
|
||||
|
||||
const sessionInfo = await gateway.resharding.getSessionInfo()
|
||||
const sessionInfo = await gateway.resharding.getSessionInfo();
|
||||
|
||||
gateway.logger.debug(`[Resharding] Session info retrieved: ${JSON.stringify(sessionInfo)}`)
|
||||
gateway.logger.debug(`[Resharding] Session info retrieved: ${JSON.stringify(sessionInfo)}`);
|
||||
|
||||
// Don't have enough identify limits to try resharding
|
||||
if (sessionInfo.sessionStartLimit.remaining < sessionInfo.shards) {
|
||||
gateway.logger.debug('[Resharding] Not enough session start limits left to reshard.')
|
||||
gateway.logger.debug('[Resharding] Not enough session start limits left to reshard.');
|
||||
|
||||
return { needed: false, info: sessionInfo }
|
||||
return { needed: false, info: sessionInfo };
|
||||
}
|
||||
|
||||
gateway.logger.debug('[Resharding] Able to reshard, checking whether necessary now.')
|
||||
gateway.logger.debug('[Resharding] Able to reshard, checking whether necessary now.');
|
||||
|
||||
// 2500 is the max amount of guilds a single shard can handle
|
||||
// 1000 is the amount of guilds discord uses to determine how many shards to recommend.
|
||||
// This algo helps check if your bot has grown enough to reshard.
|
||||
// While this is imprecise as discord changes the recommended number of shard every 1000 guilds it is good enough
|
||||
// The alternative is to store the guild count for each shard and require the Guilds intent for `GUILD_CREATE` and `GUILD_DELETE` events
|
||||
const percentage = (sessionInfo.shards / ((gateway.totalShards * 2500) / 1000)) * 100
|
||||
const percentage = (sessionInfo.shards / ((gateway.totalShards * 2500) / 1000)) * 100;
|
||||
|
||||
// Less than necessary% being used so do nothing
|
||||
if (percentage < gateway.resharding.shardsFullPercentage) {
|
||||
gateway.logger.debug('[Resharding] Resharding not needed.')
|
||||
gateway.logger.debug('[Resharding] Resharding not needed.');
|
||||
|
||||
return { needed: false, info: sessionInfo }
|
||||
return { needed: false, info: sessionInfo };
|
||||
}
|
||||
|
||||
gateway.logger.info('[Resharding] Resharding is needed.')
|
||||
gateway.logger.info('[Resharding] Resharding is needed.');
|
||||
|
||||
return { needed: true, info: sessionInfo }
|
||||
return { needed: true, info: sessionInfo };
|
||||
},
|
||||
async reshard(info) {
|
||||
gateway.logger.info(`[Resharding] Starting the reshard process. Previous total shards: ${gateway.totalShards}`)
|
||||
gateway.logger.info(`[Resharding] Starting the reshard process. Previous total shards: ${gateway.totalShards}`);
|
||||
// Set values on gateway
|
||||
gateway.totalShards = info.shards
|
||||
gateway.totalShards = info.shards;
|
||||
// Handles preparing mid sized bots for LBS
|
||||
gateway.totalShards = gateway.calculateTotalShards()
|
||||
gateway.totalShards = gateway.calculateTotalShards();
|
||||
// Set first shard id if provided in info
|
||||
if (typeof info.firstShardId === 'number') gateway.firstShardId = info.firstShardId
|
||||
if (typeof info.firstShardId === 'number') gateway.firstShardId = info.firstShardId;
|
||||
// Set last shard id if provided in info
|
||||
if (typeof info.lastShardId === 'number') gateway.lastShardId = info.lastShardId
|
||||
if (typeof info.lastShardId === 'number') gateway.lastShardId = info.lastShardId;
|
||||
// If we didn't get any lastShardId, we assume all the shards are to be used
|
||||
else gateway.lastShardId = gateway.totalShards - 1
|
||||
gateway.logger.info(`[Resharding] Starting the reshard process. New total shards: ${gateway.totalShards}`)
|
||||
else gateway.lastShardId = gateway.totalShards - 1;
|
||||
gateway.logger.info(`[Resharding] Starting the reshard process. New total shards: ${gateway.totalShards}`);
|
||||
|
||||
// Resetting buckets
|
||||
gateway.buckets.clear()
|
||||
gateway.buckets.clear();
|
||||
// Refilling buckets with new values
|
||||
gateway.prepareBuckets()
|
||||
gateway.prepareBuckets();
|
||||
|
||||
// Call all the buckets and tell their workers & shards to identify
|
||||
const promises = Array.from(gateway.buckets.entries()).map(async ([bucketId, bucket]) => {
|
||||
for (const worker of bucket.workers) {
|
||||
for (const shardId of worker.queue) {
|
||||
await gateway.resharding.tellWorkerToPrepare(worker.id, shardId, bucketId)
|
||||
await gateway.resharding.tellWorkerToPrepare(worker.id, shardId, bucketId);
|
||||
}
|
||||
}
|
||||
})
|
||||
});
|
||||
|
||||
await Promise.all(promises)
|
||||
await Promise.all(promises);
|
||||
|
||||
gateway.logger.info(`[Resharding] All shards are now online.`)
|
||||
gateway.logger.info(`[Resharding] All shards are now online.`);
|
||||
|
||||
await gateway.resharding.onReshardingSwitch()
|
||||
await gateway.resharding.onReshardingSwitch();
|
||||
},
|
||||
async tellWorkerToPrepare(workerId, shardId, bucketId) {
|
||||
gateway.logger.debug(`[Resharding] Telling worker to prepare. Worker: ${workerId} | Shard: ${shardId} | Bucket: ${bucketId}.`)
|
||||
gateway.logger.debug(`[Resharding] Telling worker to prepare. Worker: ${workerId} | Shard: ${shardId} | Bucket: ${bucketId}.`);
|
||||
const shard = new Shard({
|
||||
id: shardId,
|
||||
connection: {
|
||||
@@ -166,32 +166,32 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
await gateway.resharding.updateGuildsShardId?.(
|
||||
(payload.d as DiscordReady).guilds.map((g) => g.id),
|
||||
shardId,
|
||||
)
|
||||
);
|
||||
}
|
||||
},
|
||||
},
|
||||
logger: gateway.logger,
|
||||
requestIdentify: async () => await gateway.requestIdentify(shardId),
|
||||
makePresence: gateway.makePresence,
|
||||
})
|
||||
});
|
||||
|
||||
gateway.resharding.shards.set(shardId, shard)
|
||||
gateway.resharding.shards.set(shardId, shard);
|
||||
|
||||
await shard.identify()
|
||||
await shard.identify();
|
||||
|
||||
gateway.logger.debug(`[Resharding] Shard #${shardId} identified.`)
|
||||
gateway.logger.debug(`[Resharding] Shard #${shardId} identified.`);
|
||||
},
|
||||
async onReshardingSwitch() {
|
||||
gateway.logger.debug(`[Resharding] Making the switch from the old shards to the new ones.`)
|
||||
gateway.logger.debug(`[Resharding] Making the switch from the old shards to the new ones.`);
|
||||
|
||||
// Move the events from the old shards to the new ones
|
||||
for (const shard of gateway.resharding.shards.values()) {
|
||||
shard.events = options.events ?? {}
|
||||
shard.events = options.events ?? {};
|
||||
}
|
||||
|
||||
// Old shards stop processing events
|
||||
for (const shard of gateway.shards.values()) {
|
||||
const oldHandler = shard.events.message
|
||||
const oldHandler = shard.events.message;
|
||||
|
||||
// Change with spread operator to not affect new shards, as changing anything on shard.events will directly change options.events, which changes new shards' events
|
||||
shard.events = {
|
||||
@@ -199,29 +199,29 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
message: async function (_, message) {
|
||||
// Member checks need to continue but others can stop
|
||||
if (message.t === 'GUILD_MEMBERS_CHUNK') {
|
||||
oldHandler?.(shard, message)
|
||||
oldHandler?.(shard, message);
|
||||
}
|
||||
},
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
gateway.logger.info(`[Resharding] Shutting down old shards.`)
|
||||
await gateway.shutdown(ShardSocketCloseCodes.Resharded, 'Resharded!', false)
|
||||
gateway.logger.info(`[Resharding] Shutting down old shards.`);
|
||||
await gateway.shutdown(ShardSocketCloseCodes.Resharded, 'Resharded!', false);
|
||||
|
||||
gateway.logger.info(`[Resharding] Completed.`)
|
||||
gateway.shards = new Map(gateway.resharding.shards)
|
||||
gateway.resharding.shards.clear()
|
||||
gateway.logger.info(`[Resharding] Completed.`);
|
||||
gateway.shards = new Map(gateway.resharding.shards);
|
||||
gateway.resharding.shards.clear();
|
||||
},
|
||||
},
|
||||
|
||||
calculateTotalShards() {
|
||||
// Bots under 100k servers do not have access to LBS.
|
||||
if (gateway.totalShards < 100) {
|
||||
gateway.logger.debug(`[Gateway] Calculating total shards: ${gateway.totalShards}`)
|
||||
return gateway.totalShards
|
||||
gateway.logger.debug(`[Gateway] Calculating total shards: ${gateway.totalShards}`);
|
||||
return gateway.totalShards;
|
||||
}
|
||||
|
||||
gateway.logger.debug(`[Gateway] Calculating total shards`, gateway.totalShards, gateway.connection.sessionStartLimit.maxConcurrency)
|
||||
gateway.logger.debug(`[Gateway] Calculating total shards`, gateway.totalShards, gateway.connection.sessionStartLimit.maxConcurrency);
|
||||
// Calculate a multiple of `maxConcurrency` which can be used to connect to the gateway.
|
||||
return (
|
||||
Math.ceil(
|
||||
@@ -229,20 +229,20 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
// If `maxConcurrency` is 1, we can safely use 16 to get `totalShards` to be in a multiple of 16 so that we can prepare bots with 100k servers for LBS.
|
||||
(gateway.connection.sessionStartLimit.maxConcurrency === 1 ? 16 : gateway.connection.sessionStartLimit.maxConcurrency),
|
||||
) * (gateway.connection.sessionStartLimit.maxConcurrency === 1 ? 16 : gateway.connection.sessionStartLimit.maxConcurrency)
|
||||
)
|
||||
);
|
||||
},
|
||||
calculateWorkerId(shardId) {
|
||||
const workerId = options.spreadShardsInRoundRobin
|
||||
? shardId % gateway.totalWorkers
|
||||
: Math.min(Math.floor(shardId / gateway.shardsPerWorker), gateway.totalWorkers - 1)
|
||||
: Math.min(Math.floor(shardId / gateway.shardsPerWorker), gateway.totalWorkers - 1);
|
||||
gateway.logger.debug(
|
||||
`[Gateway] Calculating workerId: Shard: ${shardId} -> Worker: ${workerId} -> Per Worker: ${gateway.shardsPerWorker} -> Total: ${gateway.totalWorkers}`,
|
||||
)
|
||||
return workerId
|
||||
);
|
||||
return workerId;
|
||||
},
|
||||
prepareBuckets() {
|
||||
for (let i = 0; i < gateway.connection.sessionStartLimit.maxConcurrency; ++i) {
|
||||
gateway.logger.debug(`[Gateway] Preparing buckets for concurrency: ${i}`)
|
||||
gateway.logger.debug(`[Gateway] Preparing buckets for concurrency: ${i}`);
|
||||
gateway.buckets.set(i, {
|
||||
workers: [],
|
||||
leakyBucket: new LeakyBucket({
|
||||
@@ -251,95 +251,95 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
refillInterval: gateway.spawnShardDelay,
|
||||
logger: this.logger,
|
||||
}),
|
||||
})
|
||||
});
|
||||
}
|
||||
|
||||
// Organize all shards into their own buckets
|
||||
for (let shardId = gateway.firstShardId; shardId <= gateway.lastShardId; ++shardId) {
|
||||
gateway.logger.debug(`[Gateway] Preparing buckets for shard: ${shardId}`)
|
||||
gateway.logger.debug(`[Gateway] Preparing buckets for shard: ${shardId}`);
|
||||
|
||||
if (shardId >= gateway.totalShards) {
|
||||
throw new Error(`Shard (id: ${shardId}) is bigger or equal to the used amount of used shards which is ${gateway.totalShards}`)
|
||||
throw new Error(`Shard (id: ${shardId}) is bigger or equal to the used amount of used shards which is ${gateway.totalShards}`);
|
||||
}
|
||||
|
||||
const bucketId = shardId % gateway.connection.sessionStartLimit.maxConcurrency
|
||||
const bucket = gateway.buckets.get(bucketId)
|
||||
const bucketId = shardId % gateway.connection.sessionStartLimit.maxConcurrency;
|
||||
const bucket = gateway.buckets.get(bucketId);
|
||||
|
||||
if (!bucket) {
|
||||
throw new Error(
|
||||
`Shard (id: ${shardId}) got assigned to an illegal bucket id: ${bucketId}, expected a bucket id between 0 and ${
|
||||
gateway.connection.sessionStartLimit.maxConcurrency - 1
|
||||
}`,
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
// Get the worker id for this shard
|
||||
const workerId = gateway.calculateWorkerId(shardId)
|
||||
const worker = bucket.workers.find((w) => w.id === workerId)
|
||||
const workerId = gateway.calculateWorkerId(shardId);
|
||||
const worker = bucket.workers.find((w) => w.id === workerId);
|
||||
|
||||
// If this worker already exists, add the shard to its queue
|
||||
if (worker) {
|
||||
worker.queue.push(shardId)
|
||||
worker.queue.push(shardId);
|
||||
} else {
|
||||
bucket.workers.push({ id: workerId, queue: [shardId] })
|
||||
bucket.workers.push({ id: workerId, queue: [shardId] });
|
||||
}
|
||||
}
|
||||
},
|
||||
async spawnShards() {
|
||||
// Prepare the concurrency buckets
|
||||
gateway.prepareBuckets()
|
||||
gateway.prepareBuckets();
|
||||
|
||||
const promises = [...gateway.buckets.entries()].map(async ([bucketId, bucket]) => {
|
||||
for (const worker of bucket.workers) {
|
||||
for (const shardId of worker.queue) {
|
||||
await gateway.tellWorkerToIdentify(worker.id, shardId, bucketId)
|
||||
await gateway.tellWorkerToIdentify(worker.id, shardId, bucketId);
|
||||
}
|
||||
}
|
||||
})
|
||||
});
|
||||
|
||||
// We use Promise.all so we can start all buckets at the same time
|
||||
await Promise.all(promises)
|
||||
await Promise.all(promises);
|
||||
|
||||
// Check and reshard automatically if auto resharding is enabled.
|
||||
if (gateway.resharding.enabled && gateway.resharding.checkInterval !== -1) {
|
||||
// It is better to ensure there is always only one
|
||||
clearInterval(gateway.resharding.checkIntervalId)
|
||||
clearInterval(gateway.resharding.checkIntervalId);
|
||||
|
||||
if (!gateway.resharding.getSessionInfo) {
|
||||
gateway.resharding.enabled = false
|
||||
gateway.logger.warn("[Resharding] Resharding is enabled but 'resharding.getSessionInfo()' was not provided. Disabling resharding.")
|
||||
gateway.resharding.enabled = false;
|
||||
gateway.logger.warn("[Resharding] Resharding is enabled but 'resharding.getSessionInfo()' was not provided. Disabling resharding.");
|
||||
|
||||
return
|
||||
return;
|
||||
}
|
||||
|
||||
gateway.resharding.checkIntervalId = setInterval(async () => {
|
||||
const reshardingInfo = await gateway.resharding.checkIfReshardingIsNeeded()
|
||||
const reshardingInfo = await gateway.resharding.checkIfReshardingIsNeeded();
|
||||
|
||||
if (reshardingInfo.needed && reshardingInfo.info) await gateway.resharding.reshard(reshardingInfo.info)
|
||||
}, gateway.resharding.checkInterval)
|
||||
if (reshardingInfo.needed && reshardingInfo.info) await gateway.resharding.reshard(reshardingInfo.info);
|
||||
}, gateway.resharding.checkInterval);
|
||||
}
|
||||
},
|
||||
async shutdown(code, reason, clearReshardingInterval = true) {
|
||||
if (clearReshardingInterval) clearInterval(gateway.resharding.checkIntervalId)
|
||||
if (clearReshardingInterval) clearInterval(gateway.resharding.checkIntervalId);
|
||||
|
||||
await Promise.all(Array.from(gateway.shards.values()).map((shard) => shard.close(code, reason)))
|
||||
await Promise.all(Array.from(gateway.shards.values()).map((shard) => shard.close(code, reason)));
|
||||
},
|
||||
async sendPayload(shardId, payload) {
|
||||
const shard = gateway.shards.get(shardId)
|
||||
const shard = gateway.shards.get(shardId);
|
||||
|
||||
if (!shard) {
|
||||
throw new Error(`Shard (id: ${shardId} not found`)
|
||||
throw new Error(`Shard (id: ${shardId} not found`);
|
||||
}
|
||||
|
||||
await shard.send(payload)
|
||||
await shard.send(payload);
|
||||
},
|
||||
async tellWorkerToIdentify(workerId, shardId, bucketId) {
|
||||
gateway.logger.debug(`[Gateway] Tell worker #${workerId} to identify shard #${shardId} from bucket ${bucketId}`)
|
||||
await gateway.identify(shardId)
|
||||
gateway.logger.debug(`[Gateway] Tell worker #${workerId} to identify shard #${shardId} from bucket ${bucketId}`);
|
||||
await gateway.identify(shardId);
|
||||
},
|
||||
async identify(shardId: number) {
|
||||
let shard = this.shards.get(shardId)
|
||||
gateway.logger.debug(`[Gateway] Identifying ${shard ? 'existing' : 'new'} shard (${shardId})`)
|
||||
let shard = this.shards.get(shardId);
|
||||
gateway.logger.debug(`[Gateway] Identifying ${shard ? 'existing' : 'new'} shard (${shardId})`);
|
||||
|
||||
if (!shard) {
|
||||
shard = new Shard({
|
||||
@@ -358,59 +358,59 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
logger: this.logger,
|
||||
requestIdentify: async () => await gateway.requestIdentify(shardId),
|
||||
makePresence: gateway.makePresence,
|
||||
})
|
||||
});
|
||||
|
||||
this.shards.set(shardId, shard)
|
||||
this.shards.set(shardId, shard);
|
||||
}
|
||||
|
||||
await shard.identify()
|
||||
await shard.identify();
|
||||
},
|
||||
|
||||
async requestIdentify(shardId) {
|
||||
gateway.logger.debug(`[Gateway] Shard #${shardId} requested an identify.`)
|
||||
gateway.logger.debug(`[Gateway] Shard #${shardId} requested an identify.`);
|
||||
|
||||
const bucket = gateway.buckets.get(shardId % gateway.connection.sessionStartLimit.maxConcurrency)
|
||||
const bucket = gateway.buckets.get(shardId % gateway.connection.sessionStartLimit.maxConcurrency);
|
||||
|
||||
if (!bucket) {
|
||||
throw new Error("Can't request identify for a shard that is not assigned to any bucket.")
|
||||
throw new Error("Can't request identify for a shard that is not assigned to any bucket.");
|
||||
}
|
||||
|
||||
await bucket.leakyBucket.acquire()
|
||||
await bucket.leakyBucket.acquire();
|
||||
|
||||
gateway.logger.debug(`[Gateway] Approved identify request for Shard #${shardId}.`)
|
||||
gateway.logger.debug(`[Gateway] Approved identify request for Shard #${shardId}.`);
|
||||
},
|
||||
|
||||
async kill(shardId: number) {
|
||||
const shard = this.shards.get(shardId)
|
||||
const shard = this.shards.get(shardId);
|
||||
if (!shard) {
|
||||
gateway.logger.debug(`[Gateway] Shard #${shardId} was requested to be killed, but the shard could not be found.`)
|
||||
return
|
||||
gateway.logger.debug(`[Gateway] Shard #${shardId} was requested to be killed, but the shard could not be found.`);
|
||||
return;
|
||||
}
|
||||
|
||||
gateway.logger.debug(`[Gateway] Killing Shard #${shardId}`)
|
||||
this.shards.delete(shardId)
|
||||
await shard.shutdown()
|
||||
gateway.logger.debug(`[Gateway] Killing Shard #${shardId}`);
|
||||
this.shards.delete(shardId);
|
||||
await shard.shutdown();
|
||||
},
|
||||
|
||||
// Helpers methods below this
|
||||
|
||||
calculateShardId(guildId, totalShards) {
|
||||
// If none is provided, use the total shards number from gateway object.
|
||||
if (!totalShards) totalShards = gateway.totalShards
|
||||
if (!totalShards) totalShards = gateway.totalShards;
|
||||
// If it is only 1 shard, it will always be shard id 0
|
||||
if (totalShards === 1) {
|
||||
gateway.logger.debug(`[Gateway] calculateShardId (1 shard)`)
|
||||
return 0
|
||||
gateway.logger.debug(`[Gateway] calculateShardId (1 shard)`);
|
||||
return 0;
|
||||
}
|
||||
|
||||
gateway.logger.debug(`[Gateway] calculateShardId (guildId: ${guildId}, totalShards: ${totalShards})`)
|
||||
return Number((BigInt(guildId) >> 22n) % BigInt(totalShards))
|
||||
gateway.logger.debug(`[Gateway] calculateShardId (guildId: ${guildId}, totalShards: ${totalShards})`);
|
||||
return Number((BigInt(guildId) >> 22n) % BigInt(totalShards));
|
||||
},
|
||||
|
||||
async joinVoiceChannel(guildId, channelId, options) {
|
||||
const shardId = gateway.calculateShardId(guildId)
|
||||
const shardId = gateway.calculateShardId(guildId);
|
||||
|
||||
gateway.logger.debug(`[Gateway] joinVoiceChannel guildId: ${guildId} channelId: ${channelId}`)
|
||||
gateway.logger.debug(`[Gateway] joinVoiceChannel guildId: ${guildId} channelId: ${channelId}`);
|
||||
|
||||
await gateway.sendPayload(shardId, {
|
||||
op: GatewayOpcodes.VoiceStateUpdate,
|
||||
@@ -420,21 +420,21 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
self_mute: options?.selfMute ?? false,
|
||||
self_deaf: options?.selfDeaf ?? true,
|
||||
},
|
||||
})
|
||||
});
|
||||
},
|
||||
|
||||
async editBotStatus(data) {
|
||||
gateway.logger.debug(`[Gateway] editBotStatus data: ${JSON.stringify(data, jsonSafeReplacer)}`)
|
||||
gateway.logger.debug(`[Gateway] editBotStatus data: ${JSON.stringify(data, jsonSafeReplacer)}`);
|
||||
|
||||
await Promise.all(
|
||||
[...gateway.shards.values()].map(async (shard) => {
|
||||
gateway.editShardStatus(shard.id, data)
|
||||
gateway.editShardStatus(shard.id, data);
|
||||
}),
|
||||
)
|
||||
);
|
||||
},
|
||||
|
||||
async editShardStatus(shardId, data) {
|
||||
gateway.logger.debug(`[Gateway] editShardStatus shardId: ${shardId} -> data: ${JSON.stringify(data)}`)
|
||||
gateway.logger.debug(`[Gateway] editShardStatus shardId: ${shardId} -> data: ${JSON.stringify(data)}`);
|
||||
|
||||
await gateway.sendPayload(shardId, {
|
||||
op: GatewayOpcodes.PresenceUpdate,
|
||||
@@ -444,32 +444,32 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
activities: data.activities,
|
||||
status: data.status,
|
||||
},
|
||||
})
|
||||
});
|
||||
},
|
||||
|
||||
async requestMembers(guildId, options) {
|
||||
const shardId = gateway.calculateShardId(guildId)
|
||||
const shardId = gateway.calculateShardId(guildId);
|
||||
|
||||
if (gateway.intents && (!options?.limit || options.limit > 1) && !(gateway.intents & GatewayIntents.GuildMembers))
|
||||
throw new Error('Cannot fetch more then 1 member without the GUILD_MEMBERS intent')
|
||||
throw new Error('Cannot fetch more then 1 member without the GUILD_MEMBERS intent');
|
||||
|
||||
gateway.logger.debug(`[Gateway] requestMembers guildId: ${guildId} -> data: ${JSON.stringify(options)}`)
|
||||
gateway.logger.debug(`[Gateway] requestMembers guildId: ${guildId} -> data: ${JSON.stringify(options)}`);
|
||||
|
||||
if (options?.userIds?.length) {
|
||||
gateway.logger.debug(`[Gateway] requestMembers guildId: ${guildId} -> setting user limit based on userIds length: ${options.userIds.length}`)
|
||||
gateway.logger.debug(`[Gateway] requestMembers guildId: ${guildId} -> setting user limit based on userIds length: ${options.userIds.length}`);
|
||||
|
||||
options.limit = options.userIds.length
|
||||
options.limit = options.userIds.length;
|
||||
}
|
||||
|
||||
if (!options?.nonce) {
|
||||
let nonce = ''
|
||||
let nonce = '';
|
||||
|
||||
while (!nonce || gateway.cache.requestMembers.pending.has(nonce)) {
|
||||
nonce = randomBytes(16).toString('hex')
|
||||
nonce = randomBytes(16).toString('hex');
|
||||
}
|
||||
|
||||
options ??= { limit: 0 }
|
||||
options.nonce = nonce
|
||||
options ??= { limit: 0 };
|
||||
options.nonce = nonce;
|
||||
}
|
||||
|
||||
const members = !gateway.cache.requestMembers.enabled
|
||||
@@ -477,16 +477,16 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
: new Promise<Camelize<DiscordMemberWithUser[]>>((resolve, reject) => {
|
||||
// Should never happen.
|
||||
if (!gateway.cache.requestMembers.enabled || !options?.nonce) {
|
||||
reject(new Error("Can't request the members without the nonce or with the feature disabled."))
|
||||
return
|
||||
reject(new Error("Can't request the members without the nonce or with the feature disabled."));
|
||||
return;
|
||||
}
|
||||
|
||||
gateway.cache.requestMembers.pending.set(options.nonce, {
|
||||
nonce: options.nonce,
|
||||
resolve,
|
||||
members: [],
|
||||
})
|
||||
})
|
||||
});
|
||||
});
|
||||
|
||||
await gateway.sendPayload(shardId, {
|
||||
op: GatewayOpcodes.RequestGuildMembers,
|
||||
@@ -499,15 +499,15 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
user_ids: options?.userIds?.map((id) => id.toString()),
|
||||
nonce: options?.nonce,
|
||||
},
|
||||
})
|
||||
});
|
||||
|
||||
return await members
|
||||
return await members;
|
||||
},
|
||||
|
||||
async leaveVoiceChannel(guildId) {
|
||||
const shardId = gateway.calculateShardId(guildId)
|
||||
const shardId = gateway.calculateShardId(guildId);
|
||||
|
||||
gateway.logger.debug(`[Gateway] leaveVoiceChannel guildId: ${guildId} Shard ${shardId}`)
|
||||
gateway.logger.debug(`[Gateway] leaveVoiceChannel guildId: ${guildId} Shard ${shardId}`);
|
||||
|
||||
await gateway.sendPayload(shardId, {
|
||||
op: GatewayOpcodes.VoiceStateUpdate,
|
||||
@@ -517,7 +517,7 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
self_mute: false,
|
||||
self_deaf: false,
|
||||
},
|
||||
})
|
||||
});
|
||||
},
|
||||
|
||||
async requestSoundboardSounds(guildIds) {
|
||||
@@ -526,15 +526,15 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
* For this reason we need to group the ids with the shard the calculateShardId method gives
|
||||
*/
|
||||
|
||||
const map = new Map<number, BigString[]>()
|
||||
const map = new Map<number, BigString[]>();
|
||||
|
||||
for (const guildId of guildIds) {
|
||||
const shardId = gateway.calculateShardId(guildId)
|
||||
const shardId = gateway.calculateShardId(guildId);
|
||||
|
||||
const ids = map.get(shardId) ?? []
|
||||
map.set(shardId, ids)
|
||||
const ids = map.get(shardId) ?? [];
|
||||
map.set(shardId, ids);
|
||||
|
||||
ids.push(guildId)
|
||||
ids.push(guildId);
|
||||
}
|
||||
|
||||
await Promise.all(
|
||||
@@ -546,11 +546,11 @@ export function createGatewayManager(options: CreateGatewayManagerOptions): Gate
|
||||
},
|
||||
}),
|
||||
),
|
||||
)
|
||||
);
|
||||
},
|
||||
}
|
||||
};
|
||||
|
||||
return gateway
|
||||
return gateway;
|
||||
}
|
||||
|
||||
export interface CreateGatewayManagerOptions {
|
||||
@@ -558,32 +558,32 @@ export interface CreateGatewayManagerOptions {
|
||||
* Id of the first Shard which should get controlled by this manager.
|
||||
* @default 0
|
||||
*/
|
||||
firstShardId?: number
|
||||
firstShardId?: number;
|
||||
/**
|
||||
* Id of the last Shard which should get controlled by this manager.
|
||||
* @default 0
|
||||
*/
|
||||
lastShardId?: number
|
||||
lastShardId?: number;
|
||||
/**
|
||||
* Delay in milliseconds to wait before spawning next shard. OPTIMAL IS ABOVE 5100. YOU DON'T WANT TO HIT THE RATE LIMIT!!!
|
||||
* @default 5300
|
||||
*/
|
||||
spawnShardDelay?: number
|
||||
spawnShardDelay?: number;
|
||||
/**
|
||||
* Total amount of shards your bot uses. Useful for zero-downtime updates or resharding.
|
||||
* @default 1
|
||||
*/
|
||||
totalShards?: number
|
||||
totalShards?: number;
|
||||
/**
|
||||
* The amount of shards to load per worker.
|
||||
* @default 25
|
||||
*/
|
||||
shardsPerWorker?: number
|
||||
shardsPerWorker?: number;
|
||||
/**
|
||||
* The total amount of workers to use for your bot.
|
||||
* @default 4
|
||||
*/
|
||||
totalWorkers?: number
|
||||
totalWorkers?: number;
|
||||
/**
|
||||
* Whether to spread shards across workers in a round-robin manner.
|
||||
*
|
||||
@@ -596,53 +596,53 @@ export interface CreateGatewayManagerOptions {
|
||||
*
|
||||
* @default false
|
||||
*/
|
||||
spreadShardsInRoundRobin?: boolean
|
||||
spreadShardsInRoundRobin?: boolean;
|
||||
/** Important data which is used by the manager to connect shards to the gateway. */
|
||||
connection?: Camelize<DiscordGetGatewayBot>
|
||||
connection?: Camelize<DiscordGetGatewayBot>;
|
||||
/** Whether incoming payloads are compressed using zlib.
|
||||
*
|
||||
* @default false
|
||||
*/
|
||||
compress?: boolean
|
||||
compress?: boolean;
|
||||
/** What transport compression should be used */
|
||||
transportCompression?: TransportCompression | null
|
||||
transportCompression?: TransportCompression | null;
|
||||
/** The calculated intent value of the events which the shard should receive.
|
||||
*
|
||||
* @default 0
|
||||
*/
|
||||
intents?: number
|
||||
intents?: number;
|
||||
/** Identify properties to use */
|
||||
properties?: {
|
||||
/** Operating system the shard runs on.
|
||||
*
|
||||
* @default "darwin" | "linux" | "windows"
|
||||
*/
|
||||
os: string
|
||||
os: string;
|
||||
/** The "browser" where this shard is running on.
|
||||
*
|
||||
* @default "Discordeno"
|
||||
*/
|
||||
browser: string
|
||||
browser: string;
|
||||
/** The device on which the shard is running.
|
||||
*
|
||||
* @default "Discordeno"
|
||||
*/
|
||||
device: string
|
||||
}
|
||||
device: string;
|
||||
};
|
||||
/** Bot token which is used to connect to Discord */
|
||||
token: string
|
||||
token: string;
|
||||
/** The URL of the gateway which should be connected to.
|
||||
*
|
||||
* @default "wss://gateway.discord.gg"
|
||||
*/
|
||||
url?: string
|
||||
url?: string;
|
||||
/** The gateway version which should be used.
|
||||
*
|
||||
* @default 10
|
||||
*/
|
||||
version?: number
|
||||
version?: number;
|
||||
/** The events handlers */
|
||||
events?: ShardEvents
|
||||
events?: ShardEvents;
|
||||
/** This managers cache related settings. */
|
||||
cache?: {
|
||||
requestMembers?: {
|
||||
@@ -650,28 +650,28 @@ export interface CreateGatewayManagerOptions {
|
||||
* Whether or not request member requests should be cached.
|
||||
* @default false
|
||||
*/
|
||||
enabled?: boolean
|
||||
}
|
||||
}
|
||||
enabled?: boolean;
|
||||
};
|
||||
};
|
||||
/**
|
||||
* The logger that the gateway manager will use.
|
||||
* @default logger // The logger exported by `@discordeno/utils`
|
||||
*/
|
||||
logger?: Pick<typeof logger, 'debug' | 'info' | 'warn' | 'error' | 'fatal'>
|
||||
logger?: Pick<typeof logger, 'debug' | 'info' | 'warn' | 'error' | 'fatal'>;
|
||||
/**
|
||||
* Make the presence for when the bot connects to the gateway
|
||||
*
|
||||
* @remarks
|
||||
* This function will be called each time a Shard is going to identify
|
||||
*/
|
||||
makePresence?: () => Promise<DiscordUpdatePresence | undefined>
|
||||
makePresence?: () => Promise<DiscordUpdatePresence | undefined>;
|
||||
/** Options related to resharding. */
|
||||
resharding?: {
|
||||
/**
|
||||
* Whether or not automated resharding should be enabled.
|
||||
* @default true
|
||||
*/
|
||||
enabled: boolean
|
||||
enabled: boolean;
|
||||
/**
|
||||
* The % of how full a shard is when resharding should be triggered.
|
||||
*
|
||||
@@ -681,17 +681,17 @@ export interface CreateGatewayManagerOptions {
|
||||
*
|
||||
* @default 80 as in 80%
|
||||
*/
|
||||
shardsFullPercentage: number
|
||||
shardsFullPercentage: number;
|
||||
/**
|
||||
* The interval in milliseconds, of how often to check whether resharding is needed and reshard automatically. Set to -1 to disable auto resharding.
|
||||
* @default 28800000 (8 hours)
|
||||
*/
|
||||
checkInterval: number
|
||||
checkInterval: number;
|
||||
/** Handler to get shard count and other session info. */
|
||||
getSessionInfo?: () => Promise<Camelize<DiscordGetGatewayBot>>
|
||||
getSessionInfo?: () => Promise<Camelize<DiscordGetGatewayBot>>;
|
||||
/** Handler to edit the shard id on any cached guilds. */
|
||||
updateGuildsShardId?: (guildIds: string[], shardId: number) => Promise<void>
|
||||
}
|
||||
updateGuildsShardId?: (guildIds: string[], shardId: number) => Promise<void>;
|
||||
};
|
||||
}
|
||||
|
||||
export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
@@ -699,23 +699,23 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
buckets: Map<
|
||||
number,
|
||||
{
|
||||
workers: Array<{ id: number; queue: number[] }>
|
||||
workers: Array<{ id: number; queue: number[] }>;
|
||||
/** The bucket to queue the identifies. */
|
||||
leakyBucket: LeakyBucket
|
||||
leakyBucket: LeakyBucket;
|
||||
}
|
||||
>
|
||||
>;
|
||||
/** The shards that are created. */
|
||||
shards: Map<number, Shard>
|
||||
shards: Map<number, Shard>;
|
||||
/** The logger for the gateway manager. */
|
||||
logger: Pick<typeof logger, 'debug' | 'info' | 'warn' | 'error' | 'fatal'>
|
||||
logger: Pick<typeof logger, 'debug' | 'info' | 'warn' | 'error' | 'fatal'>;
|
||||
/** Everything related to resharding. */
|
||||
resharding: CreateGatewayManagerOptions['resharding'] & {
|
||||
/** The interval id of the check interval. This is used to clear the interval when the manager is shutdown. */
|
||||
checkIntervalId?: NodeJS.Timeout | undefined
|
||||
checkIntervalId?: NodeJS.Timeout | undefined;
|
||||
/** Holds the shards that resharding has created. Once resharding is done, this replaces the gateway.shards */
|
||||
shards: Map<number, Shard>
|
||||
shards: Map<number, Shard>;
|
||||
/** Handler to check if resharding is necessary. */
|
||||
checkIfReshardingIsNeeded: () => Promise<{ needed: boolean; info?: Camelize<DiscordGetGatewayBot> }>
|
||||
checkIfReshardingIsNeeded: () => Promise<{ needed: boolean; info?: Camelize<DiscordGetGatewayBot> }>;
|
||||
/**
|
||||
* Handler to begin resharding.
|
||||
*
|
||||
@@ -723,7 +723,7 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
* This function will resolve once the resharding is done.
|
||||
* So when all the calls to {@link tellWorkerToPrepare} and {@link onReshardingSwitch} are done.
|
||||
*/
|
||||
reshard: (info: Camelize<DiscordGetGatewayBot> & { firstShardId?: number; lastShardId?: number }) => Promise<void>
|
||||
reshard: (info: Camelize<DiscordGetGatewayBot> & { firstShardId?: number; lastShardId?: number }) => Promise<void>;
|
||||
/**
|
||||
* Handler to communicate to a worker that it needs to spawn a new shard and identify it for the resharding.
|
||||
*
|
||||
@@ -731,25 +731,25 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
* This handler works in the same way as the {@link tellWorkerToIdentify} handler.
|
||||
* So you should wait for the worker to have identified the shard before resolving the promise
|
||||
*/
|
||||
tellWorkerToPrepare: (workerId: number, shardId: number, bucketId: number) => Promise<void>
|
||||
tellWorkerToPrepare: (workerId: number, shardId: number, bucketId: number) => Promise<void>;
|
||||
/**
|
||||
* Handle called when all the workers have finished preparing for the resharding.
|
||||
*
|
||||
* This should make the new resharded shards become the active ones and shutdown the old ones
|
||||
*/
|
||||
onReshardingSwitch: () => Promise<void>
|
||||
}
|
||||
onReshardingSwitch: () => Promise<void>;
|
||||
};
|
||||
/** Determine max number of shards to use based upon the max concurrency. */
|
||||
calculateTotalShards: () => number
|
||||
calculateTotalShards: () => number;
|
||||
/** Determine the id of the worker which is handling a shard. */
|
||||
calculateWorkerId: (shardId: number) => number
|
||||
calculateWorkerId: (shardId: number) => number;
|
||||
/** Prepares all the buckets that are available for identifying the shards. */
|
||||
prepareBuckets: () => void
|
||||
prepareBuckets: () => void;
|
||||
/** Start identifying all the shards. */
|
||||
spawnShards: () => Promise<void>
|
||||
spawnShards: () => Promise<void>;
|
||||
/** Shutdown all shards. */
|
||||
shutdown: (code: number, reason: string, clearReshardingInterval?: boolean) => Promise<void>
|
||||
sendPayload: (shardId: number, payload: ShardSocketRequest) => Promise<void>
|
||||
shutdown: (code: number, reason: string, clearReshardingInterval?: boolean) => Promise<void>;
|
||||
sendPayload: (shardId: number, payload: ShardSocketRequest) => Promise<void>;
|
||||
/**
|
||||
* Allows users to hook in and change to communicate to different workers across different servers or anything they like.
|
||||
* For example using redis pubsub to talk to other servers.
|
||||
@@ -757,15 +757,15 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
* @remarks
|
||||
* This should wait for the worker to have identified the shard before resolving the returned promise
|
||||
*/
|
||||
tellWorkerToIdentify: (workerId: number, shardId: number, bucketId: number) => Promise<void>
|
||||
tellWorkerToIdentify: (workerId: number, shardId: number, bucketId: number) => Promise<void>;
|
||||
/** Tell the manager to identify a Shard. If this Shard is not already managed this will also add the Shard to the manager. */
|
||||
identify: (shardId: number) => Promise<void>
|
||||
identify: (shardId: number) => Promise<void>;
|
||||
/** Kill a shard. Close a shards connection to Discord's gateway (if any) and remove it from the manager. */
|
||||
kill: (shardId: number) => Promise<void>
|
||||
kill: (shardId: number) => Promise<void>;
|
||||
/** This function makes sure that the bucket is allowed to make the next identify request. */
|
||||
requestIdentify: (shardId: number) => Promise<void>
|
||||
requestIdentify: (shardId: number) => Promise<void>;
|
||||
/** Calculates the number of shards based on the guild id and total shards. */
|
||||
calculateShardId: (guildId: BigString, totalShards?: number) => number
|
||||
calculateShardId: (guildId: BigString, totalShards?: number) => number;
|
||||
/**
|
||||
* Connects the bot user to a voice or stage channel.
|
||||
*
|
||||
@@ -781,14 +781,18 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
*
|
||||
* @see {@link https://discord.com/developers/docs/topics/gateway#update-voice-state}
|
||||
*/
|
||||
joinVoiceChannel: (guildId: BigString, channelId: BigString, options?: AtLeastOne<Omit<UpdateVoiceState, 'guildId' | 'channelId'>>) => Promise<void>
|
||||
joinVoiceChannel: (
|
||||
guildId: BigString,
|
||||
channelId: BigString,
|
||||
options?: AtLeastOne<Omit<UpdateVoiceState, 'guildId' | 'channelId'>>,
|
||||
) => Promise<void>;
|
||||
/**
|
||||
* Edits the bot status in all shards that this gateway manages.
|
||||
*
|
||||
* @param data The status data to set the bots status to.
|
||||
* @returns nothing
|
||||
*/
|
||||
editBotStatus: (data: DiscordUpdatePresence) => Promise<void>
|
||||
editBotStatus: (data: DiscordUpdatePresence) => Promise<void>;
|
||||
/**
|
||||
* Edits the bot's status on one shard.
|
||||
*
|
||||
@@ -796,7 +800,7 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
* @param data The status data to set the bots status to.
|
||||
* @returns nothing
|
||||
*/
|
||||
editShardStatus: (shardId: number, data: DiscordUpdatePresence) => Promise<void>
|
||||
editShardStatus: (shardId: number, data: DiscordUpdatePresence) => Promise<void>;
|
||||
/**
|
||||
* Fetches the list of members for a guild over the gateway. If `gateway.cache.requestMembers.enabled` is not set, this function will return an empty array and you'll have to handle the `GUILD_MEMBERS_CHUNK` events yourself.
|
||||
*
|
||||
@@ -820,7 +824,7 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
*
|
||||
* @see {@link https://discord.com/developers/docs/topics/gateway#request-guild-members}
|
||||
*/
|
||||
requestMembers: (guildId: BigString, options?: Omit<RequestGuildMembers, 'guildId'>) => Promise<Camelize<DiscordMemberWithUser[]>>
|
||||
requestMembers: (guildId: BigString, options?: Omit<RequestGuildMembers, 'guildId'>) => Promise<Camelize<DiscordMemberWithUser[]>>;
|
||||
/**
|
||||
* Leaves the voice channel the bot user is currently in.
|
||||
*
|
||||
@@ -833,7 +837,7 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
*
|
||||
* @see {@link https://discord.com/developers/docs/topics/gateway#update-voice-state}
|
||||
*/
|
||||
leaveVoiceChannel: (guildId: BigString) => Promise<void>
|
||||
leaveVoiceChannel: (guildId: BigString) => Promise<void>;
|
||||
/**
|
||||
* Used to request soundboard sounds for a list of guilds.
|
||||
*
|
||||
@@ -853,7 +857,7 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
*
|
||||
* @see {@link https://discord.com/developers/docs/topics/gateway-events#request-soundboard-sounds}
|
||||
*/
|
||||
requestSoundboardSounds: (guildIds: BigString[]) => Promise<void>
|
||||
requestSoundboardSounds: (guildIds: BigString[]) => Promise<void>;
|
||||
/** This managers cache related settings. */
|
||||
cache: {
|
||||
requestMembers: {
|
||||
@@ -861,18 +865,18 @@ export interface GatewayManager extends Required<CreateGatewayManagerOptions> {
|
||||
* Whether or not request member requests should be cached.
|
||||
* @default false
|
||||
*/
|
||||
enabled: boolean
|
||||
enabled: boolean;
|
||||
/** The pending requests. */
|
||||
pending: Collection<string, RequestMemberRequest>
|
||||
}
|
||||
}
|
||||
pending: Collection<string, RequestMemberRequest>;
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
export interface RequestMemberRequest {
|
||||
/** The unique nonce for this request. */
|
||||
nonce: string
|
||||
nonce: string;
|
||||
/** The resolver handler to run when all members arrive. */
|
||||
resolve: (value: Camelize<DiscordMemberWithUser[]> | PromiseLike<Camelize<DiscordMemberWithUser[]>>) => void
|
||||
resolve: (value: Camelize<DiscordMemberWithUser[]> | PromiseLike<Camelize<DiscordMemberWithUser[]>>) => void;
|
||||
/** The members that have already arrived for this request. */
|
||||
members: DiscordMemberWithUser[]
|
||||
members: DiscordMemberWithUser[];
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import type { DiscordGatewayPayload, GatewayOpcodes } from '@discordeno/types'
|
||||
import type Shard from './Shard.js'
|
||||
import type { DiscordGatewayPayload, GatewayOpcodes } from '@discordeno/types';
|
||||
import type Shard from './Shard.js';
|
||||
|
||||
export enum ShardState {
|
||||
/** Shard is fully connected to the gateway and receiving events from Discord. */
|
||||
@@ -52,7 +52,7 @@ export interface ShardGatewayConfig {
|
||||
*
|
||||
* @see https://discord.com/developers/docs/topics/gateway#payload-compression
|
||||
*/
|
||||
compress: boolean
|
||||
compress: boolean;
|
||||
/**
|
||||
* What Transport Compression should be use
|
||||
*
|
||||
@@ -60,96 +60,96 @@ export interface ShardGatewayConfig {
|
||||
*
|
||||
* @see https://discord.com/developers/docs/topics/gateway#transport-compression
|
||||
*/
|
||||
transportCompression: TransportCompression | null
|
||||
transportCompression: TransportCompression | null;
|
||||
/** The calculated intent value of the events which the shard should receive.
|
||||
*
|
||||
* @default 0
|
||||
*/
|
||||
intents: number
|
||||
intents: number;
|
||||
/** Identify properties to use */
|
||||
properties: {
|
||||
/** Operating system the shard runs on.
|
||||
*
|
||||
* @default "darwin" | "linux" | "windows"
|
||||
*/
|
||||
os: string
|
||||
os: string;
|
||||
/** The "browser" where this shard is running on.
|
||||
*
|
||||
* @default "Discordeno"
|
||||
*/
|
||||
browser: string
|
||||
browser: string;
|
||||
/** The device on which the shard is running.
|
||||
*
|
||||
* @default "Discordeno"
|
||||
*/
|
||||
device: string
|
||||
}
|
||||
device: string;
|
||||
};
|
||||
/** Bot token which is used to connect to Discord */
|
||||
token: string
|
||||
token: string;
|
||||
/** The URL of the gateway which should be connected to.
|
||||
*
|
||||
* @default "wss://gateway.discord.gg"
|
||||
*/
|
||||
url: string
|
||||
url: string;
|
||||
/** The gateway version which should be used.
|
||||
*
|
||||
* @default 10
|
||||
*/
|
||||
version: number
|
||||
version: number;
|
||||
/**
|
||||
* The total number of shards to connect to across the entire bot.
|
||||
* @default 1
|
||||
*/
|
||||
totalShards: number
|
||||
totalShards: number;
|
||||
}
|
||||
|
||||
export interface ShardHeart {
|
||||
/** Whether or not the heartbeat was acknowledged by Discord in time. */
|
||||
acknowledged: boolean
|
||||
acknowledged: boolean;
|
||||
/** Interval between heartbeats requested by Discord. */
|
||||
interval: number
|
||||
interval: number;
|
||||
/** Id of the interval, which is used for sending the heartbeats. */
|
||||
intervalId?: NodeJS.Timeout
|
||||
intervalId?: NodeJS.Timeout;
|
||||
/** Unix (in milliseconds) timestamp when the last heartbeat ACK was received from Discord. */
|
||||
lastAck?: number
|
||||
lastAck?: number;
|
||||
/** Unix timestamp (in milliseconds) when the last heartbeat was sent. */
|
||||
lastBeat?: number
|
||||
lastBeat?: number;
|
||||
/** Round trip time (in milliseconds) from Shard to Discord and back.
|
||||
* Calculated using the heartbeat system.
|
||||
* Note: this value is undefined until the first heartbeat to Discord has happened.
|
||||
*/
|
||||
rtt?: number
|
||||
rtt?: number;
|
||||
/** Id of the timeout which is used for sending the first heartbeat to Discord since it's "special". */
|
||||
timeoutId?: NodeJS.Timeout
|
||||
timeoutId?: NodeJS.Timeout;
|
||||
}
|
||||
|
||||
export interface ShardEvents {
|
||||
/** A heartbeat has been send. */
|
||||
heartbeat?: (shard: Shard) => unknown
|
||||
heartbeat?: (shard: Shard) => unknown;
|
||||
/** A heartbeat ACK was received. */
|
||||
heartbeatAck?: (shard: Shard) => unknown
|
||||
heartbeatAck?: (shard: Shard) => unknown;
|
||||
/** Shard has received a Hello payload. */
|
||||
hello?: (shard: Shard) => unknown
|
||||
hello?: (shard: Shard) => unknown;
|
||||
/** The Shards session has been invalidated. */
|
||||
invalidSession?: (shard: Shard, resumable: boolean) => unknown
|
||||
invalidSession?: (shard: Shard, resumable: boolean) => unknown;
|
||||
/** The shard has started a resume action. */
|
||||
resuming?: (shard: Shard) => unknown
|
||||
resuming?: (shard: Shard) => unknown;
|
||||
/** The shard has successfully resumed an old session. */
|
||||
resumed?: (shard: Shard) => unknown
|
||||
resumed?: (shard: Shard) => unknown;
|
||||
/** Discord has requested the Shard to reconnect. */
|
||||
requestedReconnect?: (shard: Shard) => unknown
|
||||
requestedReconnect?: (shard: Shard) => unknown;
|
||||
/** The shard started to connect to Discord's gateway. */
|
||||
connecting?: (shard: Shard) => unknown
|
||||
connecting?: (shard: Shard) => unknown;
|
||||
/** The shard is connected with Discord's gateway. */
|
||||
connected?: (shard: Shard) => unknown
|
||||
connected?: (shard: Shard) => unknown;
|
||||
/** The shard has been disconnected from Discord's gateway. */
|
||||
disconnected?: (shard: Shard) => unknown
|
||||
disconnected?: (shard: Shard) => unknown;
|
||||
/** The shard has started to identify itself to Discord. */
|
||||
identifying?: (shard: Shard) => unknown
|
||||
identifying?: (shard: Shard) => unknown;
|
||||
/** The shard has successfully been identified itself with Discord. */
|
||||
ready?: (shard: Shard) => unknown
|
||||
ready?: (shard: Shard) => unknown;
|
||||
/** The shard has received a message from Discord. */
|
||||
message?: (shard: Shard, payload: DiscordGatewayPayload) => unknown
|
||||
message?: (shard: Shard, payload: DiscordGatewayPayload) => unknown;
|
||||
}
|
||||
|
||||
export enum ShardSocketCloseCodes {
|
||||
@@ -172,19 +172,19 @@ export enum ShardSocketCloseCodes {
|
||||
|
||||
export interface ShardSocketRequest {
|
||||
/** The OP-Code for the payload to send. */
|
||||
op: GatewayOpcodes
|
||||
op: GatewayOpcodes;
|
||||
/** Payload data. */
|
||||
d: unknown
|
||||
d: unknown;
|
||||
}
|
||||
|
||||
/** https://discord.com/developers/docs/topics/gateway#update-voice-state */
|
||||
export interface UpdateVoiceState {
|
||||
/** id of the guild */
|
||||
guildId: string
|
||||
guildId: string;
|
||||
/** id of the voice channel client wants to join (null if disconnecting) */
|
||||
channelId: string | null
|
||||
channelId: string | null;
|
||||
/** Is the client muted */
|
||||
selfMute: boolean
|
||||
selfMute: boolean;
|
||||
/** Is the client deafened */
|
||||
selfDeaf: boolean
|
||||
selfDeaf: boolean;
|
||||
}
|
||||
|
||||
@@ -1,58 +1,58 @@
|
||||
import { GatewayOpcodes, Intents } from '@discordeno/types'
|
||||
import { createGatewayManager, type GatewayManager } from '../../src/manager.js'
|
||||
import { ShardSocketCloseCodes } from '../../src/types.js'
|
||||
import { creatWSServer, heartbeatInterval } from './websocket.js'
|
||||
import { GatewayOpcodes, Intents } from '@discordeno/types';
|
||||
import { createGatewayManager, type GatewayManager } from '../../src/manager.js';
|
||||
import { ShardSocketCloseCodes } from '../../src/types.js';
|
||||
import { creatWSServer, heartbeatInterval } from './websocket.js';
|
||||
|
||||
describe('Gateway Integration', () => {
|
||||
it('Can connect to server', async () => {
|
||||
const { promise: connected, resolve: resolveConnected } = promiseWithResolvers<void>()
|
||||
const { promise: connected, resolve: resolveConnected } = promiseWithResolvers<void>();
|
||||
|
||||
const { port, close } = creatWSServer({
|
||||
onOpen: resolveConnected,
|
||||
})
|
||||
});
|
||||
|
||||
const gateway = createGatewayManagerWithPort(port)
|
||||
await gateway.spawnShards()
|
||||
await connected
|
||||
const gateway = createGatewayManagerWithPort(port);
|
||||
await gateway.spawnShards();
|
||||
await connected;
|
||||
|
||||
await gateway.shutdown(ShardSocketCloseCodes.TestingFinished, 'Testing finished')
|
||||
await gateway.shutdown(ShardSocketCloseCodes.TestingFinished, 'Testing finished');
|
||||
|
||||
// To avoid needing to wait 1m to get the bucket refil timer to fire we cancel it
|
||||
clearTimeout(gateway.shards.get(0)?.bucket.timeoutId)
|
||||
clearTimeout(gateway.shards.get(0)?.bucket.timeoutId);
|
||||
|
||||
close()
|
||||
})
|
||||
close();
|
||||
});
|
||||
|
||||
it('Can heartbeat', async () => {
|
||||
const { promise: connected, resolve: resolveConnected } = promiseWithResolvers<void>()
|
||||
const { promise: heartbeated, resolve: resolveHeartbeat, reject: rejectHeartbeat } = promiseWithResolvers<void>()
|
||||
const { promise: connected, resolve: resolveConnected } = promiseWithResolvers<void>();
|
||||
const { promise: heartbeated, resolve: resolveHeartbeat, reject: rejectHeartbeat } = promiseWithResolvers<void>();
|
||||
|
||||
const { port, close } = creatWSServer({
|
||||
onOpen: resolveConnected,
|
||||
onMessage: (message) => {
|
||||
if (message.op === GatewayOpcodes.Heartbeat) {
|
||||
resolveHeartbeat()
|
||||
resolveHeartbeat();
|
||||
}
|
||||
},
|
||||
})
|
||||
});
|
||||
|
||||
const gateway = createGatewayManagerWithPort(port)
|
||||
await gateway.spawnShards()
|
||||
await connected
|
||||
const gateway = createGatewayManagerWithPort(port);
|
||||
await gateway.spawnShards();
|
||||
await connected;
|
||||
|
||||
const timeout = setTimeout(() => rejectHeartbeat(new Error('Not heartbeat in time')), heartbeatInterval)
|
||||
await heartbeated
|
||||
const timeout = setTimeout(() => rejectHeartbeat(new Error('Not heartbeat in time')), heartbeatInterval);
|
||||
await heartbeated;
|
||||
|
||||
clearTimeout(timeout)
|
||||
clearTimeout(timeout);
|
||||
|
||||
await gateway.shutdown(ShardSocketCloseCodes.TestingFinished, 'Testing finished')
|
||||
await gateway.shutdown(ShardSocketCloseCodes.TestingFinished, 'Testing finished');
|
||||
|
||||
// To avoid needing to wait 1m to get the bucket refil timer to fire we cancel it
|
||||
clearTimeout(gateway.shards.get(0)?.bucket.timeoutId)
|
||||
clearTimeout(gateway.shards.get(0)?.bucket.timeoutId);
|
||||
|
||||
close()
|
||||
})
|
||||
})
|
||||
close();
|
||||
});
|
||||
});
|
||||
|
||||
function createGatewayManagerWithPort(port: number): GatewayManager {
|
||||
return createGatewayManager({
|
||||
@@ -74,22 +74,22 @@ function createGatewayManagerWithPort(port: number): GatewayManager {
|
||||
checkInterval: 0,
|
||||
shardsFullPercentage: 0,
|
||||
},
|
||||
})
|
||||
});
|
||||
}
|
||||
|
||||
// Polyfill for https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Promise/withResolvers
|
||||
function promiseWithResolvers<T>() {
|
||||
let resolve!: (value: T | PromiseLike<T>) => void
|
||||
let reject!: (reason?: any) => void
|
||||
let resolve!: (value: T | PromiseLike<T>) => void;
|
||||
let reject!: (reason?: any) => void;
|
||||
|
||||
const promise = new Promise<T>((_resolve, _reject) => {
|
||||
resolve = _resolve
|
||||
reject = _reject
|
||||
})
|
||||
resolve = _resolve;
|
||||
reject = _reject;
|
||||
});
|
||||
|
||||
return {
|
||||
promise,
|
||||
resolve,
|
||||
reject,
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { type DiscordGatewayPayload, GatewayOpcodes } from '@discordeno/types'
|
||||
import { type WebSocket, WebSocketServer } from 'ws'
|
||||
import { type DiscordGatewayPayload, GatewayOpcodes } from '@discordeno/types';
|
||||
import { type WebSocket, WebSocketServer } from 'ws';
|
||||
|
||||
/**
|
||||
* This value needs to be AT LEAST `1017`
|
||||
@@ -7,15 +7,15 @@ import { type WebSocket, WebSocketServer } from 'ws'
|
||||
* The reason for this is because the calculation in Shard.calculateSafeRequests will return 0 not allowing any sort of message to the websocket server.
|
||||
* Discord uses a way higher number for this value, but during this test we lower it since it would be annoying and useless make the test last 40+ seconds to test the heartbeat, but to make this work it needs to be at least 1017 so that calculateSafeRequests return 2 allowing for the shard to send messages.
|
||||
*/
|
||||
export const heartbeatInterval = 1017
|
||||
export const heartbeatInterval = 1017;
|
||||
|
||||
export function creatWSServer(options: CreateWsServerOptions) {
|
||||
// Port 0 according to the node:http docs is for requesting the OS a random unused port
|
||||
const server = new WebSocketServer({ port: 0 })
|
||||
const server = new WebSocketServer({ port: 0 });
|
||||
|
||||
const address = server.address()
|
||||
const address = server.address();
|
||||
if (typeof address !== 'object' || !address) {
|
||||
throw new TypeError('The address of the WebSocketServer should be an non-null object')
|
||||
throw new TypeError('The address of the WebSocketServer should be an non-null object');
|
||||
}
|
||||
|
||||
server.on('connection', (socket) => {
|
||||
@@ -25,22 +25,22 @@ export function creatWSServer(options: CreateWsServerOptions) {
|
||||
d: {
|
||||
heartbeat_interval: heartbeatInterval,
|
||||
},
|
||||
})
|
||||
});
|
||||
|
||||
options.onOpen?.()
|
||||
options.onOpen?.();
|
||||
|
||||
socket.on('message', (data) => {
|
||||
const msg = JSON.parse(data.toString('utf-8'))
|
||||
options.onMessage?.(msg)
|
||||
const msg = JSON.parse(data.toString('utf-8'));
|
||||
options.onMessage?.(msg);
|
||||
|
||||
switch (msg.op) {
|
||||
case GatewayOpcodes.Heartbeat: {
|
||||
send(socket, {
|
||||
op: GatewayOpcodes.HeartbeatACK,
|
||||
s: null,
|
||||
})
|
||||
});
|
||||
|
||||
break
|
||||
break;
|
||||
}
|
||||
|
||||
case GatewayOpcodes.Identify: {
|
||||
@@ -67,25 +67,25 @@ export function creatWSServer(options: CreateWsServerOptions) {
|
||||
guilds: [],
|
||||
application: { id: '0', flags: 0 },
|
||||
},
|
||||
})
|
||||
});
|
||||
|
||||
break
|
||||
break;
|
||||
}
|
||||
}
|
||||
})
|
||||
})
|
||||
});
|
||||
});
|
||||
|
||||
return {
|
||||
port: address.port,
|
||||
close: () => server.close(),
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
export interface CreateWsServerOptions {
|
||||
onOpen?: () => any
|
||||
onMessage?: (message: DiscordGatewayPayload) => any
|
||||
onOpen?: () => any;
|
||||
onMessage?: (message: DiscordGatewayPayload) => any;
|
||||
}
|
||||
|
||||
function send(ws: WebSocket, payload: object) {
|
||||
ws.send(JSON.stringify(payload))
|
||||
ws.send(JSON.stringify(payload));
|
||||
}
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { describe, it } from 'mocha'
|
||||
import { describe, it } from 'mocha';
|
||||
|
||||
describe('index.ts', () => {
|
||||
it('will import without error', async () => {
|
||||
await import('../../src/index.js')
|
||||
})
|
||||
})
|
||||
await import('../../src/index.js');
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user