use crate::row_locks::{publication::TransactionRowChange, PendingRowChangeKind};
use uqa_core::{memory::MemoryBudget, CancellationToken, DocId};
use uqa_sql::SQLError;
use uqa_storage::{
mvcc::{PrivateRecordChanges, PrivateRecordSnapshot, RecordWrite, VersionError},
read_control::StorageReadControl,
};
pub trait RewriteRowChanges {
fn storage_generation(&self, table: &str) -> Result<[u8; 16], SQLError>;
fn visit_changes(
&self,
visit: &mut dyn FnMut(&[TransactionRowChange]) -> Result<(), SQLError>,
) -> Result<(), SQLError>;
}
pub fn marker(owner: &dyn RewriteRowChanges) -> Result<usize, SQLError> {
let mut length = 0_usize;
owner.visit_changes(&mut |changes| {
length = length
.checked_add(changes.len())
.ok_or_else(|| SQLError::Internal("column rewrite change marker overflowed".into()))?;
Ok(())
})?;
Ok(length)
}
pub struct ChangedRows {
snapshot: Option<PrivateRecordSnapshot>,
control: StorageReadControl,
}
pub fn capture_since(
owner: &dyn RewriteRowChanges,
marker: usize,
generation: [u8; 16],
memory: &MemoryBudget,
cancellation: &CancellationToken,
) -> Result<ChangedRows, SQLError> {
cancellation.check()?;
let control = StorageReadControl::new(memory, cancellation);
let records = PrivateRecordChanges::new(memory);
let mut skip = marker;
owner.visit_changes(&mut |changes| {
cancellation.check()?;
let skipped = skip.min(changes.len());
skip -= skipped;
for change in &changes[skipped..] {
cancellation.check()?;
note_change(&records, change, generation, &control)?;
}
Ok(())
})?;
if skip != 0 {
return Err(SQLError::Internal(
"column rewrite change marker no longer belongs to this transaction".into(),
));
}
cancellation.check()?;
let snapshot = records
.has_written()
.then(|| records.snapshot().map_err(storage_error))
.transpose()?;
Ok(ChangedRows { snapshot, control })
}
fn note_change(
records: &PrivateRecordChanges,
change: &TransactionRowChange,
generation: [u8; 16],
control: &StorageReadControl,
) -> Result<(), SQLError> {
if change.source_generation == generation {
note(
records,
change.pending.key.doc_id,
matches!(
change.pending.kind,
PendingRowChangeKind::Insert | PendingRowChangeKind::Update
),
control,
)?;
}
if let PendingRowChangeKind::Rewrite(successor) = change.pending.kind {
if change.successor_generation == Some(generation) {
note(records, successor.doc_id, true, control)?;
}
}
Ok(())
}
fn note(
records: &PrivateRecordChanges,
id: DocId,
present: bool,
control: &StorageReadControl,
) -> Result<(), SQLError> {
records
.apply(
&[RecordWrite {
key: &id.to_be_bytes(),
expected: None,
value: Some(&[u8::from(present)]),
}],
control,
)
.map_err(storage_error)
}
impl ChangedRows {
pub fn is_empty(&self) -> bool {
self.snapshot.is_none()
}
pub fn get(&self, id: DocId) -> Result<Option<bool>, SQLError> {
self.control.cancellation().check()?;
let Some(snapshot) = &self.snapshot else {
return Ok(None);
};
snapshot
.get(&id.to_be_bytes(), &self.control)
.map_err(storage_error)?
.map(|write| presence(write.value()))
.transpose()
}
pub fn visit(
&self,
visit: &mut dyn FnMut(DocId, bool) -> Result<(), SQLError>,
) -> Result<(), SQLError> {
self.control.cancellation().check()?;
let Some(snapshot) = &self.snapshot else {
return Ok(());
};
let mut cursor = snapshot
.cursor(&[], None, &self.control)
.map_err(storage_error)?;
while let Some(entry) = cursor.next(&self.control).map_err(storage_error)? {
let id = DocId::from_be_bytes(entry.key().try_into().map_err(|_| {
SQLError::Internal("a changed rewrite row has an invalid identity".into())
})?);
let present = {
let write = entry.read(&self.control).map_err(storage_error)?;
presence(write.value())?
};
drop(entry);
visit(id, present)?;
}
Ok(())
}
}
fn presence(value: Option<&[u8]>) -> Result<bool, SQLError> {
match value {
Some([0]) => Ok(false),
Some([1]) => Ok(true),
_ => Err(SQLError::Internal(
"a changed rewrite row has an invalid presence".into(),
)),
}
}
fn storage_error(error: VersionError) -> SQLError {
crate::storage_errors::storage_error(
"retain changed column rewrite rows",
&error.into_storage_error(),
)
}
#[cfg(test)]
mod tests;