feat(lens): report each run as in flight when its review starts

This commit is contained in:
Ishaan Jaff 2026-10-03 18:01:39 -07:00
parent 45909929ff
commit 1d18be780e
No known key found for this signature in database

View file

@ -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