Merge remote-tracking branch 'origin/main' into feat/provisional-vad-v1

# Conflicts:
#	pipecat
This commit is contained in:
Abhishek Kumar 2026-06-29 17:28:42 +05:30
commit 5ea80506db
38 changed files with 1401 additions and 153 deletions

View file

@ -11,6 +11,7 @@ from typing import Any
from loguru import logger
from api.services.pipecat.realtime.static_greeting import format_static_greeting_prompt
from pipecat.frames.frames import (
BotStartedSpeakingFrame,
BotStoppedSpeakingFrame,
@ -49,6 +50,7 @@ class DograhAzureRealtimeLLMService(AzureRealtimeLLMService):
self._handled_initial_context: bool = False
self._bot_is_speaking: bool = False
self._deferred_function_calls: list[FunctionCallFromLLM] = []
self._pending_initial_greeting_text: str | None = None
async def process_frame(self, frame: Frame, direction: FrameDirection):
if isinstance(frame, UserMuteStartedFrame):
@ -61,7 +63,11 @@ class DograhAzureRealtimeLLMService(AzureRealtimeLLMService):
return
if isinstance(frame, TTSSpeakFrame):
if not self._handled_initial_context:
await self._handle_context(self._context)
greeting_text = frame.text.strip() if frame.text else ""
if greeting_text:
await self._handle_initial_greeting(self._context, greeting_text)
else:
await self._handle_context(self._context)
else:
logger.warning(
f"{self}: TTSSpeakFrame after initial context already handled — "
@ -118,6 +124,57 @@ class DograhAzureRealtimeLLMService(AzureRealtimeLLMService):
self._context = context
await self._process_completed_function_calls(send_new_results=True)
async def _handle_initial_greeting(self, context: LLMContext, greeting_text: str):
if context is None:
logger.warning(
f"{self}: received initial greeting trigger before context was set"
)
return
self._handled_initial_context = True
self._context = context
await self._create_initial_greeting_response(greeting_text)
async def _create_initial_greeting_response(self, greeting_text: str):
if self._disconnecting:
return
if not self._api_session_ready:
self._pending_initial_greeting_text = greeting_text
self._run_llm_when_api_session_ready = True
return
self._pending_initial_greeting_text = None
await self._ensure_conversation_setup()
await self._send_manual_response_create(
instructions=format_static_greeting_prompt(greeting_text),
tool_choice="none",
)
async def _ensure_conversation_setup(self):
if not self._llm_needs_conversation_setup:
return
adapter = self.get_llm_adapter()
llm_invocation_params = adapter.get_llm_invocation_params(self._context)
for item in llm_invocation_params["messages"]:
evt = events.ConversationItemCreateEvent(item=item)
self._messages_added_manually[evt.item.id] = True
await self.send_client_event(evt)
await self._send_session_update()
self._llm_needs_conversation_setup = False
async def _handle_evt_session_updated(self, evt):
self._api_session_ready = True
if self._pending_initial_greeting_text is not None:
greeting_text = self._pending_initial_greeting_text
self._run_llm_when_api_session_ready = False
await self._create_initial_greeting_response(greeting_text)
elif self._run_llm_when_api_session_ready:
self._run_llm_when_api_session_ready = False
await self._create_response()
async def _send_user_audio(self, frame):
if self._user_is_muted:
return
@ -171,14 +228,21 @@ class DograhAzureRealtimeLLMService(AzureRealtimeLLMService):
return "\n".join(parts) if parts else None
return None
async def _send_manual_response_create(self):
async def _send_manual_response_create(
self,
*,
instructions: str | None = None,
tool_choice: str | None = None,
):
await self.push_frame(LLMFullResponseStartFrame())
await self.start_processing_metrics()
await self.start_ttfb_metrics()
await self.send_client_event(
events.ResponseCreateEvent(
response=events.ResponseProperties(
output_modalities=self._get_enabled_modalities()
output_modalities=self._get_enabled_modalities(),
instructions=instructions,
tool_choice=tool_choice,
)
)
)

View file

@ -20,11 +20,13 @@ Layers Dograh engine integration quirks onto upstream-pristine
from typing import Any
from google.genai.types import Content, Part
from loguru import logger
from api.services.pipecat.gemini_json_schema_adapter import (
DograhGeminiJSONSchemaAdapter,
)
from api.services.pipecat.realtime.static_greeting import format_static_greeting_prompt
from pipecat.frames.frames import (
BotStoppedSpeakingFrame,
Frame,
@ -63,6 +65,9 @@ class DograhGeminiLiveLLMService(GeminiLiveLLMService):
# Function calls emitted by Gemini mid-bot-turn are deferred here and
# invoked when the turn ends, so they don't race the turn's audio.
self._pending_function_calls: list[FunctionCallFromLLM] = []
# Text greeting captured from the first TTSSpeakFrame while the Gemini
# session is still connecting.
self._pending_initial_greeting_text: str | None = None
# ------------------------------------------------------------------
# Hooks from upstream GeminiLiveLLMService
@ -142,10 +147,15 @@ class DograhGeminiLiveLLMService(GeminiLiveLLMService):
if isinstance(frame, TTSSpeakFrame):
# Greeting trigger: the engine queues a TTSSpeakFrame to start the
# bot's first turn after node setup. Gemini Live renders its own
# audio, so we don't pass the frame through — we re-enter
# _handle_context to kick off the initial response.
# audio, so we don't pass the frame through. For configured static
# text greetings, ask Gemini to say the exact greeting; otherwise
# re-enter _handle_context to kick off the normal initial response.
if not self._handled_initial_context:
await self._handle_context(self._context)
greeting_text = frame.text.strip() if frame.text else ""
if greeting_text:
await self._handle_initial_greeting(self._context, greeting_text)
else:
await self._handle_context(self._context)
else:
logger.warning(
f"{self}: TTSSpeakFrame after initial context already "
@ -183,6 +193,49 @@ class DograhGeminiLiveLLMService(GeminiLiveLLMService):
self._context = context
await self._process_completed_function_calls(send_new_results=True)
async def _handle_initial_greeting(self, context: LLMContext, greeting_text: str):
"""Trigger the first Gemini turn with an exact static text greeting."""
if context is None:
logger.warning(
f"{self}: received initial greeting trigger before context was set"
)
return
self._handled_initial_context = True
self._context = context
await self._create_initial_greeting_response(greeting_text)
async def _create_initial_greeting_response(self, greeting_text: str):
"""Ask Gemini Live to speak the configured greeting exactly once."""
if self._disconnecting:
return
if not self._session:
self._pending_initial_greeting_text = greeting_text
self._run_llm_when_session_ready = True
return
self._pending_initial_greeting_text = None
prompt = format_static_greeting_prompt(greeting_text)
turn = Content(role="user", parts=[Part(text=prompt)])
logger.debug("Creating Gemini Live initial response from static greeting")
await self.start_ttfb_metrics()
try:
await self._session.send_client_content(
turns=[turn],
turn_complete=True,
)
# Gemini 3.x also needs a realtime-input nudge to begin inference.
if self._is_gemini_3:
await self._session.send_realtime_input(text=" ")
except Exception as e:
await self._handle_send_error(e)
self._ready_for_realtime_input = True
# ------------------------------------------------------------------
# Session lifecycle: drop upstream's automatic reconnect-seed and
# initial-context-seed paths. The TTSSpeakFrame trigger and the
@ -201,7 +254,12 @@ class DograhGeminiLiveLLMService(GeminiLiveLLMService):
# Context arrived before session was ready — fulfil the queued
# initial response now.
self._run_llm_when_session_ready = False
await self._create_initial_response()
if self._pending_initial_greeting_text is not None:
await self._create_initial_greeting_response(
self._pending_initial_greeting_text
)
else:
await self._create_initial_response()
await self._drain_pending_tool_results()
# Otherwise: no automatic seed. Reconnect after a session-resumption
# update relies on the server-side restored state; reconnects without

View file

@ -22,6 +22,7 @@ from typing import Any
from loguru import logger
from api.services.pipecat.realtime.static_greeting import format_static_greeting_prompt
from pipecat.frames.frames import (
BotStartedSpeakingFrame,
BotStoppedSpeakingFrame,
@ -50,6 +51,7 @@ class DograhGrokRealtimeLLMService(GrokRealtimeLLMService):
self._handled_initial_context: bool = False
self._bot_is_speaking: bool = False
self._deferred_function_calls: list[FunctionCallFromLLM] = []
self._pending_initial_greeting_text: str | None = None
async def process_frame(self, frame: Frame, direction: FrameDirection):
if isinstance(frame, UserMuteStartedFrame):
@ -62,7 +64,11 @@ class DograhGrokRealtimeLLMService(GrokRealtimeLLMService):
return
if isinstance(frame, TTSSpeakFrame):
if not self._handled_initial_context:
await self._handle_context(self._context)
greeting_text = frame.text.strip() if frame.text else ""
if greeting_text:
await self._handle_initial_greeting(self._context, greeting_text)
else:
await self._handle_context(self._context)
else:
logger.warning(
f"{self}: TTSSpeakFrame after initial context already "
@ -120,6 +126,67 @@ class DograhGrokRealtimeLLMService(GrokRealtimeLLMService):
self._context = context
await self._process_completed_function_calls(send_new_results=True)
async def _handle_initial_greeting(self, context: LLMContext, greeting_text: str):
if context is None:
logger.warning(
f"{self}: received initial greeting trigger before context was set"
)
return
self._handled_initial_context = True
self._context = context
await self._create_initial_greeting_response(greeting_text)
async def _create_initial_greeting_response(self, greeting_text: str):
if self._disconnecting:
return
if not self._api_session_ready:
self._pending_initial_greeting_text = greeting_text
self._run_llm_when_api_session_ready = True
return
self._pending_initial_greeting_text = None
await self._ensure_conversation_setup()
item = events.ConversationItem(
type="message",
role="user",
content=[
events.ItemContent(
type="input_text",
text=format_static_greeting_prompt(greeting_text),
)
],
)
evt = events.ConversationItemCreateEvent(item=item)
self._messages_added_manually[evt.item.id] = True
await self.send_client_event(evt)
await self._send_manual_response_create()
async def _ensure_conversation_setup(self):
if not self._llm_needs_conversation_setup:
return
adapter = self.get_llm_adapter()
llm_invocation_params = adapter.get_llm_invocation_params(self._context)
for item in llm_invocation_params["messages"]:
evt = events.ConversationItemCreateEvent(item=item)
self._messages_added_manually[evt.item.id] = True
await self.send_client_event(evt)
await self._send_session_update()
self._llm_needs_conversation_setup = False
async def _handle_evt_session_updated(self, evt):
self._api_session_ready = True
if self._pending_initial_greeting_text is not None:
greeting_text = self._pending_initial_greeting_text
self._run_llm_when_api_session_ready = False
await self._create_initial_greeting_response(greeting_text)
elif self._run_llm_when_api_session_ready:
self._run_llm_when_api_session_ready = False
await self._create_response()
async def _send_user_audio(self, frame):
if self._user_is_muted:
return

View file

@ -22,6 +22,7 @@ from typing import Any
from loguru import logger
from api.services.pipecat.realtime.static_greeting import format_static_greeting_prompt
from pipecat.frames.frames import (
BotStartedSpeakingFrame,
BotStoppedSpeakingFrame,
@ -56,6 +57,7 @@ class DograhOpenAIRealtimeLLMService(OpenAIRealtimeLLMService):
# has finished speaking, matching Dograh's Gemini Live behavior.
self._bot_is_speaking: bool = False
self._deferred_function_calls: list[FunctionCallFromLLM] = []
self._pending_initial_greeting_text: str | None = None
# ------------------------------------------------------------------
# Frame handling: mute, TTSSpeakFrame as greeting trigger
@ -73,11 +75,16 @@ class DograhOpenAIRealtimeLLMService(OpenAIRealtimeLLMService):
if isinstance(frame, TTSSpeakFrame):
# Greeting trigger: the engine queues a TTSSpeakFrame after node
# setup. OpenAI Realtime renders its own audio, so we don't pass
# the frame to TTS. Route through _handle_context so the initial
# response and later tool-result turns share the same context
# lifecycle even when Dograh has already pre-populated self._context.
# the frame to TTS. For configured static text greetings, ask the
# model to say the exact greeting; otherwise route through
# _handle_context so the initial response and later tool-result
# turns share the same context lifecycle.
if not self._handled_initial_context:
await self._handle_context(self._context)
greeting_text = frame.text.strip() if frame.text else ""
if greeting_text:
await self._handle_initial_greeting(self._context, greeting_text)
else:
await self._handle_context(self._context)
else:
logger.warning(
f"{self}: TTSSpeakFrame after initial context already "
@ -137,6 +144,57 @@ class DograhOpenAIRealtimeLLMService(OpenAIRealtimeLLMService):
self._context = context
await self._process_completed_function_calls(send_new_results=True)
async def _handle_initial_greeting(self, context: LLMContext, greeting_text: str):
if context is None:
logger.warning(
f"{self}: received initial greeting trigger before context was set"
)
return
self._handled_initial_context = True
self._context = context
await self._create_initial_greeting_response(greeting_text)
async def _create_initial_greeting_response(self, greeting_text: str):
if self._disconnecting:
return
if not self._api_session_ready:
self._pending_initial_greeting_text = greeting_text
self._run_llm_when_api_session_ready = True
return
self._pending_initial_greeting_text = None
await self._ensure_conversation_setup()
await self._send_manual_response_create(
instructions=format_static_greeting_prompt(greeting_text),
tool_choice="none",
)
async def _ensure_conversation_setup(self):
if not self._llm_needs_conversation_setup:
return
adapter = self.get_llm_adapter()
llm_invocation_params = adapter.get_llm_invocation_params(self._context)
for item in llm_invocation_params["messages"]:
evt = events.ConversationItemCreateEvent(item=item)
self._messages_added_manually[evt.item.id] = True
await self.send_client_event(evt)
await self._send_session_update()
self._llm_needs_conversation_setup = False
async def _handle_evt_session_updated(self, evt):
self._api_session_ready = True
if self._pending_initial_greeting_text is not None:
greeting_text = self._pending_initial_greeting_text
self._run_llm_when_api_session_ready = False
await self._create_initial_greeting_response(greeting_text)
elif self._run_llm_when_api_session_ready:
self._run_llm_when_api_session_ready = False
await self._create_response()
async def _send_user_audio(self, frame):
if self._user_is_muted:
return
@ -190,7 +248,12 @@ class DograhOpenAIRealtimeLLMService(OpenAIRealtimeLLMService):
return "\n".join(parts) if parts else None
return None
async def _send_manual_response_create(self):
async def _send_manual_response_create(
self,
*,
instructions: str | None = None,
tool_choice: str | None = None,
):
"""Trigger inference after manually appending conversation items."""
await self.push_frame(LLMFullResponseStartFrame())
await self.start_processing_metrics()
@ -198,7 +261,9 @@ class DograhOpenAIRealtimeLLMService(OpenAIRealtimeLLMService):
await self.send_client_event(
events.ResponseCreateEvent(
response=events.ResponseProperties(
output_modalities=self._get_enabled_modalities()
output_modalities=self._get_enabled_modalities(),
instructions=instructions,
tool_choice=tool_choice,
)
)
)

View file

@ -0,0 +1,8 @@
def format_static_greeting_prompt(greeting_text: str) -> str:
return (
"The phone call has just connected. Greet the caller now: "
"say the following opening line out loud, exactly as written, "
"in a natural spoken voice, and then stop and wait for the "
"caller to respond. Do not add anything before or after it.\n\n"
f'"{greeting_text}"'
)

View file

@ -72,7 +72,7 @@ class RealtimeFeedbackObserver(BaseObserver):
- TTFB metrics (LLM generation time only)
Logs buffer persistence (only final data for post-call analysis):
- Complete user transcripts per turn (via on_user_turn_stopped)
- Complete user transcripts per turn (via on_user_turn_message_added)
- Complete assistant transcripts per turn (via on_assistant_turn_stopped)
- Function calls and TTFB metrics
@ -300,13 +300,13 @@ def register_turn_log_handlers(
):
"""Register event handlers on aggregators to persist final turn transcripts.
Hooks into on_user_turn_stopped and on_assistant_turn_stopped to store
Hooks into on_user_turn_message_added and on_assistant_turn_stopped to store
complete turn text in the logs buffer. Works for both WebRTC and telephony
calls independent of WebSocket availability.
"""
@user_aggregator.event_handler("on_user_turn_stopped")
async def on_user_turn_stopped(aggregator, strategy, message):
@user_aggregator.event_handler("on_user_turn_message_added")
async def on_user_turn_message_added(aggregator, message):
logs_buffer.increment_turn()
try:
await logs_buffer.append(

View file

@ -113,49 +113,53 @@ def _resolve_user_turn_stop_timeout(
def _create_realtime_user_turn_config(provider: str):
"""Return user turn strategies and optional local VAD for realtime providers."""
def external_provider_turn_config():
return (
UserTurnStrategies(
start=[ExternalUserTurnStartStrategy()],
stop=[ExternalUserTurnStopStrategy(wait_for_transcript=False)],
),
None,
)
def local_vad_turn_config(*, enable_interruptions: bool):
return (
UserTurnStrategies(
start=[
VADUserTurnStartStrategy(enable_interruptions=enable_interruptions)
],
stop=[SpeechTimeoutUserTurnStopStrategy(wait_for_transcript=False)],
),
SileroVADAnalyzer(params=VADParams(stop_secs=0.2)),
)
if provider in {
ServiceProviders.GOOGLE_REALTIME.value,
ServiceProviders.GOOGLE_VERTEX_REALTIME.value,
}:
# Let Gemini Live own barge-in via its server-side VAD, but keep local
# Silero VAD for early user-turn start and speaking-state tracking.
return (
UserTurnStrategies(
start=[VADUserTurnStartStrategy(enable_interruptions=False)],
stop=[SpeechTimeoutUserTurnStopStrategy()],
),
SileroVADAnalyzer(params=VADParams(stop_secs=0.2)),
)
return local_vad_turn_config(enable_interruptions=False)
if provider == ServiceProviders.OPENAI_REALTIME.value:
# OpenAI Realtime already emits speaking-state frames and interruption
# events from the provider, so the aggregator should follow those
# external signals rather than run its own local VAD.
return (
UserTurnStrategies(
start=[ExternalUserTurnStartStrategy()],
stop=[ExternalUserTurnStopStrategy()],
),
None,
)
if provider in {
ServiceProviders.OPENAI_REALTIME.value,
ServiceProviders.AZURE_REALTIME.value,
}:
# OpenAI-compatible Realtime services already emit speaking-state frames
# and interruption events from the provider, so the aggregator should
# follow those external signals rather than run its own local VAD.
return external_provider_turn_config()
if provider == ServiceProviders.GROK_REALTIME.value:
# Grok Voice Agent emits server-side speech-start/stop and
# interruption signals, so local VAD should stay out of the way.
return (
UserTurnStrategies(
start=[ExternalUserTurnStartStrategy()],
stop=[ExternalUserTurnStopStrategy()],
),
None,
)
return external_provider_turn_config()
if provider == ServiceProviders.ULTRAVOX_REALTIME.value:
# Ultravox does not emit user-turn frames, so local VAD supplies
# lifecycle signals for Dograh observers/controllers.
return local_vad_turn_config(enable_interruptions=True)
return (
UserTurnStrategies(
start=[VADUserTurnStartStrategy()],
stop=[SpeechTimeoutUserTurnStopStrategy()],
),
SileroVADAnalyzer(params=VADParams(stop_secs=0.2)),
)
return local_vad_turn_config(enable_interruptions=True)
async def run_pipeline_telephony(
@ -773,7 +777,10 @@ async def _run_pipeline_impl(
vad_analyzer=user_vad_analyzer,
)
context_aggregator = LLMContextAggregatorPair(
context, assistant_params=assistant_params, user_params=user_params
context,
assistant_params=assistant_params,
user_params=user_params,
realtime_service_mode=is_realtime,
)
# Create usage metrics aggregator with engine's callback

View file

@ -21,6 +21,7 @@ def _config_loader(value: Dict[str, Any]) -> Dict[str, Any]:
"private_key": value.get("private_key"),
"api_key": value.get("api_key"),
"api_secret": value.get("api_secret"),
"signature_secret": value.get("signature_secret"),
"from_numbers": value.get("from_numbers", []),
}
@ -49,6 +50,13 @@ _UI_METADATA = ProviderUIMetadata(
type="password",
sensitive=True,
),
ProviderUIField(
name="signature_secret",
label="Signature Secret",
type="password",
sensitive=True,
description="Vonage signature secret for signed webhook verification",
),
ProviderUIField(
name="from_numbers",
label="Phone Numbers",

View file

@ -1,6 +1,6 @@
"""Vonage telephony configuration schemas."""
from typing import List, Literal
from typing import List, Literal, Optional
from pydantic import BaseModel, Field
@ -13,6 +13,10 @@ class VonageConfigurationRequest(BaseModel):
api_secret: str = Field(..., description="Vonage API Secret")
application_id: str = Field(..., description="Vonage Application ID")
private_key: str = Field(..., description="Private key for JWT generation")
signature_secret: Optional[str] = Field(
None,
description="Vonage signature secret used to verify signed webhooks",
)
from_numbers: List[str] = Field(
default_factory=list,
description="List of Vonage phone numbers (without + prefix)",
@ -27,4 +31,5 @@ class VonageConfigurationResponse(BaseModel):
api_key: str # Masked
api_secret: str # Masked
private_key: str # Masked
signature_secret: Optional[str] = None # Masked
from_numbers: List[str]

View file

@ -2,6 +2,7 @@
Vonage (Nexmo) implementation of the TelephonyProvider interface.
"""
import hashlib
import json
import random
import time
@ -44,12 +45,14 @@ class VonageProvider(TelephonyProvider):
- api_secret: Vonage API Secret
- application_id: Vonage Application ID
- private_key: Private key for JWT generation
- signature_secret: Signature secret for signed webhooks
- from_numbers: List of phone numbers to use
"""
self.api_key = config.get("api_key")
self.api_secret = config.get("api_secret")
self.application_id = config.get("application_id")
self.private_key = config.get("private_key")
self.signature_secret = config.get("signature_secret")
self.from_numbers = config.get("from_numbers", [])
# Handle both single number (string) and multiple numbers (list)
@ -186,17 +189,18 @@ class VonageProvider(TelephonyProvider):
Verify Vonage webhook signature for security.
Vonage uses JWT for webhook signatures.
"""
if not self.api_secret:
logger.error("No API secret available for webhook signature verification")
if not self.signature_secret:
logger.error(
"No signature secret available for Vonage webhook verification"
)
return False
try:
# Vonage sends JWT in Authorization header. Verify the JWT signature
decoded = jwt.decode(
jwt.decode(
signature,
self.api_secret,
self.signature_secret,
algorithms=["HS256"],
options={"verify_signature": True},
options={"verify_signature": True, "verify_aud": False},
)
return True
except jwt.InvalidTokenError:
@ -295,9 +299,13 @@ class VonageProvider(TelephonyProvider):
"ringing": TelephonyCallStatus.RINGING,
"answered": TelephonyCallStatus.ANSWERED,
"complete": TelephonyCallStatus.COMPLETED,
"completed": TelephonyCallStatus.COMPLETED,
"disconnected": TelephonyCallStatus.COMPLETED,
"failed": TelephonyCallStatus.FAILED,
"busy": TelephonyCallStatus.BUSY,
"timeout": TelephonyCallStatus.NO_ANSWER,
"unanswered": TelephonyCallStatus.NO_ANSWER,
"cancelled": TelephonyCallStatus.NO_ANSWER,
"rejected": TelephonyCallStatus.BUSY,
}
@ -349,6 +357,8 @@ class VonageProvider(TelephonyProvider):
if workflow_run.gathered_context
else None
)
if not call_uuid and workflow_run.gathered_context:
call_uuid = workflow_run.gathered_context.get("call_id")
if not call_uuid:
logger.error(
@ -400,26 +410,126 @@ class VonageProvider(TelephonyProvider):
"""
Determine if this provider can handle the incoming webhook.
"""
return False
claims = cls._decode_unverified_signed_claims(headers)
if claims.get("api_key") or claims.get("application_id"):
return True
return bool(
webhook_data.get("uuid")
and webhook_data.get("conversation_uuid")
and webhook_data.get("from")
and webhook_data.get("to")
)
@staticmethod
def parse_inbound_webhook(webhook_data: Dict[str, Any]) -> NormalizedInboundData:
def parse_inbound_webhook(
webhook_data: Dict[str, Any], headers: Optional[Dict[str, str]] = None
) -> NormalizedInboundData:
"""
Parse Vonage-specific inbound webhook data into normalized format.
"""
claims = VonageProvider._decode_unverified_signed_claims(headers or {})
direction = webhook_data.get("direction") or "inbound"
status = webhook_data.get("status") or "started"
return NormalizedInboundData(
provider=VonageProvider.PROVIDER_NAME,
call_id=webhook_data.get("uuid", ""),
from_number=webhook_data.get("from", ""),
to_number=webhook_data.get("to", ""),
direction=webhook_data.get("direction", ""),
call_status=webhook_data.get("status", ""),
account_id=webhook_data.get("account_id"),
direction=direction,
call_status=status,
account_id=claims.get("api_key") or webhook_data.get("account_id"),
from_country=None,
to_country=None,
raw_data=webhook_data,
)
@staticmethod
def _header(headers: Dict[str, str], name: str) -> Optional[str]:
for key, value in headers.items():
if key.lower() == name.lower():
return value
return None
@classmethod
def _bearer_token(cls, headers: Dict[str, str]) -> Optional[str]:
auth_header = cls._header(headers, "authorization")
if not auth_header:
return None
parts = auth_header.split(None, 1)
if len(parts) != 2 or parts[0].lower() != "bearer":
return None
return parts[1].strip()
@classmethod
def _decode_unverified_signed_claims(
cls, headers: Dict[str, str]
) -> Dict[str, Any]:
token = cls._bearer_token(headers)
if not token:
return {}
try:
claims = jwt.decode(
token,
options={
"verify_signature": False,
"verify_aud": False,
"verify_exp": False,
},
)
except jwt.InvalidTokenError:
return {}
return claims if isinstance(claims, dict) else {}
def _verify_signed_claims(
self, headers: Dict[str, str], body: str = ""
) -> Optional[Dict[str, Any]]:
token = self._bearer_token(headers)
if not token:
logger.warning("Missing Vonage Authorization bearer token")
return None
if not self.signature_secret:
logger.error("Missing Vonage signature_secret for signed webhook")
return None
try:
claims = jwt.decode(
token,
self.signature_secret,
algorithms=["HS256"],
options={"verify_signature": True, "verify_aud": False},
)
except jwt.InvalidTokenError as exc:
logger.warning(f"Invalid Vonage signed webhook JWT: {exc}")
return None
if claims.get("iss") != "Vonage":
logger.warning("Vonage signed webhook JWT has unexpected issuer")
return None
if self.api_key and claims.get("api_key") != self.api_key:
logger.warning("Vonage signed webhook api_key does not match config")
return None
claim_application_id = claims.get("application_id")
if (
self.application_id
and claim_application_id
and claim_application_id != self.application_id
):
logger.warning("Vonage signed webhook application_id does not match config")
return None
payload_hash = claims.get("payload_hash")
if payload_hash:
actual_hash = hashlib.sha256(body.encode("utf-8")).hexdigest()
if actual_hash != payload_hash:
logger.warning("Vonage signed webhook payload hash mismatch")
return None
return claims
@staticmethod
def validate_account_id(config_data: dict, webhook_account_id: str) -> bool:
"""Validate Vonage account_id from webhook matches configuration"""
@ -437,9 +547,10 @@ class VonageProvider(TelephonyProvider):
body: str = "",
) -> bool:
"""
Vonage inbound signature verification - minimalist implementation.
Verify Vonage signed webhook JWT and optional payload hash.
"""
return True
claims = self._verify_signed_claims(headers, body)
return claims is not None
async def configure_inbound(
self, address: str, webhook_url: Optional[str]
@ -486,6 +597,15 @@ class VonageProvider(TelephonyProvider):
),
)
if not self.signature_secret:
return ProviderSyncResult(
ok=False,
message=(
"Vonage signature_secret is required because inbound calls "
"use signed webhook verification"
),
)
app_endpoint = f"{self.base_url}/v2/applications/{self.application_id}"
auth = aiohttp.BasicAuth(self.api_key, self.api_secret)
@ -510,12 +630,18 @@ class VonageProvider(TelephonyProvider):
capabilities = app_data.get("capabilities") or {}
voice = capabilities.get("voice") or {}
webhooks = voice.get("webhooks") or {}
backend_endpoint, _ = await get_backend_endpoints()
webhooks["answer_url"] = {
"address": webhook_url,
"http_method": "POST",
}
webhooks["event_url"] = {
"address": f"{backend_endpoint}/api/v1/telephony/vonage/events",
"http_method": "POST",
}
voice["webhooks"] = webhooks
voice["signed_callbacks"] = True
capabilities["voice"] = voice
update_body = {
@ -561,13 +687,24 @@ class VonageProvider(TelephonyProvider):
"""
Generate NCCO response for inbound Vonage webhook.
"""
# Minimalist NCCO response for interface compliance
ncco_response = [
{
"action": "talk",
"text": "Vonage inbound calls are not currently supported.",
},
{"action": "hangup"},
"action": "connect",
"eventUrl": [
f"{backend_endpoint}/api/v1/telephony/vonage/events/{workflow_run_id}"
],
"endpoint": [
{
"type": "websocket",
"uri": websocket_url,
"content-type": "audio/l16;rate=16000",
"headers": {
"workflow_run_id": str(workflow_run_id),
"call_uuid": normalized_data.call_id,
},
}
],
}
]
return Response(

View file

@ -7,16 +7,12 @@ provider registry — see ProviderSpec.router.
import json
from typing import Optional
from fastapi import APIRouter, Request
from fastapi import APIRouter, HTTPException, Request
from loguru import logger
from pipecat.utils.run_context import set_current_run_id
from api.db import db_client
from api.services.telephony.factory import get_telephony_provider_for_run
from api.services.telephony.status_processor import (
StatusCallbackRequest,
_process_status_update,
)
router = APIRouter()
@ -45,28 +41,33 @@ async def handle_ncco_webhook(
return json.loads(response_content)
@router.post("/vonage/events/{workflow_run_id}")
async def handle_vonage_events(
request: Request,
workflow_run_id: int,
):
"""Handle Vonage-specific event webhooks.
async def _read_json_body(request: Request) -> tuple[dict, str]:
body_bytes = await request.body()
try:
raw_body = body_bytes.decode("utf-8")
except UnicodeDecodeError as exc:
raise HTTPException(
status_code=400, detail="Webhook body is not valid UTF-8"
) from exc
try:
return json.loads(raw_body or "{}"), raw_body
except json.JSONDecodeError as exc:
raise HTTPException(status_code=400, detail="Webhook body is not JSON") from exc
Vonage sends all call events to a single endpoint.
Events include: started, ringing, answered, complete, failed, etc.
"""
async def _handle_vonage_event_request(request: Request, workflow_run_id: int):
set_current_run_id(workflow_run_id)
# Parse the event data
event_data = await request.json()
logger.info(f"[run {workflow_run_id}] Received Vonage event: {event_data}")
event_data, raw_body = await _read_json_body(request)
logger.info(
f"[run {workflow_run_id}] Received Vonage event "
f"uuid={event_data.get('uuid')} status={event_data.get('status')}"
)
# Get workflow run for processing
workflow_run = await db_client.get_workflow_run_by_id(workflow_run_id)
if not workflow_run:
logger.error(f"[run {workflow_run_id}] Workflow run not found")
return {"status": "error", "message": "Workflow run not found"}
# Get workflow and provider
workflow = await db_client.get_workflow_by_id(workflow_run.workflow_id)
if not workflow:
logger.error(f"[run {workflow_run_id}] Workflow not found")
@ -75,11 +76,18 @@ async def handle_vonage_events(
provider = await get_telephony_provider_for_run(
workflow_run, workflow.organization_id
)
signature_valid = await provider.verify_inbound_signature(
str(request.url), event_data, dict(request.headers), raw_body
)
if not signature_valid:
raise HTTPException(status_code=401, detail="Invalid webhook signature")
from api.services.telephony.status_processor import (
StatusCallbackRequest,
_process_status_update,
)
# Parse the event data into generic format
parsed_data = provider.parse_status_callback(event_data)
# Create StatusCallbackRequest from parsed data
status_update = StatusCallbackRequest(
call_id=parsed_data["call_id"],
status=parsed_data["status"],
@ -90,8 +98,35 @@ async def handle_vonage_events(
extra=parsed_data.get("extra", {}),
)
# Process the status update
await _process_status_update(workflow_run_id, status_update)
# Return 204 No Content as expected by Vonage
return {"status": "ok"}
@router.post("/vonage/events/{workflow_run_id}")
async def handle_vonage_events(
request: Request,
workflow_run_id: int,
):
"""Handle Vonage-specific event webhooks.
Vonage sends all call events to a single endpoint.
Events include: started, ringing, answered, complete, failed, etc.
"""
return await _handle_vonage_event_request(request, workflow_run_id)
@router.post("/vonage/events")
async def handle_vonage_events_without_run(request: Request):
"""Handle application-level events by resolving the run from call UUID."""
event_data, _ = await _read_json_body(request)
call_id = event_data.get("uuid")
if call_id:
workflow_run = await db_client.get_workflow_run_by_call_id(call_id)
if workflow_run:
return await _handle_vonage_event_request(request, workflow_run.id)
logger.info(
"Received unmatched Vonage application event "
f"uuid={event_data.get('uuid')} status={event_data.get('status')}"
)
return {"status": "ok"}

View file

@ -250,7 +250,6 @@ class _ToolDocumentRefsMixin(BaseModel):
"description": (
"Text spoken via TTS at the start of the call. Supports "
"{{template_variables}}. Leave empty to skip the greeting. "
"Not supported with realtime (speech-to-speech) models."
),
"display_options": DisplayOptions(show={"greeting_type": ["text"]}),
"placeholder": "Hi {{first_name}}, this is Sarah from Acme.",