From b3e7e2cd8973632756baf76a95d16975847bffb0 Mon Sep 17 00:00:00 2001 From: euzu Date: Sun, 3 Aug 2025 10:24:23 +0200 Subject: [PATCH] Playlist update added to web ui --- Cargo.lock | 54 +- backend/src/api/endpoints/v1_api.rs | 4 +- backend/src/api/endpoints/websocket_api.rs | 9 + backend/src/api/main_api.rs | 18 +- .../src/api/model/active_provider_manager.rs | 483 +++++++++--------- backend/src/api/model/app_state.rs | 2 +- backend/src/api/model/event_manager.rs | 5 +- backend/src/api/model/mod.rs | 2 +- backend/src/api/scheduler.rs | 42 +- backend/src/main.rs | 2 +- backend/src/processing/processor/playlist.rs | 14 +- backend/src/repository/user_repository.rs | 2 +- shared/src/model/config/mod.rs | 2 + .../src/model/config/playlist_update_state.rs | 7 + shared/src/model/web_socket.rs | 5 +- webui/Cargo.toml | 1 - webui/public/assets/i18n/en.json | 12 +- webui/scss/_theme.scss | 12 +- webui/scss/app/_component.scss | 1 + .../app/components/_loading_indicator.scss | 26 +- webui/scss/app/components/_text_button.scss | 12 + .../playlist/_playlist_processing.scss | 12 + .../playlist/_playlist_update_view.scss | 25 + .../components/playlist/_source_selector.scss | 4 - .../playlist/input/_input_table.scss | 6 + .../app/components/userlist/_proxy_type.scss | 13 + webui/src/app/components/home.rs | 41 +- webui/src/app/components/loading_indicator.rs | 44 +- webui/src/app/components/playlist/mod.rs | 2 + .../components/playlist/playlist_explorer.rs | 14 +- .../playlist/playlist_update_view.rs | 96 ++++ .../components/playlist/source_selector.rs | 43 +- .../app/components/playlist/target_table.rs | 6 +- webui/src/app/components/sidebar.rs | 2 + webui/src/app/components/toastr.rs | 2 +- .../components/userlist/proxy_type_view.rs | 8 +- webui/src/hooks/use_server_status.rs | 12 +- webui/src/hooks/use_service_context.rs | 7 +- webui/src/model/busy_status.rs | 6 + webui/src/model/event_message.rs | 14 + webui/src/model/mod.rs | 6 +- webui/src/model/view_type.rs | 3 + webui/src/services/event_service.rs | 44 ++ webui/src/services/mod.rs | 4 +- webui/src/services/toastr_service.rs | 14 +- webui/src/services/websocket_service.rs | 64 +-- 46 files changed, 750 insertions(+), 457 deletions(-) create mode 100644 shared/src/model/config/playlist_update_state.rs create mode 100644 webui/scss/app/components/playlist/_playlist_update_view.scss create mode 100644 webui/src/app/components/playlist/playlist_update_view.rs create mode 100644 webui/src/model/busy_status.rs create mode 100644 webui/src/model/event_message.rs create mode 100644 webui/src/services/event_service.rs diff --git a/Cargo.lock b/Cargo.lock index 9478135c4..d51dbd64c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -917,7 +917,6 @@ dependencies = [ "implicit-clone 0.6.0", "js-sys", "log", - "nanoid", "prost", "regex", "reqwasm", @@ -1572,7 +1571,7 @@ dependencies = [ "parking_lot", "portable-atomic", "quanta", - "rand 0.9.1", + "rand", "smallvec", "spinning_top", "web-time", @@ -2269,15 +2268,6 @@ dependencies = [ "windows-sys 0.59.0", ] -[[package]] -name = "nanoid" -version = "0.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3ffa00dec017b5b1a8b7cf5e2c008bfda1aa7e0697ac1508b491fdf2622fb4d8" -dependencies = [ - "rand 0.8.5", -] - [[package]] name = "native-tls" version = "0.2.14" @@ -2825,7 +2815,7 @@ dependencies = [ "bytes", "getrandom 0.3.3", "lru-slab", - "rand 0.9.1", + "rand", "ring", "rustc-hash", "rustls", @@ -2866,35 +2856,14 @@ version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" -[[package]] -name = "rand" -version = "0.8.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" -dependencies = [ - "libc", - "rand_chacha 0.3.1", - "rand_core 0.6.4", -] - [[package]] name = "rand" version = "0.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9fbfd9d094a40bf3ae768db9361049ace4c0e04a4fd6b359518bd7b73a73dd97" dependencies = [ - "rand_chacha 0.9.0", - "rand_core 0.9.3", -] - -[[package]] -name = "rand_chacha" -version = "0.3.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" -dependencies = [ - "ppv-lite86", - "rand_core 0.6.4", + "rand_chacha", + "rand_core", ] [[package]] @@ -2904,16 +2873,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" dependencies = [ "ppv-lite86", - "rand_core 0.9.3", -] - -[[package]] -name = "rand_core" -version = "0.6.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" -dependencies = [ - "getrandom 0.2.16", + "rand_core", ] [[package]] @@ -3913,7 +3873,7 @@ dependencies = [ "pest", "pest_derive", "quick-xml", - "rand 0.9.1", + "rand", "rayon", "regex", "reqwest", @@ -3951,7 +3911,7 @@ dependencies = [ "http 1.3.1", "httparse", "log", - "rand 0.9.1", + "rand", "sha1", "thiserror 2.0.12", "utf-8", diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index 584e83f51..a587185d5 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -142,7 +142,9 @@ async fn playlist_update( let process_targets = app_state.app_config.sources.load().validate_targets(user_targets.as_ref()); match process_targets { Ok(valid_targets) => { - tokio::spawn(playlist::exec_processing(Arc::clone(&app_state.http_client.load()), Arc::clone(&app_state.app_config), Arc::new(valid_targets))); + let app_config = Arc::clone(&app_state.app_config); + let event_manager = Arc::clone(&app_state.event_manager); + tokio::spawn(playlist::exec_processing(Arc::clone(&app_state.http_client.load()), app_config, Arc::new(valid_targets), Some(event_manager))); axum::http::StatusCode::OK.into_response() } Err(err) => { diff --git a/backend/src/api/endpoints/websocket_api.rs b/backend/src/api/endpoints/websocket_api.rs index d03e2c27f..1dd2a2a67 100644 --- a/backend/src/api/endpoints/websocket_api.rs +++ b/backend/src/api/endpoints/websocket_api.rs @@ -221,6 +221,15 @@ async fn handle_event_message(socket: &mut WebSocket, event: EventMessage, handl .await .map_err(|e| format!("Configuration files change event: {e} "))?; } + EventMessage::PlaylistUpdate(state) => { + let msg = ProtocolMessage::PlaylistUpdateResponse(state) + .to_bytes() + .map_err(|e| e.to_string())?; + socket + .send(Message::Binary(msg)) + .await + .map_err(|e| format!("Playlist update event: {e} "))?; + } } } } diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 34b332eae..a8daa34c1 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -94,15 +94,19 @@ fn create_shared_data( fn exec_update_on_boot( client: Arc, - cfg: &Arc, + app_state: &Arc, targets: &Arc, ) { - let config = cfg.config.load(); - if config.update_on_boot { - let cfg_clone = Arc::clone(cfg); + let cfg = &app_state.app_config; + let update_on_boot = { + let config = cfg.config.load(); + config.update_on_boot + }; + if update_on_boot { + let app_state_clone = Arc::clone(&app_state.app_config); let targets_clone = Arc::clone(targets); tokio::spawn( - async move { playlist::exec_processing(client, cfg_clone, targets_clone).await }, + async move { playlist::exec_processing(client, app_state_clone, targets_clone, None).await }, ); } } @@ -233,13 +237,13 @@ pub async fn start_server( exec_scheduler( &Arc::clone(&shared_data.http_client.load()), - &app_config, + &app_state, &targets, &cancel_token_scheduler, ); exec_update_on_boot( Arc::clone(&shared_data.http_client.load()), - &app_config, + &app_state, &targets, ); exec_config_watch(&app_state, &cancel_token_file_watch); diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 054e29123..dc39d2d49 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -211,11 +211,12 @@ impl SingleProviderLineup { self.provider.try_allocate(with_grace, grace_period_timeout_secs).await } - // async fn release(&self, provider_name: &str) { - // if self.provider.name == provider_name { - // self.provider.release().await; - // } - // } + #[cfg(test)] + async fn release(&self, provider_name: &str) { + if self.provider.name == provider_name { + self.provider.release().await; + } + } } @@ -462,26 +463,27 @@ impl MultiProviderLineup { } - // async fn release(&self, provider_name: &str) { - // for g in &self.providers { - // match g { - // ProviderPriorityGroup::SingleProviderGroup(pc) => { - // if pc.name == provider_name { - // pc.release().await; - // break; - // } - // } - // ProviderPriorityGroup::MultiProviderGroup(_, group) => { - // for pc in group { - // if pc.name == provider_name { - // pc.release().await; - // return; - // } - // } - // } - // } - // } - // } + #[cfg(test)] + async fn release(&self, provider_name: &str) { + for g in &self.providers { + match g { + ProviderPriorityGroup::SingleProviderGroup(pc) => { + if pc.name == provider_name { + pc.release().await; + break; + } + } + ProviderPriorityGroup::MultiProviderGroup(_, group) => { + for pc in group { + if pc.name == provider_name { + pc.release().await; + return; + } + } + } + } + } + } } @@ -910,218 +912,229 @@ mod tests { } // Test acquiring with an alias - // #[test] - // fn test_provider_with_alias() { - // let mut input = create_config_input(1, "provider1_1", 1, 1); - // let alias = create_config_input_alias(2, "http://alias1", 2, 2); - // - // // Adding alias to the provider - // input.aliases = Some(vec![alias]); - // - // let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); - // // Create MultiProviderLineup with the provider and alias - // let lineup = MultiProviderLineup::new(&input, None, /* &tokio::sync::mpsc::Sender<(std::string::String, usize)> */ &change_tx); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // // Test that the alias provider is available - // should_available!(lineup, 1, 5); - // // Try acquiring again - // should_available!(lineup, 2, 5); - // should_available!(lineup, 2, 5); - // should_grace_period!(lineup, 1, 5); - // should_grace_period!(lineup, 2, 5); - // should_exhausted!(lineup, 5); - // should_exhausted!(lineup, 5); - // }); - // } - // - // // // Test acquiring from a MultiProviderLineup where the alias has a different priority - // #[test] - // fn test_provider_with_priority_alias() { - // let mut input = create_config_input(1, "provider2_1", 1, 2); - // let alias = create_config_input_alias(2, "http://alias.com", 0, 2); - // // Adding alias with different priority - // input.aliases = Some(vec![alias]); - // let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); - // let lineup = MultiProviderLineup::new(&input, None, &change_tx); - // // The alias has a higher priority, so the alias should be acquired first - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // for _ in 0..2 { - // should_available!(lineup, 2, 5); - // } - // should_available!(lineup, 1, 5); - // }); - // } - // - // // Test provider when there are multiple aliases, all with distinct priorities - // #[test] - // fn test_provider_with_multiple_aliases() { - // let mut input = create_config_input(1, "provider3_1", 1, 1); - // let alias1 = create_config_input_alias(2, "http://alias1.com", 1, 2); - // let alias2 = create_config_input_alias(3, "http://alias2.com", 0, 1); - // - // // Adding multiple aliases - // input.aliases = Some(vec![alias1, alias2]); - // let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); - // let lineup = MultiProviderLineup::new(&input, None, &change_tx); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // // The alias with priority 0 should be acquired first (higher priority) - // should_available!(lineup, 3, 5); - // // Acquire again, and provider should still be available (with remaining capacity) - // should_available!(lineup, 1, 5); - // // // Check that the second alias with priority 2 is considered next - // should_available!(lineup, 2, 5); - // should_available!(lineup, 2, 5); - // - // should_grace_period!(lineup, 3, 5); - // should_grace_period!(lineup, 1, 5); - // should_grace_period!(lineup, 2, 5); - // - // should_exhausted!(lineup, 5); - // }); - // } - // - // - // // // Test acquiring when all aliases are exhausted - // #[test] - // fn test_provider_with_exhausted_aliases() { - // let mut input = create_config_input(1, "provider4_1", 1, 1); - // let alias1 = create_config_input_alias(2, "http://alias.com", 2, 1); - // let alias2 = create_config_input_alias(3, "http://alias.com", -2, 1); - // - // // Adding alias - // input.aliases = Some(vec![alias1, alias2]); - // let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); - // let lineup = MultiProviderLineup::new(&input, None, &change_tx); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // // Acquire connection from alias2 - // should_available!(lineup, 3, 5); - // // Acquire connection from provider1 - // should_available!(lineup, 1, 5); - // // Acquire connection from alias1 - // should_available!(lineup, 2, 5); - // - // // Acquire connection from alias2 - // should_grace_period!(lineup, 3, 5); - // // Acquire connection from provider1 - // should_grace_period!(lineup, 1, 5); - // // Acquire connection from alias1 - // should_grace_period!(lineup, 2, 5); - // - // // Now, all are exhausted - // should_exhausted!(lineup, 5); - // }); - // } - // - // // Test acquiring a connection when there is available capacity - // #[test] - // fn test_acquire_when_capacity_available() { - // let cfg = create_config_input(1, "provider5_1", 1, 2); - // let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); - // let lineup = SingleProviderLineup::new(&cfg, None, change_tx); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // // First acquire attempt should succeed - // should_available!(lineup, 1, 5); - // // Second acquire attempt should succeed as well - // should_available!(lineup, 1, 5); - // // Third with grace time - // should_grace_period!(lineup, 1, 5); - // // Fourth acquire attempt should fail as the provider is exhausted - // should_exhausted!(lineup, 5); - // }); - // } + #[test] + fn test_provider_with_alias() { + let mut input = create_config_input(1, "provider1_1", 1, 1); + let alias = create_config_input_alias(2, "http://alias1", 2, 2); + + // Adding alias to the provider + input.aliases = Some(vec![alias]); + + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + // Create MultiProviderLineup with the provider and alias + let lineup = MultiProviderLineup::new(&input, Some(dummy_get_connection), &change_tx); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + // Test that the alias provider is available + should_available!(lineup, 1, 5); + // Try acquiring again + should_available!(lineup, 2, 5); + should_available!(lineup, 2, 5); + should_grace_period!(lineup, 1, 5); + should_grace_period!(lineup, 2, 5); + should_exhausted!(lineup, 5); + should_exhausted!(lineup, 5); + }); + } + + // Test acquiring from a MultiProviderLineup where the alias has a different priority + #[test] + fn test_provider_with_priority_alias() { + let mut input = create_config_input(1, "provider2_1", 1, 2); + let alias = create_config_input_alias(2, "http://alias.com", 0, 2); + // Adding alias with different priority + input.aliases = Some(vec![alias]); + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + let lineup = MultiProviderLineup::new(&input, Some(dummy_get_connection), &change_tx); + // The alias has a higher priority, so the alias should be acquired first + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + for _ in 0..2 { + should_available!(lineup, 2, 5); + } + should_available!(lineup, 1, 5); + }); + } + + // Test provider when there are multiple aliases, all with distinct priorities + #[test] + fn test_provider_with_multiple_aliases() { + let mut input = create_config_input(1, "provider3_1", 1, 1); + let alias1 = create_config_input_alias(2, "http://alias1.com", 1, 2); + let alias2 = create_config_input_alias(3, "http://alias2.com", 0, 1); + + // Adding multiple aliases + input.aliases = Some(vec![alias1, alias2]); + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + let lineup = MultiProviderLineup::new(&input, Some(dummy_get_connection), &change_tx); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + // The alias with priority 0 should be acquired first (higher priority) + should_available!(lineup, 3, 5); + // Acquire again, and provider should still be available (with remaining capacity) + should_available!(lineup, 1, 5); + // // Check that the second alias with priority 2 is considered next + should_available!(lineup, 2, 5); + should_available!(lineup, 2, 5); + + should_grace_period!(lineup, 3, 5); + should_grace_period!(lineup, 1, 5); + should_grace_period!(lineup, 2, 5); + + should_exhausted!(lineup, 5); + }); + } - // // Test releasing a connection - // #[test] - // fn test_release_connection() { - // let cfg = create_config_input(1, "provider7_1", 1, 2); - // let lineup = SingleProviderLineup::new(&cfg, None); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // // Acquire two connections - // should_available!(lineup, 1, 5); - // should_available!(lineup, 1, 5); - // should_grace_period!(lineup, 1, 5); - // should_exhausted!(lineup, 5); - // lineup.release("provider7_1").await; - // should_grace_period!(lineup, 1, 5); - // lineup.release("provider7_1").await; - // lineup.release("provider7_1").await; - // should_available!(lineup, 1, 5); - // should_grace_period!(lineup, 1, 5); - // should_exhausted!(lineup, 5); - // }); - // } - // - // // Test acquiring with MultiProviderLineup and round-robin allocation - // #[test] - // fn test_multi_provider_acquire() { - // let mut cfg1 = create_config_input(1, "provider8_1", 1, 2); - // let alias = create_config_input_alias(2, "http://alias1", 1, 1); - // - // // Adding alias to the provider - // cfg1.aliases = Some(vec![alias]); - // - // // Create MultiProviderLineup with the provider and alias - // let lineup = MultiProviderLineup::new(&cfg1, None); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // // Test acquiring the first provider - // should_available!(lineup, 1, 5); - // - // // Test acquiring the second provider - // should_available!(lineup, 2, 5); - // - // // Test acquiring the first provider - // should_available!(lineup, 1, 5); - // - // should_grace_period!(lineup, 1, 5); - // should_grace_period!(lineup, 2, 5); - // - // lineup.release("provider8_1").await; - // lineup.release("alias_2").await; - // lineup.release("provider8_1").await; - // - // should_available!(lineup, 1, 5); - // should_grace_period!(lineup, 1, 5); - // should_grace_period!(lineup, 2, 5); - // - // should_exhausted!(lineup, 5); - // }); - // } + // Test acquiring when all aliases are exhausted + #[test] + fn test_provider_with_exhausted_aliases() { + let mut input = create_config_input(1, "provider4_1", 1, 1); + let alias1 = create_config_input_alias(2, "http://alias.com", 2, 1); + let alias2 = create_config_input_alias(3, "http://alias.com", -2, 1); + + // Adding alias + input.aliases = Some(vec![alias1, alias2]); + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + let lineup = MultiProviderLineup::new(&input, Some(dummy_get_connection), &change_tx); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + // Acquire connection from alias2 + should_available!(lineup, 3, 5); + // Acquire connection from provider1 + should_available!(lineup, 1, 5); + // Acquire connection from alias1 + should_available!(lineup, 2, 5); + + // Acquire connection from alias2 + should_grace_period!(lineup, 3, 5); + // Acquire connection from provider1 + should_grace_period!(lineup, 1, 5); + // Acquire connection from alias1 + should_grace_period!(lineup, 2, 5); + + // Now, all are exhausted + should_exhausted!(lineup, 5); + }); + } + + // Test acquiring a connection when there is available capacity + #[test] + fn test_acquire_when_capacity_available() { + let cfg = create_config_input(1, "provider5_1", 1, 2); + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + let lineup = SingleProviderLineup::new(&cfg, Some(dummy_get_connection), change_tx); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + // First acquire attempt should succeed + should_available!(lineup, 1, 5); + // Second acquire attempt should succeed as well + should_available!(lineup, 1, 5); + // Third with grace time + should_grace_period!(lineup, 1, 5); + // Fourth acquire attempt should fail as the provider is exhausted + should_exhausted!(lineup, 5); + }); + } + + + // Test releasing a connection + #[test] + fn test_release_connection() { + let cfg = create_config_input(1, "provider7_1", 1, 2); + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + + let lineup = SingleProviderLineup::new(&cfg, Some(dummy_get_connection), change_tx); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + // Acquire two connections + should_available!(lineup, 1, 5); + should_available!(lineup, 1, 5); + should_grace_period!(lineup, 1, 5); + should_exhausted!(lineup, 5); + lineup.release("provider7_1").await; + should_grace_period!(lineup, 1, 5); + lineup.release("provider7_1").await; + lineup.release("provider7_1").await; + should_available!(lineup, 1, 5); + should_grace_period!(lineup, 1, 5); + should_exhausted!(lineup, 5); + }); + } + + // Test acquiring with MultiProviderLineup and round-robin allocation + #[test] + fn test_multi_provider_acquire() { + let mut cfg1 = create_config_input(1, "provider8_1", 1, 2); + let alias = create_config_input_alias(2, "http://alias1", 1, 1); + + // Adding alias to the provider + cfg1.aliases = Some(vec![alias]); + + // Create MultiProviderLineup with the provider and alias + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + let lineup = MultiProviderLineup::new(&cfg1, Some(dummy_get_connection), &change_tx); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + // Test acquiring the first provider + should_available!(lineup, 1, 5); + + // Test acquiring the second provider + should_available!(lineup, 2, 5); + + // Test acquiring the first provider + should_available!(lineup, 1, 5); + + should_grace_period!(lineup, 1, 5); + should_grace_period!(lineup, 2, 5); + + lineup.release("provider8_1").await; + lineup.release("alias_2").await; + lineup.release("provider8_1").await; + + should_available!(lineup, 1, 5); + should_grace_period!(lineup, 1, 5); + should_grace_period!(lineup, 2, 5); + + should_exhausted!(lineup, 5); + }); + } // Test concurrent access to `acquire` using multiple threads - // #[test] - // fn test_concurrent_acquire() { - // let cfg = create_config_input(1, "provider9_1", 1, 2); - // let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); - // let lineup = Arc::new(SingleProviderLineup::new(&cfg, None, change_tx)); - // - // let available_count = Arc::new(AtomicU16::new(2)); - // let grace_period_count = Arc::new(AtomicU16::new(1)); - // let exhausted_count = Arc::new(AtomicU16::new(2)); - // - // for _ in 0..5 { - // let lineup_clone = Arc::clone(&lineup); - // let available = Arc::clone(&available_count); - // let grace_period = Arc::clone(&grace_period_count); - // let exhausted = Arc::clone(&exhausted_count); - // let rt = tokio::runtime::Runtime::new().unwrap(); - // rt.block_on(async move { - // match lineup_clone.acquire(true, 5).await { - // ProviderAllocation::Exhausted => exhausted.fetch_sub(1, Ordering::SeqCst), - // ProviderAllocation::Available(_, _) => available.fetch_sub(1, Ordering::SeqCst), - // ProviderAllocation::GracePeriod(_, _) => grace_period.fetch_sub(1, Ordering::SeqCst), - // } - // }); - // } - // assert_eq!(exhausted_count.load(Ordering::SeqCst), 0); - // assert_eq!(available_count.load(Ordering::SeqCst), 0); - // assert_eq!(grace_period_count.load(Ordering::SeqCst), 0); - // } + #[test] + fn test_concurrent_acquire() { + let cfg = create_config_input(1, "provider9_1", 1, 2); + let (change_tx, _) = tokio::sync::mpsc::channel::<(String, usize)>(1); + let dummy_get_connection = |_s: &str| -> Option<&ProviderConfigConnection> { None }; + let lineup = Arc::new(SingleProviderLineup::new(&cfg, Some(dummy_get_connection), change_tx)); + + let available_count = Arc::new(AtomicU16::new(2)); + let grace_period_count = Arc::new(AtomicU16::new(1)); + let exhausted_count = Arc::new(AtomicU16::new(2)); + + for _ in 0..5 { + let lineup_clone = Arc::clone(&lineup); + let available = Arc::clone(&available_count); + let grace_period = Arc::clone(&grace_period_count); + let exhausted = Arc::clone(&exhausted_count); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async move { + match lineup_clone.acquire(true, 5).await { + ProviderAllocation::Exhausted => exhausted.fetch_sub(1, Ordering::SeqCst), + ProviderAllocation::Available(_, _) => available.fetch_sub(1, Ordering::SeqCst), + ProviderAllocation::GracePeriod(_, _) => grace_period.fetch_sub(1, Ordering::SeqCst), + } + }); + } + assert_eq!(exhausted_count.load(Ordering::SeqCst), 0); + assert_eq!(available_count.load(Ordering::SeqCst), 0); + assert_eq!(grace_period_count.load(Ordering::SeqCst), 0); + } } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index ed4985494..296638b74 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -88,7 +88,7 @@ fn start_services(app_state: &Arc, changes: &UpdateChanges) { return; } if changes.scheduler { - exec_scheduler(&Arc::clone(&app_state.http_client.load()), &app_state.app_config, + exec_scheduler(&Arc::clone(&app_state.http_client.load()), app_state, &app_state.forced_targets.load(), &app_state.cancel_tokens.load().scheduler); } diff --git a/backend/src/api/model/event_manager.rs b/backend/src/api/model/event_manager.rs index a4b3eab6c..4dec7b578 100644 --- a/backend/src/api/model/event_manager.rs +++ b/backend/src/api/model/event_manager.rs @@ -1,6 +1,6 @@ use log::{error, info}; use tokio::task; -use shared::model::ConfigType; +use shared::model::{ConfigType, PlaylistUpdateState}; use crate::api::model::{ActiveUserConnectionChangeReceiver}; use crate::api::model::{ProviderConnectionChangeReceiver}; @@ -8,7 +8,8 @@ use crate::api::model::{ProviderConnectionChangeReceiver}; pub enum EventMessage { ActiveUser(usize, usize), // user_count, connection count ActiveProvider(String, usize), // provider name, connections - ConfigChange(ConfigType) + ConfigChange(ConfigType), + PlaylistUpdate(PlaylistUpdateState) } pub struct EventManager { diff --git a/backend/src/api/model/mod.rs b/backend/src/api/model/mod.rs index 2e8eb2a55..13d8a136f 100644 --- a/backend/src/api/model/mod.rs +++ b/backend/src/api/model/mod.rs @@ -22,4 +22,4 @@ pub(in crate::api) use self::active_user_manager::*; pub(in crate::api) use self::active_provider_manager::*; pub(in crate::api) use self::stream::*; pub(in crate::api) use self::provider_config::*; -pub(in crate::api) use self::event_manager::*; \ No newline at end of file +pub use self::event_manager::*; \ No newline at end of file diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index bea87ebd4..508de9c5b 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -1,13 +1,14 @@ -use std::sync::Arc; -use std::str::FromStr; -use std::time::{Duration, Instant, SystemTime}; -use chrono::{DateTime, FixedOffset, Local}; -use cron::Schedule; -use crate::utils::{exit}; -use log::{error}; -use tokio_util::sync::CancellationToken; +use crate::api::model::AppState; use crate::model::{AppConfig, ProcessTargets, ScheduleConfig}; use crate::processing::processor::playlist::exec_processing; +use crate::utils::exit; +use chrono::{DateTime, FixedOffset, Local}; +use cron::Schedule; +use log::error; +use std::str::FromStr; +use std::sync::Arc; +use std::time::{Duration, Instant, SystemTime}; +use tokio_util::sync::CancellationToken; pub fn datetime_to_instant(datetime: DateTime) -> Instant { // Convert DateTime to SystemTime @@ -25,8 +26,9 @@ pub fn datetime_to_instant(datetime: DateTime) -> Instant { Instant::now() + duration_until } -pub fn exec_scheduler(client: &Arc, cfg: &Arc, targets: &Arc, - cancel: &CancellationToken) { +pub fn exec_scheduler(client: &Arc, app_state: &Arc, targets: &Arc, + cancel: &CancellationToken) { + let cfg = &app_state.app_config; let config = cfg.config.load(); let schedules: Vec = if let Some(schedules) = &config.schedules { schedules.clone() @@ -36,17 +38,17 @@ pub fn exec_scheduler(client: &Arc, cfg: &Arc, targe for schedule in schedules { let expression = schedule.schedule.to_string(); let exec_targets = get_process_targets(cfg, targets, schedule.targets.as_ref()); - let cfg_clone = Arc::clone(cfg); + let app_state_clone = Arc::clone(app_state); let http_client = Arc::clone(client); let cancel_token = cancel.clone(); tokio::spawn(async move { - start_scheduler(http_client, expression.as_str(), cfg_clone, exec_targets, cancel_token).await; + start_scheduler(http_client, expression.as_str(), app_state_clone, exec_targets, cancel_token).await; }); } } -async fn start_scheduler(client: Arc, expression: &str, config: Arc, - targets: Arc, cancel: CancellationToken) { +async fn start_scheduler(client: Arc, expression: &str, app_state: Arc, + targets: Arc, cancel: CancellationToken) { match Schedule::from_str(expression) { Ok(schedule) => { let offset = *Local::now().offset(); @@ -55,7 +57,9 @@ async fn start_scheduler(client: Arc, expression: &str, config: if let Some(datetime) = upcoming.next() { tokio::select! { () = tokio::time::sleep_until(tokio::time::Instant::from(datetime_to_instant(datetime))) => { - exec_processing(Arc::clone(&client), Arc::clone(&config), Arc::clone(&targets)).await; + let app_config = Arc::clone(&app_state.app_config); + let event_manager = Arc::clone(&app_state.event_manager); + exec_processing(Arc::clone(&client), app_config, Arc::clone(&targets), Some(event_manager)).await; } () = cancel.cancelled() => { break; @@ -93,7 +97,7 @@ fn get_process_targets(cfg: &Arc, process_targets: &Arc, process_targets: &Arc, targets: Arc) { error!("Failed to build client {err}"); reqwest::Client::new() }); - playlist::exec_processing(Arc::new(client), cfg, targets).await; + playlist::exec_processing(Arc::new(client), cfg, targets, None).await; } async fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 0ba79152c..b6fc07294 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -30,9 +30,10 @@ use log::{debug, error, info, log_enabled, trace, warn, Level}; use reqwest::Client; use shared::error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind}; use shared::foundation::filter::{get_field_value, set_field_value, ValueAccessor, ValueProvider}; -use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, ProcessingOrder, UUIDType, XtreamCluster}; +use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster}; use shared::utils::default_as_default; use std::time::Instant; +use crate::api::model::{EventManager, EventMessage}; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { let provider = ValueProvider { pli }; @@ -552,14 +553,14 @@ fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTar } } -pub async fn exec_processing(client: Arc, cfg: Arc, targets: Arc) { +pub async fn exec_processing(client: Arc, app_config: Arc, targets: Arc, event_manager: Option>) { let start_time = Instant::now(); - let (stats, errors) = process_sources(Arc::clone(&client), &cfg, targets.clone()).await; + let (stats, errors) = process_sources(Arc::clone(&client), &app_config, targets.clone()).await; // log errors for err in &errors { error!("{}", err.message); } - let config = cfg.config.load(); + let config = app_config.config.load(); let messaging = config.messaging.as_ref(); if let Ok(stats_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("stats".to_string(), serde_json::to_value(stats).unwrap())]))) { // print stats @@ -569,9 +570,14 @@ pub async fn exec_processing(client: Arc, cfg: Arc, } // send errors if let Some(message) = get_errors_notify_message!(errors, 255) { + 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(&client, &MsgKind::Error, messaging, error_msg.as_str()); } + } 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."); diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index 3a1a33161..c6a4c690d 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -400,7 +400,7 @@ mod tests { use shared::model::{ConfigPaths, ProxyType, ProxyUserStatus}; use std::env::temp_dir; use std::sync::Arc; - use arc_swap::{ArcSwap, ArcSwapAny, ArcSwapOption}; + use arc_swap::{ArcSwap, ArcSwapAny}; use crate::utils::FileLockManager; #[test] diff --git a/shared/src/model/config/mod.rs b/shared/src/model/config/mod.rs index d4341d546..b45c5e46e 100644 --- a/shared/src/model/config/mod.rs +++ b/shared/src/model/config/mod.rs @@ -28,6 +28,7 @@ pub mod macros; mod paths; mod app_config; mod config_type; +mod playlist_update_state; pub use base::*; pub use api::*; @@ -58,4 +59,5 @@ pub use pattern_template::*; pub use paths::*; pub use app_config::*; pub use config_type::*; +pub use playlist_update_state::*; pub use crate::apply_batch_aliases; \ No newline at end of file diff --git a/shared/src/model/config/playlist_update_state.rs b/shared/src/model/config/playlist_update_state.rs new file mode 100644 index 000000000..0c66913df --- /dev/null +++ b/shared/src/model/config/playlist_update_state.rs @@ -0,0 +1,7 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +pub enum PlaylistUpdateState { + Success, + Failure, +} \ No newline at end of file diff --git a/shared/src/model/web_socket.rs b/shared/src/model/web_socket.rs index 4a4be9d65..50a93b494 100644 --- a/shared/src/model/web_socket.rs +++ b/shared/src/model/web_socket.rs @@ -1,6 +1,6 @@ use std::io; use bytes::Bytes; -use crate::model::{ConfigType, StatusCheck}; +use crate::model::{ConfigType, PlaylistUpdateState, StatusCheck}; use serde::{Deserialize, Serialize}; pub const PROTOCOL_VERSION: u8 = 1; @@ -80,7 +80,8 @@ pub enum ProtocolMessage { StatusResponse(StatusCheck), ActiveUserResponse(usize, usize), // user_count, connection count ActiveProviderResponse(String, usize), - ConfigChangeResponse(ConfigType) + ConfigChangeResponse(ConfigType), + PlaylistUpdateResponse(PlaylistUpdateState) } impl ProtocolMessage { diff --git a/webui/Cargo.toml b/webui/Cargo.toml index d5950b969..a7f6b5d01 100644 --- a/webui/Cargo.toml +++ b/webui/Cargo.toml @@ -27,7 +27,6 @@ prost = "0" wasm-bindgen-futures = "0" bytes = "1" regex = "1.11.1" -nanoid = "0" [dependencies.web-sys] version = "0.3" diff --git a/webui/public/assets/i18n/en.json b/webui/public/assets/i18n/en.json index e7d337782..434ce6d27 100644 --- a/webui/public/assets/i18n/en.json +++ b/webui/public/assets/i18n/en.json @@ -1,8 +1,12 @@ { "LABEL": { + "ALL": "All", "LIVE": "Live", "VOD": "VOD", "SERIES": "Series", + "LIVE_SHORT": "L", + "VOD_SHORT": "V", + "SERIES_SHORT": "S", "SAVE": "Save", "SUBMIT": "Submit", "OK": "Ok", @@ -218,8 +222,8 @@ "PROXY": "Proxy", "SERVER": "Server", "EPG_TIMESHIFT": "Epg Timeshift", - "CREATED_AT": "Created_at", - "EXP_DATE": "Exp_date", + "CREATED_AT": "Created at", + "EXP_DATE": "Exp. date", "STATUS": "Status", "UI_ENABLED": "Ui", "COMMENT": "Comment", @@ -267,7 +271,9 @@ }, "PLAYLIST_UPDATE": { "SUCCESS": "Successfully started playlist update!", - "FAIL": "Playlist update failed!" + "FAIL": "Playlist update failed!", + "SUCCESS_FINISH": "Successfully updated playlist!", + "FAIL_FINISH": "Playlist update failed!" }, "SOURCE_SELECTOR": { "MISSING_SELECTION": "No entry selected" diff --git a/webui/scss/_theme.scss b/webui/scss/_theme.scss index 8a3cbe6ea..0f379e348 100644 --- a/webui/scss/_theme.scss +++ b/webui/scss/_theme.scss @@ -66,6 +66,10 @@ --text-button-secondary-color: #ffffff; --text-button-secondary-hover-background-color: var(--secondary-color); --text-button-secondary-hover-color: #ffffff; + --text-button-active-background-color: #1e2d25; + --text-button-active-color: #49dc7f; + --text-button-active-hover-background-color: #1e2d25; + --text-button-active-hover-color: #ffffff; --toggle-switch-on-background-color: #36c55a; --toggle-switch-off-background-color: #3f3f45; @@ -147,13 +151,13 @@ --tag-series-border-color: #6366f1; --toastr-success-color: #ffffff; - --toastr-success-background-color: #71b570; + --toastr-success-background-color: #53a453; --toastr-error-color: #ffffff; - --toastr-error-background-color: #c85a52; + --toastr-error-background-color: #bc362f; --toastr-info-color: #ffffff; - --toastr-info-background-color: #51a8c0; + --toastr-info-background-color: #2f95b3; --toastr-warning-color: #ffffff; - --toastr-warning-background-color: #f7a820; + --toastr-warning-background-color: #f69306; --hint-color: #888888; --label-color: #73d6fe; diff --git a/webui/scss/app/_component.scss b/webui/scss/app/_component.scss index 9282d852e..9726b5ee1 100644 --- a/webui/scss/app/_component.scss +++ b/webui/scss/app/_component.scss @@ -41,6 +41,7 @@ @forward "components/playlist/playlist_mappings"; @forward "components/playlist/playlist_processing"; @forward "components/playlist/source_selector"; +@forward "components/playlist/playlist_update_view"; @forward "components/tag_list"; @forward "components/chip"; @forward "components/popup_menu"; diff --git a/webui/scss/app/components/_loading_indicator.scss b/webui/scss/app/components/_loading_indicator.scss index e6ab0975f..e544a1620 100644 --- a/webui/scss/app/components/_loading_indicator.scss +++ b/webui/scss/app/components/_loading_indicator.scss @@ -1,11 +1,31 @@ +$indicator-height: 5px; + +.tp__busy-indicator { + position: fixed; + z-index: 1000; + top: var(--app-header-height); + left: 0; + display: inline-block; + width: 100%; // Full width of the container + height: $indicator-height; // Thin height for the container + min-height: $indicator-height; + max-height: $indicator-height; + .tp__loading-bar-placeholder, .tp__loading-bar-container { + position: fixed; + } +} + .tp__loading-bar-placeholder, .tp__loading-bar-container { display: inline-block; width: 100%; // Full width of the container - height: 5px; // Thin height for the container - min-height: 5px; - max-height: 5px; + height: $indicator-height; // Thin height for the container + min-height: $indicator-height; + max-height: $indicator-height; + margin: 0; + padding: 0; } + // Container for the loading bar .tp__loading-bar-container { background-color: var(--spinner-color); // Light gray background for the container diff --git a/webui/scss/app/components/_text_button.scss b/webui/scss/app/components/_text_button.scss index 5b37f2f4d..518f4e4d0 100644 --- a/webui/scss/app/components/_text_button.scss +++ b/webui/scss/app/components/_text_button.scss @@ -55,4 +55,16 @@ fill : var(--text-button-secondary-hover-color); outline: none; } +} +.tp__text-button.tp__button-active { + background-color: var(--text-button-active-background-color); + color: var(--text-button-active-color); + fill : var(--text-button-active-color); + &:focus, + &:hover { + background-color: var(--text-button-active-hover-background-color); + color: var(--text-button-active-hover-color); + fill : var(--text-button-active-hover-color); + outline: none; + } } \ No newline at end of file diff --git a/webui/scss/app/components/playlist/_playlist_processing.scss b/webui/scss/app/components/playlist/_playlist_processing.scss index 97113abda..8801973f0 100644 --- a/webui/scss/app/components/playlist/_playlist_processing.scss +++ b/webui/scss/app/components/playlist/_playlist_processing.scss @@ -3,4 +3,16 @@ flex-flow: row; justify-content: center; align-items: center; + + .tp__tag_list { + gap: 0; + span:not(:first-child) { + border-top-left-radius: 0; + border-bottom-left-radius: 0; + } + span:not(:last-child) { + border-top-right-radius: 0; + border-bottom-right-radius: 0; + } + } } diff --git a/webui/scss/app/components/playlist/_playlist_update_view.scss b/webui/scss/app/components/playlist/_playlist_update_view.scss new file mode 100644 index 000000000..eb4ca5d76 --- /dev/null +++ b/webui/scss/app/components/playlist/_playlist_update_view.scss @@ -0,0 +1,25 @@ +.tp__playlist-update-view { + display: flex; + flex-flow: column; + width: 100%; + gap: var(--gap-default); + padding: var(--padding-default) 0; + box-sizing: border-box; + &__header { + display: flex; + flex-flow: row nowrap; + box-sizing: border-box; + width: 100%; + justify-content: space-between; + align-items: center; + } + + &__body { + display: grid; + grid-template-columns: repeat(auto-fit, minmax(300px, 1fr)); + gap: var(--gap-default); + box-sizing: border-box; + overflow-y: auto; + width: 100%; + } +} \ No newline at end of file diff --git a/webui/scss/app/components/playlist/_source_selector.scss b/webui/scss/app/components/playlist/_source_selector.scss index 12aa44ff6..7adf4af04 100644 --- a/webui/scss/app/components/playlist/_source_selector.scss +++ b/webui/scss/app/components/playlist/_source_selector.scss @@ -1,9 +1,5 @@ .tp__playlist-source-selector { - &__loading { - transform: translateY(-10px); - } - &__source-picker .tp__collapse-panel__body { display: flex; flex-flow: column; diff --git a/webui/scss/app/components/playlist/input/_input_table.scss b/webui/scss/app/components/playlist/input/_input_table.scss index 8919f448f..722eca778 100644 --- a/webui/scss/app/components/playlist/input/_input_table.scss +++ b/webui/scss/app/components/playlist/input/_input_table.scss @@ -14,4 +14,10 @@ background-color: var(--tag-alt-active-background-color); border-color: var(--tag-alt-active-border-color); } + + .tp__table { + .tp__table__container { + max-height: 400px; + } + } } \ No newline at end of file diff --git a/webui/scss/app/components/userlist/_proxy_type.scss b/webui/scss/app/components/userlist/_proxy_type.scss index 1a7c91340..17f33e92f 100644 --- a/webui/scss/app/components/userlist/_proxy_type.scss +++ b/webui/scss/app/components/userlist/_proxy_type.scss @@ -14,4 +14,17 @@ background-color: var(--tag-alt-active-background-color); border-color: var(--tag-alt-active-border-color); } + + &__mixed { + display: flex; + flex-flow: row nowrap; + span:not(:first-child) { + border-top-left-radius: 0; + border-bottom-left-radius: 0; + } + span:not(:last-child) { + border-top-right-radius: 0; + border-bottom-right-radius: 0; + } + } } \ No newline at end of file diff --git a/webui/src/app/components/home.rs b/webui/src/app/components/home.rs index ca87e12d9..8c2f9b8ac 100644 --- a/webui/src/app/components/home.rs +++ b/webui/src/app/components/home.rs @@ -1,14 +1,14 @@ -use crate::app::components::{AppIcon, DashboardView, IconButton, InputRow, Panel, PlaylistEditorView, PlaylistExplorerView, Sidebar, StatsView, ToastrView, UserlistView}; +use crate::app::components::{AppIcon, DashboardView, IconButton, InputRow, Panel, PlaylistEditorView, PlaylistExplorerView, PlaylistUpdateView, Sidebar, StatsView, ToastrView, UserlistView}; use crate::app::context::{ConfigContext, PlaylistContext, StatusContext}; use crate::hooks::{use_server_status, use_service_context}; -use crate::model::ViewType; -use crate::services::WsMessage; -use shared::model::{AppConfigDto, StatusCheck}; +use crate::model::{EventMessage, ViewType}; +use shared::model::{AppConfigDto, PlaylistUpdateState, StatusCheck}; use std::future; use std::rc::Rc; use yew::prelude::*; use yew::suspense::use_future; use yew_i18n::use_translation; +use crate::app::components::loading_indicator::{BusyIndicator}; #[function_component] pub fn Home() -> Html { @@ -17,7 +17,7 @@ pub fn Home() -> Html { let config = use_state(|| None::>); let status = use_state(|| None::>); - let view_visible = use_state(|| ViewType::PlaylistExplorer); + let view_visible = use_state(|| ViewType::Dashboard); let handle_logout = { let services_ctx = services.clone(); @@ -31,18 +31,24 @@ pub fn Home() -> Html { let services_ctx = services_ctx.clone(); let services_ctx_clone = services_ctx.clone(); let translate_clone = translate_clone.clone(); - let subid = services_ctx.websocket.subscribe(move |msg| { + let subid = services_ctx.event.subscribe(move |msg| { match msg { - WsMessage::Unauthorized => { + EventMessage::Unauthorized => { services_ctx_clone.auth.logout() }, - WsMessage::ConfigChange(config_type) => { + EventMessage::ConfigChange(config_type) => { services_ctx_clone.toastr.warning(format!("{}: {config_type}", translate_clone.t("MESSAGES.CONFIG_CHANGED"))); }, + EventMessage::PlaylistUpdate(update_state) => { + match update_state { + PlaylistUpdateState::Success => services_ctx_clone.toastr.success(translate_clone.t("MESSAGES.PLAYLIST_UPDATE.SUCCESS_FINISH")), + PlaylistUpdateState::Failure => services_ctx_clone.toastr.error(translate_clone.t("MESSAGES.PLAYLIST_UPDATE.FAIL_FINISH")), + } + }, _=> {} } }); - move || services_ctx.websocket.unsubscribe(subid) + move || services_ctx.event.unsubscribe(subid) }); } @@ -120,6 +126,7 @@ pub fn Home() -> Html { context={playlist_context}>
+
@@ -132,19 +139,6 @@ pub fn Home() -> Html { } }
- // - // - //
@@ -155,6 +149,9 @@ pub fn Home() -> Html { + + + diff --git a/webui/src/app/components/loading_indicator.rs b/webui/src/app/components/loading_indicator.rs index 9c7443f5f..37c46038a 100644 --- a/webui/src/app/components/loading_indicator.rs +++ b/webui/src/app/components/loading_indicator.rs @@ -1,4 +1,6 @@ -use yew::{classes, function_component, html, Html, Properties}; +use crate::hooks::use_service_context; +use crate::model::{BusyStatus, EventMessage}; +use yew::{classes, function_component, html, use_effect_with, use_mut_ref, use_state, Html, Properties}; #[derive(Properties, Clone, PartialEq, Debug)] pub struct LoadingIndicatorProps { @@ -18,4 +20,44 @@ pub fn LoadingIndicator(props: &LoadingIndicatorProps) -> Html {
} } +} + +#[function_component] +pub fn BusyIndicator() -> Html { + let counter = use_mut_ref(|| 0u32); + let loading = use_state(|| false); + let service_ctx = use_service_context(); + { + let services = service_ctx.clone(); + let loading = loading.clone(); + let counter = counter.clone(); + use_effect_with((), move |_| { + let sub_id = services.event.subscribe(move |msg| { + if let EventMessage::Busy(status) = msg { + match status { + BusyStatus::Show => { + *counter.borrow_mut() += 1; + loading.set(true); + } + BusyStatus::Hide => { + let mut cnt = counter.borrow_mut(); + if *cnt > 0 { + *cnt -= 1; + } + if *cnt == 0 { + loading.set(false); + } + } + } + } + }); + move || services.event.unsubscribe(sub_id) + }); + } + + html! { +
+ +
+ } } \ No newline at end of file diff --git a/webui/src/app/components/playlist/mod.rs b/webui/src/app/components/playlist/mod.rs index ff6f5c459..724f6b826 100644 --- a/webui/src/app/components/playlist/mod.rs +++ b/webui/src/app/components/playlist/mod.rs @@ -15,6 +15,7 @@ mod input_table; mod input; mod source_selector; mod playlist_explorer; +mod playlist_update_view; use std::rc::Rc; use yew_i18n::YewI18n; @@ -35,6 +36,7 @@ pub use self::input::*; pub use self::source_selector::*; pub use self::playlist_editor_view::*; pub use self::playlist_explorer_view::*; +pub use self::playlist_update_view::*; pub fn make_tags(data: &[(bool, &str)], translate: &YewI18n) -> Vec> { data.iter() diff --git a/webui/src/app/components/playlist/playlist_explorer.rs b/webui/src/app/components/playlist/playlist_explorer.rs index 743fcc413..f9a8aa1e3 100644 --- a/webui/src/app/components/playlist/playlist_explorer.rs +++ b/webui/src/app/components/playlist/playlist_explorer.rs @@ -5,7 +5,8 @@ use crate::app::context::PlaylistExplorerContext; use yew::prelude::*; use shared::model::{CommonPlaylistItem, SearchRequest, UiPlaylistGroup, XtreamCluster}; use crate::app::components::{IconButton, NoContent, Search}; -use crate::app::components::loading_indicator::LoadingIndicator; +use crate::hooks::use_service_context; +use crate::model::{BusyStatus, EventMessage}; enum ExplorerLevel { Categories, @@ -15,7 +16,7 @@ enum ExplorerLevel { #[function_component] pub fn PlaylistExplorer() -> Html { let context = use_context::().expect("PlaylistExplorer context not found"); - let loading = use_state(|| false); + let service_ctx = use_service_context(); let current_item = use_state(|| ExplorerLevel::Categories); let playlist = use_state(|| (*context.playlist).clone()); @@ -40,7 +41,7 @@ pub fn PlaylistExplorer() -> Html { }; let handle_search = { - let set_loading = loading.clone(); + let services = service_ctx.clone(); let set_playlist = playlist.clone(); let set_current_item = current_item.clone(); let context = context.clone(); @@ -49,11 +50,11 @@ pub fn PlaylistExplorer() -> Html { SearchRequest::Clear => set_playlist.set((*context.playlist).clone()), SearchRequest::Text(_) | SearchRequest::Regexp(_) => { - set_loading.set(true); - let set_loading = set_loading.clone(); + services.event.broadcast(EventMessage::Busy(BusyStatus::Show)); let set_playlist = set_playlist.clone(); let set_current_item = set_current_item.clone(); let context = context.clone(); + let services = services.clone(); spawn_local(async move { let filtered = context .playlist @@ -62,7 +63,7 @@ pub fn PlaylistExplorer() -> Html { .map(Rc::new); set_playlist.set(filtered); set_current_item.set(ExplorerLevel::Categories); - set_loading.set(false); + services.event.broadcast(EventMessage::Busy(BusyStatus::Hide)); }); } } @@ -166,7 +167,6 @@ pub fn PlaylistExplorer() -> Html { html! {
-