o
    î6Wj ã                   @  s:  U d dl mZ d dlZd dlZd dlZd dlmZmZm	Z	m
Z
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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 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dl,m-Z- ddl.m/Z/ ddl0m1Z1 ddl2m3Z3 ddl4m5Z5 ddl6m7Z7 ddgZ8e 9d¡Z:de;d< dZ<dZ=de;d< 	 G dd„ dƒZ>eddd �Z?d!e;d"< ed#d$d �Z@d%e;d&< G d'd„ de*ƒZAG d(d)„ d)eƒZBG d*d+„ d+eƒZCG d,d-„ d-eƒZDG d.d/„ d/eƒZEG d0d1„ d1ƒZFd¦d7d8„ZGd§d:d;„ZHG d<d=„ d=ƒZId¨dBdC„ZJd©dGdH„ZKG dIdJ„ dJƒZLdªdQdR„ZMd«dUdV„ZNd¬dXdY„ZOd­d]d^„ZPd®dbdc„ZQd®ddde„ZRd¯didj„ZSd°dndo„ZTd±dqdr„ZUd²dtdu„ZVd³dydz„ZWd´d}d~„ZXdµd€d�„ZYd¶d„d…„ZZd·d†d‡„Z[d·dˆd‰„Z\d¸d�dŽ„Z]d¹d‘d’„Z^dºd”d•„Z_d»d™dš„Z`d¼dœd�„Zad½d d¡„ZbG d¢d£„ d£ejcƒZdG d¤d¥„ d¥ejeƒZfdS )¾é    )ÚannotationsN)ÚAnyÚDictÚListÚCallableÚIterableÚIteratorÚOptionalÚ	GeneratorÚAsyncIteratorÚAsyncGeneratorÚcast)ÚTokenÚ
ContextVar)ÚLiteralÚoverrideé   )Úis_dict)Ú	BaseModel)Ú
APIRequest)ÚAPIResponseÚAsyncAPIResponse)ÚStreamÚAsyncStreamÚServerSentEvent)ÚAnthropicError)ÚCallNextÚ
MiddlewareÚAsyncCallNext)Úmerge_headers)ÚMessageé   )Úhelper_header)ÚBetaMessage)ÚAnthropicBetaParam)ÚBetaFallbackParamÚBetaFallbackStateÚBetaRefusalFallbackMiddlewarezanthropic.lib.middlewarezlogging.LoggerÚlogz/v1/messages)zfallback-credit-2026-06-01útuple[AnthropicBetaParam, ...]ÚDEFAULT_BETASc                   @  s:   e Zd ZU dZded< 	 ddd„Zddd	„Zddd„ZdS )r&   u{  Tracks which fallback a sequence of requests is pinned to.

    Create one and enter it (`with state:` â€” the same context manager works for
    both clients) around every request that should share the pin â€” the turns of
    one conversation, or any wider scope the stickiness should apply to;
    `BetaRefusalFallbackMiddleware` mutates it in place when a model refuses.
    z
int | NoneÚindexÚreturnÚNonec                 C  s
   d | _ d S ©N©r+   ©Úself© r2   úr/home/esfera/Documents/content_generation/venv/lib/python3.10/site-packages/anthropic/lib/middleware/_fallbacks.pyÚ__init__D   ó   
zBetaFallbackState.__init__c                 C  s&   t  | ¡}t g t ¡ ¢|‘R ¡ | S r.   )Ú_fallback_stateÚsetÚ_fallback_state_tokensÚget)r1   Útokenr2   r2   r3   Ú	__enter__G   s   
zBetaFallbackState.__enter__Úexc_infoÚobjectc                 G  s,   t  ¡ }t  |d d… ¡ t |d ¡ d S )Néÿÿÿÿ)r8   r9   r7   r6   Úreset)r1   r<   Útokensr2   r2   r3   Ú__exit__L   s   zBetaFallbackState.__exit__N©r,   r-   )r,   r&   )r<   r=   r,   r-   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú__annotations__r4   r;   rA   r2   r2   r2   r3   r&   4   s   
 

Úanthropic_beta_fallback_state)Údefaultz$ContextVar[BetaFallbackState | None]r6   Ú$anthropic_beta_fallback_state_tokensr2   z7ContextVar[tuple[Token[BetaFallbackState | None], ...]]r8   c                   @  s‚   e Zd ZdZddœd1d
d„Zed2dd„ƒZed3dd„ƒZd4dd„Zd5dd„Z	d6d d!„Z
d7d'd(„Zd8d)d*„Zd9d,d-„Zd:d/d0„ZdS );r'   uÈ  Middleware that retries refused beta `/v1/messages` requests down a
    fallback chain, reproducing the server-side `fallbacks` wire shape
    client-side.

    Only `client.beta.messages` requests are handled â€” refusals minted by the
    first-party `client.messages` surface carry no `fallback_credit_token`, so
    those requests pass through untouched.

    Non-streaming: when a response comes back with `stop_reason: "refusal"`, the
    request is retried with each entry of `fallbacks` merged over the original
    params â€” passing along the refusal's `fallback_credit_token` when it minted
    one â€” until a model accepts or the chain is exhausted. A `fallback` seam
    block per model boundary is prepended to the served message's content â€”
    the same block shape the streaming splice emits. The served hop's `usage`
    is left verbatim (streaming rewrites it to per-hop `usage.iterations`).

    Streaming: when the stream ends in `stop_reason: "refusal"`, a second
    request is issued to the fallback model. It carries the refusal's
    `fallback_credit_token`, plus the refused model's partial output as a
    trailing assistant prefill when the refusal grants one
    (`fallback_has_prefill_claim`). The fallback's events are then spliced onto
    the still-open stream, so the client sees one continuous message in the
    server-side `fallbacks` wire shape: a `fallback` content block at each
    model boundary, monotonic block indices, and per-hop `usage.iterations` on
    the final `message_delta`. A refusal before any output streamed retries
    even without a credit token, and the serving hop's `message_start` opens
    the wire carrying the primary's message id.

    The fallback-credit beta the credit tokens require is sent by default on
    every request the middleware handles; the `betas` option controls this.

    In both modes a fallback that itself refuses with a fresh credit token
    continues down the chain. A streaming fallback whose appended prefill the
    server rejects (HTTP 400 body mismatch) is retried once without it; a
    fallback whose request fails outright is skipped â€” its token was never
    redeemed, so it carries to the next entry. When every remaining entry fails
    over HTTP, the suppressed refusal is replayed to the client with
    `recommended_model` stamped from the final failure (the failed model for
    capacity errors, `null` otherwise). A refusal surfaced to the client rather
    than retried is reported through the `anthropic.lib.middleware` logger.

    To keep later requests on the model that accepted, run them inside a shared
    `BetaFallbackState` context; requests sharing that state start directly at
    the pinned fallback. Reuse one state across whatever scope the pin should
    apply to â€” typically a conversation. The state is the only pin: `fallback`
    seam blocks replayed in the request history are stripped from the outgoing
    request (an assistant turn left empty by the strip is dropped whole), never
    read back as a pin.

    ```py
    client = Anthropic(middleware=[BetaRefusalFallbackMiddleware([{"model": "claude-opus-4-8"}])])

    state = BetaFallbackState()
    with state:
        message = client.beta.messages.create(**params)
    ```
    N)ÚbetasÚ	fallbacksúIterable[BetaFallbackParam]rK   ú#Iterable[AnthropicBetaParam] | Noner,   r-   c                C  s*   t |ƒ| _|du rtnt |ƒ| _d| _dS )ué  
        Args:
            fallbacks: The fallback chain, tried in order. An empty chain disables
                the middleware.

            betas: Betas added to the `anthropic-beta` header of every `/v1/messages`
                request this middleware handles â€” the original request included, since
                refusals only carry a `fallback_credit_token` when the beta is enabled.
                Defaults to `("fallback-credit-2026-06-01",)`; pass `()` to send none.
        NF)ÚtupleÚ
