diff --git a/Cargo.lock b/Cargo.lock index 8a7890b16..2fb40753b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,15 +2,6 @@ # It is not intended for manual editing. version = 4 -[[package]] -name = "addr2line" -version = "0.25.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1b5d307320b3181d6d7954e663bd7c774a838b8220fe0593c86d9fb09f498b4b" -dependencies = [ - "gimli", -] - [[package]] name = "adler2" version = "2.0.1" @@ -239,21 +230,6 @@ dependencies = [ "syn 2.0.106", ] -[[package]] -name = "backtrace" -version = "0.3.76" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bb531853791a215d7c62a30daf0dde835f381ab5de4589cfe7c649d2cbe92bd6" -dependencies = [ - "addr2line", - "cfg-if", - "libc", - "miniz_oxide", - "object", - "rustc-demangle", - "windows-link 0.2.1", -] - [[package]] name = "base16ct" version = "0.2.0" @@ -1022,7 +998,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -1337,12 +1313,6 @@ dependencies = [ "wasm-bindgen", ] -[[package]] -name = "gimli" -version = "0.32.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e629b9b98ef3dd8afe6ca2bd0f89306cec16d43d907889945bc5d6687f2f13c7" - [[package]] name = "gloo" version = "0.8.1" @@ -2248,17 +2218,6 @@ dependencies = [ "libc", ] -[[package]] -name = "io-uring" -version = "0.7.10" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "046fa2d4d00aea763528b4950358d0ead425372445dc8ff86312b3c69ff7727b" -dependencies = [ - "bitflags 2.9.4", - "cfg-if", - "libc", -] - [[package]] name = "ipnet" version = "2.11.0" @@ -2342,9 +2301,9 @@ dependencies = [ [[package]] name = "jsonwebtoken" -version = "10.0.0" +version = "10.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f1417155a38e99d7704ddb3ea7445fe57fdbd5d756d727740a9ed8b9ebaed6e1" +checksum = "3d119c6924272d16f0ab9ce41f7aa0bfef9340c00b0bb7ca3dd3b263d4a9150b" dependencies = [ "base64", "ed25519-dalek", @@ -2682,15 +2641,6 @@ dependencies = [ "libc", ] -[[package]] -name = "object" -version = "0.37.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ff76201f031d8863c38aa7f905eca4f53abbfa15f609db4277d44cd8938f33fe" -dependencies = [ - "memchr", -] - [[package]] name = "once_cell" version = "1.21.3" @@ -3365,9 +3315,9 @@ dependencies = [ [[package]] name = "regex" -version = "1.12.1" +version = "1.12.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4a52d8d02cacdb176ef4678de6c052efb4b3da14b78e4db683a4252762be5433" +checksum = "843bc0191f75f3e22651ae5f1e72939ab2f72a4bc30fa80a066bd66edefc24d4" dependencies = [ "aho-corasick", "memchr", @@ -3547,12 +3497,6 @@ dependencies = [ "crossbeam-utils", ] -[[package]] -name = "rustc-demangle" -version = "0.1.26" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "56f7d92ca342cea22a06f2121d944b4fd82af56988c270852495420f961d4ace" - [[package]] name = "rustc-hash" version = "2.1.1" @@ -3578,7 +3522,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -4057,7 +4001,7 @@ dependencies = [ "getrandom 0.3.3", "once_cell", "rustix", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -4160,28 +4104,25 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" [[package]] name = "tokio" -version = "1.47.1" +version = "1.48.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "89e49afdadebb872d3145a5638b59eb0691ea23e46ca484037cfab3b76b95038" +checksum = "ff360e02eab121e0bc37a2d3b4d4dc622e6eda3a8e5253d5435ecf5bd4c68408" dependencies = [ - "backtrace", "bytes", - "io-uring", "libc", "mio", "parking_lot", "pin-project-lite", - "slab", "socket2", "tokio-macros", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] name = "tokio-macros" -version = "2.5.0" +version = "2.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6e06d43f1345a3bcd39f6a56dbb7dcab2ba47e68e8ac134855e7e2bdbaf8cab8" +checksum = "af407857209536a95c8e56f8231ef2c2e2aff839b22e07a1ffcbc617e9db9fa5" dependencies = [ "proc-macro2", "quote", @@ -4840,7 +4781,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 2d7f7ad1d..813eca38d 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -22,7 +22,7 @@ axum = { version = "0.8" , features = ["macros", "default", "ws"]} tower = "0.5" tower-http = { version = "0.6", features = ["cors", "auth", "fs", "compression-full", "trace"] } tower_governor = { version = "0.8", features = ["axum"] } -jsonwebtoken = { version = "10.0", features = ["rust_crypto"] } +jsonwebtoken = { version = "10.1", features = ["rust_crypto"] } rust-argon2 = "3" futures = "0.3" path-clean = "1.0" @@ -41,7 +41,7 @@ flate2 = "1.1" blake3 = "1.8" bytes = "1.10" tokio-stream = { version = "0.1", features = ["sync"] } -tokio = { version = "1.47", features = ["rt-multi-thread", "parking_lot", "fs"] } +tokio = { version = "1.48", features = ["rt-multi-thread", "parking_lot", "fs"] } #tokio = { version = "1.46", features = ["rt-multi-thread", "parking_lot", "fs", "tracing"] } #console-subscriber = "0" #tracing = "0.1" diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 1a22a50c2..1c9cc921b 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -1056,6 +1056,8 @@ fn get_add_cache_content( let cache = Arc::clone(cache); let add_cache_content: Arc = Arc::new(move |size| { let res_url = resource_url.clone(); + + // todo spawn, replace with unboundchannel let cache = Arc::clone(&cache); tokio::spawn(async move { if let Some(cache) = cache.load().as_ref() { diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index e040f5263..f73642ce0 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -63,7 +63,7 @@ async fn playlist_update( let app_config = Arc::clone(&app_state.app_config); let event_manager = Arc::clone(&app_state.event_manager); let playlist_state = Arc::clone(&app_state.playlists); - tokio::spawn(playlist::exec_processing(Arc::clone(&app_state.http_client.load()), app_config, Arc::new(valid_targets), Some(event_manager), Some(playlist_state))); + playlist::exec_processing(Arc::clone(&app_state.http_client.load()), app_config, Arc::new(valid_targets), Some(event_manager), Some(playlist_state)).await; axum::http::StatusCode::ACCEPTED.into_response() } Err(err) => { diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index d582b5430..9f3a53ca1 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -72,12 +72,12 @@ fn create_shared_data( let (provider_change_tx, provider_change_rx) = tokio::sync::mpsc::unbounded_channel(); let active_provider = Arc::new(ActiveProviderManager::new(app_config, provider_change_tx)); let (active_user_change_tx, active_user_change_rx) = tokio::sync::mpsc::unbounded_channel(); - let active_users = Arc::new(ActiveUserManager::new( + let active_users = ActiveUserManager::new( &config, &shared_stream_manager, &active_provider, active_user_change_tx, - )); + ); let event_manager = Arc::new(EventManager::new(active_user_change_rx, provider_change_rx, )); let client = create_http_client(app_config); diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 4c42fca18..82ea42353 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -41,6 +41,7 @@ impl ProviderConnectionGuard { fn send_release(&self, config: &Arc) { let provider_config = Arc::clone(config); if let Err(_err) = &self.release_tx.send(Arc::clone(config)) { + // Fallback tokio::spawn(async move { provider_config.release().await; }); diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index 8d79c5bb1..af67c060f 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -1,5 +1,5 @@ use std::borrow::Cow; -use crate::api::model::ActiveProviderManager; +use crate::api::model::{ActiveProviderManager}; use crate::api::model::SharedStreamManager; use crate::model::Config; use crate::model::ProxyUserCredentials; @@ -10,6 +10,7 @@ use shared::utils::{current_time_secs, default_grace_period_millis, default_grac use std::collections::HashMap; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; +use tokio::sync::mpsc::{unbounded_channel, UnboundedSender}; use tokio::sync::RwLock; @@ -40,7 +41,7 @@ macro_rules! active_user_manager_shared_impl { let user_connection_count = Self::get_active_connections(&user).await; let user_count = user.read().await.iter().filter(|(_, c)| c.connections > 0).count(); if let Err(err) = self.connection_change_tx.send(ActiveUserConnectionChange::Connections(user_count, user_connection_count)) { - error!("Failed to send active user connection change: user-count: {user_count}, user-connection-count: {user_connection_count] {err:?}"); + error!("Failed to send active user connection change: user-count: {user_count}, user-connection-count: {user_connection_count} {err:?}"); } if is_log_user_enabled { info!("Active Users: {user_count}, Active User Connections: {user_connection_count}"); @@ -69,7 +70,7 @@ macro_rules! active_user_manager_shared_impl { self.drop_connection(&addr); self.shared_stream_manager.release_connection(addr, true).await; self.provider_manager.release_connection(addr).await; - let Err(err) = self.connection_change_tx.send(ActiveUserConnectionChange::Disconnected(addr.to_string())) { + if let Err(err) = self.connection_change_tx.send(ActiveUserConnectionChange::Disconnected(addr.to_string())) { error!("Failed to send active user connection change: {err:?}"); } self.log_active_user().await; @@ -84,7 +85,7 @@ fn get_grace_options(config: &Config) -> (u64, u64) { (grace_period_millis, grace_period_timeout_secs) } -struct ConnectionGuardUserManager { +pub struct ConnectionGuardUserManager { log_active_user: bool, user: Arc>>, user_by_addr: Arc>>, @@ -105,18 +106,28 @@ pub struct UserConnectionGuard { manager: Arc, // username: String, addr: String, + release_tx: UnboundedSender, } + +impl UserConnectionGuard { + pub fn new(manager: Arc, addr: &str, release_tx: UnboundedSender) -> Self { + Self { + manager, + addr: addr.to_string(), + release_tx, + } + } +} + impl Drop for UserConnectionGuard { fn drop(&mut self) { let manager = self.manager.clone(); let addr = self.addr.clone(); - if let Ok(rt) = tokio::runtime::Handle::try_current() { - rt.spawn(async move { + if let Err(_err) = self.release_tx.send(addr.clone()) { + // fallback + tokio::spawn(async move { manager.remove_connection(&addr).await; }); - } else { - // Fallback: no runtime - error!("Runtime not available, cannot cleanly remove connection for {addr}"); } } } @@ -178,14 +189,19 @@ pub struct ActiveUserManager { shared_stream_manager: Arc, provider_manager: Arc, connection_change_tx: ActiveUserConnectionChangeSender, + release_tx: UnboundedSender, } impl ActiveUserManager { - pub fn new(config: &Config, shared_stream_manager: &Arc, provider_manager: &Arc, connection_change_tx: ActiveUserConnectionChangeSender) -> Self { - let log_active_user = config.log.as_ref().is_some_and(|l| l.log_active_user); + pub fn new(config: &Config, shared_stream_manager: &Arc, provider_manager: &Arc, connection_change_tx: ActiveUserConnectionChangeSender) -> Arc { + let log_active_user: bool = config.log.as_ref().is_some_and(|l| l.log_active_user); let (grace_period_millis, grace_period_timeout_secs) = get_grace_options(config); let (close_signal_tx, _) = tokio::sync::broadcast::channel(10); - Self { + + // Create the cleanup channel + let (cleanup_tx, mut cleanup_rx) = unbounded_channel::(); + + let active_user_manager = Arc::new(Self { grace_period_millis: AtomicU64::new(grace_period_millis), grace_period_timeout_secs: AtomicU64::new(grace_period_timeout_secs), log_active_user: AtomicBool::new(log_active_user), @@ -196,11 +212,29 @@ impl ActiveUserManager { shared_stream_manager: Arc::clone(shared_stream_manager), provider_manager: Arc::clone(provider_manager), connection_change_tx, - } + release_tx: cleanup_tx, + }); + + let active_user_manager_clone = Arc::clone(&active_user_manager); + // Spawn the async cleanup worker + tokio::spawn(async move { + while let Some(addr) = cleanup_rx.recv().await { + debug!("🧹 User manager - connection releasing {:?}", addr); + active_user_manager_clone.remove_connection(&addr).await + } + debug!("User manager - cleanup worker terminated"); + }); + + + active_user_manager } active_user_manager_shared_impl!(); + pub fn release_sender(&self) -> UnboundedSender { + self.release_tx.clone() + } + pub fn update_config(&self, config: &Config) { let log_active_user = config.log.as_ref().is_some_and(|l| l.log_active_user); let (grace_period_millis, grace_period_timeout_secs) = get_grace_options(config); @@ -316,10 +350,7 @@ impl ActiveUserManager { } self.log_active_user().await; - UserConnectionGuard { - manager: Arc::new(self.clone_inner()), - addr: addr.to_owned(), - } + UserConnectionGuard::new(Arc::new(self.clone_inner()), addr, self.release_sender()) } fn is_log_user_enabled(&self) -> bool { diff --git a/backend/src/api/model/event_manager.rs b/backend/src/api/model/event_manager.rs index 46a738d07..0275b3b72 100644 --- a/backend/src/api/model/event_manager.rs +++ b/backend/src/api/model/event_manager.rs @@ -1,5 +1,4 @@ use log::{info, trace}; -use tokio::task; use shared::model::{ActiveUserConnectionChange, ConfigType, PlaylistUpdateState}; use crate::api::model::{ActiveUserConnectionChangeReceiver}; use crate::api::model::{ProviderConnectionChangeReceiver}; @@ -27,7 +26,7 @@ impl EventManager { let (channel_tx, _channel_rx) = tokio::sync::broadcast::channel(10); let channel_tx_clone = channel_tx.clone(); - task::spawn(async move { + tokio::spawn(async move { loop { tokio::select! { Some(event) = active_user_change_rx.recv() => { diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index a31b092ff..1d02246ac 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -103,10 +103,7 @@ pub struct SharedStreamState { impl Drop for SharedStreamState { fn drop(&mut self) { if let Some(guard) = self.provider_guard.as_ref() { - let guard = guard.clone(); - tokio::spawn(async move { - guard.force_release(); - }); + guard.force_release(); } } } diff --git a/backend/src/api/serve.rs b/backend/src/api/serve.rs index caad8de58..f7b03398b 100644 --- a/backend/src/api/serve.rs +++ b/backend/src/api/serve.rs @@ -134,11 +134,7 @@ where let user_manager_clone = Arc::clone(&user_manager); let mut addr_close_rx = user_manager_clone.get_close_connection_channel(); - let connection_closed = async move || { - debug!("Connection closed: {remote_addr}"); - let addr = remote_addr.to_string(); - user_manager_clone.remove_connection(&addr).await; - }; + let connection_release = user_manager.release_sender(); debug!("Connection opened: {addr_str}"); @@ -148,11 +144,17 @@ where if let Err(err) = result { trace!("failed to serve connection: {err:#}"); } - connection_closed().await; + if let Err(_err) = connection_release.send(remote_addr.to_string()) { + let addr = remote_addr.to_string(); + user_manager_clone.remove_connection(&addr).await; + } break; } () = &mut signal_closed => { - connection_closed().await; + if let Err(_err) = connection_release.send(remote_addr.to_string()) { + let addr = remote_addr.to_string(); + user_manager_clone.remove_connection(&addr).await; + } debug!("Connection gracefully closed: {remote_addr}"); conn.as_mut().graceful_shutdown(); } diff --git a/backend/src/messaging.rs b/backend/src/messaging.rs index 5411cac44..29eed38d2 100644 --- a/backend/src/messaging.rs +++ b/backend/src/messaging.rs @@ -11,27 +11,23 @@ fn is_enabled(kind: MsgKind, cfg: &MessagingConfig) -> bool { cfg.notify_on.contains(&kind) } -fn send_http_post_request(client: &Arc, msg: &str, messaging: &MessagingConfig) { +async fn send_http_post_request(client: &Arc, msg: &str, messaging: &MessagingConfig) { if let Some(rest) = &messaging.rest { - let url = rest.url.clone(); - let data = msg.to_owned(); - let the_client = Arc::clone(client); - tokio::spawn(async move { - match the_client - .post(&url) - .header(header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) - .body(data) - .send() - .await - { - Ok(_) => debug!("Text message sent successfully to rest api"), - Err(e) => error!("Text message wasn't sent to rest api because of: {e}"), - } - }); + let data = msg.to_owned(); + match client + .post(&rest.url) + .header(header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) + .body(data) + .send() + .await + { + Ok(_) => debug!("Text message sent successfully to rest api"), + Err(e) => error!("Text message wasn't sent to rest api because of: {e}"), + } } } -fn send_telegram_message(client: &Arc, msg: &str, messaging: &MessagingConfig, json: bool) { +async fn send_telegram_message(client: &Arc, msg: &str, messaging: &MessagingConfig, json: bool) { // TODO use proxy settings if let Some(telegram) = &messaging.telegram { let (message, options) = { @@ -48,55 +44,51 @@ fn send_telegram_message(client: &Arc, msg: &str, messaging: &M for chat_id in &telegram.chat_ids { let bot = telegram_create_instance(&telegram.bot_token, chat_id); - telegram_send_message(client, &bot, &message, options.as_ref()); + telegram_send_message(client, &bot, &message, options.as_ref()).await; } } } -fn send_pushover_message(client: &Arc, msg: &str, messaging: &MessagingConfig) { +async fn send_pushover_message(client: &Arc, msg: &str, messaging: &MessagingConfig) { if let Some(pushover) = &messaging.pushover { let encoded_message: String = url::form_urlencoded::Serializer::new(String::new()) .append_pair("token", pushover.token.as_str()) .append_pair("user", pushover.user.as_str()) .append_pair("message", msg) .finish(); - let the_client = Arc::clone(client); - let pushover_url = pushover.url.clone(); - tokio::spawn(async move { - match the_client - .post(pushover_url) - .header(header::CONTENT_TYPE, mime::APPLICATION_WWW_FORM_URLENCODED.to_string()) - .body(encoded_message) - .send() - .await - { - Ok(response) => { - if response.status().is_success() { - debug!("Text message sent successfully to PUSHOVER, status code {}", response.status()); - } else { - error!("Failed to send text message to PUSHOVER, status code {}", response.status()); - } + match client + .post(&pushover.url) + .header(header::CONTENT_TYPE, mime::APPLICATION_WWW_FORM_URLENCODED.to_string()) + .body(encoded_message) + .send() + .await + { + Ok(response) => { + if response.status().is_success() { + debug!("Text message sent successfully to PUSHOVER, status code {}", response.status()); + } else { + error!("Failed to send text message to PUSHOVER, status code {}", response.status()); } - Err(e) => error!("Text message wasn't sent to PUSHOVER api because of: {e}"), } - }); - } -} - -fn dispatch_send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str, json: bool) { - if let Some(messaging) = cfg { - if is_enabled(kind, messaging) { - send_telegram_message(client, msg, messaging, json); - send_http_post_request(client, msg, messaging); - send_pushover_message(client, msg, messaging); + Err(e) => error!("Text message wasn't sent to PUSHOVER api because of: {e}"), } } } -pub fn send_message_json(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { - dispatch_send_message(client, kind, cfg, msg, true); +async fn dispatch_send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str, json: bool) { + if let Some(messaging) = cfg { + if is_enabled(kind, messaging) { + send_telegram_message(client, msg, messaging, json).await; + send_http_post_request(client, msg, messaging).await; + send_pushover_message(client, msg, messaging).await; + } + } } -pub fn send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { - dispatch_send_message(client, kind, cfg, msg, false); +pub async fn send_message_json(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { + dispatch_send_message(client, kind, cfg, msg, true).await; +} + +pub async fn send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { + dispatch_send_message(client, kind, cfg, msg, false).await; } diff --git a/backend/src/processing/playlist_watch.rs b/backend/src/processing/playlist_watch.rs index 650a207d5..3c77bbe22 100644 --- a/backend/src/processing/playlist_watch.rs +++ b/backend/src/processing/playlist_watch.rs @@ -8,7 +8,7 @@ use crate::model::Config; use crate::utils; use crate::utils::{bincode_deserialize, bincode_serialize}; -pub fn process_group_watch(client: &Arc, cfg: &Config, target_name: &str, pl: &PlaylistGroup) { +pub async fn process_group_watch(client: &Arc, cfg: &Config, target_name: &str, pl: &PlaylistGroup) { let mut new_tree = BTreeSet::new(); pl.channels.iter().for_each(|chan| { let header = &chan.header; @@ -28,7 +28,7 @@ pub fn process_group_watch(client: &Arc, cfg: &Config, target_n let removed_difference: BTreeSet = loaded_tree.difference(&new_tree).cloned().collect(); if !added_difference.is_empty() || !removed_difference.is_empty() { changed = true; - handle_watch_notification(client, cfg, &added_difference, &removed_difference, target_name, &pl.title); + handle_watch_notification(client, cfg, &added_difference, &removed_difference, target_name, &pl.title).await; } } else { error!("failed to load watch_file {}", &path.to_str().unwrap_or_default()); @@ -60,7 +60,7 @@ struct WatchChanges { pub removed: Vec, } -fn handle_watch_notification(client: &Arc, cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { +async fn handle_watch_notification(client: &Arc, cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { let added = added.iter().map(std::string::ToString::to_string).collect::>(); let removed = removed.iter().map(std::string::ToString::to_string).collect::>(); if !added.is_empty() || !removed.is_empty() { @@ -73,7 +73,7 @@ fn handle_watch_notification(client: &Arc, cfg: &Config, added: let msg = serde_json::to_string_pretty(&changes).unwrap_or_else(|_| "Error: Failed to serialize watch changes".to_string()); info!("{}", &msg); - send_message(client, MsgKind::Watch, cfg.messaging.as_ref(), &msg); + send_message(client, MsgKind::Watch, cfg.messaging.as_ref(), &msg).await; } } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 8e36a87f7..d2a573de7 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -579,7 +579,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, step.tick("assigning channel counter"); let config = app_config.config.load(); - if process_watch(&config, &client, target, &flat_new_playlist) { + if process_watch(&config, &client, target, &flat_new_playlist).await { step.tick("group watches"); } let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, playlist_state).await; @@ -620,14 +620,14 @@ async fn process_epg(processed_fetched_playlists: &mut Vec>) (new_epg, new_playlist) } -fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTarget, new_playlist: &Vec) -> bool { +async fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTarget, new_playlist: &Vec) -> bool { if let Some(watches) = &target.watch { if default_as_default().eq_ignore_ascii_case(&target.name) { error!("cant watch a target with no unique name"); } else { for pl in new_playlist { if watches.iter().any(|r| r.is_match(&pl.title)) { - process_group_watch(client, cfg, &target.name, pl); + process_group_watch(client, cfg, &target.name, pl).await; } } } @@ -657,7 +657,7 @@ pub async fn exec_processing(client: Arc, app_config: Arc error!("Failed to serialize playlist stats {err}"), } @@ -672,7 +672,7 @@ pub async fn exec_processing(client: Arc, app_config: Arc, input: &Input if let Ok(cur_status) = ProxyUserStatus::from_str(&status) { if !matches!(cur_status, ProxyUserStatus::Active | ProxyUserStatus::Trial) { warn!("User status for user {username} is {cur_status:?}"); - send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")).await; } } } @@ -175,11 +175,11 @@ async fn xtream_login(cfg: &Config, client: &Arc, input: &Input let datetime = DateTime::from_timestamp(expiration_timestamp, 0).unwrap(); let formatted = datetime.format("%Y-%m-%d %H:%M:%S").to_string(); warn!("User account for user {username} expires {formatted}"); - send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")).await; } } else { warn!("User account for user {username} is expired"); - send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")).await; } } } diff --git a/backend/src/utils/telegram.rs b/backend/src/utils/telegram.rs index 75948037e..77bdceb92 100644 --- a/backend/src/utils/telegram.rs +++ b/backend/src/utils/telegram.rs @@ -62,7 +62,7 @@ pub fn telegram_create_instance(bot_token: &str, chat_id: &str) -> BotInstance { } } -pub fn telegram_send_message( +pub async fn telegram_send_message( client: &Arc, instance: &BotInstance, msg: &str, @@ -87,27 +87,25 @@ pub fn telegram_send_message( .map(ToString::to_string), }; - let the_client = Arc::clone(client); - tokio::spawn(async move { - let result = the_client - .post(url) - .json(&request_json_obj) - .send() - .await; + let result = client + .post(url) + .json(&request_json_obj) + .send() + .await; - match result { - Ok(response) => { - if response.status().is_success() { - debug!("Message sent successfully to {chat_id} telegram api"); - } else { - match response.json::().await { - Ok(json) => error!("Message wasn't sent to {chat_id} telegram api because of: {}", json.description), - Err(_) => error!("Message wasn't sent to {chat_id} telegram api. Telegram response could not be parsed!"), - } + match result { + Ok(response) => { + if response.status().is_success() { + debug!("Message sent successfully to {chat_id} telegram api"); + } else { + match response.json::().await { + Ok(json) => error!("Message wasn't sent to {chat_id} telegram api because of: {}", json.description), + Err(_) => error!("Message wasn't sent to {chat_id} telegram api. Telegram response could not be parsed!"), } - }, - Err(e) => error!("Message wasn't sent to {chat_id} telegram api because of: {e}"), - } - }); + } + }, + Err(e) => error!("Message wasn't sent to {chat_id} telegram api because of: {e}"), + } } + diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index c144e8b76..3e4ed7db5 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -28,7 +28,7 @@ futures = "0.3" prost = "0" wasm-bindgen-futures = "0" bytes = "1" -regex = "1.12.1" +regex = "1.12.2" base64 = "0.22.1" cron = "0.15" fastrand = "2.3.0"