diff --git a/packages/@n8n/db/src/migrations/common/1766068346315-AddChatMessageIndices.ts b/packages/@n8n/db/src/migrations/common/1766068346315-AddChatMessageIndices.ts new file mode 100644 index 00000000000..a7cfb408744 --- /dev/null +++ b/packages/@n8n/db/src/migrations/common/1766068346315-AddChatMessageIndices.ts @@ -0,0 +1,44 @@ +import type { MigrationContext, ReversibleMigration } from '../migration-types'; + +export class AddChatMessageIndices1766068346315 implements ReversibleMigration { + async up({ schemaBuilder: { addNotNull }, runQuery, escape }: MigrationContext) { + const sessionsTable = escape.tableName('chat_hub_sessions'); + const idColumn = escape.columnName('id'); + const createdAtColumn = escape.columnName('createdAt'); + const ownerIdColumn = escape.columnName('ownerId'); + const lastMessageAtColumn = escape.columnName('lastMessageAt'); + + const messagesTable = escape.tableName('chat_hub_messages'); + const sessionIdColumn = escape.columnName('sessionId'); + + // Backfill lastMessageAt for existing rows to allow adding a NOT NULL constraint + await runQuery( + `UPDATE ${sessionsTable} + SET ${lastMessageAtColumn} = ${createdAtColumn} + WHERE ${lastMessageAtColumn} IS NULL`, + ); + + await addNotNull('chat_hub_sessions', 'lastMessageAt'); + + // Index intended for faster sessionRepository.getManyByUserId queries + await runQuery( + `CREATE INDEX IF NOT EXISTS ${escape.indexName('chat_hub_sessions_owner_lastmsg_id')} + ON ${sessionsTable}(${ownerIdColumn}, ${lastMessageAtColumn} DESC, ${idColumn})`, + ); + + // Index intended for faster sessionRepository.getOneByIdAndUserId queries and joins + await runQuery( + `CREATE INDEX IF NOT EXISTS ${escape.indexName('chat_hub_messages_sessionId')} + ON ${messagesTable}(${sessionIdColumn})`, + ); + } + + async down({ schemaBuilder: { dropNotNull }, runQuery, escape }: MigrationContext) { + await runQuery( + `DROP INDEX IF EXISTS ${escape.indexName('chat_hub_sessions_owner_lastmsg_id')}`, + ); + await runQuery(`DROP INDEX IF EXISTS ${escape.indexName('chat_hub_messages_sessionId')}`); + + await dropNotNull('chat_hub_sessions', 'lastMessageAt'); + } +} diff --git a/packages/@n8n/db/src/migrations/postgresdb/index.ts b/packages/@n8n/db/src/migrations/postgresdb/index.ts index a803611d14c..0034ba5b041 100644 --- a/packages/@n8n/db/src/migrations/postgresdb/index.ts +++ b/packages/@n8n/db/src/migrations/postgresdb/index.ts @@ -130,6 +130,7 @@ import { AddIconToAgentTable1765788427674 } from '../common/1765788427674-AddIco import { AddAgentIdForeignKeys1765886667897 } from '../common/1765886667897-AddAgentIdForeignKeys'; import { AddWorkflowVersionIdToExecutionData1765892199653 } from '../common/1765892199653-AddVersionIdToExecutionData'; import { AddWorkflowPublishScopeToProjectRoles1766064542000 } from '../common/1766064542000-AddWorkflowPublishScopeToProjectRoles'; +import { AddChatMessageIndices1766068346315 } from '../common/1766068346315-AddChatMessageIndices'; import type { Migration } from '../migration-types'; export const postgresMigrations: Migration[] = [ @@ -265,4 +266,5 @@ export const postgresMigrations: Migration[] = [ AddAgentIdForeignKeys1765886667897, AddWorkflowVersionIdToExecutionData1765892199653, AddWorkflowPublishScopeToProjectRoles1766064542000, + AddChatMessageIndices1766068346315, ]; diff --git a/packages/@n8n/db/src/migrations/sqlite/1766068346315-AddChatMessageIndices.ts b/packages/@n8n/db/src/migrations/sqlite/1766068346315-AddChatMessageIndices.ts new file mode 100644 index 00000000000..3ba60dc2823 --- /dev/null +++ b/packages/@n8n/db/src/migrations/sqlite/1766068346315-AddChatMessageIndices.ts @@ -0,0 +1,5 @@ +import { AddChatMessageIndices1766068346315 as BaseMigration } from '../common/1766068346315-AddChatMessageIndices'; + +export class AddChatMessageIndices1766068346315 extends BaseMigration { + transaction = false as const; +} diff --git a/packages/@n8n/db/src/migrations/sqlite/index.ts b/packages/@n8n/db/src/migrations/sqlite/index.ts index 0a5f0f65af9..12e8547d847 100644 --- a/packages/@n8n/db/src/migrations/sqlite/index.ts +++ b/packages/@n8n/db/src/migrations/sqlite/index.ts @@ -49,6 +49,7 @@ import { ChangeDependencyInfoToJson1761655473000 } from './1761655473000-ChangeD import { AddCreatorIdToProjectTable1764276827837 } from './1764276827837-AddCreatorIdToProjectTable'; import { AddResolvableFieldsToCredentials1764689448000 } from './1764689448000-AddResolvableFieldsToCredentials'; import { AddAgentIdForeignKeys1765886667897 } from './1765886667897-AddAgentIdForeignKeys'; +import { AddChatMessageIndices1766068346315 } from './1766068346315-AddChatMessageIndices'; import { UniqueWorkflowNames1620821879465 } from '../common/1620821879465-UniqueWorkflowNames'; import { UpdateWorkflowCredentials1630330987096 } from '../common/1630330987096-UpdateWorkflowCredentials'; import { AddNodeIds1658930531669 } from '../common/1658930531669-AddNodeIds'; @@ -255,6 +256,7 @@ const sqliteMigrations: Migration[] = [ AddAgentIdForeignKeys1765886667897, AddWorkflowVersionIdToExecutionData1765892199653, AddWorkflowPublishScopeToProjectRoles1766064542000, + AddChatMessageIndices1766068346315, ]; export { sqliteMigrations }; diff --git a/packages/cli/src/modules/chat-hub/__tests__/chat-hub.service.integration.test.ts b/packages/cli/src/modules/chat-hub/__tests__/chat-hub.service.integration.test.ts index e67d61d4d83..646cd965d3b 100644 --- a/packages/cli/src/modules/chat-hub/__tests__/chat-hub.service.integration.test.ts +++ b/packages/cli/src/modules/chat-hub/__tests__/chat-hub.service.integration.test.ts @@ -2,14 +2,14 @@ import { mockInstance, testDb, testModules, createActiveWorkflow } from '@n8n/ba import type { User } from '@n8n/db'; import { ProjectRepository } from '@n8n/db'; import { Container } from '@n8n/di'; +import { createAdmin, createMember } from '@test-integration/db/users'; import { BinaryDataService } from 'n8n-core'; import { CHAT_TRIGGER_NODE_TYPE } from 'n8n-workflow'; -import { createAdmin, createMember } from '@test-integration/db/users'; +import { ChatHubAgentRepository } from '../chat-hub-agent.repository'; import { ChatHubService } from '../chat-hub.service'; import { ChatHubMessageRepository } from '../chat-message.repository'; import { ChatHubSessionRepository } from '../chat-session.repository'; -import { ChatHubAgentRepository } from '../chat-hub-agent.repository'; mockInstance(BinaryDataService); @@ -318,28 +318,15 @@ describe('chatHub', () => { ).rejects.toThrow('Cursor session not found'); }); - it('should handle sessions with null lastMessageAt', async () => { - const session1 = await sessionsRepository.createChatSession({ - id: crypto.randomUUID(), - ownerId: member.id, - title: 'Session with date', - lastMessageAt: new Date('2025-01-01T00:00:00Z'), - tools: [], - }); - - const session2 = await sessionsRepository.createChatSession({ - id: crypto.randomUUID(), - ownerId: member.id, - title: 'Session without date', - lastMessageAt: null, - tools: [], - }); - - const conversations = await chatHubService.getConversations(member.id, 10); - - expect(conversations.data).toHaveLength(2); - expect(conversations.data[0].id).toBe(session1.id); - expect(conversations.data[1].id).toBe(session2.id); + it('should disallow sessions without lastMessageAt', async () => { + await expect( + sessionsRepository.createChatSession({ + id: crypto.randomUUID(), + ownerId: member.id, + title: 'Session with date', + tools: [], + }), + ).rejects.toThrow(); }); }); }); @@ -470,7 +457,7 @@ describe('chatHub', () => { crypto.randomUUID(), ]; - const msg1 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[0], sessionId: session.id, name: 'Nathan', @@ -478,31 +465,31 @@ describe('chatHub', () => { content: 'message 1', createdAt: new Date('2025-01-03T00:00:00Z'), }); - const msg2 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[1], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 2', - previousMessageId: msg1.id, + previousMessageId: ids[0], createdAt: new Date('2025-01-03T00:05:00Z'), }); - const msg3 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[2], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3', - previousMessageId: msg2.id, + previousMessageId: ids[1], createdAt: new Date('2025-01-03T00:10:00Z'), }); - const msg4 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[3], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 4', - previousMessageId: msg3.id, + previousMessageId: ids[2], createdAt: new Date('2025-01-03T00:15:00Z'), }); @@ -515,14 +502,14 @@ describe('chatHub', () => { } = response; expect(Object.keys(messages)).toHaveLength(4); - expect(messages[msg1.id].content).toBe('message 1'); - expect(messages[msg1.id].type).toBe('human'); - expect(messages[msg2.id].content).toBe('message 2'); - expect(messages[msg2.id].type).toBe('ai'); - expect(messages[msg3.id].content).toBe('message 3'); - expect(messages[msg3.id].type).toBe('human'); - expect(messages[msg4.id].content).toBe('message 4'); - expect(messages[msg4.id].type).toBe('ai'); + expect(messages[ids[0]].content).toBe('message 1'); + expect(messages[ids[0]].type).toBe('human'); + expect(messages[ids[1]].content).toBe('message 2'); + expect(messages[ids[1]].type).toBe('ai'); + expect(messages[ids[2]].content).toBe('message 3'); + expect(messages[ids[2]].type).toBe('human'); + expect(messages[ids[3]].content).toBe('message 4'); + expect(messages[ids[3]].type).toBe('ai'); }); it('should get conversation with a edit branch', async () => { @@ -542,7 +529,7 @@ describe('chatHub', () => { lastMessageAt: new Date('2025-01-03T00:00:00Z'), tools: [], }); - const msg1 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[0], sessionId: session.id, name: 'Nathan', @@ -550,51 +537,51 @@ describe('chatHub', () => { content: 'message 1', createdAt: new Date('2025-01-03T00:00:00Z'), }); - const msg2 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[1], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 2', - previousMessageId: msg1.id, + previousMessageId: ids[0], createdAt: new Date('2025-01-03T00:05:00Z'), }); - const msg3 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[2], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3a', - previousMessageId: msg2.id, + previousMessageId: ids[1], createdAt: new Date('2025-01-03T00:10:00Z'), }); - const msg4 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[3], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 4a', - previousMessageId: msg3.id, + previousMessageId: ids[2], createdAt: new Date('2025-01-03T00:15:00Z'), }); // Edit message 3 to create a branch - const msg5 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[4], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3b', - previousMessageId: msg2.id, - revisionOfMessageId: msg3.id, + previousMessageId: ids[1], + revisionOfMessageId: ids[2], createdAt: new Date('2025-01-03T00:20:00Z'), }); - const msg6 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[5], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 4b', - previousMessageId: msg5.id, + previousMessageId: ids[4], createdAt: new Date('2025-01-03T00:25:00Z'), }); @@ -607,13 +594,13 @@ describe('chatHub', () => { } = response; expect(Object.keys(messages)).toHaveLength(6); - expect(messages[msg1.id].content).toBe('message 1'); - expect(messages[msg2.id].content).toBe('message 2'); - expect(messages[msg3.id].content).toBe('message 3a'); - expect(messages[msg4.id].content).toBe('message 4a'); - expect(messages[msg5.id].content).toBe('message 3b'); - expect(messages[msg6.id].content).toBe('message 4b'); - expect(messages[msg5.id].previousMessageId).toBe(msg2.id); + expect(messages[ids[0]].content).toBe('message 1'); + expect(messages[ids[1]].content).toBe('message 2'); + expect(messages[ids[2]].content).toBe('message 3a'); + expect(messages[ids[3]].content).toBe('message 4a'); + expect(messages[ids[4]].content).toBe('message 3b'); + expect(messages[ids[5]].content).toBe('message 4b'); + expect(messages[ids[4]].previousMessageId).toBe(ids[1]); }); it('should get conversation with a edit branch at first message', async () => { @@ -631,7 +618,7 @@ describe('chatHub', () => { tools: [], }); - const msg1 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[0], sessionId: session.id, name: 'Nathan', @@ -645,17 +632,17 @@ describe('chatHub', () => { name: 'ChatGPT', type: 'ai', content: 'message 2a', - previousMessageId: msg1.id, + previousMessageId: ids[0], createdAt: new Date('2025-01-03T00:05:00Z'), }); // Edit message 1 to create a branch - const msg3 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[2], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 1b', - revisionOfMessageId: msg1.id, + revisionOfMessageId: ids[0], createdAt: new Date('2025-01-03T00:10:00Z'), }); await messagesRepository.createChatMessage({ @@ -664,7 +651,7 @@ describe('chatHub', () => { name: 'ChatGPT', type: 'ai', content: 'message 2b', - previousMessageId: msg3.id, + previousMessageId: ids[2], createdAt: new Date('2025-01-03T00:15:00Z'), }); @@ -696,7 +683,7 @@ describe('chatHub', () => { lastMessageAt: new Date('2025-01-03T00:00:00Z'), tools: [], }); - const msg1 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[0], sessionId: session.id, name: 'Nathan', @@ -704,42 +691,42 @@ describe('chatHub', () => { content: 'message 1', createdAt: new Date('2025-01-03T00:00:00Z'), }); - const msg2 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[1], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 2', - previousMessageId: msg1.id, + previousMessageId: ids[0], createdAt: new Date('2025-01-03T00:05:00Z'), }); - const msg3 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[2], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3', - previousMessageId: msg2.id, + previousMessageId: ids[1], createdAt: new Date('2025-01-03T00:10:00Z'), }); - const msg4 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[3], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 4a', - previousMessageId: msg3.id, + previousMessageId: ids[2], createdAt: new Date('2025-01-03T00:15:00Z'), }); // Retry message 4 to create a branch - const msg5 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[4], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 4b', - previousMessageId: msg3.id, - retryOfMessageId: msg4.id, + previousMessageId: ids[2], + retryOfMessageId: ids[3], createdAt: new Date('2025-01-03T00:20:00Z'), }); @@ -752,20 +739,13 @@ describe('chatHub', () => { } = response; expect(Object.keys(messages)).toHaveLength(5); - expect(messages[msg5.id].previousMessageId).toBe(msg3.id); - expect(messages[msg5.id].retryOfMessageId).toBe(msg4.id); + expect(messages[ids[4]].previousMessageId).toBe(ids[2]); + expect(messages[ids[4]].retryOfMessageId).toBe(ids[3]); }); it('should get a complex conversation with multiple branches', async () => { // This test creates a complex conversation with multiple edits and retries to ensure // the conversation tree is built correctly in all cases. - - // The structure created is as follows: - // msg1 -> msg2 -> msg3a -> msg4a - // -> msg3b (edit of msg3a) -> msg4b - // msg1b (edit of msg1) -> nothing - // msg1 -> msg2r (retry of msg2) -> msg3d -> msg4c - const ids = [ crypto.randomUUID(), crypto.randomUUID(), @@ -786,7 +766,7 @@ describe('chatHub', () => { lastMessageAt: new Date('2025-01-03T00:00:00Z'), tools: [], }); - const msg1 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[0], sessionId: session.id, name: 'Nathan', @@ -794,22 +774,22 @@ describe('chatHub', () => { content: 'message 1', createdAt: new Date('2025-01-03T00:00:00Z'), }); - const msg2 = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[1], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 2a', - previousMessageId: msg1.id, + previousMessageId: ids[0], createdAt: new Date('2025-01-03T00:05:00Z'), }); - const msg3a = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[2], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3a', - previousMessageId: msg2.id, + previousMessageId: ids[1], createdAt: new Date('2025-01-03T00:10:00Z'), }); await messagesRepository.createChatMessage({ @@ -818,17 +798,17 @@ describe('chatHub', () => { name: 'ChatGPT', type: 'ai', content: 'message 4a', - previousMessageId: msg3a.id, + previousMessageId: ids[2], createdAt: new Date('2025-01-03T00:15:00Z'), }); - const msg3b = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[4], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3b', - revisionOfMessageId: msg3a.id, - previousMessageId: msg2.id, + revisionOfMessageId: ids[2], + previousMessageId: ids[1], createdAt: new Date('2025-01-03T00:20:00Z'), }); await messagesRepository.createChatMessage({ @@ -837,7 +817,7 @@ describe('chatHub', () => { name: 'ChatGPT', type: 'ai', content: 'message 4b', - previousMessageId: msg3b.id, + previousMessageId: ids[4], createdAt: new Date('2025-01-03T00:25:00Z'), }); await messagesRepository.createChatMessage({ @@ -846,26 +826,26 @@ describe('chatHub', () => { name: 'Nathan', type: 'human', content: 'message 1b', - revisionOfMessageId: msg1.id, + revisionOfMessageId: ids[0], createdAt: new Date('2025-01-03T00:30:00Z'), }); - const msg2r = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[7], sessionId: session.id, name: 'ChatGPT', type: 'ai', content: 'message 2b', - previousMessageId: msg1.id, - retryOfMessageId: msg2.id, + previousMessageId: ids[0], + retryOfMessageId: ids[1], createdAt: new Date('2025-01-03T00:35:00Z'), }); - const msg3d = await messagesRepository.createChatMessage({ + await messagesRepository.createChatMessage({ id: ids[8], sessionId: session.id, name: 'Nathan', type: 'human', content: 'message 3d', - previousMessageId: msg2r.id, + previousMessageId: ids[7], createdAt: new Date('2025-01-03T00:40:00Z'), }); await messagesRepository.createChatMessage({ @@ -874,7 +854,7 @@ describe('chatHub', () => { name: 'ChatGPT', type: 'ai', content: 'message 4c', - previousMessageId: msg3d.id, + previousMessageId: ids[8], createdAt: new Date('2025-01-03T00:45:00Z'), }); @@ -888,8 +868,8 @@ describe('chatHub', () => { expect(Object.keys(messages)).toHaveLength(10); - expect(messages[msg2r.id].previousMessageId).toBe(msg1.id); - expect(messages[msg2r.id].retryOfMessageId).toBe(msg2.id); + expect(messages[ids[7]].previousMessageId).toBe(ids[0]); + expect(messages[ids[7]].retryOfMessageId).toBe(ids[1]); }); }); }); diff --git a/packages/cli/src/modules/chat-hub/chat-hub-agent.repository.ts b/packages/cli/src/modules/chat-hub/chat-hub-agent.repository.ts index 28eea18a0bb..e0eb66edcce 100644 --- a/packages/cli/src/modules/chat-hub/chat-hub-agent.repository.ts +++ b/packages/cli/src/modules/chat-hub/chat-hub-agent.repository.ts @@ -14,27 +14,42 @@ export class ChatHubAgentRepository extends Repository { agent: Partial & Pick, trx?: EntityManager, ) { - return await withTransaction(this.manager, trx, async (em) => { - await em.insert(ChatHubAgent, agent); - return await em.findOneOrFail(ChatHubAgent, { - where: { id: agent.id }, - }); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + await em.insert(ChatHubAgent, agent); + return await em.findOneOrFail(ChatHubAgent, { + where: { id: agent.id }, + }); + }, + false, + ); } async updateAgent(id: string, updates: Partial, trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - await em.update(ChatHubAgent, { id }, updates); - return await em.findOneOrFail(ChatHubAgent, { - where: { id }, - }); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + await em.update(ChatHubAgent, { id }, updates); + return await em.findOneOrFail(ChatHubAgent, { + where: { id }, + }); + }, + false, + ); } async deleteAgent(id: string, trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - return await em.delete(ChatHubAgent, { id }); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + return await em.delete(ChatHubAgent, { id }); + }, + false, + ); } async getManyByUserId(userId: string) { diff --git a/packages/cli/src/modules/chat-hub/chat-hub-session.entity.ts b/packages/cli/src/modules/chat-hub/chat-hub-session.entity.ts index 543bfcc0820..2a4960ba66b 100644 --- a/packages/cli/src/modules/chat-hub/chat-hub-session.entity.ts +++ b/packages/cli/src/modules/chat-hub/chat-hub-session.entity.ts @@ -27,7 +27,7 @@ export interface IChatHubSession { updatedAt: Date; title: string; ownerId: string; - lastMessageAt: Date | null; + lastMessageAt: Date; credentialId: string | null; provider: ChatHubProvider | null; model: string | null; @@ -66,8 +66,8 @@ export class ChatHubSession extends WithTimestamps { * Timestamp of the last active message in the session. * Used to sort chat sessions by recent activity. */ - @DateTimeColumn({ nullable: true }) - lastMessageAt: Date | null; + @DateTimeColumn() + lastMessageAt: Date; /* * ID of the selected credential to use by default with the selected LLM provider (if applicable). diff --git a/packages/cli/src/modules/chat-hub/chat-hub.attachment.service.ts b/packages/cli/src/modules/chat-hub/chat-hub.attachment.service.ts index a7aeffd4d39..1775599c92f 100644 --- a/packages/cli/src/modules/chat-hub/chat-hub.attachment.service.ts +++ b/packages/cli/src/modules/chat-hub/chat-hub.attachment.service.ts @@ -1,14 +1,17 @@ +import type { ChatMessageId, ChatSessionId, ChatAttachment } from '@n8n/api-types'; import { Service } from '@n8n/di'; -import { BINARY_ENCODING, type IBinaryData } from 'n8n-workflow'; +import { Not, IsNull } from '@n8n/typeorm'; +import type { EntityManager } from '@n8n/typeorm'; import { sanitizeFilename } from '@n8n/utils'; import { BinaryDataService, FileLocation } from 'n8n-core'; -import { Not, IsNull } from '@n8n/typeorm'; -import { ChatHubMessageRepository } from './chat-message.repository'; -import type { ChatMessageId, ChatSessionId, ChatAttachment } from '@n8n/api-types'; -import { NotFoundError } from '@/errors/response-errors/not-found.error'; -import { BadRequestError } from '@/errors/response-errors/bad-request.error'; +import { BINARY_ENCODING, type IBinaryData } from 'n8n-workflow'; import type Stream from 'node:stream'; +import { ChatHubMessageRepository } from './chat-message.repository'; + +import { BadRequestError } from '@/errors/response-errors/bad-request.error'; +import { NotFoundError } from '@/errors/response-errors/not-found.error'; + @Service() export class ChatHubAttachmentService { private readonly maxTotalSizeBytes = 200 * 1024 * 1024; // 200 MB @@ -97,9 +100,9 @@ export class ChatHubAttachmentService { /** * Deletes all files attached to messages in the session */ - async deleteAllBySessionId(sessionId: string): Promise { - const messages = await this.messageRepository.getManyBySessionId(sessionId); - + async deleteAllBySessionId(sessionId: string, trx?: EntityManager): Promise { + const messages = await this.messageRepository.getManyBySessionId(sessionId, trx); + // Attachment deletion cannot be rolled back, and the transaction doesn't cover it. await this.deleteAttachments(messages.flatMap((message) => message.attachments ?? [])); } diff --git a/packages/cli/src/modules/chat-hub/chat-hub.controller.ts b/packages/cli/src/modules/chat-hub/chat-hub.controller.ts index 0786e3152d2..8c4545dd281 100644 --- a/packages/cli/src/modules/chat-hub/chat-hub.controller.ts +++ b/packages/cli/src/modules/chat-hub/chat-hub.controller.ts @@ -96,7 +96,7 @@ export class ChatHubController { } // Verify user has access to this session - await this.chatService.getConversation(req.user.id, sessionId); + await this.chatService.ensureConversation(req.user.id, sessionId); const [{ mimeType, fileName }, attachmentAsStreamOrBuffer] = await this.chatAttachmentService.getAttachment(sessionId, messageId, attachmentIndex); diff --git a/packages/cli/src/modules/chat-hub/chat-hub.service.ts b/packages/cli/src/modules/chat-hub/chat-hub.service.ts index 58014988b15..834b3075f22 100644 --- a/packages/cli/src/modules/chat-hub/chat-hub.service.ts +++ b/packages/cli/src/modules/chat-hub/chat-hub.service.ts @@ -42,6 +42,13 @@ import { AGENT_LANGCHAIN_NODE_TYPE, } from 'n8n-workflow'; +import { ActiveExecutions } from '@/active-executions'; +import { BadRequestError } from '@/errors/response-errors/bad-request.error'; +import { NotFoundError } from '@/errors/response-errors/not-found.error'; +import { ExecutionService } from '@/executions/execution.service'; +import { WorkflowExecutionService } from '@/workflows/workflow-execution.service'; +import { WorkflowFinderService } from '@/workflows/workflow-finder.service'; + import { ChatHubAgentService } from './chat-hub-agent.service'; import { ChatHubCredentialsService } from './chat-hub-credentials.service'; import type { ChatHubMessage } from './chat-hub-message.entity'; @@ -57,6 +64,7 @@ import { PROVIDER_NODE_TYPE_MAP, TOOLS_AGENT_NODE_MIN_VERSION, } from './chat-hub.constants'; +import { ChatHubModelsService } from './chat-hub.models.service'; import { ChatHubSettingsService } from './chat-hub.settings.service'; import { HumanMessagePayload, @@ -69,14 +77,6 @@ import { ChatHubMessageRepository } from './chat-message.repository'; import { ChatHubSessionRepository } from './chat-session.repository'; import { interceptResponseWrites, createStructuredChunkAggregator } from './stream-capturer'; -import { ActiveExecutions } from '@/active-executions'; -import { BadRequestError } from '@/errors/response-errors/bad-request.error'; -import { NotFoundError } from '@/errors/response-errors/not-found.error'; -import { ExecutionService } from '@/executions/execution.service'; -import { WorkflowExecutionService } from '@/workflows/workflow-execution.service'; -import { WorkflowFinderService } from '@/workflows/workflow-finder.service'; -import { ChatHubModelsService } from './chat-hub.models.service'; - @Service() export class ChatHubService { constructor( @@ -668,12 +668,9 @@ export class ChatHubService { } async stopGeneration(user: User, sessionId: ChatSessionId, messageId: ChatMessageId) { - const session = await this.getChatSession(user, sessionId); - if (!session) { - throw new NotFoundError('Chat session not found'); - } + await this.ensureConversation(user.id, sessionId); - const message = await this.getChatMessage(session.id, messageId, [ + const message = await this.getChatMessage(sessionId, messageId, [ 'execution', 'execution.workflow', ]); @@ -945,7 +942,7 @@ export class ChatHubService { try { const title = await this.runTitleWorkflowAndGetTitle(user, workflowData, executionData); if (title) { - await this.sessionRepository.updateChatTitle(sessionId, title); + await this.sessionRepository.updateChatSession(sessionId, { title }); } } catch (error: unknown) { if (error instanceof Error) { @@ -1272,6 +1269,7 @@ export class ChatHubService { id: sessionId, ownerId: user.id, title: 'New Chat', + lastMessageAt: new Date(), agentName, tools, credentialId, @@ -1315,6 +1313,16 @@ export class ChatHubService { }; } + /** + * Ensures conversation exists and belongs to the user, throws otherwise + * */ + async ensureConversation(userId: string, sessionId: string, trx?: EntityManager): Promise { + const sessionExists = await this.sessionRepository.existsById(sessionId, userId, trx); + if (!sessionExists) { + throw new NotFoundError('Chat session not found'); + } + } + /** * Get a single conversation with messages and ready to render timeline of latest messages * */ @@ -1324,7 +1332,7 @@ export class ChatHubService { throw new NotFoundError('Chat session not found'); } - const messages = await this.messageRepository.getManyBySessionId(sessionId); + const messages = session.messages ?? []; return { session: this.convertSessionEntityToDto(session), @@ -1391,19 +1399,6 @@ export class ChatHubService { return result; } - /** - * Updates the title of a session - */ - async updateSessionTitle(userId: string, sessionId: ChatSessionId, title: string) { - const session = await this.sessionRepository.getOneById(sessionId, userId); - - if (!session) { - throw new NotFoundError('Session not found'); - } - - return await this.sessionRepository.updateChatTitle(sessionId, title); - } - /** * Updates a session with the provided fields */ @@ -1412,11 +1407,7 @@ export class ChatHubService { sessionId: ChatSessionId, updates: ChatHubUpdateConversationRequest, ) { - const session = await this.sessionRepository.getOneById(sessionId, user.id); - - if (!session) { - throw new NotFoundError('Session not found'); - } + await this.ensureConversation(user.id, sessionId); // Prepare the actual updates to be sent to the repository const sessionUpdates: Partial = {}; @@ -1453,14 +1444,11 @@ export class ChatHubService { * Deletes a session */ async deleteSession(userId: string, sessionId: ChatSessionId) { - const session = await this.sessionRepository.getOneById(sessionId, userId); - - if (!session) { - throw new NotFoundError('Session not found'); - } - - await this.chatHubAttachmentService.deleteAllBySessionId(sessionId); - await this.sessionRepository.deleteChatHubSession(sessionId); + await this.messageRepository.manager.transaction(async (trx) => { + await this.ensureConversation(userId, sessionId, trx); + await this.chatHubAttachmentService.deleteAllBySessionId(sessionId, trx); + await this.sessionRepository.deleteChatHubSession(sessionId, trx); + }); } private async ensureValidModel(user: User, model: ChatHubConversationModel) { diff --git a/packages/cli/src/modules/chat-hub/chat-message.repository.ts b/packages/cli/src/modules/chat-hub/chat-message.repository.ts index ab4daa173f1..f53dfcdb7df 100644 --- a/packages/cli/src/modules/chat-hub/chat-message.repository.ts +++ b/packages/cli/src/modules/chat-hub/chat-message.repository.ts @@ -2,11 +2,11 @@ import type { ChatHubMessageStatus, ChatMessageId, ChatSessionId } from '@n8n/ap import { withTransaction } from '@n8n/db'; import { Service } from '@n8n/di'; import { DataSource, EntityManager, Repository } from '@n8n/typeorm'; +import { QueryDeepPartialEntity } from '@n8n/typeorm/query-builder/QueryPartialEntity'; +import { UnexpectedError, type IBinaryData } from 'n8n-workflow'; import { ChatHubMessage } from './chat-hub-message.entity'; import { ChatHubSessionRepository } from './chat-session.repository'; -import { UnexpectedError, type IBinaryData } from 'n8n-workflow'; -import { QueryDeepPartialEntity } from '@n8n/typeorm/query-builder/QueryPartialEntity'; @Service() export class ChatHubMessageRepository extends Repository { @@ -19,23 +19,24 @@ export class ChatHubMessageRepository extends Repository { async createChatMessage(message: QueryDeepPartialEntity, trx?: EntityManager) { const messageId = message.id; - if (typeof messageId === 'function') { + const sessionId = message.sessionId; + + if (typeof messageId === 'function' || !messageId) { throw new UnexpectedError('Message ID is required and must be a string value'); } - return await withTransaction( - this.manager, - trx, - async (em) => { - await em.insert(ChatHubMessage, message); - const saved = await em.findOneOrFail(ChatHubMessage, { - where: { id: messageId }, - }); - await this.chatSessionRepository.updateLastMessageAt(saved.sessionId, saved.createdAt, em); - return saved; - }, - false, - ); + if (typeof sessionId === 'function' || !sessionId) { + throw new UnexpectedError('Session ID is required and must be a string value'); + } + + return await withTransaction(this.manager, trx, async (em) => { + await em.insert(ChatHubMessage, message); + await this.chatSessionRepository.updateChatSession( + sessionId, + { lastMessageAt: new Date() }, + em, + ); + }); } async updateChatMessage( diff --git a/packages/cli/src/modules/chat-hub/chat-session.repository.ts b/packages/cli/src/modules/chat-hub/chat-session.repository.ts index 70feed1cedf..b943b6f6607 100644 --- a/packages/cli/src/modules/chat-hub/chat-session.repository.ts +++ b/packages/cli/src/modules/chat-hub/chat-session.repository.ts @@ -16,49 +16,40 @@ export class ChatHubSessionRepository extends Repository { session: Partial & Pick, trx?: EntityManager, ) { - return await withTransaction(this.manager, trx, async (em) => { - await em.insert(ChatHubSession, session); - return await em.findOneOrFail(ChatHubSession, { - where: { id: session.id }, - relations: ['messages'], - }); - }); - } - - async updateLastMessageAt(id: string, lastMessageAt: Date, trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - await em.update(ChatHubSession, { id }, { lastMessageAt }); - return await em.findOneOrFail(ChatHubSession, { - where: { id }, - relations: ['messages'], - }); - }); - } - - async updateChatTitle(id: string, title: string, trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - await em.update(ChatHubSession, { id }, { title }); - return await em.findOneOrFail(ChatHubSession, { - where: { id }, - relations: ['messages'], - }); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + await em.insert(ChatHubSession, session); + return await em.findOneOrFail(ChatHubSession, { + where: { id: session.id }, + relations: ['messages'], + }); + }, + false, + ); } async updateChatSession(id: string, updates: Partial, trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - await em.update(ChatHubSession, { id }, updates); - return await em.findOneOrFail(ChatHubSession, { - where: { id }, - relations: ['messages'], - }); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + await em.update(ChatHubSession, { id }, updates); + }, + false, + ); } async deleteChatHubSession(id: string, trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - return await em.delete(ChatHubSession, { id }); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + return await em.delete(ChatHubSession, { id }); + }, + false, + ); } async getManyByUserId(userId: string, limit: number, cursor?: string) { @@ -69,8 +60,7 @@ export class ChatHubSessionRepository extends Repository { .leftJoinAndSelect('shared.project', 'project') .leftJoinAndSelect('workflow.activeVersion', 'activeVersion') .where('session.ownerId = :userId', { userId }) - .addSelect("COALESCE(session.lastMessageAt, '1970-01-01')", 'sortdate') - .orderBy('sortdate', 'DESC') + .orderBy('session.lastMessageAt', 'DESC') .addOrderBy('session.id', 'ASC'); if (cursor) { @@ -96,6 +86,17 @@ export class ChatHubSessionRepository extends Repository { return await queryBuilder.getMany(); } + async existsById(id: string, userId: string, trx?: EntityManager): Promise { + return await withTransaction( + this.manager, + trx, + async (em) => { + return await em.exists(ChatHubSession, { where: { id, ownerId: userId } }); + }, + false, + ); + } + async getOneById(id: string, userId: string, trx?: EntityManager) { return await withTransaction( this.manager, @@ -120,8 +121,13 @@ export class ChatHubSessionRepository extends Repository { } async deleteAll(trx?: EntityManager) { - return await withTransaction(this.manager, trx, async (em) => { - return await em.createQueryBuilder().delete().from(ChatHubSession).execute(); - }); + return await withTransaction( + this.manager, + trx, + async (em) => { + return await em.createQueryBuilder().delete().from(ChatHubSession).execute(); + }, + false, + ); } }