Merge pull request #1533 from pipecat-ai/daily_small_webrtc
Example interoping between SmallWebRTC and Daily
This commit is contained in:
61
examples/p2p-webrtc/daily-interop-bridge/README.md
Normal file
61
examples/p2p-webrtc/daily-interop-bridge/README.md
Normal file
@@ -0,0 +1,61 @@
|
|||||||
|
# SmallWebRTC and Daily
|
||||||
|
|
||||||
|
A Pipecat example demonstrating how to interoperate audio and video between `SmallWebRTCTransport` and `DailyTransport`.
|
||||||
|
|
||||||
|
## 🚀 Quick Start
|
||||||
|
|
||||||
|
### 1️⃣ Start the Bot Server
|
||||||
|
|
||||||
|
#### 🔧 Set Up the Environment
|
||||||
|
1. Create and activate a virtual environment:
|
||||||
|
```bash
|
||||||
|
python3 -m venv venv
|
||||||
|
source venv/bin/activate # On Windows: venv\Scripts\activate
|
||||||
|
```
|
||||||
|
|
||||||
|
2. Install dependencies:
|
||||||
|
```bash
|
||||||
|
pip install -r requirements.txt
|
||||||
|
```
|
||||||
|
|
||||||
|
3. Configure environment variables:
|
||||||
|
- Copy `env.example` to `.env`
|
||||||
|
```bash
|
||||||
|
cp env.example .env
|
||||||
|
```
|
||||||
|
- Add your API keys
|
||||||
|
|
||||||
|
#### ▶️ Run the Server
|
||||||
|
```bash
|
||||||
|
python server.py
|
||||||
|
```
|
||||||
|
|
||||||
|
### 1️⃣ Connect the first client using Daily Prebuilt
|
||||||
|
|
||||||
|
- Open your browser and navigate to the same URL that you configured inside your `.env` file:
|
||||||
|
- `DAILY_SAMPLE_ROOM_URL`
|
||||||
|
|
||||||
|
### 2️⃣ Connect the second client using SmallWebRTC Prebuilt UI
|
||||||
|
|
||||||
|
- Open your browser and navigate to:
|
||||||
|
👉 http://localhost:7860
|
||||||
|
- (Or use your custom port, if configured)
|
||||||
|
|
||||||
|
## ⚠️ Important Note
|
||||||
|
Ensure the bot server is running before using any client implementations.
|
||||||
|
|
||||||
|
## 📌 Requirements
|
||||||
|
|
||||||
|
- Python **3.10+**
|
||||||
|
- Node.js **16+** (for JavaScript components)
|
||||||
|
- Google API Key
|
||||||
|
- Modern web browser with WebRTC support
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### 💡 Notes
|
||||||
|
- Ensure all dependencies are installed before running the server.
|
||||||
|
- Check the `.env` file for missing configurations.
|
||||||
|
- WebRTC requires a secure environment (HTTPS) for full functionality in production.
|
||||||
|
|
||||||
|
Happy coding! 🎉
|
||||||
128
examples/p2p-webrtc/daily-interop-bridge/bot.py
Normal file
128
examples/p2p-webrtc/daily-interop-bridge/bot.py
Normal file
@@ -0,0 +1,128 @@
|
|||||||
|
#
|
||||||
|
# Copyright (c) 2025, Daily
|
||||||
|
#
|
||||||
|
# SPDX-License-Identifier: BSD 2-Clause License
|
||||||
|
#
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
|
||||||
|
from dotenv import load_dotenv
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
|
from pipecat.frames.frames import (
|
||||||
|
InputAudioRawFrame,
|
||||||
|
InputImageRawFrame,
|
||||||
|
OutputAudioRawFrame,
|
||||||
|
OutputImageRawFrame,
|
||||||
|
)
|
||||||
|
from pipecat.pipeline.parallel_pipeline import ParallelPipeline
|
||||||
|
from pipecat.pipeline.pipeline import Pipeline
|
||||||
|
from pipecat.pipeline.runner import PipelineRunner
|
||||||
|
from pipecat.pipeline.task import PipelineParams, PipelineTask
|
||||||
|
from pipecat.processors.frame_processor import Frame, FrameDirection, FrameProcessor
|
||||||
|
from pipecat.transports.base_transport import TransportParams
|
||||||
|
from pipecat.transports.network.small_webrtc import SmallWebRTCTransport
|
||||||
|
from pipecat.transports.services.daily import DailyParams, DailyTransport
|
||||||
|
|
||||||
|
load_dotenv(override=True)
|
||||||
|
|
||||||
|
logger.remove(0)
|
||||||
|
logger.add(sys.stderr, level="DEBUG")
|
||||||
|
|
||||||
|
|
||||||
|
class MirrorProcessor(FrameProcessor):
|
||||||
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
|
await super().process_frame(frame, direction)
|
||||||
|
|
||||||
|
if isinstance(frame, InputAudioRawFrame):
|
||||||
|
await self.push_frame(
|
||||||
|
OutputAudioRawFrame(
|
||||||
|
audio=frame.audio,
|
||||||
|
sample_rate=frame.sample_rate,
|
||||||
|
num_channels=frame.num_channels,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif isinstance(frame, InputImageRawFrame):
|
||||||
|
await self.push_frame(
|
||||||
|
OutputImageRawFrame(image=frame.image, size=frame.size, format=frame.format)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
await self.push_frame(frame, direction)
|
||||||
|
|
||||||
|
|
||||||
|
async def run_bot(webrtc_connection):
|
||||||
|
pipecat_transport = SmallWebRTCTransport(
|
||||||
|
webrtc_connection=webrtc_connection,
|
||||||
|
params=TransportParams(
|
||||||
|
camera_in_enabled=True,
|
||||||
|
camera_out_enabled=True,
|
||||||
|
camera_out_is_live=True,
|
||||||
|
audio_in_enabled=True,
|
||||||
|
audio_out_enabled=True,
|
||||||
|
camera_out_width=1280,
|
||||||
|
camera_out_height=720,
|
||||||
|
vad_enabled=False,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
room_url = os.getenv("DAILY_SAMPLE_ROOM_URL", "")
|
||||||
|
daily_transport = DailyTransport(
|
||||||
|
room_url,
|
||||||
|
None,
|
||||||
|
"SmallWebRTC",
|
||||||
|
params=DailyParams(
|
||||||
|
camera_in_enabled=True,
|
||||||
|
camera_out_enabled=True,
|
||||||
|
camera_out_is_live=True,
|
||||||
|
audio_in_enabled=True,
|
||||||
|
audio_out_enabled=True,
|
||||||
|
camera_out_width=1280,
|
||||||
|
camera_out_height=720,
|
||||||
|
vad_enabled=False,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
pipeline = Pipeline(
|
||||||
|
[
|
||||||
|
ParallelPipeline(
|
||||||
|
[
|
||||||
|
daily_transport.input(),
|
||||||
|
MirrorProcessor(),
|
||||||
|
pipecat_transport.output(),
|
||||||
|
],
|
||||||
|
[
|
||||||
|
pipecat_transport.input(),
|
||||||
|
MirrorProcessor(),
|
||||||
|
daily_transport.output(),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
task = PipelineTask(
|
||||||
|
pipeline,
|
||||||
|
params=PipelineParams(
|
||||||
|
allow_interruptions=False,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
@daily_transport.event_handler("on_participant_joined")
|
||||||
|
async def on_participant_joined(transport, participant):
|
||||||
|
await transport.capture_participant_video(participant["id"])
|
||||||
|
|
||||||
|
@pipecat_transport.event_handler("on_client_connected")
|
||||||
|
async def on_client_connected(transport, client):
|
||||||
|
logger.info("Pipecat Client connected")
|
||||||
|
|
||||||
|
@pipecat_transport.event_handler("on_client_disconnected")
|
||||||
|
async def on_client_disconnected(transport, client):
|
||||||
|
logger.info("Pipecat Client disconnected")
|
||||||
|
|
||||||
|
@pipecat_transport.event_handler("on_client_closed")
|
||||||
|
async def on_client_closed(transport, client):
|
||||||
|
logger.info("Pipecat Client closed")
|
||||||
|
await task.cancel()
|
||||||
|
|
||||||
|
runner = PipelineRunner(handle_sigint=False)
|
||||||
|
|
||||||
|
await runner.run(task)
|
||||||
2
examples/p2p-webrtc/daily-interop-bridge/env.example
Normal file
2
examples/p2p-webrtc/daily-interop-bridge/env.example
Normal file
@@ -0,0 +1,2 @@
|
|||||||
|
DAILY_API_KEY=
|
||||||
|
DAILY_SAMPLE_ROOM_URL=
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
python-dotenv
|
||||||
|
fastapi[all]
|
||||||
|
uvicorn
|
||||||
|
aiortc
|
||||||
|
pipecat-ai[silero, webrtc, daily]
|
||||||
89
examples/p2p-webrtc/daily-interop-bridge/server.py
Normal file
89
examples/p2p-webrtc/daily-interop-bridge/server.py
Normal file
@@ -0,0 +1,89 @@
|
|||||||
|
import argparse
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
from contextlib import asynccontextmanager
|
||||||
|
from typing import Dict
|
||||||
|
|
||||||
|
import uvicorn
|
||||||
|
from bot import run_bot
|
||||||
|
from dotenv import load_dotenv
|
||||||
|
from fastapi import BackgroundTasks, FastAPI
|
||||||
|
from fastapi.responses import RedirectResponse
|
||||||
|
from pipecat_ai_small_webrtc_prebuilt.frontend import SmallWebRTCPrebuiltUI
|
||||||
|
|
||||||
|
from pipecat.transports.network.webrtc_connection import SmallWebRTCConnection
|
||||||
|
|
||||||
|
# Load environment variables
|
||||||
|
load_dotenv(override=True)
|
||||||
|
|
||||||
|
logger = logging.getLogger("pc")
|
||||||
|
|
||||||
|
app = FastAPI()
|
||||||
|
|
||||||
|
# Store connections by pc_id
|
||||||
|
pcs_map: Dict[str, SmallWebRTCConnection] = {}
|
||||||
|
|
||||||
|
ice_servers = ["stun:stun.l.google.com:19302"]
|
||||||
|
|
||||||
|
# Mount the frontend at /
|
||||||
|
app.mount("/prebuilt", SmallWebRTCPrebuiltUI)
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/", include_in_schema=False)
|
||||||
|
async def root_redirect():
|
||||||
|
return RedirectResponse(url="/prebuilt/")
|
||||||
|
|
||||||
|
|
||||||
|
@app.post("/api/offer")
|
||||||
|
async def offer(request: dict, background_tasks: BackgroundTasks):
|
||||||
|
pc_id = request.get("pc_id")
|
||||||
|
|
||||||
|
if pc_id and pc_id in pcs_map:
|
||||||
|
pipecat_connection = pcs_map[pc_id]
|
||||||
|
logger.info(f"Reusing existing connection for pc_id: {pc_id}")
|
||||||
|
await pipecat_connection.renegotiate(
|
||||||
|
sdp=request["sdp"], type=request["type"], restart_pc=request.get("restart_pc", False)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
pipecat_connection = SmallWebRTCConnection(ice_servers)
|
||||||
|
await pipecat_connection.initialize(sdp=request["sdp"], type=request["type"])
|
||||||
|
|
||||||
|
@pipecat_connection.event_handler("closed")
|
||||||
|
async def handle_disconnected(webrtc_connection: SmallWebRTCConnection):
|
||||||
|
logger.info(f"Discarding peer connection for pc_id: {webrtc_connection.pc_id}")
|
||||||
|
pcs_map.pop(webrtc_connection.pc_id, None)
|
||||||
|
|
||||||
|
background_tasks.add_task(run_bot, pipecat_connection)
|
||||||
|
|
||||||
|
answer = pipecat_connection.get_answer()
|
||||||
|
# Updating the peer connection inside the map
|
||||||
|
pcs_map[answer["pc_id"]] = pipecat_connection
|
||||||
|
|
||||||
|
return answer
|
||||||
|
|
||||||
|
|
||||||
|
@asynccontextmanager
|
||||||
|
async def lifespan(app: FastAPI):
|
||||||
|
yield # Run app
|
||||||
|
coros = [pc.close() for pc in pcs_map.values()]
|
||||||
|
await asyncio.gather(*coros)
|
||||||
|
pcs_map.clear()
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
parser = argparse.ArgumentParser(description="WebRTC demo")
|
||||||
|
parser.add_argument(
|
||||||
|
"--host", default="localhost", help="Host for HTTP server (default: localhost)"
|
||||||
|
)
|
||||||
|
parser.add_argument(
|
||||||
|
"--port", type=int, default=7860, help="Port for HTTP server (default: 7860)"
|
||||||
|
)
|
||||||
|
parser.add_argument("--verbose", "-v", action="count")
|
||||||
|
args = parser.parse_args()
|
||||||
|
|
||||||
|
if args.verbose:
|
||||||
|
logging.basicConfig(level=logging.DEBUG)
|
||||||
|
else:
|
||||||
|
logging.basicConfig(level=logging.INFO)
|
||||||
|
|
||||||
|
uvicorn.run(app, host=args.host, port=args.port)
|
||||||
@@ -386,10 +386,13 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
async def _draw_image(self, frame: OutputImageRawFrame):
|
async def _draw_image(self, frame: OutputImageRawFrame):
|
||||||
desired_size = (self._params.camera_out_width, self._params.camera_out_height)
|
desired_size = (self._params.camera_out_width, self._params.camera_out_height)
|
||||||
|
|
||||||
|
# TODO: we should refactor in the future to support dynamic resolutions
|
||||||
|
# which is kind of what happens in P2P connections.
|
||||||
|
# We need to add support for that inside the DailyTransport
|
||||||
if frame.size != desired_size:
|
if frame.size != desired_size:
|
||||||
image = Image.frombytes(frame.format, frame.size, frame.image)
|
image = Image.frombytes(frame.format, frame.size, frame.image)
|
||||||
resized_image = image.resize(desired_size)
|
resized_image = image.resize(desired_size)
|
||||||
logger.warning(f"{frame} does not have the expected size {desired_size}, resizing")
|
# logger.warning(f"{frame} does not have the expected size {desired_size}, resizing")
|
||||||
frame = OutputImageRawFrame(
|
frame = OutputImageRawFrame(
|
||||||
resized_image.tobytes(), resized_image.size, resized_image.format
|
resized_image.tobytes(), resized_image.size, resized_image.format
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user