o
    ›¨Êh.#  ã                   @   sÀ   d 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	 zddl
Z
ddlZ
ddlZ
ddlZ
W n ey;   dZ
Y nw dZeeƒZd	Zd
ZdZdZdZdZdZdZdd„ ZG dd„ de	ƒZdS )z@Apache Cassandra result store backend using the DataStax driver.é    N)Ústates)ÚImproperlyConfigured)Ú
get_loggeré   )ÚBaseBackend)ÚCassandraBackendz
You need to install the cassandra-driver library to
use the Cassandra backend.  See https://github.com/datastax/python-driver
z�
CASSANDRA_AUTH_PROVIDER you provided is not a valid auth_provider class.
See https://datastax.github.io/python-driver/api/cassandra/auth.html.
z(Cassandra backend improperly configured.z!Cassandra backend not configured.zˆ
INSERT INTO {table} (
    task_id, status, result, date_done, traceback, children) VALUES (
        %s, %s, %s, %s, %s, %s) {expires};
z]
SELECT status, result, date_done, traceback, children
FROM {table}
WHERE task_id=%s
LIMIT 1
zà
CREATE TABLE {table} (
    task_id text,
    status text,
    result blob,
    date_done timestamp,
    traceback blob,
    children blob,
    PRIMARY KEY ((task_id), date_done)
) WITH CLUSTERING ORDER BY (date_done DESC);
z
    USING TTL {0}
c                 C   s
   t | dƒS )NÚutf8)Úbytes)Úx© r   úK/var/www/html/env/lib/python3.10/site-packages/celery/backends/cassandra.pyÚbuf_tC   s   
r   c                       sh   e Zd ZdZdZdZdZ		d‡ fdd„	Zddd	„Z	dd
d„Z	ddd„Z
dd„ Zd‡ fdd„	Z‡  ZS )r   aG  Cassandra/AstraDB backend utilizing DataStax driver.

    Raises:
        celery.exceptions.ImproperlyConfigured:
            if module :pypi:`cassandra-driver` is not available,
            or not-exactly-one of the :setting:`cassandra_servers` and
            the :setting:`cassandra_secure_bundle_path` settings is set.
    NTéR#  c                    s¨  t ƒ jdi |¤Ž tsttƒ‚| jj}|p| dd ¡| _|p#| dd ¡| _	|p,| dd ¡| _
