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,
filtered: 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("filtered", &self.filtered)
.field("poisoned", &self.poisoned)
.finish()
}
}
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);
let needs_commit = g.filtered && g.current_dim % chunk_elems != 0;
let claim = dataset.claim_for_appender(needs_commit)?;
Ok(Self {
dataset,
chunk_elems,
element_size: g.element_size,
filtered: g.filtered,
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.check_usable()?;
let avail = self.whole_elements()?;
if avail == 0 {
return Ok(());
}
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()
}
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 staged = self.filtered && before % self.chunk_elems != 0;
let head = self.pending.head(
take_elems
.to_usize()?
.saturating_mul(self.element_size.get()),
);
let result = if staged {
self.dataset.append_staged_committed(head)
} else {
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()),
);
let now = self
.dataset
.append_geometry()
.map_or(before + landed, |g| g.current_dim);
self.dataset
.set_appender_needs_commit(self.claim, self.filtered && now % self.chunk_elems != 0);
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();
}
self.dataset.release_appender_claim(self.claim);
}
}