Unified start route to make all transports available.
This commit is contained in:
@@ -19,6 +19,10 @@ All bots must implement a `bot(runner_args)` async function as the entry point.
|
|||||||
The server automatically discovers and executes this function when connections
|
The server automatically discovers and executes this function when connections
|
||||||
are established.
|
are established.
|
||||||
|
|
||||||
|
By default the runner starts a single FastAPI server that supports WebRTC, Daily,
|
||||||
|
and telephony transports simultaneously. Clients declare which transport they want
|
||||||
|
via the ``transport`` field in the ``/start`` request body (default: ``"webrtc"``).
|
||||||
|
|
||||||
Single transport example::
|
Single transport example::
|
||||||
|
|
||||||
async def bot(runner_args: RunnerArguments):
|
async def bot(runner_args: RunnerArguments):
|
||||||
@@ -55,14 +59,33 @@ Supported transports:
|
|||||||
- WebRTC - Provides local WebRTC interface with prebuilt UI
|
- WebRTC - Provides local WebRTC interface with prebuilt UI
|
||||||
- Telephony - Handles webhook and WebSocket connections for Twilio, Telnyx, Plivo, Exotel
|
- Telephony - Handles webhook and WebSocket connections for Twilio, Telnyx, Plivo, Exotel
|
||||||
|
|
||||||
|
The ``/start`` endpoint accepts::
|
||||||
|
|
||||||
|
{
|
||||||
|
"transport": "webrtc", // "webrtc" | "daily" | "twilio" | "telnyx" |
|
||||||
|
// "plivo" | "exotel" — default: "webrtc"
|
||||||
|
|
||||||
|
// WebRTC-specific
|
||||||
|
"enableDefaultIceServers": false,
|
||||||
|
"body": {...},
|
||||||
|
|
||||||
|
// Daily-specific
|
||||||
|
"createDailyRoom": true,
|
||||||
|
"dailyRoomProperties": {...},
|
||||||
|
"dailyMeetingTokenProperties": {...},
|
||||||
|
"body": {...}
|
||||||
|
}
|
||||||
|
|
||||||
To run locally:
|
To run locally:
|
||||||
|
|
||||||
- WebRTC: `python bot.py -t webrtc`
|
- All transports (default): ``python bot.py``
|
||||||
- ESP32: `python bot.py -t webrtc --esp32 --host 192.168.1.100`
|
- WebRTC only: ``python bot.py -t webrtc``
|
||||||
- Daily (server): `python bot.py -t daily`
|
- ESP32: ``python bot.py -t webrtc --esp32 --host 192.168.1.100``
|
||||||
- Daily (direct, testing only): `python bot.py -d`
|
- Daily only: ``python bot.py -t daily``
|
||||||
- Telephony: `python bot.py -t twilio -x your_username.ngrok.io`
|
- Daily (direct, testing only): ``python bot.py -d``
|
||||||
- Exotel: `python bot.py -t exotel` (no proxy needed, but ngrok connection to HTTP 7860 is required)
|
- Telephony: ``python bot.py -t twilio -x your_username.ngrok.io``
|
||||||
|
- Exotel: ``python bot.py -t exotel`` (no proxy needed, but ngrok connection to HTTP 7860 is required)
|
||||||
|
- WhatsApp: ``python bot.py --whatsapp``
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
@@ -187,7 +210,7 @@ async def _run_telephony_bot(websocket: WebSocket, args: argparse.Namespace):
|
|||||||
|
|
||||||
|
|
||||||
def _configure_server_app(args: argparse.Namespace):
|
def _configure_server_app(args: argparse.Namespace):
|
||||||
"""Configure the module-level FastAPI app with transport-specific routes."""
|
"""Configure the module-level FastAPI app with routes for all transports."""
|
||||||
app.add_middleware(
|
app.add_middleware(
|
||||||
CORSMiddleware,
|
CORSMiddleware,
|
||||||
allow_origins=["*"],
|
allow_origins=["*"],
|
||||||
@@ -196,17 +219,178 @@ def _configure_server_app(args: argparse.Namespace):
|
|||||||
allow_headers=["*"],
|
allow_headers=["*"],
|
||||||
)
|
)
|
||||||
|
|
||||||
# Set up transport-specific routes
|
# Shared session store: session_id -> body data. Used by the WebRTC /start
|
||||||
if args.transport == "webrtc":
|
# flow and the /sessions/{session_id}/... proxy routes.
|
||||||
_setup_webrtc_routes(app, args)
|
active_sessions: dict[str, dict[str, Any]] = {}
|
||||||
if args.whatsapp:
|
|
||||||
_setup_whatsapp_routes(app, args)
|
_setup_webrtc_routes(app, args, active_sessions)
|
||||||
elif args.transport == "daily":
|
_setup_daily_routes(app, args)
|
||||||
_setup_daily_routes(app, args)
|
_setup_telephony_routes(app, args)
|
||||||
elif args.transport in TELEPHONY_TRANSPORTS:
|
_setup_unified_start_route(app, args, active_sessions)
|
||||||
_setup_telephony_routes(app, args)
|
|
||||||
else:
|
if args.whatsapp:
|
||||||
logger.warning(f"Unknown transport type: {args.transport}")
|
_setup_whatsapp_routes(app, args)
|
||||||
|
|
||||||
|
|
||||||
|
def _setup_unified_start_route(
|
||||||
|
app: FastAPI, args: argparse.Namespace, active_sessions: dict[str, dict[str, Any]]
|
||||||
|
):
|
||||||
|
"""Register the unified POST /start endpoint.
|
||||||
|
|
||||||
|
Handles WebRTC, Daily, and telephony transport start flows. Clients specify
|
||||||
|
which transport they want via the ``transport`` field in the request body.
|
||||||
|
When ``-t`` was passed on the command line, requests for any other transport
|
||||||
|
are rejected with HTTP 400.
|
||||||
|
"""
|
||||||
|
|
||||||
|
class IceServer(TypedDict, total=False):
|
||||||
|
urls: str | list[str]
|
||||||
|
|
||||||
|
class IceConfig(TypedDict):
|
||||||
|
iceServers: list[IceServer]
|
||||||
|
|
||||||
|
class StartBotResult(TypedDict, total=False):
|
||||||
|
sessionId: str
|
||||||
|
iceConfig: IceConfig | None
|
||||||
|
dailyRoom: str | None
|
||||||
|
dailyToken: str | None
|
||||||
|
websocketUrl: str | None
|
||||||
|
|
||||||
|
@app.post("/start")
|
||||||
|
async def start_agent(request: Request):
|
||||||
|
"""Start a bot session.
|
||||||
|
|
||||||
|
Accepts::
|
||||||
|
|
||||||
|
{
|
||||||
|
"transport": "webrtc", // "webrtc" | "daily" | "twilio" | "telnyx" |
|
||||||
|
// "plivo" | "exotel" — default: "webrtc"
|
||||||
|
|
||||||
|
// WebRTC-specific
|
||||||
|
"enableDefaultIceServers": false,
|
||||||
|
"body": {...},
|
||||||
|
|
||||||
|
// Daily-specific
|
||||||
|
"createDailyRoom": true,
|
||||||
|
"dailyRoomProperties": {...},
|
||||||
|
"dailyMeetingTokenProperties": {...},
|
||||||
|
"body": {...}
|
||||||
|
}
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
request_data = await request.json()
|
||||||
|
logger.debug(f"Received request: {request_data}")
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Failed to parse request body: {e}")
|
||||||
|
request_data = {}
|
||||||
|
|
||||||
|
# Determine transport: explicit field → legacy Daily hint → CLI default → webrtc
|
||||||
|
transport = request_data.get("transport")
|
||||||
|
if transport is None and request_data.get("createDailyRoom", False):
|
||||||
|
transport = "daily"
|
||||||
|
if transport is None:
|
||||||
|
transport = args.transport or "webrtc"
|
||||||
|
|
||||||
|
# Enforce restriction when -t was explicitly set on the command line
|
||||||
|
if args.transport is not None and transport != args.transport:
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=400,
|
||||||
|
detail=(
|
||||||
|
f"Transport '{transport}' is not allowed. "
|
||||||
|
f"Server is configured for '{args.transport}' only (-t {args.transport})."
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
if transport == "webrtc":
|
||||||
|
# WebRTC: register the session; the bot starts when the WebRTC offer arrives.
|
||||||
|
session_id = str(uuid.uuid4())
|
||||||
|
active_sessions[session_id] = request_data.get("body", {})
|
||||||
|
|
||||||
|
result: StartBotResult = {"sessionId": session_id}
|
||||||
|
if request_data.get("enableDefaultIceServers"):
|
||||||
|
result["iceConfig"] = IceConfig(
|
||||||
|
iceServers=[IceServer(urls=["stun:stun.l.google.com:19302"])]
|
||||||
|
)
|
||||||
|
return result
|
||||||
|
|
||||||
|
elif transport == "daily":
|
||||||
|
create_daily_room = request_data.get("createDailyRoom", False)
|
||||||
|
body = request_data.get("body", {})
|
||||||
|
daily_room_properties_dict = request_data.get("dailyRoomProperties", None)
|
||||||
|
daily_token_properties_dict = request_data.get("dailyMeetingTokenProperties", None)
|
||||||
|
|
||||||
|
bot_module = _get_bot_module()
|
||||||
|
|
||||||
|
existing_room_url = os.getenv("DAILY_ROOM_URL")
|
||||||
|
session_id = str(uuid.uuid4())
|
||||||
|
result = None
|
||||||
|
|
||||||
|
if create_daily_room or existing_room_url:
|
||||||
|
from pipecat.runner.daily import configure
|
||||||
|
from pipecat.transports.daily.utils import (
|
||||||
|
DailyMeetingTokenProperties,
|
||||||
|
DailyRoomProperties,
|
||||||
|
)
|
||||||
|
|
||||||
|
async with aiohttp.ClientSession() as session:
|
||||||
|
room_properties = None
|
||||||
|
if daily_room_properties_dict:
|
||||||
|
daily_room_properties_dict.setdefault(
|
||||||
|
"exp", time.time() + PIPECAT_ROOM_EXP_HOURS * 3600
|
||||||
|
)
|
||||||
|
daily_room_properties_dict.setdefault("eject_at_room_exp", True)
|
||||||
|
try:
|
||||||
|
room_properties = DailyRoomProperties(**daily_room_properties_dict)
|
||||||
|
logger.debug(f"Using custom room properties: {room_properties}")
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Failed to parse dailyRoomProperties: {e}")
|
||||||
|
|
||||||
|
token_properties = None
|
||||||
|
if daily_token_properties_dict:
|
||||||
|
try:
|
||||||
|
token_properties = DailyMeetingTokenProperties(
|
||||||
|
**daily_token_properties_dict
|
||||||
|
)
|
||||||
|
logger.debug(f"Using custom token properties: {token_properties}")
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Failed to parse dailyMeetingTokenProperties: {e}")
|
||||||
|
|
||||||
|
room_url, token = await configure(
|
||||||
|
session,
|
||||||
|
room_exp_duration=PIPECAT_ROOM_EXP_HOURS,
|
||||||
|
room_properties=room_properties,
|
||||||
|
token_properties=token_properties,
|
||||||
|
)
|
||||||
|
runner_args = DailyRunnerArguments(
|
||||||
|
room_url=room_url, token=token, body=body, session_id=session_id
|
||||||
|
)
|
||||||
|
result = {
|
||||||
|
"dailyRoom": room_url,
|
||||||
|
"dailyToken": token,
|
||||||
|
"sessionId": session_id,
|
||||||
|
}
|
||||||
|
else:
|
||||||
|
runner_args = RunnerArguments(body=body, session_id=session_id)
|
||||||
|
|
||||||
|
runner_args.cli_args = args
|
||||||
|
asyncio.create_task(bot_module.bot(runner_args))
|
||||||
|
return result
|
||||||
|
|
||||||
|
elif transport in TELEPHONY_TRANSPORTS:
|
||||||
|
# Telephony: the bot starts when the provider connects to /ws.
|
||||||
|
# Return the WebSocket URL so the caller knows where to point their provider.
|
||||||
|
session_id = str(uuid.uuid4())
|
||||||
|
scheme = "wss" if args.host != "localhost" else "ws"
|
||||||
|
return StartBotResult(
|
||||||
|
sessionId=session_id,
|
||||||
|
websocketUrl=f"{scheme}://{args.host}:{args.port}/ws",
|
||||||
|
)
|
||||||
|
|
||||||
|
else:
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=400,
|
||||||
|
detail=f"Unknown transport '{transport}'.",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def _resolve_download_path(folder: str, filename: str) -> Path:
|
def _resolve_download_path(folder: str, filename: str) -> Path:
|
||||||
@@ -220,7 +404,9 @@ def _resolve_download_path(folder: str, filename: str) -> Path:
|
|||||||
return file_path
|
return file_path
|
||||||
|
|
||||||
|
|
||||||
def _setup_webrtc_routes(app: FastAPI, args: argparse.Namespace):
|
def _setup_webrtc_routes(
|
||||||
|
app: FastAPI, args: argparse.Namespace, active_sessions: dict[str, dict[str, Any]]
|
||||||
|
):
|
||||||
"""Set up WebRTC-specific routes."""
|
"""Set up WebRTC-specific routes."""
|
||||||
try:
|
try:
|
||||||
from pipecat_ai_small_webrtc_prebuilt.frontend import SmallWebRTCPrebuiltUI
|
from pipecat_ai_small_webrtc_prebuilt.frontend import SmallWebRTCPrebuiltUI
|
||||||
@@ -236,19 +422,6 @@ def _setup_webrtc_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
logger.error(f"WebRTC transport dependencies not installed: {e}")
|
logger.error(f"WebRTC transport dependencies not installed: {e}")
|
||||||
return
|
return
|
||||||
|
|
||||||
class IceServer(TypedDict, total=False):
|
|
||||||
urls: str | list[str]
|
|
||||||
|
|
||||||
class IceConfig(TypedDict):
|
|
||||||
iceServers: list[IceServer]
|
|
||||||
|
|
||||||
class StartBotResult(TypedDict, total=False):
|
|
||||||
sessionId: str
|
|
||||||
iceConfig: IceConfig | None
|
|
||||||
|
|
||||||
# In-memory store of active sessions: session_id -> session info
|
|
||||||
active_sessions: dict[str, dict[str, Any]] = {}
|
|
||||||
|
|
||||||
# Mount the frontend
|
# Mount the frontend
|
||||||
app.mount("/client", SmallWebRTCPrebuiltUI)
|
app.mount("/client", SmallWebRTCPrebuiltUI)
|
||||||
|
|
||||||
@@ -315,29 +488,6 @@ def _setup_webrtc_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
await small_webrtc_handler.handle_patch_request(request)
|
await small_webrtc_handler.handle_patch_request(request)
|
||||||
return {"status": "success"}
|
return {"status": "success"}
|
||||||
|
|
||||||
@app.post("/start")
|
|
||||||
async def rtvi_start(request: Request):
|
|
||||||
"""Mimic Pipecat Cloud's /start endpoint."""
|
|
||||||
# Parse the request body
|
|
||||||
try:
|
|
||||||
request_data = await request.json()
|
|
||||||
logger.debug(f"Received request: {request_data}")
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Failed to parse request body: {e}")
|
|
||||||
request_data = {}
|
|
||||||
|
|
||||||
# Store session info immediately in memory, replicate the behavior expected on Pipecat Cloud
|
|
||||||
session_id = str(uuid.uuid4())
|
|
||||||
active_sessions[session_id] = request_data.get("body", {})
|
|
||||||
|
|
||||||
result: StartBotResult = {"sessionId": session_id}
|
|
||||||
if request_data.get("enableDefaultIceServers"):
|
|
||||||
result["iceConfig"] = IceConfig(
|
|
||||||
iceServers=[IceServer(urls=["stun:stun.l.google.com:19302"])]
|
|
||||||
)
|
|
||||||
|
|
||||||
return result
|
|
||||||
|
|
||||||
@app.api_route(
|
@app.api_route(
|
||||||
"/sessions/{session_id}/{path:path}",
|
"/sessions/{session_id}/{path:path}",
|
||||||
methods=["GET", "POST", "PUT", "PATCH", "DELETE"],
|
methods=["GET", "POST", "PUT", "PATCH", "DELETE"],
|
||||||
@@ -563,12 +713,10 @@ def _setup_whatsapp_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
def _setup_daily_routes(app: FastAPI, args: argparse.Namespace):
|
def _setup_daily_routes(app: FastAPI, args: argparse.Namespace):
|
||||||
"""Set up Daily-specific routes."""
|
"""Set up Daily-specific routes."""
|
||||||
|
|
||||||
@app.get("/")
|
@app.get("/daily")
|
||||||
async def create_room_and_start_agent():
|
async def create_room_and_start_agent():
|
||||||
"""Launch a Daily bot and redirect to room."""
|
"""Launch a Daily bot and redirect to room."""
|
||||||
print("Starting bot with Daily transport and redirecting to Daily room")
|
logger.debug("Starting bot with Daily transport and redirecting to Daily room")
|
||||||
|
|
||||||
import aiohttp
|
|
||||||
|
|
||||||
from pipecat.runner.daily import configure
|
from pipecat.runner.daily import configure
|
||||||
|
|
||||||
@@ -584,105 +732,6 @@ def _setup_daily_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
asyncio.create_task(bot_module.bot(runner_args))
|
asyncio.create_task(bot_module.bot(runner_args))
|
||||||
return RedirectResponse(room_url)
|
return RedirectResponse(room_url)
|
||||||
|
|
||||||
@app.post("/start")
|
|
||||||
async def start_agent(request: Request):
|
|
||||||
"""Handler for /start endpoints.
|
|
||||||
|
|
||||||
Expects POST body like::
|
|
||||||
{
|
|
||||||
"createDailyRoom": true,
|
|
||||||
"dailyRoomProperties": { "start_video_off": true },
|
|
||||||
"dailyMeetingTokenProperties": { "is_owner": true, "user_name": "Bot" },
|
|
||||||
"body": { "custom_data": "value" }
|
|
||||||
}
|
|
||||||
"""
|
|
||||||
print("Starting bot with Daily transport")
|
|
||||||
|
|
||||||
# Parse the request body
|
|
||||||
try:
|
|
||||||
request_data = await request.json()
|
|
||||||
logger.debug(f"Received request: {request_data}")
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Failed to parse request body: {e}")
|
|
||||||
request_data = {}
|
|
||||||
|
|
||||||
create_daily_room = request_data.get("createDailyRoom", False)
|
|
||||||
body = request_data.get("body", {})
|
|
||||||
daily_room_properties_dict = request_data.get("dailyRoomProperties", None)
|
|
||||||
daily_token_properties_dict = request_data.get("dailyMeetingTokenProperties", None)
|
|
||||||
|
|
||||||
bot_module = _get_bot_module()
|
|
||||||
|
|
||||||
existing_room_url = os.getenv("DAILY_ROOM_URL")
|
|
||||||
|
|
||||||
session_id = str(uuid.uuid4())
|
|
||||||
result = None
|
|
||||||
|
|
||||||
# Configure room if:
|
|
||||||
# 1. Explicitly requested via createDailyRoom in payload
|
|
||||||
# 2. Using pre-configured room from DAILY_ROOM_URL env var
|
|
||||||
if create_daily_room or existing_room_url:
|
|
||||||
import aiohttp
|
|
||||||
|
|
||||||
from pipecat.runner.daily import configure
|
|
||||||
from pipecat.transports.daily.utils import (
|
|
||||||
DailyMeetingTokenProperties,
|
|
||||||
DailyRoomProperties,
|
|
||||||
)
|
|
||||||
|
|
||||||
async with aiohttp.ClientSession() as session:
|
|
||||||
# Parse dailyRoomProperties if provided
|
|
||||||
room_properties = None
|
|
||||||
if daily_room_properties_dict:
|
|
||||||
# Apply Pipecat Cloud's session policy if caller didn't override.
|
|
||||||
daily_room_properties_dict.setdefault(
|
|
||||||
"exp", time.time() + PIPECAT_ROOM_EXP_HOURS * 3600
|
|
||||||
)
|
|
||||||
daily_room_properties_dict.setdefault("eject_at_room_exp", True)
|
|
||||||
try:
|
|
||||||
room_properties = DailyRoomProperties(**daily_room_properties_dict)
|
|
||||||
logger.debug(f"Using custom room properties: {room_properties}")
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Failed to parse dailyRoomProperties: {e}")
|
|
||||||
# Continue without custom properties
|
|
||||||
|
|
||||||
# Parse dailyMeetingTokenProperties if provided
|
|
||||||
token_properties = None
|
|
||||||
if daily_token_properties_dict:
|
|
||||||
try:
|
|
||||||
token_properties = DailyMeetingTokenProperties(
|
|
||||||
**daily_token_properties_dict
|
|
||||||
)
|
|
||||||
logger.debug(f"Using custom token properties: {token_properties}")
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Failed to parse dailyMeetingTokenProperties: {e}")
|
|
||||||
# Continue without custom properties
|
|
||||||
|
|
||||||
room_url, token = await configure(
|
|
||||||
session,
|
|
||||||
room_exp_duration=PIPECAT_ROOM_EXP_HOURS,
|
|
||||||
room_properties=room_properties,
|
|
||||||
token_properties=token_properties,
|
|
||||||
)
|
|
||||||
runner_args = DailyRunnerArguments(
|
|
||||||
room_url=room_url, token=token, body=body, session_id=session_id
|
|
||||||
)
|
|
||||||
result = {
|
|
||||||
"dailyRoom": room_url,
|
|
||||||
"dailyToken": token,
|
|
||||||
"sessionId": session_id,
|
|
||||||
}
|
|
||||||
else:
|
|
||||||
runner_args = RunnerArguments(body=body, session_id=session_id)
|
|
||||||
|
|
||||||
# Update CLI args.
|
|
||||||
runner_args.cli_args = args
|
|
||||||
|
|
||||||
# Start the bot in the background
|
|
||||||
asyncio.create_task(bot_module.bot(runner_args))
|
|
||||||
|
|
||||||
return result
|
|
||||||
|
|
||||||
if args.dialin:
|
if args.dialin:
|
||||||
|
|
||||||
@app.post("/daily-dialin-webhook")
|
@app.post("/daily-dialin-webhook")
|
||||||
@@ -731,8 +780,6 @@ def _setup_daily_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
detail="Missing required fields: From, To, callId, callDomain",
|
detail="Missing required fields: From, To, callId, callDomain",
|
||||||
)
|
)
|
||||||
|
|
||||||
import aiohttp
|
|
||||||
|
|
||||||
from pipecat.runner.daily import configure
|
from pipecat.runner.daily import configure
|
||||||
from pipecat.runner.types import DailyDialinRequest, DialinSettings
|
from pipecat.runner.types import DailyDialinRequest, DialinSettings
|
||||||
|
|
||||||
@@ -801,44 +848,51 @@ def _setup_daily_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
|
|
||||||
|
|
||||||
def _setup_telephony_routes(app: FastAPI, args: argparse.Namespace):
|
def _setup_telephony_routes(app: FastAPI, args: argparse.Namespace):
|
||||||
"""Set up telephony-specific routes."""
|
"""Set up telephony-specific routes.
|
||||||
# XML response templates (Exotel doesn't use XML webhooks)
|
|
||||||
XML_TEMPLATES = {
|
The WebSocket endpoint (``/ws``) is always registered so providers can
|
||||||
"twilio": f"""<?xml version="1.0" encoding="UTF-8"?>
|
connect directly. The XML webhook (``POST /``) is only registered when a
|
||||||
|
specific telephony transport is chosen via ``-t`` because the XML template
|
||||||
|
is provider-specific and requires a proxy hostname (``--proxy``).
|
||||||
|
"""
|
||||||
|
if args.transport in TELEPHONY_TRANSPORTS:
|
||||||
|
# XML response templates (Exotel doesn't use XML webhooks)
|
||||||
|
XML_TEMPLATES = {
|
||||||
|
"twilio": f"""<?xml version="1.0" encoding="UTF-8"?>
|
||||||
<Response>
|
<Response>
|
||||||
<Connect>
|
<Connect>
|
||||||
<Stream url="wss://{args.proxy}/ws"></Stream>
|
<Stream url="wss://{args.proxy}/ws"></Stream>
|
||||||
</Connect>
|
</Connect>
|
||||||
<Pause length="40"/>
|
<Pause length="40"/>
|
||||||
</Response>""",
|
</Response>""",
|
||||||
"telnyx": f"""<?xml version="1.0" encoding="UTF-8"?>
|
"telnyx": f"""<?xml version="1.0" encoding="UTF-8"?>
|
||||||
<Response>
|
<Response>
|
||||||
<Connect>
|
<Connect>
|
||||||
<Stream url="wss://{args.proxy}/ws" bidirectionalMode="rtp"></Stream>
|
<Stream url="wss://{args.proxy}/ws" bidirectionalMode="rtp"></Stream>
|
||||||
</Connect>
|
</Connect>
|
||||||
<Pause length="40"/>
|
<Pause length="40"/>
|
||||||
</Response>""",
|
</Response>""",
|
||||||
"plivo": f"""<?xml version="1.0" encoding="UTF-8"?>
|
"plivo": f"""<?xml version="1.0" encoding="UTF-8"?>
|
||||||
<Response>
|
<Response>
|
||||||
<Stream bidirectional="true" keepCallAlive="true" contentType="audio/x-mulaw;rate=8000">wss://{args.proxy}/ws</Stream>
|
<Stream bidirectional="true" keepCallAlive="true" contentType="audio/x-mulaw;rate=8000">wss://{args.proxy}/ws</Stream>
|
||||||
</Response>""",
|
</Response>""",
|
||||||
}
|
}
|
||||||
|
|
||||||
@app.post("/")
|
@app.post("/")
|
||||||
async def start_call():
|
async def start_call():
|
||||||
"""Handle telephony webhook and return XML response."""
|
"""Handle telephony webhook and return XML response."""
|
||||||
if args.transport == "exotel":
|
if args.transport == "exotel":
|
||||||
# Exotel doesn't use POST webhooks - redirect to proper documentation
|
# Exotel doesn't use POST webhooks - redirect to proper documentation
|
||||||
logger.debug("POST Exotel endpoint - not used")
|
logger.debug("POST Exotel endpoint - not used")
|
||||||
return {
|
return {
|
||||||
"error": "Exotel doesn't use POST webhooks",
|
"error": "Exotel doesn't use POST webhooks",
|
||||||
"websocket_url": f"wss://{args.proxy}/ws",
|
"websocket_url": f"wss://{args.proxy}/ws",
|
||||||
"note": "Configure the WebSocket URL above in your Exotel App Bazaar Voicebot Applet",
|
"note": "Configure the WebSocket URL above in your Exotel App Bazaar Voicebot Applet",
|
||||||
}
|
}
|
||||||
else:
|
else:
|
||||||
logger.debug(f"POST {args.transport.upper()} XML")
|
logger.debug(f"POST {args.transport.upper()} XML")
|
||||||
xml_content = XML_TEMPLATES.get(args.transport, "<Response></Response>")
|
xml_content = XML_TEMPLATES.get(args.transport, "<Response></Response>")
|
||||||
return HTMLResponse(content=xml_content, media_type="application/xml")
|
return HTMLResponse(content=xml_content, media_type="application/xml")
|
||||||
|
|
||||||
@app.websocket("/ws")
|
@app.websocket("/ws")
|
||||||
async def websocket_endpoint(websocket: WebSocket):
|
async def websocket_endpoint(websocket: WebSocket):
|
||||||
@@ -847,11 +901,6 @@ def _setup_telephony_routes(app: FastAPI, args: argparse.Namespace):
|
|||||||
logger.debug("WebSocket connection accepted")
|
logger.debug("WebSocket connection accepted")
|
||||||
await _run_telephony_bot(websocket, args)
|
await _run_telephony_bot(websocket, args)
|
||||||
|
|
||||||
@app.get("/")
|
|
||||||
async def start_agent():
|
|
||||||
"""Simple status endpoint for telephony transports."""
|
|
||||||
return {"status": f"Bot started with {args.transport}"}
|
|
||||||
|
|
||||||
|
|
||||||
async def _run_daily_direct(args: argparse.Namespace):
|
async def _run_daily_direct(args: argparse.Namespace):
|
||||||
"""Run Daily bot with direct connection (no FastAPI server)."""
|
"""Run Daily bot with direct connection (no FastAPI server)."""
|
||||||
@@ -922,22 +971,27 @@ def runner_port() -> int:
|
|||||||
def main(parser: argparse.ArgumentParser | None = None):
|
def main(parser: argparse.ArgumentParser | None = None):
|
||||||
"""Start the Pipecat development runner.
|
"""Start the Pipecat development runner.
|
||||||
|
|
||||||
Parses command-line arguments and starts a FastAPI server configured
|
Parses command-line arguments and starts a FastAPI server that supports
|
||||||
for the specified transport type.
|
WebRTC, Daily, and telephony transports simultaneously. Clients declare
|
||||||
|
which transport to use via the ``transport`` field in the ``/start`` body.
|
||||||
|
|
||||||
|
When ``-t`` is provided, the server restricts ``/start`` to that transport
|
||||||
|
only and displays transport-specific startup information.
|
||||||
|
|
||||||
The runner discovers and runs any ``bot(runner_args)`` function found in the
|
The runner discovers and runs any ``bot(runner_args)`` function found in the
|
||||||
calling module.
|
calling module.
|
||||||
|
|
||||||
Command-line arguments:
|
Command-line arguments:
|
||||||
- --host: Server host address (default: localhost) 879
|
- --host: Server host address (default: localhost)
|
||||||
- --port: Server port (default: 7860)
|
- --port: Server port (default: 7860)
|
||||||
- -t/--transport: Transport type (daily, webrtc, twilio, telnyx, plivo, exotel)
|
- -t/--transport: Restrict to a single transport and set as default for /start
|
||||||
|
(daily, webrtc, twilio, telnyx, plivo, exotel). Omit to support all transports.
|
||||||
- -x/--proxy: Public proxy hostname for telephony webhooks
|
- -x/--proxy: Public proxy hostname for telephony webhooks
|
||||||
- -d/--direct: Connect directly to Daily room (automatically sets transport to daily)
|
- -d/--direct: Connect directly to Daily room (automatically sets transport to daily)
|
||||||
- -f/--folder: Path to downloads folder
|
- -f/--folder: Path to downloads folder
|
||||||
- --dialin: Enable Daily PSTN dial-in webhook handling (requires Daily transport)
|
- --dialin: Enable Daily PSTN dial-in webhook handling
|
||||||
- --esp32: Enable SDP munging for ESP32 compatibility (requires --host with IP address)
|
- --esp32: Enable SDP munging for ESP32 compatibility (requires --host with IP address)
|
||||||
- --whatsapp: Ensure requried WhatsApp environment variables are present
|
- --whatsapp: Ensure required WhatsApp environment variables are present
|
||||||
- -v/--verbose: Increase logging verbosity
|
- -v/--verbose: Increase logging verbosity
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
@@ -958,8 +1012,11 @@ def main(parser: argparse.ArgumentParser | None = None):
|
|||||||
"--transport",
|
"--transport",
|
||||||
type=str,
|
type=str,
|
||||||
choices=["daily", "webrtc", *TELEPHONY_TRANSPORTS],
|
choices=["daily", "webrtc", *TELEPHONY_TRANSPORTS],
|
||||||
default="webrtc",
|
default=None,
|
||||||
help="Transport type",
|
help=(
|
||||||
|
"Restrict the server to a single transport and set it as the default for /start. "
|
||||||
|
"Omit to support all transports simultaneously (default behaviour)."
|
||||||
|
),
|
||||||
)
|
)
|
||||||
parser.add_argument("-x", "--proxy", help="Public proxy host name")
|
parser.add_argument("-x", "--proxy", help="Public proxy host name")
|
||||||
parser.add_argument(
|
parser.add_argument(
|
||||||
@@ -977,7 +1034,7 @@ def main(parser: argparse.ArgumentParser | None = None):
|
|||||||
"--dialin",
|
"--dialin",
|
||||||
action="store_true",
|
action="store_true",
|
||||||
default=False,
|
default=False,
|
||||||
help="Enable Daily PSTN dial-in webhook handling (requires Daily transport)",
|
help="Enable Daily PSTN dial-in webhook handling",
|
||||||
)
|
)
|
||||||
parser.add_argument(
|
parser.add_argument(
|
||||||
"--esp32",
|
"--esp32",
|
||||||
@@ -989,7 +1046,7 @@ def main(parser: argparse.ArgumentParser | None = None):
|
|||||||
"--whatsapp",
|
"--whatsapp",
|
||||||
action="store_true",
|
action="store_true",
|
||||||
default=False,
|
default=False,
|
||||||
help="Ensure requried WhatsApp environment variables are present",
|
help="Ensure required WhatsApp environment variables are present",
|
||||||
)
|
)
|
||||||
|
|
||||||
args = parser.parse_args()
|
args = parser.parse_args()
|
||||||
@@ -998,12 +1055,13 @@ def main(parser: argparse.ArgumentParser | None = None):
|
|||||||
if args.proxy:
|
if args.proxy:
|
||||||
args.proxy = _validate_and_clean_proxy(args.proxy)
|
args.proxy = _validate_and_clean_proxy(args.proxy)
|
||||||
|
|
||||||
# Auto-set transport to daily if --direct is used without explicit transport
|
# --direct implies Daily transport
|
||||||
if args.direct and args.transport == "webrtc": # webrtc is the default
|
if args.direct:
|
||||||
args.transport = "daily"
|
if args.transport is None or args.transport == "daily":
|
||||||
elif args.direct and args.transport != "daily":
|
args.transport = "daily"
|
||||||
logger.error("--direct flag only works with Daily transport (-t daily)")
|
else:
|
||||||
return
|
logger.error("--direct flag only works with Daily transport (-t daily)")
|
||||||
|
return
|
||||||
|
|
||||||
# Validate ESP32 requirements
|
# Validate ESP32 requirements
|
||||||
if args.esp32 and args.host == "localhost":
|
if args.esp32 and args.host == "localhost":
|
||||||
@@ -1011,7 +1069,7 @@ def main(parser: argparse.ArgumentParser | None = None):
|
|||||||
return
|
return
|
||||||
|
|
||||||
# Validate dial-in requirements
|
# Validate dial-in requirements
|
||||||
if args.dialin and args.transport != "daily":
|
if args.dialin and args.transport is not None and args.transport != "daily":
|
||||||
logger.error("--dialin flag only works with Daily transport (-t daily)")
|
logger.error("--dialin flag only works with Daily transport (-t daily)")
|
||||||
return
|
return
|
||||||
|
|
||||||
@@ -1029,28 +1087,38 @@ def main(parser: argparse.ArgumentParser | None = None):
|
|||||||
asyncio.run(_run_daily_direct(args))
|
asyncio.run(_run_daily_direct(args))
|
||||||
return
|
return
|
||||||
|
|
||||||
# Print startup message for server-based transports
|
# Print startup message
|
||||||
if args.transport == "webrtc":
|
print()
|
||||||
print()
|
if args.transport is None:
|
||||||
|
print("🚀 Bot ready!")
|
||||||
|
print(f" → WebRTC: http://{args.host}:{args.port}/client")
|
||||||
|
print(f" → Daily: http://{args.host}:{args.port}/daily")
|
||||||
|
print(f" → Telephony: ws://{args.host}:{args.port}/ws")
|
||||||
|
elif args.transport == "webrtc":
|
||||||
if args.esp32:
|
if args.esp32:
|
||||||
print(f"🚀 Bot ready! (ESP32 mode)")
|
print("🚀 Bot ready! (ESP32 mode)")
|
||||||
elif args.whatsapp:
|
elif args.whatsapp:
|
||||||
print(f"🚀 Bot ready! (WhatsApp)")
|
print("🚀 Bot ready! (WhatsApp)")
|
||||||
else:
|
else:
|
||||||
print(f"🚀 Bot ready!")
|
print("🚀 Bot ready! (WebRTC)")
|
||||||
print(f" → Open http://{args.host}:{args.port}/client in your browser")
|
print(f" → Open http://{args.host}:{args.port}/client in your browser")
|
||||||
print()
|
|
||||||
elif args.transport == "daily":
|
elif args.transport == "daily":
|
||||||
print()
|
print("🚀 Bot ready! (Daily)")
|
||||||
print(f"🚀 Bot ready!")
|
|
||||||
if args.dialin:
|
if args.dialin:
|
||||||
print(
|
print(
|
||||||
f" → Daily dial-in webhook: http://{args.host}:{args.port}/daily-dialin-webhook"
|
f" → Daily dial-in webhook: http://{args.host}:{args.port}/daily-dialin-webhook"
|
||||||
)
|
)
|
||||||
print(f" → Configure this URL in your Daily phone number settings")
|
print(f" → Configure this URL in your Daily phone number settings")
|
||||||
else:
|
else:
|
||||||
print(f" → Open http://{args.host}:{args.port} in your browser to start a session")
|
print(
|
||||||
print()
|
f" → Open http://{args.host}:{args.port}/daily in your browser to start a session"
|
||||||
|
)
|
||||||
|
elif args.transport in TELEPHONY_TRANSPORTS:
|
||||||
|
print(f"🚀 Bot ready! ({args.transport.capitalize()})")
|
||||||
|
if args.proxy:
|
||||||
|
print(f" → XML webhook: http://{args.host}:{args.port}/")
|
||||||
|
print(f" → WebSocket: ws://{args.host}:{args.port}/ws")
|
||||||
|
print()
|
||||||
|
|
||||||
RUNNER_DOWNLOADS_FOLDER = args.folder
|
RUNNER_DOWNLOADS_FOLDER = args.folder
|
||||||
RUNNER_HOST = args.host
|
RUNNER_HOST = args.host
|
||||||
|
|||||||
Reference in New Issue
Block a user