o
    wvXjÔd  ã                   @   sT  d Z ddlmZ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mZmZ ddlmZ ddlmZmZ ddlmZ dd	lmZ dd
lmZmZmZ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dl(m)Z) ee*ƒZ+dd„ ej,D ƒZ-de-d< dZ.dd„ Z/G dd„ de0ƒZ1G dd„ de)j2ƒZ2G dd„ de)j3ƒZ3dS )aV
  Amazon SQS Transport.

Amazon SQS transport module for Kombu.  This package implements an AMQP-like
interface on top of Amazons SQS service, with the goal of being optimized for
high performance and reliability.

The default settings for this module are focused now on high performance in
task queue situations where tasks are small, idempotent and run very fast.

SQS Features supported by this transport:
  Long Polling:
    https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-long-polling.html

    Long polling is enabled by setting the `wait_time_seconds` transport
    option to a number > 1.  Amazon supports up to 20 seconds.  This is
    enabled with 10 seconds by default.

  Batch API Actions:
   https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-batch-api.html

    The default behavior of the SQS Channel.drain_events() method is to
    request up to the 'prefetch_count' messages on every request to SQS.
    These messages are stored locally in a deque object and passed back
    to the Transport until the deque is empty, before triggering a new
    API call to Amazon.

    This behavior dramatically speeds up the rate that you can pull tasks
    from SQS when you have short-running tasks (or a large number of workers).

    When a Celery worker has multiple queues to monitor, it will pull down
    up to 'prefetch_count' messages from queueA and work on them all before
    moving on to queueB.  If queueB is empty, it will wait up until
    'polling_interval' expires before moving back and checking on queueA.

Other Features supported by this transport:
  Predefined Queues:
    The default behavior of this transport is to use a single AWS credential
    pair in order to manage all SQS queues (e.g. listing queues, creating
    queues, polling queues, deleting messages).

    If it is preferable for your environment to use a single AWS credential, you
    can use the 'predefined_queues' setting inside the  'transport_options' map.
    This setting allows you to specify the SQS queue URL and AWS credentials for
    each of your queues. For example, if you have two queues which both already
    exist in AWS) you can tell this transport about them as follows:

    transport_options = {
      'predefined_queues': {
        'queue-1': {
          'url': 'https://sqs.us-east-1.amazonaws.com/xxx/aaa',
          'access_key_id': 'a',
          'secret_access_key': 'b',
        },
        'queue-2': {
          'url': 'https://sqs.us-east-1.amazonaws.com/xxx/bbb',
          'access_key_id': 'c',
          'secret_access_key': 'd',
        },
      }
    }
é    )Úabsolute_importÚunicode_literalsN)ÚClientError)Ú	transformÚensure_promiseÚpromise)Úget_event_loop)Úboto3Ú
exceptions)ÚAsyncSQSConnection)ÚAsyncMessage)ÚEmptyÚrangeÚstring_tÚtext_t)Ú
get_logger)Ú
scheduling)Úbytes_to_strÚsafe_str)ÚloadsÚdumps)Úcached_propertyé   )Úvirtualc                 C   s   i | ]}|d vrt |ƒd“qS )z-_.é_   )Úord)Ú.0Úc© r   úP/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/SQS.pyÚ
<dictcomp>Z   s    r    é-   é.   é
   c                 C   s"   zt | ƒW S  ty   |  Y S w )z5Try to convert x' to int, or return x' if that fails.)ÚintÚ
ValueError)Úxr   r   r   Ú	maybe_intc   s
   
ÿr'   c                   @   s   e Zd ZdZdS )ÚUndefinedQueueExceptionzAPredefined queues are being used and an undefined queue was used.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r   r   r(   k   s    r(   c                       s
  e Zd ZdZdZdZdZdZdZi Z	dZ
