§
    ~ŠtjÆd  ã                  ód  — d Z ddlmZ ddl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 dd	lmZ dd
lmZmZ ddlmZ ddlmZmZ ddlmZmZmZmZm Z  e!e"e"f         Z#e
rddl$m%Z%m&Z&m'Z' ddl(m)Z)  ej*        e+¦  «        Z,dZ-dZ.dZ/e G d„ d¦  «        ¦   «         Z0 G d„ de¦  «        Z1dS )ur  OTel â†’ LangSmith bridge for Pipecat.

Rewrites Pipecat's ``conversation`` / ``turn`` / ``stt`` / ``llm`` / ``tts`` OTel
spans into the ``gen_ai.*`` / ``langsmith.*`` namespaces LangSmith ingests, rolls
the whole conversation onto the root ``conversation`` span, and optionally
attaches the recorded audio there; non-Pipecat spans pass through untouched.
Shared export / ``thread_id`` / message plumbing lives in
:class:`BaseLangSmithSpanProcessor`.

Trace shape::

    conversation                 (root; whole transcript + conversation WAV)
    â””â”€â”€ turn Ã— N
        â”œâ”€â”€ stt                  (audio â†’ transcript)
        â”œâ”€â”€ llm                  (the LLM stage; kind set by llm_span_kind)
        â”œâ”€â”€ tool Ã— N             (realtime: one per tool call, args â†’ result)
        â””â”€â”€ tts                  (response text â†’ audio)

Audio uses Pipecat's ``AudioBufferProcessor``: wire it with
:meth:`attach_audio_buffer` (placed after ``transport.output()``, constructed
with ``num_channels=2``); the merged stereo (user-left / bot-right) it emits is
attached as a WAV when the conversation ends.

``llm_span_kind`` sets the kind for Pipecat's ``llm`` span â€” keep ``"llm"`` for a
service that does its own inference; pass ``"chain"`` when it only orchestrates
nested runs exported separately (else the conversation double-counts).

