Files
iistwin/server/routes/messenger.messages.routes.ts
Ильяс Султанов 9827b5df33 fix(chat): cross-tab lock реплея офлайн-очереди + идемпотентность мессенджера (conv_messages)
Клиентская часть гонки (вариант B): Web Locks API (navigator.locks) —
кросс-вкладочный мьютекс iistwin-offline-sync-replay на replay-проход в
useOfflineSync: первая вкладка под локом реплеит и удаляет записи, остальные
видят пустую очередь. Fallback для сред без Web Locks — старое поведение
(дубли гасит серверная идемпотентность). Race-тест: без locks → 2 POST,
с locks → 1 POST для одной записи очереди.

Мессенджер (conv_messages) — та же защита, что у task_messages:
- conversation_messages.client_message_id + partial unique (миграция 0087);
- POST /api/messenger/conversations/:id/messages: повтор с тем же
  clientMessageId → существующее сообщение без insert и без SSE/уведомлений,
  гонка insert'ов ловится по 23505;
- useChatController шлёт crypto.randomUUID(), attachment-реплей мессенджера
  тоже использует id записи очереди.
2026-09-29 20:31:59 +03:00

681 lines
29 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 { Router } from "express";
import { z } from "zod";
import crypto from "crypto";
import { db, openTenantCtx } from "../db";
import { conversations, conversationMembers, conversationMessages, users, bots, type Bot } from "@shared/schema";
import { eq, and, lt, gt, desc, sql, inArray } from "drizzle-orm";
import { authenticateToken, type AuthenticatedRequest } from "../middleware/auth.middleware";
import { tenantIsolation } from "../middleware/tenant.middleware";
import { eventBus } from "./shared";
import { isMember, enrichMessage, getMembersForSSE } from "./messenger.helpers";
import { formatUserName } from "../utils/formatUserName";
import { notificationService } from "../services/notification.service";
import { handleAiBotDirectMessage } from "../services/ai-bot.service";
import { checkBotAccess } from "../services/ai-bot.service";
import { handleVpnBotDirectMessage, handleVpnBotCallback } from "../vpn/vpn-bot.service";
import { generateBotServiceToken } from "../utils/jwt";
import { sendWebhook } from "../utils/webhook";
import { decrypt } from "../crypto";
export function registerMessengerMessageRoutes(router: Router): void {
// GET /api/messenger/conversations/:id/messages
router.get("/api/messenger/conversations/:id/messages",
authenticateToken,
tenantIsolation,
async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const convId = parseInt(req.params.id);
if (isNaN(convId)) return res.status(400).json({ success: false, error: "Некорректный ID" });
if (!await isMember(convId, userId)) {
return res.status(403).json({ success: false, error: "Нет доступа" });
}
const limit = Math.min(parseInt((req.query.limit as string) ?? "50"), 100);
const beforeId = req.query.beforeId ? parseInt(req.query.beforeId as string) : null;
const afterId = req.query.afterId ? parseInt(req.query.afterId as string) : null;
const conditions = [
eq(conversationMessages.conversationId, convId),
eq(conversationMessages.isDeleted, false),
];
if (beforeId !== null && !isNaN(beforeId)) {
conditions.push(lt(conversationMessages.id, beforeId));
}
if (afterId !== null && !isNaN(afterId)) {
conditions.push(gt(conversationMessages.id, afterId));
}
const msgs = await db
.select({
id: conversationMessages.id,
conversationId: conversationMessages.conversationId,
message: conversationMessages.message,
replyToId: conversationMessages.replyToId,
mentionedUserIds: conversationMessages.mentionedUserIds,
attachments: conversationMessages.attachments,
botButtons: conversationMessages.botButtons,
createdAt: conversationMessages.createdAt,
updatedAt: conversationMessages.updatedAt,
isDeleted: conversationMessages.isDeleted,
authorId: conversationMessages.authorId,
botId: conversationMessages.botId,
authorFirstName: users.firstName,
authorMiddleName: users.middleName,
authorLastName: users.lastName,
})
.from(conversationMessages)
.leftJoin(users, eq(conversationMessages.authorId, users.id))
.where(and(...conditions))
.orderBy(desc(conversationMessages.id))
.limit(limit);
// Load bot info for bot messages
const botIds = [...new Set(msgs.filter(m => m.botId).map(m => m.botId!))];
const botInfoMap = new Map<number, { id: number; name: string; avatarUrl: string | null }>();
if (botIds.length > 0) {
const botRows = await db
.select({ id: bots.id, name: bots.name, avatarUrl: bots.avatarUrl })
.from(bots)
.where(inArray(bots.id, botIds));
botRows.forEach(b => botInfoMap.set(b.id, b));
}
const enrichedMsgs = await Promise.all(
msgs.map(async m => {
let replyTo: { id: number; message: string; author: { id: number; firstName: string; middleName: string; lastName: string } } | null = null;
if (m.replyToId) {
const [r] = await db
.select({
id: conversationMessages.id,
message: conversationMessages.message,
authorId: conversationMessages.authorId,
authorFirstName: users.firstName,
authorMiddleName: users.middleName,
authorLastName: users.lastName,
})
.from(conversationMessages)
.leftJoin(users, eq(conversationMessages.authorId, users.id))
.where(and(
eq(conversationMessages.id, m.replyToId),
eq(conversationMessages.conversationId, convId),
));
if (r && r.authorId) {
replyTo = {
id: r.id,
message: r.message,
author: { id: r.authorId, firstName: r.authorFirstName ?? '', middleName: r.authorMiddleName ?? '', lastName: r.authorLastName ?? '' },
};
}
}
const botInfo = m.botId ? botInfoMap.get(m.botId) ?? null : null;
return {
...m,
author: m.authorId
? { id: m.authorId, firstName: m.authorFirstName ?? '', middleName: m.authorMiddleName ?? '', lastName: m.authorLastName ?? '' }
: null,
bot: botInfo,
replyTo,
};
})
);
res.json({ success: true, messages: enrichedMsgs.reverse() });
} catch (err) {
console.error("get messages error:", err);
res.status(500).json({ success: false, error: "Ошибка при загрузке сообщений" });
}
}
);
// POST /api/messenger/conversations/:id/messages — send message
const sendSchema = z.object({
message: z.string().max(10000),
replyToId: z.number().int().positive().optional(),
mentionedUserIds: z.array(z.number().int().positive()).optional(),
attachments: z.array(z.object({
url: z.string(),
name: z.string(),
size: z.number().default(0),
mimeType: z.string().optional(),
})).optional(),
clientMessageId: z.string().max(64).optional(),
});
router.post("/api/messenger/conversations/:id/messages",
authenticateToken,
tenantIsolation,
async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const convId = parseInt(req.params.id);
if (isNaN(convId)) return res.status(400).json({ success: false, error: "Некорректный ID" });
if (!await isMember(convId, userId)) {
return res.status(403).json({ success: false, error: "Нет доступа" });
}
const parsed = sendSchema.safeParse(req.body);
if (!parsed.success) {
return res.status(400).json({ success: false, error: "Некорректные данные", details: parsed.error.issues });
}
const { message, replyToId, mentionedUserIds, attachments } = parsed.data;
if (!message.trim() && (!attachments || attachments.length === 0)) {
return res.status(400).json({ success: false, error: "Сообщение или вложения обязательны" });
}
const [conv] = await db
.select({ organizationId: conversations.organizationId, type: conversations.type, name: conversations.name, botId: conversations.botId })
.from(conversations)
.where(eq(conversations.id, convId));
const orgId = conv?.organizationId ?? req.organizationId!;
// Идемпотентность: повтор с тем же clientMessageId — вернуть уже созданное
// сообщение БЕЗ повторного insert и без side-эффектов (SSE/уведомления
// уже отработали при первой записи). Защита от дублей при multi-tab replay
// офлайн-очереди (миграция 0087).
const clientMessageId = parsed.data.clientMessageId ?? null;
if (clientMessageId) {
const [existing] = await db
.select({ id: conversationMessages.id })
.from(conversationMessages)
.where(and(
eq(conversationMessages.conversationId, convId),
eq(conversationMessages.clientMessageId, clientMessageId)
))
.limit(1);
if (existing) {
const enrichedExisting = await enrichMessage(existing.id);
return res.status(200).json({ success: true, message: enrichedExisting, deduplicated: true });
}
}
let newMsg;
try {
[newMsg] = await db
.insert(conversationMessages)
.values({
conversationId: convId,
authorId: userId,
message: message.trim(),
replyToId: replyToId ?? null,
mentionedUserIds: mentionedUserIds?.length ? mentionedUserIds : null,
attachments: attachments?.length ? attachments : null,
clientMessageId,
})
.returning();
} catch (err) {
// Гонка: параллельный запрос с тем же clientMessageId уже вставил
// сообщение (unique index conv_messages_client_message_id_unique).
const e = err as { code?: string; cause?: { code?: string } };
if (clientMessageId && (e?.code === '23505' || e?.cause?.code === '23505')) {
const [existing] = await db
.select({ id: conversationMessages.id })
.from(conversationMessages)
.where(and(
eq(conversationMessages.conversationId, convId),
eq(conversationMessages.clientMessageId, clientMessageId)
))
.limit(1);
if (existing) {
const enrichedExisting = await enrichMessage(existing.id);
return res.status(200).json({ success: true, message: enrichedExisting, deduplicated: true });
}
}
throw err;
}
const enriched = await enrichMessage(newMsg.id);
const members = await getMembersForSSE(convId, orgId);
const [author] = await db
.select({ firstName: users.firstName, middleName: users.middleName, lastName: users.lastName })
.from(users)
.where(eq(users.id, userId));
const authorName = author ? formatUserName(author) : 'Пользователь';
const memberMuteRows = await db
.select({ userId: conversationMembers.userId, mutedAt: conversationMembers.mutedAt, externalOrgId: conversationMembers.externalOrgId })
.from(conversationMembers)
.where(eq(conversationMembers.conversationId, convId));
const chatName = conv?.type === 'group' && conv.name
? conv.name
: authorName;
const memberOrgMap = new Map(
memberMuteRows.map(r => [
r.userId,
{ orgId: r.externalOrgId ?? orgId, mutedAt: r.mutedAt ?? null },
])
);
for (const { userId: memberId, orgId: memberOrgId } of members) {
eventBus.publishEvent({
type: "conv_message_created",
data: { conversationId: convId, message: enriched },
organizationId: memberOrgId,
userId: memberId,
});
if (memberId !== userId) {
await db
.update(conversationMembers)
.set({ reactionUnreadCount: sql`COALESCE(reaction_unread_count, 0)` })
.where(and(
eq(conversationMembers.conversationId, convId),
eq(conversationMembers.userId, memberId),
));
}
}
notificationService.notifyMessengerMessage({
conversationId: convId,
authorId: userId,
authorName,
chatName,
chatType: (conv?.type ?? 'direct') as 'direct' | 'group',
messageText: message.trim(),
mentionedUserIds: mentionedUserIds ?? [],
memberOrgMap,
}).catch(err => console.error('[Messenger Notify] Error:', err));
res.status(201).json({ success: true, message: enriched });
// Trigger AI bot response for bot_direct conversations (after response sent to user)
if (conv?.type === 'bot_direct' && conv.botId) {
const botId = conv.botId;
setImmediate(async () => {
try {
const handle = await openTenantCtx(orgId);
try {
await handle.run(async () => {
const [bot] = await db
.select()
.from(bots)
.where(eq(bots.id, botId));
if (!bot || !bot.isActive) return;
if (bot.type === 'ai_assistant') {
const typedBot = bot as Bot;
if (!checkBotAccess(typedBot, req.user!)) {
const [denyMsg] = await db
.insert(conversationMessages)
.values({
conversationId: convId,
authorId: null,
botId: bot.id,
message: `⛔ У вас нет доступа к боту ${bot.name}`,
replyToId: null,
mentionedUserIds: null,
attachments: null,
})
.returning();
const enrichedDeny = await enrichMessage(denyMsg.id);
for (const { userId: memberId, orgId: memberOrgId } of members) {
eventBus.publishEvent({ type: "conv_message_created", data: { conversationId: convId, message: enrichedDeny }, organizationId: memberOrgId, userId: memberId });
}
notificationService.notifyMessengerMessage({
conversationId: convId,
authorId: -1,
authorName: bot.name,
chatName: bot.name,
chatType: 'bot_direct',
messageText: `⛔ У вас нет доступа к боту ${bot.name}`,
mentionedUserIds: [],
memberOrgMap,
}).catch(err => console.error('[Bot Deny Notify] Error:', err));
return;
}
await handleAiBotDirectMessage({
bot: typedBot,
conversationId: convId,
message: message.trim(),
user: req.user!,
organizationId: orgId,
});
const latestBotMsg = await db
.select({ id: conversationMessages.id })
.from(conversationMessages)
.where(and(
eq(conversationMessages.conversationId, convId),
eq(conversationMessages.botId, botId),
))
.orderBy(desc(conversationMessages.id))
.limit(1);
if (latestBotMsg[0]) {
const enrichedBot = await enrichMessage(latestBotMsg[0].id);
for (const { userId: memberId, orgId: memberOrgId } of members) {
eventBus.publishEvent({ type: "conv_message_created", data: { conversationId: convId, message: enrichedBot }, organizationId: memberOrgId, userId: memberId });
}
notificationService.notifyMessengerMessage({
conversationId: convId,
authorId: -1,
authorName: bot.name,
chatName: bot.name,
chatType: 'bot_direct',
messageText: (enrichedBot as any).message ?? '',
mentionedUserIds: [],
memberOrgMap,
}).catch(err => console.error('[Bot Reply Notify] Error:', err));
}
} else if (bot.type === 'vpn') {
const typedBot = bot as Bot;
if (!checkBotAccess(typedBot, req.user!)) {
const [denyMsg] = await db
.insert(conversationMessages)
.values({
conversationId: convId,
authorId: null,
botId: bot.id,
message: `⛔ У вас нет доступа к боту ${bot.name}`,
replyToId: null,
mentionedUserIds: null,
attachments: null,
})
.returning();
const enrichedDeny = await enrichMessage(denyMsg.id);
for (const { userId: memberId, orgId: memberOrgId } of members) {
eventBus.publishEvent({ type: "conv_message_created", data: { conversationId: convId, message: enrichedDeny }, organizationId: memberOrgId, userId: memberId });
}
notificationService.notifyMessengerMessage({
conversationId: convId,
authorId: -1,
authorName: bot.name,
chatName: bot.name,
chatType: 'bot_direct',
messageText: `⛔ У вас нет доступа к боту ${bot.name}`,
mentionedUserIds: [],
memberOrgMap,
}).catch(err => console.error('[VPN Bot Deny Notify] Error:', err));
return;
}
await handleVpnBotDirectMessage({
bot: typedBot,
conversationId: convId,
message: message.trim(),
user: req.user!,
organizationId: orgId,
});
}
});
} finally {
await handle.release().catch(() => {});
}
} catch (err) {
console.error('[Bot Direct] Error handling bot response:', err);
}
});
}
} catch (err) {
console.error("send message error:", err);
res.status(500).json({ success: false, error: "Ошибка при отправке" });
}
}
);
// PATCH /api/messenger/messages/:id — edit message
router.patch("/api/messenger/messages/:id",
authenticateToken,
tenantIsolation,
async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const msgId = parseInt(req.params.id);
if (isNaN(msgId)) return res.status(400).json({ success: false, error: "Некорректный ID" });
const [msg] = await db
.select({ authorId: conversationMessages.authorId, conversationId: conversationMessages.conversationId, isDeleted: conversationMessages.isDeleted })
.from(conversationMessages)
.where(eq(conversationMessages.id, msgId));
if (!msg || msg.isDeleted) return res.status(404).json({ success: false, error: "Сообщение не найдено" });
if (msg.authorId !== userId) return res.status(403).json({ success: false, error: "Нельзя редактировать чужие сообщения" });
const { message } = req.body;
if (typeof message !== "string" || !message.trim()) {
return res.status(400).json({ success: false, error: "Текст не может быть пустым" });
}
await db
.update(conversationMessages)
.set({ message: message.trim(), updatedAt: new Date() })
.where(eq(conversationMessages.id, msgId));
const enriched = await enrichMessage(msgId);
const [conv] = await db
.select({ organizationId: conversations.organizationId })
.from(conversations)
.where(eq(conversations.id, msg.conversationId));
const members = await getMembersForSSE(msg.conversationId, conv?.organizationId ?? req.organizationId!);
for (const { userId: memberId, orgId } of members) {
eventBus.publishEvent({ type: "conv_message_updated", data: { conversationId: msg.conversationId, message: enriched }, organizationId: orgId, userId: memberId });
}
res.json({ success: true, message: enriched });
} catch (err) {
console.error("edit message error:", err);
res.status(500).json({ success: false, error: "Ошибка" });
}
}
);
// DELETE /api/messenger/messages/:id — soft-delete message
router.delete("/api/messenger/messages/:id",
authenticateToken,
tenantIsolation,
async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const msgId = parseInt(req.params.id);
if (isNaN(msgId)) return res.status(400).json({ success: false, error: "Некорректный ID" });
const [msg] = await db
.select({ authorId: conversationMessages.authorId, conversationId: conversationMessages.conversationId, isDeleted: conversationMessages.isDeleted })
.from(conversationMessages)
.where(eq(conversationMessages.id, msgId));
if (!msg || msg.isDeleted) return res.status(404).json({ success: false, error: "Сообщение не найдено" });
const isAdmin = req.user?.appRole === "admin";
if (msg.authorId !== userId && !isAdmin) {
return res.status(403).json({ success: false, error: "Нельзя удалить чужое сообщение" });
}
await db
.update(conversationMessages)
.set({ isDeleted: true, message: "", updatedAt: new Date() })
.where(eq(conversationMessages.id, msgId));
const [conv] = await db
.select({ organizationId: conversations.organizationId })
.from(conversations)
.where(eq(conversations.id, msg.conversationId));
const members = await getMembersForSSE(msg.conversationId, conv?.organizationId ?? req.organizationId!);
for (const { userId: memberId, orgId } of members) {
eventBus.publishEvent({ type: "conv_message_deleted", data: { id: msgId, conversationId: msg.conversationId }, organizationId: orgId, userId: memberId });
}
res.json({ success: true });
} catch (err) {
console.error("delete message error:", err);
res.status(500).json({ success: false, error: "Ошибка" });
}
}
);
// POST /api/messenger/conversations/:id/read — mark conversation as read
router.post("/api/messenger/conversations/:id/read",
authenticateToken,
tenantIsolation,
async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const convId = parseInt(req.params.id);
if (isNaN(convId)) return res.status(400).json({ success: false, error: "Некорректный ID" });
if (!await isMember(convId, userId)) {
return res.status(403).json({ success: false, error: "Нет доступа" });
}
await db
.update(conversationMembers)
.set({ lastReadAt: new Date(), reactionUnreadCount: 0 })
.where(and(
eq(conversationMembers.conversationId, convId),
eq(conversationMembers.userId, userId),
));
const [conv] = await db
.select({ organizationId: conversations.organizationId })
.from(conversations)
.where(eq(conversations.id, convId));
const members = await getMembersForSSE(convId, conv?.organizationId ?? req.organizationId!);
for (const { userId: memberId, orgId } of members) {
eventBus.publishEvent({
type: "conv_read",
data: { conversationId: convId, userId },
organizationId: orgId,
userId: memberId,
});
}
res.json({ success: true });
} catch (err) {
console.error("mark read error:", err);
res.status(500).json({ success: false, error: "Ошибка" });
}
}
);
// POST /api/messenger/messages/:id/callback — handle bot button click in messenger
router.post("/api/messenger/messages/:id/callback",
authenticateToken,
tenantIsolation,
async (req: AuthenticatedRequest, res) => {
try {
const userId = req.user!.id;
const messageId = parseInt(req.params.id);
if (isNaN(messageId)) {
return res.status(400).json({ success: false, error: "Некорректный ID сообщения" });
}
const { callbackData } = req.body;
if (!callbackData || typeof callbackData !== "string") {
return res.status(400).json({ success: false, error: "callbackData обязателен" });
}
const [msg] = await db
.select({
id: conversationMessages.id,
conversationId: conversationMessages.conversationId,
botId: conversationMessages.botId,
})
.from(conversationMessages)
.where(eq(conversationMessages.id, messageId));
if (!msg) {
return res.status(404).json({ success: false, error: "Сообщение не найдено" });
}
if (!msg.botId) {
return res.status(400).json({ success: false, error: "Сообщение не от бота" });
}
if (!await isMember(msg.conversationId, userId)) {
return res.status(403).json({ success: false, error: "Нет доступа к диалогу" });
}
const [bot] = await db
.select()
.from(bots)
.where(and(eq(bots.id, msg.botId), eq(bots.organizationId, req.organizationId!)));
if (!bot) {
return res.status(404).json({ success: false, error: "Бот не найден" });
}
// Inline handler for built-in VPN bot (no external webhook required)
if (bot.type === 'vpn') {
await handleVpnBotCallback({
bot,
conversationId: msg.conversationId,
messageId: msg.id,
callbackData,
user: req.user!,
organizationId: req.organizationId!,
});
return res.json({ success: true, message: "Callback обработан" });
}
if (!bot.webhookUrl || !bot.webhookEnabled) {
return res.status(400).json({ success: false, error: "Бот не настроен для обработки callback" });
}
const botAccessToken = generateBotServiceToken(bot.id, req.organizationId!);
const callbackPayload = {
event: "bot_direct_button_callback",
timestamp: new Date().toISOString(),
organizationId: req.organizationId,
data: {
conversationId: msg.conversationId,
messageId: msg.id,
callbackData,
bot: {
id: bot.id,
name: bot.name,
accessToken: botAccessToken,
},
user: {
id: req.user!.id,
firstName: req.user!.firstName,
middleName: req.user!.middleName,
lastName: req.user!.lastName,
email: req.user!.email,
},
},
};
let signature = "";
if (bot.webhookSecret) {
const secret = decrypt(bot.webhookSecret);
signature = crypto
.createHmac("sha256", secret)
.update(JSON.stringify(callbackPayload))
.digest("hex");
}
await sendWebhook(
bot.webhookUrl,
callbackPayload,
{
"X-Webhook-Signature": signature,
"X-Bot-Id": bot.id.toString(),
"X-Organization-Id": req.organizationId!.toString(),
},
{
conversationId: msg.conversationId,
organizationId: req.organizationId!,
botId: bot.id,
eventType: "bot_direct_button_callback",
}
);
res.json({ success: true, message: "Callback обработан" });
} catch (err) {
console.error("Messenger callback error:", err);
res.status(500).json({ success: false, error: "Ошибка обработки callback" });
}
}
);
}