Skip to main content

kawa_storage/
memory.rs

1//! # 高性能メモリプール
2//!
3//! 1M+ events/secを達成するための高性能メモリ管理システム。
4//! プリアロケートされたバッファプールによりGC圧力を軽減し、
5//! 高速なメモリアロケーション/デアロケーションを実現。
6
7use crate::{StorageError, StorageResult};
8use std::{
9    collections::VecDeque,
10    sync::{
11        atomic::{AtomicUsize, Ordering},
12        Arc, Mutex,
13    },
14};
15use tokio::sync::Semaphore;
16
17/// プールされたバッファ
18/// 
19/// 高性能なバッファ管理のためのRAII型。
20/// ドロップ時に自動的にプールに戻される。
21pub struct PooledBuffer {
22    /// バッファデータ
23    pub buffer: Vec<u8>,
24    /// プールへの参照(戻すため)
25    pool: Arc<MemoryPoolInner>,
26    /// バッファサイズ
27    size: usize,
28}
29
30impl PooledBuffer {
31    /// バッファをクリア
32    pub fn clear(&mut self) {
33        self.buffer.clear();
34    }
35    
36    /// バッファサイズを取得
37    pub fn capacity(&self) -> usize {
38        self.buffer.capacity()
39    }
40    
41    /// データを書き込み
42    pub fn write(&mut self, data: &[u8]) -> StorageResult<()> {
43        if self.buffer.len() + data.len() > self.buffer.capacity() {
44            return Err(StorageError::InsufficientSpace {
45                required: data.len() as u64,
46                available: (self.buffer.capacity() - self.buffer.len()) as u64,
47            });
48        }
49        self.buffer.extend_from_slice(data);
50        Ok(())
51    }
52    
53    /// バッファ内のデータを取得
54    pub fn data(&self) -> &[u8] {
55        &self.buffer
56    }
57    
58    /// バッファのミュータブル参照を取得
59    pub fn as_mut(&mut self) -> &mut Vec<u8> {
60        &mut self.buffer
61    }
62}
63
64impl Drop for PooledBuffer {
65    fn drop(&mut self) {
66        // バッファをクリアしてプールに戻す
67        self.buffer.clear();
68        let buffer = std::mem::take(&mut self.buffer);
69        
70        // 容量が期待値と一致する場合のみプールに戻す
71        if buffer.capacity() >= self.size {
72            let mut pool = self.pool.buffers.lock().unwrap();
73            if pool.len() < self.pool.max_buffers {
74                pool.push_back(buffer);
75                self.pool.available_count.fetch_add(1, Ordering::Relaxed);
76            }
77        }
78    }
79}
80
81/// メモリプール内部構造
82#[derive(Debug)]
83struct MemoryPoolInner {
84    /// プールされたバッファ
85    buffers: Mutex<VecDeque<Vec<u8>>>,
86    /// 利用可能バッファ数
87    available_count: AtomicUsize,
88    /// 最大バッファ数
89    max_buffers: usize,
90    /// バッファサイズ
91    buffer_size: usize,
92    /// バッファ取得セマフォ
93    semaphore: Semaphore,
94}
95
96/// 高性能メモリプール
97/// 
98/// プリアロケートされたバッファプールによる高速メモリ管理。
99/// 
100/// # 機能
101/// - ゼロアロケーション書き込み
102/// - バッファ再利用による GC 圧力軽減
103/// - 非同期バッファ取得
104/// - 自動サイズ調整
105/// 
106/// # 性能特性
107/// - バッファ取得: O(1)
108/// - バッファ返却: O(1)
109/// - メモリ断片化: 最小限
110#[derive(Debug)]
111pub struct MemoryPool {
112    inner: Arc<MemoryPoolInner>,
113    /// プール統計
114    stats: Arc<PoolStats>,
115}
116
117/// プール統計情報
118#[derive(Debug, Default)]
119pub struct PoolStats {
120    /// 総取得回数
121    pub total_gets: AtomicUsize,
122    /// 総返却回数
123    pub total_returns: AtomicUsize,
124    /// キャッシュヒット数
125    pub cache_hits: AtomicUsize,
126    /// キャッシュミス数
127    pub cache_misses: AtomicUsize,
128    /// 現在利用中バッファ数
129    pub active_buffers: AtomicUsize,
130    /// プールサイズ
131    pub pool_size: AtomicUsize,
132}
133
134impl PoolStats {
135    /// ヒット率を計算
136    pub fn hit_rate(&self) -> f64 {
137        let hits = self.cache_hits.load(Ordering::Relaxed) as f64;
138        let total = hits + self.cache_misses.load(Ordering::Relaxed) as f64;
139        if total > 0.0 {
140            hits / total
141        } else {
142            0.0
143        }
144    }
145    
146    /// 利用率を計算
147    pub fn utilization(&self) -> f64 {
148        let active = self.active_buffers.load(Ordering::Relaxed) as f64;
149        let total = self.pool_size.load(Ordering::Relaxed) as f64;
150        if total > 0.0 {
151            active / total
152        } else {
153            0.0
154        }
155    }
156}
157
158impl MemoryPool {
159    /// 新しいメモリプールを作成
160    /// 
161    /// # Arguments
162    /// * `pool_size_bytes` - プール総サイズ(バイト)
163    /// * `buffer_size` - 個別バッファサイズ
164    /// 
165    /// # Returns
166    /// * `StorageResult<Self>` - メモリプールインスタンス
167    pub fn new(pool_size_bytes: usize, buffer_size: usize) -> StorageResult<Self> {
168        let max_buffers = pool_size_bytes / buffer_size;
169        
170        if max_buffers == 0 {
171            return Err(StorageError::configuration(
172                "Pool size too small for requested buffer size"
173            ));
174        }
175        
176        let mut buffers = VecDeque::with_capacity(max_buffers);
177        
178        // プールを事前に満たす(ウォームアップ)
179        for _ in 0..max_buffers {
180            let mut buffer = Vec::with_capacity(buffer_size);
181            buffer.reserve_exact(buffer_size);
182            buffers.push_back(buffer);
183        }
184        
185        let stats = Arc::new(PoolStats::default());
186        stats.pool_size.store(max_buffers, Ordering::Relaxed);
187        
188        let inner = Arc::new(MemoryPoolInner {
189            buffers: Mutex::new(buffers),
190            available_count: AtomicUsize::new(max_buffers),
191            max_buffers,
192            buffer_size,
193            semaphore: Semaphore::new(max_buffers),
194        });
195        
196        tracing::info!(
197            "MemoryPool initialized: {} buffers, {} bytes each, {} MB total",
198            max_buffers, 
199            buffer_size,
200            pool_size_bytes / (1024 * 1024)
201        );
202        
203        Ok(Self { inner, stats })
204    }
205    
206    /// バッファを取得(非同期)
207    /// 
208    /// # Returns
209    /// * `StorageResult<PooledBuffer>` - プールされたバッファ
210    pub async fn get_buffer(&self) -> StorageResult<PooledBuffer> {
211        // セマフォで利用可能性を確認
212        let _permit = self.inner.semaphore.acquire().await
213            .map_err(|_| StorageError::internal("Failed to acquire buffer permit"))?;
214        
215        self.stats.total_gets.fetch_add(1, Ordering::Relaxed);
216        
217        // バッファプールから取得を試行
218        if let Some(buffer) = self.try_get_from_pool() {
219            self.stats.cache_hits.fetch_add(1, Ordering::Relaxed);
220            self.stats.active_buffers.fetch_add(1, Ordering::Relaxed);
221            
222            return Ok(PooledBuffer {
223                buffer,
224                pool: Arc::clone(&self.inner),
225                size: self.inner.buffer_size,
226            });
227        }
228        
229        // プールが空の場合は新しいバッファを作成
230        self.stats.cache_misses.fetch_add(1, Ordering::Relaxed);
231        self.stats.active_buffers.fetch_add(1, Ordering::Relaxed);
232        
233        let mut buffer = Vec::with_capacity(self.inner.buffer_size);
234        buffer.reserve_exact(self.inner.buffer_size);
235        
236        Ok(PooledBuffer {
237            buffer,
238            pool: Arc::clone(&self.inner),
239            size: self.inner.buffer_size,
240        })
241    }
242    
243    /// プールから直接バッファ取得を試行
244    fn try_get_from_pool(&self) -> Option<Vec<u8>> {
245        let mut buffers = self.inner.buffers.lock().ok()?;
246        if let Some(buffer) = buffers.pop_front() {
247            self.inner.available_count.fetch_sub(1, Ordering::Relaxed);
248            Some(buffer)
249        } else {
250            None
251        }
252    }
253    
254    /// プール統計を取得
255    pub fn stats(&self) -> PoolStats {
256        PoolStats {
257            total_gets: AtomicUsize::new(self.stats.total_gets.load(Ordering::Relaxed)),
258            total_returns: AtomicUsize::new(self.stats.total_returns.load(Ordering::Relaxed)),
259            cache_hits: AtomicUsize::new(self.stats.cache_hits.load(Ordering::Relaxed)),
260            cache_misses: AtomicUsize::new(self.stats.cache_misses.load(Ordering::Relaxed)),
261            active_buffers: AtomicUsize::new(self.stats.active_buffers.load(Ordering::Relaxed)),
262            pool_size: AtomicUsize::new(self.stats.pool_size.load(Ordering::Relaxed)),
263        }
264    }
265    
266    /// プールの健全性をチェック
267    pub fn health_check(&self) -> bool {
268        let available = self.inner.available_count.load(Ordering::Relaxed);
269        let active = self.stats.active_buffers.load(Ordering::Relaxed);
270        let total = available + active;
271        
272        // 総数が期待範囲内かチェック
273        total <= self.inner.max_buffers
274    }
275    
276    /// プールをウォームアップ(事前バッファ作成)
277    pub async fn warmup(&self) -> StorageResult<()> {
278        tracing::info!("Warming up memory pool...");
279        
280        let warmup_count = self.inner.max_buffers / 2;
281        let mut buffers = Vec::new();
282        
283        // バッファを事前取得してすぐ戻す(ウォームアップ)
284        for _ in 0..warmup_count {
285            if let Ok(buffer) = self.get_buffer().await {
286                buffers.push(buffer);
287            }
288        }
289        
290        // バッファを解放(自動的にプールに戻る)
291        drop(buffers);
292        
293        tracing::info!("Memory pool warmup completed");
294        Ok(())
295    }
296}
297
298/// 並列バッファプロセッサ
299/// 
300/// 複数のワーカーでバッファ処理を並列化
301pub struct ParallelBufferProcessor {
302    pool: Arc<MemoryPool>,
303    worker_count: usize,
304}
305
306impl ParallelBufferProcessor {
307    /// 新しい並列プロセッサを作成
308    pub fn new(pool: Arc<MemoryPool>, worker_count: Option<usize>) -> Self {
309        let worker_count = worker_count.unwrap_or_else(|| {
310            std::thread::available_parallelism()
311                .map(|n| n.get() * 2)
312                .unwrap_or(8)
313        });
314        
315        Self { pool, worker_count }
316    }
317    
318    /// バッチデータを並列処理
319    pub async fn process_batch<T, F, Fut>(
320        &self,
321        items: Vec<T>,
322        processor: F,
323    ) -> StorageResult<Vec<StorageResult<()>>>
324    where
325        T: Send + 'static + Clone,
326        F: Fn(T, PooledBuffer) -> Fut + Send + Sync + 'static,
327        Fut: std::future::Future<Output = StorageResult<()>> + Send,
328    {
329        let processor = Arc::new(processor);
330        let chunk_size = (items.len() + self.worker_count - 1) / self.worker_count;
331        
332        let mut handles = Vec::new();
333        
334        for chunk in items.chunks(chunk_size) {
335            let chunk = chunk.to_vec();
336            let pool = Arc::clone(&self.pool);
337            let processor = Arc::clone(&processor);
338            
339            let handle = tokio::spawn(async move {
340                let mut results = Vec::new();
341                
342                for item in chunk {
343                    match pool.get_buffer().await {
344                        Ok(buffer) => {
345                            let result = processor(item, buffer).await;
346                            results.push(result);
347                        }
348                        Err(e) => {
349                            results.push(Err(e));
350                        }
351                    }
352                }
353                
354                results
355            });
356            
357            handles.push(handle);
358        }
359        
360        let mut all_results = Vec::new();
361        for handle in handles {
362            match handle.await {
363                Ok(results) => all_results.extend(results),
364                Err(e) => return Err(StorageError::internal(format!("Worker task failed: {}", e))),
365            }
366        }
367        
368        Ok(all_results)
369    }
370}
371
372#[cfg(test)]
373mod tests {
374    use super::*;
375    use tokio;
376    
377    #[tokio::test]
378    async fn test_memory_pool_basic_operations() {
379        let pool = MemoryPool::new(1024 * 1024, 1024).unwrap(); // 1MB pool, 1KB buffers
380        
381        // バッファ取得
382        let mut buffer = pool.get_buffer().await.unwrap();
383        assert_eq!(buffer.capacity(), 1024);
384        
385        // データ書き込み
386        let test_data = b"Hello, Memory Pool!";
387        buffer.write(test_data).unwrap();
388        assert_eq!(buffer.data(), test_data);
389        
390        // バッファ解放(ドロップ時に自動的にプールに戻る)
391        drop(buffer);
392        
393        // 統計確認
394        let stats = pool.stats();
395        assert_eq!(stats.total_gets.load(Ordering::Relaxed), 1);
396    }
397    
398    #[tokio::test]
399    async fn test_concurrent_buffer_access() {
400        let pool = Arc::new(MemoryPool::new(10 * 1024, 1024).unwrap()); // 10KB pool
401        
402        let mut handles = Vec::new();
403        
404        // 並行でバッファ取得
405        for i in 0..5 {
406            let pool_clone = Arc::clone(&pool);
407            let handle = tokio::spawn(async move {
408                let mut buffer = pool_clone.get_buffer().await.unwrap();
409                buffer.write(format!("Data {}", i).as_bytes()).unwrap();
410                // 少し待機してからドロップ
411                tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
412                buffer
413            });
414            handles.push(handle);
415        }
416        
417        // すべてのタスクが完了するまで待機
418        for handle in handles {
419            let _buffer = handle.await.unwrap();
420        }
421        
422        // プールの健全性をチェック
423        assert!(pool.health_check());
424    }
425}