Speech-to-speech (realtime) services emit operation-named spans: ``llm_response``
(â†’ ``llm`` run, usage + reply) and the ``llm_setup`` / ``llm_request`` wrappers.
Call :meth:`PipecatLangSmithSpanProcessor.instrument_user_aggregator` to supply
finalized user transcripts emitted outside those spans. A realtime tool call
arrives as two sibling spans, ``llm_tool_call`` (args) then
``llm_tool_result`` (output); the processor defers the call and merges the result
onto it, so each call renders as one ``tool`` run spanning call â†’ result.
é    )ÚannotationsN)ÚMutableMapping)Ú	dataclassÚfield)ÚTYPE_CHECKINGÚAnyÚOptional)ÚTTLCache)ÚSpanProcessor)Úget_package_version)Úbuild_assistant_messageÚbuild_user_message)Ú
pcm_to_wav)ÚBaseLangSmithSpanProcessorÚTranslatedSpan)Úbuild_completion_messageÚextract_llm_usageÚ	iso_to_nsÚparse_llm_messagesÚtool_message_key)ÚLLMContextAggregatorPairÚLLMUserAggregatorÚUserTurnMessageAddedMessage)ÚAudioBufferProcessorg      ¬@i † é,   c                  ó`   — e Zd ZU dZded<   ded<    ee¬¦  «        Zded<   dZd	ed
<   dd„Z	dS )Ú_AudioRecordu  Merged PCM accumulated for one conversation, plus its WAV parameters.

    Pipecat's ``AudioBufferProcessor`` emits already-merged audio â€” stereo
    (user-left / bot-right) when constructed with ``num_channels=2`` â€” so we just
    accumulate and wrap it.
    ÚintÚsample_rateÚnum_channels)Údefault_factoryÚ	bytearrayÚpcmFÚboolÚaudio_truncatedÚaudioÚbytesÚlimit_bytesúOptional[int]ÚreturnÚNonec                ó  — || _         || _        |�U|t          | j        ¦  «        z
  }||d|z  z  z  }|dk    r	d| _        dS t          |¦  «        |k    rd| _        |d|…         }| j                             |¦  «         dS )zFAppend PCM (truncating at ``limit_bytes``) and refresh the WAV params.Né   r   T)r   r    Úlenr#   r%   Úextend)Úselfr&   r   r    r(   Ú	remainings         úf/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langsmith/integrations/pipecat/processor.pyr/   z_AudioRecord.extendi   s—   € ð 'ˆÔØ(ˆÔØÐ"Ø#¥c¨$¬(¡m¤mÑ3ˆIØ˜ a¨,Ñ&6Ñ7Ñ7ˆIØ˜AŠ~ˆ~Ø'+�Ô$Ø�Ý�5‰zŒz˜IÒ%Ð%Ø'+�Ô$Ø˜j˜y˜jÔ)�ØŒ�Š˜ÑÔÐÐÐó    N)
r&   r'   r   r   r    r   r(   r)   r*   r+   )
Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú__annotations__r   r"   r#   r%   r/   © r3   r2   r   r   [   s}   € € € € € € ðð ð ÐÐÑØÐÐÑØ�U¨9Ð5Ñ5Ô5€CÐ5Ð5Ð5Ñ5Ø!€OÐ!Ð!Ð!Ñ!ðð ð ð ð ð r3   r   c                  óú   ‡ — e Zd ZdZ	 dEdddddedœdFˆ fd„ZdGˆ fd„ZdHd„ZdId„ZdJd$„Z	dKd(„Z
dLd*„ZdMd/„ZdNd3„ZdOd4„ZdOd5„ZdPd6„ZdOd7„ZdOd8„ZdQd:„ZdRd;„ZdRd<„ZdSd>„ZdPd?„ZdOd@„ZdTdA„ZdUdB„ZdEdVdC„ZdWˆ fdD„Zˆ xZS )XÚPipecatLangSmithSpanProcessorzCEnriches Pipecat's OTel spans with LangSmith-compatible attributes.NÚllmz	audio/wav)Úllm_span_kindÚapi_keyÚprojectÚendpointÚaudio_mime_typeÚstate_ttl_secondsÚdownstream_processorúOptional[SpanProcessor]r=   Ústrr>   úOptional[str]r?   r@   rA   rB   ÚfloatÚkwargsr   r*   r+   c               ón  •—  t          ¦   «         j        |f|||dœ|¤Ž || _        || _        t	          t
          |¬¦  «        | _        t	          t
          |¬¦  «        | _        t	          t
          |¬¦  «        | _        t	          t
          |¬¦  «        | _	        t	          t
          |¬¦  «        | _
        dS )a  Create the processor.

        Args:
            llm_span_kind: LangSmith run kind for Pipecat's ``llm`` span.
            audio_mime_type: MIME type for the attached conversation recording.
            state_ttl_seconds: lifetime for per-conversation state.
        )r>   r?   r@   )ÚmaxsizeÚttlN)ÚsuperÚ__init__Ú_llm_span_kindÚ_audio_mime_typer
   ÚDEFAULT_STATE_MAXSIZEÚ_transcript_by_traceÚ_audio_by_conversationÚ_tool_calls_by_traceÚ_trace_by_threadÚ_pending_user_transcripts)
r0   rC   r=   r>   r?   r@   rA   rB   rH   Ú	__class__s
            €r2   rM   z&PipecatLangSmithSpanProcessor.__init__‚   sø   ø€ ð& 	�‰ŒÔØ ð	
àØØð		
ð 	
ð
 ð	
ð 	
ð 	
ð ,ˆÔØ /ˆÔõ
 Õ2Ð8IÐJÑJÔJð 	Ô!õ JRÝ)Ð/@ðJ
ñ J
ô J
ˆÔ#õ
 PXÝ)Ð/@ðP
ñ P
ô P
ˆÔ!õ
 ;CÝ)Ð/@ð;
ñ ;
ô ;
ˆÔõ Õ2Ð8IÐJÑJÔJð 	Ô&Ð&Ð&r3   Útrace_idr   Ú	thread_idc                óØ   •— t          ¦   «                              ||¦  «         || j        |<   | j                             |d¦  «        }|pg D ]\  }}|                      |||¦  «         ŒdS )uG   Also index threadâ†’trace, draining user turns buffered before the map.N)rL   Ú_remember_thread_idrT   rU   ÚpopÚ_append_transcript)r0   rW   rX   ÚpendingÚsort_keyÚmessagerV   s         €r2   rZ   z1PipecatLangSmithSpanProcessor._remember_thread_id·   sƒ   ø€ å‰Œ×#Ò# H¨iÑ8Ô8Ð8Ø+3ˆÔ˜iÑ(ØÔ0×4Ò4°YÀÑEÔEˆØ!( ¨Bð 	Að 	AÑˆH�gØ×#Ò# H¨g°xÑ@Ô@Ð@Ð@ð	Að 	Ar3   Ú
aggregatorr   c                ó”   ‡ ‡— |                      ¦   «         }t          |¦  «        Š|                     d¦  «        d
ˆ ˆfd„¦   «         }d	S )a  Fold a realtime session's finalized user transcripts into its trace.

        Pipecat realtime services emit the user's finalized text through the user
        context aggregator's ``on_user_turn_message_added`` event rather than on
        an OTel span. Call this once after building the context aggregators::

            processor = configure_pipecat(...)
            aggregators = LLMContextAggregatorPair(context)
            set_thread_id(conversation_id)
            processor.instrument_user_aggregator(aggregators, conversation_id)

        The aggregator is the authoritative source of user turns: a realtime
        service transcribes the user's audio asynchronously, so the text lands
        (via this event, carrying the turn-start timestamp) after the model has
        already replied. The transcript is ordered by that timestamp, not arrival.

        Args:
            aggregator: the ``LLMContextAggregatorPair`` for the conversation.
            thread_id: the conversation id, matching :func:`set_thread_id`.
        Úon_user_turn_message_addedÚ_aggregatorr   r_   r   r*   r+   c              “  óN   •K  — ‰                      ‰|j        |j        ¦  «         d S ©N)Ú_record_user_transcriptÚcontentÚ	timestamp)rc   r_   r0   Útids     €€r2   Ú_on_user_turn_message_addedz]PipecatLangSmithSpanProcessor.instrument_user_aggregator.<locals>._on_user_turn_message_addedÛ   s,   øè è € ð ×(Ò(¨¨g¬o¸wÔ?PÑQÔQÐQÐQÐQr3   N)rc   r   r_   r   r*   r+   )ÚuserrE   Úevent_handler)r0   r`   rX   Úuser_aggregatorrj   ri   s   `    @r2   Úinstrument_user_aggregatorz8PipecatLangSmithSpanProcessor.instrument_user_aggregatorÁ   sq   øø€ ð. %Ÿ/š/Ñ+Ô+ˆÝ�)‰nŒnˆà	×	&Ò	&Ð'CÑ	DÔ	Dð	Rð 	Rð 	Rð 	Rð 	Rð 	Rñ 
EÔ	Dð	Rð 	Rð 	Rr3   Útextrh   c                ó2  — |sdS t          |¦  «        }t          |¦  «        df}| j                             |¦  «        }|€?| j                             |¦  «        pg }|                     ||f¦  «         || j        |<   dS |                      |||¦  «         dS )zDAppend or buffer a finalized user turn, keyed by when it was spoken.Nr   )r   r   rT   ÚgetrU   Úappendr\   )r0   rX   ro   rh   r_   r^   rW   r]   s           r2   rf   z5PipecatLangSmithSpanProcessor._record_user_transcriptá   s®   € ð ð 	ØˆFÝ$ TÑ*Ô*ˆÝ˜iÑ(Ô(¨!Ð,ˆØÔ(×,Ò,¨YÑ7Ô7ˆØÐØÔ4×8Ò8¸ÑCÔCÐIÀrˆGØ�NŠN˜H gÐ.Ñ/Ô/Ð/Ø8?ˆDÔ*¨9Ñ5ØˆFØ×Ò ¨'°8Ñ<Ô<Ð<Ð<Ð<r3   r_   údict[str, Any]r^   Ú_SortKeyc                ó€   — | j                              |¦  «        pg }|                     ||f¦  «         || j         |<   dS )zEAdd a message to the transcript the root renders, keyed for ordering.N)rQ   rq   rr   )r0   rW   r_   r^   Úconversations        r2   r\   z0PipecatLangSmithSpanProcessor._append_transcriptñ   sK   € ð Ô0×4Ò4°XÑ>Ô>ÐDÀ"ˆØ×Ò˜X wÐ/Ñ0Ô0Ð0à.:ˆÔ! (Ñ+Ð+Ð+r3   Úaudio_bufferr   Úconversation_idc                óT   ‡ ‡— dˆˆ fd
„} |                      d¦  «        |¦  «         dS )ug  Record this conversation's audio from a Pipecat ``AudioBufferProcessor``.

        Registers an ``on_audio_data`` handler that accumulates the merged PCM
        Pipecat emits â€” construct the buffer with ``num_channels=2`` for a stereo
        (user-left / bot-right) recording that preserves barge-in overlap â€” and
        attaches it as a WAV when the ``conversation`` span ends.
        ``conversation_id`` must match the ``PipelineTask`` id. Place the buffer
        after ``transport.output()`` and give it a non-zero ``buffer_size`` so the
        final chunk arrives before the span is exported.
        Ú_bufferr   r&   r'   r   r   r    r*   r+   c              “  ó<   •K  — ‰                      ‰|||¦  «         d S re   )Ú_accumulate_audio)rz   r&   r   r    rx   r0   s       €€r2   Ú_on_audio_datazIPipecatLangSmithSpanProcessor.attach_audio_buffer.<locals>._on_audio_data
  s)   øè è € ð ×"Ò" ?°E¸;ÈÑUÔUÐUÐUÐUr3   Úon_audio_dataN)
rz   r   r&   r'   r   r   r    r   r*   r+   )rl   )r0   rw   rx   r}   s   ` ` r2   Úattach_audio_bufferz1PipecatLangSmithSpanProcessor.attach_audio_bufferü   sT   øø€ ð	Vð 	Vð 	Vð 	Vð 	Vð 	Vð 	Vð 	4ˆ×"Ò" ?Ñ3Ô3°NÑCÔCÐCÐCÐCr3   r)   c                óN   — | j         €d S t          d| j         t          z
  ¦  «        S ©Nr   )Úaudio_size_limit_bytesÚmaxÚ_WAV_HEADER_BYTES)r0   s    r2   Ú_pcm_audio_size_limit_bytesz9PipecatLangSmithSpanProcessor._pcm_audio_size_limit_bytes  s)   € ØÔ&Ð.Ø�4Ý�1�dÔ1Õ4EÑEÑFÔFÐFr3   r&   r'   r   r    c                ó
  — | j                              |¦  «        }|€2|dk    rt                               d|¦  «         t	          ||¬¦  «        }|                     ||||                      ¦   «         ¦  «         || j         |<   d S )Nr-   z—langsmith voice: AudioBufferProcessor num_channels=%d; use num_channels=2 for a stereo (user-left/bot-right) recording that preserves barge-in overlap.)r   r    )rR   rq   ÚloggerÚwarningr   r/   r…   )r0   rx   r&   r   r    Úrecs         r2   r|   z/PipecatLangSmithSpanProcessor._accumulate_audio  s•   € ð Ô)×-Ò-¨oÑ>Ô>ˆØˆ;Ø˜qÒ Ð Ý—’ð7ð !ñ	ô ð õ ¨;À\ÐRÑRÔRˆCØ�
Š
�5˜+ |°T×5UÒ5UÑ5WÔ5WÑXÔXÐXà7:ˆÔ# OÑ4Ð4Ð4r3   Útspanr   r$   c                ó|  — |j         j        j        }|j         j        }|dk    r|                      ||¦  «         �n |dk    r|                      ||¦  «         nã|dk    r|                      |¦  «         nÇ|dk    r|                      |¦  «         n«|dk    r|                      ||¦  «         nŽ|dk    r|  	                    ||¦  «         nq|dk    r|  
                    ||¦  «         nT|dk    r|                     d	¦  «         n8|d
k    r|                      ||¦  «        S |dk    r|                      ||¦  «        S dS )NÚsttr<   ÚttsÚturnrv   Úllm_responseÚllm_requestÚ	llm_setupÚchainÚllm_tool_callÚllm_tool_resultT)ÚspanÚcontextrW   ÚnameÚ_handle_sttÚ_handle_llmÚ_handle_ttsÚ_handle_turnÚ_handle_conversationÚ_handle_realtime_responseÚ_handle_llm_requestÚset_kindÚ_handle_tool_callÚ_handle_tool_result)r0   rŠ   rW   r—   s       r2   Ú	_dispatchz'PipecatLangSmithSpanProcessor._dispatch0  sp  € Ø”:Ô%Ô.ˆØŒzŒˆà�5Š=ˆ=Ø×Ò˜U HÑ-Ô-Ð-Ñ-Ø�UŠ]ˆ]Ø×Ò˜U HÑ-Ô-Ð-Ð-Ø�UŠ]ˆ]Ø×Ò˜UÑ#Ô#Ð#Ð#Ø�VŠ^ˆ^Ø×Ò˜eÑ$Ô$Ð$Ð$Ø�^Ò#Ð#Ø×%Ò% e¨XÑ6Ô6Ð6Ð6Ø�^Ò#Ð#Ø×*Ò*¨5°(Ñ;Ô;Ð;Ð;Ø�]Ò"Ð"Ø×$Ò$ U¨HÑ5Ô5Ð5Ð5Ø�[Ò Ð Ø�NŠN˜7Ñ#Ô#Ð#Ð#Ø�_Ò$Ð$Ø×)Ò)¨%°Ñ:Ô:Ð:ØÐ&Ò&Ð&Ø×+Ò+¨E°8Ñ<Ô<Ð<Øˆtr3   c                óü  — |j                              dd¦  «        }|j                              dd¦  «        }|                     d¦  «         |                     t	          d|› d�¦  «        g¬¦  «         |rr|                     t          t          |¦  «        ¦  «        g¬	¦  «         |r?|                      |t	          t          |¦  «        ¦  «        |j        j	        pd
d
f¦  «         | 
                    ¦   «          dS )u+   STT span: audio input â†’ transcribed text.Ú
transcriptÚ Úis_finalTr<   zAudio for: "ú"©Úprompt©Ú
completionr   N)Ú
attributesrq   rŸ   Úset_messagesr   r   rE   r\   r•   Ú
start_timeÚexclude_from_message_view)r0   rŠ   rW   r¤   r¦   s        r2   r˜   z)PipecatLangSmithSpanProcessor._handle_sttL  s  € àÔ%×)Ò)¨,¸Ñ;Ô;ˆ
ØÔ#×'Ò'¨
°DÑ9Ô9ˆØ�Š�uÑÔÐØ×ÒÕ#5Ð6RÀZÐ6RÐ6RÐ6RÑ#SÔ#SÐ"TÐÑUÔUÐUØð 	Ø×ÒÕ+BÅ3ÀzÁ?Ä?Ñ+SÔ+SÐ*TÐÑUÔUÐUð
 ð Ø×'Ò'ØÝ&¥s¨:¡¤Ñ7Ô7Ø”ZÔ*Ð/¨a°Ð3ñô ð ð
 	×'Ò'Ñ)Ô)Ð)Ð)Ð)r3   c                óR  ‡— |j                              dd¦  «        }|j                              dd¦  «        }|                     | j        ¦  «         t	          |¦  «        }|r|                     |¬¦  «         |r$|                     t          |¦  «        g¬¦  «         t          |j         ¦  «        }|r |j        d
i |¤Ž d„ |D ¦   «         }|r"| 	                    t          |¦  «        ¦  «         |r3|j
        j        pdŠˆfd„t          |¦  «        D ¦   «         | j        |<   d	S d	S )zÚ``llm`` span: forward the request history and reply verbatim.

        Each request's ``input`` carries the whole history, so the latest snapshot
        is kept per trace as the conversation the root renders.
        Úinputr¥   Úoutputr¨   rª   c                óD   — g | ]}|                      d ¦  «        dk    ¯|‘ŒS )ÚroleÚsystem)rq   )Ú.0Úms     r2   ú
<listcomp>z=PipecatLangSmithSpanProcessor._handle_llm.<locals>.<listcomp>w  s,   € ÐGÐGÐG˜A¨Q¯UªU°6©]¬]¸hÒ-FÐ-F�aÐ-FÐ-FÐ-Fr3   r   c                ó"   •— g | ]\  }}‰|f|f‘ŒS r9   r9   )r¶   Úir_   Úbases      €r2   r¸   z=PipecatLangSmithSpanProcessor._handle_llm.<locals>.<listcomp>|  s3   ø€ ð 3ð 3ð 3Ù)3¨¨G�$˜�˜GÐ$ð3ð 3ð 3r3   Nr9   )r¬   rq   rŸ   rN   r   r­   r   r   Ú	set_usagerr   r•   r®   Ú	enumeraterQ   )	r0   rŠ   rW   Ú
input_dataÚoutput_dataÚmessagesÚusager¤   r»   s	           @r2   r™   z)PipecatLangSmithSpanProcessor._handle_llm`  sm  ø€ ð Ô%×)Ò)¨'°2Ñ6Ô6ˆ
ØÔ&×*Ò*¨8°RÑ8Ô8ˆØ�Š�tÔ*Ñ+Ô+Ð+å% jÑ1Ô1ˆØð 	0Ø×Ò hÐÑ/Ô/Ð/Øð 	SØ×ÒÕ+CÀKÑ+PÔ+PÐ*QÐÑRÔRÐRå! %Ô"2Ñ3Ô3ˆØð 	%ØˆEŒOÐ$Ð$˜eÐ$Ð$Ð$ð
 HÐG ÐGÑGÔGˆ
