refactor(core): Make Chat db queries faster (#23397)

This commit is contained in:
Jaakko Husso
2025-12-19 12:03:51 +02:00
committed by GitHub
parent 150d16d410
commit 6b2887c95f
12 changed files with 274 additions and 228 deletions
@@ -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');
}
}
@@ -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,
];
@@ -0,0 +1,5 @@
import { AddChatMessageIndices1766068346315 as BaseMigration } from '../common/1766068346315-AddChatMessageIndices';
export class AddChatMessageIndices1766068346315 extends BaseMigration {
transaction = false as const;
}
@@ -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 };
@@ -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]);
});
});
});
@@ -14,27 +14,42 @@ export class ChatHubAgentRepository extends Repository<ChatHubAgent> {
agent: Partial<IChatHubAgent> & Pick<IChatHubAgent, 'id'>,
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<IChatHubAgent>, 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) {
@@ -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).
@@ -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<void> {
const messages = await this.messageRepository.getManyBySessionId(sessionId);
async deleteAllBySessionId(sessionId: string, trx?: EntityManager): Promise<void> {
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 ?? []));
}
@@ -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);
@@ -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<void> {
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<IChatHubSession> = {};
@@ -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) {
@@ -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<ChatHubMessage> {
@@ -19,23 +19,24 @@ export class ChatHubMessageRepository extends Repository<ChatHubMessage> {
async createChatMessage(message: QueryDeepPartialEntity<ChatHubMessage>, 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(
@@ -16,49 +16,40 @@ export class ChatHubSessionRepository extends Repository<ChatHubSession> {
session: Partial<IChatHubSession> & Pick<IChatHubSession, 'id'>,
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<IChatHubSession>, 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<ChatHubSession> {
.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<ChatHubSession> {
return await queryBuilder.getMany();
}
async existsById(id: string, userId: string, trx?: EntityManager): Promise<boolean> {
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<ChatHubSession> {
}
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,
);
}
}