feat: add a basic_room demo (#62)

This commit is contained in:
Théo Monnom
2023-05-09 15:08:57 +02:00
committed by GitHub
parent 1beee0e1f7
commit e8d1c7ee73
12 changed files with 53 additions and 19 deletions
+511
View File
@@ -0,0 +1,511 @@
use crate::events::UiCmd;
use crate::logo_track::LogoTrack;
use crate::sine_track::SineTrack;
use crate::video_renderer::VideoRenderer;
use crate::{events::AsyncCmd, video_grid::VideoGrid};
use egui::{Rounding, Stroke};
use egui_wgpu::WgpuConfiguration;
use futures::StreamExt;
use livekit::prelude::*;
use livekit::webrtc::audio_stream::native::NativeAudioStream;
use livekit::webrtc::video_frame::native::I420BufferExt;
use livekit::SimulateScenario;
use parking_lot::Mutex;
use std::collections::HashMap;
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc,
};
use tokio::sync::{mpsc, oneshot};
// Useful default constants for developing
const DEFAULT_URL: &str = "ws://localhost:7880";
const DEFAULT_TOKEN : &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE5MDY2MTMyODgsImlzcyI6IkFQSVRzRWZpZFpqclFvWSIsIm5hbWUiOiJuYXRpdmUiLCJuYmYiOjE2NzI2MTMyODgsInN1YiI6Im5hdGl2ZSIsInZpZGVvIjp7InJvb20iOiJ0ZXN0Iiwicm9vbUFkbWluIjp0cnVlLCJyb29tQ3JlYXRlIjp0cnVlLCJyb29tSm9pbiI6dHJ1ZSwicm9vbUxpc3QiOnRydWV9fQ.uSNIangMRu8jZD5mnRYoCHjcsQWCrJXgHCs0aNIgBFY";
// eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE5MDY2MTM0MzcsImlzcyI6IkFQSVRzRWZpZFpqclFvWSIsIm5hbWUiOiJ3ZWIiLCJuYmYiOjE2NzI2MTM0MzcsInN1YiI6IndlYiIsInZpZGVvIjp7InJvb20iOiJ0ZXN0Iiwicm9vbUFkbWluIjp0cnVlLCJyb29tQ3JlYXRlIjp0cnVlLCJyb29tSm9pbiI6dHJ1ZSwicm9vbUxpc3QiOnRydWV9fQ.DFTXt60n1kzGq4cSuOhbFBTQW2nd3rlcXKQ54sXsP8s
use winit::{
event::*,
event_loop::{ControlFlow, EventLoop},
window::{WindowBuilder, WindowId},
};
struct Session {
room: Room,
logo_track: LogoTrack,
sine_track: SineTrack,
close_tx: oneshot::Sender<()>,
handle: tokio::task::JoinHandle<()>,
}
struct AppState {
session: Mutex<Option<Session>>,
connecting: AtomicBool,
}
struct App {
state: Arc<AppState>,
video_renderers: HashMap<(ParticipantSid, TrackSid), VideoRenderer>,
egui_context: egui::Context,
egui_state: egui_winit::State,
egui_painter: egui_wgpu::winit::Painter,
window: winit::window::Window,
cmd_tx: mpsc::UnboundedSender<AsyncCmd>,
cmd_rx: mpsc::UnboundedReceiver<UiCmd>,
// UI State
lk_url: String,
lk_token: String,
connection_failure: Option<String>,
room_state: ConnectionState,
}
pub fn run(rt: tokio::runtime::Runtime) {
rt.block_on(async {
let event_loop = EventLoop::new();
let window = WindowBuilder::new()
.with_title("LiveKit - NativeSDK")
.build(&event_loop)
.unwrap();
let egui_context = egui::Context::default();
let egui_state = egui_winit::State::new(&event_loop);
let mut egui_painter =
egui_wgpu::winit::Painter::new(WgpuConfiguration::default(), 1, None, false);
egui_painter.set_window(Some(&window)).await.unwrap();
let (async_cmd_tx, mut async_cmd_rx) = mpsc::unbounded_channel::<AsyncCmd>();
let (ui_cmd_tx, ui_cmd_rx) = mpsc::unbounded_channel::<UiCmd>();
let state = Arc::new(AppState {
session: Default::default(),
connecting: AtomicBool::new(false),
});
let mut app = App {
state: state.clone(),
video_renderers: HashMap::default(),
egui_context,
egui_state,
egui_painter,
window,
cmd_tx: async_cmd_tx,
cmd_rx: ui_cmd_rx,
lk_url: DEFAULT_URL.to_owned(),
lk_token: DEFAULT_TOKEN.to_owned(),
connection_failure: None,
room_state: ConnectionState::Connected,
};
// Async event loop
tokio::spawn(async move {
while let Some(event) = async_cmd_rx.recv().await {
match event {
AsyncCmd::RoomConnect { url, token } => {
state.connecting.store(true, Ordering::SeqCst);
let res = Room::connect(&url, &token).await;
if let Ok((room, room_events)) = res {
let (close_tx, close_rx) = oneshot::channel();
let logo_track = LogoTrack::new(room.session());
let sine_track = SineTrack::new(room.session());
let handle = tokio::spawn(room_task(
state.clone(),
room_events,
close_rx,
ui_cmd_tx.clone(),
));
*state.session.lock() = Some(Session {
room,
logo_track,
sine_track,
close_tx,
handle,
});
let _ = ui_cmd_tx.send(UiCmd::ConnectResult { result: Ok(()) });
} else if let Err(err) = res {
let _ = ui_cmd_tx.send(UiCmd::ConnectResult { result: Err(err) });
}
state.connecting.store(false, Ordering::SeqCst);
}
AsyncCmd::RoomDisconnect => {
if let Some(session) = state.session.lock().take() {
let _ = session.room.close().await;
let _ = session.close_tx.send(());
let _ = session.handle.await;
}
}
AsyncCmd::SimulateScenario { scenario } => {
if let Some(session) = state.session.lock().as_ref() {
let _ = session.room.session().simulate_scenario(scenario).await;
}
}
AsyncCmd::ToggleLogo => {
if let Some(session) = state.session.lock().as_mut() {
let logo_track = &mut session.logo_track;
if !logo_track.is_published() {
logo_track.publish().await.unwrap();
} else {
logo_track.unpublish().await.unwrap();
}
}
}
AsyncCmd::ToggleSine => {
if let Some(session) = state.session.lock().as_mut() {
let sine_track = &mut session.sine_track;
sine_track.publish().await.unwrap();
}
}
}
}
});
tokio::task::block_in_place(move || loop {
// ui/main thread
event_loop.run(move |event, _, control_flow| {
app.update(event, control_flow);
});
});
});
}
async fn room_task(
_app_state: Arc<AppState>,
mut room_events: mpsc::UnboundedReceiver<RoomEvent>,
mut close_rx: oneshot::Receiver<()>,
ui_cmd_tx: mpsc::UnboundedSender<UiCmd>,
) {
loop {
tokio::select! {
Some(event) = room_events.recv() => {
let _ = ui_cmd_tx.send(UiCmd::RoomEvent{event});
}
_ = &mut close_rx => {
break;
}
}
}
}
impl App {
fn update<T>(&mut self, event: Event<'_, T>, control_flow: &mut ControlFlow) {
if let Ok(cmd) = self.cmd_rx.try_recv() {
match cmd {
UiCmd::ConnectResult { result } => {
if let Err(err) = result {
self.connection_failure = Some(err.to_string());
} else {
self.connection_failure = None
}
}
UiCmd::RoomEvent { event } => {
match event {
RoomEvent::TrackSubscribed {
track, participant, ..
} => {
match track.clone() {
RemoteTrack::Video(video_track) => {
// Create a new VideoRenderer
let video_renderer = VideoRenderer::new(
self.egui_painter.render_state().clone().unwrap(),
video_track.rtc_track(),
);
self.video_renderers
.insert((participant.sid(), track.sid()), video_renderer);
}
RemoteTrack::Audio(audio_track) => {
tokio::spawn(async move {
let mut stream =
NativeAudioStream::new(audio_track.rtc_track());
while let Some(_frame) = stream.next().await {
// Received audio frames
}
});
}
};
}
RoomEvent::TrackUnsubscribed {
track, participant, ..
} => {
self.video_renderers
.remove(&(participant.sid(), track.sid()));
}
_ => {}
}
}
}
}
match event {
Event::WindowEvent { window_id, event } => {
if let Some(flow) = self.on_window_event(window_id, event) {
*control_flow = flow;
}
}
Event::RedrawRequested(window_id) if window_id == self.window.id() => {
self.render();
}
Event::RedrawEventsCleared => {
self.window.request_redraw();
}
_ => {}
};
}
fn on_window_event(
&mut self,
_window_id: WindowId,
event: WindowEvent<'_>,
) -> Option<ControlFlow> {
if self
.egui_state
.on_event(&self.egui_context, &event)
.consumed
{
return None;
}
match event {
WindowEvent::CloseRequested => Some(ControlFlow::Exit),
WindowEvent::Resized(inner_size) => {
self.egui_painter
.on_window_resized(inner_size.width, inner_size.height);
None
}
WindowEvent::ScaleFactorChanged { new_inner_size, .. } => {
self.egui_painter
.on_window_resized(new_inner_size.width, new_inner_size.height);
None
}
_ => None,
}
}
fn ui(&mut self, ui: &mut egui::Ui) {
egui::TopBottomPanel::top("top_panel").show(ui.ctx(), |ui| {
egui::menu::bar(ui, |ui| {
ui.menu_button("Tools", |ui| {
if ui.button("Logs").clicked() {}
if ui.button("Profiler").clicked() {}
if ui.button("WebRTC Stats").clicked() {}
if ui.button("Events").clicked() {}
});
ui.menu_button("Simulate", |ui| {
if ui.button("SignalReconnect").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::SignalReconnect,
});
}
if ui.button("Speaker").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::Speaker,
});
}
if ui.button("NodeFailure").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::NodeFailure,
});
}
if ui.button("ServerLeave").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::ServerLeave,
});
}
if ui.button("Migration").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::Migration,
});
}
if ui.button("ForceTcp").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::ForceTcp,
});
}
if ui.button("ForceTls").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::SimulateScenario {
scenario: SimulateScenario::ForceTls,
});
}
});
ui.menu_button("Publish", |ui| {
if ui.button("Logo").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::ToggleLogo);
}
if ui.button("SineWave").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::ToggleSine);
}
});
});
});
egui::SidePanel::right("room_panel")
.default_width(256.0)
.show(ui.ctx(), |ui| {
ui.heading("Livekit - Connect to a room");
ui.separator();
ui.horizontal(|ui| {
ui.label("URL: ");
ui.text_edit_singleline(&mut self.lk_url);
});
ui.horizontal(|ui| {
ui.label("Token: ");
ui.text_edit_singleline(&mut self.lk_token);
});
ui.horizontal(|ui| {
let connecting = self.state.connecting.load(Ordering::SeqCst);
let session = self.state.session.lock();
ui.add_enabled_ui(!connecting && session.is_none(), |ui| {
if ui.button("Connect").clicked() {
self.connection_failure = None;
let _ = self.cmd_tx.send(AsyncCmd::RoomConnect {
url: self.lk_url.clone(),
token: self.lk_token.clone(),
});
}
});
if connecting {
ui.spinner();
}
if session.is_some() {
if ui.button("Disconnect").clicked() {
let _ = self.cmd_tx.send(AsyncCmd::RoomDisconnect);
}
}
});
if let Some(err) = &self.connection_failure {
ui.colored_label(egui::Color32::RED, err);
}
ui.separator();
{
// Room Info
if let Some(session) = self.state.session.lock().as_ref() {
ui.label(format!("Name: {}", session.room.session().name()));
ui.label(format!("SID: {}", session.room.session().sid()));
ui.label(format!(
"ConnectionState: {:?}",
session.room.session().connection_state()
));
ui.label(format!(
"ParticipantCount: {:?}",
session.room.session().participants().len() + 1
));
}
}
});
egui::CentralPanel::default().show(ui.ctx(), |ui| {
egui::ScrollArea::vertical().show(ui, |ui| {
VideoGrid::new("default_grid")
.max_columns(6)
.show(ui, |ui| {
if self.room_state == ConnectionState::Disconnected {
for _ in 0..20 {
ui.video_frame(|ui| {
egui::Frame::none().fill(egui::Color32::DARK_GRAY).show(
ui,
|ui| {
ui.allocate_space(ui.available_size());
},
);
});
}
} else {
// Render participant videos
for ((participant_sid, _), video_renderer) in &self.video_renderers {
ui.video_frame(|ui| {
let rect = ui.available_rect_before_wrap();
ui.painter().rect(
rect,
Rounding::none(),
egui::Color32::DARK_GRAY,
Stroke::NONE,
);
if let Some(tex) = video_renderer.texture_id() {
ui.painter().image(
tex,
rect,
egui::Rect::from_min_max(
egui::pos2(0.0, 0.0),
egui::pos2(1.0, 1.0),
),
egui::Color32::WHITE,
);
}
let name =
self.state.session.lock().as_ref().and_then(|session| {
session
.room
.session()
.participants()
.get(participant_sid)
.map(|p| p.name())
});
if let Some(name) = name {
ui.painter().text(
egui::pos2(rect.min.x + 5.0, rect.max.y - 5.0),
egui::Align2::LEFT_BOTTOM,
name,
egui::FontId::default(),
egui::Color32::WHITE,
);
}
});
}
}
});
});
});
}
fn render(&mut self) {
self.egui_state
.set_pixels_per_point(egui_winit::native_pixels_per_point(&self.window));
let raw_inputs = self.egui_state.take_egui_input(&self.window);
let full_output = self.egui_context.clone().run(raw_inputs, |ctx| {
egui::CentralPanel::default().show(ctx, |ui| {
self.ui(ui);
});
});
let clipped_primitives = self.egui_context.tessellate(full_output.shapes);
self.egui_painter.paint_and_update_textures(
egui_winit::native_pixels_per_point(&self.window),
[0.0, 0.0, 0.0, 0.0],
&clipped_primitives,
&full_output.textures_delta,
false,
);
self.egui_state.handle_platform_output(
&self.window,
&self.egui_context,
full_output.platform_output,
);
}
}
+17
View File
@@ -0,0 +1,17 @@
use livekit::prelude::*;
use livekit::{RoomResult, SimulateScenario};
#[derive(Debug)]
pub enum AsyncCmd {
RoomConnect { url: String, token: String },
RoomDisconnect,
SimulateScenario { scenario: SimulateScenario },
ToggleLogo, // Unpublish/Publish a logo track
ToggleSine,
}
#[derive(Debug)]
pub enum UiCmd {
ConnectResult { result: RoomResult<()> },
RoomEvent { event: RoomEvent },
}
+205
View File
@@ -0,0 +1,205 @@
use image::ImageFormat;
use image::RgbaImage;
use livekit::options::{TrackPublishOptions, VideoCaptureOptions};
use livekit::prelude::*;
use livekit::webrtc::{
native::yuv_helper,
video_frame::native::I420BufferExt,
video_frame::{I420Buffer, VideoFrame, VideoRotation},
video_source::native::NativeVideoSource,
};
use parking_lot::Mutex;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::oneshot;
use tokio::task::JoinHandle;
// The logo must not be bigger than the framebuffer
const PIXEL_SIZE: usize = 4;
const FRAME_RATE: u64 = 30;
const MOVE_SPEED: i32 = 16;
const FB_WIDTH: usize = 1280;
const FB_HEIGHT: usize = 720;
#[derive(Clone)]
struct FrameData {
image: Arc<RgbaImage>,
framebuffer: Arc<Mutex<Vec<u8>>>,
video_frame: Arc<Mutex<VideoFrame<I420Buffer>>>,
pos: (u32, u32),
direction: (i32, i32),
}
struct TrackHandle {
close_tx: oneshot::Sender<()>,
track: LocalVideoTrack,
task: JoinHandle<()>,
}
pub struct LogoTrack {
rtc_source: NativeVideoSource,
session: RoomSession,
handle: Option<TrackHandle>,
}
impl LogoTrack {
pub fn new(session: RoomSession) -> Self {
Self {
rtc_source: NativeVideoSource::default(),
session,
handle: None,
}
}
pub fn is_published(&self) -> bool {
self.handle.is_some()
}
pub async fn publish(&mut self) -> Result<(), RoomError> {
self.unpublish().await?;
let (close_tx, close_rx) = oneshot::channel();
let track = LocalVideoTrack::create_video_track(
"livekit_logo",
VideoCaptureOptions::default(),
self.rtc_source.clone(),
);
let task = tokio::spawn(Self::track_task(close_rx, self.rtc_source.clone()));
self.session
.local_participant()
.publish_track(
LocalTrack::Video(track.clone()),
TrackPublishOptions {
source: TrackSource::Camera,
..Default::default()
},
)
.await?;
let handle = TrackHandle {
close_tx,
task,
track,
};
self.handle = Some(handle);
Ok(())
}
pub async fn unpublish(&mut self) -> Result<(), RoomError> {
if let Some(handle) = self.handle.take() {
let _ = handle.close_tx.send(());
let _ = handle.task.await;
self.session
.local_participant()
.unpublish_track(handle.track.sid(), true)
.await?;
}
Ok(())
}
async fn track_task(mut close_rx: oneshot::Receiver<()>, rtc_source: NativeVideoSource) {
let mut interval = tokio::time::interval(Duration::from_millis(1000 / FRAME_RATE));
let image = tokio::task::spawn_blocking(|| {
image::load_from_memory_with_format(include_bytes!("moving-logo.png"), ImageFormat::Png)
.unwrap()
.to_rgba8()
})
.await
.unwrap();
let mut data = FrameData {
image: Arc::new(image),
framebuffer: Arc::new(Mutex::new(vec![0u8; (FB_WIDTH * FB_HEIGHT * 4) as usize])),
video_frame: Arc::new(Mutex::new(VideoFrame {
rotation: VideoRotation::VideoRotation0,
buffer: I420Buffer::new(FB_WIDTH as u32, FB_HEIGHT as u32),
timestamp: 0,
})),
pos: (0, 0),
direction: (1, 1),
};
loop {
tokio::select! {
_ = &mut close_rx => {
break;
}
_ = interval.tick() => {}
}
data.pos.0 = (data.pos.0 as i32 + data.direction.0 * MOVE_SPEED) as u32;
data.pos.1 = (data.pos.1 as i32 + data.direction.1 * MOVE_SPEED) as u32;
if data.pos.0 >= (FB_WIDTH - data.image.width() as usize) as u32 {
data.direction.0 = -1;
} else if data.pos.0 <= 0 {
data.direction.0 = 1;
}
if data.pos.1 >= (FB_HEIGHT - data.image.height() as usize) as u32 {
data.direction.1 = -1;
} else if data.pos.1 <= 0 {
data.direction.1 = 1;
}
tokio::task::spawn_blocking({
let data = data.clone();
let source = rtc_source.clone();
move || {
let image = data.image.as_raw();
let mut framebuffer = data.framebuffer.lock();
let mut video_frame = data.video_frame.lock();
let i420_buffer = &mut video_frame.buffer;
let (stride_y, stride_u, stride_v) = i420_buffer.strides();
let (data_y, data_u, data_v) = i420_buffer.data_mut();
framebuffer.fill(0);
for i in 0..data.image.height() as usize {
let x = data.pos.0 as usize;
let y = data.pos.1 as usize;
let frame_width = data.image.width() as usize;
let logo_stride = frame_width * PIXEL_SIZE;
let row_start = (x + ((i + y) * FB_WIDTH)) * PIXEL_SIZE;
let row_end = row_start + logo_stride;
framebuffer[row_start..row_end].copy_from_slice(
&image[i * logo_stride..i * logo_stride + logo_stride],
);
}
yuv_helper::abgr_to_i420(
&framebuffer,
(FB_WIDTH * PIXEL_SIZE) as u32,
data_y,
stride_y,
data_u,
stride_u,
data_v,
stride_v,
FB_WIDTH as i32,
FB_HEIGHT as i32,
)
.unwrap();
source.capture_frame(&*video_frame);
}
})
.await
.unwrap();
}
}
}
impl Drop for LogoTrack {
fn drop(&mut self) {
if let Some(handle) = self.handle.take() {
let _ = handle.close_tx.send(());
}
}
}
+17
View File
@@ -0,0 +1,17 @@
mod app;
mod events;
mod logo_track;
mod sine_track;
mod video_grid;
mod video_renderer;
fn main() {
tracing_subscriber::fmt::init();
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.unwrap();
app::run(rt);
}
Binary file not shown.

