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; 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 Promise>)[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 { 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) || {}; const configFormId = config.formId; if (configFormId !== undefined && triggerData && typeof triggerData === 'object') { const data = triggerData as Record; const dataFormId = data.formId ?? (data.task as Record)?.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; 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 = {}; 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;