_fallbacksr*   Ú_betasÚ_warned_missing_state)r1   rL   rK   r2   r2   r3   r4   ™   s   

z&BetaRefusalFallbackMiddleware.__init__Úrequestr   Ú	call_nextr   úAPIResponse[Any]c                 C  sÚ  |   |¡}|d u r||ƒS t ¡ }|  |¡}|  |¡}t|| jƒ}t|ƒ}|dkr+|ni |¥| j| ¥}|j	|d�}||ƒ}	|	j
jsD|	S |jr_|d }
|
t| jƒkrT|	S | j|||	||
|d�S |}|	}t| d¡pjdƒ}g }|t| jƒd k rÐ|j
jrÐ| ¡ }t|ttfƒr‹|jdkrŒnD|d7 }||ƒ t|ƒ}t|ƒ}||j	t|| j| |ƒd�ƒ}|j
jrÃt| j| d ƒ}| t|||ƒ¡ |}|t| jƒd k rÐ|j
js{|rë|j
jrë| ¡ }t|ttfƒrë|jdkrët||ƒS |S ©Nr>   ©Úbodyé   ©rS   rX   ÚresponserT   Ú	first_hopÚpinÚmodelÚ Úrefusal)Ú_applicable_bodyr6   r9   Ú_start_indexÚ	_make_pinÚ_with_middleware_headersrQ   Ú_strip_seam_blocksrP   ÚcopyÚhttp_responseÚ
is_successÚstreamÚlenÚ_splice_fallback_streamÚstrÚparseÚ
isinstancer    r#   Ústop_reasonÚ_credit_tokenÚ_refusal_categoryÚ_merged_bodyÚappendÚ_seam_blockÚ_prepend_seam_blocks©r1   rS   rT   rX   ÚstateÚstart_indexr]   Úinitial_bodyÚinitial_requestr[   r\   r+   ÚresÚ
from_modelÚseamsÚmessager:   ÚcategoryÚto_modelÚservedr2   r2   r3   Úhandle­   s`   


ú	ó
z$BetaRefusalFallbackMiddleware.handler   úAsyncAPIResponse[Any]c                 Ã  sú  �|   |¡}|d u r||ƒI d H S t ¡ }|  |¡}|  |¡}t|| jƒ}t|ƒ}|dkr/|ni |¥| j| ¥}|j	|d�}||ƒI d H }	|	j
jsK|	S |jrf|d }
|
t| jƒkr[|	S | j|||	||
|d�S |}|	}t| d¡pqdƒ}g }|t| jƒd k rÝ|j
jrÝ| ¡ I d H }t|ttfƒr•|jdkr–nG|d7 }||ƒ t|ƒ}t|ƒ}||j	t|| j| |ƒd�ƒI d H }|j
jrÐt| j| d ƒ}| t|||ƒ¡ |}|t| jƒd k rÝ|j
js‚|rû|j
jrû| ¡ I d H }t|ttfƒrû|jdkrût||ƒS |S rV   )ra   r6   r9   rb   rc   rd   rQ   re   rP   rf   rg   rh   ri   rj   Ú_splice_fallback_stream_asyncrl   rm   rn   r    r#   ro   rp   rq   rr   rs   rt   Ú_prepend_seam_blocks_asyncrv   r2   r2   r3   Úhandle_asyncñ   sb   €


ú	$ó
z*BetaRefusalFallbackMiddleware.handle_asyncúdict[str, Any] | Nonec                 C  sj   t |jƒ}t |j¡}| jr&|j ¡ dks&|jt	ks&|j
 d¡dks&|du r(dS | d¡dur3tdƒ‚|S )zMThe request's JSON body when this middleware applies to it, `None` otherwise.ÚpostÚbetaÚtrueNrL   az  Sending the `fallbacks:` request param is not supported when using the `BetaRefusalFallbackMiddleware`. You should either remove the middleware and send `fallbacks:` with the `server-side-fallback-2026-06-01` beta header to let the API handle refusal fallbacks, or omit the `fallbacks:` param if you'd like `BetaRefusalFallbackMiddleware` to handle fallbacks on the client side.)Ú_as_dictÚjsonÚhttpxÚURLÚurlrP   ÚmethodÚlowerÚpathÚ_MESSAGES_PATHÚparamsr9   r   )r1   rS   rX   r�   r2   r2   r3   ra   5  s   
þ
ÿz.BetaRefusalFallbackMiddleware._applicable_bodyrw   úBetaFallbackState | NoneÚintc                 C  sX   |du s	|j du rdS |j }d|  krt| jƒk s*n td|› dt| jƒ› d�ƒ‚|S )zÍThe chain entry this request starts at (-1 = the original params).

        Only an explicit `BetaFallbackState` pin moves the start; without one
        the request starts at the original params.
        Nr>   zBetaFallbackState.index z! is out of bounds for a chain of z? fallback(s); was the state shared with a different middleware?)r+   rj   rP   r   )r1   rw   rx   r2   r2   r3   rb   O  s   ÿÿz*BetaRefusalFallbackMiddleware._start_indexúCallable[[int], None]c                   s   d‡ ‡fdd„}|S )	zUPin requests sharing the state to the entry being tried (or warn that there is none).r+   r–   r,   r-   c                   s0   ˆd ur	| ˆ_ d S ˆ jsdˆ _t d¡ d S d S )NTzôanthropic-sdk: BetaRefusalFallbackMiddleware fell back without an active BetaFallbackState; follow-up requests will retry models that already refused. Run them inside a shared `with BetaFallbackState():` block to pin them to the accepted model.)r+   rR   r(   Úwarningr/   ©r1   rw   r2   r3   r]   b  s   
ÿþz4BetaRefusalFallbackMiddleware._make_pin.<locals>.pinN)r+   r–   r,   r-   r2   )r1   rw   r]   r2   r™   r3   rc   _  s   z'BetaRefusalFallbackMiddleware._make_pinrX   údict[str, Any]r[   r\   r]   c             	   C  óB   | j ||||||d�}tt|jt|ƒƒ|j|jd|j|j|j	d�S )aH  Wrap the refusable stream in a response whose body passes events through
        until a retryable refusal, then splices the fallback chain's events on.

        Closing the returned response (or the `Stream` parsed from it) tears down
        whichever stream is being read and abandons any in-flight fallback request.
        rZ   T©ÚrawÚcast_toÚclientri   Ú
stream_clsÚoptionsÚretries_taken)
Ú_spliced_framesr   Ú_spliced_http_responserg   Ú_FrameByteStreamÚ_cast_toÚ_clientÚ_stream_clsÚ_optionsr¢   ©r1   rS   rX   r[   rT   r\   r]   Úframesr2   r2   r3   rk   p  s"   úùz5BetaRefusalFallbackMiddleware._splice_fallback_streamc             	   C  r›   )NrZ   Trœ   )
Ú_spliced_frames_asyncr   r¤   rg   Ú_AsyncFrameByteStreamr¦   r§   r¨   r©   r¢   rª   r2   r2   r3   r„   ’  s"   
úùz;BetaRefusalFallbackMiddleware._splice_fallback_stream_asyncúGenerator[bytes, None, None]c                c  sÒ  � | j }|j}�zU|j}	tddd dd g d�}
t|	|
ƒE d H }|jd u r0W |d ur.| ¡  d S d S |	 ¡  d }t ||¡}t|t	|ƒƒD �]}|| }t
|d ƒ}|d t	|ƒk }||ƒ | ¡ }d }d }tdƒD ]q}|j| ||¡d�}z||ƒ}W n  tyš } zt d	||¡ t|d d
�}W Y d }~ nDd }~ww |jjr¤|j} n6t|jƒ}|j ¡  |dkrÇ|jdkrÇ|rÇt dt|ƒ¡ |j}d }qht d||jt|ƒ¡ t||jd
�} |d urú|ráqC| |¡D ]}|V  qæ W |d urø| ¡  d S d S |d u�sJ ‚|j}| |¡ t|j|t|j|d�|j|j| ¡ d�}
t||
ƒE d H }|j �r-| !¡  |jd u �rB W |d u�r@| ¡  d S d S | ¡  d }| "|||¡ qCW |d u�r\| ¡  d S d S |d u�rh| ¡  w w ©Nr   TF)Ú
index_baseÚhas_nextÚspliceÚ	wire_openÚ
primary_idÚseam_framesr^   rY   r!   rW   zOanthropic-sdk: BetaRefusalFallbackMiddleware: fallback request to %s failed: %s)r^   Ústatusi�  z�anthropic-sdk: BetaRefusalFallbackMiddleware: fallback request with the partial output appended was rejected (HTTP 400: %s); retrying without itzXanthropic-sdk: BetaRefusalFallbackMiddleware: fallback request to %s failed: HTTP %s: %s)Ú
iterationsr^   )#rP   rg   Ú
_HopReaderÚ
_drive_hopÚrefusedÚcloseÚ_ChainStateÚbeginÚrangerj   rl   Úcontinuationrf   Úhop_bodyÚ	Exceptionr(   ÚerrorÚ_HopFailurerh   Ú
_read_jsonÚstatus_coder˜   Ú_json_dumpsÚbaseÚterminal_failure_framesÚ
queue_seamÚ
next_indexÚ_SpliceInfor·   r³   r´   Úpending_seam_framesÚopenedÚmark_openedÚabsorb_refusal)r1   rS   rX   r[   rT   r\   r]   rL   ÚcurrentÚstream_aÚreaderÚoutcomeÚchainÚhopÚentryr^   r±   r¿   Úres_bÚfailureÚattemptÚhop_requestÚerrÚerr_bodyÚframeÚhop_responser2   r2   r3   r£   Ê  sÌ   €
ù	
\ÿ¥ý€ù

ýüÿè
ú
	ÿø­
Uÿ
ÿz-BetaRefusalFallbackMiddleware._spliced_framesúAsyncGenerator[bytes, None]c                C sF  �| j }|j}�zŒ|j}	tddd dd g d�}
t|	|
ƒ2 z	3 d H W }|V  q6 |
 ¡ }|jd u r?W |d ur=| ¡ I d H  d S d S |	 ¡ I d H  d }t ||¡}t	|t
|ƒƒD �]-}|| }t|d ƒ}|d t
|ƒk }||ƒ | ¡ }d }d }t	dƒD ]z}|j| ||¡d�}z	||ƒI d H }W n  ty¯ } zt d	||¡ t|d d
�}W Y d }~ nJd }~ww |jjr¹|j} n<t|jƒI d H }|j ¡ I d H  |dkrâ|jdkrâ|rât dt|ƒ¡ |j}d }qzt d||jt|ƒ¡ t||jd
�} |d u�r|rýqU| |¡D ]}|V  �q W |d u�r| ¡ I d H  d S d S |d u�s"J ‚|j}| |¡ t|j|t|j|d�|j|j|  ¡ d�}
t||
ƒ2 z
3 d H W }|V  �qC6 |
 ¡ }|j!�r[| "¡  |jd u �rs W |d u�rq| ¡ I d H  d S d S | ¡ I d H  d }| #|||¡ qUW |d u�r“| ¡ I d H  d S d S |d u�r¢| ¡ I d H  w w r¯   )$rP   rg   r¸   Ú_drive_hop_asyncÚfinishrº   Úacloser¼   r½   r¾   rj   rl   r¿   rf   rÀ   rÁ   r(   rÂ   rÃ   rh   Ú_read_json_asyncrÅ   r˜   rÆ   rÇ   rÈ   rÉ   rÊ   rË   r·   r³   r´   rÌ   rÍ   rÎ   rÏ   )r1   rS   rX   r[   rT   r\   r]   rL   rÐ   rÑ   rÒ   rÝ   rÓ   rÔ   rÕ   rÖ   r^   r±   r¿   r×   rØ   rÙ   rÚ   rÛ   rÜ   rÞ   r2   r2   r3   r¬   E  sØ   €
ù	ÿ
Qÿ°ý€ùýü


ÿé
ú
ÿ
ÿû·
Kÿÿz3BetaRefusalFallbackMiddleware._spliced_frames_async)rL   rM   rK   rN   r,   r-   )rS   r   rT   r   r,   rU   )rS   r   rT   r   r,   rƒ   )rS   r   r,   r‡   )rw   r•   r,   r–   )rw   r•   r,   r—   )rS   r   rX   rš   r[   rU   rT   r   r\   r–   r]   r—   r,   rU   )rS   r   rX   rš   r[   rƒ   rT   r   r\   r–   r]   r—   r,   rƒ   )rS   r   rX   rš   r[   rU   rT   r   r\   r–   r]   r—   r,   r®   )rS   r   rX   rš   r[   rƒ   rT   r   r\   r–   r]   r—   r,   rß   )rC   rD   rE   rF   r4   r   r‚   r†   ra   rb   rc   rk   r„   r£   r¬   r2   r2   r2   r3   r'   ^   s    >üC
C



"
8{c                   @  s:   e Zd ZU ded< 	 ded< 	 ded< ded< ded< d	S )
Ú_RefusalúOptional[str]r:   r   ÚboolÚhas_prefill_claimúDict[str, Any]ÚusageÚeventN)rC   rD   rE   rG   r2   r2   r2   r3   rä   »  s   
 rä   c                   @  ó"   e Zd ZU dZded< ded< dS )rË   z6Splice context for fallback hops; `None` for stream A.úList[Dict[str, Any]]r·   rl   r^   N©rC   rD   rE   rF   rG   r2   r2   r2   r3   rË   É  s   
 rË   c                   @  rë   )rÃ   z[A fallback hop whose request failed; stamps `recommended_model` when it
    ends the chain.rl   r^   zOptional[int]r¶   Nrí   r2   r2   r2   r3   rÃ   Ð  s
   
 rÃ   c                   @  sV   e Zd ZU dZded< 	 ded< 	 ded< 	 ded	< 	 d
ed< 	 ded< 	 ded< dS )Ú_HopOutcomez*The outcome of consuming one hop's stream.zOptional[_Refusal]rº   rå   r^   zOptional[Dict[str, Any]]Ústart_eventrì   Úblocksr–   rÊ   ræ   rÍ   Úlast_fallback_toNrí   r2   r2   r2   r3   rî   Ù  s    
 rî   c                   @  s\   e Zd ZU dZded< 	 d!dd„Zed"dd„ƒZd#dd„Zd$dd„Z	d$dd„Z
d%dd„Zd S )&r¸   uÆ  Consumes one hop's SSE events, producing the frames to forward to the
    client while accumulating its content blocks.

    Events are held until the hop produces output (a content block or a
    terminal event), so a refusal that arrives before any output â€” a pre-stream
    decline â€” leaves no trace on the wire and chains silently. On open, stream
    A's held `message_start` is flushed in its original wire bytes; a spliced
    hop's is suppressed when the wire is already open, or emitted with its `id`
    rewritten to the primary's when this hop's start opens the wire. Queued
    seam frames for every hop reached so far are flushed right after.

    A spliced hop has its block indices shifted by `index_base` and its
    terminal message_delta's usage rewritten to the `usage.iterations` chain
    shape.

    A refusal that can be chained â€” an entry remains, and either a
    `fallback_credit_token` was minted or nothing has streamed yet â€” ends the
    hop early: open blocks are closed, the terminal message_delta +
    message_stop are suppressed, and the refusal is set on the outcome so the
    caller can issue the next hop. Any other refusal is logged and passes
    through to the client.
    z_HopOutcome | NonerÓ   r°   r–   r±   ræ   r²   ú_SpliceInfo | Noner³   r´   ú
str | Nonerµ   úlist[bytes]r,   r-   c                C  s\   t |ƒ| _|| _|| _|| _|| _|| _d | _d | _d | _	d | _
g | _d| _d | _d | _d S )NF)Ú_BlockTrackerÚ_trackerÚ	_has_nextÚ_spliceÚ
_wire_openÚ_primary_idÚ_seam_framesÚ_modelÚ_start_usageÚ_start_eventÚ_held_startÚ_heldÚ_openedÚ_last_fallback_torÓ   )r1   r°   r±   r²   r³   r´   rµ   r2   r2   r3   r4     s   


z_HopReader.__init__c                 C  ó   | j S r.   )r  r0   r2   r2   r3   rÍ   '  ó   z_HopReader.openedÚsser   c              	   C  s
  t t|jƒƒ}|dur| d¡nd}| j}|dkrG|durGt | d¡ƒ}|dur?| d¡}t|tƒr4|nd| _t | d¡ƒ| _|| _	|| _
g S |dkrš|durš|  ¡ }t | d¡ƒ}|dur| d¡d	krt | d
¡ƒ}	|	duru|	 d¡nd}
t|
tƒr|
| _| j |¡ |dur’g |¢td|ƒ‘S g |¢t|ƒ‘S |dkrÁ|durÁ|  ¡ }| j |¡ |dur¹g |¢td|ƒ‘S g |¢t|ƒ‘S |dkrè|durè|  ¡ }| j |¡ |duràg |¢td|ƒ‘S g |¢t|ƒ‘S |dk�rô|du�rôt | d¡ƒpúi }| d¡dk�r¥t | d¡ƒ}|du�r| d¡dk�r|nd}|du�r%| d¡nd}t|tƒ�r2|�r2|nd}| j �o>| j ¡  }| j�r—|du�sK|�r—tt | d¡ƒ| jƒ}t| j ¡ ƒ}d| _
| j ¡  tt||du�rr| d¡nd|du�o~| d¡du ||d�| j| j	| j ¡ | jj| j| jd�| _|S |�s t  d¡ nt  d¡ |  ¡ }|du�rì| d¡dk�rÈt | d¡ƒ}|du�rÈ| !dd¡ t | d¡ƒ�pÑi }g |j"¢t#||j$ƒ‘|d< ||d< g |¢td|ƒ‘S g |¢t|ƒ‘S | j�s | j %|¡ g S t|ƒgS )z:Process one event, returning the frames to forward for it.NÚtypeÚmessage_startr~   r^   ré   Úcontent_block_startÚcontent_blockÚfallbackÚtoÚcontent_block_deltaÚcontent_block_stopÚmessage_deltaÚdeltaro   r`   Ústop_detailsÚfallback_credit_tokenr   Úfallback_has_prefill_claimT)r:   r   rç   ré   rê   ©rº   r^   rï   rð   rÊ   rÍ   rñ   zvanthropic-sdk: BetaRefusalFallbackMiddleware: refusal stop_details has no fallback_credit_token; surfacing the refusalzkanthropic-sdk: BetaRefusalFallbackMiddleware: refusal but no fallback entries remain; surfacing the refusalÚrecommended_modelr·   )&r‹   Ú
_safe_jsonÚdatar9   rø   rn   rl   rü   rý   rþ   rÿ   Ú_openr  rö   ÚstartÚ_emitÚ_passthrough_sser  Ústopr  Úcontent_blocksr÷   Ú	_backfillÚlistÚclose_open_blocksr   Úclearrî   rä   rÊ   rÓ   r(   rÂ   Ú
setdefaultr·   Ú_serving_iteration_entryr^   rs   )r1   r  rê   Ú
event_typer²   r~   r^   r«   r	  r  r€   r  r  Údetailsr:   Ú
pre_streamré   r2   r2   r3   Úfeed+  sÂ   

