Ë
    _§1hæ[  ã                  ó`  — U d dl mZ d dlZd dlmZmZ d dlmZmZ ddl	m
Z
mZ ddlmZ ddlmZmZ erd d	lmZ eg ee   f   Zd
ed<   eg ef   Zd
ed<    ede¬«      Z ede¬«      Z G d„ d«      Zej4                   G d„ de«      «       Zej4                   G d„ de«      «       Zej:                  j=                  dd«      e_        ej:                  j=                  dd«      e_        ddœ	 	 	 	 	 	 	 d%d„Zd&d„Z 	 	 	 	 d'd„Z!d(d„Z" G d„ d«      Z# G d„ d e«      Z$ G d!„ d"e«      Z%d)d#„Z&d*d$„Z'y)+é    )ÚannotationsN)Ú	AwaitableÚCallable)ÚTYPE_CHECKINGÚTypeVaré   )Ú_coreÚ_util©ÚStapledStream)ÚReceiveStreamÚ
SendStream)Ú	TypeAliasr   Ú	AsyncHookÚSyncHookÚSendStreamT)ÚboundÚReceiveStreamTc                  óP   — e Zd Zd
d„Zd
d„Zd
d„Zdd„Zdd„Zdd„Zddd„Z	ddd	„Z
y)Ú_UnboundedByteQueuec                ó–   — t        «       | _        d| _        t        j                  «       | _        t        j                  d«      | _        y )NFz%another task is already fetching data)	Ú	bytearrayÚ_dataÚ_closedr	   Ú
ParkingLotÚ_lotr
   ÚConflictDetectorÚ_fetch_lock©Úselfs    úZ/var/www/skyplay_api_hub/venv/lib/python3.12/site-packages/trio/testing/_memory_streams.pyÚ__init__z_UnboundedByteQueue.__init__   s8   € Ü“[ˆŒ
ØˆŒÜ×$Ñ$Ó&ˆŒ	Ü ×1Ñ1Ø3ó
ˆÕó    c                óF   — d| _         | j                  j                  «        y ©NT)r   r   Ú
unpark_allr   s    r!   Úclosez_UnboundedByteQueue.close(   s   € ØˆŒØ�	‰	×ÑÕr#   c                óB   — t        «       | _        | j                  «        y ©N)r   r   r'   r   s    r!   Úclose_and_wipez"_UnboundedByteQueue.close_and_wipe,   s   € Ü“[ˆŒ
Ø�
‰
�r#   c                ó¤   — | j                   rt        j                  d«      ‚| xj                  |z  c_        | j                  j                  «        y )Nzvirtual connection closed)r   r	   ÚClosedResourceErrorr   r   r&   ©r    Údatas     r!   Úputz_UnboundedByteQueue.put0   s9   € Ø�<Š<Ü×+Ñ+Ð,GÓHÐHØ�
Š
�dÑ�
Ø�	‰	×ÑÕr#   c                óT   — |€y t        j                  |«      }|dk  rt        d«      ‚y )Né   úmax_bytes must be >= 1)ÚoperatorÚindexÚ
ValueError©r    Ú	max_bytess     r!   Ú_check_max_bytesz$_UnboundedByteQueue._check_max_bytes6   s1   € ØÐØÜ—N‘N 9Ó-ˆ	Ø�qŠ=ÜÐ5Ó6Ð6ð r#   c                óØ   — | j                   s| j                  sJ ‚|€t        | j                  «      }| j                  r$| j                  d | }| j                  d |…= |sJ ‚|S t        «       S r)   )r   r   Úlenr   )r    r7   Úchunks      r!   Ú	_get_implz_UnboundedByteQueue._get_impl=   sc   € Ø�|Š|˜tŸzšzÐ)Ð)ØÐÜ˜DŸJ™J›ˆIØ�:Š:Ø—J‘J˜z 	Ð*ˆEØ—
‘
˜:˜I˜:Ð&ÙˆL�5ØˆLä“;Ðr#   Nc                óÚ   — | j                   5  | j                  |«       | j                  s| j                  st        j
                  ‚| j                  |«      cd d d «       S # 1 sw Y   y xY wr)   )r   r8   r   r   r	   Ú
WouldBlockr<   r6   s     r!   Ú
get_nowaitz_UnboundedByteQueue.get_nowaitI   sS   € Ø×Ññ 	-Ø×!Ñ! )Ô,Ø—<’<¨¯
ª
Ü×&Ñ&Ð&Ø—>‘> )Ó,÷		-÷ 	-ò 	-ús   �A
A!Á!A*c              ƒ  óH  K  — | j                   5  | j                  |«       | j                  s/| j                  s#| j                  j                  «       ƒ d {  –—†  nt        j                  «       ƒ d {  –—†  | j                  |«      cd d d «       S 7 Œ;7 Œ # 1 sw Y   y xY w­wr)   )	r   r8   r   r   r   Úparkr	   Ú
checkpointr<   r6   s     r!   Úgetz_UnboundedByteQueue.getP   s   è ø€ Ø×Ññ 	-Ø×!Ñ! )Ô,Ø—<’<¨¯
ª
Ø—i‘i—n‘nÓ&×&Ñ&ä×&Ñ&Ó(×(Ð(Ø—>‘> )Ó,÷	-ñ 	-ð 'øà(ø÷	-ð 	-üsA   ‚B"�ABÁBÁBÁ3BÁ4BÂ
B"ÂBÂBÂBÂB"©ÚreturnÚNone©r.   zbytes | bytearray | memoryviewrE   rF   )r7   ú
int | NonerE   rF   ©r7   rH   rE   r   r)   )Ú__name__Ú
__module__Ú__qualname__r"   r'   r*   r/   r8   r<   r?   rC   © r#   r!   r   r      s*   „ ó
óóóó7ó
ô-õ-r#   r   c                  ób   — e Zd ZdZ	 	 	 d
	 	 	 	 	 	 	 dd„Zdd„Zdd„Zdd„Zdd„Zddd„Z	ddd	„Z
y)ÚMemorySendStreamaÆ  An in-memory :class:`~trio.abc.SendStream`.

    Args:
      send_all_hook: An async function, or None. Called from
          :meth:`send_all`. Can do whatever you like.
      wait_send_all_might_not_block_hook: An async function, or None. Called
          from :meth:`wait_send_all_might_not_block`. Can do whatever you
          like.
      close_hook: A synchronous function, or None. Called from :meth:`close`
          and :meth:`aclose`. Can do whatever you like.

    .. attribute:: send_all_hook
                   wait_send_all_might_not_block_hook
                   close_hook

       All of these hooks are also exposed as attributes on the object, and
       you can change them at any time.

    Nc                ó€   — t        j                  d«      | _        t        «       | _        || _        || _        || _        y )Nú!another task is using this stream)r
   r   Ú_conflict_detectorr   Ú	_outgoingÚsend_all_hookÚ"wait_send_all_might_not_block_hookÚ
close_hook)r    rT   rU   rV   s       r!   r"   zMemorySendStream.__init__p   s=   € ô #(×"8Ñ"8Ø/ó#
ˆÔô -Ó.ˆŒØ*ˆÔØ2TˆÔ/Ø$ˆ�r#   c              ƒ  óH  K  — | j                   5  t        j                  «       ƒ d{  –—†  t        j                  «       ƒ d{  –—†  | j                  j	                  |«       | j
                  �| j                  «       ƒ d{  –—†  ddd«       y7 Œh7 ŒN7 Œ# 1 sw Y   yxY w­w)z}Places the given data into the object's internal buffer, and then
        calls the :attr:`send_all_hook` (if any).

        N)rR   r	   rB   rS   r/   rT   r-   s     r!   Úsend_allzMemorySendStream.send_all~   sŒ   è ø€ ð ×$Ñ$ñ 	+Ü×"Ñ"Ó$×$Ð$Ü×"Ñ"Ó$×$Ð$Ø�N‰N×Ñ˜tÔ$Ø×!Ñ!Ð-Ø×(Ñ(Ó*×*Ð*÷	+ð 	+Ø$øØ$øð +ø÷	+ð 	+üóP   ‚B"�B§B¨BÁBÁ>BÂBÂBÂ	B"ÂBÂBÂBÂBÂB"c              ƒ  óH  K  — | j                   5  t        j                  «       ƒ d{  –—†  t        j                  «       ƒ d{  –—†  | j                  j	                  d«       | j
                  �| j                  «       ƒ d{  –—†  ddd«       y7 Œh7 ŒN7 Œ# 1 sw Y   yxY w­w)znCalls the :attr:`wait_send_all_might_not_block_hook` (if any), and
        then returns immediately.

        Nr#   )rR   r	   rB   rS   r/   rU   r   s    r!   Úwait_send_all_might_not_blockz.MemorySendStream.wait_send_all_might_not_blockŒ   s’   è ø€ ð ×$Ñ$ñ 	@Ü×"Ñ"Ó$×$Ð$Ü×"Ñ"Ó$×$Ð$à�N‰N×Ñ˜sÔ#Ø×6Ñ6ÐBØ×=Ñ=Ó?×?Ð?÷	@ð 	@Ø$øØ$øð @ø÷	@ð 	@ürY   c                ór   — | j                   j                  «        | j                  �| j                  «        yy)z^Marks this stream as closed, and then calls the :attr:`close_hook`
        (if any).

        N)rS   r'   rV   r   s    r!   r'   zMemorySendStream.close›   s-   € ð 	�‰×ÑÔØ�?‰?Ð&Ø�O‰OÕð 'r#   c              ƒ  óh   K  — | j                  «        t        j                  «       ƒ d{  –—†  y7 Œ­w©z!Same as :meth:`close`, but async.N©r'   r	   rB   r   s    r!   ÚaclosezMemorySendStream.aclose¬   ó!   è ø€ à�
‰
ŒÜ×ÑÓ × Ò úó   ‚(2ª0«2c              ƒ  óT   K  — | j                   j                  |«      ƒ d{  –—† S 7 Œ­w)aÀ  Retrieves data from the internal buffer, blocking if necessary.

        Args:
          max_bytes (int or None): The maximum amount of data to
              retrieve. None (the default) means to retrieve all the data
              that's present (but still blocks until at least one byte is
              available).

        Returns:
          If this stream has been closed, an empty bytearray. Otherwise, the
          requested data.

        N)rS   rC   r6   s     r!   Úget_datazMemorySendStream.get_data±   s#   è ø€ ð —^‘^×'Ñ'¨	Ó2×2Ð2Ð2úó   ‚(¡&¢(c                ó8   — | j                   j                  |«      S )zÁRetrieves data from the internal buffer, but doesn't block.

        See :meth:`get_data` for details.

        Raises:
          trio.WouldBlock: if no data is available to retrieve.

        )rS   r?   r6   s     r!   Úget_data_nowaitz MemorySendStream.get_data_nowaitÁ   s   € ð �~‰~×(Ñ(¨Ó3Ð3r#   )NNN)rT   úAsyncHook | NonerU   rh   rV   úSyncHook | NonerE   rF   rG   rD   r)   rI   )rJ   rK   rL   Ú__doc__r"   rX   r[   r'   r`   rd   rg   rM   r#   r!   rO   rO   Z   s\   „ ñð, +/Ø?CØ&*ð	%à'ð%ð -=ð%ð $ð	%ð
 
ó%ó+ó@óó"!ô
3õ 	4r#   rO   c                  óR   — e Zd ZdZ	 	 d		 	 	 	 	 d
d„Zddd„Zdd„Zdd„Zdd„Zdd„Z	y)ÚMemoryReceiveStreamað  An in-memory :class:`~trio.abc.ReceiveStream`.

    Args:
      receive_some_hook: An async function, or None. Called from
          :meth:`receive_some`. Can do whatever you like.
      close_hook: A synchronous function, or None. Called from :meth:`close`
          and :meth:`aclose`. Can do whatever you like.

    .. attribute:: receive_some_hook
                   close_hook

       Both hooks are also exposed as attributes on the object, and you can
       change them at any time.

    Nc                ó€   — t        j                  d«      | _        t        «       | _        d| _        || _        || _        y )NrQ   F)r
   r   rR   r   Ú	_incomingr   Úreceive_some_hookrV   )r    ro   rV   s      r!   r"   zMemoryReceiveStream.__init__ß   s<   € ô
 #(×"8Ñ"8Ø/ó#
