Playlist Update Optimization: Reduced Memory Usage

An optimization has been introduced to reduce memory consumption during playlist updates.

Overview
Previously, provider playlists were fully loaded into memory during the update process. For large playlists (e.g. hundreds of thousands of entries or multiple providers), this could result in significant RAM usage.
With the new implementation, it is now possible to optionally read provider playlists directly from disk instead of keeping them entirely in memory.

How It Works
 - users can configure whether:
     - Provider playlists are loaded into memory (previous behavior), or
     - Provider playlists are streamed to/from disk to minimize RAM usage.

Processing remains sequential and batch-based, ensuring identical functional behavior.
This approach significantly lowers peak memory usage, especially on systems with limited resources.

Benefits

- Reduced peak RAM consumption during playlist updates
- Better scalability for large playlists and multiple providers
- Full backward compatibility

Trade-offs / Drawbacks

 - Increased processing time due to reduced in-memory caching
 - Higher disk I/O usage, especially for large or fragmented playlists
 - Performance depends more strongly on disk speed (Nvme SSD, HDD)
 - Slightly increased CPU overhead due to repeated parsing and deserialization
 - Not optimal for environments where fast updates are more important than memory usage

Recommendation

 - Use in-memory mode for systems with sufficient RAM and a focus on update speed
 - Use disk-based mode for large playlists, multiple providers, or memory-constrained environments