Øð 	EØ×ÒÕ6°{ÑCÔCÑDÔDÐDØð 	Ø”:Ô(Ð-¨AˆDð3ð 3ð 3ð 3Ý7@ÀÑ7LÔ7Lð3ñ 3ô 3ˆDÔ% hÑ/Ð/Ð/ð	ð 	r3   c                ó”  — |j                              dd¦  «        }|                     d¦  «         |j                              d¦  «        }|r#|                     dt	          |¦  «        ¦  «         |                     t          t	          |¦  «        ¦  «        gt          d|› d�¦  «        g¬¦  «         |                     ¦   «          dS )	u=   TTS span: text â†’ audio. The voice is metadata, not content.ro   r¥   r<   Úvoice_idzGenerated audio for: "r§   )r©   r«   N)	r¬   rq   rŸ   Úset_metadatarE   r­   r   r   r¯   )r0   rŠ   ro   rÃ   s       r2   rš   z)PipecatLangSmithSpanProcessor._handle_tts€  sË   € àÔ×#Ò# F¨BÑ/Ô/ˆØ�Š�uÑÔÐØÔ#×'Ò'¨
Ñ3Ô3ˆØð 	:Ø×Ò˜z­3¨x©=¬=Ñ9Ô9Ð9Ø×ÒÝ&¥s¨4¡y¤yÑ1Ô1Ð2Ý/Ð0PÈÐ0PÐ0PÐ0PÑQÔQÐRð 	ñ 	
ô 	
ð 	
ð 	×'Ò'Ñ)Ô)Ð)Ð)Ð)r3   c                óâ  — |                      | j        ¦  «         |                      |¦  «        }|r|                     |¬¦  «         |j                             d¦  «        p|j                             dd¦  «        }|rKt          |¦  «        }|                     |g¬¦  «         |                      |||j        j	        pddf¦  «         t          |j        ¦  «        }|r |j        di |¤Ž dS dS )	u>   ``llm_response`` span: conversation so far â†’ realtime reply.r¨   r²   Útext_outputr¥   rª   r   Nr9   )rŸ   rN   Ú_render_messagesr­   r¬   rq   r   r\   r•   r®   r   r¼   )r0   rŠ   rW   rv   r¿   r«   rÁ   s          r2   r�   z7PipecatLangSmithSpanProcessor._handle_realtime_response�  s  € à�Š�tÔ*Ñ+Ô+Ð+Ø×,Ò,¨XÑ6Ô6ˆØð 	4Ø×Ò lÐÑ3Ô3Ð3àÔ&×*Ò*¨8Ñ4Ô4ð 
¸Ô8H×8LÒ8LØ˜2ñ9
ô 9
ˆð ð 	Ý1°+Ñ>Ô>ˆJØ×Ò¨:¨,ÐÑ7Ô7Ð7Ø×#Ò#Ø˜* u¤zÔ'<Ð'AÀÀ1Ð&Eñô ð õ " %Ô"2Ñ3Ô3ˆØð 	%ØˆEŒOÐ$Ð$˜eÐ$Ð$Ð$Ð$Ð$ð	%ð 	%r3   c                ó€  ‡— |                      d¦  «         t          |j                             dd¦  «        ¦  «        }ˆfd„|                      |¦  «        D ¦   «         }|j        j        pd}d}|D ]K}t          |¦  «        Š‰�‰|v rŒ|                     ‰¦  «         |  	                    ||||f¦  «         |dz  }ŒLdS )uÛ  ``llm_request`` span (OpenAI realtime): capture the tool round-trip.

        For OpenAI realtime, tool calls and results have no spans of their own â€”
        they exist only in this context snapshot. User turns come from the
        aggregator and assistant replies from ``llm_response``, so take *only* the
        tool-call / tool-result messages here, deduped by their id, keyed by this
        span's time so they sort after the user turn that triggered them.
        r’   r±   r¥   c                ó6   •— h | ]}t          |¦  «        xŠ®‰’ŒS re   )r   )r¶   r_   Úkeys     €r2   ú	<setcomp>zDPipecatLangSmithSpanProcessor._handle_llm_request.<locals>.<setcomp>­  s6   ø€ ð 
