From 5fb580f78203e17df2ac54db352dbad53e38ee02 Mon Sep 17 00:00:00 2001 From: KhooLy <73142442+KhooLy@users.noreply.github.com> Date: Sat, 27 Jun 2026 21:04:20 +0300 Subject: [PATCH] feat: update core casting and streaming support --- fluxa-streaming-engine/Cargo.toml | 12 - .../src/bin/companion_server.rs | 11 - .../src/bin/torrent_bench.rs | 150 --------- .../src/bin/torrent_serve.rs | 209 ------------- fluxa-streaming-engine/src/bindings/jni.rs | 22 ++ fluxa-streaming-engine/src/chapters.rs | 232 ++++++++++++++ .../src/companion_server.rs | 5 +- fluxa-streaming-engine/src/lib.rs | 1 + fluxa-streaming-engine/src/torrent_engine.rs | 22 +- fluxa-streaming-engine/src/transcode.rs | 47 ++- src/bindings/jni.rs | 191 +++++++----- src/bindings/uniffi.rs | 38 ++- src/cast_protocol.rs | 288 ++++++++++++++++++ src/content_identity.rs | 30 +- src/core_api.rs | 58 +++- src/headless_engine/contracts.rs | 2 + src/headless_engine/detail.rs | 6 + src/headless_engine/library.rs | 28 ++ src/headless_engine/mod.rs | 3 + src/home_ranking.rs | 58 ---- src/lib.rs | 1 + src/library_state.rs | 6 + 22 files changed, 858 insertions(+), 562 deletions(-) delete mode 100644 fluxa-streaming-engine/src/bin/companion_server.rs delete mode 100644 fluxa-streaming-engine/src/bin/torrent_bench.rs delete mode 100644 fluxa-streaming-engine/src/bin/torrent_serve.rs create mode 100644 fluxa-streaming-engine/src/chapters.rs create mode 100644 src/cast_protocol.rs diff --git a/fluxa-streaming-engine/Cargo.toml b/fluxa-streaming-engine/Cargo.toml index 173ebe6..41470cf 100644 --- a/fluxa-streaming-engine/Cargo.toml +++ b/fluxa-streaming-engine/Cargo.toml @@ -6,18 +6,6 @@ edition = "2021" [lib] crate-type = ["cdylib", "rlib"] -[[bin]] -name = "torrent_bench" -path = "src/bin/torrent_bench.rs" - -[[bin]] -name = "torrent_serve" -path = "src/bin/torrent_serve.rs" - -[[bin]] -name = "companion_server" -path = "src/bin/companion_server.rs" - [features] default = ["native"] native = [ diff --git a/fluxa-streaming-engine/src/bin/companion_server.rs b/fluxa-streaming-engine/src/bin/companion_server.rs deleted file mode 100644 index fe70309..0000000 --- a/fluxa-streaming-engine/src/bin/companion_server.rs +++ /dev/null @@ -1,11 +0,0 @@ -#[tokio::main] -async fn main() { - let port: u16 = std::env::var("FLUXA_COMPANION_PORT") - .ok() - .and_then(|p| p.parse().ok()) - .unwrap_or(48211); - - fluxa_streaming_engine::companion_server::serve(port) - .await - .expect("companion server crashed"); -} diff --git a/fluxa-streaming-engine/src/bin/torrent_bench.rs b/fluxa-streaming-engine/src/bin/torrent_bench.rs deleted file mode 100644 index 5705cb3..0000000 --- a/fluxa-streaming-engine/src/bin/torrent_bench.rs +++ /dev/null @@ -1,150 +0,0 @@ -/// rqbit bench — magnet → metadata → first 256 KB -/// Usage: cargo run --bin torrent_bench -- "" [--bytes N] -/// Prints one JSON line. - -use librqbit::{ - AddTorrent, AddTorrentOptions, Api, PeerConnectionOptions, Session, SessionOptions, -}; -use std::io::SeekFrom; -use std::time::{Duration, Instant}; -use tokio::io::{AsyncReadExt, AsyncSeekExt}; - -const DEFAULT_BYTES: usize = 256 * 1024; -const DEFAULT_MAGNET: &str = concat!( - "magnet:?xt=urn:btih:dd8255ecdc7ca55fb0bbf81323d87062db1f6d1c", - "&dn=Big+Buck+Bunny", - "&tr=udp://tracker.opentrackr.org:1337/announce", - "&tr=udp://open.stealth.si:80/announce", - "&tr=udp://tracker.openbittorrent.com:80/announce" -); - -#[tokio::main] -async fn main() { - let args: Vec = std::env::args().collect(); - let magnet = args.get(1).map(String::as_str).unwrap_or(DEFAULT_MAGNET); - let bytes_target: usize = args - .windows(2) - .find(|w| w[0] == "--bytes") - .and_then(|w| w[1].parse().ok()) - .unwrap_or(DEFAULT_BYTES); - - let cache = std::env::temp_dir().join("fluxa_bench_rqbit"); - let _ = std::fs::remove_dir_all(&cache); - std::fs::create_dir_all(&cache).unwrap(); - - eprintln!("[rqbit] starting session…"); - let t0 = Instant::now(); - - let mut session_opts = SessionOptions::default(); - session_opts.disable_dht_persistence = true; - session_opts.defer_writes_up_to = Some(64); - session_opts.listen_port_range = Some(49152..65535); - session_opts.disable_upload = true; - - let session = Session::new_with_opts(cache.clone(), session_opts) - .await - .expect("session"); - let api = Api::new(session, None); - let t_engine = t0.elapsed(); - eprintln!("[rqbit] engine ready in {}ms", t_engine.as_millis()); - - let mut add_opts = AddTorrentOptions::default(); - add_opts.overwrite = true; - add_opts.output_folder = Some(cache.to_string_lossy().into_owned()); - add_opts.peer_opts = Some(PeerConnectionOptions { - connect_timeout: Some(Duration::from_secs(10)), - read_write_timeout: Some(Duration::from_secs(20)), - ..Default::default() - }); - - eprintln!("[rqbit] adding torrent…"); - let t_add = Instant::now(); - let response = api - .api_add_torrent(AddTorrent::Url(magnet.to_string().into()), Some(add_opts)) - .await - .expect("add_torrent"); - - let torrent_id = response.id.expect("metadata must be ready after add"); - let t_metadata = t_add.elapsed(); - eprintln!("[rqbit] metadata in {}ms", t_metadata.as_millis()); - - let files = response.details.files.as_ref().expect("file list"); - let (file_id, file_name, file_len) = files - .iter() - .enumerate() - .filter(|(_, f)| { - let n = f.name.to_ascii_lowercase(); - matches!( - n.rsplit('.').next().unwrap_or(""), - "mkv" | "mp4" | "avi" | "webm" | "m4v" | "mov" - ) - }) - .max_by_key(|(_, f)| f.length) - .or_else(|| files.iter().enumerate().max_by_key(|(_, f)| f.length)) - .map(|(i, f)| (i, f.name.clone(), f.length)) - .expect("at least one file"); - - eprintln!("[rqbit] target: [{file_id}] {file_name} ({} MB)", file_len / 1024 / 1024); - - // prioritize only the target file - { - use std::collections::HashSet; - let only = HashSet::from([file_id]); - let _ = api - .api_torrent_action_update_only_files( - librqbit::api::TorrentIdOrHash::Id(torrent_id), - &only, - ) - .await; - } - - eprintln!("[rqbit] waiting for first {}KB…", bytes_target / 1024); - let t_read_start = Instant::now(); - let mut first_byte_ms: Option = None; - let mut total_read = 0usize; - - // retry until api_stream succeeds (torrent must be Live) - let mut stream = loop { - match api.api_stream( - librqbit::api::TorrentIdOrHash::Id(torrent_id), - file_id, - ) { - Ok(s) => break s, - Err(_) => tokio::time::sleep(Duration::from_millis(100)).await, - } - }; - - // seek to start (just in case) - stream.seek(SeekFrom::Start(0)).await.ok(); - - let mut buf = vec![0u8; 32 * 1024]; - while total_read < bytes_target { - let n = stream.read(&mut buf).await.expect("read"); - if n == 0 { - break; - } - if first_byte_ms.is_none() { - first_byte_ms = Some(t_read_start.elapsed().as_millis()); - eprintln!("[rqbit] first byte in {}ms", first_byte_ms.unwrap()); - } - total_read += n; - } - let t_data_done = t_read_start.elapsed(); - eprintln!( - "[rqbit] {}KB read in {}ms", - total_read / 1024, - t_data_done.as_millis() - ); - - let out = serde_json::json!({ - "engine": "rqbit", - "engineStartupMs": t_engine.as_millis(), - "metadataMs": t_metadata.as_millis(), - "firstByteMs": first_byte_ms, - "readDoneMs": t_data_done.as_millis(), - "bytesRead": total_read, - "fileName": file_name, - "fileSizeBytes": file_len - }); - println!("{}", serde_json::to_string_pretty(&out).unwrap()); -} diff --git a/fluxa-streaming-engine/src/bin/torrent_serve.rs b/fluxa-streaming-engine/src/bin/torrent_serve.rs deleted file mode 100644 index 6e16cba..0000000 --- a/fluxa-streaming-engine/src/bin/torrent_serve.rs +++ /dev/null @@ -1,209 +0,0 @@ -/// rqbit HTTP streaming server — magnet → HTTP stream URL -/// Starts a local HTTP server, adds the torrent, prints: -/// READY (once the stream is serveable) -/// then keeps running until killed. -/// -/// Usage: torrent_serve "" [--file-idx N] - -use axum::body::Body; -use axum::extract::State; -use axum::http::{HeaderMap, HeaderValue, StatusCode}; -use axum::response::{IntoResponse, Response}; -use axum::routing::get; -use axum::Router; -use std::io::Write; -use librqbit::{ - AddTorrent, AddTorrentOptions, Api, PeerConnectionOptions, Session, SessionOptions, -}; -use std::collections::HashSet; -use std::io::SeekFrom; -use std::sync::Arc; -use std::time::{Duration, Instant}; -use tokio::io::AsyncSeekExt; -use tokio::net::TcpListener; -use tokio_util::io::ReaderStream; - -#[derive(Clone)] -struct Ctx { - api: Api, - torrent_id: usize, - file_id: usize, - file_name: String, - file_len: u64, -} - -#[tokio::main] -async fn main() { - let args: Vec = std::env::args().collect(); - let magnet = args - .get(1) - .expect("Usage: torrent_serve [--file-idx N]"); - let file_idx_override: Option = args - .windows(2) - .find(|w| w[0] == "--file-idx") - .and_then(|w| w[1].parse().ok()); - - let cache = std::env::temp_dir().join("fluxa_bench_serve"); - let _ = std::fs::remove_dir_all(&cache); - std::fs::create_dir_all(&cache).unwrap(); - - eprintln!("[serve] starting rqbit session…"); - let t0 = Instant::now(); - - let mut session_opts = SessionOptions::default(); - session_opts.disable_dht_persistence = true; - session_opts.defer_writes_up_to = Some(64); - session_opts.listen_port_range = Some(49152..65535); - session_opts.disable_upload = true; - - let session = Session::new_with_opts(cache.clone(), session_opts) - .await - .expect("session"); - let api = Api::new(session, None); - eprintln!("[serve] engine up in {}ms", t0.elapsed().as_millis()); - - eprintln!("[serve] adding torrent (waiting for metadata)…"); - let t_meta = Instant::now(); - - let mut add_opts = AddTorrentOptions::default(); - add_opts.overwrite = true; - add_opts.output_folder = Some(cache.to_string_lossy().into_owned()); - add_opts.peer_opts = Some(PeerConnectionOptions { - connect_timeout: Some(Duration::from_secs(10)), - read_write_timeout: Some(Duration::from_secs(20)), - ..Default::default() - }); - - let response = api - .api_add_torrent(AddTorrent::Url(magnet.to_string().into()), Some(add_opts)) - .await - .expect("add_torrent failed"); - - let torrent_id = response.id.expect("torrent id"); - eprintln!("[serve] metadata in {}ms", t_meta.elapsed().as_millis()); - - let files = response.details.files.as_ref().expect("files"); - let (file_id, file_name, file_len) = if let Some(idx) = file_idx_override { - let f = &files[idx]; - (idx, f.name.clone(), f.length) - } else { - files - .iter() - .enumerate() - .filter(|(_, f)| { - matches!( - f.name.to_ascii_lowercase().rsplit('.').next().unwrap_or(""), - "mkv" | "mp4" | "avi" | "webm" | "m4v" | "mov" - ) - }) - .max_by_key(|(_, f)| f.length) - .or_else(|| files.iter().enumerate().max_by_key(|(_, f)| f.length)) - .map(|(i, f)| (i, f.name.clone(), f.length)) - .expect("at least one file") - }; - - eprintln!( - "[serve] target file: [{file_id}] {file_name} ({} MB)", - file_len / 1024 / 1024 - ); - - // prioritize only this file - let only = HashSet::from([file_id]); - let _ = api - .api_torrent_action_update_only_files( - librqbit::api::TorrentIdOrHash::Id(torrent_id), - &only, - ) - .await; - - let ctx = Arc::new(Ctx { - api, - torrent_id, - file_id, - file_name, - file_len, - }); - - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let port = listener.local_addr().unwrap().port(); - let stream_url = format!("http://127.0.0.1:{port}/stream"); - - let app = Router::new() - .route("/stream", get(handle_stream)) - .with_state(ctx); - - // Signal ready BEFORE accepting connections — explicit flush required when stdout is a file - println!("READY {stream_url}"); - let _ = std::io::stdout().flush(); - eprintln!("[serve] listening on {stream_url}"); - - axum::serve(listener, app).await.unwrap(); -} - -async fn handle_stream( - State(ctx): State>, - headers: HeaderMap, -) -> Response { - // retry until rqbit transitions to Live - let mut stream = loop { - match ctx - .api - .api_stream(librqbit::api::TorrentIdOrHash::Id(ctx.torrent_id), ctx.file_id) - { - Ok(s) => break s, - Err(_) => tokio::time::sleep(Duration::from_millis(100)).await, - } - }; - - let total = ctx.file_len; - let range_str = headers - .get("Range") - .and_then(|v| v.to_str().ok()) - .unwrap_or(""); - - let (status, start, length) = if let Some((s, e)) = parse_range(range_str, total) { - let end = e.unwrap_or(total - 1); - let _ = stream.seek(SeekFrom::Start(s)).await; - (StatusCode::PARTIAL_CONTENT, s, end - s + 1) - } else { - (StatusCode::OK, 0u64, total) - }; - - let ext = ctx.file_name.rsplit('.').next().unwrap_or("mp4"); - let mime = match ext { - "mkv" => "video/x-matroska", - "webm" => "video/webm", - "mp4" | "m4v" => "video/mp4", - "avi" => "video/x-msvideo", - _ => "application/octet-stream", - }; - - let mut resp_headers = HeaderMap::new(); - resp_headers.insert("Accept-Ranges", HeaderValue::from_static("bytes")); - resp_headers.insert("Content-Type", HeaderValue::from_str(mime).unwrap()); - resp_headers.insert( - "Content-Length", - HeaderValue::from_str(&length.to_string()).unwrap(), - ); - if status == StatusCode::PARTIAL_CONTENT { - let end = start + length - 1; - resp_headers.insert( - "Content-Range", - HeaderValue::from_str(&format!("bytes {start}-{end}/{total}")).unwrap(), - ); - } - - let body = Body::from_stream(ReaderStream::with_capacity(stream, 65536)); - (status, resp_headers, body).into_response() -} - -fn parse_range(value: &str, length: u64) -> Option<(u64, Option)> { - let raw = value.strip_prefix("bytes=")?; - let (start, end) = raw.split_once('-')?; - let start = start.parse::().ok()?; - if start >= length { - return None; - } - let end = end.parse::().ok().map(|e| e.min(length - 1)); - Some((start, end)) -} diff --git a/fluxa-streaming-engine/src/bindings/jni.rs b/fluxa-streaming-engine/src/bindings/jni.rs index ce9a8cb..d8cf7a1 100644 --- a/fluxa-streaming-engine/src/bindings/jni.rs +++ b/fluxa-streaming-engine/src/bindings/jni.rs @@ -1,3 +1,4 @@ +use crate::chapters::parse_mkv_chapters_json; use crate::dv_rewrite::{dv_auto_detect_was_iptpqc2, dv_get_stream_stats_json, dv_rewrite_segment_bytes, dv_rpu_self_test, start_dv_rewrite_local_stream_server}; use crate::local_stream::{start_local_stream_server, stop_local_stream_server}; use crate::torrent_engine; @@ -186,3 +187,24 @@ pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaStreamingNative_ })) .unwrap_or(ptr::null_mut()) } + +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaStreamingNative_parseMkvChaptersNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + data: JByteArray<'_>, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let len = match env.get_array_length(&data) { + Ok(l) => l as usize, + Err(_) => return write_jstring(&mut env, Some("[]".to_string())), + }; + let mut buf_i8: Vec = vec![0i8; len]; + if len > 0 && env.get_byte_array_region(&data, 0, &mut buf_i8).is_err() { + return write_jstring(&mut env, Some("[]".to_string())); + } + let input: Vec = buf_i8.into_iter().map(|b| b as u8).collect(); + write_jstring(&mut env, Some(parse_mkv_chapters_json(&input))) + })) + .unwrap_or(ptr::null_mut()) +} diff --git a/fluxa-streaming-engine/src/chapters.rs b/fluxa-streaming-engine/src/chapters.rs new file mode 100644 index 0000000..7cb4364 --- /dev/null +++ b/fluxa-streaming-engine/src/chapters.rs @@ -0,0 +1,232 @@ +// Standalone Matroska Chapters reader. +// +// ExoPlayer's MatroskaExtractor doesn't surface the Chapters element, so for the +// ExoPlayer backend we read a prefix of the file ourselves and walk the EBML tree +// far enough to find Segment -> Chapters -> EditionEntry -> ChapterAtom, without +// needing a full demuxer. Chapters always precede Cluster data in every muxer we +// care about (mkvmerge, ffmpeg), so scanning stops as soon as a Cluster is hit. + +use crate::dv_rewrite::try_parse_ebml_header; + +const ID_EBML_HEADER: u64 = 0x1A45DFA3; +const ID_SEGMENT: u64 = 0x1853_8067; +const ID_CLUSTER: u64 = 0x1F43_B675; +const ID_CHAPTERS: u64 = 0x1043_A770; +const ID_EDITION_ENTRY: u64 = 0x45B9; +const ID_CHAPTER_ATOM: u64 = 0xB6; +const ID_CHAPTER_TIME_START: u64 = 0x91; +const ID_CHAPTER_DISPLAY: u64 = 0x80; +const ID_CHAPTER_STRING: u64 = 0x85; +const UNKNOWN_SIZE: u64 = u64::MAX; + +#[derive(Debug, PartialEq)] +pub(crate) struct MkvChapter { + pub title: String, + pub start_ms: u64, +} + +/// Read `(id, header_len, content_start, content_end)` for the element at `pos`, +/// clamping an unknown-size element's content to `scan_end`. +fn read_element(buf: &[u8], pos: usize, scan_end: usize) -> Option<(u64, usize, usize, usize)> { + let (id, size, header_len) = try_parse_ebml_header(&buf[pos..scan_end.min(buf.len())])?; + let content_start = pos + header_len; + let content_end = if size == UNKNOWN_SIZE { + scan_end + } else { + (content_start + size as usize).min(scan_end) + }; + Some((id, header_len, content_start, content_end)) +} + +fn parse_chapter_display(buf: &[u8], mut pos: usize, end: usize) -> Option { + while pos < end { + let (id, _, content_start, content_end) = read_element(buf, pos, end)?; + if id == ID_CHAPTER_STRING { + return String::from_utf8(buf[content_start..content_end].to_vec()).ok(); + } + pos = content_end.max(pos + 1); + } + None +} + +fn read_uint_be(bytes: &[u8]) -> u64 { + bytes.iter().fold(0u64, |acc, &b| (acc << 8) | b as u64) +} + +fn parse_chapter_atom(buf: &[u8], mut pos: usize, end: usize) -> Option { + let mut start_ms: Option = None; + let mut title: Option = None; + while pos < end { + let (id, _, content_start, content_end) = read_element(buf, pos, end)?; + match id { + ID_CHAPTER_TIME_START => { + let raw_ns = read_uint_be(&buf[content_start..content_end]); + start_ms = Some(raw_ns / 1_000_000); + } + ID_CHAPTER_DISPLAY if title.is_none() => { + title = parse_chapter_display(buf, content_start, content_end); + } + _ => {} + } + pos = content_end.max(pos + 1); + } + start_ms.map(|ms| MkvChapter { title: title.unwrap_or_default(), start_ms: ms }) +} + +fn parse_edition_entry(buf: &[u8], mut pos: usize, end: usize, out: &mut Vec) { + while pos < end { + let Some((id, _, content_start, content_end)) = read_element(buf, pos, end) else { return }; + if id == ID_CHAPTER_ATOM { + if let Some(chapter) = parse_chapter_atom(buf, content_start, content_end) { + out.push(chapter); + } + } + pos = content_end.max(pos + 1); + } +} + +fn parse_chapters_element(buf: &[u8], mut pos: usize, end: usize) -> Vec { + let mut chapters = Vec::new(); + while pos < end { + let Some((id, _, content_start, content_end)) = read_element(buf, pos, end) else { break }; + if id == ID_EDITION_ENTRY { + parse_edition_entry(buf, content_start, content_end, &mut chapters); + } + pos = content_end.max(pos + 1); + } + chapters +} + +/// Scan the direct children of `Segment` looking for `Chapters`, stopping as soon +/// as a `Cluster` is reached (chapter metadata always precedes frame data). +fn scan_segment(buf: &[u8], mut pos: usize, end: usize) -> Vec { + while pos < end { + let Some((id, _, content_start, content_end)) = read_element(buf, pos, end) else { break }; + if id == ID_CLUSTER { + break; + } + if id == ID_CHAPTERS { + return parse_chapters_element(buf, content_start, content_end); + } + pos = content_end.max(pos + 1); + } + Vec::new() +} + +/// Parse chapters from a prefix of an MKV/WebM file. `buf` only needs to cover the +/// EBML header, Segment header, and whatever metadata elements precede the first +/// Cluster — a multi-megabyte prefix is normally more than enough. +pub(crate) fn parse_mkv_chapters(buf: &[u8]) -> Vec { + let mut pos = 0usize; + let end = buf.len(); + if let Some((id, _, _content_start, content_end)) = read_element(buf, pos, end) { + if id == ID_EBML_HEADER { + pos = content_end; + } + } + while pos < end { + let Some((id, _, content_start, content_end)) = read_element(buf, pos, end) else { break }; + if id == ID_SEGMENT { + return scan_segment(buf, content_start, content_end); + } + pos = content_end.max(pos + 1); + } + Vec::new() +} + +pub(crate) fn parse_mkv_chapters_json(buf: &[u8]) -> String { + let chapters = parse_mkv_chapters(buf); + let arr: Vec = chapters + .iter() + .map(|c| serde_json::json!({ "title": c.title, "startMs": c.start_ms })) + .collect(); + serde_json::to_string(&arr).unwrap_or_else(|_| "[]".to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::dv_rewrite::encode_ebml_element; + + fn chapter_atom(start_ms: u64, title: &str) -> Vec { + let time_start = encode_ebml_element(ID_CHAPTER_TIME_START, &(start_ms * 1_000_000).to_be_bytes()); + let chapter_string = encode_ebml_element(ID_CHAPTER_STRING, title.as_bytes()); + let display = encode_ebml_element(ID_CHAPTER_DISPLAY, &chapter_string); + let mut content = Vec::new(); + content.extend_from_slice(&time_start); + content.extend_from_slice(&display); + encode_ebml_element(ID_CHAPTER_ATOM, &content) + } + + fn segment_with_chapters(chapters: &[(u64, &str)]) -> Vec { + let atoms: Vec = chapters.iter().flat_map(|(ms, title)| chapter_atom(*ms, title)).collect(); + let edition_entry = encode_ebml_element(ID_EDITION_ENTRY, &atoms); + let chapters_elem = encode_ebml_element(ID_CHAPTERS, &edition_entry); + encode_ebml_element(ID_SEGMENT, &chapters_elem) + } + + #[test] + fn parses_single_chapter() { + let segment = segment_with_chapters(&[(0, "OP")]); + let chapters = parse_mkv_chapters(&segment); + assert_eq!(chapters, vec![MkvChapter { title: "OP".to_string(), start_ms: 0 }]); + } + + #[test] + fn parses_multiple_chapters_in_order() { + let segment = segment_with_chapters(&[(0, "OP"), (90_000, "Episode"), (1_320_000, "ED")]); + let chapters = parse_mkv_chapters(&segment); + assert_eq!( + chapters, + vec![ + MkvChapter { title: "OP".to_string(), start_ms: 0 }, + MkvChapter { title: "Episode".to_string(), start_ms: 90_000 }, + MkvChapter { title: "ED".to_string(), start_ms: 1_320_000 }, + ] + ); + } + + #[test] + fn stops_at_cluster_without_chapters() { + let cluster = encode_ebml_element(ID_CLUSTER, &[0x00, 0x01, 0x02]); + let segment = encode_ebml_element(ID_SEGMENT, &cluster); + assert!(parse_mkv_chapters(&segment).is_empty()); + } + + #[test] + fn chapters_before_cluster_are_still_found() { + let chapters_elem = { + let atom = chapter_atom(0, "Intro"); + let edition_entry = encode_ebml_element(ID_EDITION_ENTRY, &atom); + encode_ebml_element(ID_CHAPTERS, &edition_entry) + }; + let cluster = encode_ebml_element(ID_CLUSTER, &[0x00]); + let mut segment_content = Vec::new(); + segment_content.extend_from_slice(&chapters_elem); + segment_content.extend_from_slice(&cluster); + let segment = encode_ebml_element(ID_SEGMENT, &segment_content); + + let chapters = parse_mkv_chapters(&segment); + assert_eq!(chapters, vec![MkvChapter { title: "Intro".to_string(), start_ms: 0 }]); + } + + #[test] + fn handles_truncated_buffer_without_panicking() { + let segment = segment_with_chapters(&[(0, "OP")]); + for cut in 0..segment.len() { + let _ = parse_mkv_chapters(&segment[..cut]); + } + } + + #[test] + fn empty_buffer_returns_no_chapters() { + assert!(parse_mkv_chapters(&[]).is_empty()); + } + + #[test] + fn json_output_shape() { + let segment = segment_with_chapters(&[(0, "OP")]); + let json = parse_mkv_chapters_json(&segment); + assert_eq!(json, r#"[{"startMs":0,"title":"OP"}]"#); + } +} diff --git a/fluxa-streaming-engine/src/companion_server.rs b/fluxa-streaming-engine/src/companion_server.rs index 39b6ee4..df1e8b1 100644 --- a/fluxa-streaming-engine/src/companion_server.rs +++ b/fluxa-streaming-engine/src/companion_server.rs @@ -23,6 +23,7 @@ struct AppState { } #[derive(Deserialize)] +#[serde(rename_all = "camelCase")] struct StartTorrentBody { stream_json: String, title: Option, @@ -135,7 +136,9 @@ fn apply_torrent_preferences(base_url: &str, preferences: Option<&Value>) { }; let url = format!("{}/settings", base_url.trim_end_matches('/')); tokio::spawn(async move { - let client = reqwest::Client::new(); + let Ok(client) = reqwest::Client::builder().timeout(std::time::Duration::from_secs(5)).build() else { + return; + }; let _ = client.post(&url).json(&json!({ "PreloadSize": preload_size })).send().await; }); } diff --git a/fluxa-streaming-engine/src/lib.rs b/fluxa-streaming-engine/src/lib.rs index 8e2676d..d273d09 100644 --- a/fluxa-streaming-engine/src/lib.rs +++ b/fluxa-streaming-engine/src/lib.rs @@ -1,3 +1,4 @@ +mod chapters; mod dv_rewrite; mod local_stream; mod torrent_engine; diff --git a/fluxa-streaming-engine/src/torrent_engine.rs b/fluxa-streaming-engine/src/torrent_engine.rs index 57bd092..0645319 100644 --- a/fluxa-streaming-engine/src/torrent_engine.rs +++ b/fluxa-streaming-engine/src/torrent_engine.rs @@ -20,6 +20,7 @@ use std::time::Duration; use tokio::io::AsyncSeekExt; use tokio::net::TcpListener; use tokio::sync::oneshot; +use tokio::sync::Mutex as AsyncMutex; use tokio_util::io::ReaderStream; #[derive(Deserialize)] @@ -55,6 +56,10 @@ struct EngineState { output_dir: PathBuf, preload_size: Arc>, known_links: Arc>>, + // Serializes the check-then-add sequence in ensure_torrent so two + // near-simultaneous requests for the same new link (e.g. a stat poll + // racing the stream GET) can't both call api_add_torrent for it. + add_lock: Arc>, } struct TorrentServerHandle { @@ -86,7 +91,7 @@ pub fn start_torrent_server(cache_dir: &str, preferred_port: i32) -> Option(); @@ -150,6 +155,7 @@ pub fn start_torrent_server(cache_dir: &str, preferred_port: i32) -> Option); + +impl Drop for KillOnDrop { + fn drop(&mut self) { + if let Some(mut child) = self.0.take() { + let _ = child.start_kill(); + tokio::spawn(async move { + let _ = child.wait().await; + }); + } + } +} + +struct ChildStdoutGuarded { + stdout: ChildStdout, + _guard: KillOnDrop, +} + +impl AsyncRead for ChildStdoutGuarded { + fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + Pin::new(&mut self.get_mut().stdout).poll_read(cx, buf) + } +} + #[derive(Deserialize)] pub struct TranscodeQuery { url: String, @@ -131,13 +164,11 @@ pub async fn handle_transcode(Query(q): Query) -> Response { return (StatusCode::INTERNAL_SERVER_ERROR, "ffmpeg produced no stdout pipe").into_response(); }; - // Detached: the child is reaped once stdout closes (process exits) or the - // response stream is dropped (client disconnects), whichever comes first. - tokio::spawn(async move { - let _ = child.wait().await; - }); - - let body = Body::from_stream(ReaderStream::with_capacity(stdout, 65536)); + // Tied to the response body's lifetime via KillOnDrop, so a client + // disconnect kills ffmpeg immediately instead of leaving it stalled on + // a stdout pipe nobody's reading from. + let guarded = ChildStdoutGuarded { stdout, _guard: KillOnDrop(Some(child)) }; + let body = Body::from_stream(ReaderStream::with_capacity(guarded, 65536)); let mut response = (StatusCode::OK, body).into_response(); response.headers_mut().insert( axum::http::header::CONTENT_TYPE, diff --git a/src/bindings/jni.rs b/src/bindings/jni.rs index a3c51a0..f034ba3 100644 --- a/src/bindings/jni.rs +++ b/src/bindings/jni.rs @@ -25,12 +25,13 @@ use crate::stream_policy::*; use crate::watchlist_plan::*; use jni::objects::JClass; pub(crate) use jni::objects::JString; -use jni::sys::{jboolean, jfloat, jint, jlong, jstring}; +use jni::sys::{jboolean, jdouble, jfloat, jint, jlong, jstring}; pub(crate) use jni::JNIEnv; use serde_json::json; use std::ptr; pub(crate) type JBoolean = jboolean; +pub(crate) type JDouble = jdouble; pub(crate) type JFloat = jfloat; pub(crate) type JInt = jint; pub(crate) type JLong = jlong; @@ -1294,6 +1295,97 @@ pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_trakt .unwrap_or(ptr::null_mut()) } +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_simklScrobbleBodyJsonNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + ids_json: JString<'_>, + is_episode: JBoolean, + season: JLong, + ep_number: JLong, + time_pos_sec: JDouble, + duration_sec: JDouble, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let output = read_jstring(&mut env, &ids_json).and_then(|ids_json| { + simkl_scrobble_body_json( + &ids_json, + is_episode != 0, + season, + ep_number, + time_pos_sec, + duration_sec, + ) + }); + write_jstring(&mut env, output) + })) + .unwrap_or(ptr::null_mut()) + } + +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_simklMatchEpisodeJsonNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + episodes_json: JString<'_>, + target_json: JString<'_>, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let output = read_jstring(&mut env, &episodes_json).and_then(|episodes_json| { + simkl_match_episode_json(&episodes_json, &read_jstring(&mut env, &target_json)?) + }); + write_jstring(&mut env, output) + })) + .unwrap_or(ptr::null_mut()) + } + +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_simklWatchingToItemsJsonNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + shows_json: JString<'_>, + movies_json: JString<'_>, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let output = read_jstring(&mut env, &shows_json).and_then(|shows_json| { + simkl_watching_to_items_json(&shows_json, &read_jstring(&mut env, &movies_json)?) + }); + write_jstring(&mut env, output) + })) + .unwrap_or(ptr::null_mut()) + } + +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_simklWatchlistToItemsJsonNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + shows_json: JString<'_>, + movies_json: JString<'_>, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let output = read_jstring(&mut env, &shows_json).and_then(|shows_json| { + simkl_watchlist_to_items_json(&shows_json, &read_jstring(&mut env, &movies_json)?) + }); + write_jstring(&mut env, output) + })) + .unwrap_or(ptr::null_mut()) + } + +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_simklWatchedToIdsJsonNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + shows_json: JString<'_>, + movies_json: JString<'_>, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let output = read_jstring(&mut env, &shows_json).and_then(|shows_json| { + simkl_watched_to_ids_json(&shows_json, &read_jstring(&mut env, &movies_json)?) + }); + write_jstring(&mut env, output) + })) + .unwrap_or(ptr::null_mut()) + } + #[no_mangle] pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_playbackProgressItemJsonNative( mut env: JNIEnv<'_>, @@ -1572,14 +1664,28 @@ pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_conte } #[no_mangle] -pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_contentBillboardKeyNative( +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_contentTraktKeysBatchJsonNative( mut env: JNIEnv<'_>, _class: JObject<'_>, - meta_json: JString<'_>, + metas_json: JString<'_>, ) -> JStringReturn { std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - let output = - read_jstring(&mut env, &meta_json).and_then(|meta_json| content_billboard_key(&meta_json)); + let output = read_jstring(&mut env, &metas_json) + .and_then(|metas_json| content_trakt_keys_batch(&metas_json)); + write_jstring(&mut env, output) + })) + .unwrap_or(ptr::null_mut()) + } + +#[no_mangle] +pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_contentWatchedKeysBatchJsonNative( + mut env: JNIEnv<'_>, + _class: JObject<'_>, + metas_json: JString<'_>, +) -> JStringReturn { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let output = read_jstring(&mut env, &metas_json) + .and_then(|metas_json| content_watched_keys_batch(&metas_json)); write_jstring(&mut env, output) })) .unwrap_or(ptr::null_mut()) @@ -1800,81 +1906,6 @@ pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_optim .unwrap_or(ptr::null_mut()) } -#[no_mangle] -pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_billboardScoreCandidateNative( - mut env: JNIEnv<'_>, - _class: JObject<'_>, - meta_json: JString<'_>, - days_since_release: JLong, - has_days_since_release: JBoolean, -) -> JInt { - std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - read_jstring(&mut env, &meta_json) - .and_then(|meta_json| { - billboard_score_candidate_json( - &meta_json, - if has_days_since_release != 0 { - Some(days_since_release) - } else { - None - }, - ) - }) - .unwrap_or(0) - })) - .unwrap_or(0) - } - -#[no_mangle] -pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_billboardHasBackdropCandidateNative( - mut env: JNIEnv<'_>, - _class: JObject<'_>, - meta_json: JString<'_>, -) -> JBoolean { - std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - let result = read_jstring(&mut env, &meta_json) - .map(|meta_json| has_billboard_backdrop_candidate_json(&meta_json)) - .unwrap_or(false); - if result { - 1 - } else { - 0 - } - })) - .unwrap_or(0) - } - -#[no_mangle] -pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_billboardVisualScoreNative( - mut env: JNIEnv<'_>, - _class: JObject<'_>, - meta_json: JString<'_>, -) -> JInt { - std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - read_jstring(&mut env, &meta_json) - .and_then(|meta_json| billboard_visual_score_json(&meta_json)) - .unwrap_or(0) - })) - .unwrap_or(0) - } - -#[no_mangle] -pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_billboardEditorialMatchScoreNative( - mut env: JNIEnv<'_>, - _class: JObject<'_>, - meta_json: JString<'_>, - spec_json: JString<'_>, -) -> JInt { - std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - read_jstring(&mut env, &meta_json) - .and_then(|meta_json| { - billboard_editorial_match_score_json(&meta_json, &read_jstring(&mut env, &spec_json)?) - }) - .unwrap_or(0) - })) - .unwrap_or(0) - } - #[no_mangle] pub unsafe extern "system" fn Java_com_fluxa_app_core_rust_FluxaCoreNative_playerProgressPercentNative( _env: JNIEnv<'_>, diff --git a/src/bindings/uniffi.rs b/src/bindings/uniffi.rs index 1dad18a..f4a2083 100644 --- a/src/bindings/uniffi.rs +++ b/src/bindings/uniffi.rs @@ -1,5 +1,11 @@ use crate::{app_state, core_contract, headless_engine}; +// A panic anywhere below must not unwind across the UniFFI boundary into +// Swift/Kotlin — that's undefined behavior, not a catchable exception there. +fn guard(default: T, f: impl FnOnce() -> T) -> T { + std::panic::catch_unwind(std::panic::AssertUnwindSafe(f)).unwrap_or(default) +} + #[uniffi::export] pub fn fluxa_core_version() -> String { env!("CARGO_PKG_VERSION").to_string() @@ -13,12 +19,12 @@ pub fn core_invoke(method: String, args_json: String) -> String { #[uniffi::export] pub fn create_headless_engine_json(initial_json: String) -> i64 { - headless_engine::create_headless_engine(&initial_json) as i64 + guard(0, || headless_engine::create_headless_engine(&initial_json) as i64) } #[uniffi::export] pub fn destroy_headless_engine_json(handle: i64) -> bool { - handle > 0 && headless_engine::destroy_headless_engine(handle as u64) + handle > 0 && guard(false, || headless_engine::destroy_headless_engine(handle as u64)) } #[uniffi::export] @@ -26,7 +32,9 @@ pub fn headless_engine_snapshot_json(handle: i64) -> String { if handle <= 0 { return String::new(); } - headless_engine::headless_engine_snapshot_json(handle as u64).unwrap_or_default() + guard(String::new(), || { + headless_engine::headless_engine_snapshot_json(handle as u64).unwrap_or_default() + }) } #[uniffi::export] @@ -34,7 +42,9 @@ pub fn headless_engine_dispatch_json(handle: i64, action_json: String) -> String if handle <= 0 { return String::new(); } - headless_engine::headless_engine_dispatch_json(handle as u64, &action_json).unwrap_or_default() + guard(String::new(), || { + headless_engine::headless_engine_dispatch_json(handle as u64, &action_json).unwrap_or_default() + }) } #[uniffi::export] @@ -42,23 +52,25 @@ pub fn headless_engine_complete_effect_json(handle: i64, result_json: String) -> if handle <= 0 { return String::new(); } - headless_engine::headless_engine_complete_effect_json(handle as u64, &result_json) - .unwrap_or_default() + guard(String::new(), || { + headless_engine::headless_engine_complete_effect_json(handle as u64, &result_json) + .unwrap_or_default() + }) } #[uniffi::export] pub fn core_capabilities_json(portable: bool) -> String { - core_contract::core_capabilities_json(portable) + guard(String::new(), || core_contract::core_capabilities_json(portable)) } #[uniffi::export] pub fn create_app_core_state_json(initial_json: String) -> i64 { - app_state::create_app_core_state(&initial_json) as i64 + guard(0, || app_state::create_app_core_state(&initial_json) as i64) } #[uniffi::export] pub fn destroy_app_core_state_json(handle: i64) -> bool { - handle > 0 && app_state::destroy_app_core_state(handle as u64) + handle > 0 && guard(false, || app_state::destroy_app_core_state(handle as u64)) } #[uniffi::export] @@ -66,7 +78,9 @@ pub fn app_core_state_json(handle: i64) -> String { if handle <= 0 { return String::new(); } - app_state::app_core_state_json(handle as u64).unwrap_or_default() + guard(String::new(), || { + app_state::app_core_state_json(handle as u64).unwrap_or_default() + }) } #[uniffi::export] @@ -74,5 +88,7 @@ pub fn app_core_dispatch_json(handle: i64, action_json: String) -> String { if handle <= 0 { return String::new(); } - app_state::app_core_dispatch_json(handle as u64, &action_json).unwrap_or_default() + guard(String::new(), || { + app_state::app_core_dispatch_json(handle as u64, &action_json).unwrap_or_default() + }) } diff --git a/src/cast_protocol.rs b/src/cast_protocol.rs new file mode 100644 index 0000000..97fbb2a --- /dev/null +++ b/src/cast_protocol.rs @@ -0,0 +1,288 @@ +use serde_json::json; + +pub(crate) const AVTRANSPORT_URN: &str = "urn:schemas-upnp-org:service:AVTransport:1"; +pub(crate) const RENDERING_CONTROL_URN: &str = "urn:schemas-upnp-org:service:RenderingControl:1"; + +pub(crate) fn validate_stream_url(url: &str) -> bool { + let trimmed = url.trim(); + let Some(scheme_end) = trimmed.find("://") else { return false }; + let scheme = trimmed[..scheme_end].to_ascii_lowercase(); + if scheme != "http" && scheme != "https" { + return false; + } + !trimmed[scheme_end + 3..].is_empty() +} + +pub(crate) fn xml_escape(value: &str) -> String { + value.replace('&', "&").replace('<', "<").replace('>', ">") +} + +fn extract_tag(xml: &str, tag: &str) -> Option { + let open = format!("<{tag}>"); + let close = format!(""); + let start = xml.find(&open)? + open.len(); + let end = xml[start..].find(&close)? + start; + Some(xml[start..end].trim().to_string()) +} + +fn resolve_url(base_url: &str, path: &str) -> String { + if path.starts_with("http://") || path.starts_with("https://") { + return path.to_string(); + } + let base = base_url.trim_end_matches('/'); + if path.starts_with('/') { + if let Some(scheme_end) = base.find("://") { + if let Some(host_end) = base[scheme_end + 3..].find('/') { + return format!("{}{}", &base[..scheme_end + 3 + host_end], path); + } + } + return format!("{base}{path}"); + } + format!("{base}/{path}") +} + +fn extract_service_control_url(xml: &str, base_url: &str, urn: &str) -> Option { + for service_block in xml.split("").skip(1) { + let service_type = extract_tag(service_block, "serviceType")?; + if service_type != urn { + continue; + } + let control_path = extract_tag(service_block, "controlURL")?; + return Some(resolve_url(base_url, &control_path)); + } + None +} + +pub(crate) fn dlna_parse_device_description_json(xml: &str, base_url: &str) -> Option { + let name = extract_tag(xml, "friendlyName").unwrap_or_else(|| "Unknown device".to_string()); + let control_url = extract_service_control_url(xml, base_url, AVTRANSPORT_URN)?; + let rendering_control_url = extract_service_control_url(xml, base_url, RENDERING_CONTROL_URN); + Some(json!({"name": name, "controlUrl": control_url, "renderingControlUrl": rendering_control_url}).to_string()) +} + +pub(crate) fn soap_action_body(urn: &str, action: &str, args: &str) -> String { + format!( + "\ +\ +{args}" + ) +} + +fn format_didl_metadata(title: &str, subtitle_url: Option<&str>) -> String { + let subtitle_res = subtitle_url + .map(|url| format!("{}", xml_escape(url))) + .unwrap_or_default(); + let didl = format!( + "\ +{}object.item.videoItem{subtitle_res}", + xml_escape(title) + ); + xml_escape(&didl) +} + +pub(crate) fn dlna_set_av_transport_args(media_url: &str, title: &str, subtitle_url: Option<&str>) -> Option { + if !validate_stream_url(media_url) { + return None; + } + let metadata = format_didl_metadata(title, subtitle_url); + Some(format!( + "0{}{metadata}", + xml_escape(media_url) + )) +} + +pub(crate) fn format_hms(total_secs: f64) -> String { + let total = total_secs.max(0.0) as u64; + format!("{:02}:{:02}:{:02}", total / 3600, (total % 3600) / 60, total % 60) +} + +pub(crate) fn dlna_seek_args(position_secs: f64) -> String { + format!("0ABS_TIME{}", format_hms(position_secs)) +} + +pub(crate) fn dlna_set_volume_args(level: f64) -> String { + let volume = (level.clamp(0.0, 1.0) * 100.0).round() as u32; + format!("0Master{volume}") +} + +pub(crate) fn resolve_loopback_url(stream_url: &str, lan_ip: &str) -> String { + if let Some(rest) = stream_url.strip_prefix("http://127.0.0.1") { + return format!("http://{lan_ip}{rest}"); + } + stream_url.to_string() +} + +pub(crate) fn guess_cast_content_type(media_url: &str) -> &'static str { + let path = media_url.split(['?', '#']).next().unwrap_or(media_url).to_ascii_lowercase(); + if path.ends_with(".m3u8") { + "application/x-mpegurl" + } else if path.ends_with(".mkv") { + "video/x-matroska" + } else if path.ends_with(".webm") { + "video/webm" + } else { + "video/mp4" + } +} + +fn write_varint(buf: &mut Vec, mut value: u64) { + loop { + let byte = (value & 0x7F) as u8; + value >>= 7; + if value == 0 { + buf.push(byte); + break; + } + buf.push(byte | 0x80); + } +} + +fn write_tag(buf: &mut Vec, field: u32, wire_type: u8) { + write_varint(buf, ((field << 3) | wire_type as u32) as u64); +} + +fn write_string_field(buf: &mut Vec, field: u32, value: &str) { + write_tag(buf, field, 2); + write_varint(buf, value.len() as u64); + buf.extend_from_slice(value.as_bytes()); +} + +pub(crate) fn encode_cast_message(source_id: &str, destination_id: &str, namespace: &str, payload_utf8: &str) -> Vec { + let mut buf = Vec::new(); + write_tag(&mut buf, 1, 0); + write_varint(&mut buf, 0); + write_string_field(&mut buf, 2, source_id); + write_string_field(&mut buf, 3, destination_id); + write_string_field(&mut buf, 4, namespace); + write_tag(&mut buf, 5, 0); + write_varint(&mut buf, 0); + write_string_field(&mut buf, 6, payload_utf8); + buf +} + +fn read_varint(buf: &[u8], pos: &mut usize) -> Option { + let mut result = 0u64; + let mut shift = 0; + loop { + let byte = *buf.get(*pos)?; + *pos += 1; + result |= ((byte & 0x7F) as u64) << shift; + if byte & 0x80 == 0 { + break; + } + shift += 7; + } + Some(result) +} + +pub(crate) struct DecodedCastMessage { + pub namespace: String, + pub payload_utf8: String, +} + +pub(crate) fn decode_cast_message(buf: &[u8]) -> Option { + let mut pos = 0; + let mut namespace = String::new(); + let mut payload_utf8 = String::new(); + while pos < buf.len() { + let tag = read_varint(buf, &mut pos)?; + let field = (tag >> 3) as u32; + let wire_type = (tag & 0x7) as u8; + match wire_type { + 0 => { + read_varint(buf, &mut pos)?; + } + 2 => { + let len = read_varint(buf, &mut pos)? as usize; + let end = pos.checked_add(len)?; + let slice = buf.get(pos..end)?; + if field == 4 { + namespace = String::from_utf8_lossy(slice).to_string(); + } else if field == 6 { + payload_utf8 = String::from_utf8_lossy(slice).to_string(); + } + pos = end; + } + _ => return None, + } + } + Some(DecodedCastMessage { namespace, payload_utf8 }) +} + +fn roku_url_encode(value: &str) -> String { + let mut out = String::new(); + for byte in value.bytes() { + match byte { + b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => out.push(byte as char), + _ => out.push_str(&format!("%{byte:02X}")), + } + } + out +} + +pub(crate) fn roku_device_name(xml: &str) -> Option { + extract_tag(xml, "friendly-device-name") +} + +pub(crate) fn roku_launch_url(host: &str, media_url: &str, subtitle_url: Option<&str>) -> Option { + if !validate_stream_url(media_url) { + return None; + } + if let Some(sub) = subtitle_url { + if !validate_stream_url(sub) { + return None; + } + } + const ROKU_MEDIA_PLAYER_APP_ID: &str = "2213"; + let mut url = format!("http://{host}:8060/launch/{ROKU_MEDIA_PLAYER_APP_ID}?t=v&u={}", roku_url_encode(media_url)); + if let Some(sub) = subtitle_url { + url.push_str(&format!("&k={}", roku_url_encode(sub))); + } + Some(url) +} + +pub(crate) fn airplay_volume_db(level: f64) -> f64 { + if level <= 0.0 { + -30.0 + } else { + (20.0 * level.clamp(0.0, 1.0).log10()).max(-30.0) + } +} + +pub(crate) fn airplay_play_body(media_url: &str) -> Option { + if !validate_stream_url(media_url) { + return None; + } + Some(format!("Content-Location: {media_url}\nStart-Position: 0\n")) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn rejects_non_http_schemes() { + assert!(!validate_stream_url("file:///etc/passwd")); + assert!(!validate_stream_url("javascript:alert(1)")); + assert!(!validate_stream_url("not a url")); + } + + #[test] + fn accepts_http_and_https() { + assert!(validate_stream_url("http://192.168.1.5:11470/stream")); + assert!(validate_stream_url("https://example.com/movie.mkv")); + } + + #[test] + fn didl_title_with_markup_does_not_break_out_of_the_item_tag() { + let args = dlna_set_av_transport_args("http://192.168.1.5/a.mkv", "