一、客服 Agent 智能增强(正面回应"不智能、动不动就转人工")
- 决策链由 2 个出口扩到 5 个:E1 澄清 / E2 计算型 / E3 知识直返 / E4 证据约束生成 / E5 分级回退
- 转人工从"默认动作"降为最后一档 E5c,只保留 4 类白名单:
P0 反诈 / P1 账户与个人数据 / P2 写操作与争议 / 用户明确要求人工
- 46 条金标实测(修复前 → 修复后):
转人工率 43.5% → 10.9%;出口准确率 45.7% → 100%;事实正确率 69.6% → 100%
禁忌违反 1 → 0;档位越权 / 无出处数字 / 误拒 四项零容忍全 0
- 安全不变量 INV-1~INV-5;零容忍规则未删,改的是挂载点
(输出侧字面黑名单 → 检索层档位隔离 + 判定层合规词表 + 输出守护)
二、知识库:档位单点化与物理隔离
- 新增 app/core/knowledge_tier.py 作为档位规则唯一落点(G-03),
knowledge_contracts.py 原定义块改为显式再导出(X as X,非副本)
- 档位过滤由 bool 默认值(fail-open)改为 tiers 必填集合(缺参即 TypeError)
- Milvus 侧四集合按 visibility 分区键物理隔离;双 schema 收敛为一套
- 新增 app/core/actor.py:访客三元组与匿名判定的唯一构造/判定点(G-01/G-01b)
- 新增 app/core/fund_fee_rules.py:费率计算纯函数
三、前端入参边界对齐(本轮 W11 新修,4 处"校验宽于存储")
- message 加 max_length=8000(与浮窗 widget.js 的 maxlength 一致)
- session_id 加 1—64;idempotency_key 上限 128 → 64(对齐列宽 String(64))
- feedback_type 加 max_length=32(对齐列宽 String(32))
- 8 条路径参数补 min_length=1 + max_length=64 + 字符集正则
({session_id} / {run_id} / {handover_id})
- 改前超限值会落到 MySQL 才失败(500);改后一律 422 AGENT_INPUT_INVALID + 字段级定位
- 新增 tests/unit/api/test_frontend_boundaries.py(33 例),含"端点表 ↔ OpenAPI 全量对照"
四、投顾模块整体清除(D4.4 / D4.5)
- 删除投顾相关 controller / schema / model / repository / service 及门户页面
- tools/portal_api_check.py 同步作废 AD003/AD005/AD011/A047 四条用例与 advisor_t 登录
(端点与账号均已不存在,此前稳定报 3 条假红)
五、验证(提交前实测)
- pytest -q:1856 passed / 2 skipped / 0 failed
- ruff check app tools tests:19(= 基线);mypy app:2(= 基线)
- 前端接口契约体检 portal_api_check.py:38 项,通过 34,失败 0,跳过 4
- 全链路冒烟 e2e_smoke_test.py --read-only:31/31
- HTTP 全链路探针 http_probe.py:11/11 succeeded
- 跨文档一致性 _consistency.py:GATE PASS
- 真机边界复验 12 条:12/12 符合预期
六、纪律与文档
- 可改文件白名单 A-09(docs/46)与底座会签申请单 A-10(docs/47,组 1—组 4 全部受理)
- 零 DDL:未新增/修改任何表结构,89 张业务表与基线一致
- 证据留痕:docs/evidence/**(含 46 条金标 score、快照、清除与重建记录)
- 未提交(刻意排除,见提交说明):仓库内 客服agent/ 与 开发文档/ 是 2026-09-16 前的
过期副本(Todolist 440 行 vs 权威 D2.1 1167 行),权威正本在仓库外;
_chunks_report.txt 是 tools/build_knowledge_chunks.py 生成的本地产物
186 lines
6.8 KiB
Python
186 lines
6.8 KiB
Python
"""风控规则扫描定时 Worker。
|
|
|
|
该进程不依赖 Web 进程生命周期,通过 MySQL 咨询锁保证同一时刻只有一个
|
|
调度者执行扫描。默认关闭,配置开启后才执行。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import asyncio
|
|
import logging
|
|
from collections.abc import Awaitable, Callable
|
|
from contextlib import AbstractAsyncContextManager
|
|
from datetime import UTC, datetime, timedelta
|
|
from typing import Any
|
|
from uuid import uuid4
|
|
|
|
from app.core.config import get_settings
|
|
from app.core.contracts import RequestContext
|
|
from app.infrastructure.db import SessionFactory, engine, mysql_scan_lock
|
|
from app.model.audit import InteractionAudit
|
|
from app.service.risk_scan_schedule_config import (
|
|
RiskScanScheduleConfig,
|
|
load_risk_scan_schedule_config,
|
|
)
|
|
from app.service.risk_scan_service import RiskScanService
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
ConfigLoader = Callable[[], Awaitable[RiskScanScheduleConfig]]
|
|
ScanExecutor = Callable[[], Awaitable[dict[str, int | str]]]
|
|
AuditWriter = Callable[[str, dict[str, Any]], Awaitable[None]]
|
|
LockFactory = Callable[[], AbstractAsyncContextManager[bool]]
|
|
|
|
# 跨进程扫描锁已移到 `app/infrastructure/db.py`:端点与调度器**必须共用同一把锁**,
|
|
# 放在基础设施层两个入口才都能引用(service 不该反向依赖 worker)。
|
|
|
|
|
|
async def default_scan_executor() -> dict[str, int | str]:
|
|
context = RequestContext(
|
|
user_id="0",
|
|
trace_id=f"risk-scan-schedule-{uuid4()}",
|
|
roles=("system",),
|
|
permissions=("risk:alert:scan",),
|
|
data_scope="all",
|
|
portal="worker",
|
|
)
|
|
async with SessionFactory() as session:
|
|
return await RiskScanService.from_settings(session).scan(context)
|
|
|
|
|
|
async def default_audit_writer(status: str, detail: dict[str, Any]) -> None:
|
|
async with SessionFactory() as session, session.begin():
|
|
session.add(
|
|
InteractionAudit(
|
|
actor_type="system",
|
|
actor_id=None,
|
|
target_customer_id=None,
|
|
portal="worker",
|
|
action_type=f"risk_scan_scheduled_{status}",
|
|
detail=detail,
|
|
created_at=datetime.now(UTC).replace(tzinfo=None),
|
|
)
|
|
)
|
|
|
|
|
|
class RiskScanSchedulerWorker:
|
|
def __init__(
|
|
self,
|
|
*,
|
|
config_loader: ConfigLoader = load_risk_scan_schedule_config,
|
|
scan_executor: ScanExecutor = default_scan_executor,
|
|
audit_writer: AuditWriter = default_audit_writer,
|
|
lock_factory: LockFactory = mysql_scan_lock,
|
|
now: Callable[[], datetime] | None = None,
|
|
) -> None:
|
|
self.config_loader = config_loader
|
|
self.scan_executor = scan_executor
|
|
self.audit_writer = audit_writer
|
|
self.lock_factory = lock_factory
|
|
self.now = now or (lambda: datetime.now(UTC))
|
|
self.last_run_at: datetime | None = None
|
|
|
|
async def run_once(self, *, force: bool = False) -> bool:
|
|
config = await self.config_loader()
|
|
if not config.enabled:
|
|
return False
|
|
current = self.now()
|
|
if not force and not self._is_due(config, current):
|
|
return False
|
|
|
|
async with self.lock_factory() as acquired:
|
|
if not acquired:
|
|
logger.info("风控定时扫描由其他 Worker 执行,本轮跳过")
|
|
return False
|
|
return await self._execute_with_retry(config)
|
|
|
|
def _is_due(
|
|
self,
|
|
config: RiskScanScheduleConfig,
|
|
current: datetime,
|
|
) -> bool:
|
|
if self.last_run_at is None:
|
|
# 进程刚起来,不知道自己上次是什么时候跑的 —— `last_run_at` 只存在内存里。
|
|
# **保守地视为 due**:多跑一次的最坏后果是重复扫描,而扫描本身是幂等的
|
|
# (每条规则先 `_exists` 查重)并且有 MySQL 级锁;反过来"不跑"的后果可能是
|
|
# **永远不跑**:原先这里返回 `config.run_immediately`(默认 False),
|
|
# 重启之后 `_is_due` 恒为假,调度器形同虚设 —— 而且没有任何告警,
|
|
# 现场只会表现为"风控好像没在扫描"。
|
|
#
|
|
# `config.run_immediately` 因此不再承担"首次是否执行"的语义(它原本想表达的
|
|
# 是"启动后别马上跑",但那与"永远不跑"在实现上无法区分)。字段保留,
|
|
# 以免破坏既有配置。
|
|
return True
|
|
return current - self.last_run_at >= timedelta(minutes=config.interval_minutes)
|
|
|
|
async def _execute_with_retry(
|
|
self,
|
|
config: RiskScanScheduleConfig,
|
|
) -> bool:
|
|
attempts = config.retry_limit + 1
|
|
last_error: Exception | None = None
|
|
for attempt in range(1, attempts + 1):
|
|
try:
|
|
result = await self.scan_executor()
|
|
self.last_run_at = self.now()
|
|
await self.audit_writer(
|
|
"succeeded",
|
|
{
|
|
**result,
|
|
"attempt": attempt,
|
|
"trace_id": f"risk-scan-schedule-{uuid4()}",
|
|
},
|
|
)
|
|
return True
|
|
except Exception as error:
|
|
last_error = error
|
|
logger.warning(
|
|
"风控定时扫描失败 attempt=%s/%s",
|
|
attempt,
|
|
attempts,
|
|
exc_info=True,
|
|
)
|
|
await self.audit_writer(
|
|
"failed",
|
|
{
|
|
"attempts": attempts,
|
|
"error_type": type(last_error).__name__ if last_error else "unknown",
|
|
"trace_id": f"risk-scan-schedule-{uuid4()}",
|
|
},
|
|
)
|
|
return False
|
|
|
|
|
|
async def serve(*, once: bool = False, force: bool = False) -> None:
|
|
worker = RiskScanSchedulerWorker()
|
|
try:
|
|
while True:
|
|
try:
|
|
await worker.run_once(force=force)
|
|
except Exception:
|
|
logger.warning("风控定时扫描 Worker 轮次失败", exc_info=True)
|
|
if once:
|
|
raise
|
|
if once:
|
|
return
|
|
await asyncio.sleep(get_settings().risk_scan_poll_seconds)
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
|
|
def main() -> None:
|
|
parser = argparse.ArgumentParser(description="南方基金风控规则定时扫描 Worker")
|
|
parser.add_argument("--once", action="store_true", help="执行一轮后退出")
|
|
parser.add_argument("--force", action="store_true", help="忽略间隔,立即执行一次")
|
|
args = parser.parse_args()
|
|
logging.basicConfig(level=logging.INFO)
|
|
try:
|
|
asyncio.run(serve(once=args.once, force=args.force))
|
|
except KeyboardInterrupt:
|
|
pass
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|