o
    wvXj!  ã                	   @   s¼  d Z ddlmZmZ ddlZddlmZmZmZ ddl	m
Z
 ddlmZ ddlmZ zddlZdd	lmZ dd
lmZ W n eyK   d Z ZZY nw dZdZdZeeƒZG dd„ dejƒZG dd„ dejƒZedur…e ddd„ ¡ ejejdd�G dd„ de ƒƒƒZ!edkrÜe"dƒ e #¡ �AZ$e"d %ej&j'ej&j(¡ƒ e )¡ �Z*e"d %ej&j'¡ƒ e$ +e!¡Z,e* +de,¡ W d  ƒ n1 sÁw   Y  e$ -¡  W d  ƒ dS 1 sÕw   Y  dS dS )aÏ  Pyro transport, and Kombu Broker daemon.

Requires the :mod:`Pyro4` library to be installed.

To use the Pyro transport with Kombu, use an url of the form:
``pyro://localhost/kombu.broker``

The hostname is where the transport will be looking for a Pyro name server,
which is used in turn to locate the kombu.broker Pyro service.
This broker can be launched by simply executing this transport module directly,
with the command: ``python -m kombu.transport.pyro``
é    )Úabsolute_importÚunicode_literalsN)ÚreraiseÚQueueÚEmpty)Úcached_property)Ú
get_loggeré   )Úvirtual)ÚNamingError)ÚSerializerBasei‚#  z5Unable to locate pyro nameserver on host {0.hostname}zKUnable to lookup '{0.virtual_host}' in pyro nameserver on host {0.hostname}c                       s~   e Zd ZdZ‡ fdd„Zdd„ Zdd„ Zdd	„ Zddd„Zdd„ Z	dd„ Z
dd„ Zdd„ Zdd„ Zdd„ Zedd„ ƒZ‡  ZS )ÚChannelzPyro Channel.c                    s&   t t| ƒ ¡  | jr| j ¡  d S d S ©N)Úsuperr   ÚcloseÚshared_queuesÚ_pyroRelease©Úself©Ú	__class__© úQ/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/pyro.pyr   ,   s   ÿzChannel.closec                 C   s
   | j  ¡ S r   )r   Úget_queue_namesr   r   r   r   Úqueues1   ó   
zChannel.queuesc                 K   s    ||   ¡ vr| j |¡ d S d S r   ©r   r   Ú	new_queue©r   ÚqueueÚkwargsr   r   r   Ú
_new_queue4   s   ÿzChannel._new_queuec                 K   ó   | j  |¡S r   )r   Ú	has_queuer   r   r   r   Ú
_has_queue8   ó   zChannel._has_queueNc                 C   s   |   |¡}| j |¡S r   )Ú
_queue_forr   Úget)r   r   Útimeoutr   r   r   Ú_get;   s   
zChannel._getc                 C   s   ||   ¡ vr| j |¡ |S r   r   ©r   r   r   r   r   r&   ?   s   zChannel._queue_forc                 K   s   |   |¡}| j ||¡ d S r   )r&   r   Úput)r   r   Úmessager    r   r   r   Ú_putD   s   
zChannel._putc                 C   r"   r   )r   Úsizer*   r   r   r   Ú_sizeH   r%   zChannel._sizec                 O   s   | j  |¡ d S r   )r   Údelete)r   r   Úargsr    r   r   r   Ú_deleteK   s   zChannel._deletec                 C   r"   r   )r   Úpurger*   r   r   r   Ú_purgeN   r%   zChannel._purgec                 C   s   d S r   r   r*   r   r   r   Úafter_reply_message_receivedQ   s   z$Channel.after_reply_message_receivedc                 C   s   | j jS r   )Ú
connectionr   r   r   r   r   r   T   ó   zChannel.shared_queuesr   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r!   r$   r)   r&   r-   r/   r2   r4   r5   r   r   Ú__classcell__r   r   r   r   r   )   s    
r   c                   @   sD   e Zd ZdZeZe ¡ ZeZ	d Z
Zdd„ Zdd„ Zedd„ ƒZd	S )
Ú	TransportzPyro Transport.Úpyroc              	   C   s¤   t  d¡ | j}ztj|j| jd�}W n ty+   tttt	 
|¡ƒt ¡ d ƒ Y nw z| |j¡}t |¡W S  tyQ   tttt 
|¡ƒt ¡ d ƒ Y d S w )Nz0trying Pyro nameserver to find the broker daemon)ÚhostÚporté   )ÚloggerÚdebugÚclientr>   ÚlocateNSÚhostnameÚdefault_portr   r   ÚE_NAMESERVERÚformatÚsysÚexc_infoÚlookupÚvirtual_hostÚProxyÚE_LOOKUP)r   ÚconninfoÚ
nameserverÚurir   r   r   Ú_opene   s&   

