refactored xtream api

This commit is contained in:
euzu
2023-11-24 22:53:33 +01:00
parent fd69a835b6
commit 7298c28572
13 changed files with 307 additions and 199 deletions
+3
View File
@@ -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
Generated
+32
View File
@@ -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",
+1
View File
@@ -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"
+2
View File
@@ -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_`
+34 -1
View File
@@ -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<RwLock<HashMap<String, Arc<RwLock<String>>>>>,
}
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<RwLock<String>> {
{
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<Mutex<VecDeque<FileDownload>>>,
pub active: Arc<RwLock<Option<FileDownload>>>,
@@ -97,6 +125,7 @@ pub(crate) struct AppState {
pub config: Arc<Config>,
pub targets: Arc<ProcessTargets>,
pub downloads: Arc<DownloadQueue>,
pub shared_locks: Arc<SharedLocks>,
}
#[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)]
+2
View File
@@ -29,6 +29,8 @@ async fn m3u_api(
pub(crate) fn m3u_api_register() -> Vec<Resource> {
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))
]
}
+3 -2
View File
@@ -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<Config>, targets: Arc<ProcessTargets>)
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
+4 -1
View File
@@ -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,
})
}
}
+47 -21
View File
@@ -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<UserApiRequest>,
path: web::Path<(String, String, String, String, String)>,
_app_state: web::Data<AppState>,
) -> 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<UserApiRequest>,
@@ -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<Resource> {
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))
*/
]
}
+10
View File
@@ -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<String>,
#[serde(default = "default_as_true")]
pub enabled: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub options: Option<ConfigInputOptions>,
}
impl ConfigInput {
+4 -4
View File
@@ -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<String> = 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<String>, removed: BTr
}
}
fn load_tree(path: &Path) -> Option<BTreeSet<String>> {
fn load_watch_tree(path: &Path) -> Option<BTreeSet<String>> {
match std::fs::read(path) {
Ok(encoded) => {
let decoded: BTreeSet<String> = bincode::deserialize(&encoded[..]).unwrap();
@@ -86,7 +86,7 @@ fn load_tree(path: &Path) -> Option<BTreeSet<String>> {
}
}
fn save_tree(path: &Path, tree: BTreeSet<String>) -> std::io::Result<()> {
fn save_watch_tree(path: &Path, tree: BTreeSet<String>) -> std::io::Result<()> {
let encoded: Vec<u8> = bincode::serialize(&tree).unwrap();
std::fs::write(path, encoded)
}
+150 -165
View File
@@ -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<i32, (u32, u16)>;
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<PathBuf> {
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<T>(file: &Path, value: &T) -> Result<(), Error>
where
T: ?Sized + Serialize {
@@ -50,68 +78,40 @@ fn write_to_file<T>(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::<i32, (u32, u16)>::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<u8> = 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<u8> = 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<PathBuf> {
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<BTreeMap<i32, (u32, u16)>> {
match std::fs::read(path) {
fn load_index(path: &Path) -> Option<IndexTree> {
match fs::read(path) {
Ok(encoded) => {
let decoded: BTreeMap<i32, (u32, u16)> = 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<String, Error> {
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<String, Error> {
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<IndexTree> = 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<u8> = 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<String, Error> {
/*
{
"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<String, Error> {
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<String, Error> {
/*
{
"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)))
}
+15 -5
View File
@@ -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<String, String>) -> 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
}