428 lines
17 KiB
TypeScript
428 lines
17 KiB
TypeScript
import { Router } from "express";
|
||
import { Worker } from 'node:worker_threads';
|
||
import v8 from 'v8';
|
||
import { storage } from "../storage";
|
||
import { authenticateToken, requirePermission, type AuthenticatedRequest } from "../middleware/auth.middleware";
|
||
import { tenantIsolation } from "../middleware/tenant.middleware";
|
||
import { logAudit } from "../utils/audit";
|
||
import { AUTOMATION_WORKER_CODE } from "../workers/automation-worker-code";
|
||
import { sendTaskMessage } from "../services/task-message.service";
|
||
|
||
const router = Router();
|
||
|
||
const WORKER_TIMEOUT_MS = 10_000;
|
||
|
||
// Storage methods that worker threads are allowed to call via IPC
|
||
const ALLOWED_STORAGE_METHODS = new Set([
|
||
'getFormsByOrganization',
|
||
'getForm',
|
||
'getFormStatuses',
|
||
'getTasksByForm',
|
||
'getTask',
|
||
'createTask',
|
||
'updateTask',
|
||
'deleteTask',
|
||
'getTaskFieldValues',
|
||
'getFormFields',
|
||
'getTaskAssignees',
|
||
'updateTaskFieldValue',
|
||
'createTaskFieldValue',
|
||
'getUser',
|
||
'getUsersByOrganization',
|
||
'updateUser',
|
||
'setUserProfileFieldValueFromAutomation',
|
||
'recalcUserFieldAverageFromAutomation',
|
||
// Не storage-методы: обрабатываются отдельно в обработчике storage-call
|
||
'sendTaskMessageFromAutomation',
|
||
'createDelayedStatusChangeFromAutomation',
|
||
]);
|
||
|
||
interface WorkerStorageCall {
|
||
type: 'storage-call';
|
||
callId: number;
|
||
method: string;
|
||
args: unknown[];
|
||
}
|
||
|
||
interface WorkerDone {
|
||
type: 'done';
|
||
logs: string[];
|
||
error?: string;
|
||
result?: unknown;
|
||
}
|
||
|
||
type WorkerMessage = WorkerStorageCall | WorkerDone;
|
||
|
||
export function runAutomationInWorker(
|
||
code: string,
|
||
organizationId: number,
|
||
triggerData: unknown,
|
||
): Promise<{ logs: string[]; error?: string; result?: unknown }> {
|
||
return new Promise((resolve) => {
|
||
const worker = new Worker(AUTOMATION_WORKER_CODE, {
|
||
workerData: { code, organizationId, triggerData },
|
||
eval: true,
|
||
});
|
||
|
||
let settled = false;
|
||
let hardStop: ReturnType<typeof setTimeout>;
|
||
|
||
function cleanup(result: { logs: string[]; error?: string; result?: unknown }) {
|
||
if (settled) return;
|
||
settled = true;
|
||
clearTimeout(hardStop);
|
||
worker.removeAllListeners();
|
||
worker.terminate().catch(() => {});
|
||
resolve(result);
|
||
}
|
||
|
||
hardStop = setTimeout(() => {
|
||
cleanup({
|
||
logs: [`ERROR: Script execution timed out (${WORKER_TIMEOUT_MS / 1000}s) — worker terminated`],
|
||
error: `Script execution timed out (${WORKER_TIMEOUT_MS / 1000}s)`,
|
||
});
|
||
}, WORKER_TIMEOUT_MS);
|
||
|
||
worker.on('message', async (msg: WorkerMessage) => {
|
||
if (settled) return;
|
||
|
||
if (msg.type === 'storage-call') {
|
||
const { callId, method, args } = msg;
|
||
|
||
if (!ALLOWED_STORAGE_METHODS.has(method)) {
|
||
try {
|
||
worker.postMessage({
|
||
type: 'storage-response',
|
||
callId,
|
||
error: `Storage method '${method}' is not allowed`,
|
||
});
|
||
} catch { /* worker may have been terminated */ }
|
||
return;
|
||
}
|
||
|
||
// Отложенная смена статуса из автоматизации: планирует переход воркером
|
||
// delayed-status-change.worker (отменяется, если задача сменит статус раньше)
|
||
if (method === 'createDelayedStatusChangeFromAutomation') {
|
||
try {
|
||
const [taskId, orgId, targetStatusId, delayMs] = args as [number, number, number, number];
|
||
const task = await storage.getTask(Number(taskId), Number(orgId));
|
||
if (!task) throw new Error(`Задача ${taskId} не найдена`);
|
||
const statuses = await storage.getFormStatuses(task.formId, Number(orgId));
|
||
if (!statuses.some(s => s.id === Number(targetStatusId))) {
|
||
throw new Error(`Статус ${targetStatusId} не принадлежит форме ${task.formId}`);
|
||
}
|
||
const delay = Number(delayMs);
|
||
if (!Number.isFinite(delay) || delay < 0) throw new Error('delayMs должен быть неотрицательным числом');
|
||
const created = await storage.createDelayedStatusChange({
|
||
organizationId: Number(orgId),
|
||
taskId: task.id,
|
||
formId: task.formId,
|
||
expectedStatusId: task.currentStatusId,
|
||
targetStatusId: Number(targetStatusId),
|
||
runAt: new Date(Date.now() + delay),
|
||
});
|
||
if (!settled) {
|
||
try { worker.postMessage({ type: 'storage-response', callId, result: created }); } catch { /* terminated */ }
|
||
}
|
||
} catch (err: unknown) {
|
||
const message = err instanceof Error ? err.message : String(err);
|
||
if (!settled) {
|
||
try { worker.postMessage({ type: 'storage-response', callId, error: message }); } catch { /* terminated */ }
|
||
}
|
||
}
|
||
return;
|
||
}
|
||
|
||
// Сообщение в чат задачи из автоматизации: не storage-метод,
|
||
// идём через сервис сообщений (SSE, вложения, messageType 'status_change')
|
||
if (method === 'sendTaskMessageFromAutomation') {
|
||
try {
|
||
const [taskId, orgId, message] = args as [number, number, unknown];
|
||
const task = await storage.getTask(Number(taskId), Number(orgId));
|
||
if (!task) throw new Error(`Задача ${taskId} не найдена`);
|
||
const created = await sendTaskMessage({
|
||
task,
|
||
organizationId: Number(orgId),
|
||
message: String(message ?? ''),
|
||
messageType: 'status_change',
|
||
});
|
||
if (!settled) {
|
||
try { worker.postMessage({ type: 'storage-response', callId, result: created }); } catch { /* terminated */ }
|
||
}
|
||
} catch (err: unknown) {
|
||
const message = err instanceof Error ? err.message : String(err);
|
||
if (!settled) {
|
||
try { worker.postMessage({ type: 'storage-response', callId, error: message }); } catch { /* terminated */ }
|
||
}
|
||
}
|
||
return;
|
||
}
|
||
|
||
try {
|
||
const fn = (storage as unknown as Record<string, (...a: unknown[]) => Promise<unknown>>)[method];
|
||
if (typeof fn !== 'function') {
|
||
throw new Error(`Storage method '${method}' not found`);
|
||
}
|
||
const result = await fn.apply(storage, args);
|
||
if (!settled) {
|
||
try { worker.postMessage({ type: 'storage-response', callId, result }); } catch { /* terminated */ }
|
||
}
|
||
} catch (err: unknown) {
|
||
const message = err instanceof Error ? err.message : String(err);
|
||
if (!settled) {
|
||
try { worker.postMessage({ type: 'storage-response', callId, error: message }); } catch { /* terminated */ }
|
||
}
|
||
}
|
||
return;
|
||
}
|
||
|
||
if (msg.type === 'done') {
|
||
cleanup({ logs: msg.logs, error: msg.error, result: msg.result });
|
||
}
|
||
});
|
||
|
||
worker.on('error', err => {
|
||
cleanup({ logs: [`ERROR: Worker error — ${err.message}`], error: err.message });
|
||
});
|
||
|
||
worker.on('exit', code => {
|
||
cleanup({
|
||
logs: [`ERROR: Worker exited unexpectedly (code ${code})`],
|
||
error: `Worker exited unexpectedly (code ${code})`,
|
||
});
|
||
});
|
||
});
|
||
}
|
||
|
||
export interface AutomationRunResult {
|
||
automationId: number;
|
||
name: string;
|
||
logs: string[];
|
||
error?: string;
|
||
result?: unknown;
|
||
}
|
||
|
||
export async function runAutomationsByTrigger(
|
||
organizationId: number,
|
||
trigger: string,
|
||
triggerData: unknown,
|
||
): Promise<AutomationRunResult[]> {
|
||
const automations = await storage.getAutomations(organizationId);
|
||
const results: AutomationRunResult[] = [];
|
||
|
||
for (const automation of automations) {
|
||
if (!automation.isActive) continue;
|
||
if (automation.trigger !== trigger) continue;
|
||
|
||
const config = (automation.triggerConfig as Record<string, unknown>) || {};
|
||
const configFormId = config.formId;
|
||
if (configFormId !== undefined && triggerData && typeof triggerData === 'object') {
|
||
const data = triggerData as Record<string, unknown>;
|
||
const dataFormId = data.formId ?? (data.task as Record<string, unknown>)?.formId;
|
||
if (Number(configFormId) !== Number(dataFormId)) continue;
|
||
}
|
||
|
||
// Для task.status_changed: triggerConfig.statusId ограничивает срабатывание
|
||
// переходом ТОЛЬКО в указанный статус (без statusId — на любую смену статуса)
|
||
if (trigger === 'task.status_changed' && config.statusId !== undefined && triggerData && typeof triggerData === 'object') {
|
||
const data = triggerData as Record<string, unknown>;
|
||
if (Number(config.statusId) !== Number(data.newStatusId)) continue;
|
||
}
|
||
|
||
const { logs, error, result } = await runAutomationInWorker(
|
||
automation.code,
|
||
organizationId,
|
||
triggerData,
|
||
);
|
||
// Логи выполнения триггерных автоматизаций идут в консоль сервиса (иначе результаты теряются)
|
||
if (error) {
|
||
console.error(`[automation] ${automation.id} "${automation.name}" — ошибка:`, error, '| логи:', logs);
|
||
} else {
|
||
console.log(`[automation] ${automation.id} "${automation.name}" — выполнена | логи:`, logs);
|
||
}
|
||
results.push({ automationId: automation.id, name: automation.name, logs, error, result });
|
||
}
|
||
|
||
return results;
|
||
}
|
||
|
||
// =====================
|
||
// Automations
|
||
// =====================
|
||
|
||
// GET /api/automations
|
||
router.get('/api/automations', authenticateToken, tenantIsolation, requirePermission('automations.manage'), async (req: AuthenticatedRequest, res) => {
|
||
try {
|
||
const list = await storage.getAutomations(req.organizationId!);
|
||
res.json({ success: true, automations: list });
|
||
} catch (err) {
|
||
console.error('Get automations error:', err);
|
||
res.status(500).json({ error: 'Ошибка загрузки автоматизаций' });
|
||
}
|
||
});
|
||
|
||
// GET /api/automations/:id
|
||
router.get('/api/automations/:id', authenticateToken, tenantIsolation, requirePermission('automations.manage'), async (req: AuthenticatedRequest, res) => {
|
||
try {
|
||
const id = parseInt(req.params.id);
|
||
if (isNaN(id)) return res.status(400).json({ error: 'Неверный ID' });
|
||
const item = await storage.getAutomation(id, req.organizationId!);
|
||
if (!item) return res.status(404).json({ error: 'Автоматизация не найдена' });
|
||
res.json({ success: true, automation: item });
|
||
} catch (err) {
|
||
console.error('Get automation error:', err);
|
||
res.status(500).json({ error: 'Ошибка загрузки автоматизации' });
|
||
}
|
||
});
|
||
|
||
// POST /api/automations
|
||
router.post('/api/automations', authenticateToken, tenantIsolation, requirePermission('automations.manage'), async (req: AuthenticatedRequest, res) => {
|
||
try {
|
||
const { name, description, code, trigger, triggerConfig, isActive, runOffline, clientCompatible } = req.body;
|
||
if (!name?.trim()) return res.status(400).json({ error: 'Название обязательно' });
|
||
const item = await storage.createAutomation({
|
||
organizationId: req.organizationId!,
|
||
name: name.trim(),
|
||
description: description ?? null,
|
||
code: code ?? '',
|
||
trigger: trigger ?? 'manual',
|
||
triggerConfig: triggerConfig ?? null,
|
||
isActive: isActive ?? true,
|
||
runOffline: runOffline ?? false,
|
||
clientCompatible: clientCompatible ?? false,
|
||
createdBy: req.user!.id,
|
||
});
|
||
res.json({ success: true, automation: item });
|
||
} catch (err) {
|
||
console.error('Create automation error:', err);
|
||
res.status(500).json({ error: 'Ошибка создания автоматизации' });
|
||
}
|
||
});
|
||
|
||
// PUT /api/automations/:id
|
||
router.put('/api/automations/:id', authenticateToken, tenantIsolation, requirePermission('automations.manage'), async (req: AuthenticatedRequest, res) => {
|
||
try {
|
||
const id = parseInt(req.params.id);
|
||
if (isNaN(id)) return res.status(400).json({ error: 'Неверный ID' });
|
||
const existing = await storage.getAutomation(id, req.organizationId!);
|
||
if (!existing) return res.status(404).json({ error: 'Автоматизация не найдена' });
|
||
const { name, description, code, trigger, triggerConfig, isActive, runOffline, clientCompatible } = req.body;
|
||
const updates: Record<string, unknown> = {};
|
||
if (name !== undefined) updates.name = name;
|
||
if (description !== undefined) updates.description = description;
|
||
if (code !== undefined) updates.code = code;
|
||
if (trigger !== undefined) updates.trigger = trigger;
|
||
if (triggerConfig !== undefined) updates.triggerConfig = triggerConfig;
|
||
if (isActive !== undefined) updates.isActive = isActive;
|
||
if (runOffline !== undefined) updates.runOffline = runOffline;
|
||
if (clientCompatible !== undefined) updates.clientCompatible = clientCompatible;
|
||
const item = await storage.updateAutomation(id, req.organizationId!, updates as any);
|
||
res.json({ success: true, automation: item });
|
||
} catch (err) {
|
||
console.error('Update automation error:', err);
|
||
res.status(500).json({ error: 'Ошибка обновления автоматизации' });
|
||
}
|
||
});
|
||
|
||
// DELETE /api/automations/:id
|
||
router.delete('/api/automations/:id', authenticateToken, tenantIsolation, requirePermission('automations.manage'), async (req: AuthenticatedRequest, res) => {
|
||
try {
|
||
const id = parseInt(req.params.id);
|
||
if (isNaN(id)) return res.status(400).json({ error: 'Неверный ID' });
|
||
const existing = await storage.getAutomation(id, req.organizationId!);
|
||
if (!existing) return res.status(404).json({ error: 'Автоматизация не найдена' });
|
||
await storage.deleteAutomation(id, req.organizationId!);
|
||
res.json({ success: true });
|
||
} catch (err) {
|
||
console.error('Delete automation error:', err);
|
||
res.status(500).json({ error: 'Ошибка удаления автоматизации' });
|
||
}
|
||
});
|
||
|
||
// POST /api/automations/:id/run — execute automation in a Worker Thread
|
||
router.post('/api/automations/:id/run', authenticateToken, tenantIsolation, requirePermission('automations.manage'), async (req: AuthenticatedRequest, res) => {
|
||
const id = parseInt(req.params.id);
|
||
if (isNaN(id)) return res.status(400).json({ error: 'Неверный ID' });
|
||
const automation = await storage.getAutomation(id, req.organizationId!);
|
||
if (!automation) return res.status(404).json({ error: 'Автоматизация не найдена' });
|
||
|
||
const orgId = req.organizationId!;
|
||
|
||
// Memory guard: reject if process heap is > 80% of the V8 heap limit
|
||
const heapLimit = v8.getHeapStatistics().heap_size_limit;
|
||
const { heapUsed } = process.memoryUsage();
|
||
if (heapUsed > heapLimit * 0.8) {
|
||
const usedMB = Math.round(heapUsed / 1024 / 1024);
|
||
const limitMB = Math.round(heapLimit / 1024 / 1024);
|
||
console.warn(`[automation] Rejected run of automation ${automation.id} — heap ${usedMB}/${limitMB} MB`);
|
||
return res.json({
|
||
success: false,
|
||
logs: [`REJECTED: heap usage too high (${usedMB} MB / ${limitMB} MB)`],
|
||
error: 'Недостаточно памяти для запуска автоматизации',
|
||
});
|
||
}
|
||
|
||
console.log(`[automation] Starting worker for automation ${automation.id} (org=${orgId})`);
|
||
|
||
const { logs, error, result } = await runAutomationInWorker(
|
||
automation.code,
|
||
orgId,
|
||
req.body.triggerData ?? null,
|
||
);
|
||
|
||
if (error) {
|
||
const taskId: number | undefined =
|
||
typeof req.body.triggerData?.taskId === 'number' ? req.body.triggerData.taskId : undefined;
|
||
|
||
if (taskId) {
|
||
storage.addTaskAuditLog({
|
||
taskId,
|
||
organizationId: orgId,
|
||
action: 'automation.error',
|
||
fieldName: 'automation',
|
||
newValue: String(automation.id),
|
||
changedBy: req.user?.id ?? null,
|
||
changedByName: 'system',
|
||
metadata: {
|
||
automationId: automation.id,
|
||
automationName: automation.name,
|
||
error,
|
||
},
|
||
}).catch(e => console.error('[automation] Failed to log error to task_audit_log:', e));
|
||
} else {
|
||
logAudit({
|
||
action: 'automation.error',
|
||
userId: req.user?.id ?? null,
|
||
organizationId: orgId,
|
||
details: {
|
||
automationId: automation.id,
|
||
automationName: automation.name,
|
||
error,
|
||
},
|
||
});
|
||
}
|
||
|
||
return res.json({ success: false, logs, error });
|
||
}
|
||
|
||
res.json({ success: true, logs, result });
|
||
});
|
||
|
||
// POST /api/automations/run-by-trigger — execute all automations matching a trigger
|
||
router.post('/api/automations/run-by-trigger', authenticateToken, tenantIsolation, async (req: AuthenticatedRequest, res) => {
|
||
try {
|
||
const { trigger, triggerData } = req.body;
|
||
if (!trigger || typeof trigger !== 'string') {
|
||
return res.status(400).json({ error: 'Неверный триггер' });
|
||
}
|
||
|
||
const results = await runAutomationsByTrigger(req.organizationId!, trigger, triggerData ?? null);
|
||
res.json({ success: true, results });
|
||
} catch (err) {
|
||
console.error('Run by trigger error:', err);
|
||
res.status(500).json({ error: 'Ошибка выполнения автоматизаций' });
|
||
}
|
||
});
|
||
|
||
export default router;
|