ÿÿýÿÿýÿÿý"
ûóÿÿ

	ÿ
þ
z_HopReader.feedc                 C  sÆ   | j rg S d| _ g }| jdu r| jdur| t| jƒ¡ n+| jsH| jdurHt | j¡}t	| 
d¡ƒ}|dur@| jdur@| j|d< | td|ƒ¡ | | j¡ | dd„ | jD ƒ¡ d| _| j ¡  |S )z°Open the wire for this hop: its message_start (raw for stream A;
        suppressed or id-rewritten for a spliced hop), the queued seam frames,
        then anything else held.TNr~   Úidr  c                 s  ó   � | ]}t |ƒV  qd S r.   )r  )Ú.0Úheldr2   r2   r3   Ú	<genexpr>·  ó   € z#_HopReader._open.<locals>.<genexpr>)r  rø   rÿ   rs   r  rù   rþ   Ú_copyÚdeepcopyr‹   r9   rú   r  Úextendrû   r   r   )r1   r«   rï   r~   r2   r2   r3   r  ¥  s&   

€

z_HopReader._openc                 C  s   | j dus| jr
g S |  ¡ S )zCFrames still owed when the stream ends without a chainable refusal.N)rÓ   r  r  r0   r2   r2   r3   Úfinish_frames¼  s   z_HopReader.finish_framesrî   c              	   C  s8   | j dur| j S td| j| j| j ¡ | jj| j| jd�S )zCThe outcome once the hop's stream is fully consumed (or cut early).Nr  )	rÓ   rî   rü   rþ   rö   r  rÊ   r  r  r0   r2   r2   r3   rá   Â  s   
ùz_HopReader.finishN)r°   r–   r±   ræ   r²   rò   r³   ræ   r´   ró   rµ   rô   r,   r-   )r,   ræ   )r  r   r,   rô   ©r,   rô   )r,   rî   )rC   rD   rE   rF   rG   r4   ÚpropertyrÍ   r&  r  r0  rá   r2   r2   r2   r3   r¸   ó  s   
 


