o
    wvXjD  ã                   @   sÊ   d Z ddl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 ddlmZ d	d
lmZ zddlZddlmZ W n eyO   d ZZY nw dZeddƒZeeƒZG dd„ deƒZdS )z"AWS DynamoDB result store backend.é    )Úabsolute_importÚunicode_literals)Ú
namedtuple)ÚsleepÚtime)Ú
_parse_url)ÚImproperlyConfigured)Ústring)Ú
get_loggeré   )ÚKeyValueStoreBackendN)ÚClientError)ÚDynamoDBBackendÚDynamoDBAttribute©ÚnameÚ	data_typec                       s  e Zd ZdZdZdZdZdZdZdZ	dZ
eddd�Zed	d
d�Zeddd�Zeddd�ZdZd3‡ fdd„	Zd3dd„Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zd4d!d"„Zd#d$„ Zd%d&„ Zd'd(„ Zed)d*„ ƒZd+d,„ Z d-d.„ Z!d/d0„ Z"d1d2„ Z#‡  Z$S )5r   z”AWS DynamoDB result backend.

    Raises:
        celery.exceptions.ImproperlyConfigured:
            if module :pypi:`boto3` is not available.
    Úceleryr   NTÚidÚSr   ÚresultÚBÚ	timestampÚNÚttlc              
      sŒ  t t| ƒj|i |¤Ž || _|p| j| _tstdƒ‚d}d }d }|d ur­t|ƒ\}}	}
}}}}|}|}|d u}|d u}||krCtdƒ‚|}|	dkr\d |
¡| _	d| _
t d | j	¡¡ n|	| _
| jjj}|dƒ}|rm|| _	t| d	| j¡ƒ| _t| d
| j¡ƒ| _| d| j¡}|r§zt|ƒ| _W n ty¦ } z	tjd|d� |‚d }~ww |p«| j| _| j| j| jf| _d | _|rÄ| j||d� d S d S )NzBYou need to install the boto3 library to use the DynamoDB backend.Fz6You need to specify both the Access Key ID and Secret.Ú	localhostzhttp://localhost:{}z	us-east-1z*Using local-only DynamoDB endpoint URL: {}Údynamodb_endpoint_urlÚreadÚwriteÚttl_secondsz!TTL must be a number; got "{ttl}")Úexc_info)Úaccess_key_idÚsecret_access_key)Úsuperr   Ú__init__ÚurlÚ
table_nameÚboto3r   Ú	parse_urlÚformatÚendpoint_urlÚ
aws_regionÚloggerÚwarningÚappÚconfÚgetÚintÚread_capacity_unitsÚwrite_capacity_unitsÚtime_to_live_secondsÚ
ValueErrorÚerrorÚ
_key_fieldÚ_value_fieldÚ_timestamp_fieldÚ_available_fieldsÚ_clientÚ_get_client)Úselfr%   r&   ÚargsÚkwargsÚaws_credentials_givenÚaws_access_key_idÚaws_secret_access_keyÚschemeÚregionÚportÚusernameÚpasswordÚtableÚqueryÚaccess_key_givenÚsecret_key_givenÚ_getÚconfig_endpoint_urlr   Úe©Ú	__class__© úU/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/backends/dynamodb.pyr$   B   sŽ   ÿÿÿÿÿ
þÿþÿþ€ûý
þÿzDynamoDBBackend.__init__c                 C   s~   | j du r<d| ji}|dur| ||dœ¡ | jdur | j|d< tj	di |¤Ž| _ |  ¡  |  ¡ dur<|  ¡  |  	¡  | j S )zGet client connection.NÚregion_name)rA   rB   r*   Údynamodb)rT   )
r;   r+   Úupdater*   r'   ÚclientÚ_get_or_create_tableÚ_has_ttlÚ_validate_ttl_methodsÚ_set_table_ttl)r=   r!   r"   Úclient_parametersrQ   rQ   rR   r<   ›   s(   
ÿþ

ÿþzDynamoDBBackend._get_clientc                 C   s6   | j j| j jdœg| j| j jddœg| j| jdœdœS )z=Get the boto3 structure describing the DynamoDB table schema.)ÚAttributeNameÚAttributeTypeÚHASH)r\   ÚKeyType)ÚReadCapacityUnitsÚWriteCapacityUnits)ÚAttributeDefinitionsÚ	TableNameÚ	KeySchemaÚProvisionedThroughput)r7   r   r   r&   r2   r3   ©r=   rQ   rQ   rR   Ú_get_table_schema¶   s   þÿþÿþòz!DynamoDBBackend._get_table_schemac              
   C   s¢   |   ¡ }z#| jjd
