B
    à%Q[�b  ã               @   s  d Z ddlmZ ddlmZ ddlmZmZ ddlmZ ddl	m
Z
mZ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mZmZmZ ddlmZmZmZmZ ddlm Z  dZ!dZ"dZ#dZ$dZ%yddl&m'Z'm(Z( W n$ e)k
rü   ddl'm'Z'm(Z( Y nX G dd„ deƒZ*dS )Ú
é    )Údatetime)Úlinesep)ÚThreadÚLock)Úsleepé   )ÚRESTARTABLEÚget_config_parameterÚAUTO_BIND_NONEÚAUTO_BIND_NO_TLSÚAUTO_BIND_TLS_AFTER_BINDÚAUTO_BIND_TLS_BEFORE_BINDé   )ÚBaseStrategy)ÚConnectionUsage)Ú&LDAPConnectionPoolNameIsMandatoryErrorÚ!LDAPConnectionPoolNotStartedErrorÚLDAPOperationResultÚLDAPExceptionErrorÚLDAPResponseTimeoutError)ÚlogÚlog_enabledÚERRORÚBASIC)ÚLDAP_MAX_INTZTERMINATE_REUSABLE_CONNECTIONéÿÿÿÿéþÿÿÿéýÿÿÿéüÿÿÿ)ÚQueueÚEmptyc               @   s¼   e Zd ZdZeƒ Zdd„ Zdd„ Zdd„ Zdd	„ Z	d
d„ Z
G dd„ deƒZG dd„ deƒZG dd„ deƒZdd„ Zd'dd„Zdd„ Zdd„ Zd(dd„Zdd„ Zd)d!d"„Zd#d$„ Zd%d&„ ZdS )*ÚReusableStrategyaÊ  
    A pool of reusable SyncWaitRestartable connections with lazy behaviour and limited lifetime.
    The connection using this strategy presents itself as a normal connection, but internally the strategy has a pool of
    connections that can be used as needed. Each connection lives in its own thread and has a busy/available status.
    The strategy performs the requested operation on the first available connection.
    The pool of connections is instantiated at strategy initialization.
    Strategy has two customizable properties, the total number of connections in the pool and the lifetime of each connection.
    When lifetime is expired the connection is closed and will be open again when needed.
    c             C   s   t ‚d S )N)ÚNotImplementedError)Úself© r%   úZC:\Users\HIRONO~1\AppData\Local\Temp\pip-install-6i93dh7p\ldap3\ldap3\strategy\reusable.pyÚ	receivingA   s    zReusableStrategy.receivingc             C   s   t ‚d S )N)r#   )r$   r%   r%   r&   Ú_start_listenD   s    zReusableStrategy._start_listenc             C   s   t ‚d S )N)r#   )r$   Z
message_idr%   r%   r&   Ú_get_responseG   s    zReusableStrategy._get_responsec             C   s   t ‚d S )N)r#   )r$   r%   r%   r&   Ú
get_streamJ   s    zReusableStrategy.get_streamc             C   s   t ‚d S )N)r#   )r$   Úvaluer%   r%   r&   Ú
set_streamM   s    zReusableStrategy.set_streamc               @   sX   e 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
dd„ Zdd„ ZdS )zReusableStrategy.ConnectionPoolz6
        Container for the Connection Threads
        c             C   sš   |j tjkrŒtj|j  }|js2tj|j = t | ¡S |jrL|j|jkrL|j|_|jrf|j	|jkrf|j|_	|j
rˆ|j
|j
krˆ| ¡  |j
|_
|S t | ¡S d S )N)Ú	pool_namer"   ÚpoolsÚstartedÚobjectÚ__new__Úpool_keepaliveÚ	keepaliveÚpool_lifetimeÚlifetimeÚ	pool_sizeÚterminate_pool)ÚclsÚ
connectionÚpoolr%   r%   r&   r1   U   s    