z
r¸   r[   úhttpx.ResponserÒ   r,   ú#Generator[bytes, None, _HopOutcome]c                 c  sR   � t  | ¡D ]}| |¡D ]}|V  q|jdur nq| ¡ D ]}|V  q| ¡ S )zKFeed one hop's SSE events through `reader`, yielding the frames to forward.N)r   Ú
raw_eventsr&  rÓ   r0  rá   ©r[   rÒ   r  rÝ   r2   r2   r3   r¹   Ñ  s   €
ÿr¹   úAsyncIterator[bytes]c                 C sX   �t  | ¡2 z3 dH W }| |¡D ]}|V  q|jdur nq6 | ¡ D ]}|V  q$dS )u»   Feed one hop's SSE events through `reader`, yielding the frames to forward.

    The outcome is read from `reader.finish()` afterwards â€” async generators
    cannot return a value.
    N)r   r5  r&  rÓ   r0  r6  r2   r2   r3   rà   Ý  s   €
ÿýÿrà   c                   @  sÒ   e Zd ZU dZded< 	 ded< 	 ded< 	 ded	< 	 ded
< ded< ded< ded< ded< ded< ded< 	 d7dd„Zed8dd„ƒZd9d d!„Zd:d#d$„Z	d;d%d&„Z
d<d(d)„Zd=d-d.„Zd>d0d1„Zd?d4d5„Zd6S )@r¼   z9Bookkeeping shared across the hops of one spliced stream.r–   rÊ   ræ   r³   ró   r´   rl   Úprimary_modelr:   ú	List[Any]rÇ   Úpartialr|   rä   Úlast_refusalzDict[str, Any] | NoneÚlast_start_eventrì   r·   rX   rš   r,   r-   c                 C  s   || _ g | _d S r.   )Ú_bodyÚ_pending_seams)r1   rX   r2   r2   r3   r4     s   
z_ChainState.__init__rÓ   rî   c                 C  sÒ   |j dusJ ‚| |ƒ}t| d¡pdƒ|_|j|_|j|_t|jp"i  d¡ƒ}|dur0| d¡nd}t	|tƒr9|nd|_
|j j|_g |_|j jrMt|jƒng |_|jpU|j|_|j |_|j|_t|j |jƒ|_|S )z3The chain state after stream A's chainable refusal.Nr^   r_   r~   r'  )rº   rl   r9   r8  rÊ   rÍ   r³   r‹   rï   rn   r´   r:   rÇ   rç   Ú_to_prefill_blocksrð   r:  rñ   r|   r;  r<  Ú_declined_iteration_entriesr·   )ÚclsrX   rÓ   rÔ   Ústart_messageÚstart_idr2   r2   r3   r½     s    
z_ChainState.beginr€   c              
   C  sV   | j }|  j d7  _ | j tdd|t| j|| jjƒdœƒtdd|dœƒg¡ || _dS )a  Queue the reached hop's `fallback` seam block at the next monotonic index.

        Queued seams are flushed by the hop reader once output reaches the wire,
        so a chain of pre-stream declines lines its seams up after the serving
        hop's message_start.
        rY   r  )r  r+   r	  r  ©r  r+   N)rÊ   r>  r/  r  rt   r|   r;  r   )r1   r€   Ú
