use scirs2_core::ndarray::Array1;
use scirs2_core::numeric::Float;
use std::collections::VecDeque;
use std::sync::{
atomic::{AtomicUsize, Ordering},
Arc, Mutex, MutexGuard, PoisonError,
};
use std::time::{Duration, Instant};
use crate::error::{OptimError, Result};
use crate::optimizers::Optimizer;
#[cfg(test)]
mod regression_tests;
fn lock_recovered<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
#[derive(Debug, Clone)]
pub struct LowLatencyConfig {
pub target_latency_us: u64,
pub max_latency_us: u64,
pub enable_precomputation: bool,
pub precomputation_buffer_size: usize,
pub enable_lock_free: bool,
pub use_approximations: bool,
pub approximation_tolerance: f64,
pub enable_simd: bool,
pub batch_threshold: usize,
pub enable_zero_copy: bool,
pub memory_pool_size: usize,
pub enable_quantization: bool,
pub quantization_bits: u8,
}
impl Default for LowLatencyConfig {
fn default() -> Self {
Self {
target_latency_us: 100, max_latency_us: 1000, enable_precomputation: true,
precomputation_buffer_size: 64,
enable_lock_free: true,
use_approximations: true,
approximation_tolerance: 0.01,
enable_simd: true,
batch_threshold: 8,
enable_zero_copy: true,
memory_pool_size: 1024 * 1024, enable_quantization: false,
quantization_bits: 8,
}
}
}
pub struct LowLatencyOptimizer<O, A>
where
A: Float + Send + Sync + scirs2_core::ndarray::ScalarOperand + std::fmt::Debug,
O: Optimizer<A, scirs2_core::ndarray::Ix1> + Send + Sync,
{
base_optimizer: Arc<Mutex<O>>,
config: LowLatencyConfig,
parameters: Option<Array1<A>>,
precomputation_engine: Option<PrecomputationEngine<A>>,
update_buffer: LockFreeBuffer<A>,
memory_pool: FastMemoryPool<A>,
simd_processor: SIMDProcessor<A>,
quantizer: Option<GradientQuantizer<A>>,
perf_monitor: LatencyMonitor,
approximation_controller: ApproximationController<A>,
step_counter: AtomicUsize,
}
struct PrecomputationEngine<A: Float + Send + Sync> {
precomputed_updates: VecDeque<PrecomputedUpdate<A>>,
gradient_predictor: GradientPredictor<A>,
max_buffer_size: usize,
hits: usize,
misses: usize,
min_confidence: A,
}
#[derive(Debug, Clone)]
struct PrecomputedUpdate<A: Float + Send + Sync> {
gradient: Array1<A>,
update: Array1<A>,
valid_until: Instant,
confidence: A,
}
struct LockFreeBuffer<A: Float + Send + Sync> {
buffer: Vec<Option<Array1<A>>>,
write_index: AtomicUsize,
read_index: AtomicUsize,
capacity: usize,
}
struct FastMemoryPool<A> {
free_blocks: Mutex<Vec<Vec<A>>>,
elements_per_block: usize,
total_blocks: usize,
checked_out: AtomicUsize,
peak_checked_out: AtomicUsize,
misses: AtomicUsize,
}
struct SIMDProcessor<A: Float + Send + Sync> {
enabled: bool,
vector_width: usize,
_element: std::marker::PhantomData<A>,
}
struct GradientQuantizer<A: Float + Send + Sync> {
bits: u8,
scale: A,
zero_point: A,
error_accumulator: Option<Array1<A>>,
}
#[derive(Debug)]
struct LatencyMonitor {
latency_samples: VecDeque<Duration>,
maxsamples: usize,
p50_latency: Duration,
p95_latency: Duration,
p99_latency: Duration,
violations: usize,
total_operations: usize,
}
const PERFORMANCE_WINDOW_AGE: Duration = Duration::from_secs(30);
const PERFORMANCE_WINDOW_LEN: usize = 100;
struct ApproximationController<A: Float + Send + Sync> {
approximation_level: A,
performance_history: VecDeque<PerformancePoint<A>>,
adaptation_rate: A,
targetlatency: Duration,
}
#[derive(Debug, Clone)]
struct PerformancePoint<A: Float + Send + Sync> {
latency: Duration,
accuracy: A,
timestamp: Instant,
}
struct GradientPredictor<A: Float + Send + Sync> {
gradient_history: VecDeque<Array1<A>>,
trend_weights: Option<Array1<A>>,
windowsize: usize,
confidence: Option<A>,
pending_prediction: Option<Array1<A>>,
}
impl<O, A> LowLatencyOptimizer<O, A>
where
A: Float
+ Send
+ Sync
+ Default
+ Clone
+ std::fmt::Debug
+ scirs2_core::ndarray::ScalarOperand
+ 'static
+ std::iter::Sum,
O: Optimizer<A, scirs2_core::ndarray::Ix1> + Send + Sync + 'static,
{
pub fn new(_baseoptimizer: O, config: LowLatencyConfig) -> Result<Self> {
let base_optimizer = Arc::new(Mutex::new(_baseoptimizer));
let precomputation_engine = if config.enable_precomputation {
Some(PrecomputationEngine::new(config.precomputation_buffer_size))
} else {
None
};
let update_buffer = LockFreeBuffer::new(config.precomputation_buffer_size);
let memory_pool = FastMemoryPool::new(config.memory_pool_size, 4096)?; let simd_processor = SIMDProcessor::new(config.enable_simd, config.batch_threshold);
let quantizer = if config.enable_quantization {
Some(GradientQuantizer::new(config.quantization_bits))
} else {
None
};
let perf_monitor = LatencyMonitor::new(1000); let approximation_controller =
ApproximationController::new(Duration::from_micros(config.target_latency_us));
Ok(Self {
base_optimizer,
config,
parameters: None,
precomputation_engine,
update_buffer,
memory_pool,
simd_processor,
quantizer,
perf_monitor,
approximation_controller,
step_counter: AtomicUsize::new(0),
})
}
pub fn set_parameters(&mut self, parameters: Array1<A>) {
self.parameters = Some(parameters);
}
pub fn set_precomputation_min_confidence(&mut self, min_confidence: A) {
if let Some(precomp) = self.precomputation_engine.as_mut() {
precomp.set_min_confidence(min_confidence);
}
}
pub fn parameters(&self) -> Option<&Array1<A>> {
self.parameters.as_ref()
}
pub fn low_latency_step(&mut self, gradient: &Array1<A>) -> Result<Array1<A>> {
let start_time = Instant::now();
if gradient.is_empty() {
return Err(OptimError::DimensionMismatch(
"low_latency_step received an empty gradient".to_string(),
));
}
let previous_params = self.parameters.clone();
let learning_rate = self.base_learning_rate();
let tolerance = self.config.approximation_tolerance.max(0.0);
let served = self
.precomputation_engine
.as_mut()
.and_then(|precomp| precomp.try_get_precomputed(gradient, tolerance));
if let Some(precomputed) = served {
let update = precomputed.update;
self.parameters = Some(update.clone());
let latency = start_time.elapsed();
self.perf_monitor.record_latency(latency);
if self.config.enable_lock_free {
self.update_buffer.push(update.clone());
}
let validity = Duration::from_micros(self.config.max_latency_us.max(1));
if let Some(ref mut precomp) = self.precomputation_engine {
precomp.start_precomputation(gradient, &update, learning_rate, validity);
}
self.step_counter.fetch_add(1, Ordering::Relaxed);
return Ok(update);
}
let quantized = match self.quantizer.as_mut() {
Some(quantizer) => Some(quantizer.quantize(gradient)?),
None if self.config.enable_zero_copy => None,
None => Some(gradient.clone()),
};
let processed_gradient: &Array1<A> = quantized.as_ref().unwrap_or(gradient);
let approximation_level = self.approximation_controller.get_approximation_level();
let use_approximation = self.config.use_approximations && approximation_level > A::zero();
let update = if use_approximation {
let simplified = self.simplify_gradient(processed_gradient, approximation_level)?;
self.fast_path_update(&simplified, learning_rate)?
} else {
self.exact_update(processed_gradient)?
};
let latency = start_time.elapsed();
let accuracy = Self::estimate_accuracy(previous_params.as_ref(), &update, gradient);
self.approximation_controller
.record_performance(latency, approximation_level, accuracy);
self.perf_monitor.record_latency(latency);
if latency.as_micros() as u64 > self.config.max_latency_us {
self.handle_latency_violation(latency)?;
}
if self.config.enable_lock_free {
self.update_buffer.push(update.clone());
}
let validity = Duration::from_micros(self.config.max_latency_us.max(1));
if let Some(ref mut precomp) = self.precomputation_engine {
precomp.start_precomputation(gradient, &update, learning_rate, validity);
}
self.parameters = Some(update.clone());
self.step_counter.fetch_add(1, Ordering::Relaxed);
Ok(update)
}
fn base_learning_rate(&self) -> A {
lock_recovered(&self.base_optimizer).get_learning_rate()
}
fn current_parameters(&self, len: usize) -> Result<Array1<A>> {
match self.parameters.as_ref() {
Some(params) if params.len() == len => Ok(params.clone()),
Some(params) => Err(OptimError::DimensionMismatch(format!(
"gradient has {} elements but the tracked parameters have {}",
len,
params.len()
))),
None => Ok(Array1::zeros(len)),
}
}
fn exact_update(&mut self, gradient: &Array1<A>) -> Result<Array1<A>> {
let current_params = self.current_parameters(gradient.len())?;
let mut optimizer = lock_recovered(&self.base_optimizer);
optimizer.step(¤t_params, gradient)
}
fn fast_path_update(&mut self, gradient: &Array1<A>, learning_rate: A) -> Result<Array1<A>> {
if !self.simd_processor.is_active(gradient.len()) {
return self.exact_update(gradient);
}
let mut params = self.current_parameters(gradient.len())?;
self.simd_processor
.apply_scaled_subtract(&mut params, gradient, learning_rate);
Ok(params)
}
fn simplify_gradient(&self, gradient: &Array1<A>, level: A) -> Result<Array1<A>> {
let n = gradient.len();
if n == 0 {
return Ok(gradient.clone());
}
let sparsity_ratio = level.to_f64().unwrap_or(0.0).clamp(0.0, 1.0);
let keep_ratio = 1.0 - sparsity_ratio * 0.8; let keep_count = (((n as f64) * keep_ratio).round() as usize).clamp(1, n);
if keep_count == n {
return Ok(gradient.clone());
}
let mut magnitudes = self
.memory_pool
.acquire(n)
.unwrap_or_else(|| Vec::with_capacity(n));
magnitudes.clear();
magnitudes.extend(gradient.iter().map(|g| g.abs()));
let kth = keep_count - 1;
magnitudes.select_nth_unstable_by(kth, |a, b| {
b.partial_cmp(a).unwrap_or(std::cmp::Ordering::Equal)
});
let threshold = magnitudes[kth];
self.memory_pool.release(magnitudes);
let mut simplified = Array1::zeros(n);
let mut kept = 0usize;
for (i, &g) in gradient.iter().enumerate() {
if kept < keep_count && g.abs() >= threshold {
simplified[i] = g;
kept += 1;
}
}
Ok(simplified)
}
fn estimate_accuracy(
previous_params: Option<&Array1<A>>,
new_params: &Array1<A>,
gradient: &Array1<A>,
) -> A {
if new_params.len() != gradient.len() {
return A::zero();
}
let zeros = Array1::zeros(new_params.len());
let previous = match previous_params {
Some(previous) if previous.len() == new_params.len() => previous,
_ => &zeros,
};
let mut dot = A::zero();
let mut norm_delta = A::zero();
let mut norm_grad = A::zero();
for ((&p_new, &p_old), &g) in new_params.iter().zip(previous.iter()).zip(gradient.iter()) {
let delta = p_new - p_old;
dot = dot + delta * (-g);
norm_delta = norm_delta + delta * delta;
norm_grad = norm_grad + g * g;
}
let norm_delta = norm_delta.sqrt();
let norm_grad = norm_grad.sqrt();
if norm_delta == A::zero() || norm_grad == A::zero() {
A::zero()
} else {
dot / (norm_delta * norm_grad)
}
}
fn handle_latency_violation(&mut self, latency: Duration) -> Result<()> {
self.perf_monitor.violations += 1;
self.approximation_controller.increase_approximation();
if !self.config.enable_quantization
&& latency.as_micros() as u64 > self.config.max_latency_us.saturating_mul(2)
{
self.config.enable_quantization = true;
self.quantizer = Some(GradientQuantizer::new(self.config.quantization_bits));
}
Ok(())
}
pub fn try_pop_staged_update(&mut self) -> Option<Array1<A>> {
self.update_buffer.pop()
}
pub fn staged_update_count(&self) -> usize {
self.update_buffer.len()
}
pub fn get_performance_metrics(&self) -> LowLatencyMetrics {
LowLatencyMetrics {
avg_latency_us: self.perf_monitor.get_average_latency().as_micros() as u64,
p50_latency_us: self.perf_monitor.p50_latency.as_micros() as u64,
p95_latency_us: self.perf_monitor.p95_latency.as_micros() as u64,
p99_latency_us: self.perf_monitor.p99_latency.as_micros() as u64,
latency_violations: self.perf_monitor.violations,
total_operations: self.perf_monitor.total_operations,
current_approximation_level: self
.approximation_controller
.approximation_level
.to_f64()
.unwrap_or(0.0),
approximation_accuracy: self
.approximation_controller
.mean_accuracy()
.and_then(|value| value.to_f64()),
precomputation_hit_rate: self
.precomputation_engine
.as_ref()
.and_then(|pe| pe.hit_rate()),
precomputation_attempts: self
.precomputation_engine
.as_ref()
.map(|pe| pe.attempts())
.unwrap_or(0),
memory_efficiency: self.memory_pool.get_efficiency(),
memory_pool_misses: self.memory_pool.misses(),
}
}
pub fn is_meeting_latency_requirements(&self) -> bool {
let avg_latency = self.perf_monitor.get_average_latency().as_micros() as u64;
avg_latency <= self.config.target_latency_us
}
}
impl<A: Float + Send + Sync + std::iter::Sum> PrecomputationEngine<A> {
fn new(_buffersize: usize) -> Self {
let capacity = _buffersize.max(1);
Self {
precomputed_updates: VecDeque::with_capacity(capacity),
gradient_predictor: GradientPredictor::new(10), max_buffer_size: capacity,
hits: 0,
misses: 0,
min_confidence: A::zero(),
}
}
fn set_min_confidence(&mut self, min_confidence: A) {
self.min_confidence = min_confidence;
}
fn try_get_precomputed(
&mut self,
actual_gradient: &Array1<A>,
tolerance: f64,
) -> Option<PrecomputedUpdate<A>> {
let now = Instant::now();
while let Some(update) = self.precomputed_updates.front() {
if update.valid_until <= now {
self.precomputed_updates.pop_front();
} else {
break;
}
}
let candidate = self.precomputed_updates.pop_front();
self.gradient_predictor.observe(actual_gradient);
match candidate {
Some(candidate)
if candidate.confidence >= self.min_confidence
&& gradient_matches(&candidate.gradient, actual_gradient, tolerance) =>
{
self.hits += 1;
Some(candidate)
}
_ => {
self.misses += 1;
None
}
}
}
fn start_precomputation(
&mut self,
_observed_gradient: &Array1<A>,
current_params: &Array1<A>,
learning_rate: A,
validity: Duration,
) {
let Some((predicted, confidence)) = self.gradient_predictor.predict() else {
return;
};
if predicted.len() != current_params.len() {
return;
}
let mut update = current_params.clone();
for (p, &g) in update.iter_mut().zip(predicted.iter()) {
*p = *p - learning_rate * g;
}
if self.precomputed_updates.len() >= self.max_buffer_size {
self.precomputed_updates.pop_front();
}
self.precomputed_updates.push_back(PrecomputedUpdate {
gradient: predicted,
update,
valid_until: Instant::now() + validity,
confidence,
});
}
fn attempts(&self) -> usize {
self.hits + self.misses
}
fn hit_rate(&self) -> Option<f64> {
let attempts = self.attempts();
if attempts == 0 {
None
} else {
Some(self.hits as f64 / attempts as f64)
}
}
}
fn gradient_matches<A: Float>(predicted: &Array1<A>, actual: &Array1<A>, tolerance: f64) -> bool {
if predicted.len() != actual.len() || predicted.is_empty() {
return false;
}
let mut diff_sq = A::zero();
let mut actual_sq = A::zero();
for (&p, &a) in predicted.iter().zip(actual.iter()) {
let d = p - a;
diff_sq = diff_sq + d * d;
actual_sq = actual_sq + a * a;
}
let diff = diff_sq.sqrt().to_f64().unwrap_or(f64::INFINITY);
let scale = actual_sq.sqrt().to_f64().unwrap_or(0.0);
if !diff.is_finite() {
return false;
}
if scale <= f64::EPSILON {
diff <= tolerance
} else {
diff / scale <= tolerance
}
}
impl<A: Float + Send + Sync> LockFreeBuffer<A> {
fn new(capacity: usize) -> Self {
let capacity = capacity.max(1);
Self {
buffer: vec![None; capacity],
write_index: AtomicUsize::new(0),
read_index: AtomicUsize::new(0),
capacity,
}
}
fn push(&mut self, value: Array1<A>) {
let write = self.write_index.load(Ordering::Relaxed);
let read = self.read_index.load(Ordering::Relaxed);
if write - read >= self.capacity {
let slot = read % self.capacity;
self.buffer[slot] = None;
self.read_index.store(read + 1, Ordering::Relaxed);
}
let slot = write % self.capacity;
self.buffer[slot] = Some(value);
self.write_index.store(write + 1, Ordering::Relaxed);
}
fn pop(&mut self) -> Option<Array1<A>> {
let read = self.read_index.load(Ordering::Relaxed);
if read == self.write_index.load(Ordering::Relaxed) {
return None;
}
let slot = read % self.capacity;
let value = self.buffer[slot].take();
self.read_index.store(read + 1, Ordering::Relaxed);
value
}
fn len(&self) -> usize {
self.write_index.load(Ordering::Relaxed) - self.read_index.load(Ordering::Relaxed)
}
}
impl<A: Float> FastMemoryPool<A> {
fn new(_total_size: usize, block_size_bytes: usize) -> Result<Self> {
let element_size = std::mem::size_of::<A>().max(1);
let elements_per_block = (block_size_bytes / element_size).max(1);
let total_blocks = _total_size / block_size_bytes.max(1);
let mut free_blocks = Vec::with_capacity(total_blocks);
for _ in 0..total_blocks {
free_blocks.push(Vec::with_capacity(elements_per_block));
}
Ok(Self {
free_blocks: Mutex::new(free_blocks),
elements_per_block,
total_blocks,
checked_out: AtomicUsize::new(0),
peak_checked_out: AtomicUsize::new(0),
misses: AtomicUsize::new(0),
})
}
fn acquire(&self, len: usize) -> Option<Vec<A>> {
if len > self.elements_per_block {
self.misses.fetch_add(1, Ordering::Relaxed);
return None;
}
let block = lock_recovered(&self.free_blocks).pop();
match block {
Some(mut block) => {
block.clear();
let in_use = self.checked_out.fetch_add(1, Ordering::Relaxed) + 1;
self.peak_checked_out.fetch_max(in_use, Ordering::Relaxed);
Some(block)
}
None => {
self.misses.fetch_add(1, Ordering::Relaxed);
None
}
}
}
fn release(&self, mut block: Vec<A>) {
if block.capacity() < self.elements_per_block {
return;
}
block.clear();
let mut free = lock_recovered(&self.free_blocks);
if free.len() < self.total_blocks {
free.push(block);
drop(free);
let previous = self.checked_out.load(Ordering::Relaxed);
if previous > 0 {
self.checked_out.store(previous - 1, Ordering::Relaxed);
}
}
}
fn get_efficiency(&self) -> f64 {
if self.total_blocks == 0 {
return 0.0;
}
self.peak_checked_out.load(Ordering::Relaxed) as f64 / self.total_blocks as f64
}
fn misses(&self) -> usize {
self.misses.load(Ordering::Relaxed)
}
}
impl<A: Float + Send + Sync> SIMDProcessor<A> {
fn new(enabled: bool, batch_threshold: usize) -> Self {
Self {
enabled,
vector_width: batch_threshold.max(1),
_element: std::marker::PhantomData,
}
}
fn is_active(&self, len: usize) -> bool {
self.enabled && len >= self.vector_width
}
fn apply_scaled_subtract(
&self,
params: &mut Array1<A>,
gradient: &Array1<A>,
learning_rate: A,
) {
let width = self.vector_width.max(1);
match (params.as_slice_mut(), gradient.as_slice()) {
(Some(p), Some(g)) => {
for (p_chunk, g_chunk) in p.chunks_mut(width).zip(g.chunks(width)) {
for (p_value, g_value) in p_chunk.iter_mut().zip(g_chunk.iter()) {
*p_value = *p_value - learning_rate * *g_value;
}
}
}
_ => {
for (p_value, g_value) in params.iter_mut().zip(gradient.iter()) {
*p_value = *p_value - learning_rate * *g_value;
}
}
}
}
}
impl<A: Float + Send + Sync> GradientQuantizer<A> {
fn new(bits: u8) -> Self {
Self {
bits: bits.clamp(1, 24),
scale: A::one(),
zero_point: A::zero(),
error_accumulator: None,
}
}
fn quantize(&mut self, gradient: &Array1<A>) -> Result<Array1<A>> {
let n = gradient.len();
if n == 0 {
return Ok(gradient.clone());
}
let compensated = match self.error_accumulator.as_ref() {
Some(error) if error.len() == n => gradient + error,
_ => gradient.clone(),
};
if compensated.iter().any(|value| !value.is_finite()) {
return Err(OptimError::InvalidParameter(
"cannot quantize a gradient containing non-finite values".to_string(),
));
}
let max_abs = compensated.iter().fold(
A::zero(),
|acc, x| if x.abs() > acc { x.abs() } else { acc },
);
self.zero_point = A::zero(); if max_abs == A::zero() {
self.scale = A::one();
self.error_accumulator = Some(Array1::zeros(n));
return Ok(compensated);
}
let level_count = (1u32 << (self.bits.max(1) as u32 - 1))
.saturating_sub(1)
.max(1);
let levels = A::from(level_count).unwrap_or(A::one());
self.scale = max_abs / levels;
let scale = self.scale;
let zero_point = self.zero_point;
let quantized = compensated.mapv(|x| {
let mut q = (x / scale).round();
if q > levels {
q = levels;
} else if q < -levels {
q = -levels;
}
q * scale + zero_point
});
self.error_accumulator = Some(&compensated - &quantized);
Ok(quantized)
}
}
impl LatencyMonitor {
fn new(maxsamples: usize) -> Self {
Self {
latency_samples: VecDeque::with_capacity(maxsamples),
maxsamples: maxsamples.max(1),
p50_latency: Duration::from_micros(0),
p95_latency: Duration::from_micros(0),
p99_latency: Duration::from_micros(0),
violations: 0,
total_operations: 0,
}
}
fn record_latency(&mut self, latency: Duration) {
self.latency_samples.push_back(latency);
if self.latency_samples.len() > self.maxsamples {
self.latency_samples.pop_front();
}
self.total_operations += 1;
self.update_percentiles();
}
fn update_percentiles(&mut self) {
if self.latency_samples.is_empty() {
return;
}
let mut sorted: Vec<_> = self.latency_samples.iter().cloned().collect();
sorted.sort();
let last = sorted.len() - 1;
let index_for = |q: f64| ((sorted.len() as f64 * q) as usize).min(last);
self.p50_latency = sorted[index_for(0.50)];
self.p95_latency = sorted[index_for(0.95)];
self.p99_latency = sorted[index_for(0.99)];
}
fn get_average_latency(&self) -> Duration {
if self.latency_samples.is_empty() {
Duration::from_micros(0)
} else {
let total: Duration = self.latency_samples.iter().sum();
total / self.latency_samples.len() as u32
}
}
}
impl<A: Float + Send + Sync> ApproximationController<A> {
fn new(targetlatency: Duration) -> Self {
Self {
approximation_level: A::zero(),
performance_history: VecDeque::with_capacity(100),
adaptation_rate: A::from(0.1).unwrap_or_else(A::one),
targetlatency,
}
}
fn get_approximation_level(&self) -> A {
self.approximation_level
}
fn record_performance(&mut self, latency: Duration, _approximation_level: A, accuracy: A) {
let now = Instant::now();
let point = PerformancePoint {
latency,
accuracy,
timestamp: now,
};
self.performance_history.push_back(point);
while self
.performance_history
.front()
.is_some_and(|p| now.duration_since(p.timestamp) > PERFORMANCE_WINDOW_AGE)
{
self.performance_history.pop_front();
}
if self.performance_history.len() > PERFORMANCE_WINDOW_LEN {
self.performance_history.pop_front();
}
self.adapt_approximation_level();
}
fn mean_latency(&self) -> Option<Duration> {
let count = self.performance_history.len();
if count == 0 {
return None;
}
let total: Duration = self.performance_history.iter().map(|p| p.latency).sum();
Some(total / count as u32)
}
fn adapt_approximation_level(&mut self) {
let Some(latency) = self.mean_latency() else {
return;
};
let target = self.targetlatency.as_micros().max(1) as f64;
let latency_ratio = latency.as_micros() as f64 / target;
if latency_ratio > 1.1 {
self.approximation_level =
(self.approximation_level + self.adaptation_rate).min(A::one());
} else if latency_ratio < 0.8 {
self.approximation_level =
(self.approximation_level - self.adaptation_rate).max(A::zero());
}
}
fn increase_approximation(&mut self) {
let double = A::from(2.0).unwrap_or_else(A::one);
self.approximation_level =
(self.approximation_level + self.adaptation_rate * double).min(A::one());
}
fn mean_accuracy(&self) -> Option<A> {
if self.performance_history.is_empty() {
return None;
}
let count = A::from(self.performance_history.len())?;
let sum = self
.performance_history
.iter()
.fold(A::zero(), |acc, point| acc + point.accuracy);
Some(sum / count)
}
}
impl<A: Float + Send + Sync + std::iter::Sum> GradientPredictor<A> {
fn new(windowsize: usize) -> Self {
Self {
gradient_history: VecDeque::with_capacity(windowsize.max(2)),
trend_weights: None,
windowsize: windowsize.max(2),
confidence: None,
pending_prediction: None,
}
}
fn observe(&mut self, gradient: &Array1<A>) {
if let Some(prediction) = self.pending_prediction.take() {
if prediction.len() == gradient.len() {
let similarity = cosine_similarity(&prediction, gradient);
let alpha = A::from(0.2).unwrap_or_else(A::one);
self.confidence = Some(match self.confidence {
Some(previous) => previous * (A::one() - alpha) + similarity * alpha,
None => similarity,
});
}
}
self.gradient_history.push_back(gradient.clone());
while self.gradient_history.len() > self.windowsize {
self.gradient_history.pop_front();
}
self.recompute_trend();
}
fn recompute_trend(&mut self) {
let n = self.gradient_history.len();
if n < 2 {
self.trend_weights = None;
return;
}
let dim = match self.gradient_history.back() {
Some(last) => last.len(),
None => return,
};
if self.gradient_history.iter().any(|g| g.len() != dim) {
self.trend_weights = None;
return;
}
let n_f = A::from(n).unwrap_or_else(A::one);
let x_mean = A::from((n - 1) as f64 / 2.0).unwrap_or_else(A::zero);
let mut denominator = A::zero();
for i in 0..n {
let dx = A::from(i).unwrap_or_else(A::zero) - x_mean;
denominator = denominator + dx * dx;
}
if denominator == A::zero() {
self.trend_weights = None;
return;
}
let mut slopes = Array1::zeros(dim);
for coordinate in 0..dim {
let mut y_sum = A::zero();
for gradient in &self.gradient_history {
y_sum = y_sum + gradient[coordinate];
}
let y_mean = y_sum / n_f;
let mut numerator = A::zero();
for (i, gradient) in self.gradient_history.iter().enumerate() {
let dx = A::from(i).unwrap_or_else(A::zero) - x_mean;
numerator = numerator + dx * (gradient[coordinate] - y_mean);
}
slopes[coordinate] = numerator / denominator;
}
self.trend_weights = Some(slopes);
}
fn predict(&mut self) -> Option<(Array1<A>, A)> {
let last = self.gradient_history.back()?.clone();
let slopes = self.trend_weights.as_ref()?;
if slopes.len() != last.len() {
return None;
}
let mut predicted = last;
for (value, &slope) in predicted.iter_mut().zip(slopes.iter()) {
*value = *value + slope;
}
self.pending_prediction = Some(predicted.clone());
let confidence = self.confidence.unwrap_or_else(A::zero);
Some((predicted, confidence))
}
}
fn cosine_similarity<A: Float>(a: &Array1<A>, b: &Array1<A>) -> A {
if a.len() != b.len() {
return A::zero();
}
let mut dot = A::zero();
let mut norm_a = A::zero();
let mut norm_b = A::zero();
for (&x, &y) in a.iter().zip(b.iter()) {
dot = dot + x * y;
norm_a = norm_a + x * x;
norm_b = norm_b + y * y;
}
let norm_a = norm_a.sqrt();
let norm_b = norm_b.sqrt();
if norm_a == A::zero() || norm_b == A::zero() {
A::zero()
} else {
dot / (norm_a * norm_b)
}
}
#[derive(Debug, Clone)]
pub struct LowLatencyMetrics {
pub avg_latency_us: u64,
pub p50_latency_us: u64,
pub p95_latency_us: u64,
pub p99_latency_us: u64,
pub latency_violations: usize,
pub total_operations: usize,
pub current_approximation_level: f64,
pub approximation_accuracy: Option<f64>,
pub precomputation_hit_rate: Option<f64>,
pub precomputation_attempts: usize,
pub memory_efficiency: f64,
pub memory_pool_misses: usize,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::optimizers::SGD;
#[test]
fn test_low_latency_config() {
let config = LowLatencyConfig::default();
assert_eq!(config.target_latency_us, 100);
assert!(config.enable_precomputation);
assert!(config.enable_lock_free);
}
#[test]
fn test_low_latency_optimizer_creation() {
let sgd = SGD::new(0.01f64);
let config = LowLatencyConfig::default();
let result = LowLatencyOptimizer::new(sgd, config);
assert!(result.is_ok());
}
#[test]
fn test_latency_monitor() {
let mut monitor = LatencyMonitor::new(10);
for i in 1..=5 {
monitor.record_latency(Duration::from_micros(i * 100));
}
assert_eq!(monitor.total_operations, 5);
assert!(monitor.get_average_latency().as_micros() > 0);
}
#[test]
fn test_gradient_quantizer() {
let mut quantizer = GradientQuantizer::new(8);
let gradient = Array1::from_vec(vec![0.1f64, 0.5, -0.3, 0.8]);
let result = quantizer.quantize(&gradient);
assert!(result.is_ok());
let quantized = result.expect("quantization of a finite gradient must succeed");
assert_eq!(quantized.len(), gradient.len());
}
#[test]
fn test_approximation_controller() {
let mut controller = ApproximationController::new(Duration::from_micros(100));
controller.record_performance(Duration::from_micros(200), 0.0f64, 0.9f64);
assert!(controller.get_approximation_level() > 0.0);
}
#[test]
fn test_lock_free_buffer() {
let buffer = LockFreeBuffer::<f64>::new(4);
assert_eq!(buffer.capacity, 4);
assert_eq!(buffer.write_index.load(Ordering::Relaxed), 0);
assert_eq!(buffer.read_index.load(Ordering::Relaxed), 0);
}
}