z'ReusableStrategy.ConnectionPool.__new__c             C   s¸   t | dƒs´|j| _|| _g | _|jp*tdƒ| _|jp:tdƒ| _|j	| _
tƒ | _d| _d| _d| _tƒ | _d| _|jrztƒ nd | _d| _tƒ | _| tj| j< d| _ttƒr´ttd| ƒ d S )NÚworkersZREUSABLE_THREADED_POOL_SIZEZREUSABLE_THREADED_LIFETIMEFr   z!instantiated ConnectionPool: <%r>)Úhasattrr-   ÚnameÚmaster_connectionr;   r6   r
   r4   r5   r2   r3   r    Úrequest_queueÚ	open_poolÚ	bind_poolÚtls_poolÚdictÚ	_incomingÚcounterÚ_usager   Úterminated_usageÚ
terminatedr   Ú	pool_lockr"   r.   r/   r   r   r   )r$   r9   r%   r%   r&   Ú__init__f   s(    
z(ReusableStrategy.ConnectionPool.__init__c             C   s  dt | jƒ d | jrdnd }|dt t| jƒƒ 7 }|dt | jƒ 7 }|dt | jƒ 7 }|dt | jƒ 7 }|d	t | jƒ 7 }|d
t | j	ƒ 7 }|dt | j
ƒ t 7 }|dt | jƒ t 7 }|d7 }| j�rxFt| jƒD ]*\}}|tt |ƒ d¡ d t |ƒ 7 }qØW n|td 7 }|S )NzPOOL: z - status: r/   rH   z - responses in queue: z - pool size: z - lifetime: z - keepalive: z	 - open: z	 - bind: z - tls: zMASTER CONN: zWORKERS:é   z: z    no active workers in pool)Ústrr=   r/   ÚlenrD   r6   r5   r3   r@   rA   rB   r   r>   r;   Ú	enumerateÚrjust)r$   ÚsÚiÚworkerr%   r%   r&   Ú__str__|   s     (z'ReusableStrategy.ConnectionPool.__str__c             C   s   |   ¡ S )N)rS   )r$   r%   r%   r&   Ú__repr__�   s    z(ReusableStrategy.ConnectionPool.__repr__c          
   C   sH   xB| j D ]8}|j�( |jjjr(|jjjs0d|_nd|_W d Q R X qW d S )NTF)r;   Úworker_lockr9   ÚserverÚschemaÚinfoÚget_info_from_server)r$   rR   r%   r%   r&   rY   ’   s
    z4ReusableStrategy.ConnectionPool.get_info_from_serverc          
   C   sN   xH| j D ]>}|j�. |j | jj| jj| jj| jj| jj	¡ W d Q R X qW d S )N)
r;   rU   r9   Zrebindr>   ÚuserÚpasswordÚauthenticationÚsasl_mechanismÚsasl_credentials)r$   rR   r%   r%   r&   Úrebind_poolš   s    z+ReusableStrategy.ConnectionPool.rebind_poolc          
   C   sb   | j s^|  ¡  x*| jD ] }|j� |j ¡  W d Q R X qW d| _ d| _ttƒrZt	td| ƒ dS dS )NTFzworker started for pool <%s>)
