Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .serena/project.local.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
# This file allows you to locally override settings in project.yml for development purposes.
#
# Use the same keys as in project.yml here. Any setting you specify will override the corresponding
# setting in project.yml, allowing you to customise the configuration for your local development environment
# without affecting the project configuration in project.yml (which is intended to be versioned).
1 change: 1 addition & 0 deletions app/api/v1/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@
api_router.include_router(prompt_partials.router)
api_router.include_router(prompt_optimization.router)
api_router.include_router(telephony.router)
api_router.include_router(telephony.plivo_webhook_router)
api_router.include_router(vobiz_telephony.router)
api_router.include_router(call_imports.router)
api_router.include_router(call_import_schemas.router)
Expand Down
2 changes: 1 addition & 1 deletion app/api/v1/media.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,5 +6,5 @@

media_router = APIRouter()
media_router.include_router(vobiz_telephony.webhook_router)
media_router.include_router(vobiz_telephony.ws_router)
media_router.include_router(vobiz_telephony.carrier_ws_router)
media_router.include_router(voice_agent.ws_router)
38 changes: 27 additions & 11 deletions app/api/v1/routes/call_imports.py
Original file line number Diff line number Diff line change
Expand Up @@ -1360,11 +1360,14 @@ def _validate_direct_url_import_ready(
)


