use std::sync::Arc;
use ahash::AHashMap;
use crate::common::stored_bitslice::StoredBitSlice;
use crate::common::universal_io::UniversalWrite;
use itertools::Itertools;
use parking_lot::RwLock;
use crate::segment::common::Flusher;
use crate::segment::common::operation_error::{OperationError, OperationResult};
#[derive(Debug)]
pub struct BufferedUpdateBitSlice<S> {
bitslice: Arc<RwLock<StoredBitSlice<S>>>,
len: usize,
pending_updates: Arc<RwLock<AHashMap<usize, bool>>>,
is_alive_flush_lock: crate::common::is_alive_lock::IsAliveLock,
}
impl<S: UniversalWrite + Send + Sync + 'static> BufferedUpdateBitSlice<S> {
pub fn new(bitslice: StoredBitSlice<S>) -> Self {
let len = bitslice.bit_len() as usize;
Self {
bitslice: Arc::new(RwLock::new(bitslice)),
len,
pending_updates: Arc::new(RwLock::new(AHashMap::new())),
is_alive_flush_lock: crate::common::is_alive_lock::IsAliveLock::new(),
}
}
pub fn set(&self, index: usize, value: bool) {
assert!(index < self.len, "index {index} out of range: {}", self.len);
self.pending_updates.write().insert(index, value);
}
pub fn get(&self, index: usize) -> Option<bool> {
if index >= self.len {
return None;
}
if let Some(value) = self.pending_updates.read().get(&index) {
Some(*value)
} else {
self.bitslice
.read()
.get_bit(index as u64)
.unwrap_or_else(|err| {
log::error!("Error reading bit at index {index}: {err}");
debug_assert!(false, "Error reading bit at index {index}: {err}");
None
})
}
}
pub fn len(&self) -> usize {
self.len
}
pub fn is_empty(&self) -> bool {
self.len == 0
}
fn reconcile_persisted_updates(
pending_updates: &RwLock<AHashMap<usize, bool>>,
persisted: AHashMap<usize, bool>,
) {
pending_updates
.write()
.retain(|point_id, a| persisted.get(point_id).is_none_or(|b| a != b));
}
pub fn clear_cache(&self) -> OperationResult<()> {
let Self {
bitslice,
len: _,
pending_updates: _,
is_alive_flush_lock: _,
} = self;
bitslice.read().clear_ram_cache()?;
Ok(())
}
pub fn flusher(&self) -> Flusher {
let updates = {
let updates_guard = self.pending_updates.read();
if updates_guard.is_empty() {
return Box::new(|| Ok(()));
}
updates_guard.clone()
};
let bitslice = Arc::downgrade(&self.bitslice);
let pending_updates_weak = Arc::downgrade(&self.pending_updates);
let is_alive_flush_lock = self.is_alive_flush_lock.handle();
Box::new(move || {
let (Some(is_alive_flush_guard), Some(bitslice), Some(pending_updates_arc)) = (
is_alive_flush_lock.lock_if_alive(),
bitslice.upgrade(),
pending_updates_weak.upgrade(),
) else {
log::trace!("BufferedUpdateBitSlice was dropped, cancelling flush");
return Err(OperationError::cancelled(
"Aborted flushing on a dropped BufferedUpdateBitSlice instance",
));
};
let mut storage_write = bitslice.write();
storage_write.set_ascending_bits_batch(
updates
.iter()
.map(|(idx, value)| (*idx as u64, *value))
.sorted_unstable_by_key(|(idx, _)| *idx),
)?;
storage_write.flusher()()?;
drop(is_alive_flush_guard);
Self::reconcile_persisted_updates(&pending_updates_arc, updates);
Ok(())
})
}
}