ó
    ±"³jc  ã                  óú   • S r SSKJr  SSKJr  SSKJr  SSKJr  SSK	J
r
  SSKJr  SS	KJrJrJrJrJrJr  SS
KJrJr  SSKJr  SSKJr  SSKJrJrJrJr  SSK J!r!  \" SS9 " S S\!5      5       r"      SS jr#g)z(Concurrency limiting wrapper for models.é    )Úannotations)ÚAsyncGenerator)Úasynccontextmanager)Ú	dataclass)ÚAnyé   )Ú
RunContext)ÚAbstractConcurrencyLimiterÚAnyConcurrencyLimitÚConcurrencyLimitÚConcurrencyLimiterÚget_concurrency_contextÚnormalize_to_limiter)ÚModelMessageÚModelResponse)ÚModelSettings)ÚRequestUsageé   )ÚKnownModelNameÚModelÚModelRequestParametersÚStreamedResponse)ÚWrapperModelF)Úinitc                  ó¤   ^ • \ rS rSr% SrS\S'       S
U 4S jjr        SS jr        SS jr\	 S         SS jj5       r
S	rU =r$ )ÚConcurrencyLimitedModelé   a¤  A model wrapper that limits concurrent requests to the underlying model.

This wrapper applies concurrency limiting at the model level, ensuring that
the number of concurrent requests to the model does not exceed the configured
limit. This is useful for:

- Respecting API rate limits
- Managing resource usage
- Sharing a concurrency pool across multiple models

