workspace reshape
This commit is contained in:
Generated
+2230
-38
File diff suppressed because it is too large
Load Diff
+3
-18
@@ -1,18 +1,3 @@
|
|||||||
[package]
|
[workspace]
|
||||||
name = "rtmp-to-whip-simple-server"
|
members = ["crates/server", "crates/entity", "crates/migration"]
|
||||||
version = "0.1.0"
|
resolver = "2"
|
||||||
edition = "2024"
|
|
||||||
|
|
||||||
[dependencies]
|
|
||||||
async-broadcast = "0.7.2"
|
|
||||||
bytes = "1.11.1"
|
|
||||||
dashmap = "6.2.1"
|
|
||||||
rand = "0.10.1"
|
|
||||||
rml_rtmp = "0.8.0"
|
|
||||||
str0m = "0.20.0"
|
|
||||||
tokio = { version = "1", features = ["full"] }
|
|
||||||
hyper = { version = "1", features = ["full"] }
|
|
||||||
http-body-util = "0.1"
|
|
||||||
hyper-util = { version = "0.1", features = ["full"] }
|
|
||||||
serde = { version = "1.0.228", features = ["serde_derive"] }
|
|
||||||
serde_json = "1.0.150"
|
|
||||||
|
|||||||
@@ -0,0 +1,8 @@
|
|||||||
|
[package]
|
||||||
|
name = "entity"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2024"
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
sea-orm = { version = "1", features = ["sqlx-sqlite", "runtime-tokio-rustls", "macros"] }
|
||||||
|
serde = { version = "1", features = ["derive"] }
|
||||||
@@ -0,0 +1,3 @@
|
|||||||
|
pub mod prelude;
|
||||||
|
pub mod stream_key;
|
||||||
|
pub mod stream_session;
|
||||||
@@ -0,0 +1,2 @@
|
|||||||
|
pub use super::stream_key::Entity as StreamKey;
|
||||||
|
pub use super::stream_session::Entity as StreamSession;
|
||||||
@@ -0,0 +1,28 @@
|
|||||||
|
use sea_orm::entity::prelude::*;
|
||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Serialize, Deserialize)]
|
||||||
|
#[sea_orm(table_name = "stream_key")]
|
||||||
|
pub struct Model {
|
||||||
|
#[sea_orm(primary_key)]
|
||||||
|
pub id: i32,
|
||||||
|
#[sea_orm(unique)]
|
||||||
|
pub key_value: String,
|
||||||
|
pub label: String,
|
||||||
|
pub is_active: bool,
|
||||||
|
pub created_at: DateTimeUtc,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
||||||
|
pub enum Relation {
|
||||||
|
#[sea_orm(has_many = "super::stream_session::Entity")]
|
||||||
|
StreamSession,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Related<super::stream_session::Entity> for Entity {
|
||||||
|
fn to() -> RelationDef {
|
||||||
|
Relation::StreamSession.def()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ActiveModelBehavior for ActiveModel {}
|
||||||
@@ -0,0 +1,30 @@
|
|||||||
|
use sea_orm::entity::prelude::*;
|
||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Serialize, Deserialize)]
|
||||||
|
#[sea_orm(table_name = "stream_session")]
|
||||||
|
pub struct Model {
|
||||||
|
#[sea_orm(primary_key)]
|
||||||
|
pub id: i32,
|
||||||
|
pub stream_key_id: i32,
|
||||||
|
pub started_at: DateTimeUtc,
|
||||||
|
pub ended_at: Option<DateTimeUtc>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
||||||
|
pub enum Relation {
|
||||||
|
#[sea_orm(
|
||||||
|
belongs_to = "super::stream_key::Entity",
|
||||||
|
from = "Column::StreamKeyId",
|
||||||
|
to = "super::stream_key::Column::Id"
|
||||||
|
)]
|
||||||
|
StreamKey,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Related<super::stream_key::Entity> for Entity {
|
||||||
|
fn to() -> RelationDef {
|
||||||
|
Relation::StreamKey.def()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ActiveModelBehavior for ActiveModel {}
|
||||||
@@ -0,0 +1,16 @@
|
|||||||
|
[package]
|
||||||
|
name = "migration"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2024"
|
||||||
|
|
||||||
|
[lib]
|
||||||
|
name = "migration"
|
||||||
|
path = "src/lib.rs"
|
||||||
|
|
||||||
|
[[bin]]
|
||||||
|
name = "migration"
|
||||||
|
path = "src/main.rs"
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
sea-orm-migration = { version = "1", features = ["runtime-tokio-rustls", "sqlx-sqlite"] }
|
||||||
|
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
pub use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
mod m20260616_000001_create_tables;
|
||||||
|
|
||||||
|
pub struct Migrator;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigratorTrait for Migrator {
|
||||||
|
fn migrations() -> Vec<Box<dyn MigrationTrait>> {
|
||||||
|
vec![Box::new(m20260616_000001_create_tables::Migration)]
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,103 @@
|
|||||||
|
use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.create_table(
|
||||||
|
Table::create()
|
||||||
|
.table(StreamKey::Table)
|
||||||
|
.if_not_exists()
|
||||||
|
.col(
|
||||||
|
ColumnDef::new(StreamKey::Id)
|
||||||
|
.integer()
|
||||||
|
.not_null()
|
||||||
|
.auto_increment()
|
||||||
|
.primary_key(),
|
||||||
|
)
|
||||||
|
.col(
|
||||||
|
ColumnDef::new(StreamKey::KeyValue)
|
||||||
|
.string()
|
||||||
|
.not_null()
|
||||||
|
.unique_key(),
|
||||||
|
)
|
||||||
|
.col(ColumnDef::new(StreamKey::Label).string().not_null())
|
||||||
|
.col(
|
||||||
|
ColumnDef::new(StreamKey::IsActive)
|
||||||
|
.boolean()
|
||||||
|
.not_null()
|
||||||
|
.default(true),
|
||||||
|
)
|
||||||
|
.col(ColumnDef::new(StreamKey::CreatedAt).date_time().not_null())
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
manager
|
||||||
|
.create_table(
|
||||||
|
Table::create()
|
||||||
|
.table(StreamSession::Table)
|
||||||
|
.if_not_exists()
|
||||||
|
.col(
|
||||||
|
ColumnDef::new(StreamSession::Id)
|
||||||
|
.integer()
|
||||||
|
.not_null()
|
||||||
|
.auto_increment()
|
||||||
|
.primary_key(),
|
||||||
|
)
|
||||||
|
.col(ColumnDef::new(StreamSession::StreamKeyId).integer().not_null())
|
||||||
|
.col(ColumnDef::new(StreamSession::StartedAt).date_time().not_null())
|
||||||
|
.col(ColumnDef::new(StreamSession::EndedAt).date_time().null())
|
||||||
|
.foreign_key(
|
||||||
|
ForeignKey::create()
|
||||||
|
.from(StreamSession::Table, StreamSession::StreamKeyId)
|
||||||
|
.to(StreamKey::Table, StreamKey::Id),
|
||||||
|
)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
// Seed the default stream key so OBS can connect out of the box
|
||||||
|
manager
|
||||||
|
.get_connection()
|
||||||
|
.execute_unprepared(
|
||||||
|
"INSERT OR IGNORE INTO stream_key (key_value, label, is_active, created_at) \
|
||||||
|
VALUES ('test', 'Default', 1, datetime('now'))",
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.drop_table(Table::drop().table(StreamSession::Table).to_owned())
|
||||||
|
.await?;
|
||||||
|
manager
|
||||||
|
.drop_table(Table::drop().table(StreamKey::Table).to_owned())
|
||||||
|
.await?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Iden)]
|
||||||
|
enum StreamKey {
|
||||||
|
Table,
|
||||||
|
Id,
|
||||||
|
KeyValue,
|
||||||
|
Label,
|
||||||
|
IsActive,
|
||||||
|
CreatedAt,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Iden)]
|
||||||
|
enum StreamSession {
|
||||||
|
Table,
|
||||||
|
Id,
|
||||||
|
StreamKeyId,
|
||||||
|
StartedAt,
|
||||||
|
EndedAt,
|
||||||
|
}
|
||||||
@@ -0,0 +1,6 @@
|
|||||||
|
use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
#[tokio::main]
|
||||||
|
async fn main() {
|
||||||
|
cli::run_cli(migration::Migrator).await;
|
||||||
|
}
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
[package]
|
||||||
|
name = "server"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2024"
|
||||||
|
|
||||||
|
[[bin]]
|
||||||
|
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"
|
||||||
|
str0m = "0.20.0"
|
||||||
|
tokio = { version = "1", features = ["full"] }
|
||||||
|
hyper = { version = "1", features = ["full"] }
|
||||||
|
http-body-util = "0.1"
|
||||||
|
hyper-util = { version = "0.1", features = ["full"] }
|
||||||
|
serde = { version = "1.0.228", features = ["serde_derive"] }
|
||||||
|
serde_json = "1.0.150"
|
||||||
@@ -0,0 +1,86 @@
|
|||||||
|
use std::{convert::Infallible, error::Error, net::SocketAddr, sync::Arc};
|
||||||
|
|
||||||
|
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 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> {
|
||||||
|
match request.uri().path() {
|
||||||
|
"/api/catalog" => {
|
||||||
|
let catalog = catalog_from_state(appstate).await;
|
||||||
|
let body = serde_json::to_vec(&catalog).unwrap();
|
||||||
|
Ok(Response::builder()
|
||||||
|
.status(200)
|
||||||
|
.header("content-type", "application/json")
|
||||||
|
.body(Full::new(Bytes::from(body)))
|
||||||
|
.unwrap())
|
||||||
|
}
|
||||||
|
"/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 {
|
||||||
|
active_streams: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn catalog_from_state(state: &Arc<Mutex<AppState>>) -> StreamCatalog {
|
||||||
|
let app = state.lock().await;
|
||||||
|
let active_streams = app
|
||||||
|
.stream_sessions
|
||||||
|
.iter()
|
||||||
|
.map(|e| e.key().clone())
|
||||||
|
.collect();
|
||||||
|
StreamCatalog { active_streams }
|
||||||
|
}
|
||||||
@@ -0,0 +1,217 @@
|
|||||||
|
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,
|
||||||
|
sync::Mutex,
|
||||||
|
};
|
||||||
|
|
||||||
|
use crate::{
|
||||||
|
http::HttpServer,
|
||||||
|
media::{H264Parser, VideoFrame},
|
||||||
|
};
|
||||||
|
|
||||||
|
mod http;
|
||||||
|
mod media;
|
||||||
|
mod rtmp;
|
||||||
|
mod webrtc;
|
||||||
|
|
||||||
|
pub struct AppState {
|
||||||
|
pub stream_sessions: Arc<DashMap<String, StreamSession>>,
|
||||||
|
pub db: DatabaseConnection,
|
||||||
|
}
|
||||||
|
|
||||||
|
pub struct StreamSession {
|
||||||
|
pub stream_key: String,
|
||||||
|
pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>,
|
||||||
|
pub db_session_id: i32,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::main]
|
||||||
|
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 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);
|
||||||
|
|
||||||
|
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();
|
||||||
|
let db = db.clone();
|
||||||
|
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let mut stream = stream;
|
||||||
|
let mut server = Handshake::new(PeerType::Server);
|
||||||
|
|
||||||
|
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 { 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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
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()
|
||||||
|
};
|
||||||
|
|
||||||
|
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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ => {}
|
||||||
|
},
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,147 @@
|
|||||||
|
// 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) }
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
#[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
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,190 @@
|
|||||||
|
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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user