ˆÔô -Ó.ˆŒØˆŒØ!2ˆÔØ$ˆ�r#   c              ƒ  óÐ  K  — | j                   5  t        j                  «       ƒ d{  –—†  t        j                  «       ƒ d{  –—†  | j                  rt        j                  ‚| j
                  �| j                  «       ƒ d{  –—†  | j                  j                  |«      ƒ d{  –—† }| j                  rt        j                  ‚|cddd«       S 7 Œª7 Œ�7 ŒR7 Œ1# 1 sw Y   yxY w­w)zˆCalls the :attr:`receive_some_hook` (if any), and then retrieves
        data from the internal buffer, blocking if necessary.

        N)rR   r	   rB   r   r,   ro   rn   rC   )r    r7   r.   s      r!   Úreceive_somez MemoryReceiveStream.receive_someì   sÂ   è ø€ ð ×$Ñ$ñ 	Ü×"Ñ"Ó$×$Ð$Ü×"Ñ"Ó$×$Ð$Ø�|Š|Ü×/Ñ/Ð/Ø×%Ñ%Ð1Ø×,Ñ,Ó.×.Ð.ð
 Ÿ™×+Ñ+¨IÓ6×6ˆDØ�|Š|Ü×/Ñ/Ð/Ø÷	ñ 	Ø$øØ$øð /øð
 7ø÷	ð 	üsb   ‚C&�C§C¨CÁCÁ?CÂCÂ"CÂ&CÂ'!CÃ
