import { storage } from "../storage"; import type { NotificationEvent, Bot, Task, SafeUser, 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?: SafeUser; // автор сообщения; необязателен, если передан botId botId?: number | null; // сообщение от бота: authorId=null, botId, messageType='bot' organizationId: number; message: string; messageType?: string; // 'comment' (default), 'bot', 'system', 'status_change', 'reminder', ... replyToMessageId?: number | null; mentionedUserIds?: number[] | null; mentionedBotIds?: number[] | null; attachments?: MessageAttachment[] | null; bodyAuthorId?: number | null; // authorId из тела запроса (null — признак бот/системного сообщения) isBotToken?: boolean; // запрос аутентифицирован bot-service JWT visibleToUserIds?: number[] | null; // приватность: только перечисленные пользователи видят сообщение (NULL — все участники чата) clientMessageId?: string | null; // ключ идемпотентности (uuid от клиента) — защита от дублей при повторной отправке } // PG unique violation (23505): drizzle оборачивает оригинальную ошибку, поэтому // проверяем и саму ошибку, и cause. function isUniqueViolation(err: unknown): boolean { const e = err as { code?: string; cause?: { code?: string } } | null; return e?.code === '23505' || e?.cause?.code === '23505'; } // Создаёт сообщение задачи со всеми side-эффектами: // привязка файловых вложений к задаче, доступ по упоминанию, обработка ботов // (AI-ассистенты и webhook-боты), SSE-события, уведомления, исходящие вебхуки, // постановка embedding в очередь, запись взаимодействия с задачей. // Логика вынесена из POST /api/tasks/:id/messages (routes/chat.messages.routes.ts) // и переиспользуется MCP-сервером — поведение не менять. // При botId (сообщение от бота) авторство и пользовательские side-эффекты пропускаются, // как в POST /api/bot/message (bot-api.routes.ts). export async function sendTaskMessage(params: SendTaskMessageParams): Promise { const { task, user, organizationId } = params; const botId = params.botId ?? null; // Системные сообщения ('system', 'status_change', 'reminder') могут быть без автора — // их пишут автоматизации (ctx.tasks.sendMessage, напоминания), автор в чате — «Система» const isAuthorlessSystem = !user && !botId && (params.messageType === 'system' || params.messageType === 'status_change' || params.messageType === 'reminder'); if (!user && !botId && !isAuthorlessSystem) { throw new SendTaskMessageError('Не указан автор сообщения (user или botId)', 500); } const taskId = task.id; // Идемпотентность: повтор с тем же clientMessageId — вернуть уже созданное // сообщение БЕЗ повторного insert и БЕЗ side-эффектов (уведомления, SSE, // вебхуки, embedding уже отработали при первой записи). Защита от дублей // при multi-tab replay офлайн-очереди и ретраях после сетевых сбоев. const clientMessageId = params.clientMessageId ?? null; if (clientMessageId) { const existingId = await storage.getTaskMessageIdByClientMessageId(taskId, clientMessageId); if (existingId != null) { const existing = await storage.getTaskMessage(existingId, organizationId); if (existing) return existing as TaskMessage; } } 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 messageType = params.messageType || (botId ? 'bot' : 'comment'); const isBotOrSystemMessage = messageType === 'bot' || messageType === 'system' || params.bodyAuthorId === null || params.isBotToken === true; let mentionedBotIds: number[] = []; let mentionedBotsMap = new Map(); if (!isBotOrSystemMessage && user && 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: botId ? null : (user?.id ?? null), botId, replyToMessageId: params.replyToMessageId || null, message: params.message || '', messageType, mentionedUserIds: mentionedUserIds.length > 0 ? mentionedUserIds : null, attachments: params.attachments?.length ? params.attachments : null, visibleToUserIds: params.visibleToUserIds && params.visibleToUserIds.length > 0 ? [...new Set(params.visibleToUserIds)] : null, clientMessageId, }; let createdMessage: TaskMessage; try { createdMessage = await storage.createTaskMessage(messageData, organizationId); } catch (err) { // Гонка: параллельный запрос с тем же clientMessageId уже вставил сообщение // (unique index task_messages_client_message_id_unique). Возвращаем // существующее БЕЗ повторных side-эффектов. if (clientMessageId && isUniqueViolation(err)) { const existingId = await storage.getTaskMessageIdByClientMessageId(taskId, clientMessageId); if (existingId != null) { const existing = await storage.getTaskMessage(existingId, organizationId); if (existing) return existing as TaskMessage; } } throw err; } // 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 && user) { 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); } } } } // Приватные сообщения (visible_to_user_ids): SSE-событие получают только // разрешённые пользователи (перечисленные + автор). Остальные участники чата // событие не получают и сообщение не увидят ни в реальном времени, ни при загрузке // (фильтр в storage.getTaskMessages по viewerUserId). const visibleTo = messageData.visibleToUserIds; if (visibleTo && visibleTo.length > 0) { const allowed = new Set(visibleTo); if (user?.id) allowed.add(user.id); for (const uid of allowed) { eventBus.publishEvent({ type: 'message_created', data: { taskId: taskId, message: createdMessage }, organizationId, userId: uid, }); } } else { eventBus.publishEvent({ type: 'message_created', data: { taskId: taskId, message: createdMessage }, organizationId, taskId: taskId }); } // Уведомления и исходящие вебхуки — только для сообщений от пользователя // (для bot-сообщений — как в POST /api/bot/message: без notificationService). if (user) { 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; }