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
mZ ddlmZ dd	lmZ d
d„ ZG dd„ deƒZG dd„ deƒZdS )z%Generic resource pool implementation.é    )Úabsolute_importÚunicode_literalsN)Údequeé   )Ú
exceptions)ÚEmptyÚ	LifoQueue)Úregister_after_fork)Úlazyc                 C   s$   z|   ¡  W d S  ty   Y d S w ©N)Úforce_close_allÚ	Exception)Úresource© r   úK/var/www/html/myproject/venv/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„ Zd(dd„Zd)dd„Zedd „ ƒZejd!d „ ƒZej d"¡rfe
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#   2   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_empty5   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 )av  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   r/   a   s   z!Resource.acquire.<locals>.release)r   ÚRuntimeErrorr$   r    Úgetr   r,   ÚprepareÚBaseExceptionÚ
isinstancer
   r*   r/   r"   Úaddr+   )r   r-   r.   r/   r   r0   r   Úacquire=   s4   ÿ

ÿùï	zResource.acquirec                 C   s   |S r   r   ©r   r   r   r   r   r4   n   ó   zResource.preparec                 C   s   |  ¡  d S r   )Úcloser9   r   r   r   Úclose_resourceq   r   zResource.close_resourcec                 C   ó   d S r   r   r9   r   r   r   Úrelease_resourcet   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<   r9   r   r   r   Úreplacew   s   zResource.replacec                 C   s:   | j r| j |¡ | j |¡ |  |¡ d S |  |¡ d S r   )r$   r"   r?   r    r*   r>   r<   r9   r   r   r   r/   €   s
   zResource.releasec                 C   r=   r   r   r9   r   r   r   Úcollect_resourceˆ   r:   zResource.collect_resourcec                 C   s¬   | j rdS d| _ | 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)z®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).
        NT)	r   r"   r    ÚpopÚKeyErrorrA   ÚAttributeErrorr   Ú
IndexError)r   Ú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}|| _ |r7z|  ¡  W n	 ty6   Y nw |  ¡  ||k rI| j|dkd� d S d S )Nr   z.Can't shrink pool when in use: was={0} now={1}T)Úcollect)r   r"   r2   Úformatr   r   r#   Ú_shrink_down)r   r$   ÚforceÚignore_errorsÚresetÚ
prev_limitr   r   r   Úresize¬   s(   $ÿÿÿÿzResource.resizeTc                 C   s�   G dd„ dƒ}| j }t|d|ƒ ƒ�- t|jƒ| jkr6|j ¡ }|r&|  |¡ t|jƒ| jksW d   ƒ d S W d   ƒ d S 1 sAw   Y  d S )Nc                   @   s   e Zd Zdd„ Z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__c                 S   r=   r   r   )r   ÚtypeÚvalueÚ	tracebackr   r   r   Ú__exit__Ã   r:   z,Resource._shrink_down.<locals>.Noop.__exit__N)r   r   r   rQ   rU   r   r   r   r   ÚNoop¿   s    rV   Úmutex)r    Úgetattrr(   r   r$   ÚpopleftrA   )r   rI   rV   r   r1   r   r   r   rK   ¾   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   )rP   )r   r$   r   r   r   r$   Ò   s   ÚKOMBU_DEBUG_POOLr   c                 O   sz   dd l }| jd  }| _td || jj¡ƒ | j|i |¤Ž}||_td || jj¡ƒ t|dƒs3g |_	|j	 
| ¡ ¡ |S )Nr   r   z+{0} ACQUIRE {1}z-{0} ACQUIRE {1}Úacquired_by)rT   Ú_next_resource_idÚprintrJ   Ú	__class__r   Ú_orig_acquireÚ_resource_idÚhasattrr[   ÚappendÚformat_stack)r   ÚargsÚkwargsrT   ÚidÚrr   r   r   r8   Ü   s   
c                 C   sJ   |j }td || jj¡ƒ |  |¡}td || jj¡ƒ |  jd8  _|S )Nz+{0} RELEASE {1}z-{0} RELEASE {1}r   )r`   r]   rJ   r^   r   Ú_orig_releaser\   )r   r   rf   rg   r   r   r   r/   è   s   
)NNN)FN)FFF)T)r   r   r   r   r   r)   r   r%   r#   r,   r8   r4   r<   r>   r@   r/   rA   r   rP   rK   Úpropertyr$   ÚsetterÚosÚenvironr3   r_   rh   r\   r   r   r   r   r      s8    

1	
!


îr   )r   Ú
__future__r   r   rk   Úcollectionsr   Ú r   Úfiver   r   Ú
_LifoQueueÚutils.compatr	   Úutils.functionalr
   r   Úobjectr   r   r   r   r   Ú<module>   s    