Ë
    êëjÛ&  ã                  ó¶  — d 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 ddlmZ ddlmZ dd	lmZ dd
lmZmZmZmZ ddlmZmZmZmZ ddlmZmZ e	rddlm Z  ddl!m"Z" ddl#m$Z$ ejJ                  jL                  Z&ejJ                  jN                  Z'ejP                  jR                  Z) G d„ ded   «      Z* G d„ de«      Z+ G d„ de+«      Z, G d„ de,«      Z-y)z7
Objects to support the COPY protocol (async version).
é    )Úannotations)ÚABCÚabstractmethod)ÚTracebackType)ÚTYPE_CHECKINGÚAnyÚAsyncIteratorÚSequenceé   )Úerrors)Úpq)ÚSelf)ÚAQueueÚAWorkerÚagatherÚaspawn)ÚMAX_BUFFER_SIZEÚPREFER_FLUSHÚ
QUEUE_SIZEÚBaseCopy)Úcopy_endÚcopy_to)ÚBuffer)ÚAsyncCursor)ÚAsyncConnectionc                  óž   ‡ — e Zd ZU dZdZded<   dddœ	 	 	 	 	 dˆ f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ˆ xZS )Ú	AsyncCopyaj  Manage an asynchronous :sql:`COPY` operation.

    :param cursor: the cursor where the operation is performed.
    :param binary: if `!True`, write binary format.
    :param writer: the object to write to destination. If not specified, write
        to the `!cursor` connection.

    Choosing `!binary` is not necessary if the cursor has executed a
    :sql:`COPY` operation, because the operation result describes the format
    too. The parameter is useful when a `!Copy` object is created manually and
    no operation is performed on the cursor, such as when using ``writer=``\
    `~psycopg.copy.FileWriter`.
    ÚpsycopgÚAsyncWriterÚwriterN)Úbinaryr    c               ór   •— t         ‰| �  ||¬«       |st        |«      }|| _        |j                  | _        y )N)r!   )ÚsuperÚ__init__ÚAsyncLibpqWriterr    ÚwriteÚ_write)ÚselfÚcursorr!   r    Ú	__class__s       €úZ/var/www/html/kelly/kelly-backend/venv/lib/python3.12/site-packages/psycopg/_copy_async.pyr$   zAsyncCopy.__init__2   s6   ø€ ô 	‰Ñ˜¨ÐÔ/ÙÜ% fÓ-ˆFàˆŒØ—l‘lˆ�ó    c              ƒ  ó.   K  — | j                  «        | S ­w©N)Ú_enter©r(   s    r+   Ú
__aenter__zAsyncCopy.__aenter__@   s   è ø€ Ø�‰ŒØˆùs   ‚c              ƒ  óB   K  — | j                  |«      ƒ d {  –—†  y 7 Œ­wr.   )Úfinish)r(   Úexc_typeÚexc_valÚexc_tbs       r+   Ú	__aexit__zAsyncCopy.__aexit__D   s   è ø€ ð �k‰k˜'Ó"×"Ò"ús   ‚—˜c               óŠ   K  — | j                  «       ƒ d{  –—† x}r!|­–— | j                  «       ƒ d{  –—† x}rŒ yy7 Œ(7 Œ­w)z5Implement block-by-block iteration on :sql:`COPY TO`.N)Úread©r(   Údatas     r+   Ú	__aiter__zAsyncCopy.__aiter__N   s<   è ø€ à!ŸY™Y›[×(Ð)ˆdÐ)Ø‹Jð "ŸY™Y›[×(Ð)ˆdÓ)Ð(øÐ(úó#   ‚A–?—AµA¶A½AÁAc              ƒ  óp   K  — | j                   j                  | j                  «       «      ƒ d{  –—† S 7 Œ­w)zƒ
        Read an unparsed row after a :sql:`COPY TO` operation.

        Return an empty string when the data is finished.
        N)Ú
