use std::{collections::BTreeSet, sync::Arc};
use reifydb_core::interface::{
catalog::{
id::{NamespaceId, ViewId},
shape::ShapeId,
},
resolved::ResolvedView,
};
use reifydb_rql::flow::{analyzer::FlowGraphAnalyzer, loader::load_flow_dag};
use reifydb_transaction::{error::TransactionError, transaction::Transaction};
use crate::{Result, vm::services::Services};
pub mod dictionary;
pub mod index;
pub mod remote;
pub mod ringbuffer;
pub mod series;
pub mod table;
pub mod view;
pub mod vtable;
pub(crate) fn guard_view_read(view: &ResolvedView, rx: &mut Transaction<'_>, services: &Services) -> Result<()> {
if !rx.has_unprocessed_flow_changes() {
return Ok(());
}
let upstream = match services.view_lineage.upstream_of(view.def().id()) {
Some(upstream) => upstream,
None => match upstream_from_catalog(services, rx, view.def().id())? {
Some(upstream) => Arc::new(upstream),
None => return Ok(()),
},
};
let offending: Vec<ShapeId> =
rx.unprocessed_flow_change_shapes().into_iter().filter(|shape| upstream.contains(shape)).collect();
if offending.is_empty() {
return Ok(());
}
Err(TransactionError::ViewPendingUpstreamChanges {
view: view.fully_qualified_name(),
kind: view.def().kind(),
upstream: resolve_shape_names(services, rx, &offending),
fragment: view.identifier().clone(),
}
.into())
}
fn upstream_from_catalog(
services: &Services,
rx: &mut Transaction<'_>,
view: ViewId,
) -> Result<Option<BTreeSet<ShapeId>>> {
let mut dags = Vec::new();
for flow in services.catalog.list_flows_all(rx)? {
dags.push(load_flow_dag(rx, flow.id)?);
}
let mut analyzer = FlowGraphAnalyzer::new();
analyzer.add_all(dags);
Ok(analyzer.get_dependency_graph().upstream_closure().remove(&view))
}
fn resolve_shape_names(services: &Services, rx: &mut Transaction<'_>, shapes: &[ShapeId]) -> Vec<String> {
let catalog = &services.catalog;
shapes.iter()
.map(|shape| {
let named = match shape {
ShapeId::Table(id) => catalog
.find_table(rx, *id)
.ok()
.flatten()
.map(|def| ("table", def.namespace, def.name)),
ShapeId::View(id) => catalog
.find_view(rx, *id)
.ok()
.flatten()
.map(|def| ("view", def.namespace(), def.name().to_string())),
ShapeId::RingBuffer(id) => catalog
.find_ringbuffer(rx, *id)
.ok()
.flatten()
.map(|def| ("ring buffer", def.namespace, def.name)),
ShapeId::Series(id) => catalog
.find_series(rx, *id)
.ok()
.flatten()
.map(|def| ("series", def.namespace, def.name)),
ShapeId::Dictionary(id) => catalog
.find_dictionary(rx, *id)
.ok()
.flatten()
.map(|def| ("dictionary", def.namespace, def.name)),
ShapeId::TableVirtual(_) => None,
};
match named {
Some((kind, namespace, name)) => {
format!("{} '{}'", kind, qualify(services, rx, namespace, &name))
}
None => format!("shape {}", shape),
}
})
.collect()
}
fn qualify(services: &Services, rx: &mut Transaction<'_>, namespace: NamespaceId, name: &str) -> String {
match services.catalog.find_namespace(rx, namespace) {
Ok(Some(namespace)) => format!("{}::{}", namespace.name(), name),
_ => name.to_string(),
}
}