Files
tuliprox/backend/src/processing/processor/mod.rs
T
euzuandGitHub 16d32d1298 DVR Feature (#819)
feat: complete DVR and improve streaming, security, configuration, and UI

Complete the Digital Video Recorder subsystem and add a broad set of
reliability, security, streaming, configuration, processing, and Web UI
improvements across Tuliprox.

DVR:

* complete live recording and provider-aware VOD download support
* add recording queue, workers, scheduling, and recurring recording rules
* add conflict detection and capacity-aware scheduling
* add pause, resume, retry, edit, cancel, and delete workflows
* add recording quotas and configurable retention policies
* add crash recovery and startup reconciliation
* add durable lifecycle notifications with per-channel retries
* add DVR health monitoring and diagnostic tooling
* add secure access to recordings, thumbnails, and subtitles
* add WebSocket notifications for recording and rule changes
* add Web UI management for recordings, rules, progress, and task state
* add RBAC, configuration, documentation, and i18n support

Streaming and HLS:

* fix shared-stream idle handling and release dead provider streams correctly
* stop tee streams when both client and cache consumers are gone
* cancel provisioning probes when client streams terminate
* fix transient HLS origin work accounting and intermittent 503 responses
* make stream buffer byte limits configurable
* make shared subscriber idle timeout configurable
* make initial HLS manifest wait timeout configurable
* add configurable TS chunk packet count
* add configurable HLS refresh failure backoff
* centralize redirect limits and retry jitter handling
* improve provider DNS refresh behavior and failover tuning
* preserve UTF-8 characters in catchup templates
* improve stream history validation and persistence error handling

Security:

* use constant-time credential comparisons
* harden library and media path handling against traversal and symlink escapes
* only trust forwarded client IP headers from configured trusted proxies
* redact credentials and sensitive URL data from logs
* reject invalid authentication status-code configuration
* deny users with unresolved plans or invalid content filters
* improve authentication error handling across proxy and HLS endpoints

Configuration and reliability:

* prevent invalid api-proxy.yml reloads from terminating the running server
* fully validate API proxy configuration before persisting changes
* log configuration and EPG cleanup failures instead of silently discarding them
* keep the last valid configuration active after failed hot reloads
* align backend and shared media-server validation
* remove duplicated path and normalization logic
* improve DNS-store recovery and Windows rename fallback handling
* reject invalid duration, timestamp, and numeric conversions safely
* fix playlist bouquet save error handling
* fix provider record update detection
* fix cache boundary handling
* improve startup and persistence failure diagnostics

Filtering, search, sorting, and processing:

* add field-scoped playlist explorer search
* centralize shared stream-history search field definitions
* extend the filter DSL with string, set, and numeric operators
* add EPG ID, channel number, and detected quality as filterable fields
* add filter dry-run preview API with match statistics and samples
* report filter syntax errors with line and column information
* add natural numeric-aware sorting
* add quality-aware channel deduplication
* add accent-independent deduplication
* move natural sorting and quality detection helpers into shared code
* persist explorer search-field selection across reloads

User plans and content access:

* add reusable API user plans for capability tiers
* support inherited cluster and connection limits with per-user overrides
* add plan-level and user-level content filters
* enforce content filters across Xtream, M3U, direct playback, resource access,
  stream info, short EPG, categories, and XMLTV
* add trial plans with automatic expiry and Trial status
* add plan selection and content filtering to the user editor
* add full plan management to the API configuration Web UI
* migrate the API user database to schema V7 with plan and filter persistence

Web UI and accessibility:

* add live logging console to the stats page
* improve login error handling and prevent duplicate authentication requests
* add keyboard navigation to tabs, menus, tables, and search
* add ARIA roles, labels, validation state, and live-region feedback
* add confirmation dialogs for destructive actions
* add unsaved-change warnings and Ctrl/Cmd+S shortcuts
* add loading, progress, empty, and in-flight states across views
* improve dropdown and single-selection behavior
* add clipboard and credential-copy helpers
* persist table pagination and explorer search preferences
* improve error recovery when UI context providers are unavailable
* remove multiple panic-prone unwrap and browser API paths
* replace remaining hardcoded UI strings with translation keys

Maintenance:

* resolve backend and frontend compiler and Clippy warnings
* update packages and test fixtures
* consolidate duplicated helpers and validation logic
* improve documentation for configuration, filters, plans, DVR, and REST APIs
* add and update tests for migrations, filters, deduplication, sorting,
  configuration, streaming, and accessibility behavior
2026-08-21 14:50:12 +02:00

278 lines
12 KiB
Rust

mod playlist;
mod stalker;
mod stalker_refresh;
mod xtream;
// mod affix;
mod xtream_vod;
mod xtream_series;
mod deduplicate;
mod epg;
mod sort;
mod trakt;
mod library;
mod stream_probe;
mod probe_handle_guard;
mod resolve_options;
pub use self::playlist::*;
pub(crate) use self::stalker::{
download_stalker_playlist, re_resolve_stalker_url, StalkerCluster,
};
pub(crate) use self::stalker_refresh::StalkerRefreshMode;
pub use self::epg::*;
pub use self::xtream::*;
pub use self::xtream_vod::*;
pub use self::xtream_series::*;
pub use self::stream_probe::*;
pub(crate) use self::probe_handle_guard::*;
pub use self::resolve_options::*;
use crate::api::model::ProviderHandle;
use tokio_util::sync::CancellationToken;
pub(crate) const FOREGROUND_BATCH_SIZE: usize = 200;
pub(crate) const FOREGROUND_RETRY_BATCH_MAX_SIZE: usize = FOREGROUND_BATCH_SIZE * 4;
pub(crate) const FOREGROUND_MIN_RETRY_DELAY_SECS: u64 = 1;
pub(crate) fn select_cancel_token<'a>(
acquired_handle: Option<&'a ProviderHandle>,
active_handle: Option<&'a ProviderHandle>,
) -> Option<&'a CancellationToken> {
acquired_handle
.and_then(|h| h.cancel_token.as_ref())
.or_else(|| active_handle.and_then(|h| h.cancel_token.as_ref()))
}
//
// fn get_resolve_<cluster>_options(target: &ConfigTarget, fpl: &FetchedPlaylist) -> bool
//
#[macro_export]
macro_rules! create_resolve_options_function_for_xtream_target {
($cluster:ident) => {
paste::paste! {
fn [<get_resolve_ $cluster _options>](target: &ConfigTarget, fpl: &FetchedPlaylist) -> $crate::processing::processor::ResolveOptions {
match target.get_xtream_output() {
Some(_) => {
let input_options = fpl.input.options.as_ref();
let input_is_xtream = fpl.input.input_type.is_xtream();
let (
resolve_tmdb_missing,
input_resolve_enabled,
input_probe_enabled,
input_resolve_delay,
resolve_background
) = if let Some(options) = input_options {
(
options.has_flag($crate::model::ConfigInputFlags::ResolveTmdb),
options.has_flag($crate::model::ConfigInputFlags::[<Resolve $cluster:camel>]),
options.has_flag($crate::model::ConfigInputFlags::[<Probe $cluster:camel>]),
options.resolve_delay,
options.has_flag($crate::model::ConfigInputFlags::ResolveBackground)
)
} else {
(
false,
false,
false,
shared::defaults::default_resolve_delay_secs(),
shared::defaults::default_resolve_background(),
)
};
let resolve_enabled = input_resolve_enabled;
let probe_enabled = input_probe_enabled;
let resolve_delay = input_resolve_delay;
let mut flags = $crate::processing::processor::ResolveOptionsFlagsSet::new();
if resolve_enabled && input_is_xtream {
flags.set($crate::processing::processor::ResolveOptionsFlags::Resolve);
}
if resolve_tmdb_missing {
flags.set($crate::processing::processor::ResolveOptionsFlags::TmdbMissing);
}
if input_is_xtream && probe_enabled {
flags.set($crate::processing::processor::ResolveOptionsFlags::Probe);
}
if resolve_background {
flags.set($crate::processing::processor::ResolveOptionsFlags::Background);
}
$crate::processing::processor::ResolveOptions {
flags,
resolve_delay,
}
},
None => $crate::processing::processor::ResolveOptions::default(),
}
}
}
};
}
use create_resolve_options_function_for_xtream_target;
/// Foreground retry helper that retries each queued item at most once.
///
/// `retry_delay_secs` is applied sequentially per item, so the total wall-clock
/// delay is roughly `retry_delay_secs * retry_item_count` plus network/DB time.
macro_rules! process_foreground_retry_once {
(
ctx: $ctx:expr,
fpl: $fpl:expr,
filter: $filter:expr,
retry_once_ids: $retry_once_ids:ident,
retry_delay_secs: $retry_delay_secs:expr,
xtream_path: $xtream_path:ident,
db_query_holder: $db_query_holder:ident,
db_lock_holder: $db_lock_holder:ident,
batch: $batch:ident,
batch_size: $batch_size:expr,
retry_batch_max_len: $retry_batch_max_len:expr,
processed_count: $processed_count:ident,
query_error_context: $query_error_context:expr,
reasons: |$pli_reasons:ident| $reasons_expr:expr,
update: |$active_provider:ident, $pli_update:ident, $provider_id:ident, $reasons:ident, $db_query_ref:ident| $update_expr:expr,
apply_properties: |$pli_apply:ident, $updated_props:ident| $apply_expr:expr,
persist: |$updates:ident| $persist_expr:expr,
on_persist_error: |$persist_err:ident| $on_persist_error_expr:expr,
on_retry_error: |$pli_error:ident, $retry_err:ident| $on_retry_error_expr:expr,
on_after_attempt: |$pli_after:ident, $retry_succeeded:ident| $on_after_attempt_expr:expr $(,)?
) => {
for __pli in $fpl.items_mut() {
if !($filter)(__pli) {
continue;
}
let __provider_id = if let Ok(__uid) = __pli.header.id.parse::<u32>() {
crate::api::model::ProviderIdType::Id(__uid)
} else {
crate::api::model::ProviderIdType::from(&*__pli.header.id)
};
if !$retry_once_ids.remove(&__provider_id) {
continue;
}
let mut __retry_succeeded = false;
let __reasons = {
let $pli_reasons = &mut *__pli;
$reasons_expr
};
if !__reasons.is_empty() {
if let Some(__active_provider) = $ctx.provider_manager.as_ref() {
// Do not hold a read lock over the retry delay window.
if $db_query_holder.is_some() {
$db_query_holder = None;
$db_lock_holder = None;
}
tokio::time::sleep(std::time::Duration::from_secs($retry_delay_secs)).await;
if $db_query_holder.is_none() && $xtream_path.exists() {
let __file_lock = $ctx.config.file_locks.read_lock(&$xtream_path).await;
let __xtream_path = $xtream_path.clone();
let __query = match tokio::task::spawn_blocking(move || {
crate::repository::BPlusTreeQuery::<u32, shared::model::XtreamPlaylistItem>::try_new(
&__xtream_path,
)
})
.await
{
Ok(Ok(__query)) => Some((__query, __file_lock)),
Ok(Err(__err)) => {
log::error!("Failed to open BPlusTreeQuery for {}: {__err}", $query_error_context);
None
}
Err(__err) => {
log::error!("Failed to open BPlusTreeQuery for {}: {__err}", $query_error_context);
None
}
};
if let Some((__query, __guard)) = __query {
$db_query_holder = Some(std::sync::Arc::new(parking_lot::Mutex::new(__query)));
$db_lock_holder = Some(__guard);
}
}
let __db_query_ref = $db_query_holder.as_ref().map(std::sync::Arc::clone);
let __update_future = {
let $active_provider = __active_provider;
let $pli_update = &mut *__pli;
let $provider_id = __provider_id.clone();
let $reasons = &__reasons;
let $db_query_ref = __db_query_ref;
$update_expr
};
match __update_future.await {
Ok(Some(__updated_props)) => {
{
let $pli_apply = &mut *__pli;
let $updated_props = &__updated_props;
$apply_expr
}
$batch.push((__provider_id.clone(), __updated_props));
if $batch.len() >= $batch_size {
$db_query_holder = None;
$db_lock_holder = None;
let __updates: Vec<(u32, _)> = $batch
.iter()
.filter_map(|(__id, __props)| {
// Foreground retry batches can include text provider IDs (e.g. M3U).
// Persist batch functions for these paths are keyed by numeric Xtream IDs,
// so text IDs are intentionally skipped here.
if let crate::api::model::ProviderIdType::Id(__vid) = __id {
Some((*__vid, __props.clone()))
} else {
None
}
})
.collect();
if __updates.is_empty() {
$batch.clear();
} else {
let __persist_future = {
let $updates = __updates;
$persist_expr
};
match __persist_future.await {
Ok(()) => $batch.clear(),
Err($persist_err) => {
$on_persist_error_expr;
if $batch.len() > $retry_batch_max_len {
let __drop_count = $batch.len().saturating_sub($retry_batch_max_len);
if __drop_count > 0 {
$batch.drain(0..__drop_count);
}
}
}
}
}
}
$processed_count += 1;
__retry_succeeded = true;
}
Ok(None) => {}
Err($retry_err) => {
let $pli_error = &*__pli;
$on_retry_error_expr;
}
}
}
}
{
let $pli_after = &mut *__pli;
let $retry_succeeded = __retry_succeeded;
$on_after_attempt_expr;
}
}
};
}
pub(crate) use process_foreground_retry_once;