ð 
ð 
àÝ'¨Ñ0Ô0Ð0�Ð=ð à=Ð=Ð=r3   r   Né   )
rŸ   r   r¬   rq   rÇ   r•   r®   r   Úaddr\   )	r0   rŠ   rW   rÀ   Úseenr»   Úaddedr_   rÊ   s	           @r2   rž   z1PipecatLangSmithSpanProcessor._handle_llm_request¢  sï   ø€ ð 	�Š�wÑÔÐÝ% eÔ&6×&:Ò&:¸7ÀBÑ&GÔ&GÑHÔHˆð
ð 
ð 
ð 
à×0Ò0°Ñ:Ô:ð
ñ 
ô 
ˆð
 ŒzÔ$Ð)¨ˆØˆØð 	ð 	ˆGÝ" 7Ñ+Ô+ˆCØˆ{˜c T˜k˜kØØ�HŠH�S‰MŒMˆMØ×#Ò# H¨g¸¸e°}ÑEÔEÐEØ�Q‰JˆEˆEð	ð 	r3   úlist[dict[str, Any]]c                óp   — | j                              |g ¦  «        }d„ t          |d„ ¬¦  «        D ¦   «         S )z=Return the transcript as plain messages, ordered by sort key.c                ó   — g | ]\  }}|‘ŒS r9   r9   )r¶   Ú_r_   s      r2   r¸   zBPipecatLangSmithSpanProcessor._render_messages.<locals>.<listcomp>¿  s   € ÐNÐNÐN™J˜A˜w�ÐNÐNÐNr3   c                ó   — | d         S r�   r9   )Úes    r2   ú<lambda>z@PipecatLangSmithSpanProcessor._render_messages.<locals>.<lambda>¿  s
   € ÈÈ!Ì€ r3   )rÊ   )rQ   rq   Úsorted)r0   rW   Úentriess      r2   rÇ   z.PipecatLangSmithSpanProcessor._render_messages¼  s=   € àÔ+×/Ò/°¸"Ñ=Ô=ˆØNÐN­&°¸n¸nÐ*MÑ*MÔ*MÐNÑNÔNÐNr3   c                ój  — |                      d¦  «         |j                             d¦  «        }|r"|                     t	          |¦  «        ¦  «         |j                             d¦  «        }|�|                     |¦  «         | j                             |g ¦  «                             |¦  «         dS )zG``llm_tool_call`` span: defer the call (args) until its result arrives.Útoolztool.function_nameztool.argumentsNF)	rŸ   r¬   rq   Úset_namerE   Úset_tool_inputrS   Ú
