o
    —¨Êh•  ã                   @  s¦   d Z ddl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 dd	lmZ dd
lmZ ddlmZ er>ddlmZ dd„ ZG dd„ de	ƒZG dd„ dƒZdS )z%Generic resource pool implementation.é    )ÚannotationsN)Údeque)ÚEmpty)Ú	LifoQueue)ÚTYPE_CHECKINGé   )Ú
exceptions)Úregister_after_fork)Úlazy)ÚTracebackTypec                 C  s$   z|   ¡  W d S  ty   Y d S w ©N)Úforce_close_allÚ	Exception)Úresource© r   ú@/var/www/html/env/lib/python3.10/site-packages/kombu/resource.pyÚ_after_fork_cleanup_resource   s
   ÿr   c                   @  s   e Zd ZdZdd„ ZdS )r   z#Last in first out version of Queue.c                 C  s   t ƒ | _d S r   )r   Úqueue)ÚselfÚmaxsizer   r   r   Ú_init   ó   zLifoQueue._initN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r   r   r   r      s    r   c                   @  sÒ   e Zd ZdZejZdZd&dd„Zdd„ Zdd	„ Z	d'd
d„Z
dd„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zd(dd„Zd)dd„Zd(dd„Zedd „ ƒZejd!d „ ƒZej d"¡rge
ZeZd#Zd$d„ Z
d%d„ ZdS dS )*ÚResourcezPool of resources.FNc                 C  s^   || _ |pd| _d| _|d ur|n| j| _tƒ | _tƒ | _| jr)td ur)t| t	ƒ |  
¡  d S )Nr   F)Ú_limitÚpreloadÚ_closedÚclose_after_forkr   Ú	_resourceÚsetÚ_dirtyr	   r   Úsetup)r   Úlimitr   r    r   r   r   Ú__init__(   s   
ÿþ
zResource.__init__c                 C  s   t dƒ‚)Nzsubclass responsibility)ÚNotImplementedError©r   r   r   r   r$   7   s   zResource.setupc                 C  s6   | j rt| jƒ| j kr|  | j ¡‚| j |  ¡ ¡ d S r   )r%   Úlenr#   ÚLimitExceededr!   Ú
put_nowaitÚnewr(   r   r   r   Ú_add_when_empty:   s   zResource._add_when_emptyc                   sÀ   ˆj rtdƒ‚ˆjrM	 z
ˆjj||d�‰ W n ty"   ˆ ¡  Y n)w zˆ ˆ ¡‰ W n tyC   t	ˆ t
ƒr=ˆj ˆ ¡ ‚ ˆ ˆ ¡ ‚ w ˆj ˆ ¡ nqnˆ ˆ ¡ ¡‰ ‡ ‡fdd„}|ˆ _ˆ S )a–  Acquire resource.

        Arguments:
        ---------
            block (bool): If the limit is exceeded,
                then block until there is an available item.
            timeout (float): Timeout to wait
                if ``block`` is true.  Default is :const:`None` (forever).

        Raises
        ------
            LimitExceeded: if block is false and the limit has been exceeded.
        zAcquire on closed poolr   )ÚblockÚtimeoutc                     s   ˆ  ˆ ¡ dS )a'  Release resource so it can be used by another thread.

            Warnings:
            --------
                The caller is responsible for discarding the object,
                and to never use the resource again.  A new resource must
                be acquired if so needed.
            N)Úreleaser   ©ÚRr   r   r   r0   h   s   	z!Resource.acquire.<locals>.release)r   ÚRuntimeErrorr%   r!   Úgetr   r-   ÚprepareÚBaseExceptionÚ
isinstancer
   r+   r0   r#   Úaddr,   )r   r.   r/   r0   r   r1   r   ÚacquireB   s4   ÿ

