uqa-engine 0.4.8

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

//! Cross-session table-data epochs and physical cache refresh.

use super::{Engine, StorageBackendError, StorageBackendResult};

impl Engine {
    /// Publish a committed logical table-definition change to sibling
    /// sessions. Their physical stores are rebuilt lazily from their own
    /// session-bound backend on the next table lookup.
    pub(crate) fn note_table_catalog_changed(&self) {
        self.clear_regtype_output_cache();
        self.clear_bayesian_params_cache();
        if !self.session.transactions.lock().is_empty() {
            self.epochs
                .table_catalog
                .dirty
                .store(true, std::sync::atomic::Ordering::Release);
            return;
        }
        self.publish_table_catalog_changes();
    }

    pub(crate) fn publish_table_catalog_changes(&self) {
        self.clear_bayesian_params_cache();
        self.epochs
            .table_catalog
            .published
            .fetch_add(1, std::sync::atomic::Ordering::AcqRel);
        self.epochs
            .table_catalog
            .dirty
            .store(false, std::sync::atomic::Ordering::Release);
        // The writer's physical stores are current, but cached optimized and
        // prepared plans may retain a removed access path or old schema.
        // Leave `seen_table_catalog_epoch` behind so its next statement also
        // crosses the same reload/re-optimization boundary as siblings.
        self.clear_sql_statement_cache();
    }

    /// Mark table contents changed in this session. The generation is only
    /// published after the outer storage transaction commits, so sibling
    /// sessions cannot invalidate and rebuild against uncommitted data.
    pub(crate) fn note_table_data_changed(&self) {
        self.clear_bayesian_params_cache();
        self.clear_sql_statement_cache();
        // Rollback restoration replaces snapshots directly and never enters
        // this ordinary mutation hook. Therefore contention is not evidence
        // of an active transaction: wait for the stack and inspect its state.
        // This prevents an unrelated session thread from turning an
        // autocommit write into an unpublished dirty generation.
        let transaction_active = !self.session.transactions.lock().is_empty();
        if transaction_active {
            self.epochs
                .table_data
                .dirty
                .store(true, std::sync::atomic::Ordering::Release);
            return;
        }
        self.publish_table_data_changes();
    }

    pub(crate) fn publish_table_data_changes(&self) {
        self.clear_bayesian_params_cache();
        self.epochs
            .table_data
            .published
            .fetch_add(1, std::sync::atomic::Ordering::AcqRel);
        // Keep this session's observed generation behind too. Its ordinary
        // write caches were updated incrementally, but prepared/optimized
        // plans and every derived store must cross the same refresh boundary
        // as sibling sessions before the next statement.
        self.epochs
            .table_data
            .dirty
            .store(false, std::sync::atomic::Ordering::Release);
        self.clear_sql_statement_cache();
    }

    /// Refresh every session-local dependency of committed table contents.
    /// Calls made inside an already-pinned storage transaction intentionally
    /// defer the refresh: that transaction must keep using its original
    /// snapshot and will observe the new generation after it finishes.
    pub(crate) fn synchronize_table_data(&self) -> StorageBackendResult<()> {
        if self
            .storage
            .backend
            .as_ref()
            .is_some_and(|backend| backend.in_transaction())
        {
            return Ok(());
        }
        self.synchronize_external_commits()?;
        if self
            .epochs
            .table_data
            .seen
            .load(std::sync::atomic::Ordering::Acquire)
            != self
                .epochs
                .table_data
                .published
                .load(std::sync::atomic::Ordering::Acquire)
            && self.refresh_tracked_storage_snapshot()?
        {
            return Ok(());
        }
        self.refresh_table_data_cache(false)
    }