setdefaultrr   )r0   rŠ   rW   Úfunction_nameÚ	argumentss        r2   r    z/PipecatLangSmithSpanProcessor._handle_tool_callÁ  sª   € à�Š�vÑÔÐØÔ(×,Ò,Ð-AÑBÔBˆØð 	/Ø�NŠN�3˜}Ñ-Ô-Ñ.Ô.Ð.ØÔ$×(Ò(Ð)9Ñ:Ô:ˆ	ØÐ Ø× Ò  Ñ+Ô+Ð+ØÔ!×,Ò,¨X°rÑ:Ô:×AÒAÀ%ÑHÔHÐHØˆur3   c                óð  — |j                              d¦  «        }|                      |¦  «        }|€.|                     d¦  «         |�|                     |¦  «         dS |�|                     |¦  «         |j                              d¦  «        }|�#|                     dt          |¦  «        ¦  «         |j        j        �| 	                    |j        j        ¦  «         |  
                    |¦  «         dS )ub  ``llm_tool_result`` span: merge the result onto its deferred call span.

        Pairs with the oldest deferred call for the trace (Gemini Live puts no
        usable id on the result span, so pairing is by order). The result becomes
        the tool run's output and stretches the span to end at the result, so it
        spans call â†’ result.
        ztool.resultNrÚ   Tztool.result_statusÚtool_result_statusF)r¬   rq   Ú_take_pending_callrŸ   Úset_tool_outputrÄ   rE   r•   Úend_timeÚset_end_timeÚ_export)r0   rŠ   rW   ÚresultÚ
