use std::ops::Bound;
use reifydb_codec::key::encoded::EncodedKey;
use reifydb_core::{
common::CommitVersion,
interface::store::EntryKind,
key::{
row::{StoragePartitionedRowKey, StorageRowKey},
series::{StoragePartitionedSeriesKey, StorageSeriesKey},
},
};
use reifydb_runtime::shutdown::Shutdown;
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
use reifydb_sqlite::{SqliteConfig, SqliteTempPathGuard};
use reifydb_store::{coverage::cursor::Cursor, filter::KeyFilter, metrics::PageCacheMetrics};
use reifydb_store_commit::{MultiVersionScope, RangeBatch, RangeCursor, RangeStop, TierBatch, VersionedGetResult};
use reifydb_value::{Result, value::datetime::DateTime};
use crate::{
filter::MultiKeys,
tier::{TierStorage, range::NarrowLayout},
};
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
pub mod sqlite;
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
use sqlite::storage::SqlitePersistentStorage;
pub struct NarrowRangeRequest<'a, K> {
pub table: EntryKind,
pub start: Bound<&'a K>,
pub end: Bound<&'a K>,
pub scope: MultiVersionScope,
pub batch_size: usize,
pub descending: bool,
}
#[derive(Clone)]
#[cfg_attr(all(feature = "sqlite", not(target_arch = "wasm32")), repr(u8))]
pub enum MultiPersistentTier {
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
Sqlite(SqlitePersistentStorage) = 0,
}
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
impl MultiPersistentTier {
pub(crate) fn range_next_row(
&self,
cursor: &mut Cursor<RangeStop, StorageRowKey>,
request: NarrowRangeRequest<'_, StorageRowKey>,
) -> Result<RangeBatch<StorageRowKey>> {
match self {
Self::Sqlite(s) => s.range_chunk_row(cursor, request),
}
}
pub(crate) fn range_next_partitioned_row(
&self,
cursor: &mut Cursor<RangeStop, StoragePartitionedRowKey>,
request: NarrowRangeRequest<'_, StoragePartitionedRowKey>,
) -> Result<RangeBatch<StoragePartitionedRowKey>> {
match self {
Self::Sqlite(s) => s.range_chunk_partitioned(cursor, request),
}
}
pub(crate) fn range_next_series(
&self,
cursor: &mut Cursor<RangeStop, StorageSeriesKey>,
request: NarrowRangeRequest<'_, StorageSeriesKey>,
) -> Result<RangeBatch<StorageSeriesKey>> {
match self {
Self::Sqlite(s) => s.range_chunk_series(cursor, request),
}
}
pub(crate) fn range_next_partitioned_series(
&self,
cursor: &mut Cursor<RangeStop, StoragePartitionedSeriesKey>,
request: NarrowRangeRequest<'_, StoragePartitionedSeriesKey>,
) -> Result<RangeBatch<StoragePartitionedSeriesKey>> {
match self {
Self::Sqlite(s) => s.range_chunk_partitioned_series(cursor, request),
}
}
}
#[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
impl MultiPersistentTier {
pub(crate) fn range_next_row(
&self,
_cursor: &mut Cursor<RangeStop, StorageRowKey>,
_request: NarrowRangeRequest<'_, StorageRowKey>,
) -> Result<RangeBatch<StorageRowKey>> {
match *self {}
}
pub(crate) fn range_next_partitioned_row(
&self,
_cursor: &mut Cursor<RangeStop, StoragePartitionedRowKey>,
_request: NarrowRangeRequest<'_, StoragePartitionedRowKey>,
) -> Result<RangeBatch<StoragePartitionedRowKey>> {
match *self {}
}
pub(crate) fn range_next_series(
&self,
_cursor: &mut Cursor<RangeStop, StorageSeriesKey>,
_request: NarrowRangeRequest<'_, StorageSeriesKey>,
) -> Result<RangeBatch<StorageSeriesKey>> {
match *self {}
}
pub(crate) fn range_next_partitioned_series(
&self,
_cursor: &mut Cursor<RangeStop, StoragePartitionedSeriesKey>,
_request: NarrowRangeRequest<'_, StoragePartitionedSeriesKey>,
) -> Result<RangeBatch<StoragePartitionedSeriesKey>> {
match *self {}
}
}
pub trait PersistentRangeLayout: NarrowLayout {
fn range_next(
persistent: &MultiPersistentTier,
cursor: &mut Cursor<RangeStop, Self>,
request: NarrowRangeRequest<'_, Self>,
) -> Result<RangeBatch<Self>>;
}
impl PersistentRangeLayout for StorageRowKey {
fn range_next(
persistent: &MultiPersistentTier,
cursor: &mut Cursor<RangeStop, Self>,
request: NarrowRangeRequest<'_, Self>,
) -> Result<RangeBatch<Self>> {
persistent.range_next_row(cursor, request)
}
}
impl PersistentRangeLayout for StoragePartitionedRowKey {
fn range_next(
persistent: &MultiPersistentTier,
cursor: &mut Cursor<RangeStop, Self>,
request: NarrowRangeRequest<'_, Self>,
) -> Result<RangeBatch<Self>> {
persistent.range_next_partitioned_row(cursor, request)
}
}
impl PersistentRangeLayout for StorageSeriesKey {
fn range_next(
persistent: &MultiPersistentTier,
cursor: &mut Cursor<RangeStop, Self>,
request: NarrowRangeRequest<'_, Self>,
) -> Result<RangeBatch<Self>> {
persistent.range_next_series(cursor, request)
}
}
impl PersistentRangeLayout for StoragePartitionedSeriesKey {
fn range_next(
persistent: &MultiPersistentTier,
cursor: &mut Cursor<RangeStop, Self>,
request: NarrowRangeRequest<'_, Self>,
) -> Result<RangeBatch<Self>> {
persistent.range_next_partitioned_series(cursor, request)
}
}
impl Shutdown for MultiPersistentTier {
fn shutdown(&self) {
match self {
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
Self::Sqlite(s) => s.shutdown(),
#[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
_ => {}
}
}
}
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
impl MultiPersistentTier {
pub fn sqlite(config: SqliteConfig) -> Self {
Self::Sqlite(SqlitePersistentStorage::new(config))
}
pub fn sqlite_in_memory() -> (Self, SqliteTempPathGuard) {
let (storage, guard) = SqlitePersistentStorage::in_memory();
(Self::Sqlite(storage), guard)
}
pub fn sqlite_storage(&self) -> &SqlitePersistentStorage {
match self {
Self::Sqlite(storage) => storage,
}
}
pub fn filter(&self) -> &KeyFilter<MultiKeys> {
match self {
Self::Sqlite(storage) => storage.filter(),
}
}
pub fn page_cache_metrics(&self) -> PageCacheMetrics {
match self {
Self::Sqlite(storage) => storage.page_cache_metrics(),
}
}
pub fn set_checkpoint_threshold(&self, frames: u32) {
match self {
Self::Sqlite(s) => s.set_checkpoint_threshold(frames),
}
}
pub fn delete_keys(&self, table: EntryKind, keys: &[EncodedKey]) -> Result<u64> {
match self {
Self::Sqlite(s) => s.delete_keys(table, keys),
}
}
pub fn expired_keys(
&self,
table: EntryKind,
cutoff: DateTime,
cursor: Option<(DateTime, &[u8])>,
limit: usize,
) -> Result<Vec<(EncodedKey, DateTime)>> {
match self {
Self::Sqlite(s) => s.expired_keys(table, cutoff, cursor, limit),
}
}
pub fn list_current_entries(&self) -> Result<Vec<EntryKind>> {
match self {
Self::Sqlite(s) => s.list_current_entries(),
}
}
pub fn set_collecting_accepted(&self, version: CommitVersion, batches: TierBatch) -> Result<Vec<EncodedKey>> {
match self {
Self::Sqlite(s) => s.set_collecting_accepted(version, batches),
}
}
pub fn persist_sweep(&self, batches: Vec<(CommitVersion, TierBatch)>) -> Result<Vec<EncodedKey>> {
match self {
Self::Sqlite(s) => s.persist_sweep(batches),
}
}
pub fn install_floor(&self) -> Result<CommitVersion> {
match self {
Self::Sqlite(s) => s.install_floor(),
}
}
}
#[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
impl MultiPersistentTier {
pub fn filter(&self) -> &KeyFilter<MultiKeys> {
match *self {}
}
pub fn page_cache_metrics(&self) -> PageCacheMetrics {
match *self {}
}
pub fn set_checkpoint_threshold(&self, _frames: u32) {
match *self {}
}
pub fn delete_keys(&self, _table: EntryKind, _keys: &[EncodedKey]) -> Result<u64> {
match *self {}
}
pub fn expired_keys(
&self,
_table: EntryKind,
_cutoff: DateTime,
_cursor: Option<(DateTime, &[u8])>,
_limit: usize,
) -> Result<Vec<(EncodedKey, DateTime)>> {
match *self {}
}
pub fn list_current_entries(&self) -> Result<Vec<EntryKind>> {
match *self {}
}
pub fn persist_sweep(&self, _batches: Vec<(CommitVersion, TierBatch)>) -> Result<Vec<EncodedKey>> {
match *self {}
}
pub fn install_floor(&self) -> Result<CommitVersion> {
match *self {}
}
}
#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
impl TierStorage for MultiPersistentTier {
fn get(&self, table: EntryKind, key: &[u8], version: CommitVersion) -> Result<VersionedGetResult> {
match self {
Self::Sqlite(s) => s.get(table, key, version),
}
}
fn get_many(
&self,
table: EntryKind,
keys: &[&[u8]],
version: CommitVersion,
) -> Result<Vec<VersionedGetResult>> {
match self {
Self::Sqlite(s) => s.get_many(table, keys, version),
}
}
fn set(&self, version: CommitVersion, batches: TierBatch) -> Result<()> {
match self {
Self::Sqlite(s) => s.set(version, batches),
}
}
fn range_next(
&self,
table: EntryKind,
cursor: &mut RangeCursor,
start: Bound<&[u8]>,
end: Bound<&[u8]>,
scope: MultiVersionScope,
batch_size: usize,
) -> Result<RangeBatch> {
match self {
Self::Sqlite(s) => s.range_next(table, cursor, start, end, scope, batch_size),
}
}
fn range_rev_next(
&self,
table: EntryKind,
cursor: &mut RangeCursor,
start: Bound<&[u8]>,
end: Bound<&[u8]>,
scope: MultiVersionScope,
batch_size: usize,
) -> Result<RangeBatch> {
match self {
Self::Sqlite(s) => s.range_rev_next(table, cursor, start, end, scope, batch_size),
}
}
fn ensure_table(&self, table: EntryKind) -> Result<()> {
match self {
Self::Sqlite(s) => s.ensure_table(table),
}
}
fn clear_table(&self, table: EntryKind) -> Result<()> {
match self {
Self::Sqlite(s) => s.clear_table(table),
}
}
}
#[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
impl TierStorage for MultiPersistentTier {
fn get(&self, _table: EntryKind, _key: &[u8], _version: CommitVersion) -> Result<VersionedGetResult> {
match *self {}
}
fn set(&self, _version: CommitVersion, _batches: TierBatch) -> Result<()> {
match *self {}
}
fn range_next(
&self,
_table: EntryKind,
_cursor: &mut RangeCursor,
_start: Bound<&[u8]>,
_end: Bound<&[u8]>,
_scope: MultiVersionScope,
_batch_size: usize,
) -> Result<RangeBatch> {
match *self {}
}
fn range_rev_next(
&self,
_table: EntryKind,
_cursor: &mut RangeCursor,
_start: Bound<&[u8]>,
_end: Bound<&[u8]>,
_scope: MultiVersionScope,
_batch_size: usize,
) -> Result<RangeBatch> {
match *self {}
}
fn ensure_table(&self, _table: EntryKind) -> Result<()> {
match *self {}
}
fn clear_table(&self, _table: EntryKind) -> Result<()> {
match *self {}
}
}