mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-02 14:02:22 +02:00
xtrem reolve video/series wip
This commit is contained in:
@@ -1,20 +1,44 @@
|
||||
use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind};
|
||||
use crate::model::config::{Config, ConfigInput};
|
||||
use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, XtreamCluster};
|
||||
use crate::model::playlist::{PlaylistEntry, PlaylistItem, XtreamCluster};
|
||||
use crate::repository::storage::get_input_storage_path;
|
||||
use crate::repository::xtream_repository::{xtream_get_info_file_paths};
|
||||
use crate::repository::IndexedDocumentIndex;
|
||||
use crate::utils::download;
|
||||
use serde_json::Value;
|
||||
use std::collections::HashSet;
|
||||
use std::collections::{HashMap};
|
||||
use std::fs::{File, OpenOptions};
|
||||
use std::io::{BufWriter, Error, ErrorKind, Write};
|
||||
|
||||
|
||||
const FILE_SERIES_INFO:&str = "xtream_series_info";
|
||||
const FILE_VOD_INFO:&str = "xtream_vod_info";
|
||||
const FILE_SUFFIX_WAL:&str = "wal";
|
||||
|
||||
#[macro_export]
|
||||
macro_rules! handle_error {
|
||||
($stmt:expr) => {
|
||||
if let Err(err) = $stmt {
|
||||
errors.push(err);
|
||||
}
|
||||
};
|
||||
|
||||
($stmt:expr, $map_err:expr) => {
|
||||
if let Err(err) = $stmt {
|
||||
errors.push($map_err(err));
|
||||
}
|
||||
};
|
||||
}
|
||||
pub(in crate::processing) use handle_error;
|
||||
|
||||
#[macro_export]
|
||||
macro_rules! handle_error_and_return {
|
||||
($stmt:expr, $map_err:expr) => {
|
||||
if let Err(err) = $stmt {
|
||||
errors.push($map_err(err));
|
||||
return;
|
||||
}
|
||||
};
|
||||
}
|
||||
pub(in crate::processing) use handle_error_and_return;
|
||||
|
||||
#[macro_export]
|
||||
macro_rules! create_resolve_options_function_for_xtream_input {
|
||||
($cluster:ident) => {
|
||||
@@ -32,12 +56,12 @@ macro_rules! create_resolve_options_function_for_xtream_input {
|
||||
}
|
||||
|
||||
|
||||
pub fn get_u32_from_serde_value(value: &Value) -> Option<u32> {
|
||||
pub fn get_u64_from_serde_value(value: &Value) -> Option<u64> {
|
||||
match value {
|
||||
Value::Number(num_val) => num_val.as_u64().and_then(|val| u32::try_from(val).ok()),
|
||||
Value::Number(num_val) => num_val.as_u64(),
|
||||
Value::String(str_val) => {
|
||||
match str_val.parse::<u32>() {
|
||||
Ok(sid) => Some(sid),
|
||||
match str_val.parse::<u64>() {
|
||||
Ok(val) => Some(val),
|
||||
Err(_) => None
|
||||
}
|
||||
}
|
||||
@@ -45,6 +69,10 @@ pub fn get_u32_from_serde_value(value: &Value) -> Option<u32> {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn get_u32_from_serde_value(value: &Value) -> Option<u32> {
|
||||
get_u64_from_serde_value(value).and_then(|val| u32::try_from(val).ok())
|
||||
}
|
||||
|
||||
pub(in crate::processing) async fn playlist_resolve_process_playlist_item(pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec<M3uFilterError>, resolve_delay: u16, cluster: XtreamCluster) -> Option<String> {
|
||||
let mut result = None;
|
||||
let provider_id = pli.get_provider_id().unwrap_or(0);
|
||||
@@ -63,32 +91,7 @@ pub(in crate::processing) async fn playlist_resolve_process_playlist_item(pli: &
|
||||
result
|
||||
}
|
||||
|
||||
|
||||
pub(in crate::processing) async fn read_processed_info_ids(cfg: &Config, errors: &mut Vec<M3uFilterError>, fpl: &FetchedPlaylist<'_>, cluster: XtreamCluster) -> HashSet<u32> {
|
||||
let mut processed_info_ids = HashSet::new();
|
||||
{
|
||||
match get_input_storage_path(fpl.input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) {
|
||||
Ok(Some((file_path, idx_path))) => {
|
||||
match cfg.file_locks.read_lock(&file_path).await {
|
||||
Ok(file_lock) => {
|
||||
if let Ok(info_id_mapping) = IndexedDocumentIndex::<u32>::load(&idx_path) {
|
||||
info_id_mapping.traverse(|keys, _| {
|
||||
for doc_id in keys { processed_info_ids.insert(*doc_id); }
|
||||
});
|
||||
}
|
||||
drop(file_lock);
|
||||
}
|
||||
Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))),
|
||||
}
|
||||
}
|
||||
Ok(None) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &fpl.input.name.as_ref().map_or("?", |v| v)))),
|
||||
Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))),
|
||||
}
|
||||
}
|
||||
processed_info_ids
|
||||
}
|
||||
|
||||
pub(in crate::processing) fn write_info_content_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, content: &str) -> std::io::Result<()> {
|
||||
pub(in crate::processing) fn write_info_content_to_wal_file(writer: &mut BufWriter<&File>, provider_id: u32, content: &str) -> std::io::Result<()> {
|
||||
let length = u32::try_from(content.len()).map_err(|err| Error::new(ErrorKind::Other, err))?;
|
||||
if length > 0 {
|
||||
writer.write_all(&provider_id.to_le_bytes())?;
|
||||
@@ -107,10 +110,10 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input:
|
||||
XtreamCluster::Series => Some(FILE_VOD_INFO)
|
||||
} {
|
||||
let content_path = storage_path.join(format!("{file_prefix}_content.{FILE_SUFFIX_WAL}"));
|
||||
let tmdb_path = storage_path.join(format!("{file_prefix}_tmdb.{FILE_SUFFIX_WAL}"));
|
||||
let info_path = storage_path.join(format!("{file_prefix}_record.{FILE_SUFFIX_WAL}"));
|
||||
let content_file = OpenOptions::new().append(true).open(content_path).ok()?;
|
||||
let tmdb_file = OpenOptions::new().append(true).open(tmdb_path).ok()?;
|
||||
return Some((content_file, tmdb_file));
|
||||
let info_file = OpenOptions::new().append(true).open(info_path).ok()?;
|
||||
return Some((content_file, info_file));
|
||||
}
|
||||
None
|
||||
|
||||
@@ -118,3 +121,31 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input:
|
||||
Err(_) => None
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
pub(in crate::processing) fn has_different_ts(ts: &u64, pli: &PlaylistItem, field: &str) -> bool {
|
||||
pli.header
|
||||
.borrow()
|
||||
.additional_properties
|
||||
.as_ref()
|
||||
.map_or(false, |v| match v {
|
||||
Value::Object(map) => {
|
||||
if let Some(updated) = map.get(field) {
|
||||
if let Some(update_ts) = get_u64_from_serde_value(updated) {
|
||||
return update_ts != *ts;
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
_ => true,
|
||||
})
|
||||
}
|
||||
|
||||
pub(in crate::processing) fn should_update_info(pli: &PlaylistItem, processed_provider_ids: &HashMap<u32, u64>, field: &str) -> bool {
|
||||
if let Some(provider_id) = pli.header.borrow_mut().get_provider_id() {
|
||||
let timestamp = processed_provider_ids.get(&provider_id);
|
||||
timestamp.is_none() || has_different_ts(timestamp.unwrap(), pli, field)
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user