use super::*;
use crate::protocol::BackendCapabilityDescriptor;
use crate::runtime::schema_registry::{LookupError, NegotiationOutcome, SchemaRegistry};
const HEALTH_REPORT_CACHE_TTL_SECS: u64 = 3;
fn catalog_version_incompatible_status(operation: &'static str, reason: String) -> Status {
crate::runtime::executor_utils::schema_status(
tonic::Code::FailedPrecondition,
"catalog",
operation,
"catalog_version_incompatible",
reason,
)
}
fn message_schema_not_found_status(message_type: &str, project_id: &str) -> Status {
crate::runtime::executor_utils::schema_status(
tonic::Code::NotFound,
"catalog",
"LookupMessageSchema",
"message_schema_not_found",
format!("message schema '{message_type}' not found for project '{project_id}'"),
)
}
fn project_scope_mismatch_status(operation: &'static str) -> Status {
crate::runtime::executor_utils::policy_status_with_code(
tonic::Code::PermissionDenied,
operation,
"project_scope_mismatch",
"requested project_id does not match authenticated project",
)
}
#[allow(clippy::type_complexity)]
fn health_report_cache() -> &'static std::sync::Mutex<
std::collections::HashMap<(String, bool), (std::time::Instant, HealthReportResponse)>,
> {
static CACHE: std::sync::OnceLock<
std::sync::Mutex<
std::collections::HashMap<(String, bool), (std::time::Instant, HealthReportResponse)>,
>,
> = std::sync::OnceLock::new();
CACHE.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
}
impl DataBrokerService {
pub(crate) async fn get_capabilities_inner(
&self,
request: Request<CapabilitiesRequest>,
) -> Result<Response<CapabilitiesResponse>, Status> {
let (started, security) = authorized_call!(self, request, "GetCapabilities");
if let Err(err) = require_admin_scope(&security) {
return self.record_grpc("GetCapabilities", started, Err(err));
}
let mut enabled_backends = self.runtime_snapshot().enabled_backend_names();
let enabled_set: std::collections::HashSet<String> =
enabled_backends.iter().cloned().collect();
let mut degraded_backends: Vec<String> = crate::backend::all_plugins()
.into_iter()
.map(|plugin| plugin.kind().as_str().to_string())
.filter(|backend| !enabled_set.contains(backend))
.collect();
enabled_backends.sort();
enabled_backends.dedup();
degraded_backends.sort();
degraded_backends.dedup();
let manifest_checksum = if !self.catalog.active().manifest.checksum_sha256.is_empty() {
self.catalog.active().manifest.checksum_sha256.clone()
} else {
use sha2::Digest;
let mut hasher = sha2::Sha256::new();
if let Ok(bytes) = serde_json::to_vec(&self.catalog.active().manifest) {
hasher.update(bytes);
}
format!("{:x}", hasher.finalize())
};
let req = request.into_inner();
let sys_cfg = crate::runtime::system::SystemCatalogConfig::default();
let sys_schema = &sys_cfg.cdc.system_schema;
let qi = |s: &str| format!("\"{s}\"");
let qrel = |schema: &str, table: &str| format!("{}.{}", qi(schema), qi(table));
let mut system_catalog_relations: Vec<String> = vec![
qrel(sys_schema, &sys_cfg.cdc.outbox_table),
qrel(sys_schema, &sys_cfg.cdc.offsets_table),
qrel(sys_schema, &sys_cfg.cdc.lock_log_table),
qrel(sys_schema, &sys_cfg.saga_table),
qrel(&sys_cfg.abac_schema, &sys_cfg.abac_table),
];
let project_scope = req.project_id.trim().to_string();
if !project_scope.is_empty()
&& let Ok(versions) = self
.runtime_snapshot()
.get_catalog_versions(&project_scope)
.await
{
for v in &versions {
let ver = v["version"].as_str().unwrap_or("unknown");
system_catalog_relations.push(format!("project:{project_scope}:catalog:{ver}"));
}
}
let supported_rpcs: Vec<String> = SUPPORTED_RPC_NAMES
.iter()
.map(|name| (*name).to_string())
.collect();
let configured_backends: std::collections::HashSet<String> = {
let snap = self.runtime_snapshot();
snap.backend_instances()
.iter()
.map(|inst| inst.backend.clone())
.collect()
};
let backend_caps = crate::backend::capability_matrix_configured(&configured_backends);
let backend_protocol_support: Vec<crate::proto::BackendProtocolSupport> = backend_caps
.iter()
.map(|entry| {
let kind = crate::backend::BackendKind::from_token(&entry.backend);
let v2 = kind.map(|kind| kind.capabilities_v2());
let supports_streaming_reads =
v2.as_ref().is_some_and(|cap| cap.supports_streaming);
let supports_object_streaming = v2.as_ref().is_some_and(|cap| cap.is_object_store);
crate::proto::BackendProtocolSupport {
backend: entry.backend.clone(),
supports_streaming_reads,
supports_object_streaming,
encodings: Vec::new(),
}
})
.collect();
let protocol_support = Some(crate::proto::ProtocolSupport {
min_protocol_version: UDB_PROTOCOL_VERSION.to_string(),
max_protocol_version: UDB_PROTOCOL_VERSION.to_string(),
encodings: vec!["record_set_v1".to_string(), "record_batch_v2".to_string()],
compression: Vec::new(),
supports_streaming_reads: true,
supports_object_streaming: true,
max_recv_message_bytes: 0,
max_send_message_bytes: 0,
supported_rpcs: supported_rpcs.clone(),
});
let native_services: Vec<crate::proto::NativeServiceStatus> =
crate::runtime::service::native_registry::resolved_native_service_statuses(
self.runtime_snapshot().config(),
)
.into_iter()
.map(crate::runtime::service::native_registry::status_to_proto)
.collect();
let startup_summary = format!(
"[UDB] capabilities: {} table(s), {} store(s), {} backend(s) enabled, {} degraded",
self.catalog.active().manifest.tables.len(),
self.catalog.active().manifest.stores.len(),
enabled_backends.len(),
degraded_backends.len()
);
tracing::info!("{startup_summary}");
self.record_grpc(
"GetCapabilities",
started,
Ok(Response::new(CapabilitiesResponse {
schema_checksum: manifest_checksum,
protocol_version: UDB_PROTOCOL_VERSION.to_string(),
enabled_backends,
degraded_backends,
system_catalog_relations,
supported_rpcs,
backend_instances: self
.runtime_snapshot()
.backend_instances()
.iter()
.map(backend_instance_status)
.collect(),
backend_capabilities: backend_caps
.into_iter()
.map(|entry| BackendCapabilityDescriptor {
backend: entry.backend,
tier: entry.tier,
operations: entry.operations,
unsupported_error_code: entry.unsupported_error_code,
consistency_model: entry.consistency_model,
max_payload_bytes: entry.max_payload_bytes as i64,
supports_xa: entry.supports_xa,
supports_two_phase_commit: entry.supports_two_phase_commit,
})
.collect(),
protocol_support,
backend_protocol_support,
native_services,
deployment_tier: crate::runtime::core::declared_deployment_tier()
.map(|tier| tier.as_str().to_string())
.unwrap_or_default(),
})),
)
}
pub(crate) async fn get_health_report_inner(
&self,
request: Request<HealthReportRequest>,
) -> Result<Response<HealthReportResponse>, Status> {
let (started, security) = authorized_call!(self, request, "GetHealthReport");
if let Err(err) = require_admin_scope(&security) {
return self.record_grpc("GetHealthReport", started, Err(err));
}
let request = request.into_inner();
let project_scope = request.project_id.trim().to_string();
let cache_key = (project_scope.clone(), request.with_probes);
let ttl = std::time::Duration::from_secs(HEALTH_REPORT_CACHE_TTL_SECS);
if let Ok(cache) = health_report_cache().lock() {
if let Some((stored_at, cached)) = cache.get(&cache_key) {
if stored_at.elapsed() < ttl {
return self.record_grpc(
"GetHealthReport",
started,
Ok(Response::new(cached.clone())),
);
}
}
}
let init = self.runtime_snapshot().init_report().clone();
let mut errors = Vec::new();
let mut warnings = init.warnings.clone();
let priv_report = if init.postgres_configured {
let pr = self.runtime_snapshot().check_postgres_privileges().await;
if !pr.create_publication {
warnings.push("PG role lacks CREATE PUBLICATION privilege".into());
}
if !pr.replication_slot {
warnings.push("PG role lacks replication role for CDC".into());
}
Some(pr)
} else {
errors.push("PostgreSQL is required: UDB_PG_DSN / DATABASE_URL is not set".into());
None
};
let native_statuses =
crate::runtime::service::native_registry::resolved_native_service_statuses(
self.runtime_snapshot().config(),
);
let mut probes = Vec::new();
if request.with_probes {
let snapshot = self.runtime_snapshot();
let backend_probe_futs: Vec<_> = snapshot
.configured_probe_backends(false)
.into_iter()
.map(|backend| snapshot.probe_backend(backend))
.collect();
probes.extend(futures::future::join_all(backend_probe_futs).await);
#[cfg(feature = "kafka")]
probes.push(self.runtime_snapshot().probe_kafka_metadata());
let self_check_futs: Vec<_> = native_statuses
.iter()
.filter(|s| {
s.mounted
&& crate::runtime::service::native_store_binding::native_service_store_backend(
&s.service_id,
)
.is_some()
})
.map(|status| {
let service_id = status.service_id.clone();
let snapshot = self.runtime_snapshot();
let project_scope = project_scope.clone();
async move {
let outcome = snapshot
.native_store_self_check(&service_id, &project_scope)
.await;
(service_id, outcome)
}
})
.collect();
for (service_id, outcome) in futures::future::join_all(self_check_futs).await {
match outcome {
Ok(backend) => warnings.push(format!(
"native store [{service_id}]: persistence verified on '{backend}'"
)),
Err(err) => warnings.push(format!(
"native store [{service_id}]: persistence self-check failed: {}",
err.message()
)),
}
}
}
if let Some(transport) = self.runtime_snapshot().mongodb_transport_kind() {
warnings.push(format!(
"mongodb: transport={transport}; native wire-protocol not supported"
));
}
if !project_scope.is_empty() {
match self
.runtime_snapshot()
.get_catalog_versions(&project_scope)
.await
{
Ok(versions) if versions.is_empty() => {
warnings.push(format!(
"project '{project_scope}': no catalog versions found"
));
}
Ok(versions) => {
let active: Vec<_> = versions
.iter()
.filter(|v| v["status"].as_str().unwrap_or("") == "ACTIVE")
.filter_map(|v| v["version"].as_str())
.collect();
warnings.push(format!(
"project '{project_scope}': {} catalog version(s), active=[{}]",
versions.len(),
active.join(", ")
));
}
Err(_) => {
warnings.push(format!(
"project '{project_scope}': catalog version query failed"
));
}
}
}
let privileges_json = priv_report
.as_ref()
.and_then(|p| serde_json::to_vec(p).ok())
.unwrap_or_default();
let probes_json = serde_json::to_vec(&probes).unwrap_or_default();
let auth_report = crate::runtime::service::auth_service::readiness::check_auth_readiness(
&crate::runtime::security::SecurityConfig::current(),
)
.await;
let auth_triples: Vec<(String, bool, String)> = auth_report
.checks
.iter()
.map(|c| (c.name.clone(), c.ok, c.detail.clone()))
.collect();
let readiness =
crate::runtime::slo::build_readiness_facts(&init, &native_statuses, &auth_triples);
for err in readiness.errors() {
if !errors.contains(&err) {
errors.push(err);
}
}
for warn in readiness.warnings() {
if !warnings.contains(&warn) {
warnings.push(warn);
}
}
let native_services: Vec<crate::proto::NativeServiceStatus> = native_statuses
.into_iter()
.map(crate::runtime::service::native_registry::status_to_proto)
.collect();
let response = HealthReportResponse {
passed: errors.is_empty(),
postgres_configured: init.postgres_configured,
redis_configured: init.redis_configured,
qdrant_configured: init.qdrant_configured,
s3_configured: init.s3_configured,
errors,
warnings,
privileges_json,
probes_json,
backend_instances: self
.runtime_snapshot()
.backend_instances()
.iter()
.map(backend_instance_status)
.collect(),
native_services,
};
if let Ok(mut cache) = health_report_cache().lock() {
if cache.len() > 8 {
cache.clear();
}
cache.insert(cache_key, (std::time::Instant::now(), response.clone()));
}
self.record_grpc("GetHealthReport", started, Ok(Response::new(response)))
}
pub(crate) async fn lookup_message_schema_inner(
&self,
request: Request<MessageSchemaLookupRequest>,
) -> Result<Response<MessageSchemaLookupResponse>, Status> {
let (started, security) = authorized_call!(self, request, "LookupMessageSchema");
let req = request.into_inner();
if let (Some(requested), Some(bound)) =
(non_empty(&req.project_id), non_empty(&security.project_id))
{
if requested != bound && !security.has_scope("udb:admin") {
return self.record_grpc(
"LookupMessageSchema",
started,
Err(project_scope_mismatch_status("LookupMessageSchema")),
);
}
}
let project_id = non_empty(&req.project_id)
.or_else(|| non_empty(&security.project_id))
.unwrap_or("default")
.to_string();
let client_version = non_empty(&req.client_catalog_version)
.unwrap_or(&security.client_catalog_version)
.to_string();
let registry = SchemaRegistry::new(self.catalog.clone());
let descriptor =
match registry.lookup_message(&project_id, &req.message_type, &client_version) {
Ok(descriptor) => descriptor,
Err(LookupError::MessageNotFound { message_type, .. }) => {
return self.record_grpc(
"LookupMessageSchema",
started,
Err(message_schema_not_found_status(&message_type, &project_id)),
);
}
Err(LookupError::Incompatible { reason, .. }) => {
return self.record_grpc(
"LookupMessageSchema",
started,
Err(catalog_version_incompatible_status(
"LookupMessageSchema",
reason,
)),
);
}
};
self.record_grpc(
"LookupMessageSchema",
started,
Ok(Response::new(MessageSchemaLookupResponse {
schema: Some(message_descriptor_to_proto(descriptor)),
})),
)
}
pub(crate) async fn list_message_schemas_inner(
&self,
request: Request<MessageSchemaListRequest>,
) -> Result<Response<MessageSchemaListResponse>, Status> {
let (started, security) = authorized_call!(self, request, "ListMessageSchemas");
let req = request.into_inner();
if let (Some(requested), Some(bound)) =
(non_empty(&req.project_id), non_empty(&security.project_id))
{
if requested != bound && !security.has_scope("udb:admin") {
return self.record_grpc(
"ListMessageSchemas",
started,
Err(project_scope_mismatch_status("ListMessageSchemas")),
);
}
}
let project_id = non_empty(&req.project_id)
.or_else(|| non_empty(&security.project_id))
.unwrap_or("default")
.to_string();
let client_version = non_empty(&req.client_catalog_version)
.unwrap_or(&security.client_catalog_version)
.to_string();
let registry = SchemaRegistry::new(self.catalog.clone());
let outcome = registry.negotiate_version(&project_id, &client_version);
if let NegotiationOutcome::Incompatible { reason, .. } = outcome {
return self.record_grpc(
"ListMessageSchemas",
started,
Err(catalog_version_incompatible_status(
"ListMessageSchemas",
reason,
)),
);
}
let active = self.catalog.active_for(&project_id);
self.record_grpc(
"ListMessageSchemas",
started,
Ok(Response::new(MessageSchemaListResponse {
project_id,
catalog_version: active.metadata.version.clone(),
manifest_checksum: active.metadata.checksum.clone(),
message_types: registry.list_messages(&active.metadata.project_id),
})),
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HealthPlane {
DataBroker,
NativeControlPlane,
WebRtcPeer,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::proto::{ErrorDetail, ErrorKind};
use crate::runtime::executor_utils::ERROR_DETAIL_METADATA_KEY;
fn decode_detail(status: &Status) -> ErrorDetail {
let raw = status
.metadata()
.get_bin(ERROR_DETAIL_METADATA_KEY)
.expect("typed detail trailer is present");
crate::runtime::executor_utils::decode_error_detail_from_raw(&raw)
}
#[test]
fn catalog_version_incompatible_carries_schema_detail() {
let err = catalog_version_incompatible_status(
"LookupMessageSchema",
"client catalog version 1 is incompatible with active version 2".to_string(),
);
assert_eq!(err.code(), tonic::Code::FailedPrecondition);
assert_eq!(
err.message(),
"client catalog version 1 is incompatible with active version 2"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Schema as i32);
assert_eq!(detail.backend, "catalog");
assert_eq!(detail.operation, "LookupMessageSchema");
assert_eq!(detail.capability_required, "catalog_version_incompatible");
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
}
#[test]
fn message_schema_not_found_carries_schema_detail() {
let err = message_schema_not_found_status("app.Invoice", "billing");
assert_eq!(err.code(), tonic::Code::NotFound);
assert_eq!(
err.message(),
"message schema 'app.Invoice' not found for project 'billing'"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Schema as i32);
assert_eq!(detail.backend, "catalog");
assert_eq!(detail.operation, "LookupMessageSchema");
assert_eq!(detail.capability_required, "message_schema_not_found");
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
}
#[test]
fn project_scope_mismatch_carries_policy_detail() {
for operation in ["LookupMessageSchema", "ListMessageSchemas"] {
let err = project_scope_mismatch_status(operation);
assert_eq!(err.code(), tonic::Code::PermissionDenied);
assert_eq!(
err.message(),
"requested project_id does not match authenticated project"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Policy as i32);
assert_eq!(detail.backend, "");
assert_eq!(detail.operation, operation);
assert_eq!(detail.policy_decision_id, "project_scope_mismatch");
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
}
}
}
pub async fn build_listener_health_service(
plane: HealthPlane,
config: &crate::runtime::config::UdbConfig,
runtime: Option<&crate::runtime::DataBrokerRuntime>,
) -> tonic_health::pb::health_server::HealthServer<impl tonic_health::pb::health_server::Health> {
use crate::runtime::service::native_registry::NativeListenerKind;
use tonic_health::ServingStatus;
let (mut reporter, health_service) = tonic_health::server::health_reporter();
let mut any_serving = false;
let statuses =
crate::runtime::service::native_registry::resolved_native_service_statuses(config);
let readiness_passed = if let Some(runtime) = runtime {
let auth_triples = crate::runtime::service::auth_readiness_triples(
&crate::runtime::security::SecurityConfig::current(),
)
.await;
crate::runtime::slo::build_readiness_facts(runtime.init_report(), &statuses, &auth_triples)
.passed()
} else {
true
};
match plane {
HealthPlane::DataBroker => {
if readiness_passed {
reporter
.set_serving::<DataBrokerServer<DataBrokerService>>()
.await;
any_serving = true;
} else {
reporter
.set_not_serving::<DataBrokerServer<DataBrokerService>>()
.await;
}
}
HealthPlane::NativeControlPlane | HealthPlane::WebRtcPeer => {
let want_kind = match plane {
HealthPlane::NativeControlPlane => NativeListenerKind::ControlPlane.as_str(),
HealthPlane::WebRtcPeer => NativeListenerKind::WebRtcPeer.as_str(),
HealthPlane::DataBroker => unreachable!(),
};
for status in statuses {
if !status.enabled || status.listener_kind != want_kind {
continue;
}
let serving = if status.mounted && readiness_passed {
any_serving = true;
ServingStatus::Serving
} else {
ServingStatus::NotServing
};
for proto_service in &status.proto_services {
reporter
.set_service_status(proto_service.as_str(), serving)
.await;
}
}
}
}
reporter
.set_service_status(
"",
if any_serving {
ServingStatus::Serving
} else {
ServingStatus::NotServing
},
)
.await;
health_service
}
fn message_descriptor_to_proto(
descriptor: crate::runtime::schema_registry::MessageDescriptor,
) -> MessageSchemaDescriptor {
MessageSchemaDescriptor {
message_type: descriptor.message_type,
project_id: descriptor.project_id,
catalog_version: descriptor.catalog_version,
manifest_checksum: descriptor.manifest_checksum,
schema: descriptor.schema,
table: descriptor.table,
primary_key: descriptor.primary_key,
fields: descriptor
.fields
.into_iter()
.map(|field| MessageFieldDescriptor {
name: field.name,
column_name: field.column_name,
proto_type: field.proto_type,
sql_type: field.sql_type,
not_null: field.not_null,
is_primary: field.is_primary,
is_array: field.is_array,
})
.collect(),
}
}
#[cfg(test)]
mod health_cache_tests {
use super::*;
#[test]
fn health_cache_ttl_is_short_enough_to_not_mask_failures() {
assert!(
HEALTH_REPORT_CACHE_TTL_SECS <= 5,
"health cache TTL must stay tiny (capability-lie guard): {HEALTH_REPORT_CACHE_TTL_SECS}s"
);
}
#[test]
fn health_cache_key_separates_project_scope_and_probe_mode() {
let now = std::time::Instant::now();
let yes = HealthReportResponse {
passed: true,
..Default::default()
};
let no = HealthReportResponse {
passed: false,
..Default::default()
};
let k_probe = ("b2-cache-proj".to_string(), true);
let k_noprobe = ("b2-cache-proj".to_string(), false);
{
let mut cache = health_report_cache().lock().unwrap();
cache.insert(k_probe.clone(), (now, yes));
cache.insert(k_noprobe.clone(), (now, no));
}
{
let cache = health_report_cache().lock().unwrap();
assert!(
cache.get(&k_probe).unwrap().1.passed,
"with_probes=true variant"
);
assert!(
!cache.get(&k_noprobe).unwrap().1.passed,
"with_probes=false must be a distinct cache entry"
);
assert!(
cache.get(&("b2-other-proj".to_string(), true)).is_none(),
"a different project scope must not collide"
);
}
let mut cache = health_report_cache().lock().unwrap();
cache.remove(&k_probe);
cache.remove(&k_noprobe);
}
}