o
    ×èFh©c  ã                   @   sÂ  d dl Z d dlmZ d dlZd dlZd dlZd dlmZ d dlm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 ejZe e¡Zd
Zi dd “dd“dd“dd“dd“dd“dd“dd“dd“dd“dd“d d!“d"d#“d$d%“d&d'“d(d)“d*d+“d,d-i¥Zd.Ze
je
je
je
je
je
je
je
j fZ!e
jfZ"e  #d/d0d1g¡Z$G d2d3„ d3e%ƒZ&G d4d5„ d5eƒZ'G d6d7„ d7e%ƒZ(G d8d9„ d9e%ƒZ)d:d;„ Z*d<d=„ Z+d>d?„ Z,d@dA„ Z-G dBdC„ dCe%ƒZ.dS )Dé    N)ÚEnum)ÚResumableBidiRpc)ÚBackgroundConsumer)Ú
exceptions)ÚListenRequest)ÚTarget)ÚTargetChange)Ú_helpersiyP  ÚOKÚ	CANCELLEDé   ÚUNKNOWNé   ÚINVALID_ARGUMENTé   ÚDEADLINE_EXCEEDEDé   Ú	NOT_FOUNDé   ÚALREADY_EXISTSé   ÚPERMISSION_DENIEDé   ÚUNAUTHENTICATEDé   ÚRESOURCE_EXHAUSTEDé   ÚFAILED_PRECONDITIONé	   ÚABORTEDé
   ÚOUT_OF_RANGEé   ÚUNIMPLEMENTEDé   ÚINTERNALé   ÚUNAVAILABLEé   Ú	DATA_LOSSé   Ú
DO_NOT_USEéÿÿÿÿzThread-OnRpcTerminatedÚDocTreeEntryÚvalueÚindexc                   @   sT   e Z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„ Z
dd„ ZdS )ÚWatchDocTreec                 C   s   i | _ d| _d S )Nr   )Ú_dictÚ_index©Úself© r5   úX/var/www/html/loop/nvenv/lib/python3.10/site-packages/google/cloud/firestore_v1/watch.pyÚ__init__N   s   
zWatchDocTree.__init__c                 C   s   t | j ¡ ƒS ©N)Úlistr1   Úkeysr3   r5   r5   r6   r:   R   s   zWatchDocTree.keysc                 C   s"   t ƒ }| j ¡ |_| j|_|} | S r8   )r0   r1   Úcopyr2   )r4   Úwdtr5   r5   r6   Ú_copyU   s
   zWatchDocTree._copyc                 C   s,   |   ¡ } t|| jƒ| j|< |  jd7  _| S )Nr   )r=   r-   r2   r1   )r4   Úkeyr.   r5   r5   r6   Úinsert\   s   zWatchDocTree.insertc                 C   s
   | j | S r8   ©r1   ©r4   r>   r5   r5   r6   Úfindb   ó   
zWatchDocTree.findc                 C   s   |   ¡ } | j|= | S r8   )r=   r1   rA   r5   r5   r6   Úremovee   s   zWatchDocTree.removec                 c   s   � | j D ]}|V  qd S r8   r@   ©r4   Úkr5   r5   r6   Ú__iter__j   s   €
ÿzWatchDocTree.__iter__c                 C   s
   t | jƒS r8   )Úlenr1   r3   r5   r5   r6   Ú__len__n   rC   zWatchDocTree.__len__c                 C   s
   || j v S r8   r@   rE   r5   r5   r6   Ú__contains__q   rC   zWatchDocTree.__contains__N)Ú__name__Ú
__module__Ú__qualname__r7   r:   r=   r?   rB   rD   rG   rI   rJ   r5   r5   r5   r6   r0   J   s    r0   c                   @   s   e Zd ZdZdZdZdS )Ú
ChangeTyper   r   r   N)rK   rL   rM   ÚADDEDÚREMOVEDÚMODIFIEDr5   r5   r5   r6   rN   u   s    rN   c                   @   ó   e Zd Zdd„ ZdS )ÚDocumentChangec                 C   s   || _ || _|| _|| _dS )z±DocumentChange

        Args:
            type (ChangeType):
            document (document.DocumentSnapshot):
            old_index (int):
            new_index (int):
        N)ÚtypeÚdocumentÚ	old_indexÚ	new_index)r4   rT   rU   rV   rW   r5   r5   r6   r7   |   s   

