merge: 合并主干 qyqy_develop(PR #7 之后)并对齐两套投影实现

共同祖先 bbf623a;主干 54 个提交、118 个文件;本线 25 个文件;9 个冲突文件。
主干这次把 **ZSY 的整条投影实现合进来了(PR #7)**,而本线此前的提交正是
移植并修正同一套代码 —— 因此冲突的本质是"同一功能两份实现并存",取舍错了会把
已修好的缺陷又带回来。逐项取舍与理由见 `docs/39-主干合并对策记录.md`。

## 取舍(9 个冲突)

取本线:
- `app/infrastructure/milvus_profile_projection.py` —— 主干是 ZSY 原版,含两处必炸点:
  ① `customer_id` 要求 int 而本仓所有生产者都写 `str` ⇒ 每个事件必然失败;
  ② 不可投影的 `memory_key` 直接 raise ⇒ 一条 `constraint:` 记忆毒死整客户整批。
  本线版已放宽为「接受纯数字字符串」与「跳过并留痕」。
- `memory_sync_outbox_worker.py` / `conversation_privacy.py` / `risk_questionnaire.py`
  —— 代码逐行一致,仅注释与说明文字详略不同(`risk_questionnaire.py` 两边**独立做了
  完全相同的修复**,都改成 re-export `app.model.profile`)。
- 两个投影测试文件 —— 本线是他那份的**超集**(4→10、4→5 例,包含他全部用例)。

两边合并:
- `app/worker/runtime.py`:`__init__` 两边各加一个参数,都要。
- `app/service/agent/implementations/customer_service.py`:import 取并集;
  `COMPANY` 取主干的「奶龙基金责任有限公司」("奶龙"是本项目实际品牌名,主干多处出现),
  `HOTLINE`/`SERVICE_HOURS` **取本线的修复**(主干仍是占位符 `400-XXX-XXXX`,
  本线已改为引用 `customer_service_rules` 的唯一来源 —— 这是 A1 缺陷修复,
  否则同一客服给客户两个不同号码)。
- `AGENTS.md`:表数/Agent 清单取主干(90/89、7 个 Agent),本线的
  `-X utf8` 与两条 outbox 易错点保留,测试基线按合并后实测重算。

## 消费端只保留一套(本次最重要的一处)

合并后曾出现**两套消费者读同一个 `memory_sync_outbox`**:`__main__.py`(PR #7)
与 `runtime.consume_profile_projections()`(本线),而**两者的 neo4j handler 不同**
—— 前者用 ZSY 的 `Neo4jProfileProjection`(按客户各建私有节点),
后者用主干 `ProfileGraphProjectionService`(共享 tag 节点、只投影已确认事实)。
同一事件被谁领到结果不定,等于"同一事实在图里有两种说法",正是**方案 A 要避免的状态**。

现只保留 runtime 那一套(带 `memory_sources` 兜底、neo4j 复用主干服务),
删除 `__main__.py` 的重复接线;装配入口职责仍在该文件(注入 `relationships` /
`projection_cleaner`),Milvus 客户端由 `bootstrap` 工厂惰性构造、缺配置时显式降级。

副作用:`app/infrastructure/neo4j_profile_projection.py` 不再被生产代码引用,成为
**死代码**(本线未删,属架构师线,其单测仍在)—— 待架构师决定删或明确分工。

## 顺带修掉的 3 个继承缺陷(主干同样存在,PR #7 后未整套复跑故未发现)

1. `tools/seed_test_rbac.py` **少建 `review_t`(9004)账号** —— 两个集成测试都依赖它
   ("账号存在但无权限应返回 200 空集而非 404"、`PLACEHOLDER_ACCOUNTS`)。
   同时把用户↔角色绑定从 `zip(..., strict=True)` 改为**显式配对表**:原写法隐含
   "USERS 与 ROLES 一一对应",一加不绑角色的账号就 ValueError、整个种子跑不完
   (commit 在最后,外部表现是"什么都没发生")。
2. `CustomerProfileCandidateService._write_profile_snapshot` **漏写 `current_customer_id`**
   —— 该列不是生成列而是普通可空列 + 唯一键 `uk_profile_snapshot_current`,
   不写则唯一键形同虚设(多个 NULL 不冲突),且旧当前版本也没清该列、补写就会撞键。
   现旧值置 None、新值显式写入(与 `ProfileGenerationService._clear_current` 一致)。
3. 集成测试前置未记录 —— 13 个登录/RBAC 用例因 401 而红,实为"测试账号不存在",
   跑 `seed_test_rbac.py` + `set_user_password.py` 后转绿;已在 `AGENTS.md` 记明,
   避免被误判成代码缺陷。

## 文档

- 新增 `docs/39-主干合并对策记录.md`(逐文件取舍 + 理由 + 遗留)
- `docs/37` 订正一处过时说法:曾写 `current_customer_id` 无人使用且故意不映射,
  实际 `app/model/profile.py` 已映射且有人使用(详见该文档 §6.2 的订正块)
- 文档编号:主干已占 29–36,本线两份文档让号至 `docs/37`、`docs/38`

## 验证(合并后实测)

- `pytest tests`(全量)→ `2 failed, 1396 passed, 2 skipped`
- `pytest tests/integration` → `102 passed, 1 skipped`(修上述 1、2 后从 15 failed 归零)
- `mypy app` → `Success: no issues found in 245 source files`
- `tools/audit_schema.py` → 89 张业务表无缺失/意外(未改动任何表结构)
- `tools/check_authoritative_docs.py` → 52 份文档无编号冲突
- `tools/check_rbac_seed_consistency.py` → 通过

那 2 个失败是既有环境项(`test_offsite_document_recognition_adapter.py` 断言请求体
中文原文而 httpx 序列化成 `\uXXXX`),与本次合并无关。
This commit is contained in:
2026-09-12 13:14:57 +08:00
112 changed files with 11514 additions and 90 deletions
+135 -16
View File
@@ -10,12 +10,17 @@ from uuid import uuid4
from sqlalchemy import select, update
from app.core.config import Settings, get_settings
from app.core.contracts import AgentRequest, AgentResult, RequestContext
from app.core.contracts import AgentRequest, AgentRequestMetadata, AgentResult, RequestContext
from app.core.errors import AgentError, RecoverableAgentError, RunLeaseLostError
from app.infrastructure.db import SessionFactory
from app.model.audit import InteractionAudit
from app.model.conversation import ConversationMessage
from app.model.platform import AgentRun, DomainEventOutbox, RequestIdempotency
from app.model.platform import (
AgentRun,
DomainEventOutbox,
HandoverTicket,
RequestIdempotency,
)
from app.repository.agent_run_repository import AgentRunRepository
from app.service.agent.bootstrap import (
get_agent_factory,
@@ -28,6 +33,11 @@ from app.service.agent.bootstrap import (
from app.service.agent.executor import AgentExecutor
from app.service.agent.factory import AgentFactory
from app.service.agent_persistence_service import AgentPersistenceService
from app.service.customer_service_session_memory_service import (
CustomerServiceSessionMemory,
CustomerServiceSessionTurn,
build_customer_service_session_memory,
)
from app.service.identity_service import IdentityService
from app.service.memory_extraction_service import (
ExtractionEndpointResolver,
@@ -38,6 +48,7 @@ from app.service.memory_recall_service import MemoryRecallService
from app.service.memory_service import CacheDeleteAdapter, MemoryService
from app.service.memory_taxonomy import BUSINESS_EVENT_TYPES
from app.service.model_gateway import DatabaseModelEndpointResolver, ModelGenerationService
from app.worker.customer_profile_candidate_worker import CustomerProfileCandidateWorker
from app.worker.episode_worker import (
EpisodeConsumptionResult,
EpisodeExtractionConsumer,
@@ -93,6 +104,7 @@ class WorkerRuntime:
knowledge_embedder: Any = _UNSET,
knowledge_endpoint_resolver: Any = _UNSET,
profile_vector_client: Any = _UNSET,
session_memory: CustomerServiceSessionMemory | None = None,
) -> None:
self.factory = factory if factory is not None else get_agent_factory()
self.settings = settings or get_settings()
@@ -114,6 +126,11 @@ class WorkerRuntime:
# 这里不兜底。组件内部给默认实现会把"尚未装配"这一事实悄悄盖住——而"未注入即显式
# 降级并留痕"是本模块刻意保留的语义(有单测守着),因此默认值保持 None。
self.projection_cleaner = projection_cleaner
# 客服短期会话 Redis 仅用于当前会话上下文,不参与长期画像召回。
self.session_memory = (
session_memory if session_memory is not None
else build_customer_service_session_memory()
)
# 图关系服务:画像投影用它写入节点与关系(投顾的多跳推荐、风控的关系网络都读它)。
# 默认取生产装配;图库不可用时该值为 None,投影如实降级而不是失败。
if relationships is not None:
@@ -155,6 +172,48 @@ class WorkerRuntime:
# episode 聚合是低频批处理,按轮次节流而不是每轮都查。
self._episode_rounds = 0
async def restore_context(
self, *, actor_type: str, actor_id: str, trace_id: str
) -> RequestContext:
"""按受理事件中的可信身份恢复最小执行权限。"""
identity = RequestContext(user_id=actor_id, trace_id=trace_id)
if actor_type == "visitor":
return identity.model_copy(update={
"roles": ("visitor",),
"permissions": ("agent:run", "knowledge:query"),
"data_scope": "public",
})
return await self.resolve_identity(identity)
@staticmethod
def should_request_memory_extraction(
*, agent_type: str, context: RequestContext, message: str,
result: AgentResult, business_events: tuple[str, ...] | list[str],
) -> bool:
"""长期记忆抽取只接收非客服、非访客的明确业务事实。"""
if agent_type == "customer_service" or "visitor" in context.roles:
return False
return MemoryService.should_extract_memory(
conversation_content=message,
role="user",
tool_result=any(call.status == "succeeded" for call in result.result.tool_calls),
event_type=business_events[0] if business_events else None,
signals=MemoryService.detect_memory_signals(message),
)
@staticmethod
def should_request_profile_candidate(
*, agent_type: str, context: RequestContext, message: str,
) -> bool:
"""客服仅为已登录且 self 范围内的用户生成待确认画像候选。"""
if agent_type != "customer_service" or "visitor" in context.roles:
return False
if not {"customer", "authenticated_user"}.intersection(context.roles):
return False
if context.data_scope != "self":
return False
return bool(MemoryService.detect_memory_signals(message))
async def dispatch_one(self, *, run_id: str | None = None) -> bool:
# Outbox acknowledges a durable SQL queue entry, not an in-memory task.
async with SessionFactory() as session:
@@ -172,6 +231,13 @@ class WorkerRuntime:
session, extractor=self.memory_extraction, cache=self.memory_cache
).handle(payload)
async def dispatch_profile_candidate(payload: dict[str, Any]) -> None:
if "message_id" not in payload or "customer_id" not in payload:
raise OutboxHandlerError("profile candidate payload is incomplete")
await CustomerProfileCandidateWorker(
session, extractor=self.memory_extraction, cache=self.memory_cache
).handle(payload)
async def dispatch_run_completed(payload: dict[str, Any]) -> None:
# 结果消息与审计已由 complete_run 同事务落库,此事件只承担
# "运行已完成"的对外通知职责。当前没有独立外部消费者,
@@ -210,6 +276,33 @@ class WorkerRuntime:
logger.warning("graph projection degraded customer_id=%s reason=%s",
customer_id, result.reason)
async def dispatch_handover_queue_ready(payload: dict[str, Any]) -> None:
"""记录转人工队列已就绪;不向客户承诺已接单或处理时限。"""
ticket_no = str(payload.get("ticket_no", "")).strip()
if not ticket_no:
raise OutboxHandlerError(
"conversation.transfer_requested payload is incomplete"
)
ticket = await session.scalar(
select(HandoverTicket).where(HandoverTicket.ticket_no == ticket_no)
)
if ticket is None:
raise OutboxHandlerError("handover ticket not found")
session.add(InteractionAudit(
actor_type="system", actor_id=None,
target_customer_id=ticket.customer_id,
session_id=ticket.session_id, portal="worker",
action_type="handover.queue_ready",
detail={
"ticket_no": ticket.ticket_no,
"source_agent": ticket.source_agent,
"reason_code": ticket.reason_code,
"ticket_status": ticket.status,
},
created_at=datetime.now(UTC).replace(tzinfo=None),
))
await session.flush()
async def dispatch_projection_cleanup(payload: dict[str, Any]) -> None:
# memory.invalidated / memory.deleted 由 MemoryLifecycleService 按
# memory_uuid 写入,这里做幂等的投影清理(Milvus 向量、Neo4j 关系)。
@@ -234,6 +327,7 @@ class WorkerRuntime:
handlers: dict[str, Callable[[dict[str, Any]], Awaitable[None]]] = {
"agent.run_requested": dispatch,
"memory.extraction_requested": dispatch_memory_extraction,
"customer_profile.candidate_requested": dispatch_profile_candidate,
"agent.run_completed": dispatch_run_completed,
"config.cache_invalidate_requested": dispatch_cache_invalidate,
"memory.deletion_requested": dispatch_memory_deletion,
@@ -242,6 +336,7 @@ class WorkerRuntime:
"memory.deleted": dispatch_projection_cleanup,
# 画像重建:记忆写入后自动触发,使「记忆 → 画像 → 图」全链路无需手工介入
"profile.rebuild_requested": dispatch_profile_rebuild,
"conversation.transfer_requested": dispatch_handover_queue_ready,
}
# 知识向量同步/删除:Task 5 交付了 handler 与写适配器,但先前没有任何生产装配
# 调用它们 —— 事件类型不在上面的白名单里,`OutboxWorker.publish_one` 的
@@ -720,12 +815,19 @@ class WorkerRuntime:
request = AgentRequest(
agent_type=run.agent_type, message=message.content, session_id=run.session_id,
idempotency_key=idem.idempotency_key,
metadata=event.payload.get("metadata", {}) if event else {},
metadata=AgentRequestMetadata.model_validate(
event.payload.get("metadata", {}) if event else {}
),
history=history,
)
identity = RequestContext(user_id=str(run.user_id), trace_id=run.trace_id)
actor_type = (
str(event.payload.get("actor_type", "authenticated"))
if event else "authenticated"
)
# Re-check account and permissions at execution time, including delayed jobs.
context = await self.resolve_identity(identity)
context = await self.restore_context(
actor_type=actor_type, actor_id=str(run.user_id), trace_id=run.trace_id
)
result: AgentResult | None = None
async for event_data in AgentExecutor(self.factory).execute(
request.agent_type, request, context, run_id
@@ -747,19 +849,36 @@ class WorkerRuntime:
async with SessionFactory() as session:
await AgentPersistenceService(session).complete_run(
run_id, result, worker_id=worker_id,
memory_extraction_requested=MemoryService.should_extract_memory(
conversation_content=request.message,
role="user",
# 工具产出的权威事实同样构成持久记忆(工具调用记录来自终态结果)。
tool_result=any(
call.status == "succeeded" for call in result.result.tool_calls
),
# 本 run 落库的业务事件(风险评估完成、交易完成等)。
event_type=business_events[0] if business_events else None,
# 用户明确陈述的偏好/约束/身份/目标,命中才触发抽取。
signals=MemoryService.detect_memory_signals(request.message),
memory_extraction_requested=self.should_request_memory_extraction(
agent_type=run.agent_type, context=context, message=request.message,
result=result, business_events=business_events,
),
profile_candidate_requested=self.should_request_profile_candidate(
agent_type=run.agent_type, context=context, message=request.message,
),
)
await self._append_customer_service_session_memory(
agent_type=run.agent_type, actor_id=str(run.user_id), session_id=run.session_id,
request_message=request.message, response_message=result.result.text,
)
async def _append_customer_service_session_memory(
self, *, agent_type: str, actor_id: str, session_id: str,
request_message: str, response_message: str,
) -> None:
"""成功落库后追加短期会话;Redis 故障不影响主事务。"""
if agent_type != "customer_service":
return
try:
await self.session_memory.append(
actor_id=actor_id, session_id=session_id,
turns=(
CustomerServiceSessionTurn(role="user", content=request_message),
CustomerServiceSessionTurn(role="assistant", content=response_message),
),
)
except Exception:
logger.warning("客服短期会话写入降级,不影响已完成的客服运行", exc_info=True)
async def _failure(
self, run_id: str, worker_id: str, error_code: str, *, retryable: bool