diff --git a/app/infrastructure/graph.py b/app/infrastructure/graph.py new file mode 100644 index 0000000..788b719 --- /dev/null +++ b/app/infrastructure/graph.py @@ -0,0 +1,64 @@ +"""Neo4j 图驱动:把 neo4j 异步驱动封装成 `RelationshipService` 期望的 `GraphDriver` 协议。 + +为什么单独一层:`RelationshipService` 只依赖 `execute_query(query, **params)` 这个窄协议, +不关心底层是哪个驱动。有这一层,读服务、投影 worker、删除客户端就能共用同一套连接管理, +也便于测试注入替身。 + +**这层此前是空的** —— `app/service/relationship_service.py` 里的 `GraphDriver` 只有一个 +Protocol 声明,没有任何具体实现,组装层也没有装配它,所以读服务在运行期必然降级 +(`neo4j_unavailable`),这正是"Neo4j 连读适配器都没有"的根因。 + +降级语义与 Milvus 侧保持一致:本层**不吞异常**(由调用方决定降级方式),但**构造失败返回 None** +——图库不可用不该让应用起不来,也不该阻塞主链路。 +""" + +import logging +from collections.abc import Sequence +from typing import Any + +from app.core.config import get_settings + +logger = logging.getLogger(__name__) + + +class Neo4jGraphDriver: + """`GraphDriver` 协议的 neo4j 实现。""" + + def __init__(self, driver: Any, database: str) -> None: + self._driver = driver + self._database = database + + async def execute_query(self, query: str, **parameters: Any) -> Sequence[Any]: + """执行一条 Cypher 并返回记录列表。 + + 查询文本由调用方提供(`RelationshipService` 与投影 worker 各自持有白名单: + 前者限制关系类型,后者只接受 `ALLOWED_RELATIONSHIPS` 中的关系),本层不做校验, + 只负责执行与结果整形——把校验放在拥有业务语义的地方,避免两处规则漂移。 + """ + async with self._driver.session(database=self._database) as session: + result = await session.run(query, **parameters) + return list(await result.data()) + + async def close(self) -> None: + await self._driver.close() + + +def build_graph_driver() -> Neo4jGraphDriver | None: + """按配置构造图驱动;依赖缺失或配置不全时返回 None(图能力关闭)。 + + 与 `get_vector_memory_adapter` 同一取向:外部依赖不可用只关闭对应能力, + 不让整个应用起不来。 + """ + settings = get_settings() + if not settings.neo4j_uri or not settings.neo4j_password: + logger.warning("neo4j config incomplete; graph capability disabled") + return None + try: + from neo4j import AsyncGraphDatabase + except ImportError: + logger.warning("neo4j driver not installed; graph capability disabled") + return None + driver = AsyncGraphDatabase.driver( + settings.neo4j_uri, auth=(settings.neo4j_username, settings.neo4j_password) + ) + return Neo4jGraphDriver(driver, settings.neo4j_database) diff --git a/app/service/agent/bootstrap.py b/app/service/agent/bootstrap.py index 48f7ce2..76fcab7 100644 --- a/app/service/agent/bootstrap.py +++ b/app/service/agent/bootstrap.py @@ -9,6 +9,7 @@ from app.core.errors import RecoverableAgentError from app.core.fund_contracts import FundQuoteQuery from app.core.knowledge_contracts import KnowledgeSearchInput from app.infrastructure.fund_quote_cache import FundQuoteCache +from app.infrastructure.graph import build_graph_driver from app.infrastructure.memory_cache import MemoryCacheAdapter from app.infrastructure.vector_memory import VectorMemoryAdapter from app.service.agent.factory import AgentFactory @@ -27,6 +28,7 @@ from app.service.model_gateway import ( ModelEmbeddingService, ModelGenerationService, ) +from app.service.relationship_service import RelationshipService from app.service.runtime_config_service import load_active_intent_configs from app.service.suitability_service import SuitabilityToolInput, suitability_tool_handler from app.service.tool_executor import ToolDefinition, ToolExecutor, ToolRegistry @@ -126,6 +128,22 @@ def get_knowledge_search_service() -> KnowledgeSearchService: return KnowledgeSearchService(client, _embed_text) +@lru_cache(maxsize=1) +def get_relationship_service() -> RelationshipService | None: + """图关系读服务:客户 → 产品/标签/事件 的多跳查询入口。 + + 驱动构造失败时返回 None(图能力关闭),由调用方降级——图库不可用不该阻塞主链路, + 与 Milvus 侧"语义通道缺失不影响结构化召回"是同一取向。 + + 注意:本服务**只读**,且关系类型受 `RelationshipService.ALLOWED_RELATIONSHIPS` 白名单约束; + 写入走 `GraphProjectionWorker`(由领域事件驱动),这里不提供任意写接口。 + """ + driver = build_graph_driver() + if driver is None: + return None + return RelationshipService(driver) + + def build_memory_recall_service(session: AsyncSession) -> MemoryRecallService: """记忆召回组装:结构化召回始终可用,Redis 缓存与语义通道可用时叠加。