zDocumentChange.__init__N©rK   rL   rM   r7   r5   r5   r5   r6   rS   {   ó    rS   c                   @   rR   )ÚWatchResultc                 C   s   || _ || _|| _d S r8   )ÚsnapshotÚnameÚchange_type)r4   r[   r\   r]   r5   r5   r6   r7   �   s   
zWatchResult.__init__NrX   r5   r5   r5   r6   rZ   Œ   rY   rZ   c                 C   s   t | tjƒrt | ¡S | S )z(Wraps a gRPC exception class, if needed.)Ú
isinstanceÚgrpcÚRpcErrorr   Úfrom_grpc_error)Ú	exceptionr5   r5   r6   Ú_maybe_wrap_exception“   s   
rc   c                 C   s   | |ksJ dƒ‚dS )Nz+Document watches only support one document.r   r5   )Údoc1Údoc2r5   r5   r6   Údocument_watch_comparatorš   s   rf   c                 C   ó   t | ƒ}t|tƒS r8   )rc   r^   Ú_RECOVERABLE_STREAM_EXCEPTIONS©rb   Úwrappedr5   r5   r6   Ú_should_recoverŸ   ó   
rk   c                 C   rg   r8   )rc   r^   Ú_TERMINATING_STREAM_EXCEPTIONSri   r5   r5   r6   Ú_should_terminate¤   rl   rn   c                
   @   sð   e Zd Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZd	d
„ Zdd„ Z	e
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ejeejeejeejeejeiZd d!„ Zd"d#„ Zd$d%„ Zed&d'„ ƒZd(d)„ Z d*d+„ Z!d,d-„ Z"dS )/ÚWatchc                 C   sz   || _ || _|| _|| _|| _|| _|j| _t 	¡ | _
d| _|  |j¡ d| _tƒ | _i | _i | _d| _d| _|  ¡  dS )aÕ  
        Args:
            firestore:
            target:
            comparator:
            snapshot_callback: Callback method to process snapshots.
                Args:
                    docs (List(DocumentSnapshot)): A callback that returns the
                        ordered list of documents stored in this snapshot.
                    changes (List(str)): A callback that returns the list of
                        changed documents since the last snapshot delivered for
                        this watch.
                    read_time (string): The ISO 8601 time at which this
                        snapshot was obtained.

            document_snapshot_cls: factory for instances of DocumentSnapshot
        FN)Ú_document_referenceÚ
_firestoreÚ_targetsÚ_comparatorÚ_document_snapshot_clsÚ_snapshot_callbackÚ_firestore_apiÚ_apiÚ	threadingÚLockÚ_closingÚ_closedÚ_set_documents_pfxÚ_database_stringÚresume_tokenr0   Údoc_treeÚdoc_mapÚ
change_mapÚcurrentÚ
has_pushedÚ_init_stream)r4   Údocument_referenceÚ	firestoreÚtargetÚ
comparatorÚsnapshot_callbackÚdocument_snapshot_clsr5   r5   r6   r7   ª   s"   
zWatch.__init__c                 C   sP   | j }t| jjjtt|| jjd�| _	| j	 
| j¡ t| j	| jƒ| _| j ¡  d S )N)Ú	start_rpcÚshould_recoverÚshould_terminateÚinitial_requestÚmetadata)Ú_get_rpc_requestr   rw   Ú
_transportÚlistenrk   rn   rq   Ú_rpc_metadataÚ_rpcÚadd_done_callbackÚ_on_rpc_doner   Úon_snapshotÚ	_consumerÚstart)r4   Úrpc_requestr5   r5   r6   r„   è   s   ûzWatch._init_streamc                 C   s"   | ||j d|jgitdœt||ƒS )a¹  
        Creates a watch snapshot listener for a document. snapshot_callback
        receives a DocumentChange object, but may also start to get
        targetChange and such soon

        Args:
            document_ref: Reference to Document
            snapshot_callback: callback to be called on snapshot
            document_snapshot_cls: class to make snapshots with
            reference_class_instance: class make references

        Ú	documents)r›   Ú	target_id)Ú_clientÚ_document_pathÚWATCH_TARGET_IDrf   )ÚclsÚdocument_refr‰   rŠ   r5   r5   r6   Úfor_documentù   s   
