diff --git a/CHANGELOG.md b/CHANGELOG.md index 0395b7d80..186ab0629 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,9 @@ # Changelog # v1.1.4(2023-11-??) * Added regexp search in Web-UI +* Added config Web-UI +* Added xtream vod_info and series_info +* Added input options with attribute xtream_info_cache to cache get_vod_info and get_series_info on disc # v1.1.3(2023-11-08) * added new target options diff --git a/Cargo.lock b/Cargo.lock index 76047a8ac..411b2304d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -436,6 +436,12 @@ version = "3.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f30e7476521f6f8af1a1c4c0b8cc94f0bee37d91763d0ca2665f299b6cd8aec" +[[package]] +name = "byteorder" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" + [[package]] name = "bytes" version = "1.5.0" @@ -584,6 +590,21 @@ dependencies = [ "libc", ] +[[package]] +name = "crc" +version = "3.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "86ec7a15cbe22e59248fc7eadb1907dab5ba09372595da4d73dd805ed4417dfe" +dependencies = [ + "crc-catalog", +] + +[[package]] +name = "crc-catalog" +version = "2.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "19d374276b40fb8bbdee95aef7c7fa6b5316ec764510eb64b8dd0e2ed0d7e7f5" + [[package]] name = "crc32fast" version = "1.3.2" @@ -1293,6 +1314,16 @@ version = "0.4.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b5e6163cb8c49088c2c36f57875e58ccd8c87c7427f7fbd50ea6710b2f3f2e8f" +[[package]] +name = "lzma-rs" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "297e814c836ae64db86b36cf2a557ba54368d03f6afcd7d947c266692f71115e" +dependencies = [ + "byteorder", + "crc", +] + [[package]] name = "m3u-filter" version = "1.1.3" @@ -1310,6 +1341,7 @@ dependencies = [ "env_logger", "futures", "log", + "lzma-rs", "mime", "openssl", "path-absolutize", diff --git a/Cargo.toml b/Cargo.toml index 0f637c8e3..e3dba027c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -42,3 +42,4 @@ env_logger = "0.10" rustelebot = "0.3" bincode = "1.3" uuid = { version = "1.3.0", features = ["v4", "fast-rng", "macro-diagnostics"] } +lzma-rs = "0.3.0" \ No newline at end of file diff --git a/README.md b/README.md index c0b02310c..0685ca69f 100644 --- a/README.md +++ b/README.md @@ -99,6 +99,8 @@ Each input has the following attributes: - `pasword`only mandatory for type `xtream` - `prefix` is optional, it is applied to the given field with the given value - `suffix` is optional, it is applied to the given field with the given value + - `options` is optional, + + `xtream_info_cache` true or false, vod_info and series_info can be cached to disc to reduce network traffic to provider. `persist` should be different for `m3u` and `xtream` types. For `m3u` use full filename like `./playlist_{}.m3u`. For `xtream` use a prefix like `./playlist_` diff --git a/src/api/api_model.rs b/src/api/api_model.rs index 8e098644f..64c61e73e 100644 --- a/src/api/api_model.rs +++ b/src/api/api_model.rs @@ -1,4 +1,4 @@ -use std::collections::VecDeque; +use std::collections::{HashMap, VecDeque}; use std::ffi::OsStr; use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex, RwLock}; @@ -87,6 +87,34 @@ impl FileDownload { } } + +pub(crate) struct SharedLocks { + shared_locks: Arc>>>>, +} + +impl SharedLocks { + pub(crate) fn new() -> Self { + SharedLocks { + shared_locks: Arc::new(RwLock::new(HashMap::new())), + } + } + + pub(crate) fn get_lock(&self, key: &str) -> Arc> { + { + let lock = self.shared_locks.read().unwrap(); + match lock.get(key) { + None => {} + Some(file) => return file.clone(), + } + } + let mut lock = self.shared_locks.write().unwrap(); + let file = Arc::new(RwLock::new(key.to_string() )); + lock.insert(key.to_string(), file.clone()); + file + } +} + + pub(crate) struct DownloadQueue { pub queue: Arc>>, pub active: Arc>>, @@ -97,6 +125,7 @@ pub(crate) struct AppState { pub config: Arc, pub targets: Arc, pub downloads: Arc, + pub shared_locks: Arc, } #[derive(Serialize)] @@ -153,6 +182,10 @@ pub(crate) struct UserApiRequest { pub series_id: String, #[serde(default = "default_as_empty_str")] pub vod_id: String, + #[serde(default = "default_as_empty_str")] + pub stream_id: String, + #[serde(default = "default_as_empty_str")] + pub limit: String, } #[derive(Deserialize, Serialize, Debug, Clone)] diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index 9c2203a42..52cb9f263 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -29,6 +29,8 @@ async fn m3u_api( pub(crate) fn m3u_api_register() -> Vec { vec![ web::resource("/get.php").route(web::get().to(m3u_api)), + web::resource("/get.php").route(web::post().to(m3u_api)), + web::resource("/apiget").route(web::get().to(m3u_api)), web::resource("/m3u").route(web::get().to(m3u_api)) ] } \ No newline at end of file diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 4495979a9..bc18d6136 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -9,7 +9,7 @@ use actix_web::{App, get, HttpRequest, HttpServer, web}; use actix_web::middleware::Logger; use crate::api::m3u_api::{m3u_api_register}; -use crate::api::api_model::{AppState, DownloadQueue}; +use crate::api::api_model::{AppState, DownloadQueue, SharedLocks}; use crate::api::scheduler::start_scheduler; use crate::api::v1_api::{v1_api_register}; use crate::api::xmltv_api::{xmltv_api_register}; @@ -46,7 +46,8 @@ pub(crate) async fn start_server(cfg: Arc, targets: Arc) queue: Arc::from(Mutex::new(VecDeque::new())), active: Arc::from(RwLock::new(None)), finished: Arc::from(RwLock::new(Vec::new())), - }) + }), + shared_locks: Arc::new(SharedLocks::new()), }); // Scheduler diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index c91decb5d..6da7f91c4 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -3,7 +3,7 @@ use actix_web::{HttpResponse, Scope, web}; use serde_json::{json}; use crate::api::api_model::{AppState, PlaylistRequest, ServerConfig, ServerInputConfig, ServerSourceConfig, ServerTargetConfig}; use crate::download::{get_m3u_playlist, get_xtream_playlist}; -use crate::model::config::{ConfigInput, InputType, validate_targets}; +use crate::model::config::{ConfigInput, ConfigInputOptions, InputType, validate_targets}; use log::{error}; use crate::api::download_api::{download_file_info, queue_download_file}; use crate::config_reader::save_api_proxy; @@ -104,6 +104,9 @@ fn create_config_input_for_url(url: &str) -> ConfigInput { suffix: None, name: None, enabled: true, + options: Some(ConfigInputOptions { + xtream_info_cache: false, + }) } } diff --git a/src/api/xtream_player_api.rs b/src/api/xtream_player_api.rs index 7e1a9c130..5a7a4f3b1 100644 --- a/src/api/xtream_player_api.rs +++ b/src/api/xtream_player_api.rs @@ -11,7 +11,10 @@ use crate::api::api_model::{AppState, UserApiRequest, XtreamAuthorizationRespons use crate::model::api_proxy::{UserCredentials}; use crate::model::config::{Config}; use crate::model::model_config::{TargetType}; -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_series_info, xtream_get_vod_info}; +use crate::model::model_m3u::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_short_epg, xtream_get_stream_info}; use crate::utils::get_client_request; fn get_user_info(user: &UserCredentials, cfg: &Config) -> XtreamAuthorizationResponse { @@ -50,9 +53,9 @@ async fn xtream_player_api_stream( context: &str, username: &str, password: &str, - stream_id: &str, + action_path: &str, ) -> HttpResponse { - if let Some((_user, target)) = get_user_target_by_credentials(&username, &password, api_req, _app_state) { + 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) { match _app_state.config.get_xtream_input_for_target(target_name) { @@ -60,7 +63,7 @@ async fn xtream_player_api_stream( Some(input) => { let username = input.username.as_ref().unwrap().clone(); let password = input.password.as_ref().unwrap().clone(); - let stream_url = format!("{}/{}/{}/{}/{}", input.url, context, username, password, stream_id); + let stream_url = format!("{}/{}/{}/{}/{}", input.url, context, username, password, action_path); let url = reqwest::Url::parse(&stream_url).unwrap(); let client = get_client_request(input, url); if let Ok(response) = client.send().await { @@ -102,6 +105,28 @@ async fn xtream_player_api_movie_stream( xtream_player_api_stream(&api_req, &_app_state, "movie", &username, &password, &stream_id).await } +async fn xtream_player_api_timeshift_stream( + api_req: web::Query, + path: web::Path<(String, String, String, String, String)>, + _app_state: web::Data, +) -> HttpResponse { + let (username, password, duration, start, stream_id) = path.into_inner(); + let action_path = format!("{}/{}/{}", duration, start, stream_id); + xtream_player_api_stream(&api_req, &_app_state, "timeshift", &username, &password, &action_path).await +} + +async fn xtream_get_stream_info_response(app_state: &AppState, target_name: &str, stream_id: &str, cluster: XtreamCluster, + user: &UserCredentials) -> HttpResponse { + match FromStr::from_str(stream_id) { + Ok(xtream_stream_id) => { + match xtream_get_stream_info(app_state, target_name, xtream_stream_id, cluster, user).await { + Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), + Err(_) => HttpResponse::NoContent().finish() + } + } + Err(_) => HttpResponse::BadRequest().finish() + } +} async fn xtream_player_api( api_req: web::Query, @@ -119,25 +144,19 @@ async fn xtream_player_api( match action { "get_series_info" => { - match FromStr::from_str(api_req.series_id.trim()) { - Ok(stream_id) => { - match xtream_get_series_info(&_app_state.config, target_name, stream_id) { - Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), - Err(_) => HttpResponse::NoContent().finish() - } - } - Err(_) => HttpResponse::BadRequest().finish() - } + xtream_get_stream_info_response(&_app_state, target_name, + api_req.series_id.trim(), + XtreamCluster::Series, &user).await } "get_vod_info" => { - match FromStr::from_str(api_req.vod_id.trim()) { - Ok(stream_id) => { - match xtream_get_vod_info(&_app_state.config, target_name, stream_id) { - Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), - Err(_) => HttpResponse::NoContent().finish() - } - } - Err(_) => HttpResponse::BadRequest().finish() + xtream_get_stream_info_response(&_app_state, target_name, + api_req.vod_id.trim(), + XtreamCluster::Video, &user).await + } + "get_short_epg" => { + match xtream_get_short_epg(&_app_state, target_name, api_req.stream_id.trim(), api_req.limit.trim(), &user).await { + Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), + Err(_) => HttpResponse::NoContent().finish() } } _ => { @@ -187,9 +206,16 @@ async fn xtream_player_api( pub(crate) fn xtream_api_register() -> Vec { vec![ web::resource("/player_api.php").route(web::get().to(xtream_player_api)), + web::resource("/player_api.php").route(web::post().to(xtream_player_api)), web::resource("/xtream").route(web::get().to(xtream_player_api)), web::resource("/live/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_live_stream)), web::resource("/movie/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_movie_stream)), web::resource("/series/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_series_stream)), + web::resource("/timeshift/{username}/{password}/{duration}/{start}{stream_id}").route(web::get().to(xtream_player_api_timeshift_stream)), + /* TODO + web::resource("/hlsr/{token}/{username}/{password}/{channel}/{hash}/{chunk}").route(web::get().to(xtream_player_api_hlsr_stream)) + web::resource("/hls/{token}/{chunk}").route(web::get().to(xtream_player_api_hls_stream)) + web::resource("/play/{token}/{type}").route(web::get().to(xtream_player_api_play_stream)) + */ ] } \ No newline at end of file diff --git a/src/model/config.rs b/src/model/config.rs index b0542b1db..ad7f0c477 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -341,6 +341,13 @@ impl FromStr for InputType { } } +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub(crate) struct ConfigInputOptions { + #[serde(default = "default_as_false")] + pub xtream_info_cache: bool, +} + + fn default_as_type_m3u() -> InputType { InputType::M3u } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -368,6 +375,9 @@ pub(crate) struct ConfigInput { pub name: Option, #[serde(default = "default_as_true")] pub enabled: bool, + #[serde(skip_serializing_if = "Option::is_none")] + pub options: Option, + } impl ConfigInput { diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs index b873e07ab..c43e9cb64 100644 --- a/src/processing/playlist_watch.rs +++ b/src/processing/playlist_watch.rs @@ -23,7 +23,7 @@ pub(crate) fn process_group_watch(cfg: &Config, target_name: &str, pl: &Playlist let save_path = path.clone(); let mut changed = false; if path.exists() { - match load_tree(&path) { + match load_watch_tree(&path) { Some(loaded_tree) => { // Find elements in set2 but not in set1 let added_difference: BTreeSet = new_tree.difference(&loaded_tree).cloned().collect(); @@ -39,7 +39,7 @@ pub(crate) fn process_group_watch(cfg: &Config, target_name: &str, pl: &Playlist } } if changed { - match save_tree(&save_path, new_tree) { + match save_watch_tree(&save_path, new_tree) { Ok(_) => {} Err(err) => { error!("failed to write watch_file {}: {}", &save_path.to_str().unwrap(), err) @@ -76,7 +76,7 @@ fn handle_watch_notification(cfg: &Config, added: BTreeSet, removed: BTr } } -fn load_tree(path: &Path) -> Option> { +fn load_watch_tree(path: &Path) -> Option> { match std::fs::read(path) { Ok(encoded) => { let decoded: BTreeSet = bincode::deserialize(&encoded[..]).unwrap(); @@ -86,7 +86,7 @@ fn load_tree(path: &Path) -> Option> { } } -fn save_tree(path: &Path, tree: BTreeSet) -> std::io::Result<()> { +fn save_watch_tree(path: &Path, tree: BTreeSet) -> std::io::Result<()> { let encoded: Vec = bincode::serialize(&tree).unwrap(); std::fs::write(path, encoded) } diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index aa4c33bc3..e2d09f935 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,16 +1,23 @@ use std::cell::Ref; use std::collections::{BTreeMap, HashMap}; use std::fs; -use std::fs::File; +use std::fs::{File, OpenOptions}; use std::io::{BufReader, BufWriter, Error, Read, Seek, SeekFrom, Write}; use std::iter::FromIterator; use std::path::{Path, PathBuf}; +use log::error; use serde::Serialize; use serde_json::{json, Map, Value}; use crate::model::config::{Config, ConfigTarget}; use crate::model::model_m3u::{PlaylistGroup, PlaylistItemHeader, XtreamCluster}; use crate::{create_m3u_filter_error_result, utils}; +use crate::api::api_model::AppState; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::api_proxy::UserCredentials; +use crate::utils::{get_client_request}; + +type IndexTree = BTreeMap; + pub(crate) static COL_CAT_LIVE: &str = "cat_live"; pub(crate) static COL_CAT_SERIES: &str = "cat_series"; @@ -34,6 +41,27 @@ const SERIES_STREAM_FIELDS: &[&str] = &[ "stream_type", "title", "year", "youtube_trailer", ]; + +pub(crate) fn get_xtream_storage_path(cfg: &Config, target_name: &str) -> Option { + utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) +} + +pub(crate) fn get_xtream_epg_file_path(path: &Path) -> PathBuf { + path.join("epg.xml") +} + +fn get_collection_path(path: &Path, collection: &str) -> PathBuf { + path.join(format!("{}.json", collection)) +} + +fn get_info_collection_path(path: &Path, collection: &str) -> PathBuf { + path.join(format!("{}_info.db", collection)) +} + +fn get_info_idx_path(path: &Path, collection: &str) -> PathBuf { + path.join(format!("{}_info.idx", collection)) +} + fn write_to_file(file: &Path, value: &T) -> Result<(), Error> where T: ?Sized + Serialize { @@ -50,68 +78,40 @@ fn write_to_file(file: &Path, value: &T) -> Result<(), Error> } } -fn get_collection_and_idx_path(path: &Path, cluster: &XtreamCluster) -> (PathBuf, PathBuf) { +fn get_info_collection_and_idx_path(path: &Path, cluster: &XtreamCluster) -> (PathBuf, PathBuf) { let collection = match cluster { XtreamCluster::Live => COL_LIVE, XtreamCluster::Video => COL_VOD, XtreamCluster::Series => COL_SERIES, }; - (get_collection_path(path, collection), get_idx_path(path, collection)) + (get_info_collection_path(path, collection), get_info_idx_path(path, collection)) } -fn write_to_file_width_idx(path: &Path, values: &[(i32, Value)], cluster: &XtreamCluster) -> Result<(), Error> { - let (file, file_idx) = get_collection_and_idx_path(path, cluster); - match File::create(file) { - Ok(file) => { - let mut index = BTreeMap::::new(); - let mut writer = BufWriter::new(file); - writer.write_all("[".as_bytes())?; - let mut offset = 1; - let value_cnt = values.len(); - let mut value_idx = 0; - for (stream_id, data) in values { - let content = serde_json::to_string(data).unwrap(); - let bytes = content.as_bytes(); - let size = bytes.len(); - index.insert(*stream_id, (offset as u32, size as u16)); - offset += size; - let _ = writer.write_all(bytes); - value_idx += 1; - if value_idx < value_cnt { - writer.write_all(",".as_bytes())?; - offset += 1; - } - } - writer.write_all("]".as_bytes())?; - match writer.flush() { - Ok(_) => { - let encoded: Vec = bincode::serialize(&index).unwrap(); - let _ = fs::write(file_idx, encoded); - Ok(()) - } - Err(e) => Err(e) +fn write_xtream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster, content: &String, index_tree: &mut IndexTree) -> Result<(), Error> { + if let Some(path) = get_xtream_storage_path(&app_state.config, target_name) { + let (col_path, idx_path) = get_info_collection_and_idx_path(&path, cluster); + let mut comp: Vec = Vec::new(); + lzma_rs::lzma_compress(&mut BufReader::new(content.as_bytes()), &mut comp)?; + let size = comp.len(); + let lock = app_state.shared_locks.get_lock(target_name); + let shared_lock = lock.write().unwrap(); + match OpenOptions::new() + .create(true) + .write(true) + .append(true) + .open(col_path) { + Ok(mut file) => { + let offset = file.metadata().unwrap().len(); + file.write_all(comp.as_slice())?; + file.flush()?; + index_tree.insert(stream_id, (offset as u32, size as u16)); + write_index(&idx_path, index_tree)?; + drop(shared_lock); } + Err(err) => return Err(err) } - Err(e) => Err(e) } -} - - -pub(crate) fn get_xtream_storage_path(cfg: &Config, target_name: &str) -> Option { - utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) -} - -fn get_collection_path(path: &Path, collection: &str) -> PathBuf { - path.join(format!("{}.json", collection)) -} - -fn get_idx_path(path: &Path, collection: &str) -> PathBuf { - path.join(format!("{}.idx", collection)) -} - - -pub(crate) fn get_xtream_epg_file_path(path: &Path) -> PathBuf { - path.join("epg.xml") + Ok(()) } pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlist: &[PlaylistGroup]) -> Result<(), M3uFilterError> { @@ -221,7 +221,7 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis XtreamCluster::Live => &mut live_col, XtreamCluster::Series => &mut series_col, XtreamCluster::Video => &mut vod_col, - }.push((stream_id, Value::Object(document))); + }.push(Value::Object(document)); } } } @@ -230,7 +230,10 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis for (col_path, data) in [ (get_collection_path(&path, COL_CAT_LIVE), &cat_live_col), (get_collection_path(&path, COL_CAT_VOD), &cat_vod_col), - (get_collection_path(&path, COL_CAT_SERIES), &cat_series_col)] { + (get_collection_path(&path, COL_CAT_SERIES), &cat_series_col), + (get_collection_path(&path, COL_LIVE), &live_col), + (get_collection_path(&path, COL_VOD), &vod_col), + (get_collection_path(&path, COL_SERIES), &series_col)] { match write_to_file(&col_path, data) { Ok(()) => {} Err(err) => { @@ -238,17 +241,6 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis } } } - for (data, cluster) in [ - (&live_col, XtreamCluster::Live), - (&vod_col, XtreamCluster::Video), - (&series_col, XtreamCluster::Series)] { - match write_to_file_width_idx(&path, data, &cluster) { - Ok(()) => {} - Err(err) => { - errors.push(format!("Persisting collection failed: {}: {}", cluster, err)); - } - } - } if !errors.is_empty() { return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", errors.join("\n")); } @@ -315,16 +307,21 @@ pub(crate) fn xtream_get_all(cfg: &Config, target_name: &str, collection_name: & Err(Error::new(std::io::ErrorKind::Other, format!("Cant find collection: {}/{}", target_name, collection_name))) } -fn load_map(path: &Path) -> Option> { - match std::fs::read(path) { +fn load_index(path: &Path) -> Option { + match fs::read(path) { Ok(encoded) => { - let decoded: BTreeMap = bincode::deserialize(&encoded[..]).unwrap(); + let decoded: IndexTree = bincode::deserialize(&encoded[..]).unwrap(); Some(decoded) } Err(_) => None, } } +fn write_index(path: &PathBuf, index: &IndexTree) -> std::io::Result<()> { + let encoded = bincode::serialize(index).unwrap(); + fs::write(path, encoded) +} + fn seek_read( reader: &mut (impl Read + Seek), offset: u32, @@ -337,15 +334,68 @@ fn seek_read( Ok(buf) } -fn xtream_get_stream_info(cfg: &Config, target_name: &str, stream_id: i32, cluster: XtreamCluster) -> Result { - if let Some(path) = get_xtream_storage_path(cfg, target_name) { - let (col_path, idx_path) = get_collection_and_idx_path(&path, &cluster); - if idx_path.exists() && col_path.exists() { - if let Some(idx_map) = load_map(&idx_path) { - 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) { - return Ok(String::from_utf8(bytes).unwrap()); +pub(crate) async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, + cluster: XtreamCluster, + user: &UserCredentials) -> Result { + let target_input = app_state.config.get_xtream_input_for_target(target_name); + let cache_info = target_input.and_then(|i| i.options.as_ref()) + .map(|o| o.xtream_info_cache).unwrap_or(false); + let mut index_tree: Option = None; + if cache_info { + if let Some(path) = get_xtream_storage_path(&app_state.config, target_name) { + let (col_path, idx_path) = get_info_collection_and_idx_path(&path, &cluster); + let lock = app_state.shared_locks.get_lock(target_name); + let shared_lock = lock.read().unwrap(); + if idx_path.exists() && col_path.exists() { + index_tree = load_index(&idx_path); + 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) { + let mut decomp: Vec = Vec::new(); + let _ = lzma_rs::lzma_decompress(&mut bytes.as_slice(), &mut decomp); + drop(shared_lock); + return Ok(String::from_utf8(decomp).unwrap()); + } + } + } + } + drop(shared_lock); + } + } + + let (action, stream_id_field) = match cluster { + XtreamCluster::Live => ("get_live_info", "live_id"), + XtreamCluster::Video => ("get_vod_info", "vod_id"), + XtreamCluster::Series => ("get_series_info", "series_id"), + }; + + // not indexed, receive + match target_input { + None => {} + Some(input) => { + let info_url = format!("{}/player_api.php?username={}&password={}&action={}&{}={}", + input.url, &user.username, &user.password, action, + stream_id_field, stream_id); + let url = reqwest::Url::parse(&info_url).unwrap(); + let client = get_client_request(input, url); + if let Ok(response) = client.send().await { + if response.status().is_success() { + match response.text().await { + Ok(content) => { + if cache_info { + if index_tree.is_none() { + index_tree = Some(IndexTree::new()); + } + match write_xtream_info(app_state, target_name, stream_id, &cluster, &content, + index_tree.as_mut().unwrap()) { + Ok(_) => {} + Err(err) => { error!("{}", err.to_string()); } + } + } + return Ok(content); + } + Err(err) => { error!("Failed to download info {}", err.to_string()); } } } } @@ -354,93 +404,28 @@ fn xtream_get_stream_info(cfg: &Config, target_name: &str, stream_id: i32, clust Err(Error::new(std::io::ErrorKind::Other, format!("Cant find stream with id: {}/{}/{}", target_name, &cluster, stream_id))) } -pub(crate) fn xtream_get_series_info(cfg: &Config, target_name: &str, stream_id: i32) -> Result { - /* - { - "episodes": { - "": [ - { - "added": string, - "container_extension": string, - "custom_sid": string, - "direct_source": string, - "episode_num": int, - "id": string, - "info": { - "bitrate": int, - "duration": string, - "duration_secs": int, - "movie_image": string, - "name": string, - "plot": string, - "rating": float, - "releasedate": string, - "audio": FFMPEGStreamInfo, - "video": FFMPEGStreamInfo +pub(crate) async fn xtream_get_short_epg(app_state: &AppState, target_name: &str, stream_id: &str, limit: &str, user: &UserCredentials) -> Result { + match app_state.config.get_xtream_input_for_target(target_name) { + None => {} + Some(input) => { + let mut info_url = format!("{}/player_api.php?username={}&password={}&action=get_short_epg&stream_id={}", + input.url, &user.username, &user.password, stream_id); + if !(limit.is_empty() || limit.eq("0")) { + info_url = format!("{}&limit={}", info_url, limit); } - "season": int, - "title": string - } - ] - }, - "info": { - "backdrop_path: [string], - "cast": string, - "category_id": string, - "cover": string, - "director": string, - "episode_run_time": string, - "genre": string, - "last_modified": string, - "name": string, - "num": int, - "plot": string, - "rating, string, - "rating_5based": float, - "releaseDate": string, - "series_id": int, - "stream_type": string, - "youtube_trailer": string, + let url = reqwest::Url::parse(&info_url).unwrap(); + let client = get_client_request(input, url); + if let Ok(response) = client.send().await { + if response.status().is_success() { + match response.text().await { + Ok(content) => { + return Ok(content); + } + Err(err) => { error!("Failed to download epg {}", err.to_string()); } + } + } + } + } } - } - "seasons": [] -} - */ - // TODO restructure - xtream_get_stream_info(cfg, target_name, stream_id, XtreamCluster::Series) -} - -pub(crate) fn xtream_get_vod_info(cfg: &Config, target_name: &str, stream_id: i32) -> Result { - /* - { - "info": { - "backdrop_path": [string], - "bitrate": FlexInt, - "cast": string, - "director": string, - "duration": string, - "duration_secs": FlexInt, - "genre": string, - "movie_image": string, - "plot": string, - "rating": FlexFloat, - "releasedate": string, - "tmdb_id": int, - "youtube_trailer": string, - "audio": FFMPEGStreamInfo, - "video": FFMPEGStreamInfo, - } `json:"info"` - "movie_data": { - "added": string, - "category_id": string, - "container_extension": string, - "custom_sid": string, - "direct_source": string, - "name": string, - "stream_id": int - } - } - */ - // TODO restructure - xtream_get_stream_info(cfg, target_name, stream_id, XtreamCluster::Video) + Err(Error::new(std::io::ErrorKind::Other, format!("Cant find short epg with id: {}/{}", target_name, stream_id))) } diff --git a/src/utils.rs b/src/utils.rs index 19d12f42c..d82e117eb 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -124,13 +124,13 @@ pub(crate) async fn get_input_text_content(input: &ConfigInput, working_dir: &St let mut content = String::new(); match std::io::BufReader::new(file).read_to_string(&mut content) { Ok(_) => Some(content), - Err(err) => { + Err(err) => { let file_str = &filepath.to_str().unwrap_or("?"); error!("cant read file: {} {}", file_str, err); - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "Cant open file : {} => {}", file_str, err) + return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "Cant open file : {} => {}", file_str, err); } } - }, + } Err(err) => { let file_str = &filepath.to_str().unwrap_or("?"); error!("cant read file: {} {}", file_str, err); @@ -204,6 +204,16 @@ pub(crate) fn get_client_request(input: &ConfigInput, url: url::Url) -> reqwest: } request } +// +// pub(crate) fn get_client_request_sync(input: &ConfigInput, url: url::Url) -> reqwest::blocking::RequestBuilder { +// let mut request = reqwest::blocking::Client::new().get(url); +// if input.headers.is_empty() { +// let headers = get_request_headers(&input.headers); +// request = request.headers(headers); +// } +// request +// } + pub(crate) fn get_request_headers(defined_headers: &HashMap) -> HeaderMap { let mut headers = header::HeaderMap::new(); @@ -278,7 +288,7 @@ pub(crate) fn bytes_to_megabytes(bytes: u64) -> u64 { pub(crate) fn add_prefix_to_filename(path: &Path, prefix: &str, ext: Option<&str>) -> PathBuf { let file_name = path.file_name().unwrap_or_default(); let new_file_name = format!("{}{}", prefix, file_name.to_string_lossy()); - let result = path.with_file_name(new_file_name); + let result = path.with_file_name(new_file_name); match ext { None => result, Some(extension) => result.with_extension(extension) @@ -287,7 +297,7 @@ pub(crate) fn add_prefix_to_filename(path: &Path, prefix: &str, ext: Option<&str pub(crate) fn path_exists(file_path: &Path) -> bool { if let Ok(metadata) = fs::metadata(file_path) { - return metadata.is_file() + return metadata.is_file(); } false }