o
    î6Wj}Q  ã                   @  sb  d dl mZ d dlZd dl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 e	r\ddlmZmZ dd	lm Z  ed
ƒZ!G dd„ dej"ƒZ#G dd„ dee! e#d�Z$G dd„ dej"ƒZ%G dd„ dee! e%d�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_dictÚextract_type_var_from_base)Ú	AnthropicÚAsyncAnthropic)ÚFinalRequestOptionsÚ_Tc                   @  ó   e Zd Zeddd„ƒZdS )	Ú_SyncStreamMetaÚinstancer   ÚreturnÚboolc                 C  ó.   ddl m} t||ƒrtjdtdd� dS dS )Nr   )ÚMessageStreamzŽUsing `isinstance()` to check if a `MessageStream` object is an instance of `Stream` is deprecated & will be removed in the next major versioné   ©Ú
stacklevelTF)Úlib.streamingr   Ú
isinstanceÚwarningsÚwarnÚDeprecationWarning)Úselfr   r   © r)   úc/home/esfera/Documents/content_generation/venv/lib/python3.10/site-packages/anthropic/_streaming.pyÚ__instancecheck__   ó   
ýz!_SyncStreamMeta.__instancecheck__N©r   r   r   r   ©Ú__name__Ú
__module__Ú__qualname__r   r+   r)   r)   r)   r*   r      ó    r   c                   @  ó’   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d0dd„Z	e
d1dd„ƒZd/dd„Zd2d!d"„Zd3d)d*„Zd4d+d,„ZdS )5Ú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<   r   ÚNonec                C  ó0   || _ || _|| _|| _| ¡ | _|  ¡ | _d S ©N©r6   Ú_cast_toÚ_clientr8   Ú_make_sse_decoderr:   Ú
__stream__Ú	_iterator©r(   r=   r6   r?   r<   r)   r)   r*   Ú__init__4   ó   
zStream.__init__r   c                 C  s
   | j  ¡ S rB   )rH   Ú__next__©r(   r)   r)   r*   rL   C   s   
zStream.__next__úIterator[_T]c                 c  s   � | j D ]}|V  qd S rB   ©rH   ©r(   Úitemr)   r)   r*   Ú__iter__F   s   €
ÿzStream.__iter__úIterator[ServerSentEvent]c                 c  s   � | j  | j ¡ ¡E d H  d S rB   )r:   Ú
iter_bytesr6   rM   r)   r)   r*   Ú_iter_eventsJ   s   €zStream._iter_eventsc                 C  ó   t ƒ  |  ¡ ¡S ©zÂIterate the raw Server-Sent Events from `response`, before any JSON
        parsing or event-name filtering.

        This reads the response body directly, so the response is consumed.
        )Ú
SSEDecoderrT   ©r6   r)   r)   r*   Ú
raw_eventsM   ó   zStream.raw_eventsc           	   	   c  s8  � t t| jƒ}| j}| jj}|  ¡ }�zƒ|D �]x}|jdkr(|| ¡ ||d�V  |jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jd	k�s<|jd
k�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jdk�s<|jd k�s<|jd!k�s<|jd"k�s<|jd#k�s<|jd$k�s<|jd%k�s<|jd&k�s<|jd'k�s<|jd(k�s<|jd)k�s<|jd*k�s<|jd+k�s<|jd,k�s<|jd-k�s<|jd.k�s<|jd/k�s<|jd0k�rW| ¡ }t	|ƒ�rOd1|v�rO|j|d1< ||||d�V  |jd2k�r^q|jd3k�r�|j
}z	| ¡ }|› }W n t�y„   |j
�p�d4|j› �}Y nw | jj||| jd5�‚qW | ¡  d S | ¡  w ©6NÚ
completion)Údatar=   r6   Úmessage_startÚmessage_deltaÚmessage_stopÚcontent_block_startÚcontent_block_deltaÚcontent_block_stopÚmessagezuser.messagezuser.interruptzuser.tool_confirmationzuser.custom_tool_resultzuser.tool_resultzagent.messagezagent.thinkingzagent.tool_usezagent.tool_resultzagent.mcp_tool_usezagent.mcp_tool_resultzagent.custom_tool_usezagent.thread_context_compactedzsession.status_runningzsession.status_idlezsession.status_rescheduledzsession.status_terminatedzsession.errorzsession.deletedzsession.updatedzspan.model_request_startzspan.model_request_endzspan.outcome_evaluation_startzspan.outcome_evaluation_ongoingzspan.outcome_evaluation_endzuser.define_outcomezagent.thread_message_receivedzagent.thread_message_sentz%agent.session_thread_message_receivedz!agent.session_thread_message_sentzsession.thread_createdzsession.thread_status_createdzsession.thread_status_runningzsession.thread_status_idlez!session.thread_status_rescheduledz session.thread_status_terminatedÚevent_startÚevent_deltazsystem.messageÚtypeÚpingÚerrorzError code: )Úbodyr6   )r   r   rD   r6   rE   Ú_process_response_datarU   ÚeventÚjsonr   r^   Ú	ExceptionÚstatus_codeÚ_make_status_errorÚclose©	r(   r=   r6   Úprocess_dataÚiteratorÚsser^   rk   Úerr_msgr)   r)   r*   rG   V   sš   €



