From ee28fdbbbbedf2f276d60c4b0afc2c7e2ae4912b Mon Sep 17 00:00:00 2001 From: Anush008 Date: Tue, 28 Jul 2026 01:34:27 +0530 Subject: [PATCH] feat: Qdrant Vector Search --- README.md | 7 +- README.zh.md | 7 +- docs/en/API Docs/memory_server.md | 7 +- .../memory_server.md" | 7 +- jiuwen_memory/foundation/store/__init__.py | 9 +- .../store/vector/qdrant_vector_store.py | 490 ++++++++++++++++++ jiuwen_memory/server/.env.example | 10 +- jiuwen_memory/server/store_factory.py | 23 +- pyproject.toml | 3 +- .../store/vector/test_qdrant_vector_store.py | 249 +++++++++ uv.lock | 132 ++++- 11 files changed, 926 insertions(+), 18 deletions(-) create mode 100644 jiuwen_memory/foundation/store/vector/qdrant_vector_store.py create mode 100644 tests/unit_tests/core/foundation/store/vector/test_qdrant_vector_store.py diff --git a/README.md b/README.md index 578fdb54..21301f04 100644 --- a/README.md +++ b/README.md @@ -28,7 +28,7 @@ Agent conversational systems rely on limited context windows — once the Token - **Semantic Retrieval & Conflict Detection**: Unified cross-type vector semantic retrieval; `MemUpdateChecker` uses LLM to analyze semantic conflicts and intelligently decide ADD/DELETE strategies; LLM outputs UPDATE/DELETE directives validated via semantic checks before execution, ensuring memory consistency and controllable operations. -- **Full-Stack Storage Backend System**: Coverage across five storage categories — KV (InMemoryKV/ShelveStore/DbBasedKV/Redis), Vector (ChromaDB/Milvus/Elasticsearch/GaussVector), Relational (SQLite/PostgreSQL/MySQL/GaussDB), Message (SqlMessageStore), and Graph (Milvus GraphStore) — adapting from local single-node to cloud cluster scenarios. +- **Full-Stack Storage Backend System**: Coverage across five storage categories — KV (InMemoryKV/ShelveStore/DbBasedKV/Redis), Vector (ChromaDB/Milvus/Elasticsearch/GaussVector/Qdrant), Relational (SQLite/PostgreSQL/MySQL/GaussDB), Message (SqlMessageStore), and Graph (Milvus GraphStore) — adapting from local single-node to cloud cluster scenarios. - **Data Migration Framework**: Supports versioned schema migration for KV/vector/SQL/message/index stores and cross-BaseMemoryIndex batch data migration, with an operation registry for custom migration extensions. @@ -80,6 +80,9 @@ pip install JiuwenMemory[redis] # ChromaDB vector store pip install JiuwenMemory[chromadb] +# Qdrant vector store +pip install JiuwenMemory[qdrant] + # File-system memory backend (sqlite-vec + watchdog + jieba) # Required for INDEX_BACKEND=file; missing deps silently degrade: pip install JiuwenMemory[file-index] @@ -330,7 +333,7 @@ Graph Memory is an independent knowledge graph memory module. It turns input con ### **Flexible Storage Backends and Data Migration** -- **Full-Stack Storage Backends**: Coverage across five storage categories — KV (InMemoryKV/ShelveStore/DbBasedKV/Redis), Vector (ChromaDB/Milvus/Elasticsearch/GaussVector), Relational (SQLite/PostgreSQL/MySQL/GaussDB), Message (SqlMessageStore), and Graph (Milvus GraphStore) — adapting from local single-node to cloud cluster scenarios. +- **Full-Stack Storage Backends**: Coverage across five storage categories — KV (InMemoryKV/ShelveStore/DbBasedKV/Redis), Vector (ChromaDB/Milvus/Elasticsearch/GaussVector/Qdrant), Relational (SQLite/PostgreSQL/MySQL/GaussDB), Message (SqlMessageStore), and Graph (Milvus GraphStore) — adapting from local single-node to cloud cluster scenarios. - **Versioned Migration**: A complete migration framework supporting SQL schema changes, vector field renaming, KV data updates, message data transformation, and index field operations. - **Cross-index Migration**: Supports batch migration of memory data between different BaseMemoryIndex instances for smooth storage engine switching, with an operation registry for custom migration extensions. diff --git a/README.zh.md b/README.zh.md index b87d113e..c84a459e 100644 --- a/README.zh.md +++ b/README.zh.md @@ -28,7 +28,7 @@ Agent对话系统依赖有限的上下文窗口,一旦超出 Token 限制或 - **语义检索与冲突检测**:跨记忆类型统一向量语义检索;`MemUpdateChecker` 通过 LLM 分析语义冲突智能决定 ADD/DELETE 策略,LLM 输出 UPDATE/DELETE 指令经语义校验后执行,确保记忆一致与操作可控。 -- **全栈存储后端体系**:覆盖 KV(InMemoryKV/ShelveStore/DbBasedKV/Redis)、向量(ChromaDB/Milvus/Elasticsearch/GaussVector)、关系型(SQLite/PostgreSQL/MySQL/GaussDB)、消息(SqlMessageStore)、图(Milvus GraphStore)五大存储类别,适配本地单机到云端集群全场景。 +- **全栈存储后端体系**:覆盖 KV(InMemoryKV/ShelveStore/DbBasedKV/Redis)、向量(ChromaDB/Milvus/Elasticsearch/GaussVector/Qdrant)、关系型(SQLite/PostgreSQL/MySQL/GaussDB)、消息(SqlMessageStore)、图(Milvus GraphStore)五大存储类别,适配本地单机到云端集群全场景。 - **数据迁移框架**:支持 KV/向量/SQL/消息/索引的版本化 schema 迁移与跨 BaseMemoryIndex 实例批量数据迁移,操作注册表支持自定义迁移扩展。 @@ -80,6 +80,9 @@ pip install JiuwenMemory[redis] # ChromaDB 向量存储 pip install JiuwenMemory[chromadb] +# Qdrant 向量存储 +pip install JiuwenMemory[qdrant] + # 文件系统记忆后端(sqlite-vec + watchdog + jieba) # 用 INDEX_BACKEND=file 时需要,缺失则静默降级 pip install JiuwenMemory[file-index] @@ -322,7 +325,7 @@ Graph Memory 是独立的知识图谱记忆模块,可将输入内容沉淀为 ### **灵活的存储后端与数据迁移** -- **全栈存储后端**:覆盖 KV(InMemoryKV/ShelveStore/DbBasedKV/Redis)、向量(ChromaDB/Milvus/Elasticsearch/GaussVector)、关系型(SQLite/PostgreSQL/MySQL/GaussDB)、消息(SqlMessageStore)、图(Milvus GraphStore)五大存储类别,适配本地单机到云端集群全场景。 +- **全栈存储后端**:覆盖 KV(InMemoryKV/ShelveStore/DbBasedKV/Redis)、向量(ChromaDB/Milvus/Elasticsearch/GaussVector/Qdrant)、关系型(SQLite/PostgreSQL/MySQL/GaussDB)、消息(SqlMessageStore)、图(Milvus GraphStore)五大存储类别,适配本地单机到云端集群全场景。 - **版本化迁移**:提供完整的迁移框架,支持 SQL schema 变更、向量字段重命名、KV 数据更新、消息数据转换和索引字段操作等多种迁移类型。 - **跨索引迁移**:支持在不同 BaseMemoryIndex 实例之间批量迁移记忆数据,操作注册表支持自定义迁移扩展,实现存储引擎的平滑切换。 diff --git a/docs/en/API Docs/memory_server.md b/docs/en/API Docs/memory_server.md index 046a2c9d..5c758ade 100644 --- a/docs/en/API Docs/memory_server.md +++ b/docs/en/API Docs/memory_server.md @@ -105,7 +105,7 @@ If neither `.env` file is found, the service automatically creates the `~/.jiuwe | `DB_STORE_TYPE` | `default` | DB Store type. Supported values: `default` / `gauss`. | | `INDEX_BACKEND` | `simple` | Memory index backend. Supported values: `simple` (vector, KV+Vector) / `file` (markdown+SQLite). When `file`, long-term memories are persisted to markdown files; see `FILE_MEMORY_DATA_DIR`. | | `FILE_MEMORY_DATA_DIR` | `~/.jiuwenmemory/file_memory_data` | `FileMemoryIndex` data root directory. Only effective when `INDEX_BACKEND=file`. | -| `VECTOR_STORE_TYPE` | `chroma` | Vector Store type. Supported values: `chroma` / `milvus` / `elasticsearch` / `gauss`. Under `INDEX_BACKEND=file`, the vector store is only used for middle-term memory / dreaming and may be left unconfigured. | +| `VECTOR_STORE_TYPE` | `chroma` | Vector Store type. Supported values: `chroma` / `milvus` / `elasticsearch` / `gauss` / `qdrant`. Under `INDEX_BACKEND=file`, the vector store is only used for middle-term memory / dreaming and may be left unconfigured. | | `VECTOR_CHROMA_PERSIST_DIR` | `MEMORY_DATA_DIR` | Chroma persistence directory. | | `VECTOR_MILVUS_URI` | empty string | Milvus service URI. | | `VECTOR_MILVUS_TOKEN` | empty string | Milvus token; optional. | @@ -124,6 +124,11 @@ If neither `.env` file is found, the service automatically creates the `~/.jiuwe | `VECTOR_GAUSS_DATABASE` | `postgres` | Gauss vector store database name. | | `VECTOR_GAUSS_USER` | `postgres` | Gauss vector store user. | | `VECTOR_GAUSS_PASSWORD` | empty string | Gauss vector store password. | +| `VECTOR_QDRANT_URL` | `http://localhost:6333` | Qdrant HTTP endpoint. | +| `VECTOR_QDRANT_API_KEY` | empty string | Qdrant API key; optional for unsecured local deployments. | +| `VECTOR_QDRANT_COLLECTION_PREFIX` | `agent_vector` | Prefix used to isolate JiuwenMemory collections. | +| `VECTOR_QDRANT_PREFER_GRPC` | `false` | Prefer Qdrant's gRPC transport when set to `true`. | +| `VECTOR_QDRANT_TIMEOUT` | client default | Request timeout in seconds; leave empty to use the client default. | #### `INDEX_BACKEND=file` notes diff --git "a/docs/zh/API\346\226\207\346\241\243/memory_server.md" "b/docs/zh/API\346\226\207\346\241\243/memory_server.md" index 95932b3c..648620e2 100644 --- "a/docs/zh/API\346\226\207\346\241\243/memory_server.md" +++ "b/docs/zh/API\346\226\207\346\241\243/memory_server.md" @@ -105,7 +105,7 @@ IP=127.0.0.1 PORT=8000 python -m jiuwen_memory.server.memory_server | `DB_STORE_TYPE` | `default` | DB Store 类型,支持 `default` / `gauss`。 | | `INDEX_BACKEND` | `simple` | 记忆索引后端,支持 `simple`(向量,KV+Vector)/ `file`(markdown+SQLite)。`file` 时长期记忆落 markdown 文件,配合 `FILE_MEMORY_DATA_DIR`。 | | `FILE_MEMORY_DATA_DIR` | `~/.jiuwenmemory/file_memory_data` | FileMemoryIndex 数据根目录,仅 `INDEX_BACKEND=file` 时生效。 | -| `VECTOR_STORE_TYPE` | `chroma` | Vector Store 类型,支持 `chroma` / `milvus` / `elasticsearch` / `gauss`。`INDEX_BACKEND=file` 时仅用于中间记忆 / dreaming,可不配。 | +| `VECTOR_STORE_TYPE` | `chroma` | Vector Store 类型,支持 `chroma` / `milvus` / `elasticsearch` / `gauss` / `qdrant`。`INDEX_BACKEND=file` 时仅用于中间记忆 / dreaming,可不配。 | | `VECTOR_CHROMA_PERSIST_DIR` | `MEMORY_DATA_DIR` | Chroma 向量库持久化目录。 | | `VECTOR_MILVUS_URI` | 空字符串 | Milvus 服务地址。 | | `VECTOR_MILVUS_TOKEN` | 空字符串 | Milvus Token,可为空。 | @@ -124,6 +124,11 @@ IP=127.0.0.1 PORT=8000 python -m jiuwen_memory.server.memory_server | `VECTOR_GAUSS_DATABASE` | `postgres` | Gauss 向量库数据库名。 | | `VECTOR_GAUSS_USER` | `postgres` | Gauss 向量库用户名。 | | `VECTOR_GAUSS_PASSWORD` | 空字符串 | Gauss 向量库密码。 | +| `VECTOR_QDRANT_URL` | `http://localhost:6333` | Qdrant HTTP 服务地址。 | +| `VECTOR_QDRANT_API_KEY` | 空字符串 | Qdrant API Key;本地无鉴权部署可留空。 | +| `VECTOR_QDRANT_COLLECTION_PREFIX` | `agent_vector` | JiuwenMemory collection 隔离前缀。 | +| `VECTOR_QDRANT_PREFER_GRPC` | `false` | 设为 `true` 时优先使用 Qdrant gRPC 传输。 | +| `VECTOR_QDRANT_TIMEOUT` | 客户端默认值 | 请求超时秒数;留空使用客户端默认值。 | #### `INDEX_BACKEND=file` 说明 diff --git a/jiuwen_memory/foundation/store/__init__.py b/jiuwen_memory/foundation/store/__init__.py index c5236132..85aafef3 100644 --- a/jiuwen_memory/foundation/store/__init__.py +++ b/jiuwen_memory/foundation/store/__init__.py @@ -32,7 +32,7 @@ # Built-in backends. Closed to extension by design — use register_vector_store # or the entry_points mechanism for 3rd-party backends. -_BUILTIN_VECTOR_STORE_NAMES = frozenset({"chroma", "milvus", "gaussvector", "elasticsearch"}) +_BUILTIN_VECTOR_STORE_NAMES = frozenset({"chroma", "milvus", "gaussvector", "elasticsearch", "qdrant"}) # Explicit in-process registrations (register_vector_store). # Maps backend name -> factory callable (typically a class). @@ -48,7 +48,7 @@ def register_vector_store( Use this in application init code when shipping a plugin via entry_points is not practical (e.g., private in-repo backend). - Built-in names (chroma, milvus, gaussvector) cannot be overridden — a + Built-in names (chroma, milvus, gaussvector, elasticsearch, qdrant) cannot be overridden — a register call with a built-in name is kept in the registry but the built-in still wins in ``create_vector_store()`` resolution. @@ -78,6 +78,9 @@ def _resolve_builtin(store_type: str, kwargs: dict) -> "BaseVectorStore | None": if store_type == "elasticsearch": from jiuwen_memory.foundation.store.vector.es_vector_store import ElasticsearchVectorStore return ElasticsearchVectorStore(**kwargs) + if store_type == "qdrant": + from jiuwen_memory.foundation.store.vector.qdrant_vector_store import QdrantVectorStore + return QdrantVectorStore(**kwargs) return None @@ -122,7 +125,7 @@ def create_vector_store(store_type: str, **kwargs) -> "BaseVectorStore | None": """Factory for vector-store backends. Resolution order: - 1. Built-in (chroma, milvus, gaussvector) — always wins, closed set. + 1. Built-in (chroma, milvus, gaussvector, elasticsearch, qdrant) — always wins, closed set. 2. Explicit registrations via ``register_vector_store()``. 3. Entry_points in group ``openjiuwen.vector_stores``. diff --git a/jiuwen_memory/foundation/store/vector/qdrant_vector_store.py b/jiuwen_memory/foundation/store/vector/qdrant_vector_store.py new file mode 100644 index 00000000..d2dc5309 --- /dev/null +++ b/jiuwen_memory/foundation/store/vector/qdrant_vector_store.py @@ -0,0 +1,490 @@ +# coding: utf-8 +# Copyright (c) Huawei Technologies Co., Ltd. 2026. All rights reserved. +"""Qdrant vector-store adapter.""" + +from __future__ import annotations + +import copy +import uuid +from typing import Any + +from qdrant_client import AsyncQdrantClient, models +from qdrant_client.http.exceptions import UnexpectedResponse + +from jiuwen_memory.common.exception.codes import StatusCode +from jiuwen_memory.common.exception.errors import build_error +from jiuwen_memory.foundation.store.base_vector_store import ( + BaseVectorStore, + CollectionSchema, + VectorSearchResult, +) +from jiuwen_memory.foundation.store.filter_dsl import FilterGroup, FilterLogic, FilterOperator +from jiuwen_memory.foundation.store.vector.utils import ( + build_transform_func_for_operations, + compute_new_schema, + convert_cosine_similarity, + convert_ip_similarity, +) +from jiuwen_memory.memory_core.migration.operation.base_operation import BaseOperation + + +_META_KEY = "__jiuwen_collection_metadata__" +_META_ID = str(uuid.uuid5(uuid.NAMESPACE_URL, _META_KEY)) +_DISTANCES = { + "COSINE": models.Distance.COSINE, + "L2": models.Distance.EUCLID, + "EUCLID": models.Distance.EUCLID, + "EUCLIDEAN": models.Distance.EUCLID, + "IP": models.Distance.DOT, + "DOT": models.Distance.DOT, +} + + +class QdrantVectorStore(BaseVectorStore): + """Qdrant backend with persistent JiuwenMemory schema metadata.""" + + def __init__( + self, + client: AsyncQdrantClient | None = None, + url: str = "http://localhost:6333", + api_key: str | None = None, + collection_prefix: str = "agent_vector", + prefer_grpc: bool = False, + timeout: float | None = None, + **kwargs: Any, + ): + self._client = client or AsyncQdrantClient( + url=url, + api_key=api_key, + prefer_grpc=prefer_grpc, + timeout=timeout, + **kwargs, + ) + self._prefix = collection_prefix.strip("_") + self._metadata: dict[str, dict[str, Any]] = {} + + async def close(self) -> None: + await self._client.close() + self._metadata.clear() + + def _name(self, collection: str) -> str: + return f"{self._prefix}__{collection}" if self._prefix else collection + + @staticmethod + def _point_id(collection: str, value: Any) -> str: + # Qdrant only allows uint64/UUID IDs. + # Ref: https://qdrant.tech/documentation/manage-data/points/#point-ids + return str(uuid.uuid5(uuid.NAMESPACE_URL, f"{collection}\0{value}")) + + @staticmethod + def _filter(values: dict[str, Any] | FilterGroup | None = None) -> models.Filter: + if isinstance(values, FilterGroup): + conditions = [QdrantVectorStore._filter_group(values)] + else: + conditions = [QdrantVectorStore._field_condition(key, value) for key, value in (values or {}).items()] + return models.Filter( + must=conditions or None, + must_not=[models.FieldCondition(key=_META_KEY, match=models.MatchValue(value=True))], + ) + + @staticmethod + def _field_condition(key: str, value: Any) -> models.Condition: + if value is None: + return models.IsEmptyCondition(is_empty=models.PayloadField(key=key)) + match = ( + models.MatchAny(any=list(value)) + if isinstance(value, (list, tuple, set)) + else models.MatchValue(value=value) + ) + return models.FieldCondition(key=key, match=match) + + @staticmethod + def _filter_group(group: FilterGroup) -> models.Filter: + if not group.conditions: + raise build_error( + StatusCode.MEMORY_FILTER_FORMAT_ERROR, + error_msg="FilterGroup.conditions must be non-empty", + ) + conditions = [] + for condition in group.conditions: + if isinstance(condition, FilterGroup): + rendered = QdrantVectorStore._filter_group(condition) + else: + field = QdrantVectorStore._field_condition(condition.field, condition.value) + rendered = models.Filter(must_not=[field]) if condition.op == FilterOperator.NE else field + conditions.append(rendered) + if group.logic == FilterLogic.OR: + return models.Filter(should=conditions) + return models.Filter(must=conditions) + + @staticmethod + def _schema(schema: CollectionSchema | dict[str, Any]) -> CollectionSchema: + return CollectionSchema.from_dict(schema) if isinstance(schema, dict) else schema + + @staticmethod + def _distance(metric: str) -> models.Distance: + try: + return _DISTANCES[metric.upper()] + except KeyError as exc: + raise build_error( + StatusCode.STORE_VECTOR_SCHEMA_INVALID, + error_msg=f"Unsupported Qdrant distance metric: {metric}", + ) from exc + + @staticmethod + def _score(value: float, metric: str) -> float: + metric = metric.upper() + if metric == "COSINE": + return convert_cosine_similarity(value) + if metric in {"IP", "DOT"}: + return convert_ip_similarity(value) + return 1.0 / (1.0 + value) + + @staticmethod + def _vector_fields(schema: CollectionSchema) -> dict[str, int]: + return {field.name: field.dim for field in schema.get_vector_fields()} + + async def _assert_compatible( + self, + collection_name: str, + schema: CollectionSchema, + metric: str, + ) -> bool: + metadata = await self.get_collection_metadata(collection_name) + if metadata.get("schema"): + compatible = metadata["schema"] == schema.to_dict() and self._distance( + metadata.get("distance_metric", "COSINE") + ) == self._distance(metric) + else: + configured = (await self._client.get_collection(self._name(collection_name))).config.params.vectors + expected = { + name: (dimension, self._distance(metric)) for name, dimension in self._vector_fields(schema).items() + } + actual = ( + {name: (params.size, params.distance) for name, params in configured.items()} + if isinstance(configured, dict) + else {} + ) + compatible = actual == expected + if not compatible: + raise build_error( + StatusCode.STORE_VECTOR_SCHEMA_INVALID, + error_msg=f"Qdrant collection '{collection_name}' already exists with an incompatible schema", + ) + return bool(metadata.get("schema")) + + async def create_collection( + self, + collection_name: str, + schema: CollectionSchema | dict[str, Any], + **kwargs: Any, + ) -> None: + schema = self._schema(schema) + primary = schema.get_primary_key_field() + vectors = self._vector_fields(schema) + if primary is None: + raise build_error( + StatusCode.STORE_VECTOR_SCHEMA_INVALID, + error_msg="schema must contain a primary key field", + ) + if not vectors: + raise build_error( + StatusCode.STORE_VECTOR_SCHEMA_INVALID, + error_msg="schema must contain at least one FLOAT_VECTOR field", + ) + + metric = kwargs.pop("distance_metric", "COSINE").upper() + distance = self._distance(metric) + metadata = { + "schema": schema.to_dict(), + "distance_metric": metric, + "primary_key_field": primary.name, + "vector_field": next(iter(vectors)), + "vector_dim": next(iter(vectors.values())), + "schema_version": 0, + "collection_name": collection_name, + } + if not await self.collection_exists(collection_name): + try: + await self._client.create_collection( + collection_name=self._name(collection_name), + vectors_config={ + name: models.VectorParams(size=dimension, distance=distance) + for name, dimension in vectors.items() + }, + **kwargs, + ) + except UnexpectedResponse as exc: + if exc.status_code != 409: + raise + if not await self._assert_compatible(collection_name, schema, metric): + await self._write_metadata(collection_name, metadata) + + async def delete_collection(self, collection_name: str, **kwargs: Any) -> None: + if await self.collection_exists(collection_name): + await self._client.delete_collection(self._name(collection_name), **kwargs) + self._metadata.pop(collection_name, None) + + async def collection_exists(self, collection_name: str, **kwargs: Any) -> bool: + return await self._client.collection_exists(self._name(collection_name)) + + async def get_schema(self, collection_name: str, **kwargs: Any) -> CollectionSchema: + metadata = await self.get_collection_metadata(collection_name) + return self._schema_from_metadata(collection_name, metadata) + + @staticmethod + def _schema_from_metadata(collection_name: str, metadata: dict[str, Any]) -> CollectionSchema: + schema = metadata.get("schema") + if not schema: + raise build_error( + StatusCode.STORE_VECTOR_COLLECTION_NOT_FOUND, + collection_name=collection_name, + error_msg="Qdrant collection schema metadata is missing", + ) + return CollectionSchema.from_dict(schema) + + async def add_docs(self, collection_name: str, docs: list[dict[str, Any]], **kwargs: Any) -> None: + if not docs: + return + metadata = await self.get_collection_metadata(collection_name) + schema = self._schema_from_metadata(collection_name, metadata) + primary = schema.get_primary_key_field() + vector_fields = set(self._vector_fields(schema)) + points = [] + for doc in docs: + key = doc.get(primary.name) + if key is None: + if not primary.auto_id: + raise build_error( + StatusCode.STORE_VECTOR_DOC_INVALID, + collection_name=collection_name, + error_msg=f"document is missing primary key field '{primary.name}'", + ) + key = str(uuid.uuid4()) + missing = {name for name in vector_fields if doc.get(name) is None} + if missing: + raise build_error( + StatusCode.STORE_VECTOR_DOC_INVALID, + collection_name=collection_name, + error_msg=f"document is missing vector fields: {sorted(missing)}", + ) + payload = {name: value for name, value in doc.items() if name not in vector_fields and value is not None} + payload[primary.name] = key + points.append( + models.PointStruct( + id=self._point_id(collection_name, key), + vector={name: doc[name] for name in vector_fields}, + payload=payload, + ) + ) + await self._client.upsert(self._name(collection_name), points=points, wait=kwargs.get("wait", True)) + + async def search( + self, + collection_name: str, + query_vector: list[float], + vector_field: str, + top_k: int = 5, + filters: dict[str, Any] | None = None, + **kwargs: Any, + ) -> list[VectorSearchResult]: + metadata = await self.get_collection_metadata(collection_name) + output_fields = kwargs.get("output_fields") + vector_fields = set(self._vector_fields(self._schema_from_metadata(collection_name, metadata))) + include_vectors = output_fields is None or bool(vector_fields.intersection(output_fields)) + response = await self._client.query_points( + self._name(collection_name), + query=query_vector, + using=vector_field, + query_filter=self._filter(filters), + limit=top_k, + with_payload=True, + with_vectors=include_vectors, + ) + results = [] + for point in response.points: + fields = dict(point.payload or {}) + if include_vectors and isinstance(point.vector, dict): + fields.update({key: value for key, value in point.vector.items() if key in vector_fields}) + if output_fields: + fields = {key: fields[key] for key in output_fields if key in fields} + results.append( + VectorSearchResult(score=self._score(point.score, metadata["distance_metric"]), fields=fields) + ) + return results + + async def list_docs( + self, + collection_name: str, + filters: FilterGroup | None = None, + limit: int = 100, + offset: int = 0, + **kwargs: Any, + ) -> list[dict[str, Any]]: + output_fields = kwargs.get("output_fields") + metadata = await self.get_collection_metadata(collection_name) + vector_fields = set(self._vector_fields(self._schema_from_metadata(collection_name, metadata))) + with_vectors = output_fields is None or bool(vector_fields.intersection(output_fields)) + records, _ = await self._client.scroll( + self._name(collection_name), + scroll_filter=self._filter(filters), + limit=limit + offset, + with_payload=True, + with_vectors=with_vectors, + ) + documents = [] + for record in records[offset:]: + document = {**(record.payload or {}), **(record.vector or {})} + if output_fields: + document = {field: document.get(field) for field in output_fields if field in document} + documents.append(document) + return documents + + async def update_doc_fields( + self, + collection_name: str, + doc_id: str, + fields: dict[str, Any], + ) -> None: + if not doc_id: + raise build_error( + StatusCode.STORE_VECTOR_DOC_INVALID, + error_msg="doc_id is required for update_doc_fields", + ) + if not fields: + return + point_id = self._point_id(collection_name, doc_id) + records = await self._client.retrieve( + self._name(collection_name), ids=[point_id], with_payload=False, with_vectors=False + ) + if not records: + raise build_error( + StatusCode.STORE_VECTOR_DOC_INVALID, + error_msg=f"doc not found in collection, doc_id={doc_id}", + ) + await self._client.set_payload( + self._name(collection_name), + payload=fields, + points=[point_id], + wait=True, + ) + + async def delete_docs_by_ids(self, collection_name: str, ids: list[str], **kwargs: Any) -> None: + if ids: + await self._client.delete( + self._name(collection_name), + points_selector=models.PointIdsList(points=[self._point_id(collection_name, value) for value in ids]), + wait=kwargs.get("wait", True), + ) + + async def delete_docs_by_filters(self, collection_name: str, filters: dict[str, Any], **kwargs: Any) -> None: + if filters: + await self._client.delete( + self._name(collection_name), + points_selector=models.FilterSelector(filter=self._filter(filters)), + wait=kwargs.get("wait", True), + ) + + async def list_collection_names(self) -> list[str]: + prefix = f"{self._prefix}__" if self._prefix else "" + collections = (await self._client.get_collections()).collections + return [item.name.removeprefix(prefix) for item in collections if not prefix or item.name.startswith(prefix)] + + async def _write_metadata(self, collection_name: str, metadata: dict[str, Any]) -> None: + vectors = self._vector_fields(CollectionSchema.from_dict(metadata["schema"])) + await self._client.upsert( + self._name(collection_name), + points=[ + models.PointStruct( + id=_META_ID, + vector={name: [0.0] * dimension for name, dimension in vectors.items()}, + payload={_META_KEY: True, "metadata": metadata}, + ) + ], + wait=True, + ) + self._metadata[collection_name] = dict(metadata) + + async def get_collection_metadata(self, collection_name: str) -> dict[str, Any]: + if collection_name in self._metadata: + return dict(self._metadata[collection_name]) + records = await self._client.retrieve( + self._name(collection_name), ids=[_META_ID], with_payload=True, with_vectors=False + ) + metadata = dict(records[0].payload.get("metadata", {})) if records else {} + if metadata: + self._metadata[collection_name] = metadata + return dict(metadata) + + async def update_collection_metadata(self, collection_name: str, metadata: dict[str, Any]) -> None: + if not metadata: + return + version = metadata.get("schema_version") + if version is not None and (not isinstance(version, int) or version < 0): + raise build_error( + StatusCode.STORE_VECTOR_SCHEMA_INVALID, + error_msg=f"schema_version must be a non-negative integer, got {version}", + ) + current = await self.get_collection_metadata(collection_name) + current.update(metadata) + await self._write_metadata(collection_name, current) + + async def _documents(self, collection_name: str, batch_size: int = 256) -> list[dict[str, Any]]: + documents, offset = [], None + while True: + records, offset = await self._client.scroll( + self._name(collection_name), + scroll_filter=self._filter(), + limit=batch_size, + offset=offset, + with_payload=True, + with_vectors=True, + ) + for record in records: + documents.append({**(record.payload or {}), **(record.vector or {})}) + if offset is None: + return documents + + async def _replace( + self, + collection_name: str, + schema: CollectionSchema, + documents: list[dict[str, Any]], + metadata: dict[str, Any], + ) -> None: + if await self.collection_exists(collection_name): + await self.delete_collection(collection_name) + await self.create_collection(collection_name, schema, distance_metric=metadata["distance_metric"]) + await self.add_docs(collection_name, documents) + await self.update_collection_metadata( + collection_name, + {key: value for key, value in metadata.items() if key not in {"schema", "collection_name"}}, + ) + + async def update_schema(self, collection_name: str, operations: list[BaseOperation]) -> None: + if not operations: + return + old_schema = await self.get_schema(collection_name) + old_metadata = await self.get_collection_metadata(collection_name) + old_documents = await self._documents(collection_name) + new_schema = compute_new_schema(old_schema, operations) + transform = build_transform_func_for_operations(operations) + new_documents = [transform(copy.deepcopy(document)) for document in old_documents] + temp = f"{collection_name}_migration_{uuid.uuid4().hex}" + replacement_started = False + try: + await self.create_collection(temp, new_schema, distance_metric=old_metadata["distance_metric"]) + await self.add_docs(temp, new_documents) + migrated = await self._documents(temp) + replacement_started = True + await self._replace(collection_name, new_schema, migrated, old_metadata) + except Exception as migration_error: + if replacement_started: + try: + await self._replace(collection_name, old_schema, old_documents, old_metadata) + except Exception as rollback_error: + migration_error.add_note(f"Qdrant schema rollback failed: {rollback_error}") + raise + finally: + if await self.collection_exists(temp): + await self.delete_collection(temp) diff --git a/jiuwen_memory/server/.env.example b/jiuwen_memory/server/.env.example index 82ad0812..1ef93b59 100644 --- a/jiuwen_memory/server/.env.example +++ b/jiuwen_memory/server/.env.example @@ -76,7 +76,7 @@ INDEX_BACKEND=simple FILE_MEMORY_DATA_DIR= # ---------- Vector Store ---------- -# 可选值: chroma | milvus | elasticsearch | gauss (默认 chroma) +# 可选值: chroma | milvus | elasticsearch | gauss | qdrant (默认 chroma) # INDEX_BACKEND=file 时,Vector Store 仅用于中间记忆 / dreaming,可不配 VECTOR_STORE_TYPE=chroma @@ -113,3 +113,11 @@ VECTOR_GAUSS_PORT=5432 VECTOR_GAUSS_DATABASE=postgres VECTOR_GAUSS_USER=postgres VECTOR_GAUSS_PASSWORD= + +# -- qdrant --(仅 VECTOR_STORE_TYPE=qdrant 时需要) +VECTOR_QDRANT_URL=http://localhost:6333 +VECTOR_QDRANT_API_KEY= +VECTOR_QDRANT_COLLECTION_PREFIX=agent_vector +VECTOR_QDRANT_PREFER_GRPC=false +# 请求超时秒数;留空使用 qdrant-client 默认值 +VECTOR_QDRANT_TIMEOUT= diff --git a/jiuwen_memory/server/store_factory.py b/jiuwen_memory/server/store_factory.py index 7e8a32af..5bae4c36 100644 --- a/jiuwen_memory/server/store_factory.py +++ b/jiuwen_memory/server/store_factory.py @@ -52,14 +52,14 @@ def create_async_engine_from_env() -> AsyncEngine: """ data_directory = _data_dir() db_url = os.getenv("DB_URL", "").strip() - + if not db_url: db_path = Path(data_directory) / "sqlite_db.db" db_path.parent.mkdir(parents=True, exist_ok=True) db_url = f"sqlite+aiosqlite:///{db_path}" if db_url.startswith("gaussdb"): - import jiuwen_memory.foundation.store.db.gauss_dialect + import jiuwen_memory.foundation.store.db.gauss_dialect # noqa: F401 memory_logger.info("Using DB engine url=%s", db_url) return create_async_engine( @@ -124,7 +124,7 @@ def create_db_store(engine: AsyncEngine) -> BaseDbStore: def create_vector_store() -> BaseVectorStore: """根据 VECTOR_STORE_TYPE 构建 Vector store。 - 支持值: chroma | milvus | elasticsearch | gauss + 支持值: chroma | milvus | elasticsearch | gauss | qdrant 各分支懒导入对应实现,避免没装第三方包时影响其他分支启动。 """ vec_type = os.getenv("VECTOR_STORE_TYPE", "chroma").strip().lower() @@ -194,7 +194,22 @@ def create_vector_store() -> BaseVectorStore: password=os.getenv("VECTOR_GAUSS_PASSWORD", ""), ) + if vec_type == "qdrant": + from jiuwen_memory.foundation.store.vector.qdrant_vector_store import QdrantVectorStore + timeout_raw = os.getenv("VECTOR_QDRANT_TIMEOUT", "").strip() + try: + timeout = float(timeout_raw) if timeout_raw else None + except ValueError as exc: + raise ValueError(f"VECTOR_QDRANT_TIMEOUT must be a number, got {timeout_raw!r}") from exc + return QdrantVectorStore( + url=os.getenv("VECTOR_QDRANT_URL", "http://localhost:6333").strip() or "http://localhost:6333", + api_key=os.getenv("VECTOR_QDRANT_API_KEY", "").strip() or None, + collection_prefix=os.getenv("VECTOR_QDRANT_COLLECTION_PREFIX", "agent_vector").strip(), + prefer_grpc=os.getenv("VECTOR_QDRANT_PREFER_GRPC", "false").strip().lower() == "true", + timeout=timeout, + ) + raise ValueError( f"Unsupported VECTOR_STORE_TYPE={vec_type!r}, " - f"expected one of: chroma, milvus, elasticsearch, gauss" + f"expected one of: chroma, milvus, elasticsearch, gauss, qdrant" ) diff --git a/pyproject.toml b/pyproject.toml index 2eb54838..f9ac3730 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -54,6 +54,7 @@ dependencies = [ chromadb = ["chromadb>=1.5.4", "opentelemetry-proto>=1.33.0", "protobuf>=5.0,<7.0"] gaussvector = ["psycopg2-binary>=2.9.10"] elasticsearch = ["elasticsearch[async]>=8,<9"] +qdrant = ["qdrant-client>=1.10.0,<2.0.0"] postgres = ["asyncpg>=0.30.0"] mysql = ["aiomysql>=0.3.2"] gaussdb = ["async-gaussdb~=0.30.0", "asyncpg>=0.29.0"] @@ -63,7 +64,7 @@ mem0 = ["mem0ai>=1.0.0"] agentarts = ["agentarts-sdk>=0.1.2,<0.2"] server = ["uvicorn>=0.30.0", "fastapi>=0.110.0", "python-dotenv>=1.2.1", "mcp>=1.2.0"] all-storage = ["JiuwenMemory[sqlite,postgres,mysql,gaussdb,redis]"] -all-vector = ["JiuwenMemory[chromadb,gaussvector,elasticsearch]"] +all-vector = ["JiuwenMemory[chromadb,gaussvector,elasticsearch,qdrant]"] file-index = ["sqlite-vec>=0.1.9", "watchdog>=4.0", "jieba>=0.42.1"] default = [] all = ["JiuwenMemory[all-storage,all-vector,server,file-index]"] diff --git a/tests/unit_tests/core/foundation/store/vector/test_qdrant_vector_store.py b/tests/unit_tests/core/foundation/store/vector/test_qdrant_vector_store.py new file mode 100644 index 00000000..ce9da8b2 --- /dev/null +++ b/tests/unit_tests/core/foundation/store/vector/test_qdrant_vector_store.py @@ -0,0 +1,249 @@ +# coding: utf-8 +# Copyright (c) Huawei Technologies Co., Ltd. 2026. All rights reserved. + +import asyncio +from datetime import datetime, timezone + +import pytest +import pytest_asyncio + +qdrant_client = pytest.importorskip("qdrant_client") + +from jiuwen_memory.common.exception.errors import BaseError # noqa: E402 +from jiuwen_memory.foundation.store.base_vector_store import ( # noqa: E402 + CollectionSchema, + FieldSchema, + VectorDataType, +) +from jiuwen_memory.foundation.store import create_vector_store # noqa: E402 +from jiuwen_memory.foundation.store.base_memory_index import MemoryDoc # noqa: E402 +from jiuwen_memory.foundation.store.index.simple_memory_index import SimpleMemoryIndex # noqa: E402 +from jiuwen_memory.foundation.store.kv.in_memory_kv_store import InMemoryKVStore # noqa: E402 +from jiuwen_memory.foundation.store.vector.qdrant_vector_store import QdrantVectorStore # noqa: E402 +from jiuwen_memory.memory_core.migration.operation.base_operation import OperationMetadata # noqa: E402 +from jiuwen_memory.memory_core.migration.operation.operations import AddScalarFieldOperation # noqa: E402 + + +def _schema() -> CollectionSchema: + return CollectionSchema( + fields=[ + FieldSchema(name="id", dtype=VectorDataType.VARCHAR, is_primary=True), + FieldSchema(name="embedding", dtype=VectorDataType.FLOAT_VECTOR, dim=4), + FieldSchema(name="category", dtype=VectorDataType.VARCHAR), + FieldSchema(name="metadata", dtype=VectorDataType.JSON), + ] + ) + + +@pytest_asyncio.fixture(name="store") +async def _store(): + client = qdrant_client.AsyncQdrantClient(location=":memory:") + store = QdrantVectorStore(client=client, collection_prefix="pytest_agent_vector") + try: + yield store + finally: + await store.close() + + +@pytest.mark.asyncio +async def test_foundation_factory_resolves_qdrant(): + client = qdrant_client.AsyncQdrantClient(location=":memory:") + store = create_vector_store("qdrant", client=client, collection_prefix="factory_test") + try: + assert isinstance(store, QdrantVectorStore) + finally: + await store.close() + + +@pytest.mark.asyncio +async def test_qdrant_collection_creation_is_concurrent_and_schema_safe(store): + collection = "qdrant_create_race" + await asyncio.gather(*(store.create_collection(collection, _schema()) for _ in range(10))) + assert await store.collection_exists(collection) + + incompatible = _schema() + incompatible.get_vector_fields()[0].dim = 8 + with pytest.raises(BaseError): + await store.create_collection(collection, incompatible) + + +@pytest.mark.asyncio +async def test_qdrant_recovers_metadata_after_interrupted_creation(): + collection = "qdrant_metadata_recovery" + prefix = "pytest_agent_vector" + client = qdrant_client.AsyncQdrantClient(location=":memory:") + store = QdrantVectorStore(client=client, collection_prefix=prefix) + try: + await client.create_collection( + f"{prefix}__{collection}", + vectors_config={ + "embedding": qdrant_client.models.VectorParams(size=4, distance=qdrant_client.models.Distance.COSINE) + }, + ) + assert await store.get_collection_metadata(collection) == {} + + recovery_schema = CollectionSchema( + fields=[ + FieldSchema(name="id", dtype=VectorDataType.VARCHAR, is_primary=True), + FieldSchema(name="embedding", dtype=VectorDataType.FLOAT_VECTOR, dim=4), + ] + ) + await store.create_collection(collection, recovery_schema) + assert (await store.get_schema(collection)).get_vector_fields()[0].dim == 4 + finally: + await store.close() + + +@pytest.mark.asyncio +async def test_qdrant_integrates_with_simple_memory_index(store): + class Embedding: + async def embed_documents(self, texts): + return [[1.0, 0.0] if "apple" in text else [0.0, 1.0] for text in texts] + + async def embed_query(self, text): + return [1.0, 0.0] + + index = SimpleMemoryIndex(InMemoryKVStore(), store, Embedding()) + memory = MemoryDoc( + id="000000000000000000000001", + text="apple preference", + type="semantic_memory", + timestamp=datetime.now(timezone.utc), + ) + await index.add_memories("user", "scope", [memory]) + + results = await index.search("user", "scope", "fruit", ["semantic_memory"]) + assert results[0][0].text == "apple preference" + assert results[0][1] == pytest.approx(1.0) + + +@pytest.mark.asyncio +async def test_qdrant_vector_store_public_interfaces(store): + collection = "qdrant_vectors" + await store.create_collection(collection, _schema(), distance_metric="COSINE") + + assert await store.collection_exists(collection) + assert collection in await store.list_collection_names() + assert (await store.get_schema(collection)).get_vector_fields()[0].dim == 4 + + metadata = await store.get_collection_metadata(collection) + assert metadata["primary_key_field"] == "id" + assert metadata["distance_metric"] == "COSINE" + await store.update_collection_metadata(collection, {"schema_version": 1, "owner": "test"}) + assert (await store.get_collection_metadata(collection))["owner"] == "test" + + await store.add_docs( + collection, + [ + { + "id": "arbitrary/string/id", + "embedding": [1.0, 0.0, 0.0, 0.0], + "category": "fruit", + "metadata": {"source": "one"}, + }, + { + "id": "doc-2", + "embedding": [0.0, 1.0, 0.0, 0.0], + "category": "vehicle", + "metadata": {"source": "two"}, + }, + ], + ) + + results = await store.search( + collection, + [1.0, 0.0, 0.0, 0.0], + "embedding", + filters={"category": "fruit"}, + ) + assert len(results) == 1 + assert results[0].fields["id"] == "arbitrary/string/id" + assert results[0].fields["embedding"] == [1.0, 0.0, 0.0, 0.0] + assert results[0].score == pytest.approx(1.0) + + await store.delete_docs_by_ids(collection, ["arbitrary/string/id"]) + results = await store.search(collection, [1.0, 0.0, 0.0, 0.0], "embedding", top_k=5) + assert {result.fields["id"] for result in results} == {"doc-2"} + + await store.delete_docs_by_filters(collection, {"category": "vehicle"}) + assert await store.search(collection, [0.0, 1.0, 0.0, 0.0], "embedding") == [] + + await store.delete_collection(collection) + assert not await store.collection_exists(collection) + + +@pytest.mark.asyncio +async def test_qdrant_schema_migration_preserves_documents(store): + collection = "qdrant_migration" + await store.create_collection(collection, _schema()) + await store.add_docs( + collection, + [ + { + "id": "doc-1", + "embedding": [1.0, 0.0, 0.0, 0.0], + "category": "fruit", + "metadata": {}, + } + ], + ) + + await store.update_schema( + collection, + [ + AddScalarFieldOperation( + metadata=OperationMetadata(schema_version=1), + data_type="memory", + field_name="status", + field_type="string", + default_value="active", + ) + ], + ) + + schema = await store.get_schema(collection) + assert schema.get_field("status").default_value == "active" + results = await store.search(collection, [1.0, 0.0, 0.0, 0.0], "embedding") + assert results[0].fields["id"] == "doc-1" + assert results[0].fields["status"] == "active" + + +@pytest.mark.asyncio +async def test_qdrant_schema_migration_restores_original_on_rewrite_failure(store, monkeypatch): + collection = "qdrant_migration_rollback" + await store.create_collection(collection, _schema()) + await store.add_docs( + collection, + [ + { + "id": "doc-1", + "embedding": [1.0, 0.0, 0.0, 0.0], + "category": "fruit", + "metadata": {}, + } + ], + ) + original_add = store.add_docs + + async def fail_migrated_rewrite(target, docs, **kwargs): + if target == collection and docs and "status" in docs[0]: + monkeypatch.setattr(store, "add_docs", original_add) + raise RuntimeError("simulated Qdrant rewrite failure") + await original_add(target, docs, **kwargs) + + monkeypatch.setattr(store, "add_docs", fail_migrated_rewrite) + operation = AddScalarFieldOperation( + metadata=OperationMetadata(schema_version=1), + data_type="memory", + field_name="status", + field_type="string", + default_value="active", + ) + + with pytest.raises(RuntimeError, match="simulated Qdrant rewrite failure"): + await store.update_schema(collection, [operation]) + + assert (await store.get_schema(collection)).get_field("status") is None + results = await store.search(collection, [1.0, 0.0, 0.0, 0.0], "embedding") + assert results[0].fields["id"] == "doc-1" + assert "status" not in results[0].fields diff --git a/uv.lock b/uv.lock index 6c73a303..8bb04aaa 100644 --- a/uv.lock +++ b/uv.lock @@ -1131,6 +1131,15 @@ http2 = [ { name = "h2" }, ] +[[package]] +name = "httpx-sse" +version = "0.4.3" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/0f/4c/751061ffa58615a32c31b2d82e8482be8dd4a89154f003147acee90f2be9/httpx_sse-0.4.3.tar.gz", hash = "sha256:9b1ed0127459a66014aec3c56bebd93da3c1bc8bb6618c8082039a44889a755d" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/d2/fd/6668e5aec43ab844de6fc74927e155a3b37bf40d7c3790e49fc0406b6578/httpx_sse-0.4.3-py3-none-any.whl", hash = "sha256:0ac1c9fe3c0afad2e0ebb25a934a59f4c7823b60792691f779fad2c5568830fc" }, +] + [[package]] name = "huaweicloudsdkagentidentity" version = "3.1.201" @@ -1247,6 +1256,12 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/3e/95/c7c34aa53c16353c56d0b802fba48d5f5caa2cdee7958acbcb795c830416/isort-8.0.1-py3-none-any.whl", hash = "sha256:28b89bc70f751b559aeca209e6120393d43fbe2490de0559662be7a9787e3d75" }, ] +[[package]] +name = "jieba" +version = "0.42.1" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/c6/cb/18eeb235f833b726522d7ebed54f2278ce28ba9438e3135ab0278d9792a2/jieba-0.42.1.tar.gz", hash = "sha256:055ca12f62674fafed09427f176506079bc135638a14e23e25be909131928db2" } + [[package]] name = "jinja2" version = "3.1.6" @@ -1364,12 +1379,17 @@ all = [ { name = "elasticsearch", extra = ["async"] }, { name = "fastapi" }, { name = "greenlet" }, + { name = "jieba" }, + { name = "mcp" }, { name = "opentelemetry-proto" }, { name = "protobuf" }, { name = "psycopg2-binary" }, { name = "python-dotenv" }, + { name = "qdrant-client" }, { name = "redis" }, + { name = "sqlite-vec" }, { name = "uvicorn" }, + { name = "watchdog" }, ] all-storage = [ { name = "aiomysql" }, @@ -1385,6 +1405,7 @@ all-vector = [ { name = "opentelemetry-proto" }, { name = "protobuf" }, { name = "psycopg2-binary" }, + { name = "qdrant-client" }, ] chromadb = [ { name = "chromadb" }, @@ -1394,6 +1415,11 @@ chromadb = [ elasticsearch = [ { name = "elasticsearch", extra = ["async"] }, ] +file-index = [ + { name = "jieba" }, + { name = "sqlite-vec" }, + { name = "watchdog" }, +] gaussdb = [ { name = "async-gaussdb" }, { name = "asyncpg" }, @@ -1410,11 +1436,15 @@ mysql = [ postgres = [ { name = "asyncpg" }, ] +qdrant = [ + { name = "qdrant-client" }, +] redis = [ { name = "redis" }, ] server = [ { name = "fastapi" }, + { name = "mcp" }, { name = "python-dotenv" }, { name = "uvicorn" }, ] @@ -1469,10 +1499,12 @@ requires-dist = [ { name = "filelock", specifier = ">=3.20.1" }, { name = "greenlet", marker = "extra == 'sqlite'", specifier = ">=3.1.1" }, { name = "httpx", specifier = ">=0.24" }, - { name = "jiuwenmemory", extras = ["all-storage", "all-vector", "server"], marker = "extra == 'all'" }, - { name = "jiuwenmemory", extras = ["chromadb", "gaussvector", "elasticsearch"], marker = "extra == 'all-vector'" }, + { name = "jieba", marker = "extra == 'file-index'", specifier = ">=0.42.1" }, + { name = "jiuwenmemory", extras = ["all-storage", "all-vector", "server", "file-index"], marker = "extra == 'all'" }, + { name = "jiuwenmemory", extras = ["chromadb", "gaussvector", "elasticsearch", "qdrant"], marker = "extra == 'all-vector'" }, { name = "jiuwenmemory", extras = ["sqlite", "postgres", "mysql", "gaussdb", "redis"], marker = "extra == 'all-storage'" }, { name = "loguru", specifier = ">=0.7.3" }, + { name = "mcp", marker = "extra == 'server'", specifier = ">=1.2.0" }, { name = "mem0ai", marker = "extra == 'mem0'", specifier = ">=1.0.0" }, { name = "oauthlib", specifier = ">=3.3.1" }, { name = "openai", specifier = ">=1.108.0" }, @@ -1486,15 +1518,18 @@ requires-dist = [ { name = "python-dotenv", specifier = ">=1.2.1" }, { name = "python-dotenv", marker = "extra == 'server'", specifier = ">=1.2.1" }, { name = "pyyaml", specifier = ">=6.0" }, + { name = "qdrant-client", marker = "extra == 'qdrant'", specifier = ">=1.10.0,<2.0.0" }, { name = "redis", marker = "extra == 'redis'", specifier = ">=7.1.0" }, { name = "requests", specifier = ">=2.32.3" }, { name = "sqlalchemy", specifier = ">=2.0.41" }, + { name = "sqlite-vec", marker = "extra == 'file-index'", specifier = ">=0.1.9" }, { name = "sqlmodel", specifier = ">=0.0.37" }, { name = "tenacity", specifier = ">=9.1.2" }, { name = "tqdm", specifier = ">=4.0" }, { name = "uvicorn", marker = "extra == 'server'", specifier = ">=0.30.0" }, + { name = "watchdog", marker = "extra == 'file-index'", specifier = ">=4.0" }, ] -provides-extras = ["chromadb", "gaussvector", "elasticsearch", "postgres", "mysql", "gaussdb", "sqlite", "redis", "mem0", "agentarts", "server", "all-storage", "all-vector", "default", "all"] +provides-extras = ["chromadb", "gaussvector", "elasticsearch", "qdrant", "postgres", "mysql", "gaussdb", "sqlite", "redis", "mem0", "agentarts", "server", "all-storage", "all-vector", "file-index", "default", "all"] [package.metadata.requires-dev] dev = [ @@ -1715,6 +1750,31 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/27/1a/1f68f9ba0c207934b35b86a8ca3aad8395a3d6dd7921c0686e23853ff5a9/mccabe-0.7.0-py2.py3-none-any.whl", hash = "sha256:6c2d30ab6be0e4a46919781807b4f0d834ebdd6c6e3dca0bda5a15f863427b6e" }, ] +[[package]] +name = "mcp" +version = "1.28.1" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +dependencies = [ + { name = "anyio" }, + { name = "httpx" }, + { name = "httpx-sse" }, + { name = "jsonschema" }, + { name = "pydantic" }, + { name = "pydantic-settings" }, + { name = "pyjwt", extra = ["crypto"] }, + { name = "python-multipart" }, + { name = "pywin32", marker = "sys_platform == 'win32'" }, + { name = "sse-starlette" }, + { name = "starlette" }, + { name = "typing-extensions" }, + { name = "typing-inspection" }, + { name = "uvicorn", marker = "sys_platform != 'emscripten'" }, +] +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/6e/77/9450b8f251a13affb6281997d0523c4615f8a8b35d0b21ff30db3a5aac9d/mcp-1.28.1.tar.gz", hash = "sha256:d51e36a5f5644faea4f85ea649bfffa6bc6c26770d42798ad6a3de3d2ba69683" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/e2/5e/d118fce19f87a2e7d8101c35c8ae0ec289098a4df0ff244cec23e415aca0/mcp-1.28.1-py3-none-any.whl", hash = "sha256:2726bca5e7193f61c5dde8b12500a6de2d9acf6d1a1c0be9e8c2e706437991df" }, +] + [[package]] name = "mdurl" version = "0.1.2" @@ -2741,6 +2801,11 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/a3/5e/ecf12fdb62546d64385c158514e9b2b671f7832108ef2ecd2020ce0af2d1/pyjwt-2.13.0-py3-none-any.whl", hash = "sha256:66adcc2aff09b3f1bbd95fc1e1577df8ac8723c978552fd43304c8a290ac5728" }, ] +[package.optional-dependencies] +crypto = [ + { name = "cryptography" }, +] + [[package]] name = "pylint" version = "4.0.6" @@ -2957,6 +3022,15 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/0b/d7/1959b9648791274998a9c3526f6d0ec8fd2233e4d4acce81bbae76b44b2a/python_dotenv-1.2.2-py3-none-any.whl", hash = "sha256:1d8214789a24de455a8b8bd8ae6fe3c6b69a5e3d64aa8a8e5d68e694bbcb285a" }, ] +[[package]] +name = "python-multipart" +version = "0.0.32" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/5b/42/55c32bb9b12693c092ad250a0e82edb5b31ddeda6eb772de5f308b3804ad/python_multipart-0.0.32.tar.gz", hash = "sha256:be54b7f3fa167bb83e4fcd936b887b708f4e57fe75911c02aebf53efaf8d938e" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/e1/04/e8135ebd1ad02c56ec633277529b2602ff99ff634be76cdba5744cf554fd/python_multipart-0.0.32-py3-none-any.whl", hash = "sha256:ff6d3f776f16878c894e52e107296ffc890e913c611b1a4ec6c44e2821fe2e23" }, +] + [[package]] name = "pytz" version = "2026.2" @@ -3333,6 +3407,18 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/e2/22/dbf013a12ec759e54a34a119e9e217435b3f71b2dd5c61a7ade0a25dae87/sqlalchemy-2.0.51-py3-none-any.whl", hash = "sha256:bb024d8b621d0be75f4f44ecc7c950450026e76d66dc8f791bb5331d7fed59d5" }, ] +[[package]] +name = "sqlite-vec" +version = "0.1.9" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/68/85/9fad0045d8e7c8df3e0fa5a56c630e8e15ad6e5ca2e6106fceb666aa6638/sqlite_vec-0.1.9-py3-none-macosx_10_6_x86_64.whl", hash = "sha256:1b62a7f0a060d9475575d4e599bbf94a13d85af896bc1ce86ee80d1b5b48e5fb" }, + { url = "https://mirrors.aliyun.com/pypi/packages/a4/3d/3677e0cd2f92e5ebc43cd29fbf565b75582bff1ccfa0b8327c7508e1084f/sqlite_vec-0.1.9-py3-none-macosx_11_0_arm64.whl", hash = "sha256:1d52e30513bae4cc9778ddbf6145610434081be4c3afe57cd877893bad9f6b6c" }, + { url = "https://mirrors.aliyun.com/pypi/packages/00/d4/f2b936d3bdc38eadcbd2a87875815db36430fab0363182ba5d12cd8e0b51/sqlite_vec-0.1.9-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:4e921e592f24a5f9a18f590b6ddd530eb637e2d474e3b1972f9bbeb773aa3cb9" }, + { url = "https://mirrors.aliyun.com/pypi/packages/6f/ad/6afd073b0f817b3e03f9e37ad626ae341805891f23c74b5292818f49ac63/sqlite_vec-0.1.9-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.manylinux1_x86_64.whl", hash = "sha256:1515727990b49e79bcaf75fdee2ffc7d461f8b66905013231251f1c8938e7786" }, + { url = "https://mirrors.aliyun.com/pypi/packages/42/89/81b2907cda14e566b9bf215e2ad82fc9b349edf07d2010756ffdb902f328/sqlite_vec-0.1.9-py3-none-win_amd64.whl", hash = "sha256:4a28dc12fa4b53d7b1dced22da2488fade444e96b5d16fd2d698cd670675cf32" }, +] + [[package]] name = "sqlmodel" version = "0.0.39" @@ -3347,6 +3433,19 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/cf/7d/b9813a582d4eb310be35e1fc7dfaae71207d7b62e9e53be314ebd251b53b/sqlmodel-0.0.39-py3-none-any.whl", hash = "sha256:90ebe92ce5cc11d7fff8dc7cb594790a102333c8fe7c14865254f6fc5c939795" }, ] +[[package]] +name = "sse-starlette" +version = "3.4.6" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +dependencies = [ + { name = "anyio" }, + { name = "starlette" }, +] +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/6c/10/a34c656829ffc1c4b22ef36d70d9ebb6b99c020e2aeb17cee5485099f028/sse_starlette-3.4.6.tar.gz", hash = "sha256:725f8a1bd6d26ae1b2c9610c0ef5065dfdd496f3988d28adcf8c4b49dc25c627" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/49/36/e10c1d1b7ca881d2625db2ec28508578499187bb1c389952c398474e1834/sse_starlette-3.4.6-py3-none-any.whl", hash = "sha256:56217ab4c9a9f9c5db7b21e08732d3e7c2b807f45231ad23de0551a24c4a41f6" }, +] + [[package]] name = "starlette" version = "1.3.1" @@ -3557,6 +3656,33 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/f8/ba/d69adbe699b768f6b29a5eec7b47dd610bd17a69de51b251126a801369ea/uvloop-0.22.1-cp313-cp313-musllinux_1_2_x86_64.whl", hash = "sha256:1f38ec5e3f18c8a10ded09742f7fb8de0108796eb673f30ce7762ce1b8550cad" }, ] +[[package]] +name = "watchdog" +version = "6.0.0" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/db/7d/7f3d619e951c88ed75c6037b246ddcf2d322812ee8ea189be89511721d54/watchdog-6.0.0.tar.gz", hash = "sha256:9ddf7c82fda3ae8e24decda1338ede66e1c99883db93711d8fb941eaa2d8c282" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/e0/24/d9be5cd6642a6aa68352ded4b4b10fb0d7889cb7f45814fb92cecd35f101/watchdog-6.0.0-cp311-cp311-macosx_10_9_universal2.whl", hash = "sha256:6eb11feb5a0d452ee41f824e271ca311a09e250441c262ca2fd7ebcf2461a06c" }, + { url = "https://mirrors.aliyun.com/pypi/packages/63/7a/6013b0d8dbc56adca7fdd4f0beed381c59f6752341b12fa0886fa7afc78b/watchdog-6.0.0-cp311-cp311-macosx_10_9_x86_64.whl", hash = "sha256:ef810fbf7b781a5a593894e4f439773830bdecb885e6880d957d5b9382a960d2" }, + { url = "https://mirrors.aliyun.com/pypi/packages/d1/40/b75381494851556de56281e053700e46bff5b37bf4c7267e858640af5a7f/watchdog-6.0.0-cp311-cp311-macosx_11_0_arm64.whl", hash = "sha256:afd0fe1b2270917c5e23c2a65ce50c2a4abb63daafb0d419fde368e272a76b7c" }, + { url = "https://mirrors.aliyun.com/pypi/packages/39/ea/3930d07dafc9e286ed356a679aa02d777c06e9bfd1164fa7c19c288a5483/watchdog-6.0.0-cp312-cp312-macosx_10_13_universal2.whl", hash = "sha256:bdd4e6f14b8b18c334febb9c4425a878a2ac20efd1e0b231978e7b150f92a948" }, + { url = "https://mirrors.aliyun.com/pypi/packages/12/87/48361531f70b1f87928b045df868a9fd4e253d9ae087fa4cf3f7113be363/watchdog-6.0.0-cp312-cp312-macosx_10_13_x86_64.whl", hash = "sha256:c7c15dda13c4eb00d6fb6fc508b3c0ed88b9d5d374056b239c4ad1611125c860" }, + { url = "https://mirrors.aliyun.com/pypi/packages/5b/7e/8f322f5e600812e6f9a31b75d242631068ca8f4ef0582dd3ae6e72daecc8/watchdog-6.0.0-cp312-cp312-macosx_11_0_arm64.whl", hash = "sha256:6f10cb2d5902447c7d0da897e2c6768bca89174d0c6e1e30abec5421af97a5b0" }, + { url = "https://mirrors.aliyun.com/pypi/packages/68/98/b0345cabdce2041a01293ba483333582891a3bd5769b08eceb0d406056ef/watchdog-6.0.0-cp313-cp313-macosx_10_13_universal2.whl", hash = "sha256:490ab2ef84f11129844c23fb14ecf30ef3d8a6abafd3754a6f75ca1e6654136c" }, + { url = "https://mirrors.aliyun.com/pypi/packages/85/83/cdf13902c626b28eedef7ec4f10745c52aad8a8fe7eb04ed7b1f111ca20e/watchdog-6.0.0-cp313-cp313-macosx_10_13_x86_64.whl", hash = "sha256:76aae96b00ae814b181bb25b1b98076d5fc84e8a53cd8885a318b42b6d3a5134" }, + { url = "https://mirrors.aliyun.com/pypi/packages/fe/c4/225c87bae08c8b9ec99030cd48ae9c4eca050a59bf5c2255853e18c87b50/watchdog-6.0.0-cp313-cp313-macosx_11_0_arm64.whl", hash = "sha256:a175f755fc2279e0b7312c0035d52e27211a5bc39719dd529625b1930917345b" }, + { url = "https://mirrors.aliyun.com/pypi/packages/a9/c7/ca4bf3e518cb57a686b2feb4f55a1892fd9a3dd13f470fca14e00f80ea36/watchdog-6.0.0-py3-none-manylinux2014_aarch64.whl", hash = "sha256:7607498efa04a3542ae3e05e64da8202e58159aa1fa4acddf7678d34a35d4f13" }, + { url = "https://mirrors.aliyun.com/pypi/packages/5c/51/d46dc9332f9a647593c947b4b88e2381c8dfc0942d15b8edc0310fa4abb1/watchdog-6.0.0-py3-none-manylinux2014_armv7l.whl", hash = "sha256:9041567ee8953024c83343288ccc458fd0a2d811d6a0fd68c4c22609e3490379" }, + { url = "https://mirrors.aliyun.com/pypi/packages/d4/57/04edbf5e169cd318d5f07b4766fee38e825d64b6913ca157ca32d1a42267/watchdog-6.0.0-py3-none-manylinux2014_i686.whl", hash = "sha256:82dc3e3143c7e38ec49d61af98d6558288c415eac98486a5c581726e0737c00e" }, + { url = "https://mirrors.aliyun.com/pypi/packages/ab/cc/da8422b300e13cb187d2203f20b9253e91058aaf7db65b74142013478e66/watchdog-6.0.0-py3-none-manylinux2014_ppc64.whl", hash = "sha256:212ac9b8bf1161dc91bd09c048048a95ca3a4c4f5e5d4a7d1b1a7d5752a7f96f" }, + { url = "https://mirrors.aliyun.com/pypi/packages/2c/3b/b8964e04ae1a025c44ba8e4291f86e97fac443bca31de8bd98d3263d2fcf/watchdog-6.0.0-py3-none-manylinux2014_ppc64le.whl", hash = "sha256:e3df4cbb9a450c6d49318f6d14f4bbc80d763fa587ba46ec86f99f9e6876bb26" }, + { url = "https://mirrors.aliyun.com/pypi/packages/62/ae/a696eb424bedff7407801c257d4b1afda455fe40821a2be430e173660e81/watchdog-6.0.0-py3-none-manylinux2014_s390x.whl", hash = "sha256:2cce7cfc2008eb51feb6aab51251fd79b85d9894e98ba847408f662b3395ca3c" }, + { url = "https://mirrors.aliyun.com/pypi/packages/b5/e8/dbf020b4d98251a9860752a094d09a65e1b436ad181faf929983f697048f/watchdog-6.0.0-py3-none-manylinux2014_x86_64.whl", hash = "sha256:20ffe5b202af80ab4266dcd3e91aae72bf2da48c0d33bdb15c66658e685e94e2" }, + { url = "https://mirrors.aliyun.com/pypi/packages/07/f6/d0e5b343768e8bcb4cda79f0f2f55051bf26177ecd5651f84c07567461cf/watchdog-6.0.0-py3-none-win32.whl", hash = "sha256:07df1fdd701c5d4c8e55ef6cf55b8f0120fe1aef7ef39a1c6fc6bc2e606d517a" }, + { url = "https://mirrors.aliyun.com/pypi/packages/db/d9/c495884c6e548fce18a8f40568ff120bc3a4b7b99813081c8ac0c936fa64/watchdog-6.0.0-py3-none-win_amd64.whl", hash = "sha256:cbafb470cf848d93b5d013e2ecb245d4aa1c8fd0504e863ccefa32445359d680" }, + { url = "https://mirrors.aliyun.com/pypi/packages/33/e8/e40370e6d74ddba47f002a32919d91310d6074130fe4e17dabcafc15cbf1/watchdog-6.0.0-py3-none-win_ia64.whl", hash = "sha256:a1914259fa9e1454315171103c6a30961236f508b9b623eae470268bbcc6a22f" }, +] + [[package]] name = "watchfiles" version = "1.2.0"