i Zi Zeƒ Z‡ fdd„Zd	d
„ Z‡ fdd„Z‡ fdd„Zd`dd„Zdd„ Zefdd„Zdd„ Zdd„ Zdd„ Z‡ fdd„Zdd„ Zdd „ Zd!d"„ Zedfd#d$„Zd%d&„ Z dad'd(„Z!d)d*„ Z"efd+d,„Z#edfd-d.„Z$dbd0d1„Z%d2d3„ Z&	dcd4d5„Z'	6dd‡ fd7d8„	Z(de‡ fd:d;„	Z)d<d=„ Z*d>d?„ Z+‡ fd@dA„Z,dBdC„ Z-dadDdE„Z.dadFdG„Z/e0dHdI„ ƒZ1e0dJdK„ ƒZ2e3dLdM„ ƒZ4e3dNdO„ ƒZ5e3dPdQ„ ƒZ6e3dRdS„ ƒZ7e3dTdU„ ƒZ8e3dVdW„ ƒZ9e3dXdY„ ƒZ:e3dZd[„ ƒZ;e3d\d]„ ƒZ<e3d^d_„ ƒZ=‡  Z>S )fÚChannelzSQS Channel.z	us-east-1i  r#   zkombu%(vhost)sNc                    sH   t d u rtdƒ‚tt| ƒj|i |¤Ž |  | j¡ | d¡p tƒ | _	d S )Nzboto3 is not installedÚhub)
r	   ÚImportErrorÚsuperr-   Ú__init__Ú_update_queue_cacheÚqueue_name_prefixÚgetr   r.   )ÚselfÚargsÚkwargs©Ú	__class__r   r   r1   }   s
   zChannel.__init__c                 C   sj   | j r| j  ¡ D ]\}}|d | j|< qd S |  ¡ j|d�}| dg ¡D ]}| d¡d }|| j|< q$d S )NÚurl)ÚQueueNamePrefixÚ	QueueUrlsú/éÿÿÿÿ)Úpredefined_queuesÚitemsÚ_queue_cacheÚsqsÚlist_queuesr4   Úsplit)r5   r3   Ú
queue_nameÚqÚrespr:   r   r   r   r2   Š   s   þzChannel._update_queue_cachec                    s@   |r| j  |¡ | jr|  |¡ tt| ƒj||g|¢R i |¤ŽS ©N)Ú_noack_queuesÚaddr.   Ú_loop1r0   r-   Úbasic_consume)r5   ÚqueueÚno_ackr6   r7   r8   r   r   rL   •   s   

ÿÿÿzChannel.basic_consumec                    s0   || j v r| j| }| j |¡ tt| ƒ |¡S rH   )Ú
_consumersÚ_tag_to_queuerI   Údiscardr0   r-   Úbasic_cancel)r5   Úconsumer_tagrM   r8   r   r   rR   ž   s   

zChannel.basic_cancelc                 K   s,   | j r| j ¡ stƒ ‚| j| j||d� dS )z„Return a single payload message from one of our queues.

        Raises:
            Queue.Empty: if no messages available.
        )ÚtimeoutN)rO   ÚqosÚcan_consumer   Ú_pollÚcycle)r5   rT   Úcallbackr7   r   r   r   Údrain_events¤   s   zChannel.drain_eventsc                 C   s   t  | j| jt¡| _dS )a2  Reset the consume cycle.

        Returns:
            FairCycle: object that points to our _get_bulk() method
                rather than the standard _get() method.  This allows for
                multiple messages to be returned at once from SQS (
                based on the prefetch limit).
        N)r   Ú	FairCycleÚ	_get_bulkÚ_active_queuesr   Ú_cycle©r5   r   r   r   Ú_reset_cycle±   s   	

ÿzChannel._reset_cyclec                 C   sH   |  d¡r|dtdƒ … }tt|ƒƒ |¡}|d S tt|ƒƒ |¡S )z3Format AMQP queue name into a legal SQS queue name.ú.fifoN)ÚendswithÚlenr   r   Ú	translate)r5   ÚnameÚtableÚpartialr   r   r   Úentity_name¾   s
   
zChannel.entity_namec                 C   s   |   | j| ¡S rH   )rh   r3   )r5   rE   r   r   r   Úcanonical_queue_nameÇ   s   zChannel.canonical_queue_namec                 K   s¢   t |tƒs|S |  |¡}|| jvr|  |¡ z| j| W S  tyP   | jr-td |¡ƒ‚dt	| j
ƒi}| d¡r=d|d< |  ||¡}|d | j|< |d  Y S w )z-Ensure a queue with given name exists in SQS.ú<Queue with name '{}' must be defined in 'predefined_queues'.ÚVisibilityTimeoutra   ÚtrueÚ	FifoQueueÚQueueUrl)Ú
isinstancer   ri   rA   r2   ÚKeyErrorr?   r(   ÚformatÚstrÚvisibility_timeoutrb   Ú_create_queue)r5   rM   r7   Ú
attributesrG   r   r   r   Ú
_new_queueÊ   s(   



ý
ózChannel._new_queuec                 C   s6   | j rdS | | j d¡pi ¡ | j|d�j||d�S )z=Create an SQS queue with a given name and nominal attributes.Nzsqs-creation-attributes©rM   )Ú	QueueNameÚ
Attributes)r?   ÚupdateÚtransport_optionsr4   rB   Úcreate_queue)r5   rE   ru   r   r   r   rt   é   s   ÿþzChannel._create_queuec                    s,   | j rdS tt| ƒ |¡ | j |d¡ dS )zDelete queue by name.N)r?   r0   r-   Ú_deleterA   Úpop)r5   rM   r6   r7   r8   r   r   r}   ù   s   zChannel._deletec                 K   sÊ   |   |¡}|tƒ  t|ƒ¡dœ}| d¡r?d|d v r$|d d |d< nd|d< d|d v r7|d d |d< ntt ¡ ƒ|d< | j|  	|¡d�}| 
d¡r[|j||d d	 d
d� dS |jdi |¤Ž dS )zPut message onto queue.)rn   ÚMessageBodyra   ÚMessageGroupIdÚ
propertiesÚdefaultÚMessageDeduplicationIdrw   ÚredeliveredÚdelivery_tagr   )rn   ÚReceiptHandlerk   Nr   )rv   r   Úencoder   rb   rr   ÚuuidÚuuid4rB   ri   r4   Úchange_message_visibilityÚsend_message)r5   rM   Úmessager7   Úq_urlr   r   r   r   Ú_put   s*   
ÿ

ÿ
ÿ


ýzChannel._putc                 C   sÞ   zt  |d  ¡ ¡}W n ty   |d  ¡ }Y nw tt|ƒƒ}|| jv r9|  |¡}| j|d� 	||d ¡ |S z|d }|d d }W n t
y^   i }d|i}| t|ƒ|dœ¡ Y nw | ||dœ¡ |d |d< |S )	NÚBodyrw   r†   r�   Údelivery_info)Úbodyr�   ©Úsqs_messageÚ	sqs_queuer…   )Úbase64Ú	b64decoder‡   Ú	TypeErrorr   r   rI   rv   ÚasynsqsÚdelete_messagerp   rz   )r5   rŒ   rE   rM   r‘   Úpayloadr�   r�   r   r   r   Ú_message_to_python  s:   ÿ

þðþü	ÿzChannel._message_to_pythonc                    s    ˆ  ˆ¡‰ ‡ ‡‡fdd„|D ƒS )aÆ  Convert a list of SQS Message objects into Payloads.

        This method handles converting SQS Message objects into
        Payloads, and appropriately updating the queue depending on
        the 'ack' settings for that queue.

        Arguments:
            messages (SQSMessage): A list of SQS Message objects.
            queue (str): Name representing the queue they came from.

        Returns:
            List: A list of Payload objects
        c                    s   g | ]	}ˆ  |ˆˆ ¡‘qS r   )r›   )r   Úm©rF   rM   r5   r   r   Ú
<listcomp>I  s    z/Channel._messages_to_python.<locals>.<listcomp>)rv   )r5   ÚmessagesrM   r   r�   r   Ú_messages_to_python:  s   
zChannel._messages_to_pythonc           	      C   sŒ   |   ¡ }|rC|  |¡}| j|d�j||| jd�}| d¡rC|d D ]}t|d d� ¡ |d< q!|  |d |¡D ]	}| j	 
||¡ q7dS tƒ ‚)a  Try to retrieve multiple messages off ``queue``.

        Where :meth:`_get` returns a single Payload object, this method
        returns a list of Payload objects.  The number of objects returned
        is determined by the total number of messages available in the queue
        and the number of messages the QoS object allows (based on the
        prefetch_count).

        Note:
            Ignores QoS limits so caller is responsible for checking
            that we are allowed to consume at least one message from the
            queue.  get_bulk will then ask QoS for an estimate of
            the number of extra messages that we can consume.

        Arguments:
            queue (str): The queue name to pull from.

        Returns:
            List[Message]
        rw   ©rn   ÚMaxNumberOfMessagesÚWaitTimeSecondsÚMessagesr�   ©r‘   N)Ú_get_message_estimaterv   rB   Úreceive_messageÚwait_time_secondsr4   r   Údecoder    Ú
connectionÚ_deliverr   )	r5   rM   Úmax_if_unlimitedrY   Ú	max_countr�   rG   rœ   Úmsgr   r   r   r\   K  s   
þ
zChannel._get_bulkc                 C   sr   |   |¡}| j|d�j|d| jd�}| d¡r6t|d d d d� ¡ }||d d d< |  |d |¡d S tƒ ‚)z/Try to retrieve a single message off ``queue``.rw   r   r¡   r¤   r   r�   r¥   )	rv   rB   r§   r¨   r4   r   r©   r    r   )r5   rM   r�   rG   r‘   r   r   r   Ú_gett  s   
þ
zChannel._getc                 C   s   | j  | j|¡ d S rH   )r.   Ú	call_soonÚ_schedule_queue)r5   rM   Ú_r   r   r   rK   €  s   zChannel._loop1c                 C   sB   || j v r| j ¡ r| j|t| j|fƒd� d S |  |¡ d S d S ©N)rY   )r]   rU   rV   Ú_get_bulk_asyncr   rK   )r5   rM   r   r   r   r±   ƒ  s   


ÿúzChannel._schedule_queuec                 C   s*   | j  ¡ }t|d u r||ƒS t|dƒ|ƒS )Nr   )rU   Úcan_consume_max_estimateÚminÚmax)r5   r¬   Úmaxcountr   r   r   r¦   Œ  s   

þþzChannel._get_message_estimatec                 C   s0   |   ¡ }|r| j|||d�S t|ƒ}|g ƒ |S r³   )r¦   Ú
_get_asyncr   )r5   rM   r¬   rY   r¸   r   r   r   r´   “  s   zChannel._get_bulk_asyncr   c              	   C   s:   |   |¡}|  |¡}| j||| j|d�t| j|||ƒd�S )Nrw   )Úcountrª   rY   )rv   ri   Ú_get_from_sqsr˜   r   Ú_on_messages_ready)r5   rM   rº   rY   rF   Úqnamer   r   r   r¹   �  s   

þzChannel._get_asyncc                 C   sL   d|v r |d r"| j j}|d D ]}|  |||¡}|| |ƒ qd S d S d S )Nr¤   )rª   Ú
_callbacksr›   )r5   rM   r½   rŸ   Ú	callbacksr®   Ú
msg_parsedr   r   r   r¼   ¥  s   üzChannel._on_messages_readyc                 C   s\   |dur|n|j }| jr|| jvrtd |¡ƒ‚| j| }n| |¡}|j|||| j|d�S )zwRetrieve and handle messages from SQS.

        Uses long polling and returns :class:`~vine.promises.promise`.
        Nrj   )Únumber_messagesr¨   rY   )rª   r?   rA   r(   rq   Úget_queue_urlr§   r¨   )r5   rM   rº   rª   rY   Ú	queue_urlr   r   r   r»   ¬  s   
ý
ýzChannel._get_from_sqsr’   c                    s(   |D ]	}|j  |d ¡ qtt| ƒ |¡S rH   )r�   r~   r0   r-   Ú_restore)r5   rŒ   Úunwanted_delivery_infoÚunwanted_keyr8   r   r   rÄ   Â  s   zChannel._restoreFc                    s¶   z| j  |¡j}|d }W n ty   tt| ƒ |¡ Y d S w d }d|v r-|  |d ¡}z| j|d�j	|d |d d� W n t
yP   tt| ƒ |¡ Y d S w tt| ƒ |¡ d S )Nr“   Úrouting_keyrw   r”   r†   )rn   r†   )rU   r4   r�   rp   r0   r-   Ú	basic_ackri   rB   r™   r   Úbasic_reject)r5   r…   ÚmultiplerŒ   r“   rM   r8   r   r   rÈ   É  s$   ÿ
þÿzChannel.basic_ackc                 C   s<   |   |¡}| j|  |¡d�}|j|dgd�}t|d d ƒS )z)Return the number of messages in a queue.rw   ÚApproximateNumberOfMessages)rn   ÚAttributeNamesry   )rv   rB   ri   Úget_queue_attributesr$   )r5   rM   r:   r   rG   r   r   r   Ú_sizeÞ  s   
þzChannel._sizec                 C   sN   |   |¡}d}tdƒD ]}|t|  |¡ƒ7 }|s nq| j|d�j|d� |S )z'Delete all current messages in a queue.r   r#   rw   )rn   )rv   r   r$   rÎ   rB   Úpurge_queue)r5   rM   rF   ÚsizeÚir   r   r   Ú_purgeç  s   
ÿzChannel._purgec                    s   t t| ƒ ¡  d S rH   )r0   r-   Úcloser_   r8   r   r   rÓ   ô  s   zChannel.closec                 C   sR   t jj|||d�}| jd ur| jnd}d|i}| jd ur!| j|d< |jdi |¤ŽS )N)Úregion_nameÚaws_access_key_idÚaws_secret_access_keyTÚuse_sslÚendpoint_urlrB   )rB   )r	   ÚsessionÚSessionÚ	is_securerØ   Úclient)r5   ÚregionÚaccess_key_idÚsecret_access_keyrÙ   rÛ   Úclient_kwargsr   r   r   Únew_sqs_clientý  s   ýÿ

zChannel.new_sqs_clientc                 C   s¸   |d urB| j rB|| jv r| j| S || j vrtd |¡ƒ‚| j | }| j| d| j¡| d| jj¡| d| jj	¡d� }| j|< |S | j
d urJ| j
S | j| j| jj| jj	d� }| _
|S )Nrj   rÝ   rÞ   rß   )rÝ   rÞ   rß   )r?   Ú_predefined_queue_clientsr(   rq   rá   r4   rÝ   ÚconninfoÚuseridÚpasswordÚ_sqs©r5   rM   rF   r   r   r   r   rB     s.   


ý
ý
ýzChannel.sqsc                 C   s    |d ur8| j r8|| jv r| j| S || j vrtd |¡ƒ‚| j | }t| j|d�| d| j¡d� }| j|< |S | jd ur@| jS t| j|d�| jd� }| _|S )Nrj   rw   rÝ   )Úsqs_connectionrÝ   )	r?   Ú_predefined_queue_async_clientsr(   rq   r   rB   r4   rÝ   Ú_asynsqsrç   r   r   r   r˜   &  s*   


ý

þ

þzChannel.asynsqsc                 C   s   | j jS rH   )rª   rÜ   r_   r   r   r   rã   ?  s   zChannel.conninfoc                 C   s
   | j jjS rH   )rª   rÜ   r{   r_   r   r   r   r{   C  ó   
zChannel.transport_optionsc                 C   s   | j  d¡p| jS )Nrs   )r{   r4   Údefault_visibility_timeoutr_   r   r   r   rs   G  s   ÿzChannel.visibility_timeoutc                 C   ó   | j  dd¡S )z/Map of queue_name to predefined queue settings.r?   N©r{   r4   r_   r   r   r   r?   L  s   zChannel.predefined_queuesc                 C   rí   )Nr3   Ú rî   r_   r   r   r   r3   Q  s   zChannel.queue_name_prefixc                 C   s   dS )NFr   r_   r   r   r   Úsupports_fanoutU  s   zChannel.supports_fanoutc                 C   s   | j  d¡pt ¡ jp| jS )NrÝ   )r{   r4   r	   rÚ   rÔ   Údefault_regionr_   r   r   r   rÝ   Y  s
   ÿþzChannel.regionc                 C   ó   | j  d¡S )NÚ
regioninforî   r_   r   r   r   ró   _  ó   zChannel.regioninfoc                 C   rò   )NrÛ   rî   r_   r   r   r   rÛ   c  rô   zChannel.is_securec                 C   rò   ©NÚportrî   r_   r   r   r   rö   g  rô   zChannel.portc                 C   sP   | j jd ur&| jrdnd}| j jd urd | j j¡}nd}d || j j|¡S d S )NÚhttpsÚhttpz:{}rï   z	{}://{}{})rã   ÚhostnamerÛ   rö   rq   )r5   Úschemerö   r   r   r   rØ   k  s   ýúzChannel.endpoint_urlc                 C   s   | j  d| j¡S )Nr¨   )r{   r4   Údefault_wait_time_secondsr_   r   r   r   r¨   y  s   ÿzChannel.wait_time_seconds)NNrH   )r   N)r   NN)r’   )F)?r)   r*   r+   r,   rñ   rì   rû   Údomain_formatrê   ré   ræ   râ   rA   ÚsetrI   r1   r2   rL   rR   rZ   r`   ÚCHARS_REPLACE_TABLErh   ri   rv   rt   r}   rŽ   r›   r    ÚSQS_MAX_MESSAGESr\   r¯   rK   r±   r¦   r´   r¹   r¼   r»   rÄ   rÈ   rÎ   rÒ   rÓ   rá   rB   r˜   Úpropertyrã   r{   r   rs   r?   r3   rð   rÝ   ró   rÛ   rö   rØ   r¨   Ú__classcell__r   r   r8   r   r-   o   sŽ    	
	
ÿ)
	
