From af4850bd6d2f06d1f4d3768d7eca9b9b4457bb52 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=98=D0=BB=D1=8C=D1=8F=D1=81=20=D0=A1=D1=83=D0=BB=D1=82?= =?UTF-8?q?=D0=B0=D0=BD=D0=BE=D0=B2?= Date: Tue, 29 Sep 2026 17:40:47 +0300 Subject: [PATCH] =?UTF-8?q?fix(chat):=20=D0=B8=D0=B4=D0=B5=D0=BC=D0=BF?= =?UTF-8?q?=D0=BE=D1=82=D0=B5=D0=BD=D1=82=D0=BD=D0=BE=D1=81=D1=82=D1=8C=20?= =?UTF-8?q?=D1=81=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D0=B9=20=D0=BF?= =?UTF-8?q?=D0=BE=20clientMessageId=20=E2=80=94=20=D0=B7=D0=B0=D1=89=D0=B8?= =?UTF-8?q?=D1=82=D0=B0=20=D0=BE=D1=82=20=D0=B4=D1=83=D0=B1=D0=BB=D0=B5?= =?UTF-8?q?=D0=B9=20=D0=BF=D1=80=D0=B8=20=D0=BF=D0=BE=D0=B2=D1=82=D0=BE?= =?UTF-8?q?=D1=80=D0=BD=D0=BE=D0=B9=20=D0=BE=D1=82=D0=BF=D1=80=D0=B0=D0=B2?= =?UTF-8?q?=D0=BA=D0=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Первопричина дублей (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, проброс прочих ошибок). --- client/src/components/TaskChat.tsx | 7 +- client/src/hooks/useOfflineSync.ts | 6 + migrations/0086_task_message_idempotency.sql | 11 ++ server/routes/chat.messages.routes.ts | 1 + server/services/task-message.service.ts | 39 +++++- server/storage.ts | 1 + server/storage/tasks.storage.ts | 13 ++ shared/schema.ts | 5 + tests/offline-queue-multitab-race.test.tsx | 135 +++++++++++++++++++ tests/task-message-idempotency.test.ts | 70 ++++++++++ 10 files changed, 286 insertions(+), 2 deletions(-) create mode 100644 migrations/0086_task_message_idempotency.sql create mode 100644 tests/offline-queue-multitab-race.test.tsx create mode 100644 tests/task-message-idempotency.test.ts diff --git a/client/src/components/TaskChat.tsx b/client/src/components/TaskChat.tsx index 94b8dd5..0593862 100644 --- a/client/src/components/TaskChat.tsx +++ b/client/src/components/TaskChat.tsx @@ -142,6 +142,7 @@ interface SendMessagePayload { mentionedBotIds?: number[]; replyToMessageId?: number; attachments?: ChatAttachment[]; + clientMessageId?: string; } // Сообщение задачи, дополненное полем replyTo — контракт общего MessageBubble @@ -1405,13 +1406,17 @@ const TaskChat = ({ taskId }: TaskChatProps) => { return; } - // Отправляем сообщение без файлов как обычно + // Отправляем сообщение без файлов как обычно. + // clientMessageId — ключ идемпотентности: если POST сорвётся, сообщение + // попадёт в офлайн-очередь с ЭТИМ же телом, и все реплеи (в т.ч. из разных + // вкладок) пойдут на сервер с одним ключом → unique-дедуп на сервере. sendMessageMutation.mutate({ message: text, mentionedUserIds: finalMentionedIds, mentionedBotIds, replyToMessageId: replyingTo?.id, attachments: undefined, + clientMessageId: crypto.randomUUID(), }); }; diff --git a/client/src/hooks/useOfflineSync.ts b/client/src/hooks/useOfflineSync.ts index 153065d..5ea174d 100644 --- a/client/src/hooks/useOfflineSync.ts +++ b/client/src/hooks/useOfflineSync.ts @@ -256,6 +256,12 @@ async function trySendAttachmentMessage( attachments, }; + // Идемпотентность: id записи очереди стабилен между реплеями (в т.ч. из разных + // вкладок) — сервер дедупит повторы по client_message_id (миграция 0086). + if (item.type === 'task') { + payload.clientMessageId = (item as { clientMessageId?: string }).clientMessageId ?? item.id; + } + if (item.type === 'messenger' && (item.replyToId || item.replyToMessageId)) { payload.replyToId = item.replyToId || item.replyToMessageId; } diff --git a/migrations/0086_task_message_idempotency.sql b/migrations/0086_task_message_idempotency.sql new file mode 100644 index 0000000..fc630d7 --- /dev/null +++ b/migrations/0086_task_message_idempotency.sql @@ -0,0 +1,11 @@ +-- Идемпотентность сообщений чата задач: клиентский ключ (uuid). +-- Защита от дублей при повторной отправке одного сообщения: +-- multi-tab replay офлайн-очереди (каждая вкладка реплеит независимо), +-- ретраи после сетевых сбоев, двойные клики. +-- Partial unique index: NULL-значения (старые сообщения, боты, автоматизации) не конфликтуют. + +ALTER TABLE task_messages ADD COLUMN IF NOT EXISTS client_message_id varchar(64); + +CREATE UNIQUE INDEX IF NOT EXISTS task_messages_client_message_id_unique + ON task_messages (client_message_id) + WHERE client_message_id IS NOT NULL; diff --git a/server/routes/chat.messages.routes.ts b/server/routes/chat.messages.routes.ts index e57e610..5c05484 100644 --- a/server/routes/chat.messages.routes.ts +++ b/server/routes/chat.messages.routes.ts @@ -117,6 +117,7 @@ export function registerChatMessageRoutes(router: Router): void { attachments: req.body.attachments, bodyAuthorId: req.body.authorId, isBotToken: req.isBotToken === true, + clientMessageId: req.body.clientMessageId ?? null, }); res.status(201).json({ diff --git a/server/services/task-message.service.ts b/server/services/task-message.service.ts index 26aae3f..3520ce3 100644 --- a/server/services/task-message.service.ts +++ b/server/services/task-message.service.ts @@ -35,6 +35,14 @@ export interface SendTaskMessageParams { 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-эффектами: @@ -57,6 +65,19 @@ export async function sendTaskMessage(params: SendTaskMessageParams): Promise 0 ? [...new Set(params.visibleToUserIds)] : null, + clientMessageId, }; - const createdMessage = await storage.createTaskMessage(messageData, organizationId); + 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) { diff --git a/server/storage.ts b/server/storage.ts index 8625c03..8b024f0 100644 --- a/server/storage.ts +++ b/server/storage.ts @@ -227,6 +227,7 @@ export interface IStorage { getTaskMessage(messageId: number, organizationId: number): Promise; getTaskMessagesByIds(messageIds: number[], organizationId: number): Promise; getTaskMessagesByTaskIds(taskIds: number[], organizationId: number): Promise; + getTaskMessageIdByClientMessageId(taskId: number, clientMessageId: string): Promise; createTaskMessage(insertMessage: any, organizationId: number): Promise; updateTaskMessage(id: number, organizationId: number, updates: Partial): Promise; deleteTaskMessage(id: number, organizationId: number): Promise; diff --git a/server/storage/tasks.storage.ts b/server/storage/tasks.storage.ts index ee9b296..78f8e6b 100644 --- a/server/storage/tasks.storage.ts +++ b/server/storage/tasks.storage.ts @@ -410,6 +410,19 @@ export class TasksStorage extends TasksCoreStorage { }; } + /** Поиск id сообщения по клиентскому ключу идемпотентности (uuid от клиента) */ + async getTaskMessageIdByClientMessageId(taskId: number, clientMessageId: string): Promise { + const rows = await db + .select({ id: taskMessages.id }) + .from(taskMessages) + .where(and( + eq(taskMessages.taskId, taskId), + eq(taskMessages.clientMessageId, clientMessageId) + )) + .limit(1); + return rows[0]?.id ?? null; + } + /** Батч-загрузка сообщений по списку id одним IN-запросом (сырые строки, tenant-фильтр) */ async getTaskMessagesByIds(messageIds: number[], organizationId: number): Promise { if (messageIds.length === 0) return []; diff --git a/shared/schema.ts b/shared/schema.ts index dd9ee50..bd2fd8a 100644 --- a/shared/schema.ts +++ b/shared/schema.ts @@ -743,6 +743,10 @@ export const taskMessages = pgTable("task_messages", { botButtons: jsonb("bot_buttons"), // кнопки для бот-сообщений [{text, callbackData, style, row}] visibleToUserIds: jsonb("visible_to_user_ids").$type(), // NULL = видят все участники чата; массив = только перечисленные (+ автор) attachments: jsonb("attachments").$type>(), + // Идемпотентность: клиентский ключ сообщения (uuid от клиента). Защищает от + // дублей при повторной отправке (multi-tab replay офлайн-очереди, ретраи). + // Partial unique index — см. миграцию 0086. + clientMessageId: varchar("client_message_id", { length: 64 }), createdAt: timestamp("created_at").defaultNow(), }, (table) => ({ // Композитный foreign key: задача должна принадлежать той же форме @@ -1370,6 +1374,7 @@ export const createTaskMessageSchema = z.object({ messageType: z.enum(['comment', 'system', 'status_change']).default('comment'), mentionedUserIds: z.array(z.number().int().positive()).max(20).optional(), attachments: z.array(chatAttachmentSchema).max(10).optional(), + clientMessageId: z.string().max(64).optional(), }); export const searchUsersSchema = z.object({ diff --git a/tests/offline-queue-multitab-race.test.tsx b/tests/offline-queue-multitab-race.test.tsx new file mode 100644 index 0000000..05b3cea --- /dev/null +++ b/tests/offline-queue-multitab-race.test.tsx @@ -0,0 +1,135 @@ +// @vitest-environment jsdom +/** + * Воспроизведение бага «одно сообщение приходит 3 раза». + * + * Механизм: офлайн-очередь (IndexedDB) общая на устройство, а каждая открытая + * вкладка приложения реплеит её НЕЗАВИСИМО (мьютекс _syncLock в useOfflineSync — + * модульный, т.е. свой у каждой вкладки). Service Worker по background-sync шлёт + * TRIGGER_SYNC ВСЕМ клиентам (sw.js:445-447) → N вкладок одновременно читают + * одну и ту же очередь и реплеят один и тот же POST → на сервере N дублей. + * + * Тест эмулирует две вкладки: два НЕЗАВИСИМЫХ экземпляра модуля useOfflineSync + * (отдельные _syncLock, как в разных браузерных вкладках), но ОБЩУЮ очередь + * (sharedQueue — аналог общей IndexedDB). + */ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { renderHook, act } from '@testing-library/react'; +import React from 'react'; +import { QueryClient, QueryClientProvider } from '@tanstack/react-query'; + +// Общая "IndexedDB" двух вкладок — живёт вне моков, переживает resetModules +const sharedQueue: Array<{ + id: string; + method: string; + url: string; + body: string | null; + token: null; + orgId: string | null; + enqueuedAt: number; +}> = []; + +// Счётчик реальных POST-запросов (реплеев) к серверу +let replayedPosts: Array<{ url: string; body: string | null }> = []; + +vi.mock('@/lib/offlineQueue', () => ({ + getAll: vi.fn(async () => [...sharedQueue]), + remove: vi.fn(async (id: string) => { + const idx = sharedQueue.findIndex((i) => i.id === id); + if (idx >= 0) sharedQueue.splice(idx, 1); + }), + count: vi.fn(async () => sharedQueue.length), + markAsConflict: vi.fn(async () => {}), + clearConflictStatus: vi.fn(async () => {}), + updateBody: vi.fn(async () => {}), +})); + +vi.mock('@/lib/attachmentQueue', () => ({ + getAllAttachments: vi.fn(async () => []), + countAttachments: vi.fn(async () => 0), + getAttachmentById: vi.fn(), + removeAttachment: vi.fn(async () => {}), + updateAttachment: vi.fn(async () => {}), + markAttachmentError: vi.fn(async () => {}), +})); + +vi.mock('@/lib/syncEngine', () => ({ + performDeltaSync: vi.fn(async () => {}), + performInitialSync: vi.fn(async () => {}), + checkServerReachable: vi.fn(async () => true), +})); + +vi.mock('@/lib/tasksCache', () => ({ + saveInbox: vi.fn(async () => {}), + pruneStaleInboxTasks: vi.fn(async () => {}), +})); + +vi.mock('@/lib/sseManager', () => ({ + sseManager: { + setUser: vi.fn(), + on: vi.fn(() => () => {}), + onConnState: vi.fn(() => () => {}), + reconnect: vi.fn(), + }, +})); + +function makeWrapper() { + const client = new QueryClient({ defaultOptions: { queries: { retry: false } } }); + return ({ children }: { children: React.ReactNode }) => + React.createElement(QueryClientProvider, { client }, children); +} + +async function importFreshHook() { + vi.resetModules(); // новый экземпляр модуля = новая «вкладка» со своим _syncLock + return (await import('@/hooks/useOfflineSync')).useOfflineSync; +} + +describe('multi-tab offline queue replay race', () => { + beforeEach(() => { + sharedQueue.length = 0; + replayedPosts = []; + vi.stubGlobal('fetch', vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input); + if ((init?.method ?? 'GET') === 'POST') { + replayedPosts.push({ url, body: init?.body ? String(init.body) : null }); + } + return new Response( + JSON.stringify({ success: true, tasks: [] }), + { status: 200, headers: { 'Content-Type': 'application/json' } }, + ); + })); + }); + + it('две вкладки реплеят одну и ту же запись очереди дважды (дубль сообщения)', async () => { + const useOfflineSyncA = await importFreshHook(); + const useOfflineSyncB = await importFreshHook(); + + // Одна запись в общей очереди (сообщение, которое не ушло при офлайне) + sharedQueue.push({ + id: 'q1', + method: 'POST', + url: '/api/tasks/1988/messages', + body: JSON.stringify({ message: 'тестовое сообщение' }), + token: null, + orgId: '1', + enqueuedAt: Date.now(), + }); + + const tabA = renderHook(() => useOfflineSyncA(1), { wrapper: makeWrapper() }); + const tabB = renderHook(() => useOfflineSyncB(1), { wrapper: makeWrapper() }); + + // Обе «вкладки» одновременно получают TRIGGER_SYNC (background-sync от SW + // рассылается всем клиентам) и реплеят очередь + await act(async () => { + await Promise.all([ + tabA.result.current.sync(), + tabB.result.current.sync(), + ]); + }); + + const messagePosts = replayedPosts.filter((p) => p.url === '/api/tasks/1988/messages'); + console.log(`[race-test] POST /api/tasks/1988/messages выполнен ${messagePosts.length} раз(а) для 1 записи очереди`); + // БАГ: одна запись очереди уходит дважды (по разу на каждую открытую вкладку). + // После фикса (cross-tab lock / атомарное изъятие) ожидается 1. + expect(messagePosts.length).toBe(2); + }); +}); diff --git a/tests/task-message-idempotency.test.ts b/tests/task-message-idempotency.test.ts new file mode 100644 index 0000000..5ba205a --- /dev/null +++ b/tests/task-message-idempotency.test.ts @@ -0,0 +1,70 @@ +/** + * Серверная идемпотентность сообщений (вариант C, миграция 0086): + * повторная отправка с тем же clientMessageId НЕ создаёт дубль — + * sendTaskMessage возвращает уже существующее сообщение без side-эффектов. + * + * Два сценария: + * 1. Запись уже есть (быстрый путь: select до insert). + * 2. Гонка: параллельный insert упал на unique index (23505) — fallback на select. + */ +import { describe, it, expect, vi, beforeEach } from 'vitest'; + +const mockStorage = vi.hoisted(() => ({ + getTaskMessageIdByClientMessageId: vi.fn(), + getTaskMessage: vi.fn(), + createTaskMessage: vi.fn(), +})); + +vi.mock('../server/storage', () => ({ storage: mockStorage })); + +import { sendTaskMessage } from '../server/services/task-message.service'; + +const task = { id: 1988, formId: 14 } as never; +const user = { id: 1, firstName: 'Ильяс', lastName: 'Султанов' } as never; +const existingMessage = { id: 994, taskId: 1988, message: 'тест', clientMessageId: 'uuid-1' }; + +const baseParams = { + task, + user, + organizationId: 1, + message: 'тест', + clientMessageId: 'uuid-1', +}; + +describe('sendTaskMessage — идемпотентность по clientMessageId', () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it('повтор с тем же ключом: возвращает существующее сообщение, insert НЕ вызывается', async () => { + mockStorage.getTaskMessageIdByClientMessageId.mockResolvedValue(994); + mockStorage.getTaskMessage.mockResolvedValue(existingMessage); + + const result = await sendTaskMessage(baseParams); + + expect(result).toBe(existingMessage); + expect(mockStorage.createTaskMessage).not.toHaveBeenCalled(); + expect(mockStorage.getTaskMessageIdByClientMessageId).toHaveBeenCalledWith(1988, 'uuid-1'); + }); + + it('гонка: unique violation (23505) на insert → fallback возвращает существующее', async () => { + mockStorage.getTaskMessageIdByClientMessageId + .mockResolvedValueOnce(null) // select до insert — ещё пусто + .mockResolvedValueOnce(994); // select после unique-конфликта — уже вставлено параллельным запросом + const pgError = Object.assign(new Error('duplicate key value violates unique constraint'), { code: '23505' }); + mockStorage.createTaskMessage.mockRejectedValue(pgError); + mockStorage.getTaskMessage.mockResolvedValue(existingMessage); + + const result = await sendTaskMessage(baseParams); + + expect(result).toBe(existingMessage); + expect(mockStorage.createTaskMessage).toHaveBeenCalledTimes(1); + }); + + it('другие ошибки insert пробрасываются наружу', async () => { + mockStorage.getTaskMessageIdByClientMessageId.mockResolvedValue(null); + mockStorage.createTaskMessage.mockRejectedValue(new Error('connection refused')); + + await expect(sendTaskMessage(baseParams)).rejects.toThrow('connection refused'); + }); +});