Files
iistwin/server/vpn/vpn.db.ts
Ильяс Султанов aea0e9c3fe feat(vpn): монитор актуальности Яндекс Телемост + инструкция про выход из встречи
- INSTRUCTIONS_TEXT: добавлено требование выйти из встречи, а не завершать её
- vpn.db.ts: findAllActiveVpnTasks для поиска активных VPN-задач
- vpn.service.ts: checkVpnRoomsValidity — периодическая проверка комнат, concurrency=5
- vpn-bot.service.ts: findVpnBotConversationId и notifyVpnRoomExpired
- server/index.ts: worker каждые 6 часов, уведомление пользователю о протухшей ссылке
2026-07-17 13:08:24 +03:00

369 lines
10 KiB
TypeScript

import { db } from '../db';
import {
formFields,
formStatuses,
tasks,
taskFieldValues,
type Task,
type FormField,
type FormStatus,
} from '@shared/schema';
import { eq, and, sql } from 'drizzle-orm';
import { storage } from '../storage';
import { VPN_FORM_ID } from './vpn.config';
export interface VpnFormCache {
fields: FormField[];
statuses: FormStatus[];
fieldByCode: Map<string, FormField>;
statusByName: Map<string, FormStatus>;
}
const cacheByOrg = new Map<number, VpnFormCache>();
export async function getVpnFormCache(organizationId: number): Promise<VpnFormCache> {
if (cacheByOrg.has(organizationId)) {
return cacheByOrg.get(organizationId)!;
}
const [fields, statuses] = await Promise.all([
db
.select()
.from(formFields)
.where(eq(formFields.formId, VPN_FORM_ID))
.then(rows => rows as FormField[]),
db
.select()
.from(formStatuses)
.where(eq(formStatuses.formId, VPN_FORM_ID))
.then(rows => rows as FormStatus[]),
]);
const cache: VpnFormCache = {
fields,
statuses,
fieldByCode: new Map(fields.map(f => [f.code, f])),
statusByName: new Map(statuses.map(s => [s.name, s])),
};
cacheByOrg.set(organizationId, cache);
return cache;
}
export function getFieldId(cache: VpnFormCache, code: string): number {
const f = cache.fieldByCode.get(code);
if (!f) throw new Error(`VPN field not found: ${code}`);
return f.id;
}
export function getStatusId(cache: VpnFormCache, name: string): number {
const s = cache.statusByName.get(name);
if (!s) throw new Error(`VPN status not found: ${name}`);
return s.id;
}
export function clearVpnFormCache(organizationId: number): void {
cacheByOrg.delete(organizationId);
}
export interface VpnTaskData {
task: Task;
values: Record<string, unknown>;
}
export async function findVpnTaskByToken(
token: string,
organizationId: number,
): Promise<VpnTaskData | null> {
const cache = await getVpnFormCache(organizationId);
const fieldId = getFieldId(cache, 'subscription_token');
const rows = await db
.select({ task: tasks, value: taskFieldValues.value })
.from(taskFieldValues)
.innerJoin(tasks, eq(taskFieldValues.taskId, tasks.id))
.where(
and(
eq(taskFieldValues.fieldId, fieldId),
eq(taskFieldValues.formId, VPN_FORM_ID),
sql`${taskFieldValues.value} = to_jsonb(${token}::text)`,
eq(tasks.organizationId, organizationId),
),
)
.limit(1);
if (!rows.length) return null;
return loadVpnTaskValues(rows[0].task.id, organizationId);
}
export async function findVpnTaskByRoomUrl(
roomUrl: string,
organizationId: number,
): Promise<VpnTaskData | null> {
const cache = await getVpnFormCache(organizationId);
const fieldId = getFieldId(cache, 'room_url');
const rows = await db
.select({ task: tasks })
.from(taskFieldValues)
.innerJoin(tasks, eq(taskFieldValues.taskId, tasks.id))
.where(
and(
eq(taskFieldValues.fieldId, fieldId),
eq(taskFieldValues.formId, VPN_FORM_ID),
sql`${taskFieldValues.value} = to_jsonb(${roomUrl}::text)`,
eq(tasks.organizationId, organizationId),
),
)
.limit(1);
if (!rows.length) return null;
return loadVpnTaskValues(rows[0].task.id, organizationId);
}
export async function findVpnTaskByRoomAndDevice(
roomUrl: string,
deviceId: string,
organizationId: number,
): Promise<VpnTaskData | null> {
const cache = await getVpnFormCache(organizationId);
const roomFieldId = getFieldId(cache, 'room_url');
const rows = await db
.select({ task: tasks })
.from(taskFieldValues)
.innerJoin(tasks, eq(taskFieldValues.taskId, tasks.id))
.where(
and(
eq(taskFieldValues.fieldId, roomFieldId),
eq(taskFieldValues.formId, VPN_FORM_ID),
sql`${taskFieldValues.value} = to_jsonb(${roomUrl}::text)`,
eq(tasks.organizationId, organizationId),
),
)
.limit(1);
if (!rows.length) return null;
const data = await loadVpnTaskValues(rows[0].task.id, organizationId);
if (data.values.hwid !== deviceId) return null;
return data;
}
export async function findVpnTaskByRoomAndSession(
roomUrl: string,
sessionId: string,
organizationId: number,
): Promise<VpnTaskData | null> {
const cache = await getVpnFormCache(organizationId);
const roomFieldId = getFieldId(cache, 'room_url');
const rows = await db
.select({ task: tasks })
.from(taskFieldValues)
.innerJoin(tasks, eq(taskFieldValues.taskId, tasks.id))
.where(
and(
eq(taskFieldValues.fieldId, roomFieldId),
eq(taskFieldValues.formId, VPN_FORM_ID),
sql`${taskFieldValues.value} = to_jsonb(${roomUrl}::text)`,
eq(tasks.organizationId, organizationId),
),
)
.limit(1);
if (!rows.length) return null;
const data = await loadVpnTaskValues(rows[0].task.id, organizationId);
if (data.values.active_session_id !== sessionId) return null;
return data;
}
export async function loadVpnTaskValues(
taskId: number,
organizationId: number,
): Promise<VpnTaskData> {
const cache = await getVpnFormCache(organizationId);
const [task, fvs] = await Promise.all([
storage.getTask(taskId, organizationId),
storage.getTaskFieldValues(taskId, organizationId),
]);
if (!task) throw new Error(`VPN task not found: ${taskId}`);
const values: Record<string, unknown> = {};
for (const fv of fvs) {
const field = cache.fields.find(f => f.id === fv.fieldId);
if (field) {
values[field.code] = fv.value;
}
}
return { task, values };
}
export async function setVpnTaskField(
taskId: number,
fieldCode: string,
value: unknown,
organizationId: number,
): Promise<void> {
const cache = await getVpnFormCache(organizationId);
const fieldId = getFieldId(cache, fieldCode);
const existing = await db
.select({ id: taskFieldValues.id })
.from(taskFieldValues)
.where(and(eq(taskFieldValues.taskId, taskId), eq(taskFieldValues.fieldId, fieldId)))
.limit(1);
if (existing.length) {
await storage.updateTaskFieldValue(taskId, fieldId, organizationId, { value: value as any });
} else {
await storage.createTaskFieldValue({
taskId,
fieldId,
formId: VPN_FORM_ID,
value: value as any,
});
}
}
export async function setVpnTaskStatus(
taskId: number,
statusName: string,
organizationId: number,
): Promise<void> {
const cache = await getVpnFormCache(organizationId);
const statusId = getStatusId(cache, statusName);
const status = cache.statusByName.get(statusName)!;
await storage.updateTask(taskId, organizationId, {
currentStatusId: statusId,
isCompleted: status.isFinal,
});
}
export async function createVpnTask(
input: {
organizationId: number;
createdBy: number;
title: string;
statusName: string;
fields: Record<string, unknown>;
},
): Promise<Task> {
const cache = await getVpnFormCache(input.organizationId);
const statusId = getStatusId(cache, input.statusName);
return db.transaction(async (tx) => {
const [task] = await tx
.insert(tasks)
.values({
organizationId: input.organizationId,
formId: VPN_FORM_ID,
title: input.title,
currentStatusId: statusId,
createdBy: input.createdBy,
isCompleted: false,
})
.returning();
const fieldInserts = Object.entries(input.fields)
.filter(([, value]) => value !== undefined && value !== null)
.map(([code, value]) => ({
taskId: task.id,
fieldId: getFieldId(cache, code),
formId: VPN_FORM_ID,
value: value as any,
}));
if (fieldInserts.length > 0) {
await tx.insert(taskFieldValues).values(fieldInserts);
}
return task;
});
}
export async function updateVpnTaskFieldMap(
taskId: number,
organizationId: number,
fields: Record<string, unknown>,
): Promise<void> {
for (const [code, value] of Object.entries(fields)) {
if (value === undefined) continue;
await setVpnTaskField(taskId, code, value, organizationId);
}
}
export interface ActiveVpnTaskInfo {
taskId: number;
organizationId: number;
createdBy: number;
roomUrl: string;
}
export async function findAllActiveVpnTasks(): Promise<ActiveVpnTaskInfo[]> {
// 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;
}