C&ÃCÃCÃCÃCÃC#ÃC&c                ó€   — d| _         | j                  j                  «        | j                  �| j                  «        yy)zfDiscards any pending data from the internal buffer, and marks this
        stream as closed.

        TN)r   rn   r*   rV   r   s    r!   r'   zMemoryReceiveStream.close  s4   € ð
 ˆŒØ�‰×%Ñ%Ô'Ø�?‰?Ð&Ø�O‰OÕð 'r#   c              ƒ  óh   K  — | j                  «        t        j                  «       ƒ d{  –—†  y7 Œ­wr^   r_   r   s    r!   r`   zMemoryReceiveStream.aclose  ra   rb   c                ó:   — | j                   j                  |«       y)z.Appends the given data to the internal buffer.N)rn   r/   r-   s     r!   Úput_datazMemoryReceiveStream.put_data  s   € à�‰×Ñ˜4Õ r#   c                ó8   — | j                   j                  «        y)z2Adds an end-of-file marker to the internal buffer.N)rn   r'   r   s    r!   Úput_eofzMemoryReceiveStream.put_eof  s   € à�‰×ÑÕr#   )NN)ro   rh   rV   ri   rE   rF   r)   rI   rD   rG   )
rJ   rK   rL   rj   r"   rq   r'   r`   ru   rw   rM   r#   r!   rl   rl   Í   sI   „ ñð$ /3Ø&*ð%à+ð%ð $ð%ð 
ó	%ôó.ó!ó
!ôr#   rl   z._memory_streamsÚ )r7   c               ó   — 	 | j                  |«      }	 |s|j                  «        y|j	                  |«       	 y# t        j                  $ r Y yw xY w# t        j
                  $ r t        j                  d«      d‚w xY w)að  Take data out of the given :class:`MemorySendStream`'s internal buffer,
    and put it into the given :class:`MemoryReceiveStream`'s internal buffer.

    Args:
      memory_send_stream (MemorySendStream): The stream to get data from.
      memory_receive_stream (MemoryReceiveStream): The stream to put data into.
      max_bytes (int or None): The maximum amount of data to transfer in this
          call, or None to transfer all available data.

    Returns:
      True if it successfully transferred some data, or False if there was no
      data to transfer.

    This is used to implement :func:`memory_stream_one_way_pair` and
    :func:`memory_stream_pair`; see the latter's docstring for an example
    of how you might use it yourself.

    FzMemoryReceiveStream was closedNT)rg   r	   r>   rw   ru   r,   ÚBrokenResourceError)Úmemory_send_streamÚmemory_receive_streamr7   r.   s       r!   Úmemory_stream_pumpr}   $  sŒ   € ð0Ø!×1Ñ1°)Ó<ˆðTÙØ!×)Ñ)Ô+ð
 ð "×*Ñ*¨4Õ0ð øô ×Ñò Ùðûô ×$Ñ$ò TÜ×'Ñ'Ð(HÓIÈtÐSðTús   ‚: ”A §A ºAÁAÁ*A=c                 ón   ‡‡‡— t        «       Št        «       Šdˆˆfd„Šdˆfd„} | ‰_        ‰‰_        ‰‰fS )uQ  Create a connected, pure-Python, unidirectional stream with infinite
    buffering and flexible configuration options.

    You can think of this as being a no-operating-system-involved
    Trio-streamsified version of :func:`os.pipe` (except that :func:`os.pipe`
    returns the streams in the wrong order â€“ we follow the superior convention
    that data flows from left to right).

    Returns:
      A tuple (:class:`MemorySendStream`, :class:`MemoryReceiveStream`), where
      the :class:`MemorySendStream` has its hooks set up so that it calls
      :func:`memory_stream_pump` from its
      :attr:`~MemorySendStream.send_all_hook` and
      :attr:`~MemorySendStream.close_hook`.

    The end result is that data automatically flows from the
    :class:`MemorySendStream` to the :class:`MemoryReceiveStream`. But you're
    also free to rearrange things however you like. For example, you can
    temporarily set the :attr:`~MemorySendStream.send_all_hook` to None if you
    want to simulate a stall in data transmission. Or see
    :func:`memory_stream_pair` for a more elaborate example.

    c                 ó   •— t        ‰‰ «       y r)   )r}   )Úrecv_streamÚsend_streams   €€r!   Ú$pump_from_send_stream_to_recv_streamzHmemory_stream_one_way_pair.<locals>.pump_from_send_stream_to_recv_streame  s   ø€ Ü˜;¨Õ4r#   c               “  ó   •K  —  ‰ «        y ­wr)   rM   )r‚   s   €r!   Ú*async_pump_from_send_stream_to_recv_streamzNmemory_stream_one_way_pair.<locals>.async_pump_from_send_stream_to_recv_streami  s   øè ø€ Ù,Õ.ùs   ƒ	rD   )rO   rl   rT   rV   )r„   r‚   r€   r�   s    @@@r!   Úmemory_stream_one_way_pairr…   J  s=   ú€ ô0 #Ó$€KÜ%Ó'€Kö5õ/ð !K€KÔØA€KÔØ˜Ð#Ð#r#   c                ób   —  | «       \  }} | «       \  }}t        ||«      }t        ||«      }||fS r)   r   )Úone_way_pairÚ
