From d41c59288e2660e90f3750291616efb240124fab Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 6 Mar 2025 20:15:08 +0100 Subject: [PATCH] hdhomerun for empy/plex/jellyfin implemented --- CHANGELOG.md | 6 +- Cargo.lock | 1 + Cargo.toml | 1 + README.md | 49 ++++ src/api/api_utils.rs | 2 +- src/api/endpoints/download_api.rs | 4 +- src/api/endpoints/hdhomerun_api.rs | 319 +++++++++++++++++++++ src/api/endpoints/hls_api.rs | 35 ++- src/api/endpoints/m3u_api.rs | 11 +- src/api/endpoints/mod.rs | 3 +- src/api/endpoints/user_api.rs | 7 +- src/api/endpoints/v1_api.rs | 53 +--- src/api/endpoints/web_index.rs | 7 +- src/api/endpoints/xmltv_api.rs | 7 +- src/api/endpoints/xtream_api.rs | 31 +- src/api/main_api.rs | 63 +++- src/api/model/app_state.rs | 6 + src/auth/authenticator.rs | 3 +- src/main.rs | 9 +- src/model/config.rs | 131 +++++++-- src/model/hdhomerun_config.rs | 95 ++++++ src/model/mod.rs | 3 +- src/model/playlist.rs | 26 +- src/processing/parser/xtream.rs | 42 +-- src/repository/epg_repository.rs | 4 +- src/repository/kodi_repository.rs | 6 +- src/repository/m3u_playlist_iterator.rs | 72 +++-- src/repository/m3u_repository.rs | 8 +- src/repository/playlist_repository.rs | 1 + src/repository/xtream_playlist_iterator.rs | 53 +++- src/repository/xtream_repository.rs | 38 ++- src/utils/file/config_reader.rs | 8 +- src/utils/network/xtream.rs | 4 +- 33 files changed, 905 insertions(+), 203 deletions(-) create mode 100644 src/api/endpoints/hdhomerun_api.rs create mode 100644 src/model/hdhomerun_config.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 007144997..7eeb62ebd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,13 +15,15 @@ targets: - name: test ``` -- The Web UI now includes a login feature for playlist users, allowing them to set their group for filtering and managing their own boutique. - The playlist user can login with his credentials and can select the desired groups for his playlist. +- The Web UI now includes a login feature for playlist users, allowing them to set their group for filtering and managing their own bouquet of groups. + The playlist user can login with his credentials and can select the desired groups for his playlist. +Currently only xtream api is supported, m3u will follow - Added `user_config_dir` to `config.yml`. It is the storage path for user configurations (f.e. bouquets). - New Filter field `input` can be used along `name`, `group`, `title`, `url` and `type`. Input is a `regexp` filter. `input ~ "provider\-\d+"` - New option `use_user_db` in `api-proxy.yml`. The Playlist Users are stored inside the config file `api-proxy.yml`. When you set this option to `true` the user are stored in a db file. This is a better choice if you have a lot of users. If you have only a few let it default to `false` - WebUI playlist browser with tree and gallery mode. Explore self hosted and provider playlists in browser. +- Added HdHomeRun tuner target for use with Plex/Emby/Jellyfin # 2.2.1 (2025-02-14) - Added more info to `/status`. diff --git a/Cargo.lock b/Cargo.lock index 5d7187d3d..c82f38b70 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1741,6 +1741,7 @@ dependencies = [ "serde", "serde_json", "serde_yaml", + "socket2", "tempfile", "tokio", "tokio-stream", diff --git a/Cargo.toml b/Cargo.toml index 4352ff3b5..1af51a04c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -55,6 +55,7 @@ tempfile = "3.16" ruzstd = "0" filetime = "0.2" parking_lot = "0.12" +socket2 = "0.5" #[cfg(target_os = "macos")] libc = "0" #[cfg(target_os = "windows")] diff --git a/README.md b/README.md index fb1b524ab..6299afd97 100644 --- a/README.md +++ b/README.md @@ -324,6 +324,55 @@ and add it to the `config.yml`. channel_unavailable_file: /freeze_frame.ts ``` +### 1.13 `hdhomerun` + +It is possible to define `hdhomerun` target for output. To use this outputs we need to define HdHomeRun devices. + +The simplest config looks like: +```yaml +hdhomerun: + enabled: true + devices: + - name: hdhr1 + - name: hdhr2 +``` + +The `name` must be unique and is used in the target configuration in `source.yml` like. + +```yaml +sources: +- inputs: + - name: ... + ... + targets: + - name: xt_m3u + output: + - type: xtream + - type: hdhomerun + username: xtr + output: hdhr1 + filter: "!ALL_FILTER!" +``` + +The HdHomerun config has the following attribute: +`enabled`: default is `false`, you need to set it to `true` +`devices`: is a list of HdHomeRun Device configuraitons. +For each output you need to define one device with a unique name. Each output gets his own port to connect. + +HdHomeRun device config has the following attributes: + +- `name`: _mandatory_ and must be unique +- `friendly_name`: _optional_ +- `manufacturer`: _optional_ +- `model_name`: _optional_ +- `model_number`: _optional_ +- `firmware_name`: _optional_ +- `firmware_version`: _optional_ +- `device_type`: _optional_ +- `device_udn`: _optional_ +- `port`: _optional_, if not given the m3u-filter-server port is incremented for each device. + + ## 2. `source.yml` Has the following top level entries: diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index a207b9bd1..45a691f96 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -302,7 +302,7 @@ pub fn empty_json_list_response() -> HttpResponse { HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("[]") } -pub fn get_username_from_auth_header(credentials: Option, app_state: &web::Data) -> Option { +pub fn get_username_from_auth_header(credentials: Option, app_state: &web::Data>) -> Option { if let Some(bearer) = credentials { if let Some(web_auth_config) = app_state.config.web_auth.as_ref() { let secret_key = web_auth_config.secret.as_ref(); diff --git a/src/api/endpoints/download_api.rs b/src/api/endpoints/download_api.rs index 2c927c1e9..f70d7020f 100644 --- a/src/api/endpoints/download_api.rs +++ b/src/api/endpoints/download_api.rs @@ -111,7 +111,7 @@ macro_rules! download_info { pub async fn queue_download_file( req: web::Json, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { if let Some(download_cfg) = &app_state.config.video.as_ref().unwrap().download { if download_cfg.directory.is_none() { @@ -137,7 +137,7 @@ pub async fn queue_download_file( } pub async fn download_file_info( - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let finished_list: &[Value] = &app_state.downloads.finished.write().await.drain(..) .map(|fd| download_info!(fd)).collect::>(); diff --git a/src/api/endpoints/hdhomerun_api.rs b/src/api/endpoints/hdhomerun_api.rs new file mode 100644 index 000000000..2cbbe171b --- /dev/null +++ b/src/api/endpoints/hdhomerun_api.rs @@ -0,0 +1,319 @@ +use std::sync::Arc; +use crate::api::model::app_state::HdHomerunAppState; +use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; +use crate::model::config::{Config, TargetType}; +use crate::model::playlist::{M3uPlaylistItem, XtreamCluster, XtreamPlaylistItem}; +use crate::processing::parser::xtream::get_xtream_url; +use crate::utils::json_utils::get_string_from_serde_value; +// https://info.hdhomerun.com/info/http_api +use actix_web::{web, HttpResponse, Responder}; +use bytes::Bytes; +use futures::{stream, Stream, StreamExt}; +use log::{error, warn}; +use serde::{Deserialize, Serialize}; +use serde_json::{json}; +use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; +use crate::repository::xtream_playlist_iterator::{XtreamPlaylistIterator}; +// const DISCOVERY_BYTES: &[u8] = &[0, 2, 0, 12, 1, 4, 255, 255, 255, 255, 2, 4, 255, 255, 255, 255, 115, 204, 125, 143]; +// const RESPONSE_BYTES: &[u8] = &[0, 3, 0, 12, 1, 4, 255, 255, 255, 255, 2, 4, 255, 255, 255, 255, 115, 204, 125, 143]; + +#[derive(Serialize, Deserialize, Clone)] +struct Lineup { + #[serde(rename = "GuideNumber")] + guide_number: String, + #[serde(rename = "GuideName")] + guide_name: String, + #[serde(rename = "URL")] + url: String, +} + +#[derive(Serialize, Deserialize, Clone)] +struct Device { + #[serde(rename = "FriendlyName")] + friendly_name: String, + #[serde(rename = "Manufacturer")] + manufacturer: String, + // #[serde(rename = "ManufacturerURL")] + // manufacturer_url: String, + #[serde(rename = "ModelNumber")] + model_number: String, + #[serde(rename = "ModelName")] + model_name: String, + #[serde(rename = "FirmwareName")] + firmware_name: String, + #[serde(rename = "TunerCount")] + tuner_count: u8, + #[serde(rename = "FirmwareVersion")] + firmware_version: String, + #[serde(rename = "DeviceID")] + id: String, + #[serde(rename = "DeviceAuth")] + auth: String, + #[serde(rename = "BaseURL")] + base_url: String, + #[serde(rename = "LineupURL")] + lineup_url: String, + #[serde(rename = "DiscoverURL")] + discover_url: String, + +} + +impl Device { + fn as_xml(&self) -> String { + format!(r#" + +1 +0 + +{} + + urn:dial-multicast:com.silicondust.hdhomerun + {} + {} + {} + {} + {} + uuid:{} + +"#, + self.base_url, self.friendly_name, self.manufacturer, self.model_number, + self.model_number, self.id, self.id + ) + } +} + +fn xtream_item_to_lineup_stream(cfg: Arc, cluster: XtreamCluster, credentials: Arc, + base_url: Option, channels: Option) -> impl Stream> +where + I: Iterator + 'static, +{ + match channels { + Some(chans) => { + let mapped = chans.map(move |(item, has_next)| { + let input = cfg.get_input_by_name(&item.input_name); + let (live_stream_use_prefix, live_stream_without_extension) = input.map_or((true, false), |i| i.options.as_ref() + .map_or((true, false), |o| (o.xtream_live_stream_use_prefix, o.xtream_live_stream_without_extension))); + let container_extension = item.get_additional_property("container_extension").map(|v| get_string_from_serde_value(&v).unwrap_or_default()); + let stream_url = match &base_url { + None => item.url.to_string(), + Some(url) => get_xtream_url(cluster, url, &credentials.username, &credentials.password, item.virtual_id, container_extension.as_ref(), live_stream_use_prefix, live_stream_without_extension) + }; + + let lineup = Lineup { + guide_number: item.epg_channel_id.unwrap_or(item.name).to_string(), + guide_name: item.title.to_string(), + url: stream_url, + }; + match serde_json::to_string(&lineup) { + Ok(content) => { + Ok(Bytes::from(if has_next { + format!("{content},") + } else { + content + })) + } + Err(_) => Ok(Bytes::from("")), + } + }); + stream::iter(mapped).left_stream() + } + None => { + stream::once(async { Ok(Bytes::from("")) }).right_stream() + } + } +} + +fn m3u_item_to_lineup_stream(channels: Option) -> impl Stream> +where + I: Iterator + 'static, +{ + match channels { + Some(chans) => { + let mapped = chans.map(move |(item, has_next)| { + let lineup = Lineup { + guide_number: item.epg_channel_id.unwrap_or(item.name).to_string(), + guide_name: item.title.to_string(), + url: (if item.t_stream_url.is_empty() {&item.url} else {&item.t_stream_url}).to_string(), + }; + match serde_json::to_string(&lineup) { + Ok(content) => { + Ok(Bytes::from(if has_next { + format!("{content},") + } else { + content + })) + } + Err(_) => Ok(Bytes::from("")), + } + }); + stream::iter(mapped).left_stream() + } + None => { + stream::once(async { Ok(Bytes::from("")) }).right_stream() + } + } +} + +fn create_device(app_state: &web::Data) -> Option { + if let Some(credentials) = app_state.app_state.config.get_user_credentials(&app_state.device.t_username) { + let server_info = app_state.app_state.config.get_user_server_info(&credentials); + let device = &app_state.device; + let device_url = format!("{}://{}:{}", server_info.protocol, server_info.host, device.port); + Some(Device { + friendly_name: device.friendly_name.to_string(), + manufacturer: device.manufacturer.to_string(), + //manufacturer_url: "https://github.com/euzu/m3u-filter".to_string(), + model_number: device.model_number.to_string(), + model_name: device.model_name.to_string(), + firmware_name: device.firmware_name.to_string(), + tuner_count: 1, // app_state.config.hdhomerun.tuner_count, + firmware_version: device.firmware_version.to_string(), + auth: String::new(), + id: device.device_udn.to_string(), + lineup_url: format!("{device_url}/lineup.json"), + discover_url: format!("{device_url}/discover.json"), + base_url: device_url, + }) + } else { + error!("Failed to get credentials for username: {} for device: {} ", &app_state.device.t_username, &app_state.device.name); + None + } +} + +async fn device_xml(app_state: web::Data) -> impl Responder { + if let Some(device) = create_device(&app_state) { + HttpResponse::Ok().content_type("application/xml").body(device.as_xml()) + } else { + HttpResponse::InternalServerError().finish() + } +} + +async fn device_json(app_state: web::Data) -> impl Responder { + if let Some(device) = create_device(&app_state) { + HttpResponse::Ok().json(device) + } else { + HttpResponse::InternalServerError().finish() + } +} + +async fn discover_json(app_state: web::Data) -> impl Responder { + if let Some(device) = create_device(&app_state) { + HttpResponse::Ok() + .content_type("application/json") + .json(device) + } else { + HttpResponse::InternalServerError().finish() + } +} + +async fn lineup_status() -> impl Responder { + HttpResponse::Ok() + .content_type("application/json") + .json(json!({ + "ScanInProgress": 0, + "ScanPossible": 0, + "Source": "Cable", + "SourceList": ["Cable"], + })) +} + +async fn lineup_json(app_state: web::Data) -> impl Responder { + let cfg = Arc::clone(&app_state.app_state.config); + if let Some((credentials, target)) = cfg.get_target_for_username(&app_state.device.t_username) { + if target.has_output(&TargetType::Xtream) { + let server_info = app_state.app_state.config.get_user_server_info(&credentials); + let base_url = if credentials.proxy == ProxyType::Reverse { + Some(server_info.get_base_url()) + } else { + None + }; + + let live_channels = XtreamPlaylistIterator::new(XtreamCluster::Live, &cfg, target, 0, &credentials).await.ok(); + let vod_channels = XtreamPlaylistIterator::new(XtreamCluster::Video, &cfg, target, 0, &credentials).await.ok(); + // TODO include when resolved + //let series_channels = xtream_repository::iter_raw_xtream_playlist(cfg, target, XtreamCluster::Series); + let user_credentials = Arc::new(credentials); + let live_stream = xtream_item_to_lineup_stream(Arc::clone(&cfg), XtreamCluster::Live, Arc::clone(&user_credentials), base_url.clone(), live_channels); + let vod_stream = xtream_item_to_lineup_stream(Arc::clone(&cfg), XtreamCluster::Video, Arc::clone(&user_credentials), base_url.clone(), vod_channels); + let body_stream = stream::once(async { Ok(Bytes::from("[")) }) + .chain(live_stream) + .chain(stream::once(async { Ok(Bytes::from(",")) })) + .chain(vod_stream) + .chain(stream::once(async { Ok(Bytes::from("]")) })); + return HttpResponse::Ok() + .content_type("application/json") + .streaming(body_stream); + } else if target.has_output(&TargetType::M3u) { + let iterator = M3uPlaylistIterator::new(&cfg,target,&credentials).ok(); + let stream = m3u_item_to_lineup_stream(iterator); + let body_stream = stream::once(async { Ok(Bytes::from("[")) }) + .chain(stream) + .chain(stream::once(async { Ok(Bytes::from("]")) })); + return HttpResponse::Ok() + .content_type("application/json") + .streaming(body_stream); + } + } + HttpResponse::NotFound().finish() +} + +async fn auto_channel(_app_state: web::Data, path: web::Path) -> impl Responder { + let channel = path.into_inner(); + warn!("HdHomerun api not implemented for auto_channel {channel}"); + HttpResponse::NotFound().finish() +} + +pub fn hdhr_api_register(cfg: &mut web::ServiceConfig) { + cfg.service(web::resource("/device.xml").route(web::get().to(device_xml))); + cfg.service(web::resource("/device.json").route(web::get().to(device_json))); + cfg.service(web::resource("/discover.json").route(web::get().to(discover_json))); + cfg.service(web::resource("/lineup_status.json").route(web::get().to(lineup_status))); + cfg.service(web::resource("/lineup.json").route(web::get().to(lineup_json))); + // cfg.service(web::resource("/lineup.xml").route(web::get().to(lineup_xml))); + // cfg.service(web::resource("/lineup.m3u").route(web::get().to(lineup_m3u))); + cfg.service(web::resource("/auto/{channel}").route(web::get().to(auto_channel))); + cfg.service(web::resource("/tuner{tuner_num}/{channel}").route(web::get().to(auto_channel))); +} + +// fn start_hdhomerum_discovery_handler(ssdp_socket: Arc, server: String, location: String, cache_control: String, usn: String) { +// let mut buffer = [0; 4096]; +// actix_rt::spawn(async move { +// let response_bytes = RESPONSE_BYTES; +// loop { +// match ssdp_socket.recv_from(&mut buffer).await { +// Ok((size, src_addr)) => { +// let content = &buffer[..size]; +// if content == DISCOVERY_BYTES { +// match ssdp_socket.send_to(&response_bytes, src_addr).await { +// Err(err) => eprintln!("Failed to send SSDP response: {err:?}"), +// Ok(_) => println!("Sent SSDP response to: {src_addr:?}"), +// } +// } +// } +// Err(err) => eprintln!("Failed to receive SSDP request: {err:?}"), +// } +// } +// }); +// } +// +// pub async fn start_hdhomerun(/*host: &str, */port: u16) { +// let version = "2021.08.18"; +// let server_url = format!("http://10.41.41.89:{port}"); +// +// // let multicast_addr: Ipv4Addr = "255.255.255.255".parse().unwrap(); +// +// let socket_addr: SocketAddr = "0.0.0.0:65001".parse().unwrap(); +// let socket = Socket::new(Domain::IPV4, Type::DGRAM, None).unwrap(); +// // setting SO_REUSEADDR-Option if other dlna server is running +// socket.set_reuse_address(true).unwrap(); +// socket.bind(&socket_addr.into()).unwrap(); +// let udp_socket = UdpSocket::from_std(socket.into()).unwrap(); +// +// let ssdp_socket = Arc::new(udp_socket); +// // ssdp_socket.join_multicast_v4(multicast_addr, "0.0.0.0".parse().unwrap()).unwrap(); +// let server = format!("SERVER: HDHomeRun/{}", version); +// let location = format!("LOCATION: {server_url}/device.xml"); +// let cache_control = "CACHE-CONTROL: max-age=1800"; +// let usn = "USN: uuid:12345678-90ab-cdef-1234-567890abcdef::urn:dial-multicast:com.silicondust.hdhomerun"; +// start_hdhomerum_discovery_handler(Arc::clone(&ssdp_socket), server.to_string(), location.to_string(), cache_control.to_string(), usn.to_string()); +// } diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index 65da4de74..ad42327c1 100644 --- a/src/api/endpoints/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -2,6 +2,7 @@ use std::sync::Arc; use actix_web::{web, HttpRequest, HttpResponse}; use actix_web::web::Data; use log::{debug, error}; +use serde::Deserialize; use crate::api::api_utils::{get_user_target_by_credentials, stream_response}; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; @@ -15,7 +16,18 @@ use crate::repository::playlist_repository::HLS_EXT; use crate::utils::network::request; use crate::utils::network::request::{replace_extension, sanitize_sensitive_info}; -pub(in crate::api) async fn handle_hls_stream_request(app_state: &Data, user: &ProxyUserCredentials, pli: &dyn PlaylistEntry, input: &ConfigInput, target_type: TargetType) -> HttpResponse { +#[derive(Deserialize)] +#[allow(dead_code)] +struct HlsApiPathParams { + token: String, + username: String, + password: String, + channel: String, + hash: String, + chunk: String, +} + +pub(in crate::api) async fn handle_hls_stream_request(app_state: &Data>, user: &ProxyUserCredentials, pli: &dyn PlaylistEntry, input: &ConfigInput, target_type: TargetType) -> HttpResponse { let url = replace_extension(&pli.get_provider_url(), HLS_EXT); match request::download_text_content(Arc::clone(&app_state.http_client), input, &url, None).await { Ok(content) => { @@ -32,21 +44,21 @@ pub(in crate::api) async fn handle_hls_stream_request(app_state: &Data async fn hls_api_stream( req: &HttpRequest, api_req: &web::Query, - path: web::Path<(String, String, String, String, String, String)>, - app_state: &web::Data, + path: web::Path, + app_state: &web::Data>, target_type: TargetType ) -> HttpResponse { - let (_token, username, password, channel, _hash, _chunk) = path.into_inner(); + let params = path.into_inner(); let (user, target) = try_option_bad_request!( - get_user_target_by_credentials(&username, &password, api_req, app_state).await, + get_user_target_by_credentials(¶ms.username, ¶ms.password, api_req, app_state).await, false, - format!("Could not find any user {username}")); + format!("Could not find any user {}", params.username)); if !user.has_permissions(app_state) { return HttpResponse::Forbidden().finish(); } let target_name = &target.name; - let virtual_id: u32 = try_result_bad_request!(channel.parse()); + let virtual_id: u32 = try_result_bad_request!(params.channel.parse()); let (pli_url, input_name) = if target_type == TargetType::Xtream { let (pli, _ ) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); (pli.url, pli.input_name) @@ -68,8 +80,8 @@ async fn hls_api_stream( async fn hls_api_stream_xtream( req: HttpRequest, api_req: web::Query, - path: web::Path<(String, String, String, String, String, String)>, - app_state: web::Data, + path: web::Path, + app_state: web::Data>, ) -> HttpResponse { hls_api_stream(&req, &api_req, path, &app_state, TargetType::Xtream).await } @@ -77,13 +89,12 @@ async fn hls_api_stream_xtream( async fn hls_api_stream_m3u( req: HttpRequest, api_req: web::Query, - path: web::Path<(String, String, String, String, String, String)>, - app_state: web::Data, + path: web::Path, + app_state: web::Data>, ) -> HttpResponse { hls_api_stream(&req, &api_req, path, &app_state, TargetType::M3u).await } - pub fn hls_api_register(cfg: &mut web::ServiceConfig) { cfg.service(web::resource("/hlsr/{token}/{username}/{password}/{channel}/{hash}/{chunk}").route(web::get().to(hls_api_stream_xtream))); cfg.service(web::resource(format!("/{M3U_HLSR_PREFIX}/{{token}}/{{username}}/{{password}}/{{channel}}/{{hash}}/{{chunk}}")).route(web::get().to(hls_api_stream_m3u))); diff --git a/src/api/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index adc722160..933d3382b 100644 --- a/src/api/endpoints/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use actix_web::{web, HttpRequest, HttpResponse}; use bytes::Bytes; use futures::stream; @@ -27,7 +28,7 @@ async fn m3u_api( match m3u_load_rewrite_playlist(&app_state.config, target, &user).await { Ok(m3u_iter) => { // Convert the iterator into a stream of `Bytes` - let content_stream = stream::iter(m3u_iter.map(|line| Ok::(Bytes::from([line.as_bytes(), b"\n"].concat())))); + let content_stream = stream::iter(m3u_iter.map(|line| Ok::(Bytes::from([line.to_string().as_bytes(), b"\n"].concat())))); let mut builder = HttpResponse::Ok(); builder.content_type(mime::TEXT_PLAIN_UTF_8); if api_req.content_type == "m3u_plus" { @@ -46,13 +47,13 @@ async fn m3u_api( } async fn m3u_api_get(api_req: web::Query, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { m3u_api(&api_req.into_inner(), &app_state).await } async fn m3u_api_post( api_req: web::Form, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { m3u_api(&api_req.into_inner(), &app_state).await } @@ -61,7 +62,7 @@ async fn m3u_api_stream( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String)>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); let (action_stream_id, stream_ext) = separate_number_and_remainder(&stream_id); @@ -107,7 +108,7 @@ async fn m3u_api_resource( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String, String)>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let (username, password, stream_id, resource) = path.into_inner(); let Ok(m3u_stream_id) = stream_id.parse::() else { return HttpResponse::BadRequest().finish() }; diff --git a/src/api/endpoints/mod.rs b/src/api/endpoints/mod.rs index c6f8a52ca..d2926032c 100644 --- a/src/api/endpoints/mod.rs +++ b/src/api/endpoints/mod.rs @@ -5,4 +5,5 @@ pub(in crate::api) mod m3u_api; pub(in crate::api) mod xmltv_api; pub(in crate::api) mod web_index; pub(in crate::api) mod hls_api; -mod user_api; \ No newline at end of file +mod user_api; +pub(in crate::api) mod hdhomerun_api; \ No newline at end of file diff --git a/src/api/endpoints/user_api.rs b/src/api/endpoints/user_api.rs index 28fc2f1a1..d428277c5 100644 --- a/src/api/endpoints/user_api.rs +++ b/src/api/endpoints/user_api.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use crate::api::api_utils::{get_user_target_by_username, get_username_from_auth_header}; use crate::api::model::app_state::AppState; use crate::auth::authenticator::validator_user; @@ -39,7 +40,7 @@ pub(crate) async fn get_categories_content(action: Result<(Option, Opti async fn playlist_categories( credentials: Option, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { if let Some(username) = get_username_from_auth_header(credentials, &app_state) { if let Some((user, target)) = get_user_target_by_username(username.as_str(), &app_state).await { @@ -73,7 +74,7 @@ async fn playlist_categories( async fn save_playlist_bouquet( credentials: Option, - app_state: web::Data, + app_state: web::Data>, req: web::Json, ) -> HttpResponse { if let Some(username) = get_username_from_auth_header(credentials, &app_state) { @@ -96,7 +97,7 @@ async fn save_playlist_bouquet( async fn playlist_bouquet( credentials: Option, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { if let Some(username) = get_username_from_auth_header(credentials, &app_state) { if let Some((user, _target)) = get_user_target_by_username(username.as_str(), &app_state).await { diff --git a/src/api/endpoints/v1_api.rs b/src/api/endpoints/v1_api.rs index 76ee1e097..ba87d23db 100644 --- a/src/api/endpoints/v1_api.rs +++ b/src/api/endpoints/v1_api.rs @@ -6,7 +6,7 @@ use actix_web::middleware::Condition; use actix_web::{web, HttpResponse}; use actix_web_httpauth::middleware::HttpAuthentication; use bytes::Bytes; -use futures::{stream, Stream, StreamExt}; +use futures::{stream, StreamExt}; use log::error; use serde_json::json; @@ -19,12 +19,12 @@ use crate::auth::authenticator::validator_admin; use crate::m3u_filter_error::M3uFilterError; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, TargetUser}; use crate::model::config::{validate_targets, Config, ConfigDto, ConfigInput, ConfigInputOptions, ConfigSource, ConfigTarget, InputType, TargetType}; -use crate::model::playlist::{XtreamCluster, XtreamPlaylistItem}; +use crate::model::playlist::{XtreamCluster}; use crate::processing::processor::playlist; use crate::repository::user_repository::store_api_user; use crate::repository::xtream_repository; +use crate::repository::xtream_repository::playlist_iter_to_stream; use crate::utils::file::config_reader; -use crate::utils::file::file_lock_manager::FileReadGuard; use crate::utils::network::request::sanitize_sensitive_info; use crate::utils::network::{m3u, xtream}; @@ -52,7 +52,7 @@ fn intern_save_config_main(file_path: &str, backup_dir: &str, cfg: &ConfigDto) - async fn save_config_api_proxy_user( req: web::Json>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let mut users = req.0; let mut usernames = HashSet::new(); @@ -95,7 +95,7 @@ async fn save_config_api_proxy_user( async fn save_config_main( req: web::Json, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let cfg = req.0; if cfg.is_valid() { @@ -112,7 +112,7 @@ async fn save_config_main( async fn save_config_api_proxy_config( req: web::Json>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let mut req_api_proxy = req.0; for server_info in &mut req_api_proxy { @@ -132,7 +132,7 @@ async fn save_config_api_proxy_config( async fn playlist_update( req: web::Json>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let targets = req.0; let user_targets = if targets.is_empty() { None } else { Some(targets) }; @@ -206,34 +206,9 @@ async fn get_playlist(client: Arc, cfg_input: Option<&ConfigInp } } -fn to_stream(channels: Option<(FileReadGuard, I)>) -> impl Stream> -where - I: Iterator + 'static, -{ - match channels { - Some((_, chans)) => { - // Convert iterator items to Result - let mapped = chans.map(move |(item, has_next)| { - match serde_json::to_string(&item) { - Ok(content) => { - Ok(Bytes::from(if has_next { - format!("{content},") - } else { - content - })) - } - Err(_) => Ok(Bytes::from("")), - } - }); - stream::iter(mapped).left_stream() - } - None => { - stream::once(async { Ok(Bytes::from("")) }).right_stream() - } - } -} -async fn get_playlist_for_target(cfg_target: Option<&ConfigTarget>, cfg: &Config) -> HttpResponse { + +async fn get_playlist_for_target(cfg_target: Option<&ConfigTarget>, cfg: &Arc) -> HttpResponse { if let Some(target) = cfg_target { let target_name = &target.name; if target.has_output(&TargetType::Xtream) { @@ -245,9 +220,9 @@ async fn get_playlist_for_target(cfg_target: Option<&ConfigTarget>, cfg: &Config let vod_channels = xtream_repository::iter_raw_xtream_playlist(cfg, target, XtreamCluster::Video); let series_channels = xtream_repository::iter_raw_xtream_playlist(cfg, target, XtreamCluster::Series); - let live_stream = to_stream(live_channels); - let vod_stream = to_stream(vod_channels); - let series_stream = to_stream(series_channels); + let live_stream = playlist_iter_to_stream(live_channels); + let vod_stream = playlist_iter_to_stream(vod_channels); + let series_stream = playlist_iter_to_stream(series_channels); let json_stream = stream::iter(vec![ @@ -275,7 +250,7 @@ async fn get_playlist_for_target(cfg_target: Option<&ConfigTarget>, cfg: &Config async fn playlist( req: web::Json, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { match req.rtype { PlaylistRequestType::Input => { @@ -312,7 +287,7 @@ async fn playlist( } async fn config( - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let map_input = |i: &ConfigInput| ServerInputConfig { id: i.id, diff --git a/src/api/endpoints/web_index.rs b/src/api/endpoints/web_index.rs index 105f20185..af5a5335d 100644 --- a/src/api/endpoints/web_index.rs +++ b/src/api/endpoints/web_index.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use std::collections::HashMap; use std::path::{Path, PathBuf}; use actix_files::NamedFile; @@ -14,7 +15,7 @@ fn no_web_auth_token() -> HttpResponse { async fn token( mut req: web::Json, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { match &app_state.config.web_auth { None => no_web_auth_token(), @@ -53,7 +54,7 @@ async fn token( async fn token_refresh( _req: HttpRequest, credentials: Option, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { match &app_state.config.web_auth { None => { @@ -84,7 +85,7 @@ async fn token_refresh( async fn index( _req: HttpRequest, - app_state: web::Data, + app_state: web::Data>, ) -> std::io::Result { let path: PathBuf = [&app_state.config.api.web_root, "index.html"].iter().collect(); NamedFile::open(path) diff --git a/src/api/endpoints/xmltv_api.rs b/src/api/endpoints/xmltv_api.rs index d90967590..c571ed650 100644 --- a/src/api/endpoints/xmltv_api.rs +++ b/src/api/endpoints/xmltv_api.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use std::fs::File; use std::path::{Path, PathBuf}; @@ -65,8 +66,8 @@ fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option {} - } + TargetType::Strm | TargetType::HdHomeRun => {} + } } None } @@ -166,7 +167,7 @@ fn serve_epg_with_timeshift(epg_file: File, offset_minutes: i32) -> HttpResponse async fn xmltv_api( api_req: web::Query, req: HttpRequest, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { if let Some((user, target)) = get_user_target(&api_req, &app_state).await { if !user.has_permissions(&app_state) { diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index 41689843c..c88b28ce2 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -141,7 +141,7 @@ fn get_user_info(user: &ProxyUserCredentials, app_state: &AppState) -> XtreamAut async fn xtream_player_api_stream( req: &HttpRequest, api_req: &web::Query, - app_state: &web::Data, + app_state: &web::Data>, stream_req: XtreamApiStreamRequest<'_>, ) -> HttpResponse { let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state).await, false, format!("Could not find any user {}", stream_req.username)); @@ -332,7 +332,7 @@ fn xtream_get_season_resource_url(config: &Config, pli: &XtreamPlaylistItem, tar async fn xtream_player_api_resource( req: &HttpRequest, api_req: &web::Query, - app_state: &web::Data, + app_state: &web::Data>, resource_req: XtreamApiStreamRequest<'_>, ) -> HttpResponse { let (user, target) = try_option_bad_request!(get_user_target_by_credentials(resource_req.username, resource_req.password, api_req, app_state).await, false, format!("Could not find any user {}", resource_req.username)); @@ -375,7 +375,7 @@ macro_rules! create_xtream_player_api_stream { req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String)>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamRequest::from($context, &username, &password, &stream_id, "")).await @@ -389,7 +389,7 @@ macro_rules! create_xtream_player_api_resource { req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String, String)>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let (username, password, stream_id, resource) = path.into_inner(); xtream_player_api_resource(&req, &api_req, &app_state, XtreamApiStreamRequest::from($context, &username, &password, &stream_id, &resource)).await @@ -421,7 +421,7 @@ async fn xtream_player_api_timeshift_stream( api_query_req: web::Query, api_form_req: web::Form, path: web::Path<(String, String, String, String, String)>, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { let (path_username, path_password, path_duration, path_start, path_stream_id) = path.into_inner(); let username = get_non_empty(&path_username, &api_query_req.username, &api_form_req.username); @@ -579,7 +579,7 @@ macro_rules! skip_flag_optional { async fn xtream_player_api( req: &HttpRequest, api_req: UserApiRequest, - app_state: &web::Data, + app_state: &web::Data>, ) -> HttpResponse { let user_target = get_user_target(&api_req, app_state).await; if let Some((user, target)) = user_target { @@ -675,30 +675,27 @@ async fn xtream_player_api( } } -fn xtream_create_content_stream(xtream_iter: impl Iterator) -> impl Stream> { - let mut first_item = true; +fn xtream_create_content_stream(xtream_iter: impl Iterator) -> impl Stream> { stream::once(async { Ok::(Bytes::from("[")) }).chain( - stream::iter(xtream_iter.map(move |line| { - let line = if first_item { - first_item = false; - line + stream::iter(xtream_iter.map(move |(line, has_next)| { + Ok::(Bytes::from(if has_next { + format!("{line},") } else { - format!(",{line}") - }; - Ok::(Bytes::from(line)) + line.to_string() + })) })).chain(stream::once(async { Ok::(Bytes::from("]")) }))) } async fn xtream_player_api_get(req: HttpRequest, api_req: web::Query, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { xtream_player_api(&req, api_req.into_inner(), &app_state).await } async fn xtream_player_api_post(req: HttpRequest, api_req: web::Form, - app_state: web::Data, + app_state: web::Data>, ) -> HttpResponse { xtream_player_api(&req, api_req.into_inner(), &app_state).await } diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 5ceac7529..518e17f27 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -11,7 +11,7 @@ use chrono::{DateTime, Utc}; use mime::APPLICATION_JSON; use crate::api::endpoints::hls_api::hls_api_register; use crate::api::endpoints::m3u_api::m3u_api_register; -use crate::api::model::app_state::AppState; +use crate::api::model::app_state::{AppState, HdHomerunAppState}; use crate::api::model::download::DownloadQueue; use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::api::scheduler::start_scheduler; @@ -27,6 +27,7 @@ use crate::tools::lru_cache::{LRUResourceCache}; use crate::utils::size_utils::human_readable_byte_size; use crate::utils::sys_utils; use crate::{BUILD_TIMESTAMP, VERSION}; +use crate::api::endpoints::hdhomerun_api::{hdhr_api_register}; use crate::api::model::active_provider_manager::ActiveProviderManager; fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result { @@ -39,7 +40,7 @@ fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result) -> Healthcheck { +fn create_healthcheck(app_state: &web::Data>) -> Healthcheck { let server_time = chrono::offset::Local::now().with_timezone(&chrono::Local).format("%Y-%m-%d %H:%M:%S %Z").to_string(); let cache = app_state.cache.as_ref().as_ref().map(|c| c.lock().get_size_text()); let (active_clients, active_connections) = { @@ -59,11 +60,11 @@ fn create_healthcheck(app_state: &web::Data) -> Healthcheck { } } -async fn healthcheck(app_state: web::Data,) -> HttpResponse { +async fn healthcheck(app_state: web::Data>,) -> HttpResponse { HttpResponse::Ok().json(create_healthcheck(&app_state)) } -async fn status(app_state: web::Data,) -> HttpResponse { +async fn status(app_state: web::Data>,) -> HttpResponse { let status = create_healthcheck(&app_state); match serde_json::to_string_pretty(&status) { Ok(pretty_json) => HttpResponse::Ok().content_type(APPLICATION_JSON).body(pretty_json), @@ -71,7 +72,7 @@ async fn status(app_state: web::Data,) -> HttpResponse { } } -fn create_shared_data(cfg: &Arc) -> Data { +fn create_shared_data(cfg: &Arc) -> AppState { let lru_cache = cfg.reverse_proxy.as_ref().and_then(|r| r.cache.as_ref()).and_then(|c| if c.enabled { Some(PlMutex::new(LRUResourceCache::new(c.t_size, &PathBuf::from(c.dir.as_ref().unwrap())))) } else { None} ); @@ -86,7 +87,7 @@ fn create_shared_data(cfg: &Arc) -> Data { } }); let user_access_control = cfg.user_access_control; - Data::new(AppState { + AppState { config: Arc::clone(cfg), http_client: Arc::new(reqwest::Client::new()), downloads: Arc::from(DownloadQueue::new()), @@ -94,7 +95,7 @@ fn create_shared_data(cfg: &Arc) -> Data { shared_stream_manager: Arc::new(SharedStreamManager::new()), active_users: Arc::new(ActiveUserManager::new()), active_provider: Arc::new(ActiveProviderManager::new(user_access_control)), - }) + } } fn exec_update_on_boot(client: Arc, cfg: &Arc, targets: &Arc) { @@ -159,8 +160,43 @@ fn is_web_auth_enabled(cfg: &Arc, web_ui_enabled: bool) -> bool { false } +fn start_hdhomerun(cfg: &Arc, app_state: &Arc, infos: &mut Vec) { + let host = cfg.api.host.to_string(); + if let Some(hdhomerun) = &cfg.hdhomerun { + if hdhomerun.enabled { + for device in &hdhomerun.devices { + if device.t_enabled { + let app_data = Arc::clone(app_state); + let app_host = host.clone(); + let port = device.port; + let device_clone = Arc::new(device.clone()); + infos.push(format!("HdHomeRun Server '{}' running: http://{host}:{port}", device.name)); + actix_rt::spawn(async move { + HttpServer::new(move || { + App::new() + .wrap(Logger::default()) + .wrap(Cors::default() + .supports_credentials() + .allow_any_origin() + .allowed_methods(vec!["GET", "POST", "OPTIONS", "HEAD"]) + .allow_any_header() + .max_age(3600)) + .app_data(Data::new(HdHomerunAppState { + app_state: Arc::clone(&app_data), + device: Arc::clone(&device_clone), + })) + .configure(hdhr_api_register) + }).bind(format!("{}:{port}", app_host.clone()))?.run().await + }); + } + } + } + } +} + #[actix_web::main] pub async fn start_server(cfg: Arc, targets: Arc) -> futures::io::Result<()> { + let mut infos = Vec::new(); let host = cfg.api.host.to_string(); let port = cfg.api.port; let web_ui_enabled = cfg.web_ui_enabled; @@ -169,14 +205,23 @@ pub async fn start_server(cfg: Arc, targets: Arc) -> fut Err(err) => return Err(err) }; if web_ui_enabled { - info!("Web root: {:?}", &web_dir_path); + infos.push(format!("Web root: {:?}", &web_dir_path)); } - let shared_data = create_shared_data(&cfg); + let app_state = Arc::new(create_shared_data(&cfg)); + let shared_data = Data::new(Arc::clone(&app_state)); exec_scheduler(&Arc::clone(&shared_data.http_client), &cfg, &targets); exec_update_on_boot(Arc::clone(&shared_data.http_client), &cfg, &targets); let web_auth_enabled = is_web_auth_enabled(&cfg, web_ui_enabled); + if cfg.t_api_proxy.read().is_some() { + start_hdhomerun(&cfg, &app_state, &mut infos); + } + + infos.push(format!("Server running: http://{}:{}", &cfg.api.host, &cfg.api.port)); + for info in &infos { + info!("{info}"); + } // Web Server HttpServer::new(move || { App::new() diff --git a/src/api/model/app_state.rs b/src/api/model/app_state.rs index 93524093d..c56a2e6e7 100644 --- a/src/api/model/app_state.rs +++ b/src/api/model/app_state.rs @@ -5,6 +5,7 @@ use crate::api::model::active_user_manager::ActiveUserManager; use crate::api::model::download::DownloadQueue; use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::model::config::{Config}; +use crate::model::hdhomerun_config::HdHomeRunDeviceConfig; use crate::tools::lru_cache::LRUResourceCache; pub struct AppState { @@ -22,3 +23,8 @@ impl AppState { self.active_users.user_connections(username) } } + +pub struct HdHomerunAppState { + pub app_state: Arc, + pub device: Arc, +} diff --git a/src/auth/authenticator.rs b/src/auth/authenticator.rs index ee4d2054d..b62bd85a0 100644 --- a/src/auth/authenticator.rs +++ b/src/auth/authenticator.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use actix_web::{dev::ServiceRequest, Error, web}; use actix_web_httpauth::extractors::bearer::BearerAuth; use chrono::{Local, Duration}; @@ -84,7 +85,7 @@ fn validate_request( credentials: Option, verify_fn: fn(Option, &[u8]) -> bool, // Funktions-Parameter für Admin/User-Check ) -> Result { - if let Some(app_state) = req.app_data::>() { + if let Some(app_state) = req.app_data::>>() { if let Some(web_auth_config) = app_state.config.web_auth.as_ref() { let secret_key = web_auth_config.secret.as_ref(); if verify_fn(credentials, secret_key) { diff --git a/src/main.rs b/src/main.rs index b076339ca..0a99a4d4f 100644 --- a/src/main.rs +++ b/src/main.rs @@ -156,8 +156,12 @@ fn main() { } if args.server { - if let Some(api_proxy_file) = config_reader::read_api_proxy_config(args.api_proxy, &mut cfg) { - info!("Api Proxy File: {api_proxy_file:?}"); + match config_reader::read_api_proxy_config(args.api_proxy, &mut cfg) { + Ok(Some(api_proxy_file)) => { + info!("Api Proxy File: {api_proxy_file:?}"); + }, + Ok(None) => {} + Err(err) => exit!("{err}"), } start_in_server_mode(Arc::new(cfg), Arc::new(targets)); } else { @@ -200,7 +204,6 @@ fn start_in_cli_mode(cfg: Arc, targets: Arc) { } fn start_in_server_mode(cfg: Arc, targets: Arc) { - info!("Server running: http://{}:{}", &cfg.api.host, &cfg.api.port); if let Err(err) = api::main_api::start_server(cfg, targets) { exit!("Can't start server: {err}"); }; diff --git a/src/model/config.rs b/src/model/config.rs index 1bcb3cf99..81b73cf5f 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -1,5 +1,6 @@ #![allow(clippy::struct_excessive_bools)] use enum_iterator::Sequence; +use parking_lot::RwLock; use std::borrow::BorrowMut; use std::collections::{HashMap, HashSet}; use std::fmt::Display; @@ -7,7 +8,6 @@ use std::fs::File; use std::io::BufRead; use std::path::PathBuf; use std::str::FromStr; -use parking_lot::RwLock; use std::sync::Arc; use crate::auth::user::UserCredential; @@ -22,8 +22,8 @@ use crate::messaging::MsgKind; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, ProxyUserCredentials}; use crate::model::mapping::Mapping; use crate::model::mapping::Mappings; -use crate::utils::file::config_reader; use crate::utils::default_utils::{default_as_default, default_as_true, default_as_two_u16}; +use crate::utils::file::config_reader; use crate::utils::file::file_lock_manager::FileLockManager; use crate::utils::file::file_utils; use crate::utils::file::file_utils::file_reader; @@ -51,6 +51,7 @@ macro_rules! valid_property { } pub use valid_property; use crate::m3u_filter_error::{create_m3u_filter_error_result, handle_m3u_filter_error_result, handle_m3u_filter_error_result_list}; +use crate::model::hdhomerun_config::HdHomeRunConfig; use crate::utils::string_utils::get_trimmed_string; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Hash)] @@ -61,12 +62,15 @@ pub enum TargetType { Xtream, #[serde(rename = "strm")] Strm, + #[serde(rename = "hdhomerun")] + HdHomeRun, } impl TargetType { const M3U: &'static str = "M3u"; const XTREAM: &'static str = "Xtream"; const STRM: &'static str = "Strm"; + const HDHOMERUN: &'static str = "HdHomeRun"; } impl Display for TargetType { @@ -75,6 +79,7 @@ impl Display for TargetType { Self::M3u => Self::M3U, Self::Xtream => Self::XTREAM, Self::Strm => Self::STRM, + Self::HdHomeRun => Self::HDHOMERUN, }) } } @@ -305,8 +310,8 @@ pub struct ConfigTargetOptions { pub struct TargetOutput { #[serde(alias = "type")] pub target: TargetType, - #[serde(skip_serializing_if = "Option::is_none")] - pub filename: Option, + #[serde(skip_serializing_if = "Option::is_none", alias = "filename")] + pub output: Option, #[serde(skip_serializing_if = "Option::is_none")] pub username: Option, } @@ -354,26 +359,18 @@ impl ConfigTarget { let mut strm_cnt = 0; let mut xtream_cnt = 0; let mut strm_needs_xtream = false; - for format in &self.output { - let has_username = if let Some(username) = &format.username { !username.trim().is_empty() } else { false }; - let has_filename = if let Some(fname) = &format.filename { !fname.trim().is_empty() } else { false }; + let mut hdhr_cnt = 0; + for target_output in &self.output { + let has_username = if let Some(username) = &target_output.username { !username.trim().is_empty() } else { false }; + let has_filename = if let Some(fname) = &target_output.output { !fname.trim().is_empty() } else { false }; - match format.target { + match target_output.target { TargetType::M3u => { m3u_cnt += 1; if has_username { warn!("Username for target output m3u is ignored: {}", self.name); } } - TargetType::Strm => { - strm_cnt += 1; - if !has_filename { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "filename is required for strm type: {}", self.name); - } - if has_username { - strm_needs_xtream = true; - } - } TargetType::Xtream => { xtream_cnt += 1; if default_as_default().eq_ignore_ascii_case(&self.name) { @@ -386,10 +383,29 @@ impl ConfigTarget { warn!("Filename for target output xtream is ignored: {}", self.name); } } + TargetType::Strm => { + strm_cnt += 1; + if !has_filename { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "filename is required for strm type: {}", self.name); + } + if has_username { + strm_needs_xtream = true; + } + } + TargetType::HdHomeRun => { + hdhr_cnt += 1; + if !has_username { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Username is required for HdHomeRun type: {}", self.name); + } + let hdhr_name = target_output.output.as_ref().map_or("", |s| s.trim()); + if hdhr_name.is_empty() { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Output is required for HdHomeRun type: {}", self.name); + } + } } } - if m3u_cnt > 1 || strm_cnt > 1 || xtream_cnt > 1 { + if m3u_cnt > 1 || strm_cnt > 1 || xtream_cnt > 1 || hdhr_cnt > 1 { return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Multiple output formats with same type : {}", self.name); } @@ -397,6 +413,10 @@ impl ConfigTarget { return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "strm output with a username is only permitted when used in combination with xtream output: {}", self.name); } + if hdhr_cnt > 0 && (xtream_cnt == 0 && m3u_cnt == 0) { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "HdHomeRun output is only permitted when used in combination with xtream or m3u output: {}", self.name); + } + if let Some(watch) = &self.watch { let regexps: Result, _> = watch.iter().map(|s| regex::Regex::new(s)).collect(); match regexps { @@ -434,7 +454,7 @@ impl ConfigTarget { pub fn get_m3u_filename(&self) -> Option<&String> { for format in &self.output { if format.target == TargetType::M3u { - return format.filename.as_ref(); + return format.output.as_ref(); } } None @@ -781,7 +801,6 @@ pub struct VideoConfig { } impl VideoConfig { - /// # Panics /// /// Will panic if default `RegEx` gets invalid @@ -1074,6 +1093,8 @@ pub struct Config { pub messaging: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub reverse_proxy: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub hdhomerun: Option, #[serde(skip)] pub t_api_proxy: Arc>>, #[serde(skip)] @@ -1091,8 +1112,63 @@ pub struct Config { } impl Config { - pub fn set_api_proxy(&mut self, api_proxy: Option) { + pub fn set_api_proxy(&mut self, api_proxy: Option) -> Result<(), M3uFilterError> { self.t_api_proxy = Arc::new(RwLock::new(api_proxy)); + self.check_target_user() + } + + fn check_username(&self, output_username: Option<&str>, target_name: &str) -> Result<(), M3uFilterError> { + if let Some(username) = output_username { + if let Some((_, config_target)) = self.get_target_for_username(username) { + if config_target.name != target_name { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "User:{username} does not belong to target: {}", target_name); + } + } + Ok(()) + } else { + Ok(()) + } + } + fn check_target_user(&mut self) -> Result<(), M3uFilterError> { + let check_homerun = self.hdhomerun.as_ref().is_some_and(|h| h.enabled); + for source in &self.sources { + for target in &source.targets { + for output in &target.output { + match output.target { + TargetType::M3u | TargetType::Xtream => {} + TargetType::Strm => { + self.check_username(output.username.as_deref(), &target.name)?; + } + TargetType::HdHomeRun => { + if check_homerun { + if let Some(hdhr_name) = output.output.as_ref() { + self.check_username(output.username.as_deref(), &target.name)?; + if let Some(homerun) = &mut self.hdhomerun { + for device in &mut homerun.devices { + if &device.name == hdhr_name { + if let Some(username) = output.username.as_ref() { + device.t_username.clone_from(username); + device.t_enabled = true; + } + } + } + } + } + } + } + } + } + } + } + + if let Some(hdhomerun) = &self.hdhomerun { + for device in &hdhomerun.devices { + if !device.t_enabled { + debug!("HdHomeRun device '{}' has no username and will be disabled", device.name); + } + } + } + Ok(()) } pub fn is_reverse_proxy_resource_rewrite_enabled(&self) -> bool { @@ -1125,8 +1201,8 @@ impl Config { } pub fn get_target_for_username(&self, username: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { - if let Some(credentials) = self.get_user_credentials(username) { - return self.t_api_proxy.read().as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name(&credentials.username, &credentials.password))) + if let Some(credentials) = self.get_user_credentials(username) { + return self.t_api_proxy.read().as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name(&credentials.username, &credentials.password))); } None } @@ -1256,7 +1332,7 @@ impl Config { match file_utils::read_file_as_bytes(&PathBuf::from(&channel_unavailable)) { Ok(data) => { self.t_channel_unavailable_file = Some(Arc::new(data)); - }, + } Err(err) => { error!("Failed to load channel unavailable file: {channel_unavailable} {err}"); } @@ -1289,6 +1365,13 @@ impl Config { if let Some(reverse_proxy) = self.reverse_proxy.as_mut() { reverse_proxy.prepare(&self.working_dir, resolve_var); } + + if let Some(hdhomerun) = self.hdhomerun.as_mut() { + if hdhomerun.enabled { + hdhomerun.prepare(self.api.port)?; + } + } + self.api.prepare(); self.prepare_api_web_root(resolve_var); if let Some(templates) = &mut self.templates { diff --git a/src/model/hdhomerun_config.rs b/src/model/hdhomerun_config.rs new file mode 100644 index 000000000..2b49d312e --- /dev/null +++ b/src/model/hdhomerun_config.rs @@ -0,0 +1,95 @@ +use std::collections::HashSet; +use log::warn; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, create_m3u_filter_error_result}; +fn default_friendly_name() -> String { String::from("M3uFilterTV") } +fn default_manufacturer() -> String { String::from("Silicondust") } +fn default_model_name() -> String { String::from("HDTC-2US") } +fn default_firmware_name() -> String { String::from("hdhomeruntc_atsc") } +fn default_firmware_version() -> String { String::from("20170930") } +fn default_device_type() -> String { String::from("urn:schemas-upnp-org:device:MediaServer:1") } +fn default_device_udn() -> String { String::from("uuid:12345678-90ab-cdef-1234-567890abcdef::urn:dial-multicast:com.silicondust.hdhomerun") } + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] +#[serde(deny_unknown_fields)] +pub struct HdHomeRunDeviceConfig { + #[serde(default = "default_friendly_name")] + pub friendly_name: String, + #[serde(default = "default_manufacturer")] + pub manufacturer: String, + #[serde(default = "default_model_name")] + pub model_name: String, + #[serde(default = "default_model_name")] + pub model_number: String, + #[serde(default = "default_firmware_name")] + pub firmware_name: String, + #[serde(default = "default_firmware_version")] + pub firmware_version: String, + // pub device_auth: String, + #[serde(default = "default_device_type")] + pub device_type: String, + #[serde(default = "default_device_udn")] + pub device_udn: String, + pub name: String, + #[serde(default)] + pub port: u16, + #[serde(skip)] + pub t_username: String, + #[serde(skip)] + pub t_enabled: bool, +} + +impl HdHomeRunDeviceConfig { + pub fn prepare(&mut self, tuner: u8) { + self.name = self.name.trim().to_string(); + if self.name.is_empty() { + self.name = format!("tuner{tuner}"); + warn!("Tuner name empty, assigned new name: {}", self.name); + } + + if tuner > 0 && self.friendly_name == default_friendly_name() { + self.friendly_name = format!("{} {}", self.friendly_name, tuner); + } + if self.device_udn == default_device_udn() { + self.device_udn = format!("{}:{}", self.device_udn, tuner+1); + } + } +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] +#[serde(deny_unknown_fields)] +pub struct HdHomeRunConfig { + #[serde(default)] + pub enabled: bool, + pub devices: Vec, +} + +impl HdHomeRunConfig { + pub fn prepare(&mut self, api_port: u16) -> Result<(), M3uFilterError> { + let mut names = HashSet::new(); + let mut ports = HashSet::new(); + ports.insert(api_port); + for (tuner, device) in (0_u8..).zip(self.devices.iter_mut()) { + device.prepare(tuner); + if names.contains(&device.name) { + names.insert(&device.name); + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "HdHomeRun duplicate device name {}", device.name); + } + if device.port > 0 && ports.contains(&device.port) { + ports.insert(device.port); + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "HdHomeRun duplicate port {}", device.port); + } + } + let mut current_port = api_port + 1; + for device in &mut self.devices { + if device.port == 0 { + while ports.contains(¤t_port) { + current_port += 1; + } + device.port = current_port; + current_port += 1; + } + } + + Ok(()) + } +} \ No newline at end of file diff --git a/src/model/mod.rs b/src/model/mod.rs index 493b74c8b..2ce38cc60 100644 --- a/src/model/mod.rs +++ b/src/model/mod.rs @@ -6,4 +6,5 @@ pub mod stats; pub mod xmltv; pub mod xtream; pub mod healthcheck; -pub mod playlist_categories; \ No newline at end of file +pub mod playlist_categories; +pub mod hdhomerun_config; \ No newline at end of file diff --git a/src/model/playlist.rs b/src/model/playlist.rs index e521e61a5..6f7f041eb 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -284,13 +284,15 @@ pub struct M3uPlaylistItem { pub epg_channel_id: Option>, pub input_name: Rc, pub item_type: PlaylistItemType, + #[serde(skip)] + pub t_stream_url: String, + #[serde(skip)] + pub t_resource_url: Option, } impl M3uPlaylistItem { - pub fn to_m3u(&self, target_options: Option<&ConfigTargetOptions>, rewrite_urls: Option<&(String, Option)>) -> String { - let (stream_url, resource_url) = rewrite_urls - .map_or_else(|| (self.url.as_str(), None), |(su, ru)| (su.as_str(), ru.as_ref().map(String::as_str))); - + #[allow(clippy::missing_panics_doc)] + pub fn to_m3u(&self, target_options: Option<&ConfigTargetOptions>, rewrite_urls: bool) -> String { let options = target_options.as_ref(); let ignore_logo = options.is_some_and(|o| o.ignore_logo); let mut line = format!("#EXTINF:-1 tvg-id=\"{}\" tvg-name=\"{}\" group-title=\"{}\"", @@ -298,13 +300,10 @@ impl M3uPlaylistItem { self.name, self.group); if !ignore_logo { - match resource_url { - None => { - to_m3u_non_empty_fields!(self, line, (logo, "tvg-logo"), (logo_small, "tvg-logo-small");); - } - Some(res_url) => { - to_m3u_resource_non_empty_fields!(self, res_url, line, (logo, "tvg-logo"), (logo_small, "tvg-logo-small");); - } + if rewrite_urls && self.t_resource_url.is_some(){ + to_m3u_resource_non_empty_fields!(self, self.t_resource_url.as_ref().unwrap(), line, (logo, "tvg-logo"), (logo_small, "tvg-logo-small");); + } else { + to_m3u_non_empty_fields!(self, line, (logo, "tvg-logo"), (logo_small, "tvg-logo-small");); } } @@ -315,7 +314,8 @@ impl M3uPlaylistItem { (time_shift, "timeshift"), (rec, "tvg-rec");); - format!("{},{}\n{}", line, self.title, stream_url) + let url = if self.t_stream_url.is_empty() { &self.url } else { &self.t_stream_url }; + format!("{line},{}\n{url}", self.title, ) } } @@ -505,6 +505,8 @@ impl PlaylistItem { epg_channel_id: header.epg_channel_id.clone(), input_name: Rc::clone(&header.input_name), item_type: header.item_type, + t_stream_url: header.url.to_string(), + t_resource_url: None, } } diff --git a/src/processing/parser/xtream.rs b/src/processing/parser/xtream.rs index 6087d36e0..7db65e666 100644 --- a/src/processing/parser/xtream.rs +++ b/src/processing/parser/xtream.rs @@ -78,23 +78,33 @@ pub fn parse_xtream_series_info(info: &Value, group_title: &str, series_name: &s } } +#[allow(clippy::too_many_arguments)] +pub fn get_xtream_url(xtream_cluster: XtreamCluster, url: &str, + username: &str, password: &str, + stream_id: u32, container_extension: Option<&String>, + live_stream_use_prefix: bool, live_stream_without_extension: bool) -> String { + let stream_base_url = match xtream_cluster { + XtreamCluster::Live => { + let ctx_path = if live_stream_use_prefix { "live/" } else { "" }; + let suffix = if live_stream_without_extension { "" } else { ".ts" }; + format!("{url}/{ctx_path}{username}/{password}/{stream_id}{suffix}") + } + XtreamCluster::Video => { + let ext = container_extension.as_ref().map_or("mp4", |e| e.as_str()); + format!("{url}/movie/{username}/{password}/{stream_id}.{ext}") + } + XtreamCluster::Series => + format!("{}&action={ACTION_GET_SERIES_INFO}&series_id={stream_id}", get_xtream_stream_url_base(url, username, password)) + }; + stream_base_url +} + pub fn create_xtream_url(xtream_cluster: XtreamCluster, url: &str, username: &str, password: &str, - stream: &XtreamStream, live_stream_use_prefix: bool, live_stream_without_extension: bool) -> Rc { + stream: &XtreamStream, live_stream_use_prefix: bool, live_stream_without_extension: bool) -> Rc { if stream.direct_source.is_empty() { - let stream_base_url = match xtream_cluster { - XtreamCluster::Live => { - let ctx_path = if live_stream_use_prefix {"live/"} else { "" }; - let suffix = if live_stream_without_extension { "" } else { ".ts" }; - format!("{url}/{ctx_path}{username}/{password}/{}{suffix}", &stream.get_stream_id()) - }, - XtreamCluster::Video => { - let ext = stream.container_extension.as_ref().map_or("mp4", |e| e.as_str()); - format!("{url}/movie/{username}/{password}/{}.{ext}", &stream.get_stream_id()) - } - XtreamCluster::Series => - format!("{}&action={ACTION_GET_SERIES_INFO}&series_id={}", get_xtream_stream_url_base(url, username, password), &stream.get_stream_id()) - }; - Rc::new(stream_base_url) + Rc::new(get_xtream_url(xtream_cluster, url, username, password, stream.get_stream_id(), + stream.container_extension.as_ref().map(std::string::ToString::to_string).as_ref(), + live_stream_use_prefix, live_stream_without_extension)) } else { Rc::clone(&stream.direct_source) } @@ -135,7 +145,7 @@ pub fn parse_xtream(input: &ConfigInput, let item = PlaylistItem { header: RefCell::new(PlaylistItemHeader { id: Rc::new(stream.get_stream_id().to_string()), - uuid: Rc::new(generate_playlist_uuid(&input_name, &stream.get_stream_id().to_string(), item_type, &stream_url)), + uuid: Rc::new(generate_playlist_uuid(&input_name, &stream.get_stream_id().to_string(), item_type, &stream_url)), name: Rc::clone(&stream.name), logo: Rc::clone(&stream.stream_icon), group: Rc::clone(category_name), diff --git a/src/repository/epg_repository.rs b/src/repository/epg_repository.rs index b591c4dd5..727ba0177 100644 --- a/src/repository/epg_repository.rs +++ b/src/repository/epg_repository.rs @@ -54,8 +54,8 @@ pub fn epg_write(target: &ConfigTarget, cfg: &Config, target_path: &Path, epg: O None => return Err(notify_err!(format!("failed to serialize epg for target: {}, storage path not found", target.name))), } } - TargetType::Strm => {} - } + TargetType::Strm | TargetType::HdHomeRun => {} + } } Ok(()) } diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 0157df074..67d80d2f8 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -607,7 +607,7 @@ pub async fn kodi_write_strm_playlist( if new_playlist.is_empty() { return Ok(()); } - if output.filename.is_none() { + if output.output.is_none() { return Err(notify_err!( "Output directory missing. Write strm playlist failed".to_string() )); @@ -615,11 +615,11 @@ pub async fn kodi_write_strm_playlist( let Some(root_path) = file_utils::get_file_path( &cfg.working_dir, - Some(std::path::PathBuf::from(&output.filename.as_ref().map_or_else(|| "/tmp", |v| v.as_ref()))), + Some(std::path::PathBuf::from(&output.output.as_ref().map_or_else(|| "/tmp", |v| v.as_ref()))), ) else { return Err(info_err!(format!( "Failed to get file path for {}", - output.filename.as_deref().unwrap_or("") + output.output.as_deref().unwrap_or("") ))); }; diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index a8637d2a5..4edd0dc73 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -20,7 +20,6 @@ pub struct M3uPlaylistIterator { target_options: Option, mask_redirect_url: bool, include_type_in_url: bool, - started: bool, rewrite_resource: bool, proxy_type: ProxyType, _file_lock: FileReadGuard, @@ -45,6 +44,8 @@ impl M3uPlaylistIterator { let include_type_in_url = target_options.is_some_and(|opts| opts.m3u_include_type_in_url); let mask_redirect_url = target_options.is_some_and(|opts| opts.m3u_mask_redirect_url); + // TODO m3u bouquet filter + let server_info = cfg.get_user_server_info(user); Ok(Self { reader, @@ -56,7 +57,6 @@ impl M3uPlaylistIterator { mask_redirect_url, proxy_type: user.proxy.clone(), _file_lock: file_lock, // Save lock inside struct - started: false, rewrite_resource: cfg.is_reverse_proxy_resource_rewrite_enabled(), }) } @@ -92,19 +92,9 @@ impl M3uPlaylistIterator { self.get_rewritten_url(m3u_pli, false, M3U_RESOURCE_PATH) } -} - -impl Iterator for M3uPlaylistIterator { - type Item = String; - - fn next(&mut self) -> Option { - if !self.started { - self.started = true; - return Some("#EXTM3U".to_string()); - } - + fn get_next(&mut self) -> Option<(M3uPlaylistItem, bool)> { // TODO hls and unknown reverse proxy - self.reader.next().map(|(m3u_pli, _has_next)| { + self.reader.next().map(|(mut m3u_pli, _has_next)| { let rewrite_urls = match m3u_pli.item_type { PlaylistItemType::LiveHls => None, _ => if match &self.proxy_type { @@ -116,8 +106,58 @@ impl Iterator for M3uPlaylistIterator { None } }; - let target_options = self.target_options.as_ref(); - m3u_pli.to_m3u(target_options, rewrite_urls.as_ref()) + let url = m3u_pli.url.to_string(); + let (stream_url, resource_url) = rewrite_urls + .map_or_else(|| (url, None), |(su, ru)| (su, ru.as_ref().map(String::to_string))); + + m3u_pli.t_stream_url = stream_url.to_string(); + m3u_pli.t_resource_url = resource_url.map(|s| s.to_string()); + (m3u_pli, self.reader.has_next()) + }) + } +} + +impl Iterator for M3uPlaylistIterator { + type Item = (M3uPlaylistItem, bool); + + fn next(&mut self) -> Option { + self.get_next() + } +} + + +pub struct M3uPlaylistIteratorText { + inner: M3uPlaylistIterator, + started: bool, + +} + +impl M3uPlaylistIteratorText { + pub fn new( + cfg: &Config, + target: &ConfigTarget, + user: &ProxyUserCredentials, + ) -> Result { + Ok(Self { + inner: M3uPlaylistIterator::new(cfg, target, user)?, + started: false, + }) + } +} + +impl Iterator for M3uPlaylistIteratorText { + type Item = String; + + fn next(&mut self) -> Option { + if !self.started { + self.started = true; + return Some("#EXTM3U".to_string()); + } + + // TODO hls and unknown reverse proxy + self.inner.get_next().map(|(m3u_pli, _has_next)| { + let target_options = self.inner.target_options.as_ref(); + m3u_pli.to_m3u(target_options, true) }) } } diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index d550a4443..475d09f72 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -9,7 +9,7 @@ use crate::model::api_proxy::{ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, PlaylistItemType}; use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDocumentWriter}; -use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; +use crate::repository::m3u_playlist_iterator::{M3uPlaylistIteratorText}; use crate::repository::storage::{get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::utils::file::file_utils; use crate::utils::file::file_utils::file_writer; @@ -40,7 +40,7 @@ fn persist_m3u_playlist_as_text(target: &ConfigTarget, cfg: &Config, m3u_playlis let mut buf_writer = file_writer(&file); let _ = buf_writer.write(b"#EXTM3U\n"); for m3u in m3u_playlist { - let _ = buf_writer.write(m3u.to_m3u(target.options.as_ref(), None).as_bytes()); + let _ = buf_writer.write(m3u.to_m3u(target.options.as_ref(), false).to_string().as_bytes()); let _ = buf_writer.write(b"\n"); } } @@ -84,8 +84,8 @@ pub async fn m3u_load_rewrite_playlist( cfg: &Config, target: &ConfigTarget, user: &ProxyUserCredentials, -) -> Result>, M3uFilterError> { - Ok(Box::new(M3uPlaylistIterator::new(cfg, target, user)?)) +) -> Result { + M3uPlaylistIteratorText::new(cfg, target, user) } diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index ab170d256..caa15baa5 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -47,6 +47,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, TargetType::M3u => m3u_write_playlist(target, cfg, &target_path, playlist).await, TargetType::Xtream => xtream_write_playlist(target, cfg, playlist).await, TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, output).await, + TargetType::HdHomeRun => Ok(()), }; if let Err(err) = result { errors.push(err); diff --git a/src/repository/xtream_playlist_iterator.rs b/src/repository/xtream_playlist_iterator.rs index 507ec1ea6..2f81389a4 100644 --- a/src/repository/xtream_playlist_iterator.rs +++ b/src/repository/xtream_playlist_iterator.rs @@ -26,7 +26,7 @@ impl XtreamPlaylistIterator { config: &Config, target: &ConfigTarget, category_id: u32, - user: &ProxyUserCredentials + user: &ProxyUserCredentials, ) -> Result { if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { let (xtream_path, idx_path) = xtream_get_file_paths(&storage_path, cluster); @@ -56,24 +56,51 @@ impl XtreamPlaylistIterator { Err(info_err!(format!("Failed to find xtream storage for target {}", &target.name))) } } -} -impl Iterator for XtreamPlaylistIterator { - type Item = String; - - fn next(&mut self) -> Option { + fn get_next(&mut self) -> Option<(XtreamPlaylistItem, bool)> { if self.reader.has_error() { error!("Could not deserialize xtream item: {:?}", self.reader.get_path()); return None; } if let Some(set) = &self.filter { - self.reader - .find(|(pli, _has_next)| set.contains(&pli.category_id.to_string())) - .map(|(pli, _has_next)| pli.to_doc(&self.base_url, &self.options, &self.user).to_string()) + self.reader.find(|(pli, _has_next)| set.contains(&pli.category_id.to_string())) } else { - self.reader - .next() - .map(|(pli, _has_next)| pli.to_doc(&self.base_url, &self.options, &self.user).to_string()) + self.reader.next() } } -} \ No newline at end of file + +} + +impl Iterator for XtreamPlaylistIterator { + type Item = (XtreamPlaylistItem, bool); + fn next(&mut self) -> Option { + self.get_next() + } +} + + +pub struct XtreamPlaylistIteratorText { + inner: XtreamPlaylistIterator, +} + +impl XtreamPlaylistIteratorText { +pub async fn new( + cluster: XtreamCluster, + config: &Config, + target: &ConfigTarget, + category_id: u32, + user: &ProxyUserCredentials, + ) -> Result { + Ok(Self { + inner: XtreamPlaylistIterator::new(cluster, config, target, category_id, user).await? + }) + } +} + +impl Iterator for XtreamPlaylistIteratorText { + type Item = (String, bool); + fn next(&mut self) -> Option { + self.inner.get_next().map(|(pli, has_next)| (pli.to_doc(&self.inner.base_url, &self.inner.options, &self.inner.user).to_string(), has_next)) + } +} + diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index e59ebeb12..207692ef9 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -6,6 +6,9 @@ use std::fs; use std::fs::File; use std::io::{BufReader, Error, ErrorKind, Read}; use std::path::{Path, PathBuf}; +use std::sync::Arc; +use bytes::Bytes; +use futures::{stream, Stream, StreamExt}; use crate::repository::storage::hex_encode; use crate::utils::file::file_utils::file_reader; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, str_to_io_error, info_err, create_m3u_filter_error, create_m3u_filter_error_result, notify_err}; @@ -18,11 +21,11 @@ use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDo use crate::repository::playlist_repository::get_target_id_mapping; use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::repository::target_id_mapping::{VirtualIdRecord}; -use crate::repository::xtream_playlist_iterator::XtreamPlaylistIterator; use crate::utils::file::file_utils::open_readonly_file; use crate::utils::hash_utils::generate_playlist_uuid; use crate::utils::json_utils::{get_u32_from_serde_value, json_iter_array, json_write_documents_to_file}; use crate::utils::bincode_utils::{bincode_deserialize}; +use crate::repository::xtream_playlist_iterator::XtreamPlaylistIteratorText; pub static COL_CAT_LIVE: &str = "cat_live"; pub static COL_CAT_SERIES: &str = "cat_series"; @@ -429,8 +432,8 @@ pub async fn xtream_load_rewrite_playlist( target: &ConfigTarget, category_id: u32, user: &ProxyUserCredentials, -) -> Result>, M3uFilterError> { - Ok(Box::new(XtreamPlaylistIterator::new(cluster, config, target, category_id, user).await?)) +) -> Result { + XtreamPlaylistIteratorText::new(cluster, config, target, category_id, user).await } pub fn xtream_write_series_info( @@ -939,7 +942,7 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( } } -pub fn iter_raw_xtream_playlist(config: &Config, target: &ConfigTarget, cluster: XtreamCluster) -> Option<(FileReadGuard, impl Iterator)> { +pub fn iter_raw_xtream_playlist(config: &Arc, target: &ConfigTarget, cluster: XtreamCluster) -> Option<(FileReadGuard, impl Iterator)> { if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { let (xtream_path, idx_path) = xtream_get_file_paths(&storage_path, cluster); if !xtream_path.exists() || !idx_path.exists() { @@ -954,4 +957,31 @@ pub fn iter_raw_xtream_playlist(config: &Config, target: &ConfigTarget, cluster: } else { None } +} + +pub fn playlist_iter_to_stream(channels: Option<(FileReadGuard, I)>) -> impl Stream> +where + I: Iterator + 'static, +{ + match channels { + Some((_, chans)) => { + // Convert iterator items to Result + let mapped = chans.map(move |(item, has_next)| { + match serde_json::to_string(&item) { + Ok(content) => { + Ok(Bytes::from(if has_next { + format!("{content},") + } else { + content + })) + } + Err(_) => Ok(Bytes::from("")), + } + }); + stream::iter(mapped).left_stream() + } + None => { + stream::once(async { Ok(Bytes::from("")) }).right_stream() + } + } } \ No newline at end of file diff --git a/src/utils/file/config_reader.rs b/src/utils/file/config_reader.rs index ed249ba4e..a7348863e 100644 --- a/src/utils/file/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -33,18 +33,18 @@ pub fn read_mappings(args_mapping: Option, cfg: &mut Config) -> Result, cfg: &mut Config) -> Option { +pub fn read_api_proxy_config(args_api_proxy_config: Option, cfg: &mut Config) -> Result, M3uFilterError> { let api_proxy_config_file: String = args_api_proxy_config.unwrap_or_else(|| file_utils::get_default_api_proxy_config_path(cfg.t_config_path.as_str())); api_proxy_config_file.clone_into(&mut cfg.t_api_proxy_file_path); let api_proxy_config = read_api_proxy(cfg, api_proxy_config_file.as_str(), true); match api_proxy_config { None => { warn!("cant read api_proxy_config file: {}", api_proxy_config_file.as_str()); - None + Ok(None) } Some(config) => { - cfg.set_api_proxy(Some(config)); - Some(api_proxy_config_file) + cfg.set_api_proxy(Some(config))?; + Ok(Some(api_proxy_config_file)) } } } diff --git a/src/utils/network/xtream.rs b/src/utils/network/xtream.rs index 6c0fa07c4..40202da56 100644 --- a/src/utils/network/xtream.rs +++ b/src/utils/network/xtream.rs @@ -204,9 +204,7 @@ pub fn create_vod_info_from_item(user: &ProxyUserCredentials, pli: &XtreamPlayli "movie_data": {{ "added": "{added}", "category_id": {category_id}, - "category_ids": [ - {category_id} - ], + "category_ids": [{category_id}], "container_extension": "{extension}", "custom_sid": "", "direct_source": "",