i |¤Ž}t d | j¡¡ |  d¡ t d | j¡¡ |W S  tyP } z|j	d  
dd¡}|dkrJ| jj| jd�W  Y d	}~S |‚d	}~ww )z=Create table if not exists, otherwise return the description.z*DynamoDB Table {} did not exist, creating.ÚACTIVEz#DynamoDB Table {} is now available.ÚErrorÚCodeÚUnknownÚResourceInUseException©rc   NrQ   )rg   r;   Úcreate_tabler,   Úinfor)   r&   Ú_wait_for_table_statusr   Úresponser0   Údescribe_table)r=   Útable_schemaÚtable_descriptionrN   Ú
error_coderQ   rQ   rR   rW   Ì   s0   ÿÿ
ÿÿÿ€÷z$DynamoDBBackend._get_or_create_tablec                 C   s   | j du rdS | j dkS )zàReturn the desired Time to Live config.

        - True:  Enable TTL on the table; use expiry.
        - False: Disable TTL on the table; don't use expiry.
        - None:  Ignore TTL on the table; don't use expiry.
        Nr   )r4   rf   rQ   rQ   rR   rX   é   s   ÿzDynamoDBBackend._has_ttlc                 C   sb   d}g }t |ƒD ]}t| j|ƒs| |¡ q|r/t djd |¡d�¡ tdjd |¡d�ƒ‚dS )z:Verify boto support for the DynamoDB Time to Live methods.)Úupdate_time_to_liveÚdescribe_time_to_livezdboto3 method(s) {methods} not found; ensure that boto3>=1.9.178 and botocore>=1.12.178 are installedú,)Úmethodsz#boto3 method(s) {methods} not foundN)	ÚlistÚhasattrr;   Úappendr,   r6   r)   ÚjoinÚAttributeError)r=   Úrequired_methodsÚmissing_methodsÚmethodrQ   rQ   rR   rY   ó   s&   
€üÿÿÿ÷z%DynamoDBBackend._validate_ttl_methodsc                 C   s   | j |  ¡ |dœdœS )zBGet the boto3 structure describing the DynamoDB TTL specification.)ÚEnabledr\   )rc   ÚTimeToLiveSpecification)r&   rX   )r=   Úttl_attr_namerQ   rQ   rR   Ú_get_ttl_specification  s
   þþz&DynamoDBBackend._get_ttl_specificationc              
   C   sp   z| j j| jd�}W |S  ty7 } z |jd  dd¡}|jd  dd¡}t dj| j||d�¡ |‚d }~ww )Nrm   ri   rj   rk   ÚMessagezJError describing Time to Live on DynamoDB table {table}: {code}: {message})rH   ÚcodeÚmessage)	r;   rw   r&   r   rq   r0   r,   r6   r)   )r=   ÚdescriptionrN   ru   Úerror_messagerQ   rQ   rR   Ú_get_table_ttl_description  s$   ÿóú€õz*DynamoDBBackend._get_table_ttl_descriptionc           	      C   sn  |   ¡ }|d d }|dv r2|d d }|  ¡ r1|| jjkr1t dj|dkr(dnd| jd	�¡ |S n'|d
v rN|  ¡ sMt dj|dkrDdnd| jd	�¡ |S nt dj|| jd�¡ |dkr_|n| jj}z | j	j
di | j|d�¤Ž}t dj| j|  ¡ | jjd�¡ |W S  ty¶ } z'|jd  dd¡}|jd  dd¡}t dj|  ¡ r§dnd| j||d�¡ |‚d}~ww )z,Enable or disable Time to Live on the table.ÚTimeToLiveDescriptionÚTimeToLiveStatus)ÚENABLEDÚENABLINGr\   z5DynamoDB Time to Live is {situation} on table {table}rŽ   zalready enabledzcurrently being enabled)Ú	situationrH   )ÚDISABLEDÚ	DISABLINGr‘   zalready disabledzcurrently being disabledzWUnknown DynamoDB Time to Live status {status} on table {table}. Attempting to continue.)ÚstatusrH   )r„   zUDynamoDB table Time to Live updated: table={table} enabled={enabled} attribute={attr})rH   ÚenabledÚattrri   rj   rk   r†   zHError {action} Time to Live on DynamoDB table {table}: {code}: {message}ÚenablingÚ	disabling)ÚactionrH   r‡   rˆ   NrQ   )r‹   rX   Ú
_ttl_fieldr   r,   Údebugr)   r&   r-   r;   rv   r…   ro   r   rq   r0   r6   )	r=   r‰   r“   Úcur_attr_nameÚ	attr_nameÚspecificationrN   ru   rŠ   rQ   rQ   rR   rZ   /  s„   
ÿÿù	€ÿù	ôû&ÿ
ÿÿúÿ
ù	€ôzDynamoDBBackend._set_table_ttlrh   c                 C   sT   d}|s(| j j| jd�}t d | j|¡¡ |d d }||k}tdƒ |rdS dS )z#Poll for the expected table status.Frm   z+Waiting for DynamoDB table {} to become {}.ÚTableÚTableStatusr   N)rV   rr   r&   r,   rš   r)   r   )r=   ÚexpectedÚachieved_statert   Úcurrent_statusrQ   rQ   rR   rp   Ÿ  s   ÿþÿôz&DynamoDBBackend._wait_for_table_statusc                 C   s   | j | jj| jj|iidœS )z0Construct the item retrieval request parameters.)rc   ÚKey)r&   r7   r   r   )r=   ÚkeyrQ   rQ   rR   Ú_prepare_get_request°  s   ÿÿþz$DynamoDBBackend._prepare_get_requestc              	   C   s~   t ƒ }| j| jj| jj|i| jj| jj|i| jj| jjt|ƒiidœ}|  ¡ r=|d  	| j
j| j
jtt|| j ƒƒii¡ |S )z/Construct the item creation request parameters.)rc   ÚItemr¦   )r   r&   r7   r   r   r8   r9   ÚstrrX   rU   r™   r1   r4   )r=   r¤   Úvaluer   Úput_requestrQ   rQ   rR   Ú_prepare_put_request»  s*   ÿÿÿùþþÿz$DynamoDBBackend._prepare_put_requestc                    s    dˆ vri S ‡ fdd„| j D ƒS )z1Convert get_item() response to field-value pairs.r¦   c                    s$   i | ]}|j ˆ d  |j  |j “qS )r¦   r   )Ú.0Úfield©Úraw_responserQ   rR   Ú
<dictcomp>Ù  s    ÿÿz1DynamoDBBackend._item_to_dict.<locals>.<dictcomp>)r:   )r=   r®   rQ   r­   rR   Ú_item_to_dictÕ  s
   