connectionÚwaitÚ	_read_genr0   s    r+   r9   zAsyncCopy.readS   s*   è ø€ ð —_‘_×)Ñ)¨$¯.©.Ó*:Ó;×;Ð;Ð;úó   ‚-6¯4°6c               óŠ   K  — | j                  «       ƒ d{  –—† x}�!|­–— | j                  «       ƒ d{  –—† x}�Œ yy7 Œ(7 Œ­w)zé
        Iterate on the result of a :sql:`COPY TO` operation record by record.

        Note that the records returned will be tuples of unparsed strings or
        bytes, unless data types are specified using `set_types()`.
        N)Úread_row)r(   Úrecords     r+   ÚrowszAsyncCopy.rows[   s>   è ø€ ð !%§¡£×/Ð0ˆvÐ=Ø‹Lð !%§¡£×/Ð0ˆvÓ=Ð/øÐ/úr=   c              ƒ  óp   K  — | j                   j                  | j                  «       «      ƒ d{  –—† S 7 Œ­w)a  
        Read a parsed row of data from a table after a :sql:`COPY TO` operation.

        Return `!None` when the data is finished.

        Note that the records returned will be tuples of unparsed strings or
        bytes, unless data types are specified using `set_types()`.
        N)r?   r@   Ú_read_row_genr0   s    r+   rD   zAsyncCopy.read_rowe   s,   è ø€ ð —_‘_×)Ñ)¨$×*<Ñ*<Ó*>Ó?×?Ð?Ð?úrB   c              ƒ  ó~   K  — | j                   j                  |«      x}r| j                  |«      ƒ d{  –—†  yy7 Œ­w)zÜ
        Write a block of data to a table after a :sql:`COPY FROM` operation.

        If the :sql:`COPY` is in binary format `!buffer` must be `!bytes`. In
        text mode it can be either `!bytes` or `!str`.
        N)Ú	formatterr&   r'   )r(   Úbufferr;   s      r+   r&   zAsyncCopy.writep   s<   è ø€ ð —>‘>×'Ñ'¨Ó/Ð/ˆ4Ð/Ø—+‘+˜dÓ#×#Ñ#ð 0Ø#úó   ‚2=´;µ=c              ƒ  ó~   K  — | j                   j                  |«      x}r| j                  |«      ƒ d{  –—†  yy7 Œ­w)z=Write a record to a table after a :sql:`COPY FROM` operation.N)rJ   Ú	write_rowr'   )r(   Úrowr;   s      r+   rN   zAsyncCopy.write_rowz   s:   è ø€ à—>‘>×+Ñ+¨CÓ0Ð0ˆ4Ð0Ø—+‘+˜dÓ#×#Ñ#ð 1Ø#úrL   c              ƒ  óî  K  — | j                   t        k(  rb|s5| j                  j                  «       x}r| j	                  |«      ƒ d{  –—†  | j
                  j                  |«      ƒ d{  –—†  d| _        y|sy| j                  j                  t        k7  ry| j                  j                  «       ƒ d{  –—†  | j                  j                  | j                  «       «      ƒ d{  –—†  y7 Œ¤7 Œƒ7 Œ:7 Œ­w)a  Terminate the copy operation and free the resources allocated.

        You shouldn't need to call this function yourself: it is usually called
        by exit. It is available if, despite what is documented, you end up
        using the `Copy` object outside a block.
        NT)Ú
_directionÚCOPY_INrJ   Úendr'   r    r3   Ú	_finishedÚ_pgconnÚtransaction_statusÚACTIVEr?   Ú_try_cancelr@   Ú_end_copy_out_gen)r(   Úexcr;   s      r+   r3   zAsyncCopy.finish   sÉ   è ø€ ð �?‰?œgÒ%ÙØŸ>™>×-Ñ-Ó/Ð/�4Ð/ØŸ+™+ dÓ+×+Ð+Ø—+‘+×$Ñ$ SÓ)×)Ð)Ø!ˆD�NáØà�|‰|×.Ñ.´&Ò8ð ð —/‘/×-Ñ-Ó/×/Ð/Ø—/‘/×&Ñ& t×'=Ñ'=Ó'?Ó@×@Ñ@ð# ,øØ)øð 0øØ@úsI   ‚AC5ÁC-Á	"C5Á+C/Á,A
C5Â6C1Â70C5Ã'C3Ã(C5Ã/C5Ã1C5Ã3C5)r)   úAsyncCursor[Any]r!   zbool | Noner    zAsyncWriter | None)Úreturnr   )r4   ztype[BaseException] | Noner5   úBaseException | Noner6   zTracebackType | Noner\   ÚNone)r\   zAsyncIterator[Buffer])r\   r   )r\   zAsyncIterator[tuple[Any, ...]])r\   ztuple[Any, ...] | None)rK   zBuffer | strr\   r^   )rO   zSequence[Any]r\   r^   ©rZ   r]   r\   r^   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú__annotations__r$   r1   r7   r<   r9   rF   rD   r&   rN   r3   Ú__classcell__©r*   s   @r+   r   r      s“   ø… ñð €JàÓð #Ø%)ñ#à ð#ð ð	#ð
 #õ#óð#à,ð#ð &ð#ð %ð	#ð
 
ó#óó
<óó	@ó$ó$÷
Ar,   r   zAsyncConnection[Any]c                  ó,   — e Zd ZdZedd„«       Zddd„Zy)r   zG
    A class to write copy data somewhere (for async connections).
    c              ƒ  ó   K  — y­w)zWrite some data to destination.N© r:   s     r+   r&   zAsyncWriter.write¢   s   è ø€ ð 	ùó   ‚Nc              ƒ  ó   K  — y­w)z‰
        Called when write operations are finished.

        If operations finished with an error, it will be passed to ``exc``.
        Nri   )r(   rZ   s     r+   r3   zAsyncWriter.finish§   s   è ø€ ð 	ùrj   ©r;   r   r\   r^   r.   r_   )r`   ra   rb   rc   r   r&   r3   ri   r,   r+   r   r   �   s    „ ñð òó ðõr,   r   c                  ó.   — e Zd ZdZdZdd„Zdd„Zd	d
d„Zy)r%   zE
    An `AsyncWriter` to write copy data to a Postgres database.
    úpsycopg.copyc                ój   — || _         |j                  | _        | j                  j                  | _        y r.   )r)   r?   ÚpgconnrU   )r(   r)   s     r+   r$   zAsyncLibpqWriter.__init__·   s'   € ØˆŒØ ×+Ñ+ˆŒØ—‘×-Ñ-ˆ�r,   c           
   ƒ  ó€  K  — t        |«      t        k  r>| j                  j                  t	        | j
                  |t        ¬«      «      ƒ d {  –—†  y t        dt        |«      t        «      D ]I  }| j                  j                  t	        | j
                  |||t        z    t        ¬«      «      ƒ d {  –—†  ŒK y 7 Œl7 Œ	­w)N©Úflushr   )Úlenr   r?   r@   r   rU   r   Úrange©r(   r;   Úis      r+   r&   zAsyncLibpqWriter.write¼   s™   è ø€ Üˆt‹9œÒ'ð —/‘/×&Ñ&¤w¨t¯|©|¸TÌÔ'VÓW×WÑWô ˜1œc $›i¬Ö9�Ø—o‘o×*Ñ*ÜØŸ™ d¨1¨q´?Ñ/BÐ&CÌ<ôó÷ ñ ñ :ð	 Xøð
ús%   ‚AB>ÁB:ÁA$B>Â2B<Â3B>Â<B>Nc              ƒ  óh  K  — |rBdt        |«      j                  › d|› �}|j                  | j                  j                  d«      }nd }	 | j
                  j                  t        | j                  |«      «      ƒ d {  –—† }|g| j                  _	        y 7 Œ# t        j                  $ r |s‚ Y y w xY w­w)Nzerror from Python: z - Úreplace)Útyperb   ÚencoderU   Ú	_encodingr?   r@   r   r)   Ú_resultsÚeÚQueryCanceled)r(   rZ   ÚmsgÚbmsgÚress        r+   r3   zAsyncLibpqWriter.finishË   s¡   è ø€ áØ'¬¨S«	×(>Ñ(>Ð'?¸sÀ3À%ÐHˆCØ—:‘:˜dŸl™l×4Ñ4°iÓ@‰DàˆDð		)ØŸ™×,Ñ,¬X°d·l±lÀDÓ-IÓJ×JˆCð %( 5ˆD�K‰KÕ ð Kùô �‰ò 	ÙØñ ð	üs<   ‚AB2Á
2B Á<BÁ=B ÂB2ÂB ÂB/Â,B2Â.B/Â/B2©r)   r[   rl   r.   r_   )r`   ra   rb   rc   r$   r&   r3   ri   r,   r+   r%   r%   °   s   „ ñð  €Jó.ó
õ)r,   r%   c                  óF   ‡ — e Zd ZdZdZdˆ fd„Zdd„Zd	d„Zd
dˆ fd„Zˆ xZS )ÚAsyncQueuedLibpqWriterzò
    `AsyncWriter` using a buffer to queue data to write.

    `write()` returns immediately, so that the main thread can be CPU-bound
    formatting messages, while a worker thread can be IO-bound waiting to write
    on the connection.
    rn   c                ój   •— t         ‰| �  |«       t        t        ¬«      | _        d | _        d | _        y )N)Úmaxsize)r#   r$   r   r   Ú_queueÚ_workerÚ_worker_error)r(   r)   r*   s     €r+   r$   zAsyncQueuedLibpqWriter.__init__ê   s+   ø€ Ü‰Ñ˜Ô ä&,´ZÔ&@ˆŒØ'+ˆŒØ37ˆÕr,   c              ƒ  ób  K  — 	 | j                   j                  «       ƒ d{  –—† x}rc| j                  j                  t	        | j
                  |t        ¬«      «      ƒ d{  –—†  | j                   j                  «       ƒ d{  –—† x}rŒbyy7 Œj7 Œ-7 Œ# t        $ r}|| _        Y d}~yd}~ww xY w­w)zåPush data to the server when available from the copy queue.

        Terminate reading when the queue receives a false-y value, or in case
        of error.

        The function is designed to be run in a separate task.
        Nrr   )	rˆ   Úgetr?   r@   r   rU   r   ÚBaseExceptionrŠ   )r(   r;   Úexs      r+   ÚworkerzAsyncQueuedLibpqWriter.workerñ   s–   è ø€ ð	$Ø!%§¡§¡Ó!2×2Ð3�$Ð3Ø—o‘o×*Ñ*Ü˜DŸL™L¨$´lÔCó÷ ð ð "&§¡§¡Ó!2×2Ð3�$Ó3Ð2øðøð 3ùô ò 	$à!#ˆD×Ñûð	$üsb   ‚B/„B ¡B¢>B Á BÁ!!B ÂBÂB Â
B/ÂB ÂB ÂB Â	B,ÂB'Â"B/Â'B,Â,B/c              ƒ  ó”  K  — | j                   st        | j                  «      | _         | j                  r| j                  ‚t	        |«      t
        k  r$| j                  j                  |«      ƒ d {  –—†  y t        dt	        |«      t
        «      D ]/  }| j                  j                  |||t
        z    «      ƒ d {  –—†  Œ1 y 7 ŒR7 Œ	­w)Nr   )	r‰   r   r�   rŠ   rt   r   rˆ   Úputru   rv   s      r+   r&   zAsyncQueuedLibpqWriter.write  sŸ   è ø€ Ø�|Š|ä! $§+¡+Ó.ˆDŒLð ×ÒØ×$Ñ$Ð$äˆt‹9œÒ'ð —+‘+—/‘/ $Ó'×'Ñ'ô ˜1œc $›i¬Ö9�Ø—k‘k—o‘o d¨1¨q´?Ñ/BÐ&CÓD×DÑDñ :ð	 (øð
 Eús%   ‚A/CÁ1CÁ2A
CÂ<CÂ=CÃCc              ƒ  ó  •K  — | j                   j                  d«      ƒ d {  –—†  | j                  r$t        | j                  «      ƒ d {  –—†  d | _        | j                  r| j                  ‚t
        ‰| �  |«      ƒ d {  –—†  y 7 Œd7 Œ=7 Œ	­w)Nr,   )rˆ   r‘   r‰   r   rŠ   r#   r3   )r(   rZ   r*   s     €r+   r3   zAsyncQueuedLibpqWriter.finish  sw   øè ø€ Ø�k‰k�o‰o˜cÓ"×"Ð"à�<Š<Ü˜$Ÿ,™,Ó'×'Ð'ØˆDŒLð ×ÒØ×$Ñ$Ð$ä‰g‰n˜SÓ!×!Ñ!ð 	#øð (øð 	"ús3   ƒB¢B£(BÁB	Á5BÂBÂBÂ	BÂBrƒ   )r\   r^   rl   r.   r_   )	r`   ra   rb   rc   r$   r�   r&   r3   re   rf   s   @r+   r…   r…   ß   s)   ø„ ñð  €Jõ8ó$ó"E÷&"ò "r,   r…   N).rc   Ú
__future__r   Úabcr   r   Útypesr   Útypingr   r   r	   r
   Ú r   r~   r   Ú_compatr   Ú_acompatr   r   r   r   Ú
_copy_baser   r   r   r   Ú
generatorsr   r   r   Úcursor_asyncr   Úconnection_asyncr   Ú
ExecStatusrR   ÚCOPY_OUTÚTransactionStatusrW   r   r   r%   r…   ri   r,   r+   Ú<module>r¡      s­   ðñõ #ç #Ý ß >Ó >å Ý Ý ß 6Ó 6ß KÓ Kß )áÝÝ)Ý1à
�-‰-×
Ñ
€Ø�=‰=×!Ñ!€à	×	Ñ	×	$Ñ	$€ô{A�Ð/Ñ0ô {Aô|�#ô ô&,)�{ô ,)ô^A"Ð-õ A"r,   