Ë
    Ü_ËiÞI  ã                  ó°  — d Z ddlmZ ddlZddlmZmZmZmZm	Z	m
Z
mZ ddlmZm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mZmZmZmZ ddlmZ ddl m!Z!m"Z"m#Z# ddl$m%Z% ddl&m'Z'm(Z(m)Z) dZ* e+g d¢«      Z,erddl-m.Z. ddl/m0Z0 ddl1m2Z2 ddl3m4Z4 ddl5m6Z6 dd„Z7 G d„ dee(   «      Z8 G d„ de8e(   «      Z9 G d„ de8e(   «      Z: G d„ de:e(   «      Z;y) zAWatch changes on a collection, a database, or the entire cluster.é    )ÚannotationsN)ÚTYPE_CHECKINGÚAnyÚGenericÚMappingÚOptionalÚTypeÚUnion)ÚCodecOptionsÚ_bson_to_dict)ÚRawBSONDocument)Ú	Timestamp)Ú_csotÚcommon)Úvalidate_collation_or_none)ÚConnectionFailureÚCursorNotFoundÚInvalidOperationÚOperationFailureÚPyMongoError)Ú_Op)Ú_AggregationCommandÚ_CollectionAggregationCommandÚ_DatabaseAggregationCommand)ÚCommandCursor)Ú_CollationInÚ_DocumentTypeÚ	_PipelineT)é   é   éY   é[   é½   i  i)#  i{'  iP-  iR-  i{4  i|4  é?   é–   iL4  éê   é…   )ÚClientSession)Ú
Collection)ÚDatabase)ÚMongoClient)Ú
Connectionc                óú   — t        | t        t        f«      ryt        | t        «      rT| j                  €y| j                  dk\  xr | j                  d«      xs# | j                  dk  xr | j                  t        v S y)z5Return True if given a resumable change stream error.TFé	   ÚResumableChangeStreamError)Ú
isinstancer   r   r   Ú_max_wire_versionÚhas_error_labelÚcodeÚ_RESUMABLE_GETMORE_ERRORS)Úexcs    úd/var/www/pod-logistic/pod-api/venv/lib/python3.12/site-packages/pymongo/synchronous/change_stream.pyÚ
_resumabler7   M   s|   € ä�#Ô)¬>Ð:Ô;ØÜ�#Ô'Ô(Ø× Ñ Ð(Øà×!Ñ! QÑ&Ò\¨3×+>Ñ+>Ð?[Ó+\òSà×#Ñ# aÑ'ÒQ¨C¯H©HÔ8QÐ,Qð	Sð ó    c                  óN  — e Zd ZdZ	 	 	 d	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 dd„Zdd„Zedd„«       Ze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 d„Zed!d„«       Zej(                  d"d„«       ZeZed#d„«       Zej(                  d$d„«       Zd d„Zd%d„Zy)&ÚChangeStreama«  The internal abstract base class for change stream cursors.

    Should not be called directly by application developers. Use
    :meth:`pymongo.collection.Collection.watch`,
    :meth:`pymongo.database.Database.watch`, or
    :meth:`pymongo.mongo_client.MongoClient.watch` instead.

    .. versionadded:: 3.6
    .. seealso:: The MongoDB documentation on `changeStreams <https://mongodb.com/docs/manual/changeStreams/>`_.
    Nc                óô  — |€g }t        j                  d|«      }t        j                  d|«       t        |«       t        j                  d|«       d| _        |j                  | _        |j                  j                  j                  r=d| _        |j                  |j                  j                  t        ¬«      ¬«      | _        n|| _        t        j                  |«      | _        || _        || _        |