r/   Úcreate_poolr;   rU   ÚthreadÚstartrH   r   r   r   )r$   rR   r%   r%   r&   Ú
start_pool£   s    z*ReusableStrategy.ConnectionPool.start_poolc                s2   t tƒrttdˆ ƒ ‡ fdd„tˆ jƒD ƒˆ _d S )Nzcreated pool <%s>c                s   g | ]}t  ˆ jˆ j¡‘qS r%   )r"   ÚPooledConnectionWorkerr>   r?   )Ú.0Ú_)r$   r%   r&   ú
<listcomp>³   s    z?ReusableStrategy.ConnectionPool.create_pool.<locals>.<listcomp>)r   r   r   Úranger6   r;   )r$   r%   )r$   r&   r`   °   s    z+ReusableStrategy.ConnectionPool.create_poolc             C   sˆ   | j s„ttƒrttd| ƒ d| _| j ¡  x4ttdd„ | j	D ƒƒƒD ]}| j 
td d d f¡ qDW | j ¡  d| _ ttƒr„ttd| ƒ d S )Nzterminating pool <%s>Fc             S   s   g | ]}|j  ¡ r|‘qS r%   )ra   Úis_alive)re   rR   r%   r%   r&   rg   »   s    zBReusableStrategy.ConnectionPool.terminate_pool.<locals>.<listcomp>Tzpool terminated for <%s>)rH   r   r   r   r/   r?   Újoinrh   rM   r;   ÚputÚTERMINATE_REUSABLE)r$   rf   r%   r%   r&   r7   µ   s    

z.ReusableStrategy.ConnectionPool.terminate_poolN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r1   rJ   rS   rT   rY   r_   rc   r`   r7   r%   r%   r%   r&   ÚConnectionPoolQ   s   	rq   c               @   s    e Zd ZdZdd„ Zdd„ ZdS )z'ReusableStrategy.PooledConnectionThreadz®
        The thread that holds the Reusable connection and receive operation request via the queue
        Result are sent back in the pool._incoming list when ready
        c             C   s4   t  | ¡ d| _|| _|| _ttƒr0ttd| ƒ d S )NTz)instantiated PooledConnectionThread: <%r>)r   rJ   ÚdaemonrR   r>   r   r   r   )r$   rR   r>   r%   r%   r&   rJ   Ç   s    
z0ReusableStrategy.PooledConnectionThread.__init__c             C   sN  d| j _d}| jjj}�xö|�sy$|jjd| jjjjd�\}}}}W n. tk
rr   | j j	j
sl| j j	 d¡ wY nX | j j��ˆ d| j _|tkrÚd}| j j	jrÖy"| j j	 ¡  ttƒr¾ttdƒ W n tk
rÔ   Y nX �nt ¡ | j j j| jjjjk�r@y| j j	 ¡  W n tk
�r    Y nX | j  ¡  ttƒ�r@ttdƒ |dk�rà|j�r¸| j j	j
�r¸| j j	jdd� |j�r’| j j	j�s’| j j	jdd� |j �rð| j j	j�sð| j j	j!dd� n8|j�rð| j j	j
�sð|j�rð| j j	j�sð| j j	jdd� | j j"�r|�r| j j	 #¡  d| j _"d }d }d }	yR|d	k�rJ| j j	 $| j j	 %|||¡¡}n| j j	 &| j j	 %|||¡¡}| j j	j'}	W n( t(k
�rš }
 z|
}W d d }
~
X Y nX |j)�8 |�r¼|d d f|j*|< n||	t+ ,|||¡f|j*|< W d Q R X d| j _|j -¡  | j  j.d
7  _.W d Q R X qW ttƒ�r$ttdƒ | jj/�rB| j0| j j	j/7  _0d| j _d S )NTF)ÚblockÚtimeoutr   zthread terminatedzthread respawn)ÚbindRequestÚunbindRequest)Úread_server_infoZsearchRequestr   )1rR   Úrunningr>   Ústrategyr:   r?   Úgetr3   r!   r9   ÚclosedZabandonrU   Úbusyrl   ÚboundÚunbindr   r   r   r   r   ÚnowÚcreation_timeÚsecondsr5   Únew_connectionr@   ÚopenrB   Ztls_startedÚ	start_tlsrA   ÚbindrY   Z_fire_deferredÚpost_send_searchÚsendÚpost_send_single_responseÚresultr   rI   rD   r   Zdecode_requestÚ	task_doneÚtask_counterÚusagerG   )r$   Ú	terminater:   rE   Úmessage_typeÚrequestÚcontrolsÚexcÚresponser‰   Úer%   r%   r&   ÚrunÐ   s€    

$


 




$



z+ReusableStrategy.PooledConnectionThread.runN)rm   rn   ro   rp   rJ   r”   r%   r%   r%   r&   ÚPooledConnectionThreadÂ   s   	r•   c               @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	z'ReusableStrategy.PooledConnectionWorkerz�
        Container for the restartable connection. it includes a thread and a lock to execute the connection in the pool
        c             C   sh   || _ || _d| _d| _d| _d | _d | _d| _|  ¡  t	 
