合并两个不可分割的深化: 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 精神,等真实需求出现再议)。
234 lines
8.9 KiB
Python
234 lines
8.9 KiB
Python
"""Stable read projections for evaluation Campaigns."""
|
||
|
||
from typing import Any, Optional
|
||
|
||
from sqlmodel import Session
|
||
|
||
from agenteval.evaluation.campaign_scheduler import clock_offset, elapsed_seconds
|
||
from agenteval.evaluation.comparison import compute_metric_diff, resolve_auto_baseline
|
||
from agenteval.evaluation.report import (
|
||
build_campaign_timeline,
|
||
generate_campaign_report,
|
||
load_campaign_report,
|
||
summarize_campaign_progress,
|
||
)
|
||
from agenteval.exploration.summary import summarize_campaign_exploration
|
||
from agenteval.models import Campaign
|
||
from agenteval.storage.async_job_repository import CampaignAnalysisRepository, CampaignPeriodComparisonRepository
|
||
from agenteval.storage.db import iso_utc, utc_now
|
||
from agenteval.storage.model_config_repository import ModelConfigRepository
|
||
from agenteval.storage.repository import (
|
||
CampaignRepository,
|
||
RunRepository,
|
||
ScenarioRepository,
|
||
TargetRepository,
|
||
)
|
||
|
||
|
||
def load_comparison_view(session: Session, campaign: Campaign) -> dict[str, Any]:
|
||
"""周期对比读模型单一出口:返回 GET /comparison 完整响应形状。
|
||
|
||
无行时 status=none + auto_baseline;有行时含 comparison dict(含 model_name
|
||
标签)+ 对生效基线的 metric_diff。markdown 导出从同一 view 投影。
|
||
"""
|
||
payload = _build_auto_baseline_payload(campaign, session)
|
||
row = CampaignPeriodComparisonRepository(session).get_by_campaign(campaign.id)
|
||
if row is None:
|
||
return {"status": "none", "comparison": None, **payload}
|
||
|
||
effective_baseline = CampaignRepository(session).get(row.baseline_campaign_id)
|
||
metric_diff = (
|
||
compute_metric_diff(
|
||
load_campaign_report(session, effective_baseline),
|
||
load_campaign_report(session, campaign),
|
||
)
|
||
if effective_baseline is not None
|
||
else None
|
||
)
|
||
model_cfg = ModelConfigRepository(session).get(row.model_config_id) if row.model_config_id else None
|
||
model_label = (
|
||
f"{model_cfg.name}({model_cfg.model_name})"
|
||
if model_cfg and model_cfg.model_name
|
||
else (model_cfg.name if model_cfg else None)
|
||
)
|
||
comparison = {
|
||
"baseline_campaign_id": row.baseline_campaign_id,
|
||
"baseline": (
|
||
{
|
||
"id": effective_baseline.id,
|
||
"name": effective_baseline.name,
|
||
"completed_at": iso_utc(effective_baseline.completed_at),
|
||
}
|
||
if effective_baseline
|
||
else None
|
||
),
|
||
"result": row.get_result(),
|
||
"error": row.error,
|
||
"model_config_id": row.model_config_id,
|
||
"model_name": model_label,
|
||
"triggered_by": row.triggered_by,
|
||
"updated_at": iso_utc(row.updated_at),
|
||
}
|
||
return {
|
||
"status": row.status,
|
||
"comparison": comparison,
|
||
"auto_baseline": payload["auto_baseline"],
|
||
"metric_diff": metric_diff,
|
||
}
|
||
|
||
|
||
def _build_auto_baseline_payload(campaign: Campaign, session: Session) -> dict[str, Any]:
|
||
"""自动基线信息 + 机械 diff(无基线时两者均为 null)。"""
|
||
baseline = resolve_auto_baseline(campaign, session)
|
||
if baseline is None:
|
||
return {"auto_baseline": None, "metric_diff": None}
|
||
diff = compute_metric_diff(
|
||
load_campaign_report(session, baseline),
|
||
load_campaign_report(session, campaign),
|
||
)
|
||
return {
|
||
"auto_baseline": {
|
||
"id": baseline.id,
|
||
"name": baseline.name,
|
||
"completed_at": iso_utc(baseline.completed_at),
|
||
},
|
||
"metric_diff": diff,
|
||
}
|
||
|
||
|
||
class CampaignReadModel:
|
||
"""One interface for Campaign list, detail, report and export projections."""
|
||
|
||
def __init__(self, session: Session):
|
||
self._session = session
|
||
self._campaigns = CampaignRepository(session)
|
||
self._runs = RunRepository(session)
|
||
self._scenarios = ScenarioRepository(session)
|
||
self._analyses = CampaignAnalysisRepository(session)
|
||
|
||
def list_items(self) -> list[dict[str, Any]]:
|
||
campaigns = self._campaigns.list_all()
|
||
campaign_ids = [campaign.id for campaign in campaigns if campaign.id]
|
||
runs_by_campaign = self._runs.list_by_campaigns(campaign_ids)
|
||
items = []
|
||
for campaign in campaigns:
|
||
data = campaign.model_dump()
|
||
data["progress"] = summarize_campaign_progress(
|
||
campaign,
|
||
runs_by_campaign.get(campaign.id, []),
|
||
)
|
||
items.append(data)
|
||
return items
|
||
|
||
def detail(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
campaign = self._campaigns.get(campaign_id)
|
||
if campaign is None:
|
||
return None
|
||
data = campaign.model_dump()
|
||
runs = self._runs.list_by_campaign(campaign_id)
|
||
current_offset = 0.0
|
||
if campaign.started_at is not None:
|
||
current_offset = min(
|
||
clock_offset(
|
||
elapsed_seconds=elapsed_seconds(now=utc_now(), started_at=campaign.started_at),
|
||
time_scale=campaign.time_scale,
|
||
),
|
||
float(campaign.window_seconds),
|
||
)
|
||
data["progress"] = {
|
||
"current_offset_seconds": current_offset,
|
||
"spawned_runs": len(runs),
|
||
"completed_runs": sum(1 for run in runs if run.status.value == "completed"),
|
||
}
|
||
return data
|
||
|
||
def report(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
campaign = self._campaigns.get(campaign_id)
|
||
if campaign is None:
|
||
return None
|
||
report = self._core_report(campaign)
|
||
exploration = summarize_campaign_exploration(self._session, campaign_id)
|
||
if exploration is not None:
|
||
report["exploration"] = exploration
|
||
return report
|
||
|
||
def timeline(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
campaign = self._campaigns.get(campaign_id)
|
||
if campaign is None:
|
||
return None
|
||
entries = build_campaign_timeline(
|
||
campaign,
|
||
self._runs.list_by_campaign(campaign_id),
|
||
scenario_names=self._scenarios.name_map(),
|
||
)
|
||
return {"entries": entries}
|
||
|
||
def analysis(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
if self._campaigns.get(campaign_id) is None:
|
||
return None
|
||
row = self._analyses.get_by_campaign(campaign_id)
|
||
if row is None:
|
||
return {"status": "none"}
|
||
return {
|
||
"status": row.status,
|
||
"result": row.get_result(),
|
||
"error": row.error,
|
||
"model_config_id": row.model_config_id,
|
||
"triggered_by": row.triggered_by,
|
||
"updated_at": iso_utc(row.updated_at),
|
||
}
|
||
|
||
def comparison(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
campaign = self._campaigns.get(campaign_id)
|
||
return load_comparison_view(self._session, campaign) if campaign is not None else None
|
||
|
||
def full_view(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
campaign = self._campaigns.get(campaign_id)
|
||
if campaign is None:
|
||
return None
|
||
report = self._core_report(campaign)
|
||
exploration = summarize_campaign_exploration(self._session, campaign_id)
|
||
analysis = self.analysis(campaign_id)
|
||
comparison = load_comparison_view(self._session, campaign)
|
||
return {
|
||
"report": report,
|
||
"exploration": exploration,
|
||
"analysis": analysis.get("result") if analysis and analysis.get("status") == "completed" else None,
|
||
"comparison": comparison if comparison.get("status") != "none" else None,
|
||
}
|
||
|
||
def markdown_projection(self, campaign_id: str) -> Optional[dict[str, Any]]:
|
||
campaign = self._campaigns.get(campaign_id)
|
||
view = self.full_view(campaign_id)
|
||
if campaign is None or view is None:
|
||
return None
|
||
|
||
comparison = None
|
||
comparison_view = view["comparison"]
|
||
if comparison_view and comparison_view.get("status") == "completed":
|
||
row = comparison_view.get("comparison") or {}
|
||
baseline = row.get("baseline") or {}
|
||
comparison = {
|
||
"result": row.get("result"),
|
||
"baseline_name": baseline.get("name"),
|
||
"baseline_completed_at": baseline.get("completed_at"),
|
||
"model_name": row.get("model_name"),
|
||
"updated_at": row.get("updated_at"),
|
||
"metric_diff": comparison_view.get("metric_diff"),
|
||
}
|
||
|
||
target = TargetRepository(self._session).get(campaign.target_id)
|
||
return {
|
||
**view,
|
||
"comparison": comparison,
|
||
"target_name": target.name if target else None,
|
||
"scenario_names": self._scenarios.name_map(),
|
||
}
|
||
|
||
def _core_report(self, campaign) -> dict[str, Any]:
|
||
return generate_campaign_report(
|
||
campaign,
|
||
self._runs.list_by_campaign(campaign.id),
|
||
scenario_names=self._scenarios.name_map(),
|
||
)
|