From e35efb4949df897d47a5dc1ec9f433d7f57b775f Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 13 Mar 2024 12:05:51 +0100 Subject: [PATCH 01/12] refactored, moved files into util package --- src/{ => util}/config_reader.rs | 2 +- src/util/mod.rs | 3 +++ src/{ => util}/multi_file_reader.rs | 0 src/{ => util}/utils.rs | 0 4 files changed, 4 insertions(+), 1 deletion(-) rename src/{ => util}/config_reader.rs (99%) create mode 100644 src/util/mod.rs rename src/{ => util}/multi_file_reader.rs (100%) rename src/{ => util}/utils.rs (100%) diff --git a/src/config_reader.rs b/src/util/config_reader.rs similarity index 99% rename from src/config_reader.rs rename to src/util/config_reader.rs index a3a5e084d..f6555f3b4 100644 --- a/src/config_reader.rs +++ b/src/util/config_reader.rs @@ -8,7 +8,7 @@ use crate::model::config::{Config, ConfigDto}; use crate::model::mapping::Mappings; use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, utils}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::multi_file_reader::MultiFileReader; +use crate::util::multi_file_reader::MultiFileReader; pub(crate) fn read_mappings(args_mapping: Option, cfg: &mut Config) -> Result<(), M3uFilterError> { let mappings_file: String = args_mapping.unwrap_or(utils::get_default_mappings_path(cfg._config_path.as_str())); diff --git a/src/util/mod.rs b/src/util/mod.rs new file mode 100644 index 000000000..b72adc230 --- /dev/null +++ b/src/util/mod.rs @@ -0,0 +1,3 @@ +pub(crate) mod utils; +pub(crate) mod config_reader; +pub(crate) mod multi_file_reader; \ No newline at end of file diff --git a/src/multi_file_reader.rs b/src/util/multi_file_reader.rs similarity index 100% rename from src/multi_file_reader.rs rename to src/util/multi_file_reader.rs diff --git a/src/utils.rs b/src/util/utils.rs similarity index 100% rename from src/utils.rs rename to src/util/utils.rs From 22402877f27e7a04cbc97c4e8da424da8e3f35ef Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 13 Mar 2024 12:06:30 +0100 Subject: [PATCH 02/12] Case insenstive filter syntax --- src/filter.pest | 8 ++++---- src/filter.rs | 21 +++++++++++++++++++-- 2 files changed, 23 insertions(+), 6 deletions(-) diff --git a/src/filter.pest b/src/filter.pest index 704dc0854..9936f8587 100644 --- a/src/filter.pest +++ b/src/filter.pest @@ -1,8 +1,8 @@ WHITESPACE = _{ " " | "\t" } -field = { "Group" | "Title" | "Name" | "Url" } -and = {"AND" | "and"} -or = {"OR" | "or"} -not = { "NOT" | "not" } +field = { ^"group" | ^"title" | ^"name" | ^"url" } +and = { ^"and" } +or = { ^"or" } +not = { ^"not" } regexp = @{ "\"" ~ ( "\\\"" | (!"\"" ~ ANY) )* ~ "\"" } comparison_value = _{ regexp } comparison = { field ~ "~" ~ comparison_value } diff --git a/src/filter.rs b/src/filter.rs index aa00576c1..22f60453c 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -70,10 +70,27 @@ pub(crate) struct RegexWithCaptures { pub captures: Vec, } - #[derive(Parser)] //#[grammar = "filter.pest"] -#[grammar_inline = "WHITESPACE = _{ \" \" | \"\\t\" }\nfield = { \"Group\" | \"Title\" | \"Name\" | \"Url\" }\nand = {\"AND\" | \"and\"}\nor = {\"OR\" | \"or\"}\nnot = { \"NOT\" | \"not\" }\nregexp = @{ \"\\\"\" ~ ( \"\\\\\\\"\" | (!\"\\\"\" ~ ANY) )* ~ \"\\\"\" }\ncomparison_value = _{ regexp }\ncomparison = { field ~ \"~\" ~ comparison_value }\nbool_op = { and | or}\nexpr_group = { \"(\" ~ expr ~ \")\" }\nexpr = {comparison ~ (bool_op ~ expr)* | expr_group ~ (bool_op ~ expr)* | not ~ expr ~ (bool_op ~ expr)* }\nstmt = { expr ~ (bool_op ~ expr)* }\nmain = _{ SOI ~ stmt ~ EOI }"] +#[grammar_inline = r#" +WHITESPACE = _{ " " | "\t" } +field = { ^"group" | ^"title" | ^"name" | ^"url" } +and = { ^"and" } +or = { ^"or" } +not = { ^"not" } +regexp = @{ "\"" ~ ( "\\\"" | (!"\"" ~ ANY) )* ~ "\"" } +comparison_value = _{ regexp } +comparison = { field ~ "~" ~ comparison_value } +bool_op = { and | or} +expr_group = { "(" ~ expr ~ ")" } +expr = { + comparison ~ (bool_op ~ expr)* + | expr_group ~ (bool_op ~ expr)* + | not ~ expr ~ (bool_op ~ expr)* +} +stmt = { expr ~ (bool_op ~ expr)* } +main = _{ SOI ~ stmt ~ EOI } +"#] struct FilterParser; #[derive(Debug, Clone)] From 60b9551371d36341a6175aa226c62365dc86ad7e Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 15 Mar 2024 18:34:39 +0100 Subject: [PATCH 03/12] Multi Xtream Input handling --- src/api/v1_api.rs | 2 +- src/api/xtream_api.rs | 3 + src/main.rs | 7 +- src/model/config.rs | 2 + src/model/model_playlist.rs | 7 +- src/model/model_xtream.rs | 8 +- src/processing/m3u_parser.rs | 2 + src/processing/playlist_processor.rs | 108 +++++++++++++++++++-------- src/processing/xtream_parser.rs | 2 + src/repository/xtream_repository.rs | 90 ++++++++++++++++------ src/test.rs | 22 +++++- 11 files changed, 188 insertions(+), 65 deletions(-) diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index 70085cdf8..51af506be 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -6,10 +6,10 @@ use crate::download::{get_m3u_playlist, get_xtream_playlist}; use crate::model::config::{Config, ConfigDto, ConfigInput, ConfigInputOptions, ConfigSource, ConfigTarget, InputType, validate_targets}; use log::{error}; use crate::api::download_api::{download_file_info, queue_download_file}; -use crate::config_reader::{read_config, save_api_proxy, save_main_config}; use crate::m3u_filter_error::M3uFilterError; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, TargetUser}; use crate::processing::playlist_processor::exec_processing; +use crate::util::config_reader::{read_config, save_api_proxy, save_main_config}; fn _save_config_api_proxy(backup_dir: &str, api_proxy: &mut ApiProxyConfig) -> Option { match save_api_proxy(api_proxy._file_path.as_str(), backup_dir, api_proxy) { diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 2269b19c8..ddbc586cd 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -107,6 +107,9 @@ async fn xtream_player_api_stream( if let Some((user, target)) = get_user_target_by_credentials(username, password, api_req, _app_state) { let target_name = &target.name; if target.has_output(&TargetType::Xtream) { + + !! todo id mapping when multi_xtream !! + if let Some(target_input) = match _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { None => _app_state.config.get_input_for_target(target_name, &InputType::M3u), Some(inp) => Some(inp) diff --git a/src/main.rs b/src/main.rs index cc19e2d82..3853d2211 100644 --- a/src/main.rs +++ b/src/main.rs @@ -10,22 +10,21 @@ use clap::Parser; use env_logger::Builder; use log::{error, info, LevelFilter}; -use crate::config_reader::{read_api_proxy_config, read_config, read_mappings}; use crate::model::config::{Config, ProcessTargets, validate_targets}; use crate::processing::playlist_processor::exec_processing; +use crate::util::config_reader::{read_api_proxy_config, read_config, read_mappings}; +use crate::util::utils; mod m3u_filter_error; -mod config_reader; mod model; mod filter; mod repository; mod download; -mod utils; mod messaging; mod test; mod api; mod processing; -mod multi_file_reader; +mod util; #[derive(Parser)] #[command(name = "m3u-filter")] diff --git a/src/model/config.rs b/src/model/config.rs index 2b699505d..32c7b1bc8 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -301,6 +301,8 @@ impl ConfigTarget { pub(crate) struct ConfigSource { pub inputs: Vec, pub targets: Vec, + #[serde(skip_serializing, skip_deserializing)] + pub _multi_xtream_input: bool, } impl ConfigSource { diff --git a/src/model/model_playlist.rs b/src/model/model_playlist.rs index 8a766afd0..c98ceeb4e 100644 --- a/src/model/model_playlist.rs +++ b/src/model/model_playlist.rs @@ -68,6 +68,7 @@ pub(crate) trait FieldAccessor { #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct PlaylistItemHeader { pub id: Rc, + pub stream_id: Rc, pub name: Rc, pub logo: Rc, pub logo_small: Rc, @@ -86,7 +87,7 @@ pub(crate) struct PlaylistItemHeader { #[serde(skip_serializing, skip_deserializing)] pub additional_properties: Option>, #[serde(default = "default_playlist_item_type", skip_serializing, skip_deserializing)] - pub item_type: PlaylistItemType, + pub item_type: PlaylistItemType, #[serde(default = "default_as_false", skip_serializing, skip_deserializing)] pub series_fetched: bool, // only used for series_info } @@ -118,12 +119,12 @@ macro_rules! get_fields { impl FieldAccessor for PlaylistItemHeader { fn get_field(&self, field: &str) -> Option> { - get_fields!(self, field, id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url;) + get_fields!(self, field, id, stream_id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url;) } fn set_field(&mut self, field: &str, value: &str) -> bool { let val = String::from(value); - update_fields!(self, field, id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url; val) + update_fields!(self, field, id, stream_id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url; val) } } diff --git a/src/model/model_xtream.rs b/src/model/model_xtream.rs index f0b6941b8..4375a2627 100644 --- a/src/model/model_xtream.rs +++ b/src/model/model_xtream.rs @@ -330,4 +330,10 @@ impl XtreamSeriesInfoEpisode { add_str_property_if_exists!(result, series_info.info.youtube_trailer, "youtube_trailer"); if result.is_empty() { None } else { Some(result) } } -} \ No newline at end of file +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub(crate) struct MultiXtreamMapping { + pub stream_id: u32, + pub input_id: u16, +} diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 6ab8f373e..6d71808ea 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -55,6 +55,7 @@ fn skip_digit(it: &mut std::str::Chars) -> Option { fn create_empty_playlistitem_header(content: &String, url: String) -> PlaylistItemHeader { PlaylistItemHeader { id: default_as_empty_rc_str(), + stream_id: default_as_empty_rc_str(), name: default_as_empty_rc_str(), logo: default_as_empty_rc_str(), logo_small: default_as_empty_rc_str(), @@ -119,6 +120,7 @@ fn process_header(video_suffixes: &Vec<&str>, content: &String, url: String) -> if plih.group.is_empty() { plih.group = Rc::new(String::from("Unknown")); } + plih.stream_id = Rc::clone(&plih.id); plih.epg_channel_id = Some(Rc::clone(&plih.id)); } diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 6f0373f8f..e41fbd50a 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -4,7 +4,7 @@ use std::cell::RefCell; use std::collections::{HashMap, HashSet}; use std::rc::Rc; use std::sync::{Arc, Mutex}; -use std::thread; +use std::{thread}; use actix_rt::System; use log::{debug, error, info}; @@ -15,17 +15,18 @@ use crate::download::{get_m3u_playlist, get_xmltv, get_xtream_playlist, get_xtre use crate::filter::{get_field_value, MockValueProcessor, set_field_value, ValueProvider}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::messaging::{MsgKind, send_message}; -use crate::model::config::{ConfigTarget, default_as_default, InputAffix, InputType, ProcessTargets}; +use crate::model::config::{ConfigSource, ConfigTarget, default_as_default, InputAffix, InputType, ProcessTargets}; use crate::model::mapping::{Mapping, MappingValueProcessor}; use crate::model::model_config::{AFFIX_FIELDS, ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType}; use crate::model::model_playlist::{FetchedPlaylist, FieldAccessor, PlaylistGroup, PlaylistItem, PlaylistItemHeader}; +use crate::model::model_xtream::MultiXtreamMapping; use crate::model::stats::{InputStats, PlaylistStats}; use crate::model::xmltv::{Epg}; use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; use crate::repository::epg_repository::write_epg; use crate::repository::m3u_repository::{write_m3u_playlist, write_strm_playlist}; -use crate::repository::xtream_repository::write_xtream_playlist; +use crate::repository::xtream_repository::{write_xtream_mapping, write_xtream_playlist}; fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option> { debug!("Filtering {} groups", playlist.len()); @@ -44,7 +45,7 @@ fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt id: pg.id, title: pg.title.clone(), channels, - xtream_cluster: pg.xtream_cluster.clone() + xtream_cluster: pg.xtream_cluster.clone(), }); } }); @@ -248,7 +249,7 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option id: grp_id, title: Rc::clone(title), channels: vec![channel.clone()], - xtream_cluster: cluster.clone() + xtream_cluster: cluster.clone(), }) } } @@ -336,7 +337,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

{} Err(mut err) => err.drain(..).for_each(|e| errors.push(e)) } @@ -404,6 +405,7 @@ fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe { } pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], + source: &ConfigSource, target: &ConfigTarget, cfg: &Config, stats: &mut HashMap, errors: &mut Vec) -> Result<(), Vec> { @@ -425,31 +427,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], new_fpl.playlist = v; } } - let (resolve_series, resolve_series_delay) = - if let Some(options) = &target.options { - (options.xtream_resolve_series && fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u), - options.xtream_resolve_series_delay) - } else { - (false, 0) - }; - if resolve_series { - let mut series_playlist = get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; - // original content saved into original list - for plg in &series_playlist { - fpl.update_playlist(plg); - } - // run processing pipe over new items - for f in &pipe { - let r = f(&mut series_playlist, target); - if let Some(v) = r { - series_playlist = v; - } - } - // assign new items to the new playlist - for plg in &series_playlist { - new_fpl.update_playlist(plg); - } - } + playlist_resolve_series(target, errors, &pipe, fpl, &mut new_fpl).await; // stats let input_stats = stats.get_mut(&new_fpl.input.id); @@ -464,6 +442,39 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], } apply_affixes(&mut new_fetched_playlists); + + if source._multi_xtream_input { + let mut stream_id_mappings: Vec = Vec::new(); + let mut counter: u32 = 0; + new_fetched_playlists.iter() + .flat_map(|pl| { + let input_id = &pl.input.id; + pl.playlist.iter().map(move |plg| (input_id, &plg.channels)) + }) + .flat_map(|(input_id, channels)| channels.iter().map(move |chan| (input_id, chan))) + .for_each(|(input_id, chan)| { + counter += 1; + let mut header = chan.header.borrow_mut(); + header.id = Rc::new(counter.to_string()); + let xtream_mapping = MultiXtreamMapping { + stream_id: header.stream_id.parse::().unwrap(), + input_id: *input_id, + }; + stream_id_mappings.push(xtream_mapping); + }); + + match write_xtream_mapping(&stream_id_mappings, cfg, &target.name) { + Ok(_) => { + debug!("wrote multi xtream input mapping for {}", &target.name); + } + Err(err) => { + return Err(vec![M3uFilterError::new( + M3uFilterErrorKind::Notify, + format!("Write multi xtream input mapping {} failed: {}", target.name, err))]); + } + } + } + let mut new_playlist = vec![]; let mut new_epg = vec![]; let mut tv_guides = vec![]; @@ -495,7 +506,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], let watch_re = target._watch_re.as_ref().unwrap(); new_playlist.iter().for_each(|pl| { if watch_re.iter().any(|r| r.is_match(&pl.title)) { - process_group_watch(cfg, &target.name, pl) + process_group_watch(cfg, &target.name, pl) } }); } @@ -508,6 +519,37 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], } } +async fn playlist_resolve_series<'a>(target: &ConfigTarget, errors: &mut Vec, + pipe: &ProcessingPipe, + fpl: &mut FetchedPlaylist<'_>, + new_fpl: &mut FetchedPlaylist<'_>) { + let (resolve_series, resolve_series_delay) = + if let Some(options) = &target.options { + (options.xtream_resolve_series && fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u), + options.xtream_resolve_series_delay) + } else { + (false, 0) + }; + if resolve_series { + let mut series_playlist = get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; + // original content saved into original list + for plg in &series_playlist { + fpl.update_playlist(plg); + } + // run processing pipe over new items + for f in pipe { + let r = f(&mut series_playlist, target); + if let Some(v) = r { + series_playlist = v; + } + } + // assign new items to the new playlist + for plg in &series_playlist { + new_fpl.update_playlist(plg); + } + } +} + fn persist_playlist(playlist: &[PlaylistGroup], epg: Option, target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { let mut errors = vec![]; @@ -543,7 +585,7 @@ pub(crate) async fn exec_processing(cfg: Arc, targets: Arc; @@ -105,7 +106,7 @@ fn write_xtream_info(app_state: &AppState, target_name: &str, stream_id: i32, cl write_index(&idx_path, index_tree)?; } Err(err) => { - return Err(err) + return Err(err); } } } @@ -330,12 +331,12 @@ fn write_index(path: &PathBuf, index: &IndexTree) -> std::io::Result<()> { fn seek_read( reader: &mut (impl Read + Seek), - offset: u32, + offset: u64, amount_to_read: u16, ) -> Result, Error> { // A buffer filled with as many zeros as we'll read with read_exact - let mut buf = vec![0; amount_to_read as usize]; - reader.seek(SeekFrom::Start(offset as u64))?; + let mut buf = vec![0u8; amount_to_read as usize]; + reader.seek(SeekFrom::Start(offset))?; reader.read_exact(&mut buf)?; Ok(buf) } @@ -356,7 +357,7 @@ pub(crate) async fn xtream_get_stored_stream_info( if let Some(idx_map) = &index_tree { if let Some((offset, size)) = idx_map.get(&stream_id) { let mut reader = BufReader::new(File::open(&col_path).unwrap()); - if let Ok(bytes) = seek_read(&mut reader, *offset, *size) { + if let Ok(bytes) = seek_read(&mut reader, *offset as u64, *size) { let mut decomp: Vec = Vec::new(); let _ = lzma_rs::lzma_decompress(&mut bytes.as_slice(), &mut decomp); drop(shared_lock); @@ -378,22 +379,67 @@ pub(crate) async fn xtream_persist_stream_info( .map(|o| o.xtream_info_cache).unwrap_or(false); if cache_info { if let Some(path) = get_xtream_storage_path(&app_state.config, target_name) { - let lock = app_state.shared_locks.get_lock(target_name); - let shared_lock = lock.write().unwrap(); - let mut index_tree = { - let (col_path, idx_path) = get_info_collection_and_idx_path(&path, cluster); - if idx_path.exists() && col_path.exists() { - load_index(&idx_path).unwrap_or_default() - } else { - IndexTree::new() - } - }; - match write_xtream_info(app_state, target_name, stream_id, cluster, content, - &mut index_tree) { - Ok(_) => {} - Err(err) => { error!("{}", err.to_string()); } + let lock = app_state.shared_locks.get_lock(target_name); + let shared_lock = lock.write().unwrap(); + let mut index_tree = { + let (col_path, idx_path) = get_info_collection_and_idx_path(&path, cluster); + if idx_path.exists() && col_path.exists() { + load_index(&idx_path).unwrap_or_default() + } else { + IndexTree::new() } - drop(shared_lock); + }; + match write_xtream_info(app_state, target_name, stream_id, cluster, content, + &mut index_tree) { + Ok(_) => {} + Err(err) => { error!("{}", err.to_string()); } + } + drop(shared_lock); } } +} + +fn get_id_mapping_path(path: &Path) -> PathBuf { + path.join("id_mapping.db") +} + +pub(crate) fn write_xtream_mapping(mappings: &[MultiXtreamMapping], config: &Config, target_name: &str) -> io::Result<()> { + if let Some(path) = get_xtream_storage_path(config, target_name) { + let mut file = File::create(get_id_mapping_path(&path))?; + // We assume the mappings list is created with a counter as id + // and id 1 means the 0 index. We write all the data and can calculate the offset inside the + // file by (u32 size + u16 size) * index. + for mapping in mappings { + file.write_all(&mapping.stream_id.to_le_bytes())?; + file.write_all(&mapping.input_id.to_le_bytes())?; + }; + return Ok(()); + } + Err(io::Error::new(ErrorKind::Other, format!("Failed to find the xtream storage path for {}", target_name))) +} + +pub(crate) fn read_xtream_mapping(id: u32, config: &Config, target_name: &str) -> io::Result> { + if id < 1 { + return Err(io::Error::new(ErrorKind::Other, "id should start with 1")); + } + + if let Some(path) = get_xtream_storage_path(config, target_name) { + let mapping_file_path = get_id_mapping_path(&path); + if mapping_file_path.exists() { + let mut file = File::open(&mapping_file_path)?; + let index = (id - 1) as u64; + let mapping_size = 4 + 2; // u32 + u16 + let offset = mapping_size * index; + + file.seek(SeekFrom::Start(offset))?; + let mut stream_id_bytes = [0u8; 4]; + file.read_exact(&mut stream_id_bytes)?; + let stream_id = u32::from_le_bytes(stream_id_bytes); + let mut input_id_bytes = [0u8; 2]; + file.read_exact(&mut input_id_bytes)?; + let input_id = u16::from_le_bytes(input_id_bytes); + return Ok(Some(MultiXtreamMapping { stream_id, input_id })); + } + } + Ok(None) } \ No newline at end of file diff --git a/src/test.rs b/src/test.rs index d805020f8..93739355c 100644 --- a/src/test.rs +++ b/src/test.rs @@ -1,6 +1,8 @@ #[cfg(test)] mod tests { use crate::filter::get_filter; + use crate::model::model_xtream::MultiXtreamMapping; + use crate::repository::xtream_repository::{read_xtream_mapping, write_xtream_mapping}; #[test] fn test_filter() { @@ -9,7 +11,25 @@ mod tests { Ok(filter) => { assert_eq!(format!("{}", filter), flt1); }, - Err(e) => {} + Err(_e) => {} } } + + + // #[test] + // fn test_xtream_id_mapping() { + // let mappings = vec![ + // MultiXtreamMapping { stream_id: 2, input_id: 3 }, + // MultiXtreamMapping { stream_id: 4, input_id: 5 }, + // MultiXtreamMapping { stream_id: 8, input_id: 6 }, + // ]; + // + // write_xtream_mapping(&mappings).unwrap(); + // for i in 1..=mappings.len() { + // let mapping = read_xtream_mapping(i as u32).unwrap().unwrap(); + // let test_mapping = mappings.get(i-1).unwrap(); + // assert_eq!(mapping.stream_id, test_mapping.stream_id); + // assert_eq!(mapping.input_id, test_mapping.input_id); + // } + // } } \ No newline at end of file From 3527bfd9059136d495cbdd4a8a4da4c24f6fab24 Mon Sep 17 00:00:00 2001 From: euzu Date: Mon, 18 Mar 2024 17:30:52 +0100 Subject: [PATCH 04/12] added stream id mapping to xtream api --- src/api/api_utils.rs | 4 +- src/api/xmltv_api.rs | 3 +- src/api/xtream_api.rs | 88 ++++++++++++++++++++++++++++++++++--------- src/model/config.rs | 7 ++-- 4 files changed, 78 insertions(+), 24 deletions(-) diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 35ecbe0df..e89f7ed4d 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -3,7 +3,7 @@ use actix_web::http::header::{CACHE_CONTROL, HeaderValue}; use actix_web::{HttpRequest, HttpResponse, web}; use crate::api::api_model::{AppState, UserApiRequest}; use crate::model::api_proxy::{UserCredentials}; -use crate::model::config::ConfigTarget; +use crate::model::config::{ConfigTarget}; pub(crate) async fn serve_file(file_path: &Path, req: &HttpRequest) -> HttpResponse { if file_path.exists() { @@ -36,4 +36,4 @@ pub(crate) fn get_user_target<'a>(api_req: &'a UserApiRequest, app_state: &'a we let username = api_req.username.as_str().trim(); let password = api_req.password.as_str().trim(); get_user_target_by_credentials(username, password, api_req, app_state) -} \ No newline at end of file +} diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 02ac0e011..1cd4b89d9 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -54,7 +54,8 @@ async fn xmltv_api( // If you want to use xmltv then provide the url in the config to filter unnecessary content. // If you have multiple xtream sources, the first one will be used for epg let target_name = &target.name; - if let Some(input) = _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { + let inputs = _app_state.config.get_input_for_target(target_name, &InputType::Xtream); + if let Some(&input) = inputs.first() { let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); let api_url = if epg_url.is_empty() { format!("{}/xmltv.php?username={}&password={}", diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index ddbc586cd..9cc9307b4 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -14,8 +14,7 @@ use crate::model::api_proxy::{ProxyType, UserCredentials}; use crate::model::config::{Config, ConfigInput, InputType}; use crate::model::model_config::{TargetType}; use crate::model::model_playlist::XtreamCluster; -use crate::repository::xtream_repository::{COL_CAT_LIVE, COL_CAT_SERIES, COL_CAT_VOD, COL_LIVE, COL_SERIES, COL_VOD, - xtream_get_all, xtream_get_stored_stream_info, xtream_persist_stream_info}; +use crate::repository::xtream_repository::{COL_CAT_LIVE, COL_CAT_SERIES, COL_CAT_VOD, COL_LIVE, COL_SERIES, COL_VOD, read_xtream_mapping, xtream_get_all, xtream_get_stored_stream_info, xtream_persist_stream_info}; use crate::utils::{get_client_request}; fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { @@ -107,14 +106,33 @@ async fn xtream_player_api_stream( if let Some((user, target)) = get_user_target_by_credentials(username, password, api_req, _app_state) { let target_name = &target.name; if target.has_output(&TargetType::Xtream) { - - !! todo id mapping when multi_xtream !! - - if let Some(target_input) = match _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - None => _app_state.config.get_input_for_target(target_name, &InputType::M3u), - Some(inp) => Some(inp) - } { - if let Some(stream_url) = get_xtream_player_api_stream_url(target_input, context, action_path) { + let mut stream_id = action_path.to_owned(); + let mut input: Option<&ConfigInput> = None; + let mut inputs = _app_state.config.get_input_for_target(target_name, &InputType::Xtream); + if inputs.is_empty() { + inputs = _app_state.config.get_input_for_target(target_name, &InputType::M3u); + if !inputs.is_empty() { + input = inputs.first().cloned(); + } + } else if inputs.len() > 1 { + if let Ok(num) = action_path.trim().parse() { + if let Ok(Some(mapping)) = read_xtream_mapping(num, _app_state.config.as_ref(), target_name) { + match inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { + None => { + input = inputs.first().cloned(); + } + Some(cfg_input) => { + input = Some(cfg_input); + stream_id = mapping.stream_id.to_string(); + } + } + } + } + } else if !inputs.is_empty() { + input = inputs.first().cloned(); + } + if let Some(target_input) = input { + if let Some(stream_url) = get_xtream_player_api_stream_url(target_input, context, stream_id.as_str()) { if user.proxy == ProxyType::Redirect { debug!("Redirecting stream request to {}", stream_url); return HttpResponse::Found().insert_header(("Location", stream_url)).finish(); @@ -210,14 +228,41 @@ async fn xtream_player_api_timeshift_stream( xtream_player_api_stream(&req, &api_req, &_app_state, "timeshift", &username, &password, &action_path).await } + +fn get_xtream_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>){ + let mut xtream_id = stream_id; + let inputs = app_state.config.get_input_for_target(target_name, &InputType::Xtream); + let input = if inputs.len() > 1 { + if let Ok(Some(mapping)) = read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { + match inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { + None => { + inputs.first().cloned() + } + Some(cfg_input) => { + xtream_id = mapping.stream_id as i32; + Some(cfg_input) + } + } + } else { + inputs.first().cloned() + } + } else { + inputs.first().cloned() + }; + + (xtream_id, input) +} + +// TODO use u32 or i32 ?? async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { - if let Some(target_input) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - if let Ok(content) = xtream_get_stored_stream_info(app_state, target_name, stream_id, cluster, target_input).await { + let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, stream_id); + if let Some(target_input) = input { + if let Ok(content) = xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { return Ok(content); } - if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, stream_id) { + if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { if let Ok(url) = Url::parse(&info_url) { let client = get_client_request(target_input, url, None); if let Ok(response) = client.send().await { @@ -226,7 +271,7 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ match response.text().await { Ok(content) => { // TODO we are not replacing direct_source, we should add an option to do this. - xtream_persist_stream_info(app_state, target_name, stream_id, cluster, + xtream_persist_stream_info(app_state, target_name, xtream_id, cluster, target_input, content.as_str()).await; return Ok(content); } @@ -246,8 +291,9 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCreden match FromStr::from_str(stream_id) { Ok(xtream_stream_id) => { if user.proxy == ProxyType::Redirect { - if let Some(target_input) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_stream_id) { + let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); + if let Some(target_input) = input { + if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } } @@ -263,9 +309,15 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCreden } async fn xtream_get_short_epg(app_state: &AppState, user: &UserCredentials, target_name: &str, stream_id: &str, limit: &str) -> HttpResponse { - if let Some(target_input) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { + let xtream_stream_id:i32 = match FromStr::from_str(stream_id) { + Ok(id) => id, + Err(_) => return HttpResponse::BadRequest().finish() + }; + + let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); + if let Some(target_input) = input { if let Some(action_url) = get_xtream_player_api_action_url(target_input, "get_short_epg") { - let mut info_url = format!("{}&stream_id={}", action_url, stream_id); + let mut info_url = format!("{}&stream_id={}", action_url, xtream_id); if !(limit.is_empty() || limit.eq("0")) { info_url = format!("{}&limit={}", info_url, limit); } diff --git a/src/model/config.rs b/src/model/config.rs index 32c7b1bc8..e61a789ca 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -604,11 +604,12 @@ impl Config { } } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option<&ConfigInput> { + pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Vec<&ConfigInput> { + let mut result = Vec::new(); for source in &self.sources { - if let Some(cfg) = source.get_input_for_target(target_name, input_type) { return Some(cfg); } + if let Some(cfg) = source.get_input_for_target(target_name, input_type) { result.push(cfg); } } - None + result } pub fn get_target_for_user(&self, username: &str, password: &str) -> Option<(UserCredentials, &ConfigTarget)> { From 44694aa5d6e4459d41d1731a424da04a5082eb6c Mon Sep 17 00:00:00 2001 From: euzu Date: Mon, 18 Mar 2024 18:16:22 +0100 Subject: [PATCH 05/12] added stream id mapping to xtream api --- src/api/xmltv_api.rs | 42 +++++++++++++++------------- src/api/xtream_api.rs | 65 +++++++++++++++---------------------------- src/model/config.rs | 17 ++++++----- 3 files changed, 54 insertions(+), 70 deletions(-) diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 1cd4b89d9..86f58ec97 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -54,26 +54,28 @@ async fn xmltv_api( // If you want to use xmltv then provide the url in the config to filter unnecessary content. // If you have multiple xtream sources, the first one will be used for epg let target_name = &target.name; - let inputs = _app_state.config.get_input_for_target(target_name, &InputType::Xtream); - if let Some(&input) = inputs.first() { - let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); - let api_url = if epg_url.is_empty() { - format!("{}/xmltv.php?username={}&password={}", - input.url.as_str(), - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - ) - } else { epg_url.to_string() }; - if let Ok(url) = Url::parse(&api_url) { - if user.proxy == ProxyType::Redirect { - debug!("Redirecting epg request to {}", api_url); - return HttpResponse::Found().insert_header(("Location", api_url)).finish(); - } - let client = get_client_request(input, url, None); - if let Ok(response) = client.send().await { - if response.status().is_success() { - if let Ok(content) = response.text().await { - return HttpResponse::Ok().content_type(mime::TEXT_XML).body(content); + + if let Some(inputs) = _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { + if let Some(&input) = inputs.first() { + let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); + let api_url = if epg_url.is_empty() { + format!("{}/xmltv.php?username={}&password={}", + input.url.as_str(), + input.username.as_ref().unwrap_or(&"".to_string()).as_str(), + input.password.as_ref().unwrap_or(&"".to_string()).as_str(), + ) + } else { epg_url.to_string() }; + if let Ok(url) = Url::parse(&api_url) { + if user.proxy == ProxyType::Redirect { + debug!("Redirecting epg request to {}", api_url); + return HttpResponse::Found().insert_header(("Location", api_url)).finish(); + } + let client = get_client_request(input, url, None); + if let Ok(response) = client.send().await { + if response.status().is_success() { + if let Ok(content) = response.text().await { + return HttpResponse::Ok().content_type(mime::TEXT_XML).body(content); + } } } } diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 9cc9307b4..c0f3ababb 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -108,29 +108,20 @@ async fn xtream_player_api_stream( if target.has_output(&TargetType::Xtream) { let mut stream_id = action_path.to_owned(); let mut input: Option<&ConfigInput> = None; - let mut inputs = _app_state.config.get_input_for_target(target_name, &InputType::Xtream); - if inputs.is_empty() { - inputs = _app_state.config.get_input_for_target(target_name, &InputType::M3u); - if !inputs.is_empty() { - input = inputs.first().cloned(); + if let Ok(num) = action_path.trim().parse() { + let (xtream_id, cfg_input) = get_xtream_input_for_stream_id(_app_state, target_name, num); + if cfg_input.is_some() { + input = cfg_input; + stream_id = xtream_id.to_string(); } - } else if inputs.len() > 1 { - if let Ok(num) = action_path.trim().parse() { - if let Ok(Some(mapping)) = read_xtream_mapping(num, _app_state.config.as_ref(), target_name) { - match inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { - None => { - input = inputs.first().cloned(); - } - Some(cfg_input) => { - input = Some(cfg_input); - stream_id = mapping.stream_id.to_string(); - } - } - } - } - } else if !inputs.is_empty() { - input = inputs.first().cloned(); } + + if input.is_none() { + if let Some(m3u_inputs) = _app_state.config.get_input_for_target(target_name, &InputType::M3u) { + input = m3u_inputs.first().cloned(); + } + } + if let Some(target_input) = input { if let Some(stream_url) = get_xtream_player_api_stream_url(target_input, context, stream_id.as_str()) { if user.proxy == ProxyType::Redirect { @@ -229,35 +220,23 @@ async fn xtream_player_api_timeshift_stream( } -fn get_xtream_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>){ - let mut xtream_id = stream_id; - let inputs = app_state.config.get_input_for_target(target_name, &InputType::Xtream); - let input = if inputs.len() > 1 { +fn get_xtream_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>) { + if let Some(inputs) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { if let Ok(Some(mapping)) = read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { - match inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { - None => { - inputs.first().cloned() - } - Some(cfg_input) => { - xtream_id = mapping.stream_id as i32; - Some(cfg_input) - } + if let Some(cfg_input) = inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { + return (mapping.stream_id as i32, Some(cfg_input)); } - } else { - inputs.first().cloned() } - } else { - inputs.first().cloned() - }; - - (xtream_id, input) + return (stream_id, inputs.first().cloned()); + } + (stream_id, None) } // TODO use u32 or i32 ?? async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, stream_id); - if let Some(target_input) = input { + if let Some(target_input) = input { if let Ok(content) = xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { return Ok(content); } @@ -292,7 +271,7 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCreden Ok(xtream_stream_id) => { if user.proxy == ProxyType::Redirect { let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); - if let Some(target_input) = input { + if let Some(target_input) = input { if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } @@ -309,7 +288,7 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCreden } async fn xtream_get_short_epg(app_state: &AppState, user: &UserCredentials, target_name: &str, stream_id: &str, limit: &str) -> HttpResponse { - let xtream_stream_id:i32 = match FromStr::from_str(stream_id) { + let xtream_stream_id: i32 = match FromStr::from_str(stream_id) { Ok(id) => id, Err(_) => return HttpResponse::BadRequest().finish() }; diff --git a/src/model/config.rs b/src/model/config.rs index e61a789ca..723293e82 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -311,14 +311,16 @@ impl ConfigSource { Ok(index + (self.inputs.len() as u16)) } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option<&ConfigInput> { + pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option> { + let mut result = Vec::new(); for target in &self.targets { if target.name.eq(target_name) { for input in &self.inputs { - if input.input_type.eq(input_type) { - return Some(input); + if input.enabled && input.input_type.eq(input_type) { + result.push(input); } } + return Some(result); } } None @@ -604,12 +606,13 @@ impl Config { } } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Vec<&ConfigInput> { - let mut result = Vec::new(); + pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option> { for source in &self.sources { - if let Some(cfg) = source.get_input_for_target(target_name, input_type) { result.push(cfg); } + if let Some(cfg) = source.get_input_for_target(target_name, input_type) { + return Some(cfg); + } } - result + None } pub fn get_target_for_user(&self, username: &str, password: &str) -> Option<(UserCredentials, &ConfigTarget)> { From d3bdd84eda7a8e8a96ec5ccbb40803e2e00ab408 Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 19 Mar 2024 10:20:06 +0100 Subject: [PATCH 06/12] added stream id mapping to xtream api --- src/api/xtream_api.rs | 41 +++++++++++++++++++++++++---------------- 1 file changed, 25 insertions(+), 16 deletions(-) diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index c0f3ababb..b0727ddcd 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -94,6 +94,13 @@ fn get_user_info(user: &UserCredentials, cfg: &Config) -> XtreamAuthorizationRes } } +fn separate_number_and_rest(input: &str) -> (String, String) { + let dot_index = input.find('.').unwrap_or(input.len()); + let number_part = input[..dot_index].to_string(); + let rest = input[dot_index..].to_string(); + (number_part, rest) +} + async fn xtream_player_api_stream( req: &HttpRequest, api_req: &web::Query, @@ -108,11 +115,13 @@ async fn xtream_player_api_stream( if target.has_output(&TargetType::Xtream) { let mut stream_id = action_path.to_owned(); let mut input: Option<&ConfigInput> = None; - if let Ok(num) = action_path.trim().parse() { + + let (action_stream_id, action_ext) = separate_number_and_rest(action_path); + if let Ok(num) = action_stream_id.trim().parse() { let (xtream_id, cfg_input) = get_xtream_input_for_stream_id(_app_state, target_name, num); if cfg_input.is_some() { input = cfg_input; - stream_id = xtream_id.to_string(); + stream_id = format!("{}{}", xtream_id.to_string(), action_ext); } } @@ -267,23 +276,23 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCredentials, target_name: &str, stream_id: &str, cluster: &XtreamCluster) -> HttpResponse { - match FromStr::from_str(stream_id) { - Ok(xtream_stream_id) => { - if user.proxy == ProxyType::Redirect { - let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); - if let Some(target_input) = input { - if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { - return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } - } - } + let xtream_stream_id: i32 = match FromStr::from_str(stream_id) { + Ok(id) => id, + Err(_) => return HttpResponse::BadRequest().finish() + }; - match xtream_get_stream_info(app_state, target_name, xtream_stream_id, cluster).await { - Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), - Err(_) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{info:[]}"), + if user.proxy == ProxyType::Redirect { + let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); + if let Some(target_input) = input { + if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { + return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } } - Err(_) => HttpResponse::BadRequest().finish() + } + + match xtream_get_stream_info(app_state, target_name, xtream_stream_id, cluster).await { + Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), + Err(_) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{info:[]}"), } } From b8ecd8b574750120b0e02762bc3bd497f181d0f1 Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 19 Mar 2024 10:25:03 +0100 Subject: [PATCH 07/12] added stream id mapping to xtream api --- src/api/xtream_api.rs | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index b0727ddcd..5a270bbbf 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -95,10 +95,13 @@ fn get_user_info(user: &UserCredentials, cfg: &Config) -> XtreamAuthorizationRes } fn separate_number_and_rest(input: &str) -> (String, String) { - let dot_index = input.find('.').unwrap_or(input.len()); - let number_part = input[..dot_index].to_string(); - let rest = input[dot_index..].to_string(); - (number_part, rest) + if let Some(dot_index) = input.find('.') { + let number_part = input[..dot_index].to_string(); + let rest = input[dot_index..].to_string(); + (number_part, rest) + } else { + (input.to_string(), String::new()) + } } async fn xtream_player_api_stream( From 6f89a17fa2eaf07875916224944cdf3a7a42008b Mon Sep 17 00:00:00 2001 From: euzu Date: Mon, 25 Mar 2024 21:07:48 +0100 Subject: [PATCH 08/12] added stream id mapping to xtream api --- src/api/xtream_api.rs | 12 ++++++------ src/model/config.rs | 2 +- src/processing/playlist_processor.rs | 7 ++----- 3 files changed, 9 insertions(+), 12 deletions(-) diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 5a270bbbf..b6c48efc3 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -121,10 +121,10 @@ async fn xtream_player_api_stream( let (action_stream_id, action_ext) = separate_number_and_rest(action_path); if let Ok(num) = action_stream_id.trim().parse() { - let (xtream_id, cfg_input) = get_xtream_input_for_stream_id(_app_state, target_name, num); + let (xtream_id, cfg_input) = get_xtream_mapped_id_and_input_for_stream_id(_app_state, target_name, num); if cfg_input.is_some() { input = cfg_input; - stream_id = format!("{}{}", xtream_id.to_string(), action_ext); + stream_id = format!("{}{}", xtream_id, action_ext); } } @@ -232,7 +232,7 @@ async fn xtream_player_api_timeshift_stream( } -fn get_xtream_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>) { +fn get_xtream_mapped_id_and_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>) { if let Some(inputs) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { if let Ok(Some(mapping)) = read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { if let Some(cfg_input) = inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { @@ -247,7 +247,7 @@ fn get_xtream_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str // TODO use u32 or i32 ?? async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { - let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, stream_id); + let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, stream_id); if let Some(target_input) = input { if let Ok(content) = xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { return Ok(content); @@ -285,7 +285,7 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCreden }; if user.proxy == ProxyType::Redirect { - let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); + let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, xtream_stream_id); if let Some(target_input) = input { if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); @@ -305,7 +305,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &UserCredentials, targ Err(_) => return HttpResponse::BadRequest().finish() }; - let (xtream_id, input) = get_xtream_input_for_stream_id(app_state, target_name, xtream_stream_id); + let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, xtream_stream_id); if let Some(target_input) = input { if let Some(action_url) = get_xtream_player_api_action_url(target_input, "get_short_epg") { let mut info_url = format!("{}&stream_id={}", action_url, xtream_id); diff --git a/src/model/config.rs b/src/model/config.rs index 723293e82..dc5cf07e4 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -508,7 +508,7 @@ impl VideoConfig { Some(downl) => { if downl.headers.is_empty() { downl.headers.borrow_mut().insert("Accept".to_string(), "video/*".to_string()); - downl.headers.borrow_mut().insert("User-Agent".to_string(), "AppleTV/tvOS/9.1.1.".to_string()); + downl.headers.borrow_mut().insert("User-Agent".to_string(), "Apple TV; tvOS 13.3.1".to_string()); } if let Some(episode_pattern) = &downl.episode_pattern { diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index e41fbd50a..27af70776 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -188,11 +188,8 @@ fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt macro_rules! apply_pattern { ($pattern:expr, $provider:expr, $processor:expr) => {{ - match $pattern { - Some(ptrn) => { - ptrn.filter($provider, $processor); - }, - _ => {} + if let Some(ptrn) = $pattern { + ptrn.filter($provider, $processor); }; }}; } From 51609614323c249521821cc05a52888b601905fb Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 26 Mar 2024 01:21:02 +0100 Subject: [PATCH 09/12] Merge branch 'bugfi/m3u_missing_id' into bugfix/multi_xtream_input # Conflicts: # src/api/v1_api.rs # src/api/xmltv_api.rs # src/api/xtream_api.rs # src/main.rs # src/processing/playlist_processor.rs # src/repository/xtream_repository.rs # src/util/config_reader.rs # src/util/utils.rs # src/utils.rs # src/utils/request_utils.rs --- src/api/v1_api.rs | 9 ++++----- src/api/xtream_api.rs | 25 +++++++++++-------------- src/main.rs | 17 ++++++----------- src/utils/config_reader.rs | 6 ++---- src/utils/mod.rs | 4 +++- 5 files changed, 26 insertions(+), 35 deletions(-) diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index fad3f74f1..5abe5de0e 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -8,11 +8,10 @@ use crate::api::download_api::{download_file_info, queue_download_file}; use crate::m3u_filter_error::M3uFilterError; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, TargetUser}; use crate::processing::playlist_processor::exec_processing; -use crate::util::config_reader::{read_config, save_api_proxy, save_main_config}; -use crate::utils::download; +use crate::utils::{config_reader, download}; fn _save_config_api_proxy(backup_dir: &str, api_proxy: &mut ApiProxyConfig) -> Option { - match save_api_proxy(api_proxy._file_path.as_str(), backup_dir, api_proxy) { + match config_reader::save_api_proxy(api_proxy._file_path.as_str(), backup_dir, api_proxy) { Ok(_) => {} Err(err) => { error!("Failed to save api_proxy.yml {}", err.to_string()); @@ -23,7 +22,7 @@ fn _save_config_api_proxy(backup_dir: &str, api_proxy: &mut ApiProxyConfig) -> O } fn _save_config_main(file_path: &str, backup_dir: &str, cfg: &ConfigDto) -> Option { - match save_main_config(file_path, backup_dir, cfg) { + match config_reader::save_main_config(file_path, backup_dir, cfg) { Ok(_) => {} Err(err) => { error!("Failed to save config.yml {}", err.to_string()); @@ -200,7 +199,7 @@ pub(crate) async fn config( api_proxy: config._api_proxy.read().unwrap().clone(), }; - let mut result = match read_config(_app_state.config._config_path.as_str(), + let mut result = match config_reader::read_config(_app_state.config._config_path.as_str(), _app_state.config._config_file_path.as_str(), _app_state.config._sources_file_path.as_str()) { Ok(mut cfg) => { diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 49271d677..3286cc33e 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -14,11 +14,8 @@ use crate::model::api_proxy::{ProxyType, UserCredentials}; use crate::model::config::{Config, ConfigInput, InputType}; use crate::model::model_config::{TargetType}; use crate::model::model_playlist::XtreamCluster; -use crate::repository::xtream_repository::{COL_CAT_LIVE, COL_CAT_SERIES, COL_CAT_VOD, COL_LIVE, COL_SERIES, COL_VOD, read_xtream_mapping, xtream_get_all, xtream_get_stored_stream_info, xtream_persist_stream_info}; -use crate::utils::{get_client_request}; -use crate::repository::xtream_repository::{COL_CAT_LIVE, COL_CAT_SERIES, COL_CAT_VOD, COL_LIVE, COL_SERIES, COL_VOD, - xtream_get_all, xtream_get_stored_stream_info, xtream_persist_stream_info}; -use crate::utils::request_utils; +use crate::repository::xtream_repository; +use crate::utils::{request_utils}; fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { match input.input_type { @@ -237,7 +234,7 @@ async fn xtream_player_api_timeshift_stream( fn get_xtream_mapped_id_and_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>) { if let Some(inputs) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - if let Ok(Some(mapping)) = read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { + if let Ok(Some(mapping)) = xtream_repository::read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { if let Some(cfg_input) = inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { return (mapping.stream_id as i32, Some(cfg_input)); } @@ -252,7 +249,7 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ cluster: &XtreamCluster) -> Result { let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, stream_id); if let Some(target_input) = input { - if let Ok(content) = xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { + if let Ok(content) = xtream_repository::xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { return Ok(content); } @@ -265,7 +262,7 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ match response.text().await { Ok(content) => { // TODO we are not replacing direct_source, we should add an option to do this. - xtream_persist_stream_info(app_state, target_name, xtream_id, cluster, + xtream_repository::xtream_persist_stream_info(app_state, target_name, xtream_id, cluster, target_input, content.as_str()).await; return Ok(content); } @@ -374,12 +371,12 @@ async fn xtream_player_api( } _ => { match match action { - "get_live_categories" => xtream_get_all(&_app_state.config, target_name, COL_CAT_LIVE), - "get_vod_categories" => xtream_get_all(&_app_state.config, target_name, COL_CAT_VOD), - "get_series_categories" => xtream_get_all(&_app_state.config, target_name, COL_CAT_SERIES), - "get_live_streams" => xtream_get_all(&_app_state.config, target_name, COL_LIVE), - "get_vod_streams" => xtream_get_all(&_app_state.config, target_name, COL_VOD), - "get_series" => xtream_get_all(&_app_state.config, target_name, COL_SERIES), + "get_live_categories" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_CAT_LIVE), + "get_vod_categories" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_CAT_VOD), + "get_series_categories" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_CAT_SERIES), + "get_live_streams" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_LIVE), + "get_vod_streams" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_VOD), + "get_series" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_SERIES), _ => Err(Error::new(std::io::ErrorKind::Unsupported, format!("Cant find action: {}/{}", target_name, action))), } { Ok((path, content)) => { diff --git a/src/main.rs b/src/main.rs index 34fa85250..1454dcaf4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -11,22 +11,17 @@ use env_logger::Builder; use log::{error, info, LevelFilter}; use crate::model::config::{Config, ProcessTargets, validate_targets}; -use crate::processing::playlist_processor::exec_processing; -use crate::util::config_reader::{read_api_proxy_config, read_config, read_mappings}; -use crate::util::utils; -use crate::utils::file_utils; +use crate::processing::playlist_processor; +use crate::utils::{config_reader, file_utils}; mod m3u_filter_error; mod model; mod filter; mod repository; -mod download; mod messaging; mod test; mod api; mod processing; -mod util; -mod multi_file_reader; mod utils; #[derive(Parser)] @@ -81,7 +76,7 @@ fn main() { let sources_file: String = args.source_file.unwrap_or(file_utils::get_default_sources_file_path(&config_path)); - let mut cfg = read_config(config_path.as_str(), config_file.as_str(), sources_file.as_str()).unwrap_or_else(|err| exit!("{}", err)); + let mut cfg = config_reader::read_config(config_path.as_str(), config_file.as_str(), sources_file.as_str()).unwrap_or_else(|err| exit!("{}", err)); if args.log_level.is_none() { if let Some(log_level) = &cfg.log_level { @@ -99,12 +94,12 @@ fn main() { info!("Config file: {}", &config_file); info!("Source file: {}", &sources_file); - if let Err(err) = read_mappings(args.mapping_file, &mut cfg) { + if let Err(err) = config_reader::read_mappings(args.mapping_file, &mut cfg) { exit!("{}", err); } if args.server { - read_api_proxy_config(args.api_proxy, &mut cfg); + config_reader::read_api_proxy_config(args.api_proxy, &mut cfg); start_in_server_mode(Arc::new(cfg), Arc::new(targets)); } else { start_in_cli_mode(Arc::new(cfg), Arc::new(targets)) @@ -112,7 +107,7 @@ fn main() { } fn start_in_cli_mode(cfg: Arc, targets: Arc) { - System::new().block_on(async { exec_processing(cfg, targets).await }); + System::new().block_on(async { playlist_processor::exec_processing(cfg, targets).await }); } fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/src/utils/config_reader.rs b/src/utils/config_reader.rs index 52140fa36..a3afa45de 100644 --- a/src/utils/config_reader.rs +++ b/src/utils/config_reader.rs @@ -8,9 +8,7 @@ use crate::model::config::{Config, ConfigDto}; use crate::model::mapping::Mappings; use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::util::multi_file_reader::MultiFileReader; -use crate::multi_file_reader::MultiFileReader; -use crate::utils::file_utils; +use crate::utils::{file_utils, multi_file_reader}; pub(crate) fn read_mappings(args_mapping: Option, cfg: &mut Config) -> Result<(), M3uFilterError> { let mappings_file: String = args_mapping.unwrap_or(file_utils::get_default_mappings_path(cfg._config_path.as_str())); @@ -43,7 +41,7 @@ pub(crate) fn read_api_proxy_config(args_api_proxy_config: Option, cfg: pub(crate) fn read_config(config_path: &str, config_file: &str, sources_file: &str) -> Result { let files = vec![std::path::PathBuf::from(config_file), std::path::PathBuf::from(sources_file)]; - match MultiFileReader::new(&files) { + match multi_file_reader::MultiFileReader::new(&files) { Ok(file) => { match serde_yaml::from_reader::<_, Config>(file) { Ok(mut result) => { diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 79c00e8ce..2bf4fe538 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -1,4 +1,6 @@ pub (crate) mod file_utils; pub (crate) mod request_utils; pub (crate) mod download; -pub (crate) mod string_utils; \ No newline at end of file +pub (crate) mod string_utils; +pub(crate) mod config_reader; +pub(crate) mod multi_file_reader; \ No newline at end of file From 74c427765e4e41802785c134fd07dd3895207a20 Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 26 Mar 2024 10:36:05 +0100 Subject: [PATCH 10/12] Merge branch 'master' into bugfix/multi_xtream_input # Conflicts: # src/api/xtream_api.rs # src/utils/mod.rs --- src/api/api_model.rs | 2 + src/api/xtream_api.rs | 33 +++++++---- src/repository/xtream_repository.rs | 2 +- src/utils/json_utils.rs | 90 +++++++++++++++++++++++++++++ src/utils/mod.rs | 5 +- 5 files changed, 119 insertions(+), 13 deletions(-) create mode 100644 src/utils/json_utils.rs diff --git a/src/api/api_model.rs b/src/api/api_model.rs index 831ec7bae..00064b52f 100644 --- a/src/api/api_model.rs +++ b/src/api/api_model.rs @@ -210,6 +210,8 @@ pub(crate) struct UserApiRequest { #[serde(default = "default_as_empty_str")] pub stream_id: String, #[serde(default = "default_as_empty_str")] + pub category_id: String, + #[serde(default = "default_as_empty_str")] pub limit: String, } diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 3286cc33e..9894c819a 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -2,6 +2,7 @@ use std::collections::HashMap; use std::io::{Error}; +use std::path::Path; use std::str::FromStr; use actix_web::{HttpRequest, HttpResponse, web, Resource}; use chrono::{Duration, Local}; @@ -15,7 +16,12 @@ use crate::model::config::{Config, ConfigInput, InputType}; use crate::model::model_config::{TargetType}; use crate::model::model_playlist::XtreamCluster; use crate::repository::xtream_repository; -use crate::utils::{request_utils}; +use crate::utils::{json_utils, request_utils}; + +pub(crate) async fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) -> HttpResponse { + let filtered = json_utils::filter_json_file(file_path, filter); + HttpResponse::Ok().json(filtered) +} fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { match input.input_type { @@ -247,6 +253,7 @@ fn get_xtream_mapped_id_and_input_for_stream_id<'a>(app_state: &'a AppState, tar // TODO use u32 or i32 ?? async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { + let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, stream_id); if let Some(target_input) = input { if let Ok(content) = xtream_repository::xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { @@ -263,7 +270,7 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ Ok(content) => { // TODO we are not replacing direct_source, we should add an option to do this. xtream_repository::xtream_persist_stream_info(app_state, target_name, xtream_id, cluster, - target_input, content.as_str()).await; + target_input, content.as_str()).await; return Ok(content); } Err(err) => { error!("Failed to download info {}", err.to_string()); } @@ -328,7 +335,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &UserCredentials, targ error!("Failed to download epg {}", err.to_string()); HttpResponse::NoContent().finish() } - } + }; } } } @@ -371,17 +378,23 @@ async fn xtream_player_api( } _ => { match match action { - "get_live_categories" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_CAT_LIVE), - "get_vod_categories" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_CAT_VOD), - "get_series_categories" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_CAT_SERIES), - "get_live_streams" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_LIVE), - "get_vod_streams" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_VOD), - "get_series" => xtream_repository::xtream_get_all(&_app_state.config, target_name, xtream_repository::COL_SERIES), + + "get_live_categories" => xtream_repository::xtream_get_collection_path(&_app_state.config, target_name, xtream_repository::COL_CAT_LIVE), + "get_vod_categories" => xtream_repository::xtream_get_collection_path(&_app_state.config, target_name, xtream_repository::COL_CAT_VOD), + "get_series_categories" => xtream_repository::xtream_get_collection_path(&_app_state.config, target_name, xtream_repository::COL_CAT_SERIES), + "get_live_streams" => xtream_repository::xtream_get_collection_path(&_app_state.config, target_name, xtream_repository::COL_LIVE), + "get_vod_streams" => xtream_repository::xtream_get_collection_path(&_app_state.config, target_name, xtream_repository::COL_VOD), + "get_series" => xtream_repository::xtream_get_collection_path(&_app_state.config, target_name, xtream_repository::COL_SERIES), _ => Err(Error::new(std::io::ErrorKind::Unsupported, format!("Cant find action: {}/{}", target_name, action))), } { Ok((path, content)) => { if let Some(file_path) = path { - serve_file(&file_path, req).await + let category_id = api_req.category_id.trim(); + if !category_id.is_empty() { + serve_query(&file_path, &HashMap::from([("category_id", category_id)])).await + } else { + serve_file(&file_path, req).await + } } else if let Some(payload) = content { HttpResponse::Ok().body(payload) } else { diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index fa783dffa..51ffc21a1 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -308,7 +308,7 @@ fn append_mandatory_fields(document: &mut Map, fields: &[&str]) { } } -pub(crate) fn xtream_get_all(cfg: &Config, target_name: &str, collection_name: &str) -> Result<(Option, Option), Error> { +pub(crate) fn xtream_get_collection_path(cfg: &Config, target_name: &str, collection_name: &str) -> Result<(Option, Option), Error> { if let Some(path) = get_xtream_storage_path(cfg, target_name) { let col_path = get_collection_path(&path, collection_name); if col_path.exists() { diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs new file mode 100644 index 000000000..f8e6dd184 --- /dev/null +++ b/src/utils/json_utils.rs @@ -0,0 +1,90 @@ +use std::collections::HashMap; +use std::fs::File; +use serde::de::DeserializeOwned; +use serde_json::{self, Deserializer}; +use std::io::{self, BufReader, Read}; +use std::path::Path; + +fn read_skipping_ws(mut reader: impl Read) -> io::Result { + loop { + let mut byte = 0u8; + reader.read_exact(std::slice::from_mut(&mut byte))?; + if !byte.is_ascii_whitespace() { + return Ok(byte); + } + } +} + +fn invalid_data(msg: &str) -> io::Error { + io::Error::new(io::ErrorKind::InvalidData, msg) +} + +fn deserialize_single(reader: R) -> io::Result { + let next_obj = Deserializer::from_reader(reader).into_iter::().next(); + match next_obj { + Some(result) => result.map_err(Into::into), + None => Err(invalid_data("premature EOF")), + } +} + +fn yield_next_obj( + mut reader: R, + at_start: &mut bool, +) -> io::Result> { + if !*at_start { + *at_start = true; + if read_skipping_ws(&mut reader)? == b'[' { + // read the next char to see if the array is empty + let peek = read_skipping_ws(&mut reader)?; + if peek == b']' { + Ok(None) + } else { + deserialize_single(io::Cursor::new([peek]).chain(reader)).map(Some) + } + } else { + Err(invalid_data("`[` not found")) + } + } else { + match read_skipping_ws(&mut reader)? { + b',' => deserialize_single(reader).map(Some), + b']' => Ok(None), + _ => Err(invalid_data("`,` or `]` not found")), + } + } +} + +// https://stackoverflow.com/questions/68641157/how-can-i-stream-elements-from-inside-a-json-array-using-serde-json +pub(crate) fn iter_json_array( + mut reader: R, +) -> impl Iterator> { + let mut at_start = false; + std::iter::from_fn(move || yield_next_obj(&mut reader, &mut at_start).transpose()) +} + +pub(crate) fn filter_json_file(file_path: &Path, filter: &HashMap<&str, &str>) -> Vec { + let mut filtered: Vec = Vec::new(); + if file_path.exists() { + if let Ok(file) = File::open(file_path) { + let reader = BufReader::new(file); + for entry in iter_json_array::>(reader).flatten() { + if let Some(item) = entry.as_object() { + for (&key, &value) in filter { + if let Some(field_value) = item.get(key) { + if field_value.is_string() && field_value.eq(value) { + filtered.push(entry.clone()); + } else if let Some(num_val) = field_value.as_i64() { + if let Ok(filter_num_val) = value.parse::() { + if num_val == filter_num_val { + filtered.push(entry.clone()); + } + } + } + } + } + } + } + } + } + + filtered +} diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 2bf4fe538..f8658a198 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -2,5 +2,6 @@ pub (crate) mod file_utils; pub (crate) mod request_utils; pub (crate) mod download; pub (crate) mod string_utils; -pub(crate) mod config_reader; -pub(crate) mod multi_file_reader; \ No newline at end of file +pub (crate) mod json_utils; +pub (crate) mod config_reader; +pub (crate) mod multi_file_reader; From 63c3e240061bb259c07dc83315a9056b376ab0d7 Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 28 Mar 2024 16:27:39 +0100 Subject: [PATCH 11/12] Merge branch 'bugfix/multi_xtream_input' into feature/multi_xtream_input # Conflicts: # src/api/v1_api.rs # src/api/xtream_api.rs # src/processing/playlist_processor.rs # src/utils/mod.rs --- src/api/api_utils.rs | 30 +++++++++++++++++++ src/api/xmltv_api.rs | 43 +++++++++++++++------------- src/api/xtream_api.rs | 41 ++++---------------------- src/model/config.rs | 23 +++++++++++---- src/processing/playlist_processor.rs | 21 +++----------- src/processing/xtream_parser.rs | 2 +- src/repository/xtream_repository.rs | 18 ++++++++---- 7 files changed, 92 insertions(+), 86 deletions(-) diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index e89f7ed4d..7b744e4f3 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -4,6 +4,36 @@ use actix_web::{HttpRequest, HttpResponse, web}; use crate::api::api_model::{AppState, UserApiRequest}; use crate::model::api_proxy::{UserCredentials}; use crate::model::config::{ConfigTarget}; +use url::Url; + +pub (crate) struct M3uUrlInfo { + pub base_url: String, + pub username: String, + pub password: String, +} + +pub (crate) fn parse_m3u_url(url: &str) -> Option { + if let Ok(url) = Url::parse(url) { + let base_url = url.origin().ascii_serialization(); + let mut username = None; + let mut password = None; + for (key, value) in url.query_pairs() { + if key.eq("username") { + username = Some(value.into_owned()); + } else if key.eq("password") { + password = Some(value.into_owned()); + } + } + if username.is_some() || password.is_some() { + return Some(M3uUrlInfo { + base_url, + username: username.as_ref().unwrap().to_owned(), + password: username.as_ref().unwrap().to_owned(), + }); + } + } + None +} pub(crate) async fn serve_file(file_path: &Path, req: &HttpRequest) -> HttpResponse { if file_path.exists() { diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 3c4b29eda..c6037f6bd 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -51,28 +51,31 @@ async fn xmltv_api( // If no epg_url is provided for input, we did not process the xmltv for our channels. // We are now delivering the original untouched xmltv. // If you want to use xmltv then provide the url in the config to filter unnecessary content. - // If you have multiple xtream sources, the first one will be used for epg + // If you have multiple xtream sources, no response because of mapped ids + // if you want epg for multi xtream input, then provide epg_url. let target_name = &target.name; if let Some(inputs) = _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - if let Some(&input) = inputs.first() { - let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); - let api_url = if epg_url.is_empty() { - format!("{}/xmltv.php?username={}&password={}", - input.url.as_str(), - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - ) - } else { epg_url.to_string() }; - if let Ok(url) = Url::parse(&api_url) { - if user.proxy == ProxyType::Redirect { - debug!("Redirecting epg request to {}", api_url); - return HttpResponse::Found().insert_header(("Location", api_url)).finish(); - } - let client = request_utils::get_client_request(input, url, None); - if let Ok(response) = client.send().await { - if response.status().is_success() { - if let Ok(content) = response.text().await { - return HttpResponse::Ok().content_type(mime::TEXT_XML).body(content); + if inputs.len() == 1 { + if let Some(&input) = inputs.first() { + let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); + let api_url = if epg_url.is_empty() { + format!("{}/xmltv.php?username={}&password={}", + input.url.as_str(), + input.username.as_ref().unwrap_or(&"".to_string()).as_str(), + input.password.as_ref().unwrap_or(&"".to_string()).as_str(), + ) + } else { epg_url.to_string() }; + if let Ok(url) = Url::parse(&api_url) { + if user.proxy == ProxyType::Redirect { + debug!("Redirecting epg request to {}", api_url); + return HttpResponse::Found().insert_header(("Location", api_url)).finish(); + } + let client = request_utils::get_client_request(input, url, None); + if let Ok(response) = client.send().await { + if response.status().is_success() { + if let Ok(content) = response.text().await { + return HttpResponse::Ok().content_type(mime::TEXT_XML).body(content); + } } } } diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 4d272f9a7..2cdebab61 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -11,6 +11,7 @@ use url::Url; use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, serve_file}; use crate::api::api_model::{AppState, UserApiRequest, XtreamAuthorizationResponse, XtreamServerInfo, XtreamUserInfo}; +use crate::api::api_utils; use crate::model::api_proxy::{ProxyType, UserCredentials}; use crate::model::config::{Config, ConfigInput, InputType}; use crate::model::config::{TargetType}; @@ -18,35 +19,6 @@ use crate::model::playlist::XtreamCluster; use crate::repository::xtream_repository; use crate::utils::{json_utils, request_utils}; -struct M3uUrlInfo { - pub base_url: String, - pub username: String, - pub password: String, -} - -fn parse_m3u_url(url: &str) -> Option { - if let Ok(url) = Url::parse(url) { - let base_url = url.origin().ascii_serialization(); - let mut username = None; - let mut password = None; - for (key, value) in url.query_pairs() { - if key.eq("username") { - username = Some(value.into_owned()); - } else if key.eq("password") { - password = Some(value.into_owned()); - } - } - if username.is_some() || password.is_some() { - return Some(M3uUrlInfo { - base_url, - username: username.as_ref().unwrap().to_owned(), - password: username.as_ref().unwrap().to_owned(), - }); - } - } - None -} - pub(crate) async fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) -> HttpResponse { let filtered = json_utils::filter_json_file(file_path, filter); HttpResponse::Ok().json(filtered) @@ -55,7 +27,7 @@ pub(crate) async fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { match input.input_type { InputType::M3u => { - match parse_m3u_url(input.url.as_str()) { + match api_utils::parse_m3u_url(input.url.as_str()) { None => None, Some(m3u_url_info) => Some( format!("{}/player_api.php?username={}&password={}&action={}", @@ -89,7 +61,7 @@ fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: &XtreamCluster, fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &str, action_path: &str) -> Option { let ctx_path = if context.is_empty() { "".to_string() } else { format!("{}/", context) }; match input.input_type { - InputType::M3u => match parse_m3u_url(input.url.as_str()) { + InputType::M3u => match api_utils::parse_m3u_url(input.url.as_str()) { None => None, Some(m3u_url_info) => Some( format!("{}/{}{}/{}/{}", @@ -174,7 +146,6 @@ async fn xtream_player_api_stream( if target.has_output(&TargetType::Xtream) { let mut stream_id = action_path.to_owned(); let mut input: Option<&ConfigInput> = None; - let (action_stream_id, action_ext) = separate_number_and_rest(action_path); if let Ok(num) = action_stream_id.trim().parse() { let (xtream_id, cfg_input) = get_xtream_mapped_id_and_input_for_stream_id(_app_state, target_name, num); @@ -300,13 +271,11 @@ fn get_xtream_mapped_id_and_input_for_stream_id<'a>(app_state: &'a AppState, tar (stream_id, None) } -// TODO use u32 or i32 ?? async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { - let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, stream_id); if let Some(target_input) = input { - if let Ok(content) = xtream_repository::xtream_get_stored_stream_info(app_state, target_name, xtream_id, cluster, target_input).await { + if let Ok(content) = xtream_repository::xtream_get_stored_stream_info(app_state, target_name, stream_id, cluster, target_input).await { return Ok(content); } @@ -341,7 +310,7 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCreden Err(_) => return HttpResponse::BadRequest().finish() }; - if user.proxy == ProxyType::Redirect { + if user.proxy == ProxyType::Redirect && !app_state.config.is_multi_xtream_input(target_name) { let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, xtream_stream_id); if let Some(target_input) = input { if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { diff --git a/src/model/config.rs b/src/model/config.rs index eaa6cc565..b69263255 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -35,7 +35,6 @@ macro_rules! valid_property { }}; } - pub(crate) fn default_as_true() -> bool { true } pub(crate) fn default_as_false() -> bool { false } @@ -425,6 +424,7 @@ pub(crate) struct ConfigSource { impl ConfigSource { pub(crate) fn prepare(&mut self, index: u16) -> Result { handle_m3u_filter_error_result_list!(M3uFilterErrorKind::Info, self.inputs.iter_mut().enumerate().map(|(idx, i)| i.prepare(index+(idx as u16)))); + self._multi_xtream_input = self.inputs.iter().filter(|i| i.input_type == InputType::Xtream).count() > 1; Ok(index + (self.inputs.len() as u16)) } @@ -740,7 +740,7 @@ impl Config { None } - pub fn get_target_for_user(&self, username: &str, password: &str) -> Option<(UserCredentials, &ConfigTarget)> { + pub(crate) fn get_target_for_user(&self, username: &str, password: &str) -> Option<(UserCredentials, &ConfigTarget)> { match self._api_proxy.read().unwrap().as_ref() { Some(api_proxy) => { self._get_target_for_user(api_proxy.get_target_name(username, password)) @@ -749,7 +749,7 @@ impl Config { } } - pub fn get_target_for_user_by_token(&self, token: &str) -> Option<(UserCredentials, &ConfigTarget)> { + pub(crate) fn get_target_for_user_by_token(&self, token: &str) -> Option<(UserCredentials, &ConfigTarget)> { match self._api_proxy.read().unwrap().as_ref() { Some(api_proxy) => { self._get_target_for_user(api_proxy.get_target_name_by_token(token)) @@ -758,7 +758,7 @@ impl Config { } } - pub fn get_input_by_id(&self, input_id: &u16) -> Option { + pub(crate) fn get_input_by_id(&self, input_id: &u16) -> Option { for source in &self.sources { for input in &source.inputs { if input.id == *input_id { @@ -769,7 +769,18 @@ impl Config { None } - pub fn set_mappings(&mut self, mappings: Option) -> Result<(), M3uFilterError> { + pub(crate) fn is_multi_xtream_input(&self, target_name: &str) -> bool { + for source in &self.sources { + for target in &source.targets { + if target_name.eq_ignore_ascii_case(&target.name) { + return source._multi_xtream_input; + } + } + } + false + } + + pub(crate) fn set_mappings(&mut self, mappings: Option) -> Result<(), M3uFilterError> { if let Some(mapping_list) = mappings { for source in &mut self.sources { for target in &mut source.targets { @@ -789,7 +800,7 @@ impl Config { Ok(()) } - pub fn prepare(&mut self) -> Result<(), M3uFilterError> { + pub(crate) fn prepare(&mut self) -> Result<(), M3uFilterError> { self.working_dir = file_utils::get_working_path(&self.working_dir); if self.backup_dir.is_none() { self.backup_dir = Some(PathBuf::from(&self.working_dir).join(".backup").into_os_string().to_string_lossy().to_string()); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 8109f046c..ee7c04af3 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -244,39 +244,26 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option // if the group names are changed, restructure channels to the right groups // we use - let mut max_group_id = 0; let mut new_groups: Vec = Vec::new(); + let mut grp_id: u32 = 0; for playlist_group in new_playlist { - let mut group_id_used = false; for channel in &playlist_group.channels { let cluster = &channel.header.borrow().xtream_cluster; let title = &channel.header.borrow().group; match new_groups.iter_mut().find(|x| *x.title == **title) { Some(grp) => grp.channels.push(channel.clone()), _ => { - let new_group_id = if group_id_used { - 0 - } else if *title == playlist_group.title { - group_id_used = true; - max_group_id = max_group_id.max(playlist_group.id); - playlist_group.id - } else { - 0 - }; + grp_id += 1; new_groups.push(PlaylistGroup { - id: new_group_id, + id: grp_id, title: Rc::clone(title), channels: vec![channel.clone()], - xtream_cluster: cluster.clone(), + xtream_cluster: cluster.clone() }) } } } } - new_groups.iter_mut().filter(|g| g.id == 0).for_each(|grp| { - max_group_id += 1; - grp.id = max_group_id; - }); Some(new_groups) } else { None diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index a55b9e79c..2e49fe2d5 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -143,7 +143,7 @@ pub(crate) fn parse_xtream(input: &ConfigInput, Ok(Some(group_map.values().map(|category| { let cat = category.borrow(); PlaylistGroup { - id: 0, + id: cat.category_id.parse::().unwrap_or(0), xtream_cluster: xtream_cluster.clone(), title: Rc::clone(&cat.category_name), channels: cat.channels.clone() diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 7b1e44a82..301f87682 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -13,7 +13,7 @@ use crate::model::playlist::{PlaylistGroup, PlaylistItemHeader, PlaylistItemType use crate::{create_m3u_filter_error_result}; use crate::api::api_model::AppState; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::model::model_xtream::MultiXtreamMapping; +use crate::model::xtream::MultiXtreamMapping; use crate::utils::file_utils; use crate::utils::json_utils::iter_json_array; @@ -450,19 +450,25 @@ fn get_id_mapping_path(path: &Path) -> PathBuf { path.join("id_mapping.db") } -pub(crate) fn write_xtream_mapping(mappings: &[MultiXtreamMapping], config: &Config, target_name: &str) -> io::Result<()> { +pub(crate) fn write_xtream_mapping(mappings: &[MultiXtreamMapping], config: &Config, target_name: &str) -> Result<(), M3uFilterError> { if let Some(path) = get_xtream_storage_path(config, target_name) { - let mut file = File::create(get_id_mapping_path(&path))?; + if fs::create_dir_all(&path).is_err() { + let msg = format!("Failed to save, can't create directory {}", &path.to_str().unwrap()); + return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)); + } + + let err_map = |e: Error| M3uFilterError::new(M3uFilterErrorKind::Notify, e.to_string()); + let mut file = File::create(get_id_mapping_path(&path)).map_err(err_map)?; // We assume the mappings list is created with a counter as id // and id 1 means the 0 index. We write all the data and can calculate the offset inside the // file by (u32 size + u16 size) * index. for mapping in mappings { - file.write_all(&mapping.stream_id.to_le_bytes())?; - file.write_all(&mapping.input_id.to_le_bytes())?; + file.write_all(&mapping.stream_id.to_le_bytes()).map_err(err_map)?; + file.write_all(&mapping.input_id.to_le_bytes()).map_err(err_map)?; }; return Ok(()); } - Err(io::Error::new(ErrorKind::Other, format!("Failed to find the xtream storage path for {}", target_name))) + Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to find the xtream storage path for {}", target_name))) } pub(crate) fn read_xtream_mapping(id: u32, config: &Config, target_name: &str) -> io::Result> { From b7a40e1ce80aa03541ffdbd37d67c225c3dabf86 Mon Sep 17 00:00:00 2001 From: euzu Date: Mon, 1 Apr 2024 11:59:54 +0200 Subject: [PATCH 12/12] multi_xtream input first steps --- src/api/api_utils.rs | 30 -------- src/api/xmltv_api.rs | 4 +- src/api/xtream_api.rs | 107 +++++++++++---------------- src/model/config.rs | 81 +++++++++++++++----- src/model/xtream.rs | 4 +- src/processing/playlist_processor.rs | 9 +-- 6 files changed, 114 insertions(+), 121 deletions(-) diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 7b744e4f3..e89f7ed4d 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -4,36 +4,6 @@ use actix_web::{HttpRequest, HttpResponse, web}; use crate::api::api_model::{AppState, UserApiRequest}; use crate::model::api_proxy::{UserCredentials}; use crate::model::config::{ConfigTarget}; -use url::Url; - -pub (crate) struct M3uUrlInfo { - pub base_url: String, - pub username: String, - pub password: String, -} - -pub (crate) fn parse_m3u_url(url: &str) -> Option { - if let Ok(url) = Url::parse(url) { - let base_url = url.origin().ascii_serialization(); - let mut username = None; - let mut password = None; - for (key, value) in url.query_pairs() { - if key.eq("username") { - username = Some(value.into_owned()); - } else if key.eq("password") { - password = Some(value.into_owned()); - } - } - if username.is_some() || password.is_some() { - return Some(M3uUrlInfo { - base_url, - username: username.as_ref().unwrap().to_owned(), - password: username.as_ref().unwrap().to_owned(), - }); - } - } - None -} pub(crate) async fn serve_file(file_path: &Path, req: &HttpRequest) -> HttpResponse { if file_path.exists() { diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index c6037f6bd..e9d0a5e5c 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -6,7 +6,7 @@ use url::Url; use crate::api::api_utils::{get_user_target, serve_file}; use crate::api::api_model::{AppState, UserApiRequest}; use crate::model::api_proxy::ProxyType; -use crate::model::config::{Config, ConfigTarget, InputType}; +use crate::model::config::{Config, ConfigTarget}; use crate::model::config::TargetType; use crate::repository::m3u_repository::get_m3u_epg_file_path; use crate::repository::xtream_repository::{get_xtream_epg_file_path, get_xtream_storage_path}; @@ -54,7 +54,7 @@ async fn xmltv_api( // If you have multiple xtream sources, no response because of mapped ids // if you want epg for multi xtream input, then provide epg_url. let target_name = &target.name; - if let Some(inputs) = _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { + if let Some(inputs) = _app_state.config.get_inputs_for_target(target_name) { if inputs.len() == 1 { if let Some(&input) = inputs.first() { let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 2cdebab61..7b674dd0a 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -11,9 +11,8 @@ use url::Url; use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, serve_file}; use crate::api::api_model::{AppState, UserApiRequest, XtreamAuthorizationResponse, XtreamServerInfo, XtreamUserInfo}; -use crate::api::api_utils; use crate::model::api_proxy::{ProxyType, UserCredentials}; -use crate::model::config::{Config, ConfigInput, InputType}; +use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::config::{TargetType}; use crate::model::playlist::XtreamCluster; use crate::repository::xtream_repository; @@ -25,26 +24,15 @@ pub(crate) async fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) } fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { - match input.input_type { - InputType::M3u => { - match api_utils::parse_m3u_url(input.url.as_str()) { - None => None, - Some(m3u_url_info) => Some( - format!("{}/player_api.php?username={}&password={}&action={}", - m3u_url_info.base_url, - m3u_url_info.username, - m3u_url_info.password, - action - )) - } - } - InputType::Xtream => Some( - format!("{}/player_api.php?username={}&password={}&action={}", - input.url.as_str(), - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - action - )) + if let Some(user_info) = input.get_user_info() { + Some(format!("{}/player_api.php?username={}&password={}&action={}", + &user_info.base_url, + &user_info.username, + &user_info.password, + action + )) + } else { + None } } @@ -60,25 +48,16 @@ fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: &XtreamCluster, fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &str, action_path: &str) -> Option { let ctx_path = if context.is_empty() { "".to_string() } else { format!("{}/", context) }; - match input.input_type { - InputType::M3u => match api_utils::parse_m3u_url(input.url.as_str()) { - None => None, - Some(m3u_url_info) => Some( - format!("{}/{}{}/{}/{}", - m3u_url_info.base_url, - ctx_path, - m3u_url_info.username, - m3u_url_info.password, - action_path - )) - } - InputType::Xtream => Some(format!("{}/{}{}/{}/{}", - input.url.as_str(), - ctx_path, - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - action_path + if let Some(user_info) = input.get_user_info() { + Some(format!("{}/{}{}/{}/{}", + &user_info.base_url, + ctx_path, + &user_info.username, + &user_info.password, + action_path )) + } else { + None } } @@ -146,19 +125,17 @@ async fn xtream_player_api_stream( if target.has_output(&TargetType::Xtream) { let mut stream_id = action_path.to_owned(); let mut input: Option<&ConfigInput> = None; - let (action_stream_id, action_ext) = separate_number_and_rest(action_path); - if let Ok(num) = action_stream_id.trim().parse() { - let (xtream_id, cfg_input) = get_xtream_mapped_id_and_input_for_stream_id(_app_state, target_name, num); - if cfg_input.is_some() { - input = cfg_input; - stream_id = format!("{}{}", xtream_id, action_ext); - } - } - - if input.is_none() { - if let Some(m3u_inputs) = _app_state.config.get_input_for_target(target_name, &InputType::M3u) { - input = m3u_inputs.first().cloned(); + if target.is_multi_input() { + let (action_stream_id, action_ext) = separate_number_and_rest(action_path); + if let Ok(num) = action_stream_id.trim().parse() { + let (xtream_id, cfg_input) = get_xtream_mapped_id_and_input_for_stream_id(_app_state, target_name, num); + if cfg_input.is_some() { + input = cfg_input; + stream_id = format!("{}{}", xtream_id, action_ext); + } } + } else if let Some(inputs) = _app_state.config.get_inputs_for_target(target_name) { + input = inputs.first().copied(); } if let Some(target_input) = input { @@ -258,15 +235,13 @@ async fn xtream_player_api_timeshift_stream( xtream_player_api_stream(&req, &api_req, &_app_state, "timeshift", &username, &password, &action_path).await } - fn get_xtream_mapped_id_and_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>) { - if let Some(inputs) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { + if let Some(inputs) = app_state.config.get_inputs_for_target(target_name) { if let Ok(Some(mapping)) = xtream_repository::read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { if let Some(cfg_input) = inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { return (mapping.stream_id as i32, Some(cfg_input)); } } - return (stream_id, inputs.first().cloned()); } (stream_id, None) } @@ -303,23 +278,25 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ } async fn xtream_get_stream_info_response(app_state: &AppState, user: &UserCredentials, - target_name: &str, stream_id: &str, + target: &ConfigTarget, stream_id: &str, cluster: &XtreamCluster) -> HttpResponse { - let xtream_stream_id: i32 = match FromStr::from_str(stream_id) { + let req_stream_id: i32 = match FromStr::from_str(stream_id) { Ok(id) => id, Err(_) => return HttpResponse::BadRequest().finish() }; - if user.proxy == ProxyType::Redirect && !app_state.config.is_multi_xtream_input(target_name) { - let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, xtream_stream_id); - if let Some(target_input) = input { - if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { - return HttpResponse::Found().insert_header(("Location", info_url)).finish(); + if user.proxy == ProxyType::Redirect && !target.is_multi_input() { + if let Some(inputs) = app_state.config.get_inputs_for_target(&target.name) { + if let Some(&input) = inputs.first() { + if let Some(info_url) = get_xtream_player_api_info_url(input, cluster, req_stream_id) { + return HttpResponse::Found().insert_header(("Location", info_url)).finish(); + } } } + return HttpResponse::BadRequest().finish(); } - match xtream_get_stream_info(app_state, target_name, xtream_stream_id, cluster).await { + match xtream_get_stream_info(app_state, &target.name, req_stream_id, cluster).await { Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), Err(_) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{info:[]}"), } @@ -380,12 +357,12 @@ async fn xtream_player_api( match action { "get_series_info" => { - xtream_get_stream_info_response(_app_state, &user, target_name, + xtream_get_stream_info_response(_app_state, &user, target, api_req.series_id.trim(), &XtreamCluster::Series).await } "get_vod_info" => { - xtream_get_stream_info_response(_app_state, &user, target_name, + xtream_get_stream_info_response(_app_state, &user, target, api_req.vod_id.trim(), &XtreamCluster::Video).await } diff --git a/src/model/config.rs b/src/model/config.rs index b69263255..3e607383c 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -5,6 +5,7 @@ use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::str::FromStr; use std::sync::{Arc, RwLock}; +use url::Url; use log::{debug, error, warn}; use path_absolutize::*; @@ -275,6 +276,8 @@ pub(crate) struct ConfigTargetOptions { pub xtream_skip_live_direct_source: bool, #[serde(default = "default_as_true")] pub xtream_skip_video_direct_source: bool, + #[serde(default = "default_as_true")] + pub xtream_skip_series_direct_source: bool, #[serde(default = "default_as_false")] pub xtream_resolve_series: bool, #[serde(default = "default_as_two")] @@ -318,6 +321,8 @@ pub(crate) struct ConfigTarget { pub _filter: Option, #[serde(skip_serializing, skip_deserializing)] pub _mapping: Option>, + #[serde(skip_serializing, skip_deserializing)] + _multi_input: bool, } @@ -387,6 +392,11 @@ impl ConfigTarget { Err(err) => Err(err), } } + + pub(crate) fn is_multi_input(&self) -> bool { + self._multi_input + } + pub(crate) fn filter(&self, provider: &ValueProvider) -> bool { let mut processor = MockValueProcessor {}; return self._filter.as_ref().unwrap().filter(provider, &mut processor); @@ -417,23 +427,23 @@ impl ConfigTarget { pub(crate) struct ConfigSource { pub inputs: Vec, pub targets: Vec, - #[serde(skip_serializing, skip_deserializing)] - pub _multi_xtream_input: bool, } impl ConfigSource { pub(crate) fn prepare(&mut self, index: u16) -> Result { handle_m3u_filter_error_result_list!(M3uFilterErrorKind::Info, self.inputs.iter_mut().enumerate().map(|(idx, i)| i.prepare(index+(idx as u16)))); - self._multi_xtream_input = self.inputs.iter().filter(|i| i.input_type == InputType::Xtream).count() > 1; + if self.inputs.len() > 1 { + self.targets.iter_mut().for_each(|t| t._multi_input = true); + } Ok(index + (self.inputs.len() as u16)) } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option> { + pub(crate) fn get_inputs_for_target(&self, target_name: &str) -> Option> { let mut result = Vec::new(); for target in &self.targets { if target.name.eq(target_name) { for input in &self.inputs { - if input.enabled && input.input_type.eq(input_type) { + if input.enabled { result.push(input); } } @@ -494,6 +504,12 @@ pub(crate) struct ConfigInputOptions { } +pub(crate) struct InputUserInfo { + pub base_url: String, + pub username: String, + pub password: String, +} + fn default_as_type_m3u() -> InputType { InputType::M3u } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -561,6 +577,37 @@ impl ConfigInput { } Ok(()) } + + pub(crate) fn get_user_info(&self) -> Option { + if self.input_type == InputType::Xtream { + if self.username.is_some() || self.password.is_some() { + return Some(InputUserInfo { + base_url: self.url.to_owned(), + username: self.username.as_ref().unwrap().to_owned(), + password: self.password.as_ref().unwrap().to_owned(), + }); + } + } else if let Ok(url) = Url::parse(&self.url) { + let base_url = url.origin().ascii_serialization(); + let mut username = None; + let mut password = None; + for (key, value) in url.query_pairs() { + if key.eq("username") { + username = Some(value.into_owned()); + } else if key.eq("password") { + password = Some(value.into_owned()); + } + } + if username.is_some() || password.is_some() { + return Some(InputUserInfo { + base_url, + username: username.as_ref().unwrap().to_owned(), + password: password.as_ref().unwrap().to_owned(), + }); + } + } + None + } } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -731,9 +778,9 @@ impl Config { } } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option> { + pub(crate) fn get_inputs_for_target(&self, target_name: &str) -> Option> { for source in &self.sources { - if let Some(cfg) = source.get_input_for_target(target_name, input_type) { + if let Some(cfg) = source.get_inputs_for_target(target_name) { return Some(cfg); } } @@ -769,16 +816,16 @@ impl Config { None } - pub(crate) fn is_multi_xtream_input(&self, target_name: &str) -> bool { - for source in &self.sources { - for target in &source.targets { - if target_name.eq_ignore_ascii_case(&target.name) { - return source._multi_xtream_input; - } - } - } - false - } + // pub(crate) fn is_multi_input_target(&self, target_name: &str) -> bool { + // for source in &self.sources { + // for target in &source.targets { + // if target_name.eq_ignore_ascii_case(&target.name) { + // return target.is_multi_input(); + // } + // } + // } + // false + // } pub(crate) fn set_mappings(&mut self, mappings: Option) -> Result<(), M3uFilterError> { if let Some(mapping_list) = mappings { diff --git a/src/model/xtream.rs b/src/model/xtream.rs index 3298db58b..c403c379c 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -279,8 +279,8 @@ pub(crate) struct XtreamSeriesInfoEpisodeInfo { pub duration_secs: u32, pub duration: String, pub movie_image: String, - // "video": [], - // "audio": [], + pub video: Value, + pub audio: Value, pub bitrate: u32, pub rating: f64, pub season: u32, diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index ee7c04af3..508348b00 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -14,7 +14,7 @@ use crate::{Config, get_errors_notify_message, model::config, valid_property}; use crate::filter::{get_field_value, MockValueProcessor, set_field_value, ValueProvider}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::messaging::{MsgKind, send_message}; -use crate::model::config::{ConfigSource, ConfigTarget, default_as_default, InputAffix, InputType, ProcessTargets}; +use crate::model::config::{ConfigTarget, default_as_default, InputAffix, InputType, ProcessTargets}; use crate::model::mapping::{Mapping, MappingValueProcessor}; use crate::model::config::{AFFIX_FIELDS, ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType}; use crate::model::playlist::{FetchedPlaylist, FieldAccessor, PlaylistGroup, PlaylistItem, PlaylistItemHeader}; @@ -258,7 +258,7 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option id: grp_id, title: Rc::clone(title), channels: vec![channel.clone()], - xtream_cluster: cluster.clone() + xtream_cluster: cluster.clone(), }) } } @@ -350,7 +350,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

{} Err(mut err) => err.drain(..).for_each(|e| errors.push(e)) } @@ -418,7 +418,6 @@ fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe { } pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], - source: &ConfigSource, target: &ConfigTarget, cfg: &Config, stats: &mut HashMap, errors: &mut Vec) -> Result<(), Vec> { @@ -456,7 +455,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], apply_affixes(&mut new_fetched_playlists); - if source._multi_xtream_input { + if target.is_multi_input() { let mut stream_id_mappings: Vec = Vec::new(); let mut counter: u32 = 0; new_fetched_playlists.iter()