Adding telnyx serializer

This commit is contained in:
Rafal Skorski
2025-01-23 15:39:46 +01:00
parent 89b87289e2
commit 8eef21db6e
4 changed files with 125 additions and 7 deletions

View File

@@ -16,7 +16,9 @@ from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext
from pipecat.serializers.twilio import TwilioFrameSerializer
from pipecat.serializers.telnyx import TelnyxFrameSerializer
from pipecat.services.elevenlabs import ElevenLabsTTSService, Language from pipecat.services.elevenlabs import ElevenLabsTTSService, Language
from pipecat.services.deepgram import DeepgramSTTService from pipecat.services.deepgram import DeepgramSTTService
from pipecat.services.openai import OpenAILLMService from pipecat.services.openai import OpenAILLMService
@@ -25,13 +27,14 @@ from pipecat.transports.network.fastapi_websocket import (
FastAPIWebsocketTransport, FastAPIWebsocketTransport,
) )
load_dotenv(override=True) load_dotenv(override=True)
logger.remove(0) logger.remove(0)
logger.add(sys.stderr, level="DEBUG") logger.add(sys.stderr, level="DEBUG")
async def run_bot(websocket_client, stream_sid): async def run_bot(websocket_client, stream_id, encoding):
transport = FastAPIWebsocketTransport( transport = FastAPIWebsocketTransport(
websocket=websocket_client, websocket=websocket_client,
params=FastAPIWebsocketParams( params=FastAPIWebsocketParams(
@@ -40,7 +43,7 @@ async def run_bot(websocket_client, stream_sid):
vad_enabled=True, vad_enabled=True,
vad_analyzer=SileroVADAnalyzer(), vad_analyzer=SileroVADAnalyzer(),
vad_audio_passthrough=True, vad_audio_passthrough=True,
serializer=TwilioFrameSerializer(stream_sid), serializer=TelnyxFrameSerializer(stream_id, encoding),
), ),
) )

View File

@@ -30,9 +30,10 @@ async def websocket_endpoint(websocket: WebSocket):
await start_data.__anext__() await start_data.__anext__()
call_data = json.loads(await start_data.__anext__()) call_data = json.loads(await start_data.__anext__())
print(call_data, flush=True) print(call_data, flush=True)
stream_sid = call_data["stream_id"] stream_id = call_data["stream_id"]
encoding = call_data["start"]["media_format"]["encoding"]
print("WebSocket connection accepted") print("WebSocket connection accepted")
await run_bot(websocket, stream_sid) await run_bot(websocket, stream_id, encoding)
if __name__ == "__main__": if __name__ == "__main__":

View File

@@ -84,7 +84,6 @@ def ulaw_to_pcm(ulaw_bytes: bytes, in_sample_rate: int, out_sample_rate: int):
return out_pcm_bytes return out_pcm_bytes
def pcm_to_ulaw(pcm_bytes: bytes, in_sample_rate: int, out_sample_rate: int): def pcm_to_ulaw(pcm_bytes: bytes, in_sample_rate: int, out_sample_rate: int):
# Resample # Resample
in_pcm_bytes = resample_audio(pcm_bytes, in_sample_rate, out_sample_rate) in_pcm_bytes = resample_audio(pcm_bytes, in_sample_rate, out_sample_rate)
@@ -93,3 +92,22 @@ def pcm_to_ulaw(pcm_bytes: bytes, in_sample_rate: int, out_sample_rate: int):
ulaw_bytes = audioop.lin2ulaw(in_pcm_bytes, 2) ulaw_bytes = audioop.lin2ulaw(in_pcm_bytes, 2)
return ulaw_bytes return ulaw_bytes
def alaw_to_pcm(alaw_bytes: bytes, in_sample_rate: int, out_sample_rate: int) -> bytes:
# Convert a-law to PCM
in_pcm_bytes = audioop.alaw2lin(alaw_bytes, 2)
# Resample
out_pcm_bytes = resample_audio(in_pcm_bytes, in_sample_rate, out_sample_rate)
return out_pcm_bytes
def pcm_to_alaw(pcm_bytes: bytes, in_sample_rate: int, out_sample_rate: int):
# Resample
in_pcm_bytes = resample_audio(pcm_bytes, in_sample_rate, out_sample_rate)
# Convert PCM to μ-law
alaw_bytes = audioop.lin2alaw(in_pcm_bytes, 2)
return alaw_bytes

View File

@@ -0,0 +1,96 @@
#
# Copyright (c) 20242025, Daily
#
# SPDX-License-Identifier: BSD 2-Clause License
#
import base64
import json
from pydantic import BaseModel
from pipecat.audio.utils import pcm_to_ulaw, ulaw_to_pcm, pcm_to_alaw, alaw_to_pcm
from pipecat.frames.frames import (
AudioRawFrame,
Frame,
InputAudioRawFrame,
InputDTMFFrame,
KeypadEntry,
StartInterruptionFrame,
)
from pipecat.serializers.base_serializer import FrameSerializer, FrameSerializerType
class TelnyxFrameSerializer(FrameSerializer):
class InputParams(BaseModel):
telnyx_sample_rate: int = 8000
sample_rate: int = 16000
encoding: str = "PCMU"
def __init__(self, stream_id: str, encoding: str, params: InputParams = InputParams()):
self._stream_id = stream_id
params.encoding = encoding
self._params = params
@property
def type(self) -> FrameSerializerType:
return FrameSerializerType.TEXT
def serialize(self, frame: Frame) -> str | bytes | None:
if isinstance(frame, AudioRawFrame):
data = frame.audio
if self._params.encoding == "PCMU":
serialized_data = pcm_to_ulaw(data, frame.sample_rate, self._params.telnyx_sample_rate)
elif self._params.encoding == "PCMA":
serialized_data = pcm_to_alaw(data, frame.sample_rate, self._params.telnyx_sample_rate)
else:
raise ValueError(f"Unsupported encoding: {self._params.encoding}")
payload = base64.b64encode(serialized_data).decode("utf-8")
answer = {
"event": "media",
"media": {"payload": payload},
}
return json.dumps(answer)
if isinstance(frame, StartInterruptionFrame):
answer = {"event": "clear"}
return json.dumps(answer)
def deserialize(self, data: str | bytes) -> Frame | None:
message = json.loads(data)
if message["event"] == "start":
print(f"Start received encoding:{message['start']['media_format']['encoding']}")
self._params.encoding = message["start"]["media_format"]["encoding"]
if message["event"] == "media":
payload_base64 = message["media"]["payload"]
payload = base64.b64decode(payload_base64)
if self._params.encoding == "PCMU":
deserialized_data = ulaw_to_pcm(
payload, self._params.telnyx_sample_rate, self._params.sample_rate
)
elif self._params.encoding == "PCMA":
deserialized_data = alaw_to_pcm(
payload, self._params.telnyx_sample_rate, self._params.sample_rate
)
else:
raise ValueError(f"Unsupported encoding: {self._params.encoding}")
audio_frame = InputAudioRawFrame(
audio=deserialized_data, num_channels=1, sample_rate=self._params.sample_rate
)
return audio_frame
elif message["event"] == "dtmf":
digit = message.get("dtmf", {}).get("digit")
try:
return InputDTMFFrame(KeypadEntry(digit))
except ValueError as e:
# Handle case where string doesn't match any enum value
return None
else:
return None