pub(crate) mod canary;
pub(crate) mod resources;
pub(crate) mod scaling;
pub(crate) mod sourcing;
pub(crate) mod store;
pub(crate) mod subscriber;
#[cfg(test)]
mod tests;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use sqlx::PgPool;
use tokio_stream::Stream;
use tonic::{Request, Response, Status};
use crate::proto::udb::core::control::entity::v1::ResourceType;
use crate::proto::udb::core::control::services::v1 as control_pb;
use control_pb::control_plane_service_server::ControlPlaneService;
pub use crate::proto::udb::core::control::services::v1::control_plane_service_server::ControlPlaneServiceServer;
use super::DataBrokerService;
use super::mappings::{bounded_page_response, bounded_page_window, timestamp_from_unix};
use crate::runtime::metrics::MetricsRecorder;
use resources::{
ResourceModel, aggregate_version, ordered_resource_types, resource_type_from_i32,
resource_type_to_db,
};
const STREAM_POLL_INTERVAL: Duration = Duration::from_millis(1000);
#[derive(Clone)]
pub struct ControlPlaneServiceImpl {
pg_pool: Option<PgPool>,
reload: Option<Arc<tokio::sync::Notify>>,
scaler: Option<Arc<scaling::PushScaler>>,
metrics: Option<Arc<dyn MetricsRecorder>>,
}
impl ControlPlaneServiceImpl {
pub fn new() -> Self {
Self {
pg_pool: None,
reload: None,
scaler: None,
metrics: None,
}
}
pub fn with_postgres(mut self, pool: Option<PgPool>) -> Self {
self.pg_pool = pool;
self
}
pub fn with_reload(mut self, reload: Arc<tokio::sync::Notify>) -> Self {
self.reload = Some(reload);
self
}
pub fn with_scaler(mut self, scaler: Arc<scaling::PushScaler>) -> Self {
self.scaler = Some(scaler);
self
}
pub fn with_metrics(mut self, metrics: Arc<dyn MetricsRecorder>) -> Self {
self.metrics = Some(metrics);
self
}
fn require_pool(&self) -> Result<&PgPool, Status> {
self.pg_pool.as_ref().ok_or_else(|| {
crate::runtime::executor_utils::capability_status(
"control_plane",
"postgres_store",
"postgres_store",
"control-plane service requires a Postgres-backed store (no PG pool configured)",
)
})
}
fn rollback_target_required_status() -> Status {
crate::runtime::executor_utils::policy_status(
"control_plane_rollback",
"rollback_target_required",
"no retained snapshot to roll back to for this (node, resource_type, target_version)",
)
}
fn node_state_not_found_status() -> Status {
crate::runtime::executor_utils::schema_status(
tonic::Code::NotFound,
"control_plane",
"AckStatus",
"node_state_not_found",
"no node state for this (node, resource_type)",
)
}
fn internal_status(operation: impl Into<String>, message: impl Into<String>) -> Status {
crate::runtime::executor_utils::internal_status("control_plane", operation, message)
}
fn required_field(
field: &'static str,
description: &'static str,
message: &'static str,
) -> Status {
crate::runtime::executor_utils::invalid_argument_fields(message, [(field, description)])
}
fn empty_discovery_stream_status() -> Status {
Self::required_field(
"stream",
"must include an initial DiscoveryRequest",
"empty control discovery stream",
)
}
fn missing_discovery_node_id_status() -> Status {
Self::required_field(
"node_id",
"must be a non-empty node id on the first DiscoveryRequest",
"node_id is required on the first DiscoveryRequest",
)
}
fn empty_delta_stream_status() -> Status {
Self::required_field(
"stream",
"must include an initial DeltaDiscoveryRequest",
"empty control delta stream",
)
}
fn missing_delta_node_id_status() -> Status {
Self::required_field(
"node_id",
"must be a non-empty node id on the first DeltaDiscoveryRequest",
"node_id is required on the first DeltaDiscoveryRequest",
)
}
}
impl Default for ControlPlaneServiceImpl {
fn default() -> Self {
Self::new()
}
}
fn resource_to_pb(model: &ResourceModel) -> control_pb::Resource {
control_pb::Resource {
name: model.name.clone(),
version: if model.version.is_empty() {
model.content_hash.clone()
} else {
model.version.clone()
},
payload_json: model.payload_json.clone(),
resource_type: model.resource_type_enum() as i32,
}
}
fn node_state_to_pb(row: &store::NodeStateRow, current_world: &str) -> control_pb::NodeAckState {
control_pb::NodeAckState {
node_id: row.node_id.clone(),
resource_type: store::resource_type_of(&row.resource_type) as i32,
subscribed_names: store::parse_subscribed_names(&row.subscribed_names_json),
accepted_version: row.accepted_version.clone(),
last_good_version: row.last_good_version.clone(),
last_response_nonce: row.last_response_nonce.clone(),
nack_error_detail: row.nack_error_detail.clone(),
in_sync: !row.accepted_version.is_empty() && row.accepted_version == current_world,
updated_at: timestamp_from_unix(row.updated_at_unix.max(0) as u64),
}
}
type DiscoveryStream =
Pin<Box<dyn Stream<Item = Result<control_pb::DiscoveryResponse, Status>> + Send>>;
type DeltaStream =
Pin<Box<dyn Stream<Item = Result<control_pb::DeltaDiscoveryResponse, Status>> + Send>>;
#[tonic::async_trait]
impl ControlPlaneService for ControlPlaneServiceImpl {
type StreamResourcesStream = DiscoveryStream;
type DeltaResourcesStream = DeltaStream;
async fn stream_resources(
&self,
request: Request<tonic::Streaming<control_pb::DiscoveryRequest>>,
) -> Result<Response<Self::StreamResourcesStream>, Status> {
let pool = self.require_pool()?.clone();
let mut inbound = request.into_inner();
let first = inbound
.message()
.await
.map_err(|e| {
Self::internal_status("stream_resources", format!("control stream error: {e}"))
})?
.ok_or_else(Self::empty_discovery_stream_status)?;
let node_id = first.node_id.trim().to_string();
if node_id.is_empty() {
return Err(Self::missing_discovery_node_id_status());
}
for rt in ordered_resource_types() {
store::ensure_node_state(&pool, &node_id, *rt, &[]).await?;
}
apply_inbound(&pool, &node_id, &first, self.metrics.as_deref()).await?;
let inbound_pool = pool.clone();
let inbound_node = node_id.clone();
let inbound_metrics = self.metrics.clone();
tokio::spawn(async move {
while let Ok(Some(msg)) = inbound.message().await {
if !msg.node_id.trim().is_empty() && msg.node_id.trim() != inbound_node {
continue;
}
let _ = apply_inbound(
&inbound_pool,
&inbound_node,
&msg,
inbound_metrics.as_deref(),
)
.await;
}
});
let out_pool = pool.clone();
let out_node = node_id.clone();
let reload = self.reload.clone();
let scaler = self.scaler.clone();
let out = async_stream::try_stream! {
let mut last_sent: std::collections::BTreeMap<i32, String> =
std::collections::BTreeMap::new();
loop {
for rt in ordered_resource_types() {
let rt = *rt;
let state = store::get_node_state(&out_pool, &out_node, rt)
.await?
.unwrap_or_default();
let names = store::parse_subscribed_names(&state.subscribed_names_json);
let resources = store::list_resources(&out_pool, rt, None, &names).await?;
let world = aggregate_version(&resources);
let already_applied = state.accepted_version == world;
let already_sent = last_sent.get(&(rt as i32)) == Some(&world);
if already_applied || already_sent {
continue;
}
let nonce = store::next_response_nonce(&out_pool, &out_node, rt).await?;
store::retain_served_snapshot(&out_pool, &out_node, rt, &world, &resources)
.await?;
let pb_resources: Vec<control_pb::Resource> =
resources.iter().map(resource_to_pb).collect();
let _permit = match scaler.as_ref() {
Some(s) => Some(s.throttle.acquire().await),
None => None,
};
last_sent.insert(rt as i32, world.clone());
yield control_pb::DiscoveryResponse {
resource_type: rt as i32,
version_info: world,
nonce,
resources: pb_resources,
removed_resources: Vec::new(),
};
}
match reload.as_ref() {
Some(n) => {
tokio::select! {
_ = tokio::time::sleep(STREAM_POLL_INTERVAL) => {}
_ = n.notified() => {
if let Some(s) = scaler.as_ref() {
s.debouncer.notify();
s.debouncer.wait().await;
}
}
}
}
None => tokio::time::sleep(STREAM_POLL_INTERVAL).await,
}
}
};
Ok(Response::new(Box::pin(out)))
}
async fn delta_resources(
&self,
request: Request<tonic::Streaming<control_pb::DeltaDiscoveryRequest>>,
) -> Result<Response<Self::DeltaResourcesStream>, Status> {
let pool = self.require_pool()?.clone();
let mut inbound = request.into_inner();
let first = inbound
.message()
.await
.map_err(|e| {
Self::internal_status(
"delta_resources",
format!("control delta stream error: {e}"),
)
})?
.ok_or_else(Self::empty_delta_stream_status)?;
let node_id = first.node_id.trim().to_string();
if node_id.is_empty() {
return Err(Self::missing_delta_node_id_status());
}
for rt in ordered_resource_types() {
store::ensure_node_state(&pool, &node_id, *rt, &[]).await?;
}
apply_inbound_delta(&pool, &node_id, &first, self.metrics.as_deref()).await?;
let seed_rt = resource_type_from_i32(first.resource_type);
let mut known: std::collections::HashMap<String, String> = first
.initial_resource_versions
.iter()
.map(|(name, ver)| (delta_key(seed_rt, name), ver.clone()))
.collect();
let inbound_pool = pool.clone();
let inbound_node = node_id.clone();
let inbound_metrics = self.metrics.clone();
tokio::spawn(async move {
while let Ok(Some(msg)) = inbound.message().await {
if !msg.node_id.trim().is_empty() && msg.node_id.trim() != inbound_node {
continue;
}
let _ = apply_inbound_delta(
&inbound_pool,
&inbound_node,
&msg,
inbound_metrics.as_deref(),
)
.await;
}
});
let out_pool = pool.clone();
let out_node = node_id.clone();
let out = async_stream::try_stream! {
loop {
for rt in ordered_resource_types() {
let rt = *rt;
let state = store::get_node_state(&out_pool, &out_node, rt)
.await?
.unwrap_or_default();
let names = store::parse_subscribed_names(&state.subscribed_names_json);
let resources = store::list_resources(&out_pool, rt, None, &names).await?;
let world = aggregate_version(&resources);
let mut changed: Vec<control_pb::Resource> = Vec::new();
let live_names: std::collections::BTreeSet<String> =
resources.iter().map(|r| r.name.clone()).collect();
for r in &resources {
let key = delta_key(rt, &r.name);
let current = if r.version.is_empty() {
r.content_hash.clone()
} else {
r.version.clone()
};
if known.get(&key).map(String::as_str) != Some(current.as_str()) {
changed.push(resource_to_pb(r));
known.insert(key, current);
}
}
let removed: Vec<String> = known
.keys()
.filter(|k| k.starts_with(&delta_prefix(rt)))
.map(|k| delta_name(rt, k))
.filter(|name| !live_names.contains(name))
.collect();
for name in &removed {
known.remove(&delta_key(rt, name));
}
if changed.is_empty() && removed.is_empty() {
continue;
}
store::retain_served_snapshot(&out_pool, &out_node, rt, &world, &resources)
.await?;
let nonce = store::next_response_nonce(&out_pool, &out_node, rt).await?;
yield control_pb::DeltaDiscoveryResponse {
resource_type: rt as i32,
nonce,
resources: changed,
removed_resources: removed,
system_version_info: world,
};
}
tokio::time::sleep(STREAM_POLL_INTERVAL).await;
}
};
Ok(Response::new(Box::pin(out)))
}
async fn get_resources(
&self,
request: Request<control_pb::GetResourcesRequest>,
) -> Result<Response<control_pb::GetResourcesResponse>, Status> {
let req = request.into_inner();
let rt = resource_type_from_i32(req.resource_type);
if rt == ResourceType::Unspecified {
return Err(Self::required_field(
"resource_type",
"must specify a control-plane resource type",
"resource_type is required",
));
}
let pool = self.require_pool()?;
let tenant = if req.tenant_id.trim().is_empty() {
None
} else {
Some(req.tenant_id.as_str())
};
let resources = store::list_resources(pool, rt, tenant, &req.resource_names).await?;
let version = aggregate_version(&resources);
let (limit, offset, _) = bounded_page_window(req.page.as_ref());
let total = resources.len();
let windowed: Vec<control_pb::Resource> = resources
.iter()
.skip(offset)
.take(limit)
.map(resource_to_pb)
.collect();
Ok(Response::new(control_pb::GetResourcesResponse {
resources: windowed,
version_info: version,
page: Some(bounded_page_response(total, req.page.as_ref())),
}))
}
async fn list_node_states(
&self,
request: Request<control_pb::ListNodeStatesRequest>,
) -> Result<Response<control_pb::ListNodeStatesResponse>, Status> {
let pool = self.require_pool()?;
let req = request.into_inner();
let rt = resource_type_from_i32(req.resource_type);
let node = if req.node_id.trim().is_empty() {
None
} else {
Some(req.node_id.as_str())
};
let (limit, offset, _) = bounded_page_window(req.page.as_ref());
let (rows, total) =
store::list_node_states(pool, node, rt, limit as i64, offset as i64).await?;
let mut states = Vec::with_capacity(rows.len());
for row in &rows {
let row_rt = store::resource_type_of(&row.resource_type);
let names = store::parse_subscribed_names(&row.subscribed_names_json);
let world = store::world_version(pool, row_rt, None, &names).await?;
states.push(node_state_to_pb(row, &world));
}
Ok(Response::new(control_pb::ListNodeStatesResponse {
node_states: states,
page: Some(bounded_page_response(total as usize, req.page.as_ref())),
}))
}
async fn ack_status(
&self,
request: Request<control_pb::AckStatusRequest>,
) -> Result<Response<control_pb::AckStatusResponse>, Status> {
let req = request.into_inner();
if req.node_id.trim().is_empty() {
return Err(Self::required_field(
"node_id",
"must be a non-empty control-plane node id",
"node_id is required",
));
}
let rt = resource_type_from_i32(req.resource_type);
if rt == ResourceType::Unspecified {
return Err(Self::required_field(
"resource_type",
"must specify a control-plane resource type",
"resource_type is required",
));
}
let pool = self.require_pool()?;
let row = store::get_node_state(pool, req.node_id.trim(), rt)
.await?
.ok_or_else(Self::node_state_not_found_status)?;
let names = store::parse_subscribed_names(&row.subscribed_names_json);
let world = store::world_version(pool, rt, None, &names).await?;
let acknowledged = row.accepted_version == world && row.nack_error_detail.is_empty();
let nacked = !row.nack_error_detail.is_empty();
Ok(Response::new(control_pb::AckStatusResponse {
node_state: Some(node_state_to_pb(&row, &world)),
current_version: world,
acknowledged,
nacked,
}))
}
async fn rollback_resources(
&self,
request: Request<control_pb::RollbackResourcesRequest>,
) -> Result<Response<control_pb::RollbackResourcesResponse>, Status> {
let req = request.into_inner();
let node_id = req.node_id.trim();
if node_id.is_empty() {
return Err(Self::required_field(
"node_id",
"must be a non-empty control-plane node id",
"node_id is required",
));
}
let rt = resource_type_from_i32(req.resource_type);
if rt == ResourceType::Unspecified {
return Err(Self::required_field(
"resource_type",
"must specify a control-plane resource type",
"resource_type is required",
));
}
let pool = self.require_pool()?.clone();
let target = Some(req.target_version.trim()).filter(|v| !v.is_empty());
let snapshot = store::find_rollback_target(&pool, node_id, rt, target)
.await?
.ok_or_else(Self::rollback_target_required_status)?;
let mut restored = 0i32;
for r in &snapshot.resources {
store::upsert_resource(
&pool,
rt,
&r.name,
&r.tenant_id,
&r.project_id,
&r.payload_json,
"control.rollback",
)
.await?;
restored += 1;
}
if let Some(reload) = self.reload.as_ref() {
reload.notify_waiters();
}
let current = store::world_version(&pool, rt, None, &[]).await?;
Ok(Response::new(control_pb::RollbackResourcesResponse {
rolled_back_to_version: snapshot.version,
current_version: current,
resources_restored: restored,
}))
}
}
async fn apply_inbound(
pool: &PgPool,
node_id: &str,
req: &control_pb::DiscoveryRequest,
metrics: Option<&dyn MetricsRecorder>,
) -> Result<(), Status> {
let rt = resource_type_from_i32(req.resource_type);
if rt == ResourceType::Unspecified {
return Ok(());
}
if !req.resource_names.is_empty() {
store::ensure_node_state(pool, node_id, rt, &req.resource_names).await?;
} else {
store::ensure_node_state(pool, node_id, rt, &[]).await?;
}
match &req.error_detail {
Some(err) if !err.message.trim().is_empty() || err.code != 0 => {
let detail =
serde_json::json!({ "code": err.code, "message": err.message }).to_string();
store::record_nack(pool, node_id, rt, &req.response_nonce, &detail).await?;
if let Some(m) = metrics {
m.inc_control_nack(resource_type_to_db(rt));
}
}
_ if !req.version_info.trim().is_empty() => {
store::record_ack(pool, node_id, rt, &req.version_info, &req.response_nonce).await?;
}
_ => {}
}
Ok(())
}
async fn apply_inbound_delta(
pool: &PgPool,
node_id: &str,
req: &control_pb::DeltaDiscoveryRequest,
metrics: Option<&dyn MetricsRecorder>,
) -> Result<(), Status> {
let rt = resource_type_from_i32(req.resource_type);
if rt == ResourceType::Unspecified {
return Ok(());
}
let existing = store::get_node_state(pool, node_id, rt)
.await?
.map(|r| store::parse_subscribed_names(&r.subscribed_names_json))
.unwrap_or_default();
let mut set: std::collections::BTreeSet<String> = existing.into_iter().collect();
for name in &req.resource_names_subscribe {
if !name.trim().is_empty() {
set.insert(name.trim().to_string());
}
}
for name in &req.resource_names_unsubscribe {
set.remove(name.trim());
}
let names: Vec<String> = set.into_iter().collect();
store::ensure_node_state(pool, node_id, rt, &names).await?;
match &req.error_detail {
Some(err) if !err.message.trim().is_empty() || err.code != 0 => {
let detail =
serde_json::json!({ "code": err.code, "message": err.message }).to_string();
store::record_nack(pool, node_id, rt, &req.response_nonce, &detail).await?;
if let Some(m) = metrics {
m.inc_control_nack(resource_type_to_db(rt));
}
}
_ if !req.response_nonce.trim().is_empty() => {
let world = store::world_version(pool, rt, None, &names).await?;
store::record_ack(pool, node_id, rt, &world, &req.response_nonce).await?;
}
_ => {}
}
Ok(())
}
fn delta_prefix(rt: ResourceType) -> String {
format!("{}\u{1f}", rt as i32)
}
fn delta_key(rt: ResourceType, name: &str) -> String {
format!("{}\u{1f}{}", rt as i32, name)
}
fn delta_name(_rt: ResourceType, key: &str) -> String {
key.split('\u{1f}').nth(1).unwrap_or("").to_string()
}
impl DataBrokerService {
pub(crate) fn build_control_plane_service(&self) -> ControlPlaneServiceImpl {
let runtime = self.runtime.load_full();
let pg_pool = runtime
.native_store_pool_for_service("control", true, "")
.ok();
let scaler = Arc::new(scaling::PushScaler::from_env(Some(self.metrics.clone())));
let mut svc = ControlPlaneServiceImpl::new()
.with_postgres(pg_pool.clone())
.with_scaler(scaler)
.with_metrics(self.metrics.clone());
if let Some(pool) = pg_pool {
let config = Arc::new(runtime.config().clone());
let seed_pool = pool.clone();
let seed_cfg = config.clone();
tokio::spawn(async move {
if let Err(e) = sourcing::resync(&seed_pool, seed_cfg.as_ref()).await {
tracing::warn!(error = %e, "control-plane initial resource sync failed");
}
});
let authz_cell = self.authz_snapshot();
let authz_version: subscriber::AuthzVersionFn = Arc::new(move || {
let snapshot = authz_cell.load();
format!("{}:{}", snapshot.version, snapshot.relationship_version)
});
let handle = subscriber::SubscriberHandle::new(pool, config)
.with_metrics(self.metrics.clone())
.with_authz_version(authz_version);
svc = svc.with_reload(handle.reload_notify());
subscriber::spawn_control_plane_subscriber(handle);
}
svc
}
}