perf(sse+tasks): индекс SSE-соединений, getTaskTree на recursive CTE, батчи N+1

Шаги 1.4 и 1.5 плана production-готовности:
- SseConnectionIndex (byOrg/byUser), один heartbeat-таймер, протокол SSE не тронут
- getTaskTree/getTaskParentChain: WITH RECURSIVE CTE (было 1+2N запросов)
- sync initial/delta: Promise.all по формам
- embedding queue/reindex: батч-предзагрузки вместо поштучных запросов
- 8 новых тестов (96/96)
This commit is contained in:
2026-09-07 23:29:21 +03:00
parent 9528a8123a
commit 667a321341
10 changed files with 518 additions and 121 deletions

View File

@@ -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<number, Task>();
let messageById = new Map<number, TaskMessage>();
let fieldValuesByTaskId = new Map<number, TaskFieldValue[]>();
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<number, TaskFieldValue[]>();
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<number, TaskMessage[]>();
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;