两个相关的缺陷,都属于**静默失效**型。 1. `admin_service._transition_intent_config`(真正的 bug) 激活一个意图时按 `agent_type` 过滤旧 active 版本并归档,但唯一键是生成列 `active_key = concat(agent_type, ':', intent_code)` —— 同一 `agent_type` 下 **不同意图码本就允许并存**。于是激活 `general` 会把 `risk_overview` / `risk_search` / `risk_evidence` 一并归档:风控运行期只剩 1 条 active 意图,问"查看当前风险概览"被分到 `general`(confidence 0.5), 而**没有任何报错**。该函数自己的 docstring 写的正是正确行为,实现与它不符。 2. `tools/publish_risk_agent_config.py`(使脚本无法自愈) 按 `intent_code` 单键建 dict 收集现有意图,而列表**按 id 倒序**返回且同一意图码 有多个版本,于是**旧版本覆盖新版本**;取到 v1 后再"版本 +1"算出的正是已被占用的 v2,创建必然 409 IDEMPOTENCY_CONFLICT。现象是 4 条意图全部"创建失败"而库里 其实都有,`tools/seed_demo_data.py` 第 7 步因此必然失败。 修复后实测: - 库里 4 个风控意图全部 active; - 问"查看当前风险概览" → `intent=risk_overview`、`confidence=1.0000` (修复前为 `general` / 0.5); - `tools/seed_demo_data.py` 10/10 步完成、退出码 0(修复前第 7 步退出码 1)。 回归测试:`test_activating_one_intent_does_not_archive_sibling_intents`。 已确认把修复回退后该用例**确实失败**(`['beta'] != ['alpha', 'beta']`), 不是永远通过的空测试 —— 原有用例盖不住这个缺陷,因为它建的第二个意图始终停在 draft、 从未激活过。 门禁:ruff 通过;mypy 250 文件 0 错;unit+contract 1381 passed / 0 failed; integration 104 passed;e2e 冒烟 40/40。
323 lines
14 KiB
Python
323 lines
14 KiB
Python
"""发布风控 Agent 的运行期配置:意图配置 + 意图工具白名单。
|
||
|
||
配置来源:组员交付的 `20-Agent工具白名单与意图配置.md` 与两份 JSON。
|
||
工具名与权限已核对与代码一致(`risk_agent.py:24-26` 定义常量、`bootstrap.py:206-225` 注册,
|
||
`required_permission` 均为 `risk:alert:read`)。
|
||
|
||
两个必须讲清的点(与 `publish_customer_service_config.py` 同源):
|
||
|
||
1. **为什么必须发这一步**:工具白名单是**失败关闭**的——`ToolExecutor` 拿发布配置里
|
||
`agent_tools` / `risk:<intent>` 的 `allowed_tools` 与代码声明的
|
||
`AgentDefinition.allowed_tools` 取交集,缺配置时交集为空、任何工具调用都被拒。
|
||
「Agent 写好了但没发配置」的表现是"风控什么都答不了"。
|
||
|
||
2. **为什么必须继承现有配置项**:`config_release` 是**整版本替换**语义——激活新版本后,
|
||
旧版本的所有配置项都不再生效。若只发布风控自己的白名单,客服的 4 条白名单与示例 Agent
|
||
的 `fund_query_demo:fund_quote` 会被静默清空(客服表现为"一直转人工")。所以发布前先把
|
||
当前 effective 版本里的配置项原样搬进新版本,再追加本次新增项。
|
||
|
||
用法:python tools/publish_risk_agent_config.py
|
||
"""
|
||
|
||
import asyncio
|
||
import datetime as dt
|
||
import json
|
||
import sys
|
||
import uuid
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
import asyncmy
|
||
import httpx
|
||
import jwt
|
||
|
||
PROJECT_ROOT = Path(__file__).resolve().parents[1]
|
||
if str(PROJECT_ROOT) not in sys.path:
|
||
sys.path.insert(0, str(PROJECT_ROOT))
|
||
|
||
from app.core.config import get_settings # noqa: E402
|
||
from app.main import create_app # noqa: E402
|
||
|
||
ADMIN = "9003"
|
||
AGENT_TYPE = "risk"
|
||
|
||
# 工具白名单:来自 risk_agent_tools.json(release_no=risk-agent-local-v1)。
|
||
# `general` 是通用风控查询("你能做什么"),按交付文档不配任何工具。
|
||
INTENT_TOOLS: dict[str, tuple[str, ...]] = {
|
||
"risk_overview": ("get_risk_overview",),
|
||
"risk_search": ("search_risk_alerts",),
|
||
"risk_evidence": ("get_alert_evidence",),
|
||
"general": (),
|
||
}
|
||
|
||
# 意图配置:来自 risk_agent_intents.json。description 与 examples 是**给分类器看的**,
|
||
# 真正让模型分辨"风险概览"和"预警证据"的是这几个例子,所以照抄交付值、不改写。
|
||
INTENT_SPECS: tuple[dict[str, Any], ...] = (
|
||
{
|
||
"intent_code": "risk_overview", "intent_name": "风险概览",
|
||
"description": "奶龙风控智能助手:风险概览",
|
||
"examples": ["查看当前风险概览", "当前有多少高风险预警"],
|
||
},
|
||
{
|
||
"intent_code": "risk_search", "intent_name": "风险查询",
|
||
"description": "奶龙风控智能助手:风险查询",
|
||
"examples": ["查询高风险预警", "查看命中 RW-007 的预警"],
|
||
},
|
||
{
|
||
"intent_code": "risk_evidence", "intent_name": "预警证据",
|
||
"description": "奶龙风控智能助手:预警证据",
|
||
"examples": ["查询预警编号 ALERT-001 的证据", "查看这条预警的证据链"],
|
||
},
|
||
{
|
||
"intent_code": "general", "intent_name": "通用风控查询",
|
||
"description": "奶龙风控智能助手:通用风控查询",
|
||
"examples": ["你能做什么", "说明你的功能边界"],
|
||
},
|
||
)
|
||
|
||
INTENT_COMMON: dict[str, Any] = {
|
||
"classifier_instruction": "只用于只读工具查询、研判草案和边界说明,不执行人工处置。",
|
||
"confidence_threshold": "0.6500",
|
||
"max_clarification_rounds": 2,
|
||
"transfer_on_failure": True,
|
||
"priority": 100,
|
||
}
|
||
|
||
|
||
def token(subject: str) -> str:
|
||
settings = get_settings()
|
||
private_key = Path(settings.jwt_private_key_path).read_text(encoding="utf-8")
|
||
now = dt.datetime.now(dt.UTC)
|
||
return jwt.encode(
|
||
{
|
||
"sub": subject, "iss": settings.jwt_issuer, "aud": settings.jwt_audience,
|
||
"exp": now + dt.timedelta(minutes=30), "nbf": now - dt.timedelta(seconds=5),
|
||
"jti": str(uuid.uuid4()),
|
||
},
|
||
private_key,
|
||
algorithm="RS256",
|
||
)
|
||
|
||
|
||
async def active_config_items() -> list[dict[str, Any]]:
|
||
"""读取当前生效版本的全部配置项,用于在新版本里原样继承。"""
|
||
settings = get_settings()
|
||
# MYSQL_DSN 形如 mysql+asyncmy://user:pass@host:port/db
|
||
dsn = settings.mysql_dsn.split("://", 1)[1]
|
||
credentials, location = dsn.split("@", 1)
|
||
user, password = credentials.split(":", 1)
|
||
host_port, database = location.split("/", 1)
|
||
host, _, port = host_port.partition(":")
|
||
connection = await asyncmy.connect(
|
||
host=host, port=int(port or 3306), user=user, password=password, db=database
|
||
)
|
||
try:
|
||
cursor = connection.cursor()
|
||
await cursor.execute(
|
||
"""
|
||
SELECT i.namespace, i.config_key, i.value_json, i.schema_version
|
||
FROM platform_config_item i
|
||
JOIN config_release r ON r.id = i.release_id
|
||
WHERE r.status = 'active'
|
||
"""
|
||
)
|
||
rows = await cursor.fetchall()
|
||
finally:
|
||
connection.close()
|
||
items: list[dict[str, Any]] = []
|
||
for namespace, config_key, value_json, schema_version in rows:
|
||
value = json.loads(value_json) if isinstance(value_json, str) else value_json
|
||
items.append({
|
||
"namespace": namespace,
|
||
"item_key": config_key,
|
||
"value_json": value,
|
||
"schema_version": schema_version,
|
||
})
|
||
return items
|
||
|
||
|
||
async def post(
|
||
client: httpx.AsyncClient, path: str, *, auth: dict[str, str],
|
||
payload: dict[str, object] | None = None, if_match: str | None = None,
|
||
) -> httpx.Response:
|
||
headers = {**auth, "Idempotency-Key": uuid.uuid4().hex}
|
||
if if_match:
|
||
headers["If-Match"] = if_match
|
||
return await client.post(path, json=payload, headers=headers)
|
||
|
||
|
||
async def etag_of(client: httpx.AsyncClient, path: str, auth: dict[str, str]) -> str | None:
|
||
return (await client.get(path, headers=auth)).headers.get("ETag")
|
||
|
||
|
||
def _intent_payload(spec: dict[str, Any], version: int) -> dict[str, Any]:
|
||
return {
|
||
"agent_type": AGENT_TYPE,
|
||
**spec,
|
||
**INTENT_COMMON,
|
||
"allowed_tools": list(INTENT_TOOLS[str(spec["intent_code"])]),
|
||
"version": version,
|
||
}
|
||
|
||
|
||
async def ensure_risk_intents(client: httpx.AsyncClient, auth: dict[str, str]) -> int:
|
||
"""确保 4 条风控意图在运行期生效(返回 0 成功、1 失败)。
|
||
|
||
意图码要三处对齐:`AgentDefinition.supported_intents`(代码,已有)、
|
||
`agent_intent_config` 的 active 行(本函数)、发布版 `agent_tools`(下一步)。
|
||
缺任一处即失败关闭——少了这里,"查看风险概览"会被分到别的意图去。
|
||
"""
|
||
path = "/api/v1/admin/agent-intent-configs"
|
||
listed = await client.get(f"{path}?limit=100", headers=auth)
|
||
rows = listed.json().get("data", []) if listed.status_code == 200 else []
|
||
# ⚠️ 必须按 `intent_code` **聚合全部版本**,不能一个码只留一行:
|
||
# 同一意图码会有多个版本(历史版本 `status='archived'`),而列表是**按 id 倒序**返回的,
|
||
# 用单键 dict 会让**旧版本覆盖新版本**(v1 的 id 更小、排在更后面)。
|
||
# 取到 v1 之后再"版本 +1",算出来的正是已被占用的 v2,创建必然
|
||
# 409 IDEMPOTENCY_CONFLICT(撞唯一键 `uk_intent_config_version`)——
|
||
# 现象是 4 条意图全部"创建失败",而库里其实都有。
|
||
by_code: dict[str, list[dict[str, Any]]] = {}
|
||
for row in rows:
|
||
if row.get("agent_type") == AGENT_TYPE:
|
||
by_code.setdefault(str(row.get("intent_code")), []).append(row)
|
||
|
||
failures = 0
|
||
for spec in INTENT_SPECS:
|
||
code = str(spec["intent_code"])
|
||
versions = by_code.get(code, [])
|
||
active = next((r for r in versions if str(r.get("status")) == "active"), None)
|
||
if active is not None:
|
||
print(f"[意图] {code} 已生效(id={active['id']}, v{active.get('version')}),跳过")
|
||
continue
|
||
# 未生效但可推进的版本(草稿/已审)优先复用,省一次新建
|
||
reusable = next(
|
||
(
|
||
r
|
||
for r in sorted(versions, key=lambda r: int(r.get("version", 0)), reverse=True)
|
||
if str(r.get("status")) in {"draft", "approved"}
|
||
),
|
||
None,
|
||
)
|
||
if reusable is not None:
|
||
config_id = int(reusable["id"])
|
||
print(f"[意图] {code} 复用未生效版本(id={config_id}, v{reusable.get('version')})")
|
||
else:
|
||
# 新版本号取**全部版本的最大值 +1**:归档版本仍然占用着版本号
|
||
version = max((int(r.get("version", 0)) for r in versions), default=0) + 1
|
||
created = await post(client, path, auth=auth,
|
||
payload=_intent_payload(spec, version))
|
||
if created.status_code != 201:
|
||
print(f"[意图] {code} 创建失败:{created.status_code} {created.text[:200]}")
|
||
failures += 1
|
||
continue
|
||
config_id = int(created.json()["data"]["id"])
|
||
print(f"[意图] {code} 已创建(id={config_id}, v{version})")
|
||
|
||
base = f"{path}/{config_id}"
|
||
# 幂等:重跑时某一步可能已推进过("已经审过了"不该 409 让整个脚本失败)
|
||
settled = {"reviews": {"approved", "active"}, "activations": {"active"}}
|
||
for action, payload in (
|
||
("reviews", {"decision": "approved", "comment": "创建人自审"}),
|
||
# 激活端点要求 body 是对象;传 None 时 httpx 根本不发 body,会被判 422
|
||
("activations", {}),
|
||
):
|
||
current = (await client.get(base, headers=auth)).json().get("data", {})
|
||
if str(current.get("status")) in settled[action]:
|
||
continue
|
||
response = await post(client, f"{base}/{action}", auth=auth, payload=payload,
|
||
if_match=await etag_of(client, base, auth))
|
||
if response.status_code != 200:
|
||
print(f"[意图] {code} {action} 失败:{response.status_code} {response.text[:200]}")
|
||
failures += 1
|
||
break
|
||
else:
|
||
print(f"[意图] {code} 已生效")
|
||
return 1 if failures else 0
|
||
|
||
|
||
async def main() -> int:
|
||
app = create_app()
|
||
auth = {"Authorization": f"Bearer {token(ADMIN)}"}
|
||
async with httpx.AsyncClient(
|
||
transport=httpx.ASGITransport(app=app), base_url="http://test", timeout=60
|
||
) as client:
|
||
if await ensure_risk_intents(client, auth) != 0:
|
||
return 1
|
||
|
||
inherited = await active_config_items()
|
||
print(f"\n当前生效版本的配置项:{len(inherited)} 条(将原样继承)")
|
||
for item in inherited:
|
||
print(f" · {item['namespace']} / {item['item_key']}")
|
||
|
||
new_items = [
|
||
{
|
||
"namespace": "agent_tools",
|
||
"item_key": f"{AGENT_TYPE}:{intent}",
|
||
"value_json": {"allowed_tools": list(tools)},
|
||
"schema_version": "1",
|
||
}
|
||
for intent, tools in INTENT_TOOLS.items()
|
||
]
|
||
inherited_keys = {(str(i["namespace"]), str(i["item_key"])) for i in inherited}
|
||
pending = [
|
||
item for item in new_items
|
||
if (str(item["namespace"]), str(item["item_key"])) not in inherited_keys
|
||
]
|
||
if not pending:
|
||
print("\n风控白名单已存在于当前生效版本,无需发布")
|
||
return 0
|
||
|
||
created = await post(client, "/api/v1/admin/config-releases", auth=auth, payload={
|
||
"release_no": f"risk-tools-{uuid.uuid4().hex[:12]}",
|
||
"title": "风控 Agent 意图工具白名单",
|
||
"change_summary": (
|
||
"新增 risk_overview/risk_search/risk_evidence/general 的只读工具白名单"
|
||
"(general 不配工具),并继承既有配置项"
|
||
),
|
||
})
|
||
if created.status_code != 201:
|
||
print(f"创建发布版本失败:{created.status_code} {created.text[:200]}")
|
||
return 1
|
||
release_id = int(created.json()["data"]["id"])
|
||
print(f"\n发布版本 id={release_id}")
|
||
|
||
base = f"/api/v1/admin/config-releases/{release_id}/platform-config-items"
|
||
for item in [*inherited, *pending]:
|
||
response = await post(client, base, auth=auth, payload=item)
|
||
mark = "继承" if item in inherited else "新增"
|
||
print(f" [{mark}] {item['namespace']}/{item['item_key']} → {response.status_code}")
|
||
if response.status_code != 201:
|
||
print(f" 失败:{response.text[:200]}")
|
||
return 1
|
||
|
||
release_base = f"/api/v1/admin/config-releases/{release_id}"
|
||
submitted = await post(
|
||
client, f"{release_base}/validations", auth=auth, payload={},
|
||
if_match=await etag_of(client, release_base, auth),
|
||
)
|
||
print(f"\n提交复核:{submitted.status_code}")
|
||
reviewed = await post(
|
||
client, f"{release_base}/reviews", auth=auth,
|
||
payload={"decision": "approved", "comment": "风控工具白名单"},
|
||
if_match=await etag_of(client, release_base, auth),
|
||
)
|
||
print(f"审核:{reviewed.status_code}")
|
||
activated = await post(
|
||
client, f"{release_base}/activations", auth=auth, payload={},
|
||
if_match=await etag_of(client, release_base, auth),
|
||
)
|
||
print(f"激活:{activated.status_code}")
|
||
if activated.status_code not in (200, 201):
|
||
print(f" 失败:{activated.text[:200]}")
|
||
return 1
|
||
print(f"最终状态:{activated.json()['data']['status']}")
|
||
|
||
remaining = await active_config_items()
|
||
print(f"\n激活后生效版本配置项:{len(remaining)} 条")
|
||
for item in remaining:
|
||
print(f" · {item['namespace']} / {item['item_key']} = {item['value_json']}")
|
||
return 0
|
||
|
||
|
||
sys.exit(asyncio.run(main()))
|