    /// Detect commits made by independently opened engines or other
    /// processes. In-process Arc epochs only coordinate sessions derived via
    /// `new_session`; a backend commit generation closes the same visibility
    /// gap for every other writer when the backend exposes one.
    pub(super) fn synchronize_external_commits(&self) -> StorageBackendResult<()> {
        let Some(backend) = self.storage.backend.as_ref() else {
            return Ok(());
        };
        if backend.in_transaction() {
            return Ok(());
        }
        let Some(version) = backend.change_version()? else {
            return Ok(());
        };
        if self
            .epochs
            .seen_storage_change_version
            .load(std::sync::atomic::Ordering::Acquire)
            == version
        {
            return Ok(());
        }

        let _statement = self.runtime.statement_gate.lock();
        let _refresh = self.epochs.external_commit_refresh.lock();
        if backend.in_transaction() {
            return Ok(());
        }
        let Some(version) = backend.change_version()? else {
            return Ok(());
        };
        if self
            .epochs
            .seen_storage_change_version
            .load(std::sync::atomic::Ordering::Acquire)
            == version
        {
            return Ok(());
        }

        // Pin one committed snapshot for the entire restore. Merely marking
        // the observed generation is insufficient: another writer can commit
        // during restore, and a table lookup from a rule/trigger validator
        // would recursively acquire this non-reentrant refresh lock. The
        // pinned transaction both defers recursive synchronization and keeps
        // table definitions and their dependent registries consistent.
        let previous_version = self
            .epochs
            .seen_storage_change_version
            .load(std::sync::atomic::Ordering::Acquire);
        backend.begin_read_transaction()?;
        let refresh_result = self.refresh_pinned_transaction_snapshot();
        let cleanup = backend.rollback_transaction();
        let refresh_result = match (refresh_result, cleanup) {
            (Ok(()), Ok(())) => Ok(()),
            (Err(error), Ok(())) | (Ok(()), Err(error)) => Err(error),
            (Err(error), Err(cleanup)) => Err(StorageBackendError::Other(format!(
                "external catalog refresh failed: {error}; snapshot cleanup failed: {cleanup}"
            ))),
        };
        if refresh_result.is_err() {
            self.epochs
                .seen_storage_change_version
                .store(previous_version, std::sync::atomic::Ordering::Release);
        }
        refresh_result
    }

    pub(super) fn refresh_table_data_cache(&self, force: bool) -> StorageBackendResult<()> {
        let target_epoch = self
            .epochs
            .table_data
            .published
            .load(std::sync::atomic::Ordering::Acquire);
        if !force
            && self
                .epochs
                .table_data
                .seen
                .load(std::sync::atomic::Ordering::Acquire)
                == target_epoch
        {
            return Ok(());
        }

        let _refresh = self.epochs.table_data.refresh.lock();
        let target_epoch = self
            .epochs
            .table_data
            .published
            .load(std::sync::atomic::Ordering::Acquire);
        let previous_epoch = self
            .epochs
            .table_data
            .seen
            .load(std::sync::atomic::Ordering::Acquire);
        if !force && previous_epoch == target_epoch {
            return Ok(());
        }
        self.clear_bayesian_params_cache();

        let tables = self
            .storage
            .tables
            .read()
            .iter()
            .map(|(name, table)| (name.clone(), table.clone()))
            .collect::<Vec<_>>();
        for (name, table) in tables {
            let name = name.qualified_name();
            let temporary = table.persistence == uqa_sql::ast::RelationPersistence::Temporary;
            // Only this session writes a memory-only or temporary table, and every write maintains or explicitly clears its value indexes, so they stay valid across its own data generations.
            if self.storage.backend.is_some() && !temporary {
                self.rebind_persistent_table_stores(&name, &table)?;
                self.refresh_table_next_id(&name, &table)?;
            }
            table
                .doc_count_dirty
                .store(true, std::sync::atomic::Ordering::Release);
            // In-memory and temporary tables have no external writers. Their
            // mutation hooks already invalidate statistics; an epoch refresh
            // must not invalidate freshly collected ANALYZE results again.
            if let Some(catalog) = self.storage.catalog.as_ref().filter(|_| !temporary) {
                let stats = Self::load_column_stats_from_catalog(catalog.as_ref(), &name)?;
                let stats_dirty = (stats.is_empty() && !table.columns.read().is_empty())
                    || crate::statistics::MaintenanceState::load_for(
                        catalog.as_ref(),
                        &name,
                        table.object_id(),
                    )?
                    .invalidates_existing_statistics();
                *table.column_stats.write() = stats;
                table
                    .column_stats_loaded
                    .store(true, std::sync::atomic::Ordering::Release);
                table
                    .column_stats_dirty
                    .store(stats_dirty, std::sync::atomic::Ordering::Release);
            }
        }
        self.synchronize_partition_identity_watermarks()?;
        self.clear_sql_statement_cache();
        // Publish the refreshed generation and invalidate dependent executable plans.
        self.epochs
            .table_data
            .seen
            .store(target_epoch, std::sync::atomic::Ordering::Release);
        self.invalidate_prepared_plans();
        Ok(())
    }

