Имперский Стражник — большая система с кучей разных источников сообщений. Сообщение может уйти из одного из 200 классов ядра или ещё большего зоопарка внешних модулей, но есть одна важная проблема: отправку сообщений надо контролировать.
Надо следить за лимитами, чтобы не передалбывать Bot API зазря. Надо следить за ответами, чтобы, если мы вышли за лимит, подождать и повторно отправить запрос. Надо вообще убеждаться в доставке, синхронизировать состояние и ответы.
Короче, следить надо за всем. Так что давайте я расскажу, что мы сделали, чтобы всё это работало и жило с учётом нашего стека.
Начнём с учёта отправки, лимитов и структуризации. Я решил, что за отправку сообщений будет отвечать синглтон-сервис, на который все остальные сервисы будут скидывать задачи. Но как это всё соединить, учитывая необходимость быстрого масштабирования и подключения новых модулей?
Тут на помощь пришёл Bull (не реклама, если что, просто делюсь либой, на которой делал решение). Давайте посмотрим, как выглядит воркер отправки:
new Worker(process.env.DEBUG ? "guard_dev_send" : "guard_send", async (job) => { // тут я переключаю очередь в зависимости от режима, в котором запущен сервис, чтобы не задевать прод во время тестирования
let telegram = job.data.target === "media" ? mediaBot.telegram : bot.telegram; // в данный момент бот, через который будет идти отправка, выбирается немного костыльно, но мы в работе :D
try {
return await telegram.sendMessage(job.data.id, job.data.text, job.data.extra); // отправляем сообщение
} catch (e: any) {
if (e.code === 429) {
await job.moveToDelayed((Date.now() / 1000) + 30000); // словили рейтлимит, ждём 30 секунд и пробуем опять
} else {
e._sourceStack = job.data.stack; // подпихиваем стек до очереди в ошибку
e._stack = e.stack; // стек внутри очереди
e._message = e.message; // и причину
return e;
}
}
}, {
connection: {
host: process.env.REDIS_HOST,
port: Number.parseInt(process.env.REDIS_PORT ?? "6379")
},
limiter: { // режем скорость отправки
max: 2,
duration: 1000
},
concurrency: 1, // и ограничиваем параллельность
}).on("failed", (job, err) => {
Queues.logger.error({
job,
err
}, "Failed to send message"); // если что-то пошло не так, логируем
});
Небольшой такой код, но он позволяет нам чётко контролировать очередь и лимиты. А дальше посмотрим на отправку и приём)
В комментариях разрешается открыть портал в холивар))
#tech