AgentEvalTool/tests/integration/test_exploration_settlement.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

195 lines
7.4 KiB
Python

"""Integration tests for exploration settlement on campaign finalize (v0.9 票据 07).
Covers both terminal transitions: window completion (durable loop COMPLETE)
and mid-flight cancellation. Running sessions must turn expired and reject
further messages (409), leaving no dangling sessions.
"""
import asyncio
import pytest
from agenteval.evaluation import campaign_runner as runner_module
from agenteval.evaluation.campaign_runner import CampaignRuntime
from agenteval.exploration.models import ExplorationSessionStatus
from agenteval.models import (
Campaign,
CampaignPlanEntry,
Case,
CaseType,
ChannelType,
EvalTarget,
PlatformType,
Scenario,
TargetStatus,
)
from agenteval.storage.db import utc_now
from agenteval.storage.exploration_repository import ExplorationSessionRepository
from agenteval.storage.repository import (
CampaignRepository,
ScenarioRepository,
TargetRepository,
)
from agenteval.web.app import app
from httpx import ASGITransport, AsyncClient
def _session_payload(campaign_id: str) -> dict:
return {
"campaign_id": campaign_id,
"persona": {"name": "巡检用户", "traits": ["耐心"]},
"goal": "查询账单",
"triggered_by": "manual",
}
@pytest.fixture()
def seeded_db(db_session, monkeypatch):
"""Point every get_session consumer at the test session and seed the
target + scenario needed by both the loop path and the cancel path."""
from agenteval.channels import factory as factory_module
from agenteval.evaluation import engine as engine_module
from agenteval.exploration import lifecycle as exploration_lifecycle
from agenteval.storage import db as db_module
from agenteval.storage import repository as repo_module
from agenteval.web import app as app_module
monkeypatch.setattr(app_module, "init_db", lambda: None)
def _test_get_session():
return db_session
monkeypatch.setattr(db_module, "get_session", _test_get_session)
monkeypatch.setattr(repo_module, "get_session", _test_get_session)
monkeypatch.setattr(engine_module, "get_session", _test_get_session)
monkeypatch.setattr(runner_module, "get_session", _test_get_session)
monkeypatch.setattr(runner_module, "_auto_start_analysis", lambda campaign, session: None)
from tests.unit.mock_channel import MockChannel
channel = MockChannel(reply_delay=0.0)
monkeypatch.setattr(factory_module.ChannelFactory, "create", lambda target: channel)
monkeypatch.setattr(exploration_lifecycle, "start_judge_review", lambda session_id: None)
from agenteval.web.deps import get_db
def _test_get_db():
try:
yield db_session
finally:
pass
app.dependency_overrides[get_db] = _test_get_db
target = EvalTarget(
id="t-1",
name="mock-target",
platform=PlatformType.AI_DIGITAL_EMPLOYEE,
channel_type=ChannelType.TUTU_API,
channel_config={"base_url": "http://mock", "token": "x"},
status=TargetStatus.ACTIVE,
)
TargetRepository(db_session).create(target)
scenario = Scenario(
id="s-1",
name="mock-scenario",
cases=[Case(id="c1", type=CaseType.SINGLE, messages=["hi"])],
)
ScenarioRepository(db_session).create(scenario)
yield db_session
app.dependency_overrides.clear()
def _make_campaign(db_session, campaign_id: str, *, time_scale: float) -> Campaign:
return CampaignRepository(db_session).create(
Campaign(
id=campaign_id,
name=f"campaign-{campaign_id}",
target_id="t-1",
window_seconds=3600,
time_scale=time_scale,
plan=[CampaignPlanEntry(scenario_id="s-1", offset_seconds=0, count=1)],
status="running",
started_at=utc_now(),
)
)
async def _create_running_session(db_session, campaign_id: str) -> str:
async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as c:
resp = await c.post("/api/exploration/sessions", json=_session_payload(campaign_id))
assert resp.status_code == 200, resp.text
session_id = resp.json()["id"]
assert ExplorationSessionRepository(db_session).get(session_id).status == (ExplorationSessionStatus.RUNNING)
return session_id
async def test_window_completion_expires_running_sessions(seeded_db):
campaign = _make_campaign(seeded_db, "c-done", time_scale=3600.0)
session_id = await _create_running_session(seeded_db, "c-done")
# time_scale 3600 → 1 real second covers the whole window.
runtime = CampaignRuntime(session_factory=lambda: seeded_db, tick_seconds=0.05)
try:
assert runtime.start(campaign.id)
async with asyncio.timeout(5):
while CampaignRepository(seeded_db).get(campaign.id).status.value == "running":
await asyncio.sleep(0.05)
finally:
await runtime.shutdown()
assert CampaignRepository(seeded_db).get("c-done").status.value == "completed"
settled = ExplorationSessionRepository(seeded_db).get(session_id)
assert settled.status == ExplorationSessionStatus.EXPIRED
assert settled.closed_at is not None
async def test_cancel_expires_running_sessions_and_rejects_messages(seeded_db):
_make_campaign(seeded_db, "c-1", time_scale=1.0)
session_id = await _create_running_session(seeded_db, "c-1")
async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client:
resp = await client.post("/api/campaigns/c-1/cancel")
assert resp.status_code == 200, resp.text
assert resp.json()["status"] == "cancelled"
settled = ExplorationSessionRepository(seeded_db).get(session_id)
assert settled.status == ExplorationSessionStatus.EXPIRED
assert settled.closed_at is not None
resp = await client.post(f"/api/exploration/sessions/{session_id}/messages", json={"content": "还在吗"})
assert resp.status_code == 409
async def test_completed_session_survives_settlement(seeded_db):
_make_campaign(seeded_db, "c-1", time_scale=1.0)
async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client:
resp = await client.post("/api/exploration/sessions", json=_session_payload("c-1"))
session_id = resp.json()["id"]
resp = await client.post(
f"/api/exploration/sessions/{session_id}/close",
json={"experience": {"goal_achieved": True}},
)
assert resp.status_code == 200
await client.post("/api/campaigns/c-1/cancel")
assert ExplorationSessionRepository(seeded_db).get(session_id).status == (ExplorationSessionStatus.COMPLETED)
async def test_production_line_manual_trigger_also_settles(seeded_db):
_make_campaign(seeded_db, "c-1", time_scale=1.0)
async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client:
resp = await client.post(
"/api/exploration/sessions",
json={**_session_payload("c-1"), "triggered_by": "manual"},
)
assert resp.status_code == 200, resp.text
session_id = resp.json()["id"]
assert resp.json()["triggered_by"] == "manual"
# 正式线手动会话同样纳入终态结算,触发来源不影响收口。
await client.post("/api/campaigns/c-1/cancel")
assert ExplorationSessionRepository(seeded_db).get(session_id).status == (ExplorationSessionStatus.EXPIRED)