diff --git a/server/index.ts b/server/index.ts index 675f3f1..2c8146e 100644 --- a/server/index.ts +++ b/server/index.ts @@ -1179,6 +1179,35 @@ async function runStartupDataPatches() { }, msUntil23Msk); log(`Embedding queue worker scheduled at 23:00 MSK (in ${Math.round(msUntil23Msk / 1000 / 60)} minutes)`); + // VPN room validity worker — checks Yandex Telemost rooms every 6 hours + async function runVpnRoomValidityWorker() { + try { + const { checkVpnRoomsValidity } = await import('./vpn/vpn.service'); + const { notifyVpnRoomExpired } = await import('./vpn/vpn-bot.service'); + const expired = await checkVpnRoomsValidity(); + if (expired.length === 0) { + log('[VPN] Room validity check done: no expired rooms'); + return; + } + log(`[VPN] Room validity check done: ${expired.length} expired room(s)`); + for (const room of expired) { + try { + await notifyVpnRoomExpired(room.organizationId, room.createdBy, room.roomUrl); + } catch (err) { + console.error(`[VPN] Failed to notify user ${room.createdBy} about expired room:`, err); + } + } + } catch (err) { + log(`[VPN] Room validity worker error: ${err instanceof Error ? err.message : String(err)}`); + } + } + // Run once 2 minutes after startup, then every 6 hours + setTimeout(() => { + runVpnRoomValidityWorker(); + setInterval(runVpnRoomValidityWorker, 6 * 60 * 60 * 1000); + }, 120_000); + log('VPN room validity worker started (6h interval)'); + // ALWAYS serve the app on the port specified in the environment variable PORT // Other ports are firewalled. Default to 5000 if not specified. // this serves both the API and the client. diff --git a/server/vpn/vpn-bot.service.ts b/server/vpn/vpn-bot.service.ts index b5578d7..13ddba1 100644 --- a/server/vpn/vpn-bot.service.ts +++ b/server/vpn/vpn-bot.service.ts @@ -2,7 +2,7 @@ import fs from "fs"; import path from "path"; import { db } from "../db"; import { conversationMessages, conversationMembers, bots, type Bot, type User } from "@shared/schema"; -import { eq, and } from "drizzle-orm"; +import { eq, and, sql } from "drizzle-orm"; import { eventBus } from "../routes/shared"; import { enrichMessage, getMembersForSSE } from "../routes/messenger.helpers"; import { notificationService } from "../services/notification.service"; @@ -20,7 +20,8 @@ const INSTRUCTIONS_TEXT = `Привет! Для получения корпор 1. Открой Яндекс Телемост (приложение или https://telemost.yandex.ru). 2. Создай новую видеовстречу. -3. Пришли сюда ссылку на встречу. Примеры того, что подходит: +3. **Важно:** после создания **выйди** из встречи, а не **завершай** её. Тогда ссылка останется активной и будет работать для VPN. +4. Пришли сюда ссылку на встречу. Примеры того, что подходит: • https://telemost.yandex.ru/j/abc123def456 • telemost.yandex.ru/j/abc123def456 @@ -316,3 +317,46 @@ export async function handleVpnBotCallback({ organizationId, ); } + +// ── Room expiry notifications ───────────────────────────────────────────────── + +async function findVpnBotConversationId( + organizationId: number, + userId: number, +): Promise { + const botId = await getVpnBotId(organizationId); + if (!botId) return null; + + const rows = await db.execute(sql` + SELECT c.id + FROM conversations c + JOIN conversation_members cm ON cm.conversation_id = c.id + WHERE c.type = 'bot_direct' + AND c.bot_id = ${botId} + AND c.organization_id = ${organizationId} + AND cm.user_id = ${userId} + LIMIT 1 + `); + + interface Row { id: number } + const row = rows.rows[0] as unknown as Row | undefined; + return row?.id ?? null; +} + +export async function notifyVpnRoomExpired( + organizationId: number, + userId: number, + roomUrl: string, +): Promise { + const conversationId = await findVpnBotConversationId(organizationId, userId); + if (!conversationId) { + console.warn(`[VPN Bot] no bot_direct conversation found for user ${userId} in org ${organizationId}`); + return; + } + + await sendVpnBotMessage( + conversationId, + `⚠️ Ссылка на встречу Яндекс Телемост больше недействительна:\n${roomUrl}\n\nСоздайте новую встречу в Яндекс Телемост (выйдите из неё, не завершая) и пришлите новую ссылку.`, + organizationId, + ); +} diff --git a/server/vpn/vpn.db.ts b/server/vpn/vpn.db.ts index 0e5700f..f91e668 100644 --- a/server/vpn/vpn.db.ts +++ b/server/vpn/vpn.db.ts @@ -290,3 +290,79 @@ export async function updateVpnTaskFieldMap( await setVpnTaskField(taskId, code, value, organizationId); } } + +export interface ActiveVpnTaskInfo { + taskId: number; + organizationId: number; + createdBy: number; + roomUrl: string; +} + +export async function findAllActiveVpnTasks(): Promise { + // Find all organizations that have VPN tasks + const orgRows = await db + .select({ orgId: tasks.organizationId }) + .from(tasks) + .where(eq(tasks.formId, VPN_FORM_ID)) + .groupBy(tasks.organizationId); + + const results: ActiveVpnTaskInfo[] = []; + + for (const { orgId } of orgRows) { + try { + const cache = await getVpnFormCache(orgId); + const statusFieldId = getFieldId(cache, 'status'); + const activeFieldId = getFieldId(cache, 'is_active'); + const roomUrlFieldId = getFieldId(cache, 'room_url'); + + const rows = await db + .select({ + taskId: tasks.id, + createdBy: tasks.createdBy, + roomUrl: taskFieldValues.value, + }) + .from(tasks) + .innerJoin( + taskFieldValues, + and( + eq(taskFieldValues.taskId, tasks.id), + eq(taskFieldValues.formId, VPN_FORM_ID), + eq(taskFieldValues.fieldId, roomUrlFieldId), + ), + ) + .where( + and( + eq(tasks.formId, VPN_FORM_ID), + eq(tasks.organizationId, orgId), + sql`exists ( + select 1 from task_field_values tfv + where tfv.task_id = tasks.id + and tfv.field_id = ${statusFieldId} + and tfv.value = to_jsonb('Активна'::text) + )`, + sql`exists ( + select 1 from task_field_values tfv + where tfv.task_id = tasks.id + and tfv.field_id = ${activeFieldId} + and tfv.value = to_jsonb(true) + )`, + ), + ); + + for (const row of rows) { + if (typeof row.roomUrl === 'string' && row.roomUrl) { + results.push({ + taskId: row.taskId, + organizationId: orgId, + createdBy: row.createdBy, + roomUrl: row.roomUrl, + }); + } + } + } catch (err) { + console.error(`[VPN DB] findAllActiveVpnTasks error for org ${orgId}:`, err); + } + } + + return results; +} diff --git a/server/vpn/vpn.service.ts b/server/vpn/vpn.service.ts index b5d05bd..9e9fcaf 100644 --- a/server/vpn/vpn.service.ts +++ b/server/vpn/vpn.service.ts @@ -11,6 +11,7 @@ import { createVpnTask, getVpnFormCache, getFieldId, + findAllActiveVpnTasks, type VpnFormCache, } from './vpn.db'; import { @@ -310,3 +311,46 @@ export async function recordTraffic( return; } } + +export interface ExpiredVpnRoom { + taskId: number; + organizationId: number; + createdBy: number; + roomUrl: string; +} + +export async function checkVpnRoomsValidity(): Promise { + const activeTasks = await findAllActiveVpnTasks(); + const expired: ExpiredVpnRoom[] = []; + + // Limit concurrent validations to avoid overwhelming Yandex API + const concurrency = 5; + for (let i = 0; i < activeTasks.length; i += concurrency) { + const batch = activeTasks.slice(i, i + concurrency); + const results = await Promise.all( + batch.map(async (task) => { + try { + const isValid = await validateTelemostRoom(task.roomUrl); + return { task, isValid }; + } catch (err) { + console.error(`[VPN] validation error for task ${task.taskId}:`, err); + return { task, isValid: false }; + } + }), + ); + + for (const { task, isValid } of results) { + if (!isValid) { + expired.push(task); + try { + await setVpnTaskStatus(task.taskId, 'Истекла', task.organizationId); + await setVpnTaskField(task.taskId, 'last_error', 'Комната Яндекс Телемост недоступна', task.organizationId); + } catch (err) { + console.error(`[VPN] failed to mark task ${task.taskId} as expired:`, err); + } + } + } + } + + return expired; +}