Files
iistwin/server/storage/tasks-core.storage.ts
Ильяс Султанов 667a321341 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)
2026-09-07 23:29:21 +03:00

657 lines
24 KiB
TypeScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import { users, forms, formStatuses, tasks, taskFieldValues, fieldHistory, type User, type Form, type FormStatus, type Task, type TaskFieldValue, type FieldHistory, type InsertTask, type InsertTaskFieldValue, type InsertFieldHistory } from "@shared/schema";
import { taskRelations } from "@shared/schema";
import { taskAssignees } from "@shared/schema";
import { db } from "../db";
import { eq, and, or, desc, asc, exists, sql, ilike, inArray, lt, gte, isNull } from "drizzle-orm";
import { formatUserName } from "../utils/formatUserName";
import { FormsStorage } from "./forms.storage";
import { clampTasksLimit, type TasksCursor } from "../utils/task-pagination";
import { invalidateAccessibleTasksForUser } from "../utils/cache";
export class TasksCoreStorage extends FormsStorage {
// Tasks
async getTask(id: number, organizationId: number): Promise<Task | undefined> {
const results = await db
.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(tasks.id, id), eq(forms.organizationId, organizationId)));
const [result] = results;
return result?.tasks || undefined;
}
async getTasksByForm(formId: number, organizationId: number): Promise<Task[]> {
return db
.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(tasks.formId, formId), eq(forms.organizationId, organizationId)))
.orderBy(desc(tasks.createdAt))
.then(results => results.map(r => r.tasks));
}
async getTasksWithFieldsByForm(formId: number, organizationId: number): Promise<Task[]> {
const allTasks = await db
.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(tasks.formId, formId), eq(forms.organizationId, organizationId)))
.orderBy(desc(tasks.createdAt))
.then(results => results.map(r => r.tasks));
if (allTasks.length === 0) return [];
const taskIds = allTasks.map(t => t.id);
const allFieldValues = await db
.select()
.from(taskFieldValues)
.where(inArray(taskFieldValues.taskId, taskIds));
const fieldValuesByTask = new Map<number, { [fieldId: number]: unknown }>();
for (const fv of allFieldValues) {
if (!fieldValuesByTask.has(fv.taskId)) {
fieldValuesByTask.set(fv.taskId, {});
}
fieldValuesByTask.get(fv.taskId)![fv.fieldId] = fv.value;
}
return allTasks.map(task => ({
...task,
fieldValues: fieldValuesByTask.get(task.id) || {}
}));
}
async getTasksWithFieldsByFormOptimized(
formId: number,
organizationId: number,
options?: { limit?: number; offset?: number; assignedTo?: number; since?: Date }
): Promise<{ tasks: Task[]; total: number }> {
const limit = options?.limit;
const offset = options?.offset || 0;
const whereConditions = [
eq(tasks.formId, formId),
eq(forms.organizationId, organizationId),
];
if (options?.assignedTo != null) {
whereConditions.push(eq(tasks.assignedTo, options.assignedTo));
}
// Дельта-синк: фильтр по updated_at на уровне SQL (индекс tasks_form_updated_idx).
// NULL updated_at включаем для паритета со старым JS-фильтром (!t.updatedAt || >= since).
if (options?.since) {
whereConditions.push(or(
gte(tasks.updatedAt, options.since),
isNull(tasks.updatedAt)
)!);
}
let query = db
.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(...whereConditions))
.orderBy(desc(tasks.updatedAt))
.offset(offset);
const allTasks = limit
? await query.limit(limit).then(results => results.map(r => r.tasks))
: await query.then(results => results.map(r => r.tasks));
if (allTasks.length === 0) {
return { tasks: [], total: 0 };
}
const countResult = await db
.select({ count: sql<number>`count(*)` })
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(...whereConditions));
const total = Number(countResult[0]?.count || 0);
const taskIds = allTasks.map(t => t.id);
const [allFieldValues, allTaskAssignees] = await Promise.all([
db.select().from(taskFieldValues).where(inArray(taskFieldValues.taskId, taskIds)),
db
.select({
taskId: taskAssignees.taskId,
userId: taskAssignees.userId,
firstName: users.firstName,
middleName: users.middleName,
lastName: users.lastName,
email: users.email,
})
.from(taskAssignees)
.innerJoin(users, eq(taskAssignees.userId, users.id))
.where(inArray(taskAssignees.taskId, taskIds)),
]);
const fieldValuesByTask = new Map<number, { [fieldId: number]: unknown }>();
for (const fv of allFieldValues) {
if (!fieldValuesByTask.has(fv.taskId)) {
fieldValuesByTask.set(fv.taskId, {});
}
fieldValuesByTask.get(fv.taskId)![fv.fieldId] = fv.value;
}
const assigneesByTask = new Map<number, Array<{ userId: number; user: { firstName: string | null; middleName: string | null; lastName: string | null; email: string } }>>();
for (const a of allTaskAssignees) {
if (!assigneesByTask.has(a.taskId)) {
assigneesByTask.set(a.taskId, []);
}
assigneesByTask.get(a.taskId)!.push({ userId: a.userId, user: { firstName: a.firstName, middleName: a.middleName, lastName: a.lastName, email: a.email } });
}
const tasksWithFields = allTasks.map(task => ({
...task,
fieldValues: fieldValuesByTask.get(task.id) || {},
assignees: assigneesByTask.get(task.id) || [],
}));
return { tasks: tasksWithFields, total };
}
async getTaskCountsByStatus(formId: number, organizationId: number): Promise<Array<{ statusId: number; name: string; color: string; count: number }>> {
const statuses = await db
.select()
.from(formStatuses)
.where(and(eq(formStatuses.formId, formId), eq(formStatuses.isFinal, false)))
.orderBy(asc(formStatuses.position));
if (statuses.length === 0) return [];
const countsResult = await db
.select({
currentStatusId: tasks.currentStatusId,
count: sql<number>`count(*)::int`,
})
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(tasks.formId, formId), eq(forms.organizationId, organizationId)))
.groupBy(tasks.currentStatusId);
const countMap = new Map(countsResult.map(c => [c.currentStatusId, c.count]));
return statuses.map(s => ({
statusId: s.id,
name: s.name,
color: s.color || '#94a3b8',
count: countMap.get(s.id) || 0,
}));
}
async getTasksByOrganization(
organizationId: number,
options?: {
limit?: number;
offset?: number;
sortBy?: 'updatedAt' | 'createdAt';
order?: 'asc' | 'desc';
statusId?: number;
assignedTo?: number;
minimal?: boolean;
excludeDescription?: boolean;
cursor?: TasksCursor | null;
}
): Promise<Task[] | Partial<Task>[]> {
const sortBy = options?.sortBy || 'updatedAt';
const order = options?.order || 'desc';
// Дефолт 50, жёсткий кап 200 — список никогда не выгружает таблицу целиком
const limit = clampTasksLimit(options?.limit);
const offset = options?.offset;
const statusId = options?.statusId;
const assignedTo = options?.assignedTo;
const minimal = options?.minimal || false;
const sortColumn = sortBy === 'updatedAt' ? tasks.updatedAt : tasks.createdAt;
const orderFn = order === 'desc' ? desc : asc;
const conditions = [eq(forms.organizationId, organizationId)];
if (statusId) {
conditions.push(eq(tasks.currentStatusId, statusId));
}
if (assignedTo) {
conditions.push(eq(tasks.assignedTo, assignedTo));
}
// Курсорная пагинация работает только для дефолтной сортировки (updatedAt DESC, id DESC)
const cursor = options?.cursor ?? null;
const useCursor = cursor !== null && sortBy === 'updatedAt' && order === 'desc';
if (useCursor) {
conditions.push(or(
lt(tasks.updatedAt, cursor.updatedAt),
and(eq(tasks.updatedAt, cursor.updatedAt), lt(tasks.id, cursor.id))
)!);
}
// В курсорном режиме добавляем id как детерминированный тай-брейкер
const orderByClause = useCursor ? [desc(tasks.updatedAt), desc(tasks.id)] : [orderFn(sortColumn)];
if (minimal) {
return await db
.select({
id: tasks.id,
formId: tasks.formId,
title: tasks.title,
currentStatusId: tasks.currentStatusId,
assignedTo: tasks.assignedTo,
updatedAt: tasks.updatedAt,
createdAt: tasks.createdAt,
status: formStatuses.name
})
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.leftJoin(formStatuses, eq(tasks.currentStatusId, formStatuses.id))
.where(and(...conditions))
.orderBy(...orderByClause)
.limit(limit)
.offset(offset || 0);
}
// Списковый ответ без description (тяжёлая text-колонка) — только по явному флагу;
// остальные caller'ы (embedding, MCP) получают полную строку, как раньше
if (options?.excludeDescription) {
const rows = await db
.select({
id: tasks.id,
organizationId: tasks.organizationId,
formId: tasks.formId,
title: tasks.title,
currentStatusId: tasks.currentStatusId,
assignedTo: tasks.assignedTo,
createdBy: tasks.createdBy,
createdAt: tasks.createdAt,
updatedAt: tasks.updatedAt,
completedAt: tasks.completedAt,
isCompleted: tasks.isCompleted,
dueDate: tasks.dueDate,
parentTaskId: tasks.parentTaskId,
depth: tasks.depth,
position: tasks.position,
})
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(...conditions))
.orderBy(...orderByClause)
.limit(limit)
.offset(offset || 0);
return rows as unknown as Task[];
}
const results = await db
.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(...conditions))
.orderBy(...orderByClause)
.limit(limit)
.offset(offset || 0);
return results.map((r) => r.tasks);
}
/** Точный COUNT(*) задач организации — для оценки прогресса без выгрузки строк */
async getTasksCountByOrganization(organizationId: number): Promise<number> {
const result = await db
.select({ count: sql<number>`count(*)` })
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(eq(forms.organizationId, organizationId));
return Number(result[0]?.count || 0);
}
async createTask(insertTask: InsertTask): Promise<Task> {
let depth = 0;
let position = 0;
if (insertTask.parentTaskId) {
const [parentTask] = await db
.select()
.from(tasks)
.where(eq(tasks.id, insertTask.parentTaskId));
if (parentTask) {
depth = (parentTask.depth || 0) + 1;
const siblings = await db
.select({ count: sql<number>`count(*)` })
.from(tasks)
.where(eq(tasks.parentTaskId, insertTask.parentTaskId));
position = Number(siblings[0]?.count || 0);
}
}
const [task] = await db
.insert(tasks)
.values({
...(insertTask as typeof tasks.$inferInsert),
depth,
position,
})
.returning();
// Создатель и ответственный получают доступ к новой задаче — сброс их кэша
if (task.createdBy) invalidateAccessibleTasksForUser(task.createdBy, task.organizationId);
if (task.assignedTo && task.assignedTo !== task.createdBy) {
invalidateAccessibleTasksForUser(task.assignedTo, task.organizationId);
}
return task;
}
async getSubtasks(parentTaskId: number, organizationId: number): Promise<Task[]> {
return await db
.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(
eq(tasks.parentTaskId, parentTaskId),
eq(forms.organizationId, organizationId)
))
.orderBy(tasks.position)
.then(results => results.map(r => r.tasks));
}
async getTaskTree(taskId: number, organizationId: number): Promise<any> {
// Один 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 idRows = idResult.rows as Array<{ id: number; parentTaskId: number | null }>;
if (idRows.length === 0) return null;
// id'ы уже отфильтрованы по организации в CTE — дополнительный tenant-фильтр не нужен
const taskRows = await db
.select()
.from(tasks)
.where(inArray(tasks.id, idRows.map(r => r.id)));
const nodeById = new Map<number, any>();
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<number>([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<Task[]> {
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<Task>): Promise<Task> {
const [task] = await db
.update(tasks)
.set({ ...updates, updatedAt: new Date() })
.where(and(
eq(tasks.id, id),
exists(db.select().from(forms).where(and(eq(forms.id, tasks.formId), eq(forms.organizationId, organizationId))))
))
.returning();
// Переназначение ответственного даёт доступ новому исполнителю
if ('assignedTo' in updates && updates.assignedTo != null) {
invalidateAccessibleTasksForUser(updates.assignedTo, organizationId);
}
return task;
}
async updateTaskIfNotModified(
id: number,
organizationId: number,
updates: Partial<Task>,
requiredUpdatedAt: Date,
): Promise<Task | null> {
const [task] = await db
.update(tasks)
.set({ ...updates, updatedAt: new Date() })
.where(and(
eq(tasks.id, id),
eq(tasks.updatedAt, requiredUpdatedAt),
exists(db.select().from(forms).where(and(eq(forms.id, tasks.formId), eq(forms.organizationId, organizationId))))
))
.returning();
return task ?? null;
}
async deleteTask(id: number, organizationId: number): Promise<void> {
// Cleanup RAG embeddings for this task and its messages
const { deleteRagEmbeddingsByTask } = await import('../services/embedding.service.js');
await deleteRagEmbeddingsByTask(organizationId, id);
await db
.delete(tasks)
.where(and(
eq(tasks.id, id),
exists(
db.select()
.from(forms)
.where(and(eq(forms.id, tasks.formId), eq(forms.organizationId, organizationId)))
)
));
}
async searchTasksByKeyword(
organizationId: number,
query: string,
limit: number = 10
): Promise<(Pick<Task, 'id' | 'title' | 'description' | 'formId' | 'currentStatusId' | 'isCompleted'> & { statusName: string | null; statusColor: string | null; statusIsFinal: boolean | null })[]> {
const pattern = `%${query}%`;
const results = await db
.select({
id: tasks.id,
title: tasks.title,
description: tasks.description,
formId: tasks.formId,
currentStatusId: tasks.currentStatusId,
isCompleted: tasks.isCompleted,
statusName: formStatuses.name,
statusColor: formStatuses.color,
statusIsFinal: formStatuses.isFinal,
})
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.leftJoin(formStatuses, eq(tasks.currentStatusId, formStatuses.id))
.where(
and(
eq(forms.organizationId, organizationId),
or(
ilike(tasks.title, pattern),
ilike(tasks.description, pattern)
)
)
)
.orderBy(desc(tasks.updatedAt))
.limit(limit);
return results;
}
// Task Field Values
async getTaskFieldValues(taskId: number, organizationId: number): Promise<TaskFieldValue[]> {
return await db
.select()
.from(taskFieldValues)
.innerJoin(tasks, eq(taskFieldValues.taskId, tasks.id))
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(taskFieldValues.taskId, taskId), eq(forms.organizationId, organizationId)))
.orderBy(taskFieldValues.fieldId)
.then(results => results.map(r => r.task_field_values));
}
/** Батч-загрузка задач по списку id одним IN-запросом (с tenant-фильтром через форму) */
async getTasksByIds(taskIds: number[], organizationId: number): Promise<Task[]> {
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<TaskFieldValue[]> {
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<TaskFieldValue> {
const [value] = await db
.insert(taskFieldValues)
.values(insertValue)
.returning();
return value;
}
async updateTaskFieldValue(taskId: number, fieldId: number, organizationId: number, updates: Partial<TaskFieldValue>): Promise<TaskFieldValue> {
const [value] = await db
.update(taskFieldValues)
.set({ ...updates, updatedAt: new Date() })
.where(and(
eq(taskFieldValues.taskId, taskId),
eq(taskFieldValues.fieldId, fieldId),
exists(
db.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(tasks.id, taskId), eq(forms.organizationId, organizationId)))
)
))
.returning();
return value;
}
async bulkUpsertTaskFieldValues(taskId: number, formId: number, updates: Array<{ fieldId: number; value: unknown }>): Promise<TaskFieldValue[]> {
const results: TaskFieldValue[] = [];
await db.transaction(async (tx) => {
const existing = await tx
.select()
.from(taskFieldValues)
.where(eq(taskFieldValues.taskId, taskId));
const existingMap = new Map(existing.map(v => [v.fieldId, v]));
for (const { fieldId, value } of updates) {
const existingValue = existingMap.get(fieldId);
if (existingValue) {
const [updated] = await tx
.update(taskFieldValues)
.set({ value, updatedAt: new Date() })
.where(and(eq(taskFieldValues.taskId, taskId), eq(taskFieldValues.fieldId, fieldId)))
.returning();
results.push(updated);
} else {
const [inserted] = await tx
.insert(taskFieldValues)
.values({ taskId, fieldId, formId, value })
.returning();
results.push(inserted);
}
}
});
return results;
}
async deleteTaskFieldValue(taskId: number, fieldId: number, organizationId: number): Promise<void> {
await db
.delete(taskFieldValues)
.where(and(
eq(taskFieldValues.taskId, taskId),
eq(taskFieldValues.fieldId, fieldId),
exists(
db.select()
.from(tasks)
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(eq(tasks.id, taskId), eq(forms.organizationId, organizationId)))
)
));
}
// Field History
async getFieldHistory(taskId: number, fieldId: number, organizationId: number): Promise<FieldHistory[]> {
return await db
.select()
.from(fieldHistory)
.innerJoin(tasks, eq(fieldHistory.taskId, tasks.id))
.innerJoin(forms, eq(tasks.formId, forms.id))
.where(and(
eq(fieldHistory.taskId, taskId),
eq(fieldHistory.fieldId, fieldId),
eq(forms.organizationId, organizationId)
))
.orderBy(desc(fieldHistory.changedAt))
.then(results => results.map(r => r.field_history));
}
async createFieldHistory(insertHistory: InsertFieldHistory): Promise<FieldHistory> {
const [history] = await db
.insert(fieldHistory)
.values(insertHistory)
.returning();
return history;
}
}