1#![warn(missing_docs)]
27
28mod cipher;
29mod error;
30mod segment;
31
32#[cfg(feature = "encryption")]
33pub use cipher::AesGcmCipher;
34pub use cipher::SegmentCipher;
35pub use error::{Result, SegmentError};
36
37use std::fs;
38use std::path::PathBuf;
39use std::time::Instant;
40
41use parking_lot::Mutex;
42use serde::de::DeserializeOwned;
43use serde::Serialize;
44use tracing::{debug, info};
45
46use segment::SegmentRange;
47
48pub struct SegmentConfig {
50 pub max_batch_events: usize,
52 pub flush_interval_secs: u64,
55 pub max_size_bytes: u64,
57 pub compression_level: i32,
59 pub cipher: Option<Box<dyn SegmentCipher>>,
62}
63
64impl std::fmt::Debug for SegmentConfig {
65 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66 f.debug_struct("SegmentConfig")
67 .field("max_batch_events", &self.max_batch_events)
68 .field("flush_interval_secs", &self.flush_interval_secs)
69 .field("max_size_bytes", &self.max_size_bytes)
70 .field("compression_level", &self.compression_level)
71 .field("cipher", &self.cipher.as_ref().map(|_| "[set]"))
72 .finish()
73 }
74}
75
76impl Default for SegmentConfig {
77 fn default() -> Self {
78 Self {
79 max_batch_events: 256,
80 flush_interval_secs: 5,
81 max_size_bytes: 10 * 1024 * 1024 * 1024,
82 compression_level: 3,
83 cipher: None,
84 }
85 }
86}
87
88struct BufferInner<T> {
89 unflushed: Vec<T>,
93 next_seq: u64,
94 head_seq: u64,
95 last_flush: Instant,
96 approx_disk_bytes: u64,
97}
98
99pub struct SegmentBuffer<T> {
106 dir: PathBuf,
107 config: SegmentConfig,
108 inner: Mutex<BufferInner<T>>,
109}
110
111impl<T> SegmentBuffer<T>
112where
113 T: Serialize + DeserializeOwned + Clone + Send + 'static,
114{
115 pub fn open(dir: impl Into<PathBuf>, config: SegmentConfig) -> Result<Self> {
127 let dir = dir.into();
128 fs::create_dir_all(&dir)?;
129
130 let buffer = Self {
131 dir,
132 config,
133 inner: Mutex::new(BufferInner {
134 unflushed: Vec::new(),
135 next_seq: 0,
136 head_seq: 0,
137 last_flush: Instant::now(),
138 approx_disk_bytes: 0,
139 }),
140 };
141
142 buffer.recover()?;
143 Ok(buffer)
144 }
145
146 pub fn append(&self, event: T) -> Result<u64> {
155 let (should_flush, seq) = {
156 let mut inner = self.inner.lock();
157 inner.unflushed.push(event);
158 inner.next_seq += 1;
159 let seq = inner.next_seq - 1;
160
161 let batch_full = inner.unflushed.len() >= self.config.max_batch_events;
162 let interval_elapsed =
163 inner.last_flush.elapsed().as_secs() >= self.config.flush_interval_secs;
164 (batch_full || interval_elapsed, seq)
165 };
166
167 if should_flush {
168 self.flush()?;
169 }
170
171 Ok(seq)
172 }
173
174 pub fn flush(&self) -> Result<()> {
176 let (events, start_seq, end_seq) = {
177 let mut inner = self.inner.lock();
178 inner.last_flush = Instant::now();
179 if inner.unflushed.is_empty() {
180 return Ok(());
181 }
182 let events = std::mem::take(&mut inner.unflushed);
183 let count = events.len() as u64;
184 let end_seq = inner.next_seq - 1;
185 let start_seq = end_seq + 1 - count;
186 (events, start_seq, end_seq)
187 };
188
189 let compressed_len = self.write_segment(start_seq, end_seq, &events)?;
190
191 {
192 let mut inner = self.inner.lock();
193 inner.approx_disk_bytes += compressed_len;
194 }
195
196 debug!(start_seq, end_seq, count = events.len(), "Flushed segment");
197 Ok(())
198 }
199
200 pub fn read_from(&self, start_seq: u64, limit: usize) -> Result<Vec<T>> {
205 if limit == 0 {
206 return Ok(Vec::new());
207 }
208
209 let mut result: Vec<T> = Vec::with_capacity(limit.min(1024));
210
211 let segments = self.scan_segments()?;
213 for seg in &segments {
214 if result.len() >= limit {
215 break;
216 }
217 if seg.end < start_seq {
218 continue;
219 }
220
221 let events = self.read_segment(*seg)?;
222 let skip = if seg.start < start_seq {
223 (start_seq - seg.start) as usize
224 } else {
225 0
226 };
227
228 for event in events.into_iter().skip(skip) {
229 if result.len() >= limit {
230 break;
231 }
232 result.push(event);
233 }
234 }
235
236 if result.len() < limit {
238 let inner = self.inner.lock();
239 let pending_start = inner.next_seq.saturating_sub(inner.unflushed.len() as u64);
240 for (i, event) in inner.unflushed.iter().enumerate() {
241 let seq = pending_start + i as u64;
242 if seq < start_seq {
243 continue;
244 }
245 if result.len() >= limit {
246 break;
247 }
248 result.push(event.clone());
249 }
250 }
251
252 Ok(result)
253 }
254
255 pub fn delete_acked(&self, acked_seq: u64) -> Result<usize> {
270 let segments = self.scan_segments()?;
271 let mut deleted = 0;
272 let mut freed_bytes: u64 = 0;
273 let mut new_head = None;
274
275 for seg in &segments {
276 if seg.end <= acked_seq {
277 let path = self.segment_path(seg.start, seg.end);
278 if let Ok(meta) = fs::metadata(&path) {
279 freed_bytes += meta.len();
280 }
281 match fs::remove_file(&path) {
282 Ok(()) => {
283 deleted += 1;
284 debug!(start = seg.start, end = seg.end, "Deleted acked segment");
285 }
286 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
287 Err(e) => return Err(e.into()),
288 }
289 } else if new_head.is_none() {
290 new_head = Some(seg.start);
291 }
292 }
293
294 {
295 let mut inner = self.inner.lock();
296 inner.approx_disk_bytes = inner.approx_disk_bytes.saturating_sub(freed_bytes);
297 let pending_start = inner.next_seq.saturating_sub(inner.unflushed.len() as u64);
304 inner.head_seq = new_head.unwrap_or(inner.next_seq).min(pending_start);
305 }
306
307 if deleted > 0 {
308 info!(deleted, freed_bytes, acked_seq, "Deleted acked segments");
309 }
310
311 Ok(deleted)
312 }
313
314 pub fn latest_sequence(&self) -> u64 {
316 let inner = self.inner.lock();
317 if inner.next_seq == 0 {
318 0
319 } else {
320 inner.next_seq - 1
321 }
322 }
323
324 pub fn pending_count(&self) -> u64 {
326 let inner = self.inner.lock();
327 inner.next_seq.saturating_sub(inner.head_seq)
328 }
329
330 pub fn store_pressure(&self) -> f32 {
335 let inner = self.inner.lock();
336 if self.config.max_size_bytes == 0 {
337 return 0.0;
338 }
339 (inner.approx_disk_bytes as f32 / self.config.max_size_bytes as f32).min(1.0)
340 }
341
342 pub fn is_overloaded(&self) -> bool {
344 self.store_pressure() > 0.9
345 }
346
347 fn recover(&self) -> Result<()> {
352 segment::clean_tmp(&self.dir)?;
353
354 let segments = self.scan_segments()?;
355
356 let mut inner = self.inner.lock();
357 let total_bytes: u64 = segments
358 .iter()
359 .filter_map(|s| fs::metadata(self.segment_path(s.start, s.end)).ok())
360 .map(|m| m.len())
361 .sum();
362
363 match (segments.first(), segments.last()) {
364 (Some(first), Some(last)) => {
365 inner.head_seq = first.start;
366 inner.next_seq = last.end + 1;
367 }
368 _ => {
369 inner.next_seq = 0;
370 inner.head_seq = 0;
371 }
372 }
373 inner.approx_disk_bytes = total_bytes;
374
375 info!(
376 segments = segments.len(),
377 head_seq = inner.head_seq,
378 next_seq = inner.next_seq,
379 disk_bytes = total_bytes,
380 "Segment buffer recovered"
381 );
382
383 Ok(())
384 }
385
386 fn write_segment(&self, start: u64, end: u64, events: &[T]) -> Result<u64> {
387 segment::write(
388 &self.dir,
389 self.config.cipher.as_deref(),
390 self.config.compression_level,
391 SegmentRange { start, end },
392 events,
393 )
394 }
395
396 fn read_segment(&self, seg: SegmentRange) -> Result<Vec<T>> {
397 segment::read(&self.dir, self.config.cipher.as_deref(), seg)
398 }
399
400 fn scan_segments(&self) -> Result<Vec<SegmentRange>> {
401 segment::scan(&self.dir)
402 }
403
404 fn segment_path(&self, start: u64, end: u64) -> PathBuf {
405 self.dir.join(segment::filename(start, end))
406 }
407}
408
409#[cfg(test)]
410mod tests;