From d71b4052779cacf69b7243e1e8e04b2ae21e76de Mon Sep 17 00:00:00 2001 From: Vasilii Mikhailovskii Date: Sat, 14 Dec 2024 23:58:29 +0200 Subject: [PATCH] send messages to telegram from RQCP (ReadQrcodeChatPage) with link to MQCP (ManageQrcodeChatPage) --- .gitea/workflows/deploy-service-tg-bot.yaml | 2 + .../src/components/chat/ChatMessages.vue | 1 + .../src/health-server.ts | 55 ++++++--- .../DBDocuments/UserChatSettingsDocument.ts | 7 +- pnpm-lock.yaml | 22 ++++ .../package.json | 5 - .../src/dbs.ts | 4 +- .../src/service.ts | 25 ++-- .../tsconfig.json | 3 - services/tg-bot/package.json | 10 +- .../conversations/newQRCodeConversation.ts | 7 +- services/tg-bot/src/bot/index.ts | 75 ++++++------ services/tg-bot/src/config/env.ts | 6 + services/tg-bot/src/dbs.ts | 59 +++++----- services/tg-bot/src/global.d.ts | 8 ++ services/tg-bot/src/service.ts | 31 +++++ .../src/services/iterateMessagesChanges.ts | 111 ++++++++++++++++++ .../src/services/serviceChangesMessagesDb.ts | 52 ++++++++ .../src/{models/types.ts => types/models.ts} | 0 services/tg-bot/src/types/myContext.ts | 15 +-- services/tg-bot/src/utils/escapeHtml.ts | 21 ++++ services/tg-bot/tsconfig.json | 4 +- 22 files changed, 396 insertions(+), 127 deletions(-) create mode 100644 services/tg-bot/src/global.d.ts create mode 100644 services/tg-bot/src/service.ts create mode 100644 services/tg-bot/src/services/iterateMessagesChanges.ts create mode 100644 services/tg-bot/src/services/serviceChangesMessagesDb.ts rename services/tg-bot/src/{models/types.ts => types/models.ts} (100%) create mode 100644 services/tg-bot/src/utils/escapeHtml.ts diff --git a/.gitea/workflows/deploy-service-tg-bot.yaml b/.gitea/workflows/deploy-service-tg-bot.yaml index 14286ca..ab31b05 100644 --- a/.gitea/workflows/deploy-service-tg-bot.yaml +++ b/.gitea/workflows/deploy-service-tg-bot.yaml @@ -21,6 +21,8 @@ jobs: COUCHDB_PASSWORD: ${{ secrets.COUCHDB_PASSWORD }} BOT_TOKEN: ${{ secrets.BOT_TOKEN }} READ_BASE_URL: ${{ vars.READ_BASE_URL }} + APP_BASE_URL: ${{ vars.APP_BASE_URL }} + HEARTBEAT_INTERVAL: ${{ vars.HEARTBEAT_INTERVAL }} steps: - uses: actions/checkout@v4 diff --git a/apps/frontend/src/components/chat/ChatMessages.vue b/apps/frontend/src/components/chat/ChatMessages.vue index 4008039..b6b3086 100644 --- a/apps/frontend/src/components/chat/ChatMessages.vue +++ b/apps/frontend/src/components/chat/ChatMessages.vue @@ -3,6 +3,7 @@
((resolve, reject) => { - server.once("error", reject); - server.listen(port, hostname, () => { - resolve(); - console.log( - `Healthcheck server running on http://${hostname}:${port}/health`, - ); - }); + server?.unref?.(); - // Закрываем сервер при активации AbortController - controller.signal.addEventListener("abort", () => { - console.log("Healthcheck server shutting down due to abort signal."); - reject(); - server.close(); + const stop = () => { + return new Promise((resolve, reject) => { + server.close((error) => { + if (error) reject(error); + else resolve(); + }); }); - }); + }; + + const start = () => { + return new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(port, hostname, () => { + resolve(); + console.log( + `Healthcheck server running on http://${hostname}:${port}/health`, + ); + }); + + // Закрываем сервер при активации AbortController + controller.signal.addEventListener("abort", () => { + console.log("Healthcheck server shutting down due to abort signal."); + stop().then(reject, reject); + }); + }); + }; + + return { + start, + stop, + }; } + +export const startHealthServer = ( + controller: AbortController, + defaultPort: number = 8000, + defaultHostName: string = "127.0.0.1", +) => { + return createHealthServer(controller, defaultPort, defaultHostName).start(); +}; diff --git a/packages/types/src/DBDocuments/UserChatSettingsDocument.ts b/packages/types/src/DBDocuments/UserChatSettingsDocument.ts index 8fda318..87f5a24 100644 --- a/packages/types/src/DBDocuments/UserChatSettingsDocument.ts +++ b/packages/types/src/DBDocuments/UserChatSettingsDocument.ts @@ -24,7 +24,12 @@ export interface UserChatSettingsDocument extends BaseDocument { browser_uuid: string; /** Включены ли push-уведомления */ - is_push_notification_enabled: boolean; + is_push_notification_enabled?: boolean; + + /** + * Включены ли сообщения в Telegram + */ + is_telegram_messages_enabled?: boolean; /** Дата последнего обновления */ updated_at?: string; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 3770110..d48970c 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -337,6 +337,9 @@ importers: '@grammyjs/conversations': specifier: ^1.2.0 version: 1.2.0(grammy@1.33.0) + '@grammyjs/parse-mode': + specifier: ^1.10.0 + version: 1.10.0(grammy@1.33.0) '@hereconnect/constants': specifier: workspace:^ version: link:../../packages/constants @@ -425,6 +428,9 @@ importers: tsx: specifier: ^4.19.2 version: 4.19.2 + type-fest: + specifier: ^4.30.1 + version: 4.30.1 typescript: specifier: ^5.7.2 version: 5.7.2 @@ -921,6 +927,12 @@ packages: peerDependencies: grammy: ^1.20.1 + '@grammyjs/parse-mode@1.10.0': + resolution: {integrity: sha512-ZjbY2Ax0b4Nf8lPz3NV0cWDxUC10kOkzgxws+iTdgG+hAiPQVUnP/oJnAKrW1ZQOWULZYQ4GOR/+aCZHyUuLQA==} + engines: {node: '>=14.13.1'} + peerDependencies: + grammy: ^1.20.1 + '@grammyjs/types@3.17.0': resolution: {integrity: sha512-e8AR3xQwRAFX248E7Qw/7mIu1OzvoXloJzOBJVtuPKzzL7tGkn5trZAdZUBgGViVQg5ZwVS/x9N2nRrcyH/DfA==} @@ -3637,6 +3649,10 @@ packages: resolution: {integrity: sha512-Ne+eE4r0/iWnpAxD852z3A+N0Bt5RN//NjJwRd2VFHEmrywxf5vsZlh4R6lixl6B+wz/8d+maTSAkN1FIkI3LQ==} engines: {node: '>=10'} + type-fest@4.30.1: + resolution: {integrity: sha512-ojFL7eDMX2NF0xMbDwPZJ8sb7ckqtlAi1GsmgsFXvErT9kFTk1r0DuQKvrCh73M6D4nngeHJmvogF9OluXs7Hw==} + engines: {node: '>=16'} + typescript-eslint@8.16.0: resolution: {integrity: sha512-wDkVmlY6O2do4V+lZd0GtRfbtXbeD0q9WygwXXSJnC1xorE8eqyC2L1tJimqpSeFrOzRlYtWnUp/uzgHQOgfBQ==} engines: {node: ^18.18.0 || ^20.9.0 || >=21.1.0} @@ -4426,6 +4442,10 @@ snapshots: transitivePeerDependencies: - supports-color + '@grammyjs/parse-mode@1.10.0(grammy@1.33.0)': + dependencies: + grammy: 1.33.0 + '@grammyjs/types@3.17.0': {} '@humanfs/core@0.19.1': {} @@ -7380,6 +7400,8 @@ snapshots: type-fest@0.20.2: {} + type-fest@4.30.1: {} + typescript-eslint@8.16.0(eslint@9.15.0(jiti@1.21.6))(typescript@5.6.3): dependencies: '@typescript-eslint/eslint-plugin': 8.16.0(@typescript-eslint/parser@8.16.0(eslint@9.15.0(jiti@1.21.6))(typescript@5.6.3))(eslint@9.15.0(jiti@1.21.6))(typescript@5.6.3) diff --git a/services/message-delivery-method-web-push/package.json b/services/message-delivery-method-web-push/package.json index 8bdd061..ce3fbbb 100644 --- a/services/message-delivery-method-web-push/package.json +++ b/services/message-delivery-method-web-push/package.json @@ -10,11 +10,6 @@ "node": ">=v20.10.0", "pnpm": ">=8.15.0" }, - "homepage": "", - "repository": { - "type": "git", - "url": "" - }, "files": [ "dist/**/*" ], diff --git a/services/message-delivery-method-web-push/src/dbs.ts b/services/message-delivery-method-web-push/src/dbs.ts index 06fea02..89aa788 100644 --- a/services/message-delivery-method-web-push/src/dbs.ts +++ b/services/message-delivery-method-web-push/src/dbs.ts @@ -62,7 +62,7 @@ export const userChatSettingsDb = new PouchDB( }, ); -export const checkDatabases = async () => { +const checkDatabases = async () => { const entries = await Promise.all( Object.entries({ serverSettingsDb, @@ -92,3 +92,5 @@ export const checkDatabases = async () => { return Object.fromEntries(entries); }; + +export const checkDatabasesPromise = checkDatabases(); diff --git a/services/message-delivery-method-web-push/src/service.ts b/services/message-delivery-method-web-push/src/service.ts index c3fa0db..c614b08 100644 --- a/services/message-delivery-method-web-push/src/service.ts +++ b/services/message-delivery-method-web-push/src/service.ts @@ -9,7 +9,7 @@ import { messagesDbUrl, webPushSubscriptionsDb, userChatSettingsDb, - checkDatabases, + checkDatabasesPromise, serviceSeqTrackerDb, } from "./dbs"; import { setupVapidDetails } from "./setupVapidDetails"; @@ -29,7 +29,7 @@ function getTopicSafe(str: string) { } const start = async () => { - const checkedDatabases = await checkDatabases(); + const checkedDatabases = await checkDatabasesPromise; await setupVapidDetails(); const serviceSeqTracker = await ServiceSeqTracker.Init( @@ -47,10 +47,7 @@ const start = async () => { heartbeat, filter: "_selector", selector: { - $or: [ - { type: "message" }, - { type: "service", service: "message-delivery-method-web-push" }, - ], + type: "message", }, }, ); @@ -76,14 +73,14 @@ const start = async () => { for await (const change of changesStream) { const { seq, doc: msgDoc } = change; - if (!msgDoc) { - console.error("Change missing document:", change); - continue; - } - - serviceSeqTracker.checkNextSeq(seq); - try { + if (!msgDoc) { + console.error("Change missing document:", change); + continue; + } + + serviceSeqTracker.checkNextSeq(seq); + const qrDoc = await qrDb.get(`qr_code:${msgDoc.qr_code_uri}`); const timestamp = new Date(msgDoc.created_at).getTime(); @@ -264,7 +261,7 @@ const start = async () => { ); } catch (error) { console.error( - `Failed processing: ${JSON.stringify(error)} change: ${JSON.stringify({ id: msgDoc._id, rev: msgDoc._rev, seq: seq })}`, + `Failed processing: ${JSON.stringify(error)} change: ${JSON.stringify({ id: msgDoc?._id, rev: msgDoc?._rev, seq: seq })}`, ); } finally { await serviceSeqTracker.saveSeq(seq); diff --git a/services/message-delivery-method-web-push/tsconfig.json b/services/message-delivery-method-web-push/tsconfig.json index de282de..ba1efad 100644 --- a/services/message-delivery-method-web-push/tsconfig.json +++ b/services/message-delivery-method-web-push/tsconfig.json @@ -11,9 +11,6 @@ "declaration": true, // Генерация `.d.ts` файлов "sourceMap": true, // Генерация `.map` файлов "baseUrl": "./", // Базовая директория для путей - "paths": { - "*": ["src/*"] // Настройка алиасов (сопоставление путей) - } }, "include": ["src"], // Указываем папку с исходниками "exclude": ["node_modules", "dist"] // Исключаем ненужные директории diff --git a/services/tg-bot/package.json b/services/tg-bot/package.json index be1a016..ae4c6eb 100644 --- a/services/tg-bot/package.json +++ b/services/tg-bot/package.json @@ -2,7 +2,7 @@ "name": "@hereconnect/service-tg-bot", "version": "1.0.0", "description": "", - "main": "dist/server.ts", + "main": "dist/service.js", "type": "module", "license": "@hereconnect/license", "private": true, @@ -29,17 +29,18 @@ "contributors": [], "scripts": { "build": "tsc", - "dev": "tsx src/bot/index.ts", - "dev:watch": "tsx --watch src/bot/index.ts", + "dev": "tsx src/service.ts", + "dev:watch": "tsx --watch src/service.ts", "type-check": "tsc --noEmit", "clean": "rm -rf dist", "lint": "eslint src", "lint:fix": "eslint src --fix", "validate": "npm run type-check && npm run lint", - "start": "node --import=extensionless/register dist/bot/index.js" + "start": "node --import=extensionless/register dist/service.js" }, "dependencies": { "@grammyjs/conversations": "^1.2.0", + "@grammyjs/parse-mode": "^1.10.0", "@hereconnect/constants": "workspace:^", "@hereconnect/couchdb-changes-stream": "workspace:^", "@hereconnect/lib-service-health-checkserver": "workspace:^", @@ -71,6 +72,7 @@ "@typescript-eslint/parser": "^8.18.0", "eslint": "^9.16.0", "tsx": "^4.19.2", + "type-fest": "^4.30.1", "typescript": "^5.7.2" } } diff --git a/services/tg-bot/src/bot/conversations/newQRCodeConversation.ts b/services/tg-bot/src/bot/conversations/newQRCodeConversation.ts index 4edbf3c..24d0778 100644 --- a/services/tg-bot/src/bot/conversations/newQRCodeConversation.ts +++ b/services/tg-bot/src/bot/conversations/newQRCodeConversation.ts @@ -2,7 +2,7 @@ import { type QRCodeDocument } from "@hereconnect/types"; import { Conversation } from "@grammyjs/conversations"; import { InputFile } from "grammy"; import type { MyContext } from "../../types/myContext"; -import { qrCodePlacements } from "../../models/types"; +import { qrCodePlacements } from "../../types/models"; import { buildPlacementKeyboard, buildActionsKeyboard, @@ -11,11 +11,8 @@ import { createQRCode } from "../../api/qrCodes"; import { generateQRCodeImage } from "../../services/qrGenerator"; import { selectMultipleActions } from "./selectMultipleActions"; -// Тип разговора (конверсии) -type MyConversation = Conversation; - export async function newQRCodeConversation( - conversation: MyConversation, + conversation: Conversation, ctx: MyContext, ) { if (!ctx.userDoc) { diff --git a/services/tg-bot/src/bot/index.ts b/services/tg-bot/src/bot/index.ts index b26139f..54ae25b 100644 --- a/services/tg-bot/src/bot/index.ts +++ b/services/tg-bot/src/bot/index.ts @@ -1,64 +1,59 @@ import { Bot, session } from "grammy"; import { conversations, createConversation } from "@grammyjs/conversations"; -import { startHealthServer } from "@hereconnect/lib-service-health-checkserver"; import { MyContext } from "../types/myContext"; import { env } from "../config/env"; import { CouchDBStorageAdapter } from "./storageAdapters/CouchDBStorageAdapter"; import { loadUserMiddleware } from "./middleware/loadUser"; import { tgBotSessionsDb } from "../dbs"; -import { setupCommands } from "./commands/setupCommands"; import { startCommand } from "./commands/start"; import { helpCommand } from "./commands/help"; import { newCommand } from "./commands/new"; import { manageCommand } from "./commands/manage"; import { callbackQueryHandler } from "./callbackHandlers"; import { newQRCodeConversation } from "./conversations/newQRCodeConversation"; +import { hydrateReply } from "@grammyjs/parse-mode"; -const abortController = new AbortController(); +export { setupCommands } from "./commands/setupCommands"; -const storage = new CouchDBStorageAdapter(tgBotSessionsDb, { - softDelete: true, - type: "tg_bot_session", -}); +export const createBot = ({ + abortController, +}: { + abortController: AbortController; +}) => { + const storage = new CouchDBStorageAdapter(tgBotSessionsDb, { + softDelete: true, + type: "tg_bot_session", + }); -// Инициализируем бота с нужным типом контекста -const bot = new Bot(env.TELEGRAM_BOT_TOKEN, { - client: { - baseFetchConfig: { - signal: abortController.signal, + // Инициализируем бота с нужным типом контекста + const bot = new Bot(env.TELEGRAM_BOT_TOKEN, { + client: { + baseFetchConfig: { + signal: abortController.signal, + }, }, - }, -}); + }); -// Подключаем сессии. По необходимости определите initial state для сессии -bot.use(session({ storage, initial: () => ({}) })); + bot.use(hydrateReply); -// Подключаем conversations для многошаговых сценариев -bot.use(conversations()); -bot.use(createConversation(newQRCodeConversation, "newQRCodeConversation")); + // Подключаем сессии. По необходимости определите initial state для сессии + bot.use(session({ storage, initial: () => ({}) })); -// Middleware для загрузки или создания пользователя -bot.use(loadUserMiddleware); + // Подключаем conversations для многошаговых сценариев + bot.use(conversations()); + bot.use(createConversation(newQRCodeConversation)); -// Регистрация команд -bot.command("start", startCommand); -bot.command("help", helpCommand); -bot.command("new", newCommand); -bot.command("manage", manageCommand); + // Middleware для загрузки или создания пользователя + bot.use(loadUserMiddleware); -// Обработка callback_query -bot.on("callback_query:data", callbackQueryHandler); + // Регистрация команд + bot.command("start", startCommand); + bot.command("help", helpCommand); + bot.command("new", newCommand); + bot.command("manage", manageCommand); -(async () => { - // Установка команд - await setupCommands(bot); + // Обработка callback_query + bot.on("callback_query:data", callbackQueryHandler); - // Запуск бота - bot.start(); - - await startHealthServer(abortController); - - console.log("Бот запущен!"); -})(); - -export { bot }; + return bot; +}; diff --git a/services/tg-bot/src/config/env.ts b/services/tg-bot/src/config/env.ts index 92cfd57..e8ee231 100644 --- a/services/tg-bot/src/config/env.ts +++ b/services/tg-bot/src/config/env.ts @@ -10,8 +10,14 @@ if (!process.env.API_URL) { throw new Error("Environment variable API_URL is not set"); } +if (!process.env.HEARTBEAT_INTERVAL) { + throw new Error("Environment variable HEARTBEAT_INTERVAL is not set"); +} + export const env = { TELEGRAM_BOT_TOKEN: process.env.BOT_TOKEN, READ_BASE_URL: process.env.READ_BASE_URL, + APP_BASE_URL: process.env.APP_BASE_URL, API_URL: process.env.API_URL, + HEARTBEAT_INTERVAL: parseInt(process.env.HEARTBEAT_INTERVAL, 10), } as const; diff --git a/services/tg-bot/src/dbs.ts b/services/tg-bot/src/dbs.ts index 41e2de4..612770b 100644 --- a/services/tg-bot/src/dbs.ts +++ b/services/tg-bot/src/dbs.ts @@ -2,7 +2,6 @@ import PouchDB from "./PouchDB"; import { ServiceSeqTrackerDatabase } from "@hereconnect/lib-service-seq-tracker"; import type { - MessageDocument, QRCodeDocument, MessageDeliveryMethodWebPushVapid, UserChatSettingsDocument, @@ -11,25 +10,26 @@ import type { UserDocument, } from "@hereconnect/types"; +import type { ServiceSeqTrackerDocument } from "@hereconnect/lib-service-seq-tracker"; + import type { TelegramWebhookDesignDocument } from "@hereconnect/tg-bot-webhook-ddoc"; import { CouchDBSessionDocument } from "./bot/storageAdapters/CouchDBStorageAdapter"; import type { SessionData } from "./types/myContext"; import { env } from "./config/env"; -// import type { TgUpdateDocument } from "./types"; if (!env.API_URL) { throw new Error("Environment variable API_URL is not set"); } export const TG_BOT_DB_NAME = "tg_bot"; +export const MESSAGES_DB_NAME = "messages"; export const TG_BOT_DB_URL = `${env.API_URL}/${TG_BOT_DB_NAME}`; +export const MESSAGES_DB_URL = `${env.API_URL}/${MESSAGES_DB_NAME}`; -export const messagesDb = new PouchDB( - `${env.API_URL}/messages`, - { +export const messagesServiceSeqTrackerDb = + new PouchDB(MESSAGES_DB_URL, { skip_setup: true, - }, -); + }); export const serverSettingsDb = new PouchDB( `${env.API_URL}/server_settings`, @@ -49,18 +49,10 @@ export const qrDb = new PouchDB(`${env.API_URL}/qr`, { skip_setup: true, }); -export const webPushSubscriptionsDb = new PouchDB< - WebPushSubscriptionDocument ->(`${env.API_URL}/web_push_subscriptions`, { - skip_setup: true, -}); - export const tgBotDb = new PouchDB(TG_BOT_DB_URL, { skip_setup: true, }); -export const tgBotDesignDocumentsDb = - tgBotDb as PouchDB.Database; export const serviceSeqTrackerDb = tgBotDb as ServiceSeqTrackerDatabase; export const tgBotSessionsDb = new PouchDB>( @@ -81,20 +73,25 @@ export const usersDb = new PouchDB(`${env.API_URL}/_users`, { skip_setup: true, }); -export const checkDatabases = async () => { - const entries = await Promise.all( - Object.entries({ - serverSettingsDb, - publicSettingsDb, - qrDb, - webPushSubscriptionsDb, - serviceSeqTrackerDb, - userChatSettingsDb, - tgBotDb, - tgBotSessionsDb, - } as const).map( +const checkDatabases = async () => { + const databases = { + serverSettingsDb, + publicSettingsDb, + qrDb, + usersDb, + serviceSeqTrackerDb, + userChatSettingsDb, + tgBotDb, + tgBotSessionsDb, + messagesServiceSeqTrackerDb, + } as const; + + const databasesEntries = Object.entries(databases); + + const resultEntries = await Promise.all( + databasesEntries.map( async ([name, db]): Promise< - [string, Awaited>] + [typeof name, Awaited>] > => { const info = await db.info(); return [name, info]; @@ -102,7 +99,7 @@ export const checkDatabases = async () => { ), ); - for (const [name, info] of entries) { + for (const [name, info] of resultEntries) { if ("error" in info && "reason" in info) { console.error(info); throw new Error(`[${name}] ${info.error}: ${info.reason}`); @@ -112,5 +109,7 @@ export const checkDatabases = async () => { } } - return Object.fromEntries(entries); + return Object.fromEntries(resultEntries); }; + +export const checkDatabasesPromise = checkDatabases(); diff --git a/services/tg-bot/src/global.d.ts b/services/tg-bot/src/global.d.ts new file mode 100644 index 0000000..84940a6 --- /dev/null +++ b/services/tg-bot/src/global.d.ts @@ -0,0 +1,8 @@ +import type { Entries, Entry } from "type-fest"; + +declare global { + interface ObjectConstructor { + entries(o: T): Entries; + fromEntries(entries: Entry[]): T; + } +} diff --git a/services/tg-bot/src/service.ts b/services/tg-bot/src/service.ts new file mode 100644 index 0000000..6f33c26 --- /dev/null +++ b/services/tg-bot/src/service.ts @@ -0,0 +1,31 @@ +import { createHealthServer } from "@hereconnect/lib-service-health-checkserver"; +import { createBot, setupCommands } from "./bot"; +import { createServiceChangesMessagesDb } from "./services/serviceChangesMessagesDb"; +import { iterateMessagesChanges } from "./services/iterateMessagesChanges"; + +const abortController = new AbortController(); + +(async () => { + const healthServer = createHealthServer(abortController); + const bot = createBot({ abortController }); + const { changesStream, serviceSeqTracker } = + await createServiceChangesMessagesDb({ + abortController, + }); + + process.on("SIGINT", async () => { + changesStream.stop(); + bot.stop(); + }); + + await healthServer.start(); + + await bot.init(); + // Устанавливаем команд + await setupCommands(bot); + + bot.start(); + console.log("Бот запущен!"); + + iterateMessagesChanges({ changesStream, serviceSeqTracker, bot }); +})(); diff --git a/services/tg-bot/src/services/iterateMessagesChanges.ts b/services/tg-bot/src/services/iterateMessagesChanges.ts new file mode 100644 index 0000000..8400643 --- /dev/null +++ b/services/tg-bot/src/services/iterateMessagesChanges.ts @@ -0,0 +1,111 @@ +import type { MyContext } from "../types/myContext"; +import { Bot, Context, Keyboard } from "grammy"; +import type { ServiceSeqTracker } from "@hereconnect/lib-service-seq-tracker"; +import type { CouchDBChangesStream } from "@hereconnect/couchdb-changes-stream"; +import type { MessageDocument, UserDocument } from "@hereconnect/types"; +import { qrDb, userChatSettingsDb, usersDb } from "../dbs"; +// import { blockquote, fmt, italic, link } from "@grammyjs/parse-mode"; +import { env } from "../config/env"; +import { escapeHtml } from "../utils/escapeHtml"; + +async function handleChange(msgDoc: MessageDocument, bot: Bot) { + const qrDoc = await qrDb.get(`qr_code:${msgDoc.qr_code_uri}`); + if (qrDoc.messageDeliveryMethod !== "telegram") { + console.log(`Skip !telegram`); + return; + } + + const fields = ["telegram_id"] as const; + const [{ docs: telegramUsersDocs }, { rows: userChatSettingsRows }] = + await Promise.all([ + usersDb.find({ + selector: { + user_uuid: qrDoc.user_uuid, + telegram_id: { + $exists: true, + $ne: null, + }, + }, + fields: [...fields], + }) as Promise<{ + docs: Required>[]; + }>, + + userChatSettingsDb.allDocs({ + startkey: `user_chat_settings:${msgDoc.qr_code_uri}:${msgDoc.chat}:owner:`, + endkey: `user_chat_settings:${msgDoc.qr_code_uri}:${msgDoc.chat}:owner:\uffff`, + include_docs: true, + }), + ]); + + if ( + userChatSettingsRows.some( + (row) => row.doc?.is_telegram_messages_enabled === false, + ) + ) { + console.log("telegram notifications is disabled"); + return; + } + + if (telegramUsersDocs.length === 0) { + console.log("NO TELEGRAM USERS"); + } else { + console.log(`sending to ${telegramUsersDocs.length} telegram accounts`); + } + + const hashTag = `#qr_${qrDoc.uri}_c${msgDoc.chat}`; + const manageMessageUrl = `${env.APP_BASE_URL}/manage/qr/${qrDoc.uri}/chat/${msgDoc.chat}#message-${encodeURIComponent(msgDoc.created_at)}`; + // const manageMessageUrl = `tg://resolve?domain=${bot.botInfo.username}&start=message:${qrDoc.uri}:${msgDoc.created_at}`; + // const msgCaption = ` — ${link("сообщение", manageMessageUrl)} QR-кода «${qrDoc.name}»`; + + // const caption = fmt`${link("сообщение", manageMessageUrl)} QR-кода «${qrDoc.name}»`; + // const body = fmt`${msgDoc.body.trim()} — ${italic(caption)}`; + + const body = { + text: `${escapeHtml(msgDoc.body.trim())} — сообщение ${escapeHtml(hashTag)} QR-кода «${escapeHtml(qrDoc.name)}»`, + entities: undefined, + }; + + for (const telegramUserDoc of telegramUsersDocs) { + await bot.api.sendMessage(telegramUserDoc.telegram_id, body.text, { + parse_mode: "HTML", + // parse_mode: "MarkdownV2", + link_preview_options: { is_disabled: true }, + entities: body.entities, + reply_markup: { force_reply: true }, + }); + } +} + +export const iterateMessagesChanges = async ({ + changesStream, + serviceSeqTracker, + bot, +}: { + changesStream: CouchDBChangesStream; + serviceSeqTracker: ServiceSeqTracker; + bot: Bot; +}) => { + for await (const change of changesStream) { + const { id, seq, doc: msgDoc, changes } = change; + + console.log("seq", seq); + + try { + if (!msgDoc) { + console.error("Change missing document:", change); + continue; + } + + serviceSeqTracker.checkNextSeq(seq); + + await handleChange(msgDoc, bot); + } catch (error) { + console.error( + `Failed processing: ${JSON.stringify(error)} change: ${JSON.stringify({ id, rev: changes.map((c) => c.rev).join(","), seq: seq })}`, + ); + } finally { + await serviceSeqTracker.saveSeq(seq); + } + } +}; diff --git a/services/tg-bot/src/services/serviceChangesMessagesDb.ts b/services/tg-bot/src/services/serviceChangesMessagesDb.ts new file mode 100644 index 0000000..754096e --- /dev/null +++ b/services/tg-bot/src/services/serviceChangesMessagesDb.ts @@ -0,0 +1,52 @@ +import { ServiceSeqTracker } from "@hereconnect/lib-service-seq-tracker"; +import { CouchDBChangesStream } from "@hereconnect/couchdb-changes-stream"; +import type { MessageDocument } from "@hereconnect/types"; +import { + MESSAGES_DB_URL, + messagesServiceSeqTrackerDb, + checkDatabasesPromise, +} from "../dbs"; +import { env } from "../config/env"; + +export const createServiceChangesMessagesDb = async ({ + abortController, +}: { + abortController: AbortController; +}) => { + const { + messagesServiceSeqTrackerDb: { update_seq: lastSeq }, + } = await checkDatabasesPromise; + + const serviceSeqTracker = await ServiceSeqTracker.Init( + "_local/service:tg-bot", + messagesServiceSeqTrackerDb, + ); + + const changesStream = new CouchDBChangesStream( + MESSAGES_DB_URL, + { + abortController, + since: serviceSeqTracker.lastSeq, + feed: "continuous", + include_docs: true, + heartbeat: env.HEARTBEAT_INTERVAL, + filter: "_selector", + selector: { + type: "message", + }, + }, + ); + + console.info("Service started from sequence", serviceSeqTracker.lastSeq); + + if ( + parseInt(serviceSeqTracker.lastSeq.toString(), 10) === + parseInt(lastSeq.toString(), 10) + ) { + console.log("Since is the last seq"); + } else { + console.log("Last seq is", lastSeq); + } + + return { changesStream, serviceSeqTracker }; +}; diff --git a/services/tg-bot/src/models/types.ts b/services/tg-bot/src/types/models.ts similarity index 100% rename from services/tg-bot/src/models/types.ts rename to services/tg-bot/src/types/models.ts diff --git a/services/tg-bot/src/types/myContext.ts b/services/tg-bot/src/types/myContext.ts index 969a098..23228d2 100644 --- a/services/tg-bot/src/types/myContext.ts +++ b/services/tg-bot/src/types/myContext.ts @@ -1,7 +1,8 @@ // src/types/myContext.ts -import { Context, SessionFlavor } from "grammy"; -import { ConversationFlavor } from "@grammyjs/conversations"; +import type { Context, SessionFlavor } from "grammy"; +import type { ConversationFlavor, Conversation } from "@grammyjs/conversations"; +import type { ParseModeFlavor } from "@grammyjs/parse-mode"; import type { UserDocument } from "@hereconnect/types"; export interface SessionData { @@ -9,11 +10,11 @@ export interface SessionData { } // Базовый контекст без conversation -interface MyContextBase extends Context, SessionFlavor { +interface MyContextBase { userDoc?: UserDocument; } -// Теперь создаём финальный тип контекста, дополняя его ConversationFlavor -type MyContext = MyContextBase & ConversationFlavor; - -export { MyContext }; +export type MyContext = ParseModeFlavor & + SessionFlavor & + ConversationFlavor & + MyContextBase; diff --git a/services/tg-bot/src/utils/escapeHtml.ts b/services/tg-bot/src/utils/escapeHtml.ts new file mode 100644 index 0000000..d6a9244 --- /dev/null +++ b/services/tg-bot/src/utils/escapeHtml.ts @@ -0,0 +1,21 @@ +import { DOMParser, XMLSerializer } from "@xmldom/xmldom"; + +/** + * + * Экранирует HTML-код, заменяя символы "<", ">", "&" на их соответствующие entity-кодировки. + * + * @param text + */ +/** + * Экранирует HTML-код, заменяя символы "<", ">", "&", '"' и "'" на их соответствующие entity-кодировки. + * + * @param text + */ +export function escapeHtml(text: string): string { + return text + .replace(/&/g, "&") + .replace(//g, ">") + .replace(/"/g, """) + .replace(/'/g, "'"); +} diff --git a/services/tg-bot/tsconfig.json b/services/tg-bot/tsconfig.json index 3e73970..3fe645d 100644 --- a/services/tg-bot/tsconfig.json +++ b/services/tg-bot/tsconfig.json @@ -11,6 +11,6 @@ "declaration": true, // Генерация `.d.ts` файлов "sourceMap": true, // Генерация `.map` файлов }, - "include": ["src"], // Указываем папку с исходниками - "exclude": ["node_modules", "dist", "old"] // Исключаем ненужные директории + "include": ["globals.d.ts", "."], // Указываем папку с исходниками + "exclude": ["node_modules", "dist", "old"], // Исключаем ненужные директории }