o
    —¨Êhd	  ã                   @  s`   d Z ddlmZ ddlmZ ddlmZ ddlmZm	Z	 G dd„ de	j
ƒZ
G d	d
„ d
e	jƒZdS )aŽ  In-memory transport module for Kombu.

Simple transport using memory for storing messages.
Messages can be passed only between threads.

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

Connection String
=================
Connection string is in the following format:

.. code-block::

    memory://

é    )Úannotations)Údefaultdict)ÚQueueé   )ÚbaseÚvirtualc                      s�   e Zd ZdZeeƒZi ZdZdZ	dd„ Z
dd„ Zd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‡ fdd„Zdd„ Z‡  ZS )ÚChannelzIn-memory Channel.FTc                 K  s
   || j v S ©N)Úqueues©ÚselfÚqueueÚkwargs© r   úH/var/www/html/env/lib/python3.10/site-packages/kombu/transport/memory.pyÚ
_has_queue)   s   
zChannel._has_queuec                 K  s   || j vrtƒ | j |< d S d S r	   ©r
   r   r   r   r   r   Ú
_new_queue,   s   
ÿzChannel._new_queueNc                 C  s   |   |¡jdd�S )NF)Úblock)Ú
_queue_forÚget)r   r   Útimeoutr   r   r   Ú_get0   ó   zChannel._getc                 C  s    || j vrtƒ | j |< | j | S r	   r   ©r   r   r   r   r   r   3   s   

zChannel._queue_forc                 G  ó   d S r	   r   )r   Úargsr   r   r   Ú_queue_bind8   ó   zChannel._queue_bindc                 K  s&   |   ||¡D ]
}|  |¡ |¡ qd S r	   )Ú_lookupr   Úput)r   ÚexchangeÚmessageÚrouting_keyr   r   r   r   r   Ú_put_fanout;   s   ÿzChannel._put_fanoutc                 K  s   |   |¡ |¡ d S r	   )r   r    )r   r   r"   r   r   r   r   Ú_put?   s   zChannel._putc                 C  s   |   |¡ ¡ S r	   )r   Úqsizer   r   r   r   Ú_sizeB   s   zChannel._sizec                 O  s   | j  |d ¡ d S r	   )r
   Úpop)r   r   r   r   r   r   r   Ú_deleteE   r   zChannel._deletec                 C  s    |   |¡}| ¡ }|j ¡  |S r	   )r   r&   r   Úclear)r   r   ÚqÚsizer   r   r   Ú_purgeH   s   

zChannel._purgec                   s,   t ƒ  ¡  | j ¡ D ]}| ¡  q
i | _d S r	   )ÚsuperÚcloser
   ÚvaluesÚemptyr   ©Ú	__class__r   r   r/   N   s   


zChannel.closec                 C  r   r	   r   r   r   r   r   Úafter_reply_message_receivedT   r   z$Channel.after_reply_message_receivedr	   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   ÚsetÚeventsr
   Ú
do_restoreÚsupports_fanoutr   r   r   r   r   r$   r%   r'   r)   r-   r/   r4   Ú__classcell__r   r   r2   r   r   !   s$    

r   c                      sD   e Zd ZdZeZe ¡ Zej	j
Z
dZdZ‡ fdd„Zdd„ Z‡  ZS )Ú	TransportzIn-memory Transport.Úmemoryc                   s    t ƒ j|fi |¤Ž | j| _d S r	   )r.   Ú__init__Úglobal_stateÚstate)r   Úclientr   r2   r   r   r@   e   s   zTransport.__init__c                 C  s   dS )NzN/Ar   )r   r   r   r   Údriver_versioni   r   zTransport.driver_version)r5   r6   r7   r8   r   r   ÚBrokerStaterA   r   r>   Ú
implementsÚdriver_typeÚdriver_namer@   rD   r=   r   r   r2   r   r>   X   s    r>   N)r8   Ú
__future__r   Úcollectionsr   r   r   Ú r   r   r   r>   r   r   r   r   Ú<module>   s    7