ÿý÷ÃMzStream.__stream__r   c                 C  s   | S rB   r)   rM   r)   r)   r*   Ú	__enter__¬   s   zStream.__enter__Úexc_typeútype[BaseException] | NoneÚexcúBaseException | NoneÚexc_tbúTracebackType | Nonec                 C  s   |   ¡  d S rB   ©rr   ©r(   ry   r{   r}   r)   r)   r*   Ú__exit__¯   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)r6   rr   rM   r)   r)   r*   rr   ·   s   zStream.close)
r=   r>   r6   r5   r?   r   r<   r7   r   r@   ©r   r   )r   rN   )r   rS   )r6   r5   r   rS   ©r   r   ©ry   rz   r{   r|   r}   r~   r   r@   ©r   r@   )r/   r0   r1   Ú__doc__Ú__annotations__r8   rJ   rL   rR   rU   ÚstaticmethodrZ   rG   rx   r�   rr   r)   r)   r)   r*   r4   -   s    
 ú




V
r4   )Ú	metaclassc                   @  r   )	Ú_AsyncStreamMetar   r   r   r   c                 C  r   )Nr   )ÚAsyncMessageStreamz˜Using `isinstance()` to check if a `AsyncMessageStream` object is an instance of `AsyncStream` is deprecated & will be removed in the next major versionr    r!   TF)r#   rŒ   r$   r%   r&   r'   )r(   r   rŒ   r)   r)   r*   r+   Á   r,   z"_AsyncStreamMeta.__instancecheck__Nr-   r.   r)   r)   r)   r*   r‹   À   r2   r‹   c                   @  r3   )5ÚAsyncStreamzLProvides the core interface to iterate over an asynchronous stream response.r5   r6   Nr7   r8   zSSEDecoder | SSEBytesDecoderr:   r;   r=   r>   r?   r   r<   r   r@   c                C  rA   rB   rC   rI   r)   r)   r*   rJ   Ü   rK   zAsyncStream.__init__r   c                 Ã  s   �| j  ¡ I d H S rB   )rH   Ú	__anext__rM   r)   r)   r*   rŽ   ë   s   €zAsyncStream.__anext__úAsyncIterator[_T]c                 C s"   �| j 2 z	3 d H W }|V  q6 d S rB   rO   rP   r)   r)   r*   Ú	__aiter__î   s   €ÿzAsyncStream.__aiter__úAsyncIterator[ServerSentEvent]c                 C s.   �| j  | j ¡ ¡2 z	3 d H W }|V  q
6 d S rB   )r:   Úaiter_bytesr6   )r(   rv   r)   r)   r*   rU   ò   s   €ÿzAsyncStream._iter_eventsc                 C  rV   rW   )rX   r’   rY   r)   r)   r*   rZ   ö   r[   zAsyncStream.raw_eventsc           	   	   C sN  �t t| jƒ}| j}| jj}|  ¡ }�z‹|2 �z|3 d H W }|jdkr,|| ¡ ||d�V  |jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jd	k�s@|jd
k�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jdk�s@|jd k�s@|jd!k�s@|jd"k�s@|jd#k�s@|jd$k�s@|jd%k�s@|jd&k�s@|jd'k�s@|jd(k�s@|jd)k�s@|jd*k�s@|jd+k�s@|jd,k�s@|jd-k�s@|jd.k�s@|jd/k�s@|jd0k�r[| ¡ }t	|ƒ�rSd1|v�rS|j|d1< ||||d�V  |jd2k�rbq|jd3k�r“|j
}z	| ¡ }|› }W n t�yˆ   |j
�p…d4|j› �}Y nw | jj||| jd5�‚q6 W | ¡ I d H  d S | ¡ I d H  w r\   )r   r   rD   r6   rE   rl   rU   rm   rn   r   r^   ro   rp   rq   Úaclosers   r)   r)   r*   rG   ÿ   sš   €


