test(lens): cover in-flight reporting under original run ids

This commit is contained in:
Ishaan Jaff 2026-10-03 18:01:39 -07:00
parent 46f16ea03e
commit 66329e7154
No known key found for this signature in database

View file

@ -13,6 +13,7 @@ from litellm.proxy.lens.models import (
Evidence,
Execution,
ExecutionContent,
InFlight,
ModelRequest,
ModelResult,
Review,
@ -67,8 +68,10 @@ async def test_parallel_review_shares_one_model_limit_and_cleans_up(outcome: str
finally:
exited.put(request.prompt)
async def progress(stage: str, coverage: Coverage, _review: Review | None = None, /) -> None:
if stage == "Reading executions":
async def progress(
stage: str, coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
if stage == "Reading executions" and (_reading is None or _review is not None):
counts.put(coverage.screened)
claim: Final = Claim(lens_id="lens", job=queue_job(lens(), NOW, "job").jobs[0], findings=())
@ -118,7 +121,9 @@ async def test_independent_investigations_overlap_and_report_completions() -> No
async def read(_execution_id: str, _cursor: str, _offset: int) -> ExecutionContent:
pytest.fail("Inconclusive decisions must not fetch evidence")
async def progress(stage: str, coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
stage: str, coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
assert stage == "Checking original evidence"
progress_counts.put(coverage.investigated)
@ -465,7 +470,9 @@ async def test_grouping_consolidates_prior_batches_and_reports_real_progress() -
)
stages: Final = iter((0, 1))
async def progress(stage: str, coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
stage: str, coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
assert stage == "Grouping observations"
assert coverage.grouping_batches == 2
assert coverage.grouped_batches == next(stages)
@ -559,7 +566,9 @@ async def test_thousands_of_matching_runs_keep_all_members_without_a_growing_mod
cost=0,
)
async def progress(_stage: str, coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
_stage: str, coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
counts.put(coverage.grouped_batches)
batches: Final = observation_batches(observations)
@ -633,7 +642,9 @@ async def test_review_keeps_original_ids_in_per_run_assessments() -> None:
async def model(_request: ModelRequest) -> ModelResult:
return ModelResult(content='{"observations":[],"cannot_assess":false}', cost=0)
async def progress(_stage: str, _coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
_stage: str, _coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
pass
claim: Final = Claim(lens_id="lens", job=queue_job(lens(), NOW, "job").jobs[0], findings=())
@ -892,7 +903,9 @@ async def test_final_registry_reconciles_patterns_split_across_pages() -> None:
)
return ModelResult(content=Clusters(candidates=grouped).model_dump_json(), cost=0)
async def progress(_stage: str, _coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
_stage: str, _coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
return None
result: Final = await cluster_batches((observations,), model, progress, Coverage())
@ -919,7 +932,9 @@ async def test_distinct_patterns_are_consolidated_in_batches_without_losing_runs
payload: Final = json.loads(request.prompt)
return ModelResult(content=json.dumps({"candidates": payload["candidates"]}), cost=0)
async def progress(_stage: str, _coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
_stage: str, _coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
pass
result: Final = await cluster_batches(observation_batches(observations), model, progress, Coverage())
@ -951,7 +966,9 @@ async def test_invalid_candidate_response_preserves_other_findings_and_reports_i
return ModelResult(content="not JSON", cost=0)
return ModelResult(content=json.dumps({"action": "submit", "finding": finding("run").model_dump()}), cost=0)
async def progress(_stage: str, coverage: Coverage, _review: Review | None = None, /) -> None:
async def progress(
_stage: str, coverage: Coverage, _review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
counts.put(coverage.inconclusive)
claim: Final = Claim(lens_id="lens", job=queue_job(lens(), NOW, "job").jobs[0], findings=())
@ -1287,7 +1304,9 @@ async def test_each_screened_run_reports_a_review_with_the_models_reasoning() ->
cost=0,
)
async def progress(_stage: str, _coverage: Coverage, review: Review | None = None, /) -> None:
async def progress(
_stage: str, _coverage: Coverage, review: Review | None = None, _reading: tuple[InFlight, ...] | None = None, /
) -> None:
if review is not None:
reviews.put(review)
@ -1300,3 +1319,43 @@ async def test_each_screened_run_reports_a_review_with_the_models_reasoning() ->
assert review.model == claim.job.settings.model
assert tuple((s.span_id, s.cited) for s in review.spans) == (("a-root", False), ("b-tool", True))
assert tuple(v.summary for v in review.verdicts) == ("Gave up after a timeout",)
@pytest.mark.asyncio
async def test_a_run_is_reported_in_flight_under_its_original_id_until_its_review_arrives() -> None:
from litellm.proxy.lens.analysis import analyze_sample
execution: Final = Execution(
id="opaque-original", source="traces", trace_id="trace", team_id="", name="task", start_time="", span_count=1
)
reports: Final = SimpleQueue[tuple[str | None, tuple[str, ...] | None]]()
async def read(identity: str, _cursor: str, _offset: int) -> ExecutionContent:
return ExecutionContent(
execution=execution,
parts=(TracePart(execution_id=identity, span_id="s", name="agent", kind="agent", content="Hi"),),
)
async def model(request: ModelRequest) -> ModelResult:
if request.purpose == "cluster":
return ModelResult(content='{"candidates":[]}', cost=0)
return ModelResult(content='{"observations":[]}', cost=0)
async def progress(
stage: str, _coverage: Coverage, review: Review | None = None, reading: tuple[InFlight, ...] | None = None, /
) -> None:
if stage == "Reading executions":
reports.put(
(
review and review.execution_id,
None if reading is None else tuple(f"{r.execution_id}:{r.trace_id}" for r in reading),
)
)
claim: Final = Claim(lens_id="lens", job=queue_job(lens(), NOW, "job").jobs[0], findings=())
await analyze_sample(claim, Sample(executions=(execution,), eligible=1), read, model, progress)
assert tuple(reports.get_nowait() for _ in range(reports.qsize())) == (
(None, None),
(None, ("opaque-original:trace",)),
("opaque-original", ()),
)