seam_indexr2   r2   r3   rÉ   #  s   ýþ÷ÿ
z_ChainState.queue_seamrô   c                 C  s
   t | jƒS r.   )r  r>  r0   r2   r2   r3   rÌ   ;  r5   z_ChainState.pending_seam_framesc                 C  s   d| _ | j ¡  d S )NT)r³   r>  r   r0   r2   r2   r3   rÎ   >  s   z_ChainState.mark_openedú	list[Any]c                 C  s   g | j ¢| j¢S r.   )rÇ   r:  r0   r2   r2   r3   r¿   B  s   z_ChainState.continuationrÖ   r%   r¿   c                 C  sf   dd„ i | j ¥|¥ ¡ D ƒ}| jdur| j|d< |r1| d¡}g t|tƒr'|ng ¢d|dœ‘|d< |S )uz  The hop's request body: the refused request's, with the entry's
        overrides merged over it and extended by the continuation.

        The server-side `fallbacks` param is stripped â€” it is mutually exclusive
        with a credit-token retry. When the refusal granted no prefill claim the
        appended turn is omitted entirely and the same-body form is sent.
        c                 S  ó   i | ]\}}|d kr||“qS )rL   r2   ©r)  ÚkeyÚvaluer2   r2   r3   Ú
<dictcomp>M  ó    z(_ChainState.hop_body.<locals>.<dictcomp>Nr  ÚmessagesÚ	assistant)ÚroleÚcontent)r=  Úitemsr:   r9   rn   r  )r1   rÖ   r¿   rX   rM  r2   r2   r3   rÀ   E  s   


ÿþz_ChainState.hop_bodyr^   c                 C  sp   |j dusJ ‚|j j| _|| _|j jrt|jƒng | _| j t	|j |ƒ¡ |j | _
|jdur2|j| _|j| _dS )z7Fold a refused hop's outcome in, readying the next hop.N)rº   r:   rÇ   rç   r?  rð   r:  r·   r/  r@  r;  rï   r<  rÊ   )r1   rÓ   r^   r¿   r2   r2   r3   rÏ   X  s   

z_ChainState.absorb_refusalrØ   rÃ   c           	      C  s
  g }| j s-| jdur-t | j¡}t| d¡ƒ}|dur%| jdur%| j|d< | td|ƒ¡ | 	| j
