mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-04 23:12:27 +02:00
feat: prevent duplicate xtream info fetches for mapping aliases
This commit is contained in:
@@ -206,11 +206,9 @@ fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> (PlaylistItem, b
|
||||
fn map_channel_with_aliases(channel: PlaylistItem, mapping: &Mapping) -> Vec<PlaylistItem> {
|
||||
if mapping.create_alias {
|
||||
let original = channel.clone();
|
||||
let base_uuid = *original.header.get_uuid();
|
||||
let (mut mapped_channel, matched) = map_channel(channel, mapping);
|
||||
if matched {
|
||||
mapped_channel.header.uuid = create_alias_uuid(original.header.get_uuid(), &mapping.id);
|
||||
mapped_channel.header.alias_of = Some(base_uuid);
|
||||
vec![original, mapped_channel]
|
||||
} else {
|
||||
vec![mapped_channel]
|
||||
|
||||
@@ -82,9 +82,6 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input:
|
||||
pub(in crate::processing) fn should_update_info(pli: &mut PlaylistItem, processed_provider_ids: &HashMap<u32, u64>, field: &str) -> (bool, u32, u64) {
|
||||
let Some(provider_id) = pli.header.get_provider_id() else { return (false, 0, 0) };
|
||||
let last_modified = pli.header.get_additional_property_as_u64(field);
|
||||
if pli.header.alias_of.is_some() {
|
||||
return (false, provider_id, last_modified.unwrap_or(0));
|
||||
}
|
||||
let old_timestamp = processed_provider_ids.get(&provider_id);
|
||||
(old_timestamp.is_none()
|
||||
|| last_modified.is_none()
|
||||
|
||||
@@ -11,7 +11,7 @@ use crate::repository::xtream_repository::{write_series_info_to_wal_file, xtream
|
||||
use crate::repository::IndexedDocumentReader;
|
||||
use shared::error::{notify_err, info_err};
|
||||
use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target};
|
||||
use std::collections::HashMap;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::fs::File;
|
||||
use std::io::{BufWriter, Write};
|
||||
use std::sync::Arc;
|
||||
@@ -53,6 +53,7 @@ fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &Ha
|
||||
async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc<reqwest::Client>, errors: &mut Vec<TuliproxError>,
|
||||
fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool {
|
||||
let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, fpl).await;
|
||||
let mut fetched_in_run: HashSet<u32> = HashSet::new();
|
||||
// we cant write to the indexed-document directly because of the write lock and time-consuming operation.
|
||||
// All readers would be waiting for the lock and the app would be unresponsive.
|
||||
// We collect the content into a wal file and write it once we collected everything.
|
||||
@@ -81,7 +82,7 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc<reqwest::Clie
|
||||
let mut last_processed_series_info_count = 0;
|
||||
for pli in series_info_iter {
|
||||
let (should_update, provider_id, ts) = should_update_series_info(pli, &processed_info_ids);
|
||||
if should_update {
|
||||
if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) {
|
||||
if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await {
|
||||
let normalized_content = normalize_json_content(content);
|
||||
handle_error_and_return!(write_series_info_to_wal_file(provider_id, ts, &normalized_content, &mut content_writer, &mut record_writer),
|
||||
|
||||
@@ -10,7 +10,7 @@ use crate::processing::processor::{handle_error, handle_error_and_return, create
|
||||
use shared::utils::{get_u32_from_serde_value, get_u64_from_serde_value, get_string_from_serde_value};
|
||||
use crate::repository::xtream_repository::xtream_get_input_info;
|
||||
use serde_json::{from_str, Map, Value};
|
||||
use std::collections::HashMap;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::io::{Write};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
@@ -77,6 +77,7 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc<reqwest::C
|
||||
else { return; };
|
||||
|
||||
let mut processed_info_ids = read_processed_vod_info_ids(app_config, errors, fpl).await;
|
||||
let mut fetched_in_run: HashSet<u32> = HashSet::new();
|
||||
let mut content_writer = utils::file_writer(&wal_content_file);
|
||||
let mut record_writer = utils::file_writer(&wal_record_file);
|
||||
let mut content_updated = false;
|
||||
@@ -96,8 +97,8 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc<reqwest::C
|
||||
let mut last_processed_vod_info_count = 0;
|
||||
|
||||
for pli in vod_info_iter {
|
||||
let (should_update, _provider_id, _ts) = should_update_vod_info(pli, &processed_info_ids);
|
||||
if should_update {
|
||||
let (should_update, provider_id, _ts) = should_update_vod_info(pli, &processed_info_ids);
|
||||
if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) {
|
||||
if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await {
|
||||
let normalized_content = normalize_json_content(content);
|
||||
if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(&normalized_content) {
|
||||
|
||||
Reference in New Issue
Block a user