AgentEvalTool/backend/agenteval/evaluation/analysis.py
sinohqb 7eae6de52d refactor(evaluation/storage): 结算统一与 repository 拆分(Phase 2 + 3)
合并两个不可分割的深化:

Phase 2 — 智能作业结算统一(ADR-0012)
- intelligence_jobs.execute(job_kind, campaign_id, ...) 作为结算的
  唯一实现:建行 → 认领 → 校验 → generating → 落账,一处编排、
  一处截断(500 字符)。两个 executor 退化为 ensure_queued /
  validate / work_fn 三个小 adapter。
- analysis.validate_analysis_request() 共享校验入口(活动终态 →
  模型),路由捕获映射 400、executor 捕获落 failed 行,与
  validate_comparison_request 先例同构。
- campaign_runner._auto_start_analysis 的跳过守卫收敛至
  auto_intelligence_eligible 单一判断点。
- comparison.py 删除零调用的 build_comparison_payload;
  load_comparison_view 投影归位至 campaign_read_model。
- 新增 characterization 测试(认领竞争、重复触发、截断、恢复上限)。

Phase 3 — storage/repository.py 拆分
- AsyncJobRepository 及两个子类迁至
  storage/async_job_repository.py(Phase 2 的 intelligence_jobs
  与 comparison 必须 import 自该路径,故与 Phase 2 同 commit)。
- ExplorationSession / ExplorationMessage 迁至
  storage/exploration_repository.py;repository.py 由 1180 行降至
  约 814 行,grep 确认无残留符号。
- exploration 子模块与路由 import 全部更新;测试 import 跟随。

刻意不做:CAS 共享原语、app.py 五 registry 关停顺序归一
(ADR-0006 精神,等真实需求出现再议)。
2026-08-24 05:50:27 +08:00