| | j ¡| _tƒ | _ttƒrdttd| ƒ d S )NFr   z)instantiated PooledConnectionWorker: <%s>)r>   r?   rx   r|   rY   r9   r€   r‹   r‚   r"   r•   ra   r   rU   r   r   r   )r$   r9   r?   r%   r%   r&   rJ     s    z0ReusableStrategy.PooledConnectionWorker.__init__c             C   s’   dt | jƒ t d }|| jr"dnd7 }|d| jr6dnd 7 }|dd| j ¡   7 }|d	t | jjj	j
t ¡ | j j ƒ 7 }|d
t | jƒ 7 }|S )NzCONN: z       THREAD: rx   Zhaltedz - r|   Ú	availablezcreated at: z - time to live: z - requests served: )rL   r9   r   rx   r|   r€   Ú	isoformatr>   ry   r:   r5   r   r   r�   r‹   )r$   rP   r%   r%   r&   rS   *  s    (z/ReusableStrategy.PooledConnectionWorker.__str__c             C   sn  ddl m} t ¡ | _|| jjr(| jjn| jj| jj| jj	t
| jj| jjt| jj| jj| jj| jj| jj| jj| jj| jjd| jj| jj| jjd�| _| jj�rD| jjt
k�rDttƒrÄttd| jƒ | jjdd� | jjtkrî| jj dd� nV| jjt!k�r| jj"dd� | jj dd� n*| jjt#k�rD| jj dd� | jj"dd� | jj�rj| jj| j_| jj $| j¡ d S )Nr   )Ú
ConnectionF)rV   rZ   r[   Ú	auto_bindÚversionr\   Zclient_strategyÚauto_referralsÚ
auto_ranger]   r^   Úcheck_namesZcollect_usageÚ	read_onlyÚraise_exceptionsÚlazyÚfast_decoderÚreceive_timeoutZreturn_empty_attributesz"performing automatic bind for <%s>)rw   )%Zcore.connectionr˜   r   r   r€   r>   Zserver_poolrV   rZ   r[   r   rš   r\   r	   r›   rœ   r]   r^   r�   rF   rž   rŸ   r¡   r¢   Zempty_attributesr9   r™   r   r   r   rƒ   r   r…   r   r„   r   Ú
initialize)r$   r˜   r%   r%   r&   r‚   4  sH    

