Merge pull request #113 from daily-co/aleix/only-subcribe-to-participant

only subcribe to participant
This commit is contained in:
Aleix Conchillo Flaqué
2024-04-10 10:47:29 +08:00
committed by GitHub
4 changed files with 96 additions and 10 deletions

View File

@@ -5,7 +5,7 @@ import os
import tkinter as tk import tkinter as tk
from dailyai.pipeline.frames import TextFrame, EndFrame from dailyai.pipeline.frames import TextFrame
from dailyai.pipeline.pipeline import Pipeline from dailyai.pipeline.pipeline import Pipeline
from dailyai.services.fal_ai_services import FalImageGenService from dailyai.services.fal_ai_services import FalImageGenService
from dailyai.transports.local_transport import LocalTransport from dailyai.transports.local_transport import LocalTransport

View File

@@ -4,8 +4,6 @@ import logging
from typing import AsyncGenerator from typing import AsyncGenerator
from PIL import Image
from dailyai.pipeline.aggregators import FrameProcessor from dailyai.pipeline.aggregators import FrameProcessor
from dailyai.pipeline.frames import ImageFrame, Frame, UserImageFrame from dailyai.pipeline.frames import ImageFrame, Frame, UserImageFrame
@@ -25,14 +23,13 @@ logger.setLevel(logging.DEBUG)
class UserImageProcessor(FrameProcessor): class UserImageProcessor(FrameProcessor):
async def process_frame(self, frame: Frame) -> AsyncGenerator[Frame, None]: async def process_frame(self, frame: Frame) -> AsyncGenerator[Frame, None]:
print(frame)
if isinstance(frame, UserImageFrame): if isinstance(frame, UserImageFrame):
yield ImageFrame(frame.image, frame.size) yield ImageFrame(frame.image, frame.size)
else: else:
yield frame yield frame
async def main(room_url): async def main(room_url: str, token):
transport = DailyTransport( transport = DailyTransport(
room_url, room_url,
token, token,
@@ -53,4 +50,4 @@ async def main(room_url):
if __name__ == "__main__": if __name__ == "__main__":
(url, token) = configure() (url, token) = configure()
asyncio.run(main(url)) asyncio.run(main(url, token))

View File

@@ -0,0 +1,72 @@
import asyncio
import io
import logging
import tkinter as tk
from typing import AsyncGenerator
from dailyai.pipeline.aggregators import FrameProcessor
from dailyai.pipeline.frames import ImageFrame, Frame, UserImageFrame
from dailyai.pipeline.pipeline import Pipeline
from dailyai.transports.daily_transport import DailyTransport
from dailyai.transports.local_transport import LocalTransport
from runner import configure
from dotenv import load_dotenv
load_dotenv(override=True)
logging.basicConfig(format=f"%(levelno)s %(asctime)s %(message)s")
logger = logging.getLogger("dailyai")
logger.setLevel(logging.DEBUG)
class UserImageProcessor(FrameProcessor):
async def process_frame(self, frame: Frame) -> AsyncGenerator[Frame, None]:
if isinstance(frame, UserImageFrame):
yield ImageFrame(frame.image, frame.size)
else:
yield frame
async def main(room_url: str, token):
tk_root = tk.Tk()
tk_root.title("dailyai")
local_transport = LocalTransport(
tk_root=tk_root,
camera_enabled=True,
camera_width=1280,
camera_height=720
)
transport = DailyTransport(
room_url,
token,
"Render participant video",
video_rendering_enabled=True
)
@transport.event_handler("on_first_other_participant_joined")
async def on_first_other_participant_joined(transport, participant):
transport.render_participant_video(participant["id"])
async def run_tk():
while not transport._stop_threads.is_set():
tk_root.update()
tk_root.update_idletasks()
await asyncio.sleep(0.1)
local_pipeline = Pipeline([UserImageProcessor()], source=transport.receive_queue)
await asyncio.gather(
transport.run(),
local_transport.run(local_pipeline, override_pipeline_source_queue=False),
run_tk()
)
if __name__ == "__main__":
(url, token) = configure()
asyncio.run(main(url, token))

View File

@@ -162,16 +162,20 @@ class DailyTransport(ThreadedTransport, EventHandler):
return decorator return decorator
def write_frame_to_camera(self, frame: bytes): def write_frame_to_camera(self, frame: bytes):
self.camera.write_frame(frame) if self._camera_enabled:
self.camera.write_frame(frame)
def write_frame_to_mic(self, frame: bytes): def write_frame_to_mic(self, frame: bytes):
self.mic.write_frames(frame) if self._mic_enabled:
self.mic.write_frames(frame)
def send_app_message(self, message: Any, participantId: str | None): def send_app_message(self, message: Any, participantId: str | None):
self.client.send_app_message(message, participantId) self.client.send_app_message(message, participantId)
def read_audio_frames(self, desired_frame_count): def read_audio_frames(self, desired_frame_count):
bytes = self._speaker.read_frames(desired_frame_count) bytes = b""
if self._speaker_enabled or self._vad_enabled:
bytes = self._speaker.read_frames(desired_frame_count)
return bytes return bytes
def _prerun(self): def _prerun(self):
@@ -240,9 +244,12 @@ class DailyTransport(ThreadedTransport, EventHandler):
) )
self._my_participant_id = self.client.participants()["local"]["id"] self._my_participant_id = self.client.participants()["local"]["id"]
# For performance reasons, never subscribe to video streams (unless a
# video renderer is registered).
self.client.update_subscription_profiles({ self.client.update_subscription_profiles({
"base": { "base": {
"camera": "subscribed" if self._video_rendering_enabled else "unsubscribed", "camera": "unsubscribed",
"screenVideo": "unsubscribed"
} }
}) })
@@ -281,6 +288,16 @@ class DailyTransport(ThreadedTransport, EventHandler):
color_format="RGB") -> None: color_format="RGB") -> None:
if not self._video_rendering_enabled: if not self._video_rendering_enabled:
self._logger.warn("Video rendering is not enabled") self._logger.warn("Video rendering is not enabled")
return
# Only enable camera subscription on this participant
self.client.update_subscriptions(participant_settings={
participant_id: {
"media": {
video_source: "subscribed"
}
}
})
self._video_renderers[participant_id] = { self._video_renderers[participant_id] = {
"framerate": framerate, "framerate": framerate,