Integración con Pipecat

Usa una canalización de Pipecat como el cerebro LLM detrás de Speech Engine.

Esta guía muestra cómo usar Pipecat como canalización de LLM dentro de un servidor de cerebro de Speech Engine. Speech Engine gestiona el ciclo de voz —voz a texto, gestión de turnos y texto a voz—, mientras que Pipecat gestiona la generación de texto mediante una canalización componible de procesadores (llamadas a LLM, RAG, llamadas a funciones, protecciones y filtros de contenido).

Esta guía es solo para Python porque Pipecat es un framework de Python del lado del servidor. No hay un equivalente en Node para los procesadores de canalización; existe un paquete pipecat-client-js, pero es un cliente de navegador que se comunica con un servidor Pipecat, no una forma de crear canalizaciones en TypeScript.

Arquitectura

El SDK de Speech Engine funciona como capa externa: su callback on_transcript se activa cada vez que el usuario termina de hablar. Dentro del callback, creas una canalización de Pipecat, introduces el historial de la conversación como un LLMContextFrame y transmites la salida de texto de la canalización de vuelta a Speech Engine. ElevenLabs convierte el texto en voz y lo reproduce para el usuario.

User speaks (audio) on_transcript(history) LLMContextFrame(history) LLMTextFrame chunks send_response(async iterator) Agent speaks (audio) Browser ElevenLabs Brain Server (engine.serve) Pipecat Pipeline

La canalización de Pipecat se ejecuta solo durante un turno. Cuando llega una nueva transcripción, se cancela la canalización anterior antes de ejecutar la siguiente; así es como la gestión de interrupciones de Speech Engine se propaga a la canalización.

Cuándo usar este patrón

Pipecat destaca cuando tu cerebro necesita más que una única llamada a un LLM:

  • Procesadores componibles para generación aumentada por recuperación, llamadas a funciones o protecciones
  • Middleware basado en frames que puede inspeccionar, transformar o bloquear el tráfico en cada paso
  • Fragmentos de canalización reutilizables y compartidos entre varios agentes

Si tu cerebro es «entra una transcripción, sale una llamada a un LLM», la guía de inicio rápido de Speech Engine es más sencilla. Usa Pipecat cuando la propia canalización sea la parte interesante.

Requisitos previos

Instala las dependencias

pip install "pipecat-ai[openai]" "elevenlabs" "python-dotenv"

pipecat-ai[openai] incluye el servicio LLM de OpenAI. Sustituye el extra por otro proveedor (anthropic, google, etc.) si lo prefieres.

Crea el cerebro de Pipecat

El cerebro tiene dos componentes: un procesador TextSink que vuelca texto transmitido a una asyncio.Queue, y una corrutina run_pipecat_brain que crea una canalización de un turno y genera fragmentos como iterador asíncrono.

brain.py
import asyncio
import os
from typing import AsyncIterator
from dotenv import load_dotenv
from pipecat.frames.frames import (
Frame,
LLMContextFrame,
LLMFullResponseEndFrame,
LLMTextFrame,
)
from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineTask
from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
from pipecat.services.openai.llm import OpenAILLMService
load_dotenv()
SYSTEM_PROMPT = (
"You are a helpful voice assistant. Keep responses concise and conversational."
)
class TextSink(FrameProcessor):
"""Drain LLMTextFrame text into an asyncio.Queue."""
def __init__(self, queue: asyncio.Queue):
super().__init__()
self._queue = queue
async def process_frame(self, frame: Frame, direction: FrameDirection):
await super().process_frame(frame, direction)
if isinstance(frame, LLMTextFrame):
await self._queue.put(frame.text)
elif isinstance(frame, LLMFullResponseEndFrame):
await self._queue.put(None) # sentinel
await self.push_frame(frame, direction)
def build_messages(transcript: list[dict]) -> list[dict]:
messages = [{"role": "system", "content": SYSTEM_PROMPT}]
for turn in transcript:
role = "assistant" if turn["role"] == "agent" else turn["role"]
messages.append({"role": role, "content": turn["content"]})
return messages
async def run_pipecat_brain(transcript: list[dict]) -> AsyncIterator[str]:
"""Yield response text chunks from a one-turn Pipecat pipeline."""
llm = OpenAILLMService(
api_key=os.environ["OPENAI_API_KEY"],
model="gpt-4o-mini",
)
queue: asyncio.Queue[str | None] = asyncio.Queue()
sink = TextSink(queue)
task = PipelineTask(Pipeline([llm, sink]))
runner = PipelineRunner(handle_sigint=False)
async def drive():
context = LLMContext(build_messages(transcript))
await task.queue_frame(LLMContextFrame(context))
await task.stop_when_done()
run_task = asyncio.create_task(runner.run(task))
drive_task = asyncio.create_task(drive())
try:
while True:
chunk = await queue.get()
if chunk is None:
break
yield chunk
finally:
await task.cancel()
await asyncio.gather(run_task, drive_task, return_exceptions=True)

