o
    wvXjÕ  ã                   @   sÔ   d Z ddlmZm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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ejd dkr^dd„ ZneZG dd„ deƒZdS )z@Apache Cassandra result store backend using the DataStax driver.é    )Úabsolute_importÚunicode_literalsN)Ú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ˆ
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   úV/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/backends/cassandra.pyÚbuf_tA   s   
r   c                       sl   e Zd ZdZdZdZ		d‡ fdd„	Z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	   zöCassandra backend utilizing DataStax driver.

    Raises:
        celery.exceptions.ImproperlyConfigured:
            if module :pypi:`cassandra-driver` is not available,
            or if the :setting:`cassandra_servers` setting is not set.
    NTéR#  c                    sx  t t| ƒjdi |¤Ž tsttƒ‚| jj}|p| dd ¡| _	|p%| dd ¡| _
|p.| dd ¡| _|p7| dd ¡| _| di ¡| _| j	rI| jrI| jsMtdƒ‚|pT| dd ¡}|d ur^t |¡nd| _| d	¡pgd
}	| d¡pnd
}
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 | _d | _d S )NÚcassandra_serversÚcassandra_portÚcassandra_keyspaceÚcassandra_tableÚcassandra_optionsz!Cassandra backend not configured.Úcassandra_entry_ttlÚ Úcassandra_read_consistencyÚLOCAL_QUORUMÚcassandra_write_consistencyÚcassandra_auth_providerÚcassandra_auth_kwargsr   )Úsuperr	   Ú__init__Ú	cassandrar   ÚE_NO_CASSANDRAÚappÚconfÚgetÚserversÚportÚkeyspaceÚtabler   Ú	Q_EXPIRESÚformatÚ
cqlexpiresÚgetattrÚConsistencyLevelr   Úread_consistencyÚwrite_consistencyÚauth_providerÚauthÚ!E_NO_SUCH_CASSANDRA_AUTH_PROVIDERÚ_connectionÚ_sessionÚ_write_stmtÚ
_read_stmtÚ
_make_stmt)Úselfr%   r'   r(   Ú	entry_ttlr&   Úkwargsr#   ÚexpiresÚ	read_consÚ
write_consr0   Úauth_kwargsÚauth_provider_class©Ú	__class__r   r   r   U   sJ   ÿþþ
zCassandraBackend.__init__c                 C   s$   | j d ur
| j  ¡  d | _ d | _d S )N)r3   Úshutdownr4   )r8   r   r   r   Úprocess_cleanup„   s   


z CassandraBackend.process_cleanupFc                 C   s  | j durdS zltjj| jf| j| jdœ| j¤Ž| _ | j  | j	¡| _
tj tj| j| jd�¡| _| j| j_tj tj| jd�¡| _| j| j_|rqtj tj| jd�¡| _| j| j_z| j
 | j¡ W W dS  tjyp   Y W dS w W dS  tjyŒ   | j dur…| j  ¡  d| _ d| _
‚ w )zjPrepare the connection for action.

        Arguments:
            write (bool): are we a writer?
        N)r&   r0   )r(   r;   )r(   )r3   r    ÚclusterÚClusterr%   r&   r0   r   Úconnectr'   r4   ÚqueryÚSimpleStatementÚQ_INSERT_RESULTr*   r(   r+   r5   r/   Úconsistency_levelÚQ_SELECT_RESULTr6   r.   ÚQ_CREATE_RESULT_TABLEr7   ÚexecuteÚAlreadyExistsÚOperationTimedOutrB   )r8   Úwriter   r   r   Ú_get_connectionŠ   sP   
ÿþýÿÿ
ÿ
	ÿ
ÿð

ø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)rP   N)	rQ   r4   rM   r5   r   Úencoder"   ÚnowÚcurrent_task_children)r8   Útask_idÚresultÚstateÚ	tracebackÚrequestr:   r   r   r   Ú_store_resultÃ   s   

úzCassandraBackend._store_resultc                 C   s   dS )Nzcassandra://r   )r8   Ú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)ÚstatusrV   )rU   r]   rV   Ú	date_donerX   Úchildren)	rQ   r4   rM   r6   Úoner   ÚPENDINGÚmeta_from_decodedÚdecode)r8   rU   Úresr]   rV   r^   rX   r_   r   r   r   Ú_get_task_meta_forÔ   s   úz#CassandraBackend._get_task_meta_forr   c                    s6   |si n|}|  | j| j| jdœ¡ tt| ƒ ||¡S )N)r%   r'   r(   )Úupdater%   r'   r(   r   r	   Ú
__reduce__)r8   Úargsr:   r@   r   r   rg   ç   s   þÿzCassandraBackend.__reduce__)NNNNr   )F)NN)T)r   N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r%   Úsupports_autoexpirer   rC   rQ   rZ   r\   re   rg   Ú__classcell__r   r   r@   r   r	   G   s    	ÿ/
:
ÿ
r	   )rl   Ú
__future__r   r   ÚsysÚceleryr   Úcelery.exceptionsr   Úcelery.utils.logr   Úbaser   r    Úcassandra.authÚcassandra.clusterÚImportErrorÚ__all__ri   Úloggerr!   r2   rI   rK   rL   r)   Úversion_infor   Úbufferr	   r   r   r   r   Ú<module>   s4   ÿ
