o
    •ºé` :  ã                   @   s   d dl Z d dlmZ d dlZd dlmZ d dlZd dlmZ zd dlm	Z	 W n e
y5   d dlm	Z	 Y nw d dlmZmZ d dlZd dlmZ d dlmZmZ d d	lmZ d d
lmZmZ d dlmZmZmZ d dlmZm Z  d dl!m"Z" ddl#m$Z$m%Z%m&Z&m'Z'm(Z( e )e*¡Z+g d¢Z,edd„ ƒZ-G dd„ deƒZ.d%ddœde.fdd„Z/G dd„ de%ƒZ0G dd„ dƒZ1G dd „ d ƒZ2G d!d"„ d"eƒZ3e	d%ddœd#d$„ƒZ4dS )&é    N)Úcontextmanager)Úcount)ÚOptional)Úasynccontextmanager)ÚValueÚError)ÚChannel)ÚAuthenticatorÚBEGIN)Úget_bus)ÚFileDescriptorÚfds_buf_size)ÚParserÚMessageTypeÚMessage)Ú	ProxyBaseÚ
unwrap_msg)Úmessage_busé   )ÚMessageFiltersÚFilterHandleÚReplyMatcherÚRouterClosedÚcheck_replyable)Úopen_dbus_connectionÚopen_dbus_routerÚProxyc               
   c   sX   � zd V  W d S  t y+ }  z| jtjtjhv rt d¡d ‚t d | ¡¡| ‚d } ~ ww )Nzthis socket was already closedzsocket connection broken: {})ÚOSErrorÚerrnoÚEBADFÚENOTSOCKÚtrioÚClosedResourceErrorÚBrokenResourceErrorÚformat)Úexc© r&   ú1/usr/lib/python3/dist-packages/jeepney/io/trio.pyÚ)_translate_socket_errors_to_stream_errors0   s   €ÿþ€ûr(   c                   @   s~   e Zd ZdZddd„Zddœdefdd	„Zd
efdd„Zdd
e	fdd„Z
defdd„Zdd„ Zdd„ Zdd„ Zedd„ ƒZdS )ÚDBusConnectionaŽ  A plain D-Bus connection with no matching of replies.

    This doesn't run any separate tasks: sending and receiving are done in
    the task that calls those methods. It's suitable for implementing servers:
    several worker tasks can receive requests and send replies.
    For a typical client pattern, see :class:`DBusRouter`.

    Implements trio's channel interface for Message objects.
    Fc                 C   sD   || _ || _tƒ | _tdd�| _d | _t ¡ | _	t ¡ | _
d | _d S )Nr   )Ústart)ÚsocketÚ
enable_fdsr   Úparserr   Úoutgoing_serialÚunique_namer!   ÚLockÚ	send_lockÚ	recv_lockÚ_leftover_to_send)Úselfr+   r,   r&   r&   r'   Ú__init__I   s   


zDBusConnection.__init__N©ÚserialÚmessagec             	   Ã   sˆ   �| j 4 I dH š/ |du rt| jƒ}| jrt d¡nd}|j||d�}|  ||¡I dH  W d  ƒI dH  dS 1 I dH s=w   Y  dS )z.Serialise and send a :class:`~.Message` objectNÚi)Úfds)r1   Únextr.   r,   ÚarrayÚ	serialiseÚ
_send_data)r4   r8   r7   r:   Údatar&   r&   r'   ÚsendS   s   €
.ûzDBusConnection.sendr?   c              	   Ã   sà   �| j jr
t d¡‚tƒ �Y | jr|  | j¡I d H  t|ƒ�0}|r5| j  |gtj j	tj j
|fg¡I d H }n	| j  |¡I d H }|  ||¡I d H  W d   ƒ n1 sQw   Y  W d   ƒ d S W d   ƒ d S 1 siw   Y  d S )Nz!can't send data after sending EOF)r+   Údid_shutdown_SHUT_WRr!   r"   r(   r3   Ú_send_remainderÚ
memoryviewÚsendmsgÚ
SOL_SOCKETÚ
SCM_RIGHTSr@   )r4   r?   r:   Úsentr&   r&   r'   r>   ^   s"   €


ÿøû"ûzDBusConnection._send_datar   c                 Ã   sŽ   �z5|t |ƒk r1||d … �}| j |¡I d H }W d   ƒ n1 s"w   Y  ||7 }|t |ƒk sd | _W d S  tjyF   ||d … | _‚ w ©N)Úlenr+   r@   r3   r!   Ú	Cancelled)r4   r?   Úalready_sentÚ	remainingrG   r&   r&   r'   rB   q   s   €ÿýûzDBusConnection._send_remainderÚreturnc              	   Ã   sˆ   �| j 4 I dH š/ 	 | j ¡ }|dur|W  d  ƒI dH  S |  ¡ I dH \}}|s/t d¡‚| j ||¡ q
1 I dH s=w   Y  dS )z5Return the next available message from the connectionNTzSocket closed at the other end)r2   r-   Úget_next_messageÚ
_read_datar!   ÚEndOfChannelÚadd_data)r4   ÚmsgÚbr:   r&   r&   r'   Úreceive   s   €
ü
öÿzDBusConnection.receivec                 Ã   sÌ   �| j rC| j ¡ }tƒ � | j |tƒ ¡I d H \}}}}W d   ƒ n1 s&w   Y  |ttjddƒ@ r<|  	¡  t
dƒ‚|t |¡fS tƒ � | j d¡I d H }W d   ƒ |g fS 1 s]w   Y  |g fS )NÚ
MSG_CTRUNCr   z&Unable to receive all file descriptorsi   )r,   r-   Úbytes_desiredr(   r+   Úrecvmsgr   Úgetattrr!   Ú_closeÚRuntimeErrorr   Úfrom_ancdataÚrecv)r4   Únbytesr?   ÚancdataÚflagsÚ_r&   r&   r'   rO   Ž   s$   €
ÿÿ
ÿþzDBusConnection._read_datac                 C   s   | j  ¡  d | _d S rH   )r+   Úcloser3   ©r4   r&   r&   r'   rY   Ÿ   s   

zDBusConnection._closec                 Ã   s   �|   ¡  dS )zClose the D-Bus connectionN)rY   rb   r&   r&   r'   Úaclose¤   s   €zDBusConnection.aclosec              	   C  s†   �t  ¡ 4 I dH š-}t| ƒ}| |¡I dH  z|V  W | ¡ I dH  n| ¡ I dH  w W d  ƒI dH  dS 1 I dH s<w   Y  dS )aY  Temporarily wrap this connection as a :class:`DBusRouter`

        To be used like::

            async with conn.router() as req:
                reply = await req.send_and_get_reply(msg)

        While the router is running, you shouldn't use :meth:`receive`.
        Once the router is closed, you can use the plain connection again.
        N)r!   Úopen_nurseryÚ
DBusRouterr*   rc   )r4   ÚnurseryÚrouterr&   r&   r'   rg   ¨   s   €".úzDBusConnection.router)F)r   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r5   r   r@   Úbytesr>   rC   rB   rT   rO   rY   rc   r   rg   r&   r&   r&   r'   r)   ?   s    
	
r)   ÚSESSIONF©r,   rM   c          	   	   Ã   sÒ   �t | ƒ}t |¡I dH }t|d�}|D ]}| |¡I dH  | | ¡ I dH ¡ q| t¡I dH  t|j	|d�}| 
¡ 4 I dH š}| t ¡ ¡I dH }|jd |_W d  ƒI dH  |S 1 I dH sbw   Y  |S )zHOpen a plain D-Bus connection

    :return: :class:`DBusConnection`
    Nrn   r   )r   r!   Úopen_unix_socketr	   Úsend_allÚfeedÚreceive_somer
   r)   r+   rg   Úsend_and_get_replyr   ÚHelloÚbodyr/   )	Úbusr,   Úbus_addrÚsockÚauthrÚreq_dataÚconnrg   Úreplyr&   r&   r'   r   ½   s    €
þür   c                       sF   e Zd Zdef‡ fdd„Zedd„ ƒZdd„ Zdd	„ Zd
d„ Z	‡  Z
S )ÚTrioFilterHandleÚfiltersc                    s   t ƒ  |||¡ || _d S rH   )Úsuperr5   Úsend_channel)r4   r~   ÚruleÚsend_chnÚrecv_chn©Ú	__class__r&   r'   r5   Ø   s   
zTrioFilterHandle.__init__c                 C   s   | j S rH   ©Úqueuerb   r&   r&   r'   Úreceive_channelÜ   s   z TrioFilterHandle.receive_channelc                 Ã   s   �|   ¡  | j ¡ I d H  d S rH   )ra   r€   rc   rb   r&   r&   r'   rc   à   s   €zTrioFilterHandle.aclosec                 Ã   s   �| j S rH   r†   rb   r&   r&   r'   Ú
__aenter__ä   s   €zTrioFilterHandle.__aenter__c                 Ã   s   �|   ¡ I d H  d S rH   )rc   )r4   Úexc_typeÚexc_valÚexc_tbr&   r&   r'   Ú	__aexit__ç   s   €zTrioFilterHandle.__aexit__)rh   ri   rj   r   r5   Úpropertyrˆ   rc   r‰   r�   Ú__classcell__r&   r&   r„   r'   r}   ×   s    
r}   c                   @   s0   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
S )ÚFuturez4A very simple Future for trio based on `trio.Event`.c                 C   s   d | _ t ¡ | _d S rH   )Ú_outcomer!   ÚEventÚ_eventrb   r&   r&   r'   r5   í   s   zFuture.__init__c                 C   ó   t |ƒ| _| j ¡  d S rH   )r   r‘   r“   Úset)r4   Úresultr&   r&   r'   Ú
set_resultñ   ó   
zFuture.set_resultc                 C   r”   rH   )r   r‘   r“   r•   )r4   r%   r&   r&   r'   Úset_exceptionõ   r˜   zFuture.set_exceptionc                 Ã   s   �| j  ¡ I d H  | j ¡ S rH   )r“   Úwaitr‘   Úunwraprb   r&   r&   r'   Úgetù   s   €
z
Future.getN)rh   ri   rj   rk   r5   r—   r™   rœ   r&   r&   r&   r'   r�   ë   s    r�   c                   @   sž   e Zd ZdZdZdZdefdd„Zedd„ ƒZ	ddœd	d
„Z
defdd„Zdddœdeej fdd„Zdejfdd„Zdd„ Zdefdd„Zejfdd„ZdS )re   z„A client D-Bus connection which can wait for replies.

    This runs a separate receiver task and dispatches received messages.
    Nr{   c                 C   s   || _ tƒ | _tƒ | _d S rH   )Ú_connr   Ú_repliesr   Ú_filters)r4   r{   r&   r&   r'   r5     s   zDBusRouter.__init__c                 C   s   | j jS rH   )r�   r/   rb   r&   r&   r'   r/     s   zDBusRouter.unique_namer6   c                Ã   s   �| j j||d�I dH  dS )z/Send a message, don't wait for a reply
        r6   N)r�   r@   )r4   r8   r7   r&   r&   r'   r@     s   €zDBusRouter.sendrM   c                 Ã   s~   �t |ƒ | jdu rtdƒ‚t| jjƒ}| j |tƒ ¡�}| j	||d�I dH  | 
¡ I dH W  d  ƒ S 1 s8w   Y  dS )z„Send a method call message and wait for the reply

        Returns the reply message (method return or error message type).
        NzThis DBusRouter has stoppedr6   )r   Ú_rcv_cancel_scoper   r;   r�   r.   rž   Úcatchr�   r@   rœ   )r4   r8   r7   Ú	reply_futr&   r&   r'   rs     s   €
$þzDBusRouter.send_and_get_replyr   )ÚchannelÚbufsizer£   c                C   s,   |du rt  |¡\}}nd}t| j|||ƒS )a  Create a filter for incoming messages

        Usage::

            async with router.filter(rule) as receive_channel:
                matching_msg = await receive_channel.receive()

            # OR:
            send_chan, recv_chan = trio.open_memory_channel(1)
            async with router.filter(rule, channel=send_chan):
                matching_msg = await recv_chan.receive()

        If the channel fills up,
        The sending end of the channel is closed when leaving the ``async with``
        block, whether or not it was passed in.

        :param jeepney.MatchRule rule: Catch messages matching this rule
        :param trio.MemorySendChannel channel: Send matching messages here
        :param int bufsize: If no channel is passed in, create one with this size
        N)r!   Úopen_memory_channelr}   rŸ   )r4   r�   r£   r¤   Úrecv_channelr&   r&   r'   Úfilter#  s   zDBusRouter.filterrf   c                 Ã   s,   �| j d ur
tdƒ‚| | j¡I d H | _ d S )Nz+DBusRouter receiver task is already running)r    rZ   r*   Ú	_receiver)r4   rf   r&   r&   r'   r*   @  s   €
zDBusRouter.startc                 Ã   s0   �| j dur| j  ¡  d| _ t d¡I dH  dS )z Stop the sender & receiver tasksNr   )r    Úcancelr!   Úsleeprb   r&   r&   r'   rc   E  s
   €

zDBusRouter.acloserR   c              	   C   sJ   | j  |¡rdS | j |¡D ]}z|j |¡ W q tjy"   Y qw dS )zHandle one received messageN)rž   ÚdispatchrŸ   Úmatchesr€   Úsend_nowaitr!   Ú
WouldBlock)r4   rR   r§   r&   r&   r'   Ú	_dispatchR  s   ÿýzDBusRouter._dispatchc                 Ã   s´   �t  ¡ �K}d| _| |¡ z	 | j ¡ I dH }|  |¡ qd| _| j ¡  t  	d¡�}| j
j ¡ D ]}d|_|j ¡ I dH  q2W d  ƒ w 1 sJw   Y  w 1 sSw   Y  dS )z'Receiver loop - runs in a separate taskTNFé   )r!   ÚCancelScopeÚ
is_runningÚstartedr�   rT   r¯   rž   Údrop_allÚmove_on_afterrŸ   r~   ÚvaluesÚshieldr€   rc   )r4   Útask_statusÚcscoperR   Úcleanup_scoper§   r&   r&   r'   r¨   ]  s$   €


þ
þÿòzDBusRouter._receiver)rh   ri   rj   rk   Ú_nursery_mgrr    r)   r5   rŽ   r/   r@   r   rs   r   r!   ÚMemorySendChannelr§   ÚNurseryr*   rc   r¯   ÚTASK_STATUS_IGNOREDr¨   r&   r&   r&   r'   re   þ   s    
re   c                       s(   e Zd ZdZ‡ fdd„Zdd„ Z‡  ZS )r   aÓ  A trio proxy for calling D-Bus methods

    You can call methods on the proxy object, such as ``await bus_proxy.Hello()``
    to make a method call over D-Bus and wait for a reply. It will either
    return a tuple of returned data, or raise :exc:`.DBusErrorResponse`.
    The methods available are defined by the message generator you wrap.

    :param msggen: A message generator object.
    :param ~trio.DBusRouter router: Router to send and receive messages.
    c                    s(   t ƒ  |¡ t|tƒstdƒ‚|| _d S )Nz)Proxy can only be used with DBusRequester)r   r5   Ú
isinstancere   Ú	TypeErrorÚ_router)r4   Úmsggenrg   r„   r&   r'   r5   ~  s   

zProxy.__init__c                    s   ‡ ‡fdd„}|S )Nc                  Ÿ   s<   �ˆ | i |¤Ž}|j jtju sJ ‚ˆj |¡I d H }t|ƒS rH   )ÚheaderÚmessage_typer   Úmethod_callrÁ   rs   r   )ÚargsÚkwargsrR   r|   ©Úmake_msgr4   r&   r'   Úinner…  s
   €z!Proxy._method_call.<locals>.innerr&   )r4   rÉ   rÊ   r&   rÈ   r'   Ú_method_call„  s   zProxy._method_call)rh   ri   rj   rk   r5   rË   r�   r&   r&   r„   r'   r   s  s    
r   c             
   C  s”   �t | |d�I dH }|4 I dH š- | ¡ 4 I dH š}|V  W d  ƒI dH  n1 I dH s-w   Y  W d  ƒI dH  dS 1 I dH sCw   Y  dS )a§  Open a D-Bus 'router' to send and receive messages.

    Use as an async context manager::

        async with open_dbus_router() as req:
            ...

    :param str bus: 'SESSION' or 'SYSTEM' or a supported address.
    :return: :class:`DBusRouter`

    This is a shortcut for::

        conn = await open_dbus_connection()
        async with conn:
            async with conn.router() as req:
                ...
    rn   N)r   rg   )rv   r,   r{   Úrtrr&   r&   r'   r   Ž  s   €*ÿ.ÿr   )rm   )5r<   Ú
contextlibr   r   Ú	itertoolsr   ÚloggingÚtypingr   r   ÚImportErrorÚasync_generatorÚoutcomer   r   r!   Útrio.abcr   Újeepney.authr	   r
   Újeepney.busr   Újeepney.fdsr   r   Újeepney.low_levelr   r   r   Újeepney.wrappersr   r   Újeepney.bus_messagesr   Úcommonr   r   r   r   r   Ú	getLoggerrh   ÚlogÚ__all__r(   r)   r   r}   r�   re   r   r   r&   r&   r&   r'   Ú<module>   sB    ÿ

~u