ÿ
ÿÿ

ÿÿzTransport._openc                 C   s   t jS r   )r>   Ú__version__r   r   r   r   Údriver_versionv   s   zTransport.driver_versionc                 C   s   |   ¡ S r   )rS   r   r   r   r   r   y   r7   zTransport.shared_queuesN)r8   r9   r:   r;   r   r
   ÚBrokerStateÚstateÚDEFAULT_PORTrG   Údriver_typeÚdriver_namerS   rU   r   r   r   r   r   r   r=   Y   s    r=   zqueue.Emptyc                 C   s   t ƒ S r   )r   )ÚclsÚdatar   r   r   Ú<lambda>€   s    r]   Úsingle)Úinstance_modec                   @   sX   e Zd 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„ Zdd„ ZdS )ÚKombuBrokerzmKombu Broker used by the Pyro transport.

        You have to run this as a separate (Pyro) service.
        c                 C   s
   i | _ d S r   ©r   r   r   r   r   Ú__init__Š   r   zKombuBroker.__init__c                 C   s
   t | jƒS r   )Úlistr   r   r   r   r   r   �   r   zKombuBroker.get_queue_namesc                 C   s   || j v rd S tƒ | j |< d S r   )r   r   r*   r   r   r   r   �   s   
zKombuBroker.new_queuec                 C   s
   || j v S r   ra   r*   r   r   r   r#   •   r   zKombuBroker.has_queuec                 C   s   | j | jdd�S )NF)Úblock)r   r'   r*   r   r   r   r'   ˜   s   zKombuBroker.getc                 C   s   | j |  |¡ d S r   )r   r+   )r   r   r,   r   r   r   r+   ›   s   zKombuBroker.putc                 C   s   | j |  ¡ S r   )r   Úqsizer*   r   r   r   r.   ž   s   zKombuBroker.sizec                 C   s   | j |= d S r   ra   r*   r   r   r   r0   ¡   r%   zKombuBroker.deletec                 C   s0   	 z| j | jdd� W n
 ty   Y d S w q)NTF)Úblocking)r   r'   r   r*   r   r   r   r3   ¤   s   ÿýzKombuBroker.purgeN)r8   r9   r:   r;   rb   r   r   r#   r'   r+   r.   r0   r3   r   r   r   r   r`   ‚   s    r`   Ú__main__z,Launching Broker for Kombu's Pyro transport.z)(Expecting a Pyro name server at {0}:{1})zBYou can connect with Kombu using the url 'pyro://{0}/kombu.broker'zkombu.broker).r;   Ú
__future__r   r   rJ   Ú
kombu.fiver   r   r   Úkombu.utils.objectsr   Ú	kombu.logr   Ú r
   ÚPyro4r>   ÚPyro4.errorsr   Ú
Pyro4.utilr   ÚImportErrorrX   rH   rO   r8   rB   r   r=   Úregister_dict_to_classÚexposeÚbehaviorÚobjectr`   ÚprintÚDaemonÚdaemonrI   ÚconfigÚNS_HOSTÚNS_PORTrE   ÚnsÚregisterrR   ÚrequestLoopr   r   r   r   Ú<module>   sV    ÿ0%ÿ
*
ÿ

ÿ
ü
"øþ