
    xai                        d 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l
mZ ddlmZ d	Z e	e      Zej"                  Z G d
 dej$                        Zy)zWorker Task Consumer Bootstep.    )annotations)QoSignore_errors)	bootsteps)
get_logger)detect_quorum_queues   )Mingle)Tasksc                  H     e Zd ZdZefZ fdZd Zd Zd Z	d Z
ddZ xZS )	r   z,Bootstep starting the task message consumer.c                B    d x|_         |_        t        |   |fi | y )N)task_consumerqossuper__init__)selfckwargs	__class__s      n/var/www/html/BankruptcyAI-uat/bankruptcy-ai/venv/lib/python3.12/site-packages/celery/worker/consumer/tasks.pyr   zTasks.__init__   s#    "&&!%%f%    c                L  	
 j                          | j                        	j                  j                  j	                  dj
                  	       j                  j                  j                  j                  j                        _
        	fd}j                  j                  j                  }t        |j
                  |      _        j                  j                  j                  rj                  j                   j"                  dk(  }|s8t$        j'                  dj                  j                   j"                   d       ydd	lm} dd
lm
 j                  j0                  j                  }|j2                  
fd} |||      |_        yy)zStart task consumer.r   )on_decode_errorc                >    j                   j                  |       S )N)prefetch_countapply_global)r   r   )r   r   
qos_globals    r   set_prefetch_countz'Tasks.start.<locals>.set_prefetch_count,   s%    ??&&-' '  r   )max_prefetchrediszWworker_disable_prefetch is only supported for Redis brokers. Current broker transport: z$. Ignoring disable_prefetch setting.N)
MethodType)statec                    t        j                  dd       xs j                  j                  }t	        j
                        |k\  ry        S )Nmax_concurrencyF)getattr
controllerpoolnum_processeslenreserved_requests)r   limitr   original_can_consumer"   s     r   can_consumez Tasks.start.<locals>.can_consumeG   sD    .?F^!&&J^J^u../58 +--r   )update_strategiesr   
connectiondefault_channel	basic_qosinitial_prefetch_countappamqpTaskConsumerr   r   confworker_eta_task_limitr   r   worker_disable_prefetch	transportdriver_typeloggerwarningtypesr!   celery.workerr"   channelr-   )r   r   r   eta_task_limitis_redis_brokerr!   channel_qosr-   r,   r   r"   s    `      @@@r   startzTasks.start   sP   	__Q'
 	
$$..q''	
 %%**11LL!*;*; 2 
	
 99 8 8~
 55::--ll44@@GKO"1121G1G1S1S0T U9:
 (+//1155K#.#:#: . '1k&JK#1 .r   c                t    |j                   r,t        d       t        ||j                   j                         yy)zStop task consumer.zCanceling task consumer...N)r   debugr   cancelr   r   s     r   stopz
Tasks.stopP   s+    ??./!Q__334 r   c                    |j                   rD| j                  |       t        d       t        ||j                   j                         d|_         yy)zShutdown task consumer.zClosing consumer channel...N)r   rH   rE   r   closerG   s     r   shutdownzTasks.shutdownV   s=    ??IIaL/0!Q__223"AO	 r   c                P    d|j                   r|j                   j                  iS diS )zReturn task consumer info.r   zN/A)r   valuerG   s     r   infoz
Tasks.info^   s#     !%%++BBEBBr   c                   |j                   j                   }|j                  j                  j                  rPt        |j                  |j                   j                  j                        \  }}|rd}t        j                  d       |S )zDetermine if global QoS should be applied.

        Additional information:
            https://www.rabbitmq.com/docs/consumer-prefetch
            https://www.rabbitmq.com/docs/quorum-queues#global-qos
        Fz5Global QoS is disabled. Prefetch count in now static.)
r/   qos_semantics_matches_specr3   r6   worker_detect_quorum_queuesr   r9   r:   r;   rN   )r   r   r   using_quorum_queues_s        r   r   zTasks.qos_globalb   sm     @@@
55::11%9q||--99&" #"
STr   )returnbool)__name__
__module____qualname____doc__r
   requiresr   rC   rH   rK   rN   r   __classcell__)r   s   @r   r   r      s.    6yH&1Kf5#Cr   r   N)rY   
__future__r   kombu.commonr   r   celeryr   celery.utils.logr   celery.utils.quorum_queuesr   mingler
   __all__rV   r;   rE   StartStopStepr    r   r   <module>re      sH    $ " +  ' ; 
 
H	cI## cr   