mod binding;
mod fetch;
use crate::{
AnalyzerPhase, DocumentStore, Engine, EpochCoordinator, InvertedIndex, MemoryDocumentStore,
MemoryInvertedIndex, MemoryVectorIndex, PinnedPortalTransactionControl, QueryRuntime,
RelationIdentity, RuntimeExtensions, SQLError, SQLResult, SessionPortalCatalogSnapshot,
SessionPortalCommandDeclaration, SessionPortalData, SessionPortalDeclaration,
SessionPortalMaterialization, SessionPortalPosition, SessionPortalRestart,
SessionPortalSQLFunctionSnapshots, SessionPortalState, SessionPortalTableSnapshots,
SessionPortalTransactionOverlay, SessionPortalViewSnapshots, StorageContext, StoredDocument,
TableState, Value, VectorIndex,
};
use binding::bind_session_portal_function_relations;
use fetch::{
ensure_portal_rows_for_fetch, fetch_directional_query_portal, fetch_indices,
materialize_portal_to_end, select_portal_rows, uses_directional_query_execution,
};
use uqa_planner::{QueryPlan, RelationalPlan, SourcePlan};
use uqa_sql::ast::{CursorDirection, FetchCursorStmt};
struct SessionPortalTableDependencies {
tables: Option<std::collections::BTreeSet<RelationIdentity>>,
graphs: Option<std::collections::BTreeSet<String>>,
graph_catalog: bool,
}
impl SessionPortalTableDependencies {
fn empty() -> Self {
Self {
tables: Some(std::collections::BTreeSet::new()),
graphs: Some(std::collections::BTreeSet::new()),
graph_catalog: false,
}
}
fn all() -> Self {
Self {
tables: None,
graphs: None,
graph_catalog: true,
}
}
fn is_all(&self) -> bool {
self.tables.is_none() && self.graphs.is_none()
}
fn includes(&self, relation: &RelationIdentity) -> bool {
self.tables
.as_ref()
.is_none_or(|tables| tables.contains(relation))
}
fn insert(&mut self, relation: RelationIdentity) {
if let Some(relations) = &mut self.tables {
relations.insert(relation);
}
}
fn insert_graph(&mut self, graph: String) {
self.graph_catalog = true;
if let Some(graphs) = &mut self.graphs {
graphs.insert(graph);
}
}
}
type SessionPortalTableSource = (
RelationIdentity,
std::sync::Arc<TableState>,
std::sync::Arc<TableState>,
);
impl Engine {
pub(crate) fn allocate_session_portal_name(&self) -> String {
let mut next = self.session.next_portal_id.lock();
let name = format!("<unnamed portal {}>", *next);
*next += 1;
name
}
fn allocate_session_portal_transaction_origin(&self) -> u64 {
let mut next = self.session.next_portal_transaction_origin.lock();
let origin = *next;
*next = next.wrapping_add(1).max(1);
origin
}
pub(crate) fn open_pending_session_portal(
&self,
declaration: SessionPortalDeclaration,
) -> Result<(), SQLError> {
let SessionPortalDeclaration {
name,
mut query,
params,
columns,
column_types,
scrollable,
holdable,
binary,
} = declaration;
bind_session_portal_query_relations(self, &mut query, &std::collections::BTreeSet::new())?;
super::bind_query_plan_sequence_references(&mut query, &mut |reference| {
self.try_resolve_sequence_oid_reference_for_binding(reference)
.map_err(|error| {
SQLError::Internal(format!(
"bind cursor sequence `{reference}` at DECLARE: {error}"
))
})
.and_then(|bound| {
bound.ok_or_else(|| SQLError::Routine {
sqlstate: "42P01".into(),
message: format!("relation \"{reference}\" does not exist"),
})
})
})?;
let table_dependencies = session_portal_table_dependencies(self, &query)?;
let snapshot_gate = self
.row_locks
.begin_change_snapshot(&self.runtime.cancellation)?;
let transaction_overlay = self.capture_session_portal_transaction_overlay()?;
snapshot_gate.baseline()?;
drop(snapshot_gate);
let table_sources = {
let stack = self.session.transactions.lock();
let fixed_snapshot = stack
.first()
.and_then(|frame| frame.fixed_snapshot.as_ref());
self.capture_session_portal_table_sources(fixed_snapshot, &table_dependencies)
};
let table_snapshots = Self::detach_session_portal_table_snapshots(
table_sources,
transaction_overlay.as_ref(),
)?;
let mut catalog_snapshot = self.durable.snapshot();
catalog_snapshot.graphs = self.freeze_graph_read_handles(
table_dependencies.graphs.as_ref(),
table_dependencies.graph_catalog,
)?;
let catalog_snapshot = std::sync::Arc::new(catalog_snapshot);
let view_snapshots = std::sync::Arc::clone(&catalog_snapshot.views);
let sql_function_snapshots = std::sync::Arc::clone(&catalog_snapshot.sql_user_functions);
let restart = holdable.then(|| SessionPortalRestart {
query: query.clone(),
params: params.clone(),
table_snapshots: std::sync::Arc::clone(&table_snapshots),
view_snapshots: std::sync::Arc::clone(&view_snapshots),
sql_function_snapshots: std::sync::Arc::clone(&sql_function_snapshots),
catalog_snapshot: std::sync::Arc::clone(&catalog_snapshot),
});
let transaction_origin = self.allocate_session_portal_transaction_origin();
let mut portals = self.session.portals.lock();
if portals.contains_key(&name) {
return Err(cursor_error(&name, "already exists", "42P03"));
}
portals.insert(
name,
SessionPortalState {
data: SessionPortalData::Pending {
query,
params,
table_snapshots,
view_snapshots,
sql_function_snapshots,
catalog_snapshot,
restart,
},
columns,
column_types,
transaction_origin,
position: SessionPortalPosition::BeforeFirst,
scrollable,
holdable,
pinned_transaction_control: PinnedPortalTransactionControl::MakeHoldable,
pin_count: 0,
_binary: binary,
},
);
Ok(())
}
pub(crate) fn open_pending_command_session_portal(
&self,
declaration: SessionPortalCommandDeclaration,
) -> Result<(), SQLError> {
let SessionPortalCommandDeclaration {
name,
command,
params,
columns,
column_types,
scrollable,
null_returning_values,
} = declaration;
let transaction_origin = self.allocate_session_portal_transaction_origin();
let mut portals = self.session.portals.lock();
if portals.contains_key(&name) {
return Err(cursor_error(&name, "already exists", "42P03"));
}
portals.insert(
name,
SessionPortalState {
data: SessionPortalData::PendingCommand {
command,
params,
null_returning_values,
},
columns,
column_types,
transaction_origin,
position: SessionPortalPosition::BeforeFirst,
scrollable,
holdable: false,
pinned_transaction_control: PinnedPortalTransactionControl::Reject,
pin_count: 0,
_binary: false,
},
);
Ok(())
}
pub(crate) fn capture_statement_snapshot_engine(&self) -> Result<Engine, SQLError> {
let dependencies = SessionPortalTableDependencies::all();
let snapshot_gate = self
.row_locks
.begin_change_snapshot(&self.runtime.cancellation)?;
let transaction_overlay = self.capture_session_portal_transaction_overlay()?;
snapshot_gate.baseline()?;
drop(snapshot_gate);
let table_sources = {
let stack = self.session.transactions.lock();
let fixed_snapshot = stack
.first()
.and_then(|frame| frame.fixed_snapshot.as_ref());
self.capture_session_portal_table_sources(fixed_snapshot, &dependencies)
};
let table_snapshots = Self::detach_session_portal_table_snapshots(
table_sources,
transaction_overlay.as_ref(),
)?;
let mut catalog_snapshot = self.durable.snapshot();
catalog_snapshot.graphs = self.freeze_graph_read_handles(None, true)?;
let catalog_snapshot = std::sync::Arc::new(catalog_snapshot);
let view_snapshots = std::sync::Arc::clone(&catalog_snapshot.views);
let sql_function_snapshots = std::sync::Arc::clone(&catalog_snapshot.sql_user_functions);
Ok(self.session_portal_worker_engine(
table_snapshots,
view_snapshots,
sql_function_snapshots,
catalog_snapshot,
self.allocate_session_portal_transaction_origin(),
))
}
pub(crate) fn ensure_session_portal_available(&self, name: &str) -> Result<(), SQLError> {
if self.session.portals.lock().contains_key(name) {
return Err(cursor_error(name, "already exists", "42P03"));
}
Ok(())
}
pub(crate) fn pin_session_portal(&self, name: &str) -> Result<(), SQLError> {
let mut portals = self.session.portals.lock();
let portal = portals
.get_mut(name)
.ok_or_else(|| cursor_error(name, "does not exist", "34000"))?;
portal.pin_count = portal
.pin_count
.checked_add(1)
.ok_or_else(|| SQLError::Internal(format!("cursor \"{name}\" pin count overflow")))?;
Ok(())
}
pub(crate) fn unpin_session_portal(&self, name: &str) -> Result<(), SQLError> {
let mut portals = self.session.portals.lock();
let portal = portals
.get_mut(name)
.ok_or_else(|| cursor_error(name, "does not exist", "34000"))?;
portal.pin_count = portal
.pin_count
.checked_sub(1)
.ok_or_else(|| SQLError::Internal(format!("cursor \"{name}\" is not pinned")))?;
Ok(())
}
pub(crate) fn fetch_session_portal(
&self,
fetch: &FetchCursorStmt,
) -> Result<SQLResult, SQLError> {
let mut state = self
.session
.portals
.lock()
.remove(&fetch.name)
.ok_or_else(|| cursor_error(&fetch.name, "does not exist", "34000"))?;
let result = (|| {
if uses_directional_query_execution(&state) {
return fetch_directional_query_portal(
self,
&mut state,
fetch.direction,
fetch.count,
fetch.move_only,
);
}
ensure_portal_rows_for_fetch(
self,
&mut state,
fetch.direction,
fetch.count,
fetch.move_only,
)?;
let indices = fetch_indices(&mut state, fetch.direction, fetch.count, fetch.move_only)?;
if fetch.move_only {
return Ok(SQLResult::from_affected(indices.len() as u64));
}
select_portal_rows(&mut state, &indices)
})();
self.session
.portals
.lock()
.insert(fetch.name.clone(), state);
result
}
pub(crate) fn materialize_holdable_session_portals(&self) -> Result<bool, SQLError> {
let names = self
.session
.portals
.lock()
.iter()
.filter_map(|(name, portal)| {
(portal.holdable
&& matches!(
&portal.data,
SessionPortalData::Pending { .. } | SessionPortalData::Streaming { .. }
))
.then_some(name.clone())
})
.collect::<Vec<_>>();
let materialized_any = !names.is_empty();
self.materialize_named_session_portals(&names)?;
Ok(materialized_any)
}
pub(crate) fn prepare_pinned_session_portals_for_transaction_control(
&self,
) -> Result<(), SQLError> {
let names = {
let mut portals = self.session.portals.lock();
if portals.values().any(|portal| {
portal.pin_count != 0
&& portal.pinned_transaction_control == PinnedPortalTransactionControl::Reject
}) {
return Err(SQLError::Routine {
sqlstate: "55000".into(),
message: "cannot perform transaction commands inside a cursor loop that is not read-only".into(),
});
}
let names = portals
.iter()
.filter_map(|(name, portal)| (portal.pin_count != 0).then_some(name.clone()))
.collect::<Vec<_>>();
for name in &names {
if let Some(portal) = portals.get_mut(name) {
portal.holdable = true;
}
}
names
};
self.materialize_named_session_portals(&names)?;
if !names.is_empty() {
let mut stack = self.session.transactions.lock();
let frame = stack.last_mut().ok_or_else(|| {
SQLError::Internal(
"pinned PL/pgSQL portal prepared without an active transaction".into(),
)
})?;
frame.session_snapshot.portal_names.extend(names);
}
Ok(())
}
fn materialize_named_session_portals(&self, names: &[String]) -> Result<(), SQLError> {
for name in names {
let mut state = self
.session
.portals
.lock()
.remove(name)
.ok_or_else(|| cursor_error(name, "does not exist", "34000"))?;
let _row_lock_statement = self.begin_row_lock_statement();
let result = materialize_portal_to_end(self, &mut state);
self.session.portals.lock().insert(name.clone(), state);
result?;
}
Ok(())
}
pub(crate) fn close_session_portal(&self, name: &str) -> Result<(), SQLError> {
let mut portals = self.session.portals.lock();
let Some(portal) = portals.get(name) else {
return Err(cursor_error(name, "does not exist", "34000"));
};
if portal.pin_count != 0 {
return Err(SQLError::Routine {
sqlstate: "24000".into(),
message: format!("cannot drop pinned portal \"{name}\""),
});
}
portals.remove(name);
Ok(())
}
pub(crate) fn close_all_session_portals(&self) {
self.session.portals.lock().clear();
}
fn capture_session_portal_table_sources(
&self,
fixed_snapshot: Option<&crate::FixedTransactionSnapshot>,
dependencies: &SessionPortalTableDependencies,
) -> Vec<SessionPortalTableSource> {
let live_tables = self
.storage
.tables
.read()
.iter()
.filter(|(relation, _)| dependencies.includes(relation))
.map(|(relation, metadata)| (relation.clone(), std::sync::Arc::clone(metadata)))
.collect::<Vec<_>>();
live_tables
.into_iter()
.map(|(relation, metadata)| {
let data = fixed_snapshot
.and_then(|snapshot| snapshot.table_for_live_relation(&relation, &metadata))
.unwrap_or_else(|| std::sync::Arc::clone(&metadata));
(relation, data, metadata)
})
.collect()
}
fn detach_session_portal_table_snapshots(
sources: Vec<SessionPortalTableSource>,
transaction_overlay: &std::collections::BTreeMap<
String,
std::collections::BTreeMap<crate::DocId, Option<StoredDocument>>,
>,
) -> Result<SessionPortalTableSnapshots, SQLError> {
let mut snapshots = std::collections::BTreeMap::new();
for (relation, data, metadata) in sources {
let canonical = relation.qualified_name();
let changes = transaction_overlay.get(&canonical);
snapshots.insert(
relation,
Self::detach_query_table(&data, &metadata, changes)?,
);
}
Ok(std::sync::Arc::new(snapshots))
}
pub(crate) fn capture_detached_fixed_transaction_snapshot(
&self,
) -> Result<SessionPortalTableSnapshots, SQLError> {
let mut snapshots = std::collections::BTreeMap::new();
let live_tables = self
.storage
.tables
.read()
.iter()
.filter(|(_, table)| table.persistence != uqa_sql::ast::RelationPersistence::Temporary)
.map(|(relation, table)| (relation.clone(), std::sync::Arc::clone(table)))
.collect::<Vec<_>>();
for (relation, table) in live_tables {
snapshots.insert(relation, Self::detach_query_table(&table, &table, None)?);
}
Ok(std::sync::Arc::new(snapshots))
}
fn detached_documents(
data: &std::sync::Arc<TableState>,
changes: Option<&std::collections::BTreeMap<crate::DocId, Option<StoredDocument>>>,
) -> Result<std::collections::BTreeMap<crate::DocId, StoredDocument>, SQLError> {
let source_store = data.document_store.read();
let doc_ids = source_store
.doc_ids()
.map_err(|error| portal_snapshot_error("document ids", &error))?;
let mut documents = source_store
.get_stored_many(&doc_ids)
.map_err(|error| portal_snapshot_error("documents", &error))?;
drop(source_store);
if let Some(changes) = changes {
for (doc_id, document) in changes {
match document {
Some(document) => {
documents.insert(*doc_id, document.clone());
}
None => {
documents.remove(doc_id);
}
}
}
}
Ok(documents)
}
fn detached_inverted_index(
data: &std::sync::Arc<TableState>,
analyzer: &crate::Analyzer,
fts_fields: &[crate::FieldName],
documents: &std::collections::BTreeMap<crate::DocId, StoredDocument>,
) -> Result<MemoryInvertedIndex, SQLError> {
let mut inverted_index = MemoryInvertedIndex::new(analyzer.clone());
{
let source_index = data.inverted_index.read();
for field in fts_fields {
inverted_index
.set_field_analyzer(
field,
source_index.get_field_analyzer(field),
AnalyzerPhase::Index,
)
.map_err(|error| portal_snapshot_error("index analyzer", &error))?;
inverted_index
.set_field_analyzer(
field,
source_index.get_search_analyzer(field),
AnalyzerPhase::Search,
)
.map_err(|error| portal_snapshot_error("search analyzer", &error))?;
}
}
for (doc_id, document) in documents {
let fields = fts_fields
.iter()
.filter_map(|field| match document.fields().get(field) {
Some(Value::Str(value)) => Some((field.clone(), value.clone())),
_ => None,
})
.collect();
inverted_index
.add_document(*doc_id, fields)
.map_err(|error| portal_snapshot_error("inverted index", &error))?;
}
Ok(inverted_index)
}
fn detached_vector_indexes(
metadata: &std::sync::Arc<TableState>,
documents: &std::collections::BTreeMap<crate::DocId, StoredDocument>,
) -> Result<std::collections::BTreeMap<crate::FieldName, Box<dyn VectorIndex>>, SQLError> {
let mut vector_indexes: std::collections::BTreeMap<crate::FieldName, Box<dyn VectorIndex>> =
std::collections::BTreeMap::new();
for (field, source_index) in metadata.vector_indexes.read().iter() {
let mut index = MemoryVectorIndex::new(source_index.dimensions());
for (doc_id, document) in documents {
let Some(value) = document.fields().get(field) else {
continue;
};
if let Some(vectors) = Self::field_index_vectors(metadata, field, value)? {
index
.add_many(*doc_id, vectors)
.map_err(|error| portal_snapshot_error("vector index", &error))?;
}
}
vector_indexes.insert(field.clone(), Box::new(index));
}
Ok(vector_indexes)
}
pub(crate) fn detach_query_table(
data: &std::sync::Arc<TableState>,
metadata: &std::sync::Arc<TableState>,
changes: Option<&std::collections::BTreeMap<crate::DocId, Option<StoredDocument>>>,
) -> Result<std::sync::Arc<TableState>, SQLError> {
let documents = Self::detached_documents(data, changes)?;
Self::detached_query_table_from_documents(data, metadata, &documents)
}
pub(crate) fn detach_empty_query_table(
metadata: &std::sync::Arc<TableState>,
) -> Result<std::sync::Arc<TableState>, SQLError> {
Self::detached_query_table_from_documents(
metadata,
metadata,
&std::collections::BTreeMap::new(),
)
}
fn detached_query_table_from_documents(
data: &std::sync::Arc<TableState>,
metadata: &std::sync::Arc<TableState>,
documents: &std::collections::BTreeMap<crate::DocId, StoredDocument>,
) -> Result<std::sync::Arc<TableState>, SQLError> {
let data_columns = data.columns.read().clone();
let metadata_columns = metadata.columns.read().clone();
let metadata_by_id = metadata_columns
.iter()
.filter_map(|column| column.object_id.map(|object_id| (object_id, column)))
.collect::<std::collections::BTreeMap<_, _>>();
let mut adapted_documents = std::collections::BTreeMap::new();
for (doc_id, source_document) in documents {
let mut document = source_document.clone();
let fields = document.fields_mut();
for source_column in &data_columns {
let target = source_column
.object_id
.and_then(|object_id| metadata_by_id.get(&object_id).copied())
.or_else(|| {
metadata_columns.iter().find(|target| {
target.name == source_column.name
&& (source_column.object_id.is_none() || target.object_id.is_none())
})
});
let Some(target) = target else {
fields.remove(&source_column.name);
continue;
};
if target.name != source_column.name {
if let Some(value) = fields.remove(&source_column.name) {
fields.entry(target.name.clone()).or_insert(value);
}
}
}
for target in &metadata_columns {
if target.generated.is_none() && !fields.contains_key(&target.name) {
fields.insert(
target.name.clone(),
target.missing_value.clone().unwrap_or(Value::Null),
);
}
}
crate::engine_generated::materialize_missing_generated_columns(
&metadata_columns,
fields,
)?;
adapted_documents.insert(*doc_id, document);
}
let mut document_store = MemoryDocumentStore::new();
for (doc_id, document) in &adapted_documents {
document_store
.put_stored(*doc_id, document.clone())
.map_err(|error| portal_snapshot_error("memory document", &error))?;
}
let analyzer = metadata.analyzer.read().clone();
let fts_fields = metadata.fts_fields.read().clone();
let inverted_index =
Self::detached_inverted_index(metadata, &analyzer, &fts_fields, &adapted_documents)?;
let vector_indexes = Self::detached_vector_indexes(metadata, &adapted_documents)?;
Ok(std::sync::Arc::new(TableState {
lifecycle_id: std::sync::atomic::AtomicU64::new(metadata.lifecycle_id()),
object_id: metadata.object_id(),
security: crate::engine_state::CatalogCell::new(metadata.security()),
storage_generation: parking_lot::RwLock::new(metadata.storage_generation()),
document_store: parking_lot::RwLock::new(Box::new(document_store)),
inverted_index: parking_lot::RwLock::new(Box::new(inverted_index)),
vector_indexes: parking_lot::RwLock::new(vector_indexes),
fts_fields: crate::engine_state::CatalogCell::new(fts_fields),
columns: crate::engine_state::CatalogCell::new(metadata_columns),
next_id: parking_lot::Mutex::new(*metadata.next_id.lock()),
analyzer: crate::engine_state::CatalogCell::new(analyzer),
column_stats: crate::engine_state::CatalogCell::from_snapshot(
metadata.column_stats.snapshot(),
),
column_stats_loaded: std::sync::atomic::AtomicBool::new(
metadata
.column_stats_loaded
.load(std::sync::atomic::Ordering::Acquire),
),
column_stats_dirty: std::sync::atomic::AtomicBool::new(
metadata
.column_stats_dirty
.load(std::sync::atomic::Ordering::Acquire),
),
table_checks: crate::engine_state::CatalogCell::from_snapshot(
metadata.table_checks.snapshot(),
),
foreign_keys: crate::engine_state::CatalogCell::from_snapshot(
metadata.foreign_keys.snapshot(),
),
key_constraints: crate::engine_state::CatalogCell::from_snapshot(
metadata.key_constraints.snapshot(),
),
hierarchy: crate::engine_state::CatalogCell::from_snapshot(
metadata.hierarchy.snapshot(),
),
value_indexes: parking_lot::RwLock::new(std::collections::BTreeMap::new()),
doc_count_cache: std::sync::atomic::AtomicU64::new(
u64::try_from(adapted_documents.len()).unwrap_or(u64::MAX),
),
doc_count_dirty: std::sync::atomic::AtomicBool::new(false),
persistence: metadata.persistence,
on_commit: metadata.on_commit,
}))
}
fn capture_session_portal_transaction_overlay(
&self,
) -> Result<SessionPortalTransactionOverlay, SQLError> {
let relation_names = self
.storage
.tables
.read()
.iter()
.map(|(relation, table)| (table.storage_generation(), relation.qualified_name()))
.collect::<std::collections::BTreeMap<_, _>>();
let desired = {
let stack = self.session.transactions.lock();
let mut desired = std::collections::BTreeMap::<
String,
std::collections::BTreeMap<crate::DocId, bool>,
>::new();
for change in stack.iter().flat_map(|frame| frame.row_changes.iter()) {
if let Some(table) = relation_names.get(&change.source_generation) {
desired.entry(table.clone()).or_default().insert(
change.pending.key.doc_id,
!matches!(
change.pending.kind,
crate::row_locks::PendingRowChangeKind::Delete
| crate::row_locks::PendingRowChangeKind::Rewrite(_)
),
);
}
if let crate::row_locks::PendingRowChangeKind::Rewrite(successor) =
change.pending.kind
{
if let Some(table) = change
.successor_generation
.and_then(|generation| relation_names.get(&generation))
{
desired
.entry(table.clone())
.or_default()
.insert(successor.doc_id, true);
}
}
}
desired
};
let mut overlay = std::collections::BTreeMap::new();
for (table_name, desired_documents) in desired {
let table = self.require_table(&table_name)?;
let present = desired_documents
.iter()
.filter_map(|(doc_id, present)| present.then_some(*doc_id))
.collect::<Vec<_>>();
let documents = table
.document_store
.read()
.get_stored_many(&present)
.map_err(|error| portal_snapshot_error("transaction documents", &error))?;
overlay.insert(
table_name,
desired_documents
.into_iter()
.map(|(doc_id, present)| {
let document = present.then(|| documents.get(&doc_id).cloned()).flatten();
(doc_id, document)
})
.collect(),
);
}
Ok(std::sync::Arc::new(overlay))
}
}
mod dependencies;
use dependencies::{bind_session_portal_query_relations, session_portal_table_dependencies};
impl Engine {
pub(crate) fn fork_session_portal_worker_engine(&self) -> Result<Engine, SQLError> {
let table_snapshots = self.query_table_snapshots.clone().ok_or_else(|| {
SQLError::Internal("directional query branch has no table snapshot".into())
})?;
let view_snapshots = self.query_view_snapshots.clone().ok_or_else(|| {
SQLError::Internal("directional query branch has no view snapshot".into())
})?;
let sql_function_snapshots =
self.query_sql_function_snapshots.clone().ok_or_else(|| {
SQLError::Internal("directional query branch has no SQL function snapshot".into())
})?;
let catalog_snapshot = self.query_catalog_snapshot.clone().ok_or_else(|| {
SQLError::Internal("directional query branch has no catalog snapshot".into())
})?;
let transaction_origin = self.query_transaction_origin.ok_or_else(|| {
SQLError::Internal("directional query branch has no transaction origin".into())
})?;
Ok(self.session_portal_worker_engine(
table_snapshots,
view_snapshots,
sql_function_snapshots,
catalog_snapshot,
transaction_origin,
))
}
fn session_portal_worker_engine(
&self,
table_snapshots: SessionPortalTableSnapshots,
view_snapshots: SessionPortalViewSnapshots,
sql_function_snapshots: SessionPortalSQLFunctionSnapshots,
catalog_snapshot: SessionPortalCatalogSnapshot,
transaction_origin: u64,
) -> Engine {
let mut epochs = EpochCoordinator::new();
epochs.share_published_from(&self.epochs);
let mut runtime = QueryRuntime::new(self.sql_function_depth_limit());
runtime.statement_gate = std::sync::Arc::clone(&self.runtime.statement_gate);
runtime.cancellation = self.runtime.cancellation.clone();
runtime.notices = std::sync::Arc::clone(&self.runtime.notices);
runtime.notifications = std::sync::Arc::clone(&self.runtime.notifications);
Engine {
storage: StorageContext::shared_from(&self.storage),
durable: std::sync::Arc::clone(&self.durable),
session: std::sync::Arc::clone(&self.session),
extensions: RuntimeExtensions::shared_from(&self.extensions),
epochs,
runtime,
row_locks: std::sync::Arc::clone(&self.row_locks),
notification_hub: std::sync::Arc::clone(&self.notification_hub),
session_id: self.session_id,
owns_session_registration: false,
query_table_snapshots: Some(table_snapshots),
query_view_snapshots: Some(view_snapshots),
query_sql_function_snapshots: Some(sql_function_snapshots),
query_catalog_snapshot: Some(catalog_snapshot),
query_transaction_overlay: Some(std::sync::Arc::new(std::collections::BTreeMap::new())),
query_transaction_origin: Some(transaction_origin),
}
}
}
fn cursor_error(name: &str, message: &str, sqlstate: &str) -> SQLError {
SQLError::Routine {
sqlstate: sqlstate.into(),
message: format!("cursor \"{name}\" {message}"),
}
}
fn portal_snapshot_error(component: &str, error: &impl std::fmt::Display) -> SQLError {
SQLError::Internal(format!("capture cursor {component} snapshot: {error}"))
}