Vai alla navigazione

Integrazione con Pipecat

Usa una pipeline Pipecat come cervello LLM per Speech Engine.

Questa guida mostra come usare Pipecat come pipeline LLM all’interno di un server brain di Speech Engine. Speech Engine gestisce il ciclo vocale — speech-to-text, gestione dei turni e text-to-speech — mentre Pipecat gestisce la generazione di testo tramite una pipeline componibile di processor (chiamate LLM, RAG, chiamate di funzione, guardrail, filtri dei contenuti).

Questa guida è solo per Python perché Pipecat è un framework Python lato server. Non esiste un equivalente Node per i processor della pipeline; esiste un pacchetto pipecat-client-js, ma è un client per browser che comunica con un server Pipecat, non un modo per creare pipeline in TypeScript.

Architettura

L’SDK di Speech Engine funge da livello esterno: il relativo callback on_transcript si attiva ogni volta che l’utente finisce di parlare. All’interno del callback, crei una pipeline Pipecat, inserisci la cronologia della conversazione come LLMContextFrame e trasmetti l’output di testo della pipeline a Speech Engine. ElevenLabs converte il testo in parlato e lo riproduce per l’utente.

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 pipeline Pipecat viene eseguita solo per la durata di un turno. Quando arriva una nuova trascrizione, la pipeline precedente viene annullata prima dell’esecuzione della successiva: è così che la gestione delle interruzioni di Speech Engine si propaga nella pipeline.

Quando usare questo pattern

Pipecat dà il meglio quando il tuo brain richiede più di una singola chiamata LLM:

  • Processor componibili per generazione aumentata dal recupero, chiamate di funzione o guardrail
  • Middleware basato su frame che può ispezionare, trasformare o bloccare il traffico in ogni fase
  • Frammenti di pipeline riutilizzabili e condivisi tra più agenti

Se il tuo brain è “trascrizione in entrata, chiamata LLM in uscita”, il quickstart di Speech Engine è più semplice. Usa Pipecat quando la parte interessante è la pipeline stessa.

Prerequisiti

  • Un Speech Engine. Segui il quickstart di Speech Engine per crearne uno.
  • Python 3.10+ (richiesto da pipecat-ai).
  • Tunnel HTTPS pubblico per il server brain (ad es. ngrok).

Installa le dipendenze

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

pipecat-ai[openai] include il servizio LLM di OpenAI. Se preferisci, sostituisci l’extra con quello di un altro provider (anthropic, google ecc.).

Crea il brain Pipecat

Il brain ha due componenti: un processor TextSink che trasferisce il testo in streaming in una asyncio.Queue e una coroutine run_pipecat_brain che crea una pipeline per un singolo turno e produce chunk come iteratore asincrono.

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 pipeline contiene solo il servizio LLM e il sink, senza processor STT o TTS, perché Speech Engine li gestisce. LLMContextFrame è l’input; i chunk LLMTextFrame sono l’output.

run_pipecat_brain è un generatore asincrono. Ogni chunk prodotto viene inviato direttamente a Speech Engine, quindi l’agente inizia a parlare prima che la risposta completa sia pronta.

Collegalo al server Speech Engine

send_response dell’SDK di Speech Engine accetta una stringa o qualsiasi iterabile asincrono di stringhe, quindi puoi passare direttamente run_pipecat_brain(transcript). Converti gli oggetti ConversationMessage forniti da Speech Engine in semplici dict prima di passarli al brain.

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())

L’SDK di Speech Engine annulla l’attività del turno precedente quando arriva una nuova trascrizione, annullando il generatore asincrono e il PipelineTask sottostante tramite il blocco try/finally in run_pipecat_brain.

Avvia il server

ngrok http 3001
python server.py

Collegati a Speech Engine da un browser usando lo stesso endpoint token e codice client mostrati nel quickstart. La pipeline Pipecat viene eseguita lato server; il browser vede una normale conversazione di Speech Engine.

Estendi la pipeline

Una pipeline Pipecat solo testuale può includere qualsiasi frame processor che operi su LLMTextFrame o LLMContextFrame. Ecco alcune aggiunte comuni:

  • Guardrail: un FrameProcessor posizionato prima dell’LLM che ispeziona LLMContextFrame e sostituisce o blocca il contesto non sicuro.
  • Chiamate di funzione: registra strumenti su OpenAILLMService e Pipecat gestisce nativamente i frame delle chiamate agli strumenti. Il testo finale dell’assistente arriva comunque come LLMTextFrame.
  • Ragionamento multi-fase: concatena due istanze di OpenAILLMService, con un processor personalizzato intermedio che riscrive il contesto per il secondo passaggio.
  • Filtraggio dell’output: un FrameProcessor posizionato dopo l’LLM che ispeziona ogni LLMTextFrame e rimuove o riscrive i contenuti non consentiti prima che raggiungano TextSink.

La struttura della pipeline rimane la stessa — Pipeline([processor_a, llm, processor_b, sink]) — e run_pipecat_brain non cambia.

Considerazioni per la produzione

  • Sicurezza dell’annullamento: PipelineTask.cancel() può causare un deadlock se viene chiamato prima che la pipeline sia stata avviata completamente (pipecat-ai/pipecat#4276). Il pattern try/finally sopra è sicuro perché cancel() viene eseguito solo dopo che è stato inserito in coda almeno un frame.
  • Prompt injection: l’output speech-to-text è input dell’utente. Convalida o normalizza la trascrizione prima di inviarla all’LLM, soprattutto se un processor downstream usa il testo in chiamate di strumenti o query al database.
  • Autenticazione del server brain: imposta un segreto condiviso su Speech Engine e verificalo nel server brain per impedire connessioni non autorizzate al tuo endpoint /ws:
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • Provider LLM: pipecat-ai[openai] include OpenAILLMService. Per Anthropic, installa pipecat-ai[anthropic] e usa AnthropicLLMService; il resto della pipeline non cambia.

Passaggi successivi