ÿý÷Ã"MzAsyncStream.__stream__r   c                 Ã  s   �| S rB   r)   rM   r)   r)   r*   Ú
__aenter__U  s   €zAsyncStream.__aenter__ry   rz   r{   r|   r}   r~   c                 Ã  s   �|   ¡ I d H  d S rB   r   r€   r)   r)   r*   Ú	__aexit__X  s   €zAsyncStream.__aexit__c                 Ã  s   �| j  ¡ I dH  dS r‚   )r6   r“   rM   r)   r)   r*   rr   `  s   €zAsyncStream.close)
r=   r>   r6   r5   r?   r   r<   r7   r   r@   rƒ   )r   r�   )r   r‘   )r6   r5   r   r‘   r„   r…   r†   )r/   r0   r1   r‡   rˆ   r8   rJ   rŽ   r�   rU   r‰   rZ   rG   r”   r•   rr   r)   r)   r)   r*   r�   Õ   s    
 ú




V
r�   c                   @  s‚   e Zd Zd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ed$dd„ƒZ	d%dd„Z
ed#dd„ƒZdS )&ÚServerSentEventN©rm   r^   ÚidÚretryÚrawrm   ú
str | Noner^   r˜   r™   ú
int | Nonerš   úlist[str] | Noner   r@   c                C  sD   |d u rd}|| _ || _|pd | _|| _|d ur|| _d S g | _d S )NÚ )Ú_idÚ_dataÚ_eventÚ_retryÚ_raw)r(   rm   r^   r˜   r™   rš   r)   r)   r*   rJ   j  s   	
zServerSentEvent.__init__c                 C  ó   | j S rB   )r¡   rM   r)   r)   r*   rm   |  ó   zServerSentEvent.eventc                 C  r¤   rB   )rŸ   rM   r)   r)   r*   r˜   €  r¥   zServerSentEvent.idc                 C  r¤   rB   )r¢   rM   r)   r)   r*   r™   „  r¥   zServerSentEvent.retryÚstrc                 C  r¤   rB   )r    rM   r)   r)   r*   r^   ˆ  r¥   zServerSentEvent.dataú	list[str]c                 C  r¤   )zÿThe original wire lines this event was decoded from, without trailing newlines.

        Includes SSE fields the decoder does not otherwise model (comment lines,
        unknown fields). Empty for events that were constructed rather than decoded.
        )r£   rM   r)   r)   r*   rš   Œ  s   zServerSentEvent.rawr   c                 C  s   t  | j¡S rB   )rn   Úloadsr^   rM   r)   r)   r*   rn   •  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=ú))rm   r^   r˜   r™   rM   r)   r)   r*   Ú__repr__˜  s   &zServerSentEvent.__repr__)rm   r›   r^   r›   r˜   r›   r™   rœ   rš   r�   r   r@   )r   r›   )r   rœ   )r   r¦   )r   r§   )r   r   )r/   r0   r1   rJ   Úpropertyrm   r˜   r™   r^   rš   rn   r   rª   r)   r)   r)   r*   r–   i  s(    ù
r–   c                   @  sr   e Zd ZU ded< 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 )&rX   r§   r    r›   r¡   rœ   r¢   Ú_last_event_idr£   r   r@   c                 C  s"   d | _ g | _d | _d | _g | _d S rB   )r¡   r    r¬   r¢   r£   rM   r)   r)   r*   rJ   ¤  s
   
zSSEDecoder.__init__ru   úIterator[bytes]rS   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©r(   ru   ÚchunkÚraw_lineÚlinerv   r)   r)   r*   rT   «  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©r(   ru   r^   r´   r¶   r)   r)   r*   r°   µ  s   €
€ü
ÿzSSEDecoder._iter_chunksúAsyncIterator[bytes]r‘   c                 C sL   �|   |¡2 z3 dH W }| ¡ D ]}| d¡}|  |¡}|r!|V  qq6 dS )r®   Nr¯   )Ú_aiter_chunksr±   r²   r³   r)   r)   r*   r’   Á  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¾   r)   r)   r*   rÀ   Ë  s   €
€üÿ
ÿzSSEDecoder._aiter_chunksr¶   r¦   úServerSentEvent | Nonec              	   C  s*  |s4| j s| js| js| jd u rg | _d S t| j d | j¡| j| j| jd�}d | _ g | _d | _g | _|S | j |¡ | d¡rAd S | 	d¡\}}}| d¡rT|dd … }|dkr]|| _ d S |dkri| j |¡ d S |dkryd	|v rt	 d S || _d S |d
kr’zt
|ƒ| _W d S  ttfy‘   Y d S w 	 d S )NÚ
r—   ú:ú r   rm   r^   r˜   ú r™   )r¡   r    r¬   r¢   r£   r–   ÚjoinÚappendÚ
startswithÚ	partitionÚintÚ	TypeErrorÚ
ValueError)r(   r¶   rv   Ú	fieldnameÚ_Úvaluer)   r)   r*   r²   ×  sX   
û	

ñó÷	øûûzSSEDecoder.decodeNr†   ©ru   r­   r   rS   )ru   r­   r   r­   ©ru   r¿   r   r‘   )ru   r¿   r   r¿   )r¶   r¦   r   rÁ   )
r/   r0   r1   rˆ   rJ   rT   r°   r’   rÀ   r²   r)   r)   r)   r*   rX   �  s   
 






rX   c                   @  s    e Zd Zddd„Zdd	d
„ZdS )r9   ru   r­   r   rS   c                 C  ó   dS )r®   Nr)   ©r(   ru   r)   r)   r*   rT     ó   zSSEBytesDecoder.iter_bytesr¿   r‘   c                 C  rÒ   )zdGiven an async iterator that yields raw binary data, iterate over it & yield every event encounteredNr)   rÓ   r)   r)   r*   r’     rÔ   zSSEBytesDecoder.aiter_bytesNrÐ   rÑ   )r/   r0   r1   rT   r’   r)   r)   r)   r*   r9     s    
r9   Útyprh   r   ú;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Ú
issubclassr4   r�   )rÕ   Úoriginr)   r)   r*   Úis_stream_class_type  s   rÛ   )Úfailure_messageÚ
stream_clsrÜ   r›   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   )r4   r�   r   ztuple[type, ...])ÚindexÚgeneric_basesrÜ   )Ú_base_clientr4   r�   r   r   )rÝ   rÜ   r4   r�   r)   r)   r*   Úextract_stream_chunk_type  s   ürá   )rÕ   rh   r   rÖ   )rÝ   rh   rÜ   r›   r   rh   ),Ú
__future__r   Úabcrn   r×   r%   Útypesr   Útypingr   r   r   r   r   r	   r
   r   Útyping_extensionsr   r   r   r   r   r   ÚhttpxÚ_utilsr   r   rE   r   r   Ú_modelsr   r   ÚABCMetar   r4   r‹   r�   r–   rX   r9   rÛ   rá   r)   r)   r)   r*   Ú<module>   s6   (   4p

	ý