ÿùï
zResource.acquirec                 C  s   |S r   r   ©r   r   r   r   r   r5   v   ó   zResource.preparec                 C  s   |  ¡  d S r   )Úcloser:   r   r   r   Úclose_resourcey   r   zResource.close_resourcec                 C  ó   d S r   r   r:   r   r   r   Úrelease_resource|   r;   zResource.release_resourcec                 C  s    | j r	| j |¡ |  |¡ dS )zqReplace existing resource with a new instance.

        This can be used in case of defective resources.
        N)r%   r#   Údiscardr=   r:   r   r   r   Úreplace   s   zResource.replacec                 C  s:   | j r| j |¡ | j |¡ |  |¡ d S |  |¡ d S r   )r%   r#   r@   r!   r+   r?   r=   r:   r   r   r   r0   ˆ   s
   zResource.releasec                 C  r>   r   r   r:   r   r   r   Úcollect_resource�   r;   zResource.collect_resourceTc                 C  s¬   | j rdS || _ | j}| j}	 z| ¡ }W n	 ty   Y nw z|  |¡ W n	 ty/   Y nw q	 z|j ¡ }W n
 tyC   Y dS w z|  |¡ W n	 tyT   Y nw q2)aa  Close and remove all resources in the pool (also those in use).

        Used to close resources from parent processes after fork
        (e.g. sockets/connections).

        Arguments:
        ---------
            close_pool (bool): If True (default) then the pool is marked
                as closed. In case of False the pool can be reused.
        N)	r   r#   r!   ÚpopÚKeyErrorrB   ÚAttributeErrorr   Ú
IndexError)r   Ú
close_poolÚdirtyr   ÚdresÚresr   r   r   r   “   s:   ÿÿù	ÿÿözResource.force_close_allc                 C  sš   | j }| jr"d|  k r| j k r"n n|s"|s td | j |¡ƒ‚d}|| _ |r9z| jdd� W n	 ty8   Y nw |  ¡  ||k rK| j|dkd� d S d S )Nr   z,Can't shrink pool when in use: was={} now={}TF)rG   )Úcollect)r   r#   r3   Úformatr   r   r$   Ú_shrink_down)r   r%   ÚforceÚignore_errorsÚresetÚ
prev_limitr   r   r   Úresize¹   s(   $ÿÿÿÿzResource.resizec                 C  sØ   G dd„ dƒ}| j }t|d|ƒ ƒ�Q t|jƒrJt|jƒt| jƒ | jkrR|j ¡ }|r0|  |¡ t|jƒrZt|jƒt| jƒ | jks$W d   ƒ d S W d   ƒ d S W d   ƒ d S W d   ƒ d S 1 sew   Y  d S )Nc                   @  s   e Zd Zdd„ Zddd„ZdS )z#Resource._shrink_down.<locals>.Noopc                 S  r>   r   r   r(   r   r   r   Ú	__enter__Í   r;   z-Resource._shrink_down.<locals>.Noop.__enter__Úexc_typeÚtypeÚexc_valr   Úexc_tbr   ÚreturnÚNonec                 S  r>   r   r   )r   rT   rV   rW   r   r   r   Ú__exit__Ð   s   z,Resource._shrink_down.<locals>.Noop.__exit__N)rT   rU   rV   r   rW   r   rX   rY   )r   r   r   rS   rZ   r   r   r   r   ÚNoopÌ   s    r[   Úmutex)r!   Úgetattrr)   r   r#   r%   ÚpopleftrB   )r   rK   r[   r   r2   r   r   r   rM   Ë   s"   



üýþý"þzResource._shrink_downc                 C  s   | j S r   )r   r(   r   r   r   r%   â   s   zResource.limitc                 C  s   |   |¡ d S r   )rR   )r   r%   r   r   r   r%   æ   s   ÚKOMBU_DEBUG_POOLr   c                 O  s‚   dd l }| jd  }| _td|› d| jj› �ƒ | j|i |¤Ž}||_td|› d| jj› �ƒ t|dƒs7g |_|j 	| 
¡ ¡ |S )Nr   r   ú+z	 ACQUIRE ú-Úacquired_by)Ú	tracebackÚ_next_resource_idÚprintÚ	__class__r   Ú_orig_acquireÚ_resource_idÚhasattrrb   ÚappendÚformat_stack)r   ÚargsÚkwargsrc   ÚidÚrr   r   r   r9   ð   s   
c                 C  sR   |j }td|› d| jj› �ƒ |  |¡}td|› d| jj› �ƒ |  jd8  _|S )Nr`   z	 RELEASE ra   r   )rh   re   rf   r   Ú_orig_releaserd   )r   r   rn   ro   r   r   r   r0   ü   s   
)NNN)FN)T)FFF)r   r   r   r   r   r*   r    r&   r$   r-   r9   r5   r=   r?   rA   r0   rB   r   rR   rM   Úpropertyr%   ÚsetterÚosÚenvironr4   rg   rp   rd   r   r   r   r   r   !   s8    

4	

&


îr   )r   Ú
__future__r   rs   Úcollectionsr   r   r   r   Ú
_LifoQueueÚtypingr   Ú r   Úutils.compatr	   Úutils.functionalr
   Útypesr   r   r   r   r   r   r   Ú<module>   s    