use std::sync::Arc;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::time::Duration;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum PipelineMode {
#[default]
Sequential,
Concurrent,
OutOfOrder,
}
impl PipelineMode {
#[inline]
pub fn maintains_order(&self) -> bool {
matches!(self, Self::Sequential | Self::Concurrent)
}
#[inline]
pub fn is_concurrent(&self) -> bool {
matches!(self, Self::Concurrent | Self::OutOfOrder)
}
}
#[derive(Debug, Clone)]
pub struct PipelineConfig {
pub mode: PipelineMode,
pub max_concurrent: usize,
pub pipeline_flush: bool,
pub max_buffered_requests: usize,
pub keep_alive_timeout: Duration,
pub request_timeout: Option<Duration>,
pub write_timeout: Option<Duration>,
pub max_requests_per_connection: Option<u64>,
pub tcp_nodelay: bool,
pub read_buffer_size: usize,
pub write_buffer_size: usize,
pub max_header_size: usize,
}
impl Default for PipelineConfig {
fn default() -> Self {
Self {
mode: PipelineMode::Concurrent,
max_concurrent: 16,
pipeline_flush: true,
max_buffered_requests: 64,
keep_alive_timeout: Duration::from_secs(60),
request_timeout: None,
write_timeout: Some(Duration::from_secs(300)),
max_requests_per_connection: Some(10_000),
tcp_nodelay: true,
read_buffer_size: 8192,
write_buffer_size: 8192,
max_header_size: 16384,
}
}
}
impl PipelineConfig {
pub fn builder() -> PipelineConfigBuilder {
PipelineConfigBuilder::default()
}
pub fn high_performance() -> Self {
Self {
mode: PipelineMode::Concurrent,
max_concurrent: 32,
pipeline_flush: true,
max_buffered_requests: 128,
keep_alive_timeout: Duration::from_secs(120),
request_timeout: None,
write_timeout: Some(Duration::from_secs(300)),
max_requests_per_connection: Some(100_000),
tcp_nodelay: true,
read_buffer_size: 16384,
write_buffer_size: 16384,
max_header_size: 32768,
}
}
pub fn low_latency() -> Self {
Self {
mode: PipelineMode::Sequential,
max_concurrent: 1,
pipeline_flush: false,
max_buffered_requests: 16,
keep_alive_timeout: Duration::from_secs(30),
request_timeout: None,
write_timeout: Some(Duration::from_secs(60)),
max_requests_per_connection: Some(1000),
tcp_nodelay: true,
read_buffer_size: 4096,
write_buffer_size: 4096,
max_header_size: 8192,
}
}
pub fn memory_efficient() -> Self {
Self {
mode: PipelineMode::Sequential,
max_concurrent: 4,
pipeline_flush: true,
max_buffered_requests: 32,
keep_alive_timeout: Duration::from_secs(30),
request_timeout: None,
write_timeout: Some(Duration::from_secs(300)),
max_requests_per_connection: Some(1000),
tcp_nodelay: false,
read_buffer_size: 4096,
write_buffer_size: 4096,
max_header_size: 8192,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct PipelineConfigBuilder {
config: PipelineConfig,
}
impl PipelineConfigBuilder {
pub fn mode(mut self, mode: PipelineMode) -> Self {
self.config.mode = mode;
self
}
pub fn max_concurrent(mut self, max: usize) -> Self {
self.config.max_concurrent = max;
self
}
pub fn pipeline_flush(mut self, enable: bool) -> Self {
self.config.pipeline_flush = enable;
self
}
pub fn max_buffered_requests(mut self, max: usize) -> Self {
self.config.max_buffered_requests = max;
self
}
pub fn keep_alive_timeout(mut self, timeout: Duration) -> Self {
self.config.keep_alive_timeout = timeout;
self
}
pub fn max_requests_per_connection(mut self, max: Option<u64>) -> Self {
self.config.max_requests_per_connection = max;
self
}
pub fn tcp_nodelay(mut self, enable: bool) -> Self {
self.config.tcp_nodelay = enable;
self
}
pub fn read_buffer_size(mut self, size: usize) -> Self {
self.config.read_buffer_size = size;
self
}
pub fn write_buffer_size(mut self, size: usize) -> Self {
self.config.write_buffer_size = size;
self
}
pub fn max_header_size(mut self, size: usize) -> Self {
self.config.max_header_size = size;
self
}
pub fn build(self) -> PipelineConfig {
self.config
}
}
#[derive(Debug)]
pub struct ConnectionStats {
requests_processed: AtomicU64,
pending_requests: AtomicUsize,
bytes_received: AtomicU64,
bytes_sent: AtomicU64,
pipeline_depth: AtomicUsize,
}
impl Default for ConnectionStats {
fn default() -> Self {
Self::new()
}
}
impl ConnectionStats {
pub fn new() -> Self {
Self {
requests_processed: AtomicU64::new(0),
pending_requests: AtomicUsize::new(0),
bytes_received: AtomicU64::new(0),
bytes_sent: AtomicU64::new(0),
pipeline_depth: AtomicUsize::new(0),
}
}
#[inline]
pub fn request_received(&self, bytes: u64) {
self.pending_requests.fetch_add(1, Ordering::Relaxed);
self.bytes_received.fetch_add(bytes, Ordering::Relaxed);
self.pipeline_depth.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn response_sent(&self, bytes: u64) {
self.requests_processed.fetch_add(1, Ordering::Relaxed);
self.pending_requests.fetch_sub(1, Ordering::Relaxed);
self.bytes_sent.fetch_add(bytes, Ordering::Relaxed);
self.pipeline_depth.fetch_sub(1, Ordering::Relaxed);
}
#[inline]
pub fn requests_processed(&self) -> u64 {
self.requests_processed.load(Ordering::Relaxed)
}
#[inline]
pub fn pending_requests(&self) -> usize {
self.pending_requests.load(Ordering::Relaxed)
}
#[inline]
pub fn pipeline_depth(&self) -> usize {
self.pipeline_depth.load(Ordering::Relaxed)
}
#[inline]
pub fn bytes_received(&self) -> u64 {
self.bytes_received.load(Ordering::Relaxed)
}
#[inline]
pub fn bytes_sent(&self) -> u64 {
self.bytes_sent.load(Ordering::Relaxed)
}
}
#[derive(Debug, Default)]
pub struct PipelineStats {
active_connections: AtomicUsize,
total_connections: AtomicU64,
total_requests: AtomicU64,
avg_pipeline_depth: AtomicU64,
max_pipeline_depth: AtomicUsize,
}
impl PipelineStats {
pub fn new() -> Self {
Self::default()
}
#[inline]
pub fn connection_opened(&self) {
self.active_connections.fetch_add(1, Ordering::Relaxed);
self.total_connections.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn connection_closed(&self) {
self.active_connections.fetch_sub(1, Ordering::Relaxed);
}
#[inline]
pub fn request_processed(&self) {
self.total_requests.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn update_pipeline_depth(&self, depth: usize) {
self.max_pipeline_depth.fetch_max(depth, Ordering::Relaxed);
let current = self.avg_pipeline_depth.load(Ordering::Relaxed);
let new_avg = (current * 95 + (depth as u64 * 100) * 5) / 100;
self.avg_pipeline_depth.store(new_avg, Ordering::Relaxed);
}
#[inline]
pub fn active_connections(&self) -> usize {
self.active_connections.load(Ordering::Relaxed)
}
#[inline]
pub fn total_connections(&self) -> u64 {
self.total_connections.load(Ordering::Relaxed)
}
#[inline]
pub fn total_requests(&self) -> u64 {
self.total_requests.load(Ordering::Relaxed)
}
#[inline]
pub fn avg_pipeline_depth(&self) -> f64 {
self.avg_pipeline_depth.load(Ordering::Relaxed) as f64 / 100.0
}
#[inline]
pub fn max_pipeline_depth(&self) -> usize {
self.max_pipeline_depth.load(Ordering::Relaxed)
}
}
pub struct PipelinedHttp1Builder {
config: PipelineConfig,
stats: Arc<PipelineStats>,
}
impl PipelinedHttp1Builder {
pub fn new(config: PipelineConfig) -> Self {
Self {
config,
stats: Arc::new(PipelineStats::new()),
}
}
pub fn with_stats(config: PipelineConfig, stats: Arc<PipelineStats>) -> Self {
Self { config, stats }
}
pub fn config(&self) -> &PipelineConfig {
&self.config
}
pub fn stats(&self) -> Arc<PipelineStats> {
Arc::clone(&self.stats)
}
#[inline]
pub fn configure_hyper_builder(&self) -> hyper::server::conn::http1::Builder {
let mut builder = hyper::server::conn::http1::Builder::new();
builder.pipeline_flush(self.config.pipeline_flush);
builder.max_buf_size(self.config.read_buffer_size.max(8192));
builder.preserve_header_case(true);
builder.keep_alive(true);
builder
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_pipeline_mode_properties() {
assert!(PipelineMode::Sequential.maintains_order());
assert!(PipelineMode::Concurrent.maintains_order());
assert!(!PipelineMode::OutOfOrder.maintains_order());
assert!(!PipelineMode::Sequential.is_concurrent());
assert!(PipelineMode::Concurrent.is_concurrent());
assert!(PipelineMode::OutOfOrder.is_concurrent());
}
#[test]
fn test_config_builder() {
let config = PipelineConfig::builder()
.mode(PipelineMode::Concurrent)
.max_concurrent(32)
.pipeline_flush(true)
.keep_alive_timeout(Duration::from_secs(120))
.build();
assert_eq!(config.mode, PipelineMode::Concurrent);
assert_eq!(config.max_concurrent, 32);
assert!(config.pipeline_flush);
assert_eq!(config.keep_alive_timeout, Duration::from_secs(120));
}
#[test]
fn test_high_performance_config() {
let config = PipelineConfig::high_performance();
assert_eq!(config.mode, PipelineMode::Concurrent);
assert_eq!(config.max_concurrent, 32);
assert!(config.pipeline_flush);
}
#[test]
fn test_low_latency_config() {
let config = PipelineConfig::low_latency();
assert_eq!(config.mode, PipelineMode::Sequential);
assert_eq!(config.max_concurrent, 1);
assert!(!config.pipeline_flush);
}
#[test]
fn test_connection_stats() {
let stats = ConnectionStats::new();
stats.request_received(100);
assert_eq!(stats.pending_requests(), 1);
assert_eq!(stats.pipeline_depth(), 1);
assert_eq!(stats.bytes_received(), 100);
stats.request_received(200);
assert_eq!(stats.pending_requests(), 2);
assert_eq!(stats.pipeline_depth(), 2);
stats.response_sent(150);
assert_eq!(stats.pending_requests(), 1);
assert_eq!(stats.requests_processed(), 1);
assert_eq!(stats.bytes_sent(), 150);
}
#[test]
fn test_global_pipeline_stats() {
let stats = PipelineStats::new();
stats.connection_opened();
stats.connection_opened();
assert_eq!(stats.active_connections(), 2);
assert_eq!(stats.total_connections(), 2);
stats.connection_closed();
assert_eq!(stats.active_connections(), 1);
assert_eq!(stats.total_connections(), 2);
stats.request_processed();
stats.request_processed();
assert_eq!(stats.total_requests(), 2);
stats.update_pipeline_depth(5);
stats.update_pipeline_depth(10);
assert_eq!(stats.max_pipeline_depth(), 10);
}
#[test]
fn test_pipelined_builder() {
let config = PipelineConfig::default();
let builder = PipelinedHttp1Builder::new(config);
assert_eq!(builder.config().mode, PipelineMode::Concurrent);
assert_eq!(builder.stats().active_connections(), 0);
let _hyper_builder = builder.configure_hyper_builder();
}
}