ó
    °"³jå)  ã                  óZ  • % S r SSKJr  SSKJrJr  SSKJr  SSK	J
r
Jr  SSKJr  SSKJr  SSKrSS	KJrJr  SS
KJr  SSKJr  SrSS jrSS jr " S S\5      r\ " S S5      5       r " S S\5      rSrS\S'    \S S j5       r \S!S j5       r! S"     S#S jjr"SS.     S$S jjr#g)%zEConcurrency limiting infrastructure with OpenTelemetry observability.é    )Úannotations)ÚABCÚabstractmethod)ÚAsyncGenerator)ÚAbstractAsyncContextManagerÚasynccontextmanager)Ú	dataclass)Ú	TypeAliasN)ÚTracerÚ
get_tracer)ÚSelfé   ©Ú	UserError)ÚAbstractConcurrencyLimiterÚConcurrencyLimiterÚConcurrencyLimitÚAnyConcurrencyLimitc                ó.   • U S:  a  [        SU  S35      eg )Nr   zmax_running must be >= 1, got z'. Use None for no concurrency limiting.r   )Úmax_runnings    ÚT/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_ai/concurrency.pyÚ_validate_max_runningr      s$   € Ø�QƒÜÐ8¸¸ÐElÐmÓnÐnð ó    c                ó6   • U b  U S:  a  [        SU  S35      eg g )Nr   zmax_queued must be >= 0, got z. Use None for unlimited queue.r   )Ú
max_queueds    r   Ú_validate_max_queuedr      s,   € ØÑ *¨q£.ÜÐ7¸
°|ÐCbÐcÓdÐdð #1Ðr   c                  ó@   • \ rS rSrSr\SS j5       r\SS j5       rSrg)	r   é#   a‡  Abstract base class for concurrency limiters.

Subclass this to create custom concurrency limiters
(e.g., Redis-backed distributed limiters).

Example:
```python
from pydantic_ai.concurrency import AbstractConcurrencyLimiter


class RedisConcurrencyLimiter(AbstractConcurrencyLimiter):
    def __init__(self, redis_client, key: str, max_running: int):
        self._redis = redis_client
        self._key = key
        self._max_running = max_running

    async def acquire(self, source: str) -> None:
        # Implement Redis-based distributed locking
        ...

    def release(self) -> None:
        # Release the Redis lock
        ...
```
c              ƒ  ó   #   • g7f)znAcquire a slot, waiting if necessary.

Args:
    source: Identifier for observability (e.g., 'model:gpt-4o').
N© )ÚselfÚsources     r   ÚacquireÚ"AbstractConcurrencyLimiter.acquire>   s
   é € ð 	ùs   ‚c                ó   • g©zRelease a slot.Nr    ©r!   s    r   ÚreleaseÚ"AbstractConcurrencyLimiter.releaseG   s   € ð 	r   r    N©r"   ÚstrÚreturnÚNone©r,   r-   )	Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r   r#   r(   Ú__static_attributes__r    r   r   r   r   #   s/   † ñð4 óó ðð óó ór   r   c                  ó<   • \ rS rSr% SrS\S'   SrS\S'   S
S jrS	rg)r   éM   aD  Configuration for concurrency limiting with optional backpressure.

Args:
    max_running: Maximum number of concurrent operations allowed. Must be >= 1.
    max_queued: Maximum number of operations waiting in the queue. Must be >= 0.
        If None, the queue is unlimited. If exceeded, raises `ConcurrencyLimitExceeded`.
Úintr   Nú
int | Noner   c                óX   • [        U R                  5        [        U R                  5        g )N)r   r   r   r   r'   s    r   Ú__post_init__ÚConcurrencyLimit.__post_init__Z   s   € Ü˜d×.Ñ.Ô/Ü˜TŸ_™_Õ-r   r    r.   )	r/   r0   r1   r2   r3   Ú__annotations__r   r:   r4   r    r   r   r   r   M   s   ‡ ñð ÓØ!€J�
Ó!÷.r   r   c                  óê   • \ rS rSrSrSSSS.       SS jjr\SSS.       SS jj5       r\SS j5       r	\SS	 j5       r
\SS
 j5       r\SS j5       r\SS j5       rSS jrSS jrSS jrSrg)r   é_   zÿA concurrency limiter that tracks waiting operations for observability.