¡ t | jj¡}t| d¡ƒ}|dur\t| d¡ƒ}|dur\|jdv rX|jnd|d< t| d	¡ƒpdi }t | j¡|d
< ||d	< | td|ƒ¡ | tdddiƒ¡ |S )uS  Degrade to the suppressed refusal when every remaining entry failed
        over HTTP: replay its message_delta verbatim â€” `recommended_model`
        stamped from the final failure (the failed model for capacity errors,
        `null` otherwise) and `usage.iterations` carrying the recorded chain â€”
        then message_stop.
        Nr~   r'  r  r  r  )i­  i  r  ré   r·   r  Úmessage_stopr  )r³   r<  r-  r.  r‹   r9   r´   rs   r  r/  r>  r;  rê   r¶   r^   r·   )	r1   rØ   r«   rï   r~   rê   r  r$  ré   r2   r2   r3   rÈ   d  s(   
z#_ChainState.terminal_failure_framesN)rX   rš   r,   r-   )rX   rš   rÓ   rî   r,   r¼   )r€   rl   r,   r-   r1  rB   )r,   rF  )rÖ   r%   r¿   rF  r,   rš   )rÓ   rî   r^   rl   r¿   rF  r,   r-   )rØ   rÃ   r,   rô   )rC   rD   rE   rF   rG   r4   Úclassmethodr½   rÉ   rÌ   rÎ   r¿   rÀ   rÏ   rÈ   r2   r2   r2   r3   r¼   ì  s8   
 






r¼   r`   Úmodel_labelrl   úlist[dict[str, Any]]c                 C  s€   | j  d¡}t|tƒr8|r8dd„ td|ƒD ƒ}dd„ |D ƒ}|r8t|ƒdkr6t|d  d	¡tƒs6||d d	< |S td
|| j ƒgS )uS  The `usage.iterations` entries a declined hop contributes to the chain.

    The hop's refusal self-reports its iterations: adopt them, retyped to
    `message` (every one of them declined â€” a server-stitched envelope labels
    its last hop `fallback_message`). A single-entry self-report describes the
    hop itself and gains its model â€” `model_label`, the caller-spelled primary
    or the chain entry's id â€” when the wire didn't name one; a multi-entry
    self-report (a server-tool loop) is kept verbatim. Without a self-report,
    one entry is built from the refusal's usage.
    r·   c                 s  r(  r.   )r‹   ©r)  rÖ   r2   r2   r3   r+  �  r,  z._declined_iteration_entries.<locals>.<genexpr>r9  c                 S  s$   g | ]}|d uri |¥ddi¥‘qS )Nr  r~   r2   rV  r2   r2   r3   Ú
<listcomp>�  s   $ z/_declined_iteration_entries.<locals>.<listcomp>rY   r   r^   r~   )ré   r9   rn   r  r   rj   rl   Ú_to_iteration_usage)r`   rT  ÚreportedÚdict_entriesÚentriesr2   r2   r3   r@  ‚  s    r@  Údelta_usagerš   r^   c                 C  sZ   |   d¡}t|tƒr'|r'ttd|ƒd ƒ}|dur'i |¥d|  d¡p#|dœ¥S td|| ƒS )z‡The serving hop's `fallback_message` completer entry: its own
    self-reported iteration relabeled, or one built from its delta usage.r·   r9  r>   NÚfallback_messager^   )r  r^   )r9   rn   r  r‹   r   rX  )r\  r^   rY  Úlastr2   r2   r3   r"  ˜  s   
r"  c                   @  sZ   e Zd ZU dZde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
dS )rõ   a/  Block bookkeeping for one stream of the splice: accumulates each content
    block from its deltas (for the continuation prefill), shifts wire indices
    by `index_base` so they stay monotonic across hops, and tracks which blocks
    are still open so a refusal that cuts mid-block can close them.
    r–   rÊ   r   r°   r,   r-   c                 C  s   || _ || _g | _g | _d S r.   )Ú_index_baserÊ   Ú_blocksr  )r1   r°   r2   r2   r3   r4   ­  s   