def _validate_exotel_import_ready(
def _validate_credentialed_import_ready(
provider: str,
parameters: List[CallImportSchemaParameter],
parameter_mapping: Dict[str, Any],
) -> None:
"""Ensure Exotel credentialed import has a mapped recording_url column."""
"""Ensure credentialed import has a mapped recording_url column."""
provider_key = (provider or "").lower()
provider_label = provider_key.capitalize()
rec_url_param = next(
(
p
Expand All @@ -1377,7 +1380,7 @@ def _validate_exotel_import_ready(
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=(
"Exotel import requires a schema parameter of type "
f"{provider_label} import requires a schema parameter of type "
"'recording_url'."
),
)
Expand All @@ -1386,12 +1389,20 @@ def _validate_exotel_import_ready(
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=(
"Exotel import requires the 'recording_url' parameter to "
f"{provider_label} import requires the 'recording_url' parameter to "
"be mapped to a source column."
),
)


def _validate_exotel_import_ready(
parameters: List[CallImportSchemaParameter],
parameter_mapping: Dict[str, Any],
) -> None:
"""Ensure Exotel credentialed import has a mapped recording_url column."""
_validate_credentialed_import_ready("exotel", parameters, parameter_mapping)


def _is_manual_audio_call_import(call_import: CallImport) -> bool:
"""True for batches created via manual audio upload (recordings already in S3)."""
return (call_import.source_format or "").lower() == "audio"
Expand Down Expand Up @@ -2150,9 +2161,11 @@ async def start_call_import(
payload.provider or "",
)
_validate_telephony_credentials_live(db, organization_id, integration)
if (integration.provider or "").lower() == "exotel":
_validate_exotel_import_ready(
parameters, dict(call_import.parameter_mapping or {})
if (integration.provider or "").lower() in {"exotel", "plivo"}:
_validate_credentialed_import_ready(
integration.provider,
parameters,
dict(call_import.parameter_mapping or {}),
)
else:
_validate_direct_url_import_ready(
Expand Down Expand Up @@ -2345,8 +2358,10 @@ async def upload_call_import_csv(
integration = _resolve_telephony_integration(
db, organization_id, telephony_integration_id, provider or ""
)
if (integration.provider or "").lower() == "exotel":
_validate_exotel_import_ready(parameters, cleaned_mapping)
if (integration.provider or "").lower() in {"exotel", "plivo"}:
_validate_credentialed_import_ready(
integration.provider, parameters, cleaned_mapping
)
else:
_validate_direct_url_import_ready(parameters, cleaned_mapping)
integration = None
Expand Down Expand Up @@ -3527,14 +3542,15 @@ async def retry_failed_call_import_rows(
_validate_telephony_credentials_live(db, organization_id, integration)
call_import.provider = integration.provider
call_import.telephony_integration_id = integration.id
if (integration.provider or "").lower() == "exotel":
if (integration.provider or "").lower() in {"exotel", "plivo"}:
schema = _resolve_schema(
db,
organization_id,
call_import.workspace_id,
call_import.schema_id,
)
_validate_exotel_import_ready(
_validate_credentialed_import_ready(
integration.provider,
list(schema.parameters),
dict(call_import.parameter_mapping or {}),
)
Expand Down
40 changes: 34 additions & 6 deletions app/api/v1/routes/telephony.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
from app.services.telephony.platform_outbound_pool import outbound_pool_api_payload

router = APIRouter(prefix="/telephony", tags=["Telephony"])
plivo_webhook_router = APIRouter(prefix="/telephony/plivo", tags=["Plivo Telephony Webhooks"])


class TelephonyAvailableNumberResponse(BaseModel):
Expand Down Expand Up @@ -557,25 +558,52 @@ async def _read_webhook_params(request: Request) -> Dict[str, Any]:
return params


@router.post("/webhooks/answer")
async def telephony_answer_webhook(request: Request, db: Session = Depends(get_db)):
async def _handle_plivo_answer_webhook(request: Request, db: Session) -> Response:
params = await _read_webhook_params(request)
verify_plivo_webhook(request, params, "answer", db)
xml = telephony_service.handle_answer_webhook(params, db)
return Response(content=xml, media_type="application/xml")


@router.post("/webhooks/events")
async def telephony_events_webhook(request: Request, db: Session = Depends(get_db)):
async def _handle_plivo_events_webhook(request: Request, db: Session) -> Dict[str, str]:
params = await _read_webhook_params(request)
verify_plivo_webhook(request, params, "events", db)
telephony_service.handle_event_webhook(params, db)
return {"status": "ok"}


@router.post("/webhooks/masking")
async def telephony_masking_webhook(request: Request, db: Session = Depends(get_db)):
async def _handle_plivo_masking_webhook(request: Request, db: Session) -> Response:
params = await _read_webhook_params(request)
verify_plivo_webhook(request, params, "masking", db)
xml = telephony_service.handle_masking_webhook(params, db)
return Response(content=xml, media_type="application/xml")


@plivo_webhook_router.post("/webhooks/answer")
async def plivo_answer_webhook(request: Request, db: Session = Depends(get_db)):
return await _handle_plivo_answer_webhook(request, db)


@plivo_webhook_router.post("/webhooks/events")
async def plivo_events_webhook(request: Request, db: Session = Depends(get_db)):
return await _handle_plivo_events_webhook(request, db)


@plivo_webhook_router.post("/webhooks/masking")
async def plivo_masking_webhook(request: Request, db: Session = Depends(get_db)):
return await _handle_plivo_masking_webhook(request, db)


@router.post("/webhooks/answer")
async def telephony_answer_webhook(request: Request, db: Session = Depends(get_db)):
return await _handle_plivo_answer_webhook(request, db)


@router.post("/webhooks/events")
async def telephony_events_webhook(request: Request, db: Session = Depends(get_db)):
return await _handle_plivo_events_webhook(request, db)


@router.post("/webhooks/masking")
async def telephony_masking_webhook(request: Request, db: Session = Depends(get_db)):
return await _handle_plivo_masking_webhook(request, db)
52 changes: 24 additions & 28 deletions app/api/v1/routes/vobiz_telephony.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
from app.services.telephony.phone_routing import resolve_inbound_agent_for_number
from app.services.telephony.plivo_client import normalize_e164
from app.services.telephony.vobiz_agent_context import (
build_vobiz_ws_url,
build_carrier_ws_url,
extract_webhook_params,
resolve_vobiz_agent_context,
vobiz_webhook_base_url,
Expand Down Expand Up @@ -49,15 +49,19 @@
from app.services.telephony.vobiz_xml import reject_call, speak_and_hangup, stream_to_agent
from app.services.voice_agent.bot_fast_api import run_bot
from app.services.voice_agent.voice_bundle import run_voice_bundle_fastapi
from app.services.telephony.carrier_media_serializer import (
build_carrier_frame_serializer,
telephony_integration_id_from_call_row,
)
from efficientai.runner.utils import parse_telephony_websocket
from efficientai.serializers.vobiz import VobizFrameSerializer

# Exposed at module scope so tests can patch `.delay` without importing Celery tasks.
initiate_vobiz_outbound_call_task = None

router = APIRouter(prefix="/telephony/vobiz", tags=["Vobiz Telephony"])
webhook_router = APIRouter(prefix="/telephony/vobiz", tags=["Vobiz Telephony Webhooks"])
ws_router = APIRouter(prefix="/telephony/vobiz", tags=["Vobiz Telephony Media"])
carrier_ws_router = APIRouter(prefix="/telephony/carrier", tags=["Carrier Telephony Media"])
ws_router = carrier_ws_router


class VobizOutboundCallRequest(BaseModel):
Expand Down Expand Up @@ -164,7 +168,9 @@ def _resolve_agent_for_answer(
if session and session.agent_id and session.organization_id:
return UUID(session.agent_id), UUID(session.organization_id), call_ref

agent_id, organization_id = resolve_inbound_agent_for_number(db, params.get("to"))
agent_id, organization_id, _telephony_integration_id = resolve_inbound_agent_for_number(
db, params.get("to")
)
if not agent_id or not organization_id:
return None, None, None
return agent_id, organization_id, None
Expand Down Expand Up @@ -478,7 +484,7 @@ async def vobiz_answer_webhook(
if session_token and call_uuid:
link_provider_call_id(db, call_ref=session_token, provider_call_id=call_uuid)

ws_url = build_vobiz_ws_url(
ws_url = build_carrier_ws_url(
agent_id=str(agent_id),
session=session_token,
persona_id=persona_id,
Expand Down Expand Up @@ -612,8 +618,8 @@ async def vobiz_recording_ready_webhook(
return {"status": "ok"}


@ws_router.websocket("/ws")
async def vobiz_media_websocket(websocket: WebSocket):
@carrier_ws_router.websocket("/ws")
async def carrier_media_websocket(websocket: WebSocket):
agent_id = websocket.query_params.get("agent_id")
session_token = websocket.query_params.get("session")
persona_id = websocket.query_params.get("persona_id")
Expand Down Expand Up @@ -655,7 +661,7 @@ async def vobiz_media_websocket(websocket: WebSocket):
# endregion
if not call_short_id:
logger.warning(
"No CallRecording for Vobiz session {}; live transcript and recording will not be linked",
"No CallRecording for carrier media session {}; live transcript and recording will not be linked",
session_token,
)
try:
Expand All @@ -679,25 +685,13 @@ async def vobiz_media_websocket(websocket: WebSocket):
persona_id=persona_id,
scenario_id=scenario_id,
)
from app.services.telephony.vobiz_agent_context import resolve_vobiz_telephony_run_params

run_params = resolve_vobiz_telephony_run_params(
db,
context=context,
call_direction=session.direction,
persona_id=persona_id,
scenario_id=scenario_id,
evaluator_id=session.evaluator_id,
)
serializer = VobizFrameSerializer(
serializer = build_carrier_frame_serializer(
provider_platform=getattr(call_row, "provider_platform", None),
stream_id=stream_id,
call_id=call_id,
auth_id=settings.VOBIZ_AUTH_ID,
auth_token=settings.VOBIZ_AUTH_TOKEN,
params=VobizFrameSerializer.InputParams(
sample_rate=8000,
api_base=settings.VOBIZ_API_BASE,
),
organization_id=UUID(session.organization_id),
db=db,
telephony_integration_id=telephony_integration_id_from_call_row(call_row),
Comment on lines +688 to +694

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Plivo streams fall back to Vobiz

If more than 200 newer call records exist when a Plivo media WebSocket starts, the bounded call-reference lookup returns no row and this code passes None as the provider. The serializer factory then selects Vobiz, causing automatic hangups to target Vobiz while the Plivo call remains connected.

Knowledge Base Used:

)

if context.use_voice_bundle_pipeline:
Expand All @@ -717,6 +711,8 @@ async def vobiz_media_websocket(websocket: WebSocket):
stt_api_key=context.stt_api_key,
tts_api_key=context.tts_api_key,
llm_api_key=context.llm_api_key,
llm_endpoint_url=context.llm_endpoint_url,
llm_base_url=context.llm_base_url,
serializer=serializer,
telephony_mode=True,
call_short_id=call_short_id,
Expand Down Expand Up @@ -751,12 +747,12 @@ async def vobiz_media_websocket(websocket: WebSocket):
persona_speaks_via_tts=run_params.persona_speaks_via_tts,
)
except ValueError as e:
logger.error("Vobiz media websocket setup failed: {}", e)
logger.error("Carrier media websocket setup failed: {}", e)
await websocket.close(code=1011, reason=str(e))
except WebSocketDisconnect:
logger.info("Vobiz media websocket disconnected")
logger.info("Carrier media websocket disconnected")
except Exception as e:
logger.error("Vobiz media websocket error: {}", e, exc_info=True)
logger.error("Carrier media websocket error: {}", e, exc_info=True)
try:
await websocket.close(code=1011, reason="Server error")
except Exception:
Expand Down
2 changes: 1 addition & 1 deletion app/api/v1/routes/voice_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -512,7 +512,7 @@ def resolve_azure_endpoint_for_provider(provider: ModelProvider) -> str | None:

llm_base_url = (
resolve_voice_llm_base_url(db, organization_id, voice_bundle, llm_provider)
if llm_provider
if llm_provider and voice_bundle
else None
)

Expand Down
2 changes: 1 addition & 1 deletion app/app_factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,7 @@ def create_app() -> FastAPI:
)
if _service_mode() == "media":
logger.info(
"Vobiz telephony edge: /telephony/vobiz/webhooks/* and /telephony/vobiz/ws"
"Carrier telephony edge: /telephony/vobiz/webhooks/* and /telephony/carrier/ws"
)
if settings.SERVICE_MODE == "api":
logger.info(
Expand Down
Loading
Loading