From 862a50c9828e26522ad23a7defc4fb28e177496b Mon Sep 17 00:00:00 2001 From: fengmk2 Date: Mon, 23 Jun 2025 14:56:04 +0800 Subject: [PATCH] fix(server): use job queue instead event on doc indexing changes (#12893) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit close CLOUD-231 #### PR Dependency Tree * **PR #12893** 👈 This tree was auto-generated by [Charcoal](https://github.com/danerwilliams/charcoal) ## Summary by CodeRabbit - **Refactor** - Updated background processing for document indexing and deletion to use a job queue system instead of event-based triggers. - **Bug Fixes** - Improved reliability of embedding updates and deletions by ensuring tasks are properly queued and processed. - **Tests** - Adjusted tests to verify that document operations correctly trigger job queue actions. No changes to user-facing features or interface. --- .../server/src/plugins/copilot/embedding/job.ts | 12 ++++++++---- .../server/src/plugins/copilot/embedding/types.ts | 10 ++++++++++ .../src/plugins/indexer/__tests__/service.spec.ts | 8 ++++---- .../backend/server/src/plugins/indexer/index.ts | 13 ------------- .../backend/server/src/plugins/indexer/service.ts | 14 ++++++++------ 5 files changed, 30 insertions(+), 27 deletions(-) diff --git a/packages/backend/server/src/plugins/copilot/embedding/job.ts b/packages/backend/server/src/plugins/copilot/embedding/job.ts index a16272c7b..5b1adfbe6 100644 --- a/packages/backend/server/src/plugins/copilot/embedding/job.ts +++ b/packages/backend/server/src/plugins/copilot/embedding/job.ts @@ -159,8 +159,10 @@ export class CopilotEmbeddingJob { } } - @OnEvent('doc.indexer.updated') - async addDocEmbeddingQueueFromEvent(doc: Events['doc.indexer.updated']) { + @OnJob('copilot.embedding.updateDoc') + async addDocEmbeddingQueueFromEvent( + doc: Jobs['copilot.embedding.updateDoc'] + ) { if (!this.supportEmbedding || !this.embeddingClient) return; await this.queue.add( @@ -176,8 +178,10 @@ export class CopilotEmbeddingJob { ); } - @OnEvent('doc.indexer.deleted') - async deleteDocEmbeddingQueueFromEvent(doc: Events['doc.indexer.deleted']) { + @OnJob('copilot.embedding.deleteDoc') + async deleteDocEmbeddingQueueFromEvent( + doc: Jobs['copilot.embedding.deleteDoc'] + ) { await this.queue.remove( `workspace:embedding:${doc.workspaceId}:${doc.docId}`, 'copilot.embedding.docs' diff --git a/packages/backend/server/src/plugins/copilot/embedding/types.ts b/packages/backend/server/src/plugins/copilot/embedding/types.ts index 5de7b2252..afbf9d57e 100644 --- a/packages/backend/server/src/plugins/copilot/embedding/types.ts +++ b/packages/backend/server/src/plugins/copilot/embedding/types.ts @@ -43,6 +43,16 @@ declare global { docId: string; }; + 'copilot.embedding.updateDoc': { + workspaceId: string; + docId: string; + }; + + 'copilot.embedding.deleteDoc': { + workspaceId: string; + docId: string; + }; + 'copilot.embedding.files': { contextId?: string; userId: string; diff --git a/packages/backend/server/src/plugins/indexer/__tests__/service.spec.ts b/packages/backend/server/src/plugins/indexer/__tests__/service.spec.ts index f533499a5..84782dca9 100644 --- a/packages/backend/server/src/plugins/indexer/__tests__/service.spec.ts +++ b/packages/backend/server/src/plugins/indexer/__tests__/service.spec.ts @@ -1884,12 +1884,12 @@ test('should delete doc work', async t => { t.is(result4.nodes.length, 1); t.deepEqual(result4.nodes[0].fields.docId, [docId2]); - const count = module.event.count('doc.indexer.deleted'); + const count = module.queue.count('copilot.embedding.deleteDoc'); await indexerService.deleteDoc(workspaceId, docId1, { refresh: true, }); - t.is(module.event.count('doc.indexer.deleted'), count + 1); + t.is(module.queue.count('copilot.embedding.deleteDoc'), count + 1); // make sure the docId1 is deleted result1 = await indexerService.search({ @@ -2044,7 +2044,7 @@ test('should list doc ids work', async t => { // #region indexDoc() test('should index doc work', async t => { - const count = module.event.count('doc.indexer.updated'); + const count = module.queue.count('copilot.embedding.updateDoc'); const docSnapshot = await module.create(Mockers.DocSnapshot, { workspaceId: workspace.id, user, @@ -2110,7 +2110,7 @@ test('should index doc work', async t => { t.snapshot( result2.nodes.map(node => omit(node.fields, ['workspaceId', 'docId'])) ); - t.is(module.event.count('doc.indexer.updated'), count + 1); + t.is(module.queue.count('copilot.embedding.updateDoc'), count + 1); }); // #endregion diff --git a/packages/backend/server/src/plugins/indexer/index.ts b/packages/backend/server/src/plugins/indexer/index.ts index 4a012f8ca..a550d8413 100644 --- a/packages/backend/server/src/plugins/indexer/index.ts +++ b/packages/backend/server/src/plugins/indexer/index.ts @@ -27,16 +27,3 @@ export class IndexerModule {} export { IndexerService }; export type { SearchDoc } from './types'; - -declare global { - interface Events { - 'doc.indexer.updated': { - workspaceId: string; - docId: string; - }; - 'doc.indexer.deleted': { - workspaceId: string; - docId: string; - }; - } -} diff --git a/packages/backend/server/src/plugins/indexer/service.ts b/packages/backend/server/src/plugins/indexer/service.ts index a8c401369..798281c8b 100644 --- a/packages/backend/server/src/plugins/indexer/service.ts +++ b/packages/backend/server/src/plugins/indexer/service.ts @@ -2,8 +2,8 @@ import { Injectable, Logger } from '@nestjs/common'; import { camelCase, chunk, mapKeys, snakeCase } from 'lodash-es'; import { - EventBus, InvalidIndexerInput, + JobQueue, SearchProviderNotFound, } from '../../base'; import { readAllBlocksFromDocSnapshot } from '../../core/utils/blocksuite'; @@ -110,7 +110,7 @@ export class IndexerService { constructor( private readonly models: Models, private readonly factory: SearchProviderFactory, - private readonly event: EventBus + private readonly queue: JobQueue ) {} async createTables() { @@ -285,11 +285,12 @@ export class IndexerService { })), options ); - this.event.emit('doc.indexer.updated', { + + await this.queue.add('copilot.embedding.updateDoc', { workspaceId, docId, }); - this.logger.debug( + this.logger.log( `synced doc ${workspaceId}/${docId} with ${result.blocks.length} blocks` ); } @@ -319,12 +320,13 @@ export class IndexerService { }, options ); - this.logger.debug(`deleted doc ${workspaceId}/${docId}`); + await this.deleteBlocksByDocId(workspaceId, docId, options); - this.event.emit('doc.indexer.deleted', { + await this.queue.add('copilot.embedding.deleteDoc', { workspaceId, docId, }); + this.logger.log(`deleted doc ${workspaceId}/${docId}`); } async deleteBlocksByDocId(