more reshape that i missed

This commit is contained in:
2026-06-17 11:55:27 +01:00
parent b2ad096f95
commit cbfa07f8e8
10 changed files with 48 additions and 715 deletions
+1
View File
@@ -10,3 +10,4 @@ devenv.local.yaml
# pre-commit # pre-commit
.pre-commit-config.yaml .pre-commit-config.yaml
/stream.db
Generated
-6
View File
@@ -452,10 +452,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1aa79e62e7697b8e29b513a68abacf485adcd1fe8284a4316c5ae868e6633327" checksum = "1aa79e62e7697b8e29b513a68abacf485adcd1fe8284a4316c5ae868e6633327"
dependencies = [ dependencies = [
"iana-time-zone", "iana-time-zone",
"js-sys",
"num-traits", "num-traits",
"serde", "serde",
"wasm-bindgen",
"windows-link", "windows-link",
] ]
@@ -2643,16 +2641,12 @@ version = "0.1.0"
dependencies = [ dependencies = [
"async-broadcast", "async-broadcast",
"bytes", "bytes",
"chrono",
"dashmap", "dashmap",
"entity",
"http-body-util", "http-body-util",
"hyper", "hyper",
"hyper-util", "hyper-util",
"migration",
"rand 0.10.1", "rand 0.10.1",
"rml_rtmp", "rml_rtmp",
"sea-orm",
"serde", "serde",
"serde_json", "serde_json",
"str0m", "str0m",
-5
View File
@@ -8,13 +8,8 @@ name = "rtmp-to-whip"
path = "src/main.rs" path = "src/main.rs"
[dependencies] [dependencies]
entity = { path = "../entity" }
migration = { path = "../migration" }
sea-orm = { version = "1", features = ["sqlx-sqlite", "runtime-tokio-rustls", "macros"] }
async-broadcast = "0.7.2" async-broadcast = "0.7.2"
bytes = "1.11.1" bytes = "1.11.1"
chrono = "0.4"
dashmap = "6.2.1" dashmap = "6.2.1"
rand = "0.10.1" rand = "0.10.1"
rml_rtmp = "0.8.0" rml_rtmp = "0.8.0"
+10 -65
View File
@@ -1,19 +1,11 @@
use std::{error::Error, sync::Arc}; use std::{error::Error, sync::Arc};
use async_broadcast::broadcast; use async_broadcast::broadcast;
use chrono::Utc;
use dashmap::DashMap; use dashmap::DashMap;
use entity::stream_key;
use entity::stream_session;
use migration::{Migrator, MigratorTrait};
use rml_rtmp::{ use rml_rtmp::{
handshake::{Handshake, HandshakeProcessResult, PeerType}, handshake::{Handshake, HandshakeProcessResult, PeerType},
sessions::{ServerSession, ServerSessionConfig, ServerSessionEvent, ServerSessionResult}, sessions::{ServerSession, ServerSessionConfig, ServerSessionEvent, ServerSessionResult},
}; };
use sea_orm::{
ActiveModelTrait, ActiveValue::Set, ColumnTrait, Database, DatabaseConnection, EntityTrait,
QueryFilter,
};
use tokio::{ use tokio::{
io::{AsyncReadExt, AsyncWriteExt}, io::{AsyncReadExt, AsyncWriteExt},
net::TcpListener, net::TcpListener,
@@ -32,24 +24,18 @@ mod webrtc;
pub struct AppState { pub struct AppState {
pub stream_sessions: Arc<DashMap<String, StreamSession>>, pub stream_sessions: Arc<DashMap<String, StreamSession>>,
pub db: DatabaseConnection,
} }
pub struct StreamSession { pub struct StreamSession {
pub stream_key: String, pub stream_key: String,
pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>, pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>,
pub db_session_id: i32,
} }
#[tokio::main] #[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> { async fn main() -> Result<(), Box<dyn Error>> {
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 listener = TcpListener::bind("0.0.0.0:8123").await?;
let appstate = Arc::new(Mutex::new(AppState { let appstate = Arc::new(Mutex::new(AppState {
stream_sessions: Arc::new(DashMap::new()), stream_sessions: Arc::new(DashMap::new()),
db: db.clone(),
})); }));
let (offer_tx, offer_rx) = tokio::sync::mpsc::channel::<(String, String)>(4); 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 (answer_tx, answer_rx) = tokio::sync::mpsc::channel::<(String, String)>(4);
@@ -73,7 +59,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
loop { loop {
let (stream, _) = listener.accept().await?; let (stream, _) = listener.accept().await?;
let appstate = appstate.clone(); let appstate = appstate.clone();
let db = db.clone();
tokio::spawn(async move { tokio::spawn(async move {
let mut stream = stream; let mut stream = stream;
@@ -123,64 +108,24 @@ async fn main() -> Result<(), Box<dyn Error>> {
} }
ServerSessionResult::RaisedEvent(x) => match x { ServerSessionResult::RaisedEvent(x) => match x {
ServerSessionEvent::PublishStreamFinished { stream_key, .. } => { ServerSessionEvent::PublishStreamFinished { stream_key, .. } => {
let removed = { appstate.lock().await.stream_sessions.remove(&stream_key);
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();
}
} }
ServerSessionEvent::PublishStreamRequested { ServerSessionEvent::PublishStreamRequested {
request_id, request_id,
stream_key, stream_key,
.. ..
} => { } => {
let key_record = stream_key::Entity::find() let session = StreamSession {
.filter(stream_key::Column::KeyValue.eq(&stream_key)) stream_key: stream_key.clone(),
.filter(stream_key::Column::IsActive.eq(true)) frame_channel: video_channel.clone(),
.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()
}; };
appstate
.lock()
.await
.stream_sessions
.insert(stream_key.clone(), session);
let reply = rtmp_session.accept_request(request_id).unwrap();
for x in reply { for x in reply {
if let ServerSessionResult::OutboundResponse(y) = x { if let ServerSessionResult::OutboundResponse(y) = x {
stream.write_all(&y.bytes).await.unwrap(); stream.write_all(&y.bytes).await.unwrap();
+37 -8
View File
@@ -21,6 +21,7 @@
system: system:
let let
pkgs = nixpkgs.legacyPackages.${system}; pkgs = nixpkgs.legacyPackages.${system};
inherit (pkgs) lib;
craneLib = crane.mkLib pkgs; 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 commonArgs
// { // {
pname = "rtmp-to-whip";
version = "0.1.0";
cargoArtifacts = craneLib.buildDepsOnly commonArgs; cargoArtifacts = craneLib.buildDepsOnly commonArgs;
cargoExtraArgs = "-p server";
# Additional environment variables or build phases/hooks can be set src = fileSetForCrate ./crates/server;
# here *without* rebuilding all dependency crates
# MY_CUSTOM_VAR = "some value";
} }
); );
# 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 in
{ {
checks = { 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 { apps.default = flake-utils.lib.mkApp {
drv = my-crate; drv = rtmp-to-whip-simple-server;
}; };
devShells.default = craneLib.devShell { devShells.default = craneLib.devShell {
-69
View File
@@ -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<Mutex<AppState>>,
}
impl HttpServer {
pub fn start(&mut self) -> Result<(), Box<dyn Error>> {
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<impl hyper::body::Body>,
appstate: &Arc<Mutex<AppState>>,
) -> Result<Response<Full<Bytes>>, 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<Mutex<AppState>>) -> StreamCatalog {}
-173
View File
@@ -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<DashMap<String, StreamSession>>,
}
struct StreamSession {
stream_key: String,
frame_channel: async_broadcast::Sender<Arc<VideoFrame>>,
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
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::<Arc<VideoFrame>>(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();
}
}
_ => {}
},
_ => {}
}
}
}
});
}
}
-147
View File
@@ -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<Vec<u8>>,
pps: Option<Vec<u8>>,
}
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<VideoFrame> {
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<Vec<u8>> {
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) }
}
}
-52
View File
@@ -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
}
}
-190
View File
@@ -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<DashMap<String, StreamSession>>,
}
impl Webrtc {
pub fn start(mut self) -> Result<(), Box<dyn Error>> {
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<DashMap<String, StreamSession>>,
_hint_mid: Mid,
) {
let mut video_mid: Option<Mid> = None;
let mut video_pt = None;
let mut connected = false;
let mut video_stream: Option<async_broadcast::Receiver<Arc<VideoFrame>>> = 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();
}
}
}
}
}
}
}