use core::num::NonZeroUsize;
use crate::convert::TryToUsize;
use crate::edit::AppendBuilder;
use crate::element::H5Element;
use crate::error::Error;
use crate::reader::Dataset;
pub struct BufferedAppender<'a> {
dataset: &'a mut Dataset,
chunk_elems: u64,
element_size: NonZeroUsize,
lossy_filters: bool,
pending: AppendBuilder,
poisoned: bool,
finished: bool,
claim: u64,
}
impl std::fmt::Debug for BufferedAppender<'_> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("BufferedAppender")
.field("chunk_elements", &self.chunk_elems)
.field("buffered_elements", &self.buffered_elements())
.field("lossy_filters", &self.lossy_filters)
.field("poisoned", &self.poisoned)
.finish()
}
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Terminal {
No,
Yes,
}
impl<'a> BufferedAppender<'a> {
pub(crate) fn new(dataset: &'a mut Dataset) -> Result<Self, Error> {
let g = dataset.append_geometry()?;
if dataset.session_is_swmr() {
return Err(Error::SwmrAppendUnsupported(
"a SWMR writer cannot write a partial trailing chunk, which is what a buffered \
appender exists to do; append whole chunks with Dataset::append instead",
));
}
let chunk_elems = g.chunk_elems.max(1);
if g.lossy_filters && g.current_dim % chunk_elems != 0 {
return Err(Error::AppendInPlaceUnsupported(
"this dataset's filter pipeline is lossy and its length is not a whole multiple \
of the chunk length, so a buffered appender's first write would have to \
re-encode already-committed values; append whole chunks from a chunk-aligned \
length instead",
));
}
let claim = dataset.claim_for_appender()?;
Ok(Self {
dataset,
chunk_elems,
element_size: g.element_size,
lossy_filters: g.lossy_filters,
pending: AppendBuilder::new(),
poisoned: false,
finished: false,
claim,
})
}
pub fn append<T: H5Element>(&mut self, data: &[T]) -> Result<(), Error> {
self.check_usable()?;
let mark = self.pending.raw().len();
self.pending.append(data);
self.settle(mark)
}
pub fn append_raw(&mut self, bytes: &[u8]) -> Result<(), Error> {
self.check_usable()?;
let mark = self.pending.raw().len();
self.pending.append_raw(bytes);
self.settle(mark)
}
#[must_use]
pub fn buffered_elements(&self) -> u64 {
self.buffered()
}
fn buffered(&self) -> u64 {
(self.pending.raw().len() / self.element_size) as u64
}
#[must_use]
pub fn unwritten(&self) -> &[u8] {
self.pending.raw()
}
#[must_use]
pub const fn chunk_elements(&self) -> u64 {
self.chunk_elems
}
pub fn flush(&mut self) -> Result<(), Error> {
self.flush_prefix(Terminal::No)
}
fn flush_prefix(&mut self, terminal: Terminal) -> Result<(), Error> {
self.check_usable()?;
let avail = self.whole_elements()?;
if avail == 0 {
return Ok(());
}
if self.lossy_filters && terminal == Terminal::No {
let current = self.dataset.append_geometry()?.current_dim;
if current.saturating_add(avail) % self.chunk_elems != 0 {
return Err(Error::AppendInPlaceUnsupported(
"this dataset's filter pipeline is lossy, and flushing these elements would \
leave its length short of a chunk boundary — a later write would then have \
to re-encode already-committed values; buffer until the elements complete \
a chunk, or end the appender with BufferedAppender::finish, which writes \
the partial tail",
));
}
}
self.write_prefix(avail)
}
pub fn discard(mut self) -> Vec<u8> {
self.finished = true;
std::mem::replace(&mut self.pending, AppendBuilder::new()).into_raw()
}
pub fn finish(mut self) -> Result<(), Error> {
self.finished = true;
self.flush_prefix(Terminal::Yes)
}
fn settle(&mut self, mark: usize) -> Result<(), Error> {
let result = self.write_complete_chunks();
if result.is_err() && !self.poisoned {
self.pending.truncate(mark);
}
result
}
fn check_usable(&self) -> Result<(), Error> {
if self.poisoned {
return Err(Error::AppendInPlaceUnsupported(
"this buffered appender failed a write and holds unwritten elements; read them \
back with BufferedAppender::unwritten and start a new appender",
));
}
Ok(())
}
fn whole_elements(&self) -> Result<u64, Error> {
if self.pending.raw().len() % self.element_size != 0 {
return Err(Error::AppendInPlaceUnsupported(
"the buffered byte length is not a whole number of elements; append the rest of \
the trailing element before flushing",
));
}
Ok(self.buffered())
}
fn write_complete_chunks(&mut self) -> Result<(), Error> {
let avail = self.buffered();
let c = self.chunk_elems;
if avail < c {
return Ok(());
}
let partial = self.dataset.append_geometry()?.current_dim % c;
let to_boundary = (c - partial) % c;
debug_assert!(to_boundary < c && to_boundary <= avail);
let take = to_boundary + ((avail - to_boundary) / c) * c;
debug_assert!((1..=avail).contains(&take));
self.write_prefix(take)
}
fn write_prefix(&mut self, take_elems: u64) -> Result<(), Error> {
let before = self.dataset.append_geometry()?.current_dim;
let head = self.pending.head(
take_elems
.to_usize()?
.saturating_mul(self.element_size.get()),
);
let result = self.dataset.append_prebuilt(&head);
let landed = if result.is_ok() {
take_elems
} else {
self.dataset
.append_geometry()
.map_or(0, |g| g.current_dim.saturating_sub(before))
};
self.pending.drop_front(
landed
.to_usize()
.unwrap_or(usize::MAX)
.saturating_mul(self.element_size.get()),
);
if let Err(e) = result {
self.poisoned = true;
return Err(e);
}
Ok(())
}
}
impl Drop for BufferedAppender<'_> {
fn drop(&mut self) {
if !(self.finished || self.poisoned || self.pending.raw().is_empty()) {
let _ = self.flush_prefix(Terminal::Yes);
}
self.dataset.release_appender_claim(self.claim);
}
}