diff --git a/app/api/controllers/conversations.py b/app/api/controllers/conversations.py index f5375db..e123144 100644 --- a/app/api/controllers/conversations.py +++ b/app/api/controllers/conversations.py @@ -5,6 +5,7 @@ from app.api.dependencies.auth import build_request_context from app.api.dependencies.database import get_session from app.api.dependencies.rate_limit import enforce_rate_limit from app.api.schemas.conversations import FeedbackRequest +from app.api.views.envelope import envelope, list_envelope from app.core.contracts import RequestContext from app.core.cursor import parse_cursor from app.service.conversation_service import ConversationService @@ -28,9 +29,10 @@ async def list_messages( `400 INVALID_CURSOR`,而不是被静默忽略后返回第一页。 """ before = parse_cursor(cursor) - return await ConversationService(session).messages( + page = await ConversationService(session).messages( session_id, context, limit, before=before ) + return list_envelope(page, context) @router.post("/conversation-messages/{message_id}/feedback", status_code=status.HTTP_201_CREATED) @@ -40,5 +42,7 @@ async def create_feedback( context: RequestContext = Depends(build_request_context), # noqa: B008 session: AsyncSession = Depends(get_session), # noqa: B008 ) -> dict[str, object]: - return await ConversationService(session).feedback( + # 与消息列表同理:§3.3 的信封由 Controller 统一套,service 只负责业务数据。 + data = await ConversationService(session).feedback( message_id, context, payload.rating, payload.feedback_type, payload.feedback_content) + return envelope(data, context) diff --git a/app/api/controllers/knowledge.py b/app/api/controllers/knowledge.py index af6fa49..7375b0e 100644 --- a/app/api/controllers/knowledge.py +++ b/app/api/controllers/knowledge.py @@ -2,6 +2,7 @@ from fastapi import APIRouter, Depends, Path from app.api.dependencies.auth import build_request_context from app.api.dependencies.rate_limit import enforce_rate_limit +from app.api.views.envelope import envelope from app.core.contracts import RequestContext from app.service.knowledge_service import KnowledgeReferenceService @@ -14,4 +15,7 @@ async def resolve_reference( reference_token: str = Path(min_length=20, max_length=300), context: RequestContext = Depends(build_request_context), # noqa: B008 ) -> dict[str, object]: - return await KnowledgeReferenceService().resolve(context, reference_token) + # §3.3:成功响应也要有 `meta.trace_id`。此前这里直接返回资源对象,客户端拿不到 + # 本次请求的追踪标识,出问题时无法与服务端日志对上。 + data = await KnowledgeReferenceService().resolve(context, reference_token) + return envelope(data, context) diff --git a/app/api/controllers/risk.py b/app/api/controllers/risk.py index 878debe..d9f380e 100644 --- a/app/api/controllers/risk.py +++ b/app/api/controllers/risk.py @@ -24,6 +24,8 @@ from app.api.schemas.risk import ( RiskEvidenceSource, RiskNotificationPageQuery, ) +from app.api.views.envelope import envelope as _envelope +from app.api.views.envelope import list_envelope as _list_envelope from app.core.contracts import RequestContext from app.core.errors import SseNotAcceptableError from app.infrastructure.db import mysql_scan_lock @@ -304,28 +306,5 @@ async def send_risk_daily_report_mail( return _envelope(data, context) -def _envelope(data: object, context: RequestContext) -> dict[str, object]: - return { - "data": data, - "meta": {"trace_id": context.trace_id}, - } - - -def _list_envelope(page: dict[str, Any], context: RequestContext) -> dict[str, object]: - """列表资源的信封(docs/05 §3.3)。 - - §3.3 的列表样例是 `data` 为**纯数组**、游标与 `has_more` 放在 `meta` 里,并且明确 - 「业务接口不得增加其他顶层字段」。而 `RiskQueryService._page` 返回的是 - `{items, next_cursor, has_more}` —— 整体塞进 `data` 后,游标跑进了**业务数据**里、 - `meta` 只剩 trace_id,两处都不符合契约。 - - 这里统一拆包;service 侧不必改(它继续返回那个内部结构,只是不再直接当 `data` 用)。 - """ - return { - "data": page.get("items") or [], - "meta": { - "trace_id": context.trace_id, - "next_cursor": page.get("next_cursor"), - "has_more": bool(page.get("has_more")), - }, - } +# `_envelope` / `_list_envelope` 已抽到 `app/api/views/envelope.py`,与客服链路共用同一份 +# §3.3 实现 —— 两处各写一份的结果就是其中一处漏了 `meta`。 diff --git a/app/api/views/envelope.py b/app/api/views/envelope.py new file mode 100644 index 0000000..8df42e9 --- /dev/null +++ b/app/api/views/envelope.py @@ -0,0 +1,40 @@ +"""统一响应信封(`docs/05` §3.3)。 + +§3.3 规定成功响应是 `{data, meta}`,列表的 `data` 为**纯数组**、游标与 `has_more` 放在 +`meta` 里,并且明确「业务接口不得增加其他顶层字段」。 + +放在这里而不是各 Controller 各写一份:同一个偏差已经出现过两次 —— service 返回 +`{items, next_cursor, has_more}`(或干脆只有 `data`)之后被直接当响应体返回,于是 +`meta` 要么只剩 trace_id、要么整个缺失。同一份契约不该有多份实现。 +""" + +from __future__ import annotations + +from typing import Any + +from app.core.contracts import RequestContext + + +def envelope(data: object, context: RequestContext) -> dict[str, object]: + """单资源 / 单对象的信封。""" + return { + "data": data, + "meta": {"trace_id": context.trace_id}, + } + + +def list_envelope(page: dict[str, Any], context: RequestContext) -> dict[str, object]: + """列表资源的信封:`data` 只放数组,分页元数据进 `meta`。 + + `page` 是 service 的内部结构 `{items, next_cursor, has_more}` —— service 不必改, + 只是不再把它整体当 `data` 用。缺失的键按空值处理,因此只返回 `{"data": [...]}` + 的旧 service 也不会炸。 + """ + return { + "data": page.get("items") or [], + "meta": { + "trace_id": context.trace_id, + "next_cursor": page.get("next_cursor"), + "has_more": bool(page.get("has_more")), + }, + } diff --git a/app/service/agent/implementations/customer_service.py b/app/service/agent/implementations/customer_service.py index b37cb95..36518b3 100644 --- a/app/service/agent/implementations/customer_service.py +++ b/app/service/agent/implementations/customer_service.py @@ -341,7 +341,26 @@ class CustomerServiceAgent(BaseAgent): # "还没测评/已过期",不能写成"您的等级为无"这种客户看不懂的句子。 lines = ["您目前没有在有效期内的风险测评结果。"] if decision.get("allowed"): - lines.insert(0, f"{product}为 {level_name},在您的风险承受能力范围内,可以购买。") + # "在您的风险承受能力范围内"只对 C ≥ R 成立。矩阵允许的越级档(C1→R2、 + # C2→R3)和豁免档(C3→R4、C4→R5)都**超出**了客户等级,一律说成"范围内" + # 是把监管口径讲错:客户会以为自己的测评等级本来就覆盖这只产品。 + if isinstance(customer_level, int) and customer_level >= risk_level: + lines.insert( + 0, + f"{product}为 {level_name},在您的风险承受能力范围内,可以购买。", + ) + elif decision.get("reason_code") == "SUITABLE_WITH_DISCLOSURE": + lines.insert( + 0, + f"{product}为 {level_name},高于您的风险承受能力等级。" + "按照投资者适当性管理规定,签署产品风险揭示书后可以购买。", + ) + else: + lines.insert( + 0, + f"{product}为 {level_name},虽然高于您的风险测评等级," + "但仍在《个人投资者适当性管理指南》匹配矩阵允许购买的范围内。", + ) if decision.get("required_disclosure"): lines.append("购买前需签署产品风险揭示书,具体请咨询您的客户经理。") if decision.get("requires_recording"): diff --git a/app/service/conversation_service.py b/app/service/conversation_service.py index d1ecaff..03aafe0 100644 --- a/app/service/conversation_service.py +++ b/app/service/conversation_service.py @@ -22,12 +22,34 @@ class ConversationService: async def messages( self, session_id: str, context: RequestContext, limit: int, before: int | None = None ) -> dict[str, object]: - """消息列表投影;`before` 是文档 §3.8 的游标(记录 ID 边界),已由入口校验。""" + """消息列表投影;`before` 是文档 §3.8 的游标(记录 ID 边界),已由入口校验。 + + 取 `limit + 1` 行来判断"还有没有更旧的":只看"取满没取满"会把恰好等于 limit 的 + 最后一页说成还有下一页。`next_cursor` 就是本页最后一条的 `message_id` —— + 游标语义是"取更旧的一页",所以它天然可续(集成测试正是这么翻页的)。 + + 返回的是内部结构 `{items, next_cursor, has_more}`,由 Controller 用 + `list_envelope` 拆成 `docs/05` §3.3 要求的 `{data, meta}`;此前这里直接返回 + `{"data": [...]}`,成功响应因此**完全没有 `meta.trace_id`**。 + """ rows = await self.repository.messages( - session_id, int(context.user_id), limit, before=before + session_id, int(context.user_id), limit + 1, before=before ) - return {"data": [{"message_id": str(row.id), "role": row.role, "content": row.content, - "created_at": row.created_at.isoformat() + "Z"} for row in rows]} + has_more = len(rows) > limit + page_rows = rows[:limit] + return { + "items": [ + { + "message_id": str(row.id), + "role": row.role, + "content": row.content, + "created_at": row.created_at.isoformat() + "Z", + } + for row in page_rows + ], + "next_cursor": str(page_rows[-1].id) if has_more and page_rows else None, + "has_more": has_more, + } async def feedback( self, message_id: int, context: RequestContext, rating: int, @@ -50,7 +72,9 @@ class ConversationService: self.session.add(feedback) self._audit(context, message.session_id, "conversation.feedback_created", {"feedback_no": feedback.feedback_no}) - return {"data": {"feedback_no": feedback.feedback_no, "status": feedback.status}} + # 只返回业务数据;§3.3 的信封由 Controller 套(此前这里自带 {"data": ...}, + # 于是成功响应没有 meta.trace_id)。 + return {"feedback_no": feedback.feedback_no, "status": feedback.status} # 转人工申请(POST /api/v1/conversations/{id}/handover-requests)的唯一实现 # 在 PublicPlatformService.write("handover", ...):它同一事务写工单 + Outbox diff --git a/app/service/suitability_service.py b/app/service/suitability_service.py index 984fa5c..a5287bc 100644 --- a/app/service/suitability_service.py +++ b/app/service/suitability_service.py @@ -33,6 +33,31 @@ RISK_LEVEL_SOURCE = "fin_risk_assessment" PROFESSIONAL_INVESTOR_SOURCE = "sys_user" CUSTOMER_SCOPE_EXEMPT_ROLES = frozenset({"admin", "super_admin"}) +# 《个人投资者适当性管理指南》第十二条匹配矩阵(客户等级 → 可购买的产品等级)。 +# +# 矩阵与第十四条的**第 2、3 款**一致:只禁止"低两个等级及以上"的越级,低一个等级要看 +# 档位(C1→R2、C2→R3 直接可买;C3→R4、C4→R5 需签风险揭示书)。真正冲突的是第十四条 +# **第 1 款**"必须大于或等于"—— 它与同一条第 2、3 款自相矛盾。 +# +# 2026-09-11 业务裁定:**客服回答按矩阵,风控扫描保留 C ≥ R**。 +# 理由是两者的职责不同:客服要给客户一个与知识库(`POL-AST-012`)一致的"能不能买", +# 矩阵才是客户看得见的口径;风控要发现的是"越级成交且留痕不全",用更严的 C ≥ R 去 +# 事后核查。改之前两边是**同一个 Agent 自相矛盾**:问"C1 能买什么产品"答"R1、R2 可买" +# (矩阵),问"C1 能买这只 R2 吗"却答"不能购买"(C ≥ R)。 +MATRIX_ALLOWED: dict[int, frozenset[int]] = { + 1: frozenset({1, 2}), + 2: frozenset({1, 2, 3}), + 3: frozenset({1, 2, 3, 4}), + 4: frozenset({1, 2, 3, 4, 5}), + 5: frozenset({1, 2, 3, 4, 5}), +} + +# 矩阵里标"⚠️ 需签署风险揭示书"的档位,即第十五条豁免档。 +MATRIX_NEEDS_DISCLOSURE: dict[int, frozenset[int]] = { + 3: frozenset({4}), + 4: frozenset({5}), +} + _INVESTOR_TYPE_PATTERN = re.compile(r"^C([1-5])$") AuthorityReason = Literal[ @@ -192,16 +217,21 @@ class SuitabilityService: professional_investor=True, assessment_valid_until=profile.valid_until, ) - if profile.customer_risk_level < request.product_risk_level: + # 按第十二条匹配矩阵裁决(见 MATRIX_ALLOWED 的说明):不再用"C < R 即拒绝", + # 那样会把矩阵允许的 C1→R2、C2→R3 以及豁免档 C3→R4、C4→R5 一起拒掉。 + level = profile.customer_risk_level + product_level = request.product_risk_level + if product_level not in MATRIX_ALLOWED.get(level, frozenset()): return self._denied("RISK_LEVEL_MISMATCH", request, profile) - required_disclosure = request.product_requires_disclosure + needs_disclosure = product_level in MATRIX_NEEDS_DISCLOSURE.get(level, frozenset()) + required_disclosure = request.product_requires_disclosure or needs_disclosure return SuitabilityDecision( allowed=True, - reason_code="SUITABLE", + reason_code="SUITABLE_WITH_DISCLOSURE" if needs_disclosure else "SUITABLE", required_disclosure=required_disclosure, requires_confirmation=request.requires_confirmation or required_disclosure, requires_recording=required_disclosure or request.requires_confirmation, - customer_risk_level=profile.customer_risk_level, + customer_risk_level=level, professional_investor=False, assessment_valid_until=profile.valid_until, ) diff --git a/app/worker/outbox_worker.py b/app/worker/outbox_worker.py index 3973e8a..a9929bf 100644 --- a/app/worker/outbox_worker.py +++ b/app/worker/outbox_worker.py @@ -8,6 +8,41 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.model.platform import DomainEventOutbox, OutboxDelivery +class OutboxHandlerError(ValueError): + """Handler 失败,且失败原因是**可以安全落库**的固定文案。 + + 继承 `ValueError` 而不是 `Exception`:这些失败(run not found、payload 不完整、 + mode 非法)本来就是 ValueError 语义,保持继承关系才不会改动既有的 `except + ValueError` 行为与断言。 + + 为什么需要这个类型:`last_error` 默认只记异常类名,因为异常消息可能含凭据、SQL + 语句或客户标识(`tests/unit/worker/test_outbox_worker.py` 里那条 + `RuntimeError("credential=do-not-log")` 就是守这条的)。 + + 但只记类名又不够 —— `dispatch`(run not found)、`dispatch_run_completed`、 + `dispatch_memory_extraction`、`dispatch_profile_rebuild` 抛的全是 ValueError,实测 + 库里 373 条死信的 `last_error` 都是裸的 `"ValueError"`,分不清是哪一处失败的。 + + 折中办法:handler 想让人看见原因时,抛这个类型,`reason` 由**代码写死**、不含任何 + 请求数据,于是可以落库;其余异常仍然只记类名。 + """ + + def __init__(self, reason: str) -> None: + super().__init__(reason) + self.reason = reason + + +def safe_error_text(exc: BaseException) -> str: + """把异常转成可落库的失败原因。 + + - `OutboxHandlerError` → `类名: 固定文案`(文案由代码写死,安全); + - 其他异常 → **只记类名**(消息可能含凭据/请求数据,不落库)。 + """ + if isinstance(exc, OutboxHandlerError): + return f"{type(exc).__name__}: {exc.reason}"[:500] + return type(exc).__name__ + + class OutboxWorker: """One dispatcher per event type, not a fan-out transport. @@ -71,7 +106,8 @@ class OutboxWorker: await self.session.flush() except Exception as exc: event.retry_count += 1 - event.last_error = type(exc).__name__ + # 只让"代码写死的固定文案"落库;异常消息可能含凭据,见 safe_error_text。 + event.last_error = safe_error_text(exc) event.status = "dead" if event.retry_count >= 5 else "failed" event.next_retry_at = now + timedelta(seconds=min(300, 2**event.retry_count)) else: diff --git a/app/worker/runtime.py b/app/worker/runtime.py index e63858b..bdedcb6 100644 --- a/app/worker/runtime.py +++ b/app/worker/runtime.py @@ -41,7 +41,7 @@ from app.worker.episode_worker import ( EpisodeWorker, ) from app.worker.memory_extraction_worker import MemoryExtractionWorker -from app.worker.outbox_worker import OutboxWorker +from app.worker.outbox_worker import OutboxHandlerError, OutboxWorker logger = logging.getLogger(__name__) @@ -119,11 +119,11 @@ class WorkerRuntime: async def dispatch(payload: dict[str, Any]) -> None: run = await AgentRunRepository(session).get(str(payload["run_id"])) if run is None: - raise ValueError("run not found") + raise OutboxHandlerError("run not found") async def dispatch_memory_extraction(payload: dict[str, Any]) -> None: if "message_id" not in payload or "customer_id" not in payload: - raise ValueError("memory extraction payload is incomplete") + raise OutboxHandlerError("memory extraction payload is incomplete") # 幂等键只认事件 id,由 worker 自己按 payload 回查,避免调用方漏传。 # 注入召回缓存适配器:写入生效后立即失效该客户的热缓存。 await MemoryExtractionWorker( @@ -135,7 +135,7 @@ class WorkerRuntime: # "运行已完成"的对外通知职责。当前没有独立外部消费者, # 这里显式消费以免事件永久滞留;接入推送链路时在此处扩展。 if not str(payload.get("run_id", "")): - raise ValueError("agent.run_completed payload is incomplete") + raise OutboxHandlerError("agent.run_completed payload is incomplete") async def dispatch_cache_invalidate(payload: dict[str, Any]) -> None: await self._invalidate_config_cache(payload) @@ -149,7 +149,7 @@ class WorkerRuntime: """ customer_id = payload.get("customer_id") if not customer_id: - raise ValueError("profile.rebuild_requested payload is incomplete") + raise OutboxHandlerError("profile.rebuild_requested payload is incomplete") # 延迟导入:bootstrap 会间接导入本模块,模块级导入会形成循环依赖 from app.service.profile_assembly_service import ProfileAssemblyService from app.service.profile_graph_projection_service import ( @@ -181,7 +181,7 @@ class WorkerRuntime: raise ValueError("memory.deletion_requested payload is incomplete") mode = str(payload.get("mode", "invalidate")) if mode not in {"invalidate", "delete"}: - raise ValueError("memory.deletion_requested mode is invalid") + raise OutboxHandlerError("memory.deletion_requested mode is invalid") await MemoryLifecycleService(session).run( int(customer_id), mode=cast("Mode", mode), diff --git a/docs/24-客服Agent阶段性总结与下阶段计划.md b/docs/24-客服Agent阶段性总结与下阶段计划.md index 7ac5b0d..ce2db7a 100644 --- a/docs/24-客服Agent阶段性总结与下阶段计划.md +++ b/docs/24-客服Agent阶段性总结与下阶段计划.md @@ -184,9 +184,12 @@ a6c09fa 检索问句只在客户这一句说不清楚时才带上文 配置项、映射字段、写入新版本,`config_release` 的整版本替换语义不会清空工具白名单。 > 顺带记一条教训:`config_release` 是**整版本替换**,新增一项配置必须把老项带上, > 否则激活的一瞬间老配置全部消失。这条已写进 `docs/05` 与评审报告。 -2. **风险测评记录**:`fin_risk_assessment` 的补数属于演示数据准备,不是代码问题; - 风控合并后适当性链路两侧(客服 `check_suitability` / 风控 RW-007)已按同一口径 - `C ≥ R` 复核通过。 +2. **风险测评记录**:`fin_risk_assessment` 的补数属于演示数据准备,不是代码问题。 + 但**适当性口径在这一轮发现并修掉了一个真 bug**:本文第二节记着客服为 C1–C5 补的 + 「能买什么产品」问答**答案取自第十二条匹配矩阵**,而适当性出口 `check_suitability` + 用的是严格 `C ≥ R` —— 同一个 Agent 对同一个问题两种回答(问"C1 能买什么产品"答 + "R1、R2 可买",问"C1 能买这只 R2 吗"答"不能购买")。现按业务裁定**客服统一到矩阵**, + 风控扫描保留 `C ≥ R`(它要查的是越级成交且留痕不全)。详见 `docs/25` 第七节 #1。 ### 7.2 风控模块合并后的处理(详见 `docs/25`) diff --git a/docs/25-风控模块代码评审报告.md b/docs/25-风控模块代码评审报告.md index 6e9f9a6..86adbef 100644 --- a/docs/25-风控模块代码评审报告.md +++ b/docs/25-风控模块代码评审报告.md @@ -303,24 +303,43 @@ C1(1) 与 R2(2) 相比 `1 < 2`:**按矩阵可以买,按第十四条不能买 **这不是代码问题,是制度文本冲突**,需要业务方定一条为准。客服侧的适当性裁决走的是 `check_suitability`(按档案等级与匹配规则),两边口径也需要对齐。 -**业务裁定(2026-09-11):以第十四条 `C ≥ R` 为准。** +**业务裁定(2026-09-11,经一次自我更正):客服按第十二条矩阵,风控保留 `C ≥ R`。** -**客服侧口径复核(同日):本来就一致,无需改动。** 复核 `SuitabilityService._decide` -(`suitability_service.py:195`)后确认,它的判定是 -`if profile.customer_risk_level < request.product_risk_level: 拒绝` —— **同样是第十四条的 -`C ≥ R`**,并没有使用第十二条的匹配矩阵。原文"按档案等级与匹配规则"是评审时的推测, -不成立。 +第一次复核我只对比了政策文本与 `SuitabilityService`,得出"两边一致、无需改动" —— +**那个结论是错的**,因为漏了第三个东西:**客服自己的知识库**。 -两侧看上去的差异只有两点,且都不构成口径冲突: +- 知识库 `POL-AST-012` 就是第十二条矩阵原文(`个人投资者适当性管理指南.md:302-310`), + 而客服为 C1–C5 各补的「能买什么产品」问答,答案**直接取自该矩阵**(`docs/24` 第二节); +- 而 `SuitabilityService._decide` 用的是严格 `C < R 即拒绝`(第十四条**第 1 款**)。 -1. **专业投资者**:客服侧豁免等级匹配,但强制 `required_disclosure` / - `requires_confirmation` / `requires_recording`(`suitability_service.py:183-194`); - 风控扫描不做等级豁免,而是直接检查"该有的揭示、二次确认、录音留痕有没有"。 - 两者合起来是同一句话:豁免等级不等于豁免留痕。 -2. **触发条件**:客服是**事前拦截**(`C < R` 直接不许买),风控是**事后发现** - (`risk_scan_service.py:184` 的 `gap > 0 and missing_trace`)。这是职责差异,不是口径 - 差异 —— RW-007 的语义是"错配**且**留痕不全",不是"所有错配"。留痕完整却仍然成交, - 那是客服没能拦住,属另一个问题。 +于是**同一个客服 Agent 对同一个问题给出相反答案**:问"C1 能买什么产品"答"R1、R2 可买", +问"C1 能买这只 R2 吗"答"不能购买"。这不是"两侧是否需要对齐",而是客服**内部**自相矛盾, +且客户可见。 + +**冲突的精确位置**也不是原报告写的"矩阵 vs 硬匹配",而是第十四条**第 1 款与它自己的 +第 2、3 款**之间: + +| 条款 | 原文要点 | 对 C1→R2 的结论 | +|---|---|---| +| 第十四条 1 | 投资者等级**大于或等于**产品等级 | 拒绝(1 < 2) | +| 第十四条 2 | 低于**一个等级以上**的才拒绝 | 允许(只低 1 级,够不上"以上") | +| 第十四条 3 | C1 不得买 R3+、C2 不得买 R4+ | 允许(R2 不在禁止列表里) | +| 第十二条矩阵 | C1 行的 R1/R2 标 ✅ 可购买 | 允许 | + +第 2、3 款与矩阵三方一致,孤立的是第 1 款。 + +**落实**:客服侧改为按矩阵裁决 —— `suitability_service.py` 新增 `MATRIX_ALLOWED` 与 +`MATRIX_NEEDS_DISCLOSURE`,`C3→R4`、`C4→R5` 走豁免档(`SUITABLE_WITH_DISCLOSURE`, +强制揭示 + 确认 + 录音),并补了**逐格对照政策原文的 25 格用例** +(`tests/unit/service/test_suitability_service.py`)。客服话术也随之分档:越级档不再说 +"在您的风险承受能力范围内"(那会让客户以为自己的测评本来就覆盖这只产品)。 + +**风控保留 `C ≥ R`**:它要发现的是"越级成交且留痕不全",用更严的口径事后核查。 +两侧职责不同,这个差异是有意保留的。 + +**仍然成立的一点**:专业投资者上,客服侧豁免等级匹配但强制揭示/确认/录音 +(`suitability_service.py`),风控扫描不做等级豁免而是直接查留痕。两者合起来是同一句话 —— +豁免等级不等于豁免留痕。 ### 2. 第十五条豁免规则未落地 ✅(业务裁定:实现,2026-09-11 已完成) diff --git a/docs/evidence/worker-state.json b/docs/evidence/worker-state.json new file mode 100644 index 0000000..b1a3d56 --- /dev/null +++ b/docs/evidence/worker-state.json @@ -0,0 +1,135 @@ +{ + "outbox_columns": [ + "id", + "event_id", + "event_type", + "aggregate_type", + "aggregate_id", + "trace_id", + "payload", + "status", + "retry_count", + "next_retry_at", + "last_error", + "occurred_at", + "published_at", + "created_at", + "updated_at" + ], + "agent_run_columns": [ + "id", + "run_id", + "idempotency_id", + "session_id", + "user_id", + "agent_type", + "trace_id", + "request_message_id", + "result_message_id", + "status", + "attempt_count", + "worker_id", + "locked_until", + "error_code", + "result_version", + "cancel_requested_at", + "started_at", + "completed_at", + "created_at", + "updated_at" + ], + "outbox_by_status": { + "dead": 373, + "failed": 3, + "pending": 347, + "published": 410 + }, + "agent_run_by_status": { + "failed": 21, + "succeeded": 184 + }, + "outbox_created_at_range": [ + "2026-09-09 06:55:06.787006", + "2026-09-11 07:21:23.195457" + ], + "outbox_total": 1133, + "agent_run_total": 205, + "outbox_by_event_type_status": [ + { + "event_type": "agent.run_requested", + "status": "dead", + "count": 373 + }, + { + "event_type": "agent.run_requested", + "status": "pending", + "count": 132 + }, + { + "event_type": "agent.run_requested", + "status": "published", + "count": 130 + }, + { + "event_type": "agent.run_completed", + "status": "published", + "count": 118 + }, + { + "event_type": "memory.extraction_requested", + "status": "published", + "count": 106 + }, + { + "event_type": "profile.rebuild_requested", + "status": "pending", + "count": 76 + }, + { + "event_type": "agent.run_completed", + "status": "pending", + "count": 66 + }, + { + "event_type": "memory.extraction_requested", + "status": "pending", + "count": 60 + }, + { + "event_type": "config.cache_invalidate_requested", + "status": "published", + "count": 42 + }, + { + "event_type": "profile.rebuild_requested", + "status": "published", + "count": 13 + }, + { + "event_type": "config.cache_invalidate_requested", + "status": "pending", + "count": 13 + }, + { + "event_type": "agent.run_requested", + "status": "failed", + "count": 3 + }, + { + "event_type": "memory.invalidated", + "status": "published", + "count": 1 + } + ], + "dead_last_errors": [ + { + "last_error": "ValueError", + "count": 372 + }, + { + "last_error": "no handler registered", + "count": 1 + } + ], + "episode_probe_error": "ImportError: cannot import name 'MemoryEpisode' from 'app.model.memory' (C:\\Users\\Windows\\Desktop\\项目代码\\app\\model\\memory.py)" +} \ No newline at end of file diff --git a/tests/unit/api/test_response_envelope.py b/tests/unit/api/test_response_envelope.py new file mode 100644 index 0000000..9d9456b --- /dev/null +++ b/tests/unit/api/test_response_envelope.py @@ -0,0 +1,119 @@ +"""成功响应的统一信封(`docs/05` §3.3)。 + +§3.3 规定成功响应是 `{data, meta}`,列表的 `data` 为**纯数组**、`next_cursor` 与 +`has_more` 放在 `meta` 里,并且明确「业务接口不得增加其他顶层字段」。 + +此前有三个端点漏了这件事,它们都是"service 直接把内部结构当响应体返回": + +- `GET /conversations/{session_id}/messages` → 裸 `{"data": [...]}`,`meta` 整个缺失, + 游标也没有地方放; +- `POST /conversation-messages/{id}/feedback` → 同样没有 `meta`; +- `GET /knowledge-references/{token}` → 直接返回资源对象。 + +注意"缺 meta"很容易被误判成"有":`X-Trace-ID` 是**响应头**(由中间件加),和 body 里的 +`meta.trace_id` 是两件事;错误响应一直有 `meta`(异常处理器统一加),只有成功路径漏了。 +所以这里断言的是 `set(body)`,多一个或少一个顶层字段都会红。 +""" + +from typing import Any + +import pytest +from fastapi.testclient import TestClient + +from app.api.controllers import conversations as conversations_controller +from app.api.controllers import knowledge as knowledge_controller +from app.api.dependencies.auth import build_request_context +from app.api.dependencies.database import get_session +from app.core.contracts import RequestContext +from app.main import create_app + +TRACE = "trace-envelope" + + +async def resolve_context() -> RequestContext: + return RequestContext( + user_id="9001", + trace_id=TRACE, + permissions=("conversation:create", "conversation:feedback", "knowledge:reference:read"), + ) + + +class StubConversationService: + def __init__(self, _session: Any) -> None: + pass + + async def messages( + self, _session_id: str, _context: RequestContext, _limit: int, before: int | None = None + ) -> dict[str, Any]: + del before + return { + "items": [ + { + "message_id": "11", + "role": "user", + "content": "稳健型", + "created_at": "2026-09-10T00:00:00Z", + } + ], + "next_cursor": "11", + "has_more": True, + } + + async def feedback( + self, + _message_id: int, + _context: RequestContext, + _rating: int, + _feedback_type: str | None, + _feedback_content: str | None, + ) -> dict[str, Any]: + return {"feedback_no": "fb-1", "status": "open"} + + +class StubKnowledgeService: + async def resolve(self, _context: RequestContext, _token: str) -> dict[str, Any]: + return {"knowledge_id": 7, "title": "个人投资者适当性管理指南"} + + +def envelope_client(monkeypatch: pytest.MonkeyPatch) -> TestClient: + monkeypatch.setattr( + conversations_controller, "ConversationService", StubConversationService + ) + monkeypatch.setattr(knowledge_controller, "KnowledgeReferenceService", StubKnowledgeService) + application = create_app() + application.dependency_overrides[build_request_context] = resolve_context + application.dependency_overrides[get_session] = lambda: None + return TestClient(application) + + +def test_message_list_puts_cursor_in_meta(monkeypatch: pytest.MonkeyPatch) -> None: + with envelope_client(monkeypatch) as http: + response = http.get("/api/v1/conversations/session-1/messages", params={"limit": 20}) + + assert response.status_code == 200 + body = response.json() + assert set(body) == {"data", "meta"} + assert isinstance(body["data"], list) + assert body["data"][0]["message_id"] == "11" + assert body["meta"] == {"trace_id": TRACE, "next_cursor": "11", "has_more": True} + + +def test_feedback_success_response_has_meta(monkeypatch: pytest.MonkeyPatch) -> None: + with envelope_client(monkeypatch) as http: + response = http.post("/api/v1/conversation-messages/11/feedback", json={"rating": 1}) + + assert response.status_code == 201 + body = response.json() + assert set(body) == {"data", "meta"} + assert body["data"] == {"feedback_no": "fb-1", "status": "open"} + assert body["meta"] == {"trace_id": TRACE} + + +def test_knowledge_reference_success_response_has_meta(monkeypatch: pytest.MonkeyPatch) -> None: + with envelope_client(monkeypatch) as http: + response = http.get(f"/api/v1/knowledge-references/{'a' * 20}") + + assert response.status_code == 200 + body = response.json() + assert set(body) == {"data", "meta"} + assert body["meta"] == {"trace_id": TRACE} diff --git a/tests/unit/service/test_conversation_service.py b/tests/unit/service/test_conversation_service.py index c74988f..da8fb10 100644 --- a/tests/unit/service/test_conversation_service.py +++ b/tests/unit/service/test_conversation_service.py @@ -108,14 +108,20 @@ async def test_messages_are_scoped_to_current_user(monkeypatch: pytest.MonkeyPat # user_id 必须透传:Repository 靠它做归属过滤,漏传就等于不限范围。 assert captured["user_id"] == 9001 - assert captured["limit"] == 20 + # 21 = limit + 1:service 多取一行判断"还有没有更旧的",用来填 §3.3 的 has_more。 + # 只看"取满没取满"会把恰好等于 limit 的最后一页说成还有下一页。 + assert captured["limit"] == 21 assert captured["session_id"] == "session-1" # 不带游标时必须传 None,行为与加游标前一致(取最新一页)。 assert captured["before"] is None - assert result["data"] == [ + # service 返回内部结构,Controller 用 list_envelope 拆成 {data, meta}(§3.3)。 + # 一行数据小于 limit,所以没有下一页。 + assert result["items"] == [ {"message_id": "11", "role": "user", "content": "稳健型", "created_at": NOW.isoformat() + "Z"} ] + assert result["has_more"] is False + assert result["next_cursor"] is None async def test_messages_forward_cursor_to_repository(monkeypatch: pytest.MonkeyPatch) -> None: @@ -156,7 +162,7 @@ async def test_feedback_writes_feedback_and_audit_with_trace_id( result = await ConversationService(session).feedback(11, CONTEXT, 1, "rating", "很有帮助") - assert result["data"]["status"] == "open" + assert result["status"] == "open" kinds = [type(item).__name__ for item in session.added] assert kinds == ["ConversationFeedback", "InteractionAudit"] audit = session.added[1] diff --git a/tests/unit/service/test_customer_service_suitability.py b/tests/unit/service/test_customer_service_suitability.py index 4e4b625..b1b1601 100644 --- a/tests/unit/service/test_customer_service_suitability.py +++ b/tests/unit/service/test_customer_service_suitability.py @@ -63,6 +63,45 @@ def test_disclosure_and_recording_are_disclosed_when_required() -> None: assert "双录" in text +def test_one_level_up_is_not_described_as_within_capacity() -> None: + """C1 买 R2 是矩阵允许的,但不能说成"在您的风险承受能力范围内"——它超出了等级。 + + 这句错话会让客户以为自己的测评等级本来就覆盖 R2,下次买 R3 时就更难解释为什么不行。 + """ + text = CustomerServiceAgent._suitability_text( + "南方季季盈90天", 2, _decision(customer_risk_level=1) + ) + + assert "在您的风险承受能力范围内" not in text + assert "匹配矩阵允许购买的范围内" in text + + +def test_disclosure_tier_explains_the_condition() -> None: + """C3 买 R4 属豁免档:要说清"高于您的等级"以及"签揭示书后可买"这个条件。""" + text = CustomerServiceAgent._suitability_text( + "南方稳健增利180天", + 4, + _decision( + customer_risk_level=3, + reason_code="SUITABLE_WITH_DISCLOSURE", + required_disclosure=True, + requires_recording=True, + ), + ) + + assert "高于您的风险承受能力等级" in text + assert "签署产品风险揭示书后可以购买" in text + assert "在您的风险承受能力范围内" not in text + + +def test_same_level_still_uses_the_capacity_wording() -> None: + text = CustomerServiceAgent._suitability_text( + "南方季季盈90天", 2, _decision(customer_risk_level=2) + ) + + assert "在您的风险承受能力范围内" in text + + def test_rejection_does_not_assert_a_reason() -> None: """拒绝的原因可能是等级不匹配、测评过期或未测评,不能一律说成"超出承受能力"。""" text = CustomerServiceAgent._suitability_text( diff --git a/tests/unit/service/test_suitability_service.py b/tests/unit/service/test_suitability_service.py index 0c71673..89ae592 100644 --- a/tests/unit/service/test_suitability_service.py +++ b/tests/unit/service/test_suitability_service.py @@ -101,14 +101,57 @@ def test_caller_cannot_declare_risk_facts(forged: dict[str, Any]) -> None: async def test_insufficient_authority_level_is_denied() -> None: + """低两个等级及以上仍必须拒绝(第十四条第 2、3 款,矩阵里的"❌ 禁止")。""" decision = await service_with_row(authority_row(investor_type="C1")).evaluate( - query(product_risk_level=2), context(), now=NOW + query(product_risk_level=3), context(), now=NOW ) assert decision.allowed is False assert decision.reason_code == "RISK_LEVEL_MISMATCH" assert decision.requires_recording is True +@pytest.mark.parametrize( + ("investor_type", "product_level", "allowed", "reason_code"), + [ + # 逐格抄自 knowledge/policy/个人投资者适当性管理指南.md 第十二条矩阵。 + # 这张表的价值在于:任何一格被改动,都必须是有意为之并在此处说明理由。 + ("C1", 1, True, "SUITABLE"), ("C1", 2, True, "SUITABLE"), + ("C1", 3, False, "RISK_LEVEL_MISMATCH"), ("C1", 4, False, "RISK_LEVEL_MISMATCH"), + ("C1", 5, False, "RISK_LEVEL_MISMATCH"), + ("C2", 1, True, "SUITABLE"), ("C2", 2, True, "SUITABLE"), ("C2", 3, True, "SUITABLE"), + ("C2", 4, False, "RISK_LEVEL_MISMATCH"), ("C2", 5, False, "RISK_LEVEL_MISMATCH"), + ("C3", 1, True, "SUITABLE"), ("C3", 2, True, "SUITABLE"), ("C3", 3, True, "SUITABLE"), + ("C3", 4, True, "SUITABLE_WITH_DISCLOSURE"), ("C3", 5, False, "RISK_LEVEL_MISMATCH"), + ("C4", 1, True, "SUITABLE"), ("C4", 2, True, "SUITABLE"), ("C4", 3, True, "SUITABLE"), + ("C4", 4, True, "SUITABLE"), ("C4", 5, True, "SUITABLE_WITH_DISCLOSURE"), + ("C5", 1, True, "SUITABLE"), ("C5", 2, True, "SUITABLE"), ("C5", 3, True, "SUITABLE"), + ("C5", 4, True, "SUITABLE"), ("C5", 5, True, "SUITABLE"), + ], +) +async def test_full_matrix_matches_policy_document( + investor_type: str, product_level: int, allowed: bool, reason_code: str +) -> None: + """客服回答必须与知识库里的矩阵一致 —— 这是同一个 Agent 的两条出口。""" + decision = await service_with_row(authority_row(investor_type=investor_type)).evaluate( + query(product_risk_level=product_level), context(), now=NOW + ) + assert decision.allowed is allowed + assert decision.reason_code == reason_code + + +async def test_disclosure_tier_always_requires_disclosure_and_recording() -> None: + """C3→R4、C4→R5 是第十五条豁免档:可买,但必须揭示、确认、录音。""" + for investor_type, product_level in (("C3", 4), ("C4", 5)): + decision = await service_with_row(authority_row(investor_type=investor_type)).evaluate( + query(product_risk_level=product_level), context(), now=NOW + ) + assert decision.allowed is True + assert decision.reason_code == "SUITABLE_WITH_DISCLOSURE" + assert decision.required_disclosure is True + assert decision.requires_confirmation is True + assert decision.requires_recording is True + + async def test_expired_assessment_is_denied_even_for_eligible_level() -> None: row = authority_row(investor_type="C5", valid_until=NOW - timedelta(seconds=1)) decision = await service_with_row(row).evaluate(query(product_risk_level=1), context(), now=NOW) diff --git a/tests/unit/worker/test_outbox_worker.py b/tests/unit/worker/test_outbox_worker.py index 08bd270..a5f04e4 100644 --- a/tests/unit/worker/test_outbox_worker.py +++ b/tests/unit/worker/test_outbox_worker.py @@ -4,7 +4,7 @@ from unittest.mock import AsyncMock, Mock import pytest from app.model.platform import DomainEventOutbox -from app.worker.outbox_worker import OutboxWorker +from app.worker.outbox_worker import OutboxHandlerError, OutboxWorker def event() -> DomainEventOutbox: @@ -66,9 +66,32 @@ async def test_handler_failure_records_attempt(attempts: int, expected: str) -> assert row.status == expected assert row.retry_count == attempts + 1 assert "do-not-log" not in (row.last_error or "") + # 其他异常一律只记类名:消息可能含凭据/SQL/客户标识,不落库。 + assert row.last_error == "RuntimeError" session.commit.assert_awaited_once() +@pytest.mark.asyncio +async def test_handler_error_reason_is_recorded() -> None: + """`OutboxHandlerError` 的固定文案可以落库 —— 否则失败原因根本无从分辨。 + + 实测库里 373 条死信的 `last_error` 全是裸的 `"ValueError"`:`dispatch`(run not + found)、`dispatch_run_completed`、`dispatch_memory_extraction`、 + `dispatch_profile_rebuild` 抛的都是 ValueError,只记类名等于把"哪一处失败"也丢了。 + 这些 handler 现在改抛 `OutboxHandlerError`,它的 reason 由代码写死、不含请求数据。 + """ + row = event() + session = AsyncMock() + session.add = Mock() + session.scalar.side_effect = [row, None] + handler = AsyncMock(side_effect=OutboxHandlerError("run not found")) + + await OutboxWorker(session, {"test.event": handler}).publish_one() + + assert row.last_error == "OutboxHandlerError: run not found" + assert row.status == "failed" + + @pytest.mark.asyncio async def test_empty_matching_queue_releases_transaction() -> None: session = AsyncMock() diff --git a/tools/probe_worker_state.py b/tools/probe_worker_state.py new file mode 100644 index 0000000..cd349b3 --- /dev/null +++ b/tools/probe_worker_state.py @@ -0,0 +1,113 @@ +"""只读探查:Worker 的运行状态(队列有没有在消费)。 + +只做 SELECT,结果写 `docs/evidence/worker-state.json`: + + python tools/probe_worker_state.py + +三个问题: +1. `domain_event_outbox` 里各状态各有多少条、最老的一条积压了多久 + —— Worker 没在跑,事件就永远停在 pending; +2. `agent_run` 里 queued/running 有多少 + —— Worker 没在跑,用户发起的对话就一直"排队中"; +3. `memory_episode` 有没有待提取的片段(同样是 Worker 负责消费)。 +""" + +from __future__ import annotations + +import asyncio +import json +from pathlib import Path +from typing import Any + +from sqlalchemy import func, select + +from app.infrastructure.db import SessionFactory +from app.model.platform import AgentRun, DomainEventOutbox + +OUTPUT = Path("docs/evidence/worker-state.json") + + +async def _grouped(session: Any, model: Any, column_name: str) -> dict[str, int]: + column = getattr(model, column_name, None) + if column is None: + return {} + rows = await session.execute(select(column, func.count()).group_by(column)) + return {str(key): int(count) for key, count in rows.all()} + + +async def collect() -> dict[str, Any]: + report: dict[str, Any] = {} + async with SessionFactory() as session: + report["outbox_columns"] = list(DomainEventOutbox.__table__.columns.keys()) + report["agent_run_columns"] = list(AgentRun.__table__.columns.keys()) + + report["outbox_by_status"] = await _grouped( + session, DomainEventOutbox, "status" + ) + report["agent_run_by_status"] = await _grouped(session, AgentRun, "status") + + oldest = await session.scalar( + select(func.min(DomainEventOutbox.created_at)) + ) + newest = await session.scalar( + select(func.max(DomainEventOutbox.created_at)) + ) + report["outbox_created_at_range"] = [str(oldest), str(newest)] + report["outbox_total"] = await session.scalar( + select(func.count()).select_from(DomainEventOutbox) + ) + report["agent_run_total"] = await session.scalar( + select(func.count()).select_from(AgentRun) + ) + + # 哪一类事件在堆、哪一类已经进了死信 —— 死信意味着那些事件的副作用永远不会发生。 + report["outbox_by_event_type_status"] = [ + {"event_type": str(row[0]), "status": str(row[1]), "count": int(row[2])} + for row in ( + await session.execute( + select( + DomainEventOutbox.event_type, + DomainEventOutbox.status, + func.count(), + ) + .group_by(DomainEventOutbox.event_type, DomainEventOutbox.status) + .order_by(func.count().desc()) + ) + ).all() + ] + report["dead_last_errors"] = [ + {"last_error": str(row[0])[:400], "count": int(row[1])} + for row in ( + await session.execute( + select(DomainEventOutbox.last_error, func.count()) + .where(DomainEventOutbox.status == "dead") + .group_by(DomainEventOutbox.last_error) + .order_by(func.count().desc()) + ) + ).all() + ] + + try: + from app.model.memory import MemoryEpisode + + report["episode_by_status"] = await _grouped(session, MemoryEpisode, "status") + report["episode_total"] = await session.scalar( + select(func.count()).select_from(MemoryEpisode) + ) + except Exception as exc: # 模型名/字段与预期不符时只记录,不影响其余结论 + report["episode_probe_error"] = f"{type(exc).__name__}: {exc}" + return report + + +async def main() -> None: + report = await collect() + OUTPUT.parent.mkdir(parents=True, exist_ok=True) + OUTPUT.write_text( + json.dumps(report, ensure_ascii=False, indent=2, default=str), + encoding="utf-8", + ) + print(f"wrote {OUTPUT}") + + +if __name__ == "__main__": + asyncio.run(main())