Files
iistwin/server/services/task-message.service.ts
Ильяс Султанов 05e09fa08b MCP: сообщения задач, исполнители, документы, файлы, аудит — этап 2
- list_task_messages, get_task_assignees, list_document_templates, get_task_file, get_task_audit_log (read)
- send_task_message, generate_document (write); set_task_assignees (full)
- Логика POST /api/tasks/:id/messages вынесена в server/services/task-message.service.ts (переиспользуется route и MCP)
- get_task_file: нативный S3 presigned URL или локальный путь с пояснением
2026-07-21 17:19:22 +03:00

277 lines
11 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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<TaskMessage> {
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<number, Bot>();
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;
}