From 66329e715468e21835f80fbdc4d3e40375850755 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 3 Oct 2026 18:01:39 -0700 Subject: [PATCH] test(lens): cover in-flight reporting under original run ids --- tests/unit/proxy/lens/test_analysis.py | 79 ++++++++++++++++++++++---- 1 file changed, 69 insertions(+), 10 deletions(-) diff --git a/tests/unit/proxy/lens/test_analysis.py b/tests/unit/proxy/lens/test_analysis.py index ea353137bbf..832a40c0fd3 100644 --- a/tests/unit/proxy/lens/test_analysis.py +++ b/tests/unit/proxy/lens/test_analysis.py @@ -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", ()), + )