Основная сложность, с которой мы столкнулись, — необходимость рефакторить весь код. Хотелось избежать редактирования вообще всех мест, где идёт отправка сообщений. Поэтому мы вспомнили очень важную вещь, которая есть в JS: почти всё можно перезаписать по своему желанию. Этим мы и воспользовались, перезаписав метод отправки сообщения и подменив логику отправки в Telegram логикой отправки в очередь.
Получился вот такой прикольчик:
ctx.reply = async function (text: string, options?: ExtraSendMessage): Promise<Message> {
return Queues.getInstance().sendMessage("replyQueued", this.chat.id, text, options);
}
Тут мы перезаписали оригинальный метод ответа на сообщение своим собственным. Собственно, вуаля — половина работы отменена)
Дальше задача номер два, которая стала сильно интереснее: отправка сообщений разделяется на два типа — там, где нам нужен ответ, и там, где нам наплевать на ответ.
Нашей первой реализацией была идея внутри каждой операции отправки сообщения создавать листенер, который будет ждать ответа о том, что сообщение отправилось, и резолвить промис с ответом. А делать
await промиса или нет — уже решение на уровне бизнес-логики.Что могло пойти не так? Да, как оказалось, много чего. У нас утекли листенеры. Их стало так много, что Node.js перестал справляться с их контролем.
Встал вопрос: как сделать решение, которое не перегрузит среду, но сможет дожидаться ответа по отправке сотен сообщений и делать это асинхронно?
И тут я придумал прикольное решение: создаём
Map, который будет хранить ID задачи на отправку и коллбэк.
private _promisedMessages = new Map<string, { resolve: (message: Message) => void, reject: (error: Error) => void }>();
В отправке мы создаём задачу и промис, записывая в мапу ID задачи и
resolve этого промиса в качестве коллбэка.
async sendMessage(target: TargetBot, name: string, id: number, text: string, extra: ExtraSendMessage = {}, priority = 10) {
let error = new Error("Stack trace"); // получаем стектрейс до этого момента
let job = await this._sendQueue.add(name, {id, text, extra, target, stack: error.stack}, {priority});
return new Promise<Message>((resolve, reject) => {
this._promisedMessages.set(job.id!, {resolve, reject});
});
}
И вместо сотен листенеров делаем один-единственный глобальный, который будет работать до самого завершения процесса:
new QueueEvents(process.env.DEBUG ? "guard_dev_send" : "guard_send", {
connection: {
host: process.env.REDIS_HOST,
port: Number.parseInt(process.env.REDIS_PORT ?? "6379"),
},
})
.on("completed", async ({jobId, returnvalue}) => {
if (this._promisedMessages.has(jobId)) { // проверяем, что это сообщение ожидается на этом воркере
let {resolve, reject} = this._promisedMessages.get(jobId)!; // получаем коллбэки этого сообщения
this._promisedMessages.delete(jobId); // удаляем его из мапы
resolve(returnvalue); // отдаём результат промису
}
});
И тем самым получается, что в команде код выглядит как самый обычный вызов:
if (!peer.isUserKnown()) {
await ctx.reply(await ctx.lang.t("commands.general.userNotFound"));
return;
}
А на самом деле за пару миллисекунд запрос обходит две очереди и несколько сервисов, выполняя работу чётенько и безопасненько)
С первого взгляда выглядит как быстрое и логичное решение, но мы к этому шли месяц-полтора.
Вот как-то так.
#tech