This class wraps an anyio.CapacityLimiter and tracks the number of waiting operations.
When an operation has to wait to acquire a slot, a span is created for
observability purposes.
N)r   ÚnameÚtracerc               óÌ   • [        U5        [        U5        [        R                  " U5      U l        X l        X0l        X@l        [        R                  " 5       U l	        SU l
        g)aÜ  Initialize the ConcurrencyLimiter.

Args:
    max_running: Maximum number of concurrent operations. Must be >= 1.
    max_queued: Maximum queue depth before raising ConcurrencyLimitExceeded. Must be >= 0.
    name: Optional name for this limiter, used for observability when sharing
        a limiter across multiple models or agents.
    tracer: OpenTelemetry tracer for span creation.

Raises:
    UserError: If `max_running` is less than 1, or `max_queued` is less than 0.
r   N)r   r   ÚanyioÚCapacityLimiterÚ_limiterÚ_max_queuedÚ_nameÚ_tracerÚLockÚ_queue_lockÚ_waiting_count)r!   r   r   r?   r@   s        r   Ú__init__ÚConcurrencyLimiter.__init__g   sL   € ô( 	˜kÔ*Ü˜ZÔ(Ü×-Ò-¨kÓ:ˆŒØ%ÔØŒ
ØŒä Ÿ:š:›<ˆÔØˆÕr   )r?   r@   c               ót   • [        U[        5      (       a  U " XUS9$ U " UR                  UR                  UUS9$ )aC  Create a ConcurrencyLimiter from a ConcurrencyLimit configuration.

Args:
    limit: Either an int for simple limiting or a ConcurrencyLimit for full config.
    name: Optional name for this limiter, used for observability.
    tracer: OpenTelemetry tracer for span creation.

Returns:
    A configured ConcurrencyLimiter.
)r   r?   r@   )r   r   r?   r@   )Ú
isinstancer7   r   r   )ÚclsÚlimitr?   r@   s       r   Ú
from_limitÚConcurrencyLimiter.from_limit…   sC   € ô$ �eœS×!Ñ!Ù 5¸FÑCÐCáØ!×-Ñ-Ø ×+Ñ+ØØñ	ð r   c                ó   • U R                   $ )z&Name of the limiter for observability.)rF   r'   s    r   r?   ÚConcurrencyLimiter.name¡   s   € ð �z‰zÐr   c                ó   • U R                   $ )z9Number of operations currently waiting to acquire a slot.)rJ   r'   s    r   Úwaiting_countÚ ConcurrencyLimiter.waiting_count¦   s   € ð ×"Ñ"Ð"r   c                óJ   • U R                   R                  5       R                  $ )z'Number of operations currently running.)rD   Ú
statisticsÚborrowed_tokensr'   s    r   Úrunning_countÚ ConcurrencyLimiter.running_count«   s   € ð �}‰}×'Ñ'Ó)×9Ñ9Ð9r   c                ó@   • [        U R                  R                  5      $ )zNumber of slots available.)r7   rD   Úavailable_tokensr'   s    r   Úavailable_countÚ"ConcurrencyLimiter.available_count°   s   € ô �4—=‘=×1Ñ1Ó2Ð2r   c                ó@   • [        U R                  R                  5      $ )z&Maximum concurrent operations allowed.)r7   rD   Útotal_tokensr'   s    r   r   ÚConcurrencyLimiter.max_runningµ   s   € ô �4—=‘=×-Ñ-Ó.Ð.r   c                óJ   • U R                   b  U R                   $ [        S5      $ )z9Get the tracer, falling back to global tracer if not set.zpydantic-ai)rG   r   r'   s    r   Ú_get_tracerÚConcurrencyLimiter._get_tracerº   s!   € à�<‰<Ñ#Ø—<‘<ÐÜ˜-Ó(Ð(r   c              ƒ  ó:  #   • SSK Jn   U R                  R                  5         g! [        R
                   a     Of = fU R                   ISh  v•N    U R                  bj  U R                  U R                  :¼  aP  U R                  =(       d    UnU" SU R                  S-    SU R                   S3U(       a  SU 3-   5      eS-   5      eU =R                  S-  sl        SSS5      ISh  v•N    O! , ISh  v•N  (       d  f       O= f U R                  5       nU R                  =(       d    UnUU R                  [        U R                  R                  5      S	.nU R                  b  U R                  US
'   U R                  b  U R                  US'   SU S3nUR                  XeS9   U R                  R                  5       I Sh  v•N    SSS5        O! , (       d  f       O= fU =R                  S-  sl        g! U =R                  S-  sl        f = f7f)z¤Acquire a slot, creating a span if waiting is required.

Args:
    source: Identifier for the source of this acquisition (e.g., 'agent:my-agent' or 'model:gpt-4').
r   )ÚConcurrencyLimitExceededNzConcurrency queue depth (z) exceeds max_queued (Ú)z for Ú )r"   rV   r   Úlimiter_namer   zwaiting for z concurrency)Ú
attributes)Ú
exceptionsrh   rD   Úacquire_nowaitrB   Ú
WouldBlockrI   rE   rJ   rF   re   r7   rb   Ústart_as_current_spanr#   )r!   r"   rh   Údisplay_namer@   rl   Ú	span_names          r   r#   ÚConcurrencyLimiter.acquireÀ   sÛ  é € õ 	9ð	Ø�M‰M×(Ñ(Ô*ØøÜ×Ñó 	Ùð	úð ×#×#Ô#Ø×ÑÑ+°×0CÑ0CÀt×GWÑGWÓ0Wà#Ÿz™z×3¨V�Ù.Ø/°×0CÑ0CÀaÑ0GÐ/HÐH^Ð_c×_oÑ_oÐ^pÐpqÐrÞ1=˜˜|˜nÐ-ñGóð àCEñGóð ð
 ×Ò 1Ñ$Õ÷ $×#×#×#×#Ð#úð	%à×%Ñ%Ó'ˆFØŸ:™:×/¨ˆLà Ø!%×!4Ñ!4Ü" 4§=¡=×#=Ñ#=Ó>ñ0ˆJð
 �z‰zÑ%Ø-1¯Z©Z�
˜>Ñ*Ø×ÑÑ+Ø+/×+;Ñ+;�
˜<Ñ(ð ' | n°LÐAˆIØ×-Ñ-¨iÐ-ÒOØ—m‘m×+Ñ+Ó-×-Ñ-÷ P×OÖOúð ×Ò 1Ñ$ÖøˆD×Ò 1Ñ$Öüsš   ‚HŠ% ¤H¥<¹H»<¼HÁAÁHÁBC3Ã!HÃ,C/Ã-HÃ3D
Ã9C<Ã:D
ÄHÄBH Æ-GÇGÇGÇ	H Ç
G(Ç$H Ç+HÈHÈHc                ó8   • U R                   R                  5         gr&   )rD   r(   r'   s    r   r(   ÚConcurrencyLimiter.releaseõ   s   € à�‰×ÑÕr   )rD   rE   rF   rI   rG   rJ   )r   r7   r   r8   r?   ú
str | Noner@   úTracer | None)rP   zint | ConcurrencyLimitr?   rv   r@   rw   r,   r   )r,   rv   )r,   r7   )r,   r   r*   r.   )r/   r0   r1   r2   r3   rK   ÚclassmethodrQ   Úpropertyr?   rV   r[   r_   r   re   r#   r(   r4   r    r   r   r   r   _   só   † ñð "&ØØ $ñ àð ð ð	 ð
 ð ð õ ð< ð
  Ø $ñà%ðð ð	ð
 ðð 
ôó ðð6 óó ðð ó#ó ð#ð ó:ó ð:ð ó3ó ð3ð ó/ó ð/ô)ô3%÷j r   r   z:int | ConcurrencyLimit | AbstractConcurrencyLimiter | Noner
   r   c                ó   #   • S7v •  g7f)zA no-op async context manager.Nr    r    r   r   Ú_null_contextr{     s
   é € õ 
ùs   ‚	c               ó˜   #   • U R                  U5      I Sh  v•N    S7v •  U R                  5         g N! U R                  5         f = f7f)zKContext manager that acquires and releases a limiter with the given source.N)r#   r(   ©Úlimiterr"   s     r   Ú_limiter_contextr     s=   é € ð �/‰/˜&Ó
!×!Ð!ðÜà�‰Õñ	 "øð 	�‰Õüs"   ‚A
—3˜A
�5 ¢A
µAÁA
c                ó2   • U c
  [        5       $ [        X5      $ )a9  Get an async context manager for the concurrency limiter.

If limiter is None, returns a no-op context manager.

Args:
    limiter: The AbstractConcurrencyLimiter or None.
    source: Identifier for the source of this acquisition (e.g., 'agent:my-agent' or 'model:gpt-4').

Returns:
    An async context manager.
)r{   r   r}   s     r   Úget_concurrency_contextr�     s   € ð �Ü‹ÐÜ˜GÓ,Ð,r   ©r?   c               ó^   • U c  g[        U [        5      (       a  U $ [        R                  XS9$ )a  Normalize a concurrency limit configuration to an AbstractConcurrencyLimiter.

Args:
    limit: The concurrency limit configuration.
    name: Optional name for the limiter if one is created.

Returns:
    An AbstractConcurrencyLimiter if limit is not None, otherwise None.
Nr‚   )rN   r   r   rQ   )rP   r?   s     r   Únormalize_to_limiterr„   )  s3   € ð �}ØÜ	�EÔ5×	6Ñ	6Øˆä!×,Ñ,¨UÐ,Ð>Ð>r   )r   r7   r,   r-   )r   r8   r,   r-   )r,   úAsyncGenerator[None])r~   r   r"   r+   r,   r…   )Úunnamed)r~   ú!AbstractConcurrencyLimiter | Noner"   r+   r,   z!AbstractAsyncContextManager[None])rP   r   r?   rv   r,   r‡   )$r3   Ú
__future__r   Ú_annotationsÚabcr   r   Úcollections.abcr   Ú
contextlibr   r   Údataclassesr	   Útypingr
   rB   Úopentelemetry.tracer   r   Útyping_extensionsr   rm   r   Ú__all__r   r   r   r   r   r   r<   r{   r   r�   r„   r    r   r   Ú<module>r’      sô   ðÚ Kå 2ç #Ý *ß GÝ !Ý ã ß 2Ý "å !ð€ôoô
eô
' ô 'ðT ÷.ð .ó ð.ô"X Ð3ô X ðv "^Ð �YÓ ]ðð ó
ó ð
ð
 óó ðð ð-Ø.ð-àð-ð 'õ-ð. ñ?Øð?ð ð?ð 'ö	?r   