From 1d18be780e009a3d73357c97bed15ed5f9b24341 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 3 Oct 2026 18:01:39 -0700 Subject: [PATCH] feat(lens): report each run as in flight when its review starts --- litellm/proxy/lens/analysis.py | 51 ++++++++++++++++++++++++++++------ 1 file changed, 42 insertions(+), 9 deletions(-) diff --git a/litellm/proxy/lens/analysis.py b/litellm/proxy/lens/analysis.py index ef391c1e146..84061101a6b 100644 --- a/litellm/proxy/lens/analysis.py +++ b/litellm/proxy/lens/analysis.py @@ -18,6 +18,7 @@ from .models import ( Execution, ExecutionContent, FindingDraft, + InFlight, ModelRequest, ModelResult, Record, @@ -105,7 +106,14 @@ ReadContent: TypeAlias = Callable[[str, str, int], Awaitable[ExecutionContent]] class ReportProgress(Protocol): - def __call__(self, stage: str, coverage: Coverage, review: Review | None = None, /) -> Awaitable[None]: ... + def __call__( + self, + stage: str, + coverage: Coverage, + review: Review | None = None, + reading: tuple[InFlight, ...] | None = None, + /, + ) -> Awaitable[None]: ... ResponseT = TypeVar("ResponseT", bound=Record) @@ -684,11 +692,21 @@ async def analyze_sample( ) ) - async def progress_original(stage: str, coverage: Coverage, review: Review | None = None, /) -> None: + def original(identity: str) -> str: + return originals[identity].id + + async def progress_original( + stage: str, coverage: Coverage, review: Review | None = None, reading: tuple[InFlight, ...] | None = None, / + ) -> None: await progress( stage, coverage, - review and review.model_copy(update=MappingProxyType({"execution_id": originals[review.execution_id].id})), + review and review.model_copy(update=MappingProxyType({"execution_id": original(review.execution_id)})), + None + if reading is None + else tuple( + r.model_copy(update=MappingProxyType({"execution_id": original(r.execution_id)})) for r in reading + ), ) result: Final = await _analyze_sample( @@ -916,21 +934,36 @@ async def merge_candidates( async def examine_executions( claim: Claim, sample: Sample, read: ReadContent, model: ModelCall, progress: ReportProgress ) -> AsyncIterator[Examined]: + reading: tuple[InFlight, ...] = () # rebind-ok: the in-flight set changes as each read starts and finishes + screened = 0 # rebind-ok: counts finished reads for progress + reporting: Final = asyncio.Lock() + + async def report(change: Callable[[tuple[InFlight, ...]], tuple[InFlight, ...]], review: Review | None) -> None: + nonlocal reading + async with reporting: + reading = change(reading) + coverage: Final = Coverage(eligible=sample.eligible, selected=len(sample.executions), screened=screened) + await progress("Reading executions", coverage, review, reading) + async def examine(execution: Execution) -> tuple[Examined, Review]: + entry: Final = InFlight( + execution_id=execution.id, + trace_id=execution.trace_id, + agent=execution.service or execution.name, + started_at=datetime.now(timezone.utc), + ) + await report(lambda current: (*current, entry), None) started: Final = time.perf_counter() examined: Final = await extract(claim, execution, read, model) elapsed: Final = round((time.perf_counter() - started) * 1000) return examined, review_of(examined, claim.job.settings.model, elapsed, datetime.now(timezone.utc)) await progress("Reading executions", Coverage(eligible=sample.eligible, selected=len(sample.executions))) - completed: Final = iter(range(1, len(sample.executions) + 1)) async with aclosing(concurrent_results(sample.executions, examine, claim.job.settings.concurrency)) as results: async for item, review in results: - await progress( - "Reading executions", - Coverage(eligible=sample.eligible, selected=len(sample.executions), screened=next(completed)), - review, - ) + screened += 1 + done: Final = item.execution.id + await report(lambda current: tuple(r for r in current if r.execution_id != done), review) yield item