call_tspanÚstatuss         r2   r¡   z1PipecatLangSmithSpanProcessor._handle_tool_resultÍ  só   € ð Ô!×%Ò% mÑ4Ô4ˆØ×,Ò,¨XÑ6Ô6ˆ
ØÐØ�NŠN˜6Ñ"Ô"Ð"ØÐ!Ø×%Ò% fÑ-Ô-Ð-Ø�4àÐØ×&Ò& vÑ.Ô.Ð.ØÔ!×%Ò%Ð&:Ñ;Ô;ˆØÐØ×#Ò#Ð$8½#¸f¹+¼+ÑFÔFÐFØŒ:ÔÐ*Ø×#Ò# E¤JÔ$7Ñ8Ô8Ð8Ø�Š�ZÑ Ô Ð Øˆur3   úOptional[TranslatedSpan]c                ó¦   — | j                              |¦  «        }|sdS |                     d¦  «        }|s| j                              |d¦  «         |S )z=Pop the oldest deferred call span for the trace, or ``None``.Nr   )rS   rq   r[   )r0   rW   Úqueuerè   s       r2   râ   z0PipecatLangSmithSpanProcessor._take_pending_callç  s[   € àÔ)×-Ò-¨hÑ7Ô7ˆØð 	Ø�4Ø—Y’Y˜q‘\”\ˆ
Øð 	:ØÔ%×)Ò)¨(°DÑ9Ô9Ð9ØÐr3   c                óü   — |                      d¦  «         |j                             d¦  «        }|�|                     d|¦  «         |j                             d¦  «        }|�|                     d|¦  «         dS dS )zATurn span: a framework wrapper around one exchange (a ``chain``).r’   zturn.numberNÚturn_numberzturn.was_interruptedÚturn_was_interrupted)rŸ   r¬   rq   rÄ   )r0   rŠ   rî   Úwas_interrupteds       r2   r›   z*PipecatLangSmithSpanProcessor._handle_turnñ  sˆ   € à�Š�wÑÔÐØÔ&×*Ò*¨=Ñ9Ô9ˆØÐ"Ø×Ò˜}¨kÑ:Ô:Ð:ØÔ*×.Ò.Ð/EÑFÔFˆØÐ&Ø×ÒÐ5°ÑGÔGÐGÐGÐGð 'Ð&r3   c                ó  — |j                              dd¦  «        p|j                              dd¦  «        }|                     d¦  «         |                     d¦  «         |                     dd¦  «         |                     dd	¦  «         |                     d
t          d¦  «        pd¦  «         |                      |¦  «        }|r|                     |¬¦  «         |                      ||¦  «         |  	                    ||¦  «         dS )z9Conversation span: the whole session; the LangSmith root.zconversation.idr¥   rx   r’   TÚls_modalityr&   Úls_integrationÚpipecatÚls_integration_versionz
pipecat-air¨   N)
r¬   rq   rŸ   Úset_root_spanrÄ   r   rÇ   r­   Ú_attach_conversation_audioÚ_cleanup_conversation)r0   rŠ   rW   rx   rÀ   s        r2   rœ   z2PipecatLangSmithSpanProcessor._handle_conversationû  s$  € àÔ*×.Ò.Ø˜rñ
ô 
ð 9àÔ×!Ò!Ð"3°RÑ8Ô8ð 	ð 	�Š�wÑÔÐØ×Ò˜DÑ!Ô!Ð!Ø×Ò˜=¨'Ñ2Ô2Ð2Ø×ÒÐ+¨YÑ7Ô7Ð7Ø×ÒØ$Õ&9¸,Ñ&GÔ&GÐ&MÈ2ñ	
ô 	
ð 	
ð ×(Ò(¨Ñ2Ô2ˆØð 	0Ø×Ò hÐÑ/Ô/Ð/à×'Ò'¨¨Ñ?Ô?Ð?Ø×"Ò" 8¨_Ñ=Ô=Ð=Ð=Ð=r3   c                ó  — | j                              |¦  «        }|€Ht          | j         ¦  «        dk    r.t                               d|t          | j         ¦  «        ¦  «         dS |j        sdS |                      ¦   «         }|�5t          |j        ¦  «        |k    rt                               d|¦  «         dS t          t          |j        ¦  «        |j
        |j        ¦  «        }|                      |d|| j        ¬¦  «         dS )uæ   Encode the accumulated PCM and attach it, keyed strictly by id.

        Never falls back to "the only recording in flight" â€” a heuristic match
        could attach another caller's audio to this trace (a privacy leak).
        Nr   z¥langsmith voice: no recording for conversation_id=%r (%d in flight under other ids); audio not attached. Ensure attach_audio_buffer's id matches the PipelineTask id.zFlangsmith voice: skipped oversize pipecat audio for conversation_id=%rzconversation.wav)r—   ÚdataÚ	mime_type)rR   rq   r.   r‡   Údebugr#   r…   rˆ   r   r'   r   r    Ú_attach_audiorO   )r0   rŠ   rx   Ú	recordingÚpcm_limit_bytesÚwavs         r2   r÷   z8PipecatLangSmithSpanProcessor._attach_conversation_audio  s'  € ð Ô/×3Ò3°OÑDÔDˆ	ØÐÝ�4Ô.Ñ/Ô/°!Ò3Ð3Ý—’ðLð $Ý˜Ô3Ñ4Ô4ñô ð ð ˆFØŒ}ð 	ØˆFØ×:Ò:Ñ<Ô<ˆØÐ&­3¨y¬}Ñ+=Ô+=ÀÒ+OÐ+OÝ�NŠNð%àñô ð ð
 ˆFÝÝ�)”-Ñ Ô  )Ô"7¸Ô9Oñ
