Ë
    Ìi)   ã                   ó�   — d dl Z d dlmZ d dlmZmZmZmZmZ d dl	Z
d dlmZ d dlmZmZmZmZmZ d dlmZ ddlmZ  G d	„ d
«      Zy)é    N)ÚPriorityQueue)ÚDictÚ	GeneratorÚListÚOptionalÚSet)ÚManifest)ÚExposureÚGraphMemberNodeÚMetricÚ	ModelNodeÚSourceDefinition)ÚNodeTypeé   )ÚUniqueIdc                   ó`  — e Zd ZdZ	 ddej
                  dedee   de	ddf
d„Z
dee   fd	„Zd
ede	fd„Zedej
                  deee   ddf   fd„«       Zdej
                  deeef   fd„Zdde	dee   defd„Zdefd„Zde	fd„Zdede	fd„Zdd„Zd
eddfd„Zdd
ede	ddfd„Zdd„Zdefd„Z y)Ú
GraphQueuea*  A fancy queue that is backed by the dependency graph.
    Note: this will mutate input!

    This queue is thread-safe for `mark_done` calls, though you must ensure
    that separate threads do not call `.empty()` or `__len__()` and `.get()` at
    the same time, as there is an unlocked race!
    ÚgraphÚmanifestÚselectedÚpreserve_edgesÚreturnNc                 ó  — |r|n(t         j                  j                  j                  |«      | _        || _        || _        t        «       | _        t        «       | _
        t        «       | _        t        «       | _        t        j                  «       | _        | j!                  | j                  «      | _        | j%                  t'        | j                  j)                  «       «      «       t        j*                  | j                  «      | _        y ©N)ÚnxÚclassesÚfunctionÚcreate_empty_copyr   r   Ú	_selectedr   ÚinnerÚsetÚin_progressÚin_progress_microbatchÚqueuedÚ	threadingÚLockÚlockÚ_get_scoresÚ_scoresÚ_find_new_additionsÚlistÚnodesÚ	ConditionÚsome_task_done)Úselfr   r   r   r   s        úP/var/www/html/strategist-ai/venv/lib/python3.12/site-packages/dbt/graph/queue.pyÚ__init__zGraphQueue.__init__   s¶   € ñ -‘U´"·*±*×2EÑ2E×2WÑ2WÐX]Ó2^ˆŒ
Ø ˆŒØ!ˆŒä$1£OˆŒ
ô +.«%ˆÔä58³UˆÔ#ä%(£UˆŒä—N‘NÓ$ˆŒ	à×'Ñ'¨¯
©
Ó3ˆŒà× Ñ ¤ d§j¡j×&6Ñ&6Ó&8Ó!9Ô:ä'×1Ñ1°$·)±)Ó<ˆÕó    c                 ó6   — | j                   j                  «       S r   )r   Úcopy©r/   s    r0   Úget_selected_nodeszGraphQueue.get_selected_nodes:   s   € Ø�~‰~×"Ñ"Ó$Ð$r2   Únode_idc                 óÊ   — | j                   j                  |«      }|j                  t        j                  k7  ryt        |t        t        t        f«      rJ ‚|j                  ryy)NFT)
r   ÚexpectÚresource_typer   ÚModelÚ
isinstancer   r
   r   Úis_ephemeral)r/   r7   Únodes      r0   Ú_include_in_costzGraphQueue._include_in_cost=   sR   € Ø�}‰}×#Ñ# GÓ,ˆØ×Ñ¤§¡Ò/Øä˜dÔ%5´xÄÐ$HÔIÐIÐIØ×ÒØØr2   c              #   ój  K  — | j                  «       D ��ci c]  \  }}|dkD  sŒ||“Œ }}}| j                  «       D ��cg c]  \  }}|dk(  sŒ|‘Œ }}}|rP|–— g }|D ]?  }| j                  |«      D ])  \  }}||xx   dz  cc<   ||   rŒ|j                  |«       Œ+ ŒA |}|rŒOyyc c}}w c c}}w ­w)a¨  Topological sort of given graph that groups ties.

        Adapted from `nx.topological_sort`, this function returns a topo sort of a graph however
        instead of arbitrarily ordering ties in the sort order, ties are grouped together in
        lists.

        Args:
            graph: The graph to be sorted.

        Returns:
            A generator that yields lists of nodes, one list per graph depth level.
        r   r   N)Ú	in_degreeÚedgesÚappend)r   ÚvÚdÚindegree_mapÚzero_indegreeÚnew_zero_indegreeÚ_Úchilds           r0   Ú_grouped_topological_sortz$GraphQueue._grouped_topological_sortG   sÈ   è ø€ ð  */¯©Ó):×D¡  A¸aÀ!»e˜˜1™ÐDˆÑDØ',§¡Ó'8×C™t˜q !¸AÀ»FšÐCˆÑCáØÒØ "ÐØ"ò 8�Ø %§¡¨A£ò 8‘H�A�uØ  Ó'¨1Ñ,Ó'Ø'¨Ó.Ø)×0Ñ0°Õ7ñ8ð8ð
 .ˆMô ùó EùÛCùs1   ‚B3–B'¤B'©B3¿B-ÁB-Á9B3ÂB3Â%B3c                 óÜ   ‡— ˆfd„t        j                  t        j                  ‰«      «      D «       }i }|D ]2  }| j                  |«      }t	        |«      D ]  \  }}|D ]  }|||<   Œ	 Œ Œ4 |S )a  Scoring nodes for processing order.

        Scores are calculated by the graph depth level. Lowest score (0) should be processed first.

        Args:
            graph: The graph to be scored.

        Returns:
            A dictionary consisting of `node name`:`score` pairs.
        c              3   ó@   •K  — | ]  }‰j                  |«      –— Œ y ­wr   )Úsubgraph)Ú.0Úxr   s     €r0   ú	<genexpr>z)GraphQueue._get_scores.<locals>.<genexpr>p   s   øè ø€ ÒY¨1�U—^‘^ A×&ÑYùs   ƒ)r   Úconnected_componentsÚGraphrK   Ú	enumerate)	r/   r   Ú	subgraphsÚscoresrN   Úgrouped_nodesÚlevelÚgroupr>   s	    `       r0   r(   zGraphQueue._get_scoresd   s�   ø€ ó Z´×0GÑ0GÌÏÉÐQVËÓ0XÔYˆ	ð ˆØ!ò 	)ˆHØ ×:Ñ:¸8ÓDˆMÜ )¨-Ó 8ò )‘��uØ!ò )�DØ#(�F˜4’Lñ)ñ)ð	)ð ˆr2   ÚblockÚtimeoutc                 ó<  — | j                   j                  ||¬«      \  }}| j                  j                  |«      }t	        |t
        «      xr |j                  j                  dk(  }| j                  5  | j                  ||¬«       ddd«       |S # 1 sw Y   |S xY w)a¬  Get a node off the inner priority queue. By default, this blocks.

        This takes the lock, but only for part of it.

        :param block: If True, block until the inner queue has data
        :param timeout: If set, block for timeout seconds waiting for data.
        :return: The node as present in the manifest.

        See `queue.PriorityQueue` for more information on `get()` behavior and
        exceptions.
        )rZ   r[   Ú
microbatch)Úis_microbatchN)
r    Úgetr   r9   r<   r   ÚconfigÚincremental_strategyr'   Ú_mark_in_progress)r/   rZ   r[   rI   r7   r>   r^   s          r0   r_   zGraphQueue.get|   s“   € ð —Z‘Z—^‘^¨%¸�^ÓA‰
ˆˆ7Ø�}‰}×#Ñ# GÓ,ˆä�tœYÓ'Ò\¨D¯K©K×,LÑ,LÐP\Ñ,\ð 	ð �Y‰Yñ 	IØ×"Ñ" 7¸-Ð"ÔH÷	Ið ˆ÷	Ið ˆús   Á3BÂBc                 óœ   — | j                   5  t        | j                  «      t        | j                  «      z
  cddd«       S # 1 sw Y   yxY w)zÐThe length of the queue is the number of tasks left for the queue to
        give out, regardless of where they are. Incomplete tasks are not part
        of the length.

        This takes the lock.
        N)r'   Úlenr   r"   r5   s    r0   Ú__len__zGraphQueue.__len__“   s;   € ð �Y‰Yñ 	;Ü�t—z‘z“?¤S¨×)9Ñ)9Ó%:Ñ:÷	;÷ 	;ò 	;ús   �+AÁAc                 ó   — t        | «      dk(  S )z�The graph queue is 'empty' if it all remaining nodes in the graph
        are in progress.

        This takes the lock.
        r   )rd   r5   s    r0   ÚemptyzGraphQueue.empty�   s   € ô �4‹y˜A‰~Ðr2   r>   c                 ó>   — || j                   v xs || j                  v S )zðDecide if a node is already known (either handed out as a task, or
        in the queue).

        Callers must hold the lock.

        :param str node: The node ID to check
        :returns bool: If the node is in progress/queued.
        )r"   r$   )r/   r>   s     r0   Ú_already_knownzGraphQueue._already_known¥   s#   € ð �t×'Ñ'Ð'Ò>¨4°4·;±;Ð+>Ð>r2   c                 óþ   — |D ]x  }| j                   j                  |«      dk(  sŒ"| j                  |«      rŒ4| j                  j	                  | j
                  |   |f«       | j                  j                  |«       Œz y)zfFind any nodes in the graph that need to be added to the internal
        queue and add them.
        r   N)r   rA   ri   r    Úputr)   r$   Úadd)r/   Ú
