From 28b90a0a1b3210aa77740512d200df2d5d0ee69c Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:06:43 -0500 Subject: [PATCH 1/8] feat: add tls dependencies add tokio-rustls, rustls, and rustls-pemfile for upcoming tls support. uses rustls 0.23 with aws-lc-rs crypto backend by default. --- Cargo.toml | 5 +++++ crates/ember-server/Cargo.toml | 5 +++++ 2 files changed, 10 insertions(+) diff --git a/Cargo.toml b/Cargo.toml index 16771806..cf35f203 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -50,6 +50,11 @@ metrics-exporter-prometheus = { version = "0.16", features = ["http-listener"] } # benchmarking criterion = { version = "0.5", features = ["html_reports", "async_tokio"] } +# tls +tokio-rustls = "0.26" +rustls = { version = "0.23", default-features = false, features = ["std", "tls12"] } +rustls-pemfile = "2" + # internal crates (version required for crates.io publishing) emberkv-core = { version = "0.3.0", path = "crates/ember-core" } ember-protocol = { version = "0.3.0", path = "crates/ember-protocol" } diff --git a/crates/ember-server/Cargo.toml b/crates/ember-server/Cargo.toml index 7d1b81f4..49b6e7a6 100644 --- a/crates/ember-server/Cargo.toml +++ b/crates/ember-server/Cargo.toml @@ -28,5 +28,10 @@ metrics-exporter-prometheus = { workspace = true } futures = "0.3" dashmap = "6" +# tls +tokio-rustls = { workspace = true } +rustls = { workspace = true } +rustls-pemfile = { workspace = true } + # optional: better multi-threaded allocation performance tikv-jemallocator = { version = "0.6", optional = true } From 7994419a223ccede50afd555ffd235ad9994de06 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:07:40 -0500 Subject: [PATCH 2/8] feat: add tls config module new tls.rs module provides: - TlsConfig struct for cert/key paths and mTLS settings - load_tls_acceptor() to build tokio-rustls TlsAcceptor - clear error messages for common issues (file not found, invalid PEM) - optional client certificate verification via ca_cert_file --- crates/ember-server/Cargo.toml | 1 + crates/ember-server/src/main.rs | 1 + crates/ember-server/src/tls.rs | 180 ++++++++++++++++++++++++++++++++ 3 files changed, 182 insertions(+) create mode 100644 crates/ember-server/src/tls.rs diff --git a/crates/ember-server/Cargo.toml b/crates/ember-server/Cargo.toml index 49b6e7a6..6ace2686 100644 --- a/crates/ember-server/Cargo.toml +++ b/crates/ember-server/Cargo.toml @@ -25,6 +25,7 @@ tracing-subscriber = { workspace = true } clap = { workspace = true } metrics = { workspace = true } metrics-exporter-prometheus = { workspace = true } +thiserror = { workspace = true } futures = "0.3" dashmap = "6" diff --git a/crates/ember-server/src/main.rs b/crates/ember-server/src/main.rs index a50bfe33..8e25906a 100644 --- a/crates/ember-server/src/main.rs +++ b/crates/ember-server/src/main.rs @@ -12,6 +12,7 @@ mod metrics; mod pubsub; mod server; mod slowlog; +mod tls; use std::net::SocketAddr; use std::path::PathBuf; diff --git a/crates/ember-server/src/tls.rs b/crates/ember-server/src/tls.rs new file mode 100644 index 00000000..6ca142f8 --- /dev/null +++ b/crates/ember-server/src/tls.rs @@ -0,0 +1,180 @@ +//! TLS configuration and certificate loading. +//! +//! Provides utilities for setting up TLS with rustls, including loading +//! certificates and private keys from PEM files, and optionally configuring +//! client certificate verification (mTLS). + +use std::fs::File; +use std::io::BufReader; +use std::path::Path; +use std::sync::Arc; + +use rustls::pki_types::{CertificateDer, PrivateKeyDer}; +use rustls::server::WebPkiClientVerifier; +use rustls::RootCertStore; +use thiserror::Error; +use tokio_rustls::TlsAcceptor; + +/// TLS configuration loaded from command-line arguments. +#[derive(Debug, Clone)] +pub struct TlsConfig { + /// Path to the server certificate (PEM format). + pub cert_file: String, + /// Path to the server private key (PEM format). + pub key_file: String, + /// Optional path to CA certificate for client verification. + pub ca_cert_file: Option, + /// Whether to require client certificates (mTLS). + pub auth_clients: bool, +} + +/// Errors that can occur when loading TLS configuration. +#[derive(Debug, Error)] +pub enum TlsError { + #[error("certificate file not found: {0}")] + CertFileNotFound(String), + + #[error("private key file not found: {0}")] + KeyFileNotFound(String), + + #[error("CA certificate file not found: {0}")] + CaCertFileNotFound(String), + + #[error("failed to read certificate file: {0}")] + CertReadError(#[source] std::io::Error), + + #[error("failed to read private key file: {0}")] + KeyReadError(#[source] std::io::Error), + + #[error("failed to read CA certificate file: {0}")] + CaCertReadError(#[source] std::io::Error), + + #[error("no certificates found in file: {0}")] + NoCertsFound(String), + + #[error("no private key found in file: {0}")] + NoKeyFound(String), + + #[error("failed to build TLS config: {0}")] + ConfigError(#[from] rustls::Error), + + #[error("failed to build client verifier: {0}")] + VerifierError(String), +} + +/// Loads TLS configuration and creates a `TlsAcceptor`. +/// +/// Reads certificates and private key from PEM files. If `ca_cert_file` is +/// specified, sets up client certificate verification. When `auth_clients` +/// is true, clients must present valid certificates signed by the CA. +pub fn load_tls_acceptor(config: &TlsConfig) -> Result { + // load server certificates + let cert_path = Path::new(&config.cert_file); + if !cert_path.exists() { + return Err(TlsError::CertFileNotFound(config.cert_file.clone())); + } + + let cert_file = File::open(cert_path).map_err(TlsError::CertReadError)?; + let mut cert_reader = BufReader::new(cert_file); + let certs: Vec> = rustls_pemfile::certs(&mut cert_reader) + .collect::, _>>() + .map_err(TlsError::CertReadError)?; + + if certs.is_empty() { + return Err(TlsError::NoCertsFound(config.cert_file.clone())); + } + + // load private key + let key_path = Path::new(&config.key_file); + if !key_path.exists() { + return Err(TlsError::KeyFileNotFound(config.key_file.clone())); + } + + let key_file = File::open(key_path).map_err(TlsError::KeyReadError)?; + let mut key_reader = BufReader::new(key_file); + let key: PrivateKeyDer<'static> = rustls_pemfile::private_key(&mut key_reader) + .map_err(TlsError::KeyReadError)? + .ok_or_else(|| TlsError::NoKeyFound(config.key_file.clone()))?; + + // build server config + let server_config = if let Some(ref ca_path) = config.ca_cert_file { + // load CA cert for client verification + let ca_cert_path = Path::new(ca_path); + if !ca_cert_path.exists() { + return Err(TlsError::CaCertFileNotFound(ca_path.clone())); + } + + let ca_file = File::open(ca_cert_path).map_err(TlsError::CaCertReadError)?; + let mut ca_reader = BufReader::new(ca_file); + let ca_certs: Vec> = rustls_pemfile::certs(&mut ca_reader) + .collect::, _>>() + .map_err(TlsError::CaCertReadError)?; + + let mut root_store = RootCertStore::empty(); + for cert in ca_certs { + root_store + .add(cert) + .map_err(|e| TlsError::VerifierError(e.to_string()))?; + } + + let verifier = if config.auth_clients { + // require client certs + WebPkiClientVerifier::builder(Arc::new(root_store)) + .build() + .map_err(|e| TlsError::VerifierError(e.to_string()))? + } else { + // allow but don't require client certs + WebPkiClientVerifier::builder(Arc::new(root_store)) + .allow_unauthenticated() + .build() + .map_err(|e| TlsError::VerifierError(e.to_string()))? + }; + + rustls::ServerConfig::builder() + .with_client_cert_verifier(verifier) + .with_single_cert(certs, key)? + } else { + // no client verification + rustls::ServerConfig::builder() + .with_no_client_auth() + .with_single_cert(certs, key)? + }; + + Ok(TlsAcceptor::from(Arc::new(server_config))) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_missing_cert_file() { + let config = TlsConfig { + cert_file: "/nonexistent/cert.pem".into(), + key_file: "/nonexistent/key.pem".into(), + ca_cert_file: None, + auth_clients: false, + }; + let err = load_tls_acceptor(&config).unwrap_err(); + assert!(matches!(err, TlsError::CertFileNotFound(_))); + } + + #[test] + fn test_missing_key_file() { + // create temp cert file + let tmp = std::env::temp_dir().join("test_cert.pem"); + std::fs::write(&tmp, "-----BEGIN CERTIFICATE-----\ntest\n-----END CERTIFICATE-----\n") + .unwrap(); + + let config = TlsConfig { + cert_file: tmp.to_string_lossy().into(), + key_file: "/nonexistent/key.pem".into(), + ca_cert_file: None, + auth_clients: false, + }; + let err = load_tls_acceptor(&config).unwrap_err(); + assert!(matches!(err, TlsError::KeyFileNotFound(_))); + + std::fs::remove_file(tmp).ok(); + } +} From f96ced7bac67caeaf13984a959a72c9c79f374e0 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:08:23 -0500 Subject: [PATCH 3/8] feat: add tls cli arguments add redis-compatible tls flags: - --tls-port: port for tls connections - --tls-cert-file: server certificate path - --tls-key-file: server private key path - --tls-ca-cert-file: ca cert for client verification - --tls-auth-clients: require client certs (yes/no) validates that cert and key are provided when tls-port is set. --- crates/ember-server/src/main.rs | 65 +++++++++++++++++++++++++++++++++ 1 file changed, 65 insertions(+) diff --git a/crates/ember-server/src/main.rs b/crates/ember-server/src/main.rs index 8e25906a..e1aff078 100644 --- a/crates/ember-server/src/main.rs +++ b/crates/ember-server/src/main.rs @@ -82,6 +82,28 @@ struct Args { /// when set, connections must authenticate before executing any data commands. #[arg(long)] requirepass: Option, + + // -- TLS options (matching redis) -- + /// port for TLS connections. when set, enables TLS alongside plain TCP + #[arg(long)] + tls_port: Option, + + /// path to server certificate file (PEM format) + #[arg(long)] + tls_cert_file: Option, + + /// path to server private key file (PEM format) + #[arg(long)] + tls_key_file: Option, + + /// path to CA certificate for client verification (enables mTLS) + #[arg(long)] + tls_ca_cert_file: Option, + + /// require client certificates when CA cert is configured. + /// accepts: yes, no. default: no + #[arg(long, default_value = "no")] + tls_auth_clients: String, } #[tokio::main] @@ -192,6 +214,49 @@ async fn main() { info!("authentication enabled (requirepass set)"); } + // build TLS config if --tls-port is set + let tls_config = if let Some(tls_port) = args.tls_port { + let cert_file = args.tls_cert_file.unwrap_or_else(|| { + eprintln!("--tls-port requires --tls-cert-file and --tls-key-file"); + std::process::exit(1); + }); + let key_file = args.tls_key_file.unwrap_or_else(|| { + eprintln!("--tls-port requires --tls-cert-file and --tls-key-file"); + std::process::exit(1); + }); + + let auth_clients = match args.tls_auth_clients.to_lowercase().as_str() { + "yes" | "true" | "1" => true, + "no" | "false" | "0" => false, + _ => { + eprintln!("--tls-auth-clients must be 'yes' or 'no'"); + std::process::exit(1); + } + }; + + let tls_addr: SocketAddr = format!("{}:{}", args.host, tls_port) + .parse() + .expect("invalid TLS bind address"); + + info!( + tls_port = tls_port, + cert = %cert_file, + "TLS enabled" + ); + + Some(( + tls_addr, + tls::TlsConfig { + cert_file, + key_file, + ca_cert_file: args.tls_ca_cert_file, + auth_clients, + }, + )) + } else { + None + }; + let result = if args.concurrent { server::run_concurrent( addr, From e58d046a99b593c473980565e42dc89319669d30 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:10:08 -0500 Subject: [PATCH 4/8] refactor: make connection handlers generic over stream type update handle() in connection.rs and concurrent_handler.rs to accept any type implementing AsyncRead + AsyncWrite + Unpin. this enables the same handler code to work with both TcpStream and TlsStream. moved set_nodelay() call to server.rs accept loop so it applies to the underlying tcp socket before any tls handshake. --- crates/ember-server/src/concurrent_handler.rs | 17 ++++---- crates/ember-server/src/connection.rs | 30 ++++++++------ crates/ember-server/src/server.rs | 41 +++++++++++-------- 3 files changed, 52 insertions(+), 36 deletions(-) diff --git a/crates/ember-server/src/concurrent_handler.rs b/crates/ember-server/src/concurrent_handler.rs index 17438fa6..fb3d68a7 100644 --- a/crates/ember-server/src/concurrent_handler.rs +++ b/crates/ember-server/src/concurrent_handler.rs @@ -19,8 +19,7 @@ use std::time::{Duration, Instant}; use bytes::BytesMut; use ember_core::{ConcurrentKeyspace, Engine, TtlResult}; use ember_protocol::{parse_frame, Command, Frame, SetExpire}; -use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tokio::net::TcpStream; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; use crate::connection_common::{ is_allowed_before_auth, is_auth_frame, try_auth, BUF_CAPACITY, IDLE_TIMEOUT, MAX_BUF_SIZE, @@ -30,16 +29,20 @@ use crate::server::ServerContext; use crate::slowlog::SlowLog; /// Handles a connection using the concurrent keyspace for GET/SET. -pub async fn handle( - mut stream: TcpStream, +/// +/// Generic over the stream type to support both plain TCP and TLS connections. +/// Callers should set TCP_NODELAY on the underlying socket before calling. +pub async fn handle( + mut stream: S, keyspace: Arc, engine: Engine, // fallback for complex commands ctx: &Arc, slow_log: &Arc, pubsub: &Arc, -) -> Result<(), Box> { - stream.set_nodelay(true)?; - +) -> Result<(), Box> +where + S: AsyncRead + AsyncWrite + Unpin, +{ let mut authenticated = ctx.requirepass.is_none(); let mut buf = BytesMut::with_capacity(BUF_CAPACITY); diff --git a/crates/ember-server/src/connection.rs b/crates/ember-server/src/connection.rs index ce181e57..86e1e099 100644 --- a/crates/ember-server/src/connection.rs +++ b/crates/ember-server/src/connection.rs @@ -1,6 +1,6 @@ //! Per-connection handler for sharded engine mode. //! -//! Reads RESP3 frames from a TCP stream, routes them through the +//! Reads RESP3 frames from a TCP/TLS stream, routes them through the //! sharded engine, and writes responses back. Supports pipelining //! by dispatching multiple commands concurrently to shards using //! `join_all` for parallel execution. @@ -14,8 +14,7 @@ use bytes::{Bytes, BytesMut}; use ember_core::{Engine, KeyspaceStats, ShardRequest, ShardResponse, TtlResult, Value}; use ember_protocol::{parse_frame, Command, Frame, SetExpire}; use futures::future::join_all; -use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tokio::net::TcpStream; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; use tokio::sync::broadcast; use crate::connection_common::{ @@ -30,17 +29,19 @@ use crate::slowlog::SlowLog; /// Reads data into a buffer, parses complete frames, dispatches commands /// through the engine, and writes serialized responses back. The loop /// exits when the client disconnects or a protocol error occurs. -pub async fn handle( - mut stream: TcpStream, +/// +/// Generic over the stream type to support both plain TCP and TLS connections. +/// Callers should set TCP_NODELAY on the underlying socket before calling. +pub async fn handle( + mut stream: S, engine: Engine, ctx: &Arc, slow_log: &Arc, pubsub: &Arc, -) -> Result<(), Box> { - // disable Nagle's algorithm — cache servers need low-latency writes, - // and we already batch responses from pipelining into a single write - stream.set_nodelay(true)?; - +) -> Result<(), Box> +where + S: AsyncRead + AsyncWrite + Unpin, +{ // per-connection auth state. auto-authenticated when no password is set. let mut authenticated = ctx.requirepass.is_none(); @@ -173,13 +174,16 @@ fn is_subscribe_frame(frame: &Frame) -> bool { /// PSUBSCRIBE, PUNSUBSCRIBE, and PING. All other commands return an error. /// Returns to the caller when all subscriptions are removed or the client /// disconnects. -async fn handle_subscriber_mode( - stream: &mut TcpStream, +async fn handle_subscriber_mode( + stream: &mut S, buf: &mut BytesMut, out: &mut BytesMut, pubsub: &Arc, initial_frames: Vec, -) -> Result<(), Box> { +) -> Result<(), Box> +where + S: AsyncRead + AsyncWrite + Unpin, +{ // track subscriptions: channel/pattern -> receiver let mut channel_rxs: HashMap> = HashMap::new(); let mut pattern_rxs: HashMap> = HashMap::new(); diff --git a/crates/ember-server/src/server.rs b/crates/ember-server/src/server.rs index 7d39642b..3a45bd08 100644 --- a/crates/ember-server/src/server.rs +++ b/crates/ember-server/src/server.rs @@ -122,18 +122,18 @@ pub async fn run( } result = listener.accept() => { - let (mut stream, peer) = result?; + let (stream, peer) = result?; + + // disable Nagle's algorithm — cache servers need low-latency writes. + // done here before passing to handler so TLS streams also benefit. + if let Err(e) = stream.set_nodelay(true) { + warn!("failed to set TCP_NODELAY: {e}"); + } // protected mode: reject non-loopback connections when no // password is set and the server is bound to a public address if is_protected_mode_violation(&ctx, &peer) { - let msg = "-DENIED Ember is running in protected mode \ - because no password is set. In this mode \ - connections are only accepted from the loopback \ - interface. Set a password with --requirepass or \ - bind to 127.0.0.1 to resolve this.\r\n"; - let _ = stream.write_all(msg.as_bytes()).await; - let _ = stream.shutdown().await; + reject_protected_mode(stream).await; continue; } @@ -183,6 +183,17 @@ pub async fn run( Ok(()) } +/// Sends the protected mode rejection message and closes the connection. +async fn reject_protected_mode(mut stream: tokio::net::TcpStream) { + let msg = "-DENIED Ember is running in protected mode \ + because no password is set. In this mode \ + connections are only accepted from the loopback \ + interface. Set a password with --requirepass or \ + bind to 127.0.0.1 to resolve this.\r\n"; + let _ = stream.write_all(msg.as_bytes()).await; + let _ = stream.shutdown().await; +} + /// Runs the server with a concurrent keyspace (DashMap-backed). /// /// This mode bypasses shard channels for GET/SET operations, accessing @@ -253,16 +264,14 @@ pub async fn run_concurrent( } result = listener.accept() => { - let (mut stream, peer) = result?; + let (stream, peer) = result?; + + if let Err(e) = stream.set_nodelay(true) { + warn!("failed to set TCP_NODELAY: {e}"); + } if is_protected_mode_violation(&ctx, &peer) { - let msg = "-DENIED Ember is running in protected mode \ - because no password is set. In this mode \ - connections are only accepted from the loopback \ - interface. Set a password with --requirepass or \ - bind to 127.0.0.1 to resolve this.\r\n"; - let _ = stream.write_all(msg.as_bytes()).await; - let _ = stream.shutdown().await; + reject_protected_mode(stream).await; continue; } From b1d9c5220c0e9a0c929c05e24eaccae5c5786789 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:13:02 -0500 Subject: [PATCH 5/8] feat: add tls listener to server - run() and run_concurrent() now accept optional TLS config - when TLS is configured: - loads certificate and key via load_tls_acceptor() - binds separate TLS listener on --tls-port - performs TLS handshake after TCP accept - passes TlsStream to generic connection handler - both plain TCP and TLS share the same: - connection semaphore (combined limit) - ServerContext and metrics - graceful shutdown handling - TLS handshake failures logged as warnings, connection dropped --- crates/ember-server/src/main.rs | 2 + crates/ember-server/src/server.rs | 169 +++++++++++++++++++++++++++++- crates/ember-server/src/tls.rs | 14 ++- 3 files changed, 180 insertions(+), 5 deletions(-) diff --git a/crates/ember-server/src/main.rs b/crates/ember-server/src/main.rs index e1aff078..b94d57ec 100644 --- a/crates/ember-server/src/main.rs +++ b/crates/ember-server/src/main.rs @@ -268,6 +268,7 @@ async fn main() { args.metrics_port.is_some(), slowlog_config, args.requirepass, + tls_config, ) .await } else { @@ -279,6 +280,7 @@ async fn main() { args.metrics_port.is_some(), slowlog_config, args.requirepass, + tls_config, ) .await }; diff --git a/crates/ember-server/src/server.rs b/crates/ember-server/src/server.rs index 3a45bd08..a7c60528 100644 --- a/crates/ember-server/src/server.rs +++ b/crates/ember-server/src/server.rs @@ -1,4 +1,4 @@ -//! TCP server that accepts client connections and spawns handler tasks. +//! TCP/TLS server that accepts client connections and spawns handler tasks. //! //! Handles graceful shutdown on SIGINT/SIGTERM: stops accepting new //! connections and waits for in-flight requests to drain before exiting. @@ -12,11 +12,13 @@ use ember_core::{ConcurrentKeyspace, Engine, EngineConfig, EvictionPolicy}; use tokio::io::AsyncWriteExt; use tokio::net::TcpListener; use tokio::sync::Semaphore; +use tokio_rustls::TlsAcceptor; use tracing::{error, info, warn}; use crate::connection; use crate::pubsub::PubSubManager; use crate::slowlog::{SlowLog, SlowLogConfig}; +use crate::tls::TlsConfig; /// Default maximum number of concurrent client connections. const DEFAULT_MAX_CONNECTIONS: usize = 10_000; @@ -50,8 +52,12 @@ pub struct ServerContext { /// Limits concurrent connections to `max_connections` — excess clients /// are dropped immediately. /// +/// If `tls` is provided, also binds a TLS listener on the specified address. +/// Both plain TCP and TLS connections share the same engine and connection limits. +/// /// On SIGINT or SIGTERM the server stops accepting new connections, /// waits for existing handlers to finish, then exits cleanly. +#[allow(clippy::too_many_arguments)] pub async fn run( addr: SocketAddr, shard_count: usize, @@ -60,6 +66,7 @@ pub async fn run( metrics_enabled: bool, slowlog_config: SlowLogConfig, requirepass: Option, + tls: Option<(SocketAddr, TlsConfig)>, ) -> Result<(), Box> { // ensure data directory exists if persistence is configured if let Some(ref pcfg) = config.persistence { @@ -86,6 +93,17 @@ pub async fn run( let max_conn = max_connections.unwrap_or(DEFAULT_MAX_CONNECTIONS); let semaphore = Arc::new(Semaphore::new(max_conn)); + // set up TLS listener if configured + let tls_listener: Option<(TcpListener, TlsAcceptor)> = if let Some((tls_addr, tls_config)) = tls + { + let acceptor = crate::tls::load_tls_acceptor(&tls_config)?; + let tls_tcp = TcpListener::bind(tls_addr).await?; + info!("TLS listening on {tls_addr}"); + Some((tls_tcp, acceptor)) + } else { + None + }; + let ctx = Arc::new(ServerContext { start_time: Instant::now(), version: env!("CARGO_PKG_VERSION"), @@ -112,6 +130,14 @@ pub async fn run( let shutdown = tokio::signal::ctrl_c(); tokio::pin!(shutdown); + // helper to accept from TLS listener or pend forever if disabled + let tls_accept = || async { + match &tls_listener { + Some((listener, _)) => listener.accept().await, + None => std::future::pending().await, + } + }; + loop { tokio::select! { biased; @@ -121,6 +147,7 @@ pub async fn run( break; } + // plain TCP accept result = listener.accept() => { let (stream, peer) = result?; @@ -172,6 +199,64 @@ pub async fn run( drop(permit); }); } + + // TLS accept (pends forever if TLS not configured) + result = tls_accept() => { + let (stream, peer) = result?; + + if let Err(e) = stream.set_nodelay(true) { + warn!("failed to set TCP_NODELAY: {e}"); + } + + // protected mode check on TLS connections too + if is_protected_mode_violation(&ctx, &peer) { + reject_protected_mode(stream).await; + continue; + } + + let permit = match semaphore.clone().try_acquire_owned() { + Ok(permit) => permit, + Err(_) => { + warn!("connection limit reached, dropping TLS connection from {peer}"); + if metrics_enabled { + crate::metrics::on_connection_rejected(); + } + drop(stream); + continue; + } + }; + + if metrics_enabled { + crate::metrics::on_connection_accepted(); + } + ctx.connections_accepted.fetch_add(1, Ordering::Relaxed); + ctx.connections_active.fetch_add(1, Ordering::Relaxed); + + let engine = engine.clone(); + let ctx = Arc::clone(&ctx); + let slow_log = Arc::clone(&slow_log); + let pubsub = Arc::clone(&pubsub); + let acceptor = tls_listener.as_ref().map(|(_, a)| a.clone()).unwrap(); + + tokio::spawn(async move { + // perform TLS handshake + match acceptor.accept(stream).await { + Ok(tls_stream) => { + if let Err(e) = connection::handle(tls_stream, engine, &ctx, &slow_log, &pubsub).await { + error!("TLS connection error from {peer}: {e}"); + } + } + Err(e) => { + warn!("TLS handshake failed from {peer}: {e}"); + } + } + ctx.connections_active.fetch_sub(1, Ordering::Relaxed); + if ctx.metrics_enabled { + crate::metrics::on_connection_closed(); + } + drop(permit); + }); + } } } @@ -199,6 +284,8 @@ async fn reject_protected_mode(mut stream: tokio::net::TcpStream) { /// This mode bypasses shard channels for GET/SET operations, accessing /// the keyspace directly from connection handlers. Falls back to the /// sharded engine for complex commands. +/// +/// If `tls` is provided, also binds a TLS listener on the specified address. #[allow(clippy::too_many_arguments)] pub async fn run_concurrent( addr: SocketAddr, @@ -210,6 +297,7 @@ pub async fn run_concurrent( metrics_enabled: bool, slowlog_config: SlowLogConfig, requirepass: Option, + tls: Option<(SocketAddr, TlsConfig)>, ) -> Result<(), Box> { let aof_enabled = config .persistence @@ -231,6 +319,17 @@ pub async fn run_concurrent( let max_conn = max_connections.unwrap_or(DEFAULT_MAX_CONNECTIONS); let semaphore = Arc::new(Semaphore::new(max_conn)); + // set up TLS listener if configured + let tls_listener: Option<(TcpListener, TlsAcceptor)> = if let Some((tls_addr, tls_config)) = tls + { + let acceptor = crate::tls::load_tls_acceptor(&tls_config)?; + let tls_tcp = TcpListener::bind(tls_addr).await?; + info!("TLS listening on {tls_addr}"); + Some((tls_tcp, acceptor)) + } else { + None + }; + let ctx = Arc::new(ServerContext { start_time: Instant::now(), version: env!("CARGO_PKG_VERSION"), @@ -254,6 +353,14 @@ pub async fn run_concurrent( let shutdown = tokio::signal::ctrl_c(); tokio::pin!(shutdown); + // helper to accept from TLS listener or pend forever if disabled + let tls_accept = || async { + match &tls_listener { + Some((listener, _)) => listener.accept().await, + None => std::future::pending().await, + } + }; + loop { tokio::select! { biased; @@ -263,6 +370,7 @@ pub async fn run_concurrent( break; } + // plain TCP accept result = listener.accept() => { let (stream, peer) = result?; @@ -312,6 +420,65 @@ pub async fn run_concurrent( drop(permit); }); } + + // TLS accept (pends forever if TLS not configured) + result = tls_accept() => { + let (stream, peer) = result?; + + if let Err(e) = stream.set_nodelay(true) { + warn!("failed to set TCP_NODELAY: {e}"); + } + + if is_protected_mode_violation(&ctx, &peer) { + reject_protected_mode(stream).await; + continue; + } + + let permit = match semaphore.clone().try_acquire_owned() { + Ok(permit) => permit, + Err(_) => { + warn!("connection limit reached, dropping TLS connection from {peer}"); + if metrics_enabled { + crate::metrics::on_connection_rejected(); + } + drop(stream); + continue; + } + }; + + if metrics_enabled { + crate::metrics::on_connection_accepted(); + } + ctx.connections_accepted.fetch_add(1, Ordering::Relaxed); + ctx.connections_active.fetch_add(1, Ordering::Relaxed); + + let keyspace = Arc::clone(&keyspace); + let engine = engine.clone(); + let ctx = Arc::clone(&ctx); + let slow_log = Arc::clone(&slow_log); + let pubsub = Arc::clone(&pubsub); + let acceptor = tls_listener.as_ref().map(|(_, a)| a.clone()).unwrap(); + + tokio::spawn(async move { + match acceptor.accept(stream).await { + Ok(tls_stream) => { + if let Err(e) = crate::concurrent_handler::handle( + tls_stream, keyspace, engine, &ctx, &slow_log, &pubsub + ).await { + error!("TLS connection error from {peer}: {e}"); + } + } + Err(e) => { + warn!("TLS handshake failed from {peer}: {e}"); + } + } + ctx.connections_active.fetch_sub(1, Ordering::Relaxed); + if ctx.metrics_enabled { + crate::metrics::on_connection_closed(); + } + drop(permit); + }); + } } } diff --git a/crates/ember-server/src/tls.rs b/crates/ember-server/src/tls.rs index 6ca142f8..5e62575d 100644 --- a/crates/ember-server/src/tls.rs +++ b/crates/ember-server/src/tls.rs @@ -155,8 +155,11 @@ mod tests { ca_cert_file: None, auth_clients: false, }; - let err = load_tls_acceptor(&config).unwrap_err(); - assert!(matches!(err, TlsError::CertFileNotFound(_))); + match load_tls_acceptor(&config) { + Err(TlsError::CertFileNotFound(_)) => {} + Err(e) => panic!("expected CertFileNotFound, got: {e}"), + Ok(_) => panic!("expected error, got Ok"), + } } #[test] @@ -172,8 +175,11 @@ mod tests { ca_cert_file: None, auth_clients: false, }; - let err = load_tls_acceptor(&config).unwrap_err(); - assert!(matches!(err, TlsError::KeyFileNotFound(_))); + match load_tls_acceptor(&config) { + Err(TlsError::KeyFileNotFound(_)) => {} + Err(e) => panic!("expected KeyFileNotFound, got: {e}"), + Ok(_) => panic!("expected error, got Ok"), + } std::fs::remove_file(tmp).ok(); } From 7b00c79c393ef0d00b4cac08562aadad6f14cdf9 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:13:42 -0500 Subject: [PATCH 6/8] docs: add tls configuration to readme - add tls to features list - document all tls cli flags in configuration table - add tls server startup example - add redis-cli tls connection examples --- README.md | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/README.md b/README.md index 8aff7a92..343c9bd7 100644 --- a/README.md +++ b/README.md @@ -18,6 +18,7 @@ a low-latency, memory-efficient, distributed cache written in Rust. designed to - **server commands** — PING, ECHO, INFO, DBSIZE, FLUSHDB, BGSAVE, BGREWRITEAOF, AUTH, QUIT - **pub/sub** — SUBSCRIBE, UNSUBSCRIBE, PSUBSCRIBE, PUNSUBSCRIBE, PUBLISH, plus PUBSUB introspection - **authentication** — `--requirepass` for redis-compatible AUTH (legacy and username/password forms) +- **tls support** — redis-compatible TLS on a separate port, with optional mTLS for client certificates - **protected mode** — rejects non-loopback connections when no password is set on public binds - **observability** — prometheus metrics (`--metrics-port`), enriched INFO with 6 sections, SLOWLOG command - **sharded engine** — shared-nothing, thread-per-core design with no cross-shard locking @@ -46,6 +47,10 @@ cargo build --release # concurrent mode (experimental, 2x faster for GET/SET) ./target/release/ember-server --concurrent + +# with TLS (runs alongside plain TCP) +./target/release/ember-server --tls-port 6380 \ + --tls-cert-file cert.pem --tls-key-file key.pem ``` ```bash @@ -92,6 +97,11 @@ redis-cli SREM tags fast # => (integer) 1 redis-cli SCAN 0 MATCH "user:*" COUNT 100 redis-cli DBSIZE # => (integer) 6 redis-cli FLUSHDB # => OK + +# TLS connection +redis-cli -p 6380 --tls --insecure PING +# or with cert verification +redis-cli -p 6380 --tls --cacert cert.pem PING ``` ## configuration @@ -111,6 +121,11 @@ redis-cli FLUSHDB # => OK | `--slowlog-max-len` | 128 | max entries in slow log ring buffer | | `--concurrent` | false | use DashMap-backed keyspace (experimental, faster GET/SET) | | `--requirepass` | — | require AUTH with this password before running commands | +| `--tls-port` | — | port for TLS connections (enables TLS when set) | +| `--tls-cert-file` | — | path to server certificate (PEM) | +| `--tls-key-file` | — | path to server private key (PEM) | +| `--tls-ca-cert-file` | — | path to CA certificate for client verification | +| `--tls-auth-clients` | no | require client certificates (`yes` or `no`) | ## build & development From 656cda6caf3428d756c2755edb82d5f16a7d0c02 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:17:44 -0500 Subject: [PATCH 7/8] fix: correct misplaced doc comment for glob_match the doc comment for glob_match was accidentally placed above format_float. moved it to the correct location. --- crates/ember-core/src/keyspace.rs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/crates/ember-core/src/keyspace.rs b/crates/ember-core/src/keyspace.rs index 56194ff5..5443a2c3 100644 --- a/crates/ember-core/src/keyspace.rs +++ b/crates/ember-core/src/keyspace.rs @@ -1758,12 +1758,6 @@ impl Default for Keyspace { } } -/// Glob-style pattern matching for SCAN's MATCH option. -/// -/// Supports: -/// - `*` matches any sequence of characters (including empty) -/// - `?` matches exactly one character -/// - `[abc]` matches one character from the set /// Formats a float value matching Redis behavior. /// /// Uses up to 17 significant digits and strips unnecessary trailing zeros, @@ -1786,6 +1780,12 @@ fn format_float(val: f64) -> String { } } +/// Glob-style pattern matching for SCAN's MATCH option. +/// +/// Supports: +/// - `*` matches any sequence of characters (including empty) +/// - `?` matches exactly one character +/// - `[abc]` matches one character from the set /// - `[^abc]` or `[!abc]` matches one character NOT in the set /// /// Uses an iterative two-pointer algorithm with backtracking for O(n*m) From cfced606781b065119cf6f28e8e7058edc551ba4 Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Sun, 8 Feb 2026 10:19:47 -0500 Subject: [PATCH 8/8] style: apply cargo fmt to tls.rs --- crates/ember-server/src/tls.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/crates/ember-server/src/tls.rs b/crates/ember-server/src/tls.rs index 5e62575d..d642f84b 100644 --- a/crates/ember-server/src/tls.rs +++ b/crates/ember-server/src/tls.rs @@ -166,8 +166,11 @@ mod tests { fn test_missing_key_file() { // create temp cert file let tmp = std::env::temp_dir().join("test_cert.pem"); - std::fs::write(&tmp, "-----BEGIN CERTIFICATE-----\ntest\n-----END CERTIFICATE-----\n") - .unwrap(); + std::fs::write( + &tmp, + "-----BEGIN CERTIFICATE-----\ntest\n-----END CERTIFICATE-----\n", + ) + .unwrap(); let config = TlsConfig { cert_file: tmp.to_string_lossy().into(),