コンテンツへ移動

実践ガイド:オープンソースエージェントフレームワークとElevenAgents

執筆者
Akhil Chauhan
公開日
最終更新日

聴くこの記事を聴く

前回のElevenLabsの音声オーケストレーションと外部エージェントの統合に関する記事では、既存のテキストベースのエージェントオーケストレーションをカスタムLLMを介してElevenLabsに接続する方法を紹介しました。このガイドでは、その基盤をもとに、主要なオープンソースのエージェントフレームワークをカスタムLLMインターフェースの背後に適応・デプロイする方法を解説します。その結果、状態管理、ツールオーケストレーション、アプリケーション固有の制御を損なうことなく、成熟したエージェントシステムに音声を重ねられる柔軟なアーキテクチャが実現します。フレームワークにかかわらず、共通して次の3ステップに従います。生成リクエストを作成し、最終的なテキスト応答を抽出して、OpenAI互換のServer-Sent Events(SSE)形式に再構成します。ElevenLabsはChat CompletionsResponsesの両形式をサポートしています。このガイドでは広く利用されている4つのフレームワークを扱いますが、このパターンはOpenAI互換のストリーミング出力を生成できるあらゆるランタイムに応用できます。

A proxy layer translates between ElevenLabs voice orchestration and an agent framework, converting OpenAI-style messages into framework inputs and streaming SSE chunks back as agent voice output.

共通セットアップ

このセクションの例ではPythonとFastAPIを使用しますが、HTTP POSTリクエストとSSEレスポンスのストリーミングを処理できるスタックであれば使用できます。ElevenLabsの音声オーケストレーションが発話終了の可能性を検出すると、設定済みのカスタムLLMエンドポイントに生成リクエストを送信します。このセクションでは、音声オーケストレーションとエージェントフレームワークを同じ言語でつなぐブリッジ、つまりプロキシとなる変換レイヤーの主要コンポーネントを説明します。

当然ながら、各フレームワークは一般的な知名度や特定の目的を果たす能力を理由に選ばれる場合があります。たとえばLlamaIndexは、もともと検索拡張生成(RAG)のセットアップを簡単にするために開発され、CrewAIはエージェント時代における定義済みタスクの自動化を目的として構築されました。設計目標が異なれば応答構造も異なり、それぞれに固有の処理が必要です。ターン全体の完了を待つのではなく、LLMが生成したチャンクをその都度ストリーミングすることは重要です。これにより、テキスト読み上げ(TTS)モデルがより早く音声の生成を開始でき、体感レイテンシーを低減できます。ここでは、LangGraph、Google ADK、CrewAI、LlamaIndexという代表的な4つのフレームワークに焦点を当てます。

共有コードについて

各フレームワークは、OpenAI互換のSSEチャンクとして応答をストリーミングする必要があります。ここでは、これらのチャンクを構築するために各例で使用する小さなヘルパー関数を紹介します。

def sse_chunk(response_id: str, delta: dict, finish_reason=None) -> str:
    payload = {
        "id": response_id,
        "object": "chat.completion.chunk",
        "choices": [{"index": 0, "delta": delta, "finish_reason": finish_reason}],
    }
    return f"data: {json.dumps(payload)}\n\n"

準備ができたので、LangGraphから始めましょう。 

LangGraph

LangGraphでは、エージェントをグラフとしてモデル化します。ノードは個々のステップを表し、エッジはノード間の制御フローを定義します。最小限のセットアップはシンプルです。チャットモデルを初期化し、エージェントツールを定義して、エージェントグラフのランタイムを作成します。

from langchain.agents import create_agent
from langchain_openai import ChatOpenAI
from langchain_core.tools import tool
llm = ChatOpenAI(
    model=model_id,
    api_key=os.getenv("OPENAI_API_KEY"),
)
agent = create_agent(
    llm,
	tools=tool_list,
	system_prompt=system_prompt,
)

LangGraph Agentは生成リクエストごとに会話履歴全体を受け取るため、必要な状態を内部で維持できます。LangGraphはチェックポイントによるサーバー側の永続化をサポートしていますが、実装を最小限に抑えるため、ここでは扱いません。

状態管理ができたら、次にLangGraph固有の判断ポイントとなるのがストリーミングモードです。LangGraphには、用途に応じた次の2つの選択肢があります。

  • stream_mode="values"はグラフ状態のスナップショットを提供します。実装は簡単ですが、各応答により完全なメッセージ状態が含まれるため、リアルタイムの会話フローではレイテンシーが増加します。
  • stream_mode="messages"は、モデルから増分メッセージチャンクをストリーミングします。ElevenLabsのオーケストレーションレイヤーで音声の最初の出力までの時間を短縮できるため、一般にリアルタイム音声インタラクションにはこちらが推奨されます。

