fix(chat): идемпотентность сообщений по clientMessageId — защита от дублей при повторной отправке

Первопричина дублей (3 одинаковых сообщения в задаче 1988): живой POST сорвался
→ сообщение попало в общую IndexedDB-очередь → background-sync разослал
TRIGGER_SYNC всем вкладкам → каждая вкладка реплейнула одну и ту же запись
(мьютекс _syncLock модульный, кросс-вкладочной координации нет). Гонка
воспроизведена тестом tests/offline-queue-multitab-race.test.tsx (1 запись →
2 POST для двух вкладок).

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

View File

@@ -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);
});
});

View File

@@ -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');
});
});