use super::*;
const COMPENSATION_RETRY_MAX_ATTEMPTS: u32 = 3;
const COMPENSATION_RETRY_BASE_DELAY_MS: u64 = 200;
const MAX_COMPENSATIONS_PER_TX: usize = 1000;
#[derive(Debug, Clone)]
struct PendingSagaStep {
step_index: usize,
operation: String,
message_type: String,
compensation_json: String,
}
fn tx_object_invalid_field(
field: impl Into<String>,
description: impl Into<String>,
message: impl Into<String>,
) -> tonic::Status {
crate::runtime::executor_utils::invalid_argument_fields(
message,
[(field.into(), description.into())],
)
}
fn tx_object_internal_status(
operation: impl Into<String>,
message: impl Into<String>,
) -> tonic::Status {
crate::runtime::executor_utils::internal_status("tx_object", operation, message)
}
fn mysql_xa_plan_replay_status(err: impl std::fmt::Display) -> tonic::Status {
crate::runtime::executor_utils::schema_status(
tonic::Code::FailedPrecondition,
"mysql",
"xa_plan_replay",
"mysql_xa_plan_replay",
format!(
"2PC refused before side effects: plan SQL cannot be replayed on the MySQL XA participant: {err}"
),
)
}
fn xa_unsupported_participants_status(unsupported: &[String]) -> tonic::Status {
crate::runtime::executor_utils::capability_status(
"transaction",
"xa_commit",
"xa_participants",
format!(
"2PC refused before side effects: unsupported participants {}",
unsupported.join(", ")
),
)
}
fn record_encryption_key_missing_status(schema: &str, table: &str) -> tonic::Status {
crate::runtime::executor_utils::capability_status(
"encryption",
"record_encryption",
"udb_encryption_key",
format!(
"table {schema}.{table} contains encrypted columns but UDB encryption key is not configured"
),
)
}
#[cfg_attr(feature = "s3", allow(dead_code))]
fn s3_feature_disabled_status() -> tonic::Status {
crate::runtime::executor_utils::capability_status(
"s3",
"put_tx_object",
"s3_feature",
"s3/object-store feature is not enabled",
)
}
fn materialized_view_admin_scope_required_status() -> tonic::Status {
crate::runtime::executor_utils::policy_status_with_code(
tonic::Code::PermissionDenied,
"create_materialized_view",
"admin_scope_required",
"scope udb:admin is required",
)
}
fn validate_tx_identifier(value: &str, label: &str) -> Result<(), tonic::Status> {
if value.is_empty()
|| !value
.chars()
.all(|ch| ch.is_ascii_alphanumeric() || ch == '_')
|| value.starts_with(|ch: char| ch.is_ascii_digit())
{
return Err(tx_object_invalid_field(
label,
"must be a valid SQL identifier",
format!("{label} '{value}' is not a valid SQL identifier"),
));
}
Ok(())
}
fn compensation_retry_delay(attempt: u32) -> std::time::Duration {
std::time::Duration::from_millis(COMPENSATION_RETRY_BASE_DELAY_MS * u64::from(attempt))
}
fn saga_compensation_for_mutation(operation: &str, mutation: &Mutation) -> String {
let compensation = match operation {
"vector_upsert" => serde_json::json!({
"backend": "qdrant",
"operation": "delete_points",
"resource_uri": format!("qdrant://{}", mutation.collection),
"payload": {
"collection": mutation.collection,
"point_ids": mutation.vector_points.iter()
.map(|point| point.id.clone())
.filter(|id| !id.is_empty())
.collect::<Vec<_>>()
}
}),
"put_object" => serde_json::json!({
"backend": "s3",
"operation": "delete_object",
"resource_uri": format!("s3://{}/{}", mutation.bucket, mutation.object_key),
"payload": {
"bucket": mutation.bucket,
"key": mutation.object_key,
}
}),
"enqueue_outbox_event" => serde_json::json!({
"backend": "postgres",
"operation": "delete_outbox_event",
"resource_uri": "udb://outbox/event",
"payload": {
"idempotency_key": mutation.idempotency_key,
"message_type": mutation.message_type,
}
}),
_ => serde_json::json!({
"backend": "postgres",
"operation": "rollback_sql_tx",
"resource_uri": format!("postgres://{}", mutation.message_type),
"payload": {
"message_type": mutation.message_type,
"note": "rolled back by PostgreSQL transaction boundary"
}
}),
};
compensation.to_string()
}
impl DataBrokerRuntime {
pub async fn begin_tx(
&self,
manifest: &CatalogManifest,
mut stream: tonic::Streaming<Mutation>,
metadata_context: RequestContext,
) -> Vec<Result<TxStatus, tonic::Status>> {
let Some(pool) = &self.pg_pool else {
return vec![Err(crate::runtime::executor_utils::capability_status(
"postgres",
"begin_tx",
"postgres_backend",
"PostgreSQL backend is not configured",
))];
};
let mut mutations = Vec::new();
while let Some(item) = stream.next().await {
match item {
Ok(mutation) => mutations.push(mutation),
Err(err) => return vec![Err(err)],
}
}
if mutations.is_empty() {
return vec![Err(tx_object_invalid_field(
"mutations",
"transaction stream must contain at least one mutation",
"transaction stream requires at least one mutation",
))];
}
let tx_id = mutations
.iter()
.find(|mutation| !mutation.tx_id.is_empty())
.map(|mutation| mutation.tx_id.clone())
.unwrap_or_else(|| Uuid::new_v4().to_string());
if mutations.iter().any(|mutation| mutation.rollback) {
return vec![Ok(TxStatus {
state: crate::proto::tx_status::State::TxStateRolledBack as i32,
tx_id,
message: "rollback requested".to_string(),
..TxStatus::default()
})];
}
let strategy = match requested_tx_strategy(&metadata_context, &mutations) {
Ok(s) => s,
Err(err) => return vec![Err(err)],
};
if let Err(err) = validate_tx_strategy(strategy, &mutations) {
return vec![Err(err)];
}
let commit = mutations.iter().any(|mutation| mutation.commit);
let mutation_count = mutations
.iter()
.filter(|m| !m.commit && !m.rollback)
.count();
let saga_id = self
.saga_begin(
&tx_id,
mutation_count,
&metadata_context.tenant_id,
&metadata_context.correlation_id,
classify_tx_semantics(&mutations).as_str(),
tx_backend_instance(&mutations).unwrap_or_default().as_str(),
)
.await;
let mut tx = match pool.begin().await {
Ok(tx) => tx,
Err(err) => {
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "failed").await;
}
return vec![Err(crate::runtime::executor_utils::retryable_status(
"postgres",
"transaction_begin",
crate::runtime::executor_utils::HTTP_RETRYABLE_BACKOFF_MS,
format!("PostgreSQL begin failed: {err}"),
))];
}
};
let mut statuses = Vec::new();
let mut qdrant_compensations: Vec<(String, Vec<String>)> = Vec::new();
let mut s3_compensations: Vec<(String, String, Option<String>, String)> = Vec::new();
let mut compensations_spilled = false;
let projection_plans = crate::runtime::projection::ProjectionPlan::from_manifest(manifest);
let projection_config = crate::runtime::system::SystemCatalogConfig::current_arc();
let tx_mutations: Vec<&Mutation> = mutations
.iter()
.filter(|mutation| !mutation.commit)
.collect();
let mut pending_saga_steps: Vec<PendingSagaStep> = Vec::new();
#[cfg(feature = "mysql")]
let mysql_xa_capture = commit
&& strategy == TxStrategy::TwoPhase
&& super::two_phase_runtime_enabled()
&& !self.mysql_instances.is_empty();
#[cfg(feature = "mysql")]
let mut mysql_xa_statements: Vec<(String, Vec<JsonValue>)> = Vec::new();
for (mutation_index, mutation) in tx_mutations.iter().enumerate() {
let context = merge_context(mutation.context.as_ref(), metadata_context.clone());
let operation = mutation.operation.to_ascii_lowercase();
let result = if operation == "upsert" {
let record = mutation_record_json(mutation);
match record {
Ok(record) => {
let encrypted_record =
match resolve_table_for_message(manifest, &mutation.message_type) {
Ok(table) => {
let record =
crate::broker::normalize_record_keys(table, &record);
self.encrypt_record_for_table(table, &record)
}
Err(error) => Err(tx_object_invalid_field(
"message_type",
"must match exactly one manifest table message type",
error.to_string(),
)),
};
match encrypted_record {
Ok(encrypted_record) => {
let plan = build_upsert_plan(
manifest,
&UpsertPlanRequest {
context: context.clone(),
message_type: mutation.message_type.clone(),
record: record.clone(),
..UpsertPlanRequest::default()
},
);
let bind_values =
record_values(&encrypted_record, &plan.parameter_columns)
.unwrap_or_default();
let affected = execute_tx_plan(
&mut tx,
manifest,
&mutation.message_type,
&plan.sql,
&plan.parameter_columns,
&bind_values,
&plan.errors,
)
.await;
match affected {
Ok(affected) => {
#[cfg(feature = "mysql")]
if mysql_xa_capture {
mysql_xa_statements
.push((plan.sql.clone(), bind_values.clone()));
}
if affected > 0
&& let Err(err) = crate::runtime::projection::ProjectionEngine::enqueue_write_tasks_tx(
&mut tx,
&projection_config,
&context.tenant_id,
&mutation.message_type,
"upsert",
&record,
&projection_plans,
)
.await
{
Err(tx_object_internal_status(
"enqueue_projection_tasks",
format!("projection task enqueue failed: {err}"),
))
} else {
Ok(affected)
}
}
Err(err) => Err(err),
}
}
Err(err) => Err(err),
}
}
Err(err) => Err(err),
}
} else if operation == "delete" {
let filter = mutation
.filter
.as_ref()
.map(struct_to_json)
.unwrap_or(JsonValue::Null);
let plan = build_delete_plan(
manifest,
&DeletePlanRequest {
context: context.clone(),
message_type: mutation.message_type.clone(),
filter: filter.clone(),
},
);
let bind_values = filter_bind_values(&filter);
let affected = execute_tx_plan(
&mut tx,
manifest,
&mutation.message_type,
&plan.sql,
&plan.parameter_columns,
&bind_values,
&plan.errors,
)
.await;
match affected {
Ok(affected) => {
#[cfg(feature = "mysql")]
if mysql_xa_capture {
mysql_xa_statements.push((plan.sql.clone(), bind_values.clone()));
}
if affected > 0
&& let Err(err) = crate::runtime::projection::ProjectionEngine::enqueue_write_tasks_tx(
&mut tx,
&projection_config,
&context.tenant_id,
&mutation.message_type,
"delete",
&filter,
&projection_plans,
)
.await
{
Err(tx_object_internal_status(
"enqueue_projection_tasks",
format!("projection task enqueue failed: {err}"),
))
} else {
Ok(affected)
}
}
Err(err) => Err(err),
}
} else if operation == "vector_upsert" {
let request = VectorUpsertRequest {
context: mutation.context.clone(),
collection: mutation.collection.clone(),
points: mutation.vector_points.clone(),
idempotency_key: mutation.idempotency_key.clone(),
};
let point_ids = request
.points
.iter()
.map(|point| point.id.clone())
.filter(|id| !id.is_empty())
.collect::<Vec<_>>();
match self
.vector_upsert(manifest, request, metadata_context.clone())
.await
{
Ok(response) => {
if !point_ids.is_empty() {
if qdrant_compensations.len() < MAX_COMPENSATIONS_PER_TX {
qdrant_compensations.push((mutation.collection.clone(), point_ids));
} else if !compensations_spilled {
compensations_spilled = true;
tracing::warn!(
tx_id = %tx_id,
cap = MAX_COMPENSATIONS_PER_TX,
"compensation buffer cap reached; overflow spills to the durable saga ledger for recovery-worker compensation"
);
}
}
Ok(response.affected_rows as u64)
}
Err(err) => Err(err),
}
} else if operation == "put_object" {
match self
.put_tx_object(manifest, mutation, metadata_context.clone())
.await
{
Ok((affected, instance, project_id)) => {
if s3_compensations.len() < MAX_COMPENSATIONS_PER_TX {
s3_compensations.push((
mutation.bucket.clone(),
mutation.object_key.clone(),
instance,
project_id,
));
} else if !compensations_spilled {
compensations_spilled = true;
tracing::warn!(
tx_id = %tx_id,
cap = MAX_COMPENSATIONS_PER_TX,
"compensation buffer cap reached; overflow spills to the durable saga ledger for recovery-worker compensation"
);
}
Ok(affected)
}
Err(err) => Err(err),
}
} else if operation == "enqueue_outbox_event" {
let cdc_config = self.config.cdc.clone();
let topic = if mutation.collection.trim().is_empty() {
mutation.message_type.as_str()
} else {
mutation.collection.as_str()
};
match self
.topic_policy_allows(topic, &context.project_id, &context.tenant_id)
.await
{
Ok(policy_decision) => {
if policy_decision == Some(false)
|| (!cdc_config.valid_topics.is_empty()
&& !cdc_config.valid_topics.iter().any(|t| t == topic)
&& policy_decision.is_none())
{
Err(tx_object_invalid_field(
"topic",
"must be allowed by topic policy or UDB_CDC_VALID_TOPICS",
format!(
"topic '{topic}' is not in the registered topic registry; \
configure UDB_CDC_VALID_TOPICS or use an allowed topic"
),
))
} else {
let payload =
mutation.payload.as_ref().map(struct_to_json).ok_or_else(|| {
tx_object_invalid_field(
"payload",
"enqueue_outbox_event transaction mutations require payload",
"enqueue_outbox_event transaction mutation requires payload",
)
});
match payload {
Ok(payload) => {
let schema_uri = if mutation.content_type.trim().is_empty() {
None
} else {
Some(mutation.content_type.as_str())
};
let mut payload =
crate::runtime::cdc::apply_manifest_cdc_redaction(
manifest,
&mutation.message_type,
topic,
schema_uri,
payload,
cdc_config.redaction_mode,
cdc_config.redaction_version,
);
if let Some(obj) = payload.as_object_mut() {
if !context.tenant_id.trim().is_empty() {
obj.entry("tenant_id".to_string()).or_insert_with(
|| {
serde_json::Value::String(
context.tenant_id.clone(),
)
},
);
}
if !context.project_id.trim().is_empty() {
obj.entry("project_id".to_string()).or_insert_with(
|| {
serde_json::Value::String(
context.project_id.clone(),
)
},
);
}
}
match prepare_outbox_envelope(
topic,
&mutation.object_key,
payload,
schema_uri,
) {
Ok((event_id_uuid, _, enriched)) => {
let outbox_relation = cdc_config.outbox_relation();
let sql = format!(
"INSERT INTO {outbox_relation} \
(event_id, topic, partition_key, payload, created_at) \
VALUES ($1::UUID, $2, $3, $4::JSONB, NOW())"
);
sqlx::query(&sql)
.bind(event_id_uuid)
.bind(topic)
.bind(&mutation.object_key)
.bind(enriched.to_string())
.execute(&mut *tx)
.await
.map(|result| result.rows_affected())
.map_err(|err| {
tx_object_internal_status(
"enqueue_tx_event",
format!(
"failed to enqueue tx event: {err}"
),
)
})
}
Err(err) => Err(err),
}
}
Err(err) => Err(err),
}
}
}
Err(err) => Err(err),
}
} else {
Err(tx_object_invalid_field(
"operation",
"must be one of upsert, delete, vector_upsert, put_object, or enqueue_outbox_event",
format!("unsupported transaction operation {}", mutation.operation),
))
};
match result {
Ok(affected) => {
pending_saga_steps.push(PendingSagaStep {
step_index: pending_saga_steps.len(),
operation: operation.clone(),
message_type: mutation.message_type.clone(),
compensation_json: saga_compensation_for_mutation(&operation, mutation),
});
statuses.push(Ok(TxStatus {
state: crate::proto::tx_status::State::TxStateOpen as i32,
tx_id: tx_id.clone(),
mutation_id: Uuid::new_v4().to_string(),
message: format!("{affected} row(s) affected"),
}))
}
Err(err) => {
let _ = tx.rollback().await;
let compensation_message = self
.compensate_external_side_effects(&qdrant_compensations, &s3_compensations)
.await
.unwrap_or_else(|| "no external compensation was needed".to_string());
if let Some(ref sid) = saga_id {
if compensations_spilled {
self.saga_record_pending_steps(sid, &pending_saga_steps)
.await;
self.saga_set_status(sid, "in_doubt").await;
} else {
self.saga_set_status(sid, "compensated").await;
}
}
statuses.push(Err(tx_object_internal_status(
"mutation_failure_compensation",
format!("{}; {}", err.message(), compensation_message),
)));
for skipped in tx_mutations.iter().skip(mutation_index + 1) {
statuses.push(Ok(TxStatus {
state: crate::proto::tx_status::State::TxStateRolledBack as i32,
tx_id: tx_id.clone(),
mutation_id: Uuid::new_v4().to_string(),
message: format!(
"not executed because an earlier mutation failed: {} {}",
skipped.operation, skipped.message_type
),
}));
}
return statuses;
}
}
}
if commit {
if strategy == TxStrategy::TwoPhase && super::two_phase_runtime_enabled() {
let sys_config = crate::runtime::system::SystemCatalogConfig::current_arc();
if let Err(err) =
crate::runtime::xa_recovery::ensure_xa_ledger_table(pool, &sys_config).await
{
let _ = tx.rollback().await;
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "failed").await;
}
statuses.push(Err(tx_object_internal_status(
"xa_ledger_unavailable",
format!("XA ledger is unavailable: {err}"),
)));
return statuses;
}
#[cfg(feature = "mysql")]
let mysql_participants: Vec<
Box<dyn crate::runtime::xa::XaParticipant>,
> = if mysql_xa_capture {
let mut translated: Vec<(String, Vec<JsonValue>)> =
Vec::with_capacity(mysql_xa_statements.len());
let mut translation_error: Option<String> = None;
for (sql, values) in &mysql_xa_statements {
match crate::runtime::xa::translate_pg_plan_sql_to_mysql(sql) {
Ok(mysql_sql) => translated.push((mysql_sql, values.clone())),
Err(err) => {
translation_error = Some(err);
break;
}
}
}
if let Some(err) = translation_error {
let _ = tx.rollback().await;
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "failed").await;
}
statuses.push(Err(mysql_xa_plan_replay_status(err)));
return statuses;
}
let mut instance_names: Vec<String> =
self.mysql_instances.keys().cloned().collect();
instance_names.sort();
instance_names
.into_iter()
.filter_map(|name| {
self.mysql_instances.get(&name).map(|mysql_pool| {
Box::new(crate::runtime::xa::XaMysqlParticipant::new(
name.clone(),
mysql_pool.clone(),
translated.clone(),
))
as Box<dyn crate::runtime::xa::XaParticipant>
})
})
.collect()
} else {
Vec::new()
};
let participant = crate::runtime::xa_postgres::ActivePostgresXaParticipant::new(
"primary",
pool.clone(),
tx,
);
#[cfg_attr(not(feature = "mysql"), allow(unused_mut))]
let mut participants: Vec<
Box<dyn crate::runtime::xa::XaParticipant>,
> = vec![Box::new(participant)];
#[cfg(feature = "mysql")]
participants.extend(mysql_participants);
let xa_request = crate::runtime::xa::XaRequest {
tenant_id: metadata_context.tenant_id.clone(),
project_id: metadata_context.project_id.clone(),
origin_rpc: "BeginTx".to_string(),
correlation_id: metadata_context.correlation_id.clone(),
};
let write_ahead_pool = pool.clone();
let write_ahead_config = sys_config.clone();
let xa_result = crate::runtime::xa::XaCoordinator::execute_with_write_ahead(
xa_request,
participants,
&crate::backend::capability_matrix(),
move |entry| async move {
crate::runtime::xa_recovery::record_xa_ledger_entry(
&write_ahead_pool,
&write_ahead_config,
&entry,
)
.await
},
)
.await;
match xa_result {
Ok(outcome) => {
if let Err(err) = crate::runtime::xa_recovery::record_xa_ledger_entry(
pool,
&sys_config,
&outcome.ledger,
)
.await
{
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "in_doubt").await;
}
statuses.push(Err(tx_object_internal_status(
"xa_ledger_commit_record",
format!(
"XA committed but ledger write failed for xid {}: {err}",
outcome.ledger.xid
),
)));
return statuses;
}
if let Some(ref sid) = saga_id {
self.saga_record_pending_steps(sid, &pending_saga_steps)
.await;
self.saga_set_status(sid, "committed").await;
}
statuses.push(Ok(TxStatus {
state: crate::proto::tx_status::State::TxStateCommitted as i32,
tx_id,
message: format!("committed (2pc xid={})", outcome.ledger.xid),
..TxStatus::default()
}));
}
Err(crate::runtime::xa::XaError::CapabilityRefused { unsupported }) => {
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "failed").await;
}
statuses.push(Err(xa_unsupported_participants_status(&unsupported)));
}
Err(crate::runtime::xa::XaError::PrepareFailed { ledger, .. }) => {
let _ = crate::runtime::xa_recovery::record_xa_ledger_entry(
pool,
&sys_config,
&ledger,
)
.await;
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "failed").await;
}
statuses.push(Err(
crate::runtime::executor_utils::retryable_aborted_status(
"transaction",
"xa prepare",
0,
format!(
"2PC PREPARE failed for xid {}: {}",
ledger.xid, ledger.reason
),
),
));
}
Err(crate::runtime::xa::XaError::InDoubt { ledger }) => {
let ledger_err = crate::runtime::xa_recovery::record_xa_ledger_entry(
pool,
&sys_config,
&ledger,
)
.await
.err();
if let Some(ref sid) = saga_id {
self.saga_set_status(sid, "in_doubt").await;
}
let ledger_note = ledger_err
.map(|err| format!("; additionally XA ledger write failed: {err}"))
.unwrap_or_default();
statuses.push(Err(
crate::runtime::executor_utils::retryable_aborted_status(
"transaction",
"xa in-doubt recovery",
0,
format!(
"2PC xid {} is IN-DOUBT and will be resolved by XA recovery{}",
ledger.xid, ledger_note
),
),
));
}
}
} else {
match tx.commit().await {
Ok(()) => {
if let Some(ref sid) = saga_id {
self.saga_record_pending_steps(sid, &pending_saga_steps)
.await;
self.saga_set_status(sid, "committed").await;
}
statuses.push(Ok(TxStatus {
state: crate::proto::tx_status::State::TxStateCommitted as i32,
tx_id,
message: "committed".to_string(),
..TxStatus::default()
}))
}
Err(err) => {
let comp_msg = self
.compensate_external_side_effects(
&qdrant_compensations,
&s3_compensations,
)
.await
.unwrap_or_else(|| "no external compensation was needed".to_string());
if let Some(ref sid) = saga_id {
if compensations_spilled {
self.saga_record_pending_steps(sid, &pending_saga_steps)
.await;
self.saga_set_status(sid, "in_doubt").await;
} else {
self.saga_set_status(sid, "compensated").await;
}
}
statuses.push(Err(tx_object_internal_status(
"postgres_commit_compensation",
format!("PostgreSQL commit failed: {err}; {comp_msg}"),
)))
}
}
}
} else {
if let Err(err) = tx.rollback().await {
tracing::warn!("PostgreSQL rollback failed after stream close: {err}");
}
let compensation_message = self
.compensate_external_side_effects(&qdrant_compensations, &s3_compensations)
.await;
if let Some(ref sid) = saga_id {
if compensations_spilled {
self.saga_record_pending_steps(sid, &pending_saga_steps)
.await;
self.saga_set_status(sid, "in_doubt").await;
} else {
self.saga_set_status(sid, "compensated").await;
}
}
statuses.push(Ok(TxStatus {
state: crate::proto::tx_status::State::TxStateRolledBack as i32,
tx_id,
message: compensation_message.unwrap_or_else(|| {
"stream closed without commit=true; rolled back".to_string()
}),
..TxStatus::default()
}));
}
statuses
}
async fn saga_record_pending_steps(&self, saga_id: &str, steps: &[PendingSagaStep]) {
for step in steps {
self.saga_record_step(
saga_id,
step.step_index,
&step.operation,
&step.message_type,
&step.compensation_json,
)
.await;
}
}
pub async fn create_materialized_view(
&self,
manifest: &CatalogManifest,
request: ViewDefinition,
metadata_context: RequestContext,
) -> Result<MutationResponse, tonic::Status> {
let context = merge_context(request.context.as_ref(), metadata_context);
if !context
.scopes
.iter()
.any(|scope| scope == "udb:admin" || scope == "*")
{
return Err(materialized_view_admin_scope_required_status());
}
validate_tx_identifier(&request.schema, "schema")?;
validate_tx_identifier(&request.name, "name")?;
let declared = declared_materialized_view(manifest, &request.schema, &request.name)
.ok_or_else(|| {
tx_object_invalid_field(
"view",
"must be declared in the proto AST",
format!(
"materialized view {}.{} is not declared in the proto AST",
request.schema, request.name
),
)
})?;
let query = if request.query.trim().is_empty()
|| normalize_sql(&request.query) == normalize_sql(&declared.query)
{
declared.query.as_str()
} else {
return Err(tx_object_invalid_field(
"query",
"must match the proto AST declaration",
"materialized view query does not match the proto AST declaration",
));
};
let with_data = request.with_data || declared.with_data;
let sql = format!(
"CREATE MATERIALIZED VIEW IF NOT EXISTS \"{}\".\"{}\" AS {} {}",
request.schema.replace('"', "\"\""),
request.name.replace('"', "\"\""),
query,
if with_data { "" } else { "WITH NO DATA" }
);
let result = sqlx::query(&sql)
.execute(self.pg_pool()?)
.await
.map_err(|err| {
tx_object_internal_status(
"create_materialized_view",
format!("create view failed: {err}"),
)
})?;
let response = MutationResponse {
mutation_id: Uuid::new_v4().to_string(),
resource_uri: format!("materialized-view://{}/{}", request.schema, request.name),
affected_rows: result.rows_affected() as i64,
..MutationResponse::default()
};
if request.ttl_days > 0 {
self.schedule_materialized_view_refresh(
request.schema.clone(),
request.name.clone(),
request.ttl_days,
);
}
Ok(response)
}
pub fn start_materialized_view_refresh(&self, manifest: &CatalogManifest) -> usize {
let ttl_days = if self.config.materialized_view_ttl_days == 0 {
7
} else {
self.config.materialized_view_ttl_days
};
if ttl_days <= 0 {
return 0;
}
let mut scheduled = 0usize;
for view in manifest
.tables
.iter()
.flat_map(|table| table.materialized_views.iter())
{
self.schedule_materialized_view_refresh(
view.schema.clone(),
view.name.clone(),
ttl_days,
);
scheduled += 1;
}
scheduled
}
pub fn schedule_materialized_view_refresh(&self, schema: String, name: String, ttl_days: i32) {
if self.pg_pool.is_none() || ttl_days <= 0 {
return;
}
let runtime = self.clone();
let seconds = (ttl_days as u64).saturating_mul(86_400).max(60);
tokio::spawn(async move {
const STALE_ALERT_THRESHOLD: u32 = 3;
let mut interval = tokio::time::interval(Duration::from_secs(seconds));
let mut consecutive_failures: u32 = 0;
tracing::info!(
schema = %schema,
view = %name,
interval_secs = seconds,
"materialized view refresh scheduler started"
);
loop {
interval.tick().await;
match runtime.refresh_materialized_view(&schema, &name).await {
Ok(()) => consecutive_failures = 0,
Err(err) => {
consecutive_failures = consecutive_failures.saturating_add(1);
if consecutive_failures >= STALE_ALERT_THRESHOLD {
tracing::error!(
schema = %schema,
view = %name,
consecutive_failures,
error = %err,
"materialized view is staling: refresh has failed repeatedly — operator attention required"
);
} else {
tracing::warn!(
schema = %schema,
view = %name,
error = %err,
"materialized view refresh failed"
);
}
}
}
}
});
}
pub async fn refresh_materialized_view(
&self,
schema: &str,
name: &str,
) -> Result<(), tonic::Status> {
validate_tx_identifier(schema, "schema")?;
validate_tx_identifier(name, "name")?;
let sql = format!(
"REFRESH MATERIALIZED VIEW CONCURRENTLY {}.{}",
qi_runtime(schema),
qi_runtime(name)
);
if let Err(err) = sqlx::query(&sql).execute(self.pg_pool()?).await {
let error_text = err.to_string();
if !is_unpopulated_materialized_view_refresh_error(&error_text) {
return Err(tx_object_internal_status(
"refresh_materialized_view",
format!("refresh view failed: {err}"),
));
}
let sql = format!(
"REFRESH MATERIALIZED VIEW {}.{}",
qi_runtime(schema),
qi_runtime(name)
);
sqlx::query(&sql)
.execute(self.pg_pool()?)
.await
.map_err(|err| {
tx_object_internal_status(
"initial_refresh_materialized_view",
format!("initial refresh view failed: {err}"),
)
})?;
}
Ok(())
}
pub(crate) fn pg_pool(&self) -> Result<&PgPool, tonic::Status> {
self.pg_pool.as_ref().ok_or_else(|| {
crate::runtime::executor_utils::capability_status(
"postgres",
"pool",
"postgres_backend",
"PostgreSQL backend is not configured",
)
})
}
pub(crate) async fn pg_select_pool_for_table_routed(
&self,
table: &crate::generation::ManifestTable,
context: &RequestContext,
) -> Result<crate::runtime::core::RoutedReadPool, tonic::Status> {
let primary_required = table.replica_hint.eq_ignore_ascii_case("primary")
|| context.requires_primary_read()
|| super::accessors::read_fence_requires_primary(context);
self.pg_read_pool_routed(context, primary_required).await
}
#[cfg(feature = "qdrant")]
pub(crate) fn qdrant(&self) -> Result<&QdrantHttpClient, tonic::Status> {
self.qdrant.as_ref().ok_or_else(|| {
crate::runtime::executor_utils::capability_status(
"qdrant",
"client",
"qdrant_backend",
"Qdrant backend is not configured",
)
})
}
#[cfg(feature = "s3")]
pub(crate) fn s3(&self) -> Result<&aws_sdk_s3::Client, tonic::Status> {
self.s3.as_ref().ok_or_else(|| {
crate::runtime::executor_utils::capability_status(
"s3",
"client",
"s3_backend",
"S3/MinIO backend is not configured",
)
})
}
#[cfg(feature = "redis")]
pub(crate) async fn cache_get_fresh(
&self,
key: &str,
manifest_checksum: &str,
context: &RequestContext,
) -> Option<Vec<Vec<u8>>> {
let client = self.redis.as_ref()?;
let mut conn = client.get_multiplexed_async_connection().await.ok()?;
let data = match conn.get::<_, Option<Vec<u8>>>(key).await {
Ok(Some(data)) => data,
Ok(None) => {
self.cache_metrics.miss();
return None;
}
Err(_) => return None,
};
let min_lsn = serde_json::from_str::<crate::runtime::consistency::ReadFence>(
&context.read_fence_json,
)
.ok()
.and_then(|fence| crate::runtime::consistency_fence::parse_pg_lsn(&fence.min_outbox_lsn));
if let Ok(stamped) = serde_json::from_slice::<
crate::runtime::consistency_fence::StampedCacheValue<Vec<Vec<u8>>>,
>(&data)
{
if stamped.is_fresh_against(manifest_checksum, min_lsn) {
self.cache_metrics.hit();
return Some(stamped.value);
}
self.cache_metrics.miss();
return None;
}
if min_lsn.is_some() {
self.cache_metrics.miss();
return None;
}
match serde_json::from_slice(&data) {
Ok(value) => {
self.cache_metrics.hit();
Some(value)
}
Err(_) => {
self.cache_metrics.miss();
None
}
}
}
#[cfg(not(feature = "redis"))]
pub(crate) async fn cache_get_fresh(
&self,
_key: &str,
_manifest_checksum: &str,
_context: &RequestContext,
) -> Option<Vec<Vec<u8>>> {
None
}
#[cfg(feature = "redis")]
pub(crate) async fn cache_set(
&self,
key: &str,
records: &[Vec<u8>],
ttl: u64,
) -> Result<(), String> {
let Some(client) = &self.redis else {
return Ok(());
};
let mut conn = client
.get_multiplexed_async_connection()
.await
.map_err(|e| e.to_string())?;
let data = serde_json::to_vec(records).unwrap_or_default();
conn.set_ex::<_, _, ()>(key, data, ttl)
.await
.map_err(|e| e.to_string())
}
#[cfg(feature = "redis")]
pub(crate) async fn cache_set_stamped_from_pool(
&self,
key: &str,
records: &[Vec<u8>],
ttl: u64,
manifest_checksum: &str,
source_pool: &PgPool,
) -> Result<(), String> {
let source_lsn = sqlx::query_scalar::<_, String>(
"SELECT COALESCE(pg_last_wal_replay_lsn()::TEXT, pg_current_wal_lsn()::TEXT)",
)
.fetch_optional(source_pool)
.await
.ok()
.flatten()
.unwrap_or_default();
self.cache_set_stamped_with_lsn(key, records, ttl, manifest_checksum, source_lsn)
.await
}
#[cfg(not(feature = "redis"))]
pub(crate) async fn cache_set_stamped_from_pool(
&self,
_key: &str,
_records: &[Vec<u8>],
_ttl: u64,
_manifest_checksum: &str,
_source_pool: &PgPool,
) -> Result<(), String> {
Ok(())
}
#[cfg(feature = "redis")]
async fn cache_set_stamped_with_lsn(
&self,
key: &str,
records: &[Vec<u8>],
ttl: u64,
manifest_checksum: &str,
source_lsn: String,
) -> Result<(), String> {
let Some(client) = &self.redis else {
return Ok(());
};
let stamp = crate::runtime::consistency_fence::RedisCacheStamp {
manifest_checksum: manifest_checksum.to_string(),
source_lsn,
projection_version: crate::runtime::cdc::CdcConfig::current().producer_epoch,
written_at_unix_ms: unix_millis(),
};
let wrapped =
crate::runtime::consistency_fence::StampedCacheValue::new(stamp, records.to_vec());
let data = serde_json::to_vec(&wrapped).map_err(|e| e.to_string())?;
let mut conn = client
.get_multiplexed_async_connection()
.await
.map_err(|e| e.to_string())?;
conn.set_ex::<_, _, ()>(key, data, ttl)
.await
.map_err(|e| e.to_string())
}
#[cfg(not(feature = "redis"))]
pub(crate) async fn cache_set(
&self,
_key: &str,
_records: &[Vec<u8>],
_ttl: u64,
) -> Result<(), String> {
Ok(())
}
#[cfg(not(feature = "redis"))]
pub(crate) async fn cache_delete_pattern(&self, _pattern: &str) -> Result<(), String> {
Ok(())
}
#[cfg(feature = "redis")]
pub(crate) async fn cache_delete_pattern(&self, pattern: &str) -> Result<(), String> {
let Some(client) = &self.redis else {
return Ok(());
};
let mut conn = client
.get_multiplexed_async_connection()
.await
.map_err(|e| e.to_string())?;
let mut cursor = 0_u64;
loop {
let (next, keys): (u64, Vec<String>) = redis::cmd("SCAN")
.arg(cursor)
.arg("MATCH")
.arg(pattern)
.arg("COUNT")
.arg(500_u32)
.query_async(&mut conn)
.await
.map_err(|e| e.to_string())?;
if !keys.is_empty() {
let deleted = redis::cmd("DEL")
.arg(&keys)
.query_async::<u64>(&mut conn)
.await
.unwrap_or_default();
self.cache_metrics.invalidated(deleted);
}
if next == 0 {
break;
}
cursor = next;
}
Ok(())
}
pub(crate) fn encrypt_record_for_table(
&self,
table: &ManifestTable,
record: &JsonValue,
) -> Result<JsonValue, tonic::Status> {
let Some(encryption) = &self.encryption else {
if table.columns.iter().any(is_encrypted_column) {
self.encryption_metrics.record("encrypt", false);
return Err(record_encryption_key_missing_status(
&table.schema,
&table.table,
));
}
return Ok(record.clone());
};
let Some(object) = record.as_object() else {
return Ok(record.clone());
};
let mut encrypted = object.clone();
for column in table
.columns
.iter()
.filter(|column| is_encrypted_column(column))
{
let Some(value) = encrypted.get(&column.column_name).cloned() else {
continue;
};
if value.is_null() || json_is_ciphertext(&value) {
continue;
}
match encryption.encrypt_json_value(&value) {
Ok(ciphertext) => {
self.encryption_metrics.record("encrypt", true);
encrypted.insert(column.column_name.clone(), JsonValue::String(ciphertext));
}
Err(err) => {
self.encryption_metrics.record("encrypt", false);
return Err(tx_object_internal_status(
"encrypt_record_column",
format!(
"failed to encrypt {}.{}: {err}",
table.table, column.column_name
),
));
}
}
}
Ok(JsonValue::Object(encrypted))
}
pub(crate) async fn put_tx_object(
&self,
manifest: &CatalogManifest,
mutation: &Mutation,
metadata_context: RequestContext,
) -> Result<(u64, Option<String>, String), tonic::Status> {
#[cfg(not(feature = "s3"))]
{
let _ = (manifest, mutation, metadata_context);
return Err(s3_feature_disabled_status());
}
#[cfg(feature = "s3")]
{
if mutation.object_data.len() > INLINE_OBJECT_LIMIT_BYTES {
return Err(crate::runtime::executor_utils::quota_refusal_status(
"object",
"transaction inline object size",
"Use GeneratePresignedUrl for files > 1MB",
));
}
let context = merge_context(mutation.context.as_ref(), metadata_context);
let plan = build_object_stream_plan(
manifest,
&ObjectStreamPlanRequest {
context: context.clone(),
bucket: mutation.bucket.clone(),
object_key: mutation.object_key.clone(),
method: "PUT".to_string(),
chunk_count: 1,
final_chunk_seen: true,
content_type: mutation.content_type.clone(),
},
);
reject_plan(&plan.errors)?;
let target_instance = if context.target_instance.trim().is_empty() {
self.choose_instance_name_for_project("minio", true, &context.project_id)
.or_else(|| {
self.choose_instance_name_for_project("s3", true, &context.project_id)
})
} else {
Some(context.target_instance.as_str())
};
let s3 = self.s3_for_instance_for_project(target_instance, &context.project_id)?;
s3.put_object()
.bucket(&mutation.bucket)
.key(&mutation.object_key)
.set_content_type(if mutation.content_type.is_empty() {
None
} else {
Some(mutation.content_type.clone())
})
.body(ByteStream::from(mutation.object_data.clone()))
.send()
.await
.map_err(|err| {
crate::runtime::executor_utils::backend_transport_status(
"S3",
"put_object",
err,
)
})?;
Ok((
1,
target_instance.map(str::to_string),
context.project_id.clone(),
))
}
}
#[cfg(not(feature = "qdrant"))]
pub(crate) async fn compensate_qdrant_upserts(
&self,
compensations: &[(String, Vec<String>)],
) -> Option<String> {
if compensations.is_empty() {
None
} else {
Some(
"rolled back PostgreSQL; Qdrant compensation skipped because qdrant feature is not enabled"
.to_string(),
)
}
}
#[cfg(feature = "qdrant")]
pub(crate) async fn compensate_qdrant_upserts(
&self,
compensations: &[(String, Vec<String>)],
) -> Option<String> {
if compensations.is_empty() {
return None;
}
let Some(qdrant) = &self.qdrant else {
return Some(
"rolled back PostgreSQL; Qdrant compensation skipped because Qdrant is not configured"
.to_string(),
);
};
let mut deleted = 0usize;
let mut failures = Vec::new();
for (collection, point_ids) in compensations.iter().rev() {
let mut last_err = String::new();
let mut succeeded = false;
for attempt in 1..=COMPENSATION_RETRY_MAX_ATTEMPTS {
match qdrant.delete_points(collection, point_ids).await {
Ok(()) => {
deleted += point_ids.len();
succeeded = true;
break;
}
Err(err) => {
last_err = format!("{err}");
if attempt < COMPENSATION_RETRY_MAX_ATTEMPTS {
tracing::warn!(
collection = %collection,
attempt,
error = %err,
"Qdrant compensation failed; retrying"
);
tokio::time::sleep(compensation_retry_delay(attempt)).await;
}
}
}
}
if !succeeded {
failures.push(format!(
"{collection}: {last_err} (after {COMPENSATION_RETRY_MAX_ATTEMPTS} attempts)"
));
}
}
if failures.is_empty() {
Some(format!(
"rolled back PostgreSQL and compensated {deleted} Qdrant point(s)"
))
} else {
Some(format!(
"rolled back PostgreSQL; compensated {deleted} Qdrant point(s); compensation failures: {}",
failures.join("; ")
))
}
}
#[cfg(not(feature = "s3"))]
pub(crate) async fn compensate_s3_puts(
&self,
_compensations: &[(String, String, Option<String>, String)],
) -> Option<String> {
None
}
#[cfg(feature = "s3")]
pub(crate) async fn compensate_s3_puts(
&self,
compensations: &[(String, String, Option<String>, String)],
) -> Option<String> {
if compensations.is_empty() {
return None;
}
if self.s3.is_none() && self.s3_instances.is_empty() {
return Some("S3 compensation skipped because S3/MinIO is not configured".to_string());
}
let mut deleted = 0usize;
let mut failures = Vec::new();
for (bucket, object_key, instance, project_id) in compensations.iter().rev() {
let s3 = match self.s3_for_instance_for_project(instance.as_deref(), project_id) {
Ok(client) => client,
Err(err) => {
failures.push(format!("s3://{bucket}/{object_key}: {err}"));
continue;
}
};
let mut last_err = String::new();
let mut succeeded = false;
for attempt in 1..=COMPENSATION_RETRY_MAX_ATTEMPTS {
match s3
.delete_object()
.bucket(bucket)
.key(object_key)
.send()
.await
{
Ok(_) => {
deleted += 1;
succeeded = true;
break;
}
Err(err) => {
last_err = format!("{err}");
if attempt < COMPENSATION_RETRY_MAX_ATTEMPTS {
tracing::warn!(
bucket = %bucket,
key = %object_key,
attempt,
error = %err,
"S3 compensation failed; retrying"
);
tokio::time::sleep(compensation_retry_delay(attempt)).await;
}
}
}
}
if !succeeded {
failures.push(format!(
"s3://{bucket}/{object_key}: {last_err} (after {COMPENSATION_RETRY_MAX_ATTEMPTS} attempts)"
));
}
}
if failures.is_empty() {
Some(format!("compensated {deleted} S3/MinIO object(s)"))
} else {
Some(format!(
"compensated {deleted} S3/MinIO object(s); compensation failures: {}",
failures.join("; ")
))
}
}
pub(crate) async fn compensate_external_side_effects(
&self,
qdrant_compensations: &[(String, Vec<String>)],
s3_compensations: &[(String, String, Option<String>, String)],
) -> Option<String> {
let mut parts = Vec::new();
if let Some(message) = self.compensate_qdrant_upserts(qdrant_compensations).await {
parts.push(message);
}
if let Some(message) = self.compensate_s3_puts(s3_compensations).await {
parts.push(message);
}
if parts.is_empty() {
None
} else {
Some(parts.join("; "))
}
}
#[cfg(feature = "redis")]
pub async fn probe_redis_ping(&self) -> BackendProbeResult {
let start = std::time::Instant::now();
let client = match &self.redis {
Some(c) => c,
None => {
return BackendProbeResult {
backend: "redis".into(),
ok: false,
latency_ms: 0,
error: Some("Redis is not configured".into()),
};
}
};
match client.get_multiplexed_async_connection().await {
Err(e) => BackendProbeResult {
backend: "redis".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("connection failed: {e}")),
},
Ok(mut conn) => match redis::cmd("PING").query_async::<String>(&mut conn).await {
Ok(pong) if pong.eq_ignore_ascii_case("PONG") => BackendProbeResult {
backend: "redis".into(),
ok: true,
latency_ms: start.elapsed().as_millis() as u64,
error: None,
},
Ok(resp) => BackendProbeResult {
backend: "redis".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("unexpected PING response: {resp}")),
},
Err(e) => BackendProbeResult {
backend: "redis".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("PING failed: {e}")),
},
},
}
}
#[cfg(not(feature = "qdrant"))]
pub async fn probe_qdrant_collections(&self) -> BackendProbeResult {
BackendProbeResult {
backend: "qdrant".into(),
ok: false,
latency_ms: 0,
error: Some("qdrant/vector feature is not enabled".into()),
}
}
#[cfg(feature = "qdrant")]
pub async fn probe_qdrant_collections(&self) -> BackendProbeResult {
let start = std::time::Instant::now();
let qdrant = match &self.qdrant {
Some(q) => q,
None => {
return BackendProbeResult {
backend: "qdrant".into(),
ok: false,
latency_ms: 0,
error: Some("Qdrant is not configured".into()),
};
}
};
let url = format!("{}/collections", qdrant.base_url);
let mut req = qdrant.http.get(&url);
if let Some(api_key) = &qdrant.api_key {
req = req.header("api-key", api_key);
}
match req.send().await {
Err(e) => BackendProbeResult {
backend: "qdrant".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("HTTP request failed: {e}")),
},
Ok(resp) => {
let ok = resp.status().is_success();
let error = if ok {
None
} else {
Some(format!("HTTP {}", resp.status()))
};
BackendProbeResult {
backend: "qdrant".into(),
ok,
latency_ms: start.elapsed().as_millis() as u64,
error,
}
}
}
}
#[cfg(feature = "s3")]
pub async fn probe_s3_access(&self) -> BackendProbeResult {
let start = std::time::Instant::now();
let s3 = match &self.s3 {
Some(s3) => s3,
None => {
return BackendProbeResult {
backend: "s3".into(),
ok: false,
latency_ms: 0,
error: Some("S3/MinIO is not configured".into()),
};
}
};
match s3.list_buckets().send().await {
Ok(_) => BackendProbeResult {
backend: "s3".into(),
ok: true,
latency_ms: start.elapsed().as_millis() as u64,
error: None,
},
Err(e) => BackendProbeResult {
backend: "s3".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("list_buckets failed: {e}")),
},
}
}
#[cfg(feature = "kafka")]
pub fn probe_kafka_metadata(&self) -> BackendProbeResult {
let start = std::time::Instant::now();
let brokers = match self
.config
.kafka_brokers
.clone()
.filter(|v| !v.trim().is_empty())
{
Some(value) => value,
None => {
return BackendProbeResult {
backend: "kafka".into(),
ok: false,
latency_ms: 0,
error: Some(
"UDB_KAFKA_BROKERS / KAFKA_BROKERS is not set; CDC tailer disabled".into(),
),
};
}
};
let consumer: BaseConsumer = match ClientConfig::new()
.set("bootstrap.servers", &brokers)
.set("group.id", "udb-doctor")
.set("enable.auto.commit", "false")
.set("session.timeout.ms", "6000")
.create()
{
Ok(consumer) => consumer,
Err(err) => {
return BackendProbeResult {
backend: "kafka".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("Kafka consumer init failed: {err}")),
};
}
};
match consumer.fetch_metadata(None, Duration::from_secs(3)) {
Ok(metadata) => BackendProbeResult {
backend: "kafka".into(),
ok: !metadata.brokers().is_empty(),
latency_ms: start.elapsed().as_millis() as u64,
error: if metadata.brokers().is_empty() {
Some("Kafka metadata returned no brokers".into())
} else {
None
},
},
Err(err) => BackendProbeResult {
backend: "kafka".into(),
ok: false,
latency_ms: start.elapsed().as_millis() as u64,
error: Some(format!("Kafka metadata fetch failed: {err}")),
},
}
}
}
fn is_unpopulated_materialized_view_refresh_error(error_text: &str) -> bool {
let text = error_text.to_ascii_lowercase();
text.contains("concurrently")
&& (text.contains("has not been populated")
|| text.contains("not been populated")
|| text.contains("cannot refresh materialized view"))
}
#[cfg(test)]
mod materialized_view_refresh_tests {
use super::*;
use crate::generation::ManifestMaterializedView;
use crate::proto::{ErrorDetail, ErrorKind};
use crate::runtime::executor_utils::ERROR_DETAIL_METADATA_KEY;
fn decode_detail(status: &tonic::Status) -> ErrorDetail {
let raw = status
.metadata()
.get_bin(ERROR_DETAIL_METADATA_KEY)
.expect("typed error detail trailer");
crate::runtime::executor_utils::decode_error_detail_from_raw(&raw)
}
fn assert_single_field_violation(status: &tonic::Status, field: &str, description: &str) {
assert_eq!(status.code(), tonic::Code::InvalidArgument);
let detail = decode_detail(status);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert!(!detail.retryable);
assert_eq!(detail.field_violations.len(), 1);
assert_eq!(detail.field_violations[0].field, field);
assert_eq!(detail.field_violations[0].description, description);
}
fn assert_typed_detail(
status: &tonic::Status,
kind: ErrorKind,
backend: &str,
operation: &str,
capability_required: &str,
message: &str,
) {
assert_eq!(status.code(), tonic::Code::FailedPrecondition);
assert_eq!(status.message(), message);
let detail = decode_detail(status);
assert_eq!(detail.kind, kind as i32);
assert_eq!(detail.backend, backend);
assert_eq!(detail.operation, operation);
assert_eq!(detail.capability_required, capability_required);
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
assert!(detail.field_violations.is_empty());
}
fn assert_policy_detail(
status: &tonic::Status,
operation: &str,
policy_decision_id: &str,
message: &str,
) {
assert_eq!(status.code(), tonic::Code::PermissionDenied);
assert_eq!(status.message(), message);
let detail = decode_detail(status);
assert_eq!(detail.kind, ErrorKind::Policy as i32);
assert_eq!(detail.operation, operation);
assert_eq!(detail.policy_decision_id, policy_decision_id);
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
assert!(detail.field_violations.is_empty());
}
fn assert_internal_detail(status: &tonic::Status, operation: &str, message: &str) {
assert_eq!(status.code(), tonic::Code::Internal);
assert_eq!(status.message(), message);
let detail = decode_detail(status);
assert_eq!(detail.kind, ErrorKind::Internal as i32);
assert_eq!(detail.backend, "tx_object");
assert_eq!(detail.operation, operation);
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
assert!(detail.field_violations.is_empty());
}
fn admin_context() -> RequestContext {
RequestContext {
scopes: vec!["udb:admin".to_string()],
..RequestContext::default()
}
}
#[test]
fn detects_unpopulated_materialized_view_concurrent_refresh_error() {
assert!(is_unpopulated_materialized_view_refresh_error(
"ERROR: CONCURRENTLY cannot be used when the materialized view has not been populated"
));
assert!(!is_unpopulated_materialized_view_refresh_error(
"ERROR: could not create unique index for materialized view"
));
}
#[test]
fn tx_object_local_validation_carries_field_violations() {
let err = validate_tx_identifier("1bad", "schema").expect_err("invalid identifier");
assert_single_field_violation(&err, "schema", "must be a valid SQL identifier");
let err = tx_object_invalid_field(
"operation",
"must be one of upsert, delete, vector_upsert, put_object, or enqueue_outbox_event",
"unsupported transaction operation noop",
);
assert_single_field_violation(
&err,
"operation",
"must be one of upsert, delete, vector_upsert, put_object, or enqueue_outbox_event",
);
}
#[test]
fn tx_object_failed_preconditions_carry_typed_detail() {
let replay = mysql_xa_plan_replay_status("unsupported SQL");
assert_typed_detail(
&replay,
ErrorKind::Schema,
"mysql",
"xa_plan_replay",
"mysql_xa_plan_replay",
"2PC refused before side effects: plan SQL cannot be replayed on the MySQL XA participant: unsupported SQL",
);
let unsupported =
xa_unsupported_participants_status(&["mysql".to_string(), "qdrant".to_string()]);
assert_typed_detail(
&unsupported,
ErrorKind::Capability,
"transaction",
"xa_commit",
"xa_participants",
"2PC refused before side effects: unsupported participants mysql, qdrant",
);
let encryption = record_encryption_key_missing_status("private", "accounts");
assert_typed_detail(
&encryption,
ErrorKind::Capability,
"encryption",
"record_encryption",
"udb_encryption_key",
"table private.accounts contains encrypted columns but UDB encryption key is not configured",
);
let s3_feature = s3_feature_disabled_status();
assert_typed_detail(
&s3_feature,
ErrorKind::Capability,
"s3",
"put_tx_object",
"s3_feature",
"s3/object-store feature is not enabled",
);
}
#[test]
fn tx_object_internal_status_carries_typed_detail() {
let err = tx_object_internal_status(
"enqueue_projection_tasks",
"projection task enqueue failed: queue unavailable",
);
assert_internal_detail(
&err,
"enqueue_projection_tasks",
"projection task enqueue failed: queue unavailable",
);
let err = tx_object_internal_status(
"xa_ledger_commit_record",
"XA committed but ledger write failed for xid xa-1: unavailable",
);
assert_internal_detail(
&err,
"xa_ledger_commit_record",
"XA committed but ledger write failed for xid xa-1: unavailable",
);
}
#[tokio::test]
async fn create_materialized_view_admin_scope_denial_carries_policy_detail() {
let runtime = DataBrokerRuntime::default();
let err = runtime
.create_materialized_view(
&CatalogManifest::default(),
ViewDefinition::default(),
RequestContext::default(),
)
.await
.expect_err("admin scope denial must fail before validation or pg pool");
assert_policy_detail(
&err,
"create_materialized_view",
"admin_scope_required",
"scope udb:admin is required",
);
}
#[tokio::test]
async fn create_materialized_view_validation_carries_field_violations() {
let runtime = DataBrokerRuntime::default();
let manifest = CatalogManifest::default();
let request = ViewDefinition {
schema: "public".to_string(),
name: "missing_view".to_string(),
..ViewDefinition::default()
};
let err = runtime
.create_materialized_view(&manifest, request, admin_context())
.await
.expect_err("undeclared materialized view must fail before pg pool");
assert_single_field_violation(&err, "view", "must be declared in the proto AST");
let manifest = CatalogManifest {
tables: vec![ManifestTable {
materialized_views: vec![ManifestMaterializedView {
schema: "public".to_string(),
name: "orders_mv".to_string(),
query: "SELECT 1".to_string(),
with_data: false,
}],
..ManifestTable::default()
}],
..CatalogManifest::default()
};
let request = ViewDefinition {
schema: "public".to_string(),
name: "orders_mv".to_string(),
query: "SELECT 2".to_string(),
..ViewDefinition::default()
};
let err = runtime
.create_materialized_view(&manifest, request, admin_context())
.await
.expect_err("mismatched materialized view query must fail before pg pool");
assert_single_field_violation(&err, "query", "must match the proto AST declaration");
}
}