ApiWs
src/chathouse/components/ladies/ApiWs.ts
Observer на [[services/chathouse/services/lady-runner]]. Підключається до партнерського API через WebSocket (Centrifuge) від імені TU і в реальному часі обробляє вхідні події — повідомлення, typing, сповіщення, зміни лімітів.
Технологія
Використовує бібліотеку Centrifuge. Після підключення підписується на канал api#${ladyUlid}. Всі події надходять через publication handler.
Підключення (connect())
- Якщо
centrifuge.state === 'connected'— пропустити (ідемпотентність). - Отримати broadcast-токен через
authService.broadcastingToken(initiatorId). - Якщо токен не вдалося отримати — почекати 20 сек і повторити. Ліміт: 5 спроб. Після 5-ї невдачі —
disconnectOperator(примусово відключає оператора). - Створити
Centrifugeз токеном, підписатися на каналapi#${ladyUlid}. - Запустити прослуховування подій
publication.
Reconnect при розриві з’єднання обробляє сама бібліотека Centrifuge (вбудована логіка).
Оброблювані типи подій
| Тип | Що означає | Дія |
|---|---|---|
Send | RU надіслав повідомлення в чат | Форвардинг + створення таску |
InMailSend | RU надіслав mail | Форвардинг + створення таску |
Typing | RU набирає текст | Typing-таск з кеш-гардом |
NewNotification | RU онлайн або зробив покупку | → notificationLogic |
LimitChanged | Змінились ліміти в діалозі | Форвардинг якщо є хоч один ліміт |
Всі інші типи (GiftOpened, Read, DeleteMessage тощо) — ігноруються.
Обробка publication (покрокова логіка)
-
Дедуплікація. Кожне повідомлення має fingerprint (
action:messageIdабоaction:manUlid:sent_at). Якщо збігається з попереднім — ігнорується. -
Нормалізація часу. WebSocket надсилає мікросекунди (
2025-02-04T15:14:25.301443Z), HTTP — нулі (...000000Z). ApiWs обрізає до.000000Zщоб час збігався при подальшому порівнянні. -
LimitChanged: якщо є хоча б один ліміт (chat або mail) — форвардується клієнту якapiSocket-подія. Якщо лімітів нема — ігнорується. -
NewNotification→notificationLogic()(детальніше нижче). -
Фільтр вхідних повідомлень: якщо тип не в
allowedActionTypesабо повідомлення не є вхідним (isIncomingMessage) — ігнорується. -
Send+ “like message body” (повідомлення-реакція, не текст): додаткова перевірка — чи є ліміти в діалозі з цим RU (API-запит). Якщо лімітів нема — ігнорується. -
Форвардинг клієнту — raw-дані надсилаються оператору як
apiSocket-подія. -
Створення таску через
TaskWsFactory.create():- TypingTask: перевірити
cacheTypingTasksGuard(1 хв cooldown на manUlid). Якщо вже є таск в сторі — пропустити. Перевірити ліміти в діалозі. Додати до task store. - ITask (звичайний таск): якщо вже є таск з тим самим
eventTriggerTaskпо цій парі — пропустити. Зберегти в БД і додати до task store.
- TypingTask: перевірити
notificationLogic()
Спрацьовує на NewNotification з типами ONLINE_NOW або RECENT_PURCHASE і дією NEW або REPEAT.
- Отримати повідомлення діалогу з цим RU.
- Якщо є chat-ліміти → створити
NeedToWriteWsNotificationTask. - Якщо є тільки mail-ліміти → створити
NeedToWriteMailWsNotificationTask. - Якщо таск з тим самим
eventTriggerTaskвже в сторі → пропустити. - Зберегти в БД і додати до task store.
Кеш typing (cacheTypingTasksGuard)
Map<manUlid, {timestamp}> — гард проти дублікатів typing-тасків.
- Якщо typing від цього RU прийшов менше 1 хв тому — новий typing-таск не створюється.
- Очищується
setIntervalкожну 1 хв — видаляються записи старші за 1 хв.
stop()
isStopped = true.- Очищаються всі таймери (
timeout,cacheCleanupInterval). centrifuge.removeAllListeners()+centrifuge.disconnect().subscription.unsubscribe().cacheTypingTasksGuard.clear().runner.removeObserver(this).- Всі посилання на сервіси →
null(запобігання memory leaks).
Нюанси
connect()викликається зLadyRunner.start()після логіну TU — токен отримується вже під авторизованою сесією.- Centrifuge сам управляє reconnect при розриві з’єднання — ApiWs не додає власної reconnect-логіки для WebSocket-розривів (тільки для помилок отримання токена).
previousMessageзберігає fingerprint одного попереднього повідомлення — захищає тільки від миттєвих дублікатів.