o
    ØRWj9  ã                   @  s2  d dl mZ d dlZd dlZd dlmZ d dlmZmZm	Z	m
Z
mZmZmZmZ d dlmZmZmZmZmZmZ d dlZddlmZmZ ddlmZ erZdd	lmZmZ dd
lm Z  e
dƒZ!G dd„ de	e! ƒZ"G dd„ de	e! ƒZ#G dd„ dƒZ$G dd„ dƒZ%eG dd„ deƒƒZ&d"dd„Z'ddœd#d d!„Z(dS )$é    )ÚannotationsN)ÚTracebackType)ÚTYPE_CHECKINGÚAnyÚGenericÚTypeVarÚIteratorÚOptionalÚAsyncIteratorÚcast)ÚSelfÚProtocolÚ	TypeGuardÚoverrideÚ
get_originÚruntime_checkableé   )Ú
is_mappingÚextract_type_var_from_base)ÚAPIError)ÚOpenAIÚAsyncOpenAI)ÚFinalRequestOptionsÚ_Tc                   @  ó„   e Zd ZU dZded< dZded< ded< dd	œd+dd„Zd,dd„Zd-dd„Zd.dd„Z	d-dd„Z
d/dd „Zd0d'd(„Zd1d)d*„ZdS )2ÚStreamzJProvides the core interface to iterate over a synchronous stream response.úhttpx.ResponseÚresponseNúOptional[FinalRequestOptions]Ú_optionsÚSSEBytesDecoderÚ_decoder©ÚoptionsÚcast_toútype[_T]Úclientr   r#   ÚreturnÚNonec                C  ó0   || _ || _|| _|| _| ¡ | _|  ¡ | _d S ©N©r   Ú_cast_toÚ_clientr   Ú_make_sse_decoderr!   Ú
__stream__Ú	_iterator©Úselfr$   r   r&   r#   © r3   ú`/home/esfera/Documents/content_generation/venv/lib/python3.10/site-packages/openai/_streaming.pyÚ__init__   ó   
zStream.__init__r   c                 C  s
   | j  ¡ S r*   )r0   Ú__next__©r2   r3   r3   r4   r7   -   s   
zStream.__next__úIterator[_T]c                 c  s   � | j D ]}|V  qd S r*   ©r0   ©r2   Úitemr3   r3   r4   Ú__iter__0   s   €
ÿzStream.__iter__úIterator[ServerSentEvent]c                 c  s   � | j  | j ¡ ¡E d H  d S r*   )r!   Ú
iter_bytesr   r8   r3   r3   r4   Ú_iter_events4   s   €zStream._iter_eventsc           	      c  sŽ  � t t| jƒ}| j}| jj}|  ¡ }z¯|D ]ž}|j d¡r nœ|j	rk|j	 d¡rk| 
¡ }|j	dkr^t|ƒr^| d¡r^d }| d¡}t|ƒrJ| d¡}|rQt|tƒsSd}t|| jj|d d�‚|||j	dœ||d�V  q| 
¡ }t|ƒrœ| d¡rœd }| d¡}t|ƒrˆ| d¡}|r�t|tƒs‘d}t|| jj|d d�‚|| jd ur¬| jjr¬||j	dœn|||d�V  qW | ¡  d S W | ¡  d S | ¡  w ©	Nz[DONE]zthread.ÚerrorÚmessagez"An error occurred during streaming)rC   ÚrequestÚbody)ÚdataÚevent)rF   r$   r   )r   r   r,   r   r-   Ú_process_response_datar@   rF   Ú
startswithrG   Újsonr   ÚgetÚ
isinstanceÚstrr   rD   r   Úsynthesize_event_and_dataÚclose©	r2   r$   r   Úprocess_dataÚiteratorÚsserF   rC   rB   r3   r3   r4   r/   7   s`   €

ý

ýÿ
ûÙ0Ò.zStream.__stream__r   c                 C  s   | S r*   r3   r8   r3   r3   r4   Ú	__enter__p   s   zStream.__enter__Úexc_typeútype[BaseException] | NoneÚexcúBaseException | NoneÚexc_tbúTracebackType | Nonec                 C  s   |   ¡  d S r*   ©rO   ©r2   rU   rW   rY   r3   r3   r4   Ú__exit__s   s   zStream.__exit__c                 C  s   | j  ¡  dS ©zŠ
        Close the response and release the connection.

        Automatically called if the response body is read to completion.
        N)r   rO   r8   r3   r3   r4   rO   {   s   zStream.close)
r$   r%   r   r   r&   r   r#   r   r'   r(   ©r'   r   )r'   r9   )r'   r>   ©r'   r   ©rU   rV   rW   rX   rY   rZ   r'   r(   ©r'   r(   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú__annotations__r   r5   r7   r=   r@   r/   rT   r]   rO   r3   r3   r3   r4   r      s   
 ú




9
r   c                   @  r   )2ÚAsyncStreamzLProvides the core interface to iterate over an asynchronous stream response.r   r   Nr   r   zSSEDecoder | SSEBytesDecoderr!   r"   r$   r%   r&   r   r#   r'   r(   c                C  r)   r*   r+   r1   r3   r3   r4   r5   ‹   r6   zAsyncStream.__init__r   c                 Ã  s   �| j  ¡ I d H S r*   )r0   Ú	__anext__r8   r3   r3   r4   ri   š   s   €zAsyncStream.__anext__úAsyncIterator[_T]c                 C s"   �| j 2 z	3 d H W }|V  q6 d S r*   r:   r;   r3   r3   r4   Ú	__aiter__�   s   €ÿzAsyncStream.__aiter__úAsyncIterator[ServerSentEvent]c                 C s.   �| j  | j ¡ ¡2 z	3 d H W }|V  q
6 d S r*   )r!   Úaiter_bytesr   )r2   rS   r3   r3   r4   r@   ¡   s   €ÿzAsyncStream._iter_eventsc           	      C sª  �t t| jƒ}| j}| jj}|  ¡ }zº|2 z¢3 d H W }|j d¡r# n |j	ro|j	 d¡ro| 
¡ }|j	dkrbt|ƒrb| d¡rbd }| d¡}t|ƒrN| d¡}|rUt|tƒsWd}t|| jj|d d�‚|||j	dœ||d�V  q| 
¡ }t|ƒr | d¡r d }| d¡}t|ƒrŒ| d¡}|r“t|tƒs•d}t|| jj|d d�‚|| jd ur°| jjr°||j	dœn|||d�V  q6 W | ¡ I d H  d S W | ¡ I d H  d S | ¡ I d H  w rA   )r   r   r,   r   r-   rH   r@   rF   rI   rG   rJ   r   rK   rL   rM   r   rD   r   rN   ÚacloserP   r3   r3   r4   r/   ¥   s`   €

