Pipecatインテグレーション

Speech Engineの背後にあるLLMブレインとしてPipecatパイプラインを使用します。

このガイドでは、Speech Engineのブレインサーバー内でLLMパイプラインとしてPipecatを使用する方法を紹介します。Speech Engineは音声認識、ターンテイキング、テキスト読み上げといった音声ループを処理し、Pipecatはプロセッサーを組み合わせたパイプライン(LLM呼び出し、RAG、関数呼び出し、ガードレール、コンテンツフィルター)によってテキスト生成を処理します。

Pipecatはサーバー側のPythonフレームワークであるため、このガイドはPythonのみを対象としています。 パイプラインプロセッサーに対応するNode版はありません。pipecat-client-jsパッケージはありますが、 Pipecatサーバーと通信するブラウザクライアントであり、TypeScriptでパイプラインを構築するためのものではありません。

アーキテクチャ

Speech Engine SDKは外側のレイヤーとして実行され、ユーザーが話し終えるたびにon_transcriptコールバックが発火します。コールバック内でPipecatパイプラインを構築し、会話履歴をLLMContextFrameとして入力して、パイプラインのテキスト出力をSpeech Engineにストリーミングで返します。ElevenLabsがテキストを音声に変換し、ユーザーに再生します。

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

Pipecatパイプラインは1ターンの間だけ実行されます。新しいトランスクリプトが届くと、次のパイプラインを実行する前に前のパイプラインがキャンセルされます。これにより、Speech Engineの割り込み処理がパイプラインにも反映されます。

このパターンを使う場面

Pipecatは、ブレインに単一のLLM呼び出し以上の機能が必要な場合に役立ちます。

  • 検索拡張生成、関数呼び出し、ガードレール向けの組み合わせ可能なプロセッサー
  • 各ステップでトラフィックを検査、変換、ブロックできるフレームベースのミドルウェア
  • 複数のエージェントで共有できる再利用可能なパイプラインフラグメント

ブレインが「トランスクリプトを受け取り、LLMを呼び出す」だけなら、Speech Engineクイックスタートの方が簡単です。パイプラインそのものが重要な場合はPipecatを使ってください。

前提条件

  • Speech Engine。Speech Engineクイックスタートに従って作成してください。
  • Python 3.10+(pipecat-aiで必須)。
  • ブレインサーバー用のパブリックHTTPSトンネル(例:ngrok)。

依存関係をインストールする

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

pipecat-ai[openai]にはOpenAI LLMサービスが含まれます。必要に応じて、extraを別のプロバイダー(anthropic、googleなど)に置き換えてください。

Pipecatブレインを構築する

ブレインは2つの部分で構成されます。ストリーミングされたテキストをasyncio.Queueに取り込むTextSinkプロセッサーと、1ターンのパイプラインを構築して非同期イテレーターとしてチャンクを生成するrun_pipecat_brainコルーチンです。

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)

パイプラインにはLLMサービスとシンクのみが含まれ、Speech Engineが処理するため、STTまたはTTSプロセッサーは含まれません。入力はLLMContextFrameで、出力はLLMTextFrameチャンクです。

run_pipecat_brainは非同期ジェネレーターです。生成された各チャンクは直接Speech Engineに渡されるため、完全な応答が準備できる前にエージェントが話し始めます。

Speech Engineサーバーに接続する

Speech Engine SDKのsend_responseは文字列または文字列の任意の非同期イテラブルを受け取るため、run_pipecat_brain(transcript)を直接渡せます。Speech Engineが提供するConversationMessageオブジェクトは、ブレインに渡す前に通常のdictへ変換してください。

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

新しいトランスクリプトが届くと、Speech Engine SDKは前のターンのタスクをキャンセルします。これにより、run_pipecat_brain内のtry/finallyブロックを通じて、非同期ジェネレーターと基盤となるPipelineTaskがキャンセルされます。

サーバーを実行する

ngrok http 3001
python server.py

クイックスタートで紹介している同じトークンエンドポイントとクライアントコードを使って、ブラウザからSpeech Engineに接続します。Pipecatパイプラインはサーバー側で実行され、ブラウザからは通常のSpeech Engine会話として見えます。

パイプラインを拡張する

テキスト専用のPipecatパイプラインには、LLMTextFrameまたはLLMContextFrameを処理する任意のフレームプロセッサーを追加できます。一般的な追加例は次のとおりです。

  • ガードレール:LLMの前に配置するFrameProcessor。LLMContextFrameを検査し、安全でないコンテキストを置換またはブロックします。
  • 関数呼び出し:OpenAILLMServiceにツールを登録すると、Pipecatがツール呼び出しフレームをネイティブに処理します。最終的なアシスタントテキストは引き続きLLMTextFrameとして届きます。
  • マルチステージ推論:2つのOpenAILLMServiceインスタンスを連結し、その間に2回目の処理用コンテキストを書き換えるカスタムプロセッサーを配置します。
  • 出力フィルタリング:LLMの後に配置するFrameProcessor。各LLMTextFrameを検査し、許可されていないコンテンツをTextSinkに到達する前に削除または書き換えます。

パイプラインの形状はPipeline([processor_a, llm, processor_b, sink])のままで、run_pipecat_brainを変更する必要はありません。

本番環境での考慮事項

  • キャンセルの安全性:パイプラインが完全に開始される前にPipelineTask.cancel()を呼び出すと、デッドロックする場合があります(pipecat-ai/pipecat#4276)。上記のtry/finallyパターンは、少なくとも1つのフレームがキューに入った後にのみcancel()が実行されるため、安全です。
  • プロンプトインジェクション:音声認識の出力はユーザー入力です。特に下流のプロセッサーがツール呼び出しやデータベースクエリでテキストを使用する場合は、LLMに渡す前にトランスクリプトを検証または正規化してください。
  • ブレインサーバー認証:Speech Engineに共有シークレットを設定し、ブレインサーバー側で確認することで、/wsエンドポイントへの不正な接続を防止できます。
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • LLMプロバイダー:pipecat-ai[openai]にはOpenAILLMServiceが含まれます。Anthropicの場合はpipecat-ai[anthropic]をインストールしてAnthropicLLMServiceを使用します。パイプラインの残りの部分は変わりません。

次のステップ