|p5| dd ¡| _|p>| dd ¡| _| di ¡| _| jpL| j	}	|	rU| jrU| jsYttƒ‚| jrc| j	rcttƒ‚|pj| dd ¡}
|
d urtt |
¡nd| _| d	¡p}d
}| d¡p„d
}ttj|tjjƒ| _ttj|tjjƒ| _d | _| dd ¡}| dd ¡}|rÁ|rÁttj|d ƒ}|s¹ttƒ‚|di |¤Ž| _d | _d | _d | _d | _t  ¡ | _!d S )NÚcassandra_serversÚcassandra_secure_bundle_pathÚcassandra_portÚcassandra_keyspaceÚcassandra_tableÚcassandra_optionsÚcassandra_entry_ttlÚ Úcassandra_read_consistencyÚLOCAL_QUORUMÚcassandra_write_consistencyÚcassandra_auth_providerÚcassandra_auth_kwargsr   )"ÚsuperÚ__init__Ú	cassandrar   ÚE_NO_CASSANDRAÚappÚconfÚgetÚserversÚbundle_pathÚportÚkeyspaceÚtabler   ÚE_CASSANDRA_NOT_CONFIGUREDÚE_CASSANDRA_MISCONFIGUREDÚ	Q_EXPIRESÚformatÚ
cqlexpiresÚgetattrÚConsistencyLevelr   Úread_consistencyÚwrite_consistencyÚauth_providerÚauthÚ!E_NO_SUCH_CASSANDRA_AUTH_PROVIDERÚ_clusterÚ_sessionÚ_write_stmtÚ
_read_stmtÚ	threadingÚRLockÚ_lock)Úselfr#   r&   r'   Ú	entry_ttlr%   r$   Úkwargsr!   Údb_directionsÚexpiresÚ	read_consÚ
write_consr1   Úauth_kwargsÚauth_provider_class©Ú	__class__r   r   r   X   sV   ÿÿþþzCassandraBackend.__init__Fc                 C   sz  | j durdS | j ¡  zªzˆ| j durW W | j ¡  dS | jr2tjj| jf| j| j	dœ| j
¤Ž| _ntjjdd| ji| j	dœ| j
¤Ž| _| j | j¡| _ tj tj| j| jd�¡| _| j| j_tj tj| jd�¡| _| j| j_|r”tj tj| jd�¡}| j|_z| j  |¡ W n
 tjy“   Y nw W n tjy®   | jdur§| j ¡  d| _d| _ ‚ w W | j ¡  dS | j ¡  w )zjPrepare the connection for action.

        Arguments:
            write (bool): are we a writer?
        N)r%   r1   Úsecure_connect_bundle)Úcloudr1   )r'   r?   )r'   r   ) r5   r:   ÚacquireÚreleaser#   r   ÚclusterÚClusterr%   r1   r   r4   r$   Úconnectr&   ÚqueryÚSimpleStatementÚQ_INSERT_RESULTr+   r'   r,   r6   r0   Úconsistency_levelÚQ_SELECT_RESULTr7   r/   ÚQ_CREATE_RESULT_TABLEÚexecuteÚAlreadyExistsÚOperationTimedOutÚshutdown)r;   ÚwriteÚ	make_stmtr   r   r   Ú_get_connectionŽ   sl   


;Çÿþ
ýÿüûÿÿ
ÿ
	ÿÿ€

ø€
z CassandraBackend._get_connectionc                 K   sV   | j dd� | j | j||t|  |¡ƒ| j ¡ t|  |¡ƒt|  |  |¡¡ƒf¡ dS )z1Store return value and state of an executed task.T)rW   N)	rY   r5   rS   r6   r   Úencoder    ÚnowÚcurrent_task_children)r;   Útask_idÚresultÚstateÚ	tracebackÚrequestr=   r   r   r   Ú_store_resultÖ   s   

úzCassandraBackend._store_resultc                 C   s   dS )Nzcassandra://r   )r;   Úinclude_passwordr   r   r   Úas_uriä   s   zCassandraBackend.as_uric              
   C   sf   |   ¡  | j | j|f¡ ¡ }|stjddœS |\}}}}}|  |||  |¡||  |¡|  |¡dœ¡S )z$Get task meta-data for a task by id.N)Ústatusr^   )r]   re   r^   Ú	date_doner`   Úchildren)	rY   r5   rS   r7   Úoner   ÚPENDINGÚmeta_from_decodedÚdecode)r;   r]   Úresre   r^   rf   r`   rg   r   r   r   Ú_get_task_meta_forç   s   úz#CassandraBackend._get_task_meta_forr   c                    s2   |si n|}|  | j| j| jdœ¡ tƒ  ||¡S )N)r#   r&   r'   )Úupdater#   r&   r'   r   Ú
__reduce__)r;   Úargsr=   rD   r   r   ro   ú   s   þÿzCassandraBackend.__reduce__)NNNNr   N)F)NN)T)r   N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r#   r$   Úsupports_autoexpirer   rY   rb   rd   rm   ro   Ú__classcell__r   r   rD   r   r   G   s    
ÿ
6I
ÿ
r   )rt   r8   Úceleryr   Úcelery.exceptionsr   Úcelery.utils.logr   Úbaser   r   Úcassandra.authÚcassandra.clusterÚcassandra.queryÚImportErrorÚ__all__rq   Úloggerr   r3   r)   r(   rO   rQ   rR   r*   r   r   r   r   r   r   Ú<module>   s4    ÿ