diff --git a/Cargo.lock b/Cargo.lock index 368138d..fe6169f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1843,6 +1843,41 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "rustls" +version = "0.23.40" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef86cd5876211988985292b91c96a8f2d298df24e75989a43a3c73f2d4d8168b" +dependencies = [ + "log", + "once_cell", + "ring", + "rustls-pki-types", + "rustls-webpki", + "subtle", + "zeroize", +] + +[[package]] +name = "rustls-pki-types" +version = "1.14.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "30a7197ae7eb376e574fe940d068c30fe0462554a3ddbe4eca7838e049c937a9" +dependencies = [ + "zeroize", +] + +[[package]] +name = "rustls-webpki" +version = "0.103.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +dependencies = [ + "ring", + "rustls-pki-types", + "untrusted", +] + [[package]] name = "rustversion" version = "1.0.22" @@ -2256,6 +2291,16 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "tokio-rustls" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" +dependencies = [ + "rustls", + "tokio", +] + [[package]] name = "tokio-stream" version = "0.1.18" @@ -2343,6 +2388,7 @@ dependencies = [ "socket2", "sync_wrapper", "tokio", + "tokio-rustls", "tokio-stream", "tower", "tower-layer", @@ -2974,6 +3020,12 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "zeroize" +version = "1.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b97154e67e32c85465826e8bcc1c59429aaaf107c1e4a9e53c8d8ccd5eff88d0" + [[package]] name = "zmij" version = "1.0.21" diff --git a/Cargo.toml b/Cargo.toml index 453b92b..5bc9c09 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -15,7 +15,7 @@ rkyv = "0.8" bincode = "1" # --- gRPC --- -tonic = "0.14" +tonic = { version = "0.14", features = ["tls-ring"] } tonic-prost = "0.14" prost = "0.14" diff --git a/src/gateway/auth_proxy.rs b/src/gateway/auth_proxy.rs new file mode 100644 index 0000000..913270b --- /dev/null +++ b/src/gateway/auth_proxy.rs @@ -0,0 +1,162 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_auth_server::AetherAuth; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct AuthProxy { + pool: Arc>, +} +impl AuthProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result, Status> + { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.auth.clone()) + } +} + +macro_rules! unary_auth { + ($self:ident, $request:ident, $method:ident, $getter:ident) => {{ + let metadata = $request.metadata().clone(); + let req = $request.into_inner(); + let (timeout, _addr, mut client) = { + let p = $self.pool.read().await; + let c = p + .$getter() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.$method(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = $self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.$method(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + }}; +} + +#[tonic::async_trait] +impl AetherAuth for AuthProxy { + async fn authenticate( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, authenticate, get_any_auth) + } + async fn auth_enable( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, auth_enable, get_leader_auth) + } + async fn auth_disable( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, auth_disable, get_leader_auth) + } + async fn auth_status( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, auth_status, get_any_auth) + } + async fn user_add( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_add, get_leader_auth) + } + async fn user_delete( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_delete, get_leader_auth) + } + async fn user_get( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_get, get_any_auth) + } + async fn user_list( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_list, get_any_auth) + } + async fn user_change_password( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_change_password, get_leader_auth) + } + async fn user_grant_role( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_grant_role, get_leader_auth) + } + async fn user_revoke_role( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, user_revoke_role, get_leader_auth) + } + async fn role_add( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, role_add, get_leader_auth) + } + async fn role_delete( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, role_delete, get_leader_auth) + } + async fn role_get( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, role_get, get_any_auth) + } + async fn role_list( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, role_list, get_any_auth) + } + async fn role_grant_permission( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, role_grant_permission, get_leader_auth) + } + async fn role_revoke_permission( + &self, + request: Request, + ) -> Result, Status> { + unary_auth!(self, request, role_revoke_permission, get_leader_auth) + } +} diff --git a/src/gateway/barrier_proxy.rs b/src/gateway/barrier_proxy.rs new file mode 100644 index 0000000..0859e0f --- /dev/null +++ b/src/gateway/barrier_proxy.rs @@ -0,0 +1,130 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_barrier_server::AetherBarrier; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct BarrierProxy { + pool: Arc>, +} +impl BarrierProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_barrier_client::AetherBarrierClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.barrier.clone()) + } +} + +#[tonic::async_trait] +impl AetherBarrier for BarrierProxy { + async fn create( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_barrier() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.create(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.create(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn release( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_barrier() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.release(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.release(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn query( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_barrier() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.query(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.query(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/cluster_proxy.rs b/src/gateway/cluster_proxy.rs new file mode 100644 index 0000000..a661df3 --- /dev/null +++ b/src/gateway/cluster_proxy.rs @@ -0,0 +1,161 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_cluster_server::AetherCluster; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct ClusterProxy { + pool: Arc>, +} +impl ClusterProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_cluster_client::AetherClusterClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.cluster.clone()) + } +} + +#[tonic::async_trait] +impl AetherCluster for ClusterProxy { + async fn member_list( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_cluster() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.member_list(forward_request(&metadata, req))) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.member_list(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn member_add( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_cluster() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.member_add(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.member_add(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn member_remove( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_cluster() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.member_remove(forward_request(&metadata, req)), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.member_remove(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn member_promote( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_cluster() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.member_promote(forward_request(&metadata, req)), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.member_promote(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/election_proxy.rs b/src/gateway/election_proxy.rs new file mode 100644 index 0000000..2ba5c09 --- /dev/null +++ b/src/gateway/election_proxy.rs @@ -0,0 +1,159 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_election_server::AetherElection; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct ElectionProxy { + pool: Arc>, +} +impl ElectionProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_election_client::AetherElectionClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.election.clone()) + } +} + +#[tonic::async_trait] +impl AetherElection for ElectionProxy { + async fn campaign( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_election() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.campaign(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.campaign(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn leader( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_election() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.leader(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.leader(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn resign( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_election() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.resign(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.resign(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + type ObserveStream = tonic::Streaming; + + async fn observe( + &self, + request: Request, + ) -> Result, Status> { + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_election() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + let req = request.into_inner(); + match tokio::time::timeout(timeout, client.observe(Request::new(req.clone()))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + Err(Status::unavailable(format!( + "leader redirect to {leader}, please retry" + ))) + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/kv_proxy.rs b/src/gateway/kv_proxy.rs new file mode 100644 index 0000000..6afe625 --- /dev/null +++ b/src/gateway/kv_proxy.rs @@ -0,0 +1,189 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_kv_server::AetherKv; +use crate::proto::{ + DeleteRequest, DeleteResponse, GetRequest, GetResponse, PutRequest, PutResponse, RangeRequest, + RangeResponse, TxnRequest, TxnResponse, +}; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct KvProxy { + pool: Arc>, +} +impl KvProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result, Status> + { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.kv.clone()) + } +} + +#[tonic::async_trait] +impl AetherKv for KvProxy { + async fn put(&self, request: Request) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_kv() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.put(forward_request(&metadata, req.clone()))) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.put(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn get(&self, request: Request) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let serializable = req.serializable; + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = if serializable { + p.get_any_kv() + } else { + p.get_leader_kv() + } + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.get(forward_request(&metadata, req.clone()))) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.get(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn delete( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_kv() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.delete(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.delete(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn range( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let serializable = req.serializable; + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = if serializable { + p.get_any_kv() + } else { + p.get_leader_kv() + } + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.range(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.range(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn txn(&self, request: Request) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_kv() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.txn(forward_request(&metadata, req.clone()))) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.txn(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/lease_proxy.rs b/src/gateway/lease_proxy.rs new file mode 100644 index 0000000..0496867 --- /dev/null +++ b/src/gateway/lease_proxy.rs @@ -0,0 +1,202 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_lease_server::AetherLease; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct LeaseProxy { + pool: Arc>, +} +impl LeaseProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_lease_client::AetherLeaseClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.lease.clone()) + } +} + +#[tonic::async_trait] +impl AetherLease for LeaseProxy { + async fn lease_grant( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_lease() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.lease_grant(forward_request(&metadata, req))) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.lease_grant(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn lease_revoke( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_lease() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.lease_revoke(forward_request(&metadata, req)), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.lease_revoke(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + type LeaseKeepAliveStream = tonic::Streaming; + + async fn lease_keep_alive( + &self, + request: Request>, + ) -> Result, Status> { + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_lease() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + let mut client_stream = request.into_inner(); + let (tx, rx) = tokio::sync::mpsc::channel(128); + tokio::spawn(async move { + while let Ok(Some(msg)) = client_stream.message().await { + if tx.send(msg).await.is_err() { + break; + } + } + }); + let stream = tokio_stream::wrappers::ReceiverStream::new(rx); + match tokio::time::timeout(timeout, client.lease_keep_alive(Request::new(stream))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + Err(Status::unavailable(format!( + "leader redirect to {leader}, please retry" + ))) + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn lease_time_to_live( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_lease() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.lease_time_to_live(forward_request(&metadata, req)), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout( + timeout, + c.lease_time_to_live(forward_request(&metadata, req)), + ) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn lease_leases( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_lease() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.lease_leases(forward_request(&metadata, req)), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.lease_leases(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/lock_proxy.rs b/src/gateway/lock_proxy.rs new file mode 100644 index 0000000..93fbfa6 --- /dev/null +++ b/src/gateway/lock_proxy.rs @@ -0,0 +1,125 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_lock_server::AetherLock; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct LockProxy { + pool: Arc>, +} +impl LockProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result, Status> + { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.lock.clone()) + } +} + +#[tonic::async_trait] +impl AetherLock for LockProxy { + async fn lock(&self, request: Request) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_lock() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.lock(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.lock(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn unlock( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_lock() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.unlock(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.unlock(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn lock_query( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_lock() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.lock_query(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.lock_query(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/maintenance_proxy.rs b/src/gateway/maintenance_proxy.rs new file mode 100644 index 0000000..a3090ad --- /dev/null +++ b/src/gateway/maintenance_proxy.rs @@ -0,0 +1,115 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_maintenance_server::AetherMaintenance; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct MaintenanceProxy { + pool: Arc>, +} +impl MaintenanceProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_maintenance_client::AetherMaintenanceClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.maintenance.clone()) + } +} + +#[tonic::async_trait] +impl AetherMaintenance for MaintenanceProxy { + async fn defrag( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_maintenance() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.defrag(forward_request(&metadata, req))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.defrag(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn alarm( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_maintenance() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.alarm(forward_request(&metadata, req))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.alarm(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn status( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_maintenance() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.status(forward_request(&metadata, req))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.status(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs new file mode 100644 index 0000000..1c49df6 --- /dev/null +++ b/src/gateway/mod.rs @@ -0,0 +1,694 @@ +mod auth_proxy; +mod barrier_proxy; +mod cluster_proxy; +mod election_proxy; +mod kv_proxy; +mod lease_proxy; +mod lock_proxy; +mod maintenance_proxy; +mod queue_proxy; +mod session_proxy; +mod watch_proxy; + +use std::collections::HashMap; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; + +use tokio::sync::RwLock; +use tonic::Request; +use tonic::transport::Channel; +use tracing::{info, warn}; + +use crate::proto::MemberListRequest; +use crate::proto::aether_auth_client::AetherAuthClient; +use crate::proto::aether_barrier_client::AetherBarrierClient; +use crate::proto::aether_cluster_client::AetherClusterClient; +use crate::proto::aether_election_client::AetherElectionClient; +use crate::proto::aether_kv_client::AetherKvClient; +use crate::proto::aether_lease_client::AetherLeaseClient; +use crate::proto::aether_lock_client::AetherLockClient; +use crate::proto::aether_maintenance_client::AetherMaintenanceClient; +use crate::proto::aether_queue_client::AetherQueueClient; +use crate::proto::aether_session_client::AetherSessionClient; +use crate::proto::aether_watch_client::AetherWatchClient; + +pub use self::auth_proxy::AuthProxy; +pub use self::barrier_proxy::BarrierProxy; +pub use self::cluster_proxy::ClusterProxy; +pub use self::election_proxy::ElectionProxy; +pub use self::kv_proxy::KvProxy; +pub use self::lease_proxy::LeaseProxy; +pub use self::lock_proxy::LockProxy; +pub use self::maintenance_proxy::MaintenanceProxy; +pub use self::queue_proxy::QueueProxy; +pub use self::session_proxy::SessionProxy; +pub use self::watch_proxy::WatchProxy; + +const DEFAULT_REQUEST_TIMEOUT_MS: u64 = 5000; + +#[derive(Debug, Clone)] +pub struct GatewayConfig { + pub listen_addr: String, + pub backend_addrs: Vec, + pub request_timeout_ms: u64, + pub tls: Option, + pub health_addr: String, +} + +#[derive(Debug, Clone)] +pub struct TlsConfig { + pub ca_cert: String, + pub client_cert: String, + pub client_key: String, +} + +impl GatewayConfig { + pub fn new(listen_addr: String, backend_addrs: Vec) -> Self { + Self { + listen_addr, + backend_addrs, + request_timeout_ms: DEFAULT_REQUEST_TIMEOUT_MS, + tls: None, + health_addr: "127.0.0.1:9091".to_string(), + } + } +} + +/// Create a new `Request` preserving the `authorization` metadata. +pub(super) fn forward_request(metadata: &tonic::metadata::MetadataMap, inner: T) -> Request { + let mut req = Request::new(inner); + if let Some(auth) = metadata.get("authorization") { + req.metadata_mut().insert("authorization", auth.clone()); + } + req +} + +#[derive(Debug, Clone)] +pub(super) struct BackendConnection { + kv: AetherKvClient, + cluster: AetherClusterClient, + maintenance: AetherMaintenanceClient, + watch: AetherWatchClient, + lease: AetherLeaseClient, + auth: AetherAuthClient, + lock: AetherLockClient, + election: AetherElectionClient, + barrier: AetherBarrierClient, + queue: AetherQueueClient, + session: AetherSessionClient, +} + +pub struct BackendPool { + backends: Vec<(String, BackendConnection)>, + addr_to_idx: HashMap, + leader_idx: Option, + rr_counter: AtomicUsize, + request_timeout: Duration, + tls: Option, +} + +impl BackendPool { + async fn connect( + addrs: &[String], + request_timeout_ms: u64, + tls: &Option, + ) -> Result { + let mut backends = Vec::new(); + let mut addr_to_idx = HashMap::new(); + let scheme = if tls.is_some() { "https" } else { "http" }; + + for addr in addrs { + let uri = format!("{scheme}://{addr}"); + match connect_backend(&uri, tls).await { + Ok(conn) => { + info!(addr = %addr, "connected to backend"); + addr_to_idx.insert(addr.clone(), backends.len()); + backends.push((addr.clone(), conn)); + } + Err(e) => { + warn!(addr = %addr, error = %e, "failed to connect to backend"); + } + } + } + + if backends.is_empty() { + return Err(tonic::Status::unavailable( + "failed to connect to any backend node", + )); + } + + let mut pool = Self { + backends, + addr_to_idx, + leader_idx: None, + rr_counter: AtomicUsize::new(0), + tls: tls.clone(), + request_timeout: Duration::from_millis(request_timeout_ms), + }; + + pool.discover_leader().await; + Ok(pool) + } + + async fn discover_leader(&mut self) { + for (i, (addr, conn)) in self.backends.iter_mut().enumerate() { + let mut client = conn.cluster.clone(); + match tokio::time::timeout( + self.request_timeout, + client.member_list(Request::new(MemberListRequest {})), + ) + .await + { + Ok(Ok(resp)) => { + let members = &resp.get_ref().members; + info!(addr = %addr, member_count = members.len(), "discovered cluster members, setting as initial leader candidate"); + self.leader_idx = Some(i); + return; + } + Ok(Err(e)) => { + warn!(addr = %addr, error = %e, "member_list failed"); + } + Err(_) => { + warn!(addr = %addr, "member_list timed out"); + } + } + } + warn!("no backend responded to member_list, leader unknown"); + } + + pub fn update_leader(&mut self, addr: &str) { + if let Some(&idx) = self.addr_to_idx.get(addr) + && self.leader_idx != Some(idx) + { + info!(addr = %addr, "leader updated"); + self.leader_idx = Some(idx); + } + } + + pub(super) fn add_redirect_connection(&mut self, addr: String, conn: BackendConnection) { + let idx = self.backends.len(); + info!(addr = %addr, idx = idx, "caching redirect connection"); + self.addr_to_idx.insert(addr.clone(), idx); + self.leader_idx = Some(idx); + self.backends.push((addr, conn)); + } + + fn timeout(&self) -> Duration { + self.request_timeout + } + + pub(super) fn tls(&self) -> &Option { + &self.tls + } + + pub fn backends_count(&self) -> usize { + self.backends.len() + } + + // --- KV --- + pub fn get_leader_kv(&self) -> Option<(String, AetherKvClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.kv.clone())); + } + self.get_any_kv_inner() + } + pub fn get_any_kv(&self) -> Option<(String, AetherKvClient)> { + self.get_any_kv_inner() + } + fn get_any_kv_inner(&self) -> Option<(String, AetherKvClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.kv.clone())) + } + + // --- Cluster --- + pub fn get_leader_cluster(&self) -> Option<(String, AetherClusterClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.cluster.clone())); + } + self.get_any_cluster_inner() + } + pub fn get_any_cluster(&self) -> Option<(String, AetherClusterClient)> { + self.get_any_cluster_inner() + } + fn get_any_cluster_inner(&self) -> Option<(String, AetherClusterClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.cluster.clone())) + } + + // --- Maintenance --- + pub fn get_leader_maintenance(&self) -> Option<(String, AetherMaintenanceClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.maintenance.clone())); + } + self.get_any_maintenance_inner() + } + pub fn get_any_maintenance(&self) -> Option<(String, AetherMaintenanceClient)> { + self.get_any_maintenance_inner() + } + fn get_any_maintenance_inner(&self) -> Option<(String, AetherMaintenanceClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.maintenance.clone())) + } + + // --- Watch --- + pub fn get_leader_watch(&self) -> Option<(String, AetherWatchClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.watch.clone())); + } + self.get_any_watch_inner() + } + pub fn get_any_watch(&self) -> Option<(String, AetherWatchClient)> { + self.get_any_watch_inner() + } + fn get_any_watch_inner(&self) -> Option<(String, AetherWatchClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.watch.clone())) + } + + // --- Lease --- + pub fn get_leader_lease(&self) -> Option<(String, AetherLeaseClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.lease.clone())); + } + self.get_any_lease_inner() + } + pub fn get_any_lease(&self) -> Option<(String, AetherLeaseClient)> { + self.get_any_lease_inner() + } + fn get_any_lease_inner(&self) -> Option<(String, AetherLeaseClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.lease.clone())) + } + + // --- Auth --- + pub fn get_leader_auth(&self) -> Option<(String, AetherAuthClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.auth.clone())); + } + self.get_any_auth_inner() + } + pub fn get_any_auth(&self) -> Option<(String, AetherAuthClient)> { + self.get_any_auth_inner() + } + fn get_any_auth_inner(&self) -> Option<(String, AetherAuthClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.auth.clone())) + } + + // --- Lock --- + pub fn get_leader_lock(&self) -> Option<(String, AetherLockClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.lock.clone())); + } + self.get_any_lock_inner() + } + pub fn get_any_lock(&self) -> Option<(String, AetherLockClient)> { + self.get_any_lock_inner() + } + fn get_any_lock_inner(&self) -> Option<(String, AetherLockClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.lock.clone())) + } + + // --- Election --- + pub fn get_leader_election(&self) -> Option<(String, AetherElectionClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.election.clone())); + } + self.get_any_election_inner() + } + pub fn get_any_election(&self) -> Option<(String, AetherElectionClient)> { + self.get_any_election_inner() + } + fn get_any_election_inner(&self) -> Option<(String, AetherElectionClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.election.clone())) + } + + // --- Barrier --- + pub fn get_leader_barrier(&self) -> Option<(String, AetherBarrierClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.barrier.clone())); + } + self.get_any_barrier_inner() + } + pub fn get_any_barrier(&self) -> Option<(String, AetherBarrierClient)> { + self.get_any_barrier_inner() + } + fn get_any_barrier_inner(&self) -> Option<(String, AetherBarrierClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.barrier.clone())) + } + + // --- Queue --- + pub fn get_leader_queue(&self) -> Option<(String, AetherQueueClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.queue.clone())); + } + self.get_any_queue_inner() + } + pub fn get_any_queue(&self) -> Option<(String, AetherQueueClient)> { + self.get_any_queue_inner() + } + fn get_any_queue_inner(&self) -> Option<(String, AetherQueueClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.queue.clone())) + } + + // --- Session --- + pub fn get_leader_session(&self) -> Option<(String, AetherSessionClient)> { + if let Some(idx) = self.leader_idx + && let Some((addr, conn)) = self.backends.get(idx) + { + return Some((addr.clone(), conn.session.clone())); + } + self.get_any_session_inner() + } + pub fn get_any_session(&self) -> Option<(String, AetherSessionClient)> { + self.get_any_session_inner() + } + fn get_any_session_inner(&self) -> Option<(String, AetherSessionClient)> { + if self.backends.is_empty() { + return None; + } + let idx = self.rr_counter.fetch_add(1, Ordering::Relaxed); + let (addr, conn) = self.backends.get(idx % self.backends.len())?; + Some((addr.clone(), conn.session.clone())) + } +} + +/// Extract leader redirect address from a gRPC status, if present. +pub fn extract_leader_redirect(status: &tonic::Status) -> Option { + match status.metadata().get("x-aether-leader") { + Some(val) => match val.to_str() { + Ok(s) => Some(s.to_string()), + Err(e) => { + warn!(error = %e, "x-aether-leader header is not valid UTF-8"); + None + } + }, + None => None, + } +} + +/// Create a gRPC channel, optionally configured with TLS. +pub(super) async fn create_channel( + uri: &str, + tls: &Option, +) -> Result { + let mut endpoint = Channel::from_shared(uri.to_string()) + .map_err(|e| tonic::Status::unavailable(format!("invalid uri: {e}")))?; + + if let Some(tls_cfg) = tls { + let ca = std::fs::read(&tls_cfg.ca_cert) + .map_err(|e| tonic::Status::unavailable(format!("read ca cert: {e}")))?; + let cert = std::fs::read(&tls_cfg.client_cert) + .map_err(|e| tonic::Status::unavailable(format!("read client cert: {e}")))?; + let key = std::fs::read(&tls_cfg.client_key) + .map_err(|e| tonic::Status::unavailable(format!("read client key: {e}")))?; + + let ca_cert = tonic::transport::Certificate::from_pem(ca); + let identity = tonic::transport::Identity::from_pem(cert, key); + let tls_config = tonic::transport::ClientTlsConfig::new() + .ca_certificate(ca_cert) + .identity(identity); + + endpoint = endpoint + .tls_config(tls_config) + .map_err(|e| tonic::Status::unavailable(format!("tls config: {e}")))?; + } + + endpoint + .connect() + .await + .map_err(|e| tonic::Status::unavailable(format!("connect failed: {e}"))) +} + +async fn connect_backend( + uri: &str, + tls: &Option, +) -> Result { + let channel = create_channel(uri, tls).await?; + Ok(connection_from_channel(channel)) +} + +fn connection_from_channel(channel: Channel) -> BackendConnection { + BackendConnection { + kv: AetherKvClient::new(channel.clone()), + cluster: AetherClusterClient::new(channel.clone()), + maintenance: AetherMaintenanceClient::new(channel.clone()), + watch: AetherWatchClient::new(channel.clone()), + lease: AetherLeaseClient::new(channel.clone()), + auth: AetherAuthClient::new(channel.clone()), + lock: AetherLockClient::new(channel.clone()), + election: AetherElectionClient::new(channel.clone()), + barrier: AetherBarrierClient::new(channel.clone()), + queue: AetherQueueClient::new(channel.clone()), + session: AetherSessionClient::new(channel), + } +} + +/// Connect to a redirected leader, cache the connection in the pool, and return it. +pub(super) async fn redirect_connection( + pool: &Arc>, + leader: &str, +) -> Result { + let tls = { pool.read().await.tls().clone() }; + let scheme = if tls.is_some() { "https" } else { "http" }; + let channel = create_channel(&format!("{scheme}://{leader}"), &tls).await?; + let conn = connection_from_channel(channel); + pool.write() + .await + .add_redirect_connection(leader.to_string(), conn.clone()); + Ok(conn) +} + +pub async fn run_gateway(config: GatewayConfig) -> anyhow::Result<()> { + info!( + listen_addr = %config.listen_addr, + backends = ?config.backend_addrs, + timeout_ms = config.request_timeout_ms, + "starting gateway" + ); + + let pool = BackendPool::connect( + &config.backend_addrs, + config.request_timeout_ms, + &config.tls, + ) + .await?; + let pool = Arc::new(RwLock::new(pool)); + + // Spawn health check HTTP server + let health_pool = pool.clone(); + let health_addr: std::net::SocketAddr = config.health_addr.parse()?; + tokio::spawn(async move { + if let Err(e) = serve_health(health_addr, health_pool).await { + tracing::error!(error = %e, "gateway health server failed"); + } + }); + info!(addr = %config.health_addr, "gateway health server started"); + + let addr = config.listen_addr.parse()?; + info!(addr = %config.listen_addr, "gateway listening"); + + tonic::transport::Server::builder() + .add_service(crate::proto::aether_kv_server::AetherKvServer::new( + KvProxy::new(pool.clone()), + )) + .add_service( + crate::proto::aether_cluster_server::AetherClusterServer::new(ClusterProxy::new( + pool.clone(), + )), + ) + .add_service( + crate::proto::aether_maintenance_server::AetherMaintenanceServer::new( + MaintenanceProxy::new(pool.clone()), + ), + ) + .add_service(crate::proto::aether_watch_server::AetherWatchServer::new( + WatchProxy::new(pool.clone()), + )) + .add_service(crate::proto::aether_lease_server::AetherLeaseServer::new( + LeaseProxy::new(pool.clone()), + )) + .add_service(crate::proto::aether_auth_server::AetherAuthServer::new( + AuthProxy::new(pool.clone()), + )) + .add_service(crate::proto::aether_lock_server::AetherLockServer::new( + LockProxy::new(pool.clone()), + )) + .add_service( + crate::proto::aether_election_server::AetherElectionServer::new(ElectionProxy::new( + pool.clone(), + )), + ) + .add_service( + crate::proto::aether_barrier_server::AetherBarrierServer::new(BarrierProxy::new( + pool.clone(), + )), + ) + .add_service(crate::proto::aether_queue_server::AetherQueueServer::new( + QueueProxy::new(pool.clone()), + )) + .add_service( + crate::proto::aether_session_server::AetherSessionServer::new(SessionProxy::new( + pool.clone(), + )), + ) + .serve(addr) + .await?; + + info!("gateway stopped"); + Ok(()) +} + +/// Gateway health check HTTP server. +/// +/// - `GET /health/live` — always 200 if the process is running. +/// - `GET /health/ready` — 200 if at least one backend is reachable. +async fn serve_health( + addr: std::net::SocketAddr, + pool: Arc>, +) -> anyhow::Result<()> { + use http_body_util::Full; + use hyper::body::Bytes; + use hyper::server::conn::http1; + use hyper::service::service_fn; + use hyper::{Request, Response}; + use hyper_util::rt::TokioIo; + use tokio::net::TcpListener; + + let listener = TcpListener::bind(addr).await?; + info!(addr = %addr, "gateway health server listening"); + + loop { + let (stream, _) = match listener.accept().await { + Ok(conn) => conn, + Err(e) => { + warn!(error = %e, "health server accept error"); + tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; + continue; + } + }; + let io = TokioIo::new(stream); + let pool = pool.clone(); + + tokio::spawn(async move { + let service = service_fn(move |req: Request| { + let pool = pool.clone(); + async move { + if req.method() != hyper::Method::GET { + return Ok::<_, hyper::Error>( + Response::builder() + .status(405) + .body(Full::new(Bytes::from("METHOD NOT ALLOWED"))) + .unwrap(), + ); + } + let response = match req.uri().path() { + "/health/live" => Response::builder() + .status(200) + .body(Full::new(Bytes::from("OK"))) + .unwrap(), + "/health/ready" => { + let p = pool.read().await; + if p.backends_count() > 0 { + Response::builder() + .status(200) + .body(Full::new(Bytes::from("OK"))) + .unwrap() + } else { + Response::builder() + .status(503) + .body(Full::new(Bytes::from("NOT READY"))) + .unwrap() + } + } + _ => Response::builder() + .status(404) + .body(Full::new(Bytes::from("NOT FOUND"))) + .unwrap(), + }; + Ok::<_, hyper::Error>(response) + } + }); + + let conn_result = tokio::time::timeout( + std::time::Duration::from_secs(10), + http1::Builder::new().serve_connection(io, service), + ) + .await; + match conn_result { + Ok(Ok(())) => {} + Ok(Err(err)) => { + tracing::debug!(error = %err, "health connection error"); + } + Err(_) => { + tracing::debug!("health connection timed out"); + } + } + }); + } +} diff --git a/src/gateway/queue_proxy.rs b/src/gateway/queue_proxy.rs new file mode 100644 index 0000000..68a71e1 --- /dev/null +++ b/src/gateway/queue_proxy.rs @@ -0,0 +1,130 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_queue_server::AetherQueue; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct QueueProxy { + pool: Arc>, +} +impl QueueProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_queue_client::AetherQueueClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.queue.clone()) + } +} + +#[tonic::async_trait] +impl AetherQueue for QueueProxy { + async fn enqueue( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_queue() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.enqueue(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.enqueue(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn dequeue( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_queue() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.dequeue(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.dequeue(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn peek( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_queue() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout( + timeout, + client.peek(forward_request(&metadata, req.clone())), + ) + .await + { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.peek(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/session_proxy.rs b/src/gateway/session_proxy.rs new file mode 100644 index 0000000..38fe070 --- /dev/null +++ b/src/gateway/session_proxy.rs @@ -0,0 +1,153 @@ +use super::{BackendPool, extract_leader_redirect, forward_request}; +use crate::proto::aether_session_server::AetherSession; +use crate::proto::*; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct SessionProxy { + pool: Arc>, +} +impl SessionProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_session_client::AetherSessionClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.session.clone()) + } +} + +#[tonic::async_trait] +impl AetherSession for SessionProxy { + async fn create( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_session() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.create(forward_request(&metadata, req))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.create(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn close( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_session() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.close(forward_request(&metadata, req))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.close(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + type KeepAliveStream = tonic::Streaming; + + async fn keep_alive( + &self, + request: Request>, + ) -> Result, Status> { + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_session() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + let mut client_stream = request.into_inner(); + let (tx, rx) = tokio::sync::mpsc::channel(128); + tokio::spawn(async move { + while let Ok(Some(msg)) = client_stream.message().await { + if tx.send(msg).await.is_err() { + break; + } + } + }); + let stream = tokio_stream::wrappers::ReceiverStream::new(rx); + match tokio::time::timeout(timeout, client.keep_alive(Request::new(stream))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + Err(Status::unavailable(format!( + "leader redirect to {leader}, please retry" + ))) + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } + + async fn query( + &self, + request: Request, + ) -> Result, Status> { + let metadata = request.metadata().clone(); + let req = request.into_inner(); + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_any_session() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + match tokio::time::timeout(timeout, client.query(forward_request(&metadata, req))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let mut c = self.redirect_and_cache(&leader).await?; + tokio::time::timeout(timeout, c.query(forward_request(&metadata, req))) + .await + .map_err(|_| Status::deadline_exceeded("request timed out"))? + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/gateway/watch_proxy.rs b/src/gateway/watch_proxy.rs new file mode 100644 index 0000000..3ed3348 --- /dev/null +++ b/src/gateway/watch_proxy.rs @@ -0,0 +1,67 @@ +use super::{BackendPool, extract_leader_redirect}; +use crate::proto::aether_watch_server::AetherWatch; +use crate::proto::{WatchRequest, WatchResponse}; +use std::sync::Arc; +use tokio::sync::RwLock; +use tonic::{Request, Response, Status}; + +pub struct WatchProxy { + pool: Arc>, +} +impl WatchProxy { + pub fn new(pool: Arc>) -> Self { + Self { pool } + } + async fn redirect_and_cache( + &self, + leader: &str, + ) -> Result< + crate::proto::aether_watch_client::AetherWatchClient, + Status, + > { + let conn = super::redirect_connection(&self.pool, leader).await?; + Ok(conn.watch.clone()) + } +} + +#[tonic::async_trait] +impl AetherWatch for WatchProxy { + type WatchStream = tonic::Streaming; + + async fn watch( + &self, + request: Request>, + ) -> Result, Status> { + let (timeout, _addr, mut client) = { + let p = self.pool.read().await; + let c = p + .get_leader_watch() + .ok_or_else(|| Status::unavailable("no backends available"))?; + (p.timeout(), c.0, c.1) + }; + let mut client_stream = request.into_inner(); + let (tx, rx) = tokio::sync::mpsc::channel(128); + tokio::spawn(async move { + while let Ok(Some(msg)) = client_stream.message().await { + if tx.send(msg).await.is_err() { + break; + } + } + }); + let stream = tokio_stream::wrappers::ReceiverStream::new(rx); + match tokio::time::timeout(timeout, client.watch(Request::new(stream))).await { + Ok(Ok(resp)) => Ok(resp), + Ok(Err(status)) => { + if let Some(leader) = extract_leader_redirect(&status) { + let _ = self.redirect_and_cache(&leader).await; + Err(Status::unavailable(format!( + "leader redirect to {leader}, please retry" + ))) + } else { + Err(status) + } + } + Err(_) => Err(Status::deadline_exceeded("request timed out")), + } + } +} diff --git a/src/lib.rs b/src/lib.rs index 9f31c58..2978198 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -5,6 +5,7 @@ pub mod cluster; pub mod config; pub mod election; pub mod error; +pub mod gateway; pub mod lease; pub mod lock; pub mod proto; diff --git a/src/main.rs b/src/main.rs index c294f3d..376b959 100644 --- a/src/main.rs +++ b/src/main.rs @@ -141,6 +141,34 @@ struct Cli { /// DNS hosts for dns discovery (comma-separated) #[arg(long, value_delimiter = ',')] discovery_dns_hosts: Option>, + + /// Run in gateway (proxy) mode instead of Raft node mode. + #[arg(long)] + gateway: bool, + + /// Backend cluster node addresses for gateway mode (comma-separated host:port). + #[arg(long, value_delimiter = ',')] + backends: Option>, + + /// Request timeout for gateway-to-backend RPCs in milliseconds. + #[arg(long, default_value_t = 5000)] + request_timeout_ms: u64, + + /// CA certificate file for gateway TLS. + #[arg(long)] + tls_ca: Option, + + /// Client certificate file for gateway TLS. + #[arg(long)] + tls_cert: Option, + + /// Client key file for gateway TLS. + #[arg(long)] + tls_key: Option, + + /// Health check HTTP listen address for gateway mode. + #[arg(long, default_value = "127.0.0.1:9091")] + health_addr: String, } #[tokio::main] @@ -195,6 +223,27 @@ async fn main() -> anyhow::Result<()> { // Initialize tracing with config init_tracing(&config.log); + // Gateway mode: start stateless proxy instead of Raft node + if cli.gateway { + let backends = cli + .backends + .ok_or_else(|| anyhow::anyhow!("--backends is required in gateway mode"))?; + if backends.is_empty() { + anyhow::bail!("at least one backend address is required"); + } + let mut gw_config = aether::gateway::GatewayConfig::new(config.addr, backends); + gw_config.request_timeout_ms = cli.request_timeout_ms; + gw_config.health_addr = cli.health_addr; + if let (Some(ca), Some(cert), Some(key)) = (cli.tls_ca, cli.tls_cert, cli.tls_key) { + gw_config.tls = Some(aether::gateway::TlsConfig { + ca_cert: ca, + client_cert: cert, + client_key: key, + }); + } + return aether::gateway::run_gateway(gw_config).await; + } + tracing::info!("Aether starting..."); tracing::info!( node_id = config.node_id,