ó
    lQj‰1  ã                  óÒ  • S r SSKJr  SSKrSSKJs  Jr  SSK	r	SSK
r
SSKrSSKrSSKJr  SSKJrJrJrJrJrJrJr  SSKJr  SSKJrJrJrJrJrJrJ r J!r!  SSK"J#r#  \" S	5      r$\%" 5       r&\&4     SS
 jjr' " S S5      r(          SS jr) " S S\\$   5      r*\*r+S r,    SS jr-SS.       SS jjr.SS jr/  SS jr0  SS jr1\Rd                  " SS9S 5       r3g) zzAdapted.

Original source:
https://github.com/maxfischer2781/asyncstdlib/blob/master/asyncstdlib/itertools.py
MIT License
é    )ÚannotationsN)Údeque)ÚAsyncGeneratorÚAsyncIterableÚAsyncIteratorÚ	AwaitableÚ	CoroutineÚIterableÚIterator)ÚAbstractAsyncContextManager)ÚAnyÚCallableÚGenericÚOptionalÚTypeVarÚUnionÚcastÚoverload)Úget_runtime_overridesÚTc                ó  ^ ^^•  [        [        [        [           /[        [           4   [        T 5      R                  5      mT[        L a  T" T 5      $ SUUU 4S jjnU" 5       $ ! [         a    [        T < S35      ef = f)a  Pure-Python implementation of anext() for testing purposes.

Closely matches the builtin anext() C implementation.
Can be used to compare the built-in implementation of the inner
coroutines machinery to C-implementation of __anext__() and send()
or throw() on the returned generator.
z is not an async iteratorc               “  óV   >#   •  T " T5      I S h  v•N $  N! [          a    Ts $ f = f7f©N)ÚStopAsyncIteration)Ú	__anext__ÚdefaultÚiterators   €€€ÚW/home/mande/repo/quber/.venv/lib/python3.13/site-packages/langsmith/_internal/_aiter.pyÚ
anext_implÚpy_anext.<locals>.anext_implA   s1   øé € ð	ñ # 8Ó,×,Ð,Ñ,øÜ!ó 	ØŠNð	üs(   ƒ)… �‘ ”)• —&£)¥&¦))ÚreturnúUnion[T, Any])
r   r   r   r   r   Útyper   ÚAttributeErrorÚ	TypeErrorÚ_no_default)r   r   r   r   s   `` @r   Úpy_anextr'   -   s‡   ú€ ðBÜÜ”m¤AÑ&Ð'¬´1©Ð5Ñ6¼¸X»×8PÑ8Pó
ˆ	ð ”+ÒÙ˜Ó"Ð"÷	ñ 	ñ ‹<Ðøô# ó BÜ˜8™,Ð&?Ð@ÓAÐAðBús   …?A& Á&B c                  ó,   • \ rS rSrSrSS jrSS jrSrg)	ÚNoLockéO   z@Dummy lock that provides the proper interface but no protection.c              ƒ  ó   #   • g 7fr   © ©Úselfs    r   Ú
__aenter__ÚNoLock.__aenter__R   s   é € Øùó   ‚c              ƒ  ó   #   • g7f©NFr,   ©r.   Úexc_typeÚexc_valÚexc_tbs       r   Ú	__aexit__ÚNoLock.__aexit__U   s   é € Øùr1   r,   N©r!   ÚNone©r5   r   r6   r   r7   r   r!   Úbool)Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r/   r8   Ú__static_attributes__r,   ó    r   r)   r)   O   s   † ÙJô÷rD   r)   c           	    óô  #   •   U(       dh  U ISh  v•N   U(       a   SSS5      ISh  v•N   M-   U R                  5       I Sh  v•N nU H  nUR                  U5        M     SSS5      ISh  v•N   UR                  5       7v •  M„   Nv N` NG! [         a     SSS5      ISh  v•N    O f = f N@! , ISh  v•N  (       d  f       NU= fU ISh  v•N    [	        U5       H  u  peXQL d  M  UR                  U5          O   U(       d*  [        U S5      (       a  U R                  5       I Sh  v•N    SSS5      ISh  v•N    g! , ISh  v•N  (       d  f       g= f! U ISh  v•N    [	        U5       H  u  peXQL d  M  UR                  U5          O   U(       d*  [        U S5      (       a  U R                  5       I Sh  v•N    SSS5      ISh  v•N    f ! , ISh  v•N  (       d  f       f = f= f7f)zIterate over :py:func:`~.tee`.NÚaclose)r   Úappendr   ÚpopleftÚ	enumerateÚpopÚhasattrrF   )r   ÚbufferÚpeersÚlockÚitemÚpeer_bufferÚidxs          r   Útee_peerrR   Y   sX  é € ð(ØÞßš4ö Ø ÷	  Ÿ4™4ð
5Ø%-×%7Ñ%7Ó%9×9˜ó ,1˜KØ'×.Ñ.¨tÖ4ñ ,1÷  Ÿ4ð" —.‘.Ó"Ó"ñ' ãñ  :øÜ-ó Ø÷  Ÿ4™4ðú÷  Ÿ4Ÿ4˜4ú÷& “4ä$-¨eÖ$4Ñ �ØÔ(Ø—I‘I˜c”NÙñ %5ö
 œW X¨x×8Ñ8Ø—o‘oÓ'×'Ñ'÷ —4—4—4—4�4û—4“4ä$-¨eÖ$4Ñ �ØÔ(Ø—I‘I˜c”NÙñ %5ö
 œW X¨x×8Ñ8Ø—o‘oÓ'×'Ñ'÷ —4—4—4—4�4ÿsU  ‚G8„E  ’B	“E  –	B2ŸE  ªB«E  ²BÁBÁBÁ
B2Á$E  Á/B0Á0E  ÂE  ÂBÂ
B-ÂB2ÂE  Â%B(Â&E  Â,B-Â-B2Â0E  Â2C	Â8B;Â9C	ÃE  ÃG8ÃCÃG8ÃEÃ.A EÄ.D1Ä/EÄ4G8Ä?EÅ G8ÅEÅEÅEÅG8Å G5Å'E*
Å(G5Å,GÆA GÇG
ÇGÇ	G5ÇGÇG5ÇG2Ç!G$Ç"G2Ç.G5Ç5G8c                  ó¦   • \ rS rSrSr SSS.     SS jjjrSS jr\SS j5       r\SS j5       r    SS	 jrSS
 jr	SS jr
SS jrSS jrSrg)ÚTeeéƒ   aÍ  Create ``n`` separate asynchronous iterators over ``iterable``.

This splits a single ``iterable`` into multiple iterators, each providing
the same items in the same order.
All child iterators may advance separately but pare the same items
from ``iterable`` -- when the most advanced iterator retrieves an item,
it is buffered until the least advanced iterator has yielded it as well.
A ``tee`` works lazily and can handle an infinite ``iterable``, provided
that all iterators advance.

```python
async def derivative(sensor_data):
    previous, current = a.tee(sensor_data, n=2)
    await a.anext(previous)  # advance one iterator
    return a.map(operator.sub, previous, current)
```

Unlike :py:func:`itertools.tee`, :py:func:`~.tee` returns a custom type instead
of a :py:class:`tuple`. Like a tuple, it can be indexed, iterated and unpacked
to get the child iterators. In addition, its :py:meth:`~.tee.aclose` method
immediately closes all children, and it can be used in an ``async with`` context
for the same effect.

If ``iterable`` is an iterator and read elsewhere, ``tee`` will *not*
provide these items. Also, ``tee`` must internally buffer each item until the
last iterator has yielded it; if the most and least advanced iterator differ
by most data, using a :py:class:`list` is more efficient (but not lazy).

If the underlying iterable is concurrency safe (``anext`` may be awaited
concurrently) the resulting iterators are concurrency safe as well. Otherwise,
the iterators are safe if there is only ever one single "most advanced" iterator.
To enforce sequential use of ``anext``, provide a ``lock``
- e.g. an :py:class:`asyncio.Lock` instance in an :py:mod:`asyncio` application -
and access is automatically synchronised.
N)rN   c               óØ   ^ ^• UR                  5       T l        [        U5       Vs/ s H  n[        5       PM     snT l        [        UU 4S jT R                   5       5      T l        g s  snf )Nc              3  ó~   >#   • U  H2  n[        TR                  UTR                  Tb  TO	[        5       S9v •  M4     g 7f)N)r   rL   rM   rN   )rR   Ú	_iteratorÚ_buffersr)   )Ú.0rL   rN   r.   s     €€r   Ú	<genexpr>ÚTee.__init__.<locals>.<genexpr>±   s>   øé € ð 
ò (�ô ØŸ™ØØ—m‘mØ!Ñ-‘T´6³8ö	ò (ùs   ƒ:=)Ú	__aiter__rX   Úranger   rY   ÚtupleÚ	_children)r.   ÚiterableÚnrN   Ú_s   `  ` r   Ú__init__ÚTee.__init__¨   sV   ù€ ð "×+Ñ+Ó-ˆŒÜ:?À¼(Ó(Cº(°Q¬®¹(Ñ(CˆŒÜõ 
ð Ÿ-š-ó
ó 
ˆ�ùò )Ds   ¥A'c                ó,   • [        U R                  5      $ r   )Úlenr`   r-   s    r   Ú__len__ÚTee.__len__»   s   € Ü�4—>‘>Ó"Ð"rD   c                ó   • g r   r,   ©r.   rO   s     r   Ú__getitem__ÚTee.__getitem__¾   s   € Ø:=rD   c                ó   • g r   r,   rk   s     r   rl   rm   Á   s   € ØHKrD   c                ó    • U R                   U   $ r   ©r`   rk   s     r   rl   rm   Ä   s   € ð �~‰~˜dÑ#Ð#rD   c              #  ó8   #   • U R                    S h  v•N   g  N7fr   rp   r-   s    r   Ú__iter__ÚTee.__iter__É   s   é € Ø—>‘>×!Ó!ùs   ‚’“c              ƒ  ó   #   • U $ 7fr   r,   r-   s    r   r/   ÚTee.__aenter__Ì   s
   é € Øˆùs   ‚c              ƒ  ó@   #   • U R                  5       I S h  v•N   g N7fr3   )rF   r4   s       r   r8   ÚTee.__aexit__Ï   s   é € Ø�k‰k‹m×ÐØñ 	ùs   ‚–—c              ƒ  óf   #   • U R                    H  nUR                  5       I S h  v•N   M     g  N	7fr   )r`   rF   )r.   Úchilds     r   rF   Ú
Tee.acloseÓ   s%   é € Ø—^”^ˆEØ—,‘,“.× Ò ò $Ù ùs   ‚#1¥/¦
1)rY   r`   rX   )é   )ra   úAsyncIterator[T]rb   ÚintrN   z*Optional[AbstractAsyncContextManager[Any]])r!   r}   )rO   r}   r!   r|   )rO   Úslicer!   ztuple[AsyncIterator[T], ...])rO   zUnion[int, slice]r!   z5Union[AsyncIterator[T], tuple[AsyncIterator[T], ...]])r!   zIterator[AsyncIterator[T]])r!   zTee[T]r<   r:   )r>   r?   r@   rA   rB   rd   rh   r   rl   rr   r/   r8   rF   rC   r,   rD   r   rT   rT   ƒ   s…   † ñ"ðN ð
ð
 <@ñ
à"ð
ð ð
ð
 9ö
ô&#ð Û=ó Ø=àÛKó ØKð$Ø%ð$à	>ô$ô
"ôô÷!rD   rT   c                óÞ   #   • U  Vs/ s H  oR                  5       PM     nn  [        R                  " S U 5       6 I Sh  v•N n[        U5      7v •  M3  s  snf  N! [         a     gf = f7f)zAsync version of zip.c              3  ó8   #   • U  H  n[        U5      v •  M     g 7fr   )r'   )rZ   r   s     r   r[   Úasync_zip.<locals>.<genexpr>â   s   é € Ð?²Y¨”(˜8×$Ð$²Yùs   ‚N)r]   ÚasyncioÚgatherr_   r   )Úasync_iterablesra   Ú	iteratorsÚitemss       r   Ú	async_zipr‡   Û   su   é € ñ 7FÓF²o¨(×#Ñ#Ö%±o€IÐFØ
ð	Ü!Ÿ.š.Ù?±YÓ?ð÷ ˆEô ˜“,Óñ ùò Gñøô "ó 	Ùð	üsD   ‚A-‡A A-¤A ÁAÁA ÁA-ÁA Á
A*Á'A-Á)A*Á*A-c                óÆ   • [        U S5      (       a  [        [        U 5      $ [        U S5      (       a  [        [        U R                  5       5      $  " S S5      nU" U 5      $ )Nr   r]   c                  ó*   • \ rS rSrSS jrS rS rSrg)Ú3ensure_async_iterator.<locals>.AsyncIteratorWrapperéò   c                ó$   • [        U5      U l        g r   )ÚiterrX   )r.   ra   s     r   rd   Ú<ensure_async_iterator.<locals>.AsyncIteratorWrapper.__init__ó   s   € Ü!% h£�•rD   c              “  ó^   #   •  [        U R                  5      $ ! [         a    [        ef = f7fr   )ÚnextrX   ÚStopIterationr   r-   s    r   r   Ú=ensure_async_iterator.<locals>.AsyncIteratorWrapper.__anext__ö   s-   é € ð-Ü §¡Ó/Ð/øÜ$ó -Ü,Ð,ð-üs   ‚-„ ˜-™*ª-c                ó   • U $ r   r,   r-   s    r   r]   Ú=ensure_async_iterator.<locals>.AsyncIteratorWrapper.__aiter__ü   s   € Ø�rD   )rX   N)ra   r
   )r>   r?   r@   rA   rd   r   r]   rC   r,   rD   r   ÚAsyncIteratorWrapperrŠ   ò   s   † ô0ò-õrD   r•   )rK   r   r   r]   )ra   r•   s     r   Úensure_async_iteratorr–   é   sX   € ô ˆx˜×%Ñ%Ü”M 8Ó,Ð,Ü	�˜;×	'Ñ	'Ü”M 8×#5Ñ#5Ó#7Ó8Ð8÷	ñ 	ñ $ HÓ-Ð-rD   )Ú_eager_consumption_timeoutc               óØ   ^ ^^^^• T S:X  a  U4S jnU" 5       $ [        [        R                  T b  [        R                  " T 5      O	[        5       5      mSU4S jjmUUU U4S jnU" 5       $ )a½  Process async generator with max parallelism.

Args:
    n: The number of tasks to run concurrently.
    generator: The async generator to process.
    _eager_consumption_timeout: If set, check for completed tasks after
        each iteration and yield their results. This can be used to
        consume the generator eagerly while still respecting the concurrency
        limit.

Yields:
    The processed items yielded by the async generator.
r   c                óL   >#   • T  S h  v•N n U I S h  v•N 7v •  M   N N
 g 7fr   r,   )rO   Ú	generators    €r   ÚconsumeÚ'aiter_with_concurrency.<locals>.consume  s%   øé € Ù'÷ !�dØ —jÕ ñ!Ù ñ (ùs(   ƒ$†"Š‹"Ž$” •	$ž" $¢$c              “  óž   >#   • T IS h  v•N   UI S h  v•N nX4sS S S 5      IS h  v•N   $  N" N N	! , IS h  v•N  (       d  f       g = f7fr   r,   )ÚixrO   ÚresÚ	semaphores      €r   Úprocess_itemÚ,aiter_with_concurrency.<locals>.process_item   s.   øé € ß’9Ø—*ˆCØ�9÷ —9“9Ù÷ —9—9�9üsE   ƒAŠ-‹AŽ3”/•3›A§1¨A¯3±A³A
¹<ºA
ÁAc                ó,  >#   • 0 n [        5       nSnT  S h  v•N nU(       a1  [        R                  " 5       n[        R                  " T" X#5      US9nO[        R                  " T" X#5      5      nXPU'   US-  nTS:”  a>   [        R
                  " U R                  5       TS9 H  nUI S h  v•N u  pxU7v •  X	 M     Tc  M¯  [        U 5      T:¼  d  MÀ  [        R                  " U R                  5       [        R                  S9I S h  v•N u  pšU	 H  nUR                  5       u  pxU7v •  X	 M     GM    GN NŠ! [        R                   a     N‘f = f NJ
 [        R
                  " U R                  5       5       H  nUI S h  v•N  u  p¨U7v •  M     g 7f)Nr   )Úcontexté   )Útimeout)Úreturn_when)Úasyncio_accepts_contextÚcontextvarsÚcopy_contextr‚   Úcreate_taskÚas_completedÚvaluesÚTimeoutErrorrg   ÚwaitÚFIRST_COMPLETEDÚresult)ÚtasksÚaccepts_contextrž   rO   r¤   ÚtaskÚ_futÚtask_idxrŸ   Údonerc   r—   rš   rb   r¡   s              €€€€r   Úprocess_generatorÚ1aiter_with_concurrency.<locals>.process_generator%  sg  øé € ØˆÜ1Ó3ˆØˆÙ#÷ 	(�$ÞÜ%×2Ò2Ó4�Ü×*Ò*©<¸Ó+AÈ7ÑS‘ä×*Ò*©<¸Ó+AÓB�Ø�"‰IØ�!‰GˆBØ)¨AÓ-ð	Ü '× 4Ò 4ØŸ™›Ø :ô!˜ð /3¯
™˜Ø!›	Ø!šOñ!ð ‹}¤ U£¨q¥Ü '§¢Ø—L‘L“N´×0GÑ0Gñ!÷ ‘�ó !�DØ$(§K¡K£M‘M�HØ“IØšô !ò/	(ñ )3øô ×+Ñ+ó Ùðúñð) $ô8 ×(Ò(¨¯©«Ö8ˆDØ—Z�Z‰FˆAØ�Iò 9ùs€   ƒF”E˜D4™EœA$FÂ+D9Â,D7Â-D9Â>FÃFÃ4FÄEÄ	+FÄ4EÄ7D9Ä9EÅFÅEÅFÅ-FÆFÆF)rž   r}   )r   r‚   Ú	Semaphorer)   )rb   rš   r—   r›   r¸   r¡   r    s   ```  @@r   Úaiter_with_concurrencyr»     s^   ü€ ð& 	ˆAƒvõ	!ñ ‹yÐÜÜ×Ñ°1±=œ7×,Ò,¨QÔ/ÄfÃhó€I÷÷
"ð "ñH ÓÐrD   c                ó†   •  [         R                  " U 5      R                  R                  S5      SL$ ! [         a     gf = f)z/Check if a callable accepts a context argument.r¤   NF)ÚinspectÚ	signatureÚ
parametersÚgetÚ
ValueError)Úcallables    r   r³   r³   L  s@   € ðÜ× Ò  Ó*×5Ñ5×9Ñ9¸)ÓDÈDÐPÐPøÜó Ùðús   ‚03 ³
A ¿A c             �  ó´   #   • [        5       nUR                  b#  UR                  " [        X/UQ70 UD6I Sh  v•N $ [        X/UQ70 UD6I Sh  v•N $  N N7f)až  Run ``func`` in a separate thread, inside ``ctx``.

``ctx`` is the :class:`~contextvars.Context` in which ``func`` is invoked.
Callers that want default isolation should pass
``contextvars.copy_context()``; callers with a specific Context
(e.g. :func:`trace`) pass it directly so subsequent reads from that
Context see the mutations.

Return a coroutine that can be awaited to get the eventual result of ``func``.
N)r   Úaio_to_threadÚ_default_aio_to_thread)ÚctxÚfuncÚargsÚkwargsÚ	overridess        r   rÄ   rÄ   U  sm   é € ô" &Ó'€IØ×ÑÑ*Ø×,Ò,Ü" Cð
Ø04ò
Ø8>ñ
÷ 
ð 	
ô (¨ÐC°DÒC¸FÑC×CÐCñ
ñ Dùs!   ‚6A¸A¹AÁAÁAÁAc             �  ó¶   #   • [         R                  " 5       n[        R                  " U R                  U/UQ70 UD6nUR                  SU5      I Sh  v•N $  N7f)z>Default implementation of aio_to_thread using run_in_executor.N)r‚   Úget_running_loopÚ	functoolsÚpartialÚrunÚrun_in_executor)rÆ   rÇ   rÈ   rÉ   ÚloopÚ	func_calls         r   rÅ   rÅ   n  sN   é € ô ×#Ò#Ó%€DÜ×!Ò! #§'¡'¨4ÐA°$ÒA¸&ÑA€IØ×%Ñ% d¨IÓ6×6Ð6Ñ6ùs   ‚AAÁAÁAr¥   )Úmaxsizec                 ó4   • [        [        R                  5      $ )zCCheck if the current asyncio event loop accepts a context argument.)r³   r‚   r«   r,   rD   r   r¨   r¨   {  s   € ô œ7×.Ñ.Ó/Ð/rD   )r   r|   r   r"   r!   zAwaitable[Union[T, None, Any]])
r   r|   rL   zdeque[T]rM   zlist[deque[T]]rN   z AbstractAsyncContextManager[Any]r!   úAsyncGenerator[T, None])ra   zUnion[Iterable, AsyncIterable]r!   r   )rb   zOptional[int]rš   z'AsyncIterator[Coroutine[None, None, T]]r—   Úfloatr!   rÕ   )rÂ   zCallable[..., Any]r!   r=   )rÆ   zcontextvars.Context)4rB   Ú
__future__r   ÚbuiltinsÚ@py_builtinsÚ_pytest.assertion.rewriteÚ	assertionÚrewriteÚ
@pytest_arr‚   r©   rÍ   r½   Úcollectionsr   Úcollections.abcr   r   r   r   r	   r
   r   Ú
contextlibr   Útypingr   r   r   r   r   r   r   r   Úlangsmith._runtime_overridesr   r   Úobjectr&   r'   r)   rR   rT   Úateer‡   r–   r»   r³   rÄ   rÅ   Ú	lru_cacher¨   r,   rD   r   Ú<module>ræ      sf  ðñõ #ç  „ ƒÛ Û Û Ý ÷÷ ñ õ 3÷	÷ 	ó 	õ ?áˆCƒL€á‹h€ð :EðØðØ)6ðà#õ÷Dñ ð'(Øð'(ð ð'(ð
 ð'(ð +ð'(ð ô'(ôTR!ˆ'�!‰*ô R!ðj €òð.Ø,ð.àô.ð: )*ñ	GØðGà6ðGð !&ð	Gð
 õGôTðDØ	ôDð2
7Ø	ô
7ð ×Ò˜QÑñ0ó  ñ0rD   