Example usage:
```python
from pydantic_ai import Agent
from pydantic_ai.models.concurrency import ConcurrencyLimitedModel

# Limit to 5 concurrent requests
model = ConcurrencyLimitedModel('openai:gpt-4o', limiter=5)
agent = Agent(model)

# Or share a limiter across multiple models
from pydantic_ai import ConcurrencyLimiter  # noqa E402

shared_limiter = ConcurrencyLimiter(max_running=10, name='openai-pool')
model1 = ConcurrencyLimitedModel('openai:gpt-4o', limiter=shared_limiter)
model2 = ConcurrencyLimitedModel('openai:gpt-4o-mini', limiter=shared_limiter)
```
r
   Ú_limiterc                ó’   >• [         TU ]  U5        [        U[        5      (       a  X l        g[
        R                  " U5      U l        g)a°  Initialize the ConcurrencyLimitedModel.

Args:
    wrapped: The model to wrap, either a Model instance or a known model name.
    limiter: The concurrency limit configuration. Can be:
        - An `int`: Simple limit on concurrent operations (unlimited queue).
        - A `ConcurrencyLimit`: Full configuration with optional backpressure.
        - An `AbstractConcurrencyLimiter`: A pre-created limiter for sharing across models.
N)ÚsuperÚ__init__Ú
isinstancer
   r   r   Ú
from_limit)ÚselfÚwrappedÚlimiterÚ	__class__s      €Ú[/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_ai/models/concurrency.pyr!   Ú ConcurrencyLimitedModel.__init__:   s7   ø€ ô 	‰Ñ˜Ô!Ü�gÔ9×:Ñ:Ø#�Mä.×9Ò9¸'ÓBˆD�Mó    c              ƒ  ó  #   • [        U R                  SU R                   35       ISh  v•N   U R                  R	                  XU5      I Sh  v•N sSSS5      ISh  v•N   $  N9 N N	! , ISh  v•N  (       d  f       g= f7f)z6Make a request to the model with concurrency limiting.úmodel:N)r   r   Ú
model_namer%   Úrequest©r$   ÚmessagesÚmodel_settingsÚmodel_request_parameterss       r(   r.   ÚConcurrencyLimitedModel.requestN   sZ   é € ô +¨4¯=©=¸FÀ4Ç?Á?ÐBSÐ:T×UÕUØŸ™×-Ñ-¨hÐH`Óa×a÷ V×UÓUÙa÷ V×U×UÐUüóV   ‚(BªA$«B® A*ÁA&ÁA*ÁBÁA(ÁBÁ&A*Á(BÁ*BÁ0A3Á1BÁ=Bc              ƒ  ó  #   • [        U R                  SU R                   35       ISh  v•N   U R                  R	                  XU5      I Sh  v•N sSSS5      ISh  v•N   $  N9 N N	! , ISh  v•N  (       d  f       g= f7f)z'Count tokens with concurrency limiting.r,   N)r   r   r-   r%   Úcount_tokensr/   s       r(   r6   Ú$ConcurrencyLimitedModel.count_tokensX   sZ   é € ô +¨4¯=©=¸FÀ4Ç?Á?ÐBSÐ:T×UÕUØŸ™×2Ñ2°8ÐMeÓf×f÷ V×UÓUÙf÷ V×U×UÐUür4   c               óp  #   • [        U R                  SU R                   35       ISh  v•N   U R                  R	                  XX45       ISh  v•N nU7v •  SSS5      ISh  v•N   SSS5      ISh  v•N   g NO N, N! , ISh  v•N  (       d  f       N.= f N%! , ISh  v•N  (       d  f       g= f7f)z@Make a streaming request to the model with concurrency limiting.r,   N)r   r   r-   r%   Úrequest_stream)r$   r0   r1   r2   Úrun_contextÚresponse_streams         r(   r9   Ú&ConcurrencyLimitedModel.request_streamb   s„   é € ô +¨4¯=©=¸FÀ4Ç?Á?ÐBSÐ:T×UÕUØ—|‘|×2Ñ2ØÐ*B÷ô à Ø%Ó%÷÷ ÷ V×UÒU÷÷ ÷ ò ú÷ V×U×UÐUüsŒ   ‚(B6ªA:«B6®!BÁA<ÁBÁB ÁBÁ$A>Á%BÁ)B6Á4BÁ5B6Á<BÁ>BÂ B	ÂB	ÂB	ÂBÂB6ÂB3Â"B%Â#B3Â/B6)r   )r%   úModel | KnownModelNamer&   z3int | ConcurrencyLimit | AbstractConcurrencyLimiter)r0   úlist[ModelMessage]r1   úModelSettings | Noner2   r   Úreturnr   )r0   r>   r1   r?   r2   r   r@   r   )N)
r0   r>   r1   r?   r2   r   r:   zRunContext[Any] | Noner@   z AsyncGenerator[StreamedResponse])Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú__annotations__r!   r.   r6   r   r9   Ú__static_attributes__Ú__classcell__)r'   s   @r(   r   r      sÔ   ø‡ ñð8 )Ó(ðCà'ðCð E÷Cð(bà$ðbð -ðbð #9ð	bð
 
ôbðgà$ðgð -ðgð #9ð	gð
 
ôgð ð /3ð&à$ð&ð -ð&ð #9ð	&ð
 ,ð&ð 
*ô&ó ö&r*   r   c                ó~   • [        U5      nUc%  SSKJn  [        U [        5      (       a  U" U 5      $ U $ [        X5      $ )a   Wrap a model with concurrency limiting.

This is a convenience function to wrap a model with concurrency limiting.
If the limiter is None, the model is returned unchanged.

Args:
    model: The model to wrap.
    limiter: The concurrency limit configuration.

Returns:
    The wrapped model with concurrency limiting, or the original model if limiter is None.

Example:
```python
from pydantic_ai.models.concurrency import limit_model_concurrency

model = limit_model_concurrency('openai:gpt-4o', limiter=5)
```
r   )Úinfer_model)r   Ú rJ   r"   Ústrr   )Úmodelr&   Únormalized_limiterrJ   s       r(   Úlimit_model_concurrencyrO   r   s?   € ô. .¨gÓ6ÐØÑ!Ý!ä%/°´s×%;Ñ%;‰{˜5Ó!ÐFÀÐFÜ" 5Ó=Ð=r*   N)rM   r=   r&   r   r@   r   )$rE   Ú
__future__r   Úcollections.abcr   Ú
contextlibr   Údataclassesr   Útypingr   Ú_run_contextr	   Úconcurrencyr
   r   r   r   r   r   r0   r   r   Úsettingsr   Úusager   rK   r   r   r   r   Úwrapperr   r   rO   © r*   r(   Ú<module>r[      s{   ðÙ .å "å *Ý *Ý !Ý å %÷÷ ÷ 3Ý $Ý  ß MÓ MÝ !ñ �ÑôT&˜ló T&ó ðT&ðn>Ø!ð>à ð>ð õ>r*   