Add ElevenLabsHttpTTSService
This commit is contained in:
103
examples/foundational/07d-interruptible-elevenlabs-http.py
Normal file
103
examples/foundational/07d-interruptible-elevenlabs-http.py
Normal file
@@ -0,0 +1,103 @@
|
|||||||
|
#
|
||||||
|
# Copyright (c) 2024–2025, Daily
|
||||||
|
#
|
||||||
|
# SPDX-License-Identifier: BSD 2-Clause License
|
||||||
|
#
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
|
||||||
|
import aiohttp
|
||||||
|
from dotenv import load_dotenv
|
||||||
|
from loguru import logger
|
||||||
|
from runner import configure
|
||||||
|
|
||||||
|
from pipecat.audio.vad.silero import SileroVADAnalyzer
|
||||||
|
from pipecat.frames.frames import EndFrame
|
||||||
|
from pipecat.pipeline.pipeline import Pipeline
|
||||||
|
from pipecat.pipeline.runner import PipelineRunner
|
||||||
|
from pipecat.pipeline.task import PipelineParams, PipelineTask
|
||||||
|
from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext
|
||||||
|
from pipecat.services.elevenlabs import ElevenLabsHttpTTSService
|
||||||
|
from pipecat.services.openai import OpenAILLMService
|
||||||
|
from pipecat.transports.services.daily import DailyParams, DailyTransport
|
||||||
|
|
||||||
|
load_dotenv(override=True)
|
||||||
|
|
||||||
|
logger.remove(0)
|
||||||
|
logger.add(sys.stderr, level="DEBUG")
|
||||||
|
|
||||||
|
|
||||||
|
async def main():
|
||||||
|
async with aiohttp.ClientSession() as session:
|
||||||
|
(room_url, token) = await configure(session)
|
||||||
|
|
||||||
|
transport = DailyTransport(
|
||||||
|
room_url,
|
||||||
|
token,
|
||||||
|
"Respond bot",
|
||||||
|
DailyParams(
|
||||||
|
audio_out_enabled=True,
|
||||||
|
transcription_enabled=True,
|
||||||
|
vad_enabled=True,
|
||||||
|
vad_analyzer=SileroVADAnalyzer(),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
tts = ElevenLabsHttpTTSService(
|
||||||
|
api_key=os.getenv("ELEVENLABS_API_KEY", ""),
|
||||||
|
voice_id=os.getenv("ELEVENLABS_VOICE_ID", ""),
|
||||||
|
)
|
||||||
|
|
||||||
|
llm = OpenAILLMService(api_key=os.getenv("OPENAI_API_KEY"), model="gpt-4o")
|
||||||
|
|
||||||
|
messages = [
|
||||||
|
{
|
||||||
|
"role": "system",
|
||||||
|
"content": "You are a helpful LLM in a WebRTC call. Your goal is to demonstrate your capabilities in a succinct way. Your output will be converted to audio so don't include special characters in your answers. Respond to what the user said in a creative and helpful way.",
|
||||||
|
},
|
||||||
|
]
|
||||||
|
|
||||||
|
context = OpenAILLMContext(messages)
|
||||||
|
context_aggregator = llm.create_context_aggregator(context)
|
||||||
|
|
||||||
|
pipeline = Pipeline(
|
||||||
|
[
|
||||||
|
transport.input(), # Transport user input
|
||||||
|
context_aggregator.user(), # User responses
|
||||||
|
llm, # LLM
|
||||||
|
tts, # TTS
|
||||||
|
transport.output(), # Transport bot output
|
||||||
|
context_aggregator.assistant(), # Assistant spoken responses
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
task = PipelineTask(
|
||||||
|
pipeline,
|
||||||
|
PipelineParams(
|
||||||
|
allow_interruptions=True,
|
||||||
|
enable_metrics=True,
|
||||||
|
enable_usage_metrics=True,
|
||||||
|
report_only_initial_ttfb=True,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
@transport.event_handler("on_first_participant_joined")
|
||||||
|
async def on_first_participant_joined(transport, participant):
|
||||||
|
await transport.capture_participant_transcription(participant["id"])
|
||||||
|
# Kick off the conversation.
|
||||||
|
messages.append({"role": "system", "content": "Please introduce yourself to the user."})
|
||||||
|
await task.queue_frames([context_aggregator.user().get_context_frame()])
|
||||||
|
|
||||||
|
@transport.event_handler("on_participant_left")
|
||||||
|
async def on_participant_left(transport, participant, reason):
|
||||||
|
await task.queue_frame(EndFrame())
|
||||||
|
|
||||||
|
runner = PipelineRunner()
|
||||||
|
|
||||||
|
await runner.run(task)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
asyncio.run(main())
|
||||||
@@ -6,9 +6,11 @@
|
|||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import base64
|
import base64
|
||||||
|
import io
|
||||||
import json
|
import json
|
||||||
from typing import Any, AsyncGenerator, Dict, List, Literal, Mapping, Optional, Tuple
|
from typing import Any, AsyncGenerator, Dict, List, Literal, Mapping, Optional, Tuple
|
||||||
|
|
||||||
|
import aiohttp
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
from pydantic import BaseModel, model_validator
|
from pydantic import BaseModel, model_validator
|
||||||
|
|
||||||
@@ -33,6 +35,8 @@ from pipecat.transcriptions.language import Language
|
|||||||
# See .env.example for ElevenLabs configuration needed
|
# See .env.example for ElevenLabs configuration needed
|
||||||
try:
|
try:
|
||||||
import websockets
|
import websockets
|
||||||
|
from elevenlabs import Voice, VoiceSettings
|
||||||
|
from elevenlabs.client import ElevenLabs
|
||||||
except ModuleNotFoundError as e:
|
except ModuleNotFoundError as e:
|
||||||
logger.error(f"Exception: {e}")
|
logger.error(f"Exception: {e}")
|
||||||
logger.error(
|
logger.error(
|
||||||
@@ -418,3 +422,114 @@ class ElevenLabsTTSService(WordTTSService, WebsocketService):
|
|||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} exception: {e}")
|
logger.error(f"{self} exception: {e}")
|
||||||
|
|
||||||
|
|
||||||
|
class ElevenLabsHttpTTSService(WordTTSService):
|
||||||
|
class InputParams(BaseModel):
|
||||||
|
language: Optional[Language] = Language.EN
|
||||||
|
optimize_streaming_latency: Optional[str] = None
|
||||||
|
stability: Optional[float] = None
|
||||||
|
similarity_boost: Optional[float] = None
|
||||||
|
style: Optional[float] = None
|
||||||
|
use_speaker_boost: Optional[bool] = None
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
api_key: str,
|
||||||
|
voice_id: str,
|
||||||
|
model: str = "eleven_flash_v2_5",
|
||||||
|
output_format: ElevenLabsOutputFormat = "pcm_24000",
|
||||||
|
params: InputParams = InputParams(),
|
||||||
|
**kwargs,
|
||||||
|
):
|
||||||
|
sample_rate = self._sample_rate_from_output_format(output_format)
|
||||||
|
super().__init__(
|
||||||
|
aggregate_sentences=True,
|
||||||
|
push_text_frames=False,
|
||||||
|
push_stop_frames=True,
|
||||||
|
stop_frame_timeout_s=2.0,
|
||||||
|
sample_rate=sample_rate,
|
||||||
|
**kwargs,
|
||||||
|
)
|
||||||
|
|
||||||
|
self._client = ElevenLabs(api_key=api_key)
|
||||||
|
self._voice_id = voice_id
|
||||||
|
self._model = model
|
||||||
|
self._output_format = output_format
|
||||||
|
|
||||||
|
# Create voice settings if provided
|
||||||
|
self._voice_settings = None
|
||||||
|
if params.stability is not None and params.similarity_boost is not None:
|
||||||
|
self._voice_settings = VoiceSettings(
|
||||||
|
stability=params.stability,
|
||||||
|
similarity_boost=params.similarity_boost,
|
||||||
|
style=params.style or 0.0,
|
||||||
|
use_speaker_boost=params.use_speaker_boost or False,
|
||||||
|
)
|
||||||
|
|
||||||
|
logger.debug(f"Initialized with sample rate: {sample_rate}")
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _sample_rate_from_output_format(output_format: str) -> int:
|
||||||
|
return {
|
||||||
|
"pcm_16000": 16000,
|
||||||
|
"pcm_22050": 22050,
|
||||||
|
"pcm_24000": 24000,
|
||||||
|
"pcm_44100": 44100,
|
||||||
|
}[output_format]
|
||||||
|
|
||||||
|
async def start(self, frame: StartFrame):
|
||||||
|
await super().start(frame)
|
||||||
|
|
||||||
|
async def stop(self, frame: EndFrame):
|
||||||
|
await super().stop(frame)
|
||||||
|
|
||||||
|
async def cancel(self, frame: CancelFrame):
|
||||||
|
await super().cancel(frame)
|
||||||
|
|
||||||
|
async def run_tts(self, text: str) -> AsyncGenerator[Frame, None]:
|
||||||
|
def read_audio_stream(**kwargs):
|
||||||
|
# Run the streaming in a separate thread
|
||||||
|
audio_chunks = []
|
||||||
|
stream = self._client.text_to_speech.convert_as_stream(**kwargs)
|
||||||
|
for chunk in stream:
|
||||||
|
if chunk:
|
||||||
|
audio_chunks.append(chunk)
|
||||||
|
return b"".join(audio_chunks)
|
||||||
|
|
||||||
|
try:
|
||||||
|
yield TTSStartedFrame()
|
||||||
|
|
||||||
|
# Prepare parameters
|
||||||
|
params = {
|
||||||
|
"text": text,
|
||||||
|
"voice_id": self._voice_id,
|
||||||
|
"model_id": self._model,
|
||||||
|
"output_format": self._output_format,
|
||||||
|
"voice_settings": self._voice_settings,
|
||||||
|
"optimize_streaming_latency": 4, # Maximum optimization + disabled text normalizer
|
||||||
|
}
|
||||||
|
|
||||||
|
# Get audio data in a separate thread
|
||||||
|
audio_data = await asyncio.to_thread(read_audio_stream, **params)
|
||||||
|
|
||||||
|
if not audio_data:
|
||||||
|
logger.error(f"{self} No audio data returned")
|
||||||
|
yield None
|
||||||
|
return
|
||||||
|
|
||||||
|
# Stream the audio data in chunks
|
||||||
|
chunk_size = 4096 # Adjust this value as needed
|
||||||
|
for i in range(0, len(audio_data), chunk_size):
|
||||||
|
chunk = audio_data[i : i + chunk_size]
|
||||||
|
if len(chunk) > 0:
|
||||||
|
yield TTSAudioRawFrame(
|
||||||
|
chunk, self._sample_rate_from_output_format(self._output_format), 1
|
||||||
|
)
|
||||||
|
|
||||||
|
yield TTSStoppedFrame()
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error in run_tts: {e}")
|
||||||
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
Reference in New Issue
Block a user