z_BlockTracker.__init__rU  c                 C  s   dd„ | j D ƒS )z/The accumulated content blocks, in start order.c                 S  s   g | ]\}}|‘qS r2   r2   )r)  Ú_Úblockr2   r2   r3   rW  ·  ó    z0_BlockTracker.content_blocks.<locals>.<listcomp>)r`  r0   r2   r2   r3   r  µ  s   z_BlockTracker.content_blocksrê   rš   c                 C  sz   |  d¡}t|tƒsdS t|  d¡ƒ}| j ||durt|ƒni f¡ || j }||d< | j |¡ t	| j
|d ƒ| _
dS )z7Track a content_block_start, shifting `event["index"]`.r+   Nr	  rY   )r9   rn   r–   r‹   r`  rs   Údictr_  r  ÚmaxrÊ   )r1   rê   r+   r	  Úshiftedr2   r2   r3   r  ¹  s   

 
z_BlockTracker.startc                 C  sN   |  d¡}t|tƒsdS t|  d¡ƒ}|durt| j||ƒ || j |d< dS )zQApply a content_block_delta to its accumulating block, shifting `event["index"]`.r+   Nr  )r9   rn   r–   r‹   Ú_apply_deltar`  r_  )r1   rê   r+   r  r2   r2   r3   r  Å  s   

z_BlockTracker.deltac                 C  sV   |  d¡}t|tƒsdS || j }||d< || jv r | j |¡ t| j|d ƒ| _dS )z6Track a content_block_stop, shifting `event["index"]`.r+   NrY   )r9   rn   r–   r_  r  Úremovere  rÊ   )r1   rê   r+   rf  r2   r2   r3   r  Ï  s   



z_BlockTracker.stopúIterator[bytes]c                 c  s.   � | j D ]}tdd|dœƒV  q| j  ¡  dS )z4content_block_stop events for any blocks still open.r  rD  N)r  r  r   )r1   r+   r2   r2   r3   r  Ú  s   €
z_BlockTracker.close_open_blocksN)r   )r°   r–   r,   r-   )r,   rU  )rê   rš   r,   r-   ©r,   ri  )rC   rD   rE   rF   rG   r4   r  r  r  r  r  r2   r2   r2   r3   rõ   £  s   
 




rõ   rð   ú list[tuple[int, dict[str, Any]]]r+   r–   r  r-   c                   s*  t ‡ fdd„| D ƒdƒ}|du rdS | d¡}|dkr1t| d¡p"dƒt| d¡p*dƒ |d< dS |dkrKt| d	¡p<dƒt| d
¡pDdƒ |d	< dS |dkrl| d¡}t|tƒs_g }||d< td|ƒ | d¡¡ dS |dkr†t| d¡pwdƒt| d¡pdƒ |d< dS |dkr“| d¡|d< dS dS )zAApply a content_block_delta to the accumulating block at `index`.c                 3  s    � | ]\}}|ˆ kr|V  qd S r.   r2   )r)  Úblock_indexrb  r/   r2   r3   r+  æ  s   € z_apply_delta.<locals>.<genexpr>Nr  Ú
text_deltaÚtextr_   Úinput_json_deltaÚ_partial_jsonÚpartial_jsonÚcitations_deltaÚ	citationsr9  ÚcitationÚthinking_deltaÚthinkingÚsignature_deltaÚ	signature)Únextr9   rl   rn   r  r   rs   )rð   r+   r  rb  Ú
delta_typers  r2   r/   r3   rg  ä  s&   
,,

,ÿrg  Úresponse_blocksrF  c                 C  s´   g }| D ]9}|  d¡dkrq|  d¡}t|tƒs| |¡ qdd„ | ¡ D ƒ}t|ƒ}|dur1|n|  d¡|d< | |¡ q|rX|d   d¡d	v rX| ¡  |rX|d   d¡d	v sI|S )
uç  Convert a hop's accumulated response blocks to the appended assistant turn.

    A `fallback_has_prefill_claim` refusal guarantees the partial output is
    resendable, so the blocks go out near-verbatim. Three rewrites apply:
    `fallback` seam blocks (wire markers, not content) are dropped, tool inputs
    are reassembled from their accumulated `input_json_delta` JSON
    (content_block_start carries `input: {}`), and trailing thinking blocks are
    stripped â€” the server rejects an assistant turn whose final block is
    `thinking`, so a refusal that cut the stream mid-thought would otherwise
    400 the continuation. A partial that was nothing but thinking strips to
    empty, and the hop falls back to the same-body form.
    r  r
  rp  c                 S  rG  )rp  r2   rH  r2   r2   r3   rK    rL  z&_to_prefill_blocks.<locals>.<dictcomp>NÚinputr>   )rv  Úredacted_thinking)r9   rn   rl   rs   rQ  r  Úpop)r{  Úoutrb  rq  Úparsedr2   r2   r3   r?  ú  s    


ÿr?  rX   c           	      C  sæ   |   d¡}t|tƒs| S d}g }td|ƒD ]Q}t|ƒ}|dur$|  d¡nd}|du s/t|tƒs5| |¡ qtd|ƒ}dd„ |D ƒ}t|ƒt|ƒkrO| |¡ qd}|s[|  d	¡d
kr[q| i |¥d|i¥¡ q|sk| S i | ¥d|i¥S )u  A copy of `body` with `{type: "fallback"}` seam blocks filtered out of the
    replayed message history; `body` itself when there are none.

    An assistant turn left with no content by the strip â€” a seam-only turn â€” is
    dropped whole; the server rejects empty-content turns.rM  Fr9  NrP  c                 S  s&   g | ]}t |ƒr| d ¡dks|‘qS )r  r
  )r   r9   )r)  rb  r2   r2   r3   rW  -  s   & z&_strip_seam_blocks.<locals>.<listcomp>TrO  rN  )r9   rn   r  r   r‹   rs   rj   )	rX   rM  ÚstrippedÚstripped_messagesr~   Úmessage_dictrP  Úcontent_listÚkeptr2   r2   r3   re     s.   




re   rS   r   rK   c                   sš   t  dd„ | j ¡ D ƒ¡ dd¡}dd„ | d¡D ƒ‰ t ‡ fdd	„|D ƒ¡}d
d„ | j ¡ D ƒ}|s5|rBd t	d|g|¢ƒ¡|d< | j
t|tdƒƒd�S )zåA copy of `request` with `betas` appended to its `anthropic-beta` header
    (skipping values already present) and the middleware's helper-telemetry tag
    appended to `x-stainless-helper`. Single `request.copy()` for both.
    c                 S  s    i | ]\}}t |tƒr||“qS r2   )rn   rl   ©r)  ÚkÚvr2   r2   r3   rK  A  s     z,_with_middleware_headers.<locals>.<dictcomp>úanthropic-betar_   c                 S  s   h | ]}|  ¡ ’qS r2   )Ústrip)r)  rJ  r2   r2   r3   Ú	<setcomp>B  rc  z+_with_middleware_headers.<locals>.<setcomp>ú,c                 3  s$   � | ]}t |ƒˆ vrt |ƒV  qd S r.   )rl   )r)  r‰   ©Úexistingr2   r3   r+  C  s   €" z+_with_middleware_headers.<locals>.<genexpr>c                 S  s"   i | ]\}}|  ¡ d kr||“qS )r‰  )r‘   r†  r2   r2   r3   rK  D  s   " z, Nzfallback-refusal-middleware)Úheaders)r�   ÚHeadersr�  rQ  r9   Úsplitrd  ÚfromkeysÚjoinÚfilterrf   r   r"   )rS   rK   rÐ   Ú	additionsr�  r2   r�  r3   rd   :  s   "rd   r~   úMessage | BetaMessageró   c                 C  ó    t | tƒr| jdur| jjS dS )zEThe refusal's minted credit token; only the beta surface carries one.N)rn   r#   r  r  ©r~   r2   r2   r3   rp   J  ó   rp   c                 C  r—  )zOThe policy category that caused the refusal; only the beta surface carries one.N)rn   r#   r  r   r˜  r2   r2   r3   rq   Q  r™  rq   r|   r€   r   c                 C  s   dd| id|id|dœdœS )zBThe synthetic `fallback` content block marking one model boundary.r
  r^   r`   )r  r   )r  Úfromr  Útriggerr2   )r|   r€   r   r2   r2   r3   rt   X  s
   ürt   Úoriginalr}   úhttpx.Response | Nonec                 C  s”   t t| j d¡ƒƒ}|du st| d¡tƒsdS ti |¥dg |¢td|d ƒ¢i¥ƒ 	d¡}| j
 ¡ }dD ]	}||v r>||= q5tj| j||| jd�S )zÙA copy of the served hop's response with `seams` prepended to its
    `content`, or `None` when the body isn't the expected message shape.

    The caller has already parsed the response, so its body is buffered.
    úutf-8NrP  r9  ©zcontent-encodingzcontent-length)rÅ   r�  rP  rS   )r‹   r  rP  Údecodern   r9   r  rÆ   r   Úencoder�  rf   r�   ÚResponserÅ   rS   )rœ  r}   r~   rX   r�  Úheaderr2   r2   r3   Ú_seamed_http_responseb  s   ,
€ür¤  rU   c              	   C  ó8   t | j|ƒ}|d u r| S t|| j| jd| j| j| jd�S ©NFrœ   )r¤  rg   r   r¦   r§   r¨   r©   r¢   ©r[   r}   r�   r2   r2   r3   ru   x  ó   ùru   rƒ   c              	   C  r¥  r¦  )r¤  rg   r   r¦   r§   r¨   r©   r¢   r§  r2   r2   r3   r…   ‡  r¨  r…   r
  r%   Úcredit_tokenc                 C  s   i | ¥|¥}|r||d< |S )z�The non-streaming retry body: the fallback entry merged whole over the
    original params, plus the refusal's credit token when it minted one.r  r2   )rX   r
  r©  Úmergedr2   r2   r3   rr   –  s   rr   rn  r   c                 C  s"   zt  | ¡W S  ty   Y d S w r.   )rŒ   ÚloadsrÁ   )rn  r2   r2   r3   r  Ÿ  s
   ÿr  rJ  c                 C  s   t j| dd�S )N)rŒ  ú:)Ú
