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 или локальный путь с пояснением
This commit is contained in:
276
server/services/task-message.service.ts
Normal file
276
server/services/task-message.service.ts
Normal file
@@ -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<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;
|
||||
}
|
||||
Reference in New Issue
Block a user