diff --git a/client/src/App.tsx b/client/src/App.tsx index 6622591..74b8cd7 100644 --- a/client/src/App.tsx +++ b/client/src/App.tsx @@ -210,6 +210,20 @@ function ForegroundPushHandler() { const isChat = payload.data?.type === 'chat' || payload.data?.type === 'messenger'; const url = payload.url || payload.data?.url || '/'; + // Client-side suppress: don't show a foreground toast if the user is already + // looking at the corresponding chat on a visible tab. + const currentPath = window.location.pathname; + if (payload.data?.type === 'messenger' && payload.data?.conversationId) { + if (currentPath === `/chat/${payload.data.conversationId}` && document.visibilityState === 'visible') { + return; + } + } + if ((payload.type === 'chat' || payload.data?.type === 'chat') && payload.data?.taskId) { + if (currentPath.includes(`/tasks/${payload.data.taskId}`) && document.visibilityState === 'visible') { + return; + } + } + toast({ title: payload.title || 'Уведомление', description: payload.body || payload.message || '', diff --git a/client/src/components/MobileBottomNav.tsx b/client/src/components/MobileBottomNav.tsx index f8fa901..216c2f8 100644 --- a/client/src/components/MobileBottomNav.tsx +++ b/client/src/components/MobileBottomNav.tsx @@ -127,10 +127,16 @@ export function MobileBottomNav() { useEvents({ onConvMessageCreated: useCallback(() => { - queryClient.invalidateQueries({ queryKey: ['/api/messenger/unread-count'] }); + queryClient.setQueryData<{ success: boolean; unreadCount: number }>(['/api/messenger/unread-count'], (old) => { + if (!old) return old; + return { ...old, unreadCount: (old.unreadCount ?? 0) + 1 }; + }); }, []), onConvRead: useCallback(() => { - queryClient.invalidateQueries({ queryKey: ['/api/messenger/unread-count'] }); + queryClient.setQueryData<{ success: boolean; unreadCount: number }>(['/api/messenger/unread-count'], (old) => { + if (!old) return old; + return { ...old, unreadCount: 0 }; + }); }, []), onConvCreated: useCallback(() => { queryClient.invalidateQueries({ queryKey: ['/api/messenger/unread-count'] }); diff --git a/client/src/components/Sidebar.tsx b/client/src/components/Sidebar.tsx index f5b7acc..c028748 100644 --- a/client/src/components/Sidebar.tsx +++ b/client/src/components/Sidebar.tsx @@ -804,11 +804,18 @@ function MessengerNavItem({ collapsed }: { collapsed: boolean }) { useEvents({ onConvMessageCreated: useCallback(() => { - queryClient.invalidateQueries({ queryKey: ['/api/messenger/unread-count'] }); + // Optimistically bump the global unread badge immediately. + queryClient.setQueryData<{ success: boolean; unreadCount: number }>(['/api/messenger/unread-count'], (old) => { + if (!old) return old; + return { ...old, unreadCount: (old.unreadCount ?? 0) + 1 }; + }); queryClient.invalidateQueries({ queryKey: ['/api/messenger/conversations'] }); }, []), onConvRead: useCallback(() => { - queryClient.invalidateQueries({ queryKey: ['/api/messenger/unread-count'] }); + queryClient.setQueryData<{ success: boolean; unreadCount: number }>(['/api/messenger/unread-count'], (old) => { + if (!old) return old; + return { ...old, unreadCount: 0 }; + }); queryClient.invalidateQueries({ queryKey: ['/api/messenger/conversations'] }); }, []), onConvCreated: useCallback(() => { diff --git a/client/src/components/TaskChat.tsx b/client/src/components/TaskChat.tsx index 2340253..a142bc0 100644 --- a/client/src/components/TaskChat.tsx +++ b/client/src/components/TaskChat.tsx @@ -1,4 +1,4 @@ -import { useState, useRef, useEffect } from 'react'; +import { useState, useRef, useEffect, useLayoutEffect, useCallback } from 'react'; import { useQuery, useMutation } from '@tanstack/react-query'; import { useForm } from 'react-hook-form'; import { useLocation } from 'wouter'; @@ -209,6 +209,7 @@ const TaskChat = ({ taskId }: TaskChatProps) => { const [editText, setEditText] = useState(''); const textareaRef = useRef(null); const scrollAreaRef = useRef(null); + const markReadTimerRef = useRef | null>(null); const form = useForm>({ resolver: zodResolver(messageSchema), defaultValues: { @@ -258,9 +259,9 @@ const TaskChat = ({ taskId }: TaskChatProps) => { }; }); - // Автоотметка прочтения для видимого открытого чата. + // Автоотметка прочтения для видимого открытого чата (debounced). if (isVisible && data.message.id > 0) { - markAsReadMutation.mutate(data.message.id); + debouncedMarkAsRead(data.message.id); } } }, @@ -416,6 +417,13 @@ const TaskChat = ({ taskId }: TaskChatProps) => { } }); + const debouncedMarkAsRead = useCallback((messageId: number) => { + if (markReadTimerRef.current) clearTimeout(markReadTimerRef.current); + markReadTimerRef.current = setTimeout(() => { + markAsReadMutation.mutate(messageId); + }, 300); + }, []); + // Функция для проверки прочитанности сообщения const isMessageRead = (messageId: number) => { // Проверяем есть ли ID сообщения в списке прочитанных @@ -515,6 +523,15 @@ const TaskChat = ({ taskId }: TaskChatProps) => { }, [messages.length, isLoading, user, hasMarkedInitialMessages]); // Сообщаем серверу и SW, что этот task-чат открыт, чтобы подавить push. + // useLayoutEffect reports synchronously after DOM mutations so push suppression + // is active before any incoming message event is processed by the server. + useLayoutEffect(() => { + if (!taskId || taskId <= 0) return; + const visible = typeof document !== 'undefined' && document.visibilityState === 'visible'; + reportActiveChat('task', taskId, visible); + setActiveChatInSw('task', taskId, visible); + }, [taskId]); + useEffect(() => { if (!taskId || taskId <= 0) return; @@ -523,8 +540,6 @@ const TaskChat = ({ taskId }: TaskChatProps) => { setActiveChatInSw('task', taskId, visible); }; - report(document.visibilityState === 'visible'); - const handleVisibility = () => { report(document.visibilityState === 'visible'); }; diff --git a/client/src/hooks/useEvents.ts b/client/src/hooks/useEvents.ts index 036f4f4..75648c2 100644 --- a/client/src/hooks/useEvents.ts +++ b/client/src/hooks/useEvents.ts @@ -12,6 +12,8 @@ interface UseEventsOptions { onConvMessageDeleted?: (data: { conversationId: number; messageId: number }) => void; onConvRead?: (data: { conversationId: number }) => void; onConvCreated?: (data: any) => void; + onConvTypingStarted?: (data: { conversationId: number; botName: string }) => void; + onConvTypingStopped?: (data: { conversationId: number; botName: string }) => void; onReactionUpdated?: (data: { taskMessageId?: number; convMessageId?: number; conversationId?: number; reactions: any[] }) => void; onPollUpdated?: (data: { taskMessageId?: number; convMessageId?: number; conversationId?: number; poll: any }) => void; reconnectInterval?: number; @@ -58,6 +60,8 @@ export const useEvents = (options: UseEventsOptions = {}) => { if (options.onConvMessageDeleted) unsubs.push(sseManager.on('conv_message_deleted', wrap('onConvMessageDeleted'))); if (options.onConvRead) unsubs.push(sseManager.on('conv_read', wrap('onConvRead'))); if (options.onConvCreated) unsubs.push(sseManager.on('conv_created', wrap('onConvCreated'))); + if (options.onConvTypingStarted) unsubs.push(sseManager.on('conv_typing_started', wrap('onConvTypingStarted'))); + if (options.onConvTypingStopped) unsubs.push(sseManager.on('conv_typing_stopped', wrap('onConvTypingStopped'))); if (options.onReactionUpdated) unsubs.push(sseManager.on('reaction_updated', wrap('onReactionUpdated'))); if (options.onPollUpdated) unsubs.push(sseManager.on('poll_updated', wrap('onPollUpdated'))); return () => { for (const u of unsubs) u(); }; diff --git a/client/src/lib/sseManager.ts b/client/src/lib/sseManager.ts index 4b82e65..568838e 100644 --- a/client/src/lib/sseManager.ts +++ b/client/src/lib/sseManager.ts @@ -24,6 +24,8 @@ const KNOWN_EVENTS = [ 'conv_message_deleted', 'conv_read', 'conv_created', + 'conv_typing_started', + 'conv_typing_stopped', 'reaction_updated', 'poll_updated', ]; @@ -102,7 +104,12 @@ async function connect() { lastEventId = e.lastEventId; } const data = JSON.parse(e.data); - if (process.env.NODE_ENV !== 'production') { + const isChatEvent = evName === 'conv_message_created' || evName === 'conv_message_updated' || + evName === 'conv_message_deleted' || evName === 'conv_read' || evName === 'message_created' || + evName === 'message_read'; + if (isChatEvent) { + console.log(`[SSE] Received event: ${evName} id=${e.lastEventId ?? lastEventId} latency=${Date.now() - (data?.message?.createdAt ? new Date(data.message.createdAt).getTime() : 0)}ms`, data); + } else if (process.env.NODE_ENV !== 'production') { console.log(`[SSE] Received event: ${evName}`, data); } emit(evName, data); diff --git a/client/src/pages/Chat.tsx b/client/src/pages/Chat.tsx index 76757cc..3960977 100644 --- a/client/src/pages/Chat.tsx +++ b/client/src/pages/Chat.tsx @@ -1,4 +1,4 @@ -import { useState, useRef, useEffect, useCallback, useMemo } from 'react'; +import { useState, useRef, useEffect, useCallback, useMemo, useLayoutEffect } from 'react'; import { useQuery, useMutation } from '@tanstack/react-query'; import { useLocation, useParams } from 'wouter'; import { format, isToday, isYesterday } from 'date-fns'; @@ -60,6 +60,13 @@ interface ConvBot { avatarUrl?: string | null; } +interface BotButton { + text: string; + callbackData: string; + style?: 'primary' | 'secondary' | 'destructive' | 'outline'; + row?: number; +} + interface ConvMessage { id: number; conversationId: number; @@ -67,6 +74,7 @@ interface ConvMessage { replyToId?: number | null; mentionedUserIds?: number[] | null; attachments?: ChatAttachment[] | null; + botButtons?: BotButton[] | null; createdAt: string; updatedAt: string; isDeleted: boolean; @@ -124,11 +132,21 @@ function groupByDate(messages: ConvMessage[]) { return groups; } +function groupButtonsByRow(buttons: BotButton[]): BotButton[][] { + const rows = new Map(); + for (const btn of buttons) { + const rowIdx = btn.row ?? 0; + if (!rows.has(rowIdx)) rows.set(rowIdx, []); + rows.get(rowIdx)!.push(btn); + } + return Array.from(rows.entries()).sort((a, b) => a[0] - b[0]).map(([, btns]) => btns); +} + // ── Message bubble ───────────────────────────────────────────────────────────── function MessageBubble({ msg, isOwn, onReply, onEdit, onDelete, reactions, poll, onReactionsChange, onPollChange, - onOpenUserChat, onOpenBotChat, + onOpenUserChat, onOpenBotChat, onButtonClick, }: { msg: ConvMessage; isOwn: boolean; @@ -141,6 +159,7 @@ function MessageBubble({ onPollChange?: (p: PollData) => void; onOpenUserChat?: (userId: number) => void; onOpenBotChat?: (botId: number) => void; + onButtonClick?: (m: ConvMessage, button: BotButton) => void; }) { const [pickerOpen, setPickerOpen] = useState(false); const isEdited = msg.updatedAt && msg.updatedAt !== msg.createdAt; @@ -289,6 +308,27 @@ function MessageBubble({ /> )} + {/* Кнопки бота */} + {msg.botButtons && msg.botButtons.length > 0 && onButtonClick && !msg.isDeleted && ( +
+ {groupButtonsByRow(msg.botButtons).map((row, rowIdx) => ( +
+ {row.map((btn, btnIdx) => ( + + ))} +
+ ))} +
+ )} + {msg.pending && ( <> @@ -957,6 +997,9 @@ export default function ChatPage() { const [searchQuery, setSearchQuery] = useState(''); const [debouncedSearch, setDebouncedSearch] = useState(''); const searchInputRef = useRef(null); + // Typing indicator for bots + const [typingBots, setTypingBots] = useState>(new Set()); + const typingTimersRef = useRef>(new Map()); // File attachments const [pendingFiles, setPendingFiles] = useState([]); const fileInputRef = useRef(null); @@ -966,6 +1009,14 @@ export default function ChatPage() { // function to keep optimistic messages after server refetches. const activeQueueIdsRef = useRef>(new Set()); const activeAttachmentIdsRef = useRef>(new Set()); + const markReadTimerRef = useRef | null>(null); + + const debouncedMarkRead = useCallback((convId: number) => { + if (markReadTimerRef.current) clearTimeout(markReadTimerRef.current); + markReadTimerRef.current = setTimeout(() => { + markReadMutation.mutate(convId); + }, 300); + }, []); // Pagination for loading older messages const [olderMessages, setOlderMessages] = useState([]); @@ -1050,6 +1101,15 @@ export default function ChatPage() { }, [activeConvId]); // Report active chat to server and SW so push is suppressed while this chat is open. + // useLayoutEffect reports synchronously after DOM mutations so push suppression + // is active before any incoming message event is processed by the server. + useLayoutEffect(() => { + if (!activeConvId) return; + const visible = typeof document !== 'undefined' && document.visibilityState === 'visible'; + reportActiveChat('messenger', activeConvId, visible); + setActiveChatInSw('messenger', activeConvId, visible); + }, [activeConvId]); + useEffect(() => { if (!activeConvId) return; @@ -1058,8 +1118,6 @@ export default function ChatPage() { setActiveChatInSw('messenger', activeConvId, visible); }; - report(document.visibilityState === 'visible'); - const handleVisibility = () => { report(document.visibilityState === 'visible'); }; @@ -1318,6 +1376,42 @@ export default function ChatPage() { }, }); + // ── Typing indicator for bots ───────────────────────────────────────────────── + + const clearTypingBot = (botName: string) => { + const timer = typingTimersRef.current.get(botName); + if (timer) clearTimeout(timer); + typingTimersRef.current.delete(botName); + setTypingBots(prev => { + if (!prev.has(botName)) return prev; + const next = new Set(prev); + next.delete(botName); + return next; + }); + }; + + const addTypingBot = (botName: string) => { + const existing = typingTimersRef.current.get(botName); + if (existing) clearTimeout(existing); + setTypingBots(prev => { + if (prev.has(botName)) return prev; + const next = new Set(prev); + next.add(botName); + return next; + }); + typingTimersRef.current.set( + botName, + setTimeout(() => clearTypingBot(botName), 10000), + ); + }; + + useEffect(() => { + // Clear all typing indicators when switching conversations + typingTimersRef.current.forEach(timer => clearTimeout(timer)); + typingTimersRef.current.clear(); + setTypingBots(new Set()); + }, [activeConvId]); + // ── SSE ─────────────────────────────────────────────────────────────────────── useEvents({ @@ -1371,9 +1465,13 @@ export default function ChatPage() { ); // 3. Автоотметка прочтения только для активного, видимого чата и чужих сообщений. - // Используем mutation, чтобы onSuccess обновил кэш и сбросил бейджи. + // Оптимистично сбрасываем глобальный unread count, затем debounce'им запрос к серверу. if (isActiveConv && !isMine && isVisible) { - markReadMutation.mutate(conversationId); + queryClient.setQueryData<{ success: boolean; unreadCount: number }>(['/api/messenger/unread-count'], (old) => { + if (!old) return old; + return { ...old, unreadCount: 0 }; + }); + debouncedMarkRead(conversationId); } }, [activeConvId, user]), @@ -1434,6 +1532,16 @@ export default function ChatPage() { ); }, []), + onConvTypingStarted: useCallback((data: { conversationId: number; botName: string }) => { + if (data.conversationId !== activeConvId) return; + addTypingBot(data.botName); + }, [activeConvId, addTypingBot]), + + onConvTypingStopped: useCallback((data: { conversationId: number; botName: string }) => { + if (data.conversationId !== activeConvId) return; + clearTypingBot(data.botName); + }, [activeConvId, clearTypingBot]), + onReactionUpdated: useCallback((data: { taskMessageId?: number; convMessageId?: number; conversationId?: number; reactions: ReactionGroup[] }) => { if (data.convMessageId != null) { setReactionsMap(prev => ({ ...prev, [data.convMessageId!]: data.reactions })); @@ -1520,6 +1628,12 @@ export default function ChatPage() { }, }); + const botButtonCallbackMutation = useMutation({ + mutationFn: async ({ messageId, callbackData }: { messageId: number; callbackData: string }) => { + await apiRequest('POST', `/api/messenger/messages/${messageId}/callback`, { callbackData }); + }, + }); + const handleSelectConv = (id: number) => { setActiveConvId(id); setReplyTo(null); @@ -1529,6 +1643,11 @@ export default function ChatPage() { markReadMutation.mutate(id); if (isMobile) setMobileView('chat'); setLocation(`/chat/${id}`); + // Report active chat immediately so the server can suppress pushes for this chat + // before any new message event reaches the notification service. + const visible = typeof document !== 'undefined' && document.visibilityState === 'visible'; + reportActiveChat('messenger', id, visible); + setActiveChatInSw('messenger', id, visible); }; // ── Открытие/создание диалога из поиска ───────────────────────────────────── @@ -2282,12 +2401,35 @@ export default function ChatPage() { if (existing) { handleSelectConv(existing.id); return; } openBotDirectMutation.mutate(botId); }} + onButtonClick={(m, btn) => { + botButtonCallbackMutation.mutate({ messageId: m.id, callbackData: btn.callbackData }); + }} /> ))} ))} + + {/* Bot typing indicator */} + {typingBots.size > 0 && ( +
+ + 🤖 + +
+ + {Array.from(typingBots).join(', ')} + + печатает + + + + + +
+
+ )} {/* Input area */} diff --git a/migrations/0062_add_bot_buttons_to_conversation_messages.sql b/migrations/0062_add_bot_buttons_to_conversation_messages.sql new file mode 100644 index 0000000..021ddb9 --- /dev/null +++ b/migrations/0062_add_bot_buttons_to_conversation_messages.sql @@ -0,0 +1,3 @@ +-- Add bot_buttons column to conversation_messages for interactive bot buttons in messenger +ALTER TABLE conversation_messages +ADD COLUMN IF NOT EXISTS bot_buttons JSONB; diff --git a/server/routes/bot-api.routes.ts b/server/routes/bot-api.routes.ts index f30b3e7..7ee5b8c 100644 --- a/server/routes/bot-api.routes.ts +++ b/server/routes/bot-api.routes.ts @@ -506,6 +506,7 @@ export function registerBotApiRoutes(app: import("express").Express): void { replyToId: conversationMessages.replyToId, mentionedUserIds: conversationMessages.mentionedUserIds, attachments: conversationMessages.attachments, + botButtons: conversationMessages.botButtons, createdAt: conversationMessages.createdAt, updatedAt: conversationMessages.updatedAt, isDeleted: conversationMessages.isDeleted, @@ -587,9 +588,10 @@ export function registerBotApiRoutes(app: import("express").Express): void { return res.status(400).json({ success: false, error: 'Некорректный ID чата' }); } - const { message, attachments } = req.body; - if (!message?.trim() && (!attachments || attachments.length === 0)) { - return res.status(400).json({ success: false, error: 'Сообщение или вложения обязательны' }); + const { message, attachments, buttons } = req.body; + const hasButtons = Array.isArray(buttons) && buttons.length > 0; + if (!message?.trim() && (!attachments || attachments.length === 0) && !hasButtons) { + return res.status(400).json({ success: false, error: 'Сообщение, вложения или кнопки обязательны' }); } const [conv] = await db @@ -616,6 +618,7 @@ export function registerBotApiRoutes(app: import("express").Express): void { replyToId: null, mentionedUserIds: null, attachments: attachments?.length ? attachments : null, + botButtons: hasButtons ? buttons : null, }) .returning(); diff --git a/server/routes/messenger.helpers.ts b/server/routes/messenger.helpers.ts index 9f160bf..085d155 100644 --- a/server/routes/messenger.helpers.ts +++ b/server/routes/messenger.helpers.ts @@ -2,6 +2,7 @@ import { db } from "../db"; import { conversations, conversationMembers, conversationMessages, users, bots } from "@shared/schema"; import { eq, and, desc, sql } from "drizzle-orm"; import { formatUserName } from "../utils/formatUserName"; +import { eventBus } from "./shared"; export async function isMember(conversationId: number, userId: number): Promise { const [row] = await db @@ -120,6 +121,7 @@ export async function enrichMessage(msgId: number) { replyToId: conversationMessages.replyToId, mentionedUserIds: conversationMessages.mentionedUserIds, attachments: conversationMessages.attachments, + botButtons: conversationMessages.botButtons, createdAt: conversationMessages.createdAt, updatedAt: conversationMessages.updatedAt, isDeleted: conversationMessages.isDeleted, @@ -190,3 +192,21 @@ export async function getMembersForSSE(convId: number, conversationOrgId: number orgId: r.externalOrgId ?? conversationOrgId, }))); } + +export async function publishConvTypingEvent( + convId: number, + botName: string, + organizationId: number, + isTyping: boolean, +) { + const members = await getMembersForSSE(convId, organizationId); + const type = isTyping ? 'conv_typing_started' : 'conv_typing_stopped'; + for (const { userId: memberId, orgId: memberOrgId } of members) { + eventBus.publishEvent({ + type, + data: { conversationId: convId, botName }, + organizationId: memberOrgId, + userId: memberId, + }); + } +} diff --git a/server/routes/messenger.messages.routes.ts b/server/routes/messenger.messages.routes.ts index 978d7e3..475f28c 100644 --- a/server/routes/messenger.messages.routes.ts +++ b/server/routes/messenger.messages.routes.ts @@ -1,5 +1,6 @@ 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"; @@ -11,7 +12,10 @@ 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 } from "../vpn/vpn-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 @@ -50,6 +54,7 @@ export function registerMessengerMessageRoutes(router: Router): void { replyToId: conversationMessages.replyToId, mentionedUserIds: conversationMessages.mentionedUserIds, attachments: conversationMessages.attachments, + botButtons: conversationMessages.botButtons, createdAt: conversationMessages.createdAt, updatedAt: conversationMessages.updatedAt, isDeleted: conversationMessages.isDeleted, @@ -507,4 +512,125 @@ export function registerMessengerMessageRoutes(router: Router): void { } } ); + + // 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" }); + } + } + ); } diff --git a/server/routes/shared.ts b/server/routes/shared.ts index f2323d3..563ce1b 100644 --- a/server/routes/shared.ts +++ b/server/routes/shared.ts @@ -114,6 +114,15 @@ export class EventBus { timestamp: Date.now(), }; + // Логируем chat-события для диагностики задержек/пропусков + const isChatEvent = event.type === 'conv_message_created' || event.type === 'conv_message_updated' || + event.type === 'conv_message_deleted' || event.type === 'conv_read' || event.type === 'message_created' || + event.type === 'message_read'; + if (isChatEvent) { + const target = event.userId ? `user=${event.userId}` : `org=${event.organizationId}`; + console.log(`[SSE] publish ${event.type} ${target} conn=${this.connections.size} id=${bufferedEvent.id}`); + } + // Буферизуем событие для восстановления после reconnect if (event.userId) { this.pushToBuffer(this.userBuffers, event.userId, bufferedEvent); @@ -121,6 +130,7 @@ export class EventBus { this.pushToBuffer(this.orgBuffers, event.organizationId, bufferedEvent); } + let delivered = 0; this.connections.forEach((connection, id) => { // Проверяем tenant isolation if (event.organizationId && connection.organizationId !== event.organizationId) { @@ -132,8 +142,13 @@ export class EventBus { return; } - this.writeEventToConnection(connection, bufferedEvent); + const ok = this.writeEventToConnection(connection, bufferedEvent); + if (ok) delivered++; }); + + if (isChatEvent) { + console.log(`[SSE] delivered ${event.type} to ${delivered}/${this.connections.size} connection(s)`); + } } /** diff --git a/server/services/ai-bot.service.ts b/server/services/ai-bot.service.ts index e460153..9fbde90 100644 --- a/server/services/ai-bot.service.ts +++ b/server/services/ai-bot.service.ts @@ -8,6 +8,7 @@ import { decrypt } from '../crypto'; import { semanticSearch } from './embedding.service'; import { McpClient } from '../utils/mcp-client'; import { eventBus } from '../routes/shared'; +import { publishConvTypingEvent } from '../routes/messenger.helpers'; import { getOllamaNumThread, getOllamaNumCtx } from '../utils/ollama-config'; import type { Bot, Task, User } from '@shared/schema'; import { db } from '../db'; @@ -512,7 +513,10 @@ export async function handleAiBotDirectMessage(ctx: { return; } - const messages: LlmMessage[] = []; + try { + await publishConvTypingEvent(conversationId, bot.name, organizationId, true); + + const messages: LlmMessage[] = []; if (bot.sendSystemPrompt !== false) { const systemContent = bot.systemPrompt?.trim() @@ -657,7 +661,10 @@ export async function handleAiBotDirectMessage(ctx: { mentionedUserIds: null, attachments: null, }); - } catch (err) { + } finally { + await publishConvTypingEvent(conversationId, bot.name, organizationId, false).catch(() => {}); + } +} catch (err) { console.error('[AI-Bot Direct] Error:', err); await db.insert(conversationMessages).values({ conversationId, diff --git a/server/utils/webhook.ts b/server/utils/webhook.ts index f1764a8..27000db 100644 --- a/server/utils/webhook.ts +++ b/server/utils/webhook.ts @@ -3,6 +3,7 @@ import { storage } from '../storage'; export interface SendWebhookOptions { retries?: number; taskId?: number; + conversationId?: number; organizationId?: number; botId?: number; eventType?: string; diff --git a/server/vpn/vpn-bot.service.ts b/server/vpn/vpn-bot.service.ts index 0df6fae..27adb35 100644 --- a/server/vpn/vpn-bot.service.ts +++ b/server/vpn/vpn-bot.service.ts @@ -28,7 +28,12 @@ const INSTRUCTIONS_TEXT = `Привет! Для получения корпор Можно прислать и сопроводительный текст — я сам найду ссылку.`; -const PLATFORM_PROMPT = `Ссылка принята. Теперь скажи, какой у тебя телефон:\n\n• Android\n• iPhone (iOS)\n\nЯ пришлю подходящий файл приложения olcbox.`; +const PLATFORM_PROMPT = `Ссылка принята. Выбери платформу, и я пришлю подходящий файл приложения olcbox.`; + +const PLATFORM_BUTTONS: BotButton[] = [ + { text: "Android", callbackData: "platform:android", style: "primary" }, + { text: "iPhone (iOS)", callbackData: "platform:ios", style: "secondary" }, +]; const IOS_WARNING = `⚠️ Файл для iOS не подписан. Как его установить на iPhone — честно, не знаю, придётся разбираться самому (AltStore, enterprise-сертификат и т.п.).`; @@ -39,11 +44,14 @@ function formatUserName(user: User): string { type Attachment = { url: string; name: string; size: number; mimeType?: string }; +type BotButton = { text: string; callbackData: string; style?: 'primary' | 'secondary' | 'destructive' | 'outline'; row?: number }; + export async function sendVpnBotMessage( conversationId: number, text: string, organizationId: number, attachments?: Attachment[], + buttons?: BotButton[], ): Promise { const [msg] = await db .insert(conversationMessages) @@ -55,6 +63,7 @@ export async function sendVpnBotMessage( replyToId: null, mentionedUserIds: null, attachments: attachments?.length ? attachments : null, + botButtons: buttons?.length ? buttons : null, }) .returning(); @@ -166,47 +175,16 @@ export async function handleVpnBotDirectMessage({ return; } - // Continue platform-selection dialog + // Continue platform-selection dialog (text fallback) const pending = pendingPlatform.get(conversationId); if (pending) { const platform = detectPlatform(text); if (!platform) { - await sendVpnBotMessage(conversationId, PLATFORM_PROMPT, organizationId); + await sendVpnBotMessage(conversationId, PLATFORM_PROMPT, organizationId, undefined, PLATFORM_BUTTONS); return; } - pendingPlatform.delete(conversationId); - - try { - const userName = formatUserName(user); - const result = await getOrCreateVpnSubscription(organizationId, user.id, userName, pending.roomUrl); - - if (platform === "android") { - await sendVpnBotMessage( - conversationId, - "Установи приложение olcbox из прикреплённого APK, затем добавь подписку по ссылке ниже.", - organizationId, - [getFileAttachment(ANDROID_FILE)], - ); - } else { - await sendVpnBotMessage( - conversationId, - `${IOS_WARNING}\n\nПосле установки добавь подписку по ссылке ниже.`, - organizationId, - [getFileAttachment(IOS_FILE)], - ); - } - - const subscriptionText = result.isNew - ? `Подписка готова:\n${result.subscriptionUrl}\n\nTelegram и WhatsApp пойдут через латвийский сервер, российские сайты — напрямую.\nСсылка привязана к первому устройству, на котором её активируют.` - : `У вас уже есть подписка для этой комнаты:\n${result.subscriptionUrl}`; - - await sendVpnBotMessage(conversationId, subscriptionText, organizationId); - } catch (err: any) { - console.error("[VPN Bot] handle message error:", err); - const userMessage = err instanceof Error ? err.message : "Произошла ошибка. Попробуйте позже."; - await sendVpnBotMessage(conversationId, userMessage, organizationId); - } + await processPlatformSelection({ conversationId, user, organizationId, roomUrl: pending.roomUrl, platform }); return; } @@ -216,7 +194,90 @@ export async function handleVpnBotDirectMessage({ return; } - // Start platform-selection dialog + // Start platform-selection dialog with buttons pendingPlatform.set(conversationId, { roomUrl, createdAt: Date.now() }); - await sendVpnBotMessage(conversationId, PLATFORM_PROMPT, organizationId); + await sendVpnBotMessage(conversationId, PLATFORM_PROMPT, organizationId, undefined, PLATFORM_BUTTONS); +} + +async function processPlatformSelection({ + conversationId, + user, + organizationId, + roomUrl, + platform, +}: { + conversationId: number; + user: User; + organizationId: number; + roomUrl: string; + platform: "android" | "ios"; +}): Promise { + try { + const userName = formatUserName(user); + const result = await getOrCreateVpnSubscription(organizationId, user.id, userName, roomUrl); + + if (platform === "android") { + await sendVpnBotMessage( + conversationId, + "Установи приложение olcbox из прикреплённого APK, затем добавь подписку по ссылке ниже.", + organizationId, + [getFileAttachment(ANDROID_FILE)], + ); + } else { + await sendVpnBotMessage( + conversationId, + `${IOS_WARNING}\n\nПосле установки добавь подписку по ссылке ниже.`, + organizationId, + [getFileAttachment(IOS_FILE)], + ); + } + + const subscriptionText = result.isNew + ? `Подписка готова:\n${result.subscriptionUrl}\n\nTelegram и WhatsApp пойдут через латвийский сервер, российские сайты — напрямую.\nСсылка привязана к первому устройству, на котором её активируют.` + : `У вас уже есть подписка для этой комнаты:\n${result.subscriptionUrl}`; + + await sendVpnBotMessage(conversationId, subscriptionText, organizationId); + } catch (err: any) { + console.error("[VPN Bot] handle platform selection error:", err); + const userMessage = err instanceof Error ? err.message : "Произошла ошибка. Попробуйте позже."; + await sendVpnBotMessage(conversationId, userMessage, organizationId); + } +} + +export async function handleVpnBotCallback({ + bot, + conversationId, + messageId, + callbackData, + user, + organizationId, +}: { + bot: Bot; + conversationId: number; + messageId: number; + callbackData: string; + user: User; + organizationId: number; +}): Promise { + if (callbackData === "platform:android" || callbackData === "platform:ios") { + const pending = pendingPlatform.get(conversationId); + if (!pending) { + await sendVpnBotMessage( + conversationId, + "Сессия выбора платформы истекла. Пришли ссылку на встречу заново.", + organizationId, + ); + return; + } + const platform = callbackData === "platform:android" ? "android" : "ios"; + pendingPlatform.delete(conversationId); + await processPlatformSelection({ conversationId, user, organizationId, roomUrl: pending.roomUrl, platform }); + return; + } + + await sendVpnBotMessage( + conversationId, + `Неизвестное действие: ${callbackData}`, + organizationId, + ); } diff --git a/shared/schema.ts b/shared/schema.ts index 3d648ee..ff92638 100644 --- a/shared/schema.ts +++ b/shared/schema.ts @@ -2814,6 +2814,7 @@ export const conversationMessages = pgTable("conversation_messages", { replyToId: integer("reply_to_id"), mentionedUserIds: jsonb("mentioned_user_ids").$type(), attachments: jsonb("attachments").$type>(), + botButtons: jsonb("bot_buttons").$type>(), createdAt: timestamp("created_at").defaultNow(), updatedAt: timestamp("updated_at").defaultNow(), isDeleted: boolean("is_deleted").default(false),