use nodedb_types::DatabaseId;
use crate::control::security::catalog::{StoredCollection, StoredMaterializedView};
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::state::SharedState;
use super::super::super::catalog::propose_and_apply;
use super::super::super::result::{DdlError, DdlResult};
fn err(sqlstate: &str, message: String) -> DdlError {
DdlError {
sqlstate: sqlstate.to_string(),
message,
}
}
pub async fn create_materialized_view(
state: &SharedState,
identity: &AuthenticatedIdentity,
name: &str,
source: &str,
query_sql: &str,
refresh_mode: &str,
) -> Result<Vec<DdlResult>, DdlError> {
let name = name.to_string();
let source = source.to_string();
let query_sql = query_sql.to_string();
let refresh_mode = refresh_mode.to_string();
let tenant_id = identity.tenant_id;
if refresh_mode.eq_ignore_ascii_case("STREAMING") {
return create_streaming_mv(state, identity, &name, &query_sql).await;
}
let _local_lifecycle = if state.metadata_raft.get().is_none() {
Some(
state
.quiesce
.acquire_lifecycle(0, tenant_id.as_u64(), &name)
.await,
)
} else {
None
};
{
let catalog = state.credentials.catalog();
match catalog.get_collection(DatabaseId::DEFAULT, tenant_id.as_u64(), &source) {
Ok(Some(_)) => {}
_ => {
return Err(err(
"42P01",
format!("source collection '{source}' does not exist"),
));
}
}
if catalog
.get_materialized_view(tenant_id.as_u64(), &name)
.map_err(|error| err("XX000", error.to_string()))?
.is_some()
{
return Err(err(
"42P07",
format!("materialized view '{name}' already exists"),
));
}
if catalog
.get_collection(DatabaseId::DEFAULT, tenant_id.as_u64(), &name)
.map_err(|error| err("XX000", error.to_string()))?
.is_some()
{
return Err(err("42P07", format!("collection '{name}' already exists")));
}
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let view = StoredMaterializedView {
tenant_id: tenant_id.as_u64(),
name: name.clone(),
source: source.clone(),
query_sql,
refresh_mode,
owner: identity.username.clone(),
created_at: now,
descriptor_version: 0,
modification_hlc: nodedb_types::Hlc::ZERO,
};
let entry =
crate::control::catalog_entry::CatalogEntry::PutMaterializedView(Box::new(view.clone()));
propose_and_apply(state, &entry)?;
let target = StoredCollection {
tenant_id: tenant_id.as_u64(),
name: name.clone(),
owner: identity.username.clone(),
created_at: now,
descriptor_version: 0,
constraint_version: 0,
modification_hlc: nodedb_types::Hlc::ZERO,
fields: Vec::new(),
field_defs: Vec::new(),
event_defs: Vec::new(),
collection_type: nodedb_types::CollectionType::document(),
timeseries_config: None,
conflict_policy: None,
is_active: true,
append_only: false,
hash_chain: false,
balanced: None,
last_chain_hash: None,
period_lock: None,
retention_period: None,
legal_holds: Vec::new(),
state_constraints: Vec::new(),
transition_checks: Vec::new(),
type_guards: Vec::new(),
check_constraints: Vec::new(),
materialized_sums: Vec::new(),
lvc_enabled: false,
bitemporal: false,
crdt: false,
permission_tree_def: None,
indexes: Vec::new(),
size_bytes_estimate: 0,
primary: nodedb_types::PrimaryEngine::Document,
vector_primary: None,
partition_strategy: nodedb_types::PartitionStrategy::CollectionHomed,
database_id: nodedb_types::DatabaseId::DEFAULT,
cloned_from: None,
clone_status: nodedb_types::CloneStatus::default(),
has_implicit_edges: false,
declared_primary_key: None,
};
let coll_entry =
crate::control::catalog_entry::CatalogEntry::PutCollection(Box::new(target.clone()));
propose_and_apply(state, &coll_entry)?;
super::super::collection::dispatch_register_from_stored(state, &target)
.await
.map_err(|e| err("XX000", e.to_string()))?;
tracing::info!(
view = name,
source,
tenant = tenant_id.as_u64(),
"materialized view created"
);
Ok(vec![DdlResult::Status {
command: "CREATE MATERIALIZED VIEW".to_string(),
rows_affected: None,
}])
}
async fn create_streaming_mv(
state: &SharedState,
identity: &AuthenticatedIdentity,
name: &str,
query_sql: &str,
) -> Result<Vec<DdlResult>, DdlError> {
super::super::auth_support::require_tenant_admin(
identity,
"create streaming materialized views",
)?;
let parsed = super::streaming_parse::parse_streaming_mv(query_sql)?;
let tenant_id = identity.tenant_id.as_u64();
if state
.stream_registry
.get(tenant_id, &parsed.source_stream)
.is_none()
{
return Err(err(
"42704",
format!("change stream '{}' does not exist", parsed.source_stream),
));
}
if state.mv_registry.get_def(tenant_id, name).is_some() {
return Err(err(
"42710",
format!("streaming MV '{name}' already exists"),
));
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let source_stream = parsed.source_stream.clone();
let def = crate::event::streaming_mv::types::StreamingMvDef {
tenant_id,
name: name.to_string(),
source_stream: parsed.source_stream,
group_by_columns: parsed.group_by_columns,
aggregates: parsed.aggregates,
filter_expr: parsed.filter_expr,
owner: identity.username.clone(),
created_at: now,
};
let catalog = state.credentials.catalog();
catalog
.put_streaming_mv(&def)
.map_err(|e| err("XX000", format!("catalog write: {e}")))?;
state.mv_registry.register(def);
if let Some(mv_state) = state.mv_registry.get_state(tenant_id, name)
&& let Some(buffer) = state.cdc_router.get_buffer(tenant_id, &source_stream)
{
crate::event::streaming_mv::processor::backfill_from_buffer(&mv_state, &buffer);
}
state.audit_record(
crate::control::security::audit::AuditEvent::AdminAction,
Some(identity.tenant_id),
&identity.username,
&format!("CREATE MATERIALIZED VIEW {name} STREAMING"),
);
tracing::info!(
view = name,
stream = source_stream,
tenant = tenant_id,
"streaming materialized view created"
);
Ok(vec![DdlResult::Status {
command: "CREATE MATERIALIZED VIEW".to_string(),
rows_affected: None,
}])
}