separators)rŒ   Údumps©rJ  r2   r2   r3   rÆ   ¦  s   rÆ   r=   r‡   c                 C  s   t | ƒr	td| ƒS d S )Nrè   )r   r   r¯  r2   r2   r3   r‹   ª  ó   r‹   c                 C  s&   zt  |  ¡ ¡W S  ty   Y d S w r.   )rŒ   r«  ÚreadrÁ   ©r[   r2   r2   r3   rÄ   ®  s
   ÿrÄ   c                 Ã  s.   �zt  |  ¡ I d H ¡W S  ty   Y d S w r.   )rŒ   r«  ÚareadrÁ   r²  r2   r2   r3   rã   µ  s   €ÿrã   rê   ÚpayloadÚbytesc                 C  s   t | t|ƒd� d¡S )N©rê   r  rž  )Ú_serialize_sserÆ   r¡  )rê   r´  r2   r2   r3   r  ¼  r°  r  r  r   c                 C  s2   | j rd | j ¡d  d¡S t| j| jd� d¡S )zËForward a decoded event in its original wire bytes, preserving SSE fields
    beyond `event:`/`data:` (`id:`, `retry:`, comment lines). Falls back to
    re-serializing for events with no raw lines.
    Ú
z

rž  r¶  )r�   r“  r¡  r·  rê   r  )r  r2   r2   r3   r  À  s   r  r  c                 C  sD   d}| dur|d| › d�7 }|  d¡D ]
}|d|› d�7 }q|d S )zÞSerialize an event back to its SSE wire form (`event: ...\ndata: ...\n\n`).

    Multi-line `data` is emitted as one `data:` line per line, matching the
    spec. The inverse of the decoder behind `Stream.raw_events`.
    r_   Nzevent: r¸  zdata: )r‘  )rê   r  r  Úliner2   r2   r3   r·  Ê  s   r·  r  ú&Literal['message', 'fallback_message']ré   c              	   C  sJ   |pi }| ||  d¡pd|  d¡pd|  d¡pd|  d¡pd|  d¡dœS )NÚinput_tokensr   Úoutput_tokensÚcache_read_input_tokensÚcache_creation_input_tokensÚcache_creation)r  r^   r»  r¼  r½  r¾  r¿  )r9   )r  r^   ré   Úur2   r2   r3   rX  Ø  s   ùrX  Úprimaryc                 C  sP   |pi }i |¥| p
i ¥}|  ¡ D ]\}}|du r%| |¡dur%|| ||< q|S )z0Fill `None` fields on `primary` from `fallback`.N)rQ  r9   )rÁ  r
  r  rI  rJ  r2   r2   r3   r  ç  s   €r  ri   ú,httpx.SyncByteStream | httpx.AsyncByteStreamc                 C  s8   | j  ¡ }dD ]	}||v r||= qtj| j||| jd�S )zÔA synthetic response standing in for `original`, with `stream` as its body.

    The spliced frames are emitted post-decode, so the original's
    content-encoding/length headers no longer describe the body.
    rŸ  )rÅ   r�  ri   rS   )r�  rf   r�   r¢  rÅ   rS   )rœ  ri   r�  r£  r2   r2   r3   r¤   ñ  s   
€ür¤   c                   @  ó2   e Zd Zddd„Zeddd	„ƒZedd
d„ƒZdS )r¥   r«   r®   r,   r-   c                 C  ó
   || _ d S r.   ©Ú_frames©r1   r«   r2   r2   r3   r4     r5   z_FrameByteStream.__init__ri  c                 C  r  r.   rÅ  r0   r2   r2   r3   Ú__iter__	  r  z_FrameByteStream.__iter__c                 C  s   | j  ¡  d S r.   )rÆ  r»   r0   r2   r2   r3   r»     s   z_FrameByteStream.closeN)r«   r®   r,   r-   rj  rB   )rC   rD   rE   r4   r   rÈ  r»   r2   r2   r2   r3   r¥     ó    
r¥   c                   @  rÃ  )r­   r«   rß   r,   r-   c                 C  rÄ  r.   rÅ  rÇ  r2   r2   r3   r4     r5   z_AsyncFrameByteStream.__init__r7  c                 C  r  r.   rÅ  r0   r2   r2   r3   Ú	__aiter__  r  z_AsyncFrameByteStream.__aiter__c                 Ã  s   �| j  ¡ I d H  d S r.   )rÆ  râ   r0   r2   r2   r3   râ     s   €z_AsyncFrameByteStream.acloseN)r«   rß   r,   r-   )r,   r7  rB   )rC   rD   rE   r4   r   rÊ  râ   r2   r2   r2   r3   r­     rÉ  r­   )r[   r3  rÒ   r¸   r,   r4  )r[   r3  rÒ   r¸   r,   r7  )r`   rä   rT  rl   r,   rU  )r\  rš   r^   rl   r,   rš   )rð   rk  r+   r–   r  rš   r,   r-   )r{  rU  r,   rF  )rX   rš   r,   rš   )rS   r   rK   r)   r,   r   )r~   r–  r,   ró   )r|   rl   r€   rl   r   ró   r,   rš   )rœ  r3  r}   rU  r,   r�  )r[   rU   r}   rU  r,   rU   )r[   rƒ   r}   rU  r,   rƒ   )rX   rš   r
  r%   r©  ró   r,   rš   )rn  rl   r,   r   )rJ  r   r,   rl   )rJ  r=   r,   r‡   )r[   r3  r,   r   )rê   rl   r´  rš   r,   rµ  )r  r   r,   rµ  )rê   ró   r  rl   r,   rl   )r  rº  r^   rl   ré   r‡   r,   rš   )rÁ  r‡   r
  r‡   r,   rš   )rœ  r3  ri   rÂ  r,   r3  )gÚ
__future__r   rf   r-  rŒ   ÚloggingÚtypingr   r   r   r   r   r   r	   r
   r   r   r   Úcontextvarsr   r   Útyping_extensionsr   r   r�   Ú_utilsr   Ú_modelsr   Ú_requestr   Ú	_responser   r   Ú
_streamingr   r   r   Ú_exceptionsr   Ú_middlewarer   r   r   Ú_base_clientr   Útypes.messager    Ú_stainless_helpersr"   Útypes.beta.beta_messager#   Útypes.anthropic_beta_paramr$   Útypes.beta.beta_fallback_paramr%   Ú__all__Ú	getLoggerr(   rG   r“   r*   r&   r6   r8   r'   rä   rË   rÃ   rî   r¸   r¹   rà   r¼   r@  r"  rõ   rg  r?  re   rd   rp   rq   rt   r¤  ru   r…   rr   r  rÆ   r‹   rÄ   rã   r  r  r·  rX  r  r¤   ÚSyncByteStreamr¥   ÚAsyncByteStreamr­   r2   r2   r2   r3   Ú<module>   s–    4þÿ    a	 
_
 


A

!









	