より具体的には、エージェントループのmessages実装には、音声として読み上げるべきではないツール呼び出しの更新など、中間ステップが含まれます。プロキシはこれらを除外し、ユーザー向けの応答テキストだけをTTSレイヤーに渡します。ツールを使用するターンの例を示します。

[1] モデルがツールを呼び出すと判断する(tool_calls=["get_price"])
[2] ツールが実行され、データを返す(result="$24.99") 
[3] モデルが結果を使って応答を生成する(content="価格は$24.99です") 

当然、SSEストリームで転送すべきなのはステップ3のチャンクだけです。実際には、ストリーミングループ内の2つのガードチェックでこのフィルタリングを行います。1つはlanggraph_node == "model"イベントだけを保持し、もう1つは空のコンテンツをスキップします。これらのチェックにより、ユーザー向けアシスタントテキストだけがSSEとしてElevenLabsに転送されます。これらの概念を組み合わせた、軽量なリクエストプロキシの実装を紹介します。

@app.post("/chat/completions")
async def chat_completions(req: ChatCompletionRequest):
    input = {"messages": req.messages}
    async def stream():
        response_id = f"chatcmpl-{uuid.uuid4().hex[:12]}"
        sent_role = False
        async for message_chunk, metadata in agent.astream(input, stream_mode="messages"):
            # Only forward model text chunks; skip tool updates and non-text events.
            if metadata.get("langgraph_node") != "model":
                continue
            content = getattr(message_chunk, "content", None)
            if not content:
                continue
            if not sent_role:
                yield sse_chunk(response_id, {"role": "assistant"})
                sent_role = True
            # Send incremental token-like chunks to ElevenLabs in OpenAI format.
            yield sse_chunk(response_id, {"content": content})
         # Signal natural completion before using the finish_reason: "stop" [DONE]
        yield sse_chunk(response_id, {}, finish_reason="stop")
        yield "data: [DONE]\n\n"
    return StreamingResponse(stream(), media_type="text/event-stream")

これにより、ユーザー向けのモデルチャンクだけがElevenLabsに転送されます。LangGraphでは内部ツールの実行が状態ストリームを通じて可視化されるため、フィルタリングは明示的に行われ、プロキシによって制御されます。 

次に、GoogleのAgent Development Kit(ADK)を扱う際のポイントを見ていきます。

Google ADK

GoogleのADKは、ランタイムループをAgent、Runner、SessionServiceという少数のコアプリミティブの背後に抽象化します。ADKのRunnerは、HTTPレイヤーとエージェント定義の間に位置します。メッセージルーティング、ツールオーケストレーション、セッションライフサイクル、イベントストリーミングを処理します。 

from google.adk.agents import Agent
from google.adk.runners import Runner
from google.adk.agents.run_config import RunConfig, StreamingMode
from google.adk.sessions import InMemorySessionService
from google.genai import types as genai_types
agent = Agent(
    name=name,
    model=model,
    instruction=instruction,
    tools=[tool_list],
)
session_service = InMemorySessionService()
	runner = Runner(
	agent=agent,
	app_name=app_name,
	session_service=session_service
)

エージェント、セッションバックエンド、Runnerを初期化したら、プロキシは受信リクエストごとにADKセッションを取得または作成します。ADKでは、session_idがメモリの永続化を制御します。同じsession_idをターン間で再利用すると、履歴、ツール呼び出し、過去の応答が自動的に引き継がれます。会話IDはElevenLabsの上流にあるため、プロキシがこのマッピングを明示的に処理します。生成リクエストに正しいIDを渡すことで、SDKは過去のコンテキストを内部で処理できます。会話の開始時には、追加パラメータをリクエスト本文に渡して任意のIDを指定します。  

メッセージとセッションを準備すれば、Runnerを呼び出せます。実行中、ツール呼び出しとツール結果は内部ADKイベントとして引き続き表示されますが、ユーザー向け出力ではなく中間的なオーケストレーションステップとして扱われます。そのため、ツール呼び出しがユーザーに見えるテキストとして現れるフレームワークと比べ、手動フィルターは不要です。 

以下のハンドラーは、セッションの解決と取得・作成ロジックをインラインで含む簡略化した実装です。

