use uqa_core::{CancellationToken, DocId};
use uqa_sql::SQLError;
use uqa_storage::{
mvcc::{
SerializableKeySpace, SerializablePredicate, SerializableReadContext, SerializableSession,
},
read_control::StorageReadControl,
};
use crate::storage_errors::storage_error;
#[derive(Clone)]
pub struct SerializableRelationRead {
object: [u8; 16],
context: SerializableReadContext,
control: StorageReadControl,
}
impl SerializableRelationRead {
pub fn new(
object: [u8; 16],
context: SerializableReadContext,
cancellation: &CancellationToken,
) -> Self {
let control = context.read_control(cancellation);
Self {
object,
context,
control,
}
}
pub fn for_mutation(
writes: &dyn SerializableWrites,
table: &str,
) -> Result<Option<Self>, SQLError> {
let Some(session) = writes.serializable_session() else {
return Ok(None);
};
let Some(context) = session
.serializable_read_context()
.map_err(|error| storage_error("retain serializable mutation reader", &error))?
else {
return Ok(None);
};
Ok(writes
.serializable_write_object(table)?
.map(|object| Self::new(object, context, writes.serializable_cancellation())))
}
pub fn observe_scan(&self) -> Result<(), SQLError> {
self.observe(SerializablePredicate::object(self.object))
}
pub fn observe_row(&self, doc_id: DocId) -> Result<(), SQLError> {
self.observe(SerializablePredicate::point(
self.object,
SerializableKeySpace::Rows,
&doc_id.to_be_bytes(),
))
}
fn observe(&self, predicate: SerializablePredicate<'_>) -> Result<(), SQLError> {
self.context
.observe_read(predicate, &self.control)
.map_err(|error| {
storage_error("observe serializable read", &error.into_storage_error())
})
}
}
#[derive(Default)]
pub struct SerializableScan {
read: Option<SerializableRelationRead>,
relation_observed: bool,
}
impl SerializableScan {
pub fn new(read: Option<SerializableRelationRead>) -> Self {
Self {
read,
relation_observed: false,
}
}
pub fn observe_relation(&mut self) -> Result<(), SQLError> {
if !self.relation_observed {
if let Some(read) = &self.read {
read.observe_scan()?;
}
self.relation_observed = true;
}
Ok(())
}
pub fn observe_row(&self, doc_id: DocId) -> Result<(), SQLError> {
if let Some(read) = &self.read {
read.observe_row(doc_id)?;
}
Ok(())
}
}
pub trait SerializableWrites {
fn serializable_session(&self) -> Option<&dyn SerializableSession>;
fn serializable_cancellation(&self) -> &CancellationToken;
fn serializable_write_object(&self, table: &str) -> Result<Option<[u8; 16]>, SQLError>;
}
pub fn observe_row_write(
writes: &dyn SerializableWrites,
table: &str,
doc_id: DocId,
) -> Result<(), SQLError> {
let Some(session) = writes.serializable_session() else {
return Ok(());
};
if session
.serializable_read_context()
.map_err(|error| storage_error("inspect serializable writer", &error))?
.is_none()
{
return Ok(());
}
let Some(object) = writes.serializable_write_object(table)? else {
return Ok(());
};
session
.observe_serializable_write(SerializablePredicate::point(
object,
SerializableKeySpace::Rows,
&doc_id.to_be_bytes(),
))
.map_err(|error| storage_error("observe serializable write", &error))
}
pub mod column_index;
mod field;
pub mod index_key;
pub mod text;
pub mod vector;