diff --git a/Cargo.lock b/Cargo.lock index 0bf16ce65..14e14fd55 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -908,7 +908,6 @@ name = "frontend" version = "0.1.0" dependencies = [ "anyhow", - "bincode 2.0.1", "bytes", "futures", "futures-signals", @@ -3333,7 +3332,6 @@ name = "shared" version = "0.1.0" dependencies = [ "base64 0.22.1", - "bincode 2.0.1", "bitflags 2.9.1", "blake3", "bytes", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index a25fc75ce..bf862a049 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -34,7 +34,7 @@ mime = "0.3" log = "0.4" env_logger = "0.11" rustelebot = "0.3" -bincode = { version = "2", features = ["std", "serde"] } +bincode = { version = "2.0.1", features = ["std", "serde"] } rand = "0.9" rpassword = "7.4" flate2 = "1" diff --git a/backend/src/api/endpoints/websocket_api.rs b/backend/src/api/endpoints/websocket_api.rs index 8a319eef9..2f2aedef8 100644 --- a/backend/src/api/endpoints/websocket_api.rs +++ b/backend/src/api/endpoints/websocket_api.rs @@ -83,29 +83,32 @@ async fn handle_handshake( async fn handle_protocol_message( msg: Message, - socket: &mut WebSocket, app_state: &Arc, auth: bool, secret_key: Option<&Vec>, -) -> Result<(), String> { +) -> Option { if let Message::Binary(bytes) = msg { match ProtocolMessage::from_bytes(bytes) { Ok(ProtocolMessage::StatusRequest(auth_token)) => { if !auth || verify_auth_admin_token(&auth_token, secret_key) { - let status = create_status_check(app_state).await; - let response = ProtocolMessage::StatusResponse(status).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(response)).await.map_err(|e| e.to_string())?; + let status = create_status_check(app_state).await; + Some(ProtocolMessage::StatusResponse(status)) + } else { + Some(ProtocolMessage::Unauthorized) } } Ok(_) => { error!("Unexpected protocol message after handshake"); + None } Err(e) => { error!("Invalid websocket message: {e}"); + Some(ProtocolMessage::Error(format!("Invalid websocket message: {e}"))) } } + } else { + None } - Ok(()) } async fn handle_incoming_message( @@ -124,7 +127,19 @@ async fn handle_incoming_message( *handler = ProtocolHandler::Default; Ok(()) }, - ProtocolHandler::Default => handle_protocol_message(msg, socket, app_state, auth, secret_key).await, + ProtocolHandler::Default => { + let msg = handle_protocol_message(msg, app_state, auth, secret_key).await; + match msg { + None => {Ok(())}, + Some(protocol_msg) => { + let bytes = match protocol_msg.to_bytes() { + Ok(bytes) => bytes, + Err(err) => ProtocolMessage::Error(err.to_string()).to_bytes().map_err(|e| e.to_string())?, + }; + Ok(socket.send(Message::Binary(bytes)).await.map_err(|e| e.to_string())?) + } + } + }, } } @@ -132,11 +147,11 @@ async fn handle_event_message(socket: &mut WebSocket, event: EventMessage) -> Re match event { EventMessage::ActiveUserChange(users, connections) => { let msg = ProtocolMessage::ActiveUserResponse(users, connections).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(msg)).await.map_err(|e| e.to_string()) + socket.send(Message::Binary(msg)).await.map_err(|e| format!("Active user connection change event: {} ", e.to_string())) } EventMessage::ActiveProviderChange(provider, connections) => { let msg = ProtocolMessage::ActiveProviderResponse(provider, connections).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(msg)).await.map_err(|e| e.to_string()) + socket.send(Message::Binary(msg)).await.map_err(|e| format!("Provider connection change event: {} ", e.to_string())) } } @@ -165,7 +180,7 @@ async fn handle_socket(mut socket: WebSocket, app_state: Arc, auth: bo Ok(event) = event_rx.recv() => { if let Err(e) = handle_event_message(&mut socket, event).await { - error!("Failed to send active user change: {e}"); + error!("Failed to send event: {e}"); break; } } diff --git a/backend/src/processing/playlist_watch.rs b/backend/src/processing/playlist_watch.rs index 01cd46bfe..6ef0a6a90 100644 --- a/backend/src/processing/playlist_watch.rs +++ b/backend/src/processing/playlist_watch.rs @@ -3,10 +3,10 @@ use std::path::{Path}; use std::sync::Arc; use log::{error, info}; use shared::model::{MsgKind, PlaylistGroup}; -use shared::utils::{bincode_deserialize, bincode_serialize}; use crate::messaging::{send_message}; 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) { let mut new_tree = BTreeSet::new(); diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 030db216a..41402efc1 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -17,10 +17,10 @@ use std::io::{BufWriter, Write}; use std::sync::Arc; use std::time::Instant; use log::{error, info, log_enabled, Level}; -use shared::utils::bincode_serialize; use crate::model::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; use crate::utils; use crate::processing::processor::xtream::normalize_json_content; +use crate::utils::bincode_serialize; create_resolve_options_function_for_xtream_target!(series); diff --git a/backend/src/repository/bplustree.rs b/backend/src/repository/bplustree.rs index c57defeeb..b8e27cc0d 100644 --- a/backend/src/repository/bplustree.rs +++ b/backend/src/repository/bplustree.rs @@ -4,13 +4,13 @@ use ruzstd::decoding::StreamingDecoder; use ruzstd::encoding::{compress_to_vec, CompressionLevel}; use serde::{Deserialize, Serialize}; use shared::error::{str_to_io_error, to_io_error}; -use shared::utils::{bincode_deserialize, bincode_serialize}; use std::fs::File; use std::io::{self, BufReader, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::mem::size_of; use std::path::Path; use tempfile::NamedTempFile; +use crate::utils::{bincode_deserialize, bincode_serialize}; const BLOCK_SIZE: usize = 4096; const BINCODE_OVERHEAD: usize = 8; diff --git a/backend/src/repository/indexed_document.rs b/backend/src/repository/indexed_document.rs index 8c4965db8..e03a1fa22 100644 --- a/backend/src/repository/indexed_document.rs +++ b/backend/src/repository/indexed_document.rs @@ -9,8 +9,8 @@ use log::error; use serde::{Deserialize, Serialize}; use tempfile::NamedTempFile; use shared::error::{str_to_io_error, to_io_error}; -use shared::utils::{bincode_deserialize, bincode_serialize}; use crate::utils; +use crate::utils::{bincode_deserialize, bincode_serialize}; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index ec14207fe..3253fda35 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -10,7 +10,7 @@ use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_f use crate::repository::storage_const; use crate::repository::target_id_mapping::VirtualIdRecord; use crate::repository::xtream_playlist_iterator::XtreamPlaylistJsonIterator; -use crate::utils::FileReadGuard; +use crate::utils::{bincode_deserialize, FileReadGuard}; use crate::utils::file_reader; use crate::utils::open_readonly_file; use crate::utils::{json_write_documents_to_file}; @@ -25,7 +25,7 @@ use std::fs::File; use std::io::{BufReader, BufWriter, Error, ErrorKind, Read, Write}; use std::path::{Path, PathBuf}; use shared::model::{PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; -use shared::utils::{bincode_deserialize, generate_playlist_uuid, get_u32_from_serde_value, hex_encode, json_iter_array}; +use shared::utils::{generate_playlist_uuid, get_u32_from_serde_value, hex_encode, json_iter_array}; macro_rules! cant_write_result { ($path:expr, $err:expr) => { diff --git a/shared/src/utils/bincode_utils.rs b/backend/src/utils/bincode_utils.rs similarity index 76% rename from shared/src/utils/bincode_utils.rs rename to backend/src/utils/bincode_utils.rs index ae6a2b75c..76ed80ee0 100644 --- a/shared/src/utils/bincode_utils.rs +++ b/backend/src/utils/bincode_utils.rs @@ -1,5 +1,6 @@ use std::io; -use crate::error::to_io_error; +use log::error; +use shared::error::to_io_error; #[inline] pub fn bincode_serialize(value: &T) -> io::Result> @@ -16,6 +17,9 @@ where { match bincode::serde::decode_from_slice(value, bincode::config::legacy()) { Ok((instance, _size)) => Ok(instance), - Err(e) => Err(to_io_error(e)), + Err(e) => { + error!("Failed to decode {e}"); + Err(to_io_error(e)) + }, } } diff --git a/backend/src/utils/mod.rs b/backend/src/utils/mod.rs index e5174e78a..9f346cb9c 100644 --- a/backend/src/utils/mod.rs +++ b/backend/src/utils/mod.rs @@ -7,7 +7,9 @@ mod step_measure; mod logging; mod trakt; mod json_utils; +mod bincode_utils; +pub use self::bincode_utils::*; pub use self::logging::*; pub use self::trakt::*; diff --git a/shared/Cargo.toml b/shared/Cargo.toml index db7fadd55..bd07fbd04 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -21,4 +21,3 @@ fastrand = "2" zeroize = "1" chrono = "0.4.41" bytes = "1" -bincode = "2" diff --git a/shared/src/model/web_socket.rs b/shared/src/model/web_socket.rs index acfa1f9ec..6687ded8d 100644 --- a/shared/src/model/web_socket.rs +++ b/shared/src/model/web_socket.rs @@ -1,7 +1,6 @@ use std::io; use bytes::Bytes; use crate::model::StatusCheck; -use crate::utils::{bincode_deserialize, bincode_serialize}; use serde::{Deserialize, Serialize}; pub const PROTOCOL_VERSION: u8 = 1; @@ -49,6 +48,8 @@ impl WsCloseCode { #[derive(Serialize, Deserialize, Debug)] pub enum ProtocolMessage { + Unauthorized, + Error(String), Version(u8), StatusRequest(String), StatusResponse(StatusCheck), @@ -63,8 +64,10 @@ impl ProtocolMessage { Ok(Bytes::from(vec![*version])) } _ => { - let encoded = bincode_serialize(self)?; - Ok(Bytes::from(encoded)) + //let encoded = bincode_serialize(self)?; + let json = serde_json::to_string(self) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + Ok(Bytes::from(json.into_bytes())) } } } @@ -73,7 +76,11 @@ impl ProtocolMessage { if bytes.len() == 1 { Ok(ProtocolMessage::Version(bytes[0])) } else { - bincode_deserialize::(bytes.as_ref()) + //bincode_deserialize::(bytes.as_ref()) + let s = std::str::from_utf8(&bytes) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + serde_json::from_str(s) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e)) } } } \ No newline at end of file diff --git a/shared/src/utils/mod.rs b/shared/src/utils/mod.rs index debc43dc6..bb4e1efc7 100644 --- a/shared/src/utils/mod.rs +++ b/shared/src/utils/mod.rs @@ -8,7 +8,6 @@ mod directed_graph; mod hash_utils; mod json_utils; mod serde_utils; -mod bincode_utils; pub use self::default_utils::*; pub use self::time_utils::*; @@ -20,4 +19,3 @@ pub use self::directed_graph::*; pub use self::hash_utils::*; pub use self::json_utils::*; pub use self::serde_utils::*; -pub use self::bincode_utils::*; diff --git a/webui/Cargo.toml b/webui/Cargo.toml index b32e64c83..0973959e5 100644 --- a/webui/Cargo.toml +++ b/webui/Cargo.toml @@ -26,7 +26,6 @@ futures = "0.3" prost = "0" wasm-bindgen-futures = "0" bytes = "1" -bincode = { version = "2", features = ["std", "serde"] } [dependencies.web-sys] version = "0.3" diff --git a/webui/src/app/components/home.rs b/webui/src/app/components/home.rs index 5ea1488d8..6da27351c 100644 --- a/webui/src/app/components/home.rs +++ b/webui/src/app/components/home.rs @@ -1,16 +1,12 @@ -use std::cell::RefCell; -use std::collections::BTreeMap; use std::future; use std::rc::Rc; -use yew::platform::spawn_local; use yew::prelude::*; use yew::suspense::use_future; use shared::model::{AppConfigDto, StatusCheck}; use crate::app::components::{IconButton, Sidebar, DashboardView, PlaylistView, Panel, UserlistView, StatsView}; use crate::app::context::{ConfigContext, StatusContext}; use crate::model::ViewType; -use crate::hooks::use_service_context; -use crate::services::WsMessage; +use crate::hooks::{use_server_status, use_service_context}; #[function_component] pub fn Home() -> Html { @@ -18,7 +14,6 @@ pub fn Home() -> Html { let app_title = services.config.ui_config.app_title.as_ref().map_or("tuliprox", |v| v.as_str()); let config = use_state(|| None::>); let status = use_state(|| None::>); - let status_holder = use_state(|| Rc::new(RefCell::new(None::>))); let view_visible = use_state(|| ViewType::Users); @@ -32,6 +27,8 @@ pub fn Home() -> Html { Callback::from(move |view| view_vis.set(view)) }; + let _ = use_server_status(status.clone()); + { // first register for config update let services_ctx = services.clone(); @@ -53,65 +50,65 @@ pub fn Home() -> Html { }); } - { - let services_ctx = services.clone(); - let status_signal = status.clone(); - let status_holder_signal = status_holder.clone(); - - use_effect_with((), move |_| { - let subid = services_ctx.websocket.subscribe(move |msg| { - match msg { - WsMessage::ServerStatus(server_status) => { - *status_holder_signal.borrow_mut() = Some(Rc::clone(&server_status)); - status_signal.set(Some(server_status)); - } - WsMessage::ActiveUser(user_count, connections) => { - let mut server_status = { - if let Some(old_status) = status_holder_signal.borrow().as_ref() { - (**old_status).clone() - } else { - StatusCheck::default() - } - }; - server_status.active_users = user_count; - server_status.active_user_connections = connections; - let new_status = Rc::new(server_status); - *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); - status_signal.set(Some(new_status)); - } - WsMessage::ActiveProvider(provider, connections) => { - let mut server_status = { - if let Some(old_status) = status_holder_signal.borrow().as_ref() { - (**old_status).clone() - } else { - StatusCheck::default() - } - }; - if let Some(treemap) = server_status.active_provider_connections.as_mut() { - if connections == 0 { - treemap.remove(&provider); - } else { - treemap.insert(provider, connections); - } - } else if connections > 0 { - let mut treemap = BTreeMap::new(); - treemap.insert(provider, connections); - server_status.active_provider_connections = Some(treemap); - } - let new_status = Rc::new(server_status); - *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); - status_signal.set(Some(new_status)); - } - } - }); - let services_clone = services_ctx.clone(); - spawn_local(async move { - services_clone.websocket.get_server_status().await; - }); - let services_clone = services_ctx.clone(); - move || services_clone.websocket.unsubscribe(subid) - }); - } + // { + // let services_ctx = services.clone(); + // let status_signal = status.clone(); + // let status_holder_signal = status_holder.clone(); + // + // use_effect_with((), move |_| { + // let subid = services_ctx.websocket.subscribe(move |msg| { + // match msg { + // WsMessage::ServerStatus(server_status) => { + // *status_holder_signal.borrow_mut() = Some(Rc::clone(&server_status)); + // status_signal.set(Some(server_status)); + // } + // WsMessage::ActiveUser(user_count, connections) => { + // let mut server_status = { + // if let Some(old_status) = status_holder_signal.borrow().as_ref() { + // (**old_status).clone() + // } else { + // StatusCheck::default() + // } + // }; + // server_status.active_users = user_count; + // server_status.active_user_connections = connections; + // let new_status = Rc::new(server_status); + // *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); + // status_signal.set(Some(new_status)); + // } + // WsMessage::ActiveProvider(provider, connections) => { + // let mut server_status = { + // if let Some(old_status) = status_holder_signal.borrow().as_ref() { + // (**old_status).clone() + // } else { + // StatusCheck::default() + // } + // }; + // if let Some(treemap) = server_status.active_provider_connections.as_mut() { + // if connections == 0 { + // treemap.remove(&provider); + // } else { + // treemap.insert(provider, connections); + // } + // } else if connections > 0 { + // let mut treemap = BTreeMap::new(); + // treemap.insert(provider, connections); + // server_status.active_provider_connections = Some(treemap); + // } + // let new_status = Rc::new(server_status); + // *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); + // status_signal.set(Some(new_status)); + // } + // } + // }); + // let services_clone = services_ctx.clone(); + // spawn_local(async move { + // services_clone.websocket.get_server_status().await; + // }); + // let services_clone = services_ctx.clone(); + // move || services_clone.websocket.unsubscribe(subid) + // }); + // } let config_context = ConfigContext { config: (*config).clone(), diff --git a/webui/src/hooks/mod.rs b/webui/src/hooks/mod.rs index 770682e92..421fbbe58 100644 --- a/webui/src/hooks/mod.rs +++ b/webui/src/hooks/mod.rs @@ -1,5 +1,7 @@ mod use_service_context; mod use_icon_context; +mod use_server_status; pub use use_service_context::*; -pub use use_icon_context::*; \ No newline at end of file +pub use use_icon_context::*; +pub use use_server_status::*; \ No newline at end of file diff --git a/webui/src/hooks/use_server_status.rs b/webui/src/hooks/use_server_status.rs new file mode 100644 index 000000000..2a7d2adc1 --- /dev/null +++ b/webui/src/hooks/use_server_status.rs @@ -0,0 +1,89 @@ +use std::cell::RefCell; +use std::rc::Rc; +use std::collections::BTreeMap; +use gloo_timers::callback::Interval; +use yew::prelude::*; +use shared::model::StatusCheck; +use crate::hooks::use_service_context; +use crate::services::{ WsMessage}; +use yew::platform::spawn_local; + +#[hook] +pub fn use_server_status( + status: UseStateHandle>>, +) -> UseStateHandle>>> { + let services = use_service_context(); + let status_holder = use_state(|| RefCell::new(None::>)); + + { + let services_ctx = services.clone(); + let status_signal = status.clone(); + let status_holder_signal = status_holder.clone(); + + use_effect_with((), move |_| { + let subid = services_ctx.websocket.subscribe(move |msg| { + match msg { + WsMessage::ServerStatus(server_status) => { + *status_holder_signal.borrow_mut() = Some(Rc::clone(&server_status)); + status_signal.set(Some(server_status)); + } + WsMessage::ActiveUser(user_count, connections) => { + let mut server_status = { + if let Some(old_status) = status_holder_signal.borrow().as_ref() { + (**old_status).clone() + } else { + StatusCheck::default() + } + }; + server_status.active_users = user_count; + server_status.active_user_connections = connections; + let new_status = Rc::new(server_status); + *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); + status_signal.set(Some(new_status)); + } + WsMessage::ActiveProvider(provider, connections) => { + let mut server_status = { + if let Some(old_status) = status_holder_signal.borrow().as_ref() { + (**old_status).clone() + } else { + StatusCheck::default() + } + }; + if let Some(treemap) = server_status.active_provider_connections.as_mut() { + if connections == 0 { + treemap.remove(&provider); + } else { + treemap.insert(provider, connections); + } + } else if connections > 0 { + let mut treemap = BTreeMap::new(); + treemap.insert(provider, connections); + server_status.active_provider_connections = Some(treemap); + } + let new_status = Rc::new(server_status); + *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); + status_signal.set(Some(new_status)); + } + } + }); + + let fetch_status = { + let services_clone = services_ctx.clone(); + move || { + let services_clone = services_clone.clone(); + spawn_local(async move { services_clone.websocket.get_server_status().await; }); + } + }; + + fetch_status(); + let interval = Interval::new(60*1000, move || { fetch_status(); }); + + let services_clone = services_ctx.clone(); + move || { + drop(interval); + services_clone.websocket.unsubscribe(subid); + } + }); + } + status_holder +} diff --git a/webui/src/services/websocket_service.rs b/webui/src/services/websocket_service.rs index 571aa6205..4ef25539a 100644 --- a/webui/src/services/websocket_service.rs +++ b/webui/src/services/websocket_service.rs @@ -82,6 +82,12 @@ impl WebSocketService { match ProtocolMessage::from_bytes(bytes) { Ok(message) => { match message { + ProtocolMessage::Unauthorized => { + error!("Unauthorized"); + }, + ProtocolMessage::Error(err) => { + error!("{}", err); + }, ProtocolMessage::ActiveUserResponse(user_count, connections) => { broadcast(WsMessage::ActiveUser(user_count, connections)); },