@app.post("/chat/completions")
async def chat_completions(req: ChatCompletionRequest, request: Request):
    # In production, prefer a stable identifier from your upstream system.
    session_id = req.elevenlabs_extra_body.arbitrary_identifier
    session = await session_service.get_session(
        app_name="elevenlabs", user_id="user", session_id=session_id
    )
    if not session:
        session = await session_service.create_session(
            app_name="elevenlabs", user_id="user", session_id=session_id
        )
    user_text = next((m["content"] for m in reversed(req.messages) if m["role"] == "user"), "")
    content = genai_types.Content(role="user", parts=[genai_types.Part(text=user_text)])
   async def stream():
        response_id = f"chatcmpl-{uuid.uuid4().hex[:12]}"
        sent_role = False
        async for event in runner.run_async(
            user_id="user",
            session_id=session.id,
            new_message=content,
            run_config=RunConfig(streaming_mode=StreamingMode.SSE),
        ):
            if not event.content or not event.content.parts:
                continue
            # In SSE mode, ADK emits partial (incremental) and final (complete) events.
            # Forwarding only partial events avoids duplicating the full text.
            # Note: SSE streaming is experimental in ADK. For production, reconcile
            # both event types in case the model backend doesn't emit partials.
            if not getattr(event, "partial", False):
                continue
            text = "".join((getattr(p, "text", "") or "") for p in event.content.parts)
            if not text:
                continue
            if not sent_role:
                yield sse_chunk(response_id, {"role": "assistant"})
                sent_role = True
            yield sse_chunk(response_id, {"content": text})
        yield sse_chunk(response_id, {}, finish_reason="stop")
        yield "data: [DONE]\n\n"
    return StreamingResponse(stream(), media_type="text/event-stream")

次に、設計上よりタスク中心のCrewAIを見ていきます。

CrewAI

CrewAIは、自由形式の対話ループではなく、構造化されたタスク(調査、執筆、要約)を中心にマルチエージェントワークフローをオーケストレーションするために設計されています。エージェントはロール、目標、バックストーリーによって定義されます。実行は、それぞれ明確な説明と期待する出力を持つTaskオブジェクトを中心に行われます。 

from crewai import Agent, Task, Crew, Process, LLM
from crewai.tools import tool
from crewai.types.streaming import StreamChunkType
llm = LLM(
    model=model_id,
    api_key=os.getenv("OPENAI_API_KEY")
)
store_agent = Agent(
    role=role,
    goal=goal,
    backstory=backstory,
    tools=tools,
    llm=llm,
    verbose=False,
)

LangGraphやADKで使われるエージェントループモデルとは異なり、CrewAIでは通常、会話のそのターンにおける作業単位を定義するため、リクエストごとにTaskとCrewを構築します。プレースホルダーを介して前のターンを次のタスクに挿入し、会話コンテキストを引き継ぎます。{crew_chat_messages}変数にはリクエストごとに現在の会話履歴が設定され、実行時にタスクの説明へ補間されます。また、中間トレースパターン(Thought、Action、Action Input、Observation)を明示的に除外し、最終回答のテキストだけを出力することで、音声出力に適したクリーンなテキストを生成します。 

以下のハンドラーは、リクエストごとのタスク構築、履歴の補間、Crewレベルのストリーミング、トレースのフィルタリング、出力フォーマットをまとめたものです。 

@app.post("/chat/completions")
async def chat_completions(req: ChatCompletionRequest):
    # Task and Crew are assembled per request (not at startup).	
    task = Task(
        description=(
            "Conversation history:\n{crew_chat_messages}\n\n"
            "Respond to the user's latest message."
        ),
        expected_output=expected_output,
        agent=store_agent,
    )
    # stream=True returns CrewStreamingOutput instead of a single CrewOutput.
    crew = Crew(
        agents=[store_agent],
        tasks=[task],
        process=Process.sequential,
        verbose=False,
        stream=True,
    )
    async def stream():
        response_id = f"chatcmpl-{uuid.uuid4().hex[:12]}"
        sent_role = False
        final_marker = "final answer:"
        marker_buffer = ""
        marker_found = False
        emitted_any_content = False
        streaming = await crew.kickoff_async(
            inputs={"crew_chat_messages": json.dumps(req.messages)}
        )
       async for chunk in streaming:
            # Skip non-text events (e.g. tool calls).
            if chunk.chunk_type != StreamChunkType.TEXT or not chunk.content:
                continue
            # Only forward text after the "Final Answer:" marker
            if not marker_found:
                marker_buffer += chunk.content
                idx = marker_buffer.lower().find("final answer:")
                if idx == -1:
                    continue
                marker_found = True
                content = marker_buffer[idx + 13:].lstrip()
                marker_buffer = ""
            else:
                content = chunk.content
            # Clean up any trailing markdown artifacts from CrewAI output.
            content = content.rstrip("`").rstrip()
            if not content:
                continue
            if not sent_role:
                yield sse_chunk(response_id, {"role": "assistant"})
                sent_role = True
            yield sse_chunk(response_id, {"content": content})
        # Fallback to handle short responses without the "Final Answer:" marker
        if not sent_role:
            raw = getattr(streaming, "result", None)
            fallback = (raw.raw if raw else marker_buffer).strip().rstrip("`").rstrip()
            if fallback:
                yield sse_chunk(response_id, {"role": "assistant"})
                yield sse_chunk(response_id, {"content": fallback})
        yield sse_chunk(response_id, {}, finish_reason="stop")
        yield "data: [DONE]\n\n"

