Ë
    Bêñiñ!  ã                   ó  — d dl 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lmZ d dlmZ d d	lmZmZmZ d
ZdZ G d„ dee«      Z G d„ d«      Z	 ddee   dededededeeef   dz  ddfd„Z G d„ d«      Zy)é    N)Údefaultdict)ÚEnum)ÚQueueÚget_context)ÚBaseContext)ÚBaseProcess)ÚSynchronized)ÚEmpty)ÚAnyÚIterableÚTypeiX  éÈ   c                   ó   — e Zd ZdZdZdZy)ÚQueueSignalsÚstopÚconfirmÚerrorN)Ú__name__Ú
__module__Ú__qualname__r   r   r   © ó    úb/var/www/pod-logistic/pod-ai/venv/lib/python3.12/site-packages/qdrant_client/parallel_processor.pyr   r      s   „ Ø€DØ€GØ�Er   r   c                   óF   — e Zd Zedededd fd„«       Zdee   dee   fd„Zy)ÚWorkerÚargsÚkwargsÚreturnc                 ó   — t        «       ‚©N©ÚNotImplementedError)Úclsr   r   s      r   ÚstartzWorker.start   s   € ä!Ó#Ð#r   Úitemsc                 ó   — t        «       ‚r    r!   )Úselfr%   s     r   ÚprocesszWorker.process   s   € Ü!Ó#Ð#r   N)r   r   r   Úclassmethodr   r$   r   r(   r   r   r   r   r      sD   „ Øð$˜#ð $¨ð $°ò $ó ð$ð$˜X c™]ð $¨x¸©}ô $r   r   Úworker_classÚinput_queueÚoutput_queueÚnum_active_workersÚ	worker_idr   r   c                 óú  ‡— |€i }t        j                  d|› dt        j                  «       › �«       	  | j                  d	i |¤Ž}dt
        t           fˆfd„}|j                   |«       «      D ]  }|j                  |«       Œ 	 ‰j                  «        |j                  «        ‰j                  «        |j                  «        |j                  «       5  |xj                   dz  c_        ddd«       t        j                  d|› d�«       y# t        $ r>}	t        j                  |	«       |j                  t        j                  «       Y d}	~	ŒÊd}	~	ww xY w# 1 sw Y   ŒmxY w# ‰j                  «        |j                  «        ‰j                  «        |j                  «        |j                  «       5  |xj                   dz  c_        ddd«       n# 1 sw Y   nxY wt        j                  d|› d�«       w xY w)
zç
    A worker that pulls data pints off the input queue, and places the execution result on the output queue.
    When there are no data pints left on the input queue, it decrements
    num_active_workers to signal completion.
    NzReader worker: z PID: r   c               3   ó`   •K  — 	 ‰j                  «       } | t        j                  k(  ry | –— Œ)­wr    )Úgetr   r   )Úitemr+   s    €r   Úinput_queue_iterablez%_worker.<locals>.input_queue_iterable7   s1   øè ø€ ØØ"—‘Ó(�Øœ<×,Ñ,Ò,ØØ’
ð	 ùs   ƒ+.é   zReader worker z	 finishedr   )ÚloggingÚinfoÚosÚgetpidr$   r   r   r(   ÚputÚ	ExceptionÚ	exceptionr   r   ÚcloseÚjoin_threadÚget_lockÚvalue)
r*   r+   r,   r-   r.   r   Úworkerr3   Úprocessed_itemÚes
    `        r   Ú_workerrC   !   s²  ø€ ð €~Øˆä‡L�L�? 9 +¨V´B·I±I³K°=ÐAÔBð!<Ø#�×#Ñ#Ñ- fÑ-ˆð	¤h¬s¡mõ 	ð %Ÿn™nÑ-AÓ-CÓDò 	-ˆNØ×Ñ˜^Õ,ñ	-ð 	×ÑÔØ×ÑÔØ×ÑÔ!Ø× Ñ Ô"à×(Ñ(Ó*ñ 	*Ø×$Ò$¨Ñ)Õ$÷	*ô 	�‰�~ i [°	Ð:Õ;øô) ò -Ü×Ñ˜!ÔØ×Ñœ×+Ñ+×,Ñ,ûð-ú÷"	*ð 	*ûð 	×ÑÔØ×ÑÔØ×ÑÔ!Ø× Ñ Ô"à×(Ñ(Ó*ñ 	*Ø×$Ò$¨Ñ)Õ$÷	*÷ 	*ñ 	*úô 	�‰�~ i [°	Ð:Õ;úsU   ´AD ÂE$ ÃEÄ	EÄ4EÅE$ ÅEÅE$ ÅE!Å$AG:Æ5GÇ	G:ÇGÇ!G:c            	       óâ   — e Zd Zdefdedee   dedz  defd„Zde	ddfd	„Z