þzDynamoDBBackend._item_to_dictc                 C   s   |   ¡ S )N)r<   rf   rQ   rQ   rR   rV   Þ  s   zDynamoDBBackend.clientc                 C   s<   t |ƒ}|  |¡}| jjdi |¤Ž}|  |¡}| | jj¡S ©NrQ   )r	   r¥   rV   Úget_itemr°   r0   r8   r   )r=   r¤   Úrequest_parametersÚitem_responseÚitemrQ   rQ   rR   r0   â  s
   

zDynamoDBBackend.getc                 C   s*   t |ƒ}|  ||¡}| jjdi |¤Ž d S r±   )r	   rª   rV   Úput_item)r=   r¤   r¨   Ústater³   rQ   rQ   rR   Úseté  s   zDynamoDBBackend.setc                    s   ‡ fdd„|D ƒS )Nc                    s   g | ]}ˆ   |¡‘qS rQ   )r0   )r«   r¤   rf   rQ   rR   Ú
<listcomp>ï  s    z(DynamoDBBackend.mget.<locals>.<listcomp>rQ   )r=   ÚkeysrQ   rf   rR   Úmgetî  s   zDynamoDBBackend.mgetc                 C   s(   t |ƒ}|  |¡}| jjdi |¤Ž d S r±   )r	   r¥   rV   Údelete_item)r=   r¤   r³   rQ   rQ   rR   Údeleteñ  s   
zDynamoDBBackend.delete)NN)rh   )%Ú__name__Ú
__module__Ú__qualname__Ú__doc__r&   r2   r3   r+   r*   r4   Úsupports_autoexpirer   r7   r8   r9   r™   r:   r$   r<   rg   rW   rX   rY   r…   r‹   rZ   rp   r¥   rª   r°   ÚpropertyrV   r0   r¸   r»   r½   Ú__classcell__rQ   rQ   rO   rR   r      sB    
Y


p	
r   )rÁ   Ú
__future__r   r   Úcollectionsr   r   r   Úkombu.utils.urlr   r(   Úcelery.exceptionsr   Úcelery.fiver	   Úcelery.utils.logr
   Úbaser   r'   Úbotocore.exceptionsr   ÚImportErrorÚ__all__r   r¾   r,   r   rQ   rQ   rQ   rR   Ú<module>   s&   ÿ
