/
h0tnanny
/
IotPlatform
Обзор
Документация
Войти
/
h0tnanny
/
IotPlatform
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
develop
src/api/controllers/WorkflowController.ts
612 строк
18 KB
h0tnanny
Парсинг данных как строчки
20 фев 2026, 19:31
20 фев 2026, 19:31
fbb7684
Код
Авторство
О чём код?
import { Response } from 'express'; import { randomUUID } from 'crypto'; import { Main } from '../../runtime/Main'; import { Workflow } from '../../entities/Workflow'; import { ActionNode } from '../../workflow/nodes/ActionNode'; import { ConditionNode } from '../../workflow/nodes/ConditionNode'; import { LoopNode } from '../../workflow/nodes/LoopNode'; import { Field } from '../../entities/Field'; import { WorkflowRepository } from '../../repositories/WorkflowRepository'; import { PermissionRepository } from '../../repositories/PermissionRepository'; import { AuthRequest } from '../../middleware/auth'; import { LogService } from '../../services/LogService'; import { WorkflowValidationService } from '../../services/WorkflowValidationService'; export class WorkflowController { private readonly workflowRepository: WorkflowRepository; private readonly permissionRepo: PermissionRepository; constructor(private readonly main: Main) { this.workflowRepository = new WorkflowRepository(); this.permissionRepo = new PermissionRepository(); } /** * POST /workflows - Создать граф (узлы и связи) */ createWorkflow = async (req: AuthRequest, res: Response): Promise<void> => { let { id, name, description, updateInterval, isActivated, startNodeId, nodes, variableList, triggerType, triggerConfig, } = req.body; // Автоматически генерируем GUID, если ID не предоставлен if (!id) { id = randomUUID(); } if (!startNodeId || !nodes || !Array.isArray(nodes)) { res.status(400).json({ error: { code: 'VALIDATION_ERROR', message: 'Необходимы поля: startNodeId, nodes (массив)', }, }); return; } const workflow = new Workflow( id, isActivated ?? false, startNodeId, name, description, updateInterval, triggerType || 'manual', triggerConfig || {} ); // Добавляем переменные if (variableList && Array.isArray(variableList)) { for (const varData of variableList) { const field: Field = { name: varData.name, description: varData.description || '', type: varData.type, value: varData.value, }; workflow.addVariable(field); } } // Добавляем узлы for (const nodeData of nodes) { let node; switch (nodeData.type) { case 'ActionNode': node = new ActionNode( nodeData.id, nodeData.name || nodeData.id, nodeData.code || '', nodeData.nextNodeId ?? null ); break; case 'ConditionNode': node = new ConditionNode( nodeData.id, nodeData.name || nodeData.id, nodeData.fieldName, nodeData.operator, nodeData.compareValue, nodeData.trueNextId ?? null, nodeData.falseNextId ?? null ); break; case 'LoopNode': node = new LoopNode( nodeData.id, nodeData.name || nodeData.id, nodeData.fieldName, nodeData.operator, nodeData.compareValue, nodeData.bodyNodeId, nodeData.exitNodeId ?? null, nodeData.maxIterations ); break; default: res.status(400).json({ error: { code: 'INVALID_NODE_TYPE', message: `Неизвестный тип узла: ${nodeData.type}`, }, }); return; } workflow.addNode(node); } this.main.addWorkflow(workflow); // Сохраняем в БД try { await this.workflowRepository.save(workflow); // Устанавливаем владельца if (req.user) { await this.permissionRepo.setOwner('workflow', workflow.id, req.user.userId); } } catch (error) { console.error('[WorkflowController] Ошибка сохранения workflow в БД:', error); } // Валидация графа (не блокирует создание, возвращает предупреждения) const variablesForValidation = Array.isArray(variableList) ? variableList.map((v: { name: string; type: string }) => ({ name: v.name, type: v.type })) : []; const validation = WorkflowValidationService.validate(nodes, startNodeId, variablesForValidation); LogService.getInstance().info('workflow', `Workflow создан: "${workflow.name || workflow.id}"`, { userId: req.user?.userId, resourceId: workflow.id, resourceType: 'workflow', }); res.status(201).json({ data: { id: workflow.id, isActivated: workflow.isActivated, startNodeId: workflow.startNodeId, nodesCount: workflow.nodes.size, variablesCount: workflow.variableList.size, validation, }, }); }; /** * GET /workflows/:id - Получить workflow по ID */ getWorkflow = (req: AuthRequest, res: Response): void => { const { id } = req.params; const workflow = this.main.getWorkflow(id); const permission = req.resourcePermission ?? 'editor'; if (!workflow) { res.status(404).json({ error: { code: 'NOT_FOUND', message: `Workflow с ID "${id}" не найден`, }, }); return; } // Конвертируем узлы в формат для фронтенда const nodes = Array.from(workflow.nodes.values()).map((node) => { const baseNode = { id: node.id, name: node.name, type: node.constructor.name, }; if (node instanceof ActionNode) { return { ...baseNode, code: node.code, nextNodeId: node.nextNodeId, }; } if (node instanceof ConditionNode) { return { ...baseNode, fieldName: node.fieldName, operator: node.operator, compareValue: node.compareValue, trueNextId: node.trueNextId, falseNextId: node.falseNextId, }; } if (node instanceof LoopNode) { return { ...baseNode, fieldName: node.fieldName, operator: node.operator, compareValue: node.compareValue, bodyNodeId: node.bodyNodeId, exitNodeId: node.exitNodeId, maxIterations: node.maxIterations, }; } return baseNode; }); const variableList = Array.from(workflow.variableList.values()).map((field) => ({ name: field.name, description: field.description, type: field.type, value: field.value, })); res.json({ data: { id: workflow.id, name: workflow.name, description: workflow.description, updateInterval: workflow.updateInterval, isActivated: workflow.isActivated, startNodeId: workflow.startNodeId, triggerType: workflow.triggerType, triggerConfig: workflow.triggerConfig, nodes, variableList, permission, }, }); }; /** * PUT /workflows/:id - Обновить существующий workflow */ updateWorkflow = async (req: AuthRequest, res: Response): Promise<void> => { const { id } = req.params; const existingWorkflow = this.main.getWorkflow(id); if (!existingWorkflow) { res.status(404).json({ error: { code: 'NOT_FOUND', message: `Workflow с ID "${id}" не найден`, }, }); return; } const { name, description, updateInterval, isActivated, startNodeId, nodes, variableList, triggerType, triggerConfig, } = req.body; if (!startNodeId || !nodes || !Array.isArray(nodes)) { res.status(400).json({ error: { code: 'VALIDATION_ERROR', message: 'Необходимы поля: startNodeId, nodes (массив)', }, }); return; } // Удаляем старый workflow this.main.removeWorkflow(id); // Создаем новый с теми же данными const workflow = new Workflow( id, isActivated ?? false, startNodeId, name, description, updateInterval, triggerType || 'manual', triggerConfig || {} ); // Добавляем переменные if (variableList && Array.isArray(variableList)) { for (const varData of variableList) { const field: Field = { name: varData.name, description: varData.description || '', type: varData.type, value: varData.value, }; workflow.addVariable(field); } } // Добавляем узлы for (const nodeData of nodes) { let node; switch (nodeData.type) { case 'ActionNode': node = new ActionNode( nodeData.id, nodeData.name || nodeData.id, nodeData.code || '', nodeData.nextNodeId ?? null ); break; case 'ConditionNode': node = new ConditionNode( nodeData.id, nodeData.name || nodeData.id, nodeData.fieldName, nodeData.operator, nodeData.compareValue, nodeData.trueNextId ?? null, nodeData.falseNextId ?? null ); break; case 'LoopNode': node = new LoopNode( nodeData.id, nodeData.name || nodeData.id, nodeData.fieldName, nodeData.operator, nodeData.compareValue, nodeData.bodyNodeId, nodeData.exitNodeId ?? null, nodeData.maxIterations ); break; default: res.status(400).json({ error: { code: 'INVALID_NODE_TYPE', message: `Неизвестный тип узла: ${nodeData.type}`, }, }); return; } workflow.addNode(node); } this.main.addWorkflow(workflow); // Сохраняем в БД try { await this.workflowRepository.save(workflow); } catch (error) { console.error('[WorkflowController] Ошибка обновления workflow в БД:', error); // Продолжаем выполнение, так как workflow уже в памяти } // Валидация графа const variablesForValidation = Array.isArray(variableList) ? variableList.map((v: { name: string; type: string }) => ({ name: v.name, type: v.type })) : []; const validation = WorkflowValidationService.validate(nodes, startNodeId, variablesForValidation); LogService.getInstance().info('workflow', `Workflow обновлён: "${workflow.name || workflow.id}"`, { userId: req.user?.userId, resourceId: workflow.id, resourceType: 'workflow', }); res.json({ data: { id: workflow.id, isActivated: workflow.isActivated, startNodeId: workflow.startNodeId, nodesCount: workflow.nodes.size, variablesCount: workflow.variableList.size, validation, }, }); }; /** * POST /workflows/:id/execute - Ручной запуск цепочки */ executeWorkflow = async (req: AuthRequest, res: Response): Promise<void> => { const { id } = req.params; const workflow = this.main.getWorkflow(id); if (!workflow) { res.status(404).json({ error: { code: 'NOT_FOUND', message: `Workflow с ID "${id}" не найден`, }, }); return; } LogService.getInstance().info('workflow', `Workflow запущен вручную: "${workflow.name || id}"`, { userId: req.user?.userId, resourceId: id, resourceType: 'workflow', }); // Запускаем асинхронно, не блокируя ответ workflow.invoke().catch((error) => { console.error(`[Workflow ${id}] Ошибка выполнения:`, error); workflow.executionState.isRunning = false; workflow.executionState.error = error.message || 'Неизвестная ошибка'; LogService.getInstance().error('workflow', `Ошибка выполнения workflow "${id}": ${error.message}`, { resourceId: id, resourceType: 'workflow', details: { stack: error.stack }, }); }); res.json({ data: { message: `Workflow "${id}" запущен`, }, }); }; /** * GET /workflows/:id/execution - Получить состояние выполнения workflow */ getWorkflowExecution = (req: AuthRequest, res: Response): void => { const { id } = req.params; const workflow = this.main.getWorkflow(id); if (!workflow) { res.status(404).json({ error: { code: 'NOT_FOUND', message: `Workflow с ID "${id}" не найден`, }, }); return; } res.json({ data: workflow.executionState, }); }; /** * GET /state - Показать текущие значения всех переменных во всех workflow */ getState = async (req: AuthRequest, res: Response): Promise<void> => { let workflows = this.main.getAllWorkflows(); let equipment = this.main.getAllEquipment(); // Фильтрация по доступу (admin видит всё) if (req.user && req.user.role !== 'admin') { try { const accessibleWorkflowIds = await this.permissionRepo.getAccessibleResourceIds('workflow', req.user.userId); const accessibleEquipmentIds = await this.permissionRepo.getAccessibleResourceIds('equipment', req.user.userId); workflows = workflows.filter((w) => accessibleWorkflowIds.includes(w.id)); equipment = equipment.filter((e) => accessibleEquipmentIds.includes(e.id)); } catch (_err) { // При ошибке показываем всё (graceful degradation) } } const workflowIds = workflows.map((w) => w.id); const equipmentIds = equipment.map((e) => e.id); const userId = req.user?.userId ?? ''; const workflowPerms = userId ? await this.permissionRepo.getResourcePermissionsBulk('workflow', userId, workflowIds) : {}; const equipmentPerms = userId ? await this.permissionRepo.getResourcePermissionsBulk('equipment', userId, equipmentIds) : {}; const state = { workflows: workflows.map((w) => ({ id: w.id, name: w.name, description: w.description, isActivated: w.isActivated, startNodeId: w.startNodeId, permission: workflowPerms[w.id] ?? (req.user?.role === 'admin' ? 'editor' : 'viewer'), variables: Object.fromEntries( Array.from(w.variableList.entries()).map(([key, field]) => [ key, { name: field.name, description: field.description, type: field.type, value: field.value, }, ]) ), nodes: Array.from(w.nodes.keys()), })), equipment: equipment.map((e) => ({ id: e.id, name: e.name, protocol: e.protocol, dataMode: e.dataMode, endpoint: e.endpoint, pollInterval: e.pollInterval, payloadFormat: e.payloadFormat, payloadOptions: e.payloadOptions ?? undefined, mqttBrokerUrl: e.mqttBrokerUrl, mqttTopic: e.mqttTopic, permission: equipmentPerms[e.id] ?? (req.user?.role === 'admin' ? 'editor' : 'viewer'), variables: e.getFieldsCopy().map((f) => ({ name: f.name, description: f.description, type: f.type, value: f.value, })), })), }; res.json({ data: state }); }; /** * POST /workflows/validate - Валидация workflow без сохранения */ validateWorkflow = (req: AuthRequest, res: Response): void => { const { startNodeId, nodes, variableList } = req.body; if (!startNodeId || !nodes || !Array.isArray(nodes)) { res.status(400).json({ error: { code: 'VALIDATION_ERROR', message: 'Необходимы поля: startNodeId, nodes (массив)', }, }); return; } const variables = Array.isArray(variableList) ? variableList.map((v: { name: string; type: string }) => ({ name: v.name, type: v.type })) : []; const result = WorkflowValidationService.validate(nodes, startNodeId, variables); res.json({ data: result }); }; /** * DELETE /workflows/:id - Удалить workflow */ deleteWorkflow = async (req: AuthRequest, res: Response): Promise<void> => { const { id } = req.params; // Проверяем существование workflow const workflow = this.main.getWorkflow(id); if (!workflow) { res.status(404).json({ error: { code: 'NOT_FOUND', message: `Workflow с ID "${id}" не найден`, }, }); return; } // Удаляем из runtime const removedFromRuntime = this.main.removeWorkflow(id); // Удаляем из БД и разрешения try { await this.workflowRepository.delete(id); await this.permissionRepo.deleteResourcePermission('workflow', id); } catch (error) { console.error('[WorkflowController] Ошибка удаления workflow из БД:', error); } LogService.getInstance().info('workflow', `Workflow удалён: "${workflow.name || id}"`, { userId: req.user?.userId, resourceId: id, resourceType: 'workflow', }); res.json({ data: { id, removed: removedFromRuntime, message: `Workflow "${workflow.name}" успешно удален`, }, }); }; }