d
ee	   de	de	dee	   fd„Zd
ee	   de	de	dee	   fd„Zd
ee	   de	de	dee	   fd„Zdd„Zddedz  ddfd„Zdd„Zdd„Zy)ÚParallelWorkerPoolNÚnum_workersr@   Ústart_methodÚmax_internal_batch_sizec                 ó®   — || _         || _        d | _        d | _        t	        |«      | _        g | _        | j                  |z  | _        d| _        d | _	        y )NF)
r*   rF   r+   r,   r   ÚctxÚ	processesÚ
queue_sizeÚemergency_shutdownr-   )r'   rF   r@   rG   rH   s        r   Ú__init__zParallelWorkerPool.__init__X   sZ   € ð #ˆÔØ&ˆÔØ)-ˆÔØ*.ˆÔÜ +¨LÓ 9ˆŒØ,.ˆŒØ×*Ñ*Ð-DÑDˆŒØ"'ˆÔØ48ˆÕr   r   r   c                 ó   — | j                   j                  | j                  «      | _        | j                   j                  | j                  «      | _        | j                   j                  d| j                  «      }t        |t        «      sJ ‚|| _	        t        d| j                  «      D ]¢  }t        | j                   d«      sJ ‚| j                   j                  t        | j                  | j                  | j                  | j                  ||j                  «       f¬«      }|j!                  «        | j"                  j%                  |«       Œ¤ y )NÚir   ÚProcess)Útargetr   )rJ   r   rL   r+   r,   ÚValuerF   Ú
isinstanceÚ	BaseValuer-   ÚrangeÚhasattrrQ   rC   r*   Úcopyr$   rK   Úappend)r'   r   Ú	ctx_valuer.   r(   s        r   r$   zParallelWorkerPool.starti   s  € ØŸ8™8Ÿ>™>¨$¯/©/Ó:ˆÔØ ŸH™HŸN™N¨4¯?©?Ó;ˆÔà—H‘H—N‘N 3¨×(8Ñ(8Ó9ˆ	Ü˜)¤YÔ/Ð/Ð/Ø"+ˆÔä˜q $×"2Ñ"2Ó3ò 	+ˆIÜ˜4Ÿ8™8 YÔ/Ð/Ð/Ø—h‘h×&Ñ&Üà×%Ñ%Ø×$Ñ$Ø×%Ñ%Ø×+Ñ+ØØ—K‘K“Mðð 'ó 
ˆGð �M‰MŒOØ�N‰N×!Ñ! 'Õ*ñ	+r   Ústreamr   c              /   ó>  K  — 	  | j                   d	i |¤Ž | j                  €J d«       ‚| j                  €J d«       ‚d}d}|D ]º  }| j                  «        ||z
  | j                  k  r	 | j                  j                  «       }n!	 | j                  j                  t        ¬«      }|�7|t        j                  k(  r| j                  «        t        d«      ‚|–— |dz  }| j                  j                  |«       |dz  }Œ¼ t        | j                  «      D ]+  }	| j                  j                  t        j                   «       Œ- ||k  r]| j                  j                  t        ¬«      }|t        j                  k(  r| j                  «        t        d«      ‚|–— |dz  }||k  rŒ]| j                  €J d«       ‚| j                  €J d«       ‚| j#                  «        | j                  j%                  «        | j                  j%                  «        | j&                  r5| j                  j)                  «        | j                  j)                  «        y | j                  j+                  «        | j                  j+                  «        y # t        $ r d }Y �Œîw xY w# t        $ r}| j                  «        |‚d }~ww xY w# | j                  €J d«       ‚| j                  €J d«       ‚| j#                  «        | j                  j%                  «        | j                  j%                  «        | j&                  r5| j                  j)                  «        | j                  j)                  «        w | j                  j+                  «        | j                  j+                  «        w xY w­w)
