From 7827121277be8a4bc5858009146891d0cdac38cb Mon Sep 17 00:00:00 2001 From: Aminu Oluwaseun Joshua Date: Thu, 27 Nov 2025 01:02:32 +0100 Subject: [PATCH 1/4] introduces different algorithms + fix issues with request_route handling server status Signed-off-by: Aminu Oluwaseun Joshua --- src/algorithms/least_connection.rs | 2 ++ src/algorithms/mod.rs | 25 +++++++++++++++++++++ src/algorithms/resource_based.rs | 0 src/algorithms/weighted_least_connection.rs | 0 src/algorithms/weighted_response_time.rs | 0 src/config.rs | 5 +++++ src/main.rs | 9 ++++---- src/middleware/mod.rs | 17 +++++++++----- 8 files changed, 49 insertions(+), 9 deletions(-) create mode 100644 src/algorithms/least_connection.rs create mode 100644 src/algorithms/mod.rs create mode 100644 src/algorithms/resource_based.rs create mode 100644 src/algorithms/weighted_least_connection.rs create mode 100644 src/algorithms/weighted_response_time.rs diff --git a/src/algorithms/least_connection.rs b/src/algorithms/least_connection.rs new file mode 100644 index 0000000..aeceb8f --- /dev/null +++ b/src/algorithms/least_connection.rs @@ -0,0 +1,2 @@ +// Server reports its current load +// Redis collection of current state of servers or use Kafka to collect data diff --git a/src/algorithms/mod.rs b/src/algorithms/mod.rs new file mode 100644 index 0000000..e8ce92c --- /dev/null +++ b/src/algorithms/mod.rs @@ -0,0 +1,25 @@ +mod least_connection; +mod resource_based; +mod weighted_least_connection; +mod weighted_response_time; + +#[derive(Clone, Default)] +pub enum Algorithm { + #[default] + LeastConnection, + ResourceBased, + WeightedLeastConnection, + WeightedResponseTime, +} + +impl From for Algorithm { + fn from(algorithm: String) -> Self { + match algorithm.as_str() { + "least_connection" => Algorithm::LeastConnection, + "resource_based" => Algorithm::ResourceBased, + "weighted_least_connection" => Algorithm::WeightedLeastConnection, + "weighted_response_time" => Algorithm::WeightedResponseTime, + _ => Algorithm::default(), + } + } +} diff --git a/src/algorithms/resource_based.rs b/src/algorithms/resource_based.rs new file mode 100644 index 0000000..e69de29 diff --git a/src/algorithms/weighted_least_connection.rs b/src/algorithms/weighted_least_connection.rs new file mode 100644 index 0000000..e69de29 diff --git a/src/algorithms/weighted_response_time.rs b/src/algorithms/weighted_response_time.rs new file mode 100644 index 0000000..e69de29 diff --git a/src/config.rs b/src/config.rs index 977c0b7..77ec9ed 100644 --- a/src/config.rs +++ b/src/config.rs @@ -1,6 +1,7 @@ use serde::Deserialize; use crate::{ + algorithms::Algorithm, db::{self, RedisClient}, middleware::{Server, ServerClients}, }; @@ -10,6 +11,8 @@ pub struct SystemConfig { pub available_servers: String, // TODO: This should be hosted in redis pub port: u16, pub redis_url: String, + pub algorithm: String, + pub trace_level: String, } impl SystemConfig { @@ -25,6 +28,7 @@ impl SystemConfig { pub struct State { pub available_servers: ServerClients, pub redis_conn: RedisClient, + pub algorithm: Algorithm, } impl State { @@ -43,6 +47,7 @@ impl State { Ok(State { available_servers: ServerClients::new(available_servers), redis_conn, + algorithm: config.algorithm.clone().into(), }) } } diff --git a/src/main.rs b/src/main.rs index a1400f4..8b93ac5 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,6 +1,6 @@ #![deny(clippy::disallowed_methods)] -use std::net::SocketAddr; +use std::{net::SocketAddr, str::FromStr as _}; use axum::{Router, routing::get}; use tokio::task::JoinHandle; @@ -17,6 +17,7 @@ use crate::{ services::server_worker, }; +pub mod algorithms; pub mod config; pub mod db; pub mod error; @@ -26,13 +27,13 @@ mod services; #[tokio::main] async fn main() -> Result<(), Box> { + let config = SystemConfig::from_env()?; + tracing_subscriber::fmt() - .with_max_level(Level::INFO) + .with_max_level(Level::from_str(&config.trace_level)?) .pretty() .init(); - let config = SystemConfig::from_env()?; - let state = State::new(&config).await?; let server = Router::new() diff --git a/src/middleware/mod.rs b/src/middleware/mod.rs index b09a76a..83c5504 100644 --- a/src/middleware/mod.rs +++ b/src/middleware/mod.rs @@ -3,11 +3,12 @@ use axum::{ extract::State, http::Request, middleware::Next, + response::IntoResponse, }; use futures_util::stream::StreamExt; use serde_json::Value; -use crate::{config::State as AppState, error::Error, middleware::server::ApiResponse}; +use crate::{config::State as AppState, error::Error}; mod server; @@ -17,8 +18,12 @@ pub use server::{Server, ServerClients}; pub async fn request_route( State(state): State, req: Request, - _next: Next, -) -> Result { + next: Next, +) -> Result { + if req.uri().path().starts_with("/status") { + return Ok(next.run(req).await); + } + let (parts, body) = req.into_parts(); tracing::info!("New Request Received"); @@ -30,11 +35,13 @@ pub async fn request_route( let route = parts.uri.to_string(); - state + let response = state .available_servers .choiced_server() .handle_request(parts.method, route.trim_start_matches('/'), json_body) - .await + .await?; + + Ok(response.into_response()) } struct BodyBytes(Bytes); From bfe1c117d7b2716032e6b8af04946778eeebaeb4 Mon Sep 17 00:00:00 2001 From: Aminu Oluwaseun Joshua Date: Tue, 9 Dec 2025 20:55:40 +0100 Subject: [PATCH 2/4] algorithms Signed-off-by: Aminu Oluwaseun Joshua --- src/algorithms/least_connection.rs | 6 ++++++ src/algorithms/mod.rs | 19 +++++++++++++++++++ src/algorithms/resource_based.rs | 5 +++++ src/algorithms/weighted_least_connection.rs | 5 +++++ src/algorithms/weighted_response_time.rs | 5 +++++ src/middleware/mod.rs | 3 ++- src/middleware/server.rs | 7 +++---- 7 files changed, 45 insertions(+), 5 deletions(-) diff --git a/src/algorithms/least_connection.rs b/src/algorithms/least_connection.rs index aeceb8f..c14f94a 100644 --- a/src/algorithms/least_connection.rs +++ b/src/algorithms/least_connection.rs @@ -1,2 +1,8 @@ // Server reports its current load // Redis collection of current state of servers or use Kafka to collect data + +use crate::{error::Error, middleware::Server}; + +pub async fn least_connection(_available_servers: &[Server]) -> Result { + todo!() +} diff --git a/src/algorithms/mod.rs b/src/algorithms/mod.rs index e8ce92c..e0d7ae0 100644 --- a/src/algorithms/mod.rs +++ b/src/algorithms/mod.rs @@ -1,3 +1,5 @@ +use crate::{error::Error, middleware::Server}; + mod least_connection; mod resource_based; mod weighted_least_connection; @@ -23,3 +25,20 @@ impl From for Algorithm { } } } + +impl Algorithm { + pub async fn select_server(&self, available_servers: &[Server]) -> Result { + match self { + Algorithm::LeastConnection => { + least_connection::least_connection(available_servers).await + } + Algorithm::ResourceBased => resource_based::resource_based(available_servers).await, + Algorithm::WeightedLeastConnection => { + weighted_least_connection::weighted_least_connection(available_servers).await + } + Algorithm::WeightedResponseTime => { + weighted_response_time::weighted_response_time(available_servers).await + } + } + } +} diff --git a/src/algorithms/resource_based.rs b/src/algorithms/resource_based.rs index e69de29..89521d5 100644 --- a/src/algorithms/resource_based.rs +++ b/src/algorithms/resource_based.rs @@ -0,0 +1,5 @@ +use crate::{error::Error, middleware::Server}; + +pub async fn resource_based(_available_servers: &[Server]) -> Result { + todo!() +} diff --git a/src/algorithms/weighted_least_connection.rs b/src/algorithms/weighted_least_connection.rs index e69de29..b761994 100644 --- a/src/algorithms/weighted_least_connection.rs +++ b/src/algorithms/weighted_least_connection.rs @@ -0,0 +1,5 @@ +use crate::{error::Error, middleware::Server}; + +pub async fn weighted_least_connection(_available_servers: &[Server]) -> Result { + todo!() +} diff --git a/src/algorithms/weighted_response_time.rs b/src/algorithms/weighted_response_time.rs index e69de29..701004a 100644 --- a/src/algorithms/weighted_response_time.rs +++ b/src/algorithms/weighted_response_time.rs @@ -0,0 +1,5 @@ +use crate::{error::Error, middleware::Server}; + +pub async fn weighted_response_time(_available_servers: &[Server]) -> Result { + todo!() +} diff --git a/src/middleware/mod.rs b/src/middleware/mod.rs index 83c5504..1c9da3d 100644 --- a/src/middleware/mod.rs +++ b/src/middleware/mod.rs @@ -37,7 +37,8 @@ pub async fn request_route( let response = state .available_servers - .choiced_server() + .selected_server(state.algorithm) + .await? .handle_request(parts.method, route.trim_start_matches('/'), json_body) .await?; diff --git a/src/middleware/server.rs b/src/middleware/server.rs index d6352bd..31502cb 100644 --- a/src/middleware/server.rs +++ b/src/middleware/server.rs @@ -3,7 +3,7 @@ use std::str::FromStr; use axum::response::{IntoResponse, Response}; use reqwest::{Method, Response as ReqwestResponse, StatusCode, Url}; -use crate::error::Error; +use crate::{algorithms::Algorithm, error::Error}; /// Represents a collection of server clients for load balancing #[derive(Clone, Debug)] @@ -17,9 +17,8 @@ impl ServerClients { } /// Selects a server based on a load balancing algorithm - pub fn choiced_server(&self) -> Server { - // implement algorithm to select server here! - self.available_servers[0].clone() // placeholder for now + pub async fn selected_server(&self, algorithm: Algorithm) -> Result { + algorithm.select_server(&self.available_servers).await } } From 3b8e4f643af29cd5001521e6ba1ef6f76a438a12 Mon Sep 17 00:00:00 2001 From: Aminu Oluwaseun Joshua Date: Thu, 11 Dec 2025 15:47:45 +0100 Subject: [PATCH 3/4] least connection algo supported Signed-off-by: Aminu Oluwaseun Joshua --- Readme.md | 4 +- src/algorithms/least_connection.rs | 9 ++-- src/algorithms/resource_based.rs | 4 +- src/algorithms/weighted_least_connection.rs | 4 +- src/algorithms/weighted_response_time.rs | 4 +- src/middleware/server.rs | 49 ++++++++++++++++++--- 6 files changed, 55 insertions(+), 19 deletions(-) diff --git a/Readme.md b/Readme.md index f69118a..54222c1 100644 --- a/Readme.md +++ b/Readme.md @@ -12,7 +12,7 @@ Prerequisites: Create a `.env` file at the project root (example): ```bash -AVAILABLE_SERVERS=http://localhost:3001,http://localhost:3002 +AVAILABLE_SERVERS=http://localhost:3001$4,http://localhost:3002$8 PORT=8080 ``` @@ -34,7 +34,7 @@ The server binds to the configured PORT. This project reads configuration from environment variables. The important variables are: -`AVAILABLE_SERVERS` — comma-separated list of backend base URLs (e.g. http://host:port). +`AVAILABLE_SERVERS` — comma-separated list of backend base URLs and their weights (e.g. http://host:port$weight). `PORT` — port to bind the load balancer to. diff --git a/src/algorithms/least_connection.rs b/src/algorithms/least_connection.rs index c14f94a..9b34e58 100644 --- a/src/algorithms/least_connection.rs +++ b/src/algorithms/least_connection.rs @@ -1,8 +1,7 @@ -// Server reports its current load -// Redis collection of current state of servers or use Kafka to collect data - use crate::{error::Error, middleware::Server}; -pub async fn least_connection(_available_servers: &[Server]) -> Result { - todo!() +pub async fn least_connection(available_servers: &[Server]) -> Result { + let server = available_servers.iter().min_by_key(|server| server.load()); + + Ok(server.unwrap_or(&available_servers[0]).clone()) } diff --git a/src/algorithms/resource_based.rs b/src/algorithms/resource_based.rs index 89521d5..3c0d73c 100644 --- a/src/algorithms/resource_based.rs +++ b/src/algorithms/resource_based.rs @@ -1,5 +1,5 @@ use crate::{error::Error, middleware::Server}; -pub async fn resource_based(_available_servers: &[Server]) -> Result { - todo!() +pub async fn resource_based(available_servers: &[Server]) -> Result { + Ok(available_servers[0].clone()) } diff --git a/src/algorithms/weighted_least_connection.rs b/src/algorithms/weighted_least_connection.rs index b761994..97d53c5 100644 --- a/src/algorithms/weighted_least_connection.rs +++ b/src/algorithms/weighted_least_connection.rs @@ -1,5 +1,5 @@ use crate::{error::Error, middleware::Server}; -pub async fn weighted_least_connection(_available_servers: &[Server]) -> Result { - todo!() +pub async fn weighted_least_connection(available_servers: &[Server]) -> Result { + Ok(available_servers[0].clone()) } diff --git a/src/algorithms/weighted_response_time.rs b/src/algorithms/weighted_response_time.rs index 701004a..ef71c3f 100644 --- a/src/algorithms/weighted_response_time.rs +++ b/src/algorithms/weighted_response_time.rs @@ -1,5 +1,5 @@ use crate::{error::Error, middleware::Server}; -pub async fn weighted_response_time(_available_servers: &[Server]) -> Result { - todo!() +pub async fn weighted_response_time(available_servers: &[Server]) -> Result { + Ok(available_servers[0].clone()) } diff --git a/src/middleware/server.rs b/src/middleware/server.rs index 31502cb..7daa8e1 100644 --- a/src/middleware/server.rs +++ b/src/middleware/server.rs @@ -1,4 +1,7 @@ -use std::str::FromStr; +use std::{ + str::FromStr, + sync::{Arc, atomic::AtomicU32}, +}; use axum::response::{IntoResponse, Response}; use reqwest::{Method, Response as ReqwestResponse, StatusCode, Url}; @@ -18,7 +21,12 @@ impl ServerClients { /// Selects a server based on a load balancing algorithm pub async fn selected_server(&self, algorithm: Algorithm) -> Result { - algorithm.select_server(&self.available_servers).await + algorithm + .select_server(&self.available_servers) + .await + .inspect(|s| { + s.load.fetch_add(1, std::sync::atomic::Ordering::Acquire); + }) } } @@ -26,16 +34,37 @@ impl ServerClients { pub struct Server { pub url: Url, pub client: reqwest::Client, + load: Arc, + weight: u32, } impl Server { - pub fn new(url: &str) -> anyhow::Result { + pub fn new(url_and_weight: &str) -> anyhow::Result { + let (url, weight) = url_and_weight + .split_once('$') + .ok_or_else(|| anyhow::anyhow!("Invalid server format, expected 'url$weight'"))?; + let weight = weight + .parse::() + .map_err(|_| anyhow::anyhow!("Invalid weight, expected a positive integer"))?; + Ok(Self { url: Url::from_str(url)?, client: Default::default(), + load: Arc::new(AtomicU32::new(0)), + weight, }) } + /// Returns the current load of the server + pub fn load(&self) -> u32 { + self.load.load(std::sync::atomic::Ordering::Relaxed) + } + + /// Returns the weight of the server + pub fn weight(&self) -> u32 { + self.weight + } + /// Handles incoming requests and forwards them to the server pub async fn handle_request( &self, @@ -43,10 +72,18 @@ impl Server { route: &str, body: Option, ) -> Result { + // TODO: What if the request fails is the load count reduced? match method { - Method::GET => self.get_request(route, body).await, - Method::POST => self.post_request(route, body).await, - _ => Err(Error::MethodNotAllowed), + Method::GET => self.get_request(route, body).await.inspect(|_| { + self.load.fetch_sub(1, std::sync::atomic::Ordering::Release); + }), + Method::POST => self.post_request(route, body).await.inspect(|_| { + self.load.fetch_sub(1, std::sync::atomic::Ordering::Release); + }), + _ => { + self.load.fetch_sub(1, std::sync::atomic::Ordering::Release); + Err(Error::MethodNotAllowed) + } } } From 6f22ddd69989c720a7027e888a77adc5c0afb85b Mon Sep 17 00:00:00 2001 From: Aminu Oluwaseun Joshua Date: Fri, 12 Dec 2025 13:01:07 +0100 Subject: [PATCH 4/4] more algorithms added Signed-off-by: Aminu Oluwaseun Joshua --- src/algorithms/least_connection.rs | 7 ++- src/algorithms/resource_based.rs | 1 + src/algorithms/weighted_least_connection.rs | 7 ++- src/algorithms/weighted_response_time.rs | 7 ++- src/config.rs | 2 +- src/error.rs | 5 +- src/main.rs | 25 ++++++++-- src/middleware/mod.rs | 12 ++++- src/middleware/server.rs | 50 +++++++++++++++---- src/services/latency_tracker_worker.rs | 30 +++++++++++ src/services/mod.rs | 6 ++- ...ound_worker.rs => server_status_worker.rs} | 2 +- 12 files changed, 129 insertions(+), 25 deletions(-) create mode 100644 src/services/latency_tracker_worker.rs rename src/services/{background_worker.rs => server_status_worker.rs} (92%) diff --git a/src/algorithms/least_connection.rs b/src/algorithms/least_connection.rs index 9b34e58..5b9e53b 100644 --- a/src/algorithms/least_connection.rs +++ b/src/algorithms/least_connection.rs @@ -1,7 +1,10 @@ use crate::{error::Error, middleware::Server}; pub async fn least_connection(available_servers: &[Server]) -> Result { - let server = available_servers.iter().min_by_key(|server| server.load()); + let server = available_servers + .iter() + .min_by_key(|server| server.load()) + .ok_or_else(|| Error::NoServerAvailable)?; - Ok(server.unwrap_or(&available_servers[0]).clone()) + Ok(server.clone()) } diff --git a/src/algorithms/resource_based.rs b/src/algorithms/resource_based.rs index 3c0d73c..64e47ac 100644 --- a/src/algorithms/resource_based.rs +++ b/src/algorithms/resource_based.rs @@ -1,5 +1,6 @@ use crate::{error::Error, middleware::Server}; +// TODO: Implement resource-based load balancing algorithm pub async fn resource_based(available_servers: &[Server]) -> Result { Ok(available_servers[0].clone()) } diff --git a/src/algorithms/weighted_least_connection.rs b/src/algorithms/weighted_least_connection.rs index 97d53c5..84f2f87 100644 --- a/src/algorithms/weighted_least_connection.rs +++ b/src/algorithms/weighted_least_connection.rs @@ -1,5 +1,10 @@ use crate::{error::Error, middleware::Server}; pub async fn weighted_least_connection(available_servers: &[Server]) -> Result { - Ok(available_servers[0].clone()) + let server = available_servers + .iter() + .min_by_key(|server| server.load() / server.weight()) + .ok_or_else(|| Error::NoServerAvailable)?; + + Ok(server.clone()) } diff --git a/src/algorithms/weighted_response_time.rs b/src/algorithms/weighted_response_time.rs index ef71c3f..0d289b2 100644 --- a/src/algorithms/weighted_response_time.rs +++ b/src/algorithms/weighted_response_time.rs @@ -1,5 +1,10 @@ use crate::{error::Error, middleware::Server}; pub async fn weighted_response_time(available_servers: &[Server]) -> Result { - Ok(available_servers[0].clone()) + let server = available_servers + .iter() + .min_by_key(|server| server.mean_latency()) + .ok_or_else(|| Error::NoServerAvailable)?; + + Ok(server.clone()) } diff --git a/src/config.rs b/src/config.rs index 77ec9ed..8c7dbcf 100644 --- a/src/config.rs +++ b/src/config.rs @@ -6,7 +6,7 @@ use crate::{ middleware::{Server, ServerClients}, }; -#[derive(Deserialize, Debug)] +#[derive(Deserialize)] pub struct SystemConfig { pub available_servers: String, // TODO: This should be hosted in redis pub port: u16, diff --git a/src/error.rs b/src/error.rs index dc84168..4b53696 100644 --- a/src/error.rs +++ b/src/error.rs @@ -20,6 +20,8 @@ pub enum Error { InvalidUrl, #[error("Invalid Response")] InvalidResponse, + #[error("No Server Available")] + NoServerAvailable, } impl IntoResponse for Error { @@ -29,7 +31,8 @@ impl IntoResponse for Error { Error::InternalServerError | Error::Other(_) | Error::InvalidResponse - | Error::RedisError(_) => { + | Error::RedisError(_) + | Error::NoServerAvailable => { (StatusCode::INTERNAL_SERVER_ERROR, "Internal Server Error").into_response() } Error::Unauthorized => (StatusCode::UNAUTHORIZED, self).into_response(), diff --git a/src/main.rs b/src/main.rs index 8b93ac5..1f1e5f6 100644 --- a/src/main.rs +++ b/src/main.rs @@ -14,7 +14,7 @@ use crate::{ config::{State, SystemConfig}, middleware::request_route, servers::health::status, - services::server_worker, + services::server_status_worker, }; pub mod algorithms; @@ -60,12 +60,21 @@ async fn main() -> Result<(), Box> { let main = tokio::spawn(async move { axum::serve(listener, server).await }); - let background_worker = tokio::spawn(async move { - let _: () = server_worker(state.clone().available_servers.available_servers).await; + let server_status_background_worker = tokio::spawn(async move { + let _: () = server_status_worker(state.clone().available_servers.available_servers).await; Ok(()) }); - let app = App::new(main, background_worker); + // let latency_tracker_background_worker = tokio::spawn(async move { + // let _: () = latency_tracker_worker(state.clone().available_servers.available_servers).await; + // Ok(()) + // }); + + let app = App::new( + main, + server_status_background_worker, + // latency_tracker_background_worker, + ); app.start().await } @@ -76,13 +85,19 @@ type JoinHandleWrapper = JoinHandle>; struct App { main: JoinHandleWrapper, background_worker: JoinHandleWrapper, + // latency_tracker_background_worker: JoinHandleWrapper, } impl App { - fn new(main: JoinHandleWrapper, background_worker: JoinHandleWrapper) -> Self { + fn new( + main: JoinHandleWrapper, + background_worker: JoinHandleWrapper, + // latency_tracker_background_worker: JoinHandleWrapper, + ) -> Self { Self { main, background_worker, + // latency_tracker_background_worker, } } diff --git a/src/middleware/mod.rs b/src/middleware/mod.rs index 1c9da3d..3e389af 100644 --- a/src/middleware/mod.rs +++ b/src/middleware/mod.rs @@ -35,13 +35,21 @@ pub async fn request_route( let route = parts.uri.to_string(); - let response = state + let start_time = std::time::Instant::now(); + + let mut server: Server = state .available_servers .selected_server(state.algorithm) - .await? + .await?; + + let response = server .handle_request(parts.method, route.trim_start_matches('/'), json_body) .await?; + let latency = start_time.elapsed().as_millis(); + + server.update_latencies(latency); + Ok(response.into_response()) } diff --git a/src/middleware/server.rs b/src/middleware/server.rs index 7daa8e1..8ec3137 100644 --- a/src/middleware/server.rs +++ b/src/middleware/server.rs @@ -1,6 +1,9 @@ use std::{ str::FromStr, - sync::{Arc, atomic::AtomicU32}, + sync::{ + Arc, + atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering}, + }, }; use axum::response::{IntoResponse, Response}; @@ -9,7 +12,7 @@ use reqwest::{Method, Response as ReqwestResponse, StatusCode, Url}; use crate::{algorithms::Algorithm, error::Error}; /// Represents a collection of server clients for load balancing -#[derive(Clone, Debug)] +#[derive(Clone)] pub struct ServerClients { pub available_servers: Vec, } @@ -25,17 +28,20 @@ impl ServerClients { .select_server(&self.available_servers) .await .inspect(|s| { - s.load.fetch_add(1, std::sync::atomic::Ordering::Acquire); + s.load.fetch_add(1, Ordering::Acquire); }) } } -#[derive(Clone, Debug)] +#[derive(Clone)] pub struct Server { pub url: Url, pub client: reqwest::Client, load: Arc, weight: u32, + pub mean_latency: Arc, + pub latencies: Vec, + latencies_updated: Arc, } impl Server { @@ -43,21 +49,27 @@ impl Server { let (url, weight) = url_and_weight .split_once('$') .ok_or_else(|| anyhow::anyhow!("Invalid server format, expected 'url$weight'"))?; + let weight = weight .parse::() .map_err(|_| anyhow::anyhow!("Invalid weight, expected a positive integer"))?; + let url = Url::from_str(url)?; + Ok(Self { - url: Url::from_str(url)?, + url, client: Default::default(), load: Arc::new(AtomicU32::new(0)), weight, + mean_latency: Arc::new(AtomicU64::new(0)), + latencies: Vec::new(), + latencies_updated: Arc::new(AtomicBool::new(false)), }) } /// Returns the current load of the server pub fn load(&self) -> u32 { - self.load.load(std::sync::atomic::Ordering::Relaxed) + self.load.load(Ordering::Relaxed) } /// Returns the weight of the server @@ -65,6 +77,26 @@ impl Server { self.weight } + pub fn update_latencies(&mut self, latency: u128) { + if self.latencies.len() >= 20 { + // TODO: make it customisable + self.latencies.remove(0); + } + self.latencies.push(latency); + } + + pub fn latency_updated(&self) -> bool { + self.latencies_updated.load(Ordering::Relaxed) + } + + pub fn latency_update_status(&self, b: bool) { + self.latencies_updated.store(b, Ordering::Relaxed) + } + + pub fn mean_latency(&self) -> u64 { + self.mean_latency.load(Ordering::Relaxed) + } + /// Handles incoming requests and forwards them to the server pub async fn handle_request( &self, @@ -75,13 +107,13 @@ impl Server { // TODO: What if the request fails is the load count reduced? match method { Method::GET => self.get_request(route, body).await.inspect(|_| { - self.load.fetch_sub(1, std::sync::atomic::Ordering::Release); + self.load.fetch_sub(1, Ordering::Release); }), Method::POST => self.post_request(route, body).await.inspect(|_| { - self.load.fetch_sub(1, std::sync::atomic::Ordering::Release); + self.load.fetch_sub(1, Ordering::Release); }), _ => { - self.load.fetch_sub(1, std::sync::atomic::Ordering::Release); + self.load.fetch_sub(1, Ordering::Release); Err(Error::MethodNotAllowed) } } diff --git a/src/services/latency_tracker_worker.rs b/src/services/latency_tracker_worker.rs new file mode 100644 index 0000000..538e1fe --- /dev/null +++ b/src/services/latency_tracker_worker.rs @@ -0,0 +1,30 @@ +use std::sync::atomic::Ordering; + +use crate::middleware::Server; + +pub async fn _latency_tracker_worker(available_servers: Vec) { + loop { + _check(&available_servers); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + } +} + +fn _check(servers: &[Server]) { + for server in servers { + if server.latency_updated() { + server.latency_update_status(false); + let mean_latency = _mean_latency(&server.latencies); + server + .mean_latency + .store(mean_latency as u64, Ordering::Relaxed); + } + } +} + +fn _mean_latency(latencies: &[u128]) -> u128 { + if latencies.is_empty() { + 0 + } else { + latencies.iter().sum::() / latencies.len() as u128 + } +} diff --git a/src/services/mod.rs b/src/services/mod.rs index 7a2628e..533afac 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,3 +1,5 @@ -mod background_worker; +mod latency_tracker_worker; +mod server_status_worker; -pub use background_worker::server_worker; +// pub use latency_tracker_worker::latency_tracker_worker; +pub use server_status_worker::server_status_worker; diff --git a/src/services/background_worker.rs b/src/services/server_status_worker.rs similarity index 92% rename from src/services/background_worker.rs rename to src/services/server_status_worker.rs index 6a7156c..d264199 100644 --- a/src/services/background_worker.rs +++ b/src/services/server_status_worker.rs @@ -1,7 +1,7 @@ use crate::middleware::Server; /// Background worker that periodically checks the status of available servers -pub async fn server_worker(available_servers: Vec) { +pub async fn server_status_worker(available_servers: Vec) { loop { if let Err(failing_servers) = server_status(available_servers.clone()).await { // TODO: remove them from the list of available servers