pipe1_sendÚ
pipe1_recvÚ
pipe2_sendÚ
pipe2_recvÚstream1Ústream2s          r!   Ú_make_stapled_pairrŽ   q  s?   € ñ *›^Ñ€J�
Ù)›^Ñ€J�
Ü˜J¨
Ó3€GÜ˜J¨
Ó3€GØ�GÐÐr#   c                 ó    — t        t        «      S )a·  Create a connected, pure-Python, bidirectional stream with infinite
    buffering and flexible configuration options.

    This is a convenience function that creates two one-way streams using
    :func:`memory_stream_one_way_pair`, and then uses
    :class:`~trio.StapledStream` to combine them into a single bidirectional
    stream.

    This is like a no-operating-system-involved, Trio-streamsified version of
    :func:`socket.socketpair`.

    Returns:
      A pair of :class:`~trio.StapledStream` objects that are connected so
      that data automatically flows from one to the other in both directions.

    After creating a stream pair, you can send data back and forth, which is
    enough for simple tests::

       left, right = memory_stream_pair()
       await left.send_all(b"123")
       assert await right.receive_some() == b"123"
       await right.send_all(b"456")
       assert await left.receive_some() == b"456"

    But if you read the docs for :class:`~trio.StapledStream` and
    :func:`memory_stream_one_way_pair`, you'll see that all the pieces
    involved in wiring this up are public APIs, so you can adjust to suit the
    requirements of your tests. For example, here's how to tweak a stream so
    that data flowing from left to right trickles in one byte at a time (but
    data flowing from right to left proceeds at full speed)::

        left, right = memory_stream_pair()
        async def trickle():
            # left is a StapledStream, and left.send_stream is a MemorySendStream
            # right is a StapledStream, and right.recv_stream is a MemoryReceiveStream
            while memory_stream_pump(left.send_stream, right.recv_stream, max_bytes=1):
                # Pause between each byte
                await trio.sleep(1)
        # Normally this send_all_hook calls memory_stream_pump directly without
        # passing in a max_bytes. We replace it with our custom version:
        left.send_stream.send_all_hook = trickle

    And here's a simple test using our modified stream objects::

        async def sender():
            await left.send_all(b"12345")
            await left.send_eof()

        async def receiver():
            async for data in right:
                print(data)

        async with trio.open_nursery() as nursery:
            nursery.start_soon(sender)
            nursery.start_soon(receiver)

    By default, this will print ``b"12345"`` and then immediately exit; with
    our trickle stream it instead sleeps 1 second, then prints ``b"1"``, then
    sleeps 1 second, then prints ``b"2"``, etc.

    Pro-tip: you can insert sleep calls (like in our example above) to
    manipulate the flow of data across tasks... and then use
    :class:`MockClock` and its :attr:`~MockClock.autojump_threshold`
    functionality to keep your test suite running quickly.

    If you want to stress test a protocol implementation, one nice trick is to
    use the :mod:`random` module (preferably with a fixed seed) to move random
    numbers of bytes at a time, and insert random sleeps in between them. You
    can also set up a custom :attr:`~MemoryReceiveStream.receive_some_hook` if
    you want to manipulate things on the receiving side, and not just the
    sending side.

    )rŽ   r…   rM   r#   r!   Úmemory_stream_pairr�   ~  s   € ôZ Ô8Ó9Ð9r#   c                  óN   — e Zd Zd
