#![allow(clippy::result_large_err)]
use std::fs;
use std::net::SocketAddr;
use std::pin::Pin;
use std::sync::{Arc, RwLock};
use std::time::{Duration, Instant};
use arc_swap::ArcSwap;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio_stream::{Stream, StreamExt as _};
use tonic::transport::{Certificate, Identity, ServerTlsConfig};
use tonic::{Request, Response, Status};
use uuid::Uuid;
use crate::ast::ProtoSchema;
use crate::cdc::{CdcEngine, CdcRedactionMode};
use crate::engine::FsmState;
use crate::generation::CatalogManifest;
use crate::lifecycle::run_startup_lifecycle;
use crate::metrics::{MetricsRecorder, NoopMetrics, PrometheusMetrics};
use crate::proto::data_broker_server::{DataBroker, DataBrokerServer};
use crate::proto::{
AdminAuditLogRecord, AdminAuditLogRequest, AdminAuditLogResponse, AdminAuditVerifyRequest,
AdminAuditVerifyResponse, AdminBackendSummary, AdminCatalogSummary, AdminCdcSummary,
AdminSagaSummary, AdminSummaryRequest, AdminSummaryResponse, BackendInstanceStatus,
CapabilitiesRequest, CapabilitiesResponse, CatalogManifestRequest, CatalogManifestResponse,
CatalogValidationResponse, CatalogVersionListResponse, CatalogVersionRequest,
CatalogVersionResponse, CdcControlRequest, CdcEnvelope, CdcRedactionPreviewRequest,
CdcRedactionPreviewResponse, CdcStatusResponse, CdcSubscriptionRequest, Chunk, DeleteRequest,
DlqActionRequest, DlqEventRecord, DlqEventRequest, DlqEventResponse, DlqListRequest,
DlqListResponse, EnqueueOutboxEventRequest, EnqueueOutboxEventResponse, EnsureBaselineRequest,
EnsureBaselineResponse, EnsureProjectRequest, GenericDispatchRequest, GenericDispatchResponse,
HealthReportRequest, HealthReportResponse, MessageFieldDescriptor, MessageSchemaDescriptor,
MessageSchemaListRequest, MessageSchemaListResponse, MessageSchemaLookupRequest,
MessageSchemaLookupResponse, MigrationApplyRequest, MigrationPlanRequest,
MigrationPlanResponse, MigrationRunListRequest, MigrationRunListResponse, MigrationRunRequest,
MigrationStatusResponse, MultipartUploadRequest, MultipartUploadResponse, Mutation,
MutationResponse, PolicyLintResponse, PolicyListRequest, PolicyListResponse, PolicyRecord,
PolicyRequest, ProjectListRequest, ProjectListResponse, ProjectRecord,
ProjectionDriftDivergentRow, ProjectionDriftScanRequest, ProjectionDriftScanResponse,
ProjectionDriftTargetReport, PutPolicyRequest, RecordSet, ResourceAdminRequest,
ResourceListResponse, SagaListRequest, SagaListResponse, SagaRecord, SagaRequest, SagaResponse,
SelectRequest, StageCatalogRequest, TxStatus, UpsertRequest, UrlRequest, UrlResponse,
VectorHybridSearchRequest, VectorSearchRequest, VectorSet, VectorUpsertRequest, ViewDefinition,
};
use crate::runtime::DataBrokerRuntime;
use crate::runtime::authz::{AuthzQuery, AuthzSnapshot, Principal, ResourceRef};
use crate::security::{
SecurityConfig, SecurityContext, enforce_select_export_controls, ip_matches_allow_entry,
security_from_request, validate_bearer_token,
};
mod analytics_service;
mod asset_service;
pub(crate) mod auth_service;
mod backup_service;
mod cache_service;
mod config_service;
pub use auth_service::auth_readiness_triples;
pub(crate) use auth_service::events::topics::AUTH_TOPIC_PATTERNS;
pub use auth_service::{
BootstrapAdmin, bootstrap_admin_user, cli_api_key_list, cli_api_key_revoke,
migrate_service_account_grants, served_bootstrap_admin,
};
mod livequery_service;
mod lock_service;
mod metering_service;
mod embedding_service;
mod method_security;
#[cfg(feature = "bench-internals")]
pub(crate) use method_security::{
build_registry as build_method_security_registry, method_security, method_security_registry,
};
pub(crate) mod native_entity_store;
#[cfg(test)]
mod native_entity_store_tests;
mod native_helpers;
pub mod native_registry;
pub(crate) mod native_runtime;
pub(crate) mod native_store_binding;
mod notification_service;
mod scheduler_service;
mod search_service;
mod storage_service;
mod tenant_service;
mod webhook_service;
mod vault_service;
mod webrtc_service;
mod workflow_service;
const UDB_FILE_DESCRIPTOR_SET: &[u8] = tonic::include_file_descriptor_set!("udb_descriptor");
fn initial_authz_snapshot(default_allow: bool) -> AuthzSnapshot {
AuthzSnapshot {
default_allow,
..AuthzSnapshot::default()
}
}
fn startup_bool_env(key: &str) -> bool {
std::env::var(key)
.ok()
.map(|value| {
matches!(
value.trim().to_ascii_lowercase().as_str(),
"1" | "true" | "yes" | "on"
)
})
.unwrap_or(false)
}
macro_rules! authorized_call {
($self:expr, $request:expr, $method:literal) => {{
let started = Instant::now();
let security = match security_from_request(&$request) {
Ok(s) => s,
Err(e) => return $self.record_grpc($method, started, Err(e)),
};
if let Err(err) = $self.authorize(&security, "*", $method).await {
return $self.record_grpc($method, started, Err(err));
}
(started, security)
}};
}
#[derive(Debug, Clone)]
pub struct DataBrokerService {
pub catalog: Arc<crate::runtime::catalog::CatalogManager>,
pub manifest: CatalogManifest,
pub runtime: Arc<ArcSwap<DataBrokerRuntime>>,
lifecycle_state: Arc<RwLock<FsmState>>,
authz_snapshot: Arc<ArcSwap<AuthzSnapshot>>,
metrics: Arc<dyn MetricsRecorder>,
cdc_engine: Option<Arc<CdcEngine>>,
projection_engine: Option<Arc<crate::runtime::projection::ProjectionEngine>>,
#[cfg(feature = "redis")]
rate_limit_redis: Arc<tokio::sync::Mutex<Option<redis::aio::MultiplexedConnection>>>,
}
pub(crate) const UDB_PROTOCOL_VERSION: &str = "1.0.0";
pub(crate) fn admin_seed_enabled() -> bool {
std::env::var("UDB_ENABLE_ADMIN_SEED")
.map(|value| {
let value = value.trim();
value.eq_ignore_ascii_case("1") || value.eq_ignore_ascii_case("true")
})
.unwrap_or(false)
}
pub(crate) const SUPPORTED_RPC_NAMES: &[&str] = &[
"Select",
"BatchSelect",
"SelectV2",
"Upsert",
"BatchUpsert",
"Delete",
"VectorSearch",
"VectorHybridSearch",
"VectorUpsert",
"VectorBatchUpsert",
"PutObject",
"GetObject",
"GeneratePresignedUrl",
"InitiateMultipartUpload",
"BeginTx",
"PublishCDC",
"EnqueueOutboxEvent",
"StageCatalog",
"ActivateCatalog",
"RollbackCatalog",
"ValidateCatalog",
"GetCatalogVersions",
"GetCatalogVersion",
"PlanMigration",
"ApplyMigration",
"GetMigrationStatus",
"ListMigrationRuns",
"ApproveMigrationPlan",
"ListDlqEvents",
"GetDlqEvent",
"ReplayDlqEvent",
"DismissDlqEvent",
"QuarantineDlqEvent",
"GetCdcStatus",
"PauseCdc",
"ResumeCdc",
"StepDownCdcLeader",
"PreviewCdcRedaction",
"ScanProjectionDrift",
"ListSagas",
"GetSaga",
"RetrySagaCompensation",
"MarkSagaReviewed",
"ListPolicies",
"PutPolicy",
"DeletePolicy",
"ReloadPolicies",
"LintPolicies",
"GenericDispatch",
"EnsureResource",
"DropResource",
"ListResources",
"CreateMaterializedView",
"CacheGet",
"CacheSet",
"CacheDelete",
"CacheScan",
"DocumentGet",
"DocumentFind",
"DocumentUpsert",
"DocumentDelete",
"GraphQuery",
"GraphMutate",
"TimeSeriesWrite",
"TimeSeriesQuery",
"AnalyticalQuery",
"GetCapabilities",
"GetCatalogManifest",
"LookupMessageSchema",
"ListMessageSchemas",
"GetHealthReport",
"EnsureProject",
"EnsureBaseline",
"ListProjects",
"GetAdminSummary",
"ListAdminAuditLogs",
"VerifyAdminAuditLog",
];
impl DataBrokerService {
pub fn new(manifest: CatalogManifest) -> Self {
let catalog = Arc::new(crate::runtime::catalog::CatalogManager::new(
manifest.clone(),
));
Self {
catalog,
manifest,
runtime: Arc::new(ArcSwap::from_pointee(DataBrokerRuntime::planning_only())),
lifecycle_state: Arc::new(RwLock::new(FsmState::Idle)),
authz_snapshot: Arc::new(ArcSwap::from_pointee(initial_authz_snapshot(false))),
metrics: service_metrics_recorder(),
cdc_engine: None,
projection_engine: None,
#[cfg(feature = "redis")]
rate_limit_redis: Arc::new(tokio::sync::Mutex::new(None)),
}
}
pub fn with_runtime(manifest: CatalogManifest, runtime: DataBrokerRuntime) -> Self {
let catalog = Arc::new(crate::runtime::catalog::CatalogManager::new(
manifest.clone(),
));
Self {
catalog,
manifest,
runtime: Arc::new(ArcSwap::from_pointee(runtime)),
lifecycle_state: Arc::new(RwLock::new(FsmState::Completed)),
authz_snapshot: Arc::new(ArcSwap::from_pointee(initial_authz_snapshot(false))),
metrics: service_metrics_recorder(),
cdc_engine: None,
projection_engine: None,
#[cfg(feature = "redis")]
rate_limit_redis: Arc::new(tokio::sync::Mutex::new(None)),
}
}
pub fn with_runtime_and_state(
manifest: CatalogManifest,
runtime: DataBrokerRuntime,
lifecycle_state: Arc<RwLock<FsmState>>,
metrics: Arc<dyn MetricsRecorder>,
cdc_engine: Option<Arc<CdcEngine>>,
abac_default_allow: bool,
) -> Self {
let catalog = Arc::new(crate::runtime::catalog::CatalogManager::new(
manifest.clone(),
));
Self {
catalog,
manifest,
runtime: Arc::new(ArcSwap::from_pointee(runtime)),
lifecycle_state,
authz_snapshot: Arc::new(ArcSwap::from_pointee(initial_authz_snapshot(
abac_default_allow,
))),
metrics,
cdc_engine,
projection_engine: None,
#[cfg(feature = "redis")]
rate_limit_redis: Arc::new(tokio::sync::Mutex::new(None)),
}
}
pub(crate) fn authz_snapshot(&self) -> Arc<ArcSwap<AuthzSnapshot>> {
self.authz_snapshot.clone()
}
fn current_authz_snapshot(&self) -> Arc<AuthzSnapshot> {
self.authz_snapshot.load_full()
}
pub fn runtime_snapshot(&self) -> Arc<DataBrokerRuntime> {
self.runtime.load_full()
}
pub async fn reload_runtime_from_config(
&self,
config: crate::runtime::config::UdbConfig,
options: crate::runtime::ConfigReloadOptions,
) -> crate::runtime::ConfigReloadReport {
let mut next = self.runtime.load_full().as_ref().clone();
let report = next.reload_from_config(config, options).await;
if report.applied {
self.runtime.store(Arc::new(next));
}
report
}
pub async fn reload_runtime_from_env(
&self,
reason: impl Into<String>,
) -> crate::runtime::ConfigReloadReport {
self.reload_runtime_from_config(
crate::runtime::config::UdbConfig::from_merged_env(),
crate::runtime::ConfigReloadOptions {
reason: reason.into(),
require_connected_backends: true,
rollback_on_failed_health: true,
..crate::runtime::ConfigReloadOptions::default()
},
)
.await
}
pub(crate) fn ensure_ready(&self) -> Result<(), Status> {
let state = self
.lifecycle_state
.read()
.map(|state| state.clone())
.unwrap_or(FsmState::Error);
if state == FsmState::Completed {
Ok(())
} else {
Err(crate::runtime::executor_utils::retryable_status(
"data_broker",
"startup_not_ready",
crate::runtime::executor_utils::HTTP_RETRYABLE_BACKOFF_MS,
format!(
"UDB startup lifecycle is {}, DataBroker is not ready",
state.as_str()
),
))
}
}
pub(crate) async fn authorize(
&self,
security: &SecurityContext,
message_type: &str,
operation: &str,
) -> Result<String, Status> {
self.ensure_ready()?;
if let Some(detail) = self
.catalog
.compatibility_error(&security.client_catalog_version, &security.project_id)
{
let warn_only = self
.runtime_snapshot()
.config()
.service
.catalog_compat_warn_only;
let msg = format!(
"incompatible catalog version: client is '{}', active is '{}': {}",
security.client_catalog_version,
self.catalog.active().metadata.version,
detail
);
if warn_only {
tracing::warn!(
trace_id = security.trace_id,
project_id = security.project_id,
"{}",
msg
);
} else {
return Err(catalog_compatibility_status(operation, msg));
}
}
let safe = security.log_safe();
if self.runtime_snapshot().config().service.rate_limit_enabled && !safe.tenant_id.is_empty()
{
self.check_rate_limit(&safe.tenant_id, operation).await?;
}
tracing::debug!(
trace_id = security.trace_id,
correlation_id = safe.correlation_id,
tenant_id = safe.tenant_id,
purpose = safe.purpose,
service_identity = safe.service_identity,
message_type = message_type,
operation = operation,
"authorizing UDB request"
);
if message_type == "*" {
return Ok(Uuid::new_v4().to_string());
}
if security.tenant_id.trim().is_empty() {
return Err(Status::unauthenticated("tenant_id is required"));
}
if security.purpose.trim().is_empty() {
return Err(service_policy_denied(
"data_plane_authorize",
"purpose_required",
"purpose is required",
));
}
let principal = Principal::from_security_context(security, Vec::new());
let resource = ResourceRef::message(message_type);
let attributes = std::collections::BTreeMap::new();
let snapshot = self.current_authz_snapshot();
let decision = snapshot
.casbin_authorize(&AuthzQuery {
principal: &principal,
resource: &resource,
action: operation,
purpose: &security.purpose,
attributes: &attributes,
})
.await;
tracing::debug!(
trace_id = security.trace_id,
decision_id = decision.decision_id,
allowed = decision.allowed,
"authz casbin decision"
);
if decision.allowed {
Ok(decision.decision_id)
} else {
Err(service_policy_denied(
"data_plane_authorize",
decision.decision_id,
decision.deny_reason,
))
}
}
pub(crate) async fn authorize_message_item(
snapshot: &AuthzSnapshot,
security: &SecurityContext,
message_type: &str,
operation: &str,
) -> Result<String, Status> {
if security.tenant_id.trim().is_empty() {
return Err(Status::unauthenticated("tenant_id is required"));
}
if security.purpose.trim().is_empty() {
return Err(service_policy_denied(
"data_plane_authorize_item",
"purpose_required",
"purpose is required",
));
}
let principal = Principal::from_security_context(security, Vec::new());
let resource = ResourceRef::message(message_type);
let attributes = std::collections::BTreeMap::new();
let decision = snapshot
.casbin_authorize(&AuthzQuery {
principal: &principal,
resource: &resource,
action: operation,
purpose: &security.purpose,
attributes: &attributes,
})
.await;
if decision.allowed {
Ok(decision.decision_id)
} else {
Err(service_policy_denied(
"data_plane_authorize_item",
decision.decision_id,
decision.deny_reason,
))
}
}
pub(crate) fn require_portal_permission(
&self,
security: &SecurityContext,
operation: &str,
mutation: bool,
) -> Result<(), Status> {
let allowed = if security.has_scope("udb:admin") {
true
} else if mutation {
security.has_scope("udb:portal:operator") || security.has_scope("udb:portal:admin")
} else {
security.has_scope("udb:portal:viewer")
|| security.has_scope("udb:portal:operator")
|| security.has_scope("udb:portal:admin")
};
if allowed {
Ok(())
} else {
Err(service_policy_denied(
"portal_permission",
if mutation {
"portal_operator_required"
} else {
"portal_viewer_required"
},
format!(
"scope udb:admin{} is required for {operation}",
if mutation {
" or udb:portal:operator"
} else {
" or udb:portal:viewer"
}
),
))
}
}
#[cfg(not(feature = "redis"))]
pub(crate) async fn check_rate_limit(
&self,
_tenant_id: &str,
_operation: &str,
) -> Result<(), Status> {
static WARN_ONCE: std::sync::Once = std::sync::Once::new();
WARN_ONCE.call_once(|| {
tracing::warn!("rate limiting disabled in this build because the redis feature is off");
});
Ok(())
}
#[cfg(feature = "redis")]
pub(crate) async fn check_rate_limit(
&self,
tenant_id: &str,
operation: &str,
) -> Result<(), Status> {
let Some(redis) = self.runtime_snapshot().redis_clone() else {
return Ok(());
};
let window_secs = self
.runtime_snapshot()
.config()
.service
.rate_limit_window_secs
.max(1);
let max_rps = self
.runtime_snapshot()
.config()
.service
.rate_limit_max_per_window;
let unix_epoch = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let key = format!(
"udb:ratelimit:{}:{}:{}",
tenant_id,
operation,
unix_epoch / window_secs
);
let mut conn = {
let mut guard = self.rate_limit_redis.lock().await;
if guard.is_none() {
*guard = Some(
redis
.get_multiplexed_async_connection()
.await
.map_err(|e| {
crate::runtime::executor_utils::retryable_status(
"redis",
"rate_limit_connection",
crate::runtime::executor_utils::HTTP_RETRYABLE_BACKOFF_MS,
format!("rate limit redis error: {e}"),
)
})?,
);
}
guard
.as_ref()
.expect("rate limit redis connection just initialized")
.clone()
};
const RATE_LIMIT_LUA: &str = "local c = redis.call('INCR', KEYS[1]) \
if c == 1 then redis.call('EXPIRE', KEYS[1], ARGV[1]) end \
return c";
let script = redis::Script::new(RATE_LIMIT_LUA);
let first_attempt: Result<u64, redis::RedisError> = script
.key(&key)
.arg(window_secs)
.invoke_async(&mut conn)
.await;
let count: u64 = match first_attempt {
Ok(count) => count,
Err(_) => {
{
let mut guard = self.rate_limit_redis.lock().await;
*guard = None;
}
let mut fresh = redis
.get_multiplexed_async_connection()
.await
.map_err(|e| {
crate::runtime::executor_utils::retryable_status(
"redis",
"rate_limit_connection",
crate::runtime::executor_utils::HTTP_RETRYABLE_BACKOFF_MS,
format!("rate limit redis error: {e}"),
)
})?;
let count = script
.key(&key)
.arg(window_secs)
.invoke_async(&mut fresh)
.await
.map_err(|e| {
crate::runtime::executor_utils::retryable_status(
"redis",
"rate_limit_eval",
crate::runtime::executor_utils::HTTP_RETRYABLE_BACKOFF_MS,
format!("rate limit redis error: {e}"),
)
})?;
let mut guard = self.rate_limit_redis.lock().await;
*guard = Some(fresh);
count
}
};
if count > u64::from(max_rps) {
let retry_after_ms =
(window_secs.saturating_sub(unix_epoch % window_secs).max(1) as i64) * 1_000;
return Err(crate::runtime::executor_utils::quota_status(
"data_broker",
"distributed rate limit",
retry_after_ms,
format!(
"rate limit exceeded: {}/{} requests per {}s window",
count, max_rps, window_secs
),
));
}
Ok(())
}
pub(crate) fn record_grpc<T>(
&self,
method: &'static str,
started: Instant,
result: Result<Response<T>, Status>,
) -> Result<Response<T>, Status> {
let status = result
.as_ref()
.map(|_| "ok".to_string())
.unwrap_or_else(|err| format!("{:?}", err.code()).to_ascii_lowercase());
self.metrics
.record_grpc(method, &status, started.elapsed().as_secs_f64());
result
}
pub(crate) fn with_catalog_response_headers<T>(
&self,
mut response: Response<T>,
context: &crate::RequestContext,
) -> Response<T> {
let active = self.catalog.active_for(&context.project_id);
let metadata = response.metadata_mut();
insert_ascii_header(metadata, "x-udb-project-id", &active.metadata.project_id);
insert_ascii_header(metadata, "x-udb-catalog-version", &active.metadata.version);
insert_ascii_header(
metadata,
"x-udb-manifest-checksum",
&active.metadata.checksum,
);
let mut consistency = crate::runtime::consistency::ConsistencyPolicy::from_request_context(
&context.consistency,
context.max_replica_lag_ms,
context.primary_read,
context.eventual_consistency_allowed,
);
let mut read_fence_invalid = false;
if !context.read_fence_json.trim().is_empty() {
match serde_json::from_str::<crate::runtime::consistency::ReadFence>(
&context.read_fence_json,
) {
Ok(fence) => {
consistency = consistency.with_fence(fence);
}
Err(_) => {
read_fence_invalid = true;
}
}
}
insert_ascii_header(
metadata,
"x-udb-consistency-mode",
consistency.mode.as_str(),
);
insert_ascii_header(
metadata,
"x-udb-read-fence-present",
if !consistency.fence.is_empty() {
"true"
} else {
"false"
},
);
insert_ascii_header(
metadata,
"x-udb-read-fence-honored",
if !consistency.fence.is_empty() && consistency.mode.honours_fence() {
"true"
} else {
"false"
},
);
insert_ascii_header(
metadata,
"x-udb-read-fence-invalid",
if read_fence_invalid { "true" } else { "false" },
);
insert_ascii_header(
metadata,
"x-udb-primary-read",
if context.primary_read {
"true"
} else {
"false"
},
);
insert_ascii_header(
metadata,
"x-udb-eventual-consistency-allowed",
if context.eventual_consistency_allowed {
"true"
} else {
"false"
},
);
response
}
pub(crate) async fn with_mutation_response_headers(
&self,
mut mutation: MutationResponse,
context: &crate::RequestContext,
) -> Response<MutationResponse> {
let active = self.catalog.active_for(&context.project_id);
let receipt = if let Some(receipt) = mutation.write_receipt.as_ref() {
crate::runtime::consistency::WriteReceipt::from_proto(receipt)
} else if !mutation.write_receipt_json.trim().is_empty() {
serde_json::from_str::<crate::runtime::consistency::WriteReceipt>(
&mutation.write_receipt_json,
)
.unwrap_or_else(|_| crate::runtime::consistency::WriteReceipt::empty())
} else {
self.runtime_snapshot()
.current_write_receipt(&active.metadata.checksum)
.await
};
let receipt_json = serde_json::to_string(&receipt).unwrap_or_default();
mutation.write_receipt_json = receipt_json.clone();
mutation.write_receipt = Some(receipt.to_proto());
let mut response = self.with_catalog_response_headers(Response::new(mutation), context);
insert_ascii_header(
response.metadata_mut(),
"x-udb-write-receipt",
&receipt_json,
);
response
}
pub(crate) async fn execute_with_channel<F, Fut, T>(
&self,
op: crate::runtime::channels::OperationChannel,
f: F,
) -> Result<T, Status>
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = Result<T, Status>>,
{
self.execute_with_channel_scoped(op, None, None, f).await
}
pub(crate) async fn execute_with_channel_scoped<F, Fut, T>(
&self,
op: crate::runtime::channels::OperationChannel,
context: Option<&crate::RequestContext>,
backend: Option<&str>,
f: F,
) -> Result<T, Status>
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = Result<T, Status>>,
{
let runtime = self.runtime_snapshot();
let channels = runtime.channels();
let project = context
.map(|context| context.project_id.as_str())
.unwrap_or("default");
let tenant = context
.map(|context| context.tenant_id.as_str())
.unwrap_or("anonymous");
let tenant_hash = tenant_hash_label(tenant);
let backend_label = backend
.or_else(|| context.and_then(|context| non_empty(&context.target_backend)))
.unwrap_or("default");
let instance_label = context
.and_then(|context| non_empty(&context.target_instance))
.unwrap_or("default");
let cost = op.default_cost();
let _permit = match channels
.acquire_fair_with_backpressure(
op,
context.map(|context| context.tenant_id.as_str()),
context.map(|context| context.project_id.as_str()),
Some(backend_label),
context.map(|context| context.target_instance.as_str()),
cost,
)
.await
{
Ok(permit) => {
self.metrics.record_fair_admission(
project,
&tenant_hash,
backend_label,
instance_label,
op.as_str(),
"accepted",
);
self.metrics.add_fair_cost(
project,
&tenant_hash,
backend_label,
instance_label,
op.as_str(),
f64::from(cost),
);
permit
}
Err(e) => {
self.metrics.inc_channel_rejected(op.as_str());
self.metrics.record_fair_admission(
project,
&tenant_hash,
backend_label,
instance_label,
op.as_str(),
"rejected",
);
return Err(e);
}
};
self.metrics.inc_channel_inflight(op.as_str());
let start = Instant::now();
let timeout_secs = channels.deadline_secs(op, backend);
let res = tokio::time::timeout(Duration::from_secs(timeout_secs), f()).await;
self.metrics.dec_channel_inflight(op.as_str());
self.metrics
.observe_channel_latency(op.as_str(), start.elapsed().as_secs_f64());
match res {
Ok(Ok(val)) => Ok(val),
Ok(Err(e)) => Err(e),
Err(_) => {
self.metrics.inc_channel_timeout(op.as_str());
Err(crate::runtime::executor_utils::deadline_exceeded_status(
backend_label,
format!("{} channel", op.as_str()),
crate::runtime::executor_utils::HTTP_RETRYABLE_BACKOFF_MS,
format!("{} channel timeout", op.as_str()),
))
}
}
}
}
fn service_metrics_recorder() -> Arc<dyn MetricsRecorder> {
match PrometheusMetrics::new() {
Ok(metrics) => Arc::new(metrics),
Err(err) => {
tracing::warn!("prometheus metrics disabled: {err}");
Arc::new(NoopMetrics)
}
}
}
async fn admit_stream_batch_item(
channels: &crate::runtime::channels::ChannelManager,
metrics: &Arc<dyn MetricsRecorder>,
context: &crate::RequestContext,
op: crate::runtime::channels::OperationChannel,
backend: &'static str,
) -> Result<crate::runtime::channels::ChannelPermit, Status> {
let project = non_empty(&context.project_id).unwrap_or("default");
let tenant_hash = tenant_hash_label(&context.tenant_id);
let instance = non_empty(&context.target_instance).unwrap_or("default");
match channels
.acquire_fair_with_backpressure(
op,
Some(&context.tenant_id),
Some(&context.project_id),
Some(backend),
Some(&context.target_instance),
op.default_cost(),
)
.await
{
Ok(permit) => {
metrics.record_fair_admission(
project,
&tenant_hash,
backend,
instance,
op.as_str(),
"accepted",
);
metrics.add_fair_cost(
project,
&tenant_hash,
backend,
instance,
op.as_str(),
f64::from(op.default_cost()),
);
Ok(permit)
}
Err(err) => {
metrics.inc_channel_rejected(op.as_str());
metrics.record_fair_admission(
project,
&tenant_hash,
backend,
instance,
op.as_str(),
"rejected",
);
Err(err)
}
}
}
fn check_backend_capability(
backend: &str,
operation: &str,
capability_fn: impl Fn(&crate::planning::backend::BackendCapability) -> bool,
) -> Result<(), Status> {
use crate::planning::backend::BackendKind;
let backend_base = backend_selector_base(backend);
let Some(kind) = BackendKind::from_store_kind("", backend_base) else {
return Err(unknown_backend_status(backend));
};
let state = crate::backend::support_state_for_kind(&kind);
if !state.is_runtime_supported() {
return Err(backend_runtime_unsupported_status(
backend,
operation,
state.diagnostic(kind.as_str()),
));
}
let cap = kind.capabilities();
if capability_fn(&cap) {
Ok(())
} else {
Err(crate::runtime::executor_utils::capability_status(
backend,
operation,
crate::backend::UNSUPPORTED_OPERATION_CODE,
format!(
"{}: backend '{backend}' does not support operation '{operation}'",
crate::backend::UNSUPPORTED_OPERATION_CODE
),
))
}
}
fn check_generic_dispatch_operation(backend: &str, operation: &str) -> Result<(), Status> {
use crate::planning::backend::BackendKind;
let backend_base = backend_selector_base(backend);
let Some(kind) = BackendKind::from_store_kind("", backend_base) else {
return Err(unknown_backend_status(backend));
};
let state = crate::backend::support_state_for_kind(&kind);
if !state.is_runtime_supported() {
return Err(backend_runtime_unsupported_status(
backend,
operation,
state.diagnostic(kind.as_str()),
));
}
let supported = match operation {
"ping" | "probe" | "ensure_resource" | "drop_resource" | "list_resources" | "query"
| "mutate" | "transaction" | "search" | "get_object" | "put_object" | "delete_object" => {
kind.supports_operation(operation)
}
other => {
return Err(unknown_generic_operation_status(other));
}
};
if supported {
Ok(())
} else {
Err(crate::runtime::executor_utils::capability_status(
backend,
operation,
crate::backend::UNSUPPORTED_OPERATION_CODE,
format!(
"{}: backend '{backend}' does not support operation '{operation}'",
crate::backend::UNSUPPORTED_OPERATION_CODE
),
))
}
}
fn backend_runtime_unsupported_status(backend: &str, operation: &str, message: String) -> Status {
crate::runtime::executor_utils::capability_status(
backend,
operation,
"backend_runtime_support",
message,
)
}
fn unknown_backend_status(backend: &str) -> Status {
crate::runtime::executor_utils::invalid_argument_fields(
format!("unknown backend '{backend}'"),
[("backend", "must name a supported backend")],
)
}
fn unknown_generic_operation_status(operation: &str) -> Status {
crate::runtime::executor_utils::invalid_argument_fields(
format!(
"unknown operation '{operation}'; allowed: ping, probe, ensure_resource, drop_resource, list_resources, query, mutate, transaction, search, get_object, put_object, delete_object"
),
[(
"operation",
"must be a supported generic dispatch operation",
)],
)
}
fn backend_selector_base(selector: &str) -> &str {
selector
.split_once(':')
.map(|(backend, _)| backend)
.or_else(|| selector.split_once('.').map(|(backend, _)| backend))
.unwrap_or(selector)
}
fn bounded_list_limit(limit: i32) -> i32 {
if limit <= 0 { 100 } else { limit.min(1000) }
}
fn non_empty(value: &str) -> Option<&str> {
let value = value.trim();
(!value.is_empty()).then_some(value)
}
fn catalog_compatibility_status(operation: &str, message: String) -> Status {
crate::runtime::executor_utils::schema_status(
tonic::Code::FailedPrecondition,
"catalog",
operation,
"catalog_version_incompatible",
message,
)
}
fn service_policy_denied(
operation: impl Into<String>,
policy_decision_id: impl Into<String>,
message: impl Into<String>,
) -> Status {
crate::runtime::executor_utils::policy_status_with_code(
tonic::Code::PermissionDenied,
operation,
policy_decision_id,
message,
)
}
fn require_admin_scope(security: &SecurityContext) -> Result<(), Status> {
if security.has_scope("udb:admin") {
Ok(())
} else {
Err(service_policy_denied(
"admin_scope",
"admin_scope_required",
"scope udb:admin is required",
))
}
}
fn rls_bypass_ack(spec_json: &str) -> bool {
serde_json::from_str::<serde_json::Value>(spec_json)
.ok()
.and_then(|value| {
value
.get("udb_allow_rls_bypass")
.or_else(|| value.get("allow_rls_bypass"))
.and_then(|flag| flag.as_bool())
})
.unwrap_or(false)
}
fn contains_rls_bypass_sql(spec_json: &str) -> bool {
let lower = spec_json.to_ascii_lowercase();
lower.contains("truncate ")
|| lower.contains(" truncate")
|| lower.contains(" cascade")
|| lower.contains("disable row level security")
|| lower.contains("alter table")
|| lower.contains("drop table")
|| lower.contains("create unique index")
|| lower.contains(" unique ")
|| lower.contains(" primary key")
}
fn guard_rls_bypass_operation(operation: &str, spec_json: &str) -> Result<(), Status> {
let bypass_like = matches!(operation, "drop_resource")
|| (matches!(operation, "query" | "mutate" | "transaction")
&& contains_rls_bypass_sql(spec_json));
if bypass_like && !rls_bypass_ack(spec_json) {
return Err(crate::runtime::executor_utils::policy_status(
"generic_dispatch_rls_bypass",
"rls_bypass_review_required",
"operation may bypass tenant isolation/RLS; set spec_json.udb_allow_rls_bypass=true after explicit tenant-scope review",
));
}
Ok(())
}
#[derive(Clone)]
struct WebrtcPeerTokenAuth {
security: SecurityConfig,
}
impl WebrtcPeerTokenAuth {
fn new() -> Self {
Self {
security: SecurityConfig::current(),
}
}
}
impl tonic::service::Interceptor for WebrtcPeerTokenAuth {
fn call(&mut self, request: Request<()>) -> Result<Request<()>, Status> {
let metadata = request.metadata();
let auth_header = metadata
.get("authorization")
.and_then(|value| value.to_str().ok())
.unwrap_or_default();
let token = auth_header.strip_prefix("Bearer ").ok_or_else(|| {
Status::unauthenticated(
"missing or invalid authorization header (WebRTC peer bearer required)",
)
})?;
let claims =
validate_bearer_token(&self.security, token).map_err(Status::unauthenticated)?;
let scopes = claims.resolved_scopes();
let allowed = scopes.iter().any(|scope| {
matches!(
scope.as_str(),
"*" | "udb:*" | "udb:webrtc:*" | "udb:webrtc:peer" | "udb:webrtc:signal"
)
});
if !allowed {
return Err(service_policy_denied(
"webrtc_peer_token",
"webrtc_peer_scope_required",
"scope udb:webrtc:peer or udb:webrtc:signal is required",
));
}
if let Some(header_tenant) = metadata
.get("x-tenant-id")
.and_then(|value| value.to_str().ok())
&& !header_tenant.trim().is_empty()
&& claims.tenant_id.as_deref().unwrap_or_default() != header_tenant
{
return Err(service_policy_denied(
"webrtc_peer_token",
"webrtc_peer_tenant_mismatch",
"x-tenant-id must match the peer token tenant",
));
}
Ok(request)
}
}
fn tenant_hash_label(tenant: &str) -> String {
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
hasher.update(tenant.as_bytes());
let digest = hasher.finalize();
digest[..8]
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>()
}
fn page_offset(page_token: &str) -> i32 {
page_token.parse::<i32>().unwrap_or_default().max(0)
}
fn next_page_token(offset: i32, limit: i32, returned: i32) -> String {
if returned >= limit {
(offset + returned).to_string()
} else {
String::new()
}
}
pub fn context_from_metadata(metadata: &tonic::metadata::MetadataMap) -> crate::RequestContext {
let header = |name: &str| {
metadata
.get(name)
.and_then(|value| value.to_str().ok())
.unwrap_or_default()
.to_string()
};
crate::RequestContext {
tenant_id: header("x-tenant-id"),
purpose: header("x-purpose"),
correlation_id: header("x-correlation-id"),
user_id: header("x-user-id"),
project_id: header("x-udb-project-id"),
scopes: header("x-scopes")
.split(',')
.map(str::trim)
.filter(|scope| !scope.is_empty())
.map(ToString::to_string)
.collect(),
consistency: header("x-udb-consistency"),
max_replica_lag_ms: header("x-udb-max-replica-lag-ms")
.parse::<u64>()
.unwrap_or_default(),
client_catalog_version: header("x-udb-client-catalog-version"),
target_backend: header("x-udb-target-backend"),
target_instance: header("x-udb-target-instance"),
routing_policy: header("x-udb-routing-policy"),
primary_read: matches!(
header("x-udb-primary-read").to_ascii_lowercase().as_str(),
"1" | "true" | "yes" | "on"
),
eventual_consistency_allowed: matches!(
header("x-udb-eventual-consistency-allowed")
.to_ascii_lowercase()
.as_str(),
"1" | "true" | "yes" | "on"
) || matches!(
header("x-udb-consistency")
.to_ascii_lowercase()
.replace('-', "_")
.as_str(),
"eventual" | "eventual_consistency"
),
read_fence_json: header("x-udb-read-fence"),
service_identity: header("x-service-identity"),
decision_id: String::new(),
}
}
fn backend_instance_status(
instance: &crate::runtime::core::RuntimeBackendInstance,
) -> BackendInstanceStatus {
BackendInstanceStatus {
backend: instance.backend.clone(),
instance_name: instance.name.clone(),
role: instance.role.clone(),
enabled: instance.enabled,
configured: instance.configured,
connected: instance.connected,
read_weight: instance.read_weight,
write_weight: instance.write_weight,
labels: instance.labels.clone(),
capabilities: instance.capabilities.clone(),
routing_status: if !instance.enabled {
"disabled".to_string()
} else if instance.circuit_open {
"circuit_open".to_string()
} else if instance.connected {
"available".to_string()
} else if instance.configured {
"degraded".to_string()
} else {
"unconfigured".to_string()
},
healthy: instance.healthy,
circuit_open: instance.circuit_open,
}
}
fn parse_catalog_manifest_payload(bytes: &[u8]) -> Result<CatalogManifest, Status> {
if bytes.is_empty() {
return Err(crate::runtime::executor_utils::invalid_argument_fields(
"manifest_json is required",
[(
"manifest_json",
"must contain a CatalogManifest JSON payload",
)],
));
}
serde_json::from_slice::<CatalogManifest>(bytes).map_err(|err| {
crate::runtime::executor_utils::invalid_argument_fields(
format!("manifest_json is not a CatalogManifest: {err}"),
[("manifest_json", "must decode as a CatalogManifest")],
)
})
}
fn catalog_payload_version(
bytes: &[u8],
manifest: &CatalogManifest,
fallback_active_version: &str,
) -> String {
if !bytes.is_empty()
&& let Ok(value) = serde_json::from_slice::<serde_json::Value>(bytes)
&& let Some(version) = value.get("version").and_then(|v| v.as_str())
&& !version.trim().is_empty()
{
return version.trim().to_string();
}
if !fallback_active_version.trim().is_empty() {
return fallback_active_version.trim().to_string();
}
if !manifest.generator_version.trim().is_empty() {
return format!("generator-{}", manifest.generator_version.trim());
}
if !manifest.checksum_sha256.trim().is_empty() {
return manifest.checksum_sha256.chars().take(12).collect();
}
"unversioned".to_string()
}
#[cfg(unix)]
fn spawn_uds_data_plane(
service: DataBrokerService,
grpc_timeout: Duration,
grpc_max_concurrent: usize,
) -> Option<tokio::task::JoinHandle<()>> {
let path = std::env::var("UDB_DATA_UDS_PATH")
.ok()
.map(|p| p.trim().to_string())
.filter(|p| !p.is_empty())?;
Some(tokio::spawn(async move {
let _ = tokio::fs::remove_file(&path).await;
let listener = match tokio::net::UnixListener::bind(&path) {
Ok(listener) => listener,
Err(err) => {
tracing::error!(
uds_path = %path,
error = %err,
"UDB UDS data-plane bind failed; TCP listener still serving"
);
return;
}
};
tracing::info!(uds_path = %path, "UDB DataBroker UDS listener ready");
let incoming = futures::stream::unfold(listener, |listener| async move {
let conn = listener.accept().await.map(|(stream, _addr)| stream);
Some((conn, listener))
});
let layer = tower::ServiceBuilder::new()
.layer(crate::runtime::otel::TraceExtractLayer::new())
.layer(crate::runtime::credential_layer::CredentialResolveLayer::new())
.timeout(grpc_timeout)
.concurrency_limit(grpc_max_concurrent)
.into_inner();
let reflection = match tonic_reflection::server::Builder::configure()
.register_encoded_file_descriptor_set(UDB_FILE_DESCRIPTOR_SET)
.build_v1()
{
Ok(reflection) => reflection,
Err(err) => {
tracing::error!(error = %err, "UDB UDS reflection build failed; UDS listener aborted");
return;
}
};
if let Err(err) = tonic::transport::Server::builder()
.layer(layer)
.add_service(reflection)
.add_service(DataBrokerServer::new(service))
.serve_with_incoming_shutdown(incoming, shutdown_signal())
.await
{
tracing::error!(uds_path = %path, error = %err, "UDB UDS data-plane listener exited with error");
}
}))
}
pub async fn serve(
manifest: CatalogManifest,
schemas: Vec<ProtoSchema>,
addr: SocketAddr,
) -> Result<(), Box<dyn std::error::Error>> {
let runtime = DataBrokerRuntime::try_from_env().await.map_err(|err| {
std::io::Error::other(format!("UDB startup config validation failed: {err}"))
})?;
let runtime_config = runtime.config().clone();
native_registry::install_native_service_runtime_config(&runtime_config);
crate::runtime::preflight::log_findings(&crate::runtime::preflight::evaluate(
&runtime_config,
addr,
));
let transport_violation = validate_secure_transport(&runtime_config.service).err();
if let Some(err) = &transport_violation {
if runtime_config.service.require_secure_transport || runtime_config.service.mtls_required {
return Err(std::io::Error::other(format!(
"secure transport startup gate failed: {err}"
))
.into());
}
}
{
let transport = transport_violation.into_iter().collect::<Vec<_>>();
let violations = crate::runtime::security::hardened_startup_violations(&transport);
if violations.is_empty() {
if !transport.is_empty()
|| crate::runtime::security::SecurityConfig::current()
.validate_production()
.is_err()
{
tracing::warn!(
"security posture advisory (not enforced in dev mode): set UDB_ENV=production \
or UDB_FAIL_CLOSED to make these fatal"
);
}
} else {
return Err(std::io::Error::other(format!(
"production/secure-transport startup gate failed (enterprise mode refuses \
insecure transport): {}",
violations.join("; ")
))
.into());
}
}
{
let raw_profile = std::env::var("UDB_COMPLIANCE_PROFILE").unwrap_or_default();
match crate::runtime::security::selected_compliance_profile() {
Some(profile) => {
let cfg = crate::runtime::security::SecurityConfig::current();
let facts = cfg.compliance_profile_facts();
if let Err(violations) = cfg.validate_compliance_profile(profile, &facts) {
return Err(std::io::Error::other(format!(
"compliance profile '{}' startup gate failed: {}",
profile.as_str(),
violations.join("; ")
))
.into());
}
tracing::info!(profile = profile.as_str(), "compliance profile gate passed");
}
None if !raw_profile.trim().is_empty()
&& !raw_profile.trim().eq_ignore_ascii_case("none") =>
{
return Err(std::io::Error::other(format!(
"unknown UDB_COMPLIANCE_PROFILE '{}' (expected soc2 | iso27001 | pci_hipaa)",
raw_profile.trim()
))
.into());
}
None => {}
}
}
if runtime_config.has_backup() {
let backup_tenant = std::env::var("UDB_BACKUP_TENANT_ID").unwrap_or_default();
let backup_tenant = backup_tenant.trim();
let privileged = std::env::var("UDB_ALLOW_CROSS_TENANT_BACKUP")
.map(|v| matches!(v.trim().to_ascii_lowercase().as_str(), "1" | "true" | "yes"))
.unwrap_or(false);
let movement = crate::runtime::tenant_movement::TenantMovementRequest {
operation: crate::runtime::tenant_movement::TenantMovementOperation::BackupExport,
tenant_id: backup_tenant,
target_tenant_id: None,
tenant_filter_present: !backup_tenant.is_empty(),
privileged_cross_tenant: privileged,
};
match crate::runtime::tenant_movement::validate_tenant_movement_scope(&movement) {
Ok(()) => tracing::info!(
tenant_scoped = !backup_tenant.is_empty(),
privileged_cross_tenant = privileged,
"backup tenant-scope gate passed"
),
Err(violation) if crate::runtime::security::fail_closed_mode() => {
return Err(std::io::Error::other(format!(
"backup tenant-scope startup gate failed: {violation}. Set \
UDB_BACKUP_TENANT_ID=<tenant> for a tenant-scoped backup, or \
UDB_ALLOW_CROSS_TENANT_BACKUP=true to acknowledge a privileged \
broker-wide backup."
))
.into());
}
Err(violation) => tracing::warn!(
violation = %violation,
"backup tenant-scope advisory (not enforced in dev mode): set \
UDB_FAIL_CLOSED or UDB_ENV=production to make this fatal"
),
}
}
let _ = crate::runtime::service::method_security::method_security_registry();
if !runtime.postgres_configured() {
let primary = &runtime.config().primary;
match crate::runtime::core::postgres_dsn_from_config(primary) {
Some(resolved_dsn) => {
return Err(format!(
"PostgreSQL startup health gate failed: a connection string was resolved \
({}) but the database could not be reached. Verify the host/port, \
credentials and TLS settings.",
crate::generation::dsn::redact_dsn(&resolved_dsn)
)
.into());
}
None => {
return Err(
"PostgreSQL startup health gate failed: no PostgreSQL configuration found. \
Provide a connection URL via UDB_PG_DSN / DATABASE_URL, or libpq-style \
component variables (PGHOST + PGDATABASE [+ PGUSER/PGPASSWORD/PGPORT/\
PGSSLMODE])."
.into(),
);
}
}
}
crate::runtime::authz::validate_casbin_model()
.await
.map_err(|err| std::io::Error::other(format!("authz startup gate failed: {err}")))?;
if !runtime.qdrant_configured() {
tracing::warn!("Qdrant startup health gate degraded: vector RPCs will return UNAVAILABLE");
}
#[cfg(feature = "s3")]
if !runtime.s3_configured() {
tracing::warn!(
"S3/MinIO startup health gate degraded: object RPCs will return UNAVAILABLE"
);
}
let system_report = runtime.ensure_system_catalog().await?;
tracing::info!(
schema = system_report.schema,
statements_applied = system_report.statements_applied,
"UDB internal system catalog is ready"
);
let prometheus_metrics = match PrometheusMetrics::new() {
Ok(metrics) => Some(Arc::new(metrics)),
Err(err) => {
tracing::warn!("prometheus metrics disabled: {err}");
None
}
};
let metrics: Arc<dyn MetricsRecorder> = prometheus_metrics
.as_ref()
.map(|metrics| metrics.clone() as Arc<dyn MetricsRecorder>)
.unwrap_or_else(|| Arc::new(NoopMetrics));
let metrics_socket: SocketAddr = runtime_config.service.metrics_addr.parse()?;
if let Some(prometheus_metrics) = prometheus_metrics.clone() {
tokio::spawn(metrics_http_server(
prometheus_metrics,
runtime.clone(),
metrics_socket,
runtime_config.service.metrics_allowed_cidr.clone(),
));
}
tokio::spawn(cdc_metrics_poller(runtime.clone(), metrics.clone()));
let lifecycle_state = Arc::new(RwLock::new(FsmState::Initialising));
let lifecycle_started = Instant::now();
let startup_force_sync = startup_bool_env("UDB_STARTUP_FORCE_SYNC");
let startup_dry_run = startup_bool_env("UDB_STARTUP_DRY_RUN");
tracing::info!(
startup_force_sync,
startup_dry_run,
"running UDB startup lifecycle"
);
let (ready_sql_artifacts, ready_vector_collections, ready_object_buckets) =
match run_startup_lifecycle(
&runtime,
&manifest,
&schemas,
startup_force_sync,
startup_dry_run,
)
.await
{
Ok(report) => {
let elapsed = lifecycle_started.elapsed().as_secs_f64();
metrics.inc_runs_total("completed");
metrics.observe_run_duration("completed", elapsed);
metrics.set_pending_files(report.pending_migration_files);
for op in &report.migration_metric_operations {
metrics.inc_operations_total(&op.kind, &op.schema, &op.safety);
if op.safety == "blocked" || op.safety == "requires_review" {
metrics.set_blocked_operations(&op.schema, &op.kind, 1);
}
}
for _ in &report.warnings {
metrics.inc_lint_warnings("startup");
}
tracing::info!(
run_id = report.run_id,
applied_sql_artifacts = report.applied_sql_artifacts,
verified_tables = report.verified_tables,
"UDB startup lifecycle completed"
);
if let Ok(mut state) = lifecycle_state.write() {
*state = FsmState::Completed;
}
let counts = (
report.applied_sql_artifacts,
report.verified_vector_collections,
report.verified_object_buckets,
);
if !startup_dry_run {
if let Ok(pool) = runtime.pg_pool() {
if let Err(err) = auth_service::seed_system_authz_defaults(pool).await {
tracing::warn!(error = %err, "seed system authz defaults failed (non-fatal)");
}
}
}
counts
}
Err(err) => {
metrics.inc_runs_total("error");
metrics.observe_run_duration("error", lifecycle_started.elapsed().as_secs_f64());
if let Ok(mut state) = lifecycle_state.write() {
*state = FsmState::Error;
}
return Err(err.into());
}
};
runtime.mark_indeterminate_sagas().await;
if let Some(pg_pool) = runtime.pg_pool_clone() {
let sys_config = crate::runtime::system::SystemCatalogConfig::current();
let singleton_relation = runtime_config.cdc.lock_log_relation();
let recovery_config = crate::runtime::xa_recovery::RecoveryConfig::default();
let grace = std::env::var("UDB_XA_RECOVERY_GRACE_SECS")
.ok()
.and_then(|v| v.parse::<i64>().ok())
.unwrap_or(300);
#[allow(unused_mut)]
let mut registry = crate::runtime::xa_recovery::default_indoubt_registry(&pg_pool);
#[cfg(feature = "mysql")]
{
let mut instance_names: Vec<&String> = runtime.mysql_instances.keys().collect();
instance_names.sort();
for name in instance_names {
if let Some(mysql_pool) = runtime.mysql_instances.get(name) {
registry.register(std::sync::Arc::new(
crate::runtime::xa_recovery::MysqlInDoubtParticipant {
label: format!("mysql:{name}"),
pool: mysql_pool.clone(),
},
));
}
}
if let Some(mysql_pool) = runtime.mysql_pool_for_instance("primary") {
registry.register(std::sync::Arc::new(
crate::runtime::xa_recovery::MysqlInDoubtParticipant {
label: "mysql".to_string(),
pool: mysql_pool.clone(),
},
));
}
}
let xa_recovery_lease_ttl = std::cmp::max(
recovery_config.interval,
crate::runtime::singleton::WORKER_SINGLETON_LEASE_TTL,
);
tokio::spawn(async move {
let mut interval = tokio::time::interval(recovery_config.interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
interval.tick().await;
match crate::runtime::singleton::run_once(
&pg_pool,
&singleton_relation,
crate::runtime::singleton::WORKER_XA_RECOVERY,
xa_recovery_lease_ttl,
|| async {
crate::runtime::xa_recovery::run_xa_recovery_pass(
&pg_pool,
&sys_config,
®istry,
&recovery_config,
grace,
)
.await
},
)
.await
{
Ok(Some(Ok((ledger, abandoned)))) => {
if ledger > 0 {
tracing::warn!(
"XA recovery: drove {ledger} ledger in-doubt transaction(s) terminal"
);
}
if abandoned > 0 {
tracing::warn!(
"XA recovery: drove {abandoned} aged prepared transaction(s) terminal"
);
}
}
Ok(Some(Err(e))) => tracing::warn!("XA recovery sweep failed: {e}"),
Ok(None) => {
tracing::debug!("XA recovery skipped: singleton lease held by peer")
}
Err(e) => tracing::warn!("XA recovery sweep failed: {e}"),
}
}
});
tracing::info!("XA recovery worker started (periodic, lease-gated)");
}
if crate::runtime::saga::SagaRecoveryWorker::is_enabled_with_settings(&runtime_config.saga) {
if let Some(store) = runtime.default_system_stores() {
let worker = crate::runtime::saga::SagaRecoveryWorker::with_settings(
store,
&runtime_config.saga,
)
.with_compensators(runtime.saga_compensator_registry())
.with_metrics(metrics.clone());
tokio::spawn(async move { worker.run_forever().await });
tracing::info!("saga recovery worker started");
} else {
tracing::warn!("saga recovery worker disabled: no canonical store is registered");
}
}
let scheduled_views = runtime.start_materialized_view_refresh(&manifest);
if scheduled_views > 0 {
tracing::info!(
scheduled_views = scheduled_views,
"scheduled materialized view auto-refresh tasks"
);
}
let abac_default_allow = runtime_config.service.abac_default_allow;
#[cfg(feature = "kafka")]
let cdc_engine = start_cdc_engine(&runtime, metrics.clone()).await;
#[cfg(not(feature = "kafka"))]
let cdc_engine: Option<Arc<CdcEngine>> = {
let _ = &metrics;
None
};
let mut service = DataBrokerService::with_runtime_and_state(
manifest,
runtime,
lifecycle_state,
metrics.clone(),
cdc_engine,
abac_default_allow,
);
spawn_config_reload_watcher(service.clone());
if let (Some(pg_pool), Some(store)) = (
service.runtime_snapshot().pg_pool_clone(),
service.runtime_snapshot().default_system_stores(),
) {
use crate::runtime::projection::{
ProjectionEngine, ProjectionWorker, ReconciliationWorker,
};
let config = crate::runtime::system::SystemCatalogConfig::current();
let engine = Arc::new(ProjectionEngine::new(pg_pool.clone(), config));
service.projection_engine = Some(Arc::clone(&engine));
let singleton_relation = service.runtime_snapshot().config().cdc.lock_log_relation();
if ProjectionWorker::is_enabled() {
let metrics: Arc<dyn MetricsRecorder> = service.metrics.clone();
let singleton_pool = pg_pool.clone();
let singleton_relation = singleton_relation.clone();
let runtime = service.runtime_snapshot().clone();
let store = store.clone();
tokio::spawn(async move {
loop {
let metrics = metrics.clone();
let runtime = runtime.clone();
let store = store.clone();
match crate::runtime::singleton::run_while_leader(
&singleton_pool,
&singleton_relation,
crate::runtime::singleton::WORKER_PROJECTION_MATERIALIZER,
crate::runtime::singleton::WORKER_SINGLETON_LEASE_TTL,
|| async move {
ProjectionWorker::new(store, runtime, metrics)
.run_forever()
.await;
Ok::<(), String>(())
},
)
.await
{
Ok(Some(Ok(()))) => {}
Ok(Some(Err(err))) => {
tracing::warn!("projection materialization worker exited: {err}")
}
Ok(None) => tracing::debug!(
"projection materialization worker idle: singleton lease held by peer"
),
Err(err) => tracing::warn!(
"projection materialization worker singleton lease failed: {err}"
),
}
tokio::time::sleep(crate::runtime::singleton::WORKER_SINGLETON_RETRY_SLEEP)
.await;
}
});
tracing::info!("projection materialization worker started");
}
if ReconciliationWorker::is_enabled() {
let metrics: Arc<dyn MetricsRecorder> = service.metrics.clone();
let active_catalog = service.catalog.active();
let manifest = active_catalog.manifest.clone();
let project_id = active_catalog.metadata.project_id.clone();
let singleton_pool = pg_pool.clone();
let singleton_relation = singleton_relation.clone();
let worker_pool = pg_pool.clone();
let store = store.clone();
tokio::spawn(async move {
loop {
let metrics = metrics.clone();
let manifest = manifest.clone();
let project_id = project_id.clone();
let store = store.clone();
let worker_pool = worker_pool.clone();
match crate::runtime::singleton::run_while_leader(
&singleton_pool,
&singleton_relation,
crate::runtime::singleton::WORKER_PROJECTION_RECONCILIATION,
crate::runtime::singleton::WORKER_SINGLETON_LEASE_TTL,
|| async move {
ReconciliationWorker::new(
worker_pool,
store,
metrics,
manifest,
project_id,
)
.run_forever()
.await;
Ok::<(), String>(())
},
)
.await
{
Ok(Some(Ok(()))) => {}
Ok(Some(Err(err))) => {
tracing::warn!("projection reconciliation worker exited: {err}")
}
Ok(None) => tracing::debug!(
"projection reconciliation worker idle: singleton lease held by peer"
),
Err(err) => tracing::warn!(
"projection reconciliation worker singleton lease failed: {err}"
),
}
tokio::time::sleep(crate::runtime::singleton::WORKER_SINGLETON_RETRY_SLEEP)
.await;
}
});
tracing::info!("projection reconciliation worker started");
}
} else {
tracing::warn!(
"projection engine disabled: PostgreSQL pool and/or canonical store not available"
);
}
let health_runtime = service.runtime_snapshot();
let health_service = handlers_meta::build_listener_health_service(
handlers_meta::HealthPlane::DataBroker,
&runtime_config,
Some(health_runtime.as_ref()),
)
.await;
{
use sha2::Digest;
let mut hasher = sha2::Sha256::new();
for t in &service.catalog.active().manifest.tables {
hasher.update(t.message_name.as_bytes());
}
let checksum = format!("{:x}", hasher.finalize());
let mut enabled_backends = service.runtime_snapshot().enabled_backend_names();
enabled_backends.sort();
enabled_backends.dedup();
let auth_addr = {
let cp = runtime_config.native_services.control_plane_addr.trim();
if cp.is_empty() {
format!("127.0.0.1:{}", addr.port().saturating_add(10))
} else {
cp.to_string()
}
};
let table_count = service.catalog.active().manifest.tables.len();
tracing::info!(
data_addr = %addr,
auth_addr = %auth_addr,
schema_checksum = %checksum,
schemas = table_count,
table_count,
store_count = service.catalog.active().manifest.stores.len(),
sql_artifacts = ready_sql_artifacts,
vector_collections = ready_vector_collections,
object_buckets = ready_object_buckets,
enabled_backends = ?enabled_backends,
protocol_version = UDB_PROTOCOL_VERSION,
cdc_enabled = service.cdc_engine.is_some(),
supported_rpcs = SUPPORTED_RPC_NAMES.len(),
"UDB DataBroker is ready: data={addr} auth={auth_addr} schemas={table_count} \
sql_artifacts={ready_sql_artifacts} vector_collections={ready_vector_collections} \
object_buckets={ready_object_buckets}"
);
}
let grpc_timeout = Duration::from_secs(runtime_config.service.grpc_timeout_secs);
let grpc_max_concurrent: usize = runtime_config.service.grpc_max_concurrent;
let make_layer = || {
tower::ServiceBuilder::new()
.layer(crate::runtime::otel::TraceExtractLayer::new())
.layer(crate::runtime::credential_layer::CredentialResolveLayer::new())
.timeout(grpc_timeout)
.concurrency_limit(grpc_max_concurrent)
.into_inner()
};
let mut server = tonic::transport::Server::builder().layer(make_layer());
if let Some(tls) = tls_config_from_settings(&runtime_config.service.tls)? {
server = server.tls_config(tls)?;
}
let reflection_service = tonic_reflection::server::Builder::configure()
.register_encoded_file_descriptor_set(UDB_FILE_DESCRIPTOR_SET)
.build_v1()?;
let native_control_plane_enabled = native_registry::any_control_plane_enabled(&runtime_config);
let native_webrtc_peer_enabled = native_registry::any_webrtc_peer_enabled(&runtime_config);
if !native_control_plane_enabled && !native_webrtc_peer_enabled {
tracing::info!(
public_addr = %addr,
"native services disabled or no native listener selected; only the public DataBroker listener will start"
);
if let Some(pool) = service
.runtime
.load_full()
.native_store_pool_for_service("authn", true, "")
.ok()
{
crate::runtime::service::auth_service::install_data_plane_credential_resolvers(
pool,
&crate::runtime::authn::AuthnConfig::from_env(),
);
}
#[cfg(unix)]
let _uds_data_plane =
spawn_uds_data_plane(service.clone(), grpc_timeout, grpc_max_concurrent);
server
.add_service(reflection_service)
.add_service(health_service)
.add_service(DataBrokerServer::new(service))
.serve_with_shutdown(addr, shutdown_signal())
.await?;
return Ok(());
}
let (authn_service, authz_service, api_key_service) = service.build_auth_services();
let _canary_evaluator = authz_service.spawn_canary_evaluator();
if service.runtime_snapshot().pg_pool_clone().is_some() {
authz_service.warm_shared_snapshot().await;
let warmer = authz_service.clone();
let warm_interval = warmer.snapshot_ttl().max(std::time::Duration::from_secs(5));
tokio::spawn(async move {
let mut interval = tokio::time::interval(warm_interval);
interval.tick().await; loop {
interval.tick().await;
warmer.warm_shared_snapshot().await;
}
});
}
let runtime_snapshot = service.runtime_snapshot();
authn_service
.seed_signing_key_registry(runtime_snapshot.as_ref())
.await;
{
let readiness =
auth_service::readiness::check_auth_readiness(&SecurityConfig::current()).await;
for check in &readiness.checks {
if check.ok {
tracing::info!(check = %check.name, detail = %check.detail, "auth readiness ok");
} else {
tracing::error!(check = %check.name, detail = %check.detail, "auth readiness FAILED");
}
}
if !readiness.ok {
tracing::error!(
"auth-plane readiness checks failed; serving anyway — review the failed checks above"
);
let failed: Vec<serde_json::Value> = readiness
.checks
.iter()
.filter(|c| !c.ok)
.map(|c| serde_json::json!({ "check": c.name, "detail": c.detail }))
.collect();
let failed_names = readiness
.checks
.iter()
.filter(|c| !c.ok)
.map(|c| c.name.clone())
.collect::<Vec<_>>()
.join(",");
let operation_id = uuid::Uuid::new_v4().to_string();
authn_service
.emit_ops_event(
auth_service::events::AuthEvent::new(
auth_service::events::topics::OPS_READINESS_FAILURE,
operation_id.clone(),
String::new(),
serde_json::json!({
"operation_id": operation_id,
"failed_checks": failed,
}),
)
.with_correlation(operation_id.clone())
.with_compliance(
auth_service::events::ComplianceEnvelope {
actor: "udb.auth.readiness".to_string(),
target_resource: "auth-plane".to_string(),
operation: "readiness_check".to_string(),
outcome: "failure".to_string(),
reason_code: if failed_names.is_empty() {
"auth_readiness_failed".to_string()
} else {
format!("auth_readiness_failed:{failed_names}")
},
..auth_service::events::ComplianceEnvelope::default()
},
),
)
.await;
}
}
let identity_provider_service = service.build_identity_provider_service();
let scim_http_idp = std::sync::Arc::new(service.build_identity_provider_service());
let saml_http_idp = std::sync::Arc::new(service.build_identity_provider_service());
let control_plane_service = service.build_control_plane_service();
let tenant_service = service.build_tenant_service();
let notification_service = service.build_notification_service();
let analytics_service = service.build_analytics_service();
let lock_service = service.build_lock_service();
let scheduler_service = service.build_scheduler_service();
let vault_service = service.build_vault_service();
let cache_service = service.build_cache_service();
let webhook_service = service.build_webhook_service();
let backup_service = service.build_backup_service();
let search_service = service.build_search_service();
let config_service = service.build_config_service();
let metering_service = service.build_metering_service();
let livequery_service = service.build_livequery_service();
let workflow_service = service.build_workflow_service();
let embedding_service = service.build_embedding_service();
{
let vault_runtime = service.runtime.load_full();
if let Ok(vault_pool) = vault_runtime.native_store_pool_for_service("vault", true, "") {
let singleton_relation = vault_runtime.config().cdc.lock_log_relation();
let lease_pool = vault_pool.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_VAULT_LEASE_REAPER,
"vault DB credential leases revoked expired login roles",
lease_pool,
singleton_relation,
crate::runtime::service::vault_service::vault_db_lease_reaper_interval(),
move || {
let pool = vault_pool.clone();
async move {
crate::runtime::service::vault_service::run_vault_db_lease_reaper_once(
&pool,
crate::runtime::service::vault_service::VAULT_DB_LEASE_REAPER_BATCH,
)
.await
}
},
);
}
}
{
let metering_runtime = service.runtime.load_full();
if let Ok(metering_pool) =
metering_runtime.native_store_pool_for_service("metering", true, "")
{
let singleton_relation = metering_runtime.config().cdc.lock_log_relation();
let outbox_relation = metering_runtime.config().cdc.outbox_relation();
let journal_relation =
crate::runtime::system::SystemCatalogConfig::current().cdc_journal_relation();
let lease_pool = metering_pool.clone();
let metrics: Arc<dyn MetricsRecorder> = service.metrics.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_METERING_ROLLUP,
"metering rollup emitted closed usage windows",
lease_pool,
singleton_relation,
crate::runtime::service::metering_service::metering_rollup_interval(),
move || {
let pool = metering_pool.clone();
let outbox = outbox_relation.clone();
let journal = journal_relation.clone();
let metrics = metrics.clone();
async move {
crate::runtime::service::metering_service::run_metering_rollup_once(
&pool,
&outbox,
&journal,
crate::runtime::service::metering_service::METERING_ROLLUP_BATCH,
Some(&metrics),
)
.await
}
},
);
}
}
{
let analytics_runtime = service.runtime.load_full();
if let Ok(analytics_pool) =
analytics_runtime.native_store_pool_for_service("analytics", true, "")
{
let singleton_relation = analytics_runtime.config().cdc.lock_log_relation();
let lease_pool = analytics_pool.clone();
let rollup_interval = std::time::Duration::from_secs(
std::env::var("UDB_ANALYTICS_ROLLUP_INTERVAL_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|v| *v > 0)
.unwrap_or(300),
);
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_ANALYTICS_ROLLUP,
"analytics rollup refreshed percentile snapshots",
lease_pool,
singleton_relation,
rollup_interval,
move || {
let pool = analytics_pool.clone();
async move {
crate::runtime::service::analytics_service::run_analytics_rollup_once(&pool)
.await
.map(|rows| i64::try_from(rows).unwrap_or(i64::MAX))
}
},
);
}
}
{
let embedding_worker_service = std::sync::Arc::new(service.build_embedding_service());
let embedding_runtime = service.runtime.load_full();
if let Ok(embedding_pool) =
embedding_runtime.native_store_pool_for_service("embedding", true, "")
{
let singleton_relation = embedding_runtime.config().cdc.lock_log_relation();
let journal_relation =
crate::runtime::system::SystemCatalogConfig::current().cdc_journal_relation();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_EMBEDDING_WORK_EMITTER,
"embedding work emitted source changes",
embedding_pool,
singleton_relation,
crate::runtime::service::embedding_service::embedding_work_emitter_interval(),
move || {
let service = embedding_worker_service.clone();
let journal = journal_relation.clone();
async move {
crate::runtime::service::embedding_service::run_embedding_work_emitter_once(
service,
&journal,
crate::runtime::service::embedding_service::EMBEDDING_WORK_EMITTER_BATCH,
)
.await
}
},
);
}
}
{
let search_worker_service = std::sync::Arc::new(service.build_search_service());
let search_runtime = service.runtime.load_full();
if let Ok(search_pool) = search_runtime.native_store_pool_for_service("search", true, "") {
let singleton_relation = search_runtime.config().cdc.lock_log_relation();
let journal_relation =
crate::runtime::system::SystemCatalogConfig::current().cdc_journal_relation();
let freshness_service = search_worker_service.clone();
let freshness_journal = journal_relation.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_SEARCH_FRESHNESS,
"search freshness applied source changes",
search_pool.clone(),
singleton_relation.clone(),
crate::runtime::service::search_service::search_freshness_interval(),
move || {
let service = freshness_service.clone();
let journal = freshness_journal.clone();
async move {
crate::runtime::service::search_service::run_index_freshness_consumer(
service,
&journal,
crate::runtime::service::search_service::SEARCH_FRESHNESS_BATCH,
)
.await
}
},
);
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_SEARCH_REINDEX,
"search reindex rebuilt requested indexes",
search_pool,
singleton_relation,
crate::runtime::service::search_service::search_reindex_interval(),
move || {
let service = search_worker_service.clone();
let journal = journal_relation.clone();
async move {
crate::runtime::service::search_service::run_search_reindex_once(
service,
&journal,
crate::runtime::service::search_service::SEARCH_REINDEX_BATCH,
)
.await
}
},
);
}
}
{
let scheduler_runtime = service.runtime.load_full();
if let Ok(scheduler_pool) =
scheduler_runtime.native_store_pool_for_service("scheduler", true, "")
{
let singleton_relation = scheduler_runtime.config().cdc.lock_log_relation();
let outbox_relation = scheduler_runtime.config().cdc.outbox_relation();
let lease_pool = scheduler_pool.clone();
let tick_interval = std::time::Duration::from_secs(
std::env::var("UDB_SCHEDULER_TICK_INTERVAL_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|v| *v > 0)
.unwrap_or(30),
);
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_SCHEDULER_TICK,
"scheduler tick fired due jobs",
lease_pool,
singleton_relation,
tick_interval,
move || {
let pool = scheduler_pool.clone();
let outbox = outbox_relation.clone();
async move {
crate::runtime::service::scheduler_service::run_scheduler_tick_once(
&pool,
Some(&outbox),
crate::runtime::service::scheduler_service::SCHEDULER_TICK_BATCH,
)
.await
}
},
);
}
}
{
let lock_runtime = service.runtime.load_full();
if let Ok(lock_pool) = lock_runtime.native_store_pool_for_service("lock", true, "") {
let singleton_relation = lock_runtime.config().cdc.lock_log_relation();
let outbox_relation = lock_runtime.config().cdc.outbox_relation();
let lease_pool = lock_pool.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_LOCK_EXPIRY_REAPER,
"lock expiry reaper flipped lapsed leases",
lease_pool,
singleton_relation,
crate::runtime::service::lock_service::lock_expiry_interval(),
move || {
let pool = lock_pool.clone();
let outbox = outbox_relation.clone();
async move {
crate::runtime::service::lock_service::run_lock_expiry_once(
&pool,
Some(&outbox),
crate::runtime::service::lock_service::LOCK_EXPIRY_SWEEP_BATCH,
)
.await
}
},
);
}
}
{
let workflow_runtime = service.runtime.load_full();
if let Ok(workflow_pool) =
workflow_runtime.native_store_pool_for_service("workflow", true, "")
{
let singleton_relation = workflow_runtime.config().cdc.lock_log_relation();
let outbox_relation = workflow_runtime.config().cdc.outbox_relation();
let stores = workflow_runtime.default_system_stores();
let lease_pool = workflow_pool.clone();
let tick_interval = std::time::Duration::from_secs(
std::env::var("UDB_WORKFLOW_TICK_INTERVAL_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|v| *v > 0)
.unwrap_or(30),
);
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_WORKFLOW_TICK,
"workflow tick advanced due instances",
lease_pool,
singleton_relation,
tick_interval,
move || {
let pool = workflow_pool.clone();
let outbox = outbox_relation.clone();
let stores = stores.clone();
async move {
crate::runtime::service::workflow_service::run_workflow_tick_once(
&pool,
Some(&outbox),
stores,
crate::runtime::service::workflow_service::WORKFLOW_TICK_BATCH,
)
.await
}
},
);
}
}
#[cfg(feature = "http-client")]
{
let webhook_runtime = service.runtime.load_full();
if let Ok(webhook_pool) = webhook_runtime.native_store_pool_for_service("webhook", true, "")
{
let singleton_relation = webhook_runtime.config().cdc.lock_log_relation();
let outbox_relation = webhook_runtime.config().cdc.outbox_relation();
let journal_relation =
crate::runtime::system::SystemCatalogConfig::current().cdc_journal_relation();
let lease_pool = webhook_pool.clone();
let http = reqwest::Client::new();
let metrics: Arc<dyn MetricsRecorder> = service.metrics.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_WEBHOOK_DELIVERY,
"webhook delivery posted events",
lease_pool,
singleton_relation,
crate::runtime::service::webhook_service::webhook_delivery_interval(),
move || {
let pool = webhook_pool.clone();
let http = http.clone();
let outbox = outbox_relation.clone();
let journal = journal_relation.clone();
let metrics = metrics.clone();
async move {
crate::runtime::service::webhook_service::run_webhook_delivery_worker_once(
&http,
&pool,
Some(&outbox),
&journal,
crate::runtime::service::webhook_service::WEBHOOK_DELIVERY_BATCH,
Some(&metrics),
)
.await
}
},
);
}
}
#[cfg(feature = "redis")]
{
let cache_runtime = service.runtime.load_full();
if let (Some(cache_redis), Ok(cache_pool)) = (
cache_runtime.redis_clone(),
cache_runtime.native_store_pool_for_service("cache", true, ""),
) {
let singleton_relation = cache_runtime.config().cdc.lock_log_relation();
let outbox_relation = cache_runtime.config().cdc.outbox_relation();
let journal_relation =
crate::runtime::system::SystemCatalogConfig::current().cdc_journal_relation();
let lease_pool = cache_pool.clone();
let metrics: Arc<dyn MetricsRecorder> = service.metrics.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_CACHE_INVALIDATOR,
"cache invalidation swept namespaces",
lease_pool,
singleton_relation,
crate::runtime::service::cache_service::cache_invalidation_interval(),
move || {
let redis = cache_redis.clone();
let pool = cache_pool.clone();
let outbox = outbox_relation.clone();
let journal = journal_relation.clone();
let metrics = metrics.clone();
async move {
crate::runtime::service::cache_service::run_cache_invalidation_worker_once(
&redis,
&pool,
&outbox,
&journal,
crate::runtime::service::cache_service::CACHE_INVALIDATION_BATCH,
Some(&metrics),
)
.await
}
},
);
}
}
#[cfg(feature = "http-client")]
{
let notification_runtime = service.runtime.load_full();
if let Ok(notification_pool) =
notification_runtime.native_store_pool_for_service("notification", true, "")
{
let singleton_relation = notification_runtime.config().cdc.lock_log_relation();
let outbox_relation = notification_runtime.config().cdc.outbox_relation();
let http = reqwest::Client::new();
let lease_pool = notification_pool.clone();
let metrics: Arc<dyn MetricsRecorder> = service.metrics.clone();
crate::runtime::service::native_runtime::NativeWorkerHost::spawn_while_leader(
crate::runtime::singleton::WORKER_NOTIFICATION_DELIVERY,
"notification delivery processed queued intents",
lease_pool,
singleton_relation,
crate::runtime::service::notification_service::notification_delivery_interval(),
move || {
let runtime = notification_runtime.clone();
let pool = notification_pool.clone();
let outbox = outbox_relation.clone();
let http = http.clone();
let metrics = metrics.clone();
async move {
crate::runtime::service::notification_service::run_notification_delivery_worker_once(
&http,
runtime,
&pool,
Some(&outbox),
crate::runtime::service::notification_service::NOTIFICATION_DELIVERY_BATCH,
Some(&metrics),
)
.await
}
},
);
}
}
{
let evidence_runtime = service.runtime.load_full();
if let Ok(evidence_pool) =
evidence_runtime.native_store_pool_for_service("compliance", true, "")
{
let singleton_relation = evidence_runtime.config().cdc.lock_log_relation();
let config = crate::runtime::evidence_export::EvidenceExportConfig::from_env();
crate::runtime::evidence_export::spawn_evidence_export_worker(
evidence_runtime,
evidence_pool,
singleton_relation,
config,
);
}
}
let storage_service = service.build_storage_service();
let asset_service = service.build_asset_service();
let webrtc_service = service.build_webrtc_service();
#[cfg(feature = "kafka")]
if crate::runtime::cdc::cdc_delivery_enabled() {
if let Some(brokers) = runtime_config.kafka_brokers.clone() {
std::sync::Arc::new(service.build_asset_service())
.spawn_storage_finalized_consumer(brokers);
}
} else {
tracing::info!("storage→asset auto-trigger consumer disabled: UDB_CDC_ENABLED=false");
}
#[cfg(feature = "kafka")]
if crate::runtime::cdc::cdc_delivery_enabled() {
if let Some(brokers) = runtime_config.kafka_brokers.clone() {
let trigger_runtime = service.runtime.load_full();
if let Ok(trigger_pool) =
trigger_runtime.native_store_pool_for_service("asset", true, "")
{
let singleton_relation = trigger_runtime.config().cdc.lock_log_relation();
std::sync::Arc::new(service.build_asset_service()).spawn_trigger_manager(
brokers,
trigger_pool,
singleton_relation,
);
}
}
}
let auth_addr: SocketAddr = match runtime_config.native_services.control_plane_addr.as_str() {
raw if !raw.trim().is_empty() => raw
.trim()
.parse()
.map_err(|err| format!("invalid UDB_AUTH_GRPC_ADDR '{raw}': {err}"))?,
_ => SocketAddr::new(
std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
addr.port().wrapping_add(10),
),
};
tracing::info!(
%auth_addr,
public_addr = %addr,
"native auth control plane (Authn/Authz/ApiKey) bound to an internal \
listener, isolated from the public DataBroker port; set \
UDB_AUTH_GRPC_ADDR to expose it on a trusted interface"
);
let webrtc_addr: SocketAddr = match runtime_config.native_services.webrtc_peer_addr.as_str() {
raw if !raw.trim().is_empty() => raw
.trim()
.parse()
.map_err(|err| format!("invalid UDB_WEBRTC_GRPC_ADDR '{raw}': {err}"))?,
_ => SocketAddr::new(
std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
addr.port().wrapping_add(20),
),
};
tracing::info!(
%webrtc_addr,
public_addr = %addr,
"WebRTC peer listener bound with peer-token auth; set \
UDB_WEBRTC_GRPC_ADDR to expose it on a trusted interface"
);
let mut auth_server = tonic::transport::Server::builder().layer(make_layer());
if let Some(tls) = tls_config_from_settings(&runtime_config.service.tls)? {
auth_server = auth_server.tls_config(tls)?;
}
let msec = method_security::MethodSecurityLayer::new().with_metrics(metrics.clone());
let native_health_runtime = service.runtime_snapshot();
let native_health = handlers_meta::build_listener_health_service(
handlers_meta::HealthPlane::NativeControlPlane,
&runtime_config,
Some(native_health_runtime.as_ref()),
)
.await;
let (listener_shutdown_tx, listener_shutdown_rx) = tokio::sync::watch::channel(false);
tokio::spawn(async move {
shutdown_signal().await;
let _ = listener_shutdown_tx.send(true);
});
let auth_fut = auth_server
.add_service(native_health)
.add_service(msec.wrap(auth_service::AuthnServiceServer::new(authn_service)))
.add_service(msec.wrap(auth_service::AuthzServiceServer::new(authz_service)))
.add_service(msec.wrap(auth_service::ApiKeyServiceServer::new(api_key_service)))
.add_service(msec.wrap(auth_service::IdentityProviderServiceServer::new(
identity_provider_service,
)))
.add_service(msec.wrap(auth_service::ControlPlaneServiceServer::new(
control_plane_service,
)))
.add_service(msec.wrap(tenant_service::TenantServiceServer::new(tenant_service)))
.add_service(msec.wrap(lock_service::LockServiceServer::new(lock_service)))
.add_service(msec.wrap(scheduler_service::SchedulerServiceServer::new(
scheduler_service,
)))
.add_service(msec.wrap(vault_service::VaultServiceServer::new(vault_service)))
.add_service(msec.wrap(cache_service::CacheServiceServer::new(cache_service)))
.add_service(msec.wrap(webhook_service::WebhookServiceServer::new(webhook_service)))
.add_service(msec.wrap(backup_service::BackupServiceServer::new(backup_service)))
.add_service(msec.wrap(search_service::SearchServiceServer::new(search_service)))
.add_service(msec.wrap(config_service::ConfigServiceServer::new(config_service)))
.add_service(msec.wrap(metering_service::MeteringServiceServer::new(
metering_service,
)))
.add_service(msec.wrap(livequery_service::LiveQueryServiceServer::new(
livequery_service,
)))
.add_service(msec.wrap(workflow_service::WorkflowServiceServer::new(
workflow_service,
)))
.add_service(msec.wrap(embedding_service::EmbeddingServiceServer::new(
embedding_service,
)))
.add_service(
msec.wrap(notification_service::NotificationServiceServer::new(
notification_service,
)),
)
.add_service(msec.wrap(analytics_service::AnalyticsServiceServer::new(
analytics_service,
)))
.add_service(msec.wrap(storage_service::StorageServiceServer::new(storage_service)))
.add_service(msec.wrap(asset_service::AssetServiceServer::new(asset_service)))
.add_service(msec.wrap(webrtc_service::RoomServiceServer::new(
webrtc_service.clone(),
)))
.add_service(msec.wrap(webrtc_service::PeerServiceServer::new(
webrtc_service.clone(),
)))
.add_service(msec.wrap(webrtc_service::TrackServiceServer::new(
webrtc_service.clone(),
)))
.add_service(msec.wrap(webrtc_service::TurnServiceServer::new(
webrtc_service.clone(),
)))
.add_service(msec.wrap(webrtc_service::SignalingServiceServer::new(
webrtc_service.clone(),
)))
.serve_with_shutdown(
auth_addr,
listener_shutdown_signal(listener_shutdown_rx.clone()),
);
let mut webrtc_peer_server = tonic::transport::Server::builder().layer(make_layer());
if let Some(tls) = tls_config_from_settings(&runtime_config.service.tls)? {
webrtc_peer_server = webrtc_peer_server.tls_config(tls)?;
}
let peer_auth = WebrtcPeerTokenAuth::new();
let webrtc_health_runtime = service.runtime_snapshot();
let webrtc_peer_health = handlers_meta::build_listener_health_service(
handlers_meta::HealthPlane::WebRtcPeer,
&runtime_config,
Some(webrtc_health_runtime.as_ref()),
)
.await;
let webrtc_peer_fut = webrtc_peer_server
.add_service(webrtc_peer_health)
.add_service(
msec.wrap(webrtc_service::PeerServiceServer::with_interceptor(
webrtc_service.clone(),
peer_auth.clone(),
)),
)
.add_service(
msec.wrap(webrtc_service::TrackServiceServer::with_interceptor(
webrtc_service.clone(),
peer_auth.clone(),
)),
)
.add_service(
msec.wrap(webrtc_service::TurnServiceServer::with_interceptor(
webrtc_service.clone(),
peer_auth.clone(),
)),
)
.add_service(
msec.wrap(webrtc_service::SignalingServiceServer::with_interceptor(
webrtc_service,
peer_auth,
)),
)
.serve_with_shutdown(
webrtc_addr,
listener_shutdown_signal(listener_shutdown_rx.clone()),
);
#[cfg(unix)]
let _uds_data_plane = spawn_uds_data_plane(service.clone(), grpc_timeout, grpc_max_concurrent);
let main_fut = server
.add_service(reflection_service)
.add_service(health_service)
.add_service(DataBrokerServer::new(service))
.serve_with_shutdown(addr, listener_shutdown_signal(listener_shutdown_rx.clone()));
#[cfg(feature = "ws-signalling")]
let _ws_signalling = crate::runtime::signalling::SignalingServer::spawn_from_env_with_shutdown(
shutdown_signal(),
);
let _scim_http = auth_service::spawn_scim_http_from_env(scim_http_idp, shutdown_signal());
let _saml_http = auth_service::spawn_saml_http_from_env(saml_http_idp, shutdown_signal());
match (native_control_plane_enabled, native_webrtc_peer_enabled) {
(true, true) => {
tokio::pin!(main_fut);
tokio::pin!(auth_fut);
tokio::pin!(webrtc_peer_fut);
tokio::select! {
result = &mut main_fut => listener_finished("data-plane", result, *listener_shutdown_rx.borrow())?,
result = &mut auth_fut => listener_finished("native-control-plane", result, *listener_shutdown_rx.borrow())?,
result = &mut webrtc_peer_fut => listener_finished("webrtc-peer", result, *listener_shutdown_rx.borrow())?,
}
}
(true, false) => {
tokio::pin!(main_fut);
tokio::pin!(auth_fut);
tokio::select! {
result = &mut main_fut => listener_finished("data-plane", result, *listener_shutdown_rx.borrow())?,
result = &mut auth_fut => listener_finished("native-control-plane", result, *listener_shutdown_rx.borrow())?,
}
}
(false, true) => {
tokio::pin!(main_fut);
tokio::pin!(webrtc_peer_fut);
tokio::select! {
result = &mut main_fut => listener_finished("data-plane", result, *listener_shutdown_rx.borrow())?,
result = &mut webrtc_peer_fut => listener_finished("webrtc-peer", result, *listener_shutdown_rx.borrow())?,
}
}
(false, false) => unreachable!("handled before native service construction"),
}
Ok(())
}
async fn listener_shutdown_signal(mut shutdown: tokio::sync::watch::Receiver<bool>) {
if *shutdown.borrow() {
return;
}
let _ = shutdown.changed().await;
}
fn listener_finished(
listener: &'static str,
result: Result<(), tonic::transport::Error>,
shutdown_requested: bool,
) -> Result<(), Box<dyn std::error::Error>> {
match result {
Ok(()) if shutdown_requested => {
tracing::info!(listener, "gRPC listener exited after shutdown signal");
Ok(())
}
Ok(()) => {
let err =
std::io::Error::other(format!("{listener} gRPC listener exited unexpectedly"));
tracing::error!(
listener,
error = %err,
"gRPC listener exited without shutdown; shutting down sibling listeners"
);
Err(err.into())
}
Err(err) => {
tracing::error!(
listener,
error = %err,
"gRPC listener failed; shutting down sibling listeners"
);
Err(Box::new(err))
}
}
}
#[cfg(feature = "kafka")]
async fn start_cdc_engine(
runtime: &DataBrokerRuntime,
metrics: Arc<dyn MetricsRecorder>,
) -> Option<Arc<CdcEngine>> {
if !crate::runtime::cdc::cdc_delivery_enabled() {
tracing::info!("CDC tailer disabled: UDB_CDC_ENABLED=false");
return None;
}
let Some(kafka_brokers) = runtime.config().kafka_brokers.clone() else {
tracing::info!("CDC tailer disabled: kafka_brokers is not configured");
return None;
};
let Some(pg_pool) = runtime.pg_pool_clone() else {
tracing::warn!("CDC tailer disabled: PostgreSQL pool is not configured");
return None;
};
let pg_dsn = if !runtime.config().primary.direct_dsn.trim().is_empty() {
runtime.config().primary.direct_dsn.trim().to_string()
} else {
runtime.config().primary.pooler_dsn.trim().to_string()
};
if pg_dsn.is_empty() {
tracing::warn!("CDC tailer disabled: primary PostgreSQL DSN is not configured");
return None;
};
let singleton_pool = pg_pool.clone();
let singleton_relation = runtime.config().cdc.lock_log_relation();
#[cfg(feature = "redis")]
let engine = CdcEngine::new(
pg_pool,
runtime.redis_clone(),
&kafka_brokers,
pg_dsn,
metrics,
runtime.config().cdc.clone(),
);
#[cfg(not(feature = "redis"))]
let engine = CdcEngine::new(
pg_pool,
&kafka_brokers,
pg_dsn,
metrics,
runtime.config().cdc.clone(),
);
match engine {
Ok(mut engine) => {
if let Err(err) = engine.load_topic_policies().await {
tracing::warn!("CDC topic policy load failed: {err}");
}
let engine = Arc::new(engine);
if let Err(err) = engine.run_indoubt_recovery_on_startup().await {
tracing::warn!("CDC in-doubt recovery on startup failed: {err}");
}
tokio::spawn({
let engine = engine.clone();
async move {
engine.run_advisory_lock_loop().await;
}
});
if let Ok(dsn) = std::env::var("UDB_CDC_POSTGRES_SOURCE_DSN") {
if !dsn.trim().is_empty() {
let relation = std::env::var("UDB_CDC_POSTGRES_SOURCE_TABLE")
.unwrap_or_else(|_| "udb_system.udb_cdc_outbox".to_string());
let source: std::sync::Arc<dyn crate::runtime::cdc::CdcSource> =
std::sync::Arc::new(crate::runtime::cdc::source::PostgresCdcSource {
dsn,
publication: relation,
slot: "udb-postgres-source".to_string(),
});
let engine = engine.clone();
let singleton_pool = singleton_pool.clone();
let singleton_relation = singleton_relation.clone();
tokio::spawn(async move {
loop {
let engine = engine.clone();
let source = source.clone();
match crate::runtime::singleton::run_while_leader(
&singleton_pool,
&singleton_relation,
crate::runtime::singleton::WORKER_CDC_POSTGRES_SOURCE,
crate::runtime::singleton::WORKER_SINGLETON_LEASE_TTL,
|| async move { engine.tail_source(source).await },
)
.await
{
Ok(Some(Ok(()))) => {}
Ok(Some(Err(err))) => {
tracing::warn!("CDC Postgres source tailer exited: {err}")
}
Ok(None) => tracing::debug!(
"CDC Postgres source tailer idle: singleton lease held by peer"
),
Err(err) => {
tracing::warn!("CDC Postgres source tailer lease failed: {err}")
}
}
tokio::time::sleep(
crate::runtime::singleton::WORKER_SINGLETON_RETRY_SLEEP,
)
.await;
}
});
tracing::info!("CDC Postgres table source tailer started");
}
}
#[cfg(feature = "mysql")]
if let Ok(dsn) = std::env::var("UDB_CDC_MYSQL_SOURCE_DSN") {
if !dsn.trim().is_empty() {
let server_id = std::env::var("UDB_CDC_MYSQL_SOURCE_SERVER_ID")
.ok()
.and_then(|v| v.parse::<u32>().ok())
.unwrap_or(5401);
let source: std::sync::Arc<dyn crate::runtime::cdc::CdcSource> =
std::sync::Arc::new(crate::runtime::cdc::source::MysqlBinlogSource {
dsn,
server_id,
});
let engine = engine.clone();
let singleton_pool = singleton_pool.clone();
let singleton_relation = singleton_relation.clone();
tokio::spawn(async move {
loop {
let engine = engine.clone();
let source = source.clone();
match crate::runtime::singleton::run_while_leader(
&singleton_pool,
&singleton_relation,
crate::runtime::singleton::WORKER_CDC_MYSQL_SOURCE,
crate::runtime::singleton::WORKER_SINGLETON_LEASE_TTL,
|| async move { engine.tail_source(source).await },
)
.await
{
Ok(Some(Ok(()))) => {}
Ok(Some(Err(err))) => {
tracing::warn!("CDC MySQL source tailer exited: {err}")
}
Ok(None) => tracing::debug!(
"CDC MySQL source tailer idle: singleton lease held by peer"
),
Err(err) => {
tracing::warn!("CDC MySQL source tailer lease failed: {err}")
}
}
tokio::time::sleep(
crate::runtime::singleton::WORKER_SINGLETON_RETRY_SLEEP,
)
.await;
}
});
tracing::info!("CDC MySQL table source tailer started");
}
}
#[cfg(feature = "mongodb-native")]
if let Ok(uri) = std::env::var("UDB_CDC_MONGO_SOURCE_URI") {
let database = std::env::var("UDB_CDC_MONGO_SOURCE_DB").unwrap_or_default();
if !uri.trim().is_empty() && !database.trim().is_empty() {
let collection = std::env::var("UDB_CDC_MONGO_SOURCE_COLLECTION")
.ok()
.filter(|v| !v.trim().is_empty());
let source: std::sync::Arc<dyn crate::runtime::cdc::CdcSource> =
std::sync::Arc::new(crate::runtime::cdc::source::MongoCdcSource {
uri,
database,
collection,
});
let engine = engine.clone();
let singleton_pool = singleton_pool.clone();
let singleton_relation = singleton_relation.clone();
tokio::spawn(async move {
loop {
let engine = engine.clone();
let source = source.clone();
match crate::runtime::singleton::run_while_leader(
&singleton_pool,
&singleton_relation,
crate::runtime::singleton::WORKER_CDC_MONGODB_SOURCE,
crate::runtime::singleton::WORKER_SINGLETON_LEASE_TTL,
|| async move { engine.tail_source(source).await },
)
.await
{
Ok(Some(Ok(()))) => {}
Ok(Some(Err(err))) => {
tracing::warn!("CDC MongoDB source tailer exited: {err}")
}
Ok(None) => tracing::debug!(
"CDC MongoDB source tailer idle: singleton lease held by peer"
),
Err(err) => {
tracing::warn!("CDC MongoDB source tailer lease failed: {err}")
}
}
tokio::time::sleep(
crate::runtime::singleton::WORKER_SINGLETON_RETRY_SLEEP,
)
.await;
}
});
tracing::info!("CDC MongoDB change-stream source tailer started");
}
}
Some(engine)
}
Err(err) => {
tracing::warn!("CDC tailer disabled: Kafka producer initialization failed: {err}");
None
}
}
}
async fn metrics_http_server(
metrics: Arc<PrometheusMetrics>,
runtime: DataBrokerRuntime,
addr: SocketAddr,
allowed_cidr: Option<String>,
) {
let listener = match tokio::net::TcpListener::bind(addr).await {
Ok(listener) => listener,
Err(err) => {
tracing::warn!("metrics endpoint disabled: {err}");
return;
}
};
let allowed_cidr: Option<String> = allowed_cidr
.map(|cidr| cidr.trim().to_string())
.filter(|cidr| !cidr.is_empty());
loop {
let Ok((mut socket, peer)) = listener.accept().await else {
continue;
};
if let Some(allow) = &allowed_cidr
&& !ip_matches_allow_entry(peer.ip(), allow)
{
tracing::debug!(peer = %peer.ip(), "metrics scrape rejected: not in allowed CIDR");
continue;
}
refresh_native_service_degraded_metrics(&runtime, metrics.as_ref());
let text = metrics.gather_text(&format!(
"{}{}",
runtime.cache_metrics_text(),
format!(
"{}{}",
runtime.encryption_metrics_text(),
runtime.pg_pool_metrics_text()
)
));
let ready_runtime = runtime.clone();
let ready_metrics = metrics.clone();
tokio::spawn(async move {
let mut buf = [0u8; 256];
let n = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
.await
.ok()
.and_then(|r| r.ok())
.unwrap_or(0);
let request = String::from_utf8_lossy(&buf[..n]);
let path = request
.lines()
.next()
.and_then(|line| line.split_whitespace().nth(1))
.unwrap_or("/");
let response = if path == "/healthz" {
let body = "ok\n";
format!(
"HTTP/1.1 200 OK\r\ncontent-type: text/plain; charset=utf-8\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}",
body.len(),
body
)
} else if path == "/readyz" {
let (status, body) =
metrics_readiness_response(&ready_runtime, ready_metrics.as_ref()).await;
format!(
"HTTP/1.1 {status}\r\ncontent-type: text/plain; charset=utf-8\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}",
body.len(),
body
)
} else {
format!(
"HTTP/1.1 200 OK\r\ncontent-type: text/plain; version=0.0.4; charset=utf-8\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}",
text.len(),
text
)
};
let _ = socket.write_all(response.as_bytes()).await;
});
}
}
async fn metrics_readiness_response(
runtime: &DataBrokerRuntime,
metrics: &dyn MetricsRecorder,
) -> (&'static str, String) {
let native_statuses =
crate::runtime::service::native_registry::resolved_native_service_statuses(
runtime.config(),
);
refresh_native_service_degraded_metrics(runtime, metrics);
let auth_triples = auth_readiness_triples(&SecurityConfig::current()).await;
let readiness = crate::runtime::slo::build_readiness_facts(
runtime.init_report(),
&native_statuses,
&auth_triples,
);
if readiness.passed() {
return ("200 OK", "ready\n".to_string());
}
let mut body = String::from("not ready\n");
for err in readiness.errors() {
body.push_str("error: ");
body.push_str(&err);
body.push('\n');
}
for warn in readiness.warnings() {
body.push_str("warning: ");
body.push_str(&warn);
body.push('\n');
}
("503 Service Unavailable", body)
}
fn refresh_native_service_degraded_metrics(
runtime: &DataBrokerRuntime,
metrics: &dyn MetricsRecorder,
) {
for status in
crate::runtime::service::native_registry::resolved_native_service_statuses(runtime.config())
{
metrics.set_native_service_degraded(
&status.service_id,
status.degraded || (status.enabled && !status.mounted),
);
}
}
async fn cdc_metrics_poller(runtime: DataBrokerRuntime, metrics: Arc<dyn MetricsRecorder>) {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(5));
loop {
interval.tick().await;
if let Ok((lag, depth)) = runtime.cdc_outbox_metrics().await {
metrics.set_cdc_lag_seconds(lag);
metrics.set_cdc_outbox_depth(depth);
}
}
}
fn validate_secure_transport(
service: &crate::runtime::config::ServiceSettings,
) -> Result<(), String> {
if service.require_secure_transport && !service.tls.has_server_identity() {
return Err(
"UDB_REQUIRE_SECURE_TRANSPORT/UDB_TLS_REQUIRED is enabled but TLS cert/key is not configured"
.to_string(),
);
}
let any_mtls_required = service.mtls_required
|| service.broker_to_broker_mtls_required
|| service.internal_control_mtls_required;
if any_mtls_required {
if !service.tls.has_server_identity() {
return Err("mTLS is enabled but TLS cert/key is not configured".to_string());
}
if !service.tls.has_client_ca() {
return Err(
"mTLS is enabled but UDB_MTLS_CLIENT_CA_PEM/PATH is not configured".to_string(),
);
}
}
Ok(())
}
fn tls_config_from_settings(
settings: &crate::runtime::config::TlsSettings,
) -> Result<Option<ServerTlsConfig>, Box<dyn std::error::Error>> {
let Some(cert) = config_bytes(settings.cert_pem.as_deref(), settings.cert_path.as_deref())?
else {
return Ok(None);
};
let Some(key) = config_bytes(settings.key_pem.as_deref(), settings.key_path.as_deref())? else {
return Ok(None);
};
let mut tls = ServerTlsConfig::new().identity(Identity::from_pem(cert, key));
if let Some(ca) = config_bytes(
settings.client_ca_pem.as_deref(),
settings.client_ca_path.as_deref(),
)? {
tls = tls.client_ca_root(Certificate::from_pem(ca));
}
Ok(Some(tls))
}
fn config_bytes(
pem: Option<&str>,
path: Option<&str>,
) -> Result<Option<Vec<u8>>, Box<dyn std::error::Error>> {
if let Some(value) = pem.filter(|value| !value.trim().is_empty()) {
return Ok(Some(value.as_bytes().to_vec()));
}
if let Some(value) = path.filter(|value| !value.trim().is_empty()) {
return Ok(Some(fs::read(value)?));
}
Ok(None)
}
fn proto_cdc_envelope(envelope: crate::cdc::CdcEnvelope) -> CdcEnvelope {
CdcEnvelope {
event_id: envelope.event_id,
topic: envelope.topic,
partition_key: envelope.partition_key,
payload_json: envelope.payload_json,
published_at: Some(prost_types::Timestamp {
seconds: envelope.published_at.timestamp(),
nanos: envelope.published_at.timestamp_subsec_nanos() as i32,
}),
}
}
type ResponseStream<T> = Pin<Box<dyn Stream<Item = Result<T, Status>> + Send + 'static>>;
async fn shutdown_signal() {
#[cfg(unix)]
{
let mut terminate =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.unwrap_or_else(|err| {
tracing::warn!(
"failed to install SIGTERM handler: {}; only SIGINT (ctrl-c) will trigger graceful shutdown",
err
);
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::hangup())
.expect("install SIGHUP fallback handler")
});
tokio::select! {
_ = tokio::signal::ctrl_c() => {}
_ = terminate.recv() => {}
}
}
#[cfg(not(unix))]
{
let _ = tokio::signal::ctrl_c().await;
}
}
fn spawn_config_reload_watcher(service: DataBrokerService) {
#[cfg(unix)]
tokio::spawn(async move {
let mut hangup = match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::hangup())
{
Ok(signal) => signal,
Err(err) => {
tracing::warn!("failed to install SIGHUP reload handler: {err}");
return;
}
};
while hangup.recv().await.is_some() {
let report = service.reload_runtime_from_env("sighup").await;
if report.applied {
tracing::info!(
reload_id = report.reload_id,
previous_generation = report.previous_generation,
new_generation = report.new_generation,
previous_client_count = report.previous_client_count,
active_client_count = report.active_client_count,
changed_instances = ?report.changed_instances,
"UDB runtime config reload applied"
);
} else {
tracing::warn!(
reload_id = report.reload_id,
accepted = report.accepted,
rolled_back = report.rolled_back,
validation_errors = ?report.validation_errors,
failed_health_checks = ?report.failed_health_checks,
warnings = ?report.warnings,
"UDB runtime config reload rejected"
);
}
}
});
#[cfg(not(unix))]
{
let _ = service;
}
}
fn saga_record_to_proto(r: crate::runtime::saga::SagaAdminRecord) -> SagaRecord {
SagaRecord {
saga_id: r.saga_id,
tx_id: r.tx_id,
tenant_id: r.tenant_id,
correlation_id: r.correlation_id,
status: r.status,
current_step: r.current_step,
steps_json: r.steps_json.to_string().into_bytes(),
compensations_json: r.compensations_json.to_string().into_bytes(),
last_error: r.last_error,
created_at_unix: r.created_at.timestamp(),
updated_at_unix: r.updated_at.timestamp(),
}
}
fn insert_ascii_header(
metadata: &mut tonic::metadata::MetadataMap,
key: &'static str,
value: &str,
) {
if value.trim().is_empty() {
return;
}
if let Ok(parsed) = value.parse() {
metadata.insert(key, parsed);
}
}
mod handlers_admin;
mod handlers_catalog;
pub(crate) mod handlers_data;
mod handlers_meta;
mod handlers_object;
mod handlers_policy;
mod handlers_resource;
mod handlers_stores;
mod handlers_tx;
mod handlers_vector;
#[tonic::async_trait]
impl DataBroker for DataBrokerService {
type BatchSelectStream = ResponseStream<RecordSet>;
type SelectV2Stream = ResponseStream<crate::proto::RecordBatchV2>;
type BatchUpsertStream = ResponseStream<MutationResponse>;
type VectorBatchUpsertStream = ResponseStream<MutationResponse>;
type GetObjectStream = ResponseStream<Chunk>;
type BeginTxStream = ResponseStream<TxStatus>;
type PublishCDCStream = ResponseStream<CdcEnvelope>;
async fn select(&self, request: Request<SelectRequest>) -> Result<Response<RecordSet>, Status> {
self.select_inner(request).await
}
async fn select_v2(
&self,
request: Request<SelectRequest>,
) -> Result<Response<Self::SelectV2Stream>, Status> {
self.select_v2_inner(request).await
}
async fn batch_select(
&self,
request: Request<tonic::Streaming<SelectRequest>>,
) -> Result<Response<Self::BatchSelectStream>, Status> {
self.batch_select_inner(request).await
}
async fn upsert(
&self,
request: Request<UpsertRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.upsert_inner(request).await
}
async fn batch_upsert(
&self,
request: Request<tonic::Streaming<UpsertRequest>>,
) -> Result<Response<Self::BatchUpsertStream>, Status> {
self.batch_upsert_inner(request).await
}
async fn vector_search(
&self,
request: Request<VectorSearchRequest>,
) -> Result<Response<VectorSet>, Status> {
self.vector_search_inner(request).await
}
async fn vector_hybrid_search(
&self,
request: Request<VectorHybridSearchRequest>,
) -> Result<Response<VectorSet>, Status> {
self.vector_hybrid_search_inner(request).await
}
async fn vector_upsert(
&self,
request: Request<VectorUpsertRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.vector_upsert_inner(request).await
}
async fn vector_batch_upsert(
&self,
request: Request<tonic::Streaming<VectorUpsertRequest>>,
) -> Result<Response<Self::VectorBatchUpsertStream>, Status> {
self.vector_batch_upsert_inner(request).await
}
async fn put_object(
&self,
request: Request<tonic::Streaming<Chunk>>,
) -> Result<Response<MutationResponse>, Status> {
self.put_object_inner(request).await
}
async fn get_object(
&self,
request: Request<crate::proto::ObjectRequest>,
) -> Result<Response<Self::GetObjectStream>, Status> {
self.get_object_inner(request).await
}
async fn generate_presigned_url(
&self,
request: Request<UrlRequest>,
) -> Result<Response<UrlResponse>, Status> {
self.generate_presigned_url_inner(request).await
}
async fn initiate_multipart_upload(
&self,
request: Request<MultipartUploadRequest>,
) -> Result<Response<MultipartUploadResponse>, Status> {
self.initiate_multipart_upload_inner(request).await
}
async fn begin_tx(
&self,
request: Request<tonic::Streaming<Mutation>>,
) -> Result<Response<Self::BeginTxStream>, Status> {
self.begin_tx_inner(request).await
}
async fn publish_cdc(
&self,
request: Request<CdcSubscriptionRequest>,
) -> Result<Response<Self::PublishCDCStream>, Status> {
self.publish_cdc_inner(request).await
}
async fn create_materialized_view(
&self,
request: Request<ViewDefinition>,
) -> Result<Response<MutationResponse>, Status> {
self.create_materialized_view_inner(request).await
}
async fn enqueue_outbox_event(
&self,
request: Request<EnqueueOutboxEventRequest>,
) -> Result<Response<EnqueueOutboxEventResponse>, Status> {
self.enqueue_outbox_event_inner(request).await
}
async fn get_capabilities(
&self,
request: Request<CapabilitiesRequest>,
) -> Result<Response<CapabilitiesResponse>, Status> {
self.get_capabilities_inner(request).await
}
async fn get_catalog_manifest(
&self,
request: Request<CatalogManifestRequest>,
) -> Result<Response<CatalogManifestResponse>, Status> {
self.get_catalog_manifest_inner(request).await
}
async fn lookup_message_schema(
&self,
request: Request<MessageSchemaLookupRequest>,
) -> Result<Response<MessageSchemaLookupResponse>, Status> {
self.lookup_message_schema_inner(request).await
}
async fn list_message_schemas(
&self,
request: Request<MessageSchemaListRequest>,
) -> Result<Response<MessageSchemaListResponse>, Status> {
self.list_message_schemas_inner(request).await
}
async fn get_health_report(
&self,
request: Request<HealthReportRequest>,
) -> Result<Response<HealthReportResponse>, Status> {
self.get_health_report_inner(request).await
}
async fn delete(
&self,
request: Request<DeleteRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.delete_inner(request).await
}
async fn generic_dispatch(
&self,
request: Request<GenericDispatchRequest>,
) -> Result<Response<GenericDispatchResponse>, Status> {
self.generic_dispatch_inner(request).await
}
async fn ensure_resource(
&self,
request: Request<ResourceAdminRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.ensure_resource_inner(request).await
}
async fn drop_resource(
&self,
request: Request<ResourceAdminRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.drop_resource_inner(request).await
}
async fn list_resources(
&self,
request: Request<ResourceAdminRequest>,
) -> Result<Response<ResourceListResponse>, Status> {
self.list_resources_inner(request).await
}
async fn stage_catalog(
&self,
request: Request<StageCatalogRequest>,
) -> Result<Response<CatalogVersionResponse>, Status> {
self.stage_catalog_inner(request).await
}
async fn activate_catalog(
&self,
request: Request<CatalogVersionRequest>,
) -> Result<Response<CatalogVersionResponse>, Status> {
self.activate_catalog_inner(request).await
}
async fn rollback_catalog(
&self,
request: Request<CatalogVersionRequest>,
) -> Result<Response<CatalogVersionResponse>, Status> {
self.rollback_catalog_inner(request).await
}
async fn validate_catalog(
&self,
request: Request<StageCatalogRequest>,
) -> Result<Response<CatalogValidationResponse>, Status> {
self.validate_catalog_inner(request).await
}
async fn get_catalog_versions(
&self,
request: Request<CatalogManifestRequest>,
) -> Result<Response<CatalogVersionListResponse>, Status> {
self.get_catalog_versions_inner(request).await
}
async fn get_catalog_version(
&self,
request: Request<CatalogVersionRequest>,
) -> Result<Response<CatalogVersionResponse>, Status> {
self.get_catalog_version_inner(request).await
}
async fn plan_migration(
&self,
request: Request<MigrationPlanRequest>,
) -> Result<Response<MigrationPlanResponse>, Status> {
self.plan_migration_inner(request).await
}
async fn apply_migration(
&self,
request: Request<MigrationApplyRequest>,
) -> Result<Response<MigrationStatusResponse>, Status> {
self.apply_migration_inner(request).await
}
async fn get_migration_status(
&self,
request: Request<MigrationRunRequest>,
) -> Result<Response<MigrationStatusResponse>, Status> {
self.get_migration_status_inner(request).await
}
async fn list_migration_runs(
&self,
request: Request<MigrationRunListRequest>,
) -> Result<Response<MigrationRunListResponse>, Status> {
self.list_migration_runs_inner(request).await
}
async fn approve_migration_plan(
&self,
request: Request<MigrationRunRequest>,
) -> Result<Response<MigrationStatusResponse>, Status> {
self.approve_migration_plan_inner(request).await
}
async fn list_dlq_events(
&self,
request: Request<DlqListRequest>,
) -> Result<Response<DlqListResponse>, Status> {
self.list_dlq_events_inner(request).await
}
async fn get_dlq_event(
&self,
request: Request<DlqEventRequest>,
) -> Result<Response<DlqEventResponse>, Status> {
self.get_dlq_event_inner(request).await
}
async fn replay_dlq_event(
&self,
request: Request<DlqActionRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.replay_dlq_event_inner(request).await
}
async fn dismiss_dlq_event(
&self,
request: Request<DlqActionRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.dismiss_dlq_event_inner(request).await
}
async fn quarantine_dlq_event(
&self,
request: Request<DlqActionRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.quarantine_dlq_event_inner(request).await
}
async fn get_cdc_status(
&self,
request: Request<CdcControlRequest>,
) -> Result<Response<CdcStatusResponse>, Status> {
self.get_cdc_status_inner(request).await
}
async fn pause_cdc(
&self,
request: Request<CdcControlRequest>,
) -> Result<Response<CdcStatusResponse>, Status> {
self.pause_cdc_inner(request).await
}
async fn resume_cdc(
&self,
request: Request<CdcControlRequest>,
) -> Result<Response<CdcStatusResponse>, Status> {
self.resume_cdc_inner(request).await
}
async fn step_down_cdc_leader(
&self,
request: Request<CdcControlRequest>,
) -> Result<Response<CdcStatusResponse>, Status> {
self.step_down_cdc_leader_inner(request).await
}
async fn preview_cdc_redaction(
&self,
request: Request<CdcRedactionPreviewRequest>,
) -> Result<Response<CdcRedactionPreviewResponse>, Status> {
self.preview_cdc_redaction_inner(request).await
}
async fn scan_projection_drift(
&self,
request: Request<ProjectionDriftScanRequest>,
) -> Result<Response<ProjectionDriftScanResponse>, Status> {
self.scan_projection_drift_inner(request).await
}
async fn list_sagas(
&self,
request: Request<SagaListRequest>,
) -> Result<Response<SagaListResponse>, Status> {
self.list_sagas_inner(request).await
}
async fn get_saga(
&self,
request: Request<SagaRequest>,
) -> Result<Response<SagaResponse>, Status> {
self.get_saga_inner(request).await
}
async fn retry_saga_compensation(
&self,
request: Request<SagaRequest>,
) -> Result<Response<SagaResponse>, Status> {
self.retry_saga_compensation_inner(request).await
}
async fn mark_saga_reviewed(
&self,
request: Request<SagaRequest>,
) -> Result<Response<SagaResponse>, Status> {
self.mark_saga_reviewed_inner(request).await
}
async fn ensure_baseline(
&self,
request: Request<EnsureBaselineRequest>,
) -> Result<Response<EnsureBaselineResponse>, Status> {
self.ensure_baseline_inner(request).await
}
async fn list_policies(
&self,
request: Request<PolicyListRequest>,
) -> Result<Response<PolicyListResponse>, Status> {
self.list_policies_inner(request).await
}
async fn put_policy(
&self,
request: Request<PutPolicyRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.put_policy_inner(request).await
}
async fn delete_policy(
&self,
request: Request<PolicyRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.delete_policy_inner(request).await
}
async fn reload_policies(
&self,
request: Request<CapabilitiesRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.reload_policies_inner(request).await
}
async fn lint_policies(
&self,
request: Request<CapabilitiesRequest>,
) -> Result<Response<PolicyLintResponse>, Status> {
self.lint_policies_inner(request).await
}
async fn ensure_project(
&self,
request: Request<EnsureProjectRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.ensure_project_inner(request).await
}
async fn list_projects(
&self,
request: Request<ProjectListRequest>,
) -> Result<Response<ProjectListResponse>, Status> {
self.list_projects_inner(request).await
}
async fn get_admin_summary(
&self,
request: Request<AdminSummaryRequest>,
) -> Result<Response<AdminSummaryResponse>, Status> {
self.get_admin_summary_inner(request).await
}
async fn list_admin_audit_logs(
&self,
request: Request<AdminAuditLogRequest>,
) -> Result<Response<AdminAuditLogResponse>, Status> {
self.list_admin_audit_logs_inner(request).await
}
async fn verify_admin_audit_log(
&self,
request: Request<AdminAuditVerifyRequest>,
) -> Result<Response<AdminAuditVerifyResponse>, Status> {
self.verify_admin_audit_log_inner(request).await
}
async fn cache_get(
&self,
request: Request<crate::proto::CacheGetRequest>,
) -> Result<Response<crate::proto::CacheGetResponse>, Status> {
self.cache_get_inner(request).await
}
async fn cache_set(
&self,
request: Request<crate::proto::CacheSetRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.cache_set_inner(request).await
}
async fn cache_delete(
&self,
request: Request<crate::proto::CacheDeleteRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.cache_delete_inner(request).await
}
async fn cache_scan(
&self,
request: Request<crate::proto::CacheScanRequest>,
) -> Result<Response<crate::proto::CacheScanResponse>, Status> {
self.cache_scan_inner(request).await
}
async fn document_get(
&self,
request: Request<crate::proto::DocumentGetRequest>,
) -> Result<Response<crate::proto::DocumentSet>, Status> {
self.document_get_inner(request).await
}
async fn document_find(
&self,
request: Request<crate::proto::DocumentFindRequest>,
) -> Result<Response<crate::proto::DocumentSet>, Status> {
self.document_find_inner(request).await
}
async fn document_upsert(
&self,
request: Request<crate::proto::DocumentUpsertRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.document_upsert_inner(request).await
}
async fn document_delete(
&self,
request: Request<crate::proto::DocumentDeleteRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.document_delete_inner(request).await
}
async fn graph_query(
&self,
request: Request<crate::proto::GraphQueryRequest>,
) -> Result<Response<crate::proto::GraphResultSet>, Status> {
self.graph_query_inner(request).await
}
async fn graph_mutate(
&self,
request: Request<crate::proto::GraphMutationRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.graph_mutate_inner(request).await
}
async fn time_series_write(
&self,
request: Request<crate::proto::TimeSeriesWriteRequest>,
) -> Result<Response<MutationResponse>, Status> {
self.time_series_write_inner(request).await
}
async fn time_series_query(
&self,
request: Request<crate::proto::TimeSeriesQueryRequest>,
) -> Result<Response<crate::proto::TimeSeriesQueryResponse>, Status> {
self.time_series_query_inner(request).await
}
async fn analytical_query(
&self,
request: Request<crate::proto::AnalyticalQueryRequest>,
) -> Result<Response<crate::proto::AnalyticalQueryResponse>, Status> {
self.analytical_query_inner(request).await
}
}
#[cfg(test)]
mod live_tests;
#[cfg(test)]
mod tests;