import { create } from 'zustand'; import { acknowledgeChannel, addMessageReaction, createNonce, deleteChannelMessage, editChannelMessage, extractMentions, fetchChannelMessages, fetchChannelPins, pinChannelMessage, removeMessageReaction, searchChannelMessages, sendChannelMessage, sendTypingIndicator, unpinChannelMessage, MESSAGE_PAGE_SIZE, MESSAGE_SEARCH_LIMIT, type Message, type MessageReaction, type PendingMessage, } from '@/api/messages'; import type { UploadedFile } from '@/api/files'; import { useSessionStore } from '@/stores/session'; /** * Сообщения комнат (Фаза 2). * * Стор — единственный источник правды по ленте: REST-ответы и события шлюза * складываются в один список на комнату (старые → новые), поэтому рендер не * зависит от порядка, в котором пришли данные. Оптимистичные сообщения живут * в отдельном массиве `pending` и исчезают, как только сервер подтвердил * отправку (или вернул ошибку). */ /** Сколько сообщений храним на комнату, чтобы память не росла бесконечно. */ export const MESSAGE_CACHE_LIMIT = 200; /** Через сколько гаснет «печатает…» без новых событий. */ export const TYPING_TTL_MS = 6_000; /** Как часто клиент шлёт `POST /typing`, пока идёт ввод. */ export const TYPING_THROTTLE_MS = 3_000; export type { PendingMessage }; export interface TypingEntry { user_id: string; at: number; } export type ChannelStatus = 'idle' | 'loading' | 'ready' | 'error'; export type SearchStatus = 'idle' | 'searching' | 'ready' | 'error'; export interface ChannelMessagesState { /** История: старые → новые, не больше `MESSAGE_CACHE_LIMIT`. */ messages: Message[]; /** Оптимистичные сообщения, ещё не подтверждённые сервером. */ pending: PendingMessage[]; status: ChannelStatus; /** Ошибка первичной загрузки либо подгрузки истории. */ error: unknown; loadingMore: boolean; hasMore: boolean; /** * Курсор подгрузки вверх: id самого старого сообщения последней страницы. * Хранится отдельно от `messages`, потому что кэш ограничен и «самое старое» * в нём перестаёт двигаться, когда окно упирается в лимит. */ pagingCursor: string | null; pins: Message[]; pinsLoaded: boolean; pinsError: unknown; searchQuery: string; searchStatus: SearchStatus; searchResults: Message[]; searchError: unknown; /** id последнего прочитанного сообщения (read state). */ lastReadMessageId: string | null; /** Пришла ли отметка прочтения (из снапшота READY или событием). */ readStateKnown: boolean; mentionCount: number; typing: TypingEntry[]; } export interface SendMessageInput { channelId: string; content: string; replyToId?: string; /** Уже загруженные вложения (`POST /channels/{id}/files`). */ files?: UploadedFile[]; nonce?: string; } interface MessagesState { channels: Record; /** Применён ли снапшот READY (после него отсутствие read state — это факт). */ snapshotApplied: boolean; /** Первичная загрузка (идемпотентно). */ ensureChannel: (channelId: string) => Promise; /** Принудительная перезагрузка последней страницы истории. */ reloadChannel: (channelId: string) => Promise; /** Подгрузка вверх: `before` = id самого старого сообщения. */ loadMore: (channelId: string) => Promise; loadPins: (channelId: string) => Promise; search: (channelId: string, query: string) => Promise; clearSearch: (channelId: string) => void; send: (input: SendMessageInput) => Promise; edit: (channelId: string, messageId: string, content: string) => Promise; remove: (channelId: string, messageId: string) => Promise; toggleReaction: (channelId: string, messageId: string, emoji: string) => Promise; setPinned: (channelId: string, messageId: string, pinned: boolean) => Promise; notifyTyping: (channelId: string) => void; /** Отметка прочтения: локально сразу, на сервере — фоном. */ acknowledge: (channelId: string, lastMessageId?: string) => Promise; applyMessageCreate: (message: Message) => void; applyMessageUpdate: (message: Message) => void; applyMessageDelete: (channelId: string, messageId: string) => void; applyReaction: ( channelId: string, messageId: string, emoji: string, userId: string, added: boolean, ) => void; applyPinsUpdate: (channelId: string, messageId: string, pinned: boolean) => void; applyTypingStart: (channelId: string, userId: string) => void; applyReadState: (event: { channel_id: string; last_message_id?: string; mention_count: number; }) => void; /** Комната удалена — выбрасываем её состояние. */ dropChannel: (channelId: string) => void; /** Снапшот READY применён: отметки прочтения для комнат уже известны. */ markSnapshotApplied: () => void; pruneTyping: (now?: number) => void; reset: () => void; } function emptyChannel(): ChannelMessagesState { return { messages: [], pending: [], status: 'idle', error: null, loadingMore: false, hasMore: false, pagingCursor: null, pins: [], pinsLoaded: false, pinsError: null, searchQuery: '', searchStatus: 'idle', searchResults: [], searchError: null, lastReadMessageId: null, readStateKnown: false, mentionCount: 0, typing: [], }; } /** * Заглушка для незагруженных комнат. Объект один на всех: селектор не должен * возвращать новую ссылку на каждом вызове, иначе React перерисовывается зря. */ const EMPTY_CHANNEL: ChannelMessagesState = emptyChannel(); export function channelState(state: MessagesState, channelId: string): ChannelMessagesState { return state.channels[channelId] ?? EMPTY_CHANNEL; } /** Сравнение идентификаторов: снежинки — возрастающие числа в виде строк. */ function compareIds(left: string, right: string): number { try { const a = BigInt(left); const b = BigInt(right); return a === b ? 0 : a < b ? -1 : 1; } catch { // Оптимистичные id (`pending:…`) не числа — сравниваем как строки. return left === right ? 0 : left < right ? -1 : 1; } } /** Порядок ленты: по времени создания, при равенстве — по id. */ export function compareMessages(left: Message, right: Message): number { const a = Date.parse(left.created_at); const b = Date.parse(right.created_at); if (Number.isFinite(a) && Number.isFinite(b) && a !== b) { return a < b ? -1 : 1; } return compareIds(left.id, right.id); } function sortMessages(messages: Message[]): Message[] { return [...messages].sort(compareMessages); } function capMessages(messages: Message[]): Message[] { return messages.length <= MESSAGE_CACHE_LIMIT ? messages : messages.slice(messages.length - MESSAGE_CACHE_LIMIT); } function upsertMessage(messages: Message[], message: Message): Message[] { const index = messages.findIndex((item) => item.id === message.id); if (index === -1) { return capMessages(sortMessages([...messages, message])); } const next = [...messages]; next[index] = message; return next; } function removeMessage(messages: Message[], messageId: string): Message[] { return messages.filter((item) => item.id !== messageId); } /** * Слияние страницы истории с текущим списком. * * `replace` — первичная загрузка: выборка заменяет историю, но сообщения, * пришедшие событиями шлюза во время запроса (они новее выборки), остаются. * `prepend` — подгрузка вверх: старые добавляются перед текущими. */ function mergeHistory( current: ChannelMessagesState, page: Message[], mode: 'replace' | 'prepend', ): Partial { const pageIds = new Set(page.map((message) => message.id)); const kept = current.messages.filter((message) => !pageIds.has(message.id)); let merged: Message[]; if (mode === 'replace') { const newest = page.reduce( (acc, message) => (acc === null || compareMessages(message, acc) > 0 ? message : acc), null, ); // Пустая выборка — сервер действительно не отдал сообщений (всё удалено). // Совсем свежие записи всё же оставляем: они могли прийти событием шлюза // ровно в момент запроса. const fresh = newest === null ? kept.filter((message) => { const created = Date.parse(message.created_at); return Number.isFinite(created) && Date.now() - created < 60_000; }) : kept.filter((message) => compareMessages(message, newest) > 0); merged = [...page, ...fresh]; } else { merged = [...page, ...kept]; } const sorted = sortMessages(merged); // Кэш ограничен: при подгрузке вверх вытесняем самые новые сообщения, // иначе только что загруженная страница была бы сразу отброшена и // «Загрузить ещё» перестало бы двигать историю. const capped = sorted.length <= MESSAGE_CACHE_LIMIT ? sorted : mode === 'prepend' ? sorted.slice(0, MESSAGE_CACHE_LIMIT) : capMessages(sorted); // Страница приходит от новых к старым, поэтому курсор — самый старый её элемент // (порядок не полагаем на удачу: берём минимум по компаратору). const oldest = mode === 'prepend' && page.length > 0 ? page.reduce((acc, message) => (compareMessages(message, acc) < 0 ? message : acc)) : undefined; return { messages: capped, status: 'ready', error: null, loadingMore: false, hasMore: page.length >= MESSAGE_PAGE_SIZE, ...(oldest === undefined ? {} : { pagingCursor: oldest.id }), }; } /** * Локальное изменение реакций: `actorId` — кто поставил/снял реакцию, * `meId` — текущий пользователь (от него зависит подсветка «моя реакция»). */ function applyReactionLocally( reactions: MessageReaction[], emoji: string, actorId: string | null, meId: string | null, added: boolean, ): MessageReaction[] { const index = reactions.findIndex((reaction) => reaction.emoji === emoji); const current = index === -1 ? null : (reactions[index] ?? null); const isMine = actorId !== null && actorId === meId; const alreadyListed = actorId !== null && current?.user_ids !== undefined && current.user_ids.includes(actorId); if (added) { if (current === null) { const created: MessageReaction = { emoji, count: 1, me: isMine }; if (actorId !== null) { created.user_ids = [actorId]; } return [...reactions, created]; } // Дубликат события (свой же PUT и эхо шлюза): счётчик не растёт. // Когда `user_ids` сервер не прислал, ориентируемся на локальный флаг `me`. if (alreadyListed || (isMine && current.me && current.user_ids === undefined)) { return reactions; } const next: MessageReaction = { ...current, count: current.count + 1, me: current.me || isMine, }; if (actorId !== null && current.user_ids !== undefined) { next.user_ids = [...current.user_ids, actorId]; } const result = [...reactions]; result[index] = next; return result; } if (current === null) { return reactions; } if (actorId !== null && current.user_ids !== undefined && !alreadyListed) { return reactions; } // Своё эхо на уже снятую реакцию (см. выше про отсутствие `user_ids`). if (current.user_ids === undefined && isMine && !current.me) { return reactions; } const count = Math.max(0, current.count - 1); if (count === 0) { return reactions.filter((reaction) => reaction.emoji !== emoji); } const next: MessageReaction = { ...current, count, me: isMine ? false : current.me, }; if (actorId !== null && current.user_ids !== undefined) { next.user_ids = current.user_ids.filter((id) => id !== actorId); } const result = [...reactions]; result[index] = next; return result; } function toAttachment(file: UploadedFile) { return { file_id: file.file_id, filename: file.filename, ...(file.content_type === '' ? {} : { content_type: file.content_type }), ...(file.size_bytes === undefined ? {} : { size_bytes: file.size_bytes }), }; } /** Последняя отметка отправки `POST /typing` по комнате (троттлинг 3 с). */ const typingSentAt = new Map(); export const useMessagesStore = create((set, get) => { const patch = ( channelId: string, updater: (current: ChannelMessagesState) => Partial, ): void => { const current = get().channels[channelId] ?? emptyChannel(); set({ channels: { ...get().channels, [channelId]: { ...current, ...updater(current) } } }); }; return { channels: {}, snapshotApplied: false, ensureChannel: async (channelId) => { const current = get().channels[channelId] ?? EMPTY_CHANNEL; if (current.status === 'loading' || current.status === 'ready') { return; } await get().reloadChannel(channelId); }, reloadChannel: async (channelId) => { patch(channelId, () => ({ status: 'loading', error: null, loadingMore: false })); try { const page = await fetchChannelMessages(channelId, { limit: MESSAGE_PAGE_SIZE }); patch(channelId, (current) => mergeHistory(current, page, 'replace')); } catch (error) { patch(channelId, () => ({ status: 'error', error, loadingMore: false })); } }, loadMore: async (channelId) => { const current = get().channels[channelId] ?? EMPTY_CHANNEL; // Курсор берём из состояния: массив сообщений ограничен 200 записями. const cursor = current.pagingCursor ?? current.messages[0]?.id; if (current.loadingMore || !current.hasMore || cursor === undefined) { return; } patch(channelId, () => ({ loadingMore: true, error: null })); try { const page = await fetchChannelMessages(channelId, { before: cursor, limit: MESSAGE_PAGE_SIZE, }); patch(channelId, (state) => mergeHistory(state, page, 'prepend')); } catch (error) { patch(channelId, () => ({ loadingMore: false, error })); } }, loadPins: async (channelId) => { patch(channelId, () => ({ pinsError: null })); try { const pins = await fetchChannelPins(channelId); patch(channelId, () => ({ pins, pinsLoaded: true })); } catch (error) { patch(channelId, () => ({ pinsError: error, pinsLoaded: true })); } }, search: async (channelId, query) => { const trimmed = query.trim(); if (trimmed === '') { patch(channelId, () => ({ searchQuery: '', searchStatus: 'idle', searchResults: [], searchError: null, })); return; } patch(channelId, () => ({ searchQuery: trimmed, searchStatus: 'searching', searchError: null, })); try { const results = await searchChannelMessages(channelId, trimmed, MESSAGE_SEARCH_LIMIT); // Пока запрос летел, пользователь мог изменить строку поиска. if ((get().channels[channelId] ?? EMPTY_CHANNEL).searchQuery !== trimmed) { return; } patch(channelId, () => ({ searchStatus: 'ready', searchResults: results })); } catch (error) { if ((get().channels[channelId] ?? EMPTY_CHANNEL).searchQuery !== trimmed) { return; } patch(channelId, () => ({ searchStatus: 'error', searchError: error })); } }, clearSearch: (channelId) => { patch(channelId, () => ({ searchQuery: '', searchStatus: 'idle', searchResults: [], searchError: null, })); }, send: async (input) => { const { channelId } = input; const content = input.content; const nonce = input.nonce ?? createNonce(); const files = input.files ?? []; const author = useSessionStore.getState().user; const optimistic: PendingMessage = { id: `pending:${nonce}`, channel_id: channelId, content, ...(author === null ? {} : { author_id: author.id }), ...(input.replyToId === undefined ? {} : { reply_to_id: input.replyToId }), type: 'default', pinned: false, attachments: files.map(toAttachment), mentions: extractMentions(content), reactions: [], created_at: new Date().toISOString(), pending: true, nonce, }; patch(channelId, (current) => ({ pending: [...current.pending, optimistic].slice(-MESSAGE_CACHE_LIMIT), error: null, })); try { const message = await sendChannelMessage(channelId, { content, ...(input.replyToId === undefined ? {} : { reply_to_id: input.replyToId }), ...(files.length === 0 ? {} : { attachment_ids: files.map((file) => file.file_id) }), nonce, }); patch(channelId, (current) => ({ // Событие шлюза могло прийти раньше ответа REST — тогда id уже в ленте. pending: current.pending.filter((item) => item.nonce !== nonce), messages: upsertMessage(current.messages, message), })); return message; } catch (error) { // Текст не теряем: композер вернёт его пользователю сам. patch(channelId, (current) => ({ pending: current.pending.filter((item) => item.nonce !== nonce), })); throw error; } }, edit: async (channelId, messageId, content) => { const before = (get().channels[channelId] ?? EMPTY_CHANNEL).messages.find( (message) => message.id === messageId, ); patch(channelId, (current) => ({ messages: current.messages.map((message) => message.id === messageId ? { ...message, content } : message, ), })); try { const updated = await editChannelMessage(channelId, messageId, content); patch(channelId, (current) => ({ messages: upsertMessage(current.messages, updated), pins: current.pins.map((item) => (item.id === messageId ? updated : item)), })); } catch (error) { if (before !== undefined) { patch(channelId, (current) => ({ messages: upsertMessage(current.messages, before), })); } throw error; } }, remove: async (channelId, messageId) => { const before = (get().channels[channelId] ?? EMPTY_CHANNEL).messages.find( (message) => message.id === messageId, ); patch(channelId, (current) => ({ messages: removeMessage(current.messages, messageId), pins: current.pins.filter((item) => item.id !== messageId), })); try { await deleteChannelMessage(channelId, messageId); } catch (error) { if (before !== undefined) { patch(channelId, (current) => ({ messages: upsertMessage(current.messages, before), })); } throw error; } }, toggleReaction: async (channelId, messageId, emoji) => { const current = (get().channels[channelId] ?? EMPTY_CHANNEL).messages.find( (message) => message.id === messageId, ); if (current === undefined) { return; } const mine = current.reactions.find((reaction) => reaction.emoji === emoji)?.me === true; const me = useSessionStore.getState().user?.id ?? null; patch(channelId, (state) => ({ messages: state.messages.map((message) => message.id === messageId ? { ...message, reactions: applyReactionLocally(message.reactions, emoji, me, me, !mine), } : message, ), })); try { if (mine) { await removeMessageReaction(channelId, messageId, emoji); } else { await addMessageReaction(channelId, messageId, emoji); } } catch (error) { patch(channelId, (state) => ({ messages: state.messages.map((message) => message.id === messageId ? { ...message, reactions: current.reactions } : message, ), })); throw error; } }, setPinned: async (channelId, messageId, pinned) => { const current = (get().channels[channelId] ?? EMPTY_CHANNEL).messages.find( (message) => message.id === messageId, ); const previousPins = (get().channels[channelId] ?? EMPTY_CHANNEL).pins; patch(channelId, (state) => ({ messages: state.messages.map((message) => message.id === messageId ? { ...message, pinned } : message, ), pins: current === undefined ? state.pins : pinned ? upsertMessage(state.pins, { ...current, pinned: true }) : state.pins.filter((item) => item.id !== messageId), })); try { if (pinned) { await pinChannelMessage(channelId, messageId); } else { await unpinChannelMessage(channelId, messageId); } } catch (error) { patch(channelId, (state) => ({ messages: state.messages.map((message) => message.id === messageId && current !== undefined ? { ...message, pinned: current.pinned } : message, ), pins: previousPins, })); throw error; } }, notifyTyping: (channelId) => { const now = Date.now(); const last = typingSentAt.get(channelId) ?? 0; if (now - last < TYPING_THROTTLE_MS) { return; } typingSentAt.set(channelId, now); // Ошибку «печатает…» показывать не нужно: это необязательный сигнал. void sendTypingIndicator(channelId).catch(() => undefined); }, acknowledge: async (channelId, lastMessageId) => { patch(channelId, () => ({ ...(lastMessageId === undefined ? {} : { lastReadMessageId: lastMessageId }), mentionCount: 0, })); try { await acknowledgeChannel(channelId, lastMessageId); } catch { // Отметка прочтения не критична: следующая попытка придёт с новым сообщением. } }, applyMessageCreate: (message) => { if (message.channel_id === '') { return; } const me = useSessionStore.getState().user?.id ?? null; patch(message.channel_id, (current) => { // Своё сообщение могло прийти событием раньше ответа REST: снимаем // соответствующее оптимистичное, чтобы лента не показала дубль. const echo = me !== null && message.author_id === me ? current.pending.find( (item) => item.content === message.content && item.attachments.length === message.attachments.length, ) : undefined; return { messages: upsertMessage(current.messages, message), pending: echo === undefined ? current.pending : current.pending.filter((item) => item.nonce !== echo.nonce), }; }); }, applyMessageUpdate: (message) => { patch(message.channel_id, (current) => { // Payload события собран для автора правки, поэтому чужой флаг `me` // не должен затирать локальный: иначе «моя реакция» переедет на другое // устройство вместе с событием. const previous = current.messages.find((item) => item.id === message.id); const merged = previous === undefined || previous.reactions.length === 0 ? message : { ...message, reactions: message.reactions.map((reaction) => { const local = previous.reactions.find((item) => item.emoji === reaction.emoji); return local === undefined ? reaction : { ...reaction, me: local.me }; }), }; return { messages: upsertMessage(current.messages, merged), pins: current.pins.map((item) => (item.id === message.id ? merged : item)), }; }); }, applyMessageDelete: (channelId, messageId) => { patch(channelId, (current) => ({ messages: removeMessage(current.messages, messageId), pins: current.pins.filter((item) => item.id !== messageId), pending: current.pending.filter((item) => item.id !== messageId), })); }, applyReaction: (channelId, messageId, emoji, userId, added) => { const me = useSessionStore.getState().user?.id ?? null; patch(channelId, (current) => ({ messages: current.messages.map((message) => message.id === messageId ? { ...message, reactions: applyReactionLocally(message.reactions, emoji, userId, me, added), } : message, ), })); }, applyPinsUpdate: (channelId, messageId, pinned) => { patch(channelId, (current) => { const target = current.messages.find((message) => message.id === messageId); return { messages: current.messages.map((message) => message.id === messageId ? { ...message, pinned } : message, ), pins: pinned ? target === undefined ? // Сообщение вне загруженного окна: добавить его в список нечем, // но и удалять уже закреплённое нельзя. current.pins : upsertMessage(current.pins, { ...target, pinned: true }) : current.pins.filter((item) => item.id !== messageId), }; }); }, applyTypingStart: (channelId, userId) => { const me = useSessionStore.getState().user?.id ?? null; if (userId === me) { return; } const now = Date.now(); patch(channelId, (current) => ({ typing: [ ...current.typing.filter((entry) => entry.user_id !== userId), { user_id: userId, at: now, }, ], })); }, applyReadState: (event) => { patch(event.channel_id, () => ({ ...(event.last_message_id === undefined ? {} : { lastReadMessageId: event.last_message_id }), readStateKnown: true, mentionCount: event.mention_count, })); }, dropChannel: (channelId) => { const channels = { ...get().channels }; delete channels[channelId]; set({ channels }); }, markSnapshotApplied: () => { set({ snapshotApplied: true }); }, pruneTyping: (now = Date.now()) => { const channels = get().channels; let changed = false; const next: Record = {}; for (const [channelId, channel] of Object.entries(channels)) { if (channel.typing.length === 0) { next[channelId] = channel; continue; } const fresh = channel.typing.filter((entry) => now - entry.at < TYPING_TTL_MS); if (fresh.length === channel.typing.length) { next[channelId] = channel; continue; } changed = true; next[channelId] = { ...channel, typing: fresh }; } if (changed) { set({ channels: next }); } }, reset: () => { typingSentAt.clear(); set({ channels: {}, snapshotApplied: false }); }, }; }); /** Сколько сообщений комнаты ещё не прочитано (по отметке read state). */ export function unreadCount(channel: ChannelMessagesState): number { const marker = channel.lastReadMessageId; if (marker === null) { return 0; } const index = channel.messages.findIndex((message) => message.id === marker); if (index >= 0) { return channel.messages.length - index - 1; } // Отметка не попала в загруженную историю — считаем по порядку снежинок. const firstUnread = channel.messages.findIndex((message) => compareIds(message.id, marker) > 0); return firstUnread < 0 ? 0 : channel.messages.length - firstUnread; } /** Селектор записи комнаты (стабильная заглушка для незагруженных). */ export function selectChannelState(channelId: string) { return (state: MessagesState): ChannelMessagesState => channelState(state, channelId); }