d„Zd
d„Zdd„Zd
d„Zd
d„Zdd„Zd
d„Z	ddd	„Z
y)Ú_LockstepByteQueuec                óæ   — t        «       | _        d| _        d| _        d| _        t        j                  «       | _        t        j                  d«      | _
        t        j                  d«      | _        y )NFzanother task is already sendingz!another task is already receiving)r   r   Ú_sender_closedÚ_receiver_closedÚ_receiver_waitingr	   r   Ú_waitersr
   r   Ú_send_conflict_detectorÚ_receive_conflict_detectorr   s    r!   r"   z_LockstepByteQueue.__init__Ô  sa   € Ü“[ˆŒ
Ø#ˆÔØ %ˆÔØ!&ˆÔÜ×(Ñ(Ó*ˆŒÜ',×'=Ñ'=Ø-ó(
ˆÔ$ô +0×*@Ñ*@Ø/ó+
ˆÕ'r#   c                ó8   — | j                   j                  «        y r)   )r—   r&   r   s    r!   Ú_something_happenedz&_LockstepByteQueue._something_happenedá  s   € Ø�‰× Ñ Õ"r#   c              ƒ  óÖ   K  — 	  |«       rn<| j                   s| j                  rn#| j                  j                  «       ƒ d {  –—†  ŒDt	        j
                  «       ƒ d {  –—†  y 7 Œ"7 Œ­wr)   )r”   r•   r—   rA   r	   rB   )r    Úfns     r!   Ú	_wait_forz_LockstepByteQueue._wait_foræ  s]   è ø€ ØÙŒtØØ×"Ò" d×&;Ò&;ØØ—-‘-×$Ñ$Ó&×&Ð&ð ô ×ÑÓ × Ñ ð 'øØ ús$   ‚A A)ÁA%ÁA)ÁA'Á A)Á'A)c                ó2   — d| _         | j                  «        y r%   )r”   r›   r   s    r!   Úclose_senderz_LockstepByteQueue.close_senderï  s   € Ø"ˆÔØ× Ñ Õ"r#   c                ó2   — d| _         | j                  «        y r%   )r•   r›   r   s    r!   Úclose_receiverz!_LockstepByteQueue.close_receiveró  s   € Ø $ˆÔØ× Ñ Õ"r#   c              ƒ  óê  ‡ K  — ‰ j                   5  ‰ j                  rt        j                  ‚‰ j                  rt        j
                  ‚‰ j                  rJ ‚‰ xj                  |z  c_        ‰ j                  «        ‰ j                  ˆ fd„«      ƒ d {  –—†  ‰ j                  rt        j                  ‚‰ j                  r‰ j                  rt        j
                  ‚d d d «       y 7 ŒQ# 1 sw Y   y xY w­w)Nc                 ó"   •— ‰ j                   dk(  S ©Nr#   ©r   r   s   €r!   ú<lambda>z-_LockstepByteQueue.send_all.<locals>.<lambda>   s   ø€ ¨¯©°sÑ):€ r#   )	r˜   r”   r	   r,   r•   rz   r   r›   rž   r-   s   ` r!   rX   z_LockstepByteQueue.send_all÷  sÂ   øè ø€ Ø×)Ñ)ñ 	0Ø×"Ò"Ü×/Ñ/Ð/Ø×$Ò$Ü×/Ñ/Ð/Ø—z’zÐ!�>Ø�JŠJ˜$Ñ�JØ×$Ñ$Ô&Ø—.‘.Ó!:Ó;×;Ð;Ø×"Ò"Ü×/Ñ/Ð/Ø�zŠz˜d×3Ò3Ü×/Ñ/Ð/÷	0ð 	0ð <ø÷	0ð 	0üs0   ƒC3�BC'ÂC%ÂAC'Ã	C3Ã%C'Ã'C0Ã,C3c              ƒ  óf  ‡ K  — ‰ j                   5  ‰ j                  rt        j                  ‚‰ j                  r&t        j
                  «       ƒ d {  –—†  	 d d d «       y ‰ j                  ˆ fd„«      ƒ d {  –—†  ‰ j                  rt        j                  ‚	 d d d «       y 7 ŒP7 Œ,# 1 sw Y   y xY w­w)Nc                 ó   •— ‰ j                   S r)   )r–   r   s   €r!   r§   zB_LockstepByteQueue.wait_send_all_might_not_block.<locals>.<lambda>  s   ø€ ¨×)?Ñ)?€ r#   )r˜   r”   r	   r,   r•   rB   rž   r   s   `r!   r[   z0_LockstepByteQueue.wait_send_all_might_not_block  sŸ   øè ø€ Ø×)Ñ)ñ 	0Ø×"Ò"Ü×/Ñ/Ð/Ø×$Ò$Ü×&Ñ&Ó(×(Ð(Ø÷	0ð 	0ð —.‘.Ó!?Ó@×@Ð@Ø×"Ò"Ü×/Ñ/Ð/ð #÷	0ð 	0ð )øà@ø÷	0ð 	0üsM   ƒB1�A B%ÁB!ÁB%Á	B1ÁB%Á6B#Á7 B%Â	B1Â!B%Â#B%Â%B.Â*B1Nc              ƒ  óH  ‡ K  — ‰ j                   5  |�%t        j                  |«      }|dk  rt        d«      ‚‰ j                  rt
        j                  ‚d‰ _        ‰ j                  «        	 ‰ j                  ˆ fd„«      ƒ d {  –—†  d‰ _        ‰ j                  rt
        j                  ‚‰ j                  r9‰ j                  d | }‰ j                  d |…= ‰ j                  «        |cd d d «       S ‰ j                  sJ ‚	 d d d «       y7 Œ„# d‰ _        w xY w# 1 sw Y   y xY w­w)Nr1   r2   Tc                 ó"   •— ‰ j                   dk7  S r¥   r¦   r   s   €r!   r§   z1_LockstepByteQueue.receive_some.<locals>.<lambda>  s   ø€ ¨T¯Z©Z¸3Ñ->€ r#   Fr#   )r™   r3   r4   r5   r•   r	   r,   r–   r›   rž   r   r”   )r    r7   Úgots   `  r!   rq   z_LockstepByteQueue.receive_some  s  øè ø€ Ø×,Ñ,ñ 	àÐ$Ü$ŸN™N¨9Ó5�	Ø˜q’=Ü$Ð%=Ó>Ð>à×$Ò$Ü×/Ñ/Ð/à%)ˆDÔ"Ø×$Ñ$Ô&ð/Ø—n‘nÓ%>Ó?×?Ð?à).�Ô&Ø×$Ò$Ü×/Ñ/Ð/à�zŠzð —j‘j  )Ð,�Ø—J‘J˜z 	˜zÐ*Ø×(Ñ(Ô*Ø÷3	ñ 	ð6 ×*Ò*Ð*Ð*Ø÷9	ð 	ð @ùà).�Õ&ú÷	ð 	üsT   ƒD"�ADÁ,D
ÂDÂD
ÂADÃ&
D"Ã0DÃ?	D"ÄD
Ä
	DÄDÄDÄD"rD   )r�   zCallable[[], bool]rE   rF   rG   r)   ©r7   rH   rE   zbytes | bytearray)rJ   rK   rL   r"   r›   rž   r    r¢   rX   r[   rq   rM   r#   r!   r’   r’   Ó  s*   „ ó
ó#ó
!ó#ó#ó0ó	0õr#   r’   c                  ó4   — e Zd Zdd„Zdd„Zdd„Zd	d„Zdd„Zy)
Ú_LockstepSendStreamc                ó   — || _         y r)   ©Ú_lbq©r    Úlbqs     r!   r"   z_LockstepSendStream.__init__2  ó	   € Øˆ�	r#   c                ó8   — | j                   j                  «        y r)   )r²   r    r   s    r!   r'   z_LockstepSendStream.close5  s   € Ø�	‰	×ÑÕ r#   c              ƒ  óh   K  — | j                  «        t        j                  «       ƒ d {  –—†  y 7 Œ­wr)   r_   r   s    r!   r`   z_LockstepSendStream.aclose8  ó!   è ø€ Ø�
‰
ŒÜ×ÑÓ × Ò úrb   c              ƒ  óV   K  — | j                   j                  |«      ƒ d {  –—†  y 7 Œ­wr)   )r²   rX   r-   s     r!   rX   z_LockstepSendStream.send_all<  s   è ø€ Ø�i‰i× Ñ  Ó&×&Ò&ús   ‚)¡'¢)c              ƒ  óT   K  — | j                   j                  «       ƒ d {  –—†  y 7 Œ­wr)   )r²   r[   r   s    r!   r[   z1_LockstepSendStream.wait_send_all_might_not_block?  s   è ø€ Ø�i‰i×5Ñ5Ó7×7Ò7ús   ‚( &¡(N©r´   r’   rE   rF   rD   rG   )rJ   rK   rL   r"   r'   r`   rX   r[   rM   r#   r!   r¯   r¯   1  s   „ óó!ó!ó'ô8r#   r¯   c                  ó.   — e Zd Zdd„Zdd„Zdd„Zdd	d„Zy)
Ú_LockstepReceiveStreamc                ó   — || _         y r)   r±   r³   s     r!   r"   z_LockstepReceiveStream.__init__D  rµ   r#   c                ó8   — | j                   j                  «        y r)   )r²   r¢   r   s    r!   r'   z_LockstepReceiveStream.closeG  s   € Ø�	‰	× Ñ Õ"r#   c              ƒ  óh   K  — | j                  «        t        j                  «       ƒ d {  –—†  y 7 Œ­wr)   r_   r   s    r!   r`   z_LockstepReceiveStream.acloseJ  r¸   rb   Nc              ƒ  óT   K  — | j                   j                  |«      ƒ d {  –—† S 7 Œ­wr)   )r²   rq   r6   s     r!   rq   z#_LockstepReceiveStream.receive_someN  s!   è ø€ Ø—Y‘Y×+Ñ+¨IÓ6×6Ð6Ð6úre   r»   rD   r)   r­   )rJ   rK   rL   r"   r'   r`   rq   rM   r#   r!   r½   r½   C  s   „ óó#ó!õ7r#   r½   c                 óB   — t        «       } t        | «      t        | «      fS )a  Create a connected, pure Python, unidirectional stream where data flows
    in lockstep.

    Returns:
      A tuple
      (:class:`~trio.abc.SendStream`, :class:`~trio.abc.ReceiveStream`).

    This stream has *absolutely no* buffering. Each call to
    :meth:`~trio.abc.SendStream.send_all` will block until all the given data
    has been returned by a call to
    :meth:`~trio.abc.ReceiveStream.receive_some`.

    This can be useful for testing flow control mechanisms in an extreme case,
    or for setting up "clogged" streams to use with
    :func:`check_one_way_stream` and friends.

    In addition to fulfilling the :class:`~trio.abc.SendStream` and
    :class:`~trio.abc.ReceiveStream` interfaces, the return objects
    also have a synchronous ``close`` method.

    )r’   r¯   r½   )r´   s    r!   Úlockstep_stream_one_way_pairrÃ   R  s"   € ô. Ó
