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

1"""Pod-local queue of shadow-eval funnel increments, drained by the spend-update job. 

2 

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

10 

11from typing import TYPE_CHECKING, Final, Literal 

12 

13from litellm._logging import verbose_proxy_logger 

14 

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 

17 

18ShadowEvalFunnelStage = Literal["not_sampled", "unjudgeable", "shed", "withheld"] 

19 

20FUNNEL_STAGES: Final[tuple[ShadowEvalFunnelStage, ...]] = ("not_sampled", "unjudgeable", "shed", "withheld") 

21 

22_pending: dict[str, dict[ShadowEvalFunnelStage, int]] = {} # mutable-ok: module-level queue, single event loop 

23 

24_FUNNEL_PLACEHOLDERS: Final = ", ".join(f"${n + 2}" for n in range(len(FUNNEL_STAGES))) 

25 

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

32 

33 

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()) 

38 

39 

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 

45 

46 

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 )