diff --git a/CHANGELOG.md b/CHANGELOG.md index dcd7b50b6..403e5cf23 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -111,8 +111,16 @@ ## 🌟 New Features -- **Smart Connection Priority**: Introduced a priority system for connections. Users with higher priority can preempt (kick) lower priority - connections (e.g., background tasks or standard users) when provider slots are full. +- **User Connection Priority**: API users now carry a `priority` field (type `i8`, nice-style: lower value = higher priority, default `0`, probe `127`). + When all provider slots are occupied and a higher-priority user connects, the lowest-priority active connection on that provider is + evicted (oldest first when tied). Only connections with exactly one active listener are eligible for eviction; shared connections + with multiple listeners are not interrupted. Equal priority never evicts equal priority β€” the new connection is rejected normally + (with grace-period rules applied as before). User `max_connections` limits are unaffected. +- **Configurable Probe Priority**: Stream-probe tasks (`probe_live`, `probe_vod`, `probe_series`) now run with a configurable priority + instead of a fixed internal constant. Set `metadata_update.probe.user_priority` (default `127`, i.e. lowest priority) to control how + aggressively active users can preempt probe connections. +- **User DB Schema Migration V3**: The `api_user.db` file is automatically upgraded to V3 format (adds `priority` field) on first startup. + A `.userdb_mergeto_v3` guard file is created so config-driven user merges are skipped while the DB is the authoritative source. - **Background Metadata Queue**: Metadata resolution (VOD/Series) and stream analysis are now queued per input and processed in the background when provider connections are idle. This prevents "No Connections" errors for active users during playlist updates. - **Stream Probing**: Added support for probing streams (`probe_live|vod|series`) to determine codecs and resolution. This runs as a low-priority @@ -161,7 +169,12 @@ active URL of the specified provider. ## βš™οΈ New Settings +- **api-proxy.yml / Web UI (user)**: + - Added `priority` (`i8`, default `0`) to user credentials. Lower value = higher priority (nice-style). + Configurable via Web UI user editor. Negative values are valid and represent higher-than-default priority. - **config.yml**: + - Added `metadata_update.probe.user_priority` (`i8`, default `127`): priority assigned to probe connections. + Probe tasks run at the lowest priority by default; reduce this value to give probes more connection access. - Added `metadata_update` (optional) with grouped sections: `log`, `resolve`, `probe`, `ffprobe`, `tmdb`. - Added `metadata_update.cache_path` (default `metadata`): shared storage directory for TMDB cache and metadata files (moved from `library.metadata.path`). diff --git a/README.md b/README.md index c72a838d9..74a5aaad1 100644 --- a/README.md +++ b/README.md @@ -469,6 +469,7 @@ metadata_update: retry_backoff_step_3: 1h max_attempts: 3 backoff_jitter_percent: 20 + user_priority: 127 # Connection priority for probe tasks (default 127 = lowest priority) tmdb: enabled: true # api_key: "..." # Optional, fallback is internal default placeholder @@ -516,6 +517,9 @@ metadata_update: - `probe.retry_backoff_step_2` (default `30m`): Probe backoff delay for attempt 2. - `probe.retry_backoff_step_3` (default `1h`): Probe backoff delay for attempt 3 and higher. - `probe.backoff_jitter_percent` (default `20`): Random jitter percentage applied to resolve/probe retry backoff to avoid synchronized retries. +- `probe.user_priority` (default `127`): Connection priority assigned to probe tasks. Uses the same nice-style scale as user priorities + (lower value = higher priority). At `127` (lowest priority), probe connections are always the first to be evicted when a regular user connects. + Reduce this value (e.g. `64`) to give probes more access to provider slots. - `tmdb.cooldown` (default `7d`): Cooldown duration after a TMDB lookup completed successfully but returned no match. - `tmdb.enabled` / `tmdb.api_key` / `tmdb.rate_limit_ms` / `tmdb.cache_duration_days` / `tmdb.language`: TMDB resolver settings. - `tmdb.match_threshold` (default `86`): TMDB match threshold for search results for TMDB ID resolution. @@ -2675,6 +2679,27 @@ Reverse Proxy mode for user can be a subset: - `user_ui_enabled` is _optional_. If defined it can be `true` or `false`. Default is `true`. Disable/enable web_ui for user - `user_access_control` is _optional_. If defined it can be `true` or `false`. Default is `false`. +### User Priority + +Each user credential accepts an optional `priority` field (`i8`, default `0`). +The priority uses a **nice-style scale**: a **lower value means higher priority**. Negative values are allowed and represent even higher priority. +Priority range: `-128` - `127`, where `-128` has highest priority. + +When all provider connection slots are occupied and a new user with **strictly higher priority** (lower value) connects, the lowest-priority +active connection on that provider is evicted (oldest first when priorities are tied). Only connections with exactly one active listener are +eligible for eviction β€” shared connections with multiple listeners are not interrupted. A user with equal or lower priority than all existing +connections is rejected normally β€” existing grace-period rules still apply. + +| Priority | Meaning | +|:--------:|--------------------------------------------| +| `-128` | highest possible priority | +| `-10` | Very high β€” almost always preempts others | +| `0` | Default β€” standard user | +| `64` | Reduced β€” yields to default users | +| `127` | Lowest β€” default priority for probe tasks | + +`max_connections` per user is independent of priority and unaffected by eviction. + If you have a lot of users and dont want to keep them in `api-proxy.yml`, you can set the option - `use_user_db` to true to store the user information inside a db-file. @@ -2742,6 +2767,7 @@ user: exp_date: 1672705545 max_connections: 1 status: Active + priority: 0 # optional, default 0; lower = higher priority (nice-style, negative allowed) ``` If you use a reverse proxy in front of Tuliprox, don’t forget to forward: diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 91a671cf2..a3fc62169 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -346,6 +346,7 @@ async fn resolve_streaming_strategy( input: &ConfigInput, force_provider: Option<&Arc>, allow_provider_grace: bool, + user_priority: i8, ) -> StreamingStrategy { // allocate a provider connection let mut forced_provider_allocated = false; @@ -355,7 +356,7 @@ async fn resolve_streaming_strategy( // If that account is no longer available, fall back to any available account in the same lineup. if let Some(handle) = app_state .active_provider - .acquire_exact_connection_with_grace(provider, &fingerprint.addr, allow_provider_grace) + .acquire_exact_connection_with_grace(provider, &fingerprint.addr, allow_provider_grace, user_priority) .await { forced_provider_allocated = true; @@ -368,14 +369,14 @@ async fn resolve_streaming_strategy( ); app_state .active_provider - .acquire_connection_with_grace(&input.name, &fingerprint.addr, allow_provider_grace) + .acquire_connection_with_grace(&input.name, &fingerprint.addr, allow_provider_grace, user_priority) .await } } None => { app_state .active_provider - .acquire_connection_with_grace(&input.name, &fingerprint.addr, allow_provider_grace) + .acquire_connection_with_grace(&input.name, &fingerprint.addr, allow_provider_grace, user_priority) .await } }; @@ -462,9 +463,10 @@ async fn create_stream_response_details( force_provider: Option<&Arc>, allow_provider_grace: bool, virtual_id: VirtualId, + user_priority: i8, ) -> Result { let mut streaming_strategy = - resolve_streaming_strategy(app_state, stream_url, fingerprint, input, force_provider, allow_provider_grace) + resolve_streaming_strategy(app_state, stream_url, fingerprint, input, force_provider, allow_provider_grace, user_priority) .await; let mut grace_period_options = app_state.get_grace_options(); grace_period_options.period_millis = get_grace_period_millis( @@ -783,6 +785,7 @@ pub async fn force_provider_stream_response( preferred_provider, allow_provider_grace, stream_channel.virtual_id, + user.priority, ) .await { @@ -934,6 +937,7 @@ pub async fn stream_response( None, true, stream_channel.virtual_id, + user.priority, ) .await { @@ -1737,6 +1741,7 @@ pub fn create_api_proxy_user(app_state: &Arc) -> ProxyUserCredentials status: None, ui_enabled: false, comment: None, + priority: 0, t_is_api_user: true, } } diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index e66f1f0ea..dc7761c51 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -107,7 +107,7 @@ pub(in crate::api) async fn handle_hls_stream_request( Some(session) => { let handle = if let Some(handle) = app_state .active_provider - .acquire_exact_connection_with_grace(&session.provider, &fingerprint.addr, false) + .acquire_exact_connection_with_grace(&session.provider, &fingerprint.addr, false, user.priority) .await { Some(handle) @@ -117,7 +117,7 @@ pub(in crate::api) async fn handle_hls_stream_request( sanitize_sensitive_info(&session.provider), sanitize_sensitive_info(&fingerprint.addr.to_string()) ); - app_state.active_provider.acquire_connection_with_grace(&input.name, &fingerprint.addr, false).await + app_state.active_provider.acquire_connection_with_grace(&input.name, &fingerprint.addr, false, user.priority).await }; match handle { diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 29c31937d..f0fee65dc 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -21,8 +21,6 @@ use std::{ use tokio::sync::{RwLock, Semaphore}; use tokio_util::sync::CancellationToken; -const DEFAULT_USER_PRIORITY: i8 = -1; -const DEFAULT_PROBE_PRIORITY: i8 = 1; const PREEMPTED_PROBE_CANCEL_GRACE: Duration = Duration::from_secs(2); const PREEMPTED_GRACE_MAX_PENDING: usize = 64; static DUMMY_ADDR: LazyLock = LazyLock::new(|| "127.0.0.1:0".parse::().unwrap()); @@ -393,26 +391,28 @@ impl ActiveProviderManager { provider_name: &Arc, addr: &SocketAddr, allow_grace: bool, + priority: i8, ) -> Option { let allocation = self.providers.acquire_exact_connection_with_grace_override(provider_name, allow_grace).await; if matches!(allocation, ProviderAllocation::Exhausted) { return None; } - self.register_allocation(allocation, addr, DEFAULT_USER_PRIORITY, false).await + self.register_allocation(allocation, addr, priority, false).await } pub async fn force_exact_acquire_connection( &self, provider_name: &Arc, addr: &SocketAddr, + priority: i8, ) -> Option { // Compatibility wrapper: keep the exact-provider behavior but do not over-allocate exhausted accounts. - self.acquire_exact_connection_with_grace(provider_name, addr, false).await + self.acquire_exact_connection_with_grace(provider_name, addr, false, priority).await } // Returns the next available provider connection - pub async fn acquire_connection(&self, input_name: &Arc, addr: &SocketAddr) -> Option { - self.acquire_connection_inner(input_name, addr, false, None, DEFAULT_USER_PRIORITY, false).await + pub async fn acquire_connection(&self, input_name: &Arc, addr: &SocketAddr, priority: i8) -> Option { + self.acquire_connection_inner(input_name, addr, false, None, priority, false).await } /// Acquire a provider connection while explicitly controlling provider-side grace allocations. @@ -421,14 +421,15 @@ impl ActiveProviderManager { input_name: &Arc, addr: &SocketAddr, allow_grace: bool, + priority: i8, ) -> Option { - self.acquire_connection_inner(input_name, addr, false, Some(allow_grace), DEFAULT_USER_PRIORITY, false).await + self.acquire_connection_inner(input_name, addr, false, Some(allow_grace), priority, false).await } - /// Acquire a provider connection while optionally disabling provider grace allocations. - pub async fn acquire_connection_for_probe(&self, input_name: &Arc) -> Option { - // Probe is strictly low-priority and must never consume grace capacity. - self.acquire_connection_inner(input_name, &DUMMY_ADDR, false, Some(false), DEFAULT_PROBE_PRIORITY, true).await + /// Acquire a provider connection for probe tasks with configurable priority. + /// Probes never consume grace capacity. + pub async fn acquire_connection_for_probe(&self, input_name: &Arc, priority: i8) -> Option { + self.acquire_connection_inner(input_name, &DUMMY_ADDR, false, Some(false), priority, true).await } // This method is used for redirects to cycle through the provider @@ -681,6 +682,7 @@ mod tests { utils::Internable, }; use std::{collections::HashMap, net::SocketAddr, sync::Arc, time::Duration}; + use shared::utils::{default_probe_user_priority, default_user_priority}; fn build_test_app_config(aliases: Option>, max_connections: u16) -> AppConfig { let input = Arc::new(ConfigInput { @@ -757,12 +759,12 @@ mod tests { let client_2_addr: SocketAddr = "127.0.0.1:40002".parse().unwrap(); let first_alloc = - manager.acquire_connection(&input_name, &client_1_addr).await.expect("client1 initial allocation"); + manager.acquire_connection(&input_name, &client_1_addr, default_user_priority()).await.expect("client1 initial allocation"); let pinned_provider = first_alloc.allocation.get_provider_name().expect("provider name expected"); assert_eq!(pinned_provider.as_ref(), "provider_1"); // provider_1 has max_connections=1 and is already in use by client1 - let forced = manager.force_exact_acquire_connection(&pinned_provider, &client_2_addr).await; + let forced = manager.force_exact_acquire_connection(&pinned_provider, &client_2_addr, default_user_priority()).await; assert!(forced.is_none(), "forced exact acquire must not over-allocate busy provider"); manager.release_connection(&client_1_addr).await; @@ -781,7 +783,7 @@ mod tests { // Step 1: Client1 starts movie -> provider_1 let first_alloc = - manager.acquire_connection(&input_name, &client_1_addr).await.expect("client1 initial allocation"); + manager.acquire_connection(&input_name, &client_1_addr, default_user_priority()).await.expect("client1 initial allocation"); assert_eq!(first_alloc.allocation.get_provider_name().as_deref(), Some(input_name.as_ref())); // Step 2: Client1 stops -> release provider_1 @@ -789,7 +791,7 @@ mod tests { // Step 3: Client2 starts live -> provider_1 let live_alloc = - manager.acquire_connection(&input_name, &client_2_addr).await.expect("client2 live allocation"); + manager.acquire_connection(&input_name, &client_2_addr, default_user_priority()).await.expect("client2 live allocation"); let busy_provider = live_alloc.allocation.get_provider_name().expect("provider name expected"); assert_eq!(busy_provider.as_ref(), input_name.as_ref()); assert!(manager.is_exhausted(&busy_provider).await); @@ -797,7 +799,7 @@ mod tests { // Step 4: Client1 restarts same movie. // This emulates force-session fallback path by acquiring without provider grace. let fallback_alloc = manager - .acquire_connection_with_grace(&input_name, &client_1_addr, false) + .acquire_connection_with_grace(&input_name, &client_1_addr, false, 0) .await .expect("client1 fallback allocation without grace"); let fallback_provider = fallback_alloc.allocation.get_provider_name().expect("fallback provider expected"); @@ -821,12 +823,12 @@ mod tests { // Initial playback for client1. let first_alloc = - manager.acquire_connection(&input_name, &client_1_addr).await.expect("client1 initial allocation"); + manager.acquire_connection(&input_name, &client_1_addr, default_user_priority()).await.expect("client1 initial allocation"); let pinned_provider = first_alloc.allocation.get_provider_name().expect("provider name expected"); assert_eq!(pinned_provider.as_ref(), "provider_1"); // Another client occupies the alternate account while client1 keeps seeking. - let second_alloc = manager.acquire_connection(&input_name, &client_2_addr).await.expect("client2 allocation"); + let second_alloc = manager.acquire_connection(&input_name, &client_2_addr, default_user_priority()).await.expect("client2 allocation"); let second_provider = second_alloc.allocation.get_provider_name().expect("provider name expected"); assert_eq!(second_provider.as_ref(), "provider_2"); @@ -835,7 +837,7 @@ mod tests { for _ in 0..3 { manager.release_connection(&client_1_addr).await; let seek_alloc = manager - .force_exact_acquire_connection(&pinned_provider, &client_1_addr) + .force_exact_acquire_connection(&pinned_provider, &client_1_addr, default_user_priority()) .await .expect("seek reacquire should stay on pinned provider"); let seek_provider = seek_alloc.allocation.get_provider_name().expect("provider name expected"); @@ -857,12 +859,12 @@ mod tests { let user_addr: SocketAddr = "127.0.0.1:43001".parse().unwrap(); let probe_handle = - manager.acquire_connection_for_probe(&input_name).await.expect("probe allocation should succeed"); + manager.acquire_connection_for_probe(&input_name, default_probe_user_priority()).await.expect("probe allocation should succeed"); let probe_token = probe_handle.cancel_token.clone().expect("probe handle must carry cancel token"); // User request should preempt probe and immediately acquire released capacity. let user_alloc = manager - .acquire_connection_with_grace(&input_name, &user_addr, false) + .acquire_connection_with_grace(&input_name, &user_addr, false, default_user_priority()) .await .expect("user allocation should preempt probe"); assert_eq!(user_alloc.allocation.get_provider_name().as_deref(), Some(input_name.as_ref())); @@ -887,4 +889,92 @@ mod tests { manager.release_connection(&user_addr).await; } + + #[tokio::test] + async fn test_higher_priority_user_preempts_lower_priority_user() { + // User with priority 5 (low) is connected; user with priority -1 (high) arrives. + // The low-priority user should be preempted. + let app_cfg = create_test_app_config_single_provider_pool(); + let event_manager = Arc::new(EventManager::new()); + let manager = ActiveProviderManager::new(&app_cfg, &event_manager); + + let input_name = "provider_1".intern(); + let low_prio_addr: SocketAddr = "127.0.0.1:44001".parse().unwrap(); + let high_prio_addr: SocketAddr = "127.0.0.1:44002".parse().unwrap(); + + // Low-priority user connects (priority 5 = lower importance) + let low_alloc = manager + .acquire_connection(&input_name, &low_prio_addr, 5) + .await + .expect("low-priority user should get connection"); + assert_eq!(low_alloc.allocation.get_provider_name().as_deref(), Some(input_name.as_ref())); + + // Provider is now exhausted + assert!(manager.is_exhausted(&input_name).await); + + // High-priority user arrives (priority -1 = higher importance), should preempt low-priority user + let high_alloc = manager + .acquire_connection_with_grace(&input_name, &high_prio_addr, false, -1) + .await + .expect("high-priority user should preempt low-priority user and get connection"); + assert_eq!(high_alloc.allocation.get_provider_name().as_deref(), Some(input_name.as_ref())); + + manager.release_connection(&high_prio_addr).await; + } + + #[tokio::test] + async fn test_same_priority_user_does_not_preempt() { + // Two users with the same priority β€” new one should NOT preempt the existing one. + let app_cfg = create_test_app_config_single_provider_pool(); + let event_manager = Arc::new(EventManager::new()); + let manager = ActiveProviderManager::new(&app_cfg, &event_manager); + + let input_name = "provider_1".intern(); + let user_1_addr: SocketAddr = "127.0.0.1:45001".parse().unwrap(); + let user_2_addr: SocketAddr = "127.0.0.1:45002".parse().unwrap(); + + // User 1 connects with priority 0 + let alloc1 = manager + .acquire_connection(&input_name, &user_1_addr, default_user_priority()) + .await + .expect("user1 should get connection"); + assert_eq!(alloc1.allocation.get_provider_name().as_deref(), Some(input_name.as_ref())); + + // Provider is now exhausted + assert!(manager.is_exhausted(&input_name).await); + + // User 2 arrives with the same priority 0 β€” should NOT preempt user 1 + let alloc2 = manager.acquire_connection_with_grace(&input_name, &user_2_addr, false, default_user_priority()).await; + assert!(alloc2.is_none(), "same-priority user should not preempt existing user"); + + manager.release_connection(&user_1_addr).await; + } + + #[tokio::test] + async fn test_lower_priority_user_does_not_preempt_higher_priority_user() { + // User with high priority is connected; user with low priority arrives β€” should NOT preempt. + let app_cfg = create_test_app_config_single_provider_pool(); + let event_manager = Arc::new(EventManager::new()); + let manager = ActiveProviderManager::new(&app_cfg, &event_manager); + + let input_name = "provider_1".intern(); + let high_prio_addr: SocketAddr = "127.0.0.1:46001".parse().unwrap(); + let low_prio_addr: SocketAddr = "127.0.0.1:46002".parse().unwrap(); + + // High-priority user connects (priority -10) + let alloc1 = manager + .acquire_connection(&input_name, &high_prio_addr, -10) + .await + .expect("high-priority user should get connection"); + assert_eq!(alloc1.allocation.get_provider_name().as_deref(), Some(input_name.as_ref())); + + // Provider is now exhausted + assert!(manager.is_exhausted(&input_name).await); + + // Low-priority user arrives (priority 10) β€” should NOT preempt high-priority user + let alloc2 = manager.acquire_connection_with_grace(&input_name, &low_prio_addr, false, 10).await; + assert!(alloc2.is_none(), "low-priority user should not preempt high-priority user"); + + manager.release_connection(&high_prio_addr).await; + } } diff --git a/backend/src/api/model/metadata_update_manager.rs b/backend/src/api/model/metadata_update_manager.rs index de0f3bbae..10053de43 100644 --- a/backend/src/api/model/metadata_update_manager.rs +++ b/backend/src/api/model/metadata_update_manager.rs @@ -38,6 +38,7 @@ use std::{ }; use tokio::sync::{mpsc, RwLock, Semaphore}; use tokio_util::sync::CancellationToken; +use shared::utils::default_probe_user_priority; const METADATA_RETRY_STATE_FILE: &str = "metadata_retry_state.db"; const TASK_ERR_NO_CONNECTION: &str = "No connection available"; @@ -2748,9 +2749,17 @@ impl InputWorker { let needs_probe_connection = Self::task_needs_provider_connection(task, input_base.input_type); + let probe_priority = app_state + .app_config + .config + .load() + .metadata_update + .as_ref() + .map_or(default_probe_user_priority(), |cfg| cfg.probe.user_priority); + // Reserve provider capacity only for actual probe work (ffprobe paths). let provider_handle = if needs_probe_connection { - let Some(handle) = app_state.active_provider.acquire_connection_for_probe(input_name).await else { + let Some(handle) = app_state.active_provider.acquire_connection_for_probe(input_name, probe_priority).await else { debug_if_enabled!("No provider connection available for background task {}, skipping...", task); return Err(shared::error::info_err!("{}", TASK_ERR_NO_CONNECTION)); }; @@ -2800,7 +2809,7 @@ impl InputWorker { Err(shared::error::info_err!("{}", TASK_ERR_PREEMPTED)) } - res = Self::execute_task_inner_static(&app_state, &client, &input_to_use, task, item_title.as_deref(), Some(handle), collector, db_handles, failed_clusters) => { + res = Self::execute_task_inner_static(&app_state, &client, &input_to_use, task, item_title.as_deref(), Some(handle), probe_priority, collector, db_handles, failed_clusters) => { res } } @@ -2812,6 +2821,7 @@ impl InputWorker { task, item_title.as_deref(), Some(handle), + probe_priority, collector, db_handles, failed_clusters, @@ -2826,6 +2836,7 @@ impl InputWorker { task, item_title.as_deref(), None, + probe_priority, collector, db_handles, failed_clusters, @@ -2880,6 +2891,7 @@ impl InputWorker { task: &UpdateTask, item_title: Option<&str>, active_handle: Option<&ProviderHandle>, + probe_priority: i8, collector: &mut BatchResultCollector, db_handles: &mut HashMap, failed_clusters: &mut HashSet, @@ -3021,6 +3033,7 @@ impl InputWorker { *item_type, &app_state.active_provider, active_handle, + probe_priority, ) .await?; diff --git a/backend/src/main.rs b/backend/src/main.rs index c5608e9b2..87f4553be 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -21,7 +21,7 @@ use crate::{ library::LibraryProcessor, model::{AppConfig, Config, Healthcheck, HealthcheckConfig, ProcessTargets, SourcesConfig}, processing::processor::exec_processing, - repository::migrate_bplustree_databases_with_marker, + repository::run_startup_migrations, utils::{config_file_reader, db_viewer, init_logger, request::create_client, resolve_env_var}, }; use arc_swap::{access::Access, ArcSwap}; @@ -162,7 +162,7 @@ async fn main() { std::process::exit(i32::from(!healthy)); } - run_startup_bplustree_migration(&config_paths); + run_startup_migrations(&config_paths); // Handle Library scan before starting main application if args.scan_library || args.force_library_rescan { @@ -215,42 +215,6 @@ async fn main() { } } -fn run_startup_bplustree_migration(config_paths: &ConfigPaths) { - let config_file_path = Path::new(config_paths.config_file_path.as_str()); - if !config_file_path.exists() { - return; - } - - let mut roots: Vec = vec![PathBuf::from(&config_paths.config_path)]; - let storage_dir = PathBuf::from(&config_paths.storage_path); - if !roots.iter().any(|root| root == &storage_dir) { - roots.push(storage_dir.clone()); - } - - match migrate_bplustree_databases_with_marker(&roots, &storage_dir) { - Ok(stats) => { - if stats.skipped_by_marker { - info!("B+Tree startup migration skipped (marker already present)"); - } else if stats.migrated_files > 0 { - info!( - "B+Tree startup migration completed: migrated {} file(s) ({} B+Tree files checked, {} .db files scanned)", - stats.migrated_files, - stats.bplustree_files, - stats.scanned_files - ); - //} else { - // info!( - // "B+Tree startup migration check completed: {} B+Tree files already current ({} .db files scanned)", - // stats.bplustree_files, stats.scanned_files - // ); - } - } - Err(err) => { - exit!("B+Tree startup migration failed: {err}"); - } - } -} - fn print_info(app_config: &AppConfig) { let config = > as Access>::load(&app_config.config); let paths = > as Access>::load(&app_config.paths); diff --git a/backend/src/model/config/api_user.rs b/backend/src/model/config/api_user.rs index ef6a33cf3..6fa99d803 100644 --- a/backend/src/model/config/api_user.rs +++ b/backend/src/model/config/api_user.rs @@ -23,6 +23,7 @@ pub struct ProxyUserCredentials { pub status: Option, pub ui_enabled: bool, pub comment: Option, + pub priority: i8, pub t_is_api_user: bool, } @@ -43,6 +44,7 @@ impl From<&ProxyUserCredentialsDto> for ProxyUserCredentials { status: dto.status, ui_enabled: dto.ui_enabled, comment: dto.comment.clone(), + priority: dto.priority, t_is_api_user: false, } } @@ -64,6 +66,7 @@ impl From<&ProxyUserCredentials> for ProxyUserCredentialsDto { status: instance.status, ui_enabled: instance.ui_enabled, comment: instance.comment.clone(), + priority: instance.priority, } } } diff --git a/backend/src/model/config/metadata_update.rs b/backend/src/model/config/metadata_update.rs index a7e053c22..64ac8e374 100644 --- a/backend/src/model/config/metadata_update.rs +++ b/backend/src/model/config/metadata_update.rs @@ -54,6 +54,7 @@ pub struct ProbeConfig { pub retry_backoff_step_3_secs: u64, pub max_attempts: u8, pub backoff_jitter_percent: u8, + pub user_priority: i8, } #[derive(Debug, Clone)] @@ -269,6 +270,7 @@ impl From<&ProbeConfigDto> for ProbeConfig { retry_backoff_step_3: dto.retry_backoff_step_3.clone(), max_attempts: dto.max_attempts.max(1), backoff_jitter_percent: dto.backoff_jitter_percent.min(95), + user_priority: dto.user_priority, } } } @@ -283,6 +285,7 @@ impl From<&ProbeConfig> for ProbeConfigDto { retry_backoff_step_3: instance.retry_backoff_step_3.clone(), max_attempts: instance.max_attempts, backoff_jitter_percent: instance.backoff_jitter_percent, + user_priority: instance.user_priority, } } } diff --git a/backend/src/processing/processor/stream_probe.rs b/backend/src/processing/processor/stream_probe.rs index 7ddd4b447..767a82fec 100644 --- a/backend/src/processing/processor/stream_probe.rs +++ b/backend/src/processing/processor/stream_probe.rs @@ -40,6 +40,7 @@ pub async fn update_generic_stream_metadata( item_type: PlaylistItemType, active_provider: &Arc, active_handle: Option<&crate::api::model::ProviderHandle>, + probe_priority: i8, ) -> Result { let storage_dir = &app_config.config.load().storage_dir; @@ -89,7 +90,7 @@ pub async fn update_generic_stream_metadata( let acquired_handle = if !needs_provider_connection || active_handle.is_some() { None } else { - active_provider.acquire_connection_for_probe(&input.name).await + active_provider.acquire_connection_for_probe(&input.name, probe_priority).await }; if needs_provider_connection && active_handle.is_none() && acquired_handle.is_none() { diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 7aa13aecd..21f1bf8ad 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -33,6 +33,7 @@ use std::collections::{HashMap, HashSet}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::time::{Duration, Instant}; +use shared::utils::default_probe_user_priority; create_resolve_options_function_for_xtream_target!(series); @@ -848,6 +849,7 @@ pub async fn update_series_metadata( if let Some(details) = properties.details.as_mut() { if let Some(episodes) = details.episodes.as_mut() { let config = app_config.config.load(); + let probe_priority = config.metadata_update.as_ref().map_or(default_probe_user_priority(), |cfg| cfg.probe.user_priority); let user_agent = config.default_user_agent.clone(); let input_url = input.url.as_str(); @@ -868,7 +870,7 @@ pub async fn update_series_metadata( let temp_handle = if active_handle.is_some() { None } else { - active_provider.acquire_connection_for_probe(&input.name).await + active_provider.acquire_connection_for_probe(&input.name, probe_priority).await }; if active_handle.is_some() || temp_handle.is_some() { diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 526c6e966..b3bc3fa88 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -30,6 +30,7 @@ use std::collections::{HashMap, HashSet}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::time::{Duration, Instant}; +use shared::utils::default_probe_user_priority; create_resolve_options_function_for_xtream_target!(vod); @@ -814,11 +815,13 @@ pub async fn update_vod_metadata( let analyze_duration = metadata_update.ffprobe.analyze_duration_micros; let probe_size = metadata_update.ffprobe.probe_size_bytes; + let probe_priority = config.metadata_update.as_ref().map_or(default_probe_user_priority(), |cfg| cfg.probe.user_priority); + // Acquire Connection logic let temp_handle = if active_handle.is_some() { None // No new handle needed } else { - active_provider.acquire_connection_for_probe(&input.name).await + active_provider.acquire_connection_for_probe(&input.name, probe_priority).await }; if active_handle.is_some() || temp_handle.is_some() { diff --git a/backend/src/repository/bplustree_migration.rs b/backend/src/repository/bplustree_migration.rs index c776b4e25..eaa50c614 100644 --- a/backend/src/repository/bplustree_migration.rs +++ b/backend/src/repository/bplustree_migration.rs @@ -1,6 +1,8 @@ -use super::bplustree::{MAGIC, STORAGE_VERSION}; +use super::bplustree::{BPlusTree, MAGIC, STORAGE_VERSION}; +use super::storage_const; use fs2::FileExt as _; -use log::warn; +use log::{info, warn}; +use shared::model::{ConfigPaths, ProxyType, ProxyUserStatus}; use std::{ collections::{HashSet, VecDeque}, ffi::OsStr, @@ -15,8 +17,9 @@ const METADATA_MAX_SIZE: u32 = 4000; const HEADER_FLAG_HAS_METADATA_FLAGS: u32 = 1 << 31; const HEADER_FLAG_HAS_TOMBSTONES: u32 = 1 << 30; const HEADER_METADATA_LEN_MASK: u32 = !(HEADER_FLAG_HAS_METADATA_FLAGS | HEADER_FLAG_HAS_TOMBSTONES); -const MARKER_FILE_PREFIX: &str = ".db_mergeto"; -const LEGACY_MARKER_FILE_PREFIX_ALT: &str = ".db_mergedto"; +const MARKER_FILE_GUARD_PREFIX: &str = ".db_mergeto_v"; +const MARKER_FILE_GUARD_PREFIX_LEGACY_ALT: &str = ".db_mergedto"; +const MARKER_FILE_API_USER_GUARD: &str = ".userdb_mergeto_v3"; const MARKER_VERSION_KEY: &str = "migrated_to"; const MARKER_ROOTS_FINGERPRINT_KEY: &str = "roots_fingerprint"; @@ -29,7 +32,7 @@ pub struct BPlusTreeMigrationStats { } #[derive(Debug)] -pub struct BPlusTreeStartupMigrator { +struct BPlusTreeStartupMigrator { roots: Vec, migration_marker_path: Option, } @@ -77,48 +80,9 @@ impl BPlusTreeStartupMigrator { Ok(stats) } - fn resolve_scan_roots(roots: &[PathBuf]) -> Vec { - let mut resolved: Vec = Vec::new(); - let mut seen: HashSet = HashSet::new(); - - for root in roots { - if !root.exists() || !root.is_dir() { - continue; - } - let resolved_root = std::fs::canonicalize(root).unwrap_or_else(|_| root.clone()); - let key = resolved_root.to_string_lossy().into_owned(); - if seen.insert(key) { - resolved.push(resolved_root); - } - } - - resolved - } - - fn roots_fingerprint(roots: &[PathBuf]) -> String { - const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325; - const FNV_PRIME: u64 = 0x0000_0100_0000_01b3; - - let mut canonical_entries: Vec = roots.iter().map(|path| path.to_string_lossy().into_owned()).collect(); - canonical_entries.sort_unstable(); - canonical_entries.dedup(); - - let mut hash = FNV_OFFSET_BASIS; - for entry in &canonical_entries { - for byte in entry.as_bytes() { - hash ^= u64::from(*byte); - hash = hash.wrapping_mul(FNV_PRIME); - } - hash ^= 0xff; - hash = hash.wrapping_mul(FNV_PRIME); - } - - format!("{hash:016x}") - } - fn cleanup_legacy_root_markers(&self, keep_marker_path: Option<&Path>) -> io::Result<()> { let marker_name = marker_file_name(); - let marker_name_alt = marker_file_name_alt(); + let marker_name_alt = format!("{MARKER_FILE_GUARD_PREFIX_LEGACY_ALT}{STORAGE_VERSION}"); let mut visited_roots: HashSet = HashSet::new(); for root in &self.roots { @@ -172,6 +136,53 @@ impl BPlusTreeStartupMigrator { } } + fn resolve_scan_roots(roots: &[PathBuf]) -> Vec { + let mut resolved: Vec = Vec::new(); + + // Normalize paths (canonicalize where possible) + for root in roots { + if !root.exists() || !root.is_dir() { + continue; + } + // Try to resolve the absolute/real path, fall back to the original path on failure + let canon = std::fs::canonicalize(root).unwrap_or_else(|_| root.clone()); + resolved.push(canon); + } + + // Sort paths so parent directories come before child directories + resolved.sort(); + resolved.dedup(); + + // Keep only top-level directories, remove descendants + let mut final_roots: Vec = Vec::new(); + for path in resolved { + // If this path is already covered by a parent in `final_roots`, skip it + if final_roots.iter().any(|parent| path.starts_with(parent)) { + continue; + } + final_roots.push(path); + } + + final_roots + } + + fn roots_fingerprint(roots: &[PathBuf]) -> String { + const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325; + const FNV_PRIME: u64 = 0x0000_0100_0000_01b3; + + let mut hash = FNV_OFFSET_BASIS; + for path in roots { + for byte in path.to_string_lossy().as_bytes() { + hash ^= u64::from(*byte); + hash = hash.wrapping_mul(FNV_PRIME); + } + hash ^= 0xff; + hash = hash.wrapping_mul(FNV_PRIME); + } + + format!("{hash:016x}") + } + fn collect_db_files_for_root(root: &Path) -> io::Result> { let mut files: Vec = Vec::new(); let mut queue: VecDeque = VecDeque::new(); @@ -377,9 +388,243 @@ pub fn migrate_bplustree_databases_with_marker( BPlusTreeStartupMigrator::new_with_marker(roots.to_vec(), marker_path).run() } -fn marker_file_name() -> String { format!("{MARKER_FILE_PREFIX}{STORAGE_VERSION}") } +fn marker_file_name() -> String { format!("{MARKER_FILE_GUARD_PREFIX}{STORAGE_VERSION}") } -fn marker_file_name_alt() -> String { format!("{LEGACY_MARKER_FILE_PREFIX_ALT}{STORAGE_VERSION}") } +// ─── User DB schema migration ───────────────────────────────────────────────── +// +// The user database has gone through three serialization schemas (MessagePack, +// positional/sequence encoding via rmp_serde): +// +// V1 (Deprecated) – original format, 13 fields, no epg_request_timeshift +// V2 – 14 fields, added epg_request_timeshift +// V3 (current) – 15 fields, added priority +// +// On first startup after an upgrade the file is still in V1 or V2 format. +// `migrate_user_db_schema` detects this, converts every record in-place, and +// writes a merge-guard marker so that config-driven user merges cannot +// overwrite the freshly migrated data. + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +struct StoredApiUserV1 { + pub target: String, + pub username: String, + pub password: String, + pub token: Option, + pub proxy: ProxyType, + pub server: Option, + pub epg_timeshift: Option, + pub created_at: Option, + pub exp_date: Option, + pub max_connections: Option, + pub status: Option, + pub ui_enabled: bool, + pub comment: Option, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +struct StoredApiUserV2 { + pub target: String, + pub username: String, + pub password: String, + pub token: Option, + pub proxy: ProxyType, + pub server: Option, + pub epg_timeshift: Option, + pub epg_request_timeshift: Option, + pub created_at: Option, + pub exp_date: Option, + pub max_connections: Option, + pub status: Option, + pub ui_enabled: bool, + pub comment: Option, +} + +// V3 mirror β€” same layout as user_repository::StoredProxyUserCredentials. +// Defined here so the migration has no dependency on user_repository internals. +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +struct StoredApiUserV3 { + pub target: String, + pub username: String, + pub password: String, + pub token: Option, + pub proxy: ProxyType, + pub server: Option, + pub epg_timeshift: Option, + pub epg_request_timeshift: Option, + pub created_at: Option, + pub exp_date: Option, + pub max_connections: Option, + pub status: Option, + pub ui_enabled: bool, + pub comment: Option, + pub priority: Option, +} + +impl StoredApiUserV3 { + fn from_v2(v2: &StoredApiUserV2) -> Self { + Self { + target: v2.target.clone(), + username: v2.username.clone(), + password: v2.password.clone(), + token: v2.token.clone(), + proxy: v2.proxy, + server: v2.server.clone(), + epg_timeshift: v2.epg_timeshift.clone(), + epg_request_timeshift: v2.epg_request_timeshift.clone(), + created_at: v2.created_at, + exp_date: v2.exp_date, + max_connections: v2.max_connections, + status: v2.status, + ui_enabled: v2.ui_enabled, + comment: v2.comment.clone(), + priority: None, + } + } + + fn from_v1(v1: &StoredApiUserV1) -> Self { + Self { + target: v1.target.clone(), + username: v1.username.clone(), + password: v1.password.clone(), + token: v1.token.clone(), + proxy: v1.proxy, + server: v1.server.clone(), + epg_timeshift: v1.epg_timeshift.clone(), + epg_request_timeshift: None, + created_at: v1.created_at, + exp_date: v1.exp_date, + max_connections: v1.max_connections, + status: v1.status, + ui_enabled: v1.ui_enabled, + comment: v1.comment.clone(), + priority: None, + } + } +} + +fn create_user_db_merge_guard(merge_guard_path: &Path) -> io::Result<()> { + if !merge_guard_path.exists() { + std::fs::write(merge_guard_path, b"")?; + } + Ok(()) +} + +pub(crate) fn user_db_merge_guard_path(config_dir: &Path) -> PathBuf { + config_dir.join(MARKER_FILE_API_USER_GUARD) +} + +/// Migrates the user database file from V1 or V2 schema to V3 (current) in +/// place and creates a merge-guard file so config-driven merges are skipped +/// until the operator explicitly removes it. +/// +/// Returns `true` when a migration was performed, `false` when the file was +/// already in V3 format or did not exist. +fn migrate_user_db_schema(db_path: &Path, merge_guard_path: &Path) -> io::Result { + if !db_path.exists() { + return Ok(false); + } + + // Already V3? (15-field sequence β€” rmp_serde fails on shorter sequences) + if BPlusTree::::load(db_path).is_ok() { + return Ok(false); + } + + // Try V2 β†’ V3 + if let Ok(tree) = BPlusTree::::load(db_path) { + let mut v3_tree: BPlusTree = BPlusTree::new(); + for (key, v2) in &tree { + v3_tree.insert(key.clone(), StoredApiUserV3::from_v2(v2)); + } + v3_tree.store(db_path)?; + create_user_db_merge_guard(merge_guard_path)?; + return Ok(true); + } + + // Try V1 β†’ V3 + if let Ok(tree) = BPlusTree::::load(db_path) { + let mut v3_tree: BPlusTree = BPlusTree::new(); + for (key, v1) in &tree { + v3_tree.insert(key.clone(), StoredApiUserV3::from_v1(v1)); + } + v3_tree.store(db_path)?; + create_user_db_merge_guard(merge_guard_path)?; + return Ok(true); + } + + Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("User DB at '{}' exists but could not be read as V1, V2, or V3 format", db_path.display()), + )) +} + +// ─── Combined startup migration ─────────────────────────────────────────────── + +#[derive(Debug, Clone, Copy, Default)] +pub struct AllStartupMigrationStats { + pub bplustree: BPlusTreeMigrationStats, + pub user_db_migrated: bool, +} + +/// Runs all startup migrations in sequence: +/// 1. B+Tree storage-format migration (V1 β†’ current binary format) +/// 2. User DB schema migration (V1/V2 β†’ V3 `MessagePack` layout) +/// +/// `config_dir` is the directory that contains `api_user.db` and the merge-guard +/// marker. `storage_dir` is used for the B+Tree migration marker. +fn run_all_startup_migrations( + roots: &[PathBuf], + storage_dir: &Path, + config_dir: &Path, +) -> io::Result { + let marker_path = bplustree_migration_marker_path(storage_dir); + let bplustree = BPlusTreeStartupMigrator::new_with_marker(roots.to_vec(), marker_path).run()?; + + let user_db_path = config_dir.join(storage_const::API_USER_DB_FILE); + let merge_guard_path = user_db_merge_guard_path(config_dir); + let user_db_migrated = migrate_user_db_schema(&user_db_path, &merge_guard_path)?; + + Ok(AllStartupMigrationStats { bplustree, user_db_migrated }) +} + + +pub fn run_startup_migrations(config_paths: &ConfigPaths) { + let config_file_path = Path::new(config_paths.config_file_path.as_str()); + if !config_file_path.exists() { + return; + } + + let config_dir = PathBuf::from(&config_paths.config_path); + let storage_dir = if config_paths.storage_path.trim().is_empty() { + config_dir.clone() + } else { + PathBuf::from(&config_paths.storage_path) + }; + let mut roots: Vec = vec![config_dir.clone()]; + if storage_dir != config_dir { + roots.push(storage_dir.clone()); + } + + match run_all_startup_migrations(&roots, &storage_dir, &config_dir) { + Ok(stats) => { + if stats.bplustree.skipped_by_marker { + info!("B+Tree startup migration skipped (marker already present)"); + } else if stats.bplustree.migrated_files > 0 { + info!( + "B+Tree startup migration completed: migrated {} file(s) ({} B+Tree files checked, {} .db files scanned)", + stats.bplustree.migrated_files, + stats.bplustree.bplustree_files, + stats.bplustree.scanned_files + ); + } + if stats.user_db_migrated { + info!("User DB schema migrated to V3"); + } + } + Err(err) => { + crate::utils::exit!("Startup migration failed: {err}"); + } + } +} #[cfg(test)] mod tests { @@ -526,7 +771,7 @@ mod tests { let fingerprint = test_roots_fingerprint(&roots); BPlusTreeStartupMigrator::write_migration_marker(&global_marker, &fingerprint)?; - let legacy_per_root_marker = temp_other.path().join(format!("{MARKER_FILE_PREFIX}{STORAGE_VERSION}")); + let legacy_per_root_marker = temp_other.path().join(format!("{MARKER_FILE_GUARD_PREFIX}{STORAGE_VERSION}")); BPlusTreeStartupMigrator::write_migration_marker(&legacy_per_root_marker, "legacy")?; assert!(legacy_per_root_marker.exists()); diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index d00d828cf..ba9d06415 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -14,48 +14,8 @@ use std::io::Error; use std::path::{Path, PathBuf}; use tokio::task; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -struct StoredProxyUserCredentialsDeprecated { - pub target: String, - pub username: String, - pub password: String, - pub token: Option, - pub proxy: ProxyType, - pub server: Option, - pub epg_timeshift: Option, - pub created_at: Option, - pub exp_date: Option, - pub max_connections: Option, - pub status: Option, - pub ui_enabled: bool, - pub comment: Option, -} - -impl StoredProxyUserCredentialsDeprecated { - fn to(stored: &StoredProxyUserCredentialsDeprecated) -> ProxyUserCredentials { - ProxyUserCredentials { - username: stored.username.clone(), - password: stored.password.clone(), - token: stored.token.clone(), - proxy: stored.proxy, - server: stored.server.clone(), - epg_timeshift: stored.epg_timeshift.clone(), - epg_request_timeshift: None, - created_at: stored.created_at, - exp_date: stored.exp_date, - max_connections: stored.max_connections.unwrap_or_default(), - status: stored.status, - ui_enabled: stored.ui_enabled, - comment: stored.comment.clone(), - t_is_api_user: false, - } - } -} - -// This is a Helper class to store all user into one Database file. -// For the Config files we keep the old structure where a user is assigned to a target. -// But for storing inside one db file it is easier to store the target next to the user. -// due to known issue with bincode and skip_serialization_if we have to list all fields and can't use ProxyUserCredentials +// V3 (current): added priority field. V1 and V2 are migrated to V3 at startup +// by `bplustree_migration::run_all_startup_migrations`. #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] struct StoredProxyUserCredentials { pub target: String, @@ -72,6 +32,7 @@ struct StoredProxyUserCredentials { pub status: Option, pub ui_enabled: bool, pub comment: Option, + pub priority: Option, } impl StoredProxyUserCredentials { @@ -91,6 +52,7 @@ impl StoredProxyUserCredentials { status: proxy.status, ui_enabled: proxy.ui_enabled, comment: proxy.comment.clone(), + priority: if proxy.priority != 0 { Some(proxy.priority) } else { None }, } } @@ -109,6 +71,7 @@ impl StoredProxyUserCredentials { status: stored.status, ui_enabled: stored.ui_enabled, comment: stored.comment.clone(), + priority: stored.priority.unwrap_or(0), t_is_api_user: false, } } @@ -177,56 +140,29 @@ pub async fn store_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Res result } -// TODO remove me if we get stable on user_db -pub async fn load_api_user_deprecated(cfg: &AppConfig) -> Result, Error> { - let path = get_api_user_db_path(cfg); - let lock = cfg.file_locks.read_lock(&path).await; - let user_tree = BPlusTree::::load(&path)?; - drop(lock); +fn collect_target_users(user_tree: &BPlusTree) -> Vec { let mut target_users: HashMap = HashMap::new(); - for (_uname, stored_user) in &user_tree { - let proxy_user: ProxyUserCredentials = StoredProxyUserCredentialsDeprecated::to(stored_user); - let target_name = stored_user.target.clone(); - match target_users.entry(target_name) { - std::collections::hash_map::Entry::Occupied(mut entry) => { - let target = entry.get_mut(); - target.credentials.push(proxy_user); - } - std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert(TargetUser { - target: stored_user.target.clone(), - credentials: vec![proxy_user], - }); - } - } - } - Ok(target_users.into_values().collect()) -} - - -pub async fn load_api_user(cfg: &AppConfig) -> Result, Error> { - let path = get_api_user_db_path(cfg); - let lock = cfg.file_locks.read_lock(&path).await; - let Ok(user_tree) = BPlusTree::::load(&path) else { return load_api_user_deprecated(cfg).await }; - drop(lock); - let mut target_users: HashMap = HashMap::new(); - for (_uname, stored_user) in &user_tree { + for (_uname, stored_user) in user_tree { let proxy_user: ProxyUserCredentials = StoredProxyUserCredentials::to(stored_user); let target_name = stored_user.target.clone(); match target_users.entry(target_name) { std::collections::hash_map::Entry::Occupied(mut entry) => { - let target = entry.get_mut(); - target.credentials.push(proxy_user); + entry.get_mut().credentials.push(proxy_user); } std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert(TargetUser { - target: stored_user.target.clone(), - credentials: vec![proxy_user], - }); + entry.insert(TargetUser { target: stored_user.target.clone(), credentials: vec![proxy_user] }); } } } - Ok(target_users.into_values().collect()) + target_users.into_values().collect() +} + +pub async fn load_api_user(cfg: &AppConfig) -> Result, Error> { + let path = get_api_user_db_path(cfg); + let lock = cfg.file_locks.read_lock(&path).await; + let result = BPlusTree::::load(&path); + drop(lock); + result.map(|tree| collect_target_users(&tree)).map_err(|err| Error::other(format!("Failed to load user db: {err}"))) } pub fn get_user_storage_path(cfg: &Config, username: &str) -> Option { @@ -458,6 +394,7 @@ mod tests { status: Some(ProxyUserStatus::Active), ui_enabled: true, comment: None, + priority: 0, t_is_api_user: false, }, ProxyUserCredentials { @@ -474,6 +411,7 @@ mod tests { status: Some(ProxyUserStatus::Expired), ui_enabled: true, comment: None, + priority: 0, t_is_api_user: false, }, ProxyUserCredentials { @@ -490,6 +428,7 @@ mod tests { status: Some(ProxyUserStatus::Expired), ui_enabled: true, comment: None, + priority: 0, t_is_api_user: false, }, ProxyUserCredentials { @@ -506,6 +445,7 @@ mod tests { status: Some(ProxyUserStatus::Expired), ui_enabled: true, comment: None, + priority: -10, // non-zero priority to verify round-trip serialization t_is_api_user: false, } ], @@ -541,6 +481,10 @@ mod tests { let user_list = load_api_user(&cfg).await; assert!(user_list.is_ok()); assert_eq!(user_list.as_ref().unwrap().len(), 1); - assert_eq!(user_list.as_ref().unwrap().first().unwrap().credentials.len(), 4); + let loaded = user_list.as_ref().unwrap().first().unwrap(); + assert_eq!(loaded.credentials.len(), 4); + // Verify non-zero priority survives the store/load round-trip. + let test4 = loaded.credentials.iter().find(|c| c.username == "Test4").unwrap(); + assert_eq!(test4.priority, -10); } } diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index d8cfce852..01a8a684e 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -492,7 +492,8 @@ "RETRY_BACKOFF_STEP_1": "Wait time after the first probe failure.\n\nUse this as a gentle first retry step.\n\nExample: `10m`.", "RETRY_BACKOFF_STEP_2": "Wait time after the second probe failure.\n\nUsually set higher than step 1.\n\nExample: `30m`.", "RETRY_BACKOFF_STEP_3": "Wait time after the third and later probe failures.\n\nThis is the longest regular probe delay before cooldown takes over.\n\nExample: `1h`.", - "RETRY_LOAD_RETRY_DELAY": "Wait time before retrying to load saved retry-state data after a read failure.\n\nThis helps when storage is temporarily busy or unavailable.\n\nExamples: `30s`, `1m`, `5m`." + "RETRY_LOAD_RETRY_DELAY": "Wait time before retrying to load saved retry-state data after a read failure.\n\nThis helps when storage is temporarily busy or unavailable.\n\nExamples: `30s`, `1m`, `5m`.", + "USER_PRIORITY": "Connection priority assigned to probe tasks. lowest priority = 127 (default), highest priority = -128. Do not set less than 1 because default user priority is 0." }, "PROXY_CONFIG": { "PASSWORD": "Password or credential used for this account or service.", @@ -511,7 +512,8 @@ "SERVER": "Target proxy server address or definition.", "STATUS": "The active or disabled state of the user proxy credential.", "TOKEN": "Unique API token that can be used instead of username/password.", - "UI_ENABLED": "If true, this user can log into the simplified WebUI bouquet editor (default true)." + "UI_ENABLED": "If true, this user can log into the simplified WebUI bouquet editor (default true).", + "PRIORITY": "Connection priority assigned to this user proxy context. Default = 0, lowest priority = 127, highest priority = -128." }, "RATE_LIMIT_CONFIG": { "BURST_SIZE": "Defines the initial number of available connections before throttling applies (e.g. 10).", @@ -918,6 +920,7 @@ "METADATA_TMDB_COOLDOWN": "TMDB cooldown", "METADATA_UPDATE_CONFIG": "Metadata Update", "METADATA_WORKER_IDLE_TIMEOUT": "Worker idle timeout", + "METADATA_PROBE_USER_PRIORITY": "Probe user priority", "METHOD": "Method", "MODEL_NAME": "Model Name", "MODEL_NUMBER": "Model Number", diff --git a/frontend/scss/app/components/_table.scss b/frontend/scss/app/components/_table.scss index 729b1cc79..ec381350b 100644 --- a/frontend/scss/app/components/_table.scss +++ b/frontend/scss/app/components/_table.scss @@ -69,4 +69,8 @@ } } } + + &__number-cell { + text-align: right; + } } diff --git a/frontend/src/app/components/cell_value.rs b/frontend/src/app/components/cell_value.rs index 0d54f920b..bd25e8485 100644 --- a/frontend/src/app/components/cell_value.rs +++ b/frontend/src/app/components/cell_value.rs @@ -7,6 +7,7 @@ pub enum CellValue<'a> { Proxy(ProxyType), Text(&'a str), Date(i64), + I8(i8), } impl<'a> PartialOrd for CellValue<'a> { @@ -32,9 +33,12 @@ impl<'a> Ord for CellValue<'a> { std::cmp::Ordering::Greater } (CellValue::Proxy(_), _) => std::cmp::Ordering::Less, - (CellValue::Text(_), CellValue::Date(_)) => std::cmp::Ordering::Less, + (CellValue::I8(a), CellValue::I8(b)) => a.cmp(b), + (CellValue::Text(_), CellValue::Date(_) | CellValue::I8(_)) => std::cmp::Ordering::Less, (CellValue::Text(_), _) => std::cmp::Ordering::Greater, + (CellValue::Date(_), CellValue::I8(_)) => std::cmp::Ordering::Less, (CellValue::Date(_), _) => std::cmp::Ordering::Greater, + (CellValue::I8(_), _) => std::cmp::Ordering::Greater, } } } diff --git a/frontend/src/app/components/config/macros.rs b/frontend/src/app/components/config/macros.rs index d9564a54d..18c82eae6 100644 --- a/frontend/src/app/components/config/macros.rs +++ b/frontend/src/app/components/config/macros.rs @@ -410,6 +410,32 @@ macro_rules! edit_field_number_u16 { }}; } +#[macro_export] +macro_rules! edit_field_number_i8 { + ($instance:expr, $label:expr, $field:ident, $action:path) => {{ + let instance = $instance.clone(); + html! { +
+ <$crate::app::components::number_input::NumberInput + label={$label} + name={stringify!($field)} + field_id={Some($crate::app::components::dto_field_id(&instance.form, stringify!($field)))} + value={instance.form.$field as i64} + on_change={Callback::from(move |value: Option| { + match value { + Some(value) => match i8::try_from(value) { + Ok(val) => instance.dispatch($action(val)), + Err(_) => return, // keep the existing value + }, + None => {} // cleared input β€” keep the existing value + } + })} + /> +
+ } + }}; +} + #[macro_export] macro_rules! edit_field_number_i16 { ($instance:expr, $label:expr, $field:ident, $action:path) => {{ diff --git a/frontend/src/app/components/config/metadata_update_config_view.rs b/frontend/src/app/components/config/metadata_update_config_view.rs index 742ee979a..a84e02bcb 100644 --- a/frontend/src/app/components/config/metadata_update_config_view.rs +++ b/frontend/src/app/components/config/metadata_update_config_view.rs @@ -11,7 +11,7 @@ use crate::{ }, context::ConfigContext, }, - config_field, config_field_bool, config_field_optional, edit_field_bool, edit_field_number, + config_field, config_field_bool, config_field_optional, edit_field_bool, edit_field_number, edit_field_number_i8, edit_field_number_option_u64, edit_field_number_u16, edit_field_number_u64, edit_field_number_u8, edit_field_number_usize, edit_field_text, edit_field_text_option, generate_form_reducer, i18n::use_translation, @@ -39,6 +39,7 @@ const LABEL_PROBE_RETRY_BACKOFF_STEP_3: &str = "LABEL.METADATA_PROBE_RETRY_BACKO const LABEL_MAX_ATTEMPTS_RESOLVE: &str = "LABEL.METADATA_MAX_ATTEMPTS_RESOLVE"; const LABEL_MAX_ATTEMPTS_PROBE: &str = "LABEL.METADATA_MAX_ATTEMPTS_PROBE"; const LABEL_BACKOFF_JITTER_PERCENT: &str = "LABEL.METADATA_BACKOFF_JITTER_PERCENT"; +const LABEL_PROBE_USER_PRIORITY: &str = "LABEL.METADATA_PROBE_USER_PRIORITY"; const LABEL_MAX_QUEUE_SIZE: &str = "LABEL.METADATA_MAX_QUEUE_SIZE"; const LABEL_METADATA_PATH: &str = "LABEL.METADATA_PATH"; const LABEL_FFPROBE_ENABLED: &str = "LABEL.FFPROBE_ENABLED"; @@ -105,6 +106,7 @@ generate_form_reducer!( RetryBackoffStep3 => retry_backoff_step_3: String, MaxAttempts => max_attempts: u8, BackoffJitterPercent => backoff_jitter_percent: u8, + UserPriority => user_priority: i8, } ); @@ -244,6 +246,7 @@ pub fn MetadataUpdateConfigView() -> Html {

{translate.t(LABEL_PROBE)}

{ config_field!(probe, translate.t(LABEL_MAX_ATTEMPTS_PROBE), max_attempts) } { config_field!(probe, translate.t(LABEL_BACKOFF_JITTER_PERCENT), backoff_jitter_percent) } + { config_field!(probe, translate.t(LABEL_PROBE_USER_PRIORITY), user_priority) } { config_field!(probe, translate.t(LABEL_PROBE_RETRY_BACKOFF_STEP_1), retry_backoff_step_1) } { config_field!(probe, translate.t(LABEL_PROBE_RETRY_BACKOFF_STEP_2), retry_backoff_step_2) } { config_field!(probe, translate.t(LABEL_PROBE_RETRY_BACKOFF_STEP_3), retry_backoff_step_3) } @@ -315,6 +318,7 @@ pub fn MetadataUpdateConfigView() -> Html {

{translate.t(LABEL_PROBE)}

{ edit_field_number_u8!(probe_state, translate.t(LABEL_MAX_ATTEMPTS_PROBE), max_attempts, ProbeConfigFormAction::MaxAttempts) } + { edit_field_number_i8!(probe_state, translate.t(LABEL_PROBE_USER_PRIORITY), user_priority, ProbeConfigFormAction::UserPriority) }
server: Option, Status => status: Option, MaxConnections => max_connections: u32, + Priority => priority: i8, ExpDate => exp_date: Option, UiEnabled => ui_enabled: bool, EpgTimeshift => epg_timeshift: Option, @@ -251,6 +252,7 @@ pub fn ProxyUserCredentialsForm(props: &ProxyUserCredentialsFormProps) -> Html { /> }})} { edit_field_number!(form_state, translate.t("LABEL.MAX_CONNECTIONS"), max_connections, UserFormAction::MaxConnections) } + { edit_field_number_i8!(form_state, translate.t("LABEL.PRIORITY"), priority, UserFormAction::Priority) } { edit_field_date!(form_state, translate.t("LABEL.EXP_DATE"), exp_date, UserFormAction::ExpDate) } { edit_field_text_option!(form_state, translate.t("LABEL.EPG_TIMESHIFT"), epg_timeshift, UserFormAction::EpgTimeshift) } { edit_field_text_option!(form_state, translate.t("LABEL.EPG_REQUEST_TIMESHIFT"), epg_request_timeshift, UserFormAction::EpgRequestTimeshift) } diff --git a/frontend/src/app/components/userlist/user_table.rs b/frontend/src/app/components/userlist/user_table.rs index 7c523b3e5..144bb3df4 100644 --- a/frontend/src/app/components/userlist/user_table.rs +++ b/frontend/src/app/components/userlist/user_table.rs @@ -22,7 +22,7 @@ use std::{cmp::Ordering, collections::HashSet, fmt::Display, rc::Rc, str::FromSt use yew::{platform::spawn_local, prelude::*}; use yew_hooks::use_clipboard; -const HEADERS: [&str; 16] = [ +const HEADERS: [&str; 17] = [ "LABEL.EMPTY", "LABEL.ENABLED", "LABEL.STATUS", @@ -33,6 +33,7 @@ const HEADERS: [&str; 16] = [ "LABEL.PROXY", "LABEL.SERVER", "LABEL.MAX_CON", + "LABEL.PRIORITY", "LABEL.UI_ENABLED", "LABEL.EPG_TIMESHIFT", "LABEL.EPG_REQUEST_TIMESHIFT", @@ -49,13 +50,14 @@ fn get_cell_value(user: &TargetUser, col: usize) -> CellValue<'_> { 4 => CellValue::Text(user.credentials.username.as_str()), 7 => CellValue::Proxy(user.credentials.proxy), 8 => user.credentials.server.as_ref().map_or(CellValue::Empty, |s| CellValue::Text(s)), - 13 => user.credentials.created_at.as_ref().map_or(CellValue::Empty, |d| CellValue::Date(*d)), - 14 => user.credentials.exp_date.as_ref().map_or(CellValue::Empty, |d| CellValue::Date(*d)), + 10 => CellValue::I8(user.credentials.priority), + 14 => user.credentials.created_at.as_ref().map_or(CellValue::Empty, |d| CellValue::Date(*d)), + 15 => user.credentials.exp_date.as_ref().map_or(CellValue::Empty, |d| CellValue::Date(*d)), _ => CellValue::Empty, } } -fn is_col_sortable(col: usize) -> bool { matches!(col, 1 | 2 | 3 | 4 | 7 | 8 | 13 | 14) } +fn is_col_sortable(col: usize) -> bool { matches!(col, 1 | 2 | 3 | 4 | 7 | 8 | 10 | 14 | 15) } #[derive(Debug, Clone, Eq, PartialEq)] enum TableAction { @@ -202,17 +204,18 @@ pub fn UserTable(props: &UserTableProps) -> Html { 7 => html! { }, 8 => dto.credentials.server.as_ref().map_or_else(|| html! {}, |s| html! { s }), 9 => html! { }, - 10 => html! { html! { { dto.credentials.priority } }, + 11 => html! { }, - 11 => dto.credentials.epg_timeshift.as_ref().map_or_else(|| html! {}, |s| html! { s }), - 12 => dto.credentials.epg_request_timeshift.as_ref().map_or_else(|| html! {}, |s| html! { s }), - 13 => dto.credentials.created_at.as_ref().and_then(|ts| unix_ts_to_str(*ts)) + 12 => dto.credentials.epg_timeshift.as_ref().map_or_else(|| html! {}, |s| html! { s }), + 13 => dto.credentials.epg_request_timeshift.as_ref().map_or_else(|| html! {}, |s| html! { s }), + 14 => dto.credentials.created_at.as_ref().and_then(|ts| unix_ts_to_str(*ts)) .map(|s| html! { { s } }).unwrap_or_else(|| html! { }), - 14 => dto.credentials.exp_date.as_ref().and_then(|ts| unix_ts_to_str(*ts)) + 15 => dto.credentials.exp_date.as_ref().and_then(|ts| unix_ts_to_str(*ts)) .map(|s| html! { { s } }) .unwrap_or_else(|| html! { }), - 15 => dto.credentials.comment.as_ref() + 16 => dto.credentials.comment.as_ref() .map_or_else(|| html! {}, |comment| html! { {comment} }), _ => html! {""}, diff --git a/shared/src/model/config/api_user.rs b/shared/src/model/config/api_user.rs index 1bc72fd63..eb83a31eb 100644 --- a/shared/src/model/config/api_user.rs +++ b/shared/src/model/config/api_user.rs @@ -1,7 +1,10 @@ use crate::{ error::{TuliproxError, TuliproxErrorKind}, model::{ProxyType, ProxyUserStatus}, - utils::{default_as_true, deserialize_timestamp, is_blank_optional_string, is_true}, + utils::{ + default_as_true, default_user_priority, deserialize_timestamp, is_blank_optional_string, + is_default_user_priority, is_true, + }, }; #[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] @@ -38,6 +41,8 @@ pub struct ProxyUserCredentialsDto { pub ui_enabled: bool, #[serde(default, skip_serializing_if = "is_blank_optional_string")] pub comment: Option, + #[serde(default = "default_user_priority", skip_serializing_if = "is_default_user_priority")] + pub priority: i8, } impl ProxyUserCredentialsDto { diff --git a/shared/src/model/config/metadata_update.rs b/shared/src/model/config/metadata_update.rs index 747d1c3a6..21618bb2a 100644 --- a/shared/src/model/config/metadata_update.rs +++ b/shared/src/model/config/metadata_update.rs @@ -13,19 +13,20 @@ use crate::{ default_metadata_progress_log_interval, default_metadata_queue_log_interval, default_metadata_resolve_exhaustion_reset_gap, default_metadata_resolve_min_retry_base, default_metadata_retry_delay, default_metadata_tmdb_cooldown, default_metadata_worker_idle_timeout, - default_tmdb_api_key, default_tmdb_cache_duration_days, default_tmdb_language, default_tmdb_match_threshold, - default_tmdb_rate_limit_ms, deserialize_as_string, is_default_metadata_backoff_jitter_percent, - is_default_metadata_ffprobe_analyze_duration, is_default_metadata_ffprobe_live_analyze_duration, - is_default_metadata_ffprobe_live_probe_size, is_default_metadata_ffprobe_probe_size, - is_default_metadata_max_attempts_probe, is_default_metadata_max_attempts_resolve, - is_default_metadata_max_queue_size, is_default_metadata_max_resolve_retry_backoff, - is_default_metadata_no_change_cache_ttl_secs, is_default_metadata_path, is_default_metadata_probe_cooldown, - is_default_metadata_probe_fairness_resolve_burst, is_default_metadata_probe_retry_backoff_step_1, - is_default_metadata_probe_retry_backoff_step_2, is_default_metadata_probe_retry_backoff_step_3, - is_default_metadata_probe_retry_load_retry_delay, is_default_metadata_progress_log_interval, - is_default_metadata_queue_log_interval, is_default_metadata_resolve_exhaustion_reset_gap, - is_default_metadata_resolve_min_retry_base, is_default_metadata_retry_delay, is_default_metadata_tmdb_cooldown, - is_default_metadata_worker_idle_timeout, is_default_tmdb_cache_duration_days, is_default_tmdb_language, + default_probe_user_priority, default_tmdb_api_key, default_tmdb_cache_duration_days, default_tmdb_language, + default_tmdb_match_threshold, default_tmdb_rate_limit_ms, deserialize_as_string, + is_default_metadata_backoff_jitter_percent, is_default_metadata_ffprobe_analyze_duration, + is_default_metadata_ffprobe_live_analyze_duration, is_default_metadata_ffprobe_live_probe_size, + is_default_metadata_ffprobe_probe_size, is_default_metadata_max_attempts_probe, + is_default_metadata_max_attempts_resolve, is_default_metadata_max_queue_size, + is_default_metadata_max_resolve_retry_backoff, is_default_metadata_no_change_cache_ttl_secs, + is_default_metadata_path, is_default_metadata_probe_cooldown, is_default_metadata_probe_fairness_resolve_burst, + is_default_metadata_probe_retry_backoff_step_1, is_default_metadata_probe_retry_backoff_step_2, + is_default_metadata_probe_retry_backoff_step_3, is_default_metadata_probe_retry_load_retry_delay, + is_default_metadata_progress_log_interval, is_default_metadata_queue_log_interval, + is_default_metadata_resolve_exhaustion_reset_gap, is_default_metadata_resolve_min_retry_base, + is_default_metadata_retry_delay, is_default_metadata_tmdb_cooldown, is_default_metadata_worker_idle_timeout, + is_default_probe_user_priority, is_default_tmdb_cache_duration_days, is_default_tmdb_language, is_default_tmdb_match_threshold, is_default_tmdb_rate_limit_ms, is_false, is_tmdb_default_api_key, parse_duration_seconds, parse_size_base_2, TMDB_API_KEY, }, @@ -252,6 +253,8 @@ pub struct ProbeConfigDto { skip_serializing_if = "is_default_metadata_backoff_jitter_percent" )] pub backoff_jitter_percent: u8, + #[serde(default = "default_probe_user_priority", skip_serializing_if = "is_default_probe_user_priority")] + pub user_priority: i8, } impl Default for ProbeConfigDto { @@ -264,6 +267,7 @@ impl Default for ProbeConfigDto { retry_backoff_step_3: default_metadata_probe_retry_backoff_step_3(), max_attempts: default_metadata_max_attempts_probe(), backoff_jitter_percent: default_metadata_backoff_jitter_percent(), + user_priority: default_probe_user_priority(), } } } @@ -277,6 +281,7 @@ impl ProbeConfigDto { && self.retry_backoff_step_3 == default_metadata_probe_retry_backoff_step_3() && self.max_attempts == default_metadata_max_attempts_probe() && self.backoff_jitter_percent == default_metadata_backoff_jitter_percent() + && self.user_priority == default_probe_user_priority() } fn prepare(&mut self) -> Result<(), TuliproxError> { diff --git a/shared/src/utils/default_utils.rs b/shared/src/utils/default_utils.rs index 75debc98c..db038a3e5 100644 --- a/shared/src/utils/default_utils.rs +++ b/shared/src/utils/default_utils.rs @@ -228,6 +228,10 @@ pub fn default_metadata_ffprobe_live_probe_size() -> String { "5MB".to_string() pub fn is_default_metadata_ffprobe_live_probe_size(v: &String) -> bool { *v == default_metadata_ffprobe_live_probe_size() } +pub fn default_probe_user_priority() -> i8 { 127 } +pub fn is_default_probe_user_priority(v: &i8) -> bool { *v == default_probe_user_priority() } +pub fn default_user_priority() -> i8 { 0 } +pub fn is_default_user_priority(v: &i8) -> bool { *v == default_user_priority() } pub fn get_default_web_root() -> String { DEFAULT_WEB_DIR.to_string() } pub fn is_blank_or_default_web_root(value: &str) -> bool {