o
    —¨ÊhÔ  ã                	   @  s¾  d Z ddlmZ ddlZddlmZmZ ddl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yM   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ƒƒƒ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 module for kombu.

Pyro transport, and Kombu Broker daemon.

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

Features
========
* Type: Virtual
* Supports Direct: Yes
* Supports Topic: Yes
* Supports Fanout: No
* Supports Priority: No
* Supports TTL: No

Connection String
=================

To use the Pyro transport with Kombu, use an url of the form:

.. code-block::

    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``

Transport Options
=================
é    )ÚannotationsN)ÚEmptyÚQueue)Úreraise)Ú
get_logger)Úcached_propertyé   )Ú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 ƒ  ¡  | jr| j ¡  d S d S ©N)ÚsuperÚcloseÚshared_queuesÚ_pyroRelease©Úself©Ú	__class__© úF/var/www/html/env/lib/python3.10/site-packages/kombu/transport/pyro.pyr   C   s   
ÿzChannel.closec                 C  s
   | j  ¡ S r   )r   Úget_queue_namesr   r   r   r   ÚqueuesH   ó   
zChannel.queuesc                 K  s    ||   ¡ vr| j |¡ d S d S r   ©r   r   Ú	new_queue©r   ÚqueueÚkwargsr   r   r   Ú
_new_queueK   s   ÿzChannel._new_queuec                 K  ó   | j  |¡S r   )r   Ú	has_queuer   r   r   r   Ú
_has_queueO   ó   zChannel._has_queueNc                 C  s   |   |¡}| j |¡S r   )Ú
_queue_forr   Úget)r   r   Útimeoutr   r   r   Ú_getR   s   
zChannel._getc                 C  s   ||   ¡ vr| j |¡ |S r   r   ©r   r   r   r   r   r%   V   s   zChannel._queue_forc                 K  s   |   |¡}| j ||¡ d S r   )r%   r   Úput)r   r   Úmessager   r   r   r   Ú_put[   s   
zChannel._putc                 C  r!   r   )r   Úsizer)   r   r   r   Ú_size_   r$   zChannel._sizec                 O  s   | j  |¡ d S r   )r   Údelete)r   r   Úargsr   r   r   r   Ú_deleteb   s   zChannel._deletec                 C  r!   r   )r   Úpurger)   r   r   r   Ú_purgee   r$   zChannel._purgec                 C  s   d S r   r   r)   r   r   r   Úafter_reply_message_receivedh   s   z$Channel.after_reply_message_receivedc                 C  s   | j jS r   )Ú
connectionr   r   r   r   r   r   k   ó   zChannel.shared_queuesr   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r    r#   r(   r%   r,   r.   r1   r3   r4   r   r   Ú__classcell__r   r   r   r   r   @   s    
r   c                      sT   e Zd ZdZeZe ¡ ZeZ	d Z
Z‡ fdd„Zdd„ Zdd„ Zed	d
„ ƒZ‡  ZS )Ú	TransportzPyro Transport.Úpyroc                   s    t ƒ j|fi |¤Ž | j| _d S r   )r   Ú__init__Úglobal_stateÚstate)r   Úclientr   r   r   r   r>   }   s   zTransport.__init__c              	   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ÚdebugrA   r=   ÚlocateNSÚhostnameÚdefault_portr
   r   ÚE_NAMESERVERÚformatÚsysÚexc_infoÚlookupÚvirtual_hostÚProxyÚE_LOOKUP)r   ÚconninfoÚ
nameserverÚurir   r   r   Ú_open�   s&   

ÿ
ÿÿ

ÿÿzTransport._openc                 C  s   t jS r   )r=   Ú__version__r   r   r   r   Údriver_version’   s   zTransport.driver_versionc                 C  s   |   ¡ S r   )rU   r   r   r   r   r   •   r6   zTransport.shared_queues)r7   r8   r9   r:   r   r	   ÚBrokerStater?   ÚDEFAULT_PORTrI   Údriver_typeÚdriver_namer>   rU   rW   r   r   r;   r   r   r   r   r<   p   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   r>   ¦   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   rb   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   rb   r)   r   r   r   r/   ½   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   r2   À   s   ÿýzKombuBroker.purgeN)r7   r8   r9   r:   r>   r   r   r"   r&   r*   r-   r/   r2   r   r   r   r   ra   ž   s    ra   Ú__main__z,Launching Broker for Kombu's Pyro transport.z'(Expecting a Pyro name server at {}:{})zAYou can connect with Kombu using the url 'pyro://{}/kombu.broker'zkombu.broker)-r:   Ú
__future__r   rL   r   r   r   Úkombu.exceptionsr   Ú	kombu.logr   Úkombu.utils.objectsr   Ú r	   ÚPyro4r=   ÚPyro4.errorsr
   Ú
Pyro4.utilr   ÚImportErrorrY   rJ   rQ   r7   rE   r   r<   Úregister_dict_to_classÚexposeÚbehaviorra   ÚprintÚDaemonÚdaemonrK   ÚconfigÚNS_HOSTÚNS_PORTrG   ÚnsÚregisterrT   ÚrequestLoopr   r   r   r   Ú<module>   sX    "ÿ0*ÿ
*
ÿ

ÿ
ü
"øþ