uqa-engine 0.4.8

Engine: schema-aware table store, catalog restore, transactions
//
// Unified Query Algebra
//
// Copyright (c) 2023-2026 Cognica, Inc.
//

//! Bind physical publication to the active storage and transaction generation.
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> {
        // A member of a partition hierarchy draws its identities from the hierarchy's owner, whose watermark is not this table's.
        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()))
    }
}