candidatesr>   s      r0   r*   zGraphQueue._find_new_additions°   se   € ð ò 	&ˆDØ�z‰z×#Ñ# DÓ)¨QÓ.°t×7JÑ7JÈ4Õ7PØ—
‘
—‘ §¡¨TÑ 2°DÐ9Ô:Ø—‘—‘ Õ%ñ	&r2   c                 óÖ  — | j                   5  | j                  j                  |«       || j                  v r| j                  j                  |«       t	        | j
                  j                  |«      «      }| j
                  j                  |«       | j                  |«       | j                  j                  «        | j                  j                  «        ddd«       y# 1 sw Y   yxY w)z–Given a node's unique ID, mark it as done.

        This method takes the lock.

        :param str node_id: The node ID to mark as complete.
        N)r'   r"   Úremover#   r+   r   Ú
successorsÚremove_noder*   r    Ú	task_doner.   Ú
notify_all)r/   r7   rp   s      r0   Ú	mark_donezGraphQueue.mark_done¹   s±   € ð �Y‰Yñ 	-Ø×Ñ×#Ñ# GÔ,Ø˜$×5Ñ5Ñ5Ø×+Ñ+×2Ñ2°7Ô;Ü˜dŸj™j×3Ñ3°GÓ<Ó=ˆJØ�J‰J×"Ñ" 7Ô+Ø×$Ñ$ ZÔ0Ø�J‰J× Ñ Ô"Ø×Ñ×*Ñ*Ô,÷	-÷ 	-ñ 	-ús   �C	CÃC(r^   c                 ó¬   — | j                   j                  |«       | j                  j                  |«       |r| j                  j                  |«       yy)zÙMark the node as 'in progress'.

        Callers must hold the lock.

        :param str node_id: The node ID to mark as in progress.
        :param bool is_microbatch: Whether the node is a microbatch model.
        N)r$   ro   r"   rl   r#   )r/   r7   r^   s      r0   rb   zGraphQueue._mark_in_progressÊ   sF   € ð 	�‰×Ñ˜7Ô#Ø×Ñ×Ñ˜WÔ%ÙØ×'Ñ'×+Ñ+¨GÕ4ð r2   c                 ó8   — | j                   j                  «        y)z’Join the queue. Blocks until all tasks are marked as done.

        Make sure not to call this before the queue reports that it is empty.
        N)r    Újoinr5   s    r0   rw   zGraphQueue.join×   s   € ð
 	�
‰
�‰Õr2   c                 ó¦   — | j                   5  | j                  j                  «        | j                  j                  cddd«       S # 1 sw Y   yxY w)zXBlock until a task is done, then return the number of unfinished
        tasks.
        N)r'   r.   Úwaitr    Úunfinished_tasksr5   s    r0   Úwait_until_something_was_donez(GraphQueue.wait_until_something_was_doneÞ   s?   € ð �Y‰Yñ 	/Ø×Ñ×$Ñ$Ô&Ø—:‘:×.Ñ.÷	/÷ 	/ò 	/ús   �0AÁA)T)TN)r   N)F)!Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   ÚDiGraphr	   r   r   Úboolr1   r6   r?   Ústaticmethodr   r   ÚstrrK   r   Úintr(   r   Úfloatr   r_   re   rg   ri   r*   rt   rb   rw   r{   © r2   r0   r   r      sP  „ ñð  $ñ=à�z‰zð=ð ð=ð �h‘-ð	=ð
 ð=ð 
ó=ð:% C¨¡Mó %ð¨ð °Tó ð ð.Ø�z‰zð.à	�4˜‘9˜d DÐ(Ñ	)ò.ó ð.ð8 §¡ð °°S¸#°X±ó ñ0˜ð ¨x¸©ð È/ó ð.;˜ó ;ð�tó ð	? 8ð 	?°ó 	?ó&ð- ð -¨dó -ñ"5¨ð 5À$ð 5ÐSWó 5óð/¨sô /r2   r   )r%   Úqueuer   Útypingr   r   r   r   r   Únetworkxr   Údbt.contracts.graph.manifestr	   Údbt.contracts.graph.nodesr
   r   r   r   r   Údbt.node_typesr   r   r   r   r†   r2   r0   ú<module>r�      s5   ðÛ Ý ß 7Õ 7ã å 1÷õ õ $å ÷P/ò P/r2   