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 { 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 { 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 { 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(); 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`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(); 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>(); 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> { 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`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[]> { 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 { const result = await db .select({ count: sql`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 { 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`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 { 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 { // Один 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(); 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 { 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, requiredUpdatedAt: Date, ): Promise { 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 { // 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 & { 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 { 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 { 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) .values(insertValue) .returning(); return value; } async updateTaskFieldValue(taskId: number, fieldId: number, organizationId: number, updates: Partial): Promise { 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 { 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 { 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 { 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 { const [history] = await db .insert(fieldHistory) .values(insertHistory) .returning(); return history; } }