ô 
ˆð 	×ÒØÐ*°ÀÔ@Uð 	ñ 	
ô 	
ð 	
ð 	
ð 	
r3   c                ój  — | j                              |¦  «        }| j                             |d ¦  «         | j                             |d ¦  «         |�6| j                             |d ¦  «         | j                             |d ¦  «         |                      |¦  «         |                      |¦  «         d S re   )	Ú_thread_id_by_tracerq   rQ   r[   rR   rT   rU   Ú_forget_thread_idÚ_flush_tool_calls)r0   rW   rx   rX   s       r2   rø   z3PipecatLangSmithSpanProcessor._cleanup_conversation5  s°   € ØÔ,×0Ò0°Ñ:Ô:ˆ	ØÔ!×%Ò% h°Ñ5Ô5Ð5ØÔ#×'Ò'¨¸Ñ>Ô>Ð>ØÐ ØÔ!×%Ò% i°Ñ6Ô6Ð6ØÔ*×.Ò.¨y¸$Ñ?Ô?Ð?Ø×Ò˜xÑ(Ô(Ð(Ø×Ò˜xÑ(Ô(Ð(Ð(Ð(r3   c                ó¬   — |�|gnt          | j        ¦  «        }|D ]7}| j                             |d¦  «        pg D ]}|                      |¦  «         ŒŒ8dS )z�Export held tool-call spans whose result never arrived (args only).

        Scoped to one trace at conversation end, or all held spans at shutdown.
        N)ÚlistrS   r[   ræ   )r0   rW   Ú	trace_idsri   rŠ   s        r2   r  z/PipecatLangSmithSpanProcessor._flush_tool_calls?  s~   € ð #Ð.ˆXˆJˆJµD¸Ô9RÑ4SÔ4Sð 	ð ð 	$ð 	$ˆCØÔ2×6Ò6°s¸DÑAÔAÐGÀRð $ð $�Ø—’˜UÑ#Ô#Ð#Ð#ð$ð	$ð 	$r3   c                óp   •— |                       ¦   «          t          ¦   «                              ¦   «          dS )z@Flush any still-held tool-call spans, then shut down downstream.N)r  rL   Úshutdown)r0   rV   s    €r2   r	  z&PipecatLangSmithSpanProcessor.shutdownK  s1   ø€ à×ÒÑ Ô Ð Ý‰Œ×ÒÑÔÐÐÐr3   re   )rC   rD   r=   rE   r>   rF   r?   rF   r@   rF   rA   rE   rB   rG   rH   r   r*   r+   )rW   r   rX   rE   r*   r+   )r`   r   rX   rE   r*   r+   )rX   rE   ro   rE   rh   rE   r*   r+   )rW   r   r_   rs   r^   rt   r*   r+   )rw   r   rx   rE   r*   r+   )r*   r)   )
