use std::ops::Bound;
use std::sync::Arc;
use uqa_core::memory::{BudgetedSharedMap, BudgetedVec};
use uqa_core::DocId;
use uqa_storage::mvcc::PrivateRecordSnapshot;
use uqa_storage::read_control::StorageReadControl;
use uqa_storage::StorageBackendResult;
use super::spilled::{self, RowPage, SpilledRows, StagedRow};
const ROW_OVERHEAD: usize = 192;
const PRESSURE_DIVISOR: usize = 2;
const RESIDENT_DIVISOR: usize = 16;
pub(super) type MemoryRows = BudgetedSharedMap<DocId, StagedRow>;
#[derive(Clone)]
pub(super) struct StagedRows {
pub(super) memory: MemoryRows,
resident: usize,
pub(super) spilled: Option<SpilledRows>,
spilled_rows: u64,
}
impl StagedRows {
pub(super) fn new(control: &StorageReadControl) -> Self {
Self {
memory: MemoryRows::new(control.memory()),
resident: 0,
spilled: None,
spilled_rows: 0,
}
}
pub(super) fn count_bound(&self) -> u64 {
self.spilled_rows
.saturating_add(u64::try_from(self.memory.len()).unwrap_or(u64::MAX))
}
pub(super) fn is_empty(&self) -> bool {
self.memory.is_empty() && self.spilled.is_none()
}
pub(super) fn needs_room(&self, control: &StorageReadControl) -> bool {
let budget = control.memory();
!self.memory.is_empty()
&& self.resident >= budget.limit() / RESIDENT_DIVISOR
&& budget.used() > budget.limit() / PRESSURE_DIVISOR
}
pub(super) fn clear_memory(&mut self) {
self.spilled_rows = self
.spilled_rows
.saturating_add(u64::try_from(self.memory.len()).unwrap_or(u64::MAX));
self.memory = MemoryRows::new(self.memory.budget());
self.resident = 0;
}
pub(super) fn insert(&mut self, id: DocId, row: StagedRow) -> StorageBackendResult<()> {
let bytes = row.as_ref().map_or(Ok(0), |row| {
let size = |fields: &uqa_storage::document_store::Document| {
crate::spill::encoded_document_size(fields)
.map_err(|error| uqa_storage::StorageBackendError::Other(error.to_string()))
};
Ok::<_, uqa_storage::StorageBackendError>(
size(row.fields.as_ref())?.saturating_add(
row.index_values
.as_ref()
.map_or(Ok(0), |values| size(values.as_ref()))?,
),
)
})?;
self.memory.try_insert(id, row)?;
self.resident = self
.resident
.saturating_add(bytes)
.saturating_add(ROW_OVERHEAD);
Ok(())
}
pub(super) fn get(
&self,
id: DocId,
control: &StorageReadControl,
) -> StorageBackendResult<Option<StagedRow>> {
if let Some(row) = self.memory.get(&id) {
return Ok(Some(row.clone()));
}
match &self.spilled {
Some(spilled) => spilled::row(spilled.view(), id, control),
None => Ok(None),
}
}
pub(super) fn contains(
&self,
id: DocId,
control: &StorageReadControl,
) -> StorageBackendResult<bool> {
if self.memory.get(&id).is_some() {
return Ok(true);
}
match &self.spilled {
Some(spilled) => spilled::holds_row(spilled.view(), id, control),
None => Ok(false),
}
}
pub(crate) fn view(&self) -> StagedRowsView {
StagedRowsView {
memory: self.memory.clone(),
spilled: self
.spilled
.as_ref()
.map(|spilled| Arc::clone(spilled.view())),
}
}
}
#[derive(Clone)]
pub(crate) struct StagedRowsView {
memory: MemoryRows,
spilled: Option<Arc<PrivateRecordSnapshot>>,
}
impl StagedRowsView {
pub(crate) fn get(
&self,
id: DocId,
control: &StorageReadControl,
) -> StorageBackendResult<Option<StagedRow>> {
if let Some(row) = self.memory.get(&id) {
return Ok(Some(row.clone()));
}
match &self.spilled {
Some(spilled) => spilled::row(spilled, id, control),
None => Ok(None),
}
}
pub(crate) fn contains(
&self,
id: DocId,
control: &StorageReadControl,
) -> StorageBackendResult<bool> {
if self.memory.get(&id).is_some() {
return Ok(true);
}
match &self.spilled {
Some(spilled) => spilled::holds_row(spilled, id, control),
None => Ok(false),
}
}
pub(crate) fn rows(&self, after: Option<DocId>) -> StagedCursor<StagedRow> {
StagedCursor::new(
self,
after,
|row: &StagedRow| row.clone(),
spilled::row_page,
)
}
pub(crate) fn presence(&self, after: Option<DocId>) -> StagedCursor<bool> {
StagedCursor::new(
self,
after,
|row: &StagedRow| row.is_some(),
spilled::presence_page,
)
}
}
type PageReader<T> = fn(
&PrivateRecordSnapshot,
Option<DocId>,
&StorageReadControl,
) -> StorageBackendResult<RowPage<T>>;
pub(crate) struct StagedCursor<T> {
view: StagedRowsView,
from_memory: fn(&StagedRow) -> T,
read_page: PageReader<T>,
memory_after: Option<DocId>,
memory_done: bool,
page: Option<BudgetedVec<(DocId, T)>>,
position: usize,
resume: Option<DocId>,
spilled_done: bool,
}
impl<T: Default> StagedCursor<T> {
fn new(
view: &StagedRowsView,
after: Option<DocId>,
from_memory: fn(&StagedRow) -> T,
read_page: PageReader<T>,
) -> Self {
Self {
view: view.clone(),
from_memory,
read_page,
memory_after: after,
memory_done: view.memory.is_empty(),
page: None,
position: 0,
resume: after,
spilled_done: view.spilled.is_none(),
}
}
fn peek_memory(&mut self) -> Option<DocId> {
if self.memory_done {
return None;
}
let start = self
.memory_after
.as_ref()
.map_or(Bound::Unbounded, Bound::Excluded);
let next = self.view.memory.range_from(start).next().map(|(id, _)| *id);
self.memory_done = next.is_none();
next
}
fn peek_spilled(
&mut self,
control: &StorageReadControl,
) -> StorageBackendResult<Option<DocId>> {
loop {
if let Some((id, _)) = self.page.as_ref().and_then(|page| page.get(self.position)) {
return Ok(Some(*id));
}
if self.spilled_done {
return Ok(None);
}
let view = self
.view
.spilled
.as_deref()
.expect("an unfinished spilled tier");
self.page = None;
let page = (self.read_page)(view, self.resume, control)?;
self.spilled_done = page.resume.is_none();
self.resume = page.resume;
self.page = Some(page.rows);
self.position = 0;
}
}
fn take_spilled(&mut self) -> (DocId, T) {
let page = self.page.as_mut().expect("a peeked page");
let (id, row) = &mut page[self.position];
self.position += 1;
(*id, std::mem::take(row))
}
pub(crate) fn next(
&mut self,
control: &StorageReadControl,
) -> StorageBackendResult<Option<(DocId, T)>> {
control.check()?;
let spilled = self.peek_spilled(control)?;
let memory = self.peek_memory();
match (memory, spilled) {
(None, None) => Ok(None),
(Some(memory), spilled) if spilled.is_none_or(|spilled| memory <= spilled) => {
if spilled == Some(memory) {
self.take_spilled();
}
self.memory_after = Some(memory);
let row = self.view.memory.get(&memory).expect("a peeked row");
Ok(Some((memory, (self.from_memory)(row))))
}
_ => Ok(Some(self.take_spilled())),
}
}
}