    /// Bring every session-local cache onto the outer transaction's pinned
    /// database snapshot. A stable backend change version closes the gap
    /// between a physical commit and publication of the matching in-process
    /// epochs, while allowing unchanged statements to retain their caches.
    pub(crate) fn refresh_pinned_transaction_snapshot(&self) -> StorageBackendResult<()> {
        use std::sync::atomic::Ordering;

        let (storage_snapshot_unchanged, stable_storage_version) =
            if let Some(backend) = self.storage.backend.as_ref() {
                if backend.change_version_monitor_is_nonblocking()? {
                    let before = backend.change_version()?;
                    backend.pin_transaction_snapshot()?;
                    let after = backend.change_version()?;
                    let stable = before == after;
                    (
                        stable
                            && after.is_some_and(|version| {
                                self.epochs
                                    .seen_storage_change_version
                                    .load(Ordering::Acquire)
                                    == version
                            }),
                        stable.then_some(after).flatten(),
                    )
                } else {
                    // A writer or an existing reader may prevent an independent monitor from acquiring its read lock. Refresh through the pinned session instead of waiting on a lock cycle that includes this session.
                    backend.pin_transaction_snapshot()?;
                    (false, None)
                }
            } else {
                (true, None)
            };
        let table_catalog_epoch = self.epochs.table_catalog.published.load(Ordering::Acquire);
        let table_data_epoch = self.epochs.table_data.published.load(Ordering::Acquire);
        let catalog_registry_epoch = self
            .epochs
            .catalog_registry
            .published
            .load(Ordering::Acquire);
        let read_view = self
            .storage
            .backend
            .as_ref()
            .map(|backend| backend.read_view_revision())
            .transpose()?
            .flatten();
        // A complete command-root identity distinguishes writes and undo without scanning catalog records. Providers lacking it retain the conservative private-view refresh. A view is recorded only once the caches reflect it, so an unchanged view proves them current whether or not the catalog reports cache revisions.
        let same_read_view = read_view.as_ref().is_some_and(|current| {
            self.epochs.seen_storage_read_view.lock().as_ref() == Some(current)
        });
        let private_view = self.storage.backend.as_ref().is_some_and(|backend| {
            backend.transaction_model().is_versioned() && backend.in_transaction()
        });
        let epochs = [
            table_catalog_epoch,
            table_data_epoch,
            catalog_registry_epoch,
        ];
        if (same_read_view || (read_view.is_none() && !private_view && storage_snapshot_unchanged))
            && self.observes_epochs(epochs)
        {
            return Ok(());
        }
        // With the committed state unchanged since the last refresh, every private generation comes from this session's own transaction, whose writes its caches already include.
        let committed_unchanged = read_view.as_ref().is_some_and(|current| {
            self.epochs
                .seen_storage_read_view
                .lock()
                .as_ref()
                .is_some_and(|seen| current.same_committed_state(seen))
        });
        // A failed partial restoration must not leave an older successful token eligible for reuse after undo.
        *self.epochs.seen_storage_read_view.lock() = None;
        if self.refresh_tracked_pinned_snapshot(
            table_catalog_epoch,
            table_data_epoch,
            catalog_registry_epoch,
            committed_unchanged,
        )? {
            self.observe_read_view(read_view, stable_storage_version);
            return Ok(());
        }

        // A catalog without cache revisions cannot say which tables a commit changed. With the committed state the one the caches reflect, the view differs only by this transaction's own writes, which its caches already include unless it changed definitions.
        if committed_unchanged
            && !self.epochs.table_catalog.dirty.load(Ordering::Acquire)
            && !self.epochs.catalog_registry.dirty.load(Ordering::Acquire)
            && self.observes_epochs(epochs)
        {
            self.observe_read_view(read_view, stable_storage_version);
            return Ok(());
        }
        self.clear_persistent_table_bindings_for_catalog_reload();
        self.reload_table_catalog(table_catalog_epoch)?;
        // Newly restored table handles already include their statistics and
        // physical data snapshot. Do not decode the same statistics twice.
        self.epochs
            .table_data
            .seen
            .store(table_data_epoch, Ordering::Release);
        self.synchronize_partition_identity_watermarks()?;
        self.reload_catalog_registries(catalog_registry_epoch)?;
        self.observe_read_view(read_view, stable_storage_version);
        Ok(())
    }

    /// Whether this session has observed the published table catalog, table data and catalog registry epochs `[catalog, data, registry]`.
    fn observes_epochs(&self, [catalog, data, registry]: [u64; 3]) -> bool {
        use std::sync::atomic::Ordering;
        self.epochs.table_catalog.seen.load(Ordering::Acquire) == catalog
            && self.epochs.table_data.seen.load(Ordering::Acquire) == data
            && self.epochs.catalog_registry.seen.load(Ordering::Acquire) == registry
    }

    /// Record the read view the caches now reflect, and the commit version when it stood still around the pin.
    fn observe_read_view(
        &self,
        read_view: Option<uqa_storage::key_value::KeyValueReadRevision>,
        stable_version: Option<u64>,
    ) {
        *self.epochs.seen_storage_read_view.lock() = read_view;
        if let Some(version) = stable_version {
            self.epochs
                .seen_storage_change_version
                .store(version, std::sync::atomic::Ordering::Release);
        }
    }
}