La canalización contiene solo el servicio LLM y el receptor, sin procesadores de STT ni TTS, porque Speech Engine se encarga de ellos. LLMContextFrame es la entrada; los fragmentos de LLMTextFrame son la salida.

run_pipecat_brain es un generador asíncrono. Cada fragmento generado va directamente a Speech Engine, por lo que el agente empieza a hablar antes de que esté lista la respuesta completa.

Conéctalo al servidor de Speech Engine

send_response del SDK de Speech Engine acepta una cadena o cualquier iterable asíncrono de cadenas, así que puedes pasar run_pipecat_brain(transcript) directamente. Convierte los objetos ConversationMessage que proporciona Speech Engine en diccionarios simples antes de pasarlos al cerebro.

server.py
import asyncio
import os
from dotenv import load_dotenv
from elevenlabs import AsyncElevenLabs
from brain import run_pipecat_brain
load_dotenv()
elevenlabs = AsyncElevenLabs(api_key=os.environ["ELEVENLABS_API_KEY"])
SPEECH_ENGINE_ID = os.environ["SPEECH_ENGINE_ID"]
async def on_transcript(transcript, session):
history = [{"role": m.role, "content": m.content} for m in transcript]
await session.send_response(run_pipecat_brain(history))
async def main():
engine = await elevenlabs.speech_engine.get(SPEECH_ENGINE_ID)
await engine.serve(
port=3001,
path="/ws",
debug=True,
on_transcript=on_transcript,
)
if __name__ == "__main__":
asyncio.run(main())

El SDK de Speech Engine cancela la tarea del turno anterior cuando llega una nueva transcripción, lo que cancela el generador asíncrono y el PipelineTask subyacente mediante el bloque try/finally de run_pipecat_brain.

Ejecuta el servidor

ngrok http 3001
python server.py

Conéctate a Speech Engine desde un navegador usando la misma ruta de API de token y código de cliente que se muestra en la guía de inicio rápido. La canalización de Pipecat se ejecuta en el servidor; el navegador ve una conversación normal de Speech Engine.

Amplía la canalización

Una canalización de Pipecat solo de texto puede incluir cualquier procesador de frames que opere con LLMTextFrame o LLMContextFrame. Algunas adiciones habituales:

  • Protecciones: un FrameProcessor situado antes del LLM que inspecciona LLMContextFrame y sustituye o bloquea contexto no seguro.
  • Llamadas a funciones: registra herramientas en OpenAILLMService y Pipecat gestiona los frames de llamadas a herramientas de forma nativa. El texto final del asistente sigue llegando como LLMTextFrame.
  • Razonamiento en varias etapas: encadena dos instancias de OpenAILLMService, con un procesador personalizado entre ellas que reescribe el contexto para la segunda pasada.
  • Filtrado de salida: un FrameProcessor situado después del LLM que inspecciona cada LLMTextFrame y elimina o reescribe contenido no permitido antes de que llegue a TextSink.

La estructura de la canalización sigue siendo la misma: Pipeline([processor_a, llm, processor_b, sink]); run_pipecat_brain no cambia.

Consideraciones para producción

  • Seguridad de cancelación: PipelineTask.cancel() puede bloquearse indefinidamente si se llama antes de que la canalización se haya iniciado por completo (pipecat-ai/pipecat#4276). El patrón try/finally anterior es seguro porque cancel() solo se ejecuta después de que se haya encolado al menos un frame.
  • Inyección de prompts: la salida de voz a texto es entrada del usuario. Valida o normaliza la transcripción antes de proporcionársela al LLM, especialmente si algún procesador posterior usa el texto en llamadas a herramientas o consultas a bases de datos.
  • Autenticación del servidor de cerebro: configura un secreto compartido en Speech Engine y compruébalo en el servidor de cerebro para evitar conexiones no autorizadas a tu ruta de API /ws:
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • Proveedor de LLM: pipecat-ai[openai] incluye OpenAILLMService. Para Anthropic, instala pipecat-ai[anthropic] y usa AnthropicLLMService; el resto de la canalización no cambia.

Siguientes pasos