This commit is contained in:
euzu
2026-01-08 15:11:12 +01:00
parent 4aface6a31
commit 8cc9283a31
122 changed files with 3765 additions and 1397 deletions
+23 -11
View File
@@ -1,8 +1,8 @@
use crate::model::{Epg, TVGuide, XmlTag, XmlTagIcon, EPG_ATTRIB_ID};
use crate::model::{EpgConfig, EpgSmartMatchConfig};
use crate::model::{FetchedPlaylist};
use crate::model::FetchedPlaylist;
use crate::processing::parser::xmltv::normalize_channel_name;
use log::{debug, trace};
use log::{debug, trace, warn};
use rphonetic::{DoubleMetaphone, Encoder};
use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
@@ -109,10 +109,12 @@ impl EpgIdCache<'_> {
let smart_match_enabled = self.smart_match_enabled;
let fuzzy_matching = self.fuzzy_match_enabled;
for channel in fp.playlist_groups.iter().flat_map(|g| &g.channels) {
// Helper closure to process a single item
// We use a closure here to capture `self` and avoid code duplication
let mut process_item = |name: &str, epg_channel_id: Option<&str>| {
let mut missing_epg_id = true;
// insert epg_id to known channel epg_ids
if let Some(id) = channel.header.epg_channel_id.as_deref() {
if let Some(id) = epg_channel_id {
if !id.is_empty() {
missing_epg_id = false;
self.channel_epg_id.insert(Cow::Owned(id.to_string()));
@@ -124,8 +126,13 @@ impl EpgIdCache<'_> {
let needs_normalization = smart_match_enabled && (fuzzy_matching || missing_epg_id);
if needs_normalization {
let name = &channel.header.name;
self.normalize_and_store(name, channel.header.epg_channel_id.as_deref());
self.normalize_and_store(name, epg_channel_id);
}
};
for channel in fp.items() {
if channel.header.xtream_cluster == XtreamCluster::Live && channel.header.item_type.is_live() {
process_item(&channel.header.name, channel.header.epg_channel_id.as_deref());
}
}
}
@@ -211,11 +218,13 @@ async fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist<'_>
}
};
let filter_live = |c: &&mut PlaylistItem| c.header.xtream_cluster == XtreamCluster::Live;
fp.playlist_groups.iter_mut()
.flat_map(|g| &mut g.channels)
.filter(filter_live)
.for_each(assign_values);
let filter_live = |c: &&mut PlaylistItem| c.header.xtream_cluster == XtreamCluster::Live && c.header.item_type.is_live();
if fp.is_memory() {
fp.items_mut().filter(filter_live).for_each(assign_values);
} else {
warn!("Disk based playlist modification is not supported!");
}
processed_epgs.push(epg_source);
}
}
@@ -238,6 +247,9 @@ async fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist<'_>
/// process_playlist_epg(&mut playlist, &mut epg_data);
/// ```
pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec<Epg>) {
if fp.input.epg.is_none() {
return;
}
// collect all epg_channel ids
let mut id_cache = EpgIdCache::new(fp.input.epg.as_ref());
id_cache.collect_epg_id(fp);
+16 -20
View File
@@ -2,7 +2,7 @@ use crate::library::{MediaMetadata, MetadataAsyncIter, MetadataCacheEntry};
use crate::model::{AppConfig, ConfigInput};
use shared::error::TuliproxError;
use shared::model::{EpisodeStreamProperties, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, SeriesStreamDetailEpisodeProperties, SeriesStreamDetailProperties, SeriesStreamProperties, StreamProperties, UUIDType, VideoStreamDetailProperties, VideoStreamProperties, XtreamCluster};
use shared::utils::{generate_playlist_uuid};
use shared::utils::{generate_playlist_uuid, StringInterner};
use std::path::Path;
use std::sync::Arc;
@@ -25,15 +25,14 @@ pub async fn download_library_playlist(_client: &reqwest::Client, app_config: &A
channels: vec![],
xtream_cluster: XtreamCluster::Series,
};
let mut interner = StringInterner::new();
while let Some(entry) = metadata_iter.next().await {
match entry.metadata {
MediaMetadata::Movie(_) => {
let pli = to_playlist_item(&entry, &input.name, &library_config.playlist.movie_category);
group_movies.channels.extend(pli);
to_playlist_item(&mut interner, &entry, &input.name, &library_config.playlist.movie_category, &mut group_movies.channels);
}
MediaMetadata::Series(_) => {
let pli = to_playlist_item(&entry, &input.name, &library_config.playlist.series_category);
group_series.channels.extend(pli);
to_playlist_item(&mut interner, &entry, &input.name, &library_config.playlist.series_category, &mut group_series.channels);
}
}
}
@@ -49,30 +48,30 @@ pub async fn download_library_playlist(_client: &reqwest::Client, app_config: &A
(groups, vec![])
}
fn to_playlist_item(entry: &MetadataCacheEntry, input_name: &str, group_name: &str) -> Vec<PlaylistItem> {
fn to_playlist_item(interner: &mut StringInterner, entry: &MetadataCacheEntry, input_name: &str, group_name: &str, channels: &mut Vec<PlaylistItem>) {
let metadata = &entry.metadata;
match metadata {
MediaMetadata::Movie(_) => {
let additional_properties = metadata_cache_entry_to_xtream_movie_info(entry);
vec![PlaylistItem {
channels.push(PlaylistItem {
header: PlaylistItemHeader {
uuid: UUIDType::from_valid_uuid(&entry.uuid),
name: metadata.title().to_string(),
group: group_name.to_string(),
group: interner.intern(group_name),
title: metadata.title().to_string(),
logo: metadata.poster().map_or_else(String::new, ToString::to_string),
url: format!("file://{}", entry.file_path),
xtream_cluster: XtreamCluster::Video,
additional_properties,
item_type: PlaylistItemType::LocalVideo,
input_name: input_name.to_string(),
input_name: interner.intern(input_name),
..PlaylistItemHeader::default()
}
}]
});
}
MediaMetadata::Series(_series) => {
let mut items = vec![];
if let Some(additional_properties) = metadata_cache_entry_to_xtream_series_info(entry) {
let mut episodes = vec![];
if let StreamProperties::Series(series_properties) = &additional_properties {
@@ -92,13 +91,13 @@ fn to_playlist_item(entry: &MetadataCacheEntry, input_name: &str, group_name: &s
uuid: generate_playlist_uuid(input_name, &episode.id.to_string(), PlaylistItemType::LocalSeries, &episode.direct_source),
logo: logo.clone(),
name: episode.title.clone(),
group: group_name.to_string(),
group: interner.intern(group_name),
title: episode.title.clone(),
url: episode.direct_source.clone(),
xtream_cluster: XtreamCluster::Series,
item_type: PlaylistItemType::LocalSeries,
category_id: 0,
input_name: input_name.to_string(),
input_name: interner.intern(input_name),
additional_properties: Some(StreamProperties::Episode(EpisodeStreamProperties {
episode_id: episode.id,
episode: episode.episode_num,
@@ -124,23 +123,20 @@ fn to_playlist_item(entry: &MetadataCacheEntry, input_name: &str, group_name: &s
uuid: UUIDType::from_valid_uuid(&entry.uuid),
id: entry.uuid.clone(),
name: metadata.title().to_string(),
group: group_name.to_string(),
group: interner.intern(group_name),
title: metadata.title().to_string(),
logo: metadata.poster().map_or_else(String::new, ToString::to_string),
url: format!("file://{}", entry.file_path),
xtream_cluster: XtreamCluster::Series,
item_type: PlaylistItemType::LocalSeriesInfo,
input_name: input_name.to_string(),
input_name: interner.intern(input_name),
additional_properties: Some(additional_properties),
..PlaylistItemHeader::default()
}
};
items.push(series_info);
items.extend(episodes);
channels.push(series_info);
channels.extend(episodes);
}
items
}
}
}
+282 -174
View File
@@ -1,16 +1,18 @@
use crate::model::{AppConfig, Config, ConfigFavourites, ConfigInput, ConfigRename};
use crate::model::{AppConfig, Config, ConfigFavourites, ConfigInput, ConfigRename, TVGuide};
use crate::utils::m3u;
use crate::utils::xtream;
use crate::utils::{epg, StepMeasureCallback};
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::Mutex;
use std::sync::{Arc, Weak};
use tokio::sync::{Mutex, OwnedRwLockWriteGuard, RwLock};
use tokio::task::JoinSet;
use crate::api::model::{EventManager, EventMessage, PlaylistStorageState, UpdateGuard};
use crate::messaging::send_message_json;
use crate::model::Epg;
use crate::model::FetchedPlaylist;
use crate::model::Mapping;
use crate::model::{ConfigTarget, ProcessTargets};
@@ -26,16 +28,20 @@ use crate::processing::processor::trakt::process_trakt_categories_for_target;
use crate::processing::processor::xtream_series::playlist_resolve_series;
use crate::processing::processor::xtream_vod::playlist_resolve_vod;
use crate::repository::playlist_repository::{load_input_playlist, persist_input_playlist, persist_playlist};
use crate::repository::xtream_repository::CategoryKey;
use crate::repository::{MemoryPlaylistSource, PlaylistSource};
use crate::utils::StepMeasure;
use crate::utils::{debug_if_enabled, trace_if_enabled};
use futures::StreamExt;
use indexmap::IndexMap;
use log::{debug, error, info, log_enabled, trace, warn, Level};
use log::{debug, error, info, log_enabled, warn, Level};
use shared::error::{get_errors_notify_message, notify_err, TuliproxError};
use shared::foundation::filter::{get_field_value, set_field_value, Filter, ValueAccessor, ValueProvider};
use shared::model::xtream_const::XTREAM_CLUSTER;
use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster};
use shared::utils::{create_alias_uuid, default_as_default};
use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind,
PlaylistGroup, PlaylistItem, PlaylistItemType, PlaylistUpdateState,
ProcessingOrder, UUIDType};
use shared::utils::{create_alias_uuid, default_as_default, StringInterner};
use std::time::Instant;
fn is_valid(pli: &PlaylistItem, filter: &Filter, match_as_ascii: bool) -> bool {
@@ -43,14 +49,36 @@ fn is_valid(pli: &PlaylistItem, filter: &Filter, match_as_ascii: bool) -> bool {
filter.filter(&provider)
}
#[allow(clippy::unnecessary_wraps)]
pub fn apply_filter_to_source(source: &mut dyn PlaylistSource, filter: &Filter) -> Option<Vec<PlaylistGroup>> {
let mut groups: IndexMap<Arc<str>, PlaylistGroup> = IndexMap::new();
for pli in source.into_items() {
if is_valid(&pli, filter, false) {
let group_title = pli.header.group.clone();
let cluster = pli.header.xtream_cluster;
let cat_id = pli.header.category_id;
groups.entry(group_title.clone())
.or_insert_with(|| PlaylistGroup {
id: cat_id,
title: group_title,
channels: vec![],
xtream_cluster: cluster,
})
.channels.push(pli);
}
}
if groups.is_empty() { None } else { Some(groups.into_values().collect()) }
}
fn filter_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget, _interner: &mut StringInterner) -> Option<Vec<PlaylistGroup>> {
apply_filter_to_source(source, &target.filter)
}
pub fn apply_filter_to_playlist(playlist: &mut [PlaylistGroup], filter: &Filter) -> Option<Vec<PlaylistGroup>> {
debug!("Filtering {} groups", playlist.len());
let mut new_playlist = Vec::with_capacity(128);
for pg in playlist.iter_mut() {
let channels = pg.channels.iter()
.filter(|&pli| is_valid(pli, filter, false)).cloned().collect::<Vec<PlaylistItem>>();
trace!("Filtered group {} has now {}/{} items", pg.title, channels.len(), pg.channels.len());
if !channels.is_empty() {
new_playlist.push(PlaylistGroup {
id: pg.id,
@@ -60,11 +88,7 @@ pub fn apply_filter_to_playlist(playlist: &mut [PlaylistGroup], filter: &Filter)
});
}
}
Some(new_playlist)
}
fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
apply_filter_to_playlist(playlist, &target.filter)
if new_playlist.is_empty() { None } else { Some(new_playlist) }
}
fn assign_channel_no_playlist(new_playlist: &mut [PlaylistGroup]) {
@@ -86,7 +110,7 @@ fn assign_channel_no_playlist(new_playlist: &mut [PlaylistGroup]) {
}
}
fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec<ConfigRename>>) {
fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec<ConfigRename>>, interner: &mut StringInterner) {
if let Some(renames) = rename {
if !renames.is_empty() {
let result = pli;
@@ -97,33 +121,41 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec<ConfigRename>>) {
trace_if_enabled!("Renamed {}={value} to {cap}", &r.field);
}
let value = cap.into_owned();
set_field_value(result, r.field, value);
set_field_value(result, r.field, value, interner);
}
}
}
}
fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
fn rename_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget, interner: &mut StringInterner) -> Option<Vec<PlaylistGroup>> {
match &target.rename {
Some(renames) => {
if !renames.is_empty() {
let mut new_playlist: Vec<PlaylistGroup> = Vec::with_capacity(playlist.len());
for g in playlist {
let mut grp = g.clone();
for r in renames {
if matches!(r.field, ItemField::Group) {
let cap = r.pattern.replace_all(&grp.title, &r.new_name);
trace_if_enabled!("Renamed group {} to {cap} for {}", &grp.title, target.name);
grp.title = cap.into_owned();
Some(renames) if !renames.is_empty() => {
let mut groups: IndexMap<Arc<str>, PlaylistGroup> = IndexMap::new();
for mut pli in source.into_items() {
// Handle group rename first if it's in the renames
for r in renames {
if matches!(r.field, ItemField::Group) {
let value = &*pli.header.group;
let cap = r.pattern.replace_all(value, &r.new_name);
if *value != cap {
pli.header.group = interner.intern(&cap);
}
}
grp.channels.iter_mut().for_each(|pli| exec_rename(pli, target.rename.as_ref()));
new_playlist.push(grp);
}
return Some(new_playlist);
exec_rename(&mut pli, Some(renames), interner);
let group_title = pli.header.group.clone();
let cluster = pli.header.xtream_cluster;
let cat_id = pli.header.category_id;
groups.entry(group_title.clone())
.or_insert_with(|| PlaylistGroup {
id: cat_id,
title: group_title,
channels: vec![],
xtream_cluster: cluster,
})
.channels.push(pli);
}
None
Some(groups.into_values().collect())
}
_ => None
}
@@ -164,41 +196,34 @@ fn map_channel_and_flatten(channel: PlaylistItem, mapping: &Mapping) -> Vec<Play
result
}
fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
if let Some(mappings) = target.mapping.load().as_ref() {
let new_playlist: Vec<PlaylistGroup> = playlist.iter().map(|playlist_group| {
let mut grp = playlist_group.clone();
mappings.iter().filter(|&mapping| mapping.mapper.as_ref().is_some_and(|v| !v.is_empty()))
.for_each(|mapping|
grp.channels = grp.channels.drain(..).flat_map(|chan| map_channel_and_flatten(chan, mapping)).collect());
grp
}).collect();
// if the group names are changed, restructure channels to the right groups
// we use
let mut new_groups: Vec<PlaylistGroup> = Vec::with_capacity(128);
let mut grp_id: u32 = 0;
for playlist_group in new_playlist {
for channel in &playlist_group.channels {
let cluster = &channel.header.xtream_cluster;
let title = &channel.header.group;
if let Some(grp) = new_groups.iter_mut().find(|x| *x.title == **title) {
grp.channels.push(channel.clone());
} else {
grp_id += 1;
new_groups.push(PlaylistGroup {
id: grp_id,
title: title.clone(),
channels: vec![channel.clone()],
xtream_cluster: *cluster,
});
fn map_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget, _interner: &mut StringInterner) -> Option<Vec<PlaylistGroup>> {
let mapping_binding = target.mapping.load();
let mappings = mapping_binding.as_ref()?;
let valid_mappings = mappings.iter().filter(|m| m.mapper.as_ref().is_some_and(|v| !v.is_empty()));
let iter: Box<dyn Iterator<Item=PlaylistItem>> = Box::new(source.into_items());
let mapped_iter = valid_mappings.fold(iter, |iter, mapping| {
Box::new(iter.flat_map(move |chan| map_channel_and_flatten(chan, mapping)))
as Box<dyn Iterator<Item=PlaylistItem>>
});
let mut next_groups: IndexMap<Arc<str>, PlaylistGroup> = IndexMap::new();
let mut grp_id: u32 = 0;
for channel in mapped_iter {
let group_title = channel.header.group.clone();
let cluster = channel.header.xtream_cluster;
next_groups.entry(group_title.clone())
.or_insert_with(|| {
grp_id += 1;
PlaylistGroup {
id: grp_id,
title: group_title,
channels: Vec::new(),
xtream_cluster: cluster,
}
}
}
Some(new_groups)
} else {
None
})
.channels.push(channel);
}
Some(next_groups.into_values().collect())
}
fn map_playlist_counter(target: &ConfigTarget, playlist: &mut [PlaylistGroup]) {
@@ -247,7 +272,7 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo
(!user_targets.enabled && target.enabled) || (user_targets.enabled && user_targets.has_target(target.id))
}
async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc<AppConfig>, input: &ConfigInput) -> (Vec<PlaylistGroup>, Vec<TuliproxError>, bool) {
async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc<AppConfig>, input: &ConfigInput) -> (Vec<PlaylistGroup>, Vec<TuliproxError>, bool, bool) {
let config = &*app_config.config.load();
let working_dir = &config.working_dir;
@@ -286,14 +311,20 @@ async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc
};
if fully_cached {
return (vec![], vec![], true);
return (vec![], vec![], true, false);
}
let (playlist, errors) = match input.input_type {
InputType::M3u => m3u::download_m3u_playlist(client, config, input).await,
let (playlist, errors, persisted) = match input.input_type {
InputType::M3u => {
let (p, e) = m3u::download_m3u_playlist(client, config, input).await;
(p, e, false)
}
InputType::Xtream => xtream::download_xtream_playlist(config, client, input, clusters_to_download.as_deref()).await,
InputType::M3uBatch | InputType::XtreamBatch => (vec![], vec![]),
InputType::Library => library::download_library_playlist(client, app_config, input).await,
InputType::M3uBatch | InputType::XtreamBatch => (vec![], vec![], false),
InputType::Library => {
let (p, e) = library::download_library_playlist(client, app_config, input).await;
(p, e, false)
}
};
// Update Status
@@ -330,20 +361,18 @@ async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc
}
}
(playlist, errors, false)
(playlist, errors, false, persisted)
}
async fn process_source(client: &reqwest::Client, app_config: Arc<AppConfig>, source_idx: usize,
user_targets: Arc<ProcessTargets>, event_manager: Option<Arc<EventManager>>,
playlist_state: Option<&Arc<PlaylistStorageState>>,
) -> (Vec<InputStats>, Vec<TargetStats>, Vec<TuliproxError>) {
let sources = app_config.sources.load();
async fn process_source(source_idx: usize, ctx: &PlaylistProcessingContext) -> (Vec<InputStats>, Vec<TargetStats>, Vec<TuliproxError>) {
let sources = ctx.config.sources.load();
let mut errors = vec![];
let mut input_stats = HashMap::<String, InputStats>::new();
let mut target_stats = Vec::<TargetStats>::new();
if let Some(source) = sources.get_source_at(source_idx) {
let mut interner = StringInterner::default();
let mut source_playlists = Vec::with_capacity(128);
let broadcast_step = create_broadcast_callback(event_manager.as_ref());
let broadcast_step = create_broadcast_callback(ctx.event_manager.as_ref());
// Download the sources
let mut source_downloaded = false;
for input_name in &source.inputs {
@@ -351,25 +380,15 @@ async fn process_source(client: &reqwest::Client, app_config: Arc<AppConfig>, so
error!("Input {input_name} referenced by source {source_idx} does not exist");
continue;
};
if is_input_enabled(input, &user_targets) {
if is_input_enabled(input, &ctx.user_targets) {
source_downloaded = true;
let start_time = Instant::now();
// Download the playlist for input
let (playlist_groups, mut error_list) = {
let (mut playlist_groups, mut error_list) = {
broadcast_step("Playlist download", &format!("Downloading input '{}'", input.name));
// Caching Logic
let (downloaded_playlist, mut download_err, was_cached) = playlist_download_from_input(client, &app_config, input).await;
let (playlist, error) = if was_cached {
match load_input_playlist(&app_config, input, None).await {
Ok(pl) => (pl, None),
Err(e) => (vec![], Some(e)),
}
} else {
broadcast_step("Playlist download", &format!("Persisting input '{}' playlist", input.name));
persist_input_playlist(&app_config, input, downloaded_playlist).await
};
let (mut download_err, playlist, error) = download_input(ctx, input).await;
if let Some(err) = error {
broadcast_step("Playlist download", &format!("Failed to persist/load input '{}' playlist", input.name));
@@ -379,60 +398,46 @@ async fn process_source(client: &reqwest::Client, app_config: Arc<AppConfig>, so
(playlist, download_err)
};
// Download epg for input
let (tvguide, mut tvguide_errors) = if error_list.is_empty() {
broadcast_step("Playlist download", &format!("Downloading epg for input '{}'", input.name));
let working_dir = &app_config.config.load().working_dir;
epg::get_xmltv(client, input, working_dir).await
let (tvguide, mut tvguide_errors) = if input.input_type == InputType::Library {
(None, vec!())
} else {
(None, vec![])
download_input_epg(ctx,input, &mut errors).await
};
errors.append(&mut error_list);
errors.append(&mut tvguide_errors);
let group_count = playlist_groups.len();
let channel_count = playlist_groups.iter().map(|group| group.channels.len()).sum();
let group_count = playlist_groups.get_group_count();
let channel_count = playlist_groups.get_channel_count();
let input_name = &input.name;
if playlist_groups.is_empty() {
broadcast_step("Playlist download", &format!("Input '{}' playlist is empty", input.name));
info!("Source is empty {input_name}");
errors.push(notify_err!(format!("Source is empty {input_name}")));
errors.push(notify_err!("Source is empty {input_name}"));
} else {
source_playlists.push(
FetchedPlaylist {
input,
// If I create Arc here, it drops at end of loop iteration!
// SAFETY ISSUE.
// `source_playlists` stores `FetchedPlaylist`.
// `FetchedPlaylist` struct definition:
// pub struct FetchedPlaylist<'a> { pub input: &'a ConfigInput, ... }
// It holds a REFERENCE.
// If `input` comes from `sources` (guard), it lives as long as `sources` guard lives.
// `sources` guard is alive in this function scope.
// So `&input` is valid.
// `playlist_download_from_input` takes `&Arc<ConfigInput>`. I pass formatted Arc.
// `FetchedPlaylist` needs reference.
playlist_groups,
source: playlist_groups,
epg: tvguide,
}
);
}
let elapsed = start_time.elapsed().as_secs();
input_stats.insert(input_name.clone(), create_input_stat(group_count, channel_count, error_list.len(),
input_stats.insert(input_name.clone(), create_input_stat(group_count, channel_count, errors.len(),
input.input_type, input_name, elapsed));
}
}
if source_downloaded {
if source_playlists.is_empty() {
debug!("Source at index {source_idx} is empty");
errors.push(notify_err!(format!("Source at index {source_idx} is empty: {}", source.inputs.iter().map(std::string::String::as_str).collect::<Vec<&str>>().join(", "))));
errors.push(notify_err!("Source at index {source_idx} is empty: {}", source.inputs.iter().map(std::string::String::as_str).collect::<Vec<&str>>().join(", ")));
} else {
debug_if_enabled!("Source has {} groups", source_playlists.iter().map(|fpl| fpl.playlist_groups.len()).sum::<usize>());
let event_manager_clone = event_manager.clone();
debug_if_enabled!("Source has {} groups", source_playlists.iter_mut().map(FetchedPlaylist::get_channel_count).sum::<usize>());
for target in &source.targets {
let event_manager_clone = event_manager_clone.clone();
if is_target_enabled(target, &user_targets) {
match process_playlist_for_target(&app_config, client, &mut source_playlists, target, &mut input_stats, &mut errors, event_manager_clone, playlist_state).await {
if is_target_enabled(target, &ctx.user_targets) {
match process_playlist_for_target(ctx, &mut source_playlists, target,
&mut input_stats, &mut errors,
&mut interner).await {
Ok(()) => {
target_stats.push(TargetStats::success(&target.name));
}
@@ -449,6 +454,58 @@ async fn process_source(client: &reqwest::Client, app_config: Arc<AppConfig>, so
(input_stats.into_values().collect(), target_stats, errors)
}
async fn download_input_epg(ctx: &PlaylistProcessingContext, input: &Arc<ConfigInput>,
error_list: &mut [TuliproxError]) -> (Option<TVGuide>, Vec<TuliproxError>) {
// Download epg for input
let (tvguide, tvguide_errors) = if error_list.is_empty() {
debug!("Downloading epg for input '{}'", input.name);
let working_dir = &ctx.config.config.load().working_dir;
epg::get_xmltv(ctx, input, working_dir).await
} else {
(None, vec![])
};
(tvguide, tvguide_errors)
}
async fn download_input(ctx: &PlaylistProcessingContext, input: &Arc<ConfigInput>)
-> (Vec<TuliproxError>, Box<dyn PlaylistSource>, Option<TuliproxError>) {
// Coordination Logic
let need_download = !ctx.is_input_downloaded(&input.name).await;
let (downloaded_playlist, download_err, was_cached, persisted) = if need_download {
// Acquire named lock to prevent thundering herd on same input
let _input_lock = ctx.get_input_lock(&input.name).await;
// Check again after lock
let already_processed = ctx.is_input_downloaded(&input.name).await;
if already_processed {
// Use empty results, will load from disk below
(vec![], vec![], true, false)
} else {
let res = playlist_download_from_input(&ctx.client, &ctx.config, input).await;
// Mark as processed if NO critical errors?
// playlist_download_from_input returns errors but also potentially a partial playlist.
// If it attempted download, we consider it processed for this session.
ctx.mark_input_downloaded(input.name.clone()).await;
res
}
} else {
(vec![], vec![], true, false)
};
let (playlist, error) = if was_cached || persisted {
match load_input_playlist(ctx, input, None).await {
Ok(pl_source) => (pl_source, None),
Err(e) => (MemoryPlaylistSource::default().boxed(), Some(e)),
}
} else {
debug!("Persisting input '{}' playlist", input.name);
let (pl, err) = persist_input_playlist(&ctx.config, input, downloaded_playlist).await;
(MemoryPlaylistSource::new(pl).boxed(), err)
};
(download_err, playlist, error)
}
fn create_broadcast_callback(event_manager: Option<&Arc<EventManager>>) -> StepMeasureCallback {
if let Some(event_mgr) = event_manager {
let events = event_mgr.clone();
@@ -477,42 +534,80 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz
}
}
async fn process_sources(client: &reqwest::Client, config: &Arc<AppConfig>, user_targets: Arc<ProcessTargets>,
event_manager: Option<Arc<EventManager>>, playlist_state: Option<&Arc<PlaylistStorageState>>,
) -> (Vec<SourceStats>, Vec<TuliproxError>) {
#[derive(Clone)]
pub struct PlaylistProcessingContext {
pub client: reqwest::Client,
pub config: Arc<AppConfig>,
pub user_targets: Arc<ProcessTargets>,
pub event_manager: Option<Arc<EventManager>>,
pub playlist_state: Option<Arc<PlaylistStorageState>>,
// Coordination
processed_inputs: Arc<Mutex<HashSet<String>>>,
input_locks: Arc<Mutex<HashMap<String, Weak<RwLock<()>>>>>,
}
impl PlaylistProcessingContext {
pub async fn is_input_downloaded(&self, input_name: &str) -> bool {
let processed = self.processed_inputs.lock().await;
processed.contains(input_name)
}
pub async fn mark_input_downloaded(&self, input_name: String) -> bool {
let mut processed = self.processed_inputs.lock().await;
processed.insert(input_name)
}
pub async fn get_input_lock(&self, input_name: &str) -> OwnedRwLockWriteGuard<()> {
let mut locks = self.input_locks.lock().await;
// Clean up stale weak references
locks.retain(|_, weak| weak.strong_count() > 0);
if let Some(weak) = locks.get(input_name) {
if let Some(strong) = weak.upgrade() {
return strong.write_owned().await;
}
}
let lock = Arc::new(RwLock::new(()));
locks.insert(input_name.to_string(), Arc::downgrade(&lock));
lock.write_owned().await
}
}
async fn process_sources(processing_ctx: &PlaylistProcessingContext) -> (Vec<SourceStats>, Vec<TuliproxError>) {
let mut async_tasks = JoinSet::new();
let sources = config.sources.load();
let process_parallel = config.config.load().process_parallel && sources.sources.len() > 1;
let sources = processing_ctx.config.sources.load();
let process_parallel = processing_ctx.config.config.load().process_parallel && sources.sources.len() > 1;
if process_parallel && log_enabled!(Level::Debug) {
debug!("Parallel processing enabled");
}
let errors = Arc::new(Mutex::<Vec<TuliproxError>>::new(vec![]));
let stats = Arc::new(Mutex::<Vec<SourceStats>>::new(vec![]));
for (index, source) in sources.sources.iter().enumerate() {
if !source.should_process_for_user_targets(&user_targets) {
if !source.should_process_for_user_targets(&processing_ctx.user_targets) {
continue;
}
// We're using the file lock this way on purpose
let source_lock_path = PathBuf::from(format!("source_{index}"));
let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path).await else {
let Ok(update_lock) = processing_ctx.config.file_locks.try_write_lock(&source_lock_path).await else {
warn!("The update operation for the source at index {index} was skipped because an update is already in progress.");
continue;
};
let shared_errors = errors.clone();
let shared_stats = stats.clone();
let cfg = config.clone();
let usr_trgts = user_targets.clone();
let event_manager = event_manager.clone();
let ctx = processing_ctx.clone();
if process_parallel {
let http_client = client.clone();
let playlist_state = playlist_state.cloned();
async_tasks.spawn(async move {
// Hold the per-source lock for the full duration of this update.
let current_update_lock = update_lock;
let (input_stats, target_stats, mut res_errors) =
process_source(&http_client, cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await;
process_source(index, &ctx).await;
shared_errors.lock().await.append(&mut res_errors);
if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) {
shared_stats.lock().await.push(process_stats);
@@ -521,7 +616,7 @@ async fn process_sources(client: &reqwest::Client, config: &Arc<AppConfig>, user
});
} else {
let (input_stats, target_stats, mut res_errors) =
process_source(client, cfg, index, usr_trgts, event_manager, playlist_state).await;
process_source(index, &ctx).await;
shared_errors.lock().await.append(&mut res_errors);
if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) {
shared_stats.lock().await.push(process_stats);
@@ -541,7 +636,7 @@ async fn process_sources(client: &reqwest::Client, config: &Arc<AppConfig>, user
}
}
pub type ProcessingPipe = Vec<fn(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>>>;
pub type ProcessingPipe = Vec<fn(source: &mut dyn PlaylistSource, target: &ConfigTarget, interner: &mut StringInterner) -> Option<Vec<PlaylistGroup>>>;
fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe {
match &target.processing_order {
@@ -554,29 +649,27 @@ fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe {
}
}
fn duplicate_hash(item: &PlaylistItem) -> UUIDType {
item.get_uuid()
}
fn execute_pipe<'a>(target: &ConfigTarget, pipe: &ProcessingPipe, fpl: &FetchedPlaylist<'a>,
duplicates: &mut HashSet<UUIDType>) -> FetchedPlaylist<'a> {
duplicates: &mut HashSet<UUIDType>,
interner: &mut StringInterner) -> FetchedPlaylist<'a> {
let mut new_fpl = FetchedPlaylist {
input: fpl.input,
playlist_groups: fpl.playlist_groups.clone(), // we need to clone, because of multiple target definitions, we cant change the initial playlist.
source: fpl.clone_source(),
epg: fpl.epg.clone(),
};
if target.options.as_ref().is_some_and(|opt| opt.remove_duplicates) {
for group in &mut new_fpl.playlist_groups {
// `HashSet::insert` returns true for first insert, otherwise false
group.channels.retain(|item| duplicates.insert(duplicate_hash(item)));
}
new_fpl.deduplicate(duplicates);
}
for f in pipe {
if let Some(groups) = f(&mut new_fpl.playlist_groups, target) {
new_fpl.playlist_groups = groups;
if let Some(groups) = f(new_fpl.source.as_mut(), target, interner) {
new_fpl.source = MemoryPlaylistSource::new(groups).boxed();
}
}
// Ensure source is memory-based for downstream mutable processing (VOD/series resolution)
if !new_fpl.is_memory() {
new_fpl.source = MemoryPlaylistSource::new(new_fpl.source.take_groups()).boxed();
}
new_fpl
}
@@ -585,9 +678,9 @@ fn execute_pipe<'a>(target: &ConfigTarget, pipe: &ProcessingPipe, fpl: &FetchedP
fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
let mut sort_order: Vec<PlaylistGroup> = vec![];
let mut idx: usize = 0;
let mut group_map: HashMap<(String, XtreamCluster), usize> = HashMap::new();
let mut group_map: HashMap<CategoryKey, usize> = HashMap::new();
for group in playlistgroups {
let key = (group.title.clone(), group.xtream_cluster);
let key = (group.xtream_cluster, group.title.clone());
match group_map.entry(key) {
std::collections::hash_map::Entry::Vacant(v) => {
v.insert(idx);
@@ -605,14 +698,12 @@ fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
}
#[allow(clippy::too_many_arguments)]
async fn process_playlist_for_target(app_config: &Arc<AppConfig>,
client: &reqwest::Client,
async fn process_playlist_for_target(ctx: &PlaylistProcessingContext,
playlists: &mut [FetchedPlaylist<'_>],
target: &ConfigTarget,
stats: &mut HashMap<String, InputStats>,
errors: &mut Vec<TuliproxError>,
event_manager: Option<Arc<EventManager>>,
playlist_state: Option<&Arc<PlaylistStorageState>>,
interner: &mut StringInterner,
) -> Result<(), Vec<TuliproxError>> {
debug_if_enabled!("Processing order is {}", &target.processing_order);
@@ -620,22 +711,23 @@ async fn process_playlist_for_target(app_config: &Arc<AppConfig>,
let mut processed_fetched_playlists: Vec<FetchedPlaylist> = vec![];
debug!("Executing processing pipes");
let broadcast_step = create_broadcast_callback(event_manager.as_ref());
let broadcast_step = create_broadcast_callback(ctx.event_manager.as_ref());
let pipe = get_processing_pipe(target);
let mut step = StepMeasure::new(&target.name, broadcast_step);
for provider_fpl in playlists.iter_mut() {
step.broadcast("Executing transformations on '{}' playlist", &target.name);
let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates);
playlist_resolve_series(app_config, client, target, errors, &pipe, provider_fpl, &mut processed_fpl).await;
playlist_resolve_vod(app_config, client, target, errors, &mut processed_fpl).await;
let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates, interner);
processed_fpl.sort_by_provider_ordinal();
playlist_resolve_series(&ctx.config, &ctx.client, target, errors, &pipe, provider_fpl, &mut processed_fpl, interner).await;
playlist_resolve_vod(&ctx.config, &ctx.client, target, errors, provider_fpl, &mut processed_fpl).await;
// stats
let input_stats = stats.get_mut(&processed_fpl.input.name);
if let Some(stat) = input_stats {
stat.processed_stats.group_count = processed_fpl.playlist_groups.len();
stat.processed_stats.channel_count = processed_fpl.playlist_groups.iter()
.map(|group| group.channels.len())
.sum();
let input_entry_name = processed_fpl.input.name.clone();
let group_count = processed_fpl.get_group_count();
let channel_count = processed_fpl.get_channel_count();
if let Some(stat) = stats.get_mut(&input_entry_name) {
stat.processed_stats.group_count = group_count;
stat.processed_stats.channel_count = channel_count;
}
processed_fetched_playlists.push(processed_fpl);
}
@@ -654,7 +746,7 @@ async fn process_playlist_for_target(app_config: &Arc<AppConfig>,
Ok(())
} else {
// Process Trakt categories
if trakt_playlist(client, target, errors, &mut new_playlist).await {
if trakt_playlist(&ctx.client, target, errors, &mut new_playlist).await {
step.tick("trakt categories");
}
@@ -669,11 +761,11 @@ async fn process_playlist_for_target(app_config: &Arc<AppConfig>,
map_playlist_counter(target, &mut flat_new_playlist);
step.tick("assigning channel counter");
let config = app_config.config.load();
if process_watch(&config, client, target, &flat_new_playlist).await {
let config = ctx.config.config.load();
if process_watch(&config, &ctx.client, target, &flat_new_playlist).await {
step.tick("group watches");
}
let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, playlist_state).await;
let result = persist_playlist(&ctx.config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, ctx.playlist_state.as_ref(), interner).await;
step.stop("Persisting playlists");
result
}
@@ -681,7 +773,7 @@ async fn process_playlist_for_target(app_config: &Arc<AppConfig>,
pub fn process_favourites(playlist: &mut Vec<PlaylistGroup>, favourites_cfg: Option<&[ConfigFavourites]>) {
if let Some(favourites) = favourites_cfg {
let mut fav_groups: IndexMap<String, Vec<PlaylistItem>> = IndexMap::new();
let mut fav_groups: IndexMap<Arc<str>, Vec<PlaylistItem>> = IndexMap::new();
for pg in playlist.iter() {
for pli in &pg.channels {
// series episodes cant be included in favourites
@@ -737,14 +829,14 @@ async fn trakt_playlist(client: &reqwest::Client, target: &ConfigTarget, errors:
}
async fn process_epg(processed_fetched_playlists: &mut Vec<FetchedPlaylist<'_>>) -> (Vec<Epg>, Vec<PlaylistGroup>) {
let mut new_playlist = vec![];
let mut new_playlist: Vec<PlaylistGroup> = vec![];
let mut new_epg = vec![];
// each fetched playlist can have its own epgl url.
// we need to process each input epg.
for fp in processed_fetched_playlists {
process_playlist_epg(fp, &mut new_epg).await;
new_playlist.append(&mut fp.playlist_groups);
new_playlist.extend(fp.source.take_groups());
}
(new_epg, new_playlist)
}
@@ -786,9 +878,19 @@ pub async fn exec_processing(client: &reqwest::Client, app_config: Arc<AppConfig
None
};
let event_manager_clone = event_manager.clone();
// Initialize Context
let ctx = PlaylistProcessingContext {
client: client.clone(),
config: app_config.clone(),
user_targets: targets.clone(),
event_manager: event_manager.clone(),
playlist_state: playlist_state.clone(),
processed_inputs: Arc::new(Mutex::new(HashSet::new())),
input_locks: Arc::new(Mutex::new(HashMap::new())),
};
let start_time = Instant::now();
let (stats, errors) = process_sources(client, &app_config, targets.clone(), event_manager_clone, playlist_state.as_ref()).await;
let (stats, errors) = process_sources(&ctx).await;
// log errors
for err in &errors {
error!("{}", err.message);
@@ -816,17 +918,23 @@ pub async fn exec_processing(client: &reqwest::Client, app_config: Arc<AppConfig
// send errors
if let Some(message) = get_errors_notify_message!(errors, 255) {
if let Some(events) = event_manager {
if let Some(events) = &event_manager {
events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Failure));
}
if let Ok(error_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("errors".to_string(), serde_json::Value::String(message))]))) {
send_message_json(client, MsgKind::Error, messaging, error_msg.as_str()).await;
}
} else if let Some(events) = event_manager {
} else if let Some(events) = &event_manager {
events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Success));
}
let elapsed = start_time.elapsed().as_secs();
info!("🌷 Update process finished! Took {elapsed} secs.");
let update_finished_message = format!("🌷 Update process finished! Took {elapsed} secs.");
if let Some(events) = &event_manager {
events.send_event(EventMessage::PlaylistUpdateProgress("Playlist Update".to_string(), update_finished_message.clone()));
}
info!("{update_finished_message}");
}
// #[cfg(test)]
+3 -3
View File
@@ -109,8 +109,8 @@ fn playlist_comparator(
}
fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &ConfigSortGroup, match_as_ascii: bool) -> Ordering {
let value_a = if match_as_ascii { deunicode(&a.title) } else { a.title.clone() };
let value_b = if match_as_ascii { deunicode(&b.title) } else { b.title.clone() };
let value_a = if match_as_ascii { deunicode(&a.title) } else { a.title.to_string() };
let value_b = if match_as_ascii { deunicode(&b.title) } else { b.title.to_string() };
playlist_comparator(group_sort.sequence.as_ref(), group_sort.order, &value_a, &value_b)
}
@@ -144,7 +144,7 @@ pub(in crate::processing::processor) fn sort_playlist(target: &ConfigTarget, new
}
let regexp = &channel_sort.group_pattern;
for group in new_playlist.iter_mut() {
let group_title = if match_as_ascii { deunicode(&group.title) } else { group.title.clone() };
let group_title = if match_as_ascii { deunicode(&group.title) } else { group.title.to_string() };
if regexp.is_match(group_title.as_str()) {
group.channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, channel_sort, match_as_ascii));
}
+9 -4
View File
@@ -4,6 +4,7 @@ use crate::utils::{extract_year_from_title, normalize_title_for_matching, TraktC
use crate::utils::{trace_if_enabled, with};
use log::{debug, info, trace, warn};
use shared::error::TuliproxError;
use shared::utils::StringInterner;
use shared::model::{FieldGetAccessor, FieldSetAccessor, PlaylistGroup, PlaylistItem, TraktContentType, XtreamCluster};
use shared::utils::CONSTANTS;
use std::borrow::Cow;
@@ -128,6 +129,7 @@ fn find_best_match_for_item<'a>(
fn create_category_from_matches<'a>(
matches: Vec<TraktMatchResult<'a>>,
list_config: &'a TraktListConfig,
interner: &mut StringInterner,
) -> Option<PlaylistGroup> {
if matches.is_empty() { return None; }
@@ -164,7 +166,7 @@ fn create_category_from_matches<'a>(
header.set_field("caption", &caption);
}
}
header.group = String::from(group_title);
header.group = interner.intern(group_title);
header.gen_uuid();
});
matched_items.push(modified_item);
@@ -184,7 +186,7 @@ fn create_category_from_matches<'a>(
Some(PlaylistGroup {
id: 0,
title: String::from(group_title),
title: interner.intern(group_title),
channels: matched_items,
xtream_cluster: cluster,
})
@@ -194,6 +196,7 @@ fn match_trakt_items_with_playlist<'a>(
trakt_items: &'a [TraktListItem],
playlist: &'a [PlaylistGroup],
list_config: &'a TraktListConfig,
interner: &mut StringInterner,
) -> Option<PlaylistGroup> {
let trakt_match_items: Vec<TraktMatchItem<'a>> = trakt_items
.iter()
@@ -217,7 +220,7 @@ fn match_trakt_items_with_playlist<'a>(
}
}
create_category_from_matches(matches, list_config)
create_category_from_matches(matches, list_config, interner)
}
pub struct TraktCategoriesProcessor {
@@ -244,6 +247,8 @@ impl TraktCategoriesProcessor {
info!("Processing {} Trakt lists for target {}", trakt_config.lists.len(), target.name);
let mut new_categories = Vec::new();
let mut total_matches = 0;
let mut interner = StringInterner::new();
for list_config in &trakt_config.lists {
let cache_key = format!("{}:{}", list_config.user, list_config.list_slug);
@@ -251,7 +256,7 @@ impl TraktCategoriesProcessor {
Ok(trakt_items) => {
debug!("Processing Trakt list {cache_key} with {} items", trakt_items.len());
if let Some(category) = match_trakt_items_with_playlist(&trakt_items, playlist, list_config) {
if let Some(category) = match_trakt_items_with_playlist(&trakt_items, playlist, list_config, &mut interner) {
if !category.channels.is_empty() {
total_matches += category.channels.len();
let category_len = category.channels.len();
+2 -2
View File
@@ -9,10 +9,10 @@ pub(in crate::processing) async fn playlist_resolve_download_playlist_item(clien
let provider_id = pli.get_provider_id()?;
if let Some(info_url) = xtream::get_xtream_player_api_info_url(input, cluster, provider_id) {
let input_source = InputSource::from(input).with_url(info_url);
result = match xtream::get_xtream_stream_info_content(client, &input_source).await {
result = match xtream::get_xtream_stream_info_content(client, &input_source, true).await {
Ok(content) => Some(content),
Err(err) => {
errors.push(info_err!(format!("{err}")));
errors.push(info_err!("{err}"));
None
}
};
@@ -5,21 +5,28 @@ use crate::processing::processor::create_resolve_options_function_for_xtream_tar
use crate::processing::processor::playlist::ProcessingPipe;
use crate::processing::processor::xtream::playlist_resolve_download_playlist_item;
use crate::repository::storage::get_input_storage_path;
use crate::repository::xtream_repository::persists_input_series_info;
use crate::repository::xtream_repository::persist_input_series_info_batch;
use crate::repository::{MemoryPlaylistSource, PlaylistSource};
use log::{error, info, log_enabled, Level};
use shared::error::TuliproxError;
use shared::model::{InputType, PlaylistEntry, SeriesStreamProperties, StreamProperties, XtreamSeriesInfo};
use shared::model::{PlaylistGroup, PlaylistItemType, XtreamCluster};
use shared::utils::StringInterner;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Instant;
create_resolve_options_function_for_xtream_target!(series);
const BATCH_SIZE: usize = 100;
#[allow(clippy::too_many_lines)]
async fn playlist_resolve_series_info(app_config: &Arc<AppConfig>, client: &reqwest::Client,
errors: &mut Vec<TuliproxError>,
fpl: &mut FetchedPlaylist<'_>,
resolve_series: bool,
resolve_delay: u16) -> Vec<PlaylistGroup> {
resolve_delay: u16,
interner: &mut StringInterner) -> Vec<PlaylistGroup> {
let input = fpl.input;
let working_dir = &app_config.config.load().working_dir;
let storage_path = match get_input_storage_path(&input.name, working_dir) {
@@ -30,88 +37,96 @@ async fn playlist_resolve_series_info(app_config: &Arc<AppConfig>, client: &reqw
}
};
let series_info_count = if resolve_series {
let count = fpl.playlist_groups.iter()
.flat_map(|plg| &plg.channels)
.filter(|&pli| pli.header.xtream_cluster == XtreamCluster::Series
&& pli.header.item_type == PlaylistItemType::SeriesInfo
&& !pli.has_details()).count();
info!("Found {count} series info to resolve");
count
} else {
0
};
let series_info_count = fpl.get_missing_series_info_count();
if series_info_count > 0 {
info!("Found {series_info_count} series info to resolve");
}
let mut last_log_time = Instant::now();
let mut processed_series_info_count = 0;
let mut result: Vec<PlaylistGroup> = vec![];
let mut group_series: HashMap<u32, PlaylistGroup> = HashMap::new();
let mut batch = Vec::with_capacity(BATCH_SIZE);
for plg in &mut fpl.playlist_groups {
let mut group_series = vec![];
for pli in &mut plg.channels {
if pli.header.xtream_cluster != XtreamCluster::Series
|| pli.header.item_type != PlaylistItemType::SeriesInfo {
continue;
}
let input = fpl.input;
for pli in fpl.items_mut() {
if pli.header.xtream_cluster != XtreamCluster::Series
|| pli.header.item_type != PlaylistItemType::SeriesInfo {
continue;
}
let Some(provider_id) = pli.get_provider_id() else { continue; };
if provider_id == 0 {
continue;
}
let Some(provider_id) = pli.get_provider_id() else { continue; };
if provider_id == 0 {
continue;
}
let should_download = resolve_series && !pli.has_details();
if should_download {
processed_series_info_count += 1;
if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await {
if !content.is_empty() {
match serde_json::from_str::<XtreamSeriesInfo>(&content) {
Ok(info) => {
let series_stream_props = SeriesStreamProperties::from_info(&info, pli);
// the input db needs to be updated
let _ = persists_input_series_info(app_config, &storage_path, pli.header.xtream_cluster, &input.name, provider_id, &series_stream_props).await;
// Update in-memory playlist items with the newly fetched vod info.
// This makes the data available for later processing steps like STRM export.
pli.header.additional_properties = Some(StreamProperties::Series(Box::new(series_stream_props)));
}
Err(err) => {
error!("Failed to parse series info for provider_id {provider_id}: {err}");
let should_download = resolve_series && !pli.has_details();
if should_download {
processed_series_info_count += 1;
if let Some(content) = playlist_resolve_download_playlist_item(client, pli, input, errors, resolve_delay, XtreamCluster::Series).await {
if !content.is_empty() {
match serde_json::from_str::<XtreamSeriesInfo>(&content) {
Ok(info) => {
let series_stream_props = SeriesStreamProperties::from_info(&info, pli);
batch.push((provider_id, series_stream_props.clone()));
if batch.len() >= BATCH_SIZE {
if let Err(err) = persist_input_series_info_batch(app_config, &storage_path, XtreamCluster::Series, &input.name, std::mem::take(&mut batch)).await {
error!("Failed to persist batch series info: {err}");
}
}
// Update in-memory playlist items with the newly fetched vod info.
// This makes the data available for later processing steps like STRM export.
pli.header.additional_properties = Some(StreamProperties::Series(Box::new(series_stream_props)));
}
Err(err) => {
error!("Failed to parse series info for provider_id {provider_id}: {err}");
}
}
}
}
}
// extract episodes from info
if let Some(StreamProperties::Series(properties)) = pli.header.additional_properties.as_ref() {
let (group, series_name) = {
let header = &pli.header;
(header.group.clone(), if header.name.is_empty() { header.title.clone() } else { header.name.clone() })
};
if let Some(episodes) = parse_xtream_series_info(&pli.get_uuid(), properties, &group, &series_name, input) {
group_series.extend(episodes.into_iter());
}
}
if resolve_series && log_enabled!(Level::Info) && last_log_time.elapsed().as_secs() >= 30 {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
last_log_time = Instant::now();
// extract episodes from info
if let Some(StreamProperties::Series(properties)) = pli.header.additional_properties.as_ref() {
let (group, series_name) = {
let header = &pli.header;
(header.group.clone(), if header.name.is_empty() { header.title.clone() } else { header.name.clone() })
};
if let Some(episodes) = parse_xtream_series_info(&pli.get_uuid(), properties, &group, &series_name, input, interner) {
let group = group_series.entry(pli.header.category_id)
.or_insert_with(|| {
PlaylistGroup {
id: pli.header.category_id,
title: pli.header.group.clone(),
channels: Vec::new(),
xtream_cluster: XtreamCluster::Series,
}
});
group.channels.extend(episodes.into_iter());
}
}
if !group_series.is_empty() {
result.push(PlaylistGroup {
id: plg.id,
title: plg.title.clone(),
channels: group_series,
xtream_cluster: XtreamCluster::Series,
});
if resolve_series && log_enabled!(Level::Info) && last_log_time.elapsed().as_secs() >= 30 {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
last_log_time = Instant::now();
}
}
if !batch.is_empty() {
if let Err(err) = persist_input_series_info_batch(app_config, &storage_path, XtreamCluster::Series, &input.name, batch).await {
error!("Failed to persist final batch series info: {err}");
}
}
if resolve_series {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
}
result
group_series.into_values().collect()
}
#[allow(clippy::too_many_arguments)]
pub async fn playlist_resolve_series(cfg: &Arc<AppConfig>,
client: &reqwest::Client,
target: &ConfigTarget,
@@ -119,25 +134,34 @@ pub async fn playlist_resolve_series(cfg: &Arc<AppConfig>,
pipe: &ProcessingPipe,
provider_fpl: &mut FetchedPlaylist<'_>,
processed_fpl: &mut FetchedPlaylist<'_>,
interner: &mut StringInterner,
) {
let (resolve_series, resolve_delay) = get_resolve_series_options(target, processed_fpl);
let series_playlist = playlist_resolve_series_info(cfg, client, errors, processed_fpl, resolve_series, resolve_delay).await;
provider_fpl.source.release_resources(XtreamCluster::Series);
let series_playlist = playlist_resolve_series_info(cfg, client, errors, processed_fpl, resolve_series, resolve_delay, interner).await;
provider_fpl.source.obtain_resources().await;
if series_playlist.is_empty() { return; }
// original content saved into original list
for plg in &series_playlist {
provider_fpl.update_playlist(plg);
}
// run processing pipe over new items
let mut new_playlist = series_playlist;
for f in pipe {
if let Some(v) = f(&mut new_playlist, target) {
new_playlist = v;
if provider_fpl.is_memory() {
// original content saved into original list
for plg in &series_playlist {
provider_fpl.update_playlist(plg).await;
}
}
// run the processing pipe over new items
let mut new_playlist = series_playlist;
for f in pipe {
let mut source = MemoryPlaylistSource::new(new_playlist);
if let Some(v) = f(&mut source, target, interner) {
new_playlist = v;
} else {
new_playlist = source.take_groups();
}
}
// assign new items to the new playlist
for plg in &new_playlist {
processed_fpl.update_playlist(plg);
processed_fpl.update_playlist(plg).await;
}
}
+52 -37
View File
@@ -3,7 +3,7 @@ use crate::model::{AppConfig, ConfigTarget};
use crate::processing::processor::create_resolve_options_function_for_xtream_target;
use crate::processing::processor::xtream::playlist_resolve_download_playlist_item;
use crate::repository::storage::get_input_storage_path;
use crate::repository::xtream_repository::persist_input_vod_info;
use crate::repository::xtream_repository::persist_input_vod_info_batch;
use log::{error, info, log_enabled, Level};
use shared::error::TuliproxError;
use shared::model::{InputType, PlaylistEntry, StreamProperties, VideoStreamProperties, XtreamVideoInfo};
@@ -11,10 +11,16 @@ use shared::model::{PlaylistItemType, XtreamCluster};
use std::sync::Arc;
use std::time::Instant;
create_resolve_options_function_for_xtream_target!(vod);
pub async fn playlist_resolve_vod(app_config: &Arc<AppConfig>, client: &reqwest::Client,
target: &ConfigTarget, errors: &mut Vec<TuliproxError>,
const BATCH_SIZE: usize = 100;
pub async fn playlist_resolve_vod(app_config: &Arc<AppConfig>,
client: &reqwest::Client,
target: &ConfigTarget,
errors: &mut Vec<TuliproxError>,
provider_fpl: &mut FetchedPlaylist<'_>,
fpl: &mut FetchedPlaylist<'_>) {
let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl);
if !resolve_movies { return; }
@@ -29,50 +35,59 @@ pub async fn playlist_resolve_vod(app_config: &Arc<AppConfig>, client: &reqwest:
}
};
// LocalVideo entries are not resolved!
let vod_info_count = fpl.playlist_groups.iter()
.flat_map(|plg| &plg.channels)
.filter(|pli| pli.header.xtream_cluster == XtreamCluster::Video
&& pli.header.item_type == PlaylistItemType::Video
&& !pli.has_details()).count();
let vod_info_count = fpl.get_missing_vod_info_count();
info!("Found missing {vod_info_count} vod info to resolve");
info!("Found {vod_info_count} vod info to resolve");
let mut last_log_time = Instant::now();
let mut processed_vod_info_count = 0;
let mut batch = Vec::with_capacity(BATCH_SIZE);
for plg in &mut fpl.playlist_groups {
for pli in &mut plg.channels {
if pli.header.xtream_cluster != XtreamCluster::Video
|| pli.header.item_type != PlaylistItemType::Video
|| pli.has_details() {
continue;
}
let Some(provider_id) = pli.get_provider_id() else { continue; };
processed_vod_info_count += 1;
if provider_id != 0 {
if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await {
if content.is_empty() { continue; }
match serde_json::from_str::<XtreamVideoInfo>(&content) {
Ok(info) => {
let video_stream_props = VideoStreamProperties::from_info(&info, pli);
if let Err(err) = persist_input_vod_info(app_config, &storage_path, pli.header.xtream_cluster, &input.name, provider_id, &video_stream_props).await {
error!("Failed to persist VOD info for provider_id {provider_id}: {err}");
provider_fpl.source.release_resources(XtreamCluster::Video);
let input = fpl.input;
for pli in fpl.items_mut() {
if pli.header.xtream_cluster != XtreamCluster::Video
|| pli.header.item_type != PlaylistItemType::Video
|| pli.has_details() {
continue;
}
let Some(provider_id) = pli.get_provider_id() else { continue; };
processed_vod_info_count += 1;
if provider_id != 0 {
if let Some(content) = playlist_resolve_download_playlist_item(client, pli, input, errors, resolve_delay, XtreamCluster::Video).await {
if content.is_empty() { continue; }
match serde_json::from_str::<XtreamVideoInfo>(&content) {
Ok(info) => {
let video_stream_props = VideoStreamProperties::from_info(&info, pli);
batch.push((provider_id, video_stream_props.clone()));
if batch.len() >= BATCH_SIZE {
if let Err(err) = persist_input_vod_info_batch(app_config, &storage_path, XtreamCluster::Video, &input.name, std::mem::take(&mut batch)).await {
error!("Failed to persist batch VOD info: {err}");
}
// This makes the data available for subsequent processing steps like STRM export.
pli.header.additional_properties = Some(StreamProperties::Video(Box::new(video_stream_props)));
}
Err(err) => {
error!("Failed to parse video info for provider_id {provider_id}: {err}");
}
// This makes the data available for subsequent processing steps like STRM export.
pli.header.additional_properties = Some(StreamProperties::Video(Box::new(video_stream_props)));
}
Err(err) => {
error!("Failed to parse video info for provider {} stream_id {provider_id}: {err} {content}", input.name);
}
}
}
if log_enabled!(Level::Info) && last_log_time.elapsed().as_secs() >= 30 {
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
last_log_time = Instant::now();
}
}
if log_enabled!(Level::Info) && last_log_time.elapsed().as_secs() >= 30 {
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
last_log_time = Instant::now();
}
}
if !batch.is_empty() {
if let Err(err) = persist_input_vod_info_batch(app_config, &storage_path, XtreamCluster::Video, &input.name, batch).await {
error!("Failed to persist final batch VOD info: {err}");
}
}
provider_fpl.source.obtain_resources().await;
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
}