o
    —¨Êh°(  ã                   @  sv  d Z ddlmZ ddlZddlZddlZddlZddlmZ ddl	m
Z
 ddlmZ ddlmZ ddlmZ dd	lmZ dd
lmZmZ ddlmZmZ ddlmZ dZd eeeƒ¡Zej dkr‚ddl!Z!ddl"Z"ddl#Z#e"j$Z%dZ&e"j'Z(e! )¡ Z*dd„ Z+dd„ Z,nej dkrœddl-Z-ddl-m%Z%m&Z& dd„ Z+dd„ Z,ne.dƒ‚edg d¢ƒZ/G dd„ dej0ƒZ0G dd„ dej1ƒZ1dS )a†	  File-system Transport module for kombu.

Transport using the file-system as the message store. Messages written to the
queue are stored in `data_folder_in` directory and
messages read from the queue are read from `data_folder_out` directory. Both
directories must be created manually. Simple example:

* Producer:

.. code-block:: python

    import kombu

    conn = kombu.Connection(
        'filesystem://', transport_options={
            'data_folder_in': 'data_in', 'data_folder_out': 'data_out'
        }
    )
    conn.connect()

    test_queue = kombu.Queue('test', routing_key='test')

    with conn as conn:
        with conn.default_channel as channel:
            producer = kombu.Producer(channel)
            producer.publish(
                        {'hello': 'world'},
                        retry=True,
                        exchange=test_queue.exchange,
                        routing_key=test_queue.routing_key,
                        declare=[test_queue],
                        serializer='pickle'
            )

* Consumer:

.. code-block:: python

    import kombu

    conn = kombu.Connection(
        'filesystem://', transport_options={
            'data_folder_in': 'data_out', 'data_folder_out': 'data_in'
        }
    )
    conn.connect()

    def callback(body, message):
        print(body, message)
        message.ack()

    test_queue = kombu.Queue('test', routing_key='test')

    with conn as conn:
        with conn.default_channel as channel:
            consumer = kombu.Consumer(
                conn, [test_queue], accept=['pickle']
            )
            consumer.register_callback(callback)
            with consumer:
                conn.drain_events(timeout=1)

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

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

.. code-block::

    filesystem://

Transport Options
=================
* ``data_folder_in`` - directory where are messages stored when written
  to queue.
* ``data_folder_out`` - directory from which are messages read when read from
  queue.
* ``store_processed`` - if set to True, all processed messages are backed up to
  ``processed_folder``.
* ``processed_folder`` - directory where are backed up processed files.
* ``control_folder`` - directory where are exchange-queue table stored.
é    )ÚannotationsN)Ú
namedtuple)ÚPath)ÚEmpty)Ú	monotonic)ÚChannelError)Úvirtual)Úbytes_to_strÚstr_to_bytes)ÚdumpsÚloads)Úcached_property)é   r   r   Ú.Úntc                 C  s$   t  |  ¡ ¡}t  ||ddt¡ dS )úCreate file lock.r   ì     þ N)Ú	win32fileÚ_get_osfhandleÚfilenoÚ
LockFileExÚ__overlapped)ÚfileÚflagsÚhfile© r   úL/var/www/html/env/lib/python3.10/site-packages/kombu/transport/filesystem.pyÚlock}   s   r   c                 C  s"   t  |  ¡ ¡}t  |ddt¡ dS )úRemove file lock.r   r   N)r   r   r   ÚUnlockFileExr   )r   r   r   r   r   Úunlock‚   s   r    Úposix)ÚLOCK_EXÚLOCK_SHc                 C  s   t  |  ¡ |¡ dS )r   N)ÚfcntlÚflockr   )r   r   r   r   r   r   �   s   c                 C  s   t  |  ¡ t j¡ dS )r   N)r$   r%   r   ÚLOCK_UN)r   r   r   r   r    ‘   s   z9Filesystem plugin only defined for NT and POSIX platformsÚexchange_queue_t)Úrouting_keyÚpatternÚqueuec                   @  s”   e Zd 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edd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZdS )ÚChannelzFilesystem Channel.Tc                 C  sœ   | j |› d� }z-| d¡}zt|tƒ tt| ¡ ƒƒ}dd„ |D ƒW t|ƒ | ¡  W S t|ƒ | ¡  w  t	y@   g  Y S  t
yM   td|› �ƒ‚w )Nú	.exchangeÚrc                 S  ó   g | ]}t |Ž ‘qS r   ©r'   ©Ú.0Úqr   r   r   Ú
<listcomp>«   ó    z%Channel.get_table.<locals>.<listcomp>zCannot open )Úcontrol_folderÚopenr   r#   r   r	   Úreadr    ÚcloseÚFileNotFoundErrorÚOSErrorr   )ÚselfÚexchanger   Úf_objÚexchange_tabler   r   r   Ú	get_table¤   s    

ÿ
ÿzChannel.get_tablec           
      C  s  | j |› d� }| j jdd� t|pd|pd|pdƒ}zf| ¡ rT|jddd�}t|tƒ tt| 	¡ ƒƒ}dd	„ |D ƒ}	||	vrS|	 
d|¡ | d¡ | tt|	ƒƒ¡ n#|jd
dd�}t|tƒ |g}	| tt|	ƒƒ¡ W t|ƒ | ¡  d S W t|ƒ | ¡  d S t|ƒ | ¡  w )Nr,   T)Úexist_okÚ zrb+r   ©Ú	bufferingc                 S  r.   r   r/   r0   r   r   r   r3   ¾   r4   z'Channel._queue_bind.<locals>.<listcomp>Úwb)r5   Úmkdirr'   Úexistsr6   r   r"   r   r	   r7   ÚinsertÚseekÚwriter
   r   r    r8   )
