Уведомления в реальном времени
itd.notifications.events возвращает стабильный канал новых уведомлений:
import { formatNotificationText, resolveNotificationUrl } from 'itd-api';
const stream = itd.notifications.events;
stream.on('notification', ({ notification, sound }) => {
console.log(sound ? '🔔' : '🔕');
console.log(formatNotificationText(notification));
console.log(resolveNotificationUrl(notification));
});
await stream.connect();Нормализованные обновления
Промежуточные и асинхронные обработчики получают не более одного логического обновления из одного транспортного кадра:
import { type NotificationEvent, NotificationUpdateType } from 'itd-api';
type NotificationEventsUpdate =
| { type: typeof NotificationUpdateType.Notification; data: NotificationEvent }
| { type: typeof NotificationUpdateType.UnreadCount; data: number }
| { type: typeof NotificationUpdateType.Unknown; name: string; data: unknown };Обновление с типом NotificationUpdateType.Notification содержит нормализованное уведомление. Тип NotificationUpdateType.UnreadCount создаётся для отдельного серверного события или начальной REST-синхронизации. NotificationUpdateType.Unknown используется только для неизвестных библиотеке событий.
Контекст обработчика содержит:
update— нормализованные данные обновления;stream— текущий потокNotificationEvents;raw— исходный транспортный кадр илиundefinedдля REST-синхронизации;origin— источник изNotificationUpdateOrigin: поток или начальная синхронизация.
Промежуточные обработчики
use() добавляет промежуточный обработчик. Вызов next() передаёт обновление дальше, а await next() позволяет выполнить код после завершения остальной цепочки:
const removeLogging = stream.use(async (context, next) => {
const startedAt = Date.now();
await next();
console.log(context.update.type, Date.now() - startedAt);
});Если не вызвать next(), обновление не попадёт к следующим промежуточным и асинхронным обработчикам, а также к слушателям on():
stream.use(async (context, next) => {
if (
context.update.type === NotificationUpdateType.Notification &&
context.update.data.notification.isRead
) {
return;
}
await next();
});use() принимает функцию либо объект с методом middleware(), например EventRouter. Функция, возвращённая use(), снимает обработчик. Для каждого обновления используется снимок цепочки и маршрутов на момент получения, поэтому изменение подписок через use(), route() или otherwise() влияет только на следующие обновления.
Исключение не закрывает соединение. Оно приходит через событие middlewareError; если на событие никто не подписан, ошибка записывается через настроенный logger или в консоль.
Асинхронные обработчики и фильтры
onUpdate() подписывает обработчик на определённый тип обновления и сужает тип контекста:
const off = stream.onUpdate(NotificationUpdateType.UnreadCount, async ({ update }) => {
await saveBadge(update.data); // number
});Если передать только обработчик, он получит все нормализованные обновления:
stream.onUpdate(async ({ update }) => {
await saveNotificationEventsUpdate(update);
});Для уведомлений доступны краткие и объектные фильтры:
import { NotificationType } from 'itd-api';
stream.onNotification(NotificationType.PostComment, async ({ update }) => {
await handleComment(update.data.notification);
});
stream.onNotification(
{
type: [NotificationType.PostComment, NotificationType.CommentReply],
actorId: authorId,
entityId: commentId,
parentEntityId: postId,
},
async ({ update }) => {
await handleThreadUpdate(update.data.notification);
},
);Все указанные поля объектного фильтра объединяются через И. actorId совпадает, если идентификатор пользователя есть хотя бы у одного элемента notification.actors. Дополнительную проверку можно передать в predicate. onUpdate() и onNotification() также принимают собственную функцию проверки, включая функцию сужения типа.
Асинхронные обработчики выполняются после промежуточной цепочки и до слушателей on(). drain() дожидается их завершения. Ошибка одного обработчика не мешает остальным и передаётся событием handlerError.
Для низкоуровневого наблюдения stream.on('message', listener) получает каждый исходный кадр транспорта, включая известные события и подтверждение подключения. Такой слушатель выполняется синхронно и не учитывается drain().
Для большой цепочки связанные промежуточные обработчики можно собрать через EventComposer. Он поддерживает фильтры, статические маршруты и локальные границы ошибок; точный контракт и семантика вложенных цепочек описаны в справочнике событий.
Маршрутизация
EventRouter выбирает цепочку промежуточных обработчиков по ключу:
import { NotificationType, EventRouter, NotificationUpdateType } from 'itd-api';
const router = new EventRouter((context) => {
if (context.update.type !== NotificationUpdateType.Notification) return 'other';
return context.update.data.notification.type;
});
const removeCommentRoute = router.route(NotificationType.PostComment, async (context, next) => {
if (context.update.type === NotificationUpdateType.Notification) {
await handleComment(context.update.data.notification);
}
await next();
});
router.otherwise(async (_context, next) => {
await next();
});
const removeRouter = stream.use(router);Функция выбора маршрута может быть асинхронной. Если ключ не зарегистрирован, используется цепочка otherwise; без неё обновление передаётся следующему внешнему обработчику. route() и otherwise() возвращают функции удаления своих регистраций.
Порядок и конкурентность
По умолчанию обновления обрабатываются последовательно в порядке получения. Транспорт при этом продолжает принимать данные и складывает их во внутреннюю очередь.
import { NotificationUpdateType } from 'itd-api';
const itd = new ItdClient({
events: {
notifications: {
concurrency: 4,
sequentialize: (context) => {
if (context.update.type !== NotificationUpdateType.Notification) return undefined;
return context.update.data.notification.parentEntityId ?? undefined;
},
},
},
});
const stream = itd.notifications.events;concurrency задаёт общий предел одновременно обрабатываемых обновлений. sequentialize() возвращает ключ или массив ключей: обновления с хотя бы одним общим ключом выполняются последовательно в порядке получения. Независимые обновления могут выполняться параллельно.
REST и поток
Уведомления из itd.notifications.list() и потока приведены к общей форме, поэтому их можно хранить в одном массиве:
const history = await itd.notifications.list({ limit: 20 });
stream.on('notification', ({ notification }) => {
history.items.unshift(notification);
});Сервер использует короткие типы вроде like, comment и repost. Библиотека приводит их к однозначным post_reaction, post_comment, post_repost, сохраняя исходное значение в rawType, а исходный объект — в raw.
resolveNotificationUrl() учитывает смысл идентификаторов конкретного типа и строит ссылку на профиль, пост или комментарий.
Переподключение
Поток самостоятельно обрабатывает:
- обрыв соединения;
- обновление токена доступа;
- восстановление сети;
- возвращение браузерной вкладки из фона;
- отсутствие данных дольше
idleTimeout.
По умолчанию используются задержки [1, 2, 4, 8, 16, 30] секунд со случайным разбросом ±30% и не более 15 последовательных попыток. Сервер не гарантирует служебные сообщения для поддержания соединения, поэтому клиент считает молчащее соединение мёртвым через 90 секунд. handshakeTimeout ограничивает установку SSE-соединения 20 секундами.
Состояние можно отслеживать:
stream.on('status', (status) => {
console.log(status); // connecting, connected, disconnected, error
});Завершение:
stream.disconnect();
await stream.drain();
// либо закрыть все потоки клиента и дождаться активных обработчиков
await itd.close();disconnect() закрывает транспорт, отменяет переподключение и отбрасывает обновления, обработка которых ещё не началась. Уже активные обработчики завершаются; drain() позволяет их дождаться. После ручного connect() тот же экземпляр снова участвует в жизненном цикле клиента и будет закрыт следующим itd.close().
Счётчик непрочитанных
При connect() поток по умолчанию запрашивает начальный счётчик через REST и отправляет событие unreadCount. Последующие уведомления обычно не содержат актуального счётчика, поэтому увеличивайте его локально:
let unread = 0;
stream.on('unreadCount', (count) => {
unread = count;
});
stream.on('notification', (event) => {
unread = event.unreadCount ?? unread + 1;
});
await stream.connect();Начальную синхронизацию можно отключить через syncCount: false. После массовой отметки о прочтении запросите актуальное значение через itd.notifications.count().
Резервный опрос
В средах без потокового чтения ответа, например в некоторых версиях React Native, клиент автоматически переключается на периодический опрос. Интервал настраивается через pollInterval.
Можно выбрать транспорт явно:
import { NotificationEventsTransport } from 'itd-api';
const itd = new ItdClient({
events: {
notifications: {
transport: NotificationEventsTransport.Poll,
pollInterval: 5_000,
},
},
});
const stream = itd.notifications.events;Несколько аккаунтов
У каждого клиента один стабильный itd.notifications.events. Для десяти аккаунтов это не более десяти SSE-соединений, поэтому подключайте канал только там, где он действительно нужен.
Запускаемый пример
Пример принимает новые уведомления, отслеживает состояние соединения и корректно завершает активные обработчики по SIGINT или SIGTERM.
ITD_TOKEN=<токен> node guides/events/examples/notifications.mjsИсходник: examples/notifications.mjs.