use std::sync::Arc;
use crate::Result;
use crate::column_family::{ColumnFamilyHandle, DEFAULT_CF_ID, cf_upper_bound, prefix_key};
use crate::engine::RegolithEngine;
use crate::engine::iterator::RegolithIterator;
use crate::slice::DbSlice;
use crate::statistics::{Histogram, Statistics, Ticker, TimeScope};
pub struct TailingIter {
engine: Arc<RegolithEngine>,
inner: RegolithIterator,
last_returned: Option<Vec<u8>>,
cf_id: u32,
cf_upper: Vec<u8>,
valid_cf: bool,
invalid_cf: Option<Box<str>>,
stats: Option<Arc<Statistics>>,
}
impl TailingIter {
pub(crate) fn new(engine: Arc<RegolithEngine>, cf_id: u32) -> Self {
let inner = engine.new_iter_at(u64::MAX);
let stats = engine.statistics_arc();
Self {
engine,
inner,
last_returned: None,
cf_id,
cf_upper: cf_upper_bound(cf_id),
valid_cf: true,
invalid_cf: None,
stats,
}
}
pub(crate) fn empty(engine: Arc<RegolithEngine>, cf: &crate::ColumnFamilyHandle) -> Self {
let mut iter = Self::new(engine, DEFAULT_CF_ID);
iter.valid_cf = false;
iter.invalid_cf = Some(
format!(
"column family handle '{}' with id {} is not live",
cf.name(),
cf.id()
)
.into_boxed_str(),
);
iter
}
fn tick_seek(&self) {
if let Some(s) = self.stats.as_deref() {
s.add(Ticker::IterSeekCount, 1);
}
}
fn tick_next(&self) {
if let Some(s) = self.stats.as_deref() {
s.add(Ticker::IterNextCount, 1);
}
}
pub fn seek(&mut self, target: &[u8]) {
if !self.valid_cf {
return;
}
self.tick_seek();
{
let _t = TimeScope::new(self.stats.as_deref(), Histogram::DbIterSeek);
self.last_returned = None;
self.inner.seek(&prefix_key(self.cf_id, target));
}
self.record_current();
}
pub fn seek_to_first(&mut self) {
if !self.valid_cf {
return;
}
self.tick_seek();
{
let _t = TimeScope::new(self.stats.as_deref(), Histogram::DbIterSeek);
self.last_returned = None;
self.inner.seek(&self.cf_id.to_be_bytes());
}
self.record_current();
}
pub fn next(&mut self) {
if !self.valid_cf {
return;
}
{
let _t = TimeScope::new(self.stats.as_deref(), Histogram::DbIterNext);
if self.inner.valid() {
self.inner.next();
}
}
if !self.inner.valid() || !self.within_cf() {
self.refresh_and_reseek();
}
self.record_current();
if self.inner.valid() && self.within_cf() {
self.tick_next();
}
}
pub fn refresh(&mut self) {
if !self.valid_cf {
return;
}
self.refresh_and_reseek();
self.record_current();
}
fn refresh_and_reseek(&mut self) {
self.inner = self.engine.new_iter_at(u64::MAX);
match self.last_returned.clone() {
Some(k) => {
self.inner.seek(&k);
if self.inner.valid() && self.inner.key() == Some(k.as_slice()) {
self.inner.next();
}
}
None => {
self.inner.seek(&self.cf_id.to_be_bytes());
}
}
}
fn record_current(&mut self) {
if self.inner.valid()
&& self.within_cf()
&& let Some(k) = self.inner.key()
{
self.last_returned = Some(k.to_vec());
}
}
fn within_cf(&self) -> bool {
if !self.valid_cf {
return false;
}
match self.inner.key() {
Some(k) => k >= self.cf_id.to_be_bytes().as_slice() && k < self.cf_upper.as_slice(),
None => false,
}
}
pub fn valid(&self) -> bool {
self.inner.valid() && self.within_cf()
}
pub fn key(&self) -> Option<&[u8]> {
if !self.valid() {
return None;
}
self.inner.key().and_then(|k| k.get(4..))
}
pub fn value(&self) -> Option<&[u8]> {
if !self.valid() {
return None;
}
self.inner.value()
}
pub fn value_slice(&self) -> Option<DbSlice> {
if !self.valid() {
return None;
}
self.inner.value_slice()
}
pub fn status(&self) -> Result<()> {
if let Some(reason) = &self.invalid_cf {
return Err(crate::Error::invalid_column_family(reason.to_string()));
}
self.inner.status().map_err(crate::Error::from)
}
}
pub(crate) fn new_default(engine: Arc<RegolithEngine>) -> TailingIter {
TailingIter::new(engine, DEFAULT_CF_ID)
}
pub(crate) fn new_for_cf(engine: Arc<RegolithEngine>, cf: &ColumnFamilyHandle) -> TailingIter {
TailingIter::new(engine, cf.id())
}
pub(crate) fn new_empty(engine: Arc<RegolithEngine>, cf: &ColumnFamilyHandle) -> TailingIter {
TailingIter::empty(engine, cf)
}