Files
iistwin/server/routes/sync.routes.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

381 lines
15 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 { Router } from "express";
import crypto from "crypto";
import { storage } from "../storage";
import { authenticateToken, type AuthenticatedRequest } from "../middleware/auth.middleware";
import { tenantIsolation } from "../middleware/tenant.middleware";
import { logAudit } from "../utils/audit";
import { canManageUser } from "../utils/user-access";
import { z } from "zod";
const OFFLINE_TASK_LIMIT = 500;
const initialSyncSchema = z.object({
formIds: z.array(z.number()).optional(),
});
const deltaSyncSchema = z.object({
formIds: z.array(z.number()).optional(),
since: z.string().datetime(),
});
const pushSyncSchema = z.object({
mutations: z.array(z.record(z.any())).default([]),
});
export function registerSyncRoutes(app: import("express").Express): void {
const router = Router();
// Все endpoint'ы синхронизации требуют аутентификации и tenant-изоляции
router.use(authenticateToken, tenantIsolation);
/**
* GET /api/sync/config
* Конфигурация офлайн-режима для текущего пользователя:
* - формы с включённым глобальным offlineCache
* - персональные подписки пользователя
* - автоматизации, разрешённые для офлайн-выполнения
* ?userId=N — конфигурация другого пользователя (только при canManageUser: админ/руководитель)
*/
router.get('/config', async (req: AuthenticatedRequest, res) => {
try {
let userId = req.user!.id;
const organizationId = req.organizationId!;
if (req.query.userId !== undefined) {
const queryUserId = parseInt(String(req.query.userId));
if (isNaN(queryUserId)) {
return res.status(400).json({ success: false, error: 'Некорректный ID пользователя' });
}
if (queryUserId !== userId) {
if (!(await canManageUser(req.user, queryUserId, organizationId))) {
return res.status(403).json({ success: false, error: 'Нет прав на просмотр конфигурации этого пользователя' });
}
userId = queryUserId;
}
}
const [globalForms, subscriptions, automations] = await Promise.all([
storage.getFormsWithOfflineCache(organizationId),
storage.getUserOfflineSubscriptions(userId, organizationId),
storage.getOfflineAutomations(organizationId),
]);
const subscribedFormIds = new Set(subscriptions.map(s => s.formId));
const forms = [
...globalForms.filter(f => !subscribedFormIds.has(f.id)),
...subscriptions.filter(s => s.isActive).map(s => s.formId),
];
const uniqueFormIds = Array.from(new Set(forms.map(f => typeof f === 'number' ? f : f.id)));
const formList = uniqueFormIds.map(id => {
const global = globalForms.find(f => f.id === id);
if (global) return global;
return { id } as const;
});
res.json({
success: true,
forms: formList,
subscriptions,
automations: automations.map(a => ({
id: a.id,
name: a.name,
trigger: a.trigger,
triggerConfig: a.triggerConfig,
runOffline: a.runOffline,
clientCompatible: a.clientCompatible,
})),
});
} catch (err: any) {
console.error('[SYNC] config error:', err);
res.status(500).json({ success: false, error: 'Ошибка получения конфигурации синхронизации' });
}
});
/**
* POST /api/sync/initial
* Начальная полная синхронизация данных для офлайн-работы.
* Возвращает метаданные форм, задачи (до 500 на форму), значения полей,
* справочники и offline-автоматизации.
*/
router.post('/initial', async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const organizationId = req.organizationId!;
const body = initialSyncSchema.parse(req.body);
let formIds = body.formIds;
if (!formIds || formIds.length === 0) {
const [globalForms, subscriptions] = await Promise.all([
storage.getFormsWithOfflineCache(organizationId),
storage.getUserOfflineSubscriptions(userId, organizationId),
]);
const ids = new Set<number>();
globalForms.forEach(f => ids.add(f.id));
subscriptions.filter(s => s.isActive).forEach(s => ids.add(s.formId));
formIds = Array.from(ids);
}
// Проверяем доступ пользователя к каждой форме (параллельно, порядок сохраняется)
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<number>();
// Загружаем формы и их задачи параллельно; слияние результатов — в исходном порядке 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),
storage.getFormStatuses(formId, organizationId),
storage.getFormTabs(formId, organizationId),
storage.getStatusTransitions(formId, organizationId),
]);
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) {
const f = field as any;
if (f.type === 'table' && f.linkedFormId) {
tableIds.add(f.linkedFormId);
}
if ((f.type === 'select' || f.type === 'radio-group') && f.options) {
// options может быть строкой JSON
const opts = typeof f.options === 'string' ? JSON.parse(f.options) : f.options;
if (Array.isArray(opts)) {
opts.forEach((o: any) => {
if (o.tableId) tableIds.add(o.tableId);
});
}
}
}
forms.push({
...form,
fields: formFields,
statuses: formStatuses,
tabs: formTabs,
transitions,
});
for (const task of formTasks as any[]) {
const fieldValues = task.fieldValues || {};
delete task.fieldValues;
tasks.push(task);
Object.entries(fieldValues).forEach(([fieldId, value]) => {
taskFieldValues.push({ taskId: task.id, fieldId: Number(fieldId), value });
});
}
}
// Загружаем все доступные справочники организации (строки — параллельно)
const dataTables = await storage.getDataTablesByOrganization(organizationId);
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);
logAudit({
action: 'sync.initial',
userId,
organizationId,
details: { formCount: forms.length, taskCount: tasks.length },
});
res.json({
success: true,
forms,
tasks,
taskFieldValues,
dataTables,
dataTableRows,
automations: automations.map(a => ({
id: a.id,
name: a.name,
trigger: a.trigger,
triggerConfig: a.triggerConfig,
code: a.code,
runOffline: a.runOffline,
clientCompatible: a.clientCompatible,
})),
users,
});
} catch (err: any) {
console.error('[SYNC] initial error:', err);
res.status(500).json({ success: false, error: 'Ошибка начальной синхронизации' });
}
});
/**
* POST /api/sync/delta
* Дельта-синхронизация: изменения с момента `since`.
*/
router.post('/delta', async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const organizationId = req.organizationId!;
const body = deltaSyncSchema.parse(req.body);
const since = new Date(body.since);
let formIds = body.formIds;
if (!formIds || formIds.length === 0) {
const [globalForms, subscriptions] = await Promise.all([
storage.getFormsWithOfflineCache(organizationId),
storage.getUserOfflineSubscriptions(userId, organizationId),
]);
const ids = new Set<number>();
globalForms.forEach(f => ids.add(f.id));
subscriptions.filter(s => s.isActive).forEach(s => ids.add(s.formId));
formIds = Array.from(ids);
}
// Проверяем доступ пользователя к каждой форме (параллельно, порядок сохраняется)
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[] = [];
// Дельта по формам — параллельно; слияние результатов — в исходном порядке formIds
const perFormResults = await Promise.all(accessibleFormIds.map(async (formId) => {
const form = await storage.getForm(formId, organizationId);
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;
const taskOptions: { assignedTo?: number; since?: Date } = { since };
if (offlineCache?.strategy === 'assigned') {
taskOptions.assignedTo = userId;
}
// Фильтр по 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;
changedTasks.push(task);
Object.entries(fieldValues).forEach(([fieldId, value]) => {
changedFieldValues.push({ taskId: task.id, fieldId: Number(fieldId), value });
});
}
// TODO: deleted tasks tracking через task_audit_log
}
// Фильтр по updated_at — на уровне SQL (раньше: полная выгрузка + JS-фильтр)
const changedUsers = await storage.getUsersByOrganization(organizationId, { since });
res.json({
success: true,
since: since.toISOString(),
forms: changedForms,
tasks: changedTasks,
taskFieldValues: changedFieldValues,
deletedTaskIds,
users: changedUsers,
});
} catch (err: any) {
console.error('[SYNC] delta error:', err);
res.status(500).json({ success: false, error: 'Ошибка дельта-синхронизации' });
}
});
/**
* POST /api/sync/push
* Приём offline-мутаций, накопленных клиентом без интернета.
* В текущей версии возвращает accepted-список и conflicts.
*/
router.post('/push', async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const organizationId = req.organizationId!;
const body = pushSyncSchema.parse(req.body);
// LEGACY / NO-OP: the offline queue now replays mutations directly via
// authenticated fetch requests when the client comes back online.
// Keeping this endpoint for backward compatibility with older PWA builds.
const accepted: string[] = [];
const conflicts: Array<{ id: string; reason: string }> = [];
for (const mutation of body.mutations) {
const id = (mutation as any)?.id || crypto.randomUUID();
accepted.push(id);
}
logAudit({
action: 'sync.push',
userId,
organizationId,
details: { accepted: accepted.length, conflicts: conflicts.length },
});
res.json({ success: true, accepted, conflicts });
} catch (err: any) {
console.error('[SYNC] push error:', err);
res.status(500).json({ success: false, error: 'Ошибка приёма offline-мутаций' });
}
});
/**
* GET /api/sync/status
* Служебный статус синхронизации (online/offline, время сервера).
*/
router.get('/status', async (req: AuthenticatedRequest, res) => {
res.json({
success: true,
serverTime: new Date().toISOString(),
online: true,
});
});
app.use('/api/sync', router);
}