d u| _        |d u| _        t        j                  |
xs |«      | _        || _        || _        || _        || _        |	| _        || _        d| _        | j                  j8                  | _        || _        y )NÚpipelineÚfull_documentÚ	batchSizeFT)Údocument_class)Úcodec_options)r   Úvalidate_listÚvalidate_string_or_noner   Ú%validate_non_negative_integer_or_noneÚ_decode_customr@   Ú_orig_codec_optionsÚtype_registryÚ_decoder_mapÚwith_optionsr   Ú_targetÚcopyÚdeepcopyÚ	_pipelineÚ_full_documentÚ_full_document_before_changeÚ_uses_start_afterÚ_uses_resume_afterÚ_resume_tokenÚ_max_await_time_msÚ_batch_sizeÚ
_collationÚ_start_at_operation_timeÚ_sessionÚ_commentÚ_closedÚ_timeoutÚ_show_expanded_events)ÚselfÚtargetr<   r=   Úresume_afterÚmax_await_time_msÚ
batch_sizeÚ	collationÚstart_at_operation_timeÚsessionÚstart_afterÚcommentÚfull_document_before_changeÚshow_expanded_eventss                 r6   Ú__init__zChangeStream.__init__f   sR  € ð( ÐØˆHÜ×'Ñ'¨
°HÓ=ˆÜ×&Ñ& ¸ÔFÜ" 9Ô-Ü×4Ñ4°[À*ÔMà#ˆÔØ@F×@TÑ@TˆÔ Ø×Ñ×-Ñ-×:Ò:Ø"&ˆDÔð "×.Ñ.Ø$×2Ñ2×?Ñ?ÌÐ?Ó_ð /ó ˆD�Lð "ˆDŒLäŸ™ xÓ0ˆŒØ+ˆÔØ,GˆÔ)Ø!,°DÐ!8ˆÔØ".°dÐ":ˆÔÜ!Ÿ]™]¨;Ò+F¸,ÓGˆÔØ"3ˆÔØ%ˆÔØ#ˆŒØ(?ˆÔ%ØˆŒØˆŒØˆŒØŸ™×-Ñ-ˆŒØ%9ˆÕ"r8   c                ó.   — | j                  «       | _        y ©N)Ú_create_cursorÚ_cursor©r[   s    r6   Ú_initialize_cursorzChangeStream._initialize_cursor�   s   € à×*Ñ*Ó,ˆ�r8   c                ó   — t         ‚)z)The aggregation command class to be used.©ÚNotImplementedErrorrl   s    r6   Ú_aggregation_command_classz'ChangeStream._aggregation_command_class¡   s
   € ô "Ð!r8   c                ó   — t         ‚)zeThe client against which the aggregation commands for
        this ChangeStream will be run.
        ro   rl   s    r6   Ú_clientzChangeStream._client¦   s
   € ô
 "Ð!r8   c                ó.  — i }| j                   �| j                   |d<   | j                  �| j                  |d<   | j                  }|�| j                  r||d<   n!||d<   n| j                  �| j                  |d<   | j
                  r| j
                  |d<   |S )z=Return the options dict for the $changeStream pipeline stage.ÚfullDocumentÚfullDocumentBeforeChangeÚ
