use crate::Engine;
use std::collections::BTreeMap;
use uqa_core::{DocId, FieldName};
use uqa_execution::mutation::{
identity::MutationIdentifiers,
publication::{
DocumentVectors, MutationConstraintDeferrals, MutationHistory, MutationStorage,
MutationTextIndex, PublicationContext,
},
};
use uqa_sql::SQLError;
use uqa_storage::document_store::Document;
impl Engine {
pub(crate) fn mutation_publication_context(&self) -> PublicationContext<'_> {
PublicationContext {
observations: self,
storage: self,
text: self,
history: self,
deferrals: self,
identifiers: self,
catalog: self,
}
}
}
impl MutationIdentifiers for Engine {
fn allocate_next_id(&self, table: &str) -> Result<DocId, SQLError> {
Engine::allocate_next_id(self, table)
}
fn allocate_unmapped_id(&self, table: &str) -> Result<DocId, SQLError> {
Engine::allocate_unmapped_id(self, table)
}
fn maps_integer_keys(&self, table: &str) -> Result<bool, SQLError> {
Ok(self
.require_table(table)?
.maps_integer_keys
.load(std::sync::atomic::Ordering::Acquire))
}
fn advance_next_id(&self, table: &str, doc_id: DocId) -> uqa_storage::StorageBackendResult<()> {
Engine::advance_next_id(self, table, doc_id)
}
fn persist_next_id(&self, table: &str) -> uqa_storage::StorageBackendResult<()> {
Engine::persist_next_id(self, table)
}
fn generates_unused_identities(&self, table: &str) -> Result<bool, SQLError> {
if uqa_sql::semantics::partition::partition_hierarchy_root(self, table)?.is_some() {
return Ok(false);
}
let state = self
.try_table(table)
.map_err(|error| SQLError::Internal(error.to_string()))?
.ok_or_else(|| SQLError::UnknownTable(table.into()))?;
self.table_identifier_allocator(&state)
.map(|allocator| allocator.is_durable())
.map_err(|error| {
uqa_execution::storage_errors::storage_error("inspect document identities", &error)
})
}
}
impl MutationStorage for Engine {
fn can_defer_document_text(&self, table: &str) -> Result<bool, SQLError> {
let table = self
.try_resolve_table_name(table)
.map_err(|error| SQLError::Internal(error.to_string()))?
.ok_or_else(|| SQLError::UnknownTable(table.into()))?;
Ok(!self
.physical_index_definitions()
.map_err(|error| SQLError::Internal(error.to_string()))?
.row_publication_uses_expressions(&table))
}
fn delete_document(&self, table: &str, doc_id: DocId) -> Result<(), SQLError> {
Engine::delete_document(self, table, doc_id)
}
fn observe_document_identity(
&self,
table: &str,
doc_id: DocId,
) -> Result<uqa_storage::mvcc::ObservedIdentifier, SQLError> {
let state = self
.try_table(table)
.map_err(|error| SQLError::Internal(error.to_string()))?
.ok_or_else(|| SQLError::UnknownTable(table.into()))?;
self.table_identifier_allocator(&state)
.and_then(|allocator| allocator.observe_durably(doc_id))
.map_err(|error| {
uqa_execution::storage_errors::storage_error("observe document identity", &error)
})
}
fn delete_document_deferred_text(&self, table: &str, doc_id: DocId) -> Result<(), SQLError> {
self.delete_prepared_document_deferred_fts(table, doc_id)
}
fn insert_document(
&self,
table: &str,
doc_id: DocId,
document: Document,
vectors: DocumentVectors,
inserted: uqa_execution::mutation::publication::InsertedIdentity,
) -> Result<(), SQLError> {
self.add_prepared_document_with_vector_values(table, doc_id, document, vectors, inserted)
}
fn insert_document_deferred_text(
&self,
table: &str,
doc_id: DocId,
document: Document,
vectors: DocumentVectors,
inserted: uqa_execution::mutation::publication::InsertedIdentity,
) -> Result<(), SQLError> {
self.add_prepared_document_with_vector_values_deferred_fts(
table, doc_id, document, vectors, inserted,
)
}
fn rewrite_document(
&self,
table: &str,
doc_id: DocId,
document: Document,
) -> Result<(), SQLError> {
self.rewrite_prepared_document(table, doc_id, document)
}
fn rewrite_document_deferred_text(
&self,
table: &str,
doc_id: DocId,
document: Document,
) -> Result<(), SQLError> {
self.rewrite_prepared_document_deferred_fts(table, doc_id, document)
}
}
impl MutationTextIndex for Engine {
fn text_fields(
&self,
table: &str,
document: &Document,
) -> Result<BTreeMap<FieldName, String>, SQLError> {
self.prepared_document_text_fields(table, document)
}
fn add_documents(
&self,
table: &str,
documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
) -> Result<(), SQLError> {
self.add_prepared_fts_documents(table, documents)
}
}
impl MutationHistory for Engine {
fn note_rewrite(
&self,
old_table: &str,
old_doc_id: DocId,
new_table: &str,
new_doc_id: DocId,
) -> Result<(), SQLError> {
self.note_row_rewritten_between_tables(old_table, old_doc_id, new_table, new_doc_id)
}
}
impl MutationConstraintDeferrals for Engine {
fn inserted(&self, table: &str, doc_id: DocId) -> Result<(), SQLError> {
self.defer_inserted_foreign_key_checks(table, doc_id)
}
fn rewritten(
&self,
table: &str,
doc_id: DocId,
old: Option<&Document>,
new: &Document,
) -> Result<(), SQLError> {
self.defer_rewritten_foreign_key_checks(table, doc_id, old, new)
}
}
impl uqa_execution::mutation::staging::MutationCommandRows for Engine {
fn stage_shared_command_document(
&self,
table: &str,
doc_id: DocId,
document: Option<std::sync::Arc<Document>>,
) -> Result<(), SQLError> {
Engine::stage_shared_command_document(self, table, doc_id, document)
}
fn stage_command_document(
&self,
table: &str,
doc_id: DocId,
document: Option<Document>,
) -> Result<(), SQLError> {
Engine::stage_command_document(self, table, doc_id, document)
}
}
impl Engine {
pub(crate) fn mutation_staging_context(
&self,
) -> uqa_execution::mutation::staging::MutationStagingContext<'_> {
uqa_execution::mutation::staging::MutationStagingContext {
commands: self,
constraints: self.constraint_execution_context(),
triggers: self.trigger_execution_context(),
}
}
}
use uqa_execution::serializable::SerializableWrites;
use uqa_sql::ast::RelationPersistence;
use uqa_storage::mvcc::SerializableSession;
impl SerializableWrites for Engine {
fn serializable_session(&self) -> Option<&dyn SerializableSession> {
self.storage.backend.as_ref()?.serializable_session()
}
fn serializable_cancellation(&self) -> &uqa_core::CancellationToken {
&self.runtime.cancellation
}
fn serializable_write_object(&self, table: &str) -> Result<Option<[u8; 16]>, SQLError> {
let table = self.require_table(table)?;
Ok((table.persistence != RelationPersistence::Temporary).then_some(table.object_id()))
}
}