次に、ネイティブのイベント駆動型ストリーミングモデルに焦点を当てた、異なるアプローチのLlamaIndexを見ていきます。

LlamaIndex

この投稿で扱うほかのフレームワークとは異なり、LlamaIndexはLLMを外部データソース(ドキュメントストア、インデックス、検索パイプライン)に接続するために設計されています。そのエージェントレイヤーであるFunctionAgentは、この基盤の上で構造化されたコンテキストを取得・推論するものであり、自由な対話やタスク実行を目的とするものではありません。

from llama_index.llms.openai import OpenAI
from llama_index.core.agent.workflow import FunctionAgent, AgentStream
from llama_index.core.base.llms.types import ChatMessage, MessageRole
llm = OpenAI(
    model=model,
    api_key=os.getenv("OPENAI_API_KEY")
)
agent = FunctionAgent(
    tools=[list_inventory, get_item_price],
    llm=llm,
    system_prompt=system_prompt,
)

会話の連続性を保つため、プロキシは受信メッセージをLlamaIndexのチャットメッセージに変換し、最新のユーザーターン(user_msg)と過去のターン(chat_history)に分割します。各AgentStreamイベントのevent.deltaフィールドには次のテキスト断片が含まれ、OpenAI形式のdelta.contentチャンクに直接マッピングされます。空でないdeltaはそのまま転送できるため、これが本ガイドで最もシンプルなストリーミングブリッジです。ストリームには、オーケストレーションイベント(ツール呼び出し、結果)と音声イベント(アシスタントのテキストdelta)の両方が含まれます。音声出力をクリーンに保つため、プロキシはAgentStreamイベントだけを保持し、空のdeltaをスキップします。

[1] AgentStream(delta='')       ← 無視
[2] ToolCall                     ← 無視
[3] ToolCallResult               ← 無視
[4] AgentStream(delta='It')     ← 転送 ✓
[5] AgentStream(delta=' costs')← 転送 ✓
[6] AgentStream(delta=' $49.99')← 転送 ✓

この分離により、中間的なツールの仕組みを読み上げ出力から除外しながら、低レイテンシーの増分音声生成を維持できます。以下のそのまま使えるハンドラーは、これらの手順をまとめたものです。

@app.post("/chat/completions")
async def chat_completions(req: ChatCompletionRequest):
    # This assumes the last message is always a user turn with string content.
    # For production, add defensive role/content handling for non-text payloads.
    chat_history = [
        ChatMessage(role=MessageRole(m["role"]), content=m.get("content") or "")
        for m in req.messages
    ]
    user_text = chat_history.pop().content
    async def stream():
        response_id = f"chatcmpl-{uuid.uuid4().hex[:12]}"
        handler = agent.run(user_msg=user_text, chat_history=chat_history)
        async for event in handler.stream_events():
            if not isinstance(event, AgentStream):
                continue
            if not event.delta:
                continue
            yield sse_chunk(response_id, {"content": event.delta})
        yield sse_chunk(response_id, {}, finish_reason="stop")
        yield "data: [DONE]\n\n"
    return StreamingResponse(stream(), media_type="text/event-stream")

LlamaIndexは、より強力な組み込みオーケストレーションレイヤーを持つフレームワークと比べ、エンドツーエンドの会話ランタイムパターンについて規定が少なめです。本番デプロイでは通常、セッション処理、応答ガードレール、ツールオーケストレーション、トレーシングを実装する必要があります。

まとめ

このガイドの各フレームワークは、同じ契約を通じてElevenLabsに接続します。OpenAI形式のCompletionsまたはResponsesリクエストを受け取り、SSEチャンクをストリーミングで返します。これにより、既存のエージェント実装に最小限の変更で音声オーケストレーションを重ねられます。すでに構築したものを維持しつつ、リアルタイムの会話型AIを実現できます。このモジュール性はElevenAgentsプラットフォームの中核となる考え方です。既存のエージェントを拡張する場合でも、最初から音声ネイティブで構築する場合でも、ElevenAgentの音声オーケストレーションは、それぞれの状況に合わせて利用できます。

すでにオープンソースのフレームワークでエージェントを運用していて、音声を有効にしたい場合は、ぜひこのアプローチを試して感想をお聞かせください。

関連記事

最高品質のAIオーディオで創造する