264 lines
11 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Campaign intelligence analysis — the two-phase analysis agent (分析岗位).
对终态活动的聚合结果做活动级、跨场景的叙述性研判CONTEXT.md「分析岗位」
输入完全复用 ``generate_campaign_report`` 的既有聚合口径ADR-0002/0004不重算
数字),外加每场景少量代表性失败对话样例;阶段一按场景并行诊断,阶段二综合
研判产出结构化报告。LLM 调用经由 ``chat_client`` 注入,测试用假客户端替换。
"""
import asyncio
import json
from typing import Any, Awaitable, Callable, Optional
from sqlmodel import Session
from agenteval.model_gateway import ModelGateway
from agenteval.models import Campaign, CampaignStatus, ModelCapability, RunStatus
from agenteval.services.model_configs import (
ModelConfigError,
ModelConfigService,
ModelRuntimeConfig,
)
from agenteval.storage.model_config_repository import ModelConfigRepository
from agenteval.storage.repository import (
RunRepository,
)
from agenteval.utils.llm import extract_reply_text, parse_json_from_llm_text
# LLM 客户端协议:接收 chat 消息列表,返回文本内容。生产实现走 ModelGateway
# 测试注入假客户端(同 MockChannel 先例)。
ChatClient = Callable[[list[dict[str, str]]], Awaitable[str]]
MAX_SAMPLES_PER_SCENARIO = 3
SAMPLE_TEXT_LIMIT = 200
_VALID_SEVERITIES = {"high", "medium", "low"}
class AnalysisError(RuntimeError):
"""分析生成失败(数据缺失或模型输出无法解析),可重试。"""
def resolve_analysis_model(campaign: Campaign, session: Session) -> Optional[ModelRuntimeConfig]:
"""解析该活动应使用的分析模型:活动覆盖 ?? 全局分析默认;解析不到返回 None。"""
config_id = campaign.analysis_model_config_id
if config_id is None:
default = ModelConfigRepository(session).get_analysis_default()
config_id = default.id if default else None
if config_id is None:
return None
try:
return ModelConfigService(session).resolve(config_id, expected_capability=ModelCapability.CHAT)
except ModelConfigError:
return None
def validate_analysis_request(session: Session, campaign: Campaign) -> ModelRuntimeConfig:
"""共享校验入口:活动终态 → 模型。违规抛 AnalysisError。
router 触发端点捕获映射 400执行器捕获落 failed 行。校验顺序权威,
两处不再漂移validate_comparison_request 先例)。
"""
if campaign.status in (CampaignStatus.PLANNED, CampaignStatus.RUNNING):
raise AnalysisError("活动完成后才能生成智能分析")
runtime = resolve_analysis_model(campaign, session)
if runtime is None:
raise AnalysisError("未配置分析模型:请在模型配置中心将某个 chat 配置设为「分析默认」,或为该活动指定分析模型")
return runtime
def collect_failure_samples(
campaign_id: str,
session: Session,
*,
per_scenario: int = MAX_SAMPLES_PER_SCENARIO,
text_limit: int = SAMPLE_TEXT_LIMIT,
) -> dict[str, list[dict[str, str]]]:
"""每场景最多 ``per_scenario`` 条代表性失败对话(用户消息/回复/判定理由,截断)。"""
run_repo = RunRepository(session)
samples: dict[str, list[dict[str, str]]] = {}
for run in run_repo.list_by_campaign(campaign_id):
if run.status != RunStatus.COMPLETED:
continue
failed = [r for r in run_repo.get_results(run.id) if not r.passed]
if not failed:
continue
turns = {t.id: t for t in run_repo.get_turns(run.id)}
bucket = samples.setdefault(run.scenario_id, [])
for result in failed:
if len(bucket) >= per_scenario:
break
turn = turns.get(result.turn_id)
user = extract_reply_text(turn.get_sent_message().get("msgBody")) if turn else ""
reply = extract_reply_text(turn.get_reply().get("msgBody")) if turn and turn.get_reply() else ""
bucket.append(
{
"run_id": run.id or "",
"user": user[:text_limit],
"reply": reply[:text_limit],
"reason": (result.reason or "")[:text_limit],
}
)
return {sid: items for sid, items in samples.items() if items}
def _parse_stage(content: str, label: str) -> dict[str, Any]:
try:
parsed = parse_json_from_llm_text(content)
except Exception as exc:
raise AnalysisError(f"{label}输出解析失败: {exc}") from exc
if not isinstance(parsed, dict):
raise AnalysisError(f"{label}输出不是 JSON 对象")
return parsed
async def _analyze_scenario(
entry: dict[str, Any],
samples: list[dict[str, str]],
chat_client: ChatClient,
) -> dict[str, Any]:
"""阶段一:单个场景的诊断(叙述 + 问题点草稿)。"""
system_prompt = (
"你是智能客服质量评估平台的分析专家,负责对一次评估活动中某个场景的表现做诊断。"
"只输出一个 JSON 对象:"
'{"narrative": "该场景的叙述性表现分析2-4 句)", '
'"problems": [{"severity": "high|medium|low", "title": "...", "description": "...", '
'"evidence_run_ids": ["来自输入数据的真实 run_id"]}]}'
";没有问题时 problems 为空数组。全部使用中文。"
)
user_prompt = json.dumps(
{
"场景": {"id": entry["scenario_id"], "名称": entry.get("scenario_name", "")},
"聚合指标": {
"执行次数": entry.get("run_count"),
"通过率": entry.get("pass_rate"),
"可用性": entry.get("availability"),
"平均时延ms": entry.get("avg_latency_ms"),
},
"代表性失败对话": samples,
},
ensure_ascii=False,
)
parsed = _parse_stage(
await chat_client(
[
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt},
]
),
f"场景「{entry.get('scenario_name', entry['scenario_id'])}」阶段一",
)
narrative = parsed.get("narrative")
if not isinstance(narrative, str) or not narrative.strip():
raise AnalysisError(f"场景「{entry.get('scenario_name', entry['scenario_id'])}」阶段一缺少 narrative")
return {
"scenario_id": entry["scenario_id"],
"narrative": narrative,
"problems": parsed.get("problems") if isinstance(parsed.get("problems"), list) else [],
}
async def _synthesize(
campaign: Campaign,
report: dict[str, Any],
stage1: list[dict[str, Any]],
chat_client: ChatClient,
exploration_summary: Optional[dict[str, Any]] = None,
) -> dict[str, Any]:
"""阶段二:汇总各场景产出,产总体结论 + 跨场景问题 + 优先级建议。"""
system_prompt = (
"你是智能客服质量评估平台的首席分析专家,负责对整个评估活动做综合研判。"
"只输出一个 JSON 对象:"
'{"overall": "总体结论(一段话)", '
'"problems": [{"severity": "high|medium|low", "title": "...", "description": "...", '
'"scenario_ids": ["涉及场景 id"], "evidence_run_ids": ["来自输入数据的真实 run_id"]}], '
'"suggestions": [{"priority": 1, "text": "可执行的改善建议"}]}'
";问题按严重度从高到低排列,建议按优先级排列。全部使用中文。"
)
payload: dict[str, Any] = {
"活动": {
"名称": campaign.name,
"窗口秒数": campaign.window_seconds,
"总体指标": report.get("summary", {}),
},
"各场景诊断": stage1,
}
if exploration_summary is not None:
# 探索式评测证据线:只给统计与问题清单,不含全量对话
payload["探索发现"] = exploration_summary
user_prompt = json.dumps(payload, ensure_ascii=False)
parsed = _parse_stage(
await chat_client(
[
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt},
]
),
"阶段二综合研判",
)
overall = parsed.get("overall")
if not isinstance(overall, str) or not overall.strip():
raise AnalysisError("阶段二综合研判缺少 overall")
return parsed
async def analyze_campaign(
*,
campaign: Campaign,
report: dict[str, Any],
failure_samples: dict[str, list[dict[str, str]]],
valid_run_ids: set[str],
chat_client: ChatClient,
exploration_summary: Optional[dict[str, Any]] = None,
) -> dict[str, Any]:
"""两阶段编排:阶段一按场景并行诊断,阶段二综合研判。
输出遵循结构化报告 schemaoverall/problems/scenario_narratives/suggestions
模型虚构的 run_id / scenario_id 在返回前按白名单剔除;任何解析失败抛
``AnalysisError``(由调用方落 failed 状态)。
"""
capability = report.get("capability_summary") or []
if not capability:
raise AnalysisError("活动没有可分析的场景数据")
stage1 = await asyncio.gather(
*[_analyze_scenario(entry, failure_samples.get(entry["scenario_id"], []), chat_client) for entry in capability]
)
stage2 = await _synthesize(campaign, report, list(stage1), chat_client, exploration_summary=exploration_summary)
valid_scenario_ids = {entry["scenario_id"] for entry in capability}
problems = []
for p in stage2.get("problems") or []:
if not isinstance(p, dict):
continue
severity = p.get("severity")
problems.append(
{
"severity": severity if severity in _VALID_SEVERITIES else "medium",
"title": str(p.get("title", "")),
"description": str(p.get("description", "")),
"scenario_ids": [s for s in p.get("scenario_ids") or [] if s in valid_scenario_ids],
"evidence_run_ids": [r for r in p.get("evidence_run_ids") or [] if r in valid_run_ids],
}
)
suggestions = [
{"priority": int(s.get("priority", i + 1)), "text": str(s.get("text", ""))}
for i, s in enumerate(stage2.get("suggestions") or [])
if isinstance(s, dict)
]
return {
"overall": stage2["overall"],
"problems": problems,
"scenario_narratives": [{"scenario_id": s["scenario_id"], "narrative": s["narrative"]} for s in stage1],
"suggestions": suggestions,
}
def gateway_chat_client(runtime: ModelRuntimeConfig) -> ChatClient:
"""Shared ChatClient factory for analysis/comparison background executors."""
gateway = ModelGateway(timeout=180.0)
async def _chat(messages: list[dict[str, str]]) -> str:
return await gateway.chat(runtime, messages, temperature=0.2)
return _chat