ÿ


ÿÿ		












r-   c                   @   sp   e Zd ZdZeZdZdZdZej	j
ejejf Z
ej	jejf ZdZdZej	jjdedgƒd�Zed	d
„ ƒZdS )Ú	Transporta  SQS Transport.

    Additional queue attributes can be supplied to SQS during queue
    creation by passing an ``sqs-creation-attributes`` key in
    transport_options. ``sqs-creation-attributes`` must be a dict whose
    key-value pairs correspond with Attributes in the
    `CreateQueue SQS API`_.

    For example, to have SQS queues created with server-side encryption
    enabled using the default Amazon Managed Customer Master Key, you
    can set ``KmsMasterKeyId`` Attribute. When the queue is initially
    created by Kombu, encryption will be enabled.

    .. code-block:: python

        from kombu.transport.SQS import Transport

        transport = Transport(
            ...,
            transport_options={
                'sqs-creation-attributes': {
                    'KmsMasterKeyId': 'alias/aws/sqs',
                },
            }
        )

    .. _CreateQueue SQS API: https://docs.aws.amazon.com/AWSSimpleQueueService/latest/APIReference/API_CreateQueue.html#API_CreateQueue_RequestParameters
    r   r   NrB   TÚdirect)ÚasynchronousÚexchange_typec                 C   s
   d| j iS rõ   )Údefault_portr_   r   r   r   Údefault_connection_params±  rë   z#Transport.default_connection_params)r)   r*   r+   r,   r-   Úpolling_intervalr¨   r  r   r  Úconnection_errorsr
   ÚBotoCoreErrorÚsocketÚerrorÚchannel_errorsÚdriver_typeÚdriver_nameÚ
implementsÚextendÚ	frozensetr   r  r   r   r   r   r    s(    
ÿÿÿþr  )4r,   Ú
__future__r   r   r•   r  Ústringrˆ   Úbotocore.exceptionsr   Úviner   r   r   Úkombu.asynchronousr   Úkombu.asynchronous.aws.extr	   r
   Ú%kombu.asynchronous.aws.sqs.connectionr   Ú"kombu.asynchronous.aws.sqs.messager   Ú
kombu.fiver   r   r   r   Ú	kombu.logr   Úkombu.utilsr   Úkombu.utils.encodingr   r   Úkombu.utils.jsonr   r   Úkombu.utils.objectsr   rï   r   r)   ÚloggerÚpunctuationrþ   rÿ   r'   Ú	Exceptionr(   r-   r  r   r   r   r   Ú<module>   sB    >ÿ    