import { db } from '../db'; import { formFields, formStatuses, tasks, taskFieldValues, type Task, type FormField, type FormStatus, } from '@shared/schema'; import { eq, and, sql, isNull } from 'drizzle-orm'; import { storage } from '../storage'; import { VPN_FORM_ID } from './vpn.config'; export interface VpnFormCache { fields: FormField[]; statuses: FormStatus[]; fieldByCode: Map; statusByName: Map; } const cacheByOrg = new Map(); export async function getVpnFormCache(organizationId: number): Promise { if (cacheByOrg.has(organizationId)) { return cacheByOrg.get(organizationId)!; } const [fields, statuses] = await Promise.all([ db .select() .from(formFields) .where(and(eq(formFields.formId, VPN_FORM_ID), isNull(formFields.deletedAt))) .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; } export async function findVpnTaskByToken( token: string, organizationId: number, ): Promise { 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 { 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 { 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 { 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 { 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 = {}; 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 { 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 { 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; }, ): Promise { 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, ): Promise { 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 { // 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; }