From 74f6dc6f0d60906c4babf762fc8d14b7169c196c Mon Sep 17 00:00:00 2001 From: Sudhanva Bharadwaj BM <138512597+bmsvinci1729@users.noreply.github.com> Date: Thu, 30 Jul 2026 15:34:00 +0530 Subject: [PATCH] feat(ts-oss): add Qdrant server-side BM25 keywordSearch() + filter indexes (#5851) Co-authored-by: kartik-mem0 --- docs/components/vectordbs/dbs/qdrant.mdx | 6 + mem0-ts/src/oss/src/vector_stores/qdrant.ts | 183 +++++- .../oss/tests/qdrant-keyword-search.test.ts | 568 ++++++++++++++++++ 3 files changed, 739 insertions(+), 18 deletions(-) create mode 100644 mem0-ts/src/oss/tests/qdrant-keyword-search.test.ts diff --git a/docs/components/vectordbs/dbs/qdrant.mdx b/docs/components/vectordbs/dbs/qdrant.mdx index 555229e3d..75afd5da6 100644 --- a/docs/components/vectordbs/dbs/qdrant.mdx +++ b/docs/components/vectordbs/dbs/qdrant.mdx @@ -60,6 +60,12 @@ await memory.add(messages, { userId: "alice", metadata: { category: "movies" } } ``` +### Hybrid keyword search + +Mem0 blends semantic similarity with BM25 keyword scoring. On the TypeScript SDK, Qdrant computes the BM25 vectors server-side, which requires Qdrant 1.15.2 or newer with inference enabled. Qdrant Cloud enables inference by default only for clusters created after 2025-07-07; older clusters must activate it from the Cluster Detail page. The Python SDK encodes BM25 locally instead and needs the `fastembed` package, so scores are not numerically comparable between the two SDKs. + +When BM25 is unavailable, or when the collection was created before hybrid search was added, Mem0 logs a warning and falls back to semantic-only search. Writes are unaffected. To enable keyword scoring on an older collection, use a fresh collection name. + ### Config Let's see the available parameters for the `qdrant` config: diff --git a/mem0-ts/src/oss/src/vector_stores/qdrant.ts b/mem0-ts/src/oss/src/vector_stores/qdrant.ts index ddc5101b2..e84a138f1 100644 --- a/mem0-ts/src/oss/src/vector_stores/qdrant.ts +++ b/mem0-ts/src/oss/src/vector_stores/qdrant.ts @@ -4,6 +4,14 @@ import { SearchFilters, VectorStoreConfig, VectorStoreResult } from "../types"; import { loadPeer } from "../utils/load_peer"; import * as fs from "fs"; +// BM25 keyword search via Qdrant's built-in server-side inference (requires +// Qdrant >= 1.15.2). The Python adapter encodes BM25 client-side with fastembed; +// this server-side path avoids that dependency, but IDF weights, and therefore the +// scores, may differ between the two implementations. + +const BM25_VECTOR_NAME = "bm25"; +const BM25_MODEL = "Qdrant/bm25"; + interface QdrantConfig extends VectorStoreConfig { /** * Pre-configured QdrantClient instance. If using Qdrant Cloud, you must pass @@ -61,6 +69,10 @@ export class Qdrant implements VectorStore { private readonly collectionName: string; private dimension: number; private _initPromise?: Promise; + // Off for collections without the `bm25` slot (e.g. pre-hybrid-search). + private _hasBm25Slot = false; + // Payload indexes apply to a real server, not embedded/local mode. + private _isRemote = true; constructor(config: QdrantConfig) { this.config = config; @@ -73,6 +85,7 @@ export class Qdrant implements VectorStore { if (this.client) return; const config = this.config; if (config.client) { + // pre-configured client → treat as remote (mirrors Python is_local=False). this.client = config.client; return; } @@ -95,6 +108,7 @@ export class Qdrant implements VectorStore { params.port = config.port; } if (!Object.keys(params).length) { + this._isRemote = false; params.path = config.path; if (!config.onDisk && config.path) { if ( @@ -283,19 +297,98 @@ export class Qdrant implements VectorStore { payloads: Record[], ): Promise { await this.initialize(); - const points = vectors.map((vector, idx) => ({ - id: ids[idx], - vector: vector, - payload: payloads[idx] || {}, - })); - - await this.client.upsert(this.collectionName, { - points, + const points = vectors.map((vector, idx) => { + const payload = payloads[idx] || {}; + return { + id: ids[idx], + vector: this.buildPointVector(vector, payload), + payload, + }; }); + + await this.upsertPoints(points); } - async keywordSearch(): Promise { - return null; + // A server can accept a bm25 `sparse_vectors` config yet be unable to run + // inference (Qdrant < 1.15.2, or a Cloud cluster created before 2025-07-07 + // where inference was never activated), so no version check catches it and + // the first upsert failure is the only signal. Retry dense-only rather than + // lose the write. + // + // `as any`: the JS client types predate server-side sparse_vectors, so the + // named-vector and BM25-inference shapes are not typed. + private async upsertPoints( + points: { id: string | number; vector: any; payload?: any }[], + ): Promise { + try { + await this.client.upsert(this.collectionName, { points: points as any }); + } catch (error) { + if (!this._hasBm25Slot) { + throw error; + } + this._hasBm25Slot = false; + console.warn( + `Qdrant rejected the server-side BM25 vector for collection '${this.collectionName}'; ` + + "disabling hybrid keyword search. This requires Qdrant >= 1.15.2 with inference enabled. " + + "Retrying the write with a plain dense vector.", + ); + const denseOnly = points.map((p) => ({ + ...p, + vector: Array.isArray(p.vector) ? p.vector : p.vector[""], + })); + await this.client.upsert(this.collectionName, { + points: denseOnly as any, + }); + } + } + + // With a bm25 slot, return named vectors so Qdrant encodes BM25 server-side; + // otherwise return the plain dense vector (legacy behavior). + private buildPointVector( + vector: number[], + payload: Record, + ): number[] | Record { + if (!this._hasBm25Slot) { + return vector; + } + const named: Record = { "": vector }; + const text = payload?.textLemmatized || payload?.data || ""; + if (text) { + named[BM25_VECTOR_NAME] = { text, model: BM25_MODEL }; + } + return named; + } + + // BM25 keyword search; returns null (semantic-only fallback) when there is no + // bm25 slot or the query fails. + async keywordSearch( + query: string, + topK: number = 5, + filters?: SearchFilters, + ): Promise { + if (!this._hasBm25Slot) { + return null; + } + + try { + const queryFilter = this.createFilter(filters); + const response = await this.client.query(this.collectionName, { + query: { text: query, model: BM25_MODEL } as any, + using: BM25_VECTOR_NAME, + filter: queryFilter, + limit: topK, + with_payload: true, + }); + + return response.points.map((point) => ({ + id: String(point.id), + payload: (point.payload as Record) || {}, + score: point.score, + })); + } catch (error) { + console.error("Error during Qdrant keyword search:", error); + return null; + } } async search( @@ -341,13 +434,12 @@ export class Qdrant implements VectorStore { await this.initialize(); const point = { id: vectorId, - vector: vector, + // Re-encode BM25 so edited memories don't keep a stale sparse vector. + vector: this.buildPointVector(vector, payload), payload, }; - await this.client.upsert(this.collectionName, { - points: [point], - }); + await this.upsertPoints([point]); } async delete(vectorId: string): Promise { @@ -463,14 +555,31 @@ export class Qdrant implements VectorStore { } } - private async ensureCollection(name: string, size: number): Promise { + private async ensureCollection( + name: string, + size: number, + enableBm25: boolean = false, + ): Promise { try { - await this.client.createCollection(name, { + const createParams: Record = { vectors: { size, distance: "Cosine", }, - }); + }; + if (enableBm25) { + // `idf` lets Qdrant compute IDF over the live corpus at query time. + createParams.sparse_vectors = { + [BM25_VECTOR_NAME]: { modifier: "idf" }, + }; + } + await this.client.createCollection(name, createParams as any); + if (enableBm25) { + this._hasBm25Slot = true; + } + if (name === this.collectionName) { + await this.createFilterIndexes(name); + } } catch (error: any) { if ( error?.status === 409 || @@ -489,6 +598,22 @@ export class Qdrant implements VectorStore { `Expected: ${size}, got: ${vectorConfig.size}`, ); } + + if (enableBm25) { + // Existing collection: enable BM25 only if the slot is present. + const sparseConfig = (collectionInfo.config?.params as any) + ?.sparse_vectors; + this._hasBm25Slot = !!( + sparseConfig && BM25_VECTOR_NAME in sparseConfig + ); + if (!this._hasBm25Slot) { + console.warn( + `Collection '${name}' predates hybrid search (no '${BM25_VECTOR_NAME}' sparse slot). ` + + "BM25 keyword scoring is disabled for this collection; semantic search works normally. " + + "Use a fresh collection to enable hybrid keyword search.", + ); + } + } } catch (verifyError: any) { // Re-throw dimension mismatch errors if (verifyError?.message?.includes("wrong vector size")) { @@ -500,6 +625,8 @@ export class Qdrant implements VectorStore { `Collection '${name}' exists (409) but dimension verification failed: ${verifyError?.message || verifyError}. Proceeding anyway.`, ); } + // Ensure filter indexes exist even for pre-existing collections. + await this.createFilterIndexes(name); } // Otherwise collection exists and is fine — proceed } else { @@ -508,6 +635,26 @@ export class Qdrant implements VectorStore { } } + // Index the fields mem0 filters by; remote Qdrant rejects filtering on + // un-indexed fields. Mirrors Python's `_create_filter_indexes`. + private async createFilterIndexes(name: string): Promise { + if (!this._isRemote) { + return; + } + const commonFields = ["user_id", "agent_id", "run_id", "actor_id"]; + for (const field of commonFields) { + try { + await this.client.createPayloadIndex(name, { + field_name: field, + field_schema: "keyword", + }); + } catch (err) { + // Non-fatal: index likely already exists, or the server rejected it. + console.debug(`Qdrant: skipped payload index for '${field}':`, err); + } + } + } + async initialize(): Promise { if (!this._initPromise) { this._initPromise = this._doInitialize(); @@ -518,7 +665,7 @@ export class Qdrant implements VectorStore { private async _doInitialize(): Promise { try { await this.ensureClient(); - await this.ensureCollection(this.collectionName, this.dimension); + await this.ensureCollection(this.collectionName, this.dimension, true); await this.ensureCollection("memory_migrations", 1); } catch (error) { console.error("Error initializing Qdrant:", error); diff --git a/mem0-ts/src/oss/tests/qdrant-keyword-search.test.ts b/mem0-ts/src/oss/tests/qdrant-keyword-search.test.ts new file mode 100644 index 000000000..2746238cc --- /dev/null +++ b/mem0-ts/src/oss/tests/qdrant-keyword-search.test.ts @@ -0,0 +1,568 @@ +/// +/** + * Tests for Qdrant native BM25 keyword search (server-side `Qdrant/bm25` + * inference). Fully mocked — no real Qdrant server is required. + * + * Covers both the happy path and the failure modes that matter for an opt-in, + * capability-gated feature: + * - fresh collections get the `bm25` sparse slot (modifier idf); migrations don't, + * - insert/update attach the server-side BM25 inference vector, + * - keywordSearch queries the `bm25` slot and maps results to the common shape, + * - collection-init variations (slot present/absent, wrong size, 401/403/409, + * transient verify failure, fatal error), + * - malformed/hostile query responses and filter construction, + * all proving the adapter degrades to `null` (semantic-only) instead of crashing. + */ + +jest.setTimeout(15000); + +type MockClient = Record; + +function buildMockClient( + overrides: Partial = {}, + collectionInfo: any = { config: { params: { vectors: { size: 768 } } } }, +): MockClient { + return { + createCollection: jest.fn().mockResolvedValue(undefined), + createPayloadIndex: jest.fn().mockResolvedValue(undefined), + getCollection: jest.fn().mockResolvedValue(collectionInfo), + scroll: jest.fn().mockResolvedValue({ points: [] }), + upsert: jest.fn().mockResolvedValue(undefined), + retrieve: jest.fn().mockResolvedValue([]), + search: jest.fn().mockResolvedValue([]), + query: jest.fn().mockResolvedValue({ points: [] }), + delete: jest.fn().mockResolvedValue(undefined), + deleteCollection: jest.fn().mockResolvedValue(undefined), + ...overrides, + }; +} + +jest.mock("@qdrant/js-client-rest", () => ({ + QdrantClient: jest.fn().mockImplementation(() => buildMockClient()), +})); + +import { Qdrant } from "../src/vector_stores/qdrant"; + +const BASE_CONFIG = { + collectionName: "test_memories", + embeddingModelDims: 768, + dimension: 768, +}; + +// Existing collection that already carries a bm25 slot. +const INFO_WITH_SLOT = { + config: { + params: { vectors: { size: 768 }, sparse_vectors: { bm25: {} } }, + }, +}; +// Existing collection without a bm25 slot (legacy / pre-hybrid-search). +const INFO_NO_SLOT = { + config: { params: { vectors: { size: 768 } } }, +}; + +let warnSpy: jest.SpyInstance; +let errSpy: jest.SpyInstance; + +beforeEach(() => { + jest.clearAllMocks(); + warnSpy = jest.spyOn(console, "warn").mockImplementation(() => {}); + errSpy = jest.spyOn(console, "error").mockImplementation(() => {}); +}); + +afterEach(() => { + warnSpy.mockRestore(); + errSpy.mockRestore(); +}); + +function makeStore(client: MockClient): Qdrant { + return new Qdrant({ client: client as any, ...BASE_CONFIG }); +} + +describe("Qdrant BM25 keyword search — collection setup", () => { + it("creates the main collection with the bm25 sparse slot (modifier idf)", async () => { + const client = buildMockClient(); + const store = makeStore(client); + await store.initialize(); + + const createCall = client.createCollection.mock.calls.find( + (c) => c[0] === "test_memories", + ); + expect(createCall).toBeDefined(); + expect(createCall![1].sparse_vectors).toEqual({ + bm25: { modifier: "idf" }, + }); + }); + + it("does NOT add the bm25 slot to the memory_migrations collection", async () => { + const client = buildMockClient(); + const store = makeStore(client); + await store.initialize(); + + const migrationsCall = client.createCollection.mock.calls.find( + (c) => c[0] === "memory_migrations", + ); + expect(migrationsCall).toBeDefined(); + expect(migrationsCall![1].sparse_vectors).toBeUndefined(); + }); + + it("409 with an existing bm25 slot → enables keyword search, no warning", async () => { + const client = buildMockClient( + { createCollection: jest.fn().mockRejectedValue({ status: 409 }) }, + INFO_WITH_SLOT, + ); + const store = makeStore(client); + await store.initialize(); + + expect(await store.keywordSearch("anything")).not.toBeNull(); + expect(client.query).toHaveBeenCalled(); + expect(warnSpy).not.toHaveBeenCalled(); + }); + + it("409 without a bm25 slot → disabled + warns once", async () => { + const client = buildMockClient( + { createCollection: jest.fn().mockRejectedValue({ status: 409 }) }, + INFO_NO_SLOT, + ); + const store = makeStore(client); + await store.initialize(); + + expect(await store.keywordSearch("anything")).toBeNull(); + expect(client.query).not.toHaveBeenCalled(); + expect(warnSpy).toHaveBeenCalledTimes(1); + }); + + it("401 (auth-restricted existing collection) with slot → enabled", async () => { + const client = buildMockClient( + { createCollection: jest.fn().mockRejectedValue({ status: 401 }) }, + INFO_WITH_SLOT, + ); + const store = makeStore(client); + await store.initialize(); + expect(await store.keywordSearch("x")).not.toBeNull(); + }); + + it("403 with no slot → disabled (graceful)", async () => { + const client = buildMockClient( + { createCollection: jest.fn().mockRejectedValue({ status: 403 }) }, + INFO_NO_SLOT, + ); + const store = makeStore(client); + await store.initialize(); + expect(await store.keywordSearch("x")).toBeNull(); + }); + + it("existing collection with WRONG vector size → init rejects (size guard intact)", async () => { + const client = buildMockClient( + { createCollection: jest.fn().mockRejectedValue({ status: 409 }) }, + { config: { params: { vectors: { size: 1536 } } } }, // mismatch vs 768 + ); + const store = makeStore(client); + await expect(store.initialize()).rejects.toThrow(/wrong vector size/i); + }); + + it("transient getCollection failure during verify → does NOT crash init, stays disabled", async () => { + const client = buildMockClient({ + createCollection: jest.fn().mockRejectedValue({ status: 409 }), + getCollection: jest + .fn() + .mockRejectedValue({ status: 500, message: "committing" }), + }); + const store = makeStore(client); + await expect(store.initialize()).resolves.toBeUndefined(); + expect(await store.keywordSearch("x")).toBeNull(); + }); + + it("fatal createCollection error (500, not 401/403/409) → init rejects", async () => { + const client = buildMockClient({ + createCollection: jest.fn().mockRejectedValue({ status: 500 }), + }); + const store = makeStore(client); + await expect(store.initialize()).rejects.toBeDefined(); + }); +}); + +describe("Qdrant BM25 keyword search — insert / update", () => { + async function freshStore(): Promise<{ store: Qdrant; client: MockClient }> { + const client = buildMockClient(); + const store = makeStore(client); + await store.initialize(); + return { store, client }; + } + + it("attaches a server-side BM25 inference vector on insert", async () => { + const { store, client } = await freshStore(); + await store.insert( + [[0.1, 0.2, 0.3]], + ["mem-1"], + [ + { + data: "Alice reported a duplicate AWS invoice", + textLemmatized: "alice report duplicate aws invoice", + }, + ], + ); + + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.id).toBe("mem-1"); + expect(point.vector[""]).toEqual([0.1, 0.2, 0.3]); + expect(point.vector.bm25).toEqual({ + text: "alice report duplicate aws invoice", + model: "Qdrant/bm25", + }); + }); + + it("falls back to payload.data when textLemmatized is absent", async () => { + const { store, client } = await freshStore(); + await store.insert([[0.1]], ["mem-2"], [{ data: "plain text" }]); + + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.vector.bm25).toEqual({ + text: "plain text", + model: "Qdrant/bm25", + }); + }); + + it("empty textLemmatized falls through to data", async () => { + const { store, client } = await freshStore(); + await store.insert( + [[0.1]], + ["m"], + [{ textLemmatized: "", data: "fallback" }], + ); + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.vector.bm25).toEqual({ + text: "fallback", + model: "Qdrant/bm25", + }); + }); + + it("omits the bm25 vector when the point has no text (dense still present)", async () => { + const { store, client } = await freshStore(); + await store.insert([[0.1, 0.2]], ["m"], [{ textLemmatized: "", data: "" }]); + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.vector[""]).toEqual([0.1, 0.2]); + expect(point.vector.bm25).toBeUndefined(); + }); + + it("preserves the given id on insert", async () => { + const { store, client } = await freshStore(); + await store.insert([[0.1]], ["123"], [{ data: "x" }]); + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.id).toBe("123"); + }); + + it("batch insert attaches bm25 to every text-bearing point", async () => { + const { store, client } = await freshStore(); + await store.insert( + [[0.1], [0.2], [0.3]], + ["a", "b", "c"], + [{ data: "one" }, { data: "" }, { textLemmatized: "three" }], + ); + const points = client.upsert.mock.calls.at(-1)![1].points; + expect(points[0].vector.bm25.text).toBe("one"); + expect(points[1].vector.bm25).toBeUndefined(); + expect(points[2].vector.bm25.text).toBe("three"); + }); + + it("re-encodes the bm25 vector on update", async () => { + const { store, client } = await freshStore(); + await store.update("mem-1", [0.4, 0.5], { + data: "updated", + textLemmatized: "updat", + }); + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.vector[""]).toEqual([0.4, 0.5]); + expect(point.vector.bm25).toEqual({ text: "updat", model: "Qdrant/bm25" }); + }); + + it("update with no text → named dense only (bm25 not re-attached)", async () => { + const { store, client } = await freshStore(); + await store.update("m", [0.9], { user_id: "u" }); + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.vector[""]).toEqual([0.9]); + expect(point.vector.bm25).toBeUndefined(); + }); + + it("inserts a plain dense vector (no named/sparse) on legacy collections", async () => { + const client = buildMockClient( + { createCollection: jest.fn().mockRejectedValue({ status: 409 }) }, + INFO_NO_SLOT, + ); + const store = makeStore(client); + await store.initialize(); + + await store.insert( + [[0.1, 0.2]], + ["mem-1"], + [{ data: "hello", textLemmatized: "hello" }], + ); + + const point = client.upsert.mock.calls.at(-1)![1].points[0]; + expect(point.vector).toEqual([0.1, 0.2]); + }); +}); + +describe("Qdrant BM25 keyword search — querying & result mapping", () => { + it("queries the bm25 slot and maps results to {id, payload, score}", async () => { + const client = buildMockClient({ + query: jest.fn().mockResolvedValue({ + points: [ + { id: "mem-1", score: 4.2, payload: { data: "AWS invoice" } }, + { id: "mem-2", score: 1.1, payload: { data: "other" } }, + ], + }), + }); + const store = makeStore(client); + await store.initialize(); + + const results = await store.keywordSearch("aws invoice", 5, { + user_id: "u1", + }); + + expect(client.query).toHaveBeenCalledWith( + "test_memories", + expect.objectContaining({ + query: { text: "aws invoice", model: "Qdrant/bm25" }, + using: "bm25", + limit: 5, + with_payload: true, + }), + ); + expect(results).toEqual([ + { id: "mem-1", payload: { data: "AWS invoice" }, score: 4.2 }, + { id: "mem-2", payload: { data: "other" }, score: 1.1 }, + ]); + }); + + async function storeWith(queryImpl: jest.Mock): Promise { + const client = buildMockClient({ query: queryImpl }); + const store = makeStore(client); + await store.initialize(); + return store; + } + + it("query rejects → null (no throw)", async () => { + const store = await storeWith( + jest.fn().mockRejectedValue(new Error("boom")), + ); + expect(await store.keywordSearch("x")).toBeNull(); + expect(errSpy).toHaveBeenCalled(); + }); + + it("response missing the `points` array → null (no throw)", async () => { + const store = await storeWith(jest.fn().mockResolvedValue({})); + expect(await store.keywordSearch("x")).toBeNull(); + expect(errSpy).toHaveBeenCalled(); + }); + + it("response is null → null (no throw)", async () => { + const store = await storeWith(jest.fn().mockResolvedValue(null)); + expect(await store.keywordSearch("x")).toBeNull(); + }); + + it("point with numeric id → stringified", async () => { + const store = await storeWith( + jest.fn().mockResolvedValue({ + points: [{ id: 42, score: 1.0, payload: { data: "d" } }], + }), + ); + const res = await store.keywordSearch("x"); + expect(res![0].id).toBe("42"); + }); + + it("point with missing score → score undefined (pipeline coerces to 0)", async () => { + const store = await storeWith( + jest.fn().mockResolvedValue({ points: [{ id: "a", payload: {} }] }), + ); + const res = await store.keywordSearch("x"); + expect(res![0].score).toBeUndefined(); + }); + + it("point with missing payload → {}", async () => { + const store = await storeWith( + jest.fn().mockResolvedValue({ points: [{ id: "a", score: 2 }] }), + ); + const res = await store.keywordSearch("x"); + expect(res![0].payload).toEqual({}); + }); + + it("empty query string is forwarded and yields [] (graceful)", async () => { + const queryImpl = jest.fn().mockResolvedValue({ points: [] }); + const store = await storeWith(queryImpl); + const res = await store.keywordSearch(""); + expect(res).toEqual([]); + expect(queryImpl.mock.calls[0][1].query).toEqual({ + text: "", + model: "Qdrant/bm25", + }); + }); +}); + +describe("Qdrant BM25 keyword search — query construction & filters", () => { + async function freshStore(): Promise<{ store: Qdrant; client: MockClient }> { + const client = buildMockClient(); + const store = makeStore(client); + await store.initialize(); + return { store, client }; + } + + it("defaults topK to 5 when omitted", async () => { + const { store, client } = await freshStore(); + await store.keywordSearch("x"); + expect(client.query.mock.calls[0][1].limit).toBe(5); + }); + + it("passes a simple equality filter through createFilter", async () => { + const { store, client } = await freshStore(); + await store.keywordSearch("x", 3, { user_id: "u1" }); + const arg = client.query.mock.calls[0][1]; + expect(arg.limit).toBe(3); + expect(arg.using).toBe("bm25"); + expect(arg.filter).toEqual({ + must: [{ key: "user_id", match: { value: "u1" } }], + }); + }); + + it("builds an OR (should) clause for logical filters", async () => { + const { store, client } = await freshStore(); + await store.keywordSearch("x", 5, { + $or: [{ user_id: "u1" }, { agent_id: "a1" }], + }); + const filter = client.query.mock.calls[0][1].filter; + expect(filter.should).toHaveLength(2); + }); + + it("empty filter object → filter is undefined", async () => { + const { store, client } = await freshStore(); + await store.keywordSearch("x", 5, {}); + expect(client.query.mock.calls[0][1].filter).toBeUndefined(); + }); +}); + +describe("Qdrant BM25 keyword search — reactive downgrade on upsert failure", () => { + it("insert(): upsert fails once, retries with plain dense vector, disables bm25", async () => { + const upsert = jest + .fn() + .mockRejectedValueOnce({ + status: 500, + message: "InferenceService is not initialized.", + }) + .mockResolvedValueOnce(undefined); + const client = buildMockClient({ upsert }); + const store = makeStore(client); + await store.initialize(); + + await store.insert( + [[0.1, 0.2]], + ["mem-1"], + [{ data: "hello", textLemmatized: "hello" }], + ); + + expect(upsert).toHaveBeenCalledTimes(2); + const retryPoint = upsert.mock.calls[1][1].points[0]; + expect(retryPoint.vector).toEqual([0.1, 0.2]); + expect(warnSpy).toHaveBeenCalledTimes(1); + expect(await store.keywordSearch("hello")).toBeNull(); + }); + + it("update(): upsert fails once, retries with plain dense vector, disables bm25", async () => { + const upsert = jest + .fn() + .mockRejectedValueOnce({ + status: 500, + message: "InferenceService is not initialized.", + }) + .mockResolvedValueOnce(undefined); + const client = buildMockClient({ upsert }); + const store = makeStore(client); + await store.initialize(); + + await store.update("mem-1", [0.4, 0.5], { + data: "updated", + textLemmatized: "updat", + }); + + expect(upsert).toHaveBeenCalledTimes(2); + const retryPoint = upsert.mock.calls[1][1].points[0]; + expect(retryPoint.vector).toEqual([0.4, 0.5]); + expect(warnSpy).toHaveBeenCalledTimes(1); + expect(await store.keywordSearch("updat")).toBeNull(); + }); + + it("legacy collection (no bm25 slot): a genuine upsert failure still rejects, not swallowed", async () => { + const client = buildMockClient( + { + createCollection: jest.fn().mockRejectedValue({ status: 409 }), + upsert: jest.fn().mockRejectedValue(new Error("network down")), + }, + INFO_NO_SLOT, + ); + const store = makeStore(client); + await store.initialize(); + + await expect(store.insert([[0.1]], ["m"], [{ data: "x" }])).rejects.toThrow( + "network down", + ); + expect(client.upsert).toHaveBeenCalledTimes(1); + }); + + it("after a downgrade, a subsequent insert() sends a plain dense vector on the first attempt", async () => { + const upsert = jest + .fn() + .mockRejectedValueOnce({ + status: 500, + message: "InferenceService is not initialized.", + }) + .mockResolvedValueOnce(undefined) + .mockResolvedValueOnce(undefined); + const client = buildMockClient({ upsert }); + const store = makeStore(client); + await store.initialize(); + + await store.insert([[0.1]], ["mem-1"], [{ data: "first" }]); + expect(upsert).toHaveBeenCalledTimes(2); + + await store.insert([[0.2]], ["mem-2"], [{ data: "second" }]); + expect(upsert).toHaveBeenCalledTimes(3); + const secondInsertPoint = upsert.mock.calls[2][1].points[0]; + expect(secondInsertPoint.vector).toEqual([0.2]); + }); +}); + +describe("Qdrant BM25 keyword search — memory_migrations indexing", () => { + it("does not create payload indexes for memory_migrations on the existing-collection path", async () => { + const client = buildMockClient({ + createCollection: jest.fn().mockRejectedValue({ status: 409 }), + }); + const store = makeStore(client); + await store.initialize(); + + const migrationsIndexCalls = client.createPayloadIndex.mock.calls.filter( + (c) => c[0] === "memory_migrations", + ); + expect(migrationsIndexCalls).toHaveLength(0); + }); +}); + +describe("Qdrant BM25 keyword search — semantic (dense) path", () => { + it("search() still targets the unnamed dense slot after BM25 init", async () => { + const client = buildMockClient({ + search: jest + .fn() + .mockResolvedValue([ + { id: "mem-1", score: 0.9, payload: { data: "x" } }, + ]), + }); + const store = makeStore(client); + await store.initialize(); + + const res = await store.search([0.1, 0.2, 0.3], 5); + + // Plain number[] (not a named vector) keeps hitting the default dense slot. + expect(client.search).toHaveBeenCalledWith( + "test_memories", + expect.objectContaining({ vector: [0.1, 0.2, 0.3] }), + ); + expect(res[0].id).toBe("mem-1"); + }); +});