After

Width:  |  Height:  |  Size: 7.0 KiB

+127
View File
@@ -0,0 +1,127 @@
use livekit::options::{AudioCaptureOptions, TrackPublishOptions};
use livekit::webrtc::audio_frame::AudioFrame;
use livekit::{prelude::*, webrtc::audio_source::native::NativeAudioSource};
use parking_lot::Mutex;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::oneshot;
use tokio::task::JoinHandle;
#[derive(Clone)]
struct FrameData {
pub sample_rate: u32,
pub freq: f64,
pub amplitude: f64,
}
impl Default for FrameData {
fn default() -> Self {
Self {
sample_rate: 48000,
freq: 440.0,
amplitude: 1.0,
}
}
}
struct TrackHandle {
frame_data: Arc<Mutex<FrameData>>,
close_tx: oneshot::Sender<()>,
track: LocalAudioTrack,
task: JoinHandle<()>,
}
pub struct SineTrack {
rtc_source: NativeAudioSource,
session: RoomSession,
handle: Option<TrackHandle>,
}
impl SineTrack {
pub fn new(session: RoomSession) -> Self {
Self {
rtc_source: NativeAudioSource::default(),
session,
handle: None,
}
}
pub async fn publish(&mut self) -> Result<(), RoomError> {
let (close_tx, close_rx) = oneshot::channel();
let track = LocalAudioTrack::create_audio_track(
"sine_wave",
AudioCaptureOptions {
auto_gain_control: false,
echo_cancellation: false,
noise_suppression: false,
},
self.rtc_source.clone(),
);
let data = Arc::new(Mutex::new(FrameData::default()));
let task = tokio::spawn(Self::track_task(
close_rx,
self.rtc_source.clone(),
data.clone(),
));
self.session
.local_participant()
.publish_track(
LocalTrack::Audio(track.clone()),
TrackPublishOptions {
source: TrackSource::Microphone,
..Default::default()
},
)
.await?;
let handle = TrackHandle {
frame_data: data,
close_tx,
track,
task,
};
self.handle = Some(handle);
Ok(())
}
async fn track_task(
_close_rx: oneshot::Receiver<()>,
rtc_source: NativeAudioSource,
frame_options: Arc<Mutex<FrameData>>,
) {
let mut interval = tokio::time::interval(Duration::from_millis(10));
let mut samples_10ms = Vec::<i16>::new();
loop {
interval.tick().await;
let data = frame_options.lock();
let samples_count_10ms = (data.sample_rate / 100) as usize;
if samples_10ms.capacity() != samples_count_10ms {
samples_10ms.resize(samples_count_10ms, 0i16);
}
for i in 0..samples_count_10ms {
let val = data.amplitude
* f64::sin(
std::f64::consts::PI * 2.0 * data.freq * i as f64
/ samples_count_10ms as f64,
);
// WebRTC uses 16-bit signed PCM
samples_10ms[i] = (val * 32768.0) as i16;
}
rtc_source.capture_frame(&AudioFrame {
data: samples_10ms.clone(),
sample_rate: data.sample_rate as u32,
num_channels: 1,
samples_per_channel: samples_count_10ms as u32,
});
}
}
}
+175
View File
@@ -0,0 +1,175 @@
use std::cmp;
#[derive(Debug, Clone, Default, PartialEq)]
struct State {
num_videos: u32,
}
impl State {
pub fn load(ctx: &egui::Context, id: egui::Id) -> Option<Self> {
ctx.data(|i| i.get_temp(id))
}
pub fn store(self, ctx: &egui::Context, id: egui::Id) {
ctx.data_mut(|i| i.insert_temp(id, self))
}
}
pub const DEFAULT_VIDEO_SIZE: egui::Vec2 = egui::vec2(320.0, 180.0);
pub const DEFAULT_MAX_COLUMNS: u32 = 4;
pub const DEFAULT_SPACING: f32 = 16.0;
pub struct VideoGrid {
id: egui::Id,
// Current frame
available_rect: egui::Rect,
prev_state: State,
curr_state: State,
video_index: u32, // Kinda "cursor"
// Options
min_video_size: egui::Vec2,
max_columns: u32,
spacing: f32,
}
impl VideoGrid {
pub fn new(id_source: impl std::hash::Hash) -> Self {
Self {
id: egui::Id::new(id_source),
available_rect: egui::Rect::NAN,
prev_state: State::default(),
curr_state: State::default(),
video_index: 0,
min_video_size: DEFAULT_VIDEO_SIZE,
max_columns: DEFAULT_MAX_COLUMNS,
spacing: DEFAULT_SPACING,
}
}
pub fn show<R>(
mut self,
ui: &mut egui::Ui,
grid: impl FnOnce(&mut VideoGridContext) -> R,
) -> egui::InnerResponse<R> {
// TODO(theomonnom): Should I care about the current egui layout?
let prev_state = State::load(ui.ctx(), self.id);
let is_first_frame = prev_state.is_none();
self.prev_state = prev_state.unwrap_or_default();
self.available_rect = ui.available_rect_before_wrap();
ui.ctx()
.check_for_id_clash(self.id, self.available_rect, "VideoGrid");
ui.allocate_ui_at_rect(self.available_rect, |ui| {
ui.set_visible(!is_first_frame);
let mut ctx = VideoGridContext {
layout: &mut self,
ui,
};
let res = grid(&mut ctx);
// Save the new state
if self.curr_state != self.prev_state {
self.curr_state.clone().store(ui.ctx(), self.id);
ui.ctx().request_repaint();
}
res
})
}
fn next_frame_rect(&mut self) -> egui::Rect {
assert!(self.available_rect.is_finite());
assert!(self.spacing <= self.min_video_size.x);
// increment the amount of videos for the next frame
self.curr_state.num_videos += 1;
let num_videos = self.prev_state.num_videos;
if num_videos == 0 {
return egui::Rect::NOTHING;
}
let max_columns = self.max_columns;
let minimum_size = self.min_video_size;
let available_size = self.available_rect.size();
let calc_min_width =
|columns: u32| columns as f32 * minimum_size.x + (columns - 1) as f32 * self.spacing;
let total_columns = {
let mut est = (available_size.x / minimum_size.x) as u32 + 1;
if available_size.x < calc_min_width(est) {
est -= 1;
}
cmp::max(1, cmp::min(est, max_columns))
};
let aspect_ratio = minimum_size.x / minimum_size.y;
let remaining_width = available_size.x - calc_min_width(total_columns);
let w = minimum_size.x + remaining_width / total_columns as f32;
let h = w / aspect_ratio;
let x_index = self.video_index % total_columns;
let y_index = self.video_index / total_columns;
let x = {
let mut x = x_index as f32 * (w + self.spacing);
// vertically center the last row
let total_rows = num_videos / total_columns + 1;
if (y_index + 1) == total_rows {
let nb_items = num_videos - (total_rows - 1) * total_columns; // nb. of items on the last row
x += (total_columns - nb_items) as f32 * (w + self.spacing) / 2.0;
}
x
};
let y = y_index as f32 * (h + self.spacing);
let min = egui::pos2(x, y) + self.available_rect.left_top().to_vec2();
let max = egui::pos2(w, h) + min.to_vec2();
self.video_index += 1;
egui::Rect { min, max }
}
}
impl VideoGrid {
pub fn min_video_size(mut self, min_video_size: egui::Vec2) -> Self {
self.min_video_size = min_video_size;
self
}
pub fn max_columns(mut self, max_columns: u32) -> Self {
self.max_columns = max_columns;
self
}
pub fn spacing(mut self, spacing: f32) -> Self {
self.spacing = spacing;
self
}
}
pub struct VideoGridContext<'a> {
layout: &'a mut VideoGrid,
ui: &'a mut egui::Ui,
}
impl<'a> VideoGridContext<'a> {
pub fn video_frame(&mut self, add_contents: impl FnOnce(&mut egui::Ui)) -> egui::Response {
let frame_rect = self.layout.next_frame_rect();
let mut child_ui = self.ui.child_ui(frame_rect, egui::Layout::default());
add_contents(&mut child_ui);
self.ui.allocate_rect(frame_rect, egui::Sense::hover())
}
}
+179
View File
@@ -0,0 +1,179 @@
use futures::StreamExt;
use livekit::webrtc::native::yuv_helper;
use livekit::webrtc::prelude::*;
use livekit::webrtc::video_stream::native::NativeVideoStream;
use std::{
ops::DerefMut,
sync::{Arc, Mutex},
};
use tracing::debug_span;
pub struct VideoRenderer {
internal: Arc<Mutex<RendererInternal>>,
rtc_track: RtcVideoTrack,
}
struct RendererInternal {
render_state: egui_wgpu::RenderState,
width: u32,
height: u32,
rgba_data: Vec<u8>,
texture: Option<wgpu::Texture>,
texture_view: Option<wgpu::TextureView>,
egui_texture: Option<egui::TextureId>,
}
impl RendererInternal {
fn ensure_texture_size(&mut self, width: u32, height: u32) {
if self.width == width && self.height == height {
return;
}
self.width = width;
self.height = height;
self.rgba_data.resize((width * height * 4) as usize, 0);
self.texture = Some(
self.render_state
.device
.create_texture(&wgpu::TextureDescriptor {
label: Some("lk-videotexture"),
usage: wgpu::TextureUsages::TEXTURE_BINDING | wgpu::TextureUsages::COPY_DST,
dimension: wgpu::TextureDimension::D2,
size: wgpu::Extent3d {
width,
height,
..Default::default()
},
sample_count: 1,
mip_level_count: 1,
format: wgpu::TextureFormat::Rgba8UnormSrgb,
view_formats: &[wgpu::TextureFormat::Rgba8UnormSrgb],
}),
);
self.texture_view = Some(self.texture.as_mut().unwrap().create_view(
&wgpu::TextureViewDescriptor {
label: Some("lk-videotexture-view"),
format: Some(wgpu::TextureFormat::Rgba8UnormSrgb),
dimension: Some(wgpu::TextureViewDimension::D2),
mip_level_count: Some(1),
array_layer_count: Some(1),
..Default::default()
},
));
if let Some(texture_id) = self.egui_texture {
// Update the existing texture
self.render_state
.renderer
.write()
.update_egui_texture_from_wgpu_texture(
&*self.render_state.device,
self.texture_view.as_ref().unwrap(),
wgpu::FilterMode::Linear,
texture_id,
);
} else {
self.egui_texture = Some(self.render_state.renderer.write().register_native_texture(
&*self.render_state.device,
self.texture_view.as_ref().unwrap(),
wgpu::FilterMode::Linear,
));
}
}
}
impl VideoRenderer {
pub fn new(render_state: egui_wgpu::RenderState, rtc_track: RtcVideoTrack) -> Self {
let internal = Arc::new(Mutex::new(RendererInternal {
render_state,
width: 0,
height: 0,
rgba_data: Vec::default(),
texture: None,
texture_view: None,
egui_texture: None,
}));
let mut video_sink = NativeVideoStream::new(rtc_track.clone());
tokio::spawn({
let internal = internal.clone();
async move {
while let Some(frame) = video_sink.next().await {
let internal = internal.clone();
// Process the frame
let _ = tokio::task::spawn_blocking(move || {
let span = debug_span!("texture_upload");
let _enter = span.enter();
let mut internal = internal.lock().unwrap();
let buffer = frame.buffer.to_i420();
let width: u32 = buffer.width().try_into().unwrap();
let height: u32 = buffer.height().try_into().unwrap();
internal.ensure_texture_size(width, height);
let rgba_ptr = internal.rgba_data.deref_mut();
let rgba_stride = buffer.width() * 4;
let (stride_y, stride_u, stride_v) = buffer.strides();
let (data_y, data_u, data_v) = buffer.data();
yuv_helper::i420_to_abgr(
data_y,
stride_y,
data_u,
stride_u,
data_v,
stride_v,
rgba_ptr,
rgba_stride,
buffer.width() as i32,
buffer.height() as i32,
)
.unwrap();
let copy_desc = wgpu::ImageCopyTexture {
texture: internal.texture.as_ref().unwrap(),
mip_level: 0,
origin: wgpu::Origin3d::default(),
aspect: wgpu::TextureAspect::default(),
};
let copy_layout = wgpu::ImageDataLayout {
bytes_per_row: Some(width * 4),
..Default::default()
};
let copy_size = wgpu::Extent3d {
width,
height,
..Default::default()
};
internal.render_state.queue.write_texture(
copy_desc,
&internal.rgba_data,
copy_layout,
copy_size,
);
})
.await;
}
}
});
Self {
rtc_track,
internal,
}
}
pub fn texture_id(&self) -> Option<egui::TextureId> {
self.internal.lock().unwrap().egui_texture.clone()
}
}