use std::sync::Arc;
use uqa_execution::catalog::security::roles::persistence::RoleCatalogSnapshot;
use uqa_execution::catalog::sequence::snapshot::{read_sequence_snapshot, SequenceReadSnapshot};
use uqa_storage::key_value::KeyValueReadRevision;
use uqa_storage::StorageBackendResult;
use crate::Engine;
struct SequenceSnapshotSources {
commits: u64,
view: KeyValueReadRevision,
}
impl SequenceSnapshotSources {
fn unchanged(&self, later: &Self) -> bool {
self.commits == later.commits && self.view.same_private_changes(&later.view)
}
}
pub(crate) struct SequenceSnapshotMemo {
sources: SequenceSnapshotSources,
registries: SequenceReadSnapshot,
snapshot: SequenceReadSnapshot,
}
fn same_registries(left: &SequenceReadSnapshot, right: &SequenceReadSnapshot) -> bool {
Arc::ptr_eq(&left.sequences, &right.sequences)
&& Arc::ptr_eq(&left.object_ids, &right.object_ids)
&& Arc::ptr_eq(&left.persistence, &right.persistence)
&& Arc::ptr_eq(&left.security, &right.security)
&& Arc::ptr_eq(&left.roles.roles, &right.roles.roles)
&& Arc::ptr_eq(&left.roles.memberships, &right.roles.memberships)
}
impl SequenceSnapshotMemo {
fn reuse(
&self,
sources: &SequenceSnapshotSources,
registries: &SequenceReadSnapshot,
) -> Option<SequenceReadSnapshot> {
(self.sources.unchanged(sources)
&& (same_registries(&self.registries, registries)
|| same_registries(&self.snapshot, registries)))
.then(|| self.snapshot.clone())
}
}
impl Engine {
pub(crate) fn sequence_registries(&self) -> SequenceReadSnapshot {
SequenceReadSnapshot {
sequences: self.durable.sequences.snapshot(),
object_ids: self.durable.sequence_object_ids.snapshot(),
persistence: self.durable.sequence_persistence.snapshot(),
security: self.durable.sequence_security.snapshot(),
roles: RoleCatalogSnapshot {
roles: self.durable.roles.snapshot(),
memberships: self.durable.role_memberships.snapshot(),
},
}
}
fn sequence_snapshot_sources(&self) -> StorageBackendResult<Option<SequenceSnapshotSources>> {
let Some(backend) = self.storage.backend.as_ref() else {
return Ok(None);
};
if !self.versioned_backend_transactions() {
return Ok(None);
}
let Some(commits) = backend.commit_monitor_version()? else {
return Ok(None);
};
Ok(backend
.read_view_revision()?
.map(|view| SequenceSnapshotSources { commits, view }))
}
pub(crate) fn latest_sequence_snapshot(&self) -> StorageBackendResult<SequenceReadSnapshot> {
let registries = self.sequence_registries();
let sources = self.sequence_snapshot_sources()?;
if let Some(sources) = sources.as_ref() {
let remembered = self
.session
.sequence_snapshot
.lock()
.as_ref()
.and_then(|memo| memo.reuse(sources, ®istries));
if let Some(snapshot) = remembered {
return Ok(snapshot);
}
}
#[cfg(test)]
self.session
.sequence_snapshot_reads
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let session = self.open_independent_catalog_session(None)?;
let snapshot = read_sequence_snapshot(
registries.clone(),
self.storage.catalog.as_deref(),
session.as_ref(),
self.versioned_backend_transactions(),
)?;
drop(session);
let kept = match (sources, self.sequence_snapshot_sources()?) {
(Some(before), Some(after)) if before.unchanged(&after) => Some(SequenceSnapshotMemo {
sources: before,
registries,
snapshot: snapshot.clone(),
}),
_ => None,
};
*self.session.sequence_snapshot.lock() = kept;
Ok(snapshot)
}
}