use bytes::{Bytes, BytesMut};
use std::collections::VecDeque;
use std::io::IoSlice;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
#[derive(Debug)]
pub struct ResponseItem {
pub sequence: u64,
pub data: Bytes,
pub size: usize,
}
impl ResponseItem {
#[inline]
pub fn new(sequence: u64, data: Bytes) -> Self {
let size = data.len();
Self {
sequence,
data,
size,
}
}
#[inline]
pub fn from_vec(sequence: u64, data: Vec<u8>) -> Self {
Self::new(sequence, Bytes::from(data))
}
#[inline]
pub fn from_static(sequence: u64, data: &'static [u8]) -> Self {
Self::new(sequence, Bytes::from_static(data))
}
}
#[derive(Debug)]
pub struct ResponseQueue {
pending: Vec<Option<ResponseItem>>,
next_sequence: u64,
ready: VecDeque<ResponseItem>,
capacity: usize,
stats: ResponseQueueStats,
}
impl ResponseQueue {
pub fn new(capacity: usize) -> Self {
let mut pending = Vec::with_capacity(capacity);
pending.resize_with(capacity, || None);
Self {
pending,
next_sequence: 0,
ready: VecDeque::with_capacity(capacity),
capacity,
stats: ResponseQueueStats::new(),
}
}
#[inline]
pub fn push(&mut self, item: ResponseItem) -> bool {
if item.sequence < self.next_sequence {
self.stats.record_duplicate();
return false;
}
let offset = (item.sequence - self.next_sequence) as usize;
if offset >= self.capacity {
self.stats.record_overflow();
return false;
}
while self.pending.len() <= offset {
self.pending.push(None);
}
self.stats.record_push(item.size);
self.pending[offset] = Some(item);
self.promote_ready();
true
}
fn promote_ready(&mut self) {
while !self.pending.is_empty() {
if let Some(item) = self.pending[0].take() {
self.ready.push_back(item);
self.pending.remove(0);
self.next_sequence += 1;
} else {
break;
}
}
}
#[inline]
pub fn pop_ready(&mut self) -> Option<ResponseItem> {
self.ready.pop_front().inspect(|item| {
self.stats.record_pop(item.size);
})
}
#[inline]
pub fn peek_ready(&self) -> Option<&ResponseItem> {
self.ready.front()
}
#[inline]
pub fn has_ready(&self) -> bool {
!self.ready.is_empty()
}
#[inline]
pub fn ready_count(&self) -> usize {
self.ready.len()
}
#[inline]
pub fn pending_count(&self) -> usize {
self.pending.iter().filter(|x| x.is_some()).count()
}
#[inline]
pub fn next_sequence(&self) -> u64 {
self.next_sequence
}
#[inline]
pub fn ready_bytes(&self) -> usize {
self.ready.iter().map(|r| r.size).sum()
}
#[inline]
pub fn stats(&self) -> &ResponseQueueStats {
&self.stats
}
#[inline]
pub fn drain_ready(&mut self) -> ResponseBatch {
let items: Vec<_> = self.ready.drain(..).collect();
let total_size: usize = items.iter().map(|r| r.size).sum();
self.stats.record_batch_drain(items.len(), total_size);
ResponseBatch { items, total_size }
}
#[inline]
pub fn drain_ready_max(&mut self, max: usize) -> ResponseBatch {
let count = max.min(self.ready.len());
let items: Vec<_> = self.ready.drain(..count).collect();
let total_size: usize = items.iter().map(|r| r.size).sum();
self.stats.record_batch_drain(items.len(), total_size);
ResponseBatch { items, total_size }
}
pub fn clear(&mut self) {
self.pending.clear();
self.pending.resize_with(self.capacity, || None);
self.ready.clear();
}
pub fn reset(&mut self) {
self.clear();
self.next_sequence = 0;
}
}
#[derive(Debug)]
pub struct ResponseBatch {
pub items: Vec<ResponseItem>,
pub total_size: usize,
}
impl ResponseBatch {
#[inline]
pub fn is_empty(&self) -> bool {
self.items.is_empty()
}
#[inline]
pub fn len(&self) -> usize {
self.items.len()
}
#[inline]
pub fn as_io_slices(&self) -> Vec<IoSlice<'_>> {
self.items
.iter()
.map(|item| IoSlice::new(&item.data))
.collect()
}
#[inline]
pub fn concat(&self) -> Bytes {
if self.items.len() == 1 {
return self.items[0].data.clone();
}
let mut buf = BytesMut::with_capacity(self.total_size);
for item in &self.items {
buf.extend_from_slice(&item.data);
}
buf.freeze()
}
#[inline]
pub fn into_items(self) -> Vec<ResponseItem> {
self.items
}
}
#[derive(Debug, Default)]
pub struct ResponseQueueStats {
pushed: AtomicU64,
popped: AtomicU64,
bytes_pushed: AtomicU64,
bytes_popped: AtomicU64,
duplicates: AtomicU64,
overflows: AtomicU64,
batch_drains: AtomicU64,
batched_responses: AtomicU64,
}
impl ResponseQueueStats {
pub fn new() -> Self {
Self::default()
}
#[inline]
fn record_push(&self, size: usize) {
self.pushed.fetch_add(1, Ordering::Relaxed);
self.bytes_pushed.fetch_add(size as u64, Ordering::Relaxed);
}
#[inline]
fn record_pop(&self, size: usize) {
self.popped.fetch_add(1, Ordering::Relaxed);
self.bytes_popped.fetch_add(size as u64, Ordering::Relaxed);
}
#[inline]
fn record_duplicate(&self) {
self.duplicates.fetch_add(1, Ordering::Relaxed);
}
#[inline]
fn record_overflow(&self) {
self.overflows.fetch_add(1, Ordering::Relaxed);
}
#[inline]
fn record_batch_drain(&self, count: usize, _size: usize) {
self.batch_drains.fetch_add(1, Ordering::Relaxed);
self.batched_responses
.fetch_add(count as u64, Ordering::Relaxed);
}
pub fn pushed(&self) -> u64 {
self.pushed.load(Ordering::Relaxed)
}
pub fn popped(&self) -> u64 {
self.popped.load(Ordering::Relaxed)
}
pub fn bytes_pushed(&self) -> u64 {
self.bytes_pushed.load(Ordering::Relaxed)
}
pub fn bytes_popped(&self) -> u64 {
self.bytes_popped.load(Ordering::Relaxed)
}
pub fn duplicates(&self) -> u64 {
self.duplicates.load(Ordering::Relaxed)
}
pub fn overflows(&self) -> u64 {
self.overflows.load(Ordering::Relaxed)
}
pub fn batch_drains(&self) -> u64 {
self.batch_drains.load(Ordering::Relaxed)
}
pub fn avg_batch_size(&self) -> f64 {
let drains = self.batch_drains();
let batched = self.batched_responses.load(Ordering::Relaxed);
if drains > 0 {
batched as f64 / drains as f64
} else {
0.0
}
}
}
#[derive(Debug, Clone)]
pub struct ResponseWriterConfig {
pub min_batch_size: usize,
pub max_batch_size: usize,
pub min_batch_bytes: usize,
pub max_batch_bytes: usize,
pub use_tcp_cork: bool,
pub flush_timeout_us: u64,
}
impl Default for ResponseWriterConfig {
fn default() -> Self {
Self {
min_batch_size: 1,
max_batch_size: 64,
min_batch_bytes: 1024,
max_batch_bytes: 65536,
use_tcp_cork: true,
flush_timeout_us: 100,
}
}
}
impl ResponseWriterConfig {
pub fn high_throughput() -> Self {
Self {
min_batch_size: 8,
max_batch_size: 128,
min_batch_bytes: 8192,
max_batch_bytes: 131072,
use_tcp_cork: true,
flush_timeout_us: 500,
}
}
pub fn low_latency() -> Self {
Self {
min_batch_size: 1,
max_batch_size: 16,
min_batch_bytes: 512,
max_batch_bytes: 16384,
use_tcp_cork: false,
flush_timeout_us: 10,
}
}
}
#[derive(Debug, Default)]
pub struct ResponseWriterStats {
writes: AtomicU64,
bytes_written: AtomicU64,
vectored_writes: AtomicU64,
coalesced_writes: AtomicU64,
}
impl ResponseWriterStats {
pub fn new() -> Self {
Self::default()
}
#[inline]
pub fn record_write(&self, bytes: usize, vectored: bool) {
self.writes.fetch_add(1, Ordering::Relaxed);
self.bytes_written
.fetch_add(bytes as u64, Ordering::Relaxed);
if vectored {
self.vectored_writes.fetch_add(1, Ordering::Relaxed);
} else {
self.coalesced_writes.fetch_add(1, Ordering::Relaxed);
}
}
pub fn writes(&self) -> u64 {
self.writes.load(Ordering::Relaxed)
}
pub fn bytes_written(&self) -> u64 {
self.bytes_written.load(Ordering::Relaxed)
}
pub fn vectored_writes(&self) -> u64 {
self.vectored_writes.load(Ordering::Relaxed)
}
pub fn coalesced_writes(&self) -> u64 {
self.coalesced_writes.load(Ordering::Relaxed)
}
pub fn vectored_ratio(&self) -> f64 {
let total = self.writes();
let vectored = self.vectored_writes();
if total > 0 {
(vectored as f64 / total as f64) * 100.0
} else {
0.0
}
}
}
static WRITER_STATS: ResponseWriterStats = ResponseWriterStats {
writes: AtomicU64::new(0),
bytes_written: AtomicU64::new(0),
vectored_writes: AtomicU64::new(0),
coalesced_writes: AtomicU64::new(0),
};
pub fn writer_stats() -> &'static ResponseWriterStats {
&WRITER_STATS
}
#[derive(Debug)]
pub struct ConnectionPipeline {
queue: ResponseQueue,
config: ResponseWriterConfig,
sequence_counter: u64,
#[allow(dead_code)]
connection_id: u64,
}
impl ConnectionPipeline {
pub fn new(connection_id: u64, queue_capacity: usize) -> Self {
Self {
queue: ResponseQueue::new(queue_capacity),
config: ResponseWriterConfig::default(),
sequence_counter: 0,
connection_id,
}
}
pub fn with_config(
connection_id: u64,
queue_capacity: usize,
config: ResponseWriterConfig,
) -> Self {
Self {
queue: ResponseQueue::new(queue_capacity),
config,
sequence_counter: 0,
connection_id,
}
}
#[inline]
pub fn next_sequence(&mut self) -> u64 {
let seq = self.sequence_counter;
self.sequence_counter += 1;
seq
}
#[inline]
pub fn queue_response(&mut self, sequence: u64, data: Bytes) -> bool {
self.queue.push(ResponseItem::new(sequence, data))
}
#[inline]
pub fn should_flush(&self) -> bool {
let count = self.queue.ready_count();
let bytes = self.queue.ready_bytes();
count >= self.config.min_batch_size
|| bytes >= self.config.min_batch_bytes
|| count >= self.config.max_batch_size
|| bytes >= self.config.max_batch_bytes
}
#[inline]
pub fn must_flush(&self) -> bool {
let count = self.queue.ready_count();
let bytes = self.queue.ready_bytes();
count >= self.config.max_batch_size || bytes >= self.config.max_batch_bytes
}
#[inline]
pub fn drain_for_send(&mut self) -> ResponseBatch {
self.queue.drain_ready()
}
#[inline]
pub fn drain_max(&mut self, max: usize) -> ResponseBatch {
self.queue.drain_ready_max(max)
}
#[inline]
pub fn queue_stats(&self) -> &ResponseQueueStats {
self.queue.stats()
}
#[inline]
pub fn has_pending(&self) -> bool {
self.queue.pending_count() > 0
}
#[inline]
pub fn has_ready(&self) -> bool {
self.queue.has_ready()
}
#[inline]
pub fn ready_count(&self) -> usize {
self.queue.ready_count()
}
#[inline]
pub fn pending_count(&self) -> usize {
self.queue.pending_count()
}
#[inline]
pub fn ready_bytes(&self) -> usize {
self.queue.ready_bytes()
}
#[inline]
pub fn config(&self) -> &ResponseWriterConfig {
&self.config
}
pub fn reset(&mut self) {
self.queue.reset();
self.sequence_counter = 0;
}
}
#[derive(Debug, Default)]
pub struct GlobalPipelineStats {
responses_queued: AtomicU64,
responses_sent: AtomicU64,
batches_sent: AtomicU64,
bytes_sent: AtomicU64,
max_batch_size: AtomicUsize,
reordered: AtomicU64,
}
impl GlobalPipelineStats {
pub fn new() -> Self {
Self::default()
}
#[inline]
pub fn record_queued(&self, count: usize) {
self.responses_queued
.fetch_add(count as u64, Ordering::Relaxed);
}
#[inline]
pub fn record_batch_sent(&self, responses: usize, bytes: usize) {
self.responses_sent
.fetch_add(responses as u64, Ordering::Relaxed);
self.batches_sent.fetch_add(1, Ordering::Relaxed);
self.bytes_sent.fetch_add(bytes as u64, Ordering::Relaxed);
self.max_batch_size.fetch_max(responses, Ordering::Relaxed);
}
#[inline]
pub fn record_reorder(&self) {
self.reordered.fetch_add(1, Ordering::Relaxed);
}
pub fn responses_queued(&self) -> u64 {
self.responses_queued.load(Ordering::Relaxed)
}
pub fn responses_sent(&self) -> u64 {
self.responses_sent.load(Ordering::Relaxed)
}
pub fn batches_sent(&self) -> u64 {
self.batches_sent.load(Ordering::Relaxed)
}
pub fn bytes_sent(&self) -> u64 {
self.bytes_sent.load(Ordering::Relaxed)
}
pub fn max_batch_size(&self) -> usize {
self.max_batch_size.load(Ordering::Relaxed)
}
pub fn reordered(&self) -> u64 {
self.reordered.load(Ordering::Relaxed)
}
pub fn avg_batch_size(&self) -> f64 {
let batches = self.batches_sent();
let responses = self.responses_sent();
if batches > 0 {
responses as f64 / batches as f64
} else {
0.0
}
}
pub fn avg_response_size(&self) -> f64 {
let responses = self.responses_sent();
let bytes = self.bytes_sent();
if responses > 0 {
bytes as f64 / responses as f64
} else {
0.0
}
}
}
static GLOBAL_PIPELINE_STATS: GlobalPipelineStats = GlobalPipelineStats {
responses_queued: AtomicU64::new(0),
responses_sent: AtomicU64::new(0),
batches_sent: AtomicU64::new(0),
bytes_sent: AtomicU64::new(0),
max_batch_size: AtomicUsize::new(0),
reordered: AtomicU64::new(0),
};
pub fn global_pipeline_stats() -> &'static GlobalPipelineStats {
&GLOBAL_PIPELINE_STATS
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_response_item_creation() {
let item = ResponseItem::new(0, Bytes::from("hello"));
assert_eq!(item.sequence, 0);
assert_eq!(item.size, 5);
assert_eq!(&item.data[..], b"hello");
let item2 = ResponseItem::from_vec(1, vec![1, 2, 3]);
assert_eq!(item2.sequence, 1);
assert_eq!(item2.size, 3);
let item3 = ResponseItem::from_static(2, b"static");
assert_eq!(item3.sequence, 2);
assert_eq!(item3.size, 6);
}
#[test]
fn test_response_queue_in_order() {
let mut queue = ResponseQueue::new(16);
assert!(queue.push(ResponseItem::new(0, Bytes::from("first"))));
assert!(queue.push(ResponseItem::new(1, Bytes::from("second"))));
assert!(queue.push(ResponseItem::new(2, Bytes::from("third"))));
assert_eq!(queue.ready_count(), 3);
assert_eq!(queue.pending_count(), 0);
let item0 = queue.pop_ready().unwrap();
assert_eq!(item0.sequence, 0);
let item1 = queue.pop_ready().unwrap();
assert_eq!(item1.sequence, 1);
let item2 = queue.pop_ready().unwrap();
assert_eq!(item2.sequence, 2);
assert!(queue.pop_ready().is_none());
}
#[test]
fn test_response_queue_out_of_order() {
let mut queue = ResponseQueue::new(16);
assert!(queue.push(ResponseItem::new(2, Bytes::from("third"))));
assert_eq!(queue.ready_count(), 0);
assert_eq!(queue.pending_count(), 1);
assert!(queue.push(ResponseItem::new(0, Bytes::from("first"))));
assert_eq!(queue.ready_count(), 1); assert_eq!(queue.pending_count(), 1);
assert!(queue.push(ResponseItem::new(1, Bytes::from("second"))));
assert_eq!(queue.ready_count(), 3); assert_eq!(queue.pending_count(), 0);
let item0 = queue.pop_ready().unwrap();
assert_eq!(item0.sequence, 0);
let item1 = queue.pop_ready().unwrap();
assert_eq!(item1.sequence, 1);
let item2 = queue.pop_ready().unwrap();
assert_eq!(item2.sequence, 2);
}
#[test]
fn test_response_queue_duplicate_rejected() {
let mut queue = ResponseQueue::new(16);
assert!(queue.push(ResponseItem::new(0, Bytes::from("first"))));
queue.pop_ready();
assert!(!queue.push(ResponseItem::new(0, Bytes::from("duplicate"))));
assert_eq!(queue.stats().duplicates(), 1);
}
#[test]
fn test_response_queue_overflow_rejected() {
let mut queue = ResponseQueue::new(4);
assert!(!queue.push(ResponseItem::new(10, Bytes::from("too far"))));
assert_eq!(queue.stats().overflows(), 1);
}
#[test]
fn test_response_batch_drain() {
let mut queue = ResponseQueue::new(16);
queue.push(ResponseItem::new(0, Bytes::from("aaaa")));
queue.push(ResponseItem::new(1, Bytes::from("bbbb")));
queue.push(ResponseItem::new(2, Bytes::from("cccc")));
let batch = queue.drain_ready();
assert_eq!(batch.len(), 3);
assert_eq!(batch.total_size, 12);
assert_eq!(queue.ready_count(), 0);
}
#[test]
fn test_response_batch_concat() {
let mut queue = ResponseQueue::new(16);
queue.push(ResponseItem::new(0, Bytes::from("aa")));
queue.push(ResponseItem::new(1, Bytes::from("bb")));
queue.push(ResponseItem::new(2, Bytes::from("cc")));
let batch = queue.drain_ready();
let concat = batch.concat();
assert_eq!(&concat[..], b"aabbcc");
}
#[test]
fn test_response_batch_io_slices() {
let mut queue = ResponseQueue::new(16);
queue.push(ResponseItem::new(0, Bytes::from("hello")));
queue.push(ResponseItem::new(1, Bytes::from("world")));
let batch = queue.drain_ready();
let slices = batch.as_io_slices();
assert_eq!(slices.len(), 2);
assert_eq!(&*slices[0], b"hello");
assert_eq!(&*slices[1], b"world");
}
#[test]
fn test_connection_pipeline() {
let mut pipeline = ConnectionPipeline::new(1, 32);
let seq0 = pipeline.next_sequence();
let seq1 = pipeline.next_sequence();
let seq2 = pipeline.next_sequence();
assert_eq!(seq0, 0);
assert_eq!(seq1, 1);
assert_eq!(seq2, 2);
assert!(pipeline.queue_response(seq2, Bytes::from("c")));
assert!(pipeline.queue_response(seq0, Bytes::from("a")));
assert!(pipeline.queue_response(seq1, Bytes::from("b")));
assert_eq!(pipeline.ready_count(), 3);
assert_eq!(pipeline.ready_bytes(), 3);
}
#[test]
fn test_connection_pipeline_should_flush() {
let config = ResponseWriterConfig {
min_batch_size: 2,
max_batch_size: 10,
min_batch_bytes: 100,
max_batch_bytes: 1000,
use_tcp_cork: true,
flush_timeout_us: 100,
};
let mut pipeline = ConnectionPipeline::with_config(1, 32, config);
pipeline.queue_response(0, Bytes::from("hello"));
assert!(!pipeline.should_flush());
pipeline.queue_response(1, Bytes::from("world"));
assert!(pipeline.should_flush());
}
#[test]
fn test_connection_pipeline_must_flush() {
let config = ResponseWriterConfig {
min_batch_size: 1,
max_batch_size: 2,
min_batch_bytes: 1,
max_batch_bytes: 100,
use_tcp_cork: true,
flush_timeout_us: 100,
};
let mut pipeline = ConnectionPipeline::with_config(1, 32, config);
pipeline.queue_response(0, Bytes::from("a"));
assert!(!pipeline.must_flush());
pipeline.queue_response(1, Bytes::from("b"));
assert!(pipeline.must_flush()); }
#[test]
fn test_response_writer_config_presets() {
let high = ResponseWriterConfig::high_throughput();
assert_eq!(high.min_batch_size, 8);
assert!(high.use_tcp_cork);
let low = ResponseWriterConfig::low_latency();
assert_eq!(low.min_batch_size, 1);
assert!(!low.use_tcp_cork);
}
#[test]
fn test_response_queue_stats() {
let mut queue = ResponseQueue::new(16);
queue.push(ResponseItem::new(0, Bytes::from("hello")));
queue.push(ResponseItem::new(1, Bytes::from("world")));
assert_eq!(queue.stats().pushed(), 2);
assert_eq!(queue.stats().bytes_pushed(), 10);
queue.pop_ready();
assert_eq!(queue.stats().popped(), 1);
}
#[test]
fn test_global_pipeline_stats() {
let stats = global_pipeline_stats();
let _ = stats.responses_queued();
let _ = stats.responses_sent();
let _ = stats.avg_batch_size();
}
}