From 5884b0bdba0f48dc5c7dc43fcacea6032588d6a0 Mon Sep 17 00:00:00 2001 From: Kushida Date: Sun, 26 Jul 2026 20:38:32 +0300 Subject: [PATCH] fix(server): validate synced message ownership (#2054) Co-authored-by: RainbowBird Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> --- apps/server/src/services/domain/chats.test.ts | 238 +++++++++++++++++- apps/server/src/services/domain/chats.ts | 57 ++++- 2 files changed, 278 insertions(+), 17 deletions(-) diff --git a/apps/server/src/services/domain/chats.test.ts b/apps/server/src/services/domain/chats.test.ts index bb770043c..68727f87c 100644 --- a/apps/server/src/services/domain/chats.test.ts +++ b/apps/server/src/services/domain/chats.test.ts @@ -1,17 +1,20 @@ -import { describe, expect, it } from 'vitest' +import type { Database } from '../../libs/db' -import { clampLimit, resolveSenderId } from './chats' +import { eq } from 'drizzle-orm' +import { beforeEach, describe, expect, it } from 'vitest' + +import { mockDB } from '../../libs/mock-db' +import { clampLimit, createChatService, resolveSenderId } from './chats' + +import * as schema from '../../schemas' describe('resolveSenderId', () => { it('returns userId for user role', () => { - expect(resolveSenderId('user', 'user-123', 'char-456')).toBe('user-123') + expect(resolveSenderId('user', 'user-123')).toBe('user-123') }) - it('returns characterId for non-user role when available', () => { - expect(resolveSenderId('assistant', 'user-123', 'char-456')).toBe('char-456') - }) - it('returns null for non-user role without characterId', () => { - expect(resolveSenderId('assistant', 'user-123')).toBeNull() - expect(resolveSenderId('system', 'user-123', null)).toBeNull() + it('returns userId for assistant role', () => { + expect(resolveSenderId('assistant', 'user-123')).toBe('user-123') + expect(resolveSenderId('system', 'user-123')).toBeNull() }) }) @@ -33,3 +36,220 @@ describe('clampLimit', () => { expect(clampLimit(1000)).toBe(500) }) }) + +describe('pushMessages', () => { + let db: Database + + beforeEach(async () => { + db = await mockDB(schema) + }) + + it('rejects a member attempt to update another member’s message', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values([ + { chatId: 'group', memberType: 'user', userId: 'author' }, + { chatId: 'group', memberType: 'user', userId: 'member' }, + ]) + await db.insert(schema.messages).values({ + id: 'message', + chatId: 'group', + senderId: 'author', + role: 'user', + seq: 1, + content: 'original', + mediaIds: [], + stickerIds: [], + }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'group', [{ id: 'message', role: 'user', content: 'forged' }])) + .rejects + .toMatchObject({ statusCode: 403, errorCode: 'FORBIDDEN', message: 'Forbidden' }) + + const message = await db.query.messages.findFirst({ where: eq(schema.messages.id, 'message') }) + expect(message?.content).toBe('original') + expect(message?.senderId).toBe('author') + expect(message?.seq).toBe(1) + }) + + it('rejects an existing message ID from another chat', async () => { + await db.insert(schema.chats).values([ + { id: 'source', type: 'group' }, + { id: 'target', type: 'group' }, + ]) + await db.insert(schema.chatMembers).values([ + { chatId: 'source', memberType: 'user', userId: 'member' }, + { chatId: 'target', memberType: 'user', userId: 'member' }, + ]) + await db.insert(schema.messages).values({ + id: 'message', + chatId: 'source', + senderId: 'member', + role: 'user', + seq: 1, + content: 'source message', + mediaIds: [], + stickerIds: [], + }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'target', [{ id: 'message', role: 'user', content: 'target message' }])) + .rejects + .toMatchObject({ statusCode: 409, errorCode: 'CONFLICT', message: 'Message already belongs to another chat' }) + + const sourceMessage = await db.query.messages.findFirst({ where: eq(schema.messages.id, 'message') }) + const targetMessages = await db.query.messages.findMany({ where: eq(schema.messages.chatId, 'target') }) + expect(sourceMessage?.content).toBe('source message') + expect(targetMessages).toHaveLength(0) + }) + + it('allows an author to update their own message', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values({ chatId: 'group', memberType: 'user', userId: 'author' }) + await db.insert(schema.messages).values({ + id: 'message', + chatId: 'group', + senderId: 'author', + role: 'user', + seq: 1, + content: 'original', + mediaIds: [], + stickerIds: [], + }) + + const service = createChatService(db) + + await expect(service.pushMessages('author', 'group', [{ id: 'message', role: 'user', content: 'updated' }])) + .resolves + .toMatchObject({ seq: 2, fromSeq: 2, toSeq: 2 }) + + const message = await db.query.messages.findFirst({ where: eq(schema.messages.id, 'message') }) + expect(message?.content).toBe('updated') + expect(message?.senderId).toBe('author') + expect(message?.role).toBe('user') + expect(message?.chatId).toBe('group') + expect(message?.seq).toBe(2) + }) + + it('acknowledges an unchanged legacy assistant retry without mutating it', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values({ chatId: 'group', memberType: 'user', userId: 'member' }) + await db.insert(schema.messages).values({ + id: 'message', + chatId: 'group', + senderId: null, + role: 'assistant', + seq: 1, + content: 'original response', + mediaIds: [], + stickerIds: [], + }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'group', [{ id: 'message', role: 'assistant', content: 'original response' }])) + .resolves + .toMatchObject({ seq: 1, fromSeq: 2, toSeq: 1 }) + + const message = await db.query.messages.findFirst({ where: eq(schema.messages.id, 'message') }) + expect(message?.content).toBe('original response') + expect(message?.senderId).toBeNull() + expect(message?.role).toBe('assistant') + expect(message?.seq).toBe(1) + }) + + it('persists later messages batched with an unchanged legacy assistant retry', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values({ chatId: 'group', memberType: 'user', userId: 'member' }) + await db.insert(schema.messages).values({ + id: 'legacy-assistant', + chatId: 'group', + senderId: null, + role: 'assistant', + seq: 1, + content: 'original response', + mediaIds: [], + stickerIds: [], + }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'group', [ + { id: 'legacy-assistant', role: 'assistant', content: 'original response' }, + { id: 'new-user-message', role: 'user', content: 'next turn' }, + ])) + .resolves + .toMatchObject({ seq: 2, fromSeq: 2, toSeq: 2 }) + + const messages = await db.query.messages.findMany({ + where: eq(schema.messages.chatId, 'group'), + orderBy: schema.messages.seq, + }) + expect(messages).toHaveLength(2) + expect(messages[0]?.id).toBe('legacy-assistant') + expect(messages[0]?.seq).toBe(1) + expect(messages[1]?.id).toBe('new-user-message') + expect(messages[1]?.senderId).toBe('member') + expect(messages[1]?.seq).toBe(2) + }) + + it('accepts an assistant message from local-first sync', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values({ chatId: 'group', memberType: 'user', userId: 'member' }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'group', [{ id: 'message', role: 'assistant', content: 'response' }])) + .resolves + .toMatchObject({ seq: 1, fromSeq: 1, toSeq: 1 }) + + const message = await db.query.messages.findFirst({ where: eq(schema.messages.id, 'message') }) + expect(message?.role).toBe('assistant') + expect(message?.content).toBe('response') + expect(message?.senderId).toBe('member') + }) + + it('rejects updates to unowned assistant messages', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values([ + { chatId: 'group', memberType: 'user', userId: 'author' }, + { chatId: 'group', memberType: 'user', userId: 'member' }, + ]) + await db.insert(schema.messages).values({ + id: 'message', + chatId: 'group', + senderId: null, + role: 'assistant', + seq: 1, + content: 'original response', + mediaIds: [], + stickerIds: [], + }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'group', [{ id: 'message', role: 'assistant', content: 'forged response' }])) + .rejects + .toMatchObject({ statusCode: 403, errorCode: 'FORBIDDEN', message: 'Forbidden' }) + + const message = await db.query.messages.findFirst({ where: eq(schema.messages.id, 'message') }) + expect(message?.content).toBe('original response') + expect(message?.seq).toBe(1) + }) + + it('rejects roles that are not part of cloud chat sync', async () => { + await db.insert(schema.chats).values({ id: 'group', type: 'group' }) + await db.insert(schema.chatMembers).values({ chatId: 'group', memberType: 'user', userId: 'member' }) + + const service = createChatService(db) + + await expect(service.pushMessages('member', 'group', [{ id: 'message', role: 'system', content: 'local prompt' }])) + .rejects + .toMatchObject({ statusCode: 400, errorCode: 'BAD_REQUEST', message: 'Only user and assistant messages can be synchronized' }) + + const messages = await db.query.messages.findMany({ where: eq(schema.messages.chatId, 'group') }) + expect(messages).toHaveLength(0) + }) +}) diff --git a/apps/server/src/services/domain/chats.ts b/apps/server/src/services/domain/chats.ts index b3e7a5303..c28ba79c2 100644 --- a/apps/server/src/services/domain/chats.ts +++ b/apps/server/src/services/domain/chats.ts @@ -7,7 +7,7 @@ import type { ProductEventService } from './product-events' import { useLogger } from '@guiiai/logg' import { and, eq, gt, inArray, isNull, sql } from 'drizzle-orm' -import { createForbiddenError, createNotFoundError } from '../../utils/error' +import { createBadRequestError, createConflictError, createForbiddenError, createNotFoundError } from '../../utils/error' import { nanoid } from '../../utils/id' import * as schema from '../../schemas/chats' @@ -40,10 +40,10 @@ export function clampLimit(limit?: number): number { return Math.min(limit, 500) } -export function resolveSenderId(role: string, userId: string, characterId?: string | null): string | null { - if (role === 'user') +export function resolveSenderId(role: string, userId: string): string | null { + if (role === 'user' || role === 'assistant') return userId - return characterId ?? null + return null } // --------------------------------------------------------------------------- @@ -222,7 +222,10 @@ export function createChatService(db: Database, metrics?: EngagementMetrics | nu // -- Message sync (WS) -------------------------------------------------- - async pushMessages(userId: string, chatId: string, messages: PushMessage[], characterId?: string) { + async pushMessages(userId: string, chatId: string, messages: PushMessage[]) { + if (messages.some(message => message.role !== 'user' && message.role !== 'assistant')) + throw createBadRequestError('Only user and assistant messages can be synchronized') + const result = await db.transaction(async (tx) => { await verifyMembership(tx, chatId, userId) @@ -247,12 +250,50 @@ export function createChatService(db: Database, metrics?: EngagementMetrics | nu // Split into new vs existing messages const messageIds = messages.map(m => m.id) const existingMessages = messageIds.length > 0 - ? await tx.select({ id: schema.messages.id }).from(schema.messages).where(inArray(schema.messages.id, messageIds)) + ? await tx.select({ + id: schema.messages.id, + chatId: schema.messages.chatId, + senderId: schema.messages.senderId, + role: schema.messages.role, + content: schema.messages.content, + }).from(schema.messages).where(inArray(schema.messages.id, messageIds)) : [] + + if (existingMessages.some(message => message.chatId !== chatId)) + throw createConflictError('Message already belongs to another chat') + + const existingMessagesById = new Map(existingMessages.map(message => [message.id, message])) + const unchangedLegacyAssistantIds = new Set() + if (messages.some((message) => { + const existingMessage = existingMessagesById.get(message.id) + if (existingMessage == null) + return false + + if (existingMessage.senderId === resolveSenderId(message.role, userId)) + return false + + // A pre-ownership assistant row cannot be safely attributed to a user. + // An exact retry is nevertheless safe to acknowledge because it does + // not mutate the stored message or its sequence. + if ( + existingMessage.senderId == null + && existingMessage.role === 'assistant' + && message.role === 'assistant' + && existingMessage.content === message.content + ) { + unchangedLegacyAssistantIds.add(message.id) + return false + } + + return true + })) { + throw createForbiddenError() + } + const existingIds = new Set(existingMessages.map(m => m.id)) const newMsgs = messages.filter(m => !existingIds.has(m.id)) - const updateMsgs = messages.filter(m => existingIds.has(m.id)) + const updateMsgs = messages.filter(m => existingIds.has(m.id) && !unchangedLegacyAssistantIds.has(m.id)) let currentSeq = maxSeq @@ -263,7 +304,7 @@ export function createChatService(db: Database, metrics?: EngagementMetrics | nu return { id: m.id, chatId, - senderId: resolveSenderId(m.role, userId, characterId), + senderId: resolveSenderId(m.role, userId), role: m.role, seq: currentSeq, content: m.content,