diff --git a/server/mcp.ts b/server/mcp.ts index 5bec0f6..bdad8d7 100644 --- a/server/mcp.ts +++ b/server/mcp.ts @@ -43,6 +43,19 @@ import { } from "./services/embedding.service"; import { formatUserName } from "./utils/formatUserName"; import { buildDataTableTree } from "./utils/data-table-tree"; +import { sendTaskMessage, SendTaskMessageError } from "./services/task-message.service"; +import { tasksMinimalCache } from "./utils/cache"; +import { evaluateAutoTransitions } from "./utils/auto-transitions"; +import { notifyTaskAssigned } from "./utils/notifyAssignee"; +import { eventBus } from "./routes/shared"; +import { DocumentTemplateService } from "./documents/template.service"; +import { DocumentGenerationService } from "./documents/generation.service"; +import { DataResolutionService } from "./documents/data-resolution.service"; +import { AssetService } from "./documents/asset.service"; +import { isS3Enabled, getPresignedUrl } from "./utils/s3"; +import { db } from "./db"; +import { fileUploads } from "@shared/schema"; +import { eq } from "drizzle-orm"; /** Format JS page code with consistent indentation before storing in DB. */ function formatPageCode(code: string): string { @@ -119,6 +132,11 @@ const READ_TOOLS: readonly string[] = [ 'list_directory_rows', 'get_directory_row', 'get_directory_column_values', + 'list_task_messages', + 'get_task_assignees', + 'list_document_templates', + 'get_task_file', + 'get_task_audit_log', ]; // WRITE_EXTRA_TOOLS — дополнительно доступны в режимах write и full: создание данных. @@ -128,6 +146,8 @@ const WRITE_EXTRA_TOOLS: readonly string[] = [ 'link_tasks', 'create_directory_row', 'bulk_create_directory_rows', + 'send_task_message', + 'generate_document', ]; // Все остальные инструменты (изменение/удаление форм, задач, пользователей, @@ -220,6 +240,9 @@ function buildMcpServer(organizationId: number, scopes: ApiKeyScopes): McpServer isError: true, }); + // Универсальная ошибка MCP-инструмента в формате { error } (русский текст) + const mcpError = directoryError; + // Нормализация values строки справочника до длины массива columns: // лишние значения обрезаются, недостающие дополняются пустыми строками // (как normalizeValues в server/routes/data-tables-sync.routes.ts). @@ -3537,6 +3560,458 @@ To block task creation from task.before_create, set: ctx.result = { allow: false } ); + // ── Сообщения задач ──────────────────────────────────────────────────────── + + // list_task_messages + register( + "list_task_messages", + { + title: "List Task Messages", + description: "List messages (chat) of a task, oldest first. Use afterId for incremental polling.", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task"), + afterId: z.number().int().optional().describe("Only messages with ID greater than this (incremental polling)"), + limit: z.number().int().min(1).max(500).optional().describe("Max messages to return (default 50)"), + }, + }, + async ({ taskId, afterId, limit }) => { + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + + const messages = await storage.getTaskMessages(taskId, organizationId, afterId); + const limited = messages.slice(0, limit ?? 50); + return { + content: [{ + type: "text" as const, + text: JSON.stringify(limited.map((m) => ({ + id: m.id, + message: m.message, + messageType: m.messageType, + authorId: m.authorId, + author: m.author + ? (`${m.author.firstName || ''} ${m.author.middleName || ''} ${m.author.lastName || ''}`.trim() || null) + : null, + bot: m.bot ?? null, + replyToMessageId: m.replyToMessageId, + mentionedUserIds: m.mentionedUserIds, + attachments: m.attachments, + createdAt: m.createdAt, + })), null, 2), + }], + }; + } + ); + + // send_task_message + register( + "send_task_message", + { + title: "Send Task Message", + description: + "Post a comment message to a task chat. Triggers the same side effects as the web UI: " + + "notifications, SSE events, webhooks, embeddings. The message is authored by the first admin user of the organization.", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task"), + content: z.string().min(1).describe("Message text"), + }, + }, + async ({ taskId, content }) => { + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + if (!content.trim()) return mcpError("Текст сообщения обязателен"); + + const orgUsers = await storage.getUsersByOrganization(organizationId); + const adminUser = orgUsers.find((u) => u.appRole === "admin") ?? orgUsers[0]; + if (!adminUser) return mcpError("В организации нет пользователей для авторства сообщения"); + + try { + const created = await sendTaskMessage({ + task, + user: adminUser, + organizationId, + message: content, + }); + return { + content: [{ + type: "text" as const, + text: JSON.stringify({ + success: true, + message: { + id: created.id, + taskId: created.taskId, + authorId: created.authorId, + message: created.message, + messageType: created.messageType, + createdAt: created.createdAt, + }, + }, null, 2), + }], + }; + } catch (err: unknown) { + if (err instanceof SendTaskMessageError) return mcpError(err.message); + const msg = err instanceof Error ? err.message : String(err); + return mcpError(`Ошибка отправки сообщения: ${msg}`); + } + } + ); + + // ── Исполнители задачи ───────────────────────────────────────────────────── + + // get_task_assignees + register( + "get_task_assignees", + { + title: "Get Task Assignees", + description: "List assignees (исполнители) of a task with user id, name and email", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task"), + }, + }, + async ({ taskId }) => { + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + + const assignees = await storage.getTaskAssignees(taskId, organizationId); + return { + content: [{ + type: "text" as const, + text: JSON.stringify({ + taskId, + assignedTo: task.assignedTo, + assignees: assignees.map((a) => ({ + userId: a.userId, + name: (`${a.user.firstName || ''} ${a.user.lastName || ''}`.trim()) || null, + email: a.user.email, + })), + }, null, 2), + }], + }; + } + ); + + // set_task_assignees + register( + "set_task_assignees", + { + title: "Set Task Assignees", + description: + "Replace the full list of task assignees (исполнители). " + + "The legacy assignedTo field is synced to the first user in the list (or null when empty). " + + "Triggers auto-transitions, audit log entry and notifications, like the web UI.", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task"), + userIds: z.array(z.number().int()).describe("Full new list of assignee user IDs (empty array to clear all)"), + }, + }, + async ({ taskId, userIds }) => { + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + + const orgUsers = await storage.getUsersByOrganization(organizationId); + const orgUserIds = new Set(orgUsers.map((u) => u.id)); + const uniqueIds = [...new Set(userIds)]; + const invalid = uniqueIds.filter((id) => !orgUserIds.has(id)); + if (invalid.length > 0) { + return mcpError(`Пользователи не найдены в организации: ${invalid.join(', ')}`); + } + + const adminUser = orgUsers.find((u) => u.appRole === "admin") ?? orgUsers[0]; + + // Полная замена списка: удаляем лишних, добавляем недостающих + const current = await storage.getTaskAssignees(taskId, organizationId); + const currentIds = current.map((a) => a.userId); + const toAdd = uniqueIds.filter((id) => !currentIds.includes(id)); + const toRemove = currentIds.filter((id) => !uniqueIds.includes(id)); + for (const id of toRemove) { + await storage.removeTaskAssignee(taskId, id, organizationId); + } + for (const id of toAdd) { + await storage.addTaskAssignee(taskId, id, organizationId); + } + + // Синхронизация legacy-поля assignedTo (как в route assignees: первый из списка или null) + const newAssignedTo = uniqueIds[0] ?? null; + await storage.updateTask(taskId, organizationId, { assignedTo: newAssignedTo }); + tasksMinimalCache.invalidatePrefix(`tasks:${organizationId}:minimal:`); + const autoResult = await evaluateAutoTransitions(taskId, organizationId, { triggeredBy: adminUser?.id ?? null }); + + const editorName = adminUser + ? (`${adminUser.firstName || ''} ${adminUser.middleName || ''} ${adminUser.lastName || ''}`.trim() || adminUser.email || 'MCP') + : 'MCP'; + const names = uniqueIds + .map((id) => { + const u = orgUsers.find((x) => x.id === id); + return u ? (`${u.firstName || ''} ${u.lastName || ''}`.trim() || u.email) : String(id); + }) + .join(', '); + storage.addTaskAuditLog({ + taskId, + organizationId, + action: 'task.updated', + fieldName: 'Ответственные обновлены', + oldValue: null, + newValue: names || '(пусто)', + changedBy: adminUser?.id ?? null, + changedByName: editorName, + }).catch((e: unknown) => { console.error('Audit log error (MCP set_task_assignees):', e); }); + + const refreshedTask = await storage.getTask(taskId, organizationId); + + // Уведомление новому основному ответственному (если сменился) + if (refreshedTask && newAssignedTo && newAssignedTo !== task.assignedTo) { + notifyTaskAssigned(refreshedTask, newAssignedTo, adminUser?.id ?? null, organizationId) + .catch((err) => console.error('[MCP set_task_assignees] notifyTaskAssigned error:', err)); + } + + eventBus.publishEvent({ + type: 'task_updated', + organizationId, + data: { taskId, formId: refreshedTask?.formId, task: refreshedTask }, + }); + + return { + content: [{ + type: "text" as const, + text: JSON.stringify({ + success: true, + taskId, + assignedTo: newAssignedTo, + assignees: uniqueIds, + added: toAdd, + removed: toRemove, + autoTransition: autoResult.changed, + }, null, 2), + }], + }; + } + ); + + // ── Документы ────────────────────────────────────────────────────────────── + + // list_document_templates + register( + "list_document_templates", + { + title: "List Document Templates", + description: "List document templates of the organization with their variables. Optionally filter by folder.", + inputSchema: { + folderId: z.number().int().optional().describe("Filter by template folder ID (optional)"), + }, + }, + async ({ folderId }) => { + const templateService = new DocumentTemplateService(); + const templates = await templateService.list({ organizationId, folderId }); + const result = await Promise.all( + templates.map(async (t) => { + const full = await templateService.getById(t.id, organizationId); + return { + id: t.id, + name: t.name, + description: t.description, + folderId: t.folderId, + categoryId: t.categoryId, + formId: t.formId, + status: t.status, + variables: (full?.variables ?? []).map((v) => ({ + id: v.id, + code: v.code, + label: v.label, + source: v.source, + })), + }; + }) + ); + return { + content: [{ type: "text" as const, text: JSON.stringify(result, null, 2) }], + }; + } + ); + + // generate_document + register( + "generate_document", + { + title: "Generate Document", + description: + "Generate a document (docx or pdf) from a template for a given task. " + + "Returns the generation ID and a download URL (requires user authorization to download).", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task to fill the template with"), + templateId: z.number().int().describe("The numeric ID of the document template"), + format: z.enum(["docx", "pdf"]).describe("Output format"), + }, + }, + async ({ taskId, templateId, format }) => { + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + + const orgUsers = await storage.getUsersByOrganization(organizationId); + const adminUser = orgUsers.find((u) => u.appRole === "admin") ?? orgUsers[0]; + if (!adminUser) return mcpError("В организации нет пользователей для генерации документа"); + + // Сервис генерации собирается так же, как в server/documents/routes.ts + const templateService = new DocumentTemplateService(); + const dataResolution = new DataResolutionService(); + const assetService = new AssetService(); + const generationService = new DocumentGenerationService(templateService, dataResolution, assetService); + + try { + const result = await generationService.generate({ + templateId, + taskId, + organizationId, + userId: adminUser.id, + outputFormat: format, + }); + return { + content: [{ + type: "text" as const, + text: JSON.stringify({ + success: true, + generationId: result.generationId, + downloadUrl: `/api/documents/generations/${result.generationId}/download/${format}`, + pdfUrl: result.pdfUrl ?? null, + docxUrl: result.docxUrl ?? null, + status: result.status, + unresolvedVariables: result.unresolvedVariables, + }, null, 2), + }], + }; + } catch (err: unknown) { + const msg = err instanceof Error ? err.message : String(err); + return mcpError(`Ошибка генерации документа: ${msg}`); + } + } + ); + + // ── Файлы ────────────────────────────────────────────────────────────────── + + // get_task_file + register( + "get_task_file", + { + title: "Get Task File", + description: + "Get a download URL for a file attached to a task (by fileKey from message attachments or file field values). " + + "In S3/MinIO mode returns a presigned URL (valid ~5 minutes). In local mode returns a direct path that requires user authorization.", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task"), + fileKey: z.string().min(1).describe("File key (the part after /api/files/ or /uploads/ in the attachment URL)"), + }, + }, + async ({ taskId, fileKey }) => { + if (!fileKey || fileKey.includes('..') || fileKey.includes('/')) { + return mcpError("Неверный ключ файла"); + } + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + + // Файл должен отслеживаться в file_uploads и принадлежать организации (как canAccessFile в index.ts) + const [upload] = await db + .select() + .from(fileUploads) + .where(eq(fileUploads.fileKey, fileKey)) + .limit(1); + if (!upload) return mcpError("Файл не найден или не отслеживается"); + if (upload.organizationId !== organizationId) { + return mcpError("Файл принадлежит другой организации"); + } + + // Проверка принадлежности файла именно этой задаче: + // напрямую (file_uploads.taskId), через вложения сообщений или через значения file-полей + let belongsToTask = upload.taskId === taskId; + if (!belongsToTask) { + const messages = await storage.getTaskMessages(taskId, organizationId); + belongsToTask = messages.some((m) => + Array.isArray(m.attachments) && + m.attachments.some((a) => typeof a?.url === 'string' && a.url.includes(fileKey)) + ); + } + if (!belongsToTask) { + const fieldValues = await storage.getTaskFieldValues(taskId, organizationId); + belongsToTask = fieldValues.some((fv) => typeof fv.value === 'string' && fv.value.includes(fileKey)); + } + if (!belongsToTask) { + return mcpError(`Файл не относится к задаче ${taskId}`); + } + + if (isS3Enabled) { + const presignedUrl = await getPresignedUrl(fileKey, 300); + if (presignedUrl) { + return { + content: [{ + type: "text" as const, + text: JSON.stringify({ + fileKey, + originalName: upload.originalName, + downloadUrl: presignedUrl, + type: "presigned", + expiresInSeconds: 300, + }, null, 2), + }], + }; + } + } + + // Локальный режим (или ошибка presigned): отдаём прямой путь с пояснением + return { + content: [{ + type: "text" as const, + text: JSON.stringify({ + fileKey, + originalName: upload.originalName, + downloadUrl: isS3Enabled ? `/api/files/${fileKey}` : `/uploads/${fileKey}`, + type: "direct", + note: "Presigned URL недоступен (локальный режим хранения). Ссылка требует авторизации пользователя (JWT/сессия); временную ссылку выдаёт GET /api/files/:key/presigned.", + }, null, 2), + }], + }; + } + ); + + // ── Аудит задачи ─────────────────────────────────────────────────────────── + + // get_task_audit_log + register( + "get_task_audit_log", + { + title: "Get Task Audit Log", + description: "Get the audit log entries of a task (field changes, status changes, etc.), newest first", + inputSchema: { + taskId: z.number().int().describe("The numeric ID of the task"), + limit: z.number().int().min(1).max(500).optional().describe("Max entries to return (default 50)"), + }, + }, + async ({ taskId, limit }) => { + const task = await storage.getTask(taskId, organizationId); + if (!task) return mcpError(`Задача ${taskId} не найдена`); + if (!isFormAllowed(task.formId)) return formDenied(task.formId); + + const entries = await storage.getTaskAuditLog(taskId, organizationId); + return { + content: [{ + type: "text" as const, + text: JSON.stringify(entries.slice(0, limit ?? 50).map((e) => ({ + id: e.id, + action: e.action, + fieldName: e.fieldName, + oldValue: e.oldValue, + newValue: e.newValue, + changedBy: e.changedBy, + changedByName: e.changedByName, + createdAt: e.createdAt, + })), null, 2), + }], + }; + } + ); + return server; } diff --git a/server/routes/chat.messages.routes.ts b/server/routes/chat.messages.routes.ts index ee88e71..743f59a 100644 --- a/server/routes/chat.messages.routes.ts +++ b/server/routes/chat.messages.routes.ts @@ -4,20 +4,9 @@ import { storage } from "../storage"; import { authenticateToken, tryBotServiceToken, type AuthenticatedRequest } from "../middleware/auth.middleware"; import { tenantIsolation } from "../middleware/tenant.middleware"; import { validateRequest } from "../middleware/validation.middleware"; -import { - createTaskMessageSchema, - type NotificationEvent, type Bot, -} from "@shared/schema"; -import { notificationService, EVENT_TYPES } from "../services/notification.service"; -import { webhookService } from "../services/webhook.service"; -import { generateBotServiceToken } from "../utils/jwt"; -import { sendWebhook } from "../utils/webhook"; +import { createTaskMessageSchema } from "@shared/schema"; import { eventBus } from "./shared"; - -import { handleAiBotMention, checkBotAccess } from "../services/ai-bot.service"; -import { db } from "../db"; -import { fileUploads } from "@shared/schema"; -import { eq } from "drizzle-orm"; +import { sendTaskMessage, SendTaskMessageError } from "../services/task-message.service"; export function registerChatMessageRoutes(router: Router): void { // Get task messages @@ -86,238 +75,34 @@ export function registerChatMessageRoutes(router: Router): void { } } - let replyToMessage = null; - if (req.body.replyToMessageId) { - const existingMessages = await storage.getTaskMessages(taskId, req.organizationId!); - replyToMessage = existingMessages.find(m => m.id === req.body.replyToMessageId); - if (!replyToMessage) { - return res.status(400).json({ success: false, error: 'Сообщение для ответа не найдено' }); - } - } - - let mentionedUserIds: number[] = []; - if (req.body.mentionedUserIds && req.body.mentionedUserIds.length > 0) { - const orgUsers = await storage.getUsersByOrganization(req.organizationId!); - const orgUserIds = orgUsers.map(user => user.id); - mentionedUserIds = req.body.mentionedUserIds.filter((id: number) => - orgUserIds.includes(id) && id !== req.user!.id - ); - } - - // Prevent bot response loops: skip bot processing when: - // 1. messageType is 'bot' or 'system' (messages authored by the bot service) - // 2. authorId is null in the payload (another indicator of a bot/system message) - // 3. Request was authenticated via a bot-service JWT (req.isBotToken = true) - const isBotOrSystemMessage = - req.body.messageType === 'bot' || - req.body.messageType === 'system' || - req.body.authorId === null || - req.isBotToken === true; - - let mentionedBotIds: number[] = []; - let mentionedBotsMap = new Map(); - if (!isBotOrSystemMessage && req.body.mentionedBotIds && req.body.mentionedBotIds.length > 0) { - const orgBots = await storage.getBotsByOrganization(req.organizationId!); - const activeBots = orgBots.filter(b => b.isActive); - activeBots.forEach(bot => mentionedBotsMap.set(bot.id, bot)); - const orgBotIds = activeBots.map(bot => bot.id); - mentionedBotIds = req.body.mentionedBotIds.filter((id: number) => orgBotIds.includes(id)); - } - - const messageData = { - taskId, - formId: task.formId, - authorId: req.user!.id, - replyToMessageId: req.body.replyToMessageId || null, - message: req.body.message || '', - messageType: req.body.messageType || 'comment', - mentionedUserIds: mentionedUserIds.length > 0 ? mentionedUserIds : null, - attachments: req.body.attachments?.length ? req.body.attachments : null, - }; - - const createdMessage = await storage.createTaskMessage(messageData, req.organizationId!); - - // Link file uploads to this task - if (createdMessage.attachments && createdMessage.attachments.length > 0) { - try { - for (const att of createdMessage.attachments) { - const fileKey = att.url.startsWith('/api/files/') ? att.url.replace('/api/files/', '') : (att.url.startsWith('/uploads/') ? att.url.replace('/uploads/', '') : att.url); - await db.update(fileUploads).set({ taskId }).where(eq(fileUploads.fileKey, fileKey)); - } - } catch (err) { - console.warn('[Chat] Failed to link file uploads to task:', err); - } - } - - if (mentionedUserIds.length > 0) { - try { - await Promise.all( - mentionedUserIds.map(uid => - storage.upsertTaskUserAccess(taskId, uid, req.organizationId!, 'mention') - ) - ); - // Notify subscribers that the task has changed (new mention access) - eventBus.publishEvent({ - type: 'task_updated', - data: { taskId: taskId, formId: task.formId, task }, - organizationId: req.organizationId!, - taskId: taskId - }); - } catch (err) { - console.error('upsertTaskUserAccess error for mention:', err); - } - } - - if (mentionedBotIds.length > 0) { - const taskFieldValues = await storage.getTaskFieldValues(taskId, req.organizationId!); - const form = await storage.getForm(task.formId, req.organizationId!); - const taskWithFields = { ...task, fieldValues: taskFieldValues }; - - for (const botId of mentionedBotIds) { - const bot = mentionedBotsMap.get(botId); - if (!bot) continue; - - if (bot.type === 'ai_assistant') { - if (!checkBotAccess(bot, req.user!)) { - let denialMessage; - try { - denialMessage = await storage.createTaskMessage({ - taskId, - formId: task.formId, - authorId: null, - botId: bot.id, - replyToMessageId: null, - message: `⛔ У вас нет доступа к боту ${bot.name}`, - messageType: 'system', - mentionedUserIds: null, - attachments: null, - }, req.organizationId!); - } catch (err) { - console.error(`[AI-Bot] Failed to post denial message for bot ${bot.id}:`, err); - } - if (denialMessage) { - eventBus.publishEvent({ - type: 'message_created', - data: { - taskId, - message: { - ...denialMessage, - bot: bot ? { id: bot.id, name: bot.name, avatarUrl: bot.avatarUrl ?? null } : null, - }, - }, - organizationId: req.organizationId, - taskId, - }); - } - continue; - } - - handleAiBotMention({ - bot, - task, - message: createdMessage.message, - attachments: createdMessage.attachments || undefined, - user: req.user!, - organizationId: req.organizationId!, - }).catch(err => console.error(`[AI-Bot] Error for bot ${bot.name}:`, err)); - // postBotReply внутри handleAiBotMention сам публикует message_created с полным сообщением - - continue; - } - - if (bot.webhookUrl) { - try { - const botAccessToken = generateBotServiceToken(bot.id, req.organizationId!); - const webhookPayload = { - event: 'bot_mentioned', - timestamp: new Date().toISOString(), - organizationId: req.organizationId, - data: { - task: taskWithFields, - form: form ? { id: form.id, name: form.name } : null, - message: { - id: createdMessage.id, - text: createdMessage.message, - authorId: req.user!.id, - authorName: `${req.user!.firstName} ${req.user!.middleName || ''} ${req.user!.lastName}`.trim() - }, - bot: { id: bot.id, name: bot.name, accessToken: botAccessToken } - } - }; - sendWebhook( - bot.webhookUrl, - webhookPayload, - { - 'X-Bot-Id': bot.id.toString(), - 'X-Organization-Id': req.organizationId!.toString() - }, - { taskId, organizationId: req.organizationId!, botId: bot.id, eventType: 'bot_mentioned' } - ).catch(err => console.error(`Webhook to bot ${bot.name} failed:`, err)); - } catch (err) { - console.error(`Error sending webhook to bot ${bot.id}:`, err); - } - } - } - } - - eventBus.publishEvent({ - type: 'message_created', - data: { taskId: taskId, message: createdMessage }, - organizationId: req.organizationId, - taskId: taskId - }); - - const notificationEvent: NotificationEvent = { - type: replyToMessage ? EVENT_TYPES.TASK_COMMENT_REPLIED : EVENT_TYPES.TASK_COMMENT_CREATED, - organizationId: req.organizationId!, - triggeredBy: req.user!.id, - taskId: taskId, - formId: task.formId, - messageId: createdMessage.id, - payload: { - message: createdMessage.message.length > 100 - ? createdMessage.message.substring(0, 100) + '...' - : createdMessage.message, - taskTitle: task.title, - authorName: `${req.user!.firstName} ${req.user!.middleName || ''} ${req.user!.lastName}`.trim(), - originalAuthorId: replyToMessage?.authorId, - }, - mentionedUserIds: mentionedUserIds, - timestamp: new Date(), - }; - - notificationService.processEvent(notificationEvent).then(notifiedUserIds => { - notifiedUserIds.forEach(userId => { - eventBus.publishEvent({ - type: 'notification', - data: { type: notificationEvent.type, taskId, messageId: createdMessage.id }, - organizationId: req.organizationId!, - userId: userId - }); + // Основная логика (сообщение + side-эффекты) вынесена в сервис — + // переиспользуется также MCP-сервером (инструмент send_task_message). + try { + const createdMessage = await sendTaskMessage({ + task, + user: req.user!, + organizationId: req.organizationId!, + message: req.body.message || '', + messageType: req.body.messageType, + replyToMessageId: req.body.replyToMessageId, + mentionedUserIds: req.body.mentionedUserIds, + mentionedBotIds: req.body.mentionedBotIds, + attachments: req.body.attachments, + bodyAuthorId: req.body.authorId, + isBotToken: req.isBotToken === true, }); - }).catch(err => console.error('Notification processing error:', err)); - if (messageData.messageType === 'comment') { - webhookService.dispatchComment( - req.organizationId!, taskId, task.formId, createdMessage, req.user!.id - ).catch(err => console.error('Webhook dispatch error:', err)); + res.status(201).json({ + success: true, + message: 'Сообщение добавлено', + taskMessage: createdMessage + }); + } catch (err) { + if (err instanceof SendTaskMessageError) { + return res.status(err.status).json({ success: false, error: err.message }); + } + throw err; } - - if (messageData.messageType === 'comment') { - storage.enqueueEmbedding(req.organizationId!, 'task_message', createdMessage.id, 'upsert') - .catch(err => console.error('[RAG] enqueue message embedding error:', err)); - } - - if (req.user?.id) { - storage.recordTaskInteractionAuto(taskId, req.user.id, req.organizationId!) - .catch(err => console.error('recordTaskInteraction error:', err)); - } - - res.status(201).json({ - success: true, - message: 'Сообщение добавлено', - taskMessage: createdMessage - }); } catch (error) { console.error('Create task message error:', error); res.status(500).json({ success: false, error: 'Ошибка при создании сообщения' }); diff --git a/server/services/task-message.service.ts b/server/services/task-message.service.ts new file mode 100644 index 0000000..4457cc5 --- /dev/null +++ b/server/services/task-message.service.ts @@ -0,0 +1,276 @@ +import { storage } from "../storage"; +import type { NotificationEvent, Bot, Task, User, TaskMessage } from "@shared/schema"; +import { notificationService, EVENT_TYPES } from "./notification.service"; +import { webhookService } from "./webhook.service"; +import { generateBotServiceToken } from "../utils/jwt"; +import { sendWebhook } from "../utils/webhook"; +import { eventBus } from "../routes/shared"; +import { handleAiBotMention, checkBotAccess } from "./ai-bot.service"; +import { db } from "../db"; +import { fileUploads } from "@shared/schema"; +import { eq } from "drizzle-orm"; + +// Ошибка отправки сообщения с HTTP-статусом — REST-роут маппит её в ответ, +// MCP превращает в текстовую ошибку инструмента. +export class SendTaskMessageError extends Error { + constructor(message: string, public readonly status: number = 400) { + super(message); + this.name = 'SendTaskMessageError'; + } +} + +type MessageAttachment = { url: string; name: string; size: number; mimeType?: string }; + +export interface SendTaskMessageParams { + task: Task; // задача (уже загружена и проверена вызывающим кодом) + user: User; // автор сообщения + organizationId: number; + message: string; + messageType?: string; // 'comment' (default), 'bot', 'system', ... + replyToMessageId?: number | null; + mentionedUserIds?: number[] | null; + mentionedBotIds?: number[] | null; + attachments?: MessageAttachment[] | null; + bodyAuthorId?: number | null; // authorId из тела запроса (null — признак бот/системного сообщения) + isBotToken?: boolean; // запрос аутентифицирован bot-service JWT +} + +// Создаёт сообщение задачи со всеми side-эффектами: +// привязка файловых вложений к задаче, доступ по упоминанию, обработка ботов +// (AI-ассистенты и webhook-боты), SSE-события, уведомления, исходящие вебхуки, +// постановка embedding в очередь, запись взаимодействия с задачей. +// Логика вынесена из POST /api/tasks/:id/messages (routes/chat.messages.routes.ts) +// и переиспользуется MCP-сервером — поведение не менять. +export async function sendTaskMessage(params: SendTaskMessageParams): Promise { + const { task, user, organizationId } = params; + const taskId = task.id; + + let replyToMessage = null; + if (params.replyToMessageId) { + const existingMessages = await storage.getTaskMessages(taskId, organizationId); + replyToMessage = existingMessages.find(m => m.id === params.replyToMessageId); + if (!replyToMessage) { + throw new SendTaskMessageError('Сообщение для ответа не найдено'); + } + } + + let mentionedUserIds: number[] = []; + if (params.mentionedUserIds && params.mentionedUserIds.length > 0) { + const orgUsers = await storage.getUsersByOrganization(organizationId); + const orgUserIds = orgUsers.map(u => u.id); + mentionedUserIds = params.mentionedUserIds.filter(id => + orgUserIds.includes(id) && id !== user.id + ); + } + + // Prevent bot response loops: skip bot processing when: + // 1. messageType is 'bot' or 'system' (messages authored by the bot service) + // 2. authorId is null in the payload (another indicator of a bot/system message) + // 3. Request was authenticated via a bot-service JWT (isBotToken = true) + const isBotOrSystemMessage = + params.messageType === 'bot' || + params.messageType === 'system' || + params.bodyAuthorId === null || + params.isBotToken === true; + + let mentionedBotIds: number[] = []; + let mentionedBotsMap = new Map(); + if (!isBotOrSystemMessage && params.mentionedBotIds && params.mentionedBotIds.length > 0) { + const orgBots = await storage.getBotsByOrganization(organizationId); + const activeBots = orgBots.filter(b => b.isActive); + activeBots.forEach(bot => mentionedBotsMap.set(bot.id, bot)); + const orgBotIds = activeBots.map(bot => bot.id); + mentionedBotIds = params.mentionedBotIds.filter((id: number) => orgBotIds.includes(id)); + } + + const messageData = { + taskId, + formId: task.formId, + authorId: user.id, + replyToMessageId: params.replyToMessageId || null, + message: params.message || '', + messageType: params.messageType || 'comment', + mentionedUserIds: mentionedUserIds.length > 0 ? mentionedUserIds : null, + attachments: params.attachments?.length ? params.attachments : null, + }; + + const createdMessage = await storage.createTaskMessage(messageData, organizationId); + + // Link file uploads to this task + if (createdMessage.attachments && createdMessage.attachments.length > 0) { + try { + for (const att of createdMessage.attachments) { + const fileKey = att.url.startsWith('/api/files/') ? att.url.replace('/api/files/', '') : (att.url.startsWith('/uploads/') ? att.url.replace('/uploads/', '') : att.url); + await db.update(fileUploads).set({ taskId }).where(eq(fileUploads.fileKey, fileKey)); + } + } catch (err) { + console.warn('[Chat] Failed to link file uploads to task:', err); + } + } + + if (mentionedUserIds.length > 0) { + try { + await Promise.all( + mentionedUserIds.map(uid => + storage.upsertTaskUserAccess(taskId, uid, organizationId, 'mention') + ) + ); + // Notify subscribers that the task has changed (new mention access) + eventBus.publishEvent({ + type: 'task_updated', + data: { taskId: taskId, formId: task.formId, task }, + organizationId, + taskId: taskId + }); + } catch (err) { + console.error('upsertTaskUserAccess error for mention:', err); + } + } + + if (mentionedBotIds.length > 0) { + const taskFieldValues = await storage.getTaskFieldValues(taskId, organizationId); + const form = await storage.getForm(task.formId, organizationId); + const taskWithFields = { ...task, fieldValues: taskFieldValues }; + + for (const botId of mentionedBotIds) { + const bot = mentionedBotsMap.get(botId); + if (!bot) continue; + + if (bot.type === 'ai_assistant') { + if (!checkBotAccess(bot, user)) { + let denialMessage; + try { + denialMessage = await storage.createTaskMessage({ + taskId, + formId: task.formId, + authorId: null, + botId: bot.id, + replyToMessageId: null, + message: `⛔ У вас нет доступа к боту ${bot.name}`, + messageType: 'system', + mentionedUserIds: null, + attachments: null, + }, organizationId); + } catch (err) { + console.error(`[AI-Bot] Failed to post denial message for bot ${bot.id}:`, err); + } + if (denialMessage) { + eventBus.publishEvent({ + type: 'message_created', + data: { + taskId, + message: { + ...denialMessage, + bot: bot ? { id: bot.id, name: bot.name, avatarUrl: bot.avatarUrl ?? null } : null, + }, + }, + organizationId, + taskId, + }); + } + continue; + } + + handleAiBotMention({ + bot, + task, + message: createdMessage.message, + attachments: createdMessage.attachments || undefined, + user, + organizationId, + }).catch(err => console.error(`[AI-Bot] Error for bot ${bot.name}:`, err)); + // postBotReply внутри handleAiBotMention сам публикует message_created с полным сообщением + + continue; + } + + if (bot.webhookUrl) { + try { + const botAccessToken = generateBotServiceToken(bot.id, organizationId); + const webhookPayload = { + event: 'bot_mentioned', + timestamp: new Date().toISOString(), + organizationId, + data: { + task: taskWithFields, + form: form ? { id: form.id, name: form.name } : null, + message: { + id: createdMessage.id, + text: createdMessage.message, + authorId: user.id, + authorName: `${user.firstName || ''} ${user.middleName || ''} ${user.lastName || ''}`.trim() + }, + bot: { id: bot.id, name: bot.name, accessToken: botAccessToken } + } + }; + sendWebhook( + bot.webhookUrl, + webhookPayload, + { + 'X-Bot-Id': bot.id.toString(), + 'X-Organization-Id': organizationId.toString() + }, + { taskId, organizationId, botId: bot.id, eventType: 'bot_mentioned' } + ).catch(err => console.error(`Webhook to bot ${bot.name} failed:`, err)); + } catch (err) { + console.error(`Error sending webhook to bot ${bot.id}:`, err); + } + } + } + } + + eventBus.publishEvent({ + type: 'message_created', + data: { taskId: taskId, message: createdMessage }, + organizationId, + taskId: taskId + }); + + const notificationEvent: NotificationEvent = { + type: replyToMessage ? EVENT_TYPES.TASK_COMMENT_REPLIED : EVENT_TYPES.TASK_COMMENT_CREATED, + organizationId, + triggeredBy: user.id, + taskId: taskId, + formId: task.formId, + messageId: createdMessage.id, + payload: { + message: createdMessage.message.length > 100 + ? createdMessage.message.substring(0, 100) + '...' + : createdMessage.message, + taskTitle: task.title, + authorName: `${user.firstName || ''} ${user.middleName || ''} ${user.lastName || ''}`.trim(), + originalAuthorId: replyToMessage?.authorId, + }, + mentionedUserIds: mentionedUserIds, + timestamp: new Date(), + }; + + notificationService.processEvent(notificationEvent).then(notifiedUserIds => { + notifiedUserIds.forEach(userId => { + eventBus.publishEvent({ + type: 'notification', + data: { type: notificationEvent.type, taskId, messageId: createdMessage.id }, + organizationId, + userId: userId + }); + }); + }).catch(err => console.error('Notification processing error:', err)); + + if (messageData.messageType === 'comment') { + webhookService.dispatchComment( + organizationId, taskId, task.formId, createdMessage, user.id + ).catch(err => console.error('Webhook dispatch error:', err)); + } + + if (messageData.messageType === 'comment') { + storage.enqueueEmbedding(organizationId, 'task_message', createdMessage.id, 'upsert') + .catch(err => console.error('[RAG] enqueue message embedding error:', err)); + } + + if (user.id) { + storage.recordTaskInteractionAuto(taskId, user.id, organizationId) + .catch(err => console.error('recordTaskInteraction error:', err)); + } + + return createdMessage; +}