ý

ýÿ
ûÙ0Ò".zAsyncStream.__stream__r   c                 Ã  s   �| S r*   r3   r8   r3   r3   r4   Ú
__aenter__Þ   s   €zAsyncStream.__aenter__rU   rV   rW   rX   rY   rZ   c                 Ã  s   �|   ¡ I d H  d S r*   r[   r\   r3   r3   r4   Ú	__aexit__á   s   €zAsyncStream.__aexit__c                 Ã  s   �| j  ¡ I dH  dS r^   )r   rn   r8   r3   r3   r4   rO   é   s   €zAsyncStream.close)
r$   r%   r   r   r&   r   r#   r   r'   r(   r_   )r'   rj   )r'   rl   r`   ra   rb   )rc   rd   re   rf   rg   r   r5   ri   rk   r@   r/   ro   rp   rO   r3   r3   r3   r4   rh   „   s   
 ú




9
rh   c                   @  sr   e Zd Zdddddœddd„Zeddd„ƒZeddd„ƒZeddd„ƒZeddd„ƒZddd„Z	e
ddd„ƒZdS ) ÚServerSentEventN©rG   rF   ÚidÚretryrG   ú
str | NonerF   rs   rt   ú
int | Noner'   r(   c                C  s,   |d u rd}|| _ || _|pd | _|| _d S )NÚ )Ú_idÚ_dataÚ_eventÚ_retry)r2   rG   rF   rs   rt   r3   r3   r4   r5   ó   s   

zServerSentEvent.__init__c                 C  ó   | j S r*   )rz   r8   r3   r3   r4   rG     ó   zServerSentEvent.eventc                 C  r|   r*   )rx   r8   r3   r3   r4   rs     r}   zServerSentEvent.idc                 C  r|   r*   )r{   r8   r3   r3   r4   rt     r}   zServerSentEvent.retryrM   c                 C  r|   r*   )ry   r8   r3   r3   r4   rF     r}   zServerSentEvent.datar   c                 C  s   t  | j¡S r*   )rJ   ÚloadsrF   r8   r3   r3   r4   rJ     s   zServerSentEvent.jsonc              	   C  s&   d| j › d| j› d| j› d| j› d�	S )NzServerSentEvent(event=z, data=z, id=z, retry=ú)rr   r8   r3   r3   r4   Ú__repr__  s   &zServerSentEvent.__repr__)
rG   ru   rF   ru   rs   ru   rt   rv   r'   r(   )r'   ru   )r'   rv   )r'   rM   )r'   r   )rc   rd   re   r5   ÚpropertyrG   rs   rt   rF   rJ   r   r€   r3   r3   r3   r4   rq   ò   s"    ú
rq   c                   @  sj   e Zd ZU ded< ded< ded< ded< dd
d„Zd dd„Zd!dd„Zd"dd„Zd#dd„Zd$dd„Z	dS )%Ú
SSEDecoderz	list[str]ry   ru   rz   rv   r{   Ú_last_event_idr'   r(   c                 C  s   d | _ g | _d | _d | _d S r*   )rz   ry   rƒ   r{   r8   r3   r3   r4   r5   !  s   
zSSEDecoder.__init__rR   úIterator[bytes]r>   c                 c  sB   � |   |¡D ]}| ¡ D ]}| d¡}|  |¡}|r|V  qqdS )ú^Given an iterator that yields raw binary data, iterate over it & yield every event encounteredúutf-8N)Ú_iter_chunksÚ
splitlinesÚdecode©r2   rR   ÚchunkÚraw_lineÚlinerS   r3   r3   r4   r?   '  s   €