z6ReusableStrategy.PooledConnectionWorker.new_connectionN)rm   rn   ro   rp   rJ   rS   r‚   r%   r%   r%   r&   rd     s   
rd   c             C   s`   t  | |¡ d| _d| _d| _d| _t|dƒrB|jrBt 	|¡| _
nttƒrTttdƒ tdƒ‚d S )NFTr-   z)reusable connection must have a pool_name)r   rJ   ZsyncZno_real_dsaZpooledZ
can_streamr<   r-   r"   rq   r:   r   r   r   r   )r$   Zldap_connectionr%   r%   r&   rJ   _  s    
zReusableStrategy.__init__Tc             C   s@   d| j _| j  ¡  d| j_| jjr<|s0| jjjs<| jj ¡  d S )NTF)	r:   r@   rc   r9   r{   rŒ   rF   Zinitial_connection_start_timerb   )r$   Zreset_usagerw   r%   r%   r&   rƒ   l  s    
zReusableStrategy.openc             C   s6   | j  ¡  d| j _d| j_d| j_d| j _d| j _d S )NFT)r:   r7   r@   r9   r}   r{   rA   rB   )r$   r%   r%   r&   r�   u  s    
zReusableStrategy.terminatec             C   s&   d| j _| j jr"| j j jd7  _dS )z1
        Doesn't really close the socket
        Tr   N)r9   r{   rŒ   rF   Zclosed_sockets)r$   r%   r%   r&   Ú_close_socket}  s    zReusableStrategy._close_socketNc          	   C   sØ   | j jrº|dkrd| j _t}n˜|dkr4d| j _t}n‚|dkrBt}nt|dkr`| jjr`d| j _t	}nV| j j
�2 | j  jd7  _| j jtkrŽd| j _| j j}W d Q R X | j j ||||f¡ |S ttƒrÌttdƒ tdƒ‚d S )	Nru   Trv   FZabandonRequestZextendedReqr   z$reusable connection pool not started)r:   r/   rA   Ú
BOGUS_BINDÚBOGUS_UNBINDÚBOGUS_ABANDONr9   Ústarting_tlsrB   ÚBOGUS_EXTENDEDrI   rE   r   r?   rk   r   r   r   r   )r$   rŽ   r�   r�   rE   r%   r%   r&   r‡   †  s,    

zReusableStrategy.sendc             C   s"  | j j| jjjksZ| j j| jjjksZ| j j| jjjksZ| j j| jjjksZ| j j| jjjkrª| j j| jj_| j j| jj_| j j| jj_| j j| jj_| j j| jj_| j ¡  | jj	d j }d|_
| j jjrÒ| j jjsê| jj	d j j|d�}n| jj	d j j|dd�}| ¡  d|_
|�rd| j_|S )Nr   F)r�   )r�   rw   T)r9   rZ   r:   r>   r[   r\   r]   r^   r_   r;   r    rV   rW   rX   r…   r~   rA   )r$   r�   Ztemp_connectionr‰   r%   r%   r&   Úvalidate_bindŸ  s*    
zReusableStrategy.validate_bindFc          	   C   sn  t dƒ}d }|d krt dƒ}|tkrBtƒ }dd ddddd dœ}�n|tkrTd }d }nò|tkrztƒ }dd d	d
dddddœ}nÌ|tkr¨tƒ }dd d	d
dddddœ}d| j_nžd }d }xn|dk�ry4| jjj	j
� | jjj	j |¡\}}}W d Q R X W n( tk
�r   t|ƒ ||8 }w²Y nX P q²W |dk�rFttƒ�r>ttdƒ tdƒ‚t|tƒ�rV|‚|�rf|||fS ||fS )NZRESPONSE_SLEEPTIMEZRESPONSE_WAITING_TIMEOUTÚsuccessZbindResponser   Ú z<bogus Bind response>)ÚdescriptionÚ	referralsÚtyper‰   ÚdnÚmessageZ	saslCredsz1.3.6.1.4.1.1466.20037ZextendedRespÚNonez<bogus StartTls response>)r‰   r®   ZresponseNamer¯   r­   ZresponseValuer°   r±   Fz6no response from worker threads in Reusable connection)r
   r¥   Úlistr¦   r§   r©   r9   r¨   ry   r:   rI   rD   ÚpopÚKeyErrorr   r   r   r   r   Ú
isinstancer   )r$   rE   rt   Úget_requestZ	sleeptimer�   r’   r‰   r%   r%   r&   Úget_response¹  sJ    
&



zReusableStrategy.get_responsec             C   s   |S )Nr%   )r$   rE   r%   r%   r&   rˆ   å  s    z*ReusableStrategy.post_send_single_responsec             C   s   |S )Nr%   )r$   rE   r%   r%   r&   r†   è  s    z!ReusableStrategy.post_send_search)TT)N)NF)rm   rn   ro   rp   rC   r.   r'   r(   r)   r*   r,   r0   rq   r   r•   rd   rJ   rƒ   r�   r¤   r‡   rª   r¸   rˆ   r†   r%   r%   r%   r&   r"   5   s&   	qUH
		

,r"   N)+rp   r   Úosr   Ú	threadingr   r   Útimer   r¬   r	   r
   r   r   r   r   Úbaser   Z
core.usager   Zcore.exceptionsr   r   r   r   r   Z	utils.logr   r   r   r   Zprotocol.rfc4511r   rl   r¥   r¦   r©   r§   Úqueuer    r!   ÚImportErrorr"   r%   r%   r%   r&   Ú<module>   s(    