startAfterÚresumeAfterÚstartAtOperationTimeÚshowExpandedEvents)rM   rN   Úresume_tokenrO   rU   rZ   )r[   Úoptionsr{   s      r6   Ú_change_stream_optionsz#ChangeStream._change_stream_options­   sª   € à"$ˆØ×ÑÐ*Ø&*×&9Ñ&9ˆG�NÑ#à×,Ñ,Ð8Ø26×2SÑ2SˆGÐ.Ñ/à×(Ñ(ˆØÐ#Ø×%Ò%Ø(4�˜Ò%à)5�˜Ò&à×*Ñ*Ð6Ø.2×.KÑ.KˆGÐ*Ñ+à×%Ò%Ø,0×,FÑ,FˆGÐ(Ñ)àˆr8   c                óv   — i }| j                   �| j                   |d<   | j                  �| j                  |d<   |S )z4Return the options dict for the aggregation command.ÚmaxAwaitTimeMSr>   )rR   rS   )r[   r|   s     r6   Ú_command_optionszChangeStream._command_optionsÅ   sE   € àˆØ×"Ñ"Ð.Ø(,×(?Ñ(?ˆGÐ$Ñ%Ø×ÑÐ'Ø#'×#3Ñ#3ˆG�KÑ Øˆr8   c                óf   — | j                  «       }d|ig}|j                  | j                  «       |S )z;Return the full aggregation pipeline for this ChangeStream.z$changeStream)r}   ÚextendrL   )r[   r|   Úfull_pipelines      r6   Ú_aggregation_pipelinez"ChangeStream._aggregation_pipelineÎ   s5   € à×-Ñ-Ó/ˆØ0?ÀÐ/IÐ.JˆØ×Ñ˜TŸ^™^Ô,ØÐr8   c                ó  — |d   d   s�d|d   v r|d   d   | _         y| j                  €_| j                  du rP| j                  du rA|j                  dk\  r1|j                  d«      | _        | j                  €t        d|›�«      ‚yyyyyy)	aM  Callback that caches the postBatchResumeToken or
        startAtOperationTime from a changeStream aggregate command response
        containing an empty batch of change documents.

        This is implemented as a callback because we need access to the wire
        version in order to determine whether to cache this value.
        ÚcursorÚ
firstBatchÚpostBatchResumeTokenNFr    ÚoperationTimez?Expected field 'operationTime' missing from command response : )rQ   rU   rP   rO   Úmax_wire_versionÚgetr   )r[   ÚresultÚconns      r6   Ú_process_resultzChangeStream._process_resultÕ   s¾   € ð �hÑ Ò-Ø%¨°Ñ)9Ñ9Ø%+¨HÑ%5Ð6LÑ%M�Õ"à×-Ñ-Ð5Ø×+Ñ+¨uÑ4Ø×*Ñ*¨eÑ3Ø×)Ñ)¨QÒ.à06·
±
¸?Ó0K�Ô-à×0Ñ0Ð8Ü*ð&Ø&, Zð1óð ð 9ð	 /ð 4ð 5ð 6ð	 .r8   c                óL  — | j                  | j                  t        | j                  «       | j	                  «       | j
                  | j                  ¬«      }| j                  j                  |j                  | j                  j                  |«      |t        j                  ¬«      S )ztRun the full aggregation pipeline for this ChangeStream and return
        the corresponding CommandCursor.
        )Úresult_processorrd   )Ú	operation)rq   rI   r   r„   r€   rŽ   rW   rs   Ú_retryable_readÚ
