From cbfa07f8e8ee241f7c980adcc8d62a4b2fdba5af Mon Sep 17 00:00:00 2001 From: Doloro1978 Date: Wed, 17 Jun 2026 11:55:27 +0100 Subject: [PATCH] more reshape that i missed --- .gitignore | 1 + Cargo.lock | 6 -- crates/server/Cargo.toml | 5 - crates/server/src/main.rs | 75 ++------------- flake.nix | 45 +++++++-- src/http.rs | 69 -------------- src/main.rs | 173 ---------------------------------- src/media.rs | 147 ----------------------------- src/rtmp.rs | 52 ----------- src/webrtc.rs | 190 -------------------------------------- 10 files changed, 48 insertions(+), 715 deletions(-) delete mode 100644 src/http.rs delete mode 100644 src/main.rs delete mode 100644 src/media.rs delete mode 100644 src/rtmp.rs delete mode 100644 src/webrtc.rs diff --git a/.gitignore b/.gitignore index 0201000..91a15d8 100644 --- a/.gitignore +++ b/.gitignore @@ -10,3 +10,4 @@ devenv.local.yaml # pre-commit .pre-commit-config.yaml +/stream.db diff --git a/Cargo.lock b/Cargo.lock index b4e82e6..1cf9c6a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -452,10 +452,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1aa79e62e7697b8e29b513a68abacf485adcd1fe8284a4316c5ae868e6633327" dependencies = [ "iana-time-zone", - "js-sys", "num-traits", "serde", - "wasm-bindgen", "windows-link", ] @@ -2643,16 +2641,12 @@ version = "0.1.0" dependencies = [ "async-broadcast", "bytes", - "chrono", "dashmap", - "entity", "http-body-util", "hyper", "hyper-util", - "migration", "rand 0.10.1", "rml_rtmp", - "sea-orm", "serde", "serde_json", "str0m", diff --git a/crates/server/Cargo.toml b/crates/server/Cargo.toml index 72994ef..7d0b185 100644 --- a/crates/server/Cargo.toml +++ b/crates/server/Cargo.toml @@ -8,13 +8,8 @@ name = "rtmp-to-whip" path = "src/main.rs" [dependencies] -entity = { path = "../entity" } -migration = { path = "../migration" } -sea-orm = { version = "1", features = ["sqlx-sqlite", "runtime-tokio-rustls", "macros"] } - async-broadcast = "0.7.2" bytes = "1.11.1" -chrono = "0.4" dashmap = "6.2.1" rand = "0.10.1" rml_rtmp = "0.8.0" diff --git a/crates/server/src/main.rs b/crates/server/src/main.rs index cde7e76..27e1fae 100644 --- a/crates/server/src/main.rs +++ b/crates/server/src/main.rs @@ -1,19 +1,11 @@ use std::{error::Error, sync::Arc}; use async_broadcast::broadcast; -use chrono::Utc; use dashmap::DashMap; -use entity::stream_key; -use entity::stream_session; -use migration::{Migrator, MigratorTrait}; use rml_rtmp::{ handshake::{Handshake, HandshakeProcessResult, PeerType}, sessions::{ServerSession, ServerSessionConfig, ServerSessionEvent, ServerSessionResult}, }; -use sea_orm::{ - ActiveModelTrait, ActiveValue::Set, ColumnTrait, Database, DatabaseConnection, EntityTrait, - QueryFilter, -}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::TcpListener, @@ -32,24 +24,18 @@ mod webrtc; pub struct AppState { pub stream_sessions: Arc>, - pub db: DatabaseConnection, } pub struct StreamSession { pub stream_key: String, pub frame_channel: async_broadcast::Sender>, - pub db_session_id: i32, } #[tokio::main] async fn main() -> Result<(), Box> { - let db = Database::connect("sqlite://stream.db?mode=rwc").await?; - Migrator::up(&db, None).await?; - let listener = TcpListener::bind("0.0.0.0:8123").await?; let appstate = Arc::new(Mutex::new(AppState { stream_sessions: Arc::new(DashMap::new()), - db: db.clone(), })); let (offer_tx, offer_rx) = tokio::sync::mpsc::channel::<(String, String)>(4); let (answer_tx, answer_rx) = tokio::sync::mpsc::channel::<(String, String)>(4); @@ -73,7 +59,6 @@ async fn main() -> Result<(), Box> { loop { let (stream, _) = listener.accept().await?; let appstate = appstate.clone(); - let db = db.clone(); tokio::spawn(async move { let mut stream = stream; @@ -123,64 +108,24 @@ async fn main() -> Result<(), Box> { } ServerSessionResult::RaisedEvent(x) => match x { ServerSessionEvent::PublishStreamFinished { stream_key, .. } => { - let removed = { - appstate - .lock() - .await - .stream_sessions - .remove(&stream_key) - }; - if let Some((_, session)) = removed { - let mut record: stream_session::ActiveModel = - stream_session::Entity::find_by_id(session.db_session_id) - .one(&db) - .await - .unwrap() - .unwrap() - .into(); - record.ended_at = Set(Some(Utc::now())); - record.update(&db).await.unwrap(); - } + appstate.lock().await.stream_sessions.remove(&stream_key); } ServerSessionEvent::PublishStreamRequested { request_id, stream_key, .. } => { - let key_record = stream_key::Entity::find() - .filter(stream_key::Column::KeyValue.eq(&stream_key)) - .filter(stream_key::Column::IsActive.eq(true)) - .one(&db) - .await - .unwrap(); - - let reply = if let Some(key_record) = key_record { - let session_record = stream_session::ActiveModel { - stream_key_id: Set(key_record.id), - started_at: Set(Utc::now()), - ended_at: Set(None), - ..Default::default() - }; - let session_record = - session_record.insert(&db).await.unwrap(); - - let session = StreamSession { - stream_key: stream_key.clone(), - frame_channel: video_channel.clone(), - db_session_id: session_record.id, - }; - appstate - .lock() - .await - .stream_sessions - .insert(stream_key.clone(), session); - - rtmp_session.accept_request(request_id).unwrap() - } else { - println!("Rejected unknown/inactive stream key: {stream_key}"); - rtmp_session.reject_request(request_id, "", "").unwrap() + let session = StreamSession { + stream_key: stream_key.clone(), + frame_channel: video_channel.clone(), }; + appstate + .lock() + .await + .stream_sessions + .insert(stream_key.clone(), session); + let reply = rtmp_session.accept_request(request_id).unwrap(); for x in reply { if let ServerSessionResult::OutboundResponse(y) = x { stream.write_all(&y.bytes).await.unwrap(); diff --git a/flake.nix b/flake.nix index 5f262bc..c9b3514 100644 --- a/flake.nix +++ b/flake.nix @@ -21,6 +21,7 @@ system: let pkgs = nixpkgs.legacyPackages.${system}; + inherit (pkgs) lib; craneLib = crane.mkLib pkgs; @@ -39,26 +40,54 @@ ]; }; - my-crate = craneLib.buildPackage ( + fileSetForCrate = + crate: + lib.fileset.toSource { + root = ./.; + fileset = lib.fileset.unions [ + ./Cargo.toml + ./Cargo.lock + (craneLib.fileset.commonCargoSources ./crates/entity) + (craneLib.fileset.commonCargoSources ./crates/migration) + (craneLib.fileset.commonCargoSources ./crates/server) + (craneLib.fileset.commonCargoSources crate) + ]; + }; + + rtmp-to-whip-simple-server = craneLib.buildPackage ( commonArgs // { + pname = "rtmp-to-whip"; + version = "0.1.0"; cargoArtifacts = craneLib.buildDepsOnly commonArgs; - - # Additional environment variables or build phases/hooks can be set - # here *without* rebuilding all dependency crates - # MY_CUSTOM_VAR = "some value"; + cargoExtraArgs = "-p server"; + src = fileSetForCrate ./crates/server; } ); + # migration = craneLib.buildPackage ( + # commonArgs + # // { + # pname = "migration"; + # version = "0.1.0"; + # cargoArtifacts = craneLib.buildDepsOnly commonArgs; + # cargoExtraArgs = "-p migration"; + # src = fileSetForCrate ./crates/migration; + # + # # Additional environment variables or build phases/hooks can be set + # # here *without* rebuilding all dependency crates + # # MY_CUSTOM_VAR = "some value"; + # } + # ); in { checks = { - inherit my-crate; + inherit rtmp-to-whip-simple-server; }; - packages.default = my-crate; + packages.default = rtmp-to-whip-simple-server; apps.default = flake-utils.lib.mkApp { - drv = my-crate; + drv = rtmp-to-whip-simple-server; }; devShells.default = craneLib.devShell { diff --git a/src/http.rs b/src/http.rs deleted file mode 100644 index 59f7f1c..0000000 --- a/src/http.rs +++ /dev/null @@ -1,69 +0,0 @@ -use std::{convert::Infallible, error::Error, net::SocketAddr}; - -use bytes::Bytes; -use http_body_util::Full; -use hyper::{Request, Response, server::conn::http2, service::service_fn}; -use hyper_util::rt::{TokioExecutor, TokioIo}; -use serde::Serialize; -use std::sync::Arc; -use tokio::{ - net::TcpListener, - sync::{ - Mutex, - mpsc::{Receiver, Sender}, - }, -}; - -use crate::AppState; - -pub struct HttpServer { - pub offer_tx: Sender<(String, String)>, - pub accept_rx: Receiver<(String, String)>, - pub appstate: Arc>, -} - -impl HttpServer { - pub fn start(&mut self) -> Result<(), Box> { - let state = self.appstate.clone(); - tokio::spawn(async move { - let socket = SocketAddr::from(([127, 0, 0, 1], 3000)); - let listener = TcpListener::bind(socket).await.unwrap(); - loop { - let (stream, _) = listener.accept().await.unwrap(); - let tokio_io = TokioIo::new(stream); - let state = state.clone(); - tokio::spawn(async move { - let service = service_fn(move |req| { - let state = state.clone(); - async move { handle_request(req, &state).await } - }); - let hyper = http2::Builder::new(TokioExecutor::new()); - hyper.serve_connection(tokio_io, service).await.unwrap(); - }); - } - }); - Ok(()) - } -} - -async fn handle_request( - request: Request, - appstate: &Arc>, -) -> Result>, Infallible> { - return match request.uri().path() { - "/api/catalog" => {} - "/api/meow" => Ok(Response::builder() - .status(200) - .body(Full::new(Bytes::from("meow"))) - .unwrap()), - _ => Ok(Response::builder() - .status(404) - .body(Full::new(Bytes::from(""))) - .unwrap()), - }; -} - -#[derive(Serialize)] -struct StreamCatalog {} - -fn catalog_from_state(state: &Arc>) -> StreamCatalog {} diff --git a/src/main.rs b/src/main.rs deleted file mode 100644 index b606016..0000000 --- a/src/main.rs +++ /dev/null @@ -1,173 +0,0 @@ -use std::{error::Error, sync::Arc}; - -use async_broadcast::broadcast; -use dashmap::DashMap; -use rml_rtmp::{ - handshake::{Handshake, HandshakeProcessResult, PeerType}, - sessions::{ServerSession, ServerSessionConfig, ServerSessionEvent, ServerSessionResult}, -}; -use tokio::{ - io::{AsyncReadExt, AsyncWriteExt}, - net::TcpListener, - sync::Mutex, -}; - -use crate::{ - http::HttpServer, - media::{H264Parser, VideoFrame}, -}; - -mod http; -mod media; -mod rtmp; -mod webrtc; - -struct AppState { - stream_sessions: Arc>, -} - -struct StreamSession { - stream_key: String, - frame_channel: async_broadcast::Sender>, -} - -#[tokio::main] -async fn main() -> Result<(), Box> { - let listener = TcpListener::bind("0.0.0.0:8123").await?; - let appstate = Arc::new(Mutex::new(AppState { - stream_sessions: Arc::new(DashMap::new()), - })); - let (offer_tx, offer_rx) = tokio::sync::mpsc::channel::<(String, String)>(4); - let (answer_tx, answer_rx) = tokio::sync::mpsc::channel::<(String, String)>(4); - - let mut http = HttpServer { - offer_tx, - accept_rx: answer_rx, - appstate: appstate.clone(), - }; - - http.start()?; - - let app = appstate.lock().await; - - let webrtc = webrtc::Webrtc { - offer_rx, - accept_tx: answer_tx, - sessions_ref: app.stream_sessions.clone(), - }; - - webrtc.start()?; - - drop(app); - - loop { - let (stream, _) = listener.accept().await?; - let appstate = appstate.clone(); - - tokio::spawn(async move { - let mut stream = stream; - let mut server = Handshake::new(PeerType::Server); - let appstate = appstate; - - let mut c0_c1: [u8; 1537] = [0; 1537]; - stream.read_exact(&mut c0_c1).await.unwrap(); - let s0_s1_s2 = server.process_bytes(&c0_c1); - let s0_s1_s2 = match s0_s1_s2 { - Ok(HandshakeProcessResult::InProgress { - response_bytes: bytes, - }) => bytes, - _ => panic!("handshake failed"), - }; - stream.write_all(&s0_s1_s2).await.unwrap(); - let mut c2 = vec![0u8; 1536]; - stream.read_exact(&mut c2).await.unwrap(); - - match server.process_bytes(&c2[..]) { - Ok(HandshakeProcessResult::Completed { .. }) => {} - Ok(HandshakeProcessResult::InProgress { - response_bytes: meow, - }) => stream.write_all(&meow).await.unwrap(), - x => panic!("Unexpected process_bytes response: {:?}", x), - } - - let config = ServerSessionConfig::new(); - let (mut rtmp_session, bytes) = ServerSession::new(config).unwrap(); - for x in bytes { - if let ServerSessionResult::OutboundResponse(packet) = x { - stream.write_all(&packet.bytes).await.unwrap(); - } - } - let (mut video_channel, _video_rx) = broadcast::>(32); - video_channel.set_overflow(true); - let mut parser = H264Parser::new(); - - loop { - let mut buf: [u8; 4096] = [0; 4096]; - let n = stream.read(&mut buf).await.unwrap(); - - let obs_events = rtmp_session.handle_input(&buf[..n]).unwrap(); - for event in obs_events { - match event { - ServerSessionResult::OutboundResponse(packet) => { - stream.write_all(&packet.bytes).await.unwrap(); - } - ServerSessionResult::RaisedEvent(x) => match x { - ServerSessionEvent::PublishStreamFinished { .. } => { - appstate.lock().await.stream_sessions.remove("test"); - } - ServerSessionEvent::PublishStreamRequested { - request_id, - stream_key, - .. - } => { - let mut reply = rtmp_session.accept_request(request_id).unwrap(); - if stream_key != "test" { - reply = - rtmp_session.reject_request(request_id, "", "").unwrap(); - } - - let session = StreamSession { - stream_key: stream_key.clone(), - frame_channel: video_channel.clone(), - }; - - appstate - .lock() - .await - .stream_sessions - .insert(stream_key.clone(), session); - - for x in reply { - if let ServerSessionResult::OutboundResponse(y) = x { - stream.write_all(&y.bytes).await.unwrap(); - } - } - } - ServerSessionEvent::ConnectionRequested { request_id, .. } => { - let reply = rtmp_session.accept_request(request_id).unwrap(); - for x in reply { - if let ServerSessionResult::OutboundResponse(y) = x { - stream.write_all(&y.bytes).await.unwrap(); - } - } - } - ServerSessionEvent::VideoDataReceived { - data, timestamp, .. - } => { - if data.len() >= 5 && &data[1..5] == b"hvc1" { - println!("HEVC/H.265 not supported, closing connection"); - return; - } - if let Some(parsed_frame) = parser.parse(&data, timestamp.value) { - video_channel.broadcast(Arc::new(parsed_frame)).await.ok(); - } - } - _ => {} - }, - _ => {} - } - } - } - }); - } -} diff --git a/src/media.rs b/src/media.rs deleted file mode 100644 index 98ef41e..0000000 --- a/src/media.rs +++ /dev/null @@ -1,147 +0,0 @@ -// Claude slop... im not skilled amount to do this bullshit - -use bytes::Bytes; - -pub struct VideoFrame { - pub data: Bytes, - pub is_keyframe: bool, - pub timestamp_ms: u32, -} - -pub struct AudioFrame { - pub data: Bytes, - pub timestamp_ms: u32, -} - -pub struct H264Parser { - sps: Option>, - pps: Option>, -} - -impl H264Parser { - pub fn new() -> Self { - Self { - sps: None, - pps: None, - } - } - - /// Parse an RTMP VideoDataReceived payload. Returns None for sequence - /// header packets (which carry SPS/PPS but no displayable frame). - pub fn parse(&mut self, bytes: &[u8], timestamp_ms: u32) -> Option { - if bytes.len() < 5 { - return None; - } - - let frame_type = (bytes[0] >> 4) & 0x0F; - let codec_id = bytes[0] & 0x0F; - - if codec_id != 7 { - return None; // not H.264 - } - - let avc_packet_type = bytes[1]; - // bytes[2..5] are the composition time offset — not needed for sending - let payload = &bytes[5..]; - - match avc_packet_type { - 0 => { - self.parse_sequence_header(payload); - None - } - 1 => { - let is_keyframe = frame_type == 1; - let data = self.avcc_to_annexb(payload, is_keyframe)?; - Some(VideoFrame { - data: Bytes::from(data), - is_keyframe, - timestamp_ms, - }) - } - _ => None, - } - } - - fn parse_sequence_header(&mut self, payload: &[u8]) { - // AVCDecoderConfigurationRecord layout: - // [0] configurationVersion - // [1] AVCProfileIndication - // [2] profile_compatibility - // [3] AVCLevelIndication - // [4] 0xFF (lower 2 bits = lengthSizeMinusOne, always 3 meaning 4-byte lengths) - // [5] 0xE0 | numSPS - // [6..] SPS entries: 2-byte length + bytes - // then: numPPS, PPS entries: 2-byte length + bytes - if payload.len() < 7 { - return; - } - - let mut i = 5; - - let num_sps = (payload[i] & 0x1F) as usize; - i += 1; - - for _ in 0..num_sps { - if i + 2 > payload.len() { - return; - } - let len = u16::from_be_bytes([payload[i], payload[i + 1]]) as usize; - i += 2; - if i + len > payload.len() { - return; - } - self.sps = Some(payload[i..i + len].to_vec()); - i += len; - } - - if i >= payload.len() { - return; - } - - let num_pps = payload[i] as usize; - i += 1; - - for _ in 0..num_pps { - if i + 2 > payload.len() { - return; - } - let len = u16::from_be_bytes([payload[i], payload[i + 1]]) as usize; - i += 2; - if i + len > payload.len() { - return; - } - self.pps = Some(payload[i..i + len].to_vec()); - i += len; - } - } - - fn avcc_to_annexb(&self, payload: &[u8], is_keyframe: bool) -> Option> { - let mut out = Vec::new(); - - // Prepend SPS+PPS before every keyframe so str0m's packetizer - // can bundle them into a STAP-A alongside the IDR NALU. - if is_keyframe { - if let (Some(sps), Some(pps)) = (&self.sps, &self.pps) { - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(sps); - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(pps); - } - } - - // Convert each length-prefixed NALU to an Annex B start-code NALU. - let mut i = 0; - while i + 4 <= payload.len() { - let nalu_len = u32::from_be_bytes(payload[i..i + 4].try_into().unwrap()) as usize; - i += 4; - if i + nalu_len > payload.len() { - break; - } - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(&payload[i..i + nalu_len]); - i += nalu_len; - } - - if out.is_empty() { None } else { Some(out) } - } -} diff --git a/src/rtmp.rs b/src/rtmp.rs deleted file mode 100644 index 98ac91f..0000000 --- a/src/rtmp.rs +++ /dev/null @@ -1,52 +0,0 @@ -#[derive(Default, Debug)] -pub struct RtmpSession { - client: RtmpClient, -} - -#[derive(Debug)] -struct RtmpClient { - version: u8, - timestamp: u32, - magic_bytes: [u8; 1536], -} - -impl Default for RtmpClient { - fn default() -> Self { - Self { - magic_bytes: [0; 1536], - version: 0, - timestamp: 0, - } - } -} - -impl RtmpSession { - pub fn consume_handshake(&mut self, bytes: &[u8]) -> Result<(), ()> { - let version = bytes[0]; - let timestamp = u32::from_be_bytes(bytes[1..5].try_into().unwrap()); - let mut random: [u8; 1536] = [0; 1536]; - random.copy_from_slice(&bytes[1..1537]); - - println!("size of rand: {}", random.len()); - - let client = RtmpClient { - version, - timestamp, - magic_bytes: random, - }; - - self.client = client; - - Ok(()) - } - pub fn response(&self) -> [u8; 1537] { - let mut reply: [u8; 1537] = [0; 1537]; - reply[0] = 3; - // reply[1..1537].copy_from_slice(&self.client.magic_bytes); - let mut rand: [u8; 1536] = [0; 1536]; - rand.fill(1); - reply[1..1537].copy_from_slice(&rand); - - reply - } -} diff --git a/src/webrtc.rs b/src/webrtc.rs deleted file mode 100644 index 0b921cf..0000000 --- a/src/webrtc.rs +++ /dev/null @@ -1,190 +0,0 @@ -use dashmap::DashMap; -use tokio::{ - net::UdpSocket, - sync::mpsc::{Receiver, Sender}, -}; -use std::{ - error::Error, - sync::Arc, - time::{Duration, Instant}, -}; - -use str0m::{ - Candidate, Event, Input, Output, Rtc, - change::SdpOffer, - media::{MediaKind, MediaTime, Mid}, - net::{Protocol, Receive}, -}; - -use crate::{StreamSession, media::VideoFrame}; - -pub struct Webrtc { - pub offer_rx: Receiver<(String, String)>, - pub accept_tx: Sender<(String, String)>, - pub sessions_ref: Arc>, -} - -impl Webrtc { - pub fn start(mut self) -> Result<(), Box> { - tokio::spawn(async move { - while let Some(offer) = self.offer_rx.recv().await { - let (stream_key, sdp_body) = offer; - let socket = UdpSocket::bind("127.0.0.1:0").await.unwrap(); - let local_addr = socket.local_addr().unwrap(); - - let mut builder = Rtc::builder(); - { - let cc = builder.codec_config(); - cc.enable_h264(false); - cc.add_h264(102.into(), None, true, 0x42e01f); - cc.add_h264(104.into(), None, true, 0x4d001f); - cc.add_h264(106.into(), None, true, 0x64001f); - } - let mut rtc = builder.build(Instant::now()); - - let candidate = Candidate::host(local_addr, Protocol::Udp).unwrap(); - rtc.add_local_candidate(candidate); - - let offer_sdp = SdpOffer::from_sdp_string(&sdp_body).unwrap(); - let mut changes = rtc.sdp_api(); - let mid = changes.add_media( - MediaKind::Video, - str0m::media::Direction::SendOnly, - Some(stream_key.clone()), - Some("video0".to_string()), - None, - ); - let offer_answer = match changes.accept_offer(offer_sdp) { - Ok(a) => a, - Err(e) => { - println!("accept_offer failed: {:?}", e); - continue; - } - }; - let answer_sdp = offer_answer.to_sdp_string(); - - self.accept_tx - .send((stream_key.clone(), answer_sdp)) - .await - .unwrap(); - - let sessions_ref = self.sessions_ref.clone(); - tokio::spawn(async move { - Webrtc::detach_connection(socket, rtc, sessions_ref, mid).await; - }); - } - }); - Ok(()) - } - - async fn detach_connection( - socket: UdpSocket, - mut rtc: Rtc, - sessions_ref: Arc>, - _hint_mid: Mid, - ) { - let mut video_mid: Option = None; - let mut video_pt = None; - let mut connected = false; - let mut video_stream: Option>> = None; - - let mut recv_buf = vec![0u8; 65535]; - let local_addr = socket.local_addr().unwrap(); - - loop { - let deadline = loop { - match rtc.poll_output() { - Ok(Output::Timeout(t)) => break t, - Ok(Output::Transmit(t)) => { - if socket.send_to(&t.contents, t.destination).await.is_err() { - return; - } - } - Ok(Output::Event(e)) => match e { - Event::MediaAdded(ma) => { - if ma.kind == MediaKind::Video { - if let Some(writer) = rtc.writer(ma.mid) { - let best = writer - .payload_params() - .max_by_key(|p| { - p.spec().format.profile_level_id.unwrap_or(0) - }); - if let Some(params) = best { - println!("Selected PT {:?}", params.pt()); - video_pt = Some(params.pt()); - video_mid = Some(ma.mid); - } - } - } - } - Event::IceConnectionStateChange(state) => { - println!("ICE state: {:?}", state); - } - Event::Connected => { - println!("DTLS+ICE connected, ready for media"); - connected = true; - } - _ => {} - }, - Err(_) => return, - } - }; - - if connected { - if video_stream.is_none() { - if let Some(session) = sessions_ref.get("test") { - video_stream = Some(session.frame_channel.new_receiver()); - } - } - if let Some(ref mut stream) = video_stream { - for _ in 0..8 { - match stream.try_recv() { - Ok(frame) => { - let now = Instant::now(); - let rtp_time = MediaTime::from_90khz(frame.timestamp_ms as u64 * 90); - if let (Some(pt), Some(writer)) = - (video_pt, video_mid.and_then(|m| rtc.writer(m))) - { - if let Err(e) = - writer.write(pt, now, rtp_time, frame.data.to_vec()) - { - println!("write error: {:?}", e); - } - } - } - Err(async_broadcast::TryRecvError::Empty) => break, - Err(async_broadcast::TryRecvError::Closed) => return, - Err(async_broadcast::TryRecvError::Overflowed(_)) => continue, - } - } - } - } - - let wait_until = deadline.min(Instant::now() + Duration::from_millis(20)).max(Instant::now()); - let sleep = tokio::time::sleep_until(wait_until.into()); - - tokio::select! { - _ = sleep => { - rtc.handle_input(Input::Timeout(Instant::now())).ok(); - } - result = socket.recv_from(&mut recv_buf) => { - if let Ok((n, from)) = result { - let data = recv_buf[..n].to_vec(); - if let Ok(contents) = data.as_slice().try_into() { - rtc.handle_input(Input::Receive( - Instant::now(), - Receive { - proto: Protocol::Udp, - source: from, - destination: local_addr, - contents, - }, - )) - .ok(); - } - } - } - } - } - } -}