合并两个不可分割的深化: 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 精神,等真实需求出现再议)。
156 lines
5.2 KiB
Python
156 lines
5.2 KiB
Python
"""Integration tests for the runner's auto analysis hook (v0.7 ticket 05).
|
||
|
||
Only 正式线 campaigns (time_scale == 1) enqueue the analysis task on
|
||
COMPLETED; accelerated/cancelled campaigns and missing analysis models all
|
||
skip silently. The analysis service itself is spied, not executed.
|
||
"""
|
||
|
||
import pytest
|
||
from agenteval.evaluation import campaign_runner
|
||
from agenteval.evaluation.campaign_lifecycle import cancel_campaign
|
||
from agenteval.evaluation.campaign_runner import CampaignRuntime
|
||
from agenteval.models import (
|
||
Campaign,
|
||
CampaignPlanEntry,
|
||
CampaignStatus,
|
||
Case,
|
||
CaseType,
|
||
ChannelType,
|
||
EvalTarget,
|
||
PlatformType,
|
||
Scenario,
|
||
TargetStatus,
|
||
)
|
||
from agenteval.storage.repository import CampaignRepository, ScenarioRepository, TargetRepository
|
||
|
||
from tests.unit.mock_channel import MockChannel
|
||
|
||
TICK = 0.01
|
||
|
||
|
||
@pytest.fixture()
|
||
def seeded_db(db_session, monkeypatch):
|
||
from agenteval.channels import factory as factory_module
|
||
from agenteval.evaluation import engine as engine_module
|
||
from agenteval.storage import db as db_module
|
||
from agenteval.storage import repository as repo_module
|
||
|
||
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(campaign_runner, "get_session", _test_get_session)
|
||
|
||
channel = MockChannel(reply_delay=0.0)
|
||
monkeypatch.setattr(factory_module.ChannelFactory, "create", lambda target: channel)
|
||
|
||
TargetRepository(db_session).create(
|
||
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,
|
||
)
|
||
)
|
||
ScenarioRepository(db_session).create(
|
||
Scenario(
|
||
id="s-1",
|
||
name="mock-scenario",
|
||
cases=[Case(id="c1", type=CaseType.SINGLE, messages=["hi"])],
|
||
)
|
||
)
|
||
return db_session
|
||
|
||
|
||
@pytest.fixture()
|
||
def analysis_spy(monkeypatch):
|
||
"""Spy the analysis seam: resolvable model, recorded enqueue calls."""
|
||
from agenteval.evaluation import analysis as analysis_module
|
||
|
||
calls: list[tuple[str, str]] = []
|
||
monkeypatch.setattr(
|
||
campaign_runner,
|
||
"enqueue_campaign_analysis",
|
||
lambda cid, *, triggered_by: calls.append((cid, triggered_by)),
|
||
)
|
||
monkeypatch.setattr(analysis_module, "resolve_analysis_model", lambda campaign, session: object())
|
||
return calls
|
||
|
||
|
||
def _make_campaign(session, **overrides) -> Campaign:
|
||
payload = dict(
|
||
name="auto-trigger",
|
||
target_id="t-1",
|
||
window_seconds=1,
|
||
time_scale=1.0, # 正式线:1s 窗口真实耗时 ~1s
|
||
plan=[CampaignPlanEntry(scenario_id="s-1", offset_seconds=0, count=1)],
|
||
)
|
||
payload.update(overrides)
|
||
return CampaignRepository(session).create(Campaign(**payload))
|
||
|
||
|
||
@pytest.fixture()
|
||
async def runtime(seeded_db):
|
||
value = CampaignRuntime(session_factory=lambda: seeded_db, tick_seconds=TICK)
|
||
yield value
|
||
await value.shutdown()
|
||
|
||
|
||
async def _await_terminal(session, campaign_id, timeout=5.0):
|
||
import asyncio
|
||
|
||
async with asyncio.timeout(timeout):
|
||
while CampaignRepository(session).get(campaign_id).status is CampaignStatus.RUNNING:
|
||
await asyncio.sleep(TICK)
|
||
|
||
|
||
async def test_realtime_completion_auto_enqueues_analysis(seeded_db, runtime, analysis_spy):
|
||
campaign = _make_campaign(seeded_db)
|
||
assert runtime.start(campaign.id)
|
||
await _await_terminal(seeded_db, campaign.id)
|
||
|
||
final = CampaignRepository(seeded_db).get(campaign.id)
|
||
assert final.status == CampaignStatus.COMPLETED
|
||
assert analysis_spy == [(campaign.id, "auto")]
|
||
|
||
|
||
async def test_accelerated_completion_does_not_enqueue(seeded_db, runtime, analysis_spy):
|
||
campaign = _make_campaign(seeded_db, time_scale=1000.0) # 加速调试线
|
||
assert runtime.start(campaign.id)
|
||
await _await_terminal(seeded_db, campaign.id)
|
||
|
||
final = CampaignRepository(seeded_db).get(campaign.id)
|
||
assert final.status == CampaignStatus.COMPLETED
|
||
assert analysis_spy == []
|
||
|
||
|
||
async def test_cancelled_campaign_does_not_enqueue(seeded_db, runtime, analysis_spy):
|
||
import asyncio
|
||
|
||
campaign = _make_campaign(seeded_db, window_seconds=100)
|
||
assert runtime.start(campaign.id)
|
||
await asyncio.sleep(0.05)
|
||
|
||
repo = CampaignRepository(seeded_db)
|
||
cancel_campaign(seeded_db, campaign.id, stop=runtime.cancel)
|
||
|
||
assert repo.get(campaign.id).status == CampaignStatus.CANCELLED
|
||
assert analysis_spy == []
|
||
|
||
|
||
async def test_missing_analysis_model_skips_silently(seeded_db, runtime, monkeypatch, analysis_spy):
|
||
from agenteval.evaluation import analysis as analysis_module
|
||
|
||
monkeypatch.setattr(analysis_module, "resolve_analysis_model", lambda campaign, session: None)
|
||
campaign = _make_campaign(seeded_db)
|
||
assert runtime.start(campaign.id)
|
||
await _await_terminal(seeded_db, campaign.id)
|
||
|
||
final = CampaignRepository(seeded_db).get(campaign.id)
|
||
assert final.status == CampaignStatus.COMPLETED # 活动完成流程不受影响
|
||
assert analysis_spy == []
|