diff --git a/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.test.ts b/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.test.ts index fff524ded..1acb8895b 100644 --- a/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.test.ts +++ b/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.test.ts @@ -4,6 +4,7 @@ import type {SearchableMessage} from '@fluxer/schema/src/contracts/search/Search import {describe, expect, it} from 'vitest'; import type {MeilisearchClient, MeilisearchTask} from './MeilisearchClient'; import {MeilisearchMessageAdapter} from './MeilisearchDomainAdapters'; +import {MEILISEARCH_MAX_TRACKED_BULK_TASKS} from './MeilisearchIndexAdapter'; interface RecordedMeilisearchRequest { method: string; @@ -35,6 +36,9 @@ class FakeMeilisearchClient implements MeilisearchClient { if (method === 'POST' && path.endsWith('/documents')) { return this.nextTask() as TResponse; } + if (method === 'POST' && path.endsWith('/documents/delete-batch')) { + return this.nextTask() as TResponse; + } if (method === 'POST' && path.endsWith('/search')) { return { hits: [{id: 'message-1'}], @@ -120,17 +124,66 @@ describe('MeilisearchMessageAdapter', () => { }); }); - it('waits for queued write tasks when refreshed', async () => { + it('waits for queued bulk index tasks when refreshed', async () => { const client = new FakeMeilisearchClient(); client.indexExists = true; const adapter = new MeilisearchMessageAdapter({client}); await adapter.initialize(); client.clear(); - await adapter.indexDocument({id: 'message-1'} as SearchableMessage); + await adapter.bulkIndexDocuments([{id: 'message-1'} as SearchableMessage]); expect(client.waitedTaskUids).toEqual([]); await adapter.refreshIndex(); expect(client.waitedTaskUids).toEqual([5]); }); + + it('does not retain task state for per-document writes', async () => { + const client = new FakeMeilisearchClient(); + client.indexExists = true; + const adapter = new MeilisearchMessageAdapter({client}); + await adapter.initialize(); + client.clear(); + + for (let index = 0; index < 500; index++) { + await adapter.indexDocument({id: `message-${index}`} as SearchableMessage); + await adapter.updateDocument({id: `message-${index}`} as SearchableMessage); + await adapter.deleteDocument(`message-${index}`); + } + + expect(client.requests).toHaveLength(1500); + await adapter.refreshIndex(); + expect(client.waitedTaskUids).toEqual([]); + }); + + it('bounds the tracked bulk task set', async () => { + const client = new FakeMeilisearchClient(); + client.indexExists = true; + const adapter = new MeilisearchMessageAdapter({client}); + await adapter.initialize(); + client.clear(); + + const batches = MEILISEARCH_MAX_TRACKED_BULK_TASKS + 100; + for (let index = 0; index < batches; index++) { + await adapter.bulkIndexDocuments([{id: `message-${index}`} as SearchableMessage]); + } + + await adapter.refreshIndex(); + expect(client.waitedTaskUids).toHaveLength(MEILISEARCH_MAX_TRACKED_BULK_TASKS); + expect(client.waitedTaskUids[0]).toBe(105); + }); + + it('drops tracked bulk tasks on shutdown', async () => { + const client = new FakeMeilisearchClient(); + client.indexExists = true; + const adapter = new MeilisearchMessageAdapter({client}); + await adapter.initialize(); + client.clear(); + + await adapter.bulkIndexDocuments([{id: 'message-1'} as SearchableMessage]); + await adapter.shutdown(); + await adapter.refreshIndex(); + + expect(client.waitedTaskUids).toEqual([]); + }); }); diff --git a/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.ts b/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.ts index 21771526a..0293d2391 100644 --- a/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.ts +++ b/fluxer_api/src/api/search/meilisearch/MeilisearchIndexAdapter.ts @@ -9,6 +9,8 @@ import type {MeilisearchIndexDefinition} from './MeilisearchIndexDefinitions'; const MAX_SEARCH_LIMIT = 1000; const MAX_TOTAL_HITS = 10000; +export const MEILISEARCH_MAX_TRACKED_BULK_TASKS = 4096; + interface MeilisearchIndexAdapterOptions { client: MeilisearchClient; index: MeilisearchIndexDefinition; @@ -37,7 +39,7 @@ export class MeilisearchIndexAdapter< protected readonly buildFilters: (filters: TFilters) => Array; protected readonly buildSort: ((filters: TFilters) => Array | undefined) | undefined; protected readonly buildQuery: ((query: string, filters: TFilters) => string) | undefined; - private readonly pendingTaskIds = new Set(); + private readonly pendingBulkTaskIds = new Set(); private initialized = false; constructor(options: MeilisearchIndexAdapterOptions) { @@ -56,7 +58,7 @@ export class MeilisearchIndexAdapter< uid, primaryKey: this.indexDefinition.primaryKey, }); - await this.waitForTask(task.taskUid); + await this.client.waitForTask(task.taskUid); } await Promise.all([ this.applySetting('PUT', 'searchable-attributes', this.indexDefinition.searchableAttributes), @@ -68,7 +70,7 @@ export class MeilisearchIndexAdapter< } async shutdown(): Promise { - this.pendingTaskIds.clear(); + this.pendingBulkTaskIds.clear(); this.initialized = false; } @@ -97,40 +99,41 @@ export class MeilisearchIndexAdapter< return; } this.assertInitialised(); - const task = await this.client.request( + await this.client.request( 'POST', `/indexes/${encodeURIComponent(this.indexDefinition.uid)}/documents/delete-batch`, ids, ); - this.trackTask(task.taskUid); } async deleteByFilter(filter: MeilisearchFilter): Promise { this.assertInitialised(); - const task = await this.client.request( + await this.client.request( 'POST', `/indexes/${encodeURIComponent(this.indexDefinition.uid)}/documents/delete`, {filter}, ); - this.trackTask(task.taskUid); } async deleteAllDocuments(): Promise { this.assertInitialised(); - const task = await this.client.request( + await this.client.request( 'DELETE', `/indexes/${encodeURIComponent(this.indexDefinition.uid)}/documents`, ); - this.trackTask(task.taskUid); } async bulkIndexDocuments(docs: Array): Promise { - await this.addDocuments(docs); + const taskUid = await this.addDocuments(docs); + if (taskUid === undefined) { + return; + } + this.trackBulkTask(taskUid); } async refreshIndex(): Promise { - const taskIds = Array.from(this.pendingTaskIds); - this.pendingTaskIds.clear(); + const taskIds = Array.from(this.pendingBulkTaskIds); + this.pendingBulkTaskIds.clear(); await Promise.all(taskIds.map((taskId) => this.client.waitForTask(taskId))); } @@ -193,12 +196,12 @@ export class MeilisearchIndexAdapter< `/indexes/${encodeURIComponent(this.indexDefinition.uid)}/settings/${setting}`, value, ); - await this.waitForTask(task.taskUid); + await this.client.waitForTask(task.taskUid); } - private async addDocuments(docs: Array): Promise { + private async addDocuments(docs: Array): Promise { if (docs.length === 0) { - return; + return undefined; } this.assertInitialised(); const task = await this.client.request( @@ -206,15 +209,16 @@ export class MeilisearchIndexAdapter< `/indexes/${encodeURIComponent(this.indexDefinition.uid)}/documents`, docs, ); - this.trackTask(task.taskUid); + return task.taskUid; } - private trackTask(taskUid: number): void { - this.pendingTaskIds.add(taskUid); - } - - private async waitForTask(taskUid: number): Promise { - this.pendingTaskIds.delete(taskUid); - await this.client.waitForTask(taskUid); + private trackBulkTask(taskUid: number): void { + if (this.pendingBulkTaskIds.size >= MEILISEARCH_MAX_TRACKED_BULK_TASKS) { + const oldest = this.pendingBulkTaskIds.values().next().value; + if (oldest !== undefined) { + this.pendingBulkTaskIds.delete(oldest); + } + } + this.pendingBulkTaskIds.add(taskUid); } }