Files
iistwin/server/services/task-message.service.ts
Ильяс Султанов af4850bd6d fix(chat): идемпотентность сообщений по clientMessageId — защита от дублей при повторной отправке
Первопричина дублей (3 одинаковых сообщения в задаче 1988): живой POST сорвался
→ сообщение попало в общую IndexedDB-очередь → background-sync разослал
TRIGGER_SYNC всем вкладкам → каждая вкладка реплейнула одну и ту же запись
(мьютекс _syncLock модульный, кросс-вкладочной координации нет). Гонка
воспроизведена тестом tests/offline-queue-multitab-race.test.tsx (1 запись →
2 POST для двух вкладок).

Вариант C — серверная идемпотентность:
- task_messages.client_message_id + partial unique index (миграция 0086);
- sendTaskMessage: повтор с тем же clientMessageId возвращает существующее
  сообщение БЕЗ insert и БЕЗ side-эффектов (уведомления/SSE/вебхуки не дублируются);
  гонка insert'ов ловится по 23505 → fallback на select;
- клиент TaskChat шлёт crypto.randomUUID() в каждом сообщении; при офлайн-постановке
  тело с ключом сохраняется, все реплеи идут с одним ключом;
- attachment-реплей: clientMessageId = id записи очереди (стабилен между реплеями);
- юнит-тесты tests/task-message-idempotency.test.ts (3: быстрый путь, гонка 23505,
  проброс прочих ошибок).
2026-09-29 17:40:47 +03:00

353 lines
16 KiB
TypeScript
Raw Permalink 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, 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<TaskMessage> {
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<number, Bot>();
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<number>(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;
}