NzInput queue was not initializedz Output queue was not initializedr   ©ÚtimeoutzThread unexpectedly terminatedr4   zInput queue is NonezOutput queue is Noner   )r$   r+   r,   Úcheck_worker_healthrL   Ú
get_nowaitr
   r1   Úprocessing_timeoutÚjoin_or_terminater   r   ÚRuntimeErrorr9   rV   rF   r   Újoinr<   rM   Úcancel_join_threadr=   )
r'   r[   r   r   ÚpushedÚreadr2   Úout_itemrB   Ú_s
             r   Úunordered_mapz ParallelWorkerPool.unordered_map�   sA  è ø€ ð4	0ØˆD�J‰JÑ ˜Ò à×#Ñ#Ð/ÐRÐ1RÓRÐ/Ø×$Ñ$Ð0ÐTÐ2TÓTÐ0àˆFØˆDØò �Ø×(Ñ(Ô*Ø˜D‘= 4§?¡?Ò2ð(Ø#'×#4Ñ#4×#?Ñ#?Ó#A™ð Ø#'×#4Ñ#4×#8Ñ#8ÔASÐ#8Ó#T˜ð
 Ð'Ø¤<×#5Ñ#5Ò5Ø×.Ñ.Ô0Ü*Ð+KÓLÐLØ"’NØ˜A‘I�DØ× Ñ ×$Ñ$ TÔ*Ø˜!‘‘ð+ô. ˜4×+Ñ+Ó,ò 8�Ø× Ñ ×$Ñ$¤\×%6Ñ%6Õ7ð8ð ˜’-Ø×,Ñ,×0Ñ0Ô9KÐ0ÓL�Øœ|×1Ñ1Ò1Ø×*Ñ*Ô,Ü&Ð'GÓHÐHØ’Ø˜‘	�ð ˜“-ð ×#Ñ#Ð/ÐFÐ1FÓFÐ/Ø×$Ñ$Ð0ÐHÐ2HÓHÐ0Ø�I‰IŒKØ×Ñ×"Ñ"Ô$Ø×Ñ×#Ñ#Ô%Ø×&Ò&Ø× Ñ ×3Ñ3Ô5Ø×!Ñ!×4Ñ4Õ6à× Ñ ×,Ñ,Ô.Ø×!Ñ!×-Ñ-Õ/øôO !ò (Ø#'›ð(ûô
 !ò  Ø×.Ñ.Ô0Ø˜ûð ûð0 ×#Ñ#Ð/ÐFÐ1FÓFÐ/Ø×$Ñ$Ð0ÐHÐ2HÓHÐ0Ø�I‰IŒKØ×Ñ×"Ñ"Ô$Ø×Ñ×#Ñ#Ô%Ø×&Ò&Ø× Ñ ×3Ñ3Ô5Ø×!Ñ!×4Ñ4Õ6à× Ñ ×,Ñ,Ô.Ø×!Ñ!×-Ñ-Õ/üsh   ‚N„A#J9 Á(JÂJ9 Â JÂ$C?J9 Æ$C NÊJÊJ9 ÊJÊJ9 Ê	J6ÊJ1Ê1J6Ê6J9 Ê9C!NÎNc                 ó@   —  | j                   t        |«      g|¢­i |¤ŽS r    )rj   Ú	enumerate)r'   r[   r   r   s       r   Úsemi_ordered_mapz#ParallelWorkerPool.semi_ordered_map¸   s$   € Ø!ˆt×!Ñ!¤)¨FÓ"3ÐE°dÒE¸fÑEÐEr   c              /   ó¸   K  — t        t        «      }d} | j                  |g|¢­i |¤ŽD ],  \  }}|||<   ||v sŒ|j                  |«      –— |dz  }||v rŒŒ. y ­w)Nr   r4   )r   Úintrm   Úpop)r'   r[   r   r   ÚbufferÚnext_expectedÚidxr2   s           r   Úordered_mapzParallelWorkerPool.ordered_map»   ss   è ø€ ÜœSÓ!ˆØˆà.˜×.Ñ.¨vÐG¸ÒGÀÑGò 	#‰IˆC�ØˆF�3‰KØ 6Ò)Ø—j‘j Ó/Ò/Ø Ñ"�ð   6Ó)ñ	#ùs   ‚7AºAÁAc                 óÞ   — | j                   D ]^  }|j                  «       rŒ|j                  dk7  sŒ$d| _        | j	                  «        t        d|j                  › d|j                  › �«      ‚ y)zJ
        Checks if any worker process has terminated unexpectedly
        r   TzWorker PID: z# terminated unexpectedly with code N)rK   Úis_aliveÚexitcoderM   rb   rc   Úpid©r'   r(   s     r   r_   z&ParallelWorkerPool.check_worker_healthÅ   sn   € ð —~‘~ò 	ˆGØ×#Ñ#Õ%¨'×*:Ñ*:¸aÓ*?Ø*.�Ô'Ø×&Ñ&Ô(Ü"Ø" 7§;¡; -Ð/RÐSZ×ScÑScÐRdÐeóð ñ		r   r^   c                 óÎ   — d| _         | j                  D ]5  }|j                  |¬«       |j                  «       sŒ&|j	                  «        Œ7 | j                  j                  «        y)zM
        Emergency shutdown
        @param timeout:
        @return:
        Tr]   N)rM   rK   rd   rv   Ú	terminateÚclear)r'   r^   r(   s      r   rb   z$ParallelWorkerPool.join_or_terminateÑ   sW   € ð #'ˆÔØ—~‘~ò 	$ˆGØ�L‰L ˆLÔ)Ø×ÑÕ!Ø×!Ñ!Õ#ð	$ð 	�‰×ÑÕr   c                 óz   — | j                   D ]  }|j                  «        Œ | j                   j                  «        y r    )rK   rd   r|   ry   s     r   rd   zParallelWorkerPool.joinÞ   s.   € Ø—~‘~ò 	ˆGØ�L‰L�Nð	à�‰×ÑÕr   c                 óh   — | j                   D ]#  }|j                  «       sŒ|j                  «        Œ% y)a  
        Terminate processes if the user hasn't joined. This is necessary as
        leaving stray processes running can corrupt shared state. In brief,
        we've observed shared memory counters being reused (when the memory was
        free from the perspective of the parent process) while the stray
        workers still held a reference to them.
        For a discussion of using destructors in Python in this manner, see
        https://eli.thegreenplace.net/2009/06/12/safely-using-destructors-in-python/.
        N)rK   rv   r{   ry   s     r   Ú__del__zParallelWorkerPool.__del__ã   s/   € ð —~‘~ò 	$ˆGØ×ÑÕ!Ø×!Ñ!Õ#ñ	$r   )r   N)r4   )r   r   r   ÚMAX_INTERNAL_BATCH_SIZEro   r   r   ÚstrrN   r   r$   r   rj   rm   rt   r_   rb   rd   r   r   r   r   rE   rE   W   sþ   „ ð
 $(Ø'>ñ9àð9ð �V‘ð9ð ˜D‘jð	9ð
 "%ó9ð"+˜cð + dó +ð050 H¨S¡Mð 50¸#ð 50Èð 50ÐQYÐZ]ÑQ^ó 50ðnF x°¡}ð F¸Sð FÈCð FÐT\Ð]`ÑTaó Fð# (¨3¡-ð #¸ð #Àsð #ÈxÐX[É}ó #ó
ñ¨¨t©ð ¸Dó óô
$r   rE   r    )r5   r7   Úcollectionsr   Úenumr   Úmultiprocessingr   r   Úmultiprocessing.contextr   Úmultiprocessing.processr   Úmultiprocessing.sharedctypesr	   rU   Úqueuer
   Útypingr   r   r   ra   r€   r�   r   r   ro   ÚdictrC   rE   r   r   r   ú<module>r‹      s¶   ðÛ Û 	Ý #Ý ß .Ý /Ý /Ý BÝ ß &Ñ &ð Ð àÐ ô�3˜ô ÷$ñ $ð %)ñ3<Ø�v‘,ð3<àð3<ð ð3<ð "ð	3<ð
 ð3<ð ��c�‰N˜TÑ!ð3<ð 
ó3<÷lX$ò X$r   