Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/shadow_eval_funnel.py: 53%
24 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1"""Pod-local queue of shadow-eval funnel increments, drained by the spend-update job.
3The shadow-eval success hook counts the sampled-traffic outcomes that never produce an
4attempt row (a lost sampling dice roll, an unjudgeable request shape, a concurrency
5shed), so a job's results can state what share of its eligible traffic the judged rows
6represent. Counters are advisory coverage stats: a pod dying loses at most one flush
7interval, and a failed flush drops its batch because a repeated increment is worse
8than an undercount (same call as the auto-router session rollup flush).
9"""
11from typing import TYPE_CHECKING, Final, Literal
13from litellm._logging import verbose_proxy_logger
15if TYPE_CHECKING: 15 ↛ 16line 15 didn't jump to line 16 because the condition on line 15 was never true
16 from litellm.proxy.utils import PrismaClient
18ShadowEvalFunnelStage = Literal["not_sampled", "unjudgeable", "shed", "withheld"]
20FUNNEL_STAGES: Final[tuple[ShadowEvalFunnelStage, ...]] = ("not_sampled", "unjudgeable", "shed", "withheld")
22_pending: dict[str, dict[ShadowEvalFunnelStage, int]] = {} # mutable-ok: module-level queue, single event loop
24_FUNNEL_PLACEHOLDERS: Final = ", ".join(f"${n + 2}" for n in range(len(FUNNEL_STAGES)))
26_UPSERT_FUNNEL_SQL: Final = f"""
27INSERT INTO "LiteLLM_ShadowEvalFunnel" (job_id, {", ".join(FUNNEL_STAGES)})
28VALUES ($1, {_FUNNEL_PLACEHOLDERS})
29ON CONFLICT (job_id) DO UPDATE SET
30 {", ".join(f'{stage} = "LiteLLM_ShadowEvalFunnel".{stage} + EXCLUDED.{stage}' for stage in FUNNEL_STAGES)}
31"""
34def pending_shadow_eval_funnel_events() -> int:
35 """Queue census for the drain triggers: entries not yet flushed, so a funnel-only
36 batch still wakes the spend job that would otherwise skip an empty-queue run."""
37 return sum(sum(counters.values()) for counters in _pending.values())
40def record_shadow_eval_funnel_event(job_id: str, stage: ShadowEvalFunnelStage) -> None:
41 """Count one skipped request for one job leg; synchronous so the hook's read-modify-
42 write cannot interleave with the flush's snapshot on the shared event loop."""
43 counters: Final = _pending.setdefault(job_id, dict.fromkeys(FUNNEL_STAGES, 0))
44 counters[stage] += 1
47async def flush_shadow_eval_funnel(prisma_client: "PrismaClient") -> None:
48 if not _pending: 48 ↛ 50line 48 didn't jump to line 50 because the condition on line 48 was always true
49 return
50 batch: Final = dict(_pending) # mutable-ok: snapshot drained from the queue
51 _pending.clear()
52 for job_id, counters in batch.items():
53 try:
54 await prisma_client.db.execute_raw(
55 _UPSERT_FUNNEL_SQL,
56 job_id,
57 *(counters[stage] for stage in FUNNEL_STAGES),
58 )
59 except Exception as flush_err: # noqa: BLE001 # drop this leg's batch: a repeated increment is worse than an undercount
60 verbose_proxy_logger.error(
61 "Spend tracking - shadow eval funnel flush failed for job %s, %s dropped: %s",
62 job_id,
63 counters,
64 flush_err,
65 )