€CÜ˜sÓ#Ô%;¸CÓ%@Ð@Ð@r#   c                 ó    — t        t        «      S )a“  Create a connected, pure-Python, bidirectional stream where data flows
    in lockstep.

    Returns:
      A tuple (:class:`~trio.StapledStream`, :class:`~trio.StapledStream`).

    This is a convenience function that creates two one-way streams using
    :func:`lockstep_stream_one_way_pair`, and then uses
    :class:`~trio.StapledStream` to combine them into a single bidirectional
    stream.

    )rŽ   rÃ   rM   r#   r!   Úlockstep_stream_pairrÅ   m  s   € ô  Ô:Ó;Ð;r#   )r{   rO   r|   rl   r7   rH   rE   Úbool)rE   z,tuple[MemorySendStream, MemoryReceiveStream])r‡   z0Callable[[], tuple[SendStreamT, ReceiveStreamT]]rE   z]tuple[StapledStream[SendStreamT, ReceiveStreamT], StapledStream[SendStreamT, ReceiveStreamT]])rE   zqtuple[StapledStream[MemorySendStream, MemoryReceiveStream], StapledStream[MemorySendStream, MemoryReceiveStream]])rE   z tuple[SendStream, ReceiveStream])rE   zYtuple[StapledStream[SendStream, ReceiveStream], StapledStream[SendStream, ReceiveStream]])(Ú
__future__r   r3   Úcollections.abcr   r   Útypingr   r   rx   r	   r
   Ú_highlevel_genericr   Úabcr   r   Útyping_extensionsr   Úobjectr   Ú__annotations__r   r   r   r   ÚfinalrO   rl   rK   Úreplacer}   r…   rŽ   r�   r’   r¯   r½   rÃ   rÅ   rM   r#   r!   ú<module>rÑ      s}  ðÞ "ã ß /ß )ç Ý .ß +áÝ+ð    I¨fÑ$5Ð 5Ñ6€	ˆ9Ó 6à˜r 6˜zÑ*€ˆ)Ó *Ù�m¨:Ô6€ÙÐ)°Ô?€÷<-ñ <-ð~ ‡�ôo4�zó o4ó ðo4ðd ‡�ôJ˜-ó Jó ðJð\ /×9Ñ9×AÑAØ˜óÐ Ô ð "5×!?Ñ!?×!GÑ!GØ˜ó"Ð Ô ð !ñ	#Ø(ð#à.ð#ð ð	#ð
 
ó#óL$$ðN
ØBð
ðó
óM:÷j[ñ [ô|8˜*ô 8ô$7˜]ô 7óAô6<r#   