test(lens): cover in-flight runs across progress, old workers and terminal states

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

View file

@ -1,6 +1,6 @@
from datetime import datetime, timedelta, timezone
from functools import reduce
from typing import Final
from typing import Final, Literal
import pytest
@ -12,11 +12,13 @@ from litellm.proxy.lens.models import (
Evidence,
Execution,
FindingDraft,
InFlight,
IssueBrief,
Job,
Lens,
LensSettings,
MetadataFilter,
Progress,
Review,
Sample,
Scope,
@ -26,13 +28,17 @@ from litellm.proxy.lens.models import (
from litellm.proxy.lens.state import (
add_review,
add_step,
apply_progress,
can_access,
cancel_job,
claim_job,
current_job,
end_job,
merge_finding,
next_scan_start,
queue_job,
renew_budget,
replace_job,
reviews_after,
summarized,
)
@ -428,3 +434,39 @@ def test_review_polling_after_the_window_moved_on_returns_what_is_still_kept() -
f"run-{MAX_REVIEWS + 8}",
f"run-{MAX_REVIEWS + 9}",
]
def in_flight(execution: str) -> InFlight:
return InFlight(execution_id=execution, trace_id="t", agent="support", started_at=NOW)
def reading_job() -> Job:
running: Final = claim_job(queue_job(lens(), NOW, "job"), worker(), NOW).jobs[0]
return apply_progress(running, Progress(stage=running.stage, reading=(in_flight("a"), in_flight("b"))), NOW)
def test_progress_replaces_the_in_flight_runs_and_old_workers_leave_them_alone() -> None:
job: Final = reading_job()
assert [r.execution_id for r in job.reading] == ["a", "b"]
finished: Final = apply_progress(job, Progress(stage=job.stage, review=review(0), reading=(in_flight("b"),)), NOW)
assert [r.execution_id for r in finished.reading] == ["b"]
assert finished.reviewed == 1
assert apply_progress(job, Progress(stage=job.stage, review=review(1)), NOW).reading == job.reading
assert apply_progress(job, Progress(stage=job.stage, reading=()), NOW).reading == ()
@pytest.mark.parametrize("status", ("completed", "failed", "cancelled"))
def test_finished_jobs_stop_showing_runs_in_flight(status: Literal["completed", "failed", "cancelled"]) -> None:
ended: Final = end_job(reading_job(), status, NOW)
assert ended.status == status
assert ended.finished_at == NOW
assert ended.reading == ()
def test_cancel_and_repeated_disconnects_clear_runs_in_flight() -> None:
reading: Final = replace_job(queue_job(lens(), NOW, "job"), reading_job())
cancelled: Final = cancel_job(reading, NOW).jobs[0]
assert (cancelled.status, cancelled.reading) == ("cancelled", ())
abandoned: Final = reading.model_copy(update={"jobs": (reading.jobs[0].model_copy(update={"attempts": 3}),)})
expired: Final = claim_job(abandoned, worker(), NOW + timedelta(minutes=10)).jobs[0]
assert (expired.status, expired.reading) == ("failed", ())