diff --git a/IMPLEMENTATION_LOG.md b/IMPLEMENTATION_LOG.md index ea3d3a4..20c5ab0 100644 --- a/IMPLEMENTATION_LOG.md +++ b/IMPLEMENTATION_LOG.md @@ -204,3 +204,29 @@ - Что изменено: `server/routes/task-crud-list.routes.ts` — убраны 7 debug-логов на запрос + per-task лог contract-number; параметр limit (дефолт 20, кап 1000). - Как проверялось: выборка покрыта индексом tasks_org_id_form_id_idx — новый индекс не нужен; vitest 88/88. - Влияние на поиск/UX: нет. + +--- + +## [1.4] SSE: индекс соединений + +- Статус: ✅ done +- Зачем: рассылка событий не перебирает все соединения; один heartbeat вместо таймера на соединение (Фаза 1). +- Что изменено: `server/utils/sse-connection-index.ts` (новый класс: byId/byOrg/byUser); `server/routes/shared.ts` — EventBus рассылает через индекс (таргетинг org/user/user+org/broadcast с прежней tenant-семантикой), один глобальный heartbeat-таймер (unref, стоп при отсутствии соединений), убраны per-event console.log. +- Протокол НЕ тронут: id:/event:/data:, буферы 500, replay по lastEventId — клиенты совместимы. +- Как проверялось: 8 новых юнит-тестов индекса; vitest 96/96; `npm run check` чисто. +- Влияние на поиск/UX: нет. +- Подводные камни: живой прогон SSE (чат в двух вкладках, reconnect по lastEventId) — проверить на проде после деплоя (юнит-тесты покрывают маршрутизацию, не HTTP-поток). + +## [1.5] N+1: recursive CTE, Promise.all, батчи + +- Статус: ✅ done +- Зачем: убрать N+1 на горячих путях (Фаза 1). +- Что изменено: + - `getTaskTree` — один WITH RECURSIVE CTE + один IN-запрос (было 1+2N); защита от циклов (lvl<100, assembled-Set); формат результата прежний. + - Новый `getTaskParentChain` (CTE вверх, max 10) — MCP get_task_tree без цикла getTask. + - Батч-методы: getTasksByIds, getTaskFieldValuesByTaskIds, getTaskMessagesByIds/ByTaskIds. + - sync.routes.ts — /initial и /delta: Promise.all по формам (порядок ответа сохранён). + - embedding.service.ts — processEmbeddingQueue и reindexOrganization: батч-предзагрузки (3–4 запроса на org вместо 2–3 на элемент; 2 запроса на чанк 100 задач вместо ~200). +- Как проверялось: vitest 96/96; `npm run check` чисто; caller'ы дерева (роут /tree, MCP) — формат не изменился. +- Влияние на поиск/UX: нет. +- Подводные камни: raw CTE под PostgreSQL — при переименовании колонок tasks/forms обновить вручную (drizzle не проверяет); getRelatedTasks (граф task_relations) сознательно оставлен — кандидат на отдельный шаг. diff --git a/server/mcp.ts b/server/mcp.ts index 7d550d1..f8267d5 100644 --- a/server/mcp.ts +++ b/server/mcp.ts @@ -4900,14 +4900,11 @@ To block task creation from task.before_create, set: ctx.result = { allow: false }; }; - // Цепочка родителей вверх по parentTaskId (с защитой от циклов) + // Цепочка родителей вверх по parentTaskId (ближайший первым, до 10 уровней, + // один recursive CTE вместо последовательных getTask на каждый уровень) + const parentChain = await storage.getTaskParentChain(taskId, organizationId, 10); const parents: Array> = []; - const seen = new Set([task.id]); - let cursor = task.parentTaskId ?? null; - while (cursor !== null && !seen.has(cursor) && parents.length < 10) { - seen.add(cursor); - const parent = await storage.getTask(cursor, organizationId); - if (!parent) break; + for (const parent of parentChain) { if (isFormAllowed(parent.formId)) { parents.push({ id: parent.id, @@ -4917,7 +4914,6 @@ To block task creation from task.before_create, set: ctx.result = { allow: false isCompleted: parent.isCompleted, }); } - cursor = parent.parentTaskId ?? null; } return { diff --git a/server/routes/shared.ts b/server/routes/shared.ts index 8b2a4d5..00cdf4d 100644 --- a/server/routes/shared.ts +++ b/server/routes/shared.ts @@ -6,6 +6,7 @@ import { db } from '../db'; import { formatUserName } from '../utils/formatUserName'; import { tasks, forms, roles } from '@shared/schema'; import { eq, and, inArray } from 'drizzle-orm'; +import { SseConnectionIndex } from '../utils/sse-connection-index'; // ── Rate-limit helpers ─────────────────────────────────────────────────────── // In offices with a single public IP and many users, IP-only rate limiting @@ -61,7 +62,6 @@ export interface SSEConnection { userId: number; organizationId: number; res: any; // Express Response object - heartbeat?: NodeJS.Timeout; } export interface Event { @@ -78,8 +78,10 @@ export interface BufferedEvent extends Event { } export class EventBus { - private connections: Map = new Map(); + // Индекс соединений по организации/пользователю — publishEvent не перебирает все соединения + private index = new SseConnectionIndex(); private heartbeatInterval = 30000; // 30 seconds + private heartbeatTimer: NodeJS.Timeout | null = null; private eventIdCounter = 0; // Буферы последних событий для восстановления после reconnect private userBuffers = new Map(); @@ -87,32 +89,39 @@ export class EventBus { private readonly BUFFER_SIZE = 500; addConnection(connection: SSEConnection) { - this.connections.set(connection.id, connection); - - // Heartbeat для поддержания соединения - connection.heartbeat = setInterval(() => { - try { - connection.res.write(`: heartbeat\n\n`); - const resAny = connection.res as any; - if (typeof resAny.flush === 'function') { - resAny.flush(); - } - } catch (error) { - this.removeConnection(connection.id); - } - }, this.heartbeatInterval); + this.index.add(connection); + this.ensureHeartbeatTimer(); } removeConnection(connectionId: string) { - const connection = this.connections.get(connectionId); - if (connection) { - if (connection.heartbeat) { - clearInterval(connection.heartbeat); - } - this.connections.delete(connectionId); + this.index.remove(connectionId); + // Когда соединений нет, останавливаем heartbeat-таймер, чтобы не будить event loop впустую + if (this.index.size === 0 && this.heartbeatTimer) { + clearInterval(this.heartbeatTimer); + this.heartbeatTimer = null; } } + // Один глобальный heartbeat-таймер на все соединения (вместо таймера на каждое соединение) + private ensureHeartbeatTimer() { + if (this.heartbeatTimer) return; + this.heartbeatTimer = setInterval(() => { + this.index.forEach((connection) => { + try { + connection.res.write(`: heartbeat\n\n`); + const resAny = connection.res as any; + if (typeof resAny.flush === 'function') { + resAny.flush(); + } + } catch (error) { + this.removeConnection(connection.id); + } + }); + }, this.heartbeatInterval); + // unref, чтобы таймер не удерживал процесс (и не мешал тестам завершаться) + this.heartbeatTimer.unref?.(); + } + private generateEventId(): string { return `${Date.now()}-${++this.eventIdCounter}`; } @@ -163,15 +172,6 @@ export class EventBus { timestamp: Date.now(), }; - // Логируем chat-события для диагностики задержек/пропусков - const isChatEvent = event.type === 'conv_message_created' || event.type === 'conv_message_updated' || - event.type === 'conv_message_deleted' || event.type === 'conv_read' || event.type === 'message_created' || - event.type === 'message_read'; - if (isChatEvent) { - const target = event.userId ? `user=${event.userId}` : `org=${event.organizationId}`; - console.log(`[SSE] publish ${event.type} ${target} conn=${this.connections.size} id=${bufferedEvent.id}`); - } - // Буферизуем событие для восстановления после reconnect if (event.userId) { this.pushToBuffer(this.userBuffers, event.userId, bufferedEvent); @@ -179,24 +179,9 @@ export class EventBus { this.pushToBuffer(this.orgBuffers, event.organizationId, bufferedEvent); } - let delivered = 0; - this.connections.forEach((connection, id) => { - // Проверяем tenant isolation - if (event.organizationId && connection.organizationId !== event.organizationId) { - return; // Пропускаем соединения из других организаций - } - - // Если событие для конкретного пользователя - if (event.userId && connection.userId !== event.userId) { - return; - } - - const ok = this.writeEventToConnection(connection, bufferedEvent); - if (ok) delivered++; - }); - - if (isChatEvent) { - console.log(`[SSE] delivered ${event.type} to ${delivered}/${this.connections.size} connection(s)`); + // Рассылка только соединениям нужной организации/пользователя (tenant isolation в индексе) + for (const connection of this.index.getForEventTarget(event.organizationId, event.userId)) { + this.writeEventToConnection(connection, bufferedEvent); } } @@ -237,14 +222,11 @@ export class EventBus { } getActiveConnections() { - return this.connections.size; + return this.index.size; } isUserConnected(userId: number, organizationId: number): boolean { - for (const conn of this.connections.values()) { - if (conn.userId === userId && conn.organizationId === organizationId) return true; - } - return false; + return this.index.isUserConnected(userId, organizationId); } } diff --git a/server/routes/sync.routes.ts b/server/routes/sync.routes.ts index 0eec872..a8c4e07 100644 --- a/server/routes/sync.routes.ts +++ b/server/routes/sync.routes.ts @@ -115,19 +115,19 @@ export function registerSyncRoutes(app: import("express").Express): void { formIds = Array.from(ids); } - // Проверяем доступ пользователя к каждой форме - const accessibleFormIds: number[] = []; - for (const formId of formIds) { - const canAccess = await storage.canUserAccessForm(userId, formId, organizationId, 'participate'); - if (canAccess) accessibleFormIds.push(formId); - } + // Проверяем доступ пользователя к каждой форме (параллельно, порядок сохраняется) + const accessFlags = await Promise.all( + formIds.map(formId => storage.canUserAccessForm(userId, formId, organizationId, 'participate')) + ); + const accessibleFormIds = formIds.filter((_, i) => accessFlags[i]); const forms: any[] = []; const tasks: any[] = []; const taskFieldValues: any[] = []; const tableIds = new Set(); - for (const formId of accessibleFormIds) { + // Загружаем формы и их задачи параллельно; слияние результатов — в исходном порядке formIds + const perFormResults = await Promise.all(accessibleFormIds.map(async (formId) => { const [form, formFields, formStatuses, formTabs, transitions] = await Promise.all([ storage.getForm(formId, organizationId), storage.getFormFields(formId, organizationId), @@ -136,7 +136,24 @@ export function registerSyncRoutes(app: import("express").Express): void { storage.getStatusTransitions(formId, organizationId), ]); - if (!form) continue; + if (!form) return null; + + // Honor the form's offline-cache strategy: 'assigned' fetches only tasks + // assigned to the current user; 'recent' / 'all' fetch the latest/up to limit. + const offlineCache = (form as any)?.offlineCache as { enabled?: boolean; strategy?: string; maxTasks?: number } | null | undefined; + const taskLimit = offlineCache?.maxTasks ?? OFFLINE_TASK_LIMIT; + const taskOptions: { limit: number; assignedTo?: number } = { limit: taskLimit }; + if (offlineCache?.strategy === 'assigned') { + taskOptions.assignedTo = userId; + } + const { tasks: formTasks } = await storage.getTasksWithFieldsByFormOptimized(formId, organizationId, taskOptions); + + return { form, formFields, formStatuses, formTabs, transitions, formTasks }; + })); + + for (const result of perFormResults) { + if (!result) continue; + const { form, formFields, formStatuses, formTabs, transitions, formTasks } = result; // Собираем ID справочников из полей типа table / select for (const field of formFields) { @@ -163,16 +180,6 @@ export function registerSyncRoutes(app: import("express").Express): void { transitions, }); - // Honor the form's offline-cache strategy: 'assigned' fetches only tasks - // assigned to the current user; 'recent' / 'all' fetch the latest/up to limit. - const offlineCache = (form as any)?.offlineCache as { enabled?: boolean; strategy?: string; maxTasks?: number } | null | undefined; - const taskLimit = offlineCache?.maxTasks ?? OFFLINE_TASK_LIMIT; - const taskOptions: { limit: number; assignedTo?: number } = { limit: taskLimit }; - if (offlineCache?.strategy === 'assigned') { - taskOptions.assignedTo = userId; - } - const { tasks: formTasks } = await storage.getTasksWithFieldsByFormOptimized(formId, organizationId, taskOptions); - for (const task of formTasks as any[]) { const fieldValues = task.fieldValues || {}; delete task.fieldValues; @@ -183,13 +190,12 @@ export function registerSyncRoutes(app: import("express").Express): void { } } - // Загружаем все доступные справочники организации + // Загружаем все доступные справочники организации (строки — параллельно) const dataTables = await storage.getDataTablesByOrganization(organizationId); - const dataTableRows: any[] = []; - for (const table of dataTables) { - const rows = await storage.getDataTableRows(table.id, organizationId); - dataTableRows.push(...rows); - } + const rowsPerTable = await Promise.all( + dataTables.map(table => storage.getDataTableRows(table.id, organizationId)) + ); + const dataTableRows: any[] = rowsPerTable.flat(); const automations = await storage.getOfflineAutomations(organizationId); const users = await storage.getUsersByOrganization(organizationId); @@ -248,28 +254,30 @@ export function registerSyncRoutes(app: import("express").Express): void { formIds = Array.from(ids); } - const accessibleFormIds: number[] = []; - for (const formId of formIds) { - const canAccess = await storage.canUserAccessForm(userId, formId, organizationId, 'participate'); - if (canAccess) accessibleFormIds.push(formId); - } + // Проверяем доступ пользователя к каждой форме (параллельно, порядок сохраняется) + const accessFlags = await Promise.all( + formIds.map(formId => storage.canUserAccessForm(userId, formId, organizationId, 'participate')) + ); + const accessibleFormIds = formIds.filter((_, i) => accessFlags[i]); const changedForms: any[] = []; const changedTasks: any[] = []; const changedFieldValues: any[] = []; const deletedTaskIds: number[] = []; - for (const formId of accessibleFormIds) { + // Дельта по формам — параллельно; слияние результатов — в исходном порядке formIds + const perFormResults = await Promise.all(accessibleFormIds.map(async (formId) => { const form = await storage.getForm(formId, organizationId); - if (form && (!form.updatedAt || new Date(form.updatedAt) >= since)) { - const [fields, statuses, tabs, transitions] = await Promise.all([ - storage.getFormFields(formId, organizationId), - storage.getFormStatuses(formId, organizationId), - storage.getFormTabs(formId, organizationId), - storage.getStatusTransitions(formId, organizationId), - ]); - changedForms.push({ ...form, fields, statuses, tabs, transitions }); - } + const formChanged = form && (!form.updatedAt || new Date(form.updatedAt) >= since); + + const meta = formChanged + ? await Promise.all([ + storage.getFormFields(formId, organizationId), + storage.getFormStatuses(formId, organizationId), + storage.getFormTabs(formId, organizationId), + storage.getStatusTransitions(formId, organizationId), + ]) + : null; // Honor the form's offline-cache strategy for delta sync as well. const offlineCache = (form as any)?.offlineCache as { enabled?: boolean; strategy?: string; maxTasks?: number } | null | undefined; @@ -279,6 +287,17 @@ export function registerSyncRoutes(app: import("express").Express): void { } // Фильтр по updated_at — на уровне SQL (раньше грузились ВСЕ задачи формы + JS-фильтр) const { tasks: formTasks } = await storage.getTasksWithFieldsByFormOptimized(formId, organizationId, taskOptions); + + return { form, formChanged: !!formChanged, meta, formTasks }; + })); + + for (const result of perFormResults) { + const { form, formChanged, meta, formTasks } = result; + if (formChanged && meta) { + const [fields, statuses, tabs, transitions] = meta; + changedForms.push({ ...form, fields, statuses, tabs, transitions }); + } + for (const task of formTasks as any[]) { const fieldValues = task.fieldValues || {}; delete task.fieldValues; diff --git a/server/services/embedding.service.ts b/server/services/embedding.service.ts index a3702eb..6ae38a6 100644 --- a/server/services/embedding.service.ts +++ b/server/services/embedding.service.ts @@ -2,7 +2,7 @@ import { db } from "../db"; import { sql } from "drizzle-orm"; import { storage } from "../storage"; import { decrypt } from "../crypto"; -import type { Form, FormField, FormStatus, Task, TaskMessage } from "@shared/schema"; +import type { Form, FormField, FormStatus, Task, TaskFieldValue, TaskMessage } from "@shared/schema"; import type { IStorage } from "../storage"; import { generateEmbedding, @@ -523,6 +523,43 @@ export async function processEmbeddingQueue( return meta; } + // Батч-предзагрузка сущностей очереди IN-запросами, чтобы не дёргать БД на каждый элемент + const upsertItems = orgItems.filter(i => i.operation !== 'delete'); + const queueTaskIds = upsertItems.filter(i => i.entityType === 'task').map(i => i.entityId); + const queueMessageIds = upsertItems.filter(i => i.entityType === 'task_message').map(i => i.entityId); + + let taskById = new Map(); + let messageById = new Map(); + let fieldValuesByTaskId = new Map(); + try { + const [tasksBatch, messagesBatch] = await Promise.all([ + stor.getTasksByIds(queueTaskIds, organizationId), + stor.getTaskMessagesByIds(queueMessageIds, organizationId), + ]); + taskById = new Map(tasksBatch.map(t => [t.id, t])); + messageById = new Map(messagesBatch.map(m => [m.id, m])); + + // Родительские задачи сообщений могут не входить в очередь — догружаем одним запросом + const missingTaskIds = [...new Set(messagesBatch.map(m => m.taskId))].filter(id => !taskById.has(id)); + if (missingTaskIds.length > 0) { + const extraTasks = await stor.getTasksByIds(missingTaskIds, organizationId); + for (const t of extraTasks) taskById.set(t.id, t); + } + + const fieldValuesBatch = await stor.getTaskFieldValuesByTaskIds(queueTaskIds, organizationId); + for (const fv of fieldValuesBatch) { + const list = fieldValuesByTaskId.get(fv.taskId); + if (list) list.push(fv); + else fieldValuesByTaskId.set(fv.taskId, [fv]); + } + } catch (err) { + console.error(`[RAG] Failed to preload queue entities for org ${organizationId}:`, err); + attempted += orgItems.length; + failed += orgItems.length; + if (onProgress) onProgress(succeeded + failed, attempted); + continue; + } + const embeddingItems: EmbeddingItem[] = []; for (const item of orgItems) { @@ -554,7 +591,7 @@ export async function processEmbeddingQueue( metadata: { name: form.name }, }); } else if (item.entityType === 'task') { - const task = await stor.getTask(item.entityId, organizationId); + const task = taskById.get(item.entityId); if (!task) { failed++; continue; @@ -570,7 +607,7 @@ export async function processEmbeddingQueue( ? `${assigneeUser.firstName ?? ''} ${assigneeUser.middleName ?? ''} ${assigneeUser.lastName ?? ''}`.trim() || assigneeUser.email : null; const fieldMap = new Map(formFields.map(f => [f.id, f])); - const fieldValues = await stor.getTaskFieldValues(task.id, organizationId); + const fieldValues = fieldValuesByTaskId.get(task.id) ?? []; const enrichedFieldValues = fieldValues.map(fv => ({ name: fieldMap.get(fv.fieldId)?.name ?? '', value: fv.value, @@ -591,12 +628,12 @@ export async function processEmbeddingQueue( metadata: { formId: task.formId, formName: form.name, title: task.title }, }); } else if (item.entityType === 'task_message') { - const message = await stor.getTaskMessage(item.entityId, organizationId); + const message = messageById.get(item.entityId); if (!message) { failed++; continue; } - const task = await stor.getTask(message.taskId, organizationId); + const task = taskById.get(message.taskId); if (!task) { failed++; continue; @@ -734,6 +771,30 @@ export async function reindexOrganization( if (chunk.length === 0) break; + // Батч-загрузка значений полей и сообщений для всего чанка IN-запросами + // вместо поштучных getTaskFieldValues / getTaskMessages на каждую задачу + const chunkTaskIds = chunk.map(t => t.id); + const [chunkFieldValues, chunkMessages] = await Promise.all([ + entityTypes.includes("task") + ? stor.getTaskFieldValuesByTaskIds(chunkTaskIds, organizationId) + : Promise.resolve([] as TaskFieldValue[]), + entityTypes.includes("task_message") + ? stor.getTaskMessagesByTaskIds(chunkTaskIds, organizationId) + : Promise.resolve([] as TaskMessage[]), + ]); + const fieldValuesByTaskId = new Map(); + for (const fv of chunkFieldValues) { + const list = fieldValuesByTaskId.get(fv.taskId); + if (list) list.push(fv); + else fieldValuesByTaskId.set(fv.taskId, [fv]); + } + const messagesByTaskId = new Map(); + for (const msg of chunkMessages) { + const list = messagesByTaskId.get(msg.taskId); + if (list) list.push(msg); + else messagesByTaskId.set(msg.taskId, [msg]); + } + const taskItems: EmbeddingItem[] = []; const messageItems: EmbeddingItem[] = []; @@ -747,7 +808,7 @@ export async function reindexOrganization( const assigneeName = assigneeUser ? `${assigneeUser.firstName ?? ""} ${assigneeUser.middleName ?? ""} ${assigneeUser.lastName ?? ""}`.trim() || assigneeUser.email : null; - const fieldValues = await stor.getTaskFieldValues(task.id, organizationId); + const fieldValues = fieldValuesByTaskId.get(task.id) ?? []; const enrichedFieldValues = fieldValues.map(fv => ({ name: fieldMap.get(fv.fieldId)?.name ?? "", value: fv.value, @@ -765,7 +826,7 @@ export async function reindexOrganization( } if (entityTypes.includes("task_message")) { - const messages = await stor.getTaskMessages(task.id, organizationId); + const messages = messagesByTaskId.get(task.id) ?? []; for (const msg of messages) { if (msg.messageType !== "comment") continue; const authorUser = msg.authorId ? usersMap.get(msg.authorId) : null; diff --git a/server/storage.ts b/server/storage.ts index b5e3801..ac370a1 100644 --- a/server/storage.ts +++ b/server/storage.ts @@ -204,6 +204,8 @@ export interface IStorage { // Task Field Values getTaskFieldValues(taskId: number, organizationId: number): Promise; + getTasksByIds(taskIds: number[], organizationId: number): Promise; + getTaskFieldValuesByTaskIds(taskIds: number[], organizationId: number): Promise; createTaskFieldValue(insertValue: InsertTaskFieldValue): Promise; updateTaskFieldValue(taskId: number, fieldId: number, organizationId: number, updates: Partial): Promise; bulkUpsertTaskFieldValues(taskId: number, formId: number, updates: Array<{ fieldId: number; value: unknown }>): Promise; @@ -216,6 +218,8 @@ export interface IStorage { // Task Messages (Chat) getTaskMessages(taskId: number, organizationId: number, afterId?: number): Promise; getTaskMessage(messageId: number, organizationId: number): Promise; + getTaskMessagesByIds(messageIds: number[], organizationId: number): Promise; + getTaskMessagesByTaskIds(taskIds: number[], organizationId: number): Promise; createTaskMessage(insertMessage: any, organizationId: number): Promise; updateTaskMessage(id: number, organizationId: number, updates: Partial): Promise; deleteTaskMessage(id: number, organizationId: number): Promise; diff --git a/server/storage/tasks-core.storage.ts b/server/storage/tasks-core.storage.ts index 54642e9..2efb3d4 100644 --- a/server/storage/tasks-core.storage.ts +++ b/server/storage/tasks-core.storage.ts @@ -350,20 +350,92 @@ export class TasksCoreStorage extends FormsStorage { } async getTaskTree(taskId: number, organizationId: number): Promise { - const task = await this.getTask(taskId, organizationId); - if (!task) return null; + // Один recursive CTE собирает id всего поддерева (включая корень) с tenant-фильтром, + // затем полные строки задач загружаются одним IN-запросом — вместо 1+N рекурсивных запросов. + // Лимит lvl < 100 — страховка от бесконечной рекурсии при циклическом parentTaskId в данных. + const idResult = await db.execute(sql` + WITH RECURSIVE subtree AS ( + SELECT t.id, t.parent_task_id, t.position, 0 AS lvl + FROM tasks t + INNER JOIN forms f ON f.id = t.form_id + WHERE t.id = ${taskId} AND f.organization_id = ${organizationId} + UNION ALL + SELECT t.id, t.parent_task_id, t.position, s.lvl + 1 + FROM tasks t + INNER JOIN forms f ON f.id = t.form_id + INNER JOIN subtree s ON t.parent_task_id = s.id + WHERE f.organization_id = ${organizationId} AND s.lvl < 100 + ) + SELECT id, parent_task_id AS "parentTaskId" + FROM subtree + ORDER BY lvl, position + `); - const subtasks = await this.getSubtasks(taskId, organizationId); - const subtasksWithChildren = await Promise.all( - subtasks.map(async (subtask) => { - return await this.getTaskTree(subtask.id, organizationId); - }) - ); + const idRows = idResult.rows as Array<{ id: number; parentTaskId: number | null }>; + if (idRows.length === 0) return null; - return { - ...task, - subtasks: subtasksWithChildren.filter(Boolean) - }; + // id'ы уже отфильтрованы по организации в CTE — дополнительный tenant-фильтр не нужен + const taskRows = await db + .select() + .from(tasks) + .where(inArray(tasks.id, idRows.map(r => r.id))); + + const nodeById = new Map(); + for (const task of taskRows) { + nodeById.set(task.id, { ...task, subtasks: [] }); + } + + // Строки идут в порядке ORDER BY lvl, position — при добавлении в списки детей + // сохраняется сортировка по position внутри каждого родителя (как у getSubtasks). + // assembled защищает от повторного добавления узла при циклическом parentTaskId в данных. + const root = nodeById.get(taskId); + if (!root) return null; + const assembled = new Set([taskId]); + for (const row of idRows) { + if (assembled.has(row.id)) continue; + const node = nodeById.get(row.id); + const parent = row.parentTaskId != null ? nodeById.get(row.parentTaskId) : undefined; + if (node && parent) { + parent.subtasks.push(node); + assembled.add(row.id); + } + } + return root; + } + + /** + * Цепочка родителей вверх по parentTaskId (ближайший первым), максимум maxDepth уровней. + * Один recursive CTE + один IN-запрос вместо последовательных getTask на каждый уровень. + */ + async getTaskParentChain(taskId: number, organizationId: number, maxDepth = 10): Promise { + const idResult = await db.execute(sql` + WITH RECURSIVE ancestors AS ( + SELECT t.id, t.parent_task_id, 1 AS lvl + FROM tasks t + INNER JOIN forms f ON f.id = t.form_id + WHERE t.id = (SELECT parent_task_id FROM tasks WHERE id = ${taskId}) + AND f.organization_id = ${organizationId} + UNION ALL + SELECT t.id, t.parent_task_id, a.lvl + 1 + FROM tasks t + INNER JOIN forms f ON f.id = t.form_id + INNER JOIN ancestors a ON t.id = a.parent_task_id + WHERE f.organization_id = ${organizationId} AND a.lvl < ${maxDepth} + ) + SELECT id FROM ancestors ORDER BY lvl + `); + + const ids = (idResult.rows as Array<{ id: number }>).map(r => r.id); + if (ids.length === 0) return []; + + const taskRows = await db + .select() + .from(tasks) + .where(inArray(tasks.id, ids)); + + const byId = new Map(taskRows.map(t => [t.id, t])); + // Порядок — от ближайшего родителя вверх (по lvl из CTE) + return ids.map(id => byId.get(id)).filter((t): t is Task => !!t); } async updateTask(id: number, organizationId: number, updates: Partial): Promise { @@ -464,6 +536,30 @@ export class TasksCoreStorage extends FormsStorage { .then(results => results.map(r => r.task_field_values)); } + /** Батч-загрузка задач по списку id одним IN-запросом (с tenant-фильтром через форму) */ + async getTasksByIds(taskIds: number[], organizationId: number): Promise { + if (taskIds.length === 0) return []; + return await db + .select() + .from(tasks) + .innerJoin(forms, eq(tasks.formId, forms.id)) + .where(and(inArray(tasks.id, taskIds), eq(forms.organizationId, organizationId))) + .then(results => results.map(r => r.tasks)); + } + + /** Батч-загрузка значений полей для списка задач одним IN-запросом (с tenant-фильтром) */ + async getTaskFieldValuesByTaskIds(taskIds: number[], organizationId: number): Promise { + if (taskIds.length === 0) return []; + return await db + .select() + .from(taskFieldValues) + .innerJoin(tasks, eq(taskFieldValues.taskId, tasks.id)) + .innerJoin(forms, eq(tasks.formId, forms.id)) + .where(and(inArray(taskFieldValues.taskId, taskIds), eq(forms.organizationId, organizationId))) + .orderBy(taskFieldValues.taskId, taskFieldValues.fieldId) + .then(results => results.map(r => r.task_field_values)); + } + async createTaskFieldValue(insertValue: InsertTaskFieldValue): Promise { const [value] = await db .insert(taskFieldValues) diff --git a/server/storage/tasks.storage.ts b/server/storage/tasks.storage.ts index 5b5c774..c1b765d 100644 --- a/server/storage/tasks.storage.ts +++ b/server/storage/tasks.storage.ts @@ -2,7 +2,7 @@ import { users, forms, tasks, taskMessages, messageReads, bots, roles, taskViews import { taskRelations } from "@shared/schema"; import { taskAssignees, type TaskAssignee } from "@shared/schema"; import { db } from "../db"; -import { eq, and, desc, exists, sql, gt } from "drizzle-orm"; +import { eq, and, desc, exists, sql, gt, inArray } from "drizzle-orm"; import { formatUserName } from "../utils/formatUserName"; import { TasksCoreStorage } from "./tasks-core.storage"; import { fieldConditionsMatch, resolveFieldDisplayText } from "../utils/field-conditions"; @@ -398,6 +398,32 @@ export class TasksStorage extends TasksCoreStorage { }; } + /** Батч-загрузка сообщений по списку id одним IN-запросом (сырые строки, tenant-фильтр) */ + async getTaskMessagesByIds(messageIds: number[], organizationId: number): Promise { + if (messageIds.length === 0) return []; + return await db + .select() + .from(taskMessages) + .innerJoin(tasks, eq(taskMessages.taskId, tasks.id)) + .innerJoin(forms, eq(tasks.formId, forms.id)) + .where(and(inArray(taskMessages.id, messageIds), eq(forms.organizationId, organizationId))) + .orderBy(taskMessages.createdAt) + .then(results => results.map(r => r.task_messages)); + } + + /** Батч-загрузка сообщений для списка задач одним IN-запросом (сырые строки, tenant-фильтр) */ + async getTaskMessagesByTaskIds(taskIds: number[], organizationId: number): Promise { + if (taskIds.length === 0) return []; + return await db + .select() + .from(taskMessages) + .innerJoin(tasks, eq(taskMessages.taskId, tasks.id)) + .innerJoin(forms, eq(tasks.formId, forms.id)) + .where(and(inArray(taskMessages.taskId, taskIds), eq(forms.organizationId, organizationId))) + .orderBy(taskMessages.createdAt) + .then(results => results.map(r => r.task_messages)); + } + async createTaskMessage(insertMessage: any, organizationId: number): Promise { const taskExists = await db .select({ id: tasks.id }) diff --git a/server/utils/sse-connection-index.ts b/server/utils/sse-connection-index.ts new file mode 100644 index 0000000..f0a4255 --- /dev/null +++ b/server/utils/sse-connection-index.ts @@ -0,0 +1,97 @@ +// Индекс SSE-соединений по организации и пользователю. +// Позволяет publishEvent доставлять событие только соединениям нужной +// организации/пользователя, не перебирая все открытые соединения. + +export interface IndexedSSEConnection { + id: string; + userId: number; + organizationId: number; +} + +export class SseConnectionIndex { + private byId = new Map(); + private byOrg = new Map>(); + private byUser = new Map>(); + + add(connection: T): void { + this.byId.set(connection.id, connection); + + let orgSet = this.byOrg.get(connection.organizationId); + if (!orgSet) { + orgSet = new Set(); + this.byOrg.set(connection.organizationId, orgSet); + } + orgSet.add(connection); + + let userSet = this.byUser.get(connection.userId); + if (!userSet) { + userSet = new Set(); + this.byUser.set(connection.userId, userSet); + } + userSet.add(connection); + } + + remove(connectionId: string): T | undefined { + const connection = this.byId.get(connectionId); + if (!connection) return undefined; + + this.byId.delete(connectionId); + + const orgSet = this.byOrg.get(connection.organizationId); + if (orgSet) { + orgSet.delete(connection); + if (orgSet.size === 0) this.byOrg.delete(connection.organizationId); + } + + const userSet = this.byUser.get(connection.userId); + if (userSet) { + userSet.delete(connection); + if (userSet.size === 0) this.byUser.delete(connection.userId); + } + + return connection; + } + + get size(): number { + return this.byId.size; + } + + /** Все соединения (для глобального heartbeat) */ + forEach(callback: (connection: T) => void): void { + this.byId.forEach(callback); + } + + /** + * Соединения-получатели события. Семантика соответствует прежнему + * полному перебору с фильтрами tenant isolation: + * - userId + organizationId — соединения этого пользователя в этой организации; + * - только userId — все соединения пользователя (в любой организации); + * - только organizationId — все соединения организации; + * - ни того ни другого — broadcast на все соединения. + */ + getForEventTarget(organizationId?: number, userId?: number): Iterable { + if (userId !== undefined) { + const userSet = this.byUser.get(userId); + if (!userSet) return []; + if (organizationId === undefined) return userSet; + const filtered: T[] = []; + for (const connection of userSet) { + if (connection.organizationId === organizationId) filtered.push(connection); + } + return filtered; + } + if (organizationId !== undefined) { + return this.byOrg.get(organizationId) ?? []; + } + return this.byId.values(); + } + + isUserConnected(userId: number, organizationId: number): boolean { + const userSet = this.byUser.get(userId); + if (!userSet) return false; + for (const connection of userSet) { + if (connection.organizationId === organizationId) return true; + } + return false; + } +} diff --git a/tests/sse-connection-index.test.ts b/tests/sse-connection-index.test.ts new file mode 100644 index 0000000..e6a61ff --- /dev/null +++ b/tests/sse-connection-index.test.ts @@ -0,0 +1,90 @@ +import { describe, it, expect } from 'vitest'; +import { SseConnectionIndex, type IndexedSSEConnection } from '../server/utils/sse-connection-index'; + +function conn(id: string, userId: number, organizationId: number): IndexedSSEConnection { + return { id, userId, organizationId }; +} + +describe('SseConnectionIndex', () => { + it('add/remove ведут учёт размера и очищают оба индекса', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + index.add(conn('b', 1, 10)); + index.add(conn('c', 2, 20)); + expect(index.size).toBe(3); + + index.remove('a'); + expect(index.size).toBe(2); + // Соединение удалено из всех индексов — не придёт ни по org, ни по user + expect(Array.from(index.getForEventTarget(10)).map(c => c.id)).toEqual(['b']); + expect(Array.from(index.getForEventTarget(10, 1)).map(c => c.id)).toEqual(['b']); + + index.remove('b'); + index.remove('c'); + expect(index.size).toBe(0); + // Пустые Set'ы убраны из карт + expect(Array.from(index.getForEventTarget(10))).toEqual([]); + expect(Array.from(index.getForEventTarget(20, 2))).toEqual([]); + }); + + it('remove несуществующего id — no-op', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + expect(index.remove('missing')).toBeUndefined(); + expect(index.size).toBe(1); + }); + + it('событие организации уходит только соединениям этой организации', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + index.add(conn('b', 2, 10)); + index.add(conn('c', 3, 20)); + const targets = Array.from(index.getForEventTarget(10)).map(c => c.id); + expect(targets).toEqual(['a', 'b']); + }); + + it('событие пользователя без organizationId уходит всем его соединениям (любая org)', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + index.add(conn('b', 1, 20)); + index.add(conn('c', 2, 10)); + const targets = Array.from(index.getForEventTarget(undefined, 1)).map(c => c.id); + expect(targets).toEqual(['a', 'b']); + }); + + it('событие пользователя + organizationId уходит только его соединениям в этой org', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + index.add(conn('b', 1, 20)); + index.add(conn('c', 2, 10)); + const targets = Array.from(index.getForEventTarget(10, 1)).map(c => c.id); + expect(targets).toEqual(['a']); + }); + + it('событие без userId и organizationId — broadcast на все соединения', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + index.add(conn('b', 2, 20)); + const targets = Array.from(index.getForEventTarget(undefined, undefined)).map(c => c.id); + expect(targets).toEqual(['a', 'b']); + }); + + it('isUserConnected учитывает и пользователя, и организацию', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + expect(index.isUserConnected(1, 10)).toBe(true); + expect(index.isUserConnected(1, 20)).toBe(false); + expect(index.isUserConnected(2, 10)).toBe(false); + index.remove('a'); + expect(index.isUserConnected(1, 10)).toBe(false); + }); + + it('forEach обходит все соединения (глобальный heartbeat)', () => { + const index = new SseConnectionIndex(); + index.add(conn('a', 1, 10)); + index.add(conn('b', 2, 20)); + const seen: string[] = []; + index.forEach(c => seen.push(c.id)); + expect(seen).toEqual(['a', 'b']); + }); +});