€üþzSSEDecoder.iter_bytesc                 c  sP   � d}|D ]}|j dd�D ]}||7 }| d¡r|V  d}qq|r&|V  dS dS )ú^Given an iterator that yields raw binary data, iterate over it and yield individual SSE chunksó    T©Úkeepends©s   s   

s   

N©rˆ   Úendswith©r2   rR   rF   r‹   r�   r3   r3   r4   r‡   1  s   €
€ü
ÿzSSEDecoder._iter_chunksúAsyncIterator[bytes]rl   c                 C sL   �|   |¡2 z3 dH W }| ¡ D ]}| d¡}|  |¡}|r!|V  qq6 dS )r…   Nr†   )Ú_aiter_chunksrˆ   r‰   rŠ   r3   r3   r4   rm   =  s   €

€üþzSSEDecoder.aiter_bytesc                 C sZ   �d}|2 z3 dH W }|j dd�D ]}||7 }| d¡r!|V  d}qq6 |r+|V  dS dS )rŽ   r�   NTr�   r’   r“   r•   r3   r3   r4   r—   G  s   €
€üÿ
ÿzSSEDecoder._aiter_chunksr�   rM   úServerSentEvent | Nonec              	   C  s  |s,| j s| js| js| jd u rd S t| j d | j¡| j| jd�}d | _ g | _d | _|S | d¡r3d S | d¡\}}}| d¡rF|dd … }|dkrO|| _ d S |dkr[| j |¡ d S |dkrkd	|v rf	 d S || _d S |d
kr„zt	|ƒ| _W d S  t
tfyƒ   Y d S w 	 d S )NÚ
rr   ú:ú r   rG   rF   rs   ú rt   )rz   ry   rƒ   r{   rq   ÚjoinrI   Ú	partitionÚappendÚintÚ	TypeErrorÚ
ValueError)r2   r�   rS   Ú	fieldnameÚ_Úvaluer3   r3   r4   r‰   S  sP   
ü

ñó÷	øûûzSSEDecoder.decodeNrb   ©rR   r„   r'   r>   )rR   r„   r'   r„   ©rR   r–   r'   rl   )rR   r–   r'   r–   )r�   rM   r'   r˜   )
rc   rd   re   rg   r5   r?   r‡   rm   r—   r‰   r3   r3   r3   r4   r‚     s   
 






r‚   c                   @  s    e Zd Zddd„Zdd	d
„ZdS )r    rR   r„   r'   r>   c                 C  ó   dS )r…   Nr3   ©r2   rR   r3   r3   r4   r?   †  ó   zSSEBytesDecoder.iter_bytesr–   rl   c                 C  r¨   )zdGiven an async iterator that yields raw binary data, iterate over it & yield every event encounteredNr3   r©   r3   r3   r4   rm   Š  rª   zSSEBytesDecoder.aiter_bytesNr¦   r§   )rc   rd   re   r?   rm   r3   r3   r3   r4   r    „  s    
r    ÚtypÚtyper'   ú;TypeGuard[type[Stream[object]] | type[AsyncStream[object]]]c                 C  s$   t | ƒp| }t |¡ot|ttfƒS )zaTypeGuard for determining whether or not the given type is a subclass of `Stream` / `AsyncStream`)r   ÚinspectÚisclassÚ
issubclassr   rh   )r«   Úoriginr3   r3   r4   Úis_stream_class_type�  s   r²   )Úfailure_messageÚ
stream_clsr³   ru   c                C  s*   ddl m}m} t| dtd||fƒ|d�S )a  Given a type like `Stream[T]`, returns the generic type variable `T`.

    This also handles the case where a concrete subclass is given, e.g.
    ```py
    class MyStream(Stream[bytes]):
        ...

    extract_stream_chunk_type(MyStream) -> bytes
    ```
    r   )r   rh   r   ztuple[type, ...])ÚindexÚgeneric_basesr³   )Ú_base_clientr   rh   r   r   )r´   r³   r   rh   r3   r3   r4   Úextract_stream_chunk_type•  s   ür¸   )r«   r¬   r'   r­   )r´   r¬   r³   ru   r'   r¬   ))Ú
__future__r   rJ   r®   Útypesr   Útypingr   r   r   r   r   r	   r
   r   Útyping_extensionsr   r   r   r   r   r   ÚhttpxÚ_utilsr   r   Ú_exceptionsr   r-   r   r   Ú_modelsr   r   r   rh   rq   r‚   r    r²   r¸   r3   r3   r3   r4   Ú<module>   s,   ( mn)i

	ý