get_cursorÚ_read_preference_forr   Ú	AGGREGATE)r[   rb   Úcmds      r6   Ú_run_aggregation_cmdz!ChangeStream._run_aggregation_cmdî   sŒ   € ð ×-Ñ-Ø�L‰LÜØ×&Ñ&Ó(Ø×!Ñ!Ó#Ø!×1Ñ1Ø—M‘Mð .ó 
ˆð �|‰|×+Ñ+Ø�N‰NØ�L‰L×-Ñ-¨gÓ6ØÜ—m‘mð	 ,ó 
ð 	
r8   c                óœ   — | j                   j                  | j                  «      5 }| j                  |¬«      cd d d «       S # 1 sw Y   y xY w)N)rb   )rs   Ú_tmp_sessionrV   r—   )r[   Úss     r6   rj   zChangeStream._create_cursor  s@   € Ø�\‰\×&Ñ& t§}¡}Ó5ð 	8¸Ø×,Ñ,°QÐ,Ó7÷	8÷ 	8ò 	8ús   ¦AÁAc                ó‚   — 	 | j                   j                  «        | j                  «       | _         y# t        $ r Y Œ!w xY w)z7Reestablish this change stream after a resumable error.N)rk   Úcloser   rj   rl   s    r6   Ú_resumezChangeStream._resume  s=   € ð	Ø�L‰L×ÑÔ ð ×*Ñ*Ó,ˆ�øô ò 	Ùð	ús   ‚2 ²	>½>c                óF   — d| _         | j                  j                  «        y)zClose this ChangeStream.TN)rX   rk   rœ   rl   s    r6   rœ   zChangeStream.close  s   € àˆŒØ�‰×ÑÕr8   c                ó   — | S ri   © rl   s    r6   Ú__iter__zChangeStream.__iter__  ó   € Øˆr8   c                ó@   — t        j                  | j                  «      S )zŒThe cached resume token that will be used to resume after the most
        recently returned change.

        .. versionadded:: 3.9
        )rJ   rK   rQ   rl   s    r6   r{   zChangeStream.resume_token  s   € ô �}‰}˜T×/Ñ/Ó0Ð0r8   c                óh   — | j                   r!| j                  «       }|�|S | j                   rŒ!t        ‚)aû  Advance the cursor.

        This method blocks until the next change document is returned or an
        unrecoverable error is raised. This method is used when iterating over
        all changes in the cursor. For example::

            try:
                resume_token = None
                pipeline = [{'$match': {'operationType': 'insert'}}]
                with db.collection.watch(pipeline) as stream:
                    for insert_change in stream:
                        print(insert_change)
                        resume_token = stream.resume_token
            except pymongo.errors.PyMongoError:
                # The ChangeStream encountered an unrecoverable error or the
                # resume attempt failed to recreate the cursor.
                if resume_token is None:
                    # There is no usable resume token because there was a
                    # failure during ChangeStream initialization.
                    logging.error('...')
                else:
                    # Use the interrupted ChangeStream's resume token to create
                    # a new ChangeStream. The new stream will continue from the
                    # last seen insert change without missing any events.
                    with db.collection.watch(
                            pipeline, resume_after=resume_token) as stream:
                        for insert_change in stream:
                            print(insert_change)

        Raises :exc:`StopIteration` if this ChangeStream is closed.
        )ÚaliveÚtry_nextÚStopIteration)r[   Údocs     r6   ÚnextzChangeStream.next  s2   € ðB �jŠjØ—-‘-“/ˆCØˆØ�
ð �j‹jô
 Ðr8   c                ó   — | j                    S )zøDoes this cursor have the potential to return more data?

        .. note:: Even if :attr:`alive` is ``True``, :meth:`next` can raise
            :exc:`StopIteration` and :meth:`try_next` can return ``None``.

        .. versionadded:: 3.8
        )rX   rl   s    r6   r¥   zChangeStream.aliveH  s   € ð —<‘<ÐÐr8   c                ó  — | j                   s&| j                  j                  s| j                  «        	 	 | j                  j	                  d«      }| j                  j                  sd| _         |€:| j                  j                  �"| j                  j                  | _        d| _        |S 	 |d   }| j                  j                  «       s,| j                  j                  r| j                  j                  }d| _        d| _        || _        d| _        | j$                  r t'        |j(                  | j*                  «      S |S # t
        $ rB}t        |«      s‚ | j                  «        | j                  j	                  d«      }Y d}~�Œ5d}~ww xY w# t
        $ r-}t        |«      s|j                  s| j                  «        ‚ d}~wt        $ r | j                  «        ‚ w xY w# t        $ r | j                  «        t        d«      d‚w xY w)a­  Advance the cursor without blocking indefinitely.

        This method returns the next change document without waiting
        indefinitely for the next change. For example::

            with db.collection.watch() as stream:
                while stream.alive:
                    change = stream.try_next()
                    # Note that the ChangeStream's resume token may be updated
                    # even when no changes are returned.
                    print("Current resume token: %r" % (stream.resume_token,))
                    if change is not None:
                        print("Change document: %r" % (change,))
                        continue
                    # We end up here when there are no recent changes.
                    # Sleep for a while before trying again to avoid flooding
                    # the server with getMore requests when no changes are
                    # available.
                    time.sleep(10)

        If no change document is cached locally then this method runs a single
        getMore command. If the getMore yields any documents, the next
        document is returned, otherwise, if the getMore returns no documents
        (because there have been no changes) then ``None`` is returned.

        :return: The next change document or ``None`` when no document is available
          after running a single getMore or when the cursor is closed.

        .. versionadded:: 3.8
        TFNÚ_idzECannot provide resume functionality when the resume token is missing.)rX   rk   r¥   r�   Ú	_try_nextr   r7   Útimeoutrœ   ÚBaseExceptionÚ_post_batch_resume_tokenrQ   rU   ÚKeyErrorr   Ú	_has_nextrO   rP   rD   r   ÚrawrE   )r[   Úchanger5   r{   s       r6   r¦   zChangeStream.try_nextS  sµ  € ð@ �|Š| D§L¡L×$6Ò$6Ø�L‰LŒNð	ð7ØŸ™×/Ñ/°Ó5�ð" �|‰|×!Ò!ØˆDŒLð ˆ>ð
 �|‰|×4Ñ4Ð@Ø%)§\¡\×%JÑ%J�Ô"Ø04�Ô-ØˆMð	Ø! %™=ˆLð �|‰|×%Ñ%Ô'¨D¯L©L×,QÒ,QØŸ<™<×@Ñ@ˆLð "'ˆÔØ"&ˆÔð *ˆÔØ(,ˆÔ%à×ÒÜ  §¡¨T×-EÑ-EÓFÐFØˆøôm  ò 7Ü! #”ØØ—‘”ØŸ™×/Ñ/°Ó6–ûð	7ûô
 ò 	ä˜c”?¨3¯;ª;Ø—
‘
”Øûäò 	Ø�J‰JŒLØð	ûô, ò 	Ø�J‰JŒLÜ"ØWóàðð	úsA   µD? Â*G Ä?	F
Å7FÅ?F ÆF
Æ
F Æ	GÆ(F>Æ>GÇ&Hc                ó   — | S ri   r    rl   s    r6   Ú	__enter__zChangeStream.__enter__³  r¢   r8   c                ó$   — | j                  «        y ri   )rœ   )r[   Úexc_typeÚexc_valÚexc_tbs       r6   Ú__exit__zChangeStream.__exit__¶  s   € Ø�
‰
�r8   )NNN)r\   zUUnion[MongoClient[_DocumentType], Database[_DocumentType], Collection[_DocumentType]]r<   zOptional[_Pipeline]r=   úOptional[str]r]   úOptional[Mapping[str, Any]]r^   úOptional[int]r_   r¾   r`   zOptional[_CollationIn]ra   zOptional[Timestamp]rb   úOptional[ClientSession]rc   r½   rd   zOptional[Any]re   r¼   rf   zOptional[bool]ÚreturnÚNone)rÀ   rÁ   )rÀ   zType[_AggregationCommand])rÀ   r+   ©rÀ   zdict[str, Any])rÀ   zlist[dict[str, Any]])rŒ   zMapping[str, Any]r�   r,   rÀ   rÁ   )rb   r¿   rÀ   r   )rÀ   r   )rÀ   zChangeStream[_DocumentType])rÀ   r½   )rÀ   r   )rÀ   Úbool)rÀ   zOptional[_DocumentType])r¸   r   r¹   r   rº   r   rÀ   rÁ   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__rg   rm   Úpropertyrq   rs   r}   r€   r„   rŽ   r—   rj   r�   rœ   r¡   r{   r   Úapplyr©   Ú__next__r¥   r¦   r¶   r»   r    r8   r6   r:   r:   Z   si  „ ñ	ð6 "&Ø59Ø/3ð%5:ð
ð5:ð &ð5:ð %ð5:ð 2ð5:ð )ð5:ð "ð5:ð *ð5:ð "5ð5:ð )ð5:ð 1ð5:ð  ð!5:ð" &3ð#5:ð$ -ð%5:ð& 
ó'5:ón-ð ò"ó ð"ð ò"ó ð"óó0óóó2
ó&8ó-óó
ð ò1ó ð1ð ‡[�[ò%ó ð%ðN €Hàò ó ð ð ‡[�[ò]ó ð]ó~ôr8   r:   c                  ó@   — e Zd ZU dZded<   edd„«       Zedd„«       Zy)	ÚCollectionChangeStreamzåA change stream that watches changes on a single collection.

    Should not be called directly by application developers. Use
    helper method :meth:`pymongo.collection.Collection.watch` instead.

    .. versionadded:: 3.7
    zCollection[_DocumentType]rI   c                ó   — t         S ri   )r   rl   s    r6   rq   z1CollectionChangeStream._aggregation_command_classÅ  s   € ä,Ð,r8   c                óB   — | j                   j                  j                  S ri   )rI   ÚdatabaseÚclientrl   s    r6   rs   zCollectionChangeStream._clientÉ  s   € à�|‰|×$Ñ$×+Ñ+Ð+r8   N)rÀ   z#Type[_CollectionAggregationCommand]©rÀ   zMongoClient[_DocumentType]©rÄ   rÅ   rÆ   rÇ   Ú__annotations__rÈ   rq   rs   r    r8   r6   rÌ   rÌ   º  s5   … ñð 'Ó&àò-ó ð-ð ò,ó ñ,r8   rÌ   c                  ó@   — e Zd ZU dZded<   edd„«       Zedd„«       Zy)	ÚDatabaseChangeStreamzëA change stream that watches changes on all collections in a database.

    Should not be called directly by application developers. Use
    helper method :meth:`pymongo.database.Database.watch` instead.

    .. versionadded:: 3.7
    zDatabase[_DocumentType]rI   c                ó   — t         S ri   )r   rl   s    r6   rq   z/DatabaseChangeStream._aggregation_command_classÙ  s   € ä*Ð*r8   c                ó.   — | j                   j                  S ri   )rI   rÐ   rl   s    r6   rs   zDatabaseChangeStream._clientÝ  s   € à�|‰|×"Ñ"Ð"r8   N)rÀ   z!Type[_DatabaseAggregationCommand]rÑ   rÒ   r    r8   r6   rÕ   rÕ   Î  s5   … ñð %Ó$àò+ó ð+ð ò#ó ñ#r8   rÕ   c                  ó$   ‡ — e Zd ZdZdˆ fd„Zˆ xZS )ÚClusterChangeStreamzóA change stream that watches changes on all collections in the cluster.

    Should not be called directly by application developers. Use
    helper method :meth:`pymongo.mongo_client.MongoClient.watch` instead.

    .. versionadded:: 3.7
    c                ó.   •— t         ‰| �  «       }d|d<   |S )NTÚallChangesForCluster)Úsuperr}   )r[   r|   Ú	__class__s     €r6   r}   z*ClusterChangeStream._change_stream_optionsë  s    ø€ Ü‘'Ñ0Ó2ˆØ*.ˆÐ&Ñ'Øˆr8   rÂ   )rÄ   rÅ   rÆ   rÇ   r}   Ú__classcell__)rÝ   s   @r6   rÙ   rÙ   â  s   ø„ ñ÷ñ r8   rÙ   )r5   r   rÀ   rÃ   )<rÇ   Ú
__future__r   rJ   Útypingr   r   r   r   r   r	   r
   Úbsonr   r   Úbson.raw_bsonr   Úbson.timestampr   Úpymongor   r   Úpymongo.collationr   Úpymongo.errorsr   r   r   r   r   Úpymongo.operationsr   Úpymongo.synchronous.aggregationr   r   r   Ú"pymongo.synchronous.command_cursorr   Úpymongo.typingsr   r   r   Ú_IS_SYNCÚ	frozensetr4   Ú"pymongo.synchronous.client_sessionr(   Úpymongo.synchronous.collectionr)   Úpymongo.synchronous.databaser*   Ú pymongo.synchronous.mongo_clientr+   Úpymongo.synchronous.poolr,   r7   r:   rÌ   rÕ   rÙ   r    r8   r6   ú<module>rò      sÅ   ðñ HÝ "ã ß N× NÑ Nç ,Ý )Ý $ß !Ý 8÷õ õ #÷ñ õ
 =ß BÑ Bà€ñ &òóÐ ñ. Ý@Ý9Ý5Ý<Ý3ó
ô]�7˜=Ñ)ô ]ô@,˜\¨-Ñ8ô ,ô(#˜<¨Ñ6ô #ô(Ð.¨}Ñ=õ r8   