use crate::Result;
use crate::options::{DurabilityMode, WriteOptions};
use crate::{Db, WriteBatch};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct StreamOptions {
pub max_buffered_bytes: usize,
pub durability: Option<DurabilityMode>,
}
impl Default for StreamOptions {
fn default() -> Self {
Self {
max_buffered_bytes: 1 << 20,
durability: None,
}
}
}
pub struct StreamingWriter<'db> {
db: &'db Db,
batch: WriteBatch,
buffered: usize,
opts: StreamOptions,
last_sequence: u64,
}
impl<'db> StreamingWriter<'db> {
pub(crate) fn new(db: &'db Db, opts: StreamOptions) -> Self {
Self {
db,
batch: WriteBatch::new(),
buffered: 0,
opts,
last_sequence: 0,
}
}
pub fn put(&mut self, key: &[u8], value: &[u8]) -> Result<()> {
self.put_owned(key, value.to_vec())
}
pub fn put_owned(&mut self, key: &[u8], value: Vec<u8>) -> Result<()> {
self.buffered += key.len() + value.len();
self.batch.put_owned(key, value);
self.flush_if_full()
}
pub fn delete(&mut self, key: &[u8]) -> Result<()> {
self.buffered += key.len();
self.batch.delete(key);
self.flush_if_full()
}
pub fn buffered_bytes(&self) -> usize {
self.buffered
}
pub fn flush(&mut self) -> Result<()> {
if self.batch.is_empty() {
return Ok(());
}
let batch = std::mem::take(&mut self.batch);
self.buffered = 0;
self.last_sequence = match self.opts.durability {
Some(durability) => {
let opts = WriteOptions {
sync: matches!(durability, DurabilityMode::Immediate),
..WriteOptions::default()
};
self.db.write_sequenced_opt(&opts, batch)?
}
None => self.db.write_sequenced(batch)?,
};
Ok(())
}
pub fn finish(mut self) -> Result<u64> {
self.flush()?;
Ok(self.last_sequence)
}
fn flush_if_full(&mut self) -> Result<()> {
if self.buffered >= self.opts.max_buffered_bytes {
self.flush()?;
}
Ok(())
}
}
impl std::fmt::Debug for StreamingWriter<'_> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("StreamingWriter")
.field("buffered", &self.buffered)
.field("max_buffered_bytes", &self.opts.max_buffered_bytes)
.finish_non_exhaustive()
}
}