rx   rE   r&   r'   r   r   r    r   r*   r+   )rŠ   r   r*   r$   )rŠ   r   rW   r   r*   r+   )rŠ   r   r*   r+   )rW   r   r*   rÐ   )rŠ   r   rW   r   r*   r$   )rW   r   r*   rê   )rŠ   r   rx   rE   r*   r+   )rW   r   rx   rE   r*   r+   )rW   r)   r*   r+   )r*   r+   )r4   r5   r6   r7   ÚDEFAULT_STATE_TTL_SECONDSrM   rZ   rn   rf   r\   r   r…   r|   r¢   r˜   r™   rš   r�   rž   rÇ   r    r¡   râ   r›   rœ   r÷   rø   r  r	  Ú__classcell__)rV   s   @r2   r;   r;      sw  ø€ € € € € ØMÐMð 9=ð3Kð #Ø!%Ø!%Ø"&Ø*Ø#<ð3Kð 3Kð 3Kð 3Kð 3Kð 3Kð 3Kð 3KðjAð Að Að Að Að AðRð Rð Rð Rð@=ð =ð =ð =ð ;ð ;ð ;ð ;ðDð Dð Dð Dð0Gð Gð Gð Gð
;ð ;ð ;ð ;ð.ð ð ð ð8*ð *ð *ð *ð(ð ð ð ð@*ð *ð *ð *ð%ð %ð %ð %ð*ð ð ð ð4Oð Oð Oð Oð

ð 
ð 
ð 
ðð ð ð ð4ð ð ð ðHð Hð Hð Hð>ð >ð >ð >ð,"
ð "
ð "
ð "
ðH)ð )ð )ð )ð
$ð 
$ð 
$ð 
$ð 
$ðð ð ð ð ð ð ð ð ð r3   r;   )2r7   Ú
__future__r   ÚloggingÚcollections.abcr   Údataclassesr   r   Útypingr   r   r	   Ú
cachetoolsr
   Úopentelemetry.sdk.tracer   Ú$langsmith._internal._package_versionr   Ú"langsmith._internal.voice._helpersr   r   Úlangsmith._internal.voice.audior   Ú-langsmith._internal.voice.base_span_processorr   r   Ú'langsmith.integrations.pipecat._helpersr   r   r   r   r   Útupler   rt   Ú5pipecat.processors.aggregators.llm_response_universalr   r   r   Ú/pipecat.processors.audio.audio_buffer_processorr   Ú	getLoggerr4   r‡   r
  rP   r„   r   r;   r9   r3   r2   ú<module>r     sX  ðð"ð "ðH #Ð "Ð "Ð "Ð "Ð "à €€€Ø *Ð *Ð *Ð *Ð *Ð *Ø (Ð (Ð (Ð (Ð (Ð (Ð (Ð (Ø /Ð /Ð /Ð /Ð /Ð /Ð /Ð /Ð /Ð /à Ð Ð Ð Ð Ð Ø 1Ð 1Ð 1Ð 1Ð 1Ð 1à DÐ DÐ DÐ DÐ DÐ Dðð ð ð ð ð ð ð ð 7Ð 6Ð 6Ð 6Ð 6Ð 6ðð ð ð ð ð ð ð ðð ð ð ð ð ð ð ð ð ð ð ð ð ð ��c�Œ?€àð ðð ð ð ð ð ð ð ð ð ð
ð ð ð ð ð ð 
ˆÔ	˜8Ñ	$Ô	$€ð #Ð ØÐ àÐ ð ð ð  ð  ð  ð  ñ  ô  ñ „ð ðFOð Oð Oð Oð OÐ$>ñ Oô Oð Oð Oð Or3   