r;   r<   r(   r)   r*   r   Ú	queue_valr=   r>   Úqueuesr   r   r   Ú_queue_bind´   s6   ÿ

€
€ÿÿ
zChannel._queue_bindc                 K  s*   |   |¡D ]}| j|j|fi |¤Ž qd S ©N)r?   Ú_putr*   )r;   r<   Úpayloadr(   Úkwargsr2   r   r   r   Ú_put_fanoutÌ   s   ÿzChannel._put_fanoutc                 K  s¨   d  tttƒ d ƒƒt ¡ |¡}tj | j	|¡}z2zt
|ddd�}t|tƒ | tt|ƒƒ¡ W n ty?   td|›d�ƒ‚w W t|ƒ | ¡  dS t|ƒ | ¡  w )	zPut `message` onto `queue`.z{}_{}.{}.msgiè  rD   r   rB   zCannot add file z to directoryN)ÚformatÚintÚroundr   ÚuuidÚuuid4ÚosÚpathÚjoinÚdata_folder_outr6   r   r"   rI   r
   r   r:   r   r    r8   )r;   r*   rO   rP   ÚfilenameÚfr   r   r   rN   Ð   s$   ÿ

ÿÿÿÿ
zChannel._putc                 C  sú   d| d }t  | j¡}t|ƒ}t|ƒdkrz| d¡}| |¡dk r#q| jr*| j}nt	 
¡ }zt t j | j|¡|¡ W n	 tyE   Y qw t j ||¡}zt|dƒ}| ¡ }| ¡  | jsct  |¡ W n tys   td|›d�ƒ‚w tt|ƒƒS tƒ ‚)zGet next message from `queue`.r   ú.msgr   ÚrbzCannot read file z from queue.)rW   ÚlistdirÚdata_folder_inÚsortedÚlenÚpopÚfindÚstore_processedÚprocessed_folderÚtempfileÚ
gettempdirÚshutilÚmoverX   rY   r:   r6   r7   r8   Úremover   r   r	   r   )r;   r*   Ú
queue_findÚfolderr[   rf   r\   rO   r   r   r   Ú_getá   s@   
ÿþ

€
ÿÿzChannel._getc                 C  sŒ   d}d| d }t  | j¡}t|ƒdkrD| ¡ }z| |¡dk r"W qt j | j|¡}t  |¡ |d7 }W n	 t	y=   Y nw t|ƒdks|S )z!Remove all messages from `queue`.r   r   r]   r   )
rW   r_   r`   rb   rc   rd   rX   rY   rk   r:   ©r;   r*   Úcountrl   rm   r[   r   r   r   Ú_purge	  s    
ýôzChannel._purgec                 C  sX   d}d|› d�}t  | j¡}t|ƒdkr*| ¡ }| |¡dk r q|d7 }t|ƒdks|S )z<Return the number of messages in `queue` as an :class:`int`.r   r   r]   r   )rW   r_   r`   rb   rc   rd   ro   r   r   r   Ú_size"  s   ù	zChannel._sizec                 C  s
   | j jjS rM   )Ú
connectionÚclientÚtransport_options©r;   r   r   r   ru   3  s   
zChannel.transport_optionsc                 C  ó   | j  dd¡S )Nr`   Údata_in©ru   Úgetrv   r   r   r   r`   7  ó   zChannel.data_folder_inc                 C  rw   )NrZ   Údata_outry   rv   r   r   r   rZ   ;  r{   zChannel.data_folder_outc                 C  rw   )Nre   Fry   rv   r   r   r   re   ?  r{   zChannel.store_processedc                 C  rw   )Nrf   Ú	processedry   rv   r   r   r   rf   C  r{   zChannel.processed_folderc                 C  s   t | j dd¡ƒS )Nr5   Úcontrol)r   ru   rz   rv   r   r   r   r5   G  s   zChannel.control_folderN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Úsupports_fanoutr?   rL   rQ   rN   rn   rq   rr   Úpropertyru   r   r`   rZ   re   rf   r5   r   r   r   r   r+   Ÿ   s,    (




r+   c                      sZ   e Zd ZdZejjjdeg d¢ƒd�Ze	Z	e 
¡ ZdZdZdZ‡ fdd„Zd	d
„ Z‡  ZS )Ú	TransportzFilesystem Transport.F)ÚdirectÚtopicÚfanout)ÚasynchronousÚexchange_typer   Ú
filesystemc                   s    t ƒ j|fi |¤Ž | j| _d S rM   )ÚsuperÚ__init__Úglobal_stateÚstate)r;   rt   rP   ©Ú	__class__r   r   r�   [  s   zTransport.__init__c                 C  s   dS )NzN/Ar   rv   r   r   r   Údriver_version_  s   zTransport.driver_version)r   r€   r�   r‚   r   r…   Ú
implementsÚextendÚ	frozensetr+   ÚBrokerStaterŽ   Údefault_portÚdriver_typeÚdriver_namer�   r’   Ú__classcell__r   r   r�   r   r…   L  s    
þr…   )2r‚   Ú
__future__r   rW   ri   rg   rU   Úcollectionsr   Úpathlibr   r*   r   Útimer   Úkombu.exceptionsr   Úkombu.transportr   Úkombu.utils.encodingr	   r
   Úkombu.utils.jsonr   r   Úkombu.utils.objectsr   ÚVERSIONrY   ÚmapÚstrÚ__version__ÚnameÚ
pywintypesÚwin32conr   ÚLOCKFILE_EXCLUSIVE_LOCKr"   r#   ÚLOCKFILE_FAIL_IMMEDIATELYÚLOCK_NBÚ
OVERLAPPEDr   r   r    r$   ÚRuntimeErrorr'   r+   r…   r   r   r   r   Ú<module>   sR    [



ÿÿ .