þ÷zWatch.for_documentc                 C   s>   |j  ¡ \}}tj|| ¡ d�}| ||j|jtdœ|j||ƒS )N)ÚparentÚstructured_query)Úqueryrœ   )	Ú_parentÚ_parent_infor   ÚQueryTargetÚ_to_protobufr�   Ú_pbrŸ   rs   )r    r¥   r‰   rŠ   Úparent_pathÚ_Úquery_targetr5   r5   r6   Ú	for_query  s   ÿ
úzWatch.for_queryc                 C   s8   | j d ur| j | jd< n| j dd ¡ t| jj| jd�S )Nr~   )ÚdatabaseÚ
add_target)r~   rr   Úpopr   rq   r}   r3   r5   r5   r6   r�   (  s   

ÿzWatch._get_rpc_requestc                 C   s   |› d�| _ t| j ƒ| _d S )Nz/documents/)Ú_documents_pfxrH   Ú_documents_pfx_len)r4   Údatabase_stringr5   r5   r6   r|   2  s   zWatch._set_documents_pfxc                 C   s   | j duo| j jS )z¸bool: True if this manager is actively streaming.

        Note that ``False`` does not indicate this is complete shut down,
        just that it stopped getting new messages.
        N)r˜   Ú	is_activer3   r5   r5   r6   rµ   6  s   zWatch.is_activeNc                 C   sª   | j �4 | jr	 W d  ƒ dS | jrt d¡ | j ¡  d| _| j ¡  d| _d| _t d¡ W d  ƒ n1 s:w   Y  |rSt d| ¡ t	|t
ƒrO|‚t|ƒ‚dS )a  Stop consuming messages and shutdown all helper threads.

        This method is idempotent. Additional calls will have no effect.

        Args:
            reason (Any): The reason to close this. If None, this is considered
                an "intentional" shutdown.
        NzStopping consumer.TzFinished stopping manager.zreason for closing: %s)rz   r{   rµ   Ú_LOGGERÚdebugr˜   Ústopr”   Úcloser^   Ú	ExceptionÚRuntimeError)r4   Úreasonr5   r5   r6   r¹   ?  s&   	þ


ó
ûzWatch.closec                 C   s:   t  d¡ t|ƒ}tjt| jd|id�}d|_| ¡  dS )a
  Triggered whenever the underlying RPC terminates without recovery.

        This is typically triggered from one of two threads: the background
        consumer thread (when calling ``recv()`` produces a non-recoverable
        error) or the grpc management thread (when cancelling the RPC).

        This method is *non-blocking*. It will start another thread to deal
        with shutting everything down. This is to prevent blocking in the
        background consumer and preventing it from being ``joined()``.
        z.RPC termination has signaled manager shutdown.r¼   )r\   r‡   ÚkwargsTN)	r¶   Úinforc   rx   ÚThreadÚ_RPC_ERROR_THREAD_NAMEr¹   Údaemonr™   )r4   ÚfutureÚthreadr5   r5   r6   r–   ^  s   
ÿzWatch._on_rpc_donec                 C   s   |   ¡  d S r8   )r¹   r3   r5   r5   r6   Úunsubscribeq  s   zWatch.unsubscribec                 C   sR   t  d¡ |jd u pt|jƒdk}|r#|jr%| jr'|  |j|j¡ d S d S d S d S )Nz%on_snapshot: target change: NO_CHANGEr   )r¶   r·   Ú
target_idsrH   Ú	read_timer‚   Úpushr~   )r4   Útarget_changeÚno_target_idsr5   r5   r6   Ú$_on_snapshot_target_change_no_changet  s   
ÿûz*Watch._on_snapshot_target_change_no_changec                 C   s,   t  d¡ |jd }|tkrtd| ƒ‚d S )Nzon_snapshot: target change: ADDr   z&Unexpected target ID %s sent by server)r¶   r·   rÅ   rŸ   r»   )r4   rÈ   rœ   r5   r5   r6   Ú_on_snapshot_target_change_add�  s
   

ÿz$Watch._on_snapshot_target_change_addc                 C   sJ   t  d¡ |jjr|jj}|jj}nd}d}d||f }t|ƒt ||¡‚)Nz"on_snapshot: target change: REMOVEr&   zinternal errorzError %s:  %s)r¶   r·   ÚcauseÚcodeÚmessager»   r   Úfrom_grpc_status)r4   rÈ   rÍ   rÎ   Úerror_messager5   r5   r6   Ú!_on_snapshot_target_change_remove‡  s   


ÿz'Watch._on_snapshot_target_change_removec                 C   s   t  d¡ |  ¡  d S )Nz!on_snapshot: target change: RESET)r¶   r·   Ú_reset_docs©r4   rÈ   r5   r5   r6   Ú _on_snapshot_target_change_reset—  s   
z&Watch._on_snapshot_target_change_resetc                 C   s   t  d¡ d| _d S )Nz#on_snapshot: target change: CURRENTT)r¶   r·   r‚   rÓ   r5   r5   r6   Ú"_on_snapshot_target_change_currentœ  s   

z(Watch._on_snapshot_target_change_currentc                 C   s   |  | j¡r|| jd … }|S r8   )Ú
startswithr²   r³   )r4   Údocument_namer5   r5   r6   Ú_strip_document_pfx¨  s   zWatch._strip_document_pfxc              
   C   sZ  |du r
|   ¡  dS |j}| d¡}|dkr`|jj}t d|› �¡ | j |¡}|du rAd|› �}t 	d|› �¡ | j t
|ƒd� z	|| |jƒ W dS  ty_ } z	t d|› �¡ ‚ d}~ww |d	kr»t d
¡ t|jjv }t|jjv }	|jj}
|r©t d¡ t |
j| j¡}|  |
j¡}| j |¡}| j||dd|
j|
jd�}|| j|
j< dS |	r¹t d¡ tj| j|
j< dS dS |dkrÐt d¡ |jj}tj| j|< dS |dkråt d¡ |jj}tj| j|< dS |dk�rt d¡ |jj |  !¡ k�rt 	d¡ t"j#t$| j d�}| %¡  | &¡  |  '¡  |  (¡  dS dS t d¡ d|› �}| j t
|ƒd� dS )aS  Process a response from the bi-directional gRPC stream.

        Collect changes and push the changes in a batch to the customer
        when we receive 'current' from the listen response.

        Args:
            proto(`google.cloud.firestore_v1.types.ListenResponse`):
                Callback method that receives a object to
        NÚresponse_typerÈ   zon_snapshot: target change: zUnknown target change type: zon_snapshot: )r¼   zmeth(proto) exc: Údocument_changezon_snapshot: document changez%on_snapshot: document change: CHANGEDT)Ú	referenceÚdataÚexistsrÆ   Úcreate_timeÚupdate_timez%on_snapshot: document change: REMOVEDÚdocument_deletez$on_snapshot: document change: DELETEÚdocument_removez$on_snapshot: document change: REMOVEÚfilterzon_snapshot: filter updatez%Filter mismatch -- restarting stream.)r\   r‡   zUNKNOWN TYPE. UHOHzUnknown listen response type: ))r¹   rª   Ú
WhichOneofrÈ   Útarget_change_typer¶   r·   Ú_target_changetype_dispatchÚgetr¾   Ú
ValueErrorrº   rŸ   rÚ   rÅ   Úremoved_target_idsrU   r	   Údecode_dictÚfieldsrq   rØ   r\   rt   rÞ   rß   r�   rN   rP   rà   rá   râ   ÚcountÚ_current_sizerx   r¿   rÀ   r™   ÚjoinrÒ   r„   )r4   ÚprotoÚpbÚwhichrä   ÚmethrÎ   Úexc2ÚchangedÚremovedrU   rÜ   r×   r¡   r[   r\   rÃ   r5   r5   r6   r—   ­  s†   


€þ

ú
þ




þô

zWatch.on_snapshotc                 C   s’   |   | j| j|¡\}}}|  | j| j|||¡\}}}| jr!t|ƒr9t | j	¡}	t
| ¡ |	d�}
|  |
||¡ d| _|| _|| _| j ¡  || _dS )z Invoke the callback with a new snapshot

        Build the sntapshot from the current set of changes.

        Clear the current changes on completion.
        ©r>   TN)Ú_extract_changesr€   r�   Ú_compute_snapshotr   rƒ   rH   Ú	functoolsÚ
cmp_to_keyrs   Úsortedr:   ru   Úclearr~   )r4   rÆ   Únext_resume_tokenÚdeletesÚaddsÚupdatesÚupdated_treeÚupdated_mapÚappliedChangesr>   r:   r5   r5   r6   rÇ     s   

ÿ
ÿ

z
Watch.pushc                 C   s€   g }g }g }|  ¡ D ]0\}}|tjkr|| v r| |¡ q
|| v r.|d ur(||_| |¡ q
|d ur5||_| |¡ q
|||fS r8   )ÚitemsrN   rP   ÚappendrÆ   )r€   ÚchangesrÆ   rý   rþ   rÿ   r\   r.   r5   r5   r6   rö   8  s    

€
zWatch._extract_changesc                    s  |}|}t |ƒt |ƒksJ dƒ‚dd„ ‰dd„ ‰ ‡ ‡fdd„}g }	t | j¡}
t|ƒ}|D ]}ˆ|||ƒ\}}}|	 |¡ q-t||
d�}t d	¡ |D ]}t d
¡ ˆ |||ƒ\}}}|	 |¡ qKt||
d�}|D ]}||||ƒ\}}}|d ur}|	 |¡ qit |ƒt |ƒksŠJ dƒ‚|||	fS )NzJThe document tree and document map should have the same number of entries.c                 S   sP   | |v sJ dƒ‚|  | ¡}| |¡}|j}| |¡}|| = ttj||dƒ||fS )z–
            Applies a document delete to the document tree and document map.
            Returns the corresponding DocumentChange event.
            z!Document to delete does not existr,   )ræ   rB   r/   rD   rS   rN   rP   )r\   r   r  Úold_documentÚexistingrV   r5   r5   r6   Ú
delete_docX  s   


ýz+Watch._compute_snapshot.<locals>.delete_docc                 S   sN   | j j}||vsJ dƒ‚| | d¡}| | ¡j}| ||< ttj| d|ƒ||fS )z—
            Applies a document add to the document tree and the document map.
            Returns the corresponding DocumentChange event.
            zDocument to add already existsNr,   )rÛ   rž   r?   rB   r/   rS   rN   rO   )Únew_documentr   r  r\   rW   r5   r5   r6   Úadd_docj  s   ýz(Watch._compute_snapshot.<locals>.add_docc                    sv   | j j}||v sJ dƒ‚| |¡}|j| jkr6ˆ|||ƒ\}}}ˆ | ||ƒ\}}}ttj| |j|jƒ||fS d||fS )z»
            Applies a document modification to the document tree and the
            document map.
            Returns the DocumentChange event for successful modifications.
            z!Document to modify does not existN)	rÛ   rž   ræ   rß   rS   rN   rQ   rV   rW   )r	  r   r  r\   r  Úremove_changeÚ
add_change©r
  r  r5   r6   Ú
modify_docz  s(   

ÿ
ÿüø
z+Watch._compute_snapshot.<locals>.modify_docrõ   zwalk over add_changeszin add_changeszQThe update document tree and document map should have the same number of entries.)rH   rø   rù   rs   rú   r  r¶   r·   )r4   r   r€   Údelete_changesÚadd_changesÚupdate_changesr   r  r  r  r>   r\   Úchanger[   r5   r  r6   r÷   M  sH   ÿ!
ÿ


ÿ
ÿ
€ÿ
zWatch._compute_snapshotc                 C   s2   |   | j| jd¡\}}}t| jƒt|ƒ t|ƒ S )zsReturn the current count of all documents.

        Count includes the changes from the current changeMap.
        N)rö   r€   r�   rH   )r4   rý   rþ   r¬   r5   r5   r6   rì   ¾  s   zWatch._current_sizec                 C   sH   t  d¡ | j ¡  d| _| j ¡ D ]}|jj}t	j
| j|< qd| _dS )zG
        Helper to clear the docs on RESET or filter mismatch.
        zresetting documentsNF)r¶   r·   r�   rû   r~   r   r:   rÛ   rž   rN   rP   r‚   )r4   r[   r\   r5   r5   r6   rÒ   Æ  s   


zWatch._reset_docsr8   )#rK   rL   rM   r7   r„   Úclassmethodr¢   r®   r�   r|   Úpropertyrµ   r¹   r–   rÄ   rÊ   rË   rÑ   rÔ   rÕ   ÚTargetChangeTypeÚ	NO_CHANGEÚADDÚREMOVEÚRESETÚCURRENTrå   rØ   r—   rÇ   Ústaticmethodrö   r÷   rì   rÒ   r5   r5   r5   r6   ro   ©   sB    >




ûn
qro   )/ÚcollectionsÚenumr   rø   Úloggingrx   Úgoogle.api_core.bidir   r   Úgoogle.api_corer   r_   Ú)google.cloud.firestore_v1.types.firestorer   r   r   Úgoogle.cloud.firestore_v1r	   r  Ú	getLoggerrK   r¶   rŸ   ÚGRPC_STATUS_CODErÀ   ÚAbortedÚ	CancelledÚUnknownÚDeadlineExceededÚResourceExhaustedÚInternalServerErrorÚServiceUnavailableÚUnauthenticatedrh   rm   Ú
namedtupler-   Úobjectr0   rN   rS   rZ   rc   rf   rk   rn   ro   r5   r5   r5   r6   Ú<module>   s”   
ÿþýüûúùø	÷
öõôóòñðïîø
+