diff --git a/docs/components/vectordbs/dbs/upstash-vector.mdx b/docs/components/vectordbs/dbs/upstash-vector.mdx index 6b20544b8..1858516f8 100644 --- a/docs/components/vectordbs/dbs/upstash-vector.mdx +++ b/docs/components/vectordbs/dbs/upstash-vector.mdx @@ -8,6 +8,10 @@ description: "Use Upstash Vector as a serverless vector database in Mem0 with op You can enable the built-in embedding models by setting `enable_embeddings` to `True`. This allows you to use Upstash's embedding models for vectorization. + + Server-side Upstash embeddings (`enable_embeddings`) are available in the Python SDK only. The TypeScript SDK always embeds text with your configured embedder before writing to Upstash, so use the external embedding provider setup below. + + ```python import os from mem0 import Memory @@ -34,7 +38,8 @@ m.add("Likes to play cricket on weekends", user_id="alice", metadata={"category" ### Usage with external embedding providers -```python + +```python Python import os from mem0 import Memory @@ -58,6 +63,36 @@ m = Memory.from_config(config) m.add("Likes to play cricket on weekends", user_id="alice", metadata={"category": "hobbies"}) ``` +```typescript TypeScript +import { Memory } from "mem0ai/oss"; + +// Set OPENAI_API_KEY, UPSTASH_VECTOR_REST_URL, and UPSTASH_VECTOR_REST_TOKEN in your environment. +const config = { + embedder: { + provider: "openai", + config: { + apiKey: process.env.OPENAI_API_KEY, + model: "text-embedding-3-large", + }, + }, + vectorStore: { + provider: "upstash_vector", + config: { + collectionName: "memories", + url: process.env.UPSTASH_VECTOR_REST_URL, + token: process.env.UPSTASH_VECTOR_REST_TOKEN, + }, + }, +}; + +const memory = new Memory(config); +await memory.add("Likes to play cricket on weekends", { + userId: "alice", + metadata: { category: "hobbies" }, +}); +``` + + ### Config Here are the parameters available for configuring Upstash Vector: @@ -74,3 +109,7 @@ Here are the parameters available for configuring Upstash Vector: When `url` and `token` are not provided, the `UPSTASH_VECTOR_REST_URL` and `UPSTASH_VECTOR_REST_TOKEN` environment variables are used. + + + The TypeScript SDK uses camelCase config keys (`collectionName`, `url`, `token`), where `collectionName` is required. Pass `url` and `token` (or a preconfigured `client`) explicitly, since the TypeScript SDK does not read them from environment variables. `enable_embeddings` is not supported in TypeScript. + diff --git a/mem0-ts/package.json b/mem0-ts/package.json index 04a0bcaf4..54c4d5822 100644 --- a/mem0-ts/package.json +++ b/mem0-ts/package.json @@ -122,6 +122,7 @@ "@turbopuffer/turbopuffer": "^2.0.0", "@types/jest": "29.5.14", "@types/pg": "8.11.0", + "@upstash/vector": "^1.2.3", "better-sqlite3": "^12.6.2", "cassandra-driver": "4.8.0", "cloudflare": "^4.2.0", diff --git a/mem0-ts/pnpm-lock.yaml b/mem0-ts/pnpm-lock.yaml index 5fe30d89c..bafa6ba9f 100644 --- a/mem0-ts/pnpm-lock.yaml +++ b/mem0-ts/pnpm-lock.yaml @@ -74,6 +74,9 @@ importers: '@types/pg': specifier: 8.11.0 version: 8.11.0 + '@upstash/vector': + specifier: ^1.2.3 + version: 1.2.3 axios: specifier: ^1.16.0 version: 1.17.0 @@ -1366,6 +1369,9 @@ packages: resolution: {integrity: sha512-jIXhD0eWQ1JA6ln/5Dltyx22UxWNrw0hZmhy2rlv6m6KgF7kplHx3g0fzi09lNmTJQRR91OlemYp3xFnvDK9og==} engines: {node: '>=20.0.0'} + '@upstash/vector@1.2.3': + resolution: {integrity: sha512-yXsWKeuHNYyH72BcSZd3bV5ZD5MybAoTvKxkMaeV2UzuGfNzbHBVh5eO+ysTWTFAf8I9XcOueF4tZfAGjCa4Iw==} + abort-controller@3.0.0: resolution: {integrity: sha512-h8lQ8tacZYnR3vNQTgibj+tODHI5/+l06Au2Pcriv/Gmet0eaj4TwWH41sO9wnHDiQsEj19q0drzdWdeAHtweg==} engines: {node: '>=6.5'} @@ -5234,6 +5240,8 @@ snapshots: transitivePeerDependencies: - supports-color + '@upstash/vector@1.2.3': {} + abort-controller@3.0.0: dependencies: event-target-shim: 5.0.1 diff --git a/mem0-ts/src/oss/package.json b/mem0-ts/src/oss/package.json index 3ea60e431..7d97a894e 100644 --- a/mem0-ts/src/oss/package.json +++ b/mem0-ts/src/oss/package.json @@ -16,6 +16,7 @@ "@anthropic-ai/sdk": "^0.18.0", "@google/genai": "^0.7.0", "@qdrant/js-client-rest": "^1.13.0", + "@upstash/vector": "^1.2.3", "@types/node": "^20.11.19", "@types/pg": "^8.11.0", "@types/redis": "^4.0.10", diff --git a/mem0-ts/src/oss/src/index.ts b/mem0-ts/src/oss/src/index.ts index e3e15ab40..e7f7b2134 100644 --- a/mem0-ts/src/oss/src/index.ts +++ b/mem0-ts/src/oss/src/index.ts @@ -32,8 +32,11 @@ export * from "./vector_stores/langchain"; export * from "./vector_stores/vectorize"; export * from "./vector_stores/azure_ai_search"; export * from "./vector_stores/pgvector"; +export * from "./vector_stores/upstash_vector"; export * from "./vector_stores/azure_mysql"; export * from "./vector_stores/cassandra"; export * from "./vector_stores/s3_vectors"; export * from "./vector_stores/vertex_ai_vector_search"; +export * from "./vector_stores/pinecone"; +export * from "./vector_stores/turbopuffer"; export * from "./utils/factory"; diff --git a/mem0-ts/src/oss/src/utils/factory.ts b/mem0-ts/src/oss/src/utils/factory.ts index 79197c0d7..1a9526f74 100644 --- a/mem0-ts/src/oss/src/utils/factory.ts +++ b/mem0-ts/src/oss/src/utils/factory.ts @@ -42,6 +42,7 @@ import { LangchainEmbedder } from "../embeddings/langchain"; import { LangchainVectorStore } from "../vector_stores/langchain"; import { AzureAISearch } from "../vector_stores/azure_ai_search"; import { PGVector } from "../vector_stores/pgvector"; +import { UpstashVector } from "../vector_stores/upstash_vector"; import { AzureMySQLDB } from "../vector_stores/azure_mysql"; import { VertexAIVectorSearch } from "../vector_stores/vertex_ai_vector_search"; import { CassandraDB } from "../vector_stores/cassandra"; @@ -136,6 +137,8 @@ export class VectorStoreFactory { return new VertexAIVectorSearch(config as any); case "pgvector": return new PGVector(config as any); + case "upstash_vector": + return new UpstashVector(config as any); case "azure_mysql": return new AzureMySQLDB(config as any); case "cassandra": diff --git a/mem0-ts/src/oss/src/vector_stores/upstash_vector.ts b/mem0-ts/src/oss/src/vector_stores/upstash_vector.ts new file mode 100644 index 000000000..4de38e2b2 --- /dev/null +++ b/mem0-ts/src/oss/src/vector_stores/upstash_vector.ts @@ -0,0 +1,240 @@ +import { Index, QueryResult, Vector } from "@upstash/vector"; +import { VectorStore } from "./base"; +import { SearchFilters, VectorStoreConfig, VectorStoreResult } from "../types"; + +interface UpstashVectorConfig extends VectorStoreConfig { + collectionName: string; + url?: string; + token?: string; + client?: Index>; +} + +type UpstashMetadata = Record; + +export class UpstashVector implements VectorStore { + private readonly client: Index; + private readonly collectionName: string; + + constructor(config: UpstashVectorConfig) { + if (!config.collectionName) { + throw new Error("collectionName is required for Upstash Vector."); + } + + if (config.client) { + this.client = config.client; + } else if (config.url && config.token) { + this.client = new Index({ + url: config.url, + token: config.token, + }); + } else { + throw new Error("Either a client or url and token must be provided."); + } + + this.collectionName = config.collectionName; + } + + async initialize(): Promise { + return; + } + + async insert( + vectors: number[][], + ids: string[], + payloads: Record[], + ): Promise { + const upsertData = vectors.map((vector, idx) => { + return { + id: ids[idx], + vector, + metadata: payloads[idx] ?? {}, + }; + }); + + await this.client.upsert(upsertData, { namespace: this.collectionName }); + } + + async search( + query: number[], + topK: number = 5, + filters?: SearchFilters, + ): Promise { + const response = await this.client.query( + { + vector: query, + topK, + filter: this.convertFilters(filters), + includeMetadata: true, + }, + { namespace: this.collectionName }, + ); + + return response.map((result) => this.parseResult(result)); + } + + async keywordSearch( + query: string, + topK: number = 5, + filters?: SearchFilters, + ): Promise { + try { + const response = await this.client.query( + { + data: query, + topK, + filter: this.convertFilters(filters), + includeMetadata: true, + }, + { namespace: this.collectionName }, + ); + + return response.map((result) => this.parseResult(result)); + } catch (error) { + console.error(`Error during keyword search for query '${query}':`, error); + return null; + } + } + + async get(vectorId: string): Promise { + const response = await this.client.fetch([vectorId], { + includeMetadata: true, + namespace: this.collectionName, + }); + const vector = response[0]; + + if (!vector) { + return null; + } + + return { + id: String(vector.id), + payload: (vector.metadata ?? {}) as Record, + }; + } + + async update( + vectorId: string, + vector: number[], + payload: Record, + ): Promise { + // Upstash's `update` can't set the vector and metadata in one call (its + // payload is a discriminated union of vector | data | metadata), so a + // single `upsert` replaces both atomically, the same way insert() writes. + await this.client.upsert( + { + id: vectorId, + vector, + metadata: payload, + }, + { namespace: this.collectionName }, + ); + } + + async delete(vectorId: string): Promise { + await this.client.delete(vectorId, { namespace: this.collectionName }); + } + + async deleteCol(): Promise { + await this.client.reset({ namespace: this.collectionName }); + } + + async list( + filters?: SearchFilters, + topK: number = 100, + ): Promise<[VectorStoreResult[], number]> { + const results: VectorStoreResult[] = []; + let cursor = "0"; + + do { + const response = await this.client.range( + { + cursor, + limit: Math.min(100, topK - results.length), + includeMetadata: true, + }, + { namespace: this.collectionName }, + ); + + for (const vector of response.vectors) { + if (this.matchesFilters(vector, filters)) { + results.push({ + id: String(vector.id), + payload: (vector.metadata ?? {}) as Record, + }); + } + + if (results.length >= topK) { + break; + } + } + + cursor = response.nextCursor; + // Upstash returns an empty-string cursor once the scan is exhausted (it + // never comes back as "0"), so "" is the termination sentinel. Checking + // for "0" here would re-scan from the start and return duplicates. + } while (cursor !== "" && results.length < topK); + + return [results, results.length]; + } + + async getUserId(): Promise { + return "anonymous-upstash-vector"; + } + + async setUserId(): Promise { + return; + } + + async reset(): Promise { + await this.deleteCol(); + } + + private parseResult(result: QueryResult): VectorStoreResult { + return { + id: String(result.id), + payload: (result.metadata ?? {}) as Record, + score: result.score, + }; + } + + private stringifyFilterValue(value: unknown): string { + if (typeof value === "string") { + return JSON.stringify(value); + } + + if (typeof value === "boolean") { + return value ? "true" : "false"; + } + + return String(value); + } + + private convertFilters(filters?: SearchFilters): string | undefined { + if (!filters) { + return undefined; + } + + const expressions = Object.entries(filters) + .filter(([, value]) => value !== undefined && value !== null) + .map(([key, value]) => `${key} = ${this.stringifyFilterValue(value)}`); + + return expressions.length > 0 ? expressions.join(" AND ") : undefined; + } + + private matchesFilters( + vector: Vector, + filters?: SearchFilters, + ): boolean { + if (!filters) { + return true; + } + + return Object.entries(filters).every(([key, value]) => { + if (value === undefined || value === null) { + return true; + } + + return vector.metadata?.[key] === value; + }); + } +} diff --git a/mem0-ts/src/oss/tests/factory.unit.test.ts b/mem0-ts/src/oss/tests/factory.unit.test.ts index 72b18c4ee..0b4cdeeb1 100644 --- a/mem0-ts/src/oss/tests/factory.unit.test.ts +++ b/mem0-ts/src/oss/tests/factory.unit.test.ts @@ -158,6 +158,11 @@ jest.mock("../src/vector_stores/pgvector", () => ({ .fn() .mockImplementation((config) => ({ type: "pgvector", config })), })); +jest.mock("../src/vector_stores/upstash_vector", () => ({ + UpstashVector: jest + .fn() + .mockImplementation((config) => ({ type: "upstash-vector", config })), +})); jest.mock("../src/vector_stores/azure_mysql", () => ({ AzureMySQLDB: jest .fn() @@ -299,6 +304,7 @@ describe("VectorStoreFactory", () => { ["vectorize"], ["azure-ai-search"], ["pgvector"], + ["upstash_vector"], ["azure_mysql"], ["cassandra"], ["s3-vectors"], diff --git a/mem0-ts/src/oss/tests/upstash_vector.unit.test.ts b/mem0-ts/src/oss/tests/upstash_vector.unit.test.ts new file mode 100644 index 000000000..b740ab724 --- /dev/null +++ b/mem0-ts/src/oss/tests/upstash_vector.unit.test.ts @@ -0,0 +1,166 @@ +import { UpstashVector } from "../src/vector_stores/upstash_vector"; + +describe("UpstashVector", () => { + const namespace = "memories"; + + function createClient(overrides: Record = {}) { + return { + upsert: jest.fn().mockResolvedValue("Success"), + query: jest.fn().mockResolvedValue([]), + fetch: jest.fn().mockResolvedValue([]), + update: jest.fn().mockResolvedValue({ updated: 1 }), + delete: jest.fn().mockResolvedValue({ deleted: 1 }), + reset: jest.fn().mockResolvedValue("Success"), + range: jest.fn().mockResolvedValue({ vectors: [], nextCursor: "" }), + ...overrides, + }; + } + + it("upserts vectors into the collection namespace", async () => { + const client = createClient(); + const store = new UpstashVector({ + collectionName: namespace, + client: client as any, + }); + + await store.insert( + [[0.1, 0.2]], + ["memory-1"], + [{ data: "hello", user_id: "user-1" }], + ); + + expect(client.upsert).toHaveBeenCalledWith( + [ + { + id: "memory-1", + vector: [0.1, 0.2], + metadata: { data: "hello", user_id: "user-1" }, + }, + ], + { namespace }, + ); + }); + + it("queries vectors with converted filters", async () => { + const client = createClient({ + query: jest.fn().mockResolvedValue([ + { + id: "memory-1", + score: 0.9, + metadata: { data: "hello", user_id: "user-1" }, + }, + ]), + }); + const store = new UpstashVector({ + collectionName: namespace, + client: client as any, + }); + + const results = await store.search([0.1, 0.2], 3, { + user_id: "user-1", + active: true, + }); + + expect(client.query).toHaveBeenCalledWith( + { + vector: [0.1, 0.2], + topK: 3, + filter: 'user_id = "user-1" AND active = true', + includeMetadata: true, + }, + { namespace }, + ); + expect(results).toEqual([ + { + id: "memory-1", + payload: { data: "hello", user_id: "user-1" }, + score: 0.9, + }, + ]); + }); + + it("fetches, updates, deletes, resets, and lists vectors in the namespace", async () => { + const client = createClient({ + fetch: jest.fn().mockResolvedValue([ + { + id: "memory-1", + metadata: { data: "hello" }, + }, + ]), + range: jest + .fn() + .mockResolvedValueOnce({ + vectors: [ + { id: "memory-1", metadata: { user_id: "user-1" } }, + { id: "memory-2", metadata: { user_id: "user-2" } }, + ], + nextCursor: "2", + }) + .mockResolvedValueOnce({ + vectors: [{ id: "memory-3", metadata: { user_id: "user-1" } }], + nextCursor: "", + }), + }); + const store = new UpstashVector({ + collectionName: namespace, + client: client as any, + }); + + await expect(store.get("memory-1")).resolves.toEqual({ + id: "memory-1", + payload: { data: "hello" }, + }); + await store.update("memory-1", [0.3], { data: "updated" }); + await store.delete("memory-1"); + await store.deleteCol(); + await store.reset(); + await expect(store.list({ user_id: "user-1" }, 2)).resolves.toEqual([ + [ + { id: "memory-1", payload: { user_id: "user-1" } }, + { id: "memory-3", payload: { user_id: "user-1" } }, + ], + 2, + ]); + + expect(client.fetch).toHaveBeenCalledWith(["memory-1"], { + includeMetadata: true, + namespace, + }); + expect(client.upsert).toHaveBeenCalledWith( + { id: "memory-1", vector: [0.3], metadata: { data: "updated" } }, + { namespace }, + ); + expect(client.delete).toHaveBeenCalledWith("memory-1", { namespace }); + expect(client.reset).toHaveBeenCalledTimes(2); + expect(client.reset).toHaveBeenCalledWith({ namespace }); + }); + + it("stops paging when the cursor is exhausted instead of re-scanning", async () => { + // Upstash returns nextCursor "" at the end of a scan. If list() treated "" + // as "keep going" (e.g. by checking for "0"), it would re-fetch from the + // start and pile up duplicates until it hit topK. With fewer vectors than + // topK, a single page must end the scan: one range call, no duplicates. + const client = createClient({ + range: jest.fn().mockResolvedValue({ + vectors: [ + { id: "memory-1", metadata: { user_id: "user-1" } }, + { id: "memory-2", metadata: { user_id: "user-1" } }, + ], + nextCursor: "", + }), + }); + const store = new UpstashVector({ + collectionName: namespace, + client: client as any, + }); + + const [rows, count] = await store.list({ user_id: "user-1" }, 100); + + expect(client.range).toHaveBeenCalledTimes(1); + expect(rows).toEqual([ + { id: "memory-1", payload: { user_id: "user-1" } }, + { id: "memory-2", payload: { user_id: "user-1" } }, + ]); + expect(count).toBe(2); + }); +}); diff --git a/mem0-ts/tsup.config.ts b/mem0-ts/tsup.config.ts index ae8d290cc..e9189b7cf 100644 --- a/mem0-ts/tsup.config.ts +++ b/mem0-ts/tsup.config.ts @@ -20,6 +20,7 @@ const external = [ "@google-cloud/aiplatform", "@mistralai/mistralai", "@supabase/supabase-js", + "@upstash/vector", "@azure/search-documents", "@azure/identity", "cloudflare", diff --git a/mem0/configs/vector_stores/upstash_vector.py b/mem0/configs/vector_stores/upstash_vector.py index d4c3c7c3b..a382012f3 100644 --- a/mem0/configs/vector_stores/upstash_vector.py +++ b/mem0/configs/vector_stores/upstash_vector.py @@ -29,6 +29,13 @@ class UpstashVectorConfig(BaseModel): if not client and not (url and token): raise ValueError("Either a client or URL and token must be provided.") + + # Persist the env-resolved credentials so the provider constructor receives + # them; the validator used to check the env vars but drop them, so an + # env-var-only config passed validation and then raised on build. + if not client: + values["url"] = url + values["token"] = token return values model_config = ConfigDict(arbitrary_types_allowed=True) diff --git a/tests/vector_stores/test_upstash_vector.py b/tests/vector_stores/test_upstash_vector.py index 070e78bd1..5628028c7 100644 --- a/tests/vector_stores/test_upstash_vector.py +++ b/tests/vector_stores/test_upstash_vector.py @@ -4,6 +4,7 @@ from unittest.mock import MagicMock, call, patch import pytest +from mem0.configs.vector_stores.upstash_vector import UpstashVectorConfig from mem0.vector_stores.upstash_vector import UpstashVector @@ -393,3 +394,27 @@ def test_search_vectors_multi_query_namespace_at_top_level(upstash_instance): assert [r.id for r in results] == ["id1", "id2", "id3"] assert results[0].score == 0.9 assert results[1].payload == {"name": "vector2"} + + +def test_env_var_only_config_builds_provider(monkeypatch): + """Regression: an env-var-only config (no url/token/client) must build. + + VectorStoreFactory does ``UpstashVector(**config.model_dump())``. The config + validator read ``UPSTASH_VECTOR_REST_URL``/``UPSTASH_VECTOR_REST_TOKEN`` only + to pass its presence check, then returned the config unchanged, so + ``model_dump()`` still carried ``url=token=None`` and construction raised + "Either a client or URL and token must be provided." — even though the docs + advertise env-var setup. The resolved credentials must reach the ctor. + """ + monkeypatch.setenv("UPSTASH_VECTOR_REST_URL", "https://example.upstash.io") + monkeypatch.setenv("UPSTASH_VECTOR_REST_TOKEN", "tok_123") + + dumped = UpstashVectorConfig(collection_name="mem0").model_dump() + assert dumped["url"] == "https://example.upstash.io" + assert dumped["token"] == "tok_123" + + # Mirrors VectorStoreFactory.create: instance(**config.model_dump()). + with patch("mem0.vector_stores.upstash_vector.Index") as mock_index: + UpstashVector(**dumped) + + mock_index.assert_called_once_with("https://example.upstash.io", "tok_123")