Compare commits
16
Commits
aa1d4e8f67
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a9bd890d16
|
||
|
|
21a20b898c
|
||
|
|
6ac7e60a34
|
||
|
|
94daa39d87
|
||
|
|
c4a7e622b7
|
||
|
|
98dbb3d36d
|
||
|
|
0a039a97c4
|
||
|
|
f855e0e471
|
||
|
|
72cfb26e38
|
||
|
|
80368b3ea5
|
||
|
|
0c6705badd
|
||
|
|
9960434c1c
|
||
|
|
c58ca27490
|
||
|
|
98febe4bf3
|
||
|
|
f0a40efc25
|
||
|
|
7d9b3b6285
|
Generated
+4
-20
@@ -1151,7 +1151,6 @@ checksum = "8b147ee9d1f6d097cef9ce628cd2ee62288d963e16fb287bd9286455b241382d"
|
|||||||
dependencies = [
|
dependencies = [
|
||||||
"futures-channel",
|
"futures-channel",
|
||||||
"futures-core",
|
"futures-core",
|
||||||
"futures-executor",
|
|
||||||
"futures-io",
|
"futures-io",
|
||||||
"futures-sink",
|
"futures-sink",
|
||||||
"futures-task",
|
"futures-task",
|
||||||
@@ -1202,17 +1201,6 @@ version = "0.3.32"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "cecba35d7ad927e23624b22ad55235f2239cfa44fd10428eecbeba6d6a717718"
|
checksum = "cecba35d7ad927e23624b22ad55235f2239cfa44fd10428eecbeba6d6a717718"
|
||||||
|
|
||||||
[[package]]
|
|
||||||
name = "futures-macro"
|
|
||||||
version = "0.3.32"
|
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
|
||||||
checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b"
|
|
||||||
dependencies = [
|
|
||||||
"proc-macro2",
|
|
||||||
"quote",
|
|
||||||
"syn 2.0.119",
|
|
||||||
]
|
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "futures-sink"
|
name = "futures-sink"
|
||||||
version = "0.3.32"
|
version = "0.3.32"
|
||||||
@@ -1231,10 +1219,8 @@ version = "0.3.32"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6"
|
checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"futures-channel",
|
|
||||||
"futures-core",
|
"futures-core",
|
||||||
"futures-io",
|
"futures-io",
|
||||||
"futures-macro",
|
|
||||||
"futures-sink",
|
"futures-sink",
|
||||||
"futures-task",
|
"futures-task",
|
||||||
"memchr",
|
"memchr",
|
||||||
@@ -2870,9 +2856,9 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "serde_json"
|
name = "serde_json"
|
||||||
version = "1.0.150"
|
version = "1.0.151"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9"
|
checksum = "c841b55ecdae098c80dcae9cf767f6f8a0c2cdb3416bbef72181df4d0fe73f14"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"itoa",
|
"itoa",
|
||||||
"memchr",
|
"memchr",
|
||||||
@@ -2906,7 +2892,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "server"
|
name = "server"
|
||||||
version = "0.4.0"
|
version = "0.6.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"argon2",
|
"argon2",
|
||||||
"async-broadcast",
|
"async-broadcast",
|
||||||
@@ -2916,11 +2902,8 @@ dependencies = [
|
|||||||
"chrono",
|
"chrono",
|
||||||
"dashmap",
|
"dashmap",
|
||||||
"entity",
|
"entity",
|
||||||
"futures",
|
|
||||||
"migration",
|
"migration",
|
||||||
"opus",
|
"opus",
|
||||||
"rand 0.10.2",
|
|
||||||
"regex",
|
|
||||||
"rml_rtmp",
|
"rml_rtmp",
|
||||||
"rubato",
|
"rubato",
|
||||||
"sea-orm",
|
"sea-orm",
|
||||||
@@ -2930,6 +2913,7 @@ dependencies = [
|
|||||||
"symphonia",
|
"symphonia",
|
||||||
"sysinfo",
|
"sysinfo",
|
||||||
"thiserror 2.0.18",
|
"thiserror 2.0.18",
|
||||||
|
"time",
|
||||||
"tokio",
|
"tokio",
|
||||||
"tower-http",
|
"tower-http",
|
||||||
"tracing",
|
"tracing",
|
||||||
|
|||||||
+1
-1
@@ -43,6 +43,6 @@ EXPOSE 1935
|
|||||||
EXPOSE 3000
|
EXPOSE 3000
|
||||||
EXPOSE 6969/udp
|
EXPOSE 6969/udp
|
||||||
|
|
||||||
ENV RUST_LOG="info,warn"
|
ENV RUST_LOG="info"
|
||||||
|
|
||||||
CMD ["/rtmp-to-whip"]
|
CMD ["/rtmp-to-whip"]
|
||||||
|
|||||||
@@ -55,4 +55,6 @@ EXPOSE 1935
|
|||||||
EXPOSE 3000
|
EXPOSE 3000
|
||||||
EXPOSE 6969/udp
|
EXPOSE 6969/udp
|
||||||
|
|
||||||
|
ENV RUST_LOG="info"
|
||||||
|
|
||||||
CMD ["/rtmp-to-whip"]
|
CMD ["/rtmp-to-whip"]
|
||||||
|
|||||||
@@ -0,0 +1,36 @@
|
|||||||
|
# TODO
|
||||||
|
|
||||||
|
Audit findings from ponytail-audit (ranked biggest cut first). Net: ~-300 lines, -5 deps possible.
|
||||||
|
|
||||||
|
## Deletions
|
||||||
|
|
||||||
|
- [ ] Delete `Dockerfile.aarch64` — merge into `Dockerfile` with `ARG TARGETARCH` (buildx sets it); both files are 90% identical except the dep-copy block. [Dockerfile.aarch64]
|
||||||
|
- [ ] Delete entity dead methods: `update_username`, `update_password`, `change_stream_key_limit`, `find_by_user_id`, `get_stream_session`, `get_all_active_sessions`. Zero callers; scaffolding for the unbuilt admin panel. [crates/entity/src/users.rs:77, crates/entity/src/auth_session.rs:43, crates/entity/src/stream_session.rs:53]
|
||||||
|
- [ ] Delete `/api/health`, `/api/uptime`, `/api/version` + `HealthResponse`/`UptimeResponse` — all subsets of `/api/stats`; also kills `serde_json` (only used by `version_handler`). [crates/server/src/http.rs:497]
|
||||||
|
- [ ] Delete `SessionCookie` extractor — its value is only debug-logged, never used; kill the extractor + 2 logs. [crates/server/src/http.rs:213]
|
||||||
|
- [ ] Delete `WebrtcProxy::local_addr()` (no callers), dead `let _ = x_port ^ …` stmt, and single-field `WebRtcProxyConfig` struct → pass `u16` (also fixes i32 port type). [crates/server/src/webrtc_proxy.rs:167, crates/server/src/webrtc_proxy.rs:139]
|
||||||
|
- [ ] Delete `webrtc.rs` `mid` binding from `add_media` + `_hint_mid` param — passed straight into an underscore. [crates/server/src/webrtc.rs:142]
|
||||||
|
- [ ] Delete `stream_key` `is_active`/`is_unlisted` columns — never read anywhere; `create()` hardcodes `is_active=false` and `is_unlisted` is always passed `false`. Drop param + column (migration, expand-contract). [crates/entity/src/stream_key.rs:45]
|
||||||
|
- [ ] Delete `rtmp.rs` redundant `stream_id` var (warn can use `current_stream_key_id`), `warn!("")` empty-log on parse error, and `parse_video_codec`'s `Result<_, Box<dyn Error>>` → `Option`. [crates/server/src/rtmp.rs:206, crates/server/src/rtmp.rs:369]
|
||||||
|
- [ ] Delete `catalog_handler`'s `get_all_active_sessions` query — result only feeds a `debug!`; kills the entity method too. [crates/server/src/http.rs:185]
|
||||||
|
- [ ] Delete commented-out routes/code (`admin/server_stats`, duplicate whip route, `.max_age`, `fs::File`) + `#[axum::debug_handler]`. [crates/server/src/http.rs:109]
|
||||||
|
- [ ] Delete `StreamSession.active_clients` — written 0, never read; drops `AtomicU32` import. [crates/server/src/main.rs:65]
|
||||||
|
- [ ] Delete `webrtc_ingest` `_connected` flag — set true, never read. [crates/server/src/webrtc_ingest.rs:363]
|
||||||
|
- [ ] Delete `users.rs` `let user = …; user` pointless binding in `find_by_username`. [crates/entity/src/users.rs:70]
|
||||||
|
- [ ] Delete `meow_handler` + route — joke endpoint, zero consumers. [crates/server/src/http.rs:492]
|
||||||
|
- [ ] Delete `index.html` `prevStats` — assigned, never read. [index.html:57]
|
||||||
|
- [ ] Delete deps — server: `futures`, `rand`; entity: `rand`, `argon2` (hashing lives in server crate now); `serde_json` (with health/uptime/version cut). [crates/server/Cargo.toml, crates/entity/Cargo.toml]
|
||||||
|
|
||||||
|
## Shrinks
|
||||||
|
|
||||||
|
- [ ] Extract shared "remove session + finish_stream_session" fn — closure duplicated 3× (rtmp.rs cleanup, webrtc_ingest detach cleanup, whip DELETE handler). Also lets `handle_whip_injest` drop its manual `err()` closure for `Result<HttpError>` like its siblings. [crates/server/src/rtmp.rs:141, crates/server/src/webrtc_ingest.rs:338, crates/server/src/webrtc_ingest.rs:34]
|
||||||
|
- [ ] Extract `fixup_answer_sdp()` — `whip_sdp_probe.rs` re-encodes the SDP string-fixups verbatim; probe drops ~30 lines. [crates/server/tests/whip_sdp_probe.rs:66]
|
||||||
|
- [ ] Derive `thiserror` on `AudioParseError` instead of hand-rolled `Display`+`Error` impls (~15 lines). [crates/server/src/audio.rs:66]
|
||||||
|
- [ ] Collapse `webrtc.rs` four near-identical channel re-subscribe blocks (`is_none`/`is_closed` × video/audio) → one helper. [crates/server/src/webrtc.rs:336]
|
||||||
|
|
||||||
|
## Out of scope (correctness — route to normal review)
|
||||||
|
|
||||||
|
- `Local::now()` vs `Utc::now()` inconsistency in session creation
|
||||||
|
- `KEY_RE` `{1,67}` vs `MAX_LABEL_LEN = 64` mismatch
|
||||||
|
- `index.html` hardcoded port 5000
|
||||||
|
- `/api/stream/{slug}` unauthenticated
|
||||||
@@ -12,6 +12,8 @@ pub struct Model {
|
|||||||
pub label: String,
|
pub label: String,
|
||||||
pub is_active: bool,
|
pub is_active: bool,
|
||||||
pub is_unlisted: bool,
|
pub is_unlisted: bool,
|
||||||
|
pub password: Option<String>,
|
||||||
|
pub custom_id: Option<String>,
|
||||||
pub created_at: DateTimeUtc,
|
pub created_at: DateTimeUtc,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -78,15 +80,14 @@ impl Entity {
|
|||||||
.all(db)
|
.all(db)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
impl ActiveModel {
|
pub async fn find_by_custom_id(
|
||||||
pub async fn change_label_value(
|
|
||||||
mut self,
|
|
||||||
db: &DatabaseConnection,
|
db: &DatabaseConnection,
|
||||||
value: String,
|
custom_id: String,
|
||||||
) -> Result<Model, DbErr> {
|
) -> Result<Option<Model>, DbErr> {
|
||||||
self.label = Set(value);
|
Entity::find()
|
||||||
self.update(db).await
|
.filter(Column::CustomId.eq(custom_id))
|
||||||
|
.one(db)
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -50,32 +50,21 @@ impl Model {
|
|||||||
.insert(db)
|
.insert(db)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
pub async fn get_stream_session(
|
|
||||||
db: &DatabaseConnection,
|
|
||||||
stream_session_id: i32,
|
|
||||||
) -> Result<Model, DbErr> {
|
|
||||||
Entity::find_by_id(stream_session_id)
|
|
||||||
.one(db)
|
|
||||||
.await?
|
|
||||||
.ok_or(DbErr::RecordNotFound(format!(
|
|
||||||
"stream_session {stream_session_id}"
|
|
||||||
)))
|
|
||||||
}
|
|
||||||
pub async fn get_all_active_sessions(db: &DatabaseConnection) -> Result<Vec<Model>, DbErr> {
|
pub async fn get_all_active_sessions(db: &DatabaseConnection) -> Result<Vec<Model>, DbErr> {
|
||||||
Ok(Entity::find()
|
Entity::find()
|
||||||
.filter(Column::EndedAt.is_null())
|
.filter(Column::EndedAt.is_null())
|
||||||
.all(db)
|
.all(db)
|
||||||
.await?)
|
.await
|
||||||
}
|
}
|
||||||
pub async fn get_active_by_stream_key_id(
|
pub async fn get_active_by_stream_key_id(
|
||||||
db: &DatabaseConnection,
|
db: &DatabaseConnection,
|
||||||
stream_key_id: i32,
|
stream_key_id: i32,
|
||||||
) -> Result<Option<Model>, DbErr> {
|
) -> Result<Option<Model>, DbErr> {
|
||||||
Ok(Entity::find()
|
Entity::find()
|
||||||
.filter(Column::StreamKeyId.eq(stream_key_id))
|
.filter(Column::StreamKeyId.eq(stream_key_id))
|
||||||
.filter(Column::EndedAt.is_null())
|
.filter(Column::EndedAt.is_null())
|
||||||
.one(db)
|
.one(db)
|
||||||
.await?)
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn clean_unended_streams(db: &DatabaseConnection) -> Result<u64, DbErr> {
|
pub async fn clean_unended_streams(db: &DatabaseConnection) -> Result<u64, DbErr> {
|
||||||
|
|||||||
@@ -58,18 +58,18 @@ impl Entity {
|
|||||||
if let Some(x) = sessions {
|
if let Some(x) = sessions {
|
||||||
Entity::find_by_id(x.id_user).one(db).await
|
Entity::find_by_id(x.id_user).one(db).await
|
||||||
} else {
|
} else {
|
||||||
return Ok(None);
|
Ok(None)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
pub async fn find_by_username(
|
pub async fn find_by_username(
|
||||||
db: &DatabaseConnection,
|
db: &DatabaseConnection,
|
||||||
username: String,
|
username: String,
|
||||||
) -> Result<Option<Model>, DbErr> {
|
) -> Result<Option<Model>, DbErr> {
|
||||||
let user = Entity::find()
|
|
||||||
|
Entity::find()
|
||||||
.filter(Column::Username.eq(username))
|
.filter(Column::Username.eq(username))
|
||||||
.one(db)
|
.one(db)
|
||||||
.await;
|
.await
|
||||||
user
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,8 @@ mod m20260616_000001_create_users;
|
|||||||
mod m20260616_000002_create_stream_key;
|
mod m20260616_000002_create_stream_key;
|
||||||
mod m20260616_000003_create_stream_session;
|
mod m20260616_000003_create_stream_session;
|
||||||
mod m20260616_000004_create_auth_session;
|
mod m20260616_000004_create_auth_session;
|
||||||
|
mod m20260815_000005_add_password_to_stream_key;
|
||||||
|
mod m20260815_000006_set_stream_keys_unlisted;
|
||||||
|
|
||||||
pub struct Migrator;
|
pub struct Migrator;
|
||||||
|
|
||||||
@@ -15,6 +17,8 @@ impl MigratorTrait for Migrator {
|
|||||||
Box::new(m20260616_000002_create_stream_key::Migration),
|
Box::new(m20260616_000002_create_stream_key::Migration),
|
||||||
Box::new(m20260616_000003_create_stream_session::Migration),
|
Box::new(m20260616_000003_create_stream_session::Migration),
|
||||||
Box::new(m20260616_000004_create_auth_session::Migration),
|
Box::new(m20260616_000004_create_auth_session::Migration),
|
||||||
|
Box::new(m20260815_000005_add_password_to_stream_key::Migration),
|
||||||
|
Box::new(m20260815_000006_set_stream_keys_unlisted::Migration),
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -38,7 +38,7 @@ impl MigrationTrait for Migration {
|
|||||||
ColumnDef::new(StreamKey::IsUnlisted)
|
ColumnDef::new(StreamKey::IsUnlisted)
|
||||||
.boolean()
|
.boolean()
|
||||||
.not_null()
|
.not_null()
|
||||||
.default(true),
|
.default(false),
|
||||||
)
|
)
|
||||||
.col(ColumnDef::new(StreamKey::CreatedAt).date_time().not_null())
|
.col(ColumnDef::new(StreamKey::CreatedAt).date_time().not_null())
|
||||||
.foreign_key(
|
.foreign_key(
|
||||||
@@ -67,5 +67,7 @@ pub enum StreamKey {
|
|||||||
Label,
|
Label,
|
||||||
IsActive,
|
IsActive,
|
||||||
IsUnlisted,
|
IsUnlisted,
|
||||||
|
Password,
|
||||||
|
CustomId,
|
||||||
CreatedAt,
|
CreatedAt,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,47 @@
|
|||||||
|
use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
use super::m20260616_000002_create_stream_key::StreamKey;
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
Table::alter()
|
||||||
|
.table(StreamKey::Table)
|
||||||
|
.add_column(ColumnDef::new(StreamKey::Password).string().null())
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
Table::alter()
|
||||||
|
.table(StreamKey::Table)
|
||||||
|
.add_column(ColumnDef::new(StreamKey::CustomId).string().null())
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
Table::alter()
|
||||||
|
.table(StreamKey::Table)
|
||||||
|
.drop_column(StreamKey::Password)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
Table::alter()
|
||||||
|
.table(StreamKey::Table)
|
||||||
|
.drop_column(StreamKey::CustomId)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,26 @@
|
|||||||
|
use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
use super::m20260616_000002_create_stream_key::StreamKey;
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
// By accident the is_unlisted column defaulted to true; existing stream
|
||||||
|
// keys were meant to be listed. Reset all rows to false here.
|
||||||
|
manager
|
||||||
|
.exec_stmt(
|
||||||
|
Query::update()
|
||||||
|
.table(StreamKey::Table)
|
||||||
|
.values([(StreamKey::IsUnlisted, false.into())])
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, _manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "server"
|
name = "server"
|
||||||
version = "0.4.0"
|
version = "0.6.0"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
|
|
||||||
[target.x86_64-unknown-linux-gnu]
|
[target.x86_64-unknown-linux-gnu]
|
||||||
@@ -15,27 +15,25 @@ path = "src/main.rs"
|
|||||||
async-broadcast = "0.7.2"
|
async-broadcast = "0.7.2"
|
||||||
bytes = "1"
|
bytes = "1"
|
||||||
dashmap = "6.2.1"
|
dashmap = "6.2.1"
|
||||||
rand = "0.10.1"
|
|
||||||
rml_rtmp = "0.8.0"
|
rml_rtmp = "0.8.0"
|
||||||
str0m = "0.20.0"
|
str0m = "0.20.0"
|
||||||
tokio = { version = "1", features = ["full"] }
|
tokio = { version = "1", features = ["full"] }
|
||||||
axum = { version = "0.8", features = ["macros"] }
|
axum = { version = "0.8", features = ["macros"] }
|
||||||
serde = { version = "1.0.228", features = ["serde_derive"] }
|
serde = { version = "1.0.228", features = ["serde_derive"] }
|
||||||
serde_json = "1.0.150"
|
|
||||||
sea-orm = { version = "1", features = [ "sqlx-sqlite", "runtime-tokio-rustls", "macros" ] }
|
sea-orm = { version = "1", features = [ "sqlx-sqlite", "runtime-tokio-rustls", "macros" ] }
|
||||||
entity = {path = "../entity"}
|
entity = {path = "../entity"}
|
||||||
migration = {path = "../migration"}
|
migration = {path = "../migration"}
|
||||||
argon2 = "0.5.3"
|
argon2 = "0.5.3"
|
||||||
uuid = { version = "1.23.3", features = ["v4"] }
|
uuid = { version = "1.23.3", features = ["v4"] }
|
||||||
tower-http = { version = "0.6", features = ["cors"] }
|
tower-http = { version = "0.6", features = ["cors"] }
|
||||||
futures = "0.3.32"
|
|
||||||
tracing = "0.1"
|
tracing = "0.1"
|
||||||
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
||||||
symphonia = { version = "0.5", features = ["aac"] }
|
symphonia = { version = "0.5", features = ["aac"] }
|
||||||
opus = "0.3.1"
|
opus = "0.3.1"
|
||||||
rubato = "3.0.0"
|
rubato = "3.0.0"
|
||||||
chrono = "0.4.45"
|
chrono = "0.4.45"
|
||||||
regex = "1"
|
time = "0.3"
|
||||||
thiserror = "2.0.18"
|
thiserror = "2.0.18"
|
||||||
sysinfo = "0.36"
|
sysinfo = "0.36"
|
||||||
axum-extra = { version = "0.12.6", features = ["cookie"] }
|
axum-extra = { version = "0.12.6", features = ["cookie"] }
|
||||||
|
serde_json = "1.0.151"
|
||||||
|
|||||||
@@ -1,10 +1,10 @@
|
|||||||
use std::{error::Error, fmt::Display};
|
use std::error::Error;
|
||||||
|
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use rubato::{audioadapter_buffers::direct::InterleavedSlice, Fft, Resampler};
|
use rubato::{Fft, Resampler, audioadapter_buffers::direct::InterleavedSlice};
|
||||||
use symphonia::core::{
|
use symphonia::core::{
|
||||||
audio::SampleBuffer,
|
audio::SampleBuffer,
|
||||||
codecs::{CodecParameters, Decoder, DecoderOptions, CODEC_TYPE_AAC},
|
codecs::{CODEC_TYPE_AAC, CodecParameters, Decoder, DecoderOptions},
|
||||||
formats::Packet,
|
formats::Packet,
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -44,7 +44,7 @@ impl AudioProcesser {
|
|||||||
pub fn encode(&mut self, frame: AudioFrame) -> Vec<OpusAudioFrame> {
|
pub fn encode(&mut self, frame: AudioFrame) -> Vec<OpusAudioFrame> {
|
||||||
let mut samples: Vec<f32> = frame
|
let mut samples: Vec<f32> = frame
|
||||||
.data
|
.data
|
||||||
.chunks_exact(4)
|
.as_chunks::<4>().0.iter()
|
||||||
.map(|b| f32::from_le_bytes([b[0], b[1], b[2], b[3]]))
|
.map(|b| f32::from_le_bytes([b[0], b[1], b[2], b[3]]))
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
@@ -95,23 +95,6 @@ pub struct AACParser {
|
|||||||
decoder: Option<Box<dyn Decoder>>,
|
decoder: Option<Box<dyn Decoder>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug)]
|
|
||||||
enum AudioParseError {
|
|
||||||
InvalidCodec,
|
|
||||||
NoConfigPacket,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Display for AudioParseError {
|
|
||||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
||||||
match self {
|
|
||||||
Self::InvalidCodec => write!(f, "Not the right codec provided"),
|
|
||||||
Self::NoConfigPacket => write!(f, "No config packet cached"),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Error for AudioParseError {}
|
|
||||||
|
|
||||||
// Byte 0: upper nibble = sound format (10 = AAC)
|
// Byte 0: upper nibble = sound format (10 = AAC)
|
||||||
// Byte 1: 0 = AudioSpecificConfig, 1 = raw AAC frame
|
// Byte 1: 0 = AudioSpecificConfig, 1 = raw AAC frame
|
||||||
impl AACParser {
|
impl AACParser {
|
||||||
@@ -129,7 +112,7 @@ impl AACParser {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (bytes[0] >> 4) != 10 {
|
if (bytes[0] >> 4) != 10 {
|
||||||
return Err(Box::new(AudioParseError::InvalidCodec));
|
return Err("not the right codec provided".into());
|
||||||
}
|
}
|
||||||
|
|
||||||
if bytes[1] == 0 {
|
if bytes[1] == 0 {
|
||||||
@@ -142,10 +125,7 @@ impl AACParser {
|
|||||||
return Ok(None);
|
return Ok(None);
|
||||||
}
|
}
|
||||||
|
|
||||||
let decoder = self
|
let decoder = self.decoder.as_mut().ok_or("no config packet cached")?;
|
||||||
.decoder
|
|
||||||
.as_mut()
|
|
||||||
.ok_or(AudioParseError::NoConfigPacket)?;
|
|
||||||
|
|
||||||
let packet = Packet::new_from_boxed_slice(
|
let packet = Packet::new_from_boxed_slice(
|
||||||
0,
|
0,
|
||||||
|
|||||||
@@ -50,7 +50,9 @@ impl Av1CodecParser {
|
|||||||
let has_size = (header >> 1) & 1 != 0;
|
let has_size = (header >> 1) & 1 != 0;
|
||||||
i += 1;
|
i += 1;
|
||||||
if has_extension {
|
if has_extension {
|
||||||
if i >= data.len() { return false; }
|
if i >= data.len() {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
i += 1;
|
i += 1;
|
||||||
}
|
}
|
||||||
if obu_type == 1 {
|
if obu_type == 1 {
|
||||||
@@ -61,13 +63,19 @@ impl Av1CodecParser {
|
|||||||
let mut size: usize = 0;
|
let mut size: usize = 0;
|
||||||
let mut shift = 0;
|
let mut shift = 0;
|
||||||
loop {
|
loop {
|
||||||
if i >= data.len() { return false; }
|
if i >= data.len() {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
let b = data[i] as usize;
|
let b = data[i] as usize;
|
||||||
i += 1;
|
i += 1;
|
||||||
size |= (b & 0x7F) << shift;
|
size |= (b & 0x7F) << shift;
|
||||||
shift += 7;
|
shift += 7;
|
||||||
if b & 0x80 == 0 { break; }
|
if b & 0x80 == 0 {
|
||||||
if shift > 32 { return false; }
|
break;
|
||||||
|
}
|
||||||
|
if shift > 32 {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
i += size;
|
i += size;
|
||||||
} else {
|
} else {
|
||||||
@@ -78,7 +86,11 @@ impl Av1CodecParser {
|
|||||||
false
|
false
|
||||||
}
|
}
|
||||||
|
|
||||||
fn obus_for_frame(&mut self, payload: &[u8], rtmp_is_keyframe: bool) -> Option<(Vec<u8>, bool)> {
|
fn obus_for_frame(
|
||||||
|
&mut self,
|
||||||
|
payload: &[u8],
|
||||||
|
rtmp_is_keyframe: bool,
|
||||||
|
) -> Option<(Vec<u8>, bool)> {
|
||||||
if payload.is_empty() {
|
if payload.is_empty() {
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
@@ -95,13 +107,12 @@ impl Av1CodecParser {
|
|||||||
|
|
||||||
self.first_coded_frame = false;
|
self.first_coded_frame = false;
|
||||||
|
|
||||||
if is_keyframe
|
if is_keyframe && let Some(config) = &self.config_obus {
|
||||||
&& let Some(config) = &self.config_obus {
|
let mut out = Vec::with_capacity(config.len() + payload.len());
|
||||||
let mut out = Vec::with_capacity(config.len() + payload.len());
|
out.extend_from_slice(config);
|
||||||
out.extend_from_slice(config);
|
out.extend_from_slice(payload);
|
||||||
out.extend_from_slice(payload);
|
return Some((out, true));
|
||||||
return Some((out, true));
|
}
|
||||||
}
|
|
||||||
|
|
||||||
Some((payload.to_vec(), is_keyframe))
|
Some((payload.to_vec(), is_keyframe))
|
||||||
}
|
}
|
||||||
@@ -129,7 +140,11 @@ impl CodecParser for Av1CodecParser {
|
|||||||
// FourCC — enhanced RTMP only defines CTS for hvc1 CodedFrames.
|
// FourCC — enhanced RTMP only defines CTS for hvc1 CodedFrames.
|
||||||
let payload = data.get(5..)?;
|
let payload = data.get(5..)?;
|
||||||
let (obus, is_keyframe) = self.obus_for_frame(payload, rtmp_is_keyframe)?;
|
let (obus, is_keyframe) = self.obus_for_frame(payload, rtmp_is_keyframe)?;
|
||||||
Some(VideoFrame { data: Bytes::from(obus), is_keyframe, timestamp_ms })
|
Some(VideoFrame {
|
||||||
|
data: Bytes::from(obus),
|
||||||
|
is_keyframe,
|
||||||
|
timestamp_ms,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
_ => None,
|
_ => None,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -53,12 +53,20 @@ impl H264CodecParser {
|
|||||||
}
|
}
|
||||||
let pts_ms = Self::pts_ms(timestamp_ms, &bytes[5..8]);
|
let pts_ms = Self::pts_ms(timestamp_ms, &bytes[5..8]);
|
||||||
let data = self.avcc_to_annexb(bytes.get(8..)?, is_keyframe)?;
|
let data = self.avcc_to_annexb(bytes.get(8..)?, is_keyframe)?;
|
||||||
Some(VideoFrame { data: Bytes::from(data), is_keyframe, timestamp_ms: pts_ms })
|
Some(VideoFrame {
|
||||||
|
data: Bytes::from(data),
|
||||||
|
is_keyframe,
|
||||||
|
timestamp_ms: pts_ms,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
3 => {
|
3 => {
|
||||||
// CodedFramesX: no CTS field, bytes 5+ = AVCC NALUs
|
// CodedFramesX: no CTS field, bytes 5+ = AVCC NALUs
|
||||||
let data = self.avcc_to_annexb(bytes.get(5..)?, is_keyframe)?;
|
let data = self.avcc_to_annexb(bytes.get(5..)?, is_keyframe)?;
|
||||||
Some(VideoFrame { data: Bytes::from(data), is_keyframe, timestamp_ms })
|
Some(VideoFrame {
|
||||||
|
data: Bytes::from(data),
|
||||||
|
is_keyframe,
|
||||||
|
timestamp_ms,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
_ => None,
|
_ => None,
|
||||||
}
|
}
|
||||||
@@ -77,7 +85,11 @@ impl H264CodecParser {
|
|||||||
// PTS = DTS + CTS — see note on CodedFrames above.
|
// PTS = DTS + CTS — see note on CodedFrames above.
|
||||||
let pts_ms = Self::pts_ms(timestamp_ms, &bytes[2..5]);
|
let pts_ms = Self::pts_ms(timestamp_ms, &bytes[2..5]);
|
||||||
let data = self.avcc_to_annexb(&bytes[5..], is_keyframe)?;
|
let data = self.avcc_to_annexb(&bytes[5..], is_keyframe)?;
|
||||||
Some(VideoFrame { data: Bytes::from(data), is_keyframe, timestamp_ms: pts_ms })
|
Some(VideoFrame {
|
||||||
|
data: Bytes::from(data),
|
||||||
|
is_keyframe,
|
||||||
|
timestamp_ms: pts_ms,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
_ => None,
|
_ => None,
|
||||||
}
|
}
|
||||||
@@ -148,13 +160,12 @@ impl H264CodecParser {
|
|||||||
|
|
||||||
// Prepend SPS+PPS before every keyframe so str0m's packetizer
|
// Prepend SPS+PPS before every keyframe so str0m's packetizer
|
||||||
// can bundle them into a STAP-A alongside the IDR NALU.
|
// can bundle them into a STAP-A alongside the IDR NALU.
|
||||||
if is_keyframe
|
if is_keyframe && let (Some(sps), Some(pps)) = (&self.sps, &self.pps) {
|
||||||
&& let (Some(sps), Some(pps)) = (&self.sps, &self.pps) {
|
out.extend_from_slice(&[0, 0, 0, 1]);
|
||||||
out.extend_from_slice(&[0, 0, 0, 1]);
|
out.extend_from_slice(sps);
|
||||||
out.extend_from_slice(sps);
|
out.extend_from_slice(&[0, 0, 0, 1]);
|
||||||
out.extend_from_slice(&[0, 0, 0, 1]);
|
out.extend_from_slice(pps);
|
||||||
out.extend_from_slice(pps);
|
}
|
||||||
}
|
|
||||||
|
|
||||||
// Convert each length-prefixed NALU to an Annex B start-code NALU.
|
// Convert each length-prefixed NALU to an Annex B start-code NALU.
|
||||||
let mut i = 0;
|
let mut i = 0;
|
||||||
|
|||||||
@@ -74,15 +74,15 @@ impl H265CodecParser {
|
|||||||
fn hvcc_to_annexb(&self, payload: &[u8], is_keyframe: bool) -> Option<Vec<u8>> {
|
fn hvcc_to_annexb(&self, payload: &[u8], is_keyframe: bool) -> Option<Vec<u8>> {
|
||||||
let mut out = Vec::with_capacity(payload.len());
|
let mut out = Vec::with_capacity(payload.len());
|
||||||
|
|
||||||
if is_keyframe
|
if is_keyframe && let (Some(vps), Some(sps), Some(pps)) = (&self.vps, &self.sps, &self.pps)
|
||||||
&& let (Some(vps), Some(sps), Some(pps)) = (&self.vps, &self.sps, &self.pps) {
|
{
|
||||||
out.extend_from_slice(&[0, 0, 0, 1]);
|
out.extend_from_slice(&[0, 0, 0, 1]);
|
||||||
out.extend_from_slice(vps);
|
out.extend_from_slice(vps);
|
||||||
out.extend_from_slice(&[0, 0, 0, 1]);
|
out.extend_from_slice(&[0, 0, 0, 1]);
|
||||||
out.extend_from_slice(sps);
|
out.extend_from_slice(sps);
|
||||||
out.extend_from_slice(&[0, 0, 0, 1]);
|
out.extend_from_slice(&[0, 0, 0, 1]);
|
||||||
out.extend_from_slice(pps);
|
out.extend_from_slice(pps);
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut i = 0;
|
let mut i = 0;
|
||||||
while i + 4 <= payload.len() {
|
while i + 4 <= payload.len() {
|
||||||
@@ -128,12 +128,20 @@ impl CodecParser for H265CodecParser {
|
|||||||
let cts = i32::from_be_bytes([0, data[5], data[6], data[7]]) << 8 >> 8;
|
let cts = i32::from_be_bytes([0, data[5], data[6], data[7]]) << 8 >> 8;
|
||||||
let pts_ms = (timestamp_ms as i64 + cts as i64).max(0) as u32;
|
let pts_ms = (timestamp_ms as i64 + cts as i64).max(0) as u32;
|
||||||
let annexb = self.hvcc_to_annexb(data.get(8..)?, is_keyframe)?;
|
let annexb = self.hvcc_to_annexb(data.get(8..)?, is_keyframe)?;
|
||||||
Some(VideoFrame { data: Bytes::from(annexb), is_keyframe, timestamp_ms: pts_ms })
|
Some(VideoFrame {
|
||||||
|
data: Bytes::from(annexb),
|
||||||
|
is_keyframe,
|
||||||
|
timestamp_ms: pts_ms,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
3 => {
|
3 => {
|
||||||
// CodedFramesX: no CTS, bytes 5+ = HVCC
|
// CodedFramesX: no CTS, bytes 5+ = HVCC
|
||||||
let annexb = self.hvcc_to_annexb(data.get(5..)?, is_keyframe)?;
|
let annexb = self.hvcc_to_annexb(data.get(5..)?, is_keyframe)?;
|
||||||
Some(VideoFrame { data: Bytes::from(annexb), is_keyframe, timestamp_ms })
|
Some(VideoFrame {
|
||||||
|
data: Bytes::from(annexb),
|
||||||
|
is_keyframe,
|
||||||
|
timestamp_ms,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
_ => None,
|
_ => None,
|
||||||
}
|
}
|
||||||
|
|||||||
+147
-123
@@ -1,28 +1,28 @@
|
|||||||
use std::{
|
use std::{
|
||||||
net::SocketAddr,
|
net::SocketAddr,
|
||||||
sync::{
|
sync::{
|
||||||
Arc, LazyLock,
|
Arc,
|
||||||
atomic::{AtomicI32, Ordering},
|
atomic::{AtomicI32, Ordering},
|
||||||
},
|
},
|
||||||
time::Duration,
|
time::{Duration, Instant},
|
||||||
};
|
};
|
||||||
|
|
||||||
use axum_extra::extract::{CookieJar, cookie::Cookie};
|
use axum_extra::extract::{CookieJar, cookie::Cookie};
|
||||||
use regex::Regex;
|
|
||||||
|
|
||||||
// The Rust `regex` crate is guaranteed linear-time and therefore does NOT
|
/// Allowed charset for labels, passwords and custom IDs: letters, numbers,
|
||||||
// support lookaround, so the JS source pattern
|
/// dashes and apostrophes — no spaces. Mirrors the frontend
|
||||||
// /^(?=.*[A-Za-z])[A-Za-z0-9_-]{1,67}$/
|
/// `/^[A-Za-z0-9'-]{0,67}$/` used by the keys-page popups.
|
||||||
// cannot be ported verbatim. The `(?=.*[A-Za-z])` lookahead only means
|
fn valid_charset(s: &str) -> bool {
|
||||||
// "must contain at least one letter" — we drop it from the pattern and
|
s.len() <= 67
|
||||||
// enforce that condition with a separate `.chars().any(..)` check below.
|
&& s.chars()
|
||||||
static KEY_RE: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"^[A-Za-z0-9 _-]{1,67}$").unwrap());
|
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '\'')
|
||||||
|
}
|
||||||
|
|
||||||
use axum::{
|
use axum::{
|
||||||
Json, Router,
|
Json, Router,
|
||||||
extract::{FromRequestParts, Path, State},
|
extract::{FromRequestParts, Path, State},
|
||||||
http::{
|
http::{
|
||||||
HeaderName, Method, StatusCode,
|
HeaderMap, HeaderName, Method, StatusCode,
|
||||||
header::{ACCEPT, AUTHORIZATION, CONTENT_TYPE},
|
header::{ACCEPT, AUTHORIZATION, CONTENT_TYPE},
|
||||||
request::Parts,
|
request::Parts,
|
||||||
},
|
},
|
||||||
@@ -31,7 +31,7 @@ use axum::{
|
|||||||
};
|
};
|
||||||
use chrono::{DateTime, Utc};
|
use chrono::{DateTime, Utc};
|
||||||
use entity::{auth_session, stream_key, stream_session, users};
|
use entity::{auth_session, stream_key, stream_session, users};
|
||||||
use sea_orm::{DatabaseConnection, EntityTrait, IntoActiveModel};
|
use sea_orm::{ActiveModelTrait, DatabaseConnection, EntityTrait, IntoActiveModel, Set};
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use sysinfo::System;
|
use sysinfo::System;
|
||||||
use tokio::{
|
use tokio::{
|
||||||
@@ -52,23 +52,15 @@ use crate::{
|
|||||||
const MAX_USERNAME_LEN: usize = 32;
|
const MAX_USERNAME_LEN: usize = 32;
|
||||||
const MAX_LABEL_LEN: usize = 64;
|
const MAX_LABEL_LEN: usize = 64;
|
||||||
|
|
||||||
pub struct HttpServerConfig {
|
|
||||||
pub signup_code: String,
|
|
||||||
}
|
|
||||||
|
|
||||||
pub struct ServerInfo {
|
|
||||||
pub version: &'static str,
|
|
||||||
pub start_time: std::time::Instant,
|
|
||||||
}
|
|
||||||
|
|
||||||
pub struct HttpServer {
|
pub struct HttpServer {
|
||||||
pub offer_tx: Sender<(i32, i32, String)>,
|
pub offer_tx: Sender<(i32, i32, String)>,
|
||||||
pub accept_rx: async_broadcast::InactiveReceiver<(i32, Option<String>)>,
|
pub accept_rx: async_broadcast::InactiveReceiver<(i32, Result<String, String>)>,
|
||||||
pub appstate: Arc<Mutex<AppState>>,
|
pub appstate: Arc<Mutex<AppState>>,
|
||||||
pub request_count: AtomicI32,
|
pub request_count: AtomicI32,
|
||||||
pub db: DatabaseConnection,
|
pub db: DatabaseConnection,
|
||||||
pub config: Arc<HttpServerConfig>,
|
pub signup_code: String,
|
||||||
pub info: Arc<ServerInfo>,
|
pub version: &'static str,
|
||||||
|
pub start_time: Instant,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl HttpServer {
|
impl HttpServer {
|
||||||
@@ -106,8 +98,6 @@ impl HttpServer {
|
|||||||
.get(get_all_stream_keys)
|
.get(get_all_stream_keys)
|
||||||
.patch(edit_stream_key),
|
.patch(edit_stream_key),
|
||||||
)
|
)
|
||||||
// .route("/api/admin/server_stats", get(todo!()))
|
|
||||||
// .route("/api/whip", post(handle_whip_injest))
|
|
||||||
.route("/api/login", post(login_handler))
|
.route("/api/login", post(login_handler))
|
||||||
.route("/api/stream/{slug}", post(stream_handler))
|
.route("/api/stream/{slug}", post(stream_handler))
|
||||||
.route("/api/whip", post(handle_whip_injest))
|
.route("/api/whip", post(handle_whip_injest))
|
||||||
@@ -115,14 +105,8 @@ impl HttpServer {
|
|||||||
"/api/whip/{slug}",
|
"/api/whip/{slug}",
|
||||||
delete(handle_whip_injest_delete).patch(handle_whip_injest_patch),
|
delete(handle_whip_injest_delete).patch(handle_whip_injest_patch),
|
||||||
)
|
)
|
||||||
// .route(
|
|
||||||
// "/api/whip/{slug}",
|
|
||||||
// delete(handle_whip_injest_delete).patch(handle_whip_injest_patch),
|
|
||||||
// )
|
|
||||||
.route("/api/meow", get(meow_handler))
|
.route("/api/meow", get(meow_handler))
|
||||||
.route("/api/health", get(health_handler))
|
.route("/api/health", get(health_handler))
|
||||||
.route("/api/uptime", get(uptime_handler))
|
|
||||||
.route("/api/version", get(version_handler))
|
|
||||||
.route("/api/stats", get(stats_handler))
|
.route("/api/stats", get(stats_handler))
|
||||||
.layer(cors)
|
.layer(cors)
|
||||||
.with_state(state);
|
.with_state(state);
|
||||||
@@ -141,15 +125,24 @@ impl HttpServer {
|
|||||||
#[derive(Serialize)]
|
#[derive(Serialize)]
|
||||||
struct StreamListing {
|
struct StreamListing {
|
||||||
label: String,
|
label: String,
|
||||||
|
custom_url_label: String,
|
||||||
id: i32,
|
id: i32,
|
||||||
user: String,
|
user: String,
|
||||||
started_at: DateTime<Utc>, //UNIX TIMESTAMP
|
started_at: DateTime<Utc>,
|
||||||
|
is_password_protected: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Deserialize)]
|
#[derive(Deserialize)]
|
||||||
struct EditStreamKeyRequest {
|
struct EditStreamKeyRequest {
|
||||||
id: i32,
|
id: i32,
|
||||||
new: String,
|
#[serde(default)]
|
||||||
|
new: Option<String>,
|
||||||
|
#[serde(default)]
|
||||||
|
password: Option<String>,
|
||||||
|
#[serde(default)]
|
||||||
|
unlisted: Option<bool>,
|
||||||
|
#[serde(default)]
|
||||||
|
custom_id: Option<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn edit_stream_key(
|
async fn edit_stream_key(
|
||||||
@@ -158,15 +151,25 @@ async fn edit_stream_key(
|
|||||||
Json(payload): Json<EditStreamKeyRequest>,
|
Json(payload): Json<EditStreamKeyRequest>,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
// Trim surrounding whitespace so labels aren't stored with leading/trailing spaces.
|
// Trim surrounding whitespace so labels aren't stored with leading/trailing spaces.
|
||||||
let new_label = payload.new.trim();
|
let new_label = payload.new.as_deref().map(str::trim);
|
||||||
// Length (1..=67) and allowed charset ([A-Za-z0-9 _-], space included).
|
// Length (1..=67), allowed charset, and at least one letter.
|
||||||
if !KEY_RE.is_match(new_label) {
|
if let Some(label) = new_label {
|
||||||
return Err(HttpError::BadRequest("invalid label".into()));
|
if label.is_empty() || !valid_charset(label) {
|
||||||
}
|
return Err(HttpError::BadRequest("invalid label".into()));
|
||||||
// Replaces the JS lookahead: label must contain at least one letter.
|
}
|
||||||
if !new_label.chars().any(|c| c.is_ascii_alphabetic()) {
|
if !label.chars().any(|c| c.is_ascii_alphabetic()) {
|
||||||
return Err(HttpError::BadRequest("label must contain a letter".into()));
|
return Err(HttpError::BadRequest("label must contain a letter".into()));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
// Empty values are allowed (they clear the field).
|
||||||
|
if let Some(custom_id) = payload.custom_id.as_deref()
|
||||||
|
&& !custom_id.is_empty() && !valid_charset(custom_id) {
|
||||||
|
return Err(HttpError::BadRequest("invalid custom id".into()));
|
||||||
|
}
|
||||||
|
if let Some(pwd) = payload.password.as_deref()
|
||||||
|
&& !pwd.is_empty() && !valid_charset(pwd) {
|
||||||
|
return Err(HttpError::BadRequest("invalid password".into()));
|
||||||
|
}
|
||||||
|
|
||||||
let stream_key = stream_key::Entity::find_by_id(payload.id)
|
let stream_key = stream_key::Entity::find_by_id(payload.id)
|
||||||
.one(&state.db)
|
.one(&state.db)
|
||||||
@@ -176,10 +179,60 @@ async fn edit_stream_key(
|
|||||||
if stream_key.user_id != auth.0.id {
|
if stream_key.user_id != auth.0.id {
|
||||||
return Err(HttpError::Forbidden);
|
return Err(HttpError::Forbidden);
|
||||||
}
|
}
|
||||||
stream_key
|
|
||||||
.into_active_model()
|
let lock = state.appstate.lock().await;
|
||||||
.change_label_value(&state.db, new_label.to_string())
|
let mut ses = { lock.stream_sessions.get_mut(&payload.id) };
|
||||||
.await?;
|
// drop(lock);
|
||||||
|
|
||||||
|
let mut am = stream_key.into_active_model();
|
||||||
|
if let Some(custom_id) = payload.custom_id {
|
||||||
|
am.custom_id = if custom_id.is_empty() {
|
||||||
|
if let Some(ref mut ses) = ses {
|
||||||
|
ses.custom_id = None;
|
||||||
|
}
|
||||||
|
Set(None)
|
||||||
|
} else {
|
||||||
|
// Check for conflict.
|
||||||
|
if stream_key::Entity::find_by_custom_id(&state.db, custom_id.clone())
|
||||||
|
.await?
|
||||||
|
.is_some()
|
||||||
|
{
|
||||||
|
return Err(HttpError::Conflict);
|
||||||
|
}
|
||||||
|
if let Some(ref mut ses) = ses {
|
||||||
|
ses.custom_id = Some(custom_id.clone());
|
||||||
|
}
|
||||||
|
Set(Some(custom_id))
|
||||||
|
};
|
||||||
|
}
|
||||||
|
if let Some(label) = new_label {
|
||||||
|
if let Some(ref mut ses) = ses {
|
||||||
|
ses.stream_key_label = label.to_string();
|
||||||
|
}
|
||||||
|
am.label = Set(label.to_string());
|
||||||
|
}
|
||||||
|
if let Some(pwd) = payload.password {
|
||||||
|
am.password = if pwd.is_empty() {
|
||||||
|
if let Some(ref mut ses) = ses {
|
||||||
|
ses.password = None;
|
||||||
|
}
|
||||||
|
Set(None)
|
||||||
|
} else {
|
||||||
|
if let Some(ref mut ses) = ses {
|
||||||
|
ses.password = Some(pwd.clone());
|
||||||
|
}
|
||||||
|
Set(Some(pwd))
|
||||||
|
};
|
||||||
|
}
|
||||||
|
if let Some(unlisted) = payload.unlisted {
|
||||||
|
if let Some(ref mut ses) = ses {
|
||||||
|
ses.is_unlisted = unlisted;
|
||||||
|
}
|
||||||
|
am.is_unlisted = Set(unlisted);
|
||||||
|
}
|
||||||
|
// This hopefully will not fail, if it does, our values for the stream key will be mismatched,
|
||||||
|
// that would be bad
|
||||||
|
am.update(&state.db).await?;
|
||||||
|
|
||||||
Ok(StatusCode::OK)
|
Ok(StatusCode::OK)
|
||||||
}
|
}
|
||||||
@@ -195,11 +248,14 @@ async fn catalog_handler(
|
|||||||
.await
|
.await
|
||||||
.stream_sessions
|
.stream_sessions
|
||||||
.iter()
|
.iter()
|
||||||
|
.filter(|x| !x.is_unlisted)
|
||||||
.map(|x| StreamListing {
|
.map(|x| StreamListing {
|
||||||
id: x.stream_key_id,
|
id: x.stream_key_id,
|
||||||
label: x.stream_key_label.clone(),
|
label: x.stream_key_label.clone(),
|
||||||
|
custom_url_label: x.custom_id.clone().unwrap_or_default(),
|
||||||
user: x.stream_key_user.clone(),
|
user: x.stream_key_user.clone(),
|
||||||
started_at: x.started_at,
|
started_at: x.started_at,
|
||||||
|
is_password_protected: x.password.is_some(),
|
||||||
})
|
})
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
Ok(Json(catalog))
|
Ok(Json(catalog))
|
||||||
@@ -210,22 +266,6 @@ struct CreateStreamKeyBody {
|
|||||||
label: String,
|
label: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
struct SessionCookie(Option<String>);
|
|
||||||
|
|
||||||
impl<S> FromRequestParts<S> for SessionCookie
|
|
||||||
where
|
|
||||||
S: Send + Sync,
|
|
||||||
{
|
|
||||||
type Rejection = HttpError;
|
|
||||||
|
|
||||||
async fn from_request_parts(parts: &mut Parts, _state: &S) -> Result<Self, Self::Rejection> {
|
|
||||||
let jar = CookieJar::from_request_parts(parts, _state).await.unwrap();
|
|
||||||
Ok(SessionCookie(
|
|
||||||
jar.get("session").map(|c| c.value().to_string()),
|
|
||||||
))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
struct AuthUser(entity::users::Model);
|
struct AuthUser(entity::users::Model);
|
||||||
|
|
||||||
impl FromRequestParts<Arc<HttpServer>> for AuthUser {
|
impl FromRequestParts<Arc<HttpServer>> for AuthUser {
|
||||||
@@ -315,30 +355,26 @@ where
|
|||||||
|
|
||||||
fn cookie_for_token(token: &str, dev: bool) -> Cookie<'static> {
|
fn cookie_for_token(token: &str, dev: bool) -> Cookie<'static> {
|
||||||
let token = token.to_owned();
|
let token = token.to_owned();
|
||||||
match dev {
|
let mut cookie = Cookie::build(("session", token))
|
||||||
true => Cookie::build(("session", token))
|
.path("/")
|
||||||
.path("/")
|
.max_age(time::Duration::days(30))
|
||||||
.http_only(true)
|
.same_site(if dev {
|
||||||
.same_site(axum_extra::extract::cookie::SameSite::None)
|
axum_extra::extract::cookie::SameSite::None
|
||||||
.build(),
|
} else {
|
||||||
false => Cookie::build(("session", token))
|
axum_extra::extract::cookie::SameSite::Lax
|
||||||
.path("/")
|
});
|
||||||
.http_only(true)
|
if dev {
|
||||||
// .max_age(Duration::from_hours(999999))
|
cookie = cookie.secure(true);
|
||||||
.same_site(axum_extra::extract::cookie::SameSite::Lax)
|
|
||||||
.build(),
|
|
||||||
}
|
}
|
||||||
|
cookie.build()
|
||||||
}
|
}
|
||||||
|
|
||||||
#[axum::debug_handler]
|
#[axum::debug_handler]
|
||||||
async fn login_handler(
|
async fn login_handler(
|
||||||
session: SessionCookie,
|
|
||||||
State(state): State<Arc<HttpServer>>,
|
State(state): State<Arc<HttpServer>>,
|
||||||
DevFlag(dev): DevFlag,
|
DevFlag(dev): DevFlag,
|
||||||
Json(payload): Json<LoginForm>,
|
Json(payload): Json<LoginForm>,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
tracing::debug!(session = ?session.0, "login: existing session cookie");
|
|
||||||
|
|
||||||
let user = users::Entity::find_by_username(&state.db, payload.username.clone())
|
let user = users::Entity::find_by_username(&state.db, payload.username.clone())
|
||||||
.await?
|
.await?
|
||||||
.ok_or_else(|| {
|
.ok_or_else(|| {
|
||||||
@@ -370,13 +406,11 @@ struct CreateUserForm {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn create_user_handler(
|
async fn create_user_handler(
|
||||||
session: SessionCookie,
|
|
||||||
State(state): State<Arc<HttpServer>>,
|
State(state): State<Arc<HttpServer>>,
|
||||||
DevFlag(dev): DevFlag,
|
DevFlag(dev): DevFlag,
|
||||||
Json(payload): Json<CreateUserForm>,
|
Json(payload): Json<CreateUserForm>,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
tracing::debug!(session = ?session.0, "create_user: existing session cookie");
|
if state.signup_code.is_empty() || payload.ref_token != state.signup_code {
|
||||||
if state.config.signup_code.is_empty() || payload.ref_token != state.config.signup_code {
|
|
||||||
warn!(username = %payload.username, "signup rejected: invalid signup code");
|
warn!(username = %payload.username, "signup rejected: invalid signup code");
|
||||||
return Err(HttpError::Unauthorized);
|
return Err(HttpError::Unauthorized);
|
||||||
}
|
}
|
||||||
@@ -412,16 +446,18 @@ async fn create_user_handler(
|
|||||||
async fn stream_handler(
|
async fn stream_handler(
|
||||||
State(state): State<Arc<HttpServer>>,
|
State(state): State<Arc<HttpServer>>,
|
||||||
Path(slug): Path<String>,
|
Path(slug): Path<String>,
|
||||||
|
headers: HeaderMap,
|
||||||
body: String,
|
body: String,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
let request_id_clone = state.request_count.fetch_add(1, Ordering::Relaxed);
|
let request_id_clone = state.request_count.fetch_add(1, Ordering::Relaxed);
|
||||||
|
let app = state.appstate.lock().await;
|
||||||
|
|
||||||
let stream_key_id = {
|
let stream_key_id = {
|
||||||
let app = state.appstate.lock().await;
|
app.stream_sessions.iter().find(|e| {
|
||||||
app.stream_sessions
|
e.value().stream_key_id.to_string() == slug
|
||||||
.iter()
|
|| e.value().custom_id.as_deref() == Some(slug.as_str())
|
||||||
.find(|e| e.value().stream_key_id.to_string() == slug)
|
})
|
||||||
.map(|e| *e.key())
|
// .map(|e| *e.key())
|
||||||
};
|
};
|
||||||
|
|
||||||
let stream_key_id = stream_key_id.ok_or_else(|| {
|
let stream_key_id = stream_key_id.ok_or_else(|| {
|
||||||
@@ -429,12 +465,29 @@ async fn stream_handler(
|
|||||||
HttpError::NotFound
|
HttpError::NotFound
|
||||||
})?;
|
})?;
|
||||||
|
|
||||||
info!(request_id = request_id_clone, slug = %slug, stream_key_id, "WHEP offer received");
|
// Dont return stream by its id in the db if its unlisted
|
||||||
|
if stream_key_id.is_unlisted && stream_key_id.custom_id.clone().ok_or("") != Ok(slug.clone()) {
|
||||||
|
return Err(HttpError::NotFound);
|
||||||
|
}
|
||||||
|
|
||||||
|
let auth_header = headers.get("auth");
|
||||||
|
|
||||||
|
if let Some(password) = &stream_key_id.password {
|
||||||
|
if let Some(auth) = auth_header {
|
||||||
|
if auth.to_str().unwrap() != password {
|
||||||
|
return Err(HttpError::Unauthorized);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
return Err(HttpError::Unauthorized);
|
||||||
|
};
|
||||||
|
};
|
||||||
|
|
||||||
|
info!(request_id = request_id_clone, slug = %slug, stream_key_id.stream_key_id, "WHEP offer received");
|
||||||
let accept_rx = state.accept_rx.activate_cloned();
|
let accept_rx = state.accept_rx.activate_cloned();
|
||||||
// The webrtc worker owning the receiver died if this fails.
|
// The webrtc worker owning the receiver died if this fails.
|
||||||
state
|
state
|
||||||
.offer_tx
|
.offer_tx
|
||||||
.send((request_id_clone, stream_key_id, body))
|
.send((request_id_clone, stream_key_id.stream_key_id, body))
|
||||||
.await
|
.await
|
||||||
.map_err(|_| HttpError::Internal)?;
|
.map_err(|_| HttpError::Internal)?;
|
||||||
debug!(
|
debug!(
|
||||||
@@ -456,18 +509,15 @@ async fn stream_handler(
|
|||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(Some(Some(reply))) => {
|
Ok(Some(Ok(reply))) => {
|
||||||
return Ok(Response::builder()
|
return Ok(Response::builder()
|
||||||
.status(StatusCode::CREATED)
|
.status(StatusCode::CREATED)
|
||||||
.header("content-type", "application/sdp")
|
.header("content-type", "application/sdp")
|
||||||
.body(reply)
|
.body(reply)
|
||||||
.unwrap());
|
.unwrap());
|
||||||
}
|
}
|
||||||
Ok(Some(None)) => {
|
Ok(Some(Err(codec))) => {
|
||||||
return Ok(Response::builder()
|
return Err(HttpError::WhepCodecError(codec));
|
||||||
.status(StatusCode::UNSUPPORTED_MEDIA_TYPE)
|
|
||||||
.body(String::new())
|
|
||||||
.unwrap());
|
|
||||||
}
|
}
|
||||||
Ok(None) => {
|
Ok(None) => {
|
||||||
info!(
|
info!(
|
||||||
@@ -503,43 +553,17 @@ struct HealthResponse {
|
|||||||
async fn health_handler(
|
async fn health_handler(
|
||||||
State(state): State<Arc<HttpServer>>,
|
State(state): State<Arc<HttpServer>>,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
let uptime = state.info.start_time.elapsed().as_secs();
|
let uptime = state.start_time.elapsed().as_secs();
|
||||||
Ok((
|
Ok((
|
||||||
StatusCode::OK,
|
StatusCode::OK,
|
||||||
Json(HealthResponse {
|
Json(HealthResponse {
|
||||||
status: "ok",
|
status: "ok",
|
||||||
uptime_seconds: uptime,
|
uptime_seconds: uptime,
|
||||||
version: state.info.version,
|
version: state.version,
|
||||||
}),
|
}),
|
||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Serialize)]
|
|
||||||
struct UptimeResponse {
|
|
||||||
uptime_seconds: u64,
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn uptime_handler(
|
|
||||||
State(state): State<Arc<HttpServer>>,
|
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
|
||||||
let uptime = state.info.start_time.elapsed().as_secs();
|
|
||||||
Ok((
|
|
||||||
StatusCode::OK,
|
|
||||||
Json(UptimeResponse {
|
|
||||||
uptime_seconds: uptime,
|
|
||||||
}),
|
|
||||||
))
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn version_handler(
|
|
||||||
State(state): State<Arc<HttpServer>>,
|
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
|
||||||
Ok((
|
|
||||||
StatusCode::OK,
|
|
||||||
Json(serde_json::json!({ "version": state.info.version })),
|
|
||||||
))
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Serialize)]
|
#[derive(Serialize)]
|
||||||
struct StatsResponse {
|
struct StatsResponse {
|
||||||
version: &'static str,
|
version: &'static str,
|
||||||
@@ -552,7 +576,7 @@ struct StatsResponse {
|
|||||||
async fn stats_handler(
|
async fn stats_handler(
|
||||||
State(state): State<Arc<HttpServer>>,
|
State(state): State<Arc<HttpServer>>,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
let uptime = state.info.start_time.elapsed().as_secs();
|
let uptime = state.start_time.elapsed().as_secs();
|
||||||
let mut system = System::new();
|
let mut system = System::new();
|
||||||
system.refresh_cpu_all();
|
system.refresh_cpu_all();
|
||||||
let cpu_usage = system.global_cpu_usage();
|
let cpu_usage = system.global_cpu_usage();
|
||||||
@@ -564,7 +588,7 @@ async fn stats_handler(
|
|||||||
Ok((
|
Ok((
|
||||||
StatusCode::OK,
|
StatusCode::OK,
|
||||||
Json(StatsResponse {
|
Json(StatsResponse {
|
||||||
version: state.info.version,
|
version: state.version,
|
||||||
uptime_seconds: uptime,
|
uptime_seconds: uptime,
|
||||||
cpu_usage_percent: cpu_usage,
|
cpu_usage_percent: cpu_usage,
|
||||||
active_streams,
|
active_streams,
|
||||||
|
|||||||
@@ -25,6 +25,8 @@ pub enum HttpError {
|
|||||||
Unprocessable(String),
|
Unprocessable(String),
|
||||||
#[error("not acceptable: {0}")]
|
#[error("not acceptable: {0}")]
|
||||||
NotAcceptable(String),
|
NotAcceptable(String),
|
||||||
|
#[error("unsupported codec: {0}")]
|
||||||
|
WhepCodecError(String),
|
||||||
#[error("internal error")]
|
#[error("internal error")]
|
||||||
Internal,
|
Internal,
|
||||||
}
|
}
|
||||||
@@ -40,6 +42,7 @@ impl HttpError {
|
|||||||
Self::BadRequest(_) => StatusCode::BAD_REQUEST,
|
Self::BadRequest(_) => StatusCode::BAD_REQUEST,
|
||||||
Self::Unprocessable(_) => StatusCode::UNPROCESSABLE_ENTITY,
|
Self::Unprocessable(_) => StatusCode::UNPROCESSABLE_ENTITY,
|
||||||
Self::NotAcceptable(_) => StatusCode::NOT_ACCEPTABLE,
|
Self::NotAcceptable(_) => StatusCode::NOT_ACCEPTABLE,
|
||||||
|
Self::WhepCodecError(_) => StatusCode::UNSUPPORTED_MEDIA_TYPE,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -47,12 +50,14 @@ impl HttpError {
|
|||||||
impl IntoResponse for HttpError {
|
impl IntoResponse for HttpError {
|
||||||
fn into_response(self) -> Response<Body> {
|
fn into_response(self) -> Response<Body> {
|
||||||
let status = self.status();
|
let status = self.status();
|
||||||
// Only surface a message body for client (4xx) errors; keep an empty
|
let body = match self {
|
||||||
// body for 5xx so internal details aren't leaked.
|
// The WHEP client (frontend) reads this body as the rejected
|
||||||
let body = if status.is_client_error() {
|
// codec, so send it bare rather than the full error string.
|
||||||
self.to_string()
|
Self::WhepCodecError(codec) => codec,
|
||||||
} else {
|
// Only surface a message body for client (4xx) errors; keep an
|
||||||
String::new()
|
// empty body for 5xx so internal details aren't leaked.
|
||||||
|
e if status.is_client_error() => e.to_string(),
|
||||||
|
_ => String::new(),
|
||||||
};
|
};
|
||||||
(status, body).into_response()
|
(status, body).into_response()
|
||||||
}
|
}
|
||||||
|
|||||||
+33
-23
@@ -13,12 +13,10 @@ use migration::{Migrator, MigratorTrait};
|
|||||||
use sea_orm::Database;
|
use sea_orm::Database;
|
||||||
use tokio::{net::TcpListener, sync::Mutex, task::JoinSet};
|
use tokio::{net::TcpListener, sync::Mutex, task::JoinSet};
|
||||||
use tracing_subscriber::EnvFilter;
|
use tracing_subscriber::EnvFilter;
|
||||||
|
use uuid::Uuid;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
audio::OpusAudioFrame,
|
audio::OpusAudioFrame, codec::VideoFrame, http::HttpServer, webrtc_proxy::WebrtcProxy,
|
||||||
codec::VideoFrame,
|
|
||||||
http::{HttpServer, HttpServerConfig, ServerInfo},
|
|
||||||
webrtc_proxy::{WebRtcProxyConfig, WebrtcProxy},
|
|
||||||
};
|
};
|
||||||
|
|
||||||
const SERVER_VERSION: &str = env!("CARGO_PKG_VERSION");
|
const SERVER_VERSION: &str = env!("CARGO_PKG_VERSION");
|
||||||
@@ -29,11 +27,11 @@ mod hash;
|
|||||||
mod http;
|
mod http;
|
||||||
mod http_error;
|
mod http_error;
|
||||||
mod rtmp;
|
mod rtmp;
|
||||||
|
mod stream_session_2;
|
||||||
mod webrtc;
|
mod webrtc;
|
||||||
mod webrtc_ingest;
|
mod webrtc_ingest;
|
||||||
mod webrtc_proxy;
|
mod webrtc_proxy;
|
||||||
|
|
||||||
// #[derive(Debug)]
|
|
||||||
pub struct AppState {
|
pub struct AppState {
|
||||||
pub stream_sessions: Arc<DashMap<i32, StreamSession>>,
|
pub stream_sessions: Arc<DashMap<i32, StreamSession>>,
|
||||||
pub webrtc_proxy: WebrtcProxy,
|
pub webrtc_proxy: WebrtcProxy,
|
||||||
@@ -46,10 +44,26 @@ pub enum StreamCodec {
|
|||||||
AV1,
|
AV1,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl StreamCodec {
|
||||||
|
/// Map str0m's codec enum to ours; None for non-video codecs.
|
||||||
|
pub fn from_str0m(c: str0m::format::Codec) -> Option<Self> {
|
||||||
|
use str0m::format::Codec;
|
||||||
|
match c {
|
||||||
|
Codec::H264 => Some(Self::H264),
|
||||||
|
Codec::H265 => Some(Self::H265),
|
||||||
|
Codec::Av1 => Some(Self::AV1),
|
||||||
|
_ => None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub struct StreamSession {
|
pub struct StreamSession {
|
||||||
pub stream_key_id: i32,
|
pub stream_key_id: i32,
|
||||||
pub stream_key_label: String,
|
pub stream_key_label: String,
|
||||||
pub stream_key_user: String,
|
pub stream_key_user: String,
|
||||||
|
pub custom_id: Option<String>,
|
||||||
|
pub is_unlisted: bool,
|
||||||
|
pub password: Option<String>,
|
||||||
pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>,
|
pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>,
|
||||||
pub audio_channel: async_broadcast::Sender<Arc<OpusAudioFrame>>,
|
pub audio_channel: async_broadcast::Sender<Arc<OpusAudioFrame>>,
|
||||||
pub codec: Option<StreamCodec>,
|
pub codec: Option<StreamCodec>,
|
||||||
@@ -61,6 +75,9 @@ pub struct StreamSession {
|
|||||||
/// profile instead of all three H.264 profiles.
|
/// profile instead of all three H.264 profiles.
|
||||||
pub video_profile_level_id: Option<u32>,
|
pub video_profile_level_id: Option<u32>,
|
||||||
// ----
|
// ----
|
||||||
|
/// Unique id per publish. Lets a session's cleanup verify it is still the
|
||||||
|
/// live session for its key instead of clobbering a newer one.
|
||||||
|
pub session_id: Uuid,
|
||||||
pub started_at: DateTime<Utc>,
|
pub started_at: DateTime<Utc>,
|
||||||
pub active_clients: AtomicU32,
|
pub active_clients: AtomicU32,
|
||||||
}
|
}
|
||||||
@@ -85,12 +102,10 @@ async fn main() -> Result<(), Box<dyn Error>> {
|
|||||||
stream_session::Model::clean_unended_streams(&db)
|
stream_session::Model::clean_unended_streams(&db)
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let proxyconfig = WebRtcProxyConfig {
|
let proxyconfig = env::var("RTC_PORT")
|
||||||
proxy_port: env::var("RTC_PORT")
|
.unwrap_or("6969".into())
|
||||||
.unwrap_or("6969".into())
|
.parse()
|
||||||
.parse()
|
.expect("RTC_PORT needs to be a number (i32)");
|
||||||
.expect("RTC_PORT needs to be a number (i32)"),
|
|
||||||
};
|
|
||||||
let proxy = webrtc_proxy::WebrtcProxy::new(proxyconfig).await.unwrap();
|
let proxy = webrtc_proxy::WebrtcProxy::new(proxyconfig).await.unwrap();
|
||||||
|
|
||||||
let appstate = Arc::new(Mutex::new(AppState {
|
let appstate = Arc::new(Mutex::new(AppState {
|
||||||
@@ -98,11 +113,12 @@ async fn main() -> Result<(), Box<dyn Error>> {
|
|||||||
webrtc_proxy: proxy.clone(),
|
webrtc_proxy: proxy.clone(),
|
||||||
}));
|
}));
|
||||||
let (offer_tx, offer_rx) = tokio::sync::mpsc::channel::<(i32, i32, String)>(64);
|
let (offer_tx, offer_rx) = tokio::sync::mpsc::channel::<(i32, i32, String)>(64);
|
||||||
|
|
||||||
// Request_Id,
|
// Request_Id,
|
||||||
// String_Label,
|
// String_Label,
|
||||||
// Offer_body
|
// Offer_body
|
||||||
|
|
||||||
let (answer_tx, answer_rx) = broadcast::<(i32, Option<String>)>(64);
|
let (answer_tx, answer_rx) = broadcast::<(i32, Result<String, String>)>(64);
|
||||||
// Request_Id,
|
// Request_Id,
|
||||||
// Answer_body
|
// Answer_body
|
||||||
|
|
||||||
@@ -112,17 +128,12 @@ async fn main() -> Result<(), Box<dyn Error>> {
|
|||||||
appstate: appstate.clone(),
|
appstate: appstate.clone(),
|
||||||
request_count: std::sync::atomic::AtomicI32::new(0),
|
request_count: std::sync::atomic::AtomicI32::new(0),
|
||||||
db: db.clone(),
|
db: db.clone(),
|
||||||
config: HttpServerConfig {
|
signup_code: env::var("SIGNUP_CODE").unwrap_or_else(|_| {
|
||||||
signup_code: env::var("SIGNUP_CODE").unwrap_or_else(|_| {
|
warn!("SIGNUP_CODE not set; signup will be disabled");
|
||||||
warn!("SIGNUP_CODE not set; signup will be disabled");
|
String::new()
|
||||||
String::new()
|
|
||||||
}),
|
|
||||||
}
|
|
||||||
.into(),
|
|
||||||
info: Arc::new(ServerInfo {
|
|
||||||
version: SERVER_VERSION,
|
|
||||||
start_time: std::time::Instant::now(),
|
|
||||||
}),
|
}),
|
||||||
|
version: SERVER_VERSION,
|
||||||
|
start_time: std::time::Instant::now(),
|
||||||
};
|
};
|
||||||
|
|
||||||
let app = appstate.lock().await;
|
let app = appstate.lock().await;
|
||||||
@@ -156,7 +167,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
|
|||||||
res = workers.join_next() => {
|
res = workers.join_next() => {
|
||||||
if let Some(Err(e)) = res {
|
if let Some(Err(e)) = res {
|
||||||
tracing::error!("worker panicked: {:?}", e);
|
tracing::error!("worker panicked: {:?}", e);
|
||||||
// fs::File::
|
|
||||||
} else {
|
} else {
|
||||||
tracing::error!("a worker exited unexpectedly");
|
tracing::error!("a worker exited unexpectedly");
|
||||||
}
|
}
|
||||||
|
|||||||
+12
-50
@@ -15,6 +15,7 @@ use tokio::{
|
|||||||
time::timeout,
|
time::timeout,
|
||||||
};
|
};
|
||||||
use tracing::{debug, info, warn};
|
use tracing::{debug, info, warn};
|
||||||
|
use uuid::Uuid;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
StreamCodec, StreamSession,
|
StreamCodec, StreamSession,
|
||||||
@@ -133,7 +134,6 @@ impl Rtmp {
|
|||||||
// ponytail: fixed 512; per-stream dynamic sizing if memory ever matters.
|
// ponytail: fixed 512; per-stream dynamic sizing if memory ever matters.
|
||||||
let (mut video_tx, video_rx) = broadcast::<Arc<VideoFrame>>(512);
|
let (mut video_tx, video_rx) = broadcast::<Arc<VideoFrame>>(512);
|
||||||
let (mut audio_tx, audio_rx) = broadcast::<Arc<OpusAudioFrame>>(512);
|
let (mut audio_tx, audio_rx) = broadcast::<Arc<OpusAudioFrame>>(512);
|
||||||
// video_rx.cycle
|
|
||||||
|
|
||||||
video_tx.set_overflow(true);
|
video_tx.set_overflow(true);
|
||||||
audio_tx.set_overflow(true);
|
audio_tx.set_overflow(true);
|
||||||
@@ -312,11 +312,15 @@ impl Rtmp {
|
|||||||
stream_key_id: key.id,
|
stream_key_id: key.id,
|
||||||
stream_key_label: key.label,
|
stream_key_label: key.label,
|
||||||
stream_key_user: user.username,
|
stream_key_user: user.username,
|
||||||
|
custom_id: key.custom_id,
|
||||||
|
is_unlisted: key.is_unlisted,
|
||||||
|
password: key.password,
|
||||||
frame_channel: video_tx.clone(),
|
frame_channel: video_tx.clone(),
|
||||||
audio_channel: audio_tx.clone(),
|
audio_channel: audio_tx.clone(),
|
||||||
codec: None,
|
codec: None,
|
||||||
started_at: Utc::now(),
|
started_at: Utc::now(),
|
||||||
active_clients: 0.into(),
|
active_clients: 0.into(),
|
||||||
|
session_id: Uuid::new_v4(),
|
||||||
video_pt: None,
|
video_pt: None,
|
||||||
video_profile_level_id: None,
|
video_profile_level_id: None,
|
||||||
},
|
},
|
||||||
@@ -364,9 +368,6 @@ impl Rtmp {
|
|||||||
}
|
}
|
||||||
current_stream_key_id = None;
|
current_stream_key_id = None;
|
||||||
}
|
}
|
||||||
// TODO: We can totally replace the broadcast with a
|
|
||||||
// circular_buff
|
|
||||||
// Arc<Vec<ArcSwap<Frame>>>
|
|
||||||
ServerSessionEvent::VideoDataReceived {
|
ServerSessionEvent::VideoDataReceived {
|
||||||
data, timestamp, ..
|
data, timestamp, ..
|
||||||
} => match Self::parse_video_codec(&data) {
|
} => match Self::parse_video_codec(&data) {
|
||||||
@@ -380,52 +381,13 @@ impl Rtmp {
|
|||||||
}
|
}
|
||||||
codec_stamped = true;
|
codec_stamped = true;
|
||||||
}
|
}
|
||||||
match codec {
|
let p = parser.get_or_insert_with(|| match codec {
|
||||||
StreamCodec::H264 => {
|
StreamCodec::H264 => Box::new(H264CodecParser::new()),
|
||||||
let p = parser.get_or_insert_with(|| {
|
StreamCodec::H265 => Box::new(H265CodecParser::new()),
|
||||||
Box::new(H264CodecParser::new())
|
StreamCodec::AV1 => Box::new(Av1CodecParser::new()),
|
||||||
});
|
});
|
||||||
if let Some(frame) = p.parse(&data, timestamp.value)
|
if let Some(frame) = p.parse(&data, timestamp.value) {
|
||||||
{
|
video_tx.broadcast(Arc::new(frame)).await.ok();
|
||||||
video_tx.broadcast(Arc::new(frame)).await.ok();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
StreamCodec::H265 => {
|
|
||||||
let p = parser.get_or_insert_with(|| {
|
|
||||||
Box::new(H265CodecParser::new())
|
|
||||||
});
|
|
||||||
if let Some(frame) = p.parse(&data, timestamp.value)
|
|
||||||
{
|
|
||||||
video_tx.broadcast(Arc::new(frame)).await.ok();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
StreamCodec::AV1 => {
|
|
||||||
let p = parser.get_or_insert_with(|| {
|
|
||||||
Box::new(Av1CodecParser::new())
|
|
||||||
});
|
|
||||||
let pkt_type = data[0] & 0x0F;
|
|
||||||
match p.parse(&data, timestamp.value) {
|
|
||||||
Some(frame) => {
|
|
||||||
debug!(
|
|
||||||
pkt_type,
|
|
||||||
is_keyframe = frame.is_keyframe,
|
|
||||||
ts = frame.timestamp_ms,
|
|
||||||
bytes = frame.data.len(),
|
|
||||||
"AV1 frame → broadcast"
|
|
||||||
);
|
|
||||||
video_tx
|
|
||||||
.broadcast(Arc::new(frame))
|
|
||||||
.await
|
|
||||||
.ok();
|
|
||||||
}
|
|
||||||
None => {
|
|
||||||
debug!(
|
|
||||||
pkt_type,
|
|
||||||
"AV1 packet produced no frame (seq header or unknown type)"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Err(_err) => {
|
Err(_err) => {
|
||||||
|
|||||||
@@ -0,0 +1,18 @@
|
|||||||
|
use serde::Serialize;
|
||||||
|
|
||||||
|
use crate::StreamSession;
|
||||||
|
|
||||||
|
#[derive(Serialize)]
|
||||||
|
pub struct StreamUpdateData {
|
||||||
|
viewers: u32,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl StreamSession {
|
||||||
|
pub fn stream_update_data(&self) -> StreamUpdateData {
|
||||||
|
StreamUpdateData {
|
||||||
|
viewers: self
|
||||||
|
.active_clients
|
||||||
|
.load(std::sync::atomic::Ordering::Relaxed),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+76
-20
@@ -2,13 +2,18 @@ use bytes::Bytes;
|
|||||||
use dashmap::DashMap;
|
use dashmap::DashMap;
|
||||||
use entity::stream_key;
|
use entity::stream_key;
|
||||||
use sea_orm::{DatabaseConnection, EntityTrait};
|
use sea_orm::{DatabaseConnection, EntityTrait};
|
||||||
use std::{net::SocketAddr, sync::Arc, time::Instant};
|
use std::{
|
||||||
|
net::SocketAddr,
|
||||||
|
sync::Arc,
|
||||||
|
time::{Duration, Instant},
|
||||||
|
};
|
||||||
use tokio::{net::UdpSocket, sync::mpsc::Receiver};
|
use tokio::{net::UdpSocket, sync::mpsc::Receiver};
|
||||||
use tracing::{debug, error, info, warn};
|
use tracing::{debug, error, info, trace, warn};
|
||||||
|
|
||||||
use str0m::{
|
use str0m::{
|
||||||
Candidate, Event, Input, Output, Rtc,
|
Candidate, Event, Input, Output, Rtc,
|
||||||
change::SdpOffer,
|
change::SdpOffer,
|
||||||
|
channel::ChannelId,
|
||||||
media::{Frequency, MediaKind, MediaTime, Mid, Pt},
|
media::{Frequency, MediaKind, MediaTime, Mid, Pt},
|
||||||
net::{Protocol, Receive},
|
net::{Protocol, Receive},
|
||||||
};
|
};
|
||||||
@@ -18,11 +23,11 @@ use crate::{
|
|||||||
};
|
};
|
||||||
|
|
||||||
pub struct Webrtc {
|
pub struct Webrtc {
|
||||||
pub offer_rx: Receiver<(i32, i32, String)>,
|
pub offer_rx: Receiver<(i32, i32, String)>,
|
||||||
pub accept_tx: async_broadcast::Sender<(i32, Option<String>)>,
|
pub accept_tx: async_broadcast::Sender<(i32, Result<String, String>)>,
|
||||||
pub sessions_ref: Arc<DashMap<i32, StreamSession>>,
|
pub sessions_ref: Arc<DashMap<i32, StreamSession>>,
|
||||||
pub proxy: Arc<WebrtcProxy>,
|
pub proxy: Arc<WebrtcProxy>,
|
||||||
pub db: DatabaseConnection,
|
pub db: DatabaseConnection,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Webrtc {
|
impl Webrtc {
|
||||||
@@ -40,7 +45,10 @@ impl Webrtc {
|
|||||||
request_id,
|
request_id,
|
||||||
stream_id, "stream key not found in DB, rejecting offer"
|
stream_id, "stream key not found in DB, rejecting offer"
|
||||||
);
|
);
|
||||||
self.accept_tx.broadcast((request_id, None)).await.unwrap();
|
self.accept_tx
|
||||||
|
.broadcast((request_id, Err(String::new())))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if let Err(ref e) = stream_key {
|
if let Err(ref e) = stream_key {
|
||||||
@@ -48,13 +56,13 @@ impl Webrtc {
|
|||||||
request_id,
|
request_id,
|
||||||
stream_id, "DB error looking up stream key: {:?}", e
|
stream_id, "DB error looking up stream key: {:?}", e
|
||||||
);
|
);
|
||||||
self.accept_tx.broadcast((request_id, None)).await.unwrap();
|
self.accept_tx
|
||||||
|
.broadcast((request_id, Err(String::new())))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
// let socket = UdpSocket::bind("127.0.0.1:0").await.unwrap();
|
|
||||||
// let local_addr = socket.local_addr().unwrap();
|
|
||||||
|
|
||||||
let local_addr = self.proxy.public_addr();
|
let local_addr = self.proxy.public_addr();
|
||||||
|
|
||||||
let session = self.sessions_ref.get(&stream_id);
|
let session = self.sessions_ref.get(&stream_id);
|
||||||
@@ -134,7 +142,10 @@ impl Webrtc {
|
|||||||
Ok(sdp) => sdp,
|
Ok(sdp) => sdp,
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
warn!(request_id, stream_id, "malformed SDP offer: {:?}", e);
|
warn!(request_id, stream_id, "malformed SDP offer: {:?}", e);
|
||||||
self.accept_tx.broadcast((request_id, None)).await.unwrap();
|
self.accept_tx
|
||||||
|
.broadcast((request_id, Err(String::new())))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@@ -162,7 +173,7 @@ impl Webrtc {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
let answer_sdp = offer_answer.to_sdp_string();
|
let answer_sdp = offer_answer.to_sdp_string();
|
||||||
info!(request_id, "SDP answer:\n{}", answer_sdp);
|
trace!(request_id, "SDP answer:\n{}", answer_sdp);
|
||||||
|
|
||||||
// Detect the case where str0m couldn't match any video codec.
|
// Detect the case where str0m couldn't match any video codec.
|
||||||
// str0m serialises the m-line with an empty PT list, which is invalid SDP
|
// str0m serialises the m-line with an empty PT list, which is invalid SDP
|
||||||
@@ -179,7 +190,10 @@ impl Webrtc {
|
|||||||
"no video codec negotiated — browser likely doesn't support {:?}; rejecting offer",
|
"no video codec negotiated — browser likely doesn't support {:?}; rejecting offer",
|
||||||
stream_codec
|
stream_codec
|
||||||
);
|
);
|
||||||
self.accept_tx.broadcast((request_id, None)).await.unwrap();
|
self.accept_tx
|
||||||
|
.broadcast((request_id, Err(format!("{:?}", stream_codec))))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -193,7 +207,7 @@ impl Webrtc {
|
|||||||
|
|
||||||
debug!(request_id, "sending answer back");
|
debug!(request_id, "sending answer back");
|
||||||
self.accept_tx
|
self.accept_tx
|
||||||
.broadcast((request_id, Some(answer_sdp)))
|
.broadcast((request_id, Ok(answer_sdp)))
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
@@ -227,10 +241,34 @@ impl Webrtc {
|
|||||||
let mut video_pt = None;
|
let mut video_pt = None;
|
||||||
let mut audio_mid: Option<Mid> = None;
|
let mut audio_mid: Option<Mid> = None;
|
||||||
let mut audio_pt = None;
|
let mut audio_pt = None;
|
||||||
|
let mut channel_id = None;
|
||||||
let mut connected = false;
|
let mut connected = false;
|
||||||
let mut video_stream: Option<async_broadcast::Receiver<Arc<VideoFrame>>> = None;
|
let mut video_stream: Option<async_broadcast::Receiver<Arc<VideoFrame>>> = None;
|
||||||
let mut audio_stream: Option<async_broadcast::Receiver<Arc<OpusAudioFrame>>> = None;
|
let mut audio_stream: Option<async_broadcast::Receiver<Arc<OpusAudioFrame>>> = None;
|
||||||
let mut saw_keyframe = false;
|
let mut saw_keyframe = false;
|
||||||
|
// Update client with stream info through webrtc data channel ()
|
||||||
|
let mut tick = tokio::time::interval(Duration::from_millis(2000));
|
||||||
|
|
||||||
|
if let Some(ses) = sessions_ref.get(&stream_id) {
|
||||||
|
ses.active_clients
|
||||||
|
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
struct ActiveClientGuard {
|
||||||
|
sessions: Arc<DashMap<i32, StreamSession>>,
|
||||||
|
id: i32,
|
||||||
|
}
|
||||||
|
impl Drop for ActiveClientGuard {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if let Some(ses) = self.sessions.get(&self.id) {
|
||||||
|
ses.active_clients
|
||||||
|
.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let _active_client_guard = ActiveClientGuard {
|
||||||
|
sessions: sessions_ref.clone(),
|
||||||
|
id: stream_id,
|
||||||
|
};
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let deadline = loop {
|
let deadline = loop {
|
||||||
@@ -243,6 +281,9 @@ impl Webrtc {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
Ok(Output::Event(e)) => match e {
|
Ok(Output::Event(e)) => match e {
|
||||||
|
Event::ChannelOpen(channelId, _name) => {
|
||||||
|
channel_id = Some(channelId);
|
||||||
|
}
|
||||||
Event::MediaAdded(ma) => {
|
Event::MediaAdded(ma) => {
|
||||||
info!(stream_id, kind = ?ma.kind, mid = ?ma.mid, "MediaAdded");
|
info!(stream_id, kind = ?ma.kind, mid = ?ma.mid, "MediaAdded");
|
||||||
if ma.kind == MediaKind::Video {
|
if ma.kind == MediaKind::Video {
|
||||||
@@ -310,8 +351,8 @@ impl Webrtc {
|
|||||||
}
|
}
|
||||||
_ => {}
|
_ => {}
|
||||||
},
|
},
|
||||||
Err(e) => {
|
Err(_e) => {
|
||||||
error!("poll_output error (connection closing): {:?}", e);
|
// error!("poll_output error (connection closing): {:?}", e);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -422,6 +463,14 @@ impl Webrtc {
|
|||||||
Err(async_broadcast::RecvError::Overflowed(_)) => {}
|
Err(async_broadcast::RecvError::Overflowed(_)) => {}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
_interval = tick.tick() => {
|
||||||
|
if let Some(id) = channel_id
|
||||||
|
&& let Some(session) = sessions_ref.get(&stream_id)
|
||||||
|
&& let Ok(json) = serde_json::to_vec(&session.stream_update_data())
|
||||||
|
&& let Some(mut ch) = rtc.channel(id) {
|
||||||
|
let _ = ch.write(false, &json);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -487,6 +536,13 @@ impl Webrtc {
|
|||||||
{
|
{
|
||||||
warn!("RTP write error: {:?}", e);
|
warn!("RTP write error: {:?}", e);
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn write_channel_data(rtc: &mut Rtc, channel_id: ChannelId, data: &[u8]) {
|
||||||
|
if let Some(mut channel) = rtc.channel(channel_id)
|
||||||
|
&& let Err(e) = channel.write(false, data)
|
||||||
|
{
|
||||||
|
warn!("Channel write error: {:?}", e)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,4 +1,8 @@
|
|||||||
use std::{net::SocketAddr, sync::Arc, time::Instant};
|
use std::{
|
||||||
|
net::SocketAddr,
|
||||||
|
sync::Arc,
|
||||||
|
time::{Duration, Instant},
|
||||||
|
};
|
||||||
|
|
||||||
use async_broadcast::broadcast;
|
use async_broadcast::broadcast;
|
||||||
use axum::{
|
use axum::{
|
||||||
@@ -18,25 +22,32 @@ use str0m::{
|
|||||||
net::{Protocol, Receive},
|
net::{Protocol, Receive},
|
||||||
};
|
};
|
||||||
use tokio::{net::UdpSocket, sync::mpsc::Receiver};
|
use tokio::{net::UdpSocket, sync::mpsc::Receiver};
|
||||||
use tracing::{debug, error, info, warn};
|
use tracing::{debug, error, info, trace, warn};
|
||||||
|
use uuid::Uuid;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
StreamSession, audio::OpusAudioFrame, codec::VideoFrame, http::HttpServer,
|
StreamSession, audio::OpusAudioFrame, codec::VideoFrame, http::HttpServer,
|
||||||
http_error::HttpError,
|
http_error::HttpError,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/// Kill a WHIP publish if it delivers no media (video or audio) this long.
|
||||||
|
const NO_MEDIA_TIMEOUT: Duration = Duration::from_secs(30);
|
||||||
|
|
||||||
|
fn bearer_token(headers: &HeaderMap) -> Option<String> {
|
||||||
|
headers
|
||||||
|
.get(header::AUTHORIZATION)
|
||||||
|
.and_then(|v| v.to_str().ok())
|
||||||
|
.and_then(|v| v.strip_prefix("Bearer "))
|
||||||
|
.map(|s| s.trim().to_string())
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn handle_whip_injest_delete(
|
pub async fn handle_whip_injest_delete(
|
||||||
State(state): State<Arc<HttpServer>>,
|
State(state): State<Arc<HttpServer>>,
|
||||||
ConnectInfo(remote): ConnectInfo<SocketAddr>,
|
ConnectInfo(remote): ConnectInfo<SocketAddr>,
|
||||||
Path(_slug): Path<String>,
|
Path(_slug): Path<String>,
|
||||||
headers: HeaderMap,
|
headers: HeaderMap,
|
||||||
) -> Result<impl IntoResponse, HttpError> {
|
) -> Result<impl IntoResponse, HttpError> {
|
||||||
let token = headers
|
let Some(token) = bearer_token(&headers) else {
|
||||||
.get(header::AUTHORIZATION)
|
|
||||||
.and_then(|v| v.to_str().ok())
|
|
||||||
.and_then(|v| v.strip_prefix("Bearer "))
|
|
||||||
.map(|s| s.trim().to_string());
|
|
||||||
let Some(token) = token else {
|
|
||||||
warn!("Whip: missing bearer token from {}", remote);
|
warn!("Whip: missing bearer token from {}", remote);
|
||||||
return Err(HttpError::Unauthorized);
|
return Err(HttpError::Unauthorized);
|
||||||
};
|
};
|
||||||
@@ -58,6 +69,15 @@ pub async fn handle_whip_injest_delete(
|
|||||||
}
|
}
|
||||||
|
|
||||||
state.appstate.lock().await.stream_sessions.remove(&key.id);
|
state.appstate.lock().await.stream_sessions.remove(&key.id);
|
||||||
|
// Drop the trickle channel too, so PATCHes to a deleted session 404 and
|
||||||
|
// the lingering detach task's own cleanup can't reach a newer session.
|
||||||
|
state
|
||||||
|
.appstate
|
||||||
|
.lock()
|
||||||
|
.await
|
||||||
|
.webrtc_proxy
|
||||||
|
.trickle_tx
|
||||||
|
.remove(&key.id);
|
||||||
|
|
||||||
Ok(StatusCode::OK)
|
Ok(StatusCode::OK)
|
||||||
}
|
}
|
||||||
@@ -79,10 +99,13 @@ pub async fn handle_whip_injest_patch(
|
|||||||
|
|
||||||
// Look up the trickle sender for this session.
|
// Look up the trickle sender for this session.
|
||||||
let trickle_map = state.appstate.lock().await.webrtc_proxy.trickle_tx.clone();
|
let trickle_map = state.appstate.lock().await.webrtc_proxy.trickle_tx.clone();
|
||||||
let tx = trickle_map.get(&stream_key_id).ok_or_else(|| {
|
let tx = trickle_map
|
||||||
warn!(stream_key_id, "Whip PATCH: no trickle channel for session");
|
.get(&stream_key_id)
|
||||||
HttpError::NotFound
|
.map(|e| e.value().1.clone())
|
||||||
})?;
|
.ok_or_else(|| {
|
||||||
|
warn!(stream_key_id, "Whip PATCH: no trickle channel for session");
|
||||||
|
HttpError::NotFound
|
||||||
|
})?;
|
||||||
|
|
||||||
debug!(stream_key_id, %body, "Whip PATCH: forwarding trickle candidate");
|
debug!(stream_key_id, %body, "Whip PATCH: forwarding trickle candidate");
|
||||||
tx.send(body).ok();
|
tx.send(body).ok();
|
||||||
@@ -106,12 +129,7 @@ pub async fn handle_whip_injest(
|
|||||||
};
|
};
|
||||||
let public_addr = state.appstate.lock().await.webrtc_proxy.public_addr();
|
let public_addr = state.appstate.lock().await.webrtc_proxy.public_addr();
|
||||||
|
|
||||||
let token = headers
|
let Some(token) = bearer_token(&headers) else {
|
||||||
.get(header::AUTHORIZATION)
|
|
||||||
.and_then(|v| v.to_str().ok())
|
|
||||||
.and_then(|v| v.strip_prefix("Bearer "))
|
|
||||||
.map(|s| s.trim().to_string());
|
|
||||||
let Some(token) = token else {
|
|
||||||
warn!("Whip: missing bearer token from {}", remote);
|
warn!("Whip: missing bearer token from {}", remote);
|
||||||
return err(StatusCode::UNAUTHORIZED, "missing bearer token");
|
return err(StatusCode::UNAUTHORIZED, "missing bearer token");
|
||||||
};
|
};
|
||||||
@@ -173,17 +191,13 @@ pub async fn handle_whip_injest(
|
|||||||
cc.clear();
|
cc.clear();
|
||||||
cc.enable_opus(true);
|
cc.enable_opus(true);
|
||||||
match &stream_codec {
|
match &stream_codec {
|
||||||
Some(crate::StreamCodec::H264) => {
|
Some(codec) => {
|
||||||
info!("Whip: enabling H.264 codec");
|
info!("Whip: enabling {codec:?} codec");
|
||||||
cc.enable_h264(true);
|
match codec {
|
||||||
}
|
crate::StreamCodec::H264 => cc.enable_h264(true),
|
||||||
Some(crate::StreamCodec::H265) => {
|
crate::StreamCodec::H265 => cc.enable_h265(true),
|
||||||
info!("Whip: enabling H.265 codec");
|
crate::StreamCodec::AV1 => cc.enable_av1(true),
|
||||||
cc.enable_h265(true);
|
}
|
||||||
}
|
|
||||||
Some(crate::StreamCodec::AV1) => {
|
|
||||||
info!("Whip: enabling AV1 codec");
|
|
||||||
cc.enable_av1(true);
|
|
||||||
}
|
}
|
||||||
None => {
|
None => {
|
||||||
warn!("Whip: no video codec detected in offer, enabling H.264 as fallback");
|
warn!("Whip: no video codec detected in offer, enabling H.264 as fallback");
|
||||||
@@ -276,6 +290,9 @@ pub async fn handle_whip_injest(
|
|||||||
let (socket, rx) = state.appstate.lock().await.webrtc_proxy.add_client(ufrag);
|
let (socket, rx) = state.appstate.lock().await.webrtc_proxy.add_client(ufrag);
|
||||||
|
|
||||||
let stream_sessions = state.appstate.lock().await.stream_sessions.clone();
|
let stream_sessions = state.appstate.lock().await.stream_sessions.clone();
|
||||||
|
// Unique per publish: cleanup only removes the session this task created,
|
||||||
|
// never a newer one that re-published on the same stream key.
|
||||||
|
let session_id = Uuid::new_v4();
|
||||||
let (mut video_tx, video_rx) = broadcast::<Arc<VideoFrame>>(32);
|
let (mut video_tx, video_rx) = broadcast::<Arc<VideoFrame>>(32);
|
||||||
let (mut audio_tx, audio_rx) = broadcast::<Arc<OpusAudioFrame>>(32);
|
let (mut audio_tx, audio_rx) = broadcast::<Arc<OpusAudioFrame>>(32);
|
||||||
// Never block the ingest loop on slow/missing viewers: overwrite the
|
// Never block the ingest loop on slow/missing viewers: overwrite the
|
||||||
@@ -295,9 +312,13 @@ pub async fn handle_whip_injest(
|
|||||||
stream_key_id: key.id,
|
stream_key_id: key.id,
|
||||||
stream_key_label: key.label,
|
stream_key_label: key.label,
|
||||||
stream_key_user: user.username,
|
stream_key_user: user.username,
|
||||||
|
custom_id: key.custom_id,
|
||||||
|
is_unlisted: key.is_unlisted,
|
||||||
|
password: key.password,
|
||||||
frame_channel: video_tx,
|
frame_channel: video_tx,
|
||||||
audio_channel: audio_tx,
|
audio_channel: audio_tx,
|
||||||
codec: negotiated_codec,
|
codec: negotiated_codec,
|
||||||
|
session_id,
|
||||||
started_at: Utc::now(),
|
started_at: Utc::now(),
|
||||||
active_clients: 0.into(),
|
active_clients: 0.into(),
|
||||||
video_pt,
|
video_pt,
|
||||||
@@ -313,7 +334,7 @@ pub async fn handle_whip_injest(
|
|||||||
// initial offer. We forward them to the Rtc task for add_remote_candidate.
|
// initial offer. We forward them to the Rtc task for add_remote_candidate.
|
||||||
let (trickle_tx, trickle_rx) = tokio::sync::mpsc::unbounded_channel();
|
let (trickle_tx, trickle_rx) = tokio::sync::mpsc::unbounded_channel();
|
||||||
let trickle_map = state.appstate.lock().await.webrtc_proxy.trickle_tx.clone();
|
let trickle_map = state.appstate.lock().await.webrtc_proxy.trickle_tx.clone();
|
||||||
trickle_map.insert(key.id, trickle_tx);
|
trickle_map.insert(key.id, (session_id, trickle_tx));
|
||||||
|
|
||||||
let db = state.db.clone();
|
let db = state.db.clone();
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
@@ -328,10 +349,12 @@ pub async fn handle_whip_injest(
|
|||||||
public_addr,
|
public_addr,
|
||||||
video_rx,
|
video_rx,
|
||||||
audio_rx,
|
audio_rx,
|
||||||
|
session_id,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
// Clean up trickle channel when done.
|
// Drop only our own trickle channel; a newer session on the same key
|
||||||
trickle_map.remove(&key.id);
|
// must keep its own.
|
||||||
|
trickle_map.remove_if(&key.id, |_, (sid, _)| *sid == session_id);
|
||||||
});
|
});
|
||||||
|
|
||||||
stream_session::Model::create_stream_session(&state.db, key.id, Local::now().into())
|
stream_session::Model::create_stream_session(&state.db, key.id, Local::now().into())
|
||||||
@@ -343,12 +366,6 @@ pub async fn handle_whip_injest(
|
|||||||
.status(StatusCode::CREATED)
|
.status(StatusCode::CREATED)
|
||||||
.header(header::CONTENT_TYPE, "application/sdp")
|
.header(header::CONTENT_TYPE, "application/sdp")
|
||||||
.header(header::LOCATION, &location)
|
.header(header::LOCATION, &location)
|
||||||
// Provide a STUN server via Link header so OBS can gather
|
|
||||||
// ICE candidates even without explicit STUN configuration.
|
|
||||||
.header(
|
|
||||||
header::LINK,
|
|
||||||
"<stun:stun.l.google.com:19302>; rel=\"ice-server\"",
|
|
||||||
)
|
|
||||||
.body(axum::body::Body::from(answer_sdp))
|
.body(axum::body::Body::from(answer_sdp))
|
||||||
.unwrap()
|
.unwrap()
|
||||||
}
|
}
|
||||||
@@ -364,11 +381,26 @@ async fn detach_inject_rtc(
|
|||||||
local_addr: SocketAddr,
|
local_addr: SocketAddr,
|
||||||
video_rx: async_broadcast::Receiver<Arc<VideoFrame>>,
|
video_rx: async_broadcast::Receiver<Arc<VideoFrame>>,
|
||||||
audio_rx: async_broadcast::Receiver<Arc<OpusAudioFrame>>,
|
audio_rx: async_broadcast::Receiver<Arc<OpusAudioFrame>>,
|
||||||
|
session_id: Uuid,
|
||||||
) {
|
) {
|
||||||
let cleanup = || {
|
let cleanup = || {
|
||||||
let sessions_ref = sessions_ref.clone();
|
let sessions_ref = sessions_ref.clone();
|
||||||
let db = db.clone();
|
let db = db.clone();
|
||||||
async move {
|
async move {
|
||||||
|
// Only clean up the session this task created. If the user has
|
||||||
|
// already opened another stream on the same key (e.g. DELETE then
|
||||||
|
// immediate republish), that newer session must not be touched.
|
||||||
|
let is_ours = sessions_ref
|
||||||
|
.get(&stream_key_id)
|
||||||
|
.is_some_and(|s| s.session_id == session_id);
|
||||||
|
if !is_ours {
|
||||||
|
debug!(
|
||||||
|
stream_key_id,
|
||||||
|
"skipping cleanup: session replaced by a newer stream on this key"
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
debug!("Cleaning up {:?}", &stream_key_id);
|
||||||
sessions_ref.remove(&stream_key_id);
|
sessions_ref.remove(&stream_key_id);
|
||||||
if let Ok(Some(s)) =
|
if let Ok(Some(s)) =
|
||||||
stream_session::Model::get_active_by_stream_key_id(&db, stream_key_id).await
|
stream_session::Model::get_active_by_stream_key_id(&db, stream_key_id).await
|
||||||
@@ -384,7 +416,7 @@ async fn detach_inject_rtc(
|
|||||||
let mut audio_mid: Option<Mid> = None;
|
let mut audio_mid: Option<Mid> = None;
|
||||||
let mut video_tx: Option<async_broadcast::Sender<Arc<VideoFrame>>> = None;
|
let mut video_tx: Option<async_broadcast::Sender<Arc<VideoFrame>>> = None;
|
||||||
let mut audio_tx: Option<async_broadcast::Sender<Arc<OpusAudioFrame>>> = None;
|
let mut audio_tx: Option<async_broadcast::Sender<Arc<OpusAudioFrame>>> = None;
|
||||||
let mut _connected = false;
|
let mut disconnect_timer: Option<Instant> = None;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// Blankly using them so they don't drop (like RTMP).
|
// Blankly using them so they don't drop (like RTMP).
|
||||||
@@ -395,7 +427,8 @@ async fn detach_inject_rtc(
|
|||||||
match rtc.poll_output() {
|
match rtc.poll_output() {
|
||||||
Ok(Output::Timeout(t)) => break t,
|
Ok(Output::Timeout(t)) => break t,
|
||||||
Ok(Output::Transmit(t)) => {
|
Ok(Output::Transmit(t)) => {
|
||||||
debug!(
|
// The keep alive loop, by default is every 1 second.
|
||||||
|
trace!(
|
||||||
"Whip TX: {} bytes → {}:{}",
|
"Whip TX: {} bytes → {}:{}",
|
||||||
t.contents.len(),
|
t.contents.len(),
|
||||||
t.destination.ip(),
|
t.destination.ip(),
|
||||||
@@ -406,6 +439,22 @@ async fn detach_inject_rtc(
|
|||||||
cleanup().await;
|
cleanup().await;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// 30 Sec time out if disconnected, then we clean up and disconnect.
|
||||||
|
if let Some(instant_since_disconnet) = disconnect_timer
|
||||||
|
&& Instant::now()
|
||||||
|
.duration_since(instant_since_disconnet)
|
||||||
|
.as_secs()
|
||||||
|
> NO_MEDIA_TIMEOUT.as_secs()
|
||||||
|
{
|
||||||
|
info!(
|
||||||
|
"WHIP connection on stream_key_id: {:?}, has been disconnected for {:?} seconds. Cleaning up and destroying the connection,",
|
||||||
|
stream_key_id,
|
||||||
|
NO_MEDIA_TIMEOUT.as_secs()
|
||||||
|
);
|
||||||
|
cleanup().await;
|
||||||
|
rtc.disconnect();
|
||||||
|
return;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Ok(Output::Event(e)) => match e {
|
Ok(Output::Event(e)) => match e {
|
||||||
Event::MediaAdded(ma) => {
|
Event::MediaAdded(ma) => {
|
||||||
@@ -445,10 +494,19 @@ async fn detach_inject_rtc(
|
|||||||
}
|
}
|
||||||
Event::IceConnectionStateChange(state) => {
|
Event::IceConnectionStateChange(state) => {
|
||||||
info!(stream_key_id, ?state, "Whip ICE state change");
|
info!(stream_key_id, ?state, "Whip ICE state change");
|
||||||
if matches!(state, str0m::IceConnectionState::Disconnected) {
|
match state {
|
||||||
info!("Whip ICE disconnected, closing connection");
|
str0m::IceConnectionState::Disconnected => {
|
||||||
cleanup().await;
|
info!(
|
||||||
return;
|
"Whip ICE disconnected... (State changed to Disconnected for {:?}) ((This is usually due to network jitter))",
|
||||||
|
&stream_key_id
|
||||||
|
);
|
||||||
|
disconnect_timer = Some(Instant::now());
|
||||||
|
}
|
||||||
|
str0m::IceConnectionState::Connected
|
||||||
|
| str0m::IceConnectionState::Completed => {
|
||||||
|
disconnect_timer = None;
|
||||||
|
}
|
||||||
|
_ => {}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Event::Connected => {
|
Event::Connected => {
|
||||||
@@ -460,7 +518,6 @@ async fn detach_inject_rtc(
|
|||||||
video_tx = Some(session.frame_channel.clone());
|
video_tx = Some(session.frame_channel.clone());
|
||||||
audio_tx = Some(session.audio_channel.clone());
|
audio_tx = Some(session.audio_channel.clone());
|
||||||
}
|
}
|
||||||
_connected = true;
|
|
||||||
}
|
}
|
||||||
_ => {}
|
_ => {}
|
||||||
},
|
},
|
||||||
@@ -498,13 +555,21 @@ async fn detach_inject_rtc(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
result = rx.recv() => {
|
// No UDP input at all for 30s: closing connection.
|
||||||
let Some((data, from)) = result else {
|
// (ICE keepalives still arrive when only the encoder is stalled,
|
||||||
info!(stream_key_id, "Whip proxy channel closed, cleaning up");
|
// so this fires on a dead peer, not a paused one.)
|
||||||
cleanup().await;
|
result = tokio::time::timeout(NO_MEDIA_TIMEOUT, rx.recv()) => {
|
||||||
return;
|
let Ok(Some((data, from))) = result else {
|
||||||
};
|
if result.is_err() {
|
||||||
debug!(
|
warn!(stream_key_id, "Whip: no input for {NO_MEDIA_TIMEOUT:?}, closing connection");
|
||||||
|
rtc.disconnect();
|
||||||
|
} else {
|
||||||
|
info!(stream_key_id, "Whip proxy channel closed, cleaning up");
|
||||||
|
}
|
||||||
|
cleanup().await;
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
trace!(
|
||||||
stream_key_id,
|
stream_key_id,
|
||||||
len = data.len(),
|
len = data.len(),
|
||||||
%from,
|
%from,
|
||||||
@@ -537,22 +602,14 @@ async fn detach_inject_rtc(
|
|||||||
pub fn extract_negotiated_codec_info(
|
pub fn extract_negotiated_codec_info(
|
||||||
answer: &str0m::change::SdpAnswer,
|
answer: &str0m::change::SdpAnswer,
|
||||||
) -> (Option<crate::StreamCodec>, Option<u8>, Option<u32>) {
|
) -> (Option<crate::StreamCodec>, Option<u8>, Option<u32>) {
|
||||||
use str0m::format::Codec;
|
|
||||||
|
|
||||||
for line in answer.media_lines.iter() {
|
for line in answer.media_lines.iter() {
|
||||||
// Iterate rtp_params on every m-line; video check via
|
|
||||||
// codec.is_video() avoids needing str0m's private MediaType.
|
|
||||||
for p in line.rtp_params() {
|
for p in line.rtp_params() {
|
||||||
if p.spec().codec.is_video() {
|
if p.spec().codec.is_video() {
|
||||||
let pt = Some(*p.pt());
|
return (
|
||||||
let profile = p.spec().format.profile_level_id;
|
crate::StreamCodec::from_str0m(p.spec().codec),
|
||||||
let codec = match p.spec().codec {
|
Some(*p.pt()),
|
||||||
Codec::H264 => Some(crate::StreamCodec::H264),
|
p.spec().format.profile_level_id,
|
||||||
Codec::H265 => Some(crate::StreamCodec::H265),
|
);
|
||||||
Codec::Av1 => Some(crate::StreamCodec::AV1),
|
|
||||||
_ => None,
|
|
||||||
};
|
|
||||||
return (codec, pt, profile);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -561,24 +618,14 @@ pub fn extract_negotiated_codec_info(
|
|||||||
|
|
||||||
/// Extract the video codec from an SDP offer's media lines.
|
/// Extract the video codec from an SDP offer's media lines.
|
||||||
///
|
///
|
||||||
/// Walks every m= line's `rtp_params()`, maps str0m's `Codec` enum
|
/// Walks every m= line's `rtp_params()`; returns the first video codec found.
|
||||||
/// to our `StreamCodec`. Returns the first video codec found.
|
|
||||||
pub fn video_codec_from_sdp_offer(
|
pub fn video_codec_from_sdp_offer(
|
||||||
sdp_offer: &str0m::change::SdpOffer,
|
sdp_offer: &str0m::change::SdpOffer,
|
||||||
) -> Option<crate::StreamCodec> {
|
) -> Option<crate::StreamCodec> {
|
||||||
use str0m::format::Codec;
|
|
||||||
|
|
||||||
// SdpOffer derefs to Sdp, which has pub media_lines.
|
|
||||||
// MediaLine is pub(crate) in str0m, but we can call its pub methods.
|
|
||||||
for line in sdp_offer.media_lines.iter() {
|
for line in sdp_offer.media_lines.iter() {
|
||||||
for p in line.rtp_params() {
|
for p in line.rtp_params() {
|
||||||
if p.spec().codec.is_video() {
|
if p.spec().codec.is_video() {
|
||||||
return match p.spec().codec {
|
return crate::StreamCodec::from_str0m(p.spec().codec);
|
||||||
Codec::H264 => Some(crate::StreamCodec::H264),
|
|
||||||
Codec::H265 => Some(crate::StreamCodec::H265),
|
|
||||||
Codec::Av1 => Some(crate::StreamCodec::AV1),
|
|
||||||
_ => None,
|
|
||||||
};
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,11 +11,8 @@ use tokio::{
|
|||||||
net::UdpSocket,
|
net::UdpSocket,
|
||||||
sync::mpsc::{self, Receiver},
|
sync::mpsc::{self, Receiver},
|
||||||
};
|
};
|
||||||
use tracing::{debug, error, info, warn};
|
use tracing::{debug, error, info, trace, warn};
|
||||||
|
use uuid::Uuid;
|
||||||
pub struct WebRtcProxyConfig {
|
|
||||||
pub proxy_port: i32,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
pub struct WebrtcProxy {
|
pub struct WebrtcProxy {
|
||||||
@@ -24,22 +21,23 @@ pub struct WebrtcProxy {
|
|||||||
socket: Arc<UdpSocket>,
|
socket: Arc<UdpSocket>,
|
||||||
public_addr: SocketAddr,
|
public_addr: SocketAddr,
|
||||||
/// Trickle-ICE candidate channels for WHIP ingest.
|
/// Trickle-ICE candidate channels for WHIP ingest.
|
||||||
/// Keyed by stream_key_id; sender stored here so PATCH handler
|
/// Keyed by stream_key_id; each entry is tagged with the owning session's
|
||||||
/// can forward candidates to the detach task.
|
/// id so cleanup can remove only its own entry and never a newer session
|
||||||
pub trickle_tx: Arc<DashMap<i32, tokio::sync::mpsc::UnboundedSender<String>>>,
|
/// that re-published on the same key.
|
||||||
|
pub trickle_tx: Arc<DashMap<i32, (Uuid, tokio::sync::mpsc::UnboundedSender<String>)>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
const STUN_MAGIC: u32 = 0x2112A442;
|
const STUN_MAGIC: u32 = 0x2112A442;
|
||||||
|
|
||||||
impl WebrtcProxy {
|
impl WebrtcProxy {
|
||||||
pub async fn new(config: WebRtcProxyConfig) -> Result<Self, Box<dyn Error>> {
|
pub async fn new(proxy_port: i32) -> Result<Self, Box<dyn Error>> {
|
||||||
let sock = UdpSocket::bind(format!("0.0.0.0:{}", config.proxy_port)).await?;
|
let sock = UdpSocket::bind(format!("0.0.0.0:{}", proxy_port)).await?;
|
||||||
let port = sock.local_addr()?.port();
|
let port = sock.local_addr()?.port();
|
||||||
|
|
||||||
let public_ip = match env::var("PUBLIC_DOMAIN")
|
let public_ip = match env::var("PUBLIC_DOMAIN")
|
||||||
.ok()
|
.ok()
|
||||||
.filter(|s| !s.trim().is_empty())
|
.filter(|s| !s.trim().is_empty())
|
||||||
{
|
{
|
||||||
Some(domain) => {
|
Some(domain) => {
|
||||||
let ip = resolve_domain(&domain).await?;
|
let ip = resolve_domain(&domain).await?;
|
||||||
info!(%domain, %ip, "resolved PUBLIC_DOMAIN for WebRTC candidates");
|
info!(%domain, %ip, "resolved PUBLIC_DOMAIN for WebRTC candidates");
|
||||||
@@ -53,12 +51,14 @@ impl WebrtcProxy {
|
|||||||
// IP_UNICAST_IF), which makes checks to 127.0.0.1 vanish.
|
// IP_UNICAST_IF), which makes checks to 127.0.0.1 vanish.
|
||||||
match default_iface_ipv4() {
|
match default_iface_ipv4() {
|
||||||
Some(ip) => {
|
Some(ip) => {
|
||||||
info!(%ip, "using default interface IP for WebRTC candidates (debug build)");
|
info!(%ip, "using default interface IP for WebRTC candidates (debug build)");
|
||||||
ip
|
ip
|
||||||
}
|
}
|
||||||
None => {
|
None => {
|
||||||
warn!("no non-loopback IPv4 interface found, falling back to 127.0.0.1");
|
warn!(
|
||||||
IpAddr::from([127, 0, 0, 1])
|
"no non-loopback IPv4 interface found, falling back to 127.0.0.1"
|
||||||
|
);
|
||||||
|
IpAddr::from([127, 0, 0, 1])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
@@ -97,7 +97,7 @@ impl WebrtcProxy {
|
|||||||
|
|
||||||
// By addr
|
// By addr
|
||||||
if let Some(tx) = by_addr.get(&from) {
|
if let Some(tx) = by_addr.get(&from) {
|
||||||
debug!(
|
trace!(
|
||||||
"proxy: routing {} bytes by addr {}:{} → channel",
|
"proxy: routing {} bytes by addr {}:{} → channel",
|
||||||
b,
|
b,
|
||||||
from.ip(),
|
from.ip(),
|
||||||
@@ -120,7 +120,7 @@ impl WebrtcProxy {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let Some((part1, part2)) = self::WebrtcProxy::ufrag_pair(&data) else {
|
let Some((part1, part2)) = self::WebrtcProxy::ufrag_pair(&data) else {
|
||||||
debug!("huh, packet isnt stun or added as client.");
|
trace!("huh, packet isnt stun or added as client.");
|
||||||
continue;
|
continue;
|
||||||
};
|
};
|
||||||
// Try both parts of the STUN username — the first packet
|
// Try both parts of the STUN username — the first packet
|
||||||
@@ -163,10 +163,6 @@ impl WebrtcProxy {
|
|||||||
self.clients_ufrag.insert(ufrag, tx);
|
self.clients_ufrag.insert(ufrag, tx);
|
||||||
(self.socket.clone(), rx)
|
(self.socket.clone(), rx)
|
||||||
}
|
}
|
||||||
pub fn local_addr(&self) -> SocketAddr {
|
|
||||||
self.socket.local_addr().unwrap()
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn public_addr(&self) -> SocketAddr {
|
pub fn public_addr(&self) -> SocketAddr {
|
||||||
self.public_addr
|
self.public_addr
|
||||||
}
|
}
|
||||||
@@ -178,33 +174,37 @@ impl WebrtcProxy {
|
|||||||
if magic != STUN_MAGIC {
|
if magic != STUN_MAGIC {
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
// attribies start at 20
|
let value = stun_attributes(b).find(|(t, _)| *t == 0x0006)?.1;
|
||||||
let mut pos = 20usize;
|
let value = std::str::from_utf8(value).ok()?;
|
||||||
while (pos + 4) <= b.len() {
|
let mut parts = value.split(':');
|
||||||
let attr_type: u16 = u16::from_be_bytes(b[pos..pos + 2].try_into().ok()?);
|
Some((
|
||||||
let attr_len: u16 = u16::from_be_bytes(b[pos + 2..pos + 4].try_into().ok()?);
|
parts.next()?.to_string(),
|
||||||
pos += 4;
|
parts.next().map(|s| s.to_string()),
|
||||||
if attr_type == 0x0006 {
|
))
|
||||||
let value =
|
|
||||||
std::str::from_utf8(b[pos..pos + (attr_len as usize)].try_into().ok()?).ok()?;
|
|
||||||
let mut parts = value.split(':');
|
|
||||||
let first = parts.next()?.to_string();
|
|
||||||
let second = parts.next().map(|s| s.to_string());
|
|
||||||
return Some((first, second));
|
|
||||||
}
|
|
||||||
pos += (attr_len as usize + 3) & !3;
|
|
||||||
}
|
|
||||||
None
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Iterate `(attr_type, attr_value)` pairs over a STUN message's attributes
|
||||||
|
/// (RFC 5389 §15: header is 20 bytes, each attribute padded to a 4-byte boundary).
|
||||||
|
fn stun_attributes(data: &[u8]) -> impl Iterator<Item = (u16, &[u8])> {
|
||||||
|
let mut pos = 20usize;
|
||||||
|
std::iter::from_fn(move || {
|
||||||
|
let attr_type = u16::from_be_bytes(data.get(pos..pos + 2)?.try_into().ok()?);
|
||||||
|
let attr_len = u16::from_be_bytes(data.get(pos + 2..pos + 4)?.try_into().ok()?) as usize;
|
||||||
|
pos += 4;
|
||||||
|
let value = data.get(pos..pos + attr_len)?;
|
||||||
|
pos += (attr_len + 3) & !3;
|
||||||
|
Some((attr_type, value))
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
/// IP of the interface holding the default route, via a UDP connect() trick:
|
/// IP of the interface holding the default route, via a UDP connect() trick:
|
||||||
/// connect() only does a route lookup (no packets sent), so the kernel binds
|
/// connect() only does a route lookup (no packets sent), so the kernel binds
|
||||||
/// the source address the OS would use for outbound traffic.
|
/// the source address the OS would use for outbound traffic.
|
||||||
fn default_iface_ipv4() -> Option<IpAddr> {
|
fn default_iface_ipv4() -> Option<IpAddr> {
|
||||||
let sock = std::net::UdpSocket::bind("0.0.0.0:0").ok()?;
|
let sock = std::net::UdpSocket::bind("0.0.0.0:0").ok()?;
|
||||||
sock.connect("8.8.8.8:9").ok()?;
|
sock.connect("8.8.8.8:9").ok()?;
|
||||||
sock.local_addr().ok().map(|a| a.ip())
|
sock.local_addr().ok().map(|a| a.ip())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn resolve_domain(domain: &str) -> Result<std::net::IpAddr, Box<dyn Error>> {
|
async fn resolve_domain(domain: &str) -> Result<std::net::IpAddr, Box<dyn Error>> {
|
||||||
@@ -249,27 +249,17 @@ fn parse_xor_mapped_address(data: &[u8]) -> Option<std::net::IpAddr> {
|
|||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut pos = 20usize;
|
let value = stun_attributes(data).find(|(t, _)| *t == 0x0020)?.1;
|
||||||
while pos + 4 <= data.len() {
|
if value.len() < 8 {
|
||||||
let attr_type = u16::from_be_bytes(data[pos..pos + 2].try_into().ok()?);
|
return None;
|
||||||
let attr_len = u16::from_be_bytes(data[pos + 2..pos + 4].try_into().ok()?) as usize;
|
}
|
||||||
pos += 4;
|
// byte 0: reserved, byte 1: family (0x01=IPv4, 0x02=IPv6)
|
||||||
if pos + attr_len > data.len() {
|
if value[1] == 0x01 {
|
||||||
break;
|
let x_addr = u32::from_be_bytes(value[4..8].try_into().ok()?);
|
||||||
}
|
Some(std::net::IpAddr::V4(std::net::Ipv4Addr::from(
|
||||||
if attr_type == 0x0020 && attr_len >= 8 {
|
x_addr ^ magic,
|
||||||
// byte 0: reserved, byte 1: family (0x01=IPv4, 0x02=IPv6)
|
)))
|
||||||
let family = data[pos + 1];
|
} else {
|
||||||
let x_port = u16::from_be_bytes(data[pos + 2..pos + 4].try_into().ok()?);
|
None
|
||||||
let _ = x_port ^ (STUN_MAGIC >> 16) as u16; // port (unused here)
|
|
||||||
|
|
||||||
if family == 0x01 {
|
|
||||||
let x_addr = u32::from_be_bytes(data[pos + 4..pos + 8].try_into().ok()?);
|
|
||||||
let addr = x_addr ^ STUN_MAGIC;
|
|
||||||
return Some(std::net::IpAddr::V4(std::net::Ipv4Addr::from(addr)));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
pos += (attr_len + 3) & !3;
|
|
||||||
}
|
}
|
||||||
None
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,7 +3,7 @@
|
|||||||
//! whether the string-fixups in webrtc_ingest.rs actually match.
|
//! whether the string-fixups in webrtc_ingest.rs actually match.
|
||||||
use std::{net::SocketAddr, time::Instant};
|
use std::{net::SocketAddr, time::Instant};
|
||||||
|
|
||||||
use str0m::{change::SdpOffer, net::Protocol, Candidate, Rtc};
|
use str0m::{Candidate, Rtc, change::SdpOffer, net::Protocol};
|
||||||
|
|
||||||
fn obs_like_offer() -> String {
|
fn obs_like_offer() -> String {
|
||||||
let mut s = String::new();
|
let mut s = String::new();
|
||||||
|
|||||||
Generated
+28
-7
@@ -2,11 +2,11 @@
|
|||||||
"nodes": {
|
"nodes": {
|
||||||
"crane": {
|
"crane": {
|
||||||
"locked": {
|
"locked": {
|
||||||
"lastModified": 1785284101,
|
"lastModified": 1788465171,
|
||||||
"narHash": "sha256-ghcXEpYEM4a7pbEkoqbn8c0ptJJqgGzuFiG3T6W5g4I=",
|
"narHash": "sha256-Y1/TTVXjYXGF068IThQH9fPSZ0SIE74PABlUxnWTUH0=",
|
||||||
"owner": "ipetkov",
|
"owner": "ipetkov",
|
||||||
"repo": "crane",
|
"repo": "crane",
|
||||||
"rev": "756d6d07c3818ea95d1e2cdac63fa7d02fe3e61b",
|
"rev": "eb35abda9f232cc6610b1d1e3200d15c49b7ac54",
|
||||||
"type": "github"
|
"type": "github"
|
||||||
},
|
},
|
||||||
"original": {
|
"original": {
|
||||||
@@ -35,11 +35,11 @@
|
|||||||
},
|
},
|
||||||
"nixpkgs": {
|
"nixpkgs": {
|
||||||
"locked": {
|
"locked": {
|
||||||
"lastModified": 1785301185,
|
"lastModified": 1789012029,
|
||||||
"narHash": "sha256-eoS3KQTO0aPWXZvIaRbRAzSSHW3l5wdMFXtT1ISfoKA=",
|
"narHash": "sha256-1CBkBf+Nhggykzlx0jvXrj5rk20btl+tK8Fa8Ml/BL4=",
|
||||||
"owner": "NixOS",
|
"owner": "NixOS",
|
||||||
"repo": "nixpkgs",
|
"repo": "nixpkgs",
|
||||||
"rev": "9bc02893134c733dd85de46ee4fb2fac696b5529",
|
"rev": "d5dfd8e6716dde34398bc14bc87c10dece9c8c68",
|
||||||
"type": "github"
|
"type": "github"
|
||||||
},
|
},
|
||||||
"original": {
|
"original": {
|
||||||
@@ -53,7 +53,28 @@
|
|||||||
"inputs": {
|
"inputs": {
|
||||||
"crane": "crane",
|
"crane": "crane",
|
||||||
"flake-utils": "flake-utils",
|
"flake-utils": "flake-utils",
|
||||||
"nixpkgs": "nixpkgs"
|
"nixpkgs": "nixpkgs",
|
||||||
|
"rust-overlay": "rust-overlay"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"rust-overlay": {
|
||||||
|
"inputs": {
|
||||||
|
"nixpkgs": [
|
||||||
|
"nixpkgs"
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"locked": {
|
||||||
|
"lastModified": 1789024335,
|
||||||
|
"narHash": "sha256-kCy/MVLRIr95DJ4vspVzWj+kO/x+JuwSHmYZCvShg8w=",
|
||||||
|
"owner": "oxalica",
|
||||||
|
"repo": "rust-overlay",
|
||||||
|
"rev": "577bb1e1fc5af0713169176c5c76622c21fa3ec0",
|
||||||
|
"type": "github"
|
||||||
|
},
|
||||||
|
"original": {
|
||||||
|
"owner": "oxalica",
|
||||||
|
"repo": "rust-overlay",
|
||||||
|
"type": "github"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"systems": {
|
"systems": {
|
||||||
|
|||||||
@@ -7,6 +7,12 @@
|
|||||||
crane.url = "github:ipetkov/crane";
|
crane.url = "github:ipetkov/crane";
|
||||||
|
|
||||||
flake-utils.url = "github:numtide/flake-utils";
|
flake-utils.url = "github:numtide/flake-utils";
|
||||||
|
|
||||||
|
rust-overlay = {
|
||||||
|
url = "github:oxalica/rust-overlay";
|
||||||
|
inputs.nixpkgs.follows = "nixpkgs";
|
||||||
|
};
|
||||||
|
|
||||||
};
|
};
|
||||||
|
|
||||||
outputs =
|
outputs =
|
||||||
@@ -16,14 +22,27 @@
|
|||||||
crane,
|
crane,
|
||||||
flake-utils,
|
flake-utils,
|
||||||
...
|
...
|
||||||
}:
|
}@inputs:
|
||||||
flake-utils.lib.eachDefaultSystem (
|
flake-utils.lib.eachDefaultSystem (
|
||||||
system:
|
system:
|
||||||
let
|
let
|
||||||
pkgs = nixpkgs.legacyPackages.${system};
|
pkgs = import inputs.nixpkgs {
|
||||||
|
inherit system;
|
||||||
|
overlays = [ (import inputs.rust-overlay) ];
|
||||||
|
};
|
||||||
|
|
||||||
inherit (pkgs) lib;
|
inherit (pkgs) lib;
|
||||||
|
|
||||||
craneLib = crane.mkLib pkgs;
|
craneLib = (inputs.crane.mkLib pkgs).overrideToolchain (
|
||||||
|
p:
|
||||||
|
p.rust-bin.nightly.latest.default.override {
|
||||||
|
extensions = [
|
||||||
|
"rustc-codegen-cranelift-preview"
|
||||||
|
"rust-analyzer"
|
||||||
|
"rust-src"
|
||||||
|
];
|
||||||
|
}
|
||||||
|
);
|
||||||
|
|
||||||
# Common arguments can be set here to avoid repeating them later
|
# Common arguments can be set here to avoid repeating them later
|
||||||
# Note: changes here will rebuild all dependency crates
|
# Note: changes here will rebuild all dependency crates
|
||||||
@@ -56,7 +75,7 @@
|
|||||||
commonArgs
|
commonArgs
|
||||||
// {
|
// {
|
||||||
pname = "rtmp-to-whip";
|
pname = "rtmp-to-whip";
|
||||||
version = "0.1.0";
|
version = "0.4.0";
|
||||||
cargoArtifacts = craneLib.buildDepsOnly commonArgs;
|
cargoArtifacts = craneLib.buildDepsOnly commonArgs;
|
||||||
cargoExtraArgs = "-p server";
|
cargoExtraArgs = "-p server";
|
||||||
src = fileSetForCrate ./crates/server;
|
src = fileSetForCrate ./crates/server;
|
||||||
|
|||||||
Reference in New Issue
Block a user