use super::governance_logic::PolicyDocument;
use super::governance_store::version_state_to_db;
use super::*;
use crate::proto::udb::core::authz::entity::v1 as authz_entity_pb;
use crate::proto::udb::core::authz::services::v1 as authz_pb;
impl AuthzServiceImpl {
pub(super) async fn activate_policy_version_impl(
&self,
request: Request<authz_pb::ActivatePolicyVersionRequest>,
) -> Result<Response<authz_pb::ActivationResponse>, Status> {
let now = now_unix() as i64;
let req = request.into_inner();
let actor = self
.authorize_governance(
req.actor.as_ref(),
"ActivatePolicyVersion",
"governance.version.activate",
&[super::governance::SCOPE_AUTHZ_ADMIN],
now,
)
.await?;
if req.policy_version_id.trim().is_empty() {
return Err(Status::invalid_argument("policy_version_id is required"));
}
let version = self.load_version(&req.policy_version_id).await?;
let state =
authz_entity_pb::PolicyVersionState::try_from(version.state).unwrap_or_default();
if !matches!(
state,
authz_entity_pb::PolicyVersionState::Approved
| authz_entity_pb::PolicyVersionState::Superseded
| authz_entity_pb::PolicyVersionState::RolledBack
) {
return Err(Status::failed_precondition(format!(
"only an approved (or previously-active) version may be activated; version is {state:?}"
)));
}
if req.expected_revision != 0 && req.expected_revision != version.revision {
return Err(Status::aborted(
"policy version changed concurrently (expected_revision mismatch)",
));
}
let (cur_policy, cur_rel, _) = self
.current_authz_revision(&version.tenant_id, &version.project_id)
.await?;
if req.expected_policy_revision != 0 && req.expected_policy_revision != cur_policy {
return Err(Status::aborted(
"live authz policy revision changed concurrently",
));
}
if req.expected_relationship_revision != 0 && req.expected_relationship_revision != cur_rel
{
return Err(Status::aborted(
"live authz relationship revision changed concurrently",
));
}
let document = self.load_version_document(&req.policy_version_id).await?;
let (policy_rev, rel_rev) = self
.apply_document_and_activate(&version, &document, &actor, false)
.await?;
self.emit_event(AuthEvent::new(
topics::POLICY_VERSION_ACTIVATED,
version.policy_version_id.clone(),
version.tenant_id.clone(),
serde_json::json!({
"policy_version_id": version.policy_version_id,
"policy_set_id": version.policy_set_id,
"version_number": version.version_number,
"activated_by": actor,
"policy_revision": policy_rev,
}),
))
.await;
self.emit_bundle_invalidation(
&version.tenant_id,
&version.project_id,
policy_rev,
rel_rev,
"policy version activated",
&actor,
)
.await;
let activated = self.load_version(&req.policy_version_id).await?;
Ok(Response::new(authz_pb::ActivationResponse {
version: Some(activated),
policy_set: self
.load_policy_set(&version.policy_set_id)
.await
.ok()
.flatten(),
policy_revision: policy_rev,
relationship_revision: rel_rev,
content_hash: document.content_hash(),
}))
}
pub(super) async fn rollback_policy_version_impl(
&self,
request: Request<authz_pb::RollbackPolicyVersionRequest>,
) -> Result<Response<authz_pb::ActivationResponse>, Status> {
let now = now_unix() as i64;
let req = request.into_inner();
let actor = self
.authorize_governance(
req.actor.as_ref(),
"RollbackPolicyVersion",
"governance.version.rollback",
&[super::governance::SCOPE_AUTHZ_ADMIN],
now,
)
.await?;
if req.policy_set_id.trim().is_empty() {
return Err(Status::invalid_argument("policy_set_id is required"));
}
let policy_set = self
.load_policy_set(&req.policy_set_id)
.await?
.ok_or_else(|| Status::not_found("policy set not found"))?;
let target_id = if req.target_version_id.trim().is_empty() {
policy_set.rollback_version_id.clone()
} else {
req.target_version_id.clone()
};
if target_id.trim().is_empty() {
return Err(Status::failed_precondition(
"no rollback target: supply target_version_id or activate a second version first",
));
}
let target = self.load_version(&target_id).await?;
let document = self.load_version_document(&target_id).await?;
let (policy_rev, rel_rev) = self
.apply_document_and_activate(&target, &document, &actor, true)
.await?;
self.emit_event(AuthEvent::new(
topics::POLICY_VERSION_ROLLED_BACK,
target.policy_version_id.clone(),
target.tenant_id.clone(),
serde_json::json!({
"policy_set_id": req.policy_set_id,
"restored_version_id": target.policy_version_id,
"rolled_back_by": actor,
"reason": req.change_reason,
"policy_revision": policy_rev,
}),
))
.await;
self.emit_bundle_invalidation(
&target.tenant_id,
&target.project_id,
policy_rev,
rel_rev,
"policy version rolled back",
&actor,
)
.await;
Ok(Response::new(authz_pb::ActivationResponse {
version: Some(self.load_version(&target_id).await?),
policy_set: self
.load_policy_set(&req.policy_set_id)
.await
.ok()
.flatten(),
policy_revision: policy_rev,
relationship_revision: rel_rev,
content_hash: document.content_hash(),
}))
}
async fn apply_document_and_activate(
&self,
version: &authz_entity_pb::PolicyVersion,
document: &PolicyDocument,
actor: &str,
is_rollback: bool,
) -> Result<(i64, i64), Status> {
let pool = self.require_pool()?;
let tenant = version.tenant_id.clone();
let project = version.project_id.clone();
let prior_active = self
.policy_set_active_version(&version.policy_set_id)
.await?;
let mut tx = pool
.begin()
.await
.map_err(|err| Status::internal(format!("activation tx begin failed: {err}")))?;
crate::runtime::core::set_request_local_settings(
&mut tx,
&crate::RequestContext {
tenant_id: tenant.clone(),
project_id: project.clone(),
..crate::RequestContext::default()
},
)
.await?;
let pm = self.policies_model();
sqlx::query(&format!(
"DELETE FROM {rel} WHERE {tenant_id} = $1 AND COALESCE({project_id}, '') = $2",
rel = pm.relation,
tenant_id = pm.q("tenant_id"),
project_id = pm.q("project_id"),
))
.bind(&tenant)
.bind(&project)
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("clear policies failed: {err}")))?;
let mut seen_policy_ids = std::collections::HashSet::new();
for p in &document.policies {
if !p.id.trim().is_empty() && !seen_policy_ids.insert(p.id.as_str()) {
return Err(Status::invalid_argument(format!(
"duplicate policy id in version document: {}",
p.id
)));
}
}
for p in &document.policies {
let mut attrs = serde_json::Map::new();
for (k, v) in &p.conditions {
attrs.insert(k.clone(), serde_json::Value::String(v.clone()));
}
attrs.insert(
"priority".into(),
serde_json::Value::String(p.priority.to_string()),
);
attrs.insert("role".into(), serde_json::Value::String(p.role.clone()));
attrs.insert(
"purpose".into(),
serde_json::Value::String(p.purpose.clone()),
);
attrs.insert(
"relationship".into(),
serde_json::Value::String(p.relationship.clone()),
);
attrs.insert(
"required_scopes".into(),
serde_json::Value::String(
crate::runtime::service::auth_service::mappings::scopes_to_db(
&p.required_scopes,
),
),
);
sqlx::query(&format!(
"INSERT INTO {rel} ({policy_id}, {subject}, {domain_col}, {object_col}, {action_col}, {effect_col}, {condition}, {description}, {is_active}, {tenant_id}, {project_id}, {attributes_json}) \
VALUES (COALESCE(NULLIF($1,'')::uuid, gen_random_uuid()), $2, $3, $4, $5, $6, '', '', $7, $3, $8, $9::JSONB)",
rel = pm.relation,
policy_id = pm.q("policy_id"),
subject = pm.q("subject"),
domain_col = pm.q("domain"),
object_col = pm.q("object"),
action_col = pm.q("action"),
effect_col = pm.q("effect"),
condition = pm.q("condition"),
description = pm.q("description"),
is_active = pm.q("is_active"),
tenant_id = pm.q("tenant_id"),
project_id = pm.q("project_id"),
attributes_json = pm.q("attributes_json"),
))
.bind(&p.id)
.bind(&p.subject)
.bind(if p.tenant.trim().is_empty() { &tenant } else { &p.tenant })
.bind(&p.resource)
.bind(&p.action)
.bind(effect_to_db(p.effect))
.bind(p.enabled)
.bind(if p.project.trim().is_empty() { &project } else { &p.project })
.bind(serde_json::Value::Object(attrs))
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("insert policy failed: {err}")))?;
}
let tm = self.relationship_tuples_model();
sqlx::query(&format!(
"DELETE FROM {rel} WHERE {tenant_id} = $1 AND COALESCE({project_id}, '') = $2",
rel = tm.relation,
tenant_id = tm.q("tenant_id"),
project_id = tm.q("project_id"),
))
.bind(&tenant)
.bind(&project)
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("clear tuples failed: {err}")))?;
for b in &document.role_bindings {
let scope_tenant = if b.tenant.trim().is_empty() {
&tenant
} else {
&b.tenant
};
sqlx::query(&format!(
"INSERT INTO {rel} ({tuple_kind}, {subject}, {domain_col}, {object_col}, {action_col}, {effect_col}, {condition}, {tenant_id}, {project_id}) \
VALUES ('grouping', $1, $2, '', $3, '', '{{}}', $2, $4) \
ON CONFLICT ({tuple_kind}, {subject}, {domain_col}, {object_col}, {action_col}, {effect_col}) DO NOTHING",
rel = tm.relation,
tuple_kind = tm.q("tuple_kind"),
subject = tm.q("subject"),
domain_col = tm.q("domain"),
object_col = tm.q("object"),
action_col = tm.q("action"),
effect_col = tm.q("effect"),
condition = tm.q("condition"),
tenant_id = tm.q("tenant_id"),
project_id = tm.q("project_id"),
))
.bind(&b.subject)
.bind(scope_tenant)
.bind(&b.role)
.bind(if b.project.trim().is_empty() { &project } else { &b.project })
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("insert grouping tuple failed: {err}")))?;
}
for t in &document.tuples {
let scope_tenant = if t.tenant.trim().is_empty() {
&tenant
} else {
&t.tenant
};
sqlx::query(&format!(
"INSERT INTO {rel} ({tuple_kind}, {subject}, {domain_col}, {object_col}, {action_col}, {effect_col}, {condition}, {tenant_id}, {project_id}) \
VALUES ('relationship', $1, $2, $3, $4, '', '{{}}', $2, $5) \
ON CONFLICT ({tuple_kind}, {subject}, {domain_col}, {object_col}, {action_col}, {effect_col}) DO NOTHING",
rel = tm.relation,
tuple_kind = tm.q("tuple_kind"),
subject = tm.q("subject"),
domain_col = tm.q("domain"),
object_col = tm.q("object"),
action_col = tm.q("action"),
effect_col = tm.q("effect"),
condition = tm.q("condition"),
tenant_id = tm.q("tenant_id"),
project_id = tm.q("project_id"),
))
.bind(&t.subject)
.bind(scope_tenant)
.bind(&t.object)
.bind(&t.relation)
.bind(if t.project.trim().is_empty() { &project } else { &t.project })
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("insert relationship tuple failed: {err}")))?;
}
let vm = self.policy_versions_model();
if let Some(prior) = &prior_active {
if prior != &version.policy_version_id {
sqlx::query(&format!(
"UPDATE {rel} SET {state} = '{superseded}' WHERE {policy_version_id} = $1::UUID",
rel = vm.relation,
state = vm.q("state"),
superseded = version_state_to_db(authz_entity_pb::PolicyVersionState::Superseded),
policy_version_id = vm.q("policy_version_id"),
))
.bind(prior)
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("supersede prior version failed: {err}")))?;
}
}
let new_state = if is_rollback {
authz_entity_pb::PolicyVersionState::RolledBack
} else {
authz_entity_pb::PolicyVersionState::Active
};
sqlx::query(&format!(
"UPDATE {rel} SET {state} = $2, {activated_by} = $3, {activated_at} = NOW(), {revision} = {revision} + 1, {rollback_of} = $4::UUID \
WHERE {policy_version_id} = $1::UUID",
rel = vm.relation,
state = vm.q("state"),
activated_by = vm.q("activated_by"),
activated_at = vm.q("activated_at"),
revision = vm.q("revision"),
rollback_of = vm.q("rollback_of"),
policy_version_id = vm.q("policy_version_id"),
))
.bind(&version.policy_version_id)
.bind(version_state_to_db(new_state))
.bind(actor)
.bind(prior_active.as_deref())
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("activate version failed: {err}")))?;
let sm = self.policy_sets_model();
sqlx::query(&format!(
"UPDATE {rel} SET {active_version_id} = $2::UUID, {rollback_version_id} = $3::UUID WHERE {policy_set_id} = $1::UUID",
rel = sm.relation,
active_version_id = sm.q("active_version_id"),
rollback_version_id = sm.q("rollback_version_id"),
policy_set_id = sm.q("policy_set_id"),
))
.bind(&version.policy_set_id)
.bind(&version.policy_version_id)
.bind(prior_active.as_deref())
.execute(&mut *tx)
.await
.map_err(|err| Status::internal(format!("update policy set pointers failed: {err}")))?;
tx.commit()
.await
.map_err(|err| Status::internal(format!("activation tx commit failed: {err}")))?;
let change_type = if is_rollback {
authz_entity_pb::AuthzChangeType::Rollback
} else {
authz_entity_pb::AuthzChangeType::Activation
};
let (policy_rev, rel_rev) = self
.bump_authz_revision(
&tenant,
&project,
change_type,
&document.content_hash(),
actor,
)
.await?;
self.invalidate_snapshot_cache();
Ok((policy_rev, rel_rev))
}
async fn policy_set_active_version(
&self,
policy_set_id: &str,
) -> Result<Option<String>, Status> {
let pool = self.require_pool()?;
let m = self.policy_sets_model();
let rel = m.relation.clone();
let row = sqlx::query(&format!(
"SELECT COALESCE({active_version_id}::TEXT, '') AS active FROM {rel} WHERE {policy_set_id} = $1::UUID",
active_version_id = m.q("active_version_id"),
policy_set_id = m.q("policy_set_id"),
))
.bind(policy_set_id)
.fetch_optional(pool)
.await
.map_err(|err| Status::internal(format!("read active version failed: {err}")))?;
Ok(row
.and_then(|r| r.try_get::<String, _>("active").ok())
.filter(|s| !s.trim().is_empty()))
}
}
use crate::runtime::service::auth_service::control_plane::canary as canarylib;
use canarylib::{CanaryExecutor, CanaryMetricSource, CanaryScope};
impl AuthzServiceImpl {
pub(crate) fn canary_metrics(&self) -> Arc<dyn crate::metrics::MetricsRecorder> {
self.metrics.clone()
}
pub(crate) fn build_canary_metric_source(&self) -> Option<Arc<dyn CanaryMetricSource>> {
self.pg_pool
.clone()
.map(|pool| Arc::new(NackRateMetricSource::new(pool)) as Arc<dyn CanaryMetricSource>)
}
pub(crate) fn spawn_canary_evaluator(&self) -> tokio::task::JoinHandle<()> {
let executor: Arc<dyn CanaryExecutor> = Arc::new(self.clone());
let recorder = self.canary_metrics();
let source = self
.build_canary_metric_source()
.unwrap_or_else(|| Arc::new(canarylib::NoSignalSource) as Arc<dyn CanaryMetricSource>);
canarylib::spawn_canary_evaluator(
executor,
source,
recorder,
canarylib::CANARY_EVAL_INTERVAL,
Arc::new(|| {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}),
)
}
pub(super) fn policy_canaries_model(&self) -> NativeModel {
native_model(
"udb.core.authz.entity.v1.PolicyCanary",
&[
"canary_id",
"policy_set_id",
"policy_version_id",
"scope_kind",
"scope_values",
"state",
"success_window_secs",
"metric_threshold",
"created_by",
"tenant_id",
"project_id",
"min_samples",
"rollback_version_id",
"outcome_reason",
"revision",
],
)
}
pub(super) async fn activate_canary_impl(
&self,
request: Request<authz_pb::ActivateCanaryRequest>,
) -> Result<Response<authz_pb::CanaryResponse>, Status> {
let now = now_unix() as i64;
let req = request.into_inner();
let actor = self
.authorize_governance(
req.actor.as_ref(),
"ActivateCanary",
"governance.canary.activate",
&[super::governance::SCOPE_AUTHZ_ADMIN],
now,
)
.await?;
if req.policy_version_id.trim().is_empty() {
return Err(Status::invalid_argument("policy_version_id is required"));
}
let version = self.load_version(&req.policy_version_id).await?;
let state =
authz_entity_pb::PolicyVersionState::try_from(version.state).unwrap_or_default();
if !matches!(
state,
authz_entity_pb::PolicyVersionState::Approved
| authz_entity_pb::PolicyVersionState::Superseded
| authz_entity_pb::PolicyVersionState::RolledBack
) {
return Err(Status::failed_precondition(format!(
"only an approved (or previously-active) version may be canaried; version is {state:?}"
)));
}
if req.expected_revision != 0 && req.expected_revision != version.revision {
return Err(Status::aborted(
"policy version changed concurrently (expected_revision mismatch)",
));
}
let scope_kind = authz_entity_pb::CanaryScopeKind::try_from(req.scope_kind)
.unwrap_or(authz_entity_pb::CanaryScopeKind::Unspecified);
let scope = CanaryScope::from_row(scope_kind, &req.scope_values);
match &scope {
CanaryScope::Nodes(ids) | CanaryScope::Tenants(ids) if ids.is_empty() => {
return Err(Status::invalid_argument(
"canary scope_values must be non-empty for NODE/TENANT scope",
));
}
CanaryScope::Percent(0) => {
return Err(Status::invalid_argument(
"canary PERCENT scope must be 1..=100 (0 includes nobody)",
));
}
_ => {}
}
let rollback_target = self
.policy_set_active_version(&version.policy_set_id)
.await?
.unwrap_or_default();
let window = if req.success_window_secs > 0 {
req.success_window_secs
} else {
DEFAULT_CANARY_WINDOW_SECS
};
let threshold = if req.metric_threshold.is_finite() && req.metric_threshold >= 0.0 {
req.metric_threshold
} else {
DEFAULT_CANARY_THRESHOLD
};
let min_samples = req.min_samples.max(1);
let scope_values_json =
serde_json::to_string(&scope.values()).unwrap_or_else(|_| "[]".to_string());
let pool = self.require_pool()?;
let m = self.policy_canaries_model();
let canary_id: String = sqlx::query(&format!(
"INSERT INTO {rel} \
({policy_set_id}, {policy_version_id}, {scope_kind}, {scope_values}, {state}, \
{success_window_secs}, {metric_threshold}, {created_by}, {tenant_id}, {project_id}, \
{min_samples}, {rollback_version_id}, {outcome_reason}) \
VALUES ($1::UUID, $2::UUID, $3, $4::JSONB, $5, $6, $7, $8, $9, $10, $11, $12::UUID, '') \
RETURNING {canary_id}::TEXT AS canary_id",
rel = m.relation,
policy_set_id = m.q("policy_set_id"),
policy_version_id = m.q("policy_version_id"),
scope_kind = m.q("scope_kind"),
scope_values = m.q("scope_values"),
state = m.q("state"),
success_window_secs = m.q("success_window_secs"),
metric_threshold = m.q("metric_threshold"),
created_by = m.q("created_by"),
tenant_id = m.q("tenant_id"),
project_id = m.q("project_id"),
min_samples = m.q("min_samples"),
rollback_version_id = m.q("rollback_version_id"),
outcome_reason = m.q("outcome_reason"),
canary_id = m.q("canary_id"),
))
.bind(&version.policy_set_id)
.bind(&version.policy_version_id)
.bind(canary_scope_kind_to_db(scope.kind()))
.bind(scope_values_json)
.bind(canary_state_to_db(authz_entity_pb::CanaryState::Active))
.bind(window)
.bind(threshold)
.bind(&actor)
.bind(&version.tenant_id)
.bind(&version.project_id)
.bind(min_samples)
.bind(if rollback_target.is_empty() {
None
} else {
Some(rollback_target.clone())
})
.fetch_one(pool)
.await
.map_err(|err| Status::internal(format!("create canary failed: {err}")))?
.try_get("canary_id")
.map_err(|err| Status::internal(format!("create canary returned no id: {err}")))?;
let canary = self.load_canary(&canary_id).await?;
self.emit_canary_event(
topics::POLICY_CANARY_ACTIVATED,
&canary,
&actor,
"canary activated to subset scope",
)
.await;
Ok(Response::new(authz_pb::CanaryResponse {
canary: Some(canary),
version: Some(self.load_version(&req.policy_version_id).await?),
policy_set: self
.load_policy_set(&version.policy_set_id)
.await
.ok()
.flatten(),
}))
}
pub(super) async fn promote_canary_impl(
&self,
request: Request<authz_pb::PromoteCanaryRequest>,
) -> Result<Response<authz_pb::CanaryResponse>, Status> {
let now = now_unix() as i64;
let req = request.into_inner();
let actor = self
.authorize_governance(
req.actor.as_ref(),
"PromoteCanary",
"governance.canary.promote",
&[super::governance::SCOPE_AUTHZ_ADMIN],
now,
)
.await?;
if req.canary_id.trim().is_empty() {
return Err(Status::invalid_argument("canary_id is required"));
}
let canary = self.load_canary(&req.canary_id).await?;
if req.expected_revision != 0 && req.expected_revision != canary.revision {
return Err(Status::aborted(
"canary changed concurrently (expected_revision mismatch)",
));
}
if canary.state != authz_entity_pb::CanaryState::Active as i32 {
let st = authz_entity_pb::CanaryState::try_from(canary.state).unwrap_or_default();
return Err(Status::failed_precondition(format!(
"only an ACTIVE canary can be promoted; canary is {st:?}"
)));
}
if !canarylib::promote_eligible(&canary, now) {
return Err(Status::failed_precondition(
"canary is not promote-eligible yet: its success window has not elapsed",
));
}
let version = self.load_version(&canary.policy_version_id).await?;
let document = self
.load_version_document(&canary.policy_version_id)
.await?;
let (policy_rev, rel_rev) = self
.apply_document_and_activate(&version, &document, &actor, false)
.await?;
self.set_canary_state(
&canary.canary_id,
authz_entity_pb::CanaryState::Promoted,
"promoted fleet-wide after healthy bake",
)
.await?;
let promoted = self.load_canary(&canary.canary_id).await?;
self.emit_canary_event(
topics::POLICY_CANARY_PROMOTED,
&promoted,
&actor,
"canary promoted fleet-wide",
)
.await;
self.emit_bundle_invalidation(
&version.tenant_id,
&version.project_id,
policy_rev,
rel_rev,
"canary promoted fleet-wide",
&actor,
)
.await;
Ok(Response::new(authz_pb::CanaryResponse {
canary: Some(promoted),
version: Some(self.load_version(&canary.policy_version_id).await?),
policy_set: self
.load_policy_set(&version.policy_set_id)
.await
.ok()
.flatten(),
}))
}
pub(super) async fn get_canary_status_impl(
&self,
request: Request<authz_pb::GetCanaryStatusRequest>,
) -> Result<Response<authz_pb::GetCanaryStatusResponse>, Status> {
let now = now_unix() as i64;
let req = request.into_inner();
self.authorize_governance(
req.actor.as_ref(),
"GetCanaryStatus",
"governance.canary.read",
&[
super::governance::SCOPE_POLICY_READ,
super::governance::SCOPE_AUTHZ_ADMIN,
],
now,
)
.await?;
if req.canary_id.trim().is_empty() {
return Err(Status::invalid_argument("canary_id is required"));
}
let canary = self.load_canary(&req.canary_id).await?;
let promote_eligible = canarylib::promote_eligible(&canary, now);
let window_remaining_secs = canarylib::window_remaining_secs(&canary, now);
Ok(Response::new(authz_pb::GetCanaryStatusResponse {
canary: Some(canary),
promote_eligible,
window_remaining_secs,
}))
}
pub(super) async fn load_canary(
&self,
canary_id: &str,
) -> Result<authz_entity_pb::PolicyCanary, Status> {
let pool = self.require_pool()?;
let m = self.policy_canaries_model();
let rel = m.relation.clone();
let row = sqlx::query(&format!(
"SELECT {canary_id}::TEXT AS canary_id, {policy_set_id}::TEXT AS policy_set_id, \
{policy_version_id}::TEXT AS policy_version_id, {scope_kind}, \
{scope_values}::TEXT AS scope_values, {state}, {success_window_secs}, \
{metric_threshold}, COALESCE({created_by}, '') AS created_by, \
{tenant_id}::TEXT AS tenant_id, COALESCE({project_id}, '') AS project_id, \
{min_samples}, COALESCE({rollback_version_id}::TEXT, '') AS rollback_version_id, \
COALESCE({outcome_reason}, '') AS outcome_reason, {revision}, \
EXTRACT(EPOCH FROM {started_at})::BIGINT AS started_at_unix \
FROM {rel} WHERE {canary_id} = $1::UUID",
canary_id = m.q("canary_id"),
policy_set_id = m.q("policy_set_id"),
policy_version_id = m.q("policy_version_id"),
scope_kind = m.q("scope_kind"),
scope_values = m.q("scope_values"),
state = m.q("state"),
success_window_secs = m.q("success_window_secs"),
metric_threshold = m.q("metric_threshold"),
created_by = m.q("created_by"),
tenant_id = m.q("tenant_id"),
project_id = m.q("project_id"),
min_samples = m.q("min_samples"),
rollback_version_id = m.q("rollback_version_id"),
outcome_reason = m.q("outcome_reason"),
revision = m.q("revision"),
started_at = m.q("started_at"),
rel = rel,
))
.bind(canary_id)
.fetch_optional(pool)
.await
.map_err(|err| Status::internal(format!("load canary failed: {err}")))?
.ok_or_else(|| Status::not_found("canary not found"))?;
Ok(canary_from_row(&row))
}
pub(super) async fn load_active_canaries(
&self,
) -> Result<Vec<authz_entity_pb::PolicyCanary>, Status> {
let pool = self.require_pool()?;
let m = self.policy_canaries_model();
let rel = m.relation.clone();
let rows = sqlx::query(&format!(
"SELECT {canary_id}::TEXT AS canary_id, {policy_set_id}::TEXT AS policy_set_id, \
{policy_version_id}::TEXT AS policy_version_id, {scope_kind}, \
{scope_values}::TEXT AS scope_values, {state}, {success_window_secs}, \
{metric_threshold}, COALESCE({created_by}, '') AS created_by, \
{tenant_id}::TEXT AS tenant_id, COALESCE({project_id}, '') AS project_id, \
{min_samples}, COALESCE({rollback_version_id}::TEXT, '') AS rollback_version_id, \
COALESCE({outcome_reason}, '') AS outcome_reason, {revision}, \
EXTRACT(EPOCH FROM {started_at})::BIGINT AS started_at_unix \
FROM {rel} WHERE {state} = $1 ORDER BY {started_at} ASC",
canary_id = m.q("canary_id"),
policy_set_id = m.q("policy_set_id"),
policy_version_id = m.q("policy_version_id"),
scope_kind = m.q("scope_kind"),
scope_values = m.q("scope_values"),
state = m.q("state"),
success_window_secs = m.q("success_window_secs"),
metric_threshold = m.q("metric_threshold"),
created_by = m.q("created_by"),
tenant_id = m.q("tenant_id"),
project_id = m.q("project_id"),
min_samples = m.q("min_samples"),
rollback_version_id = m.q("rollback_version_id"),
outcome_reason = m.q("outcome_reason"),
revision = m.q("revision"),
started_at = m.q("started_at"),
rel = rel,
))
.bind(canary_state_to_db(authz_entity_pb::CanaryState::Active))
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("list active canaries failed: {err}")))?;
Ok(rows.iter().map(canary_from_row).collect())
}
pub(super) async fn set_canary_state(
&self,
canary_id: &str,
new_state: authz_entity_pb::CanaryState,
reason: &str,
) -> Result<bool, Status> {
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition("native authz requires runtime-backed canary persistence")
})?;
let mut assignments = std::collections::BTreeMap::new();
assignments.insert(
"state".to_string(),
LogicalAssignment::Set {
value: LogicalValue::String(canary_state_to_db(new_state).to_string()),
},
);
assignments.insert(
"outcome_reason".to_string(),
LogicalAssignment::Set {
value: LogicalValue::String(reason.to_string()),
},
);
assignments.insert(
"revision".to_string(),
LogicalAssignment::Increment {
by: LogicalValue::Int(1),
},
);
let (affected, _) = runtime
.native_entity_update_for_service(
"authz",
&crate::RequestContext::default(),
LogicalUpdate {
message_type: "udb.core.authz.entity.v1.PolicyCanary".to_string(),
filter: LogicalFilter::And(vec![
LogicalFilter::Comparison {
field: "canary_id".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String(canary_id.to_string()),
},
LogicalFilter::Comparison {
field: "state".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String(
canary_state_to_db(authz_entity_pb::CanaryState::Active)
.to_string(),
),
},
]),
assignments,
return_fields: Vec::new(),
require_affected: false,
},
)
.await
.map_err(|err| Status::internal(format!("update canary state failed: {err}")))?;
Ok(affected > 0)
}
pub(super) async fn emit_canary_event(
&self,
topic: &'static str,
canary: &authz_entity_pb::PolicyCanary,
actor: &str,
reason: &str,
) {
let severity = if topic == topics::POLICY_CANARY_ROLLED_BACK {
"high"
} else {
"info"
};
self.emit_event(AuthEvent::new(
topic,
canary.canary_id.clone(),
canary.tenant_id.clone(),
serde_json::json!({
"canary_id": canary.canary_id,
"policy_set_id": canary.policy_set_id,
"policy_version_id": canary.policy_version_id,
"rollback_version_id": canary.rollback_version_id,
"state": canary.state,
"scope_kind": canary.scope_kind,
"tenant_id": canary.tenant_id,
"project_id": canary.project_id,
"actor": actor,
"reason": reason,
"severity": severity,
}),
))
.await;
}
pub(super) async fn canary_auto_rollback(
&self,
canary: &authz_entity_pb::PolicyCanary,
reason: &str,
) -> Result<(), Status> {
if !self
.set_canary_state(
&canary.canary_id,
authz_entity_pb::CanaryState::RolledBack,
reason,
)
.await?
{
return Ok(());
}
let actor = "system:canary-evaluator";
if !canary.rollback_version_id.trim().is_empty() {
let target = self.load_version(&canary.rollback_version_id).await?;
let document = self
.load_version_document(&canary.rollback_version_id)
.await?;
let (policy_rev, rel_rev) = self
.apply_document_and_activate(&target, &document, actor, true)
.await?;
self.emit_event(AuthEvent::new(
topics::POLICY_VERSION_ROLLED_BACK,
target.policy_version_id.clone(),
target.tenant_id.clone(),
serde_json::json!({
"policy_set_id": canary.policy_set_id,
"restored_version_id": target.policy_version_id,
"rolled_back_by": actor,
"reason": reason,
"policy_revision": policy_rev,
"trigger": "canary_auto_rollback",
}),
))
.await;
self.emit_bundle_invalidation(
&target.tenant_id,
&target.project_id,
policy_rev,
rel_rev,
"canary auto-rollback",
actor,
)
.await;
}
let rolled = self.load_canary(&canary.canary_id).await?;
self.emit_canary_event(topics::POLICY_CANARY_ROLLED_BACK, &rolled, actor, reason)
.await;
Ok(())
}
pub(super) async fn canary_pause(
&self,
canary: &authz_entity_pb::PolicyCanary,
reason: &str,
) -> Result<(), Status> {
if !self
.set_canary_state(
&canary.canary_id,
authz_entity_pb::CanaryState::Paused,
reason,
)
.await?
{
return Ok(());
}
let paused = self.load_canary(&canary.canary_id).await?;
self.emit_canary_event(
topics::POLICY_CANARY_PAUSED,
&paused,
"system:canary-evaluator",
reason,
)
.await;
Ok(())
}
}
const DEFAULT_CANARY_WINDOW_SECS: i64 = 300;
const DEFAULT_CANARY_THRESHOLD: f64 = 0.05;
#[async_trait::async_trait]
impl CanaryExecutor for AuthzServiceImpl {
async fn list_active_canaries(&self) -> Vec<authz_entity_pb::PolicyCanary> {
match self.load_active_canaries().await {
Ok(canaries) => canaries,
Err(err) => {
tracing::warn!(error = %err, "canary evaluator: list_active_canaries failed");
Vec::new()
}
}
}
async fn auto_rollback(&self, canary: &authz_entity_pb::PolicyCanary, reason: &str) {
if let Err(err) = self.canary_auto_rollback(canary, reason).await {
tracing::error!(
canary_id = %canary.canary_id,
error = %err,
"canary auto-rollback failed"
);
}
}
async fn pause(&self, canary: &authz_entity_pb::PolicyCanary, reason: &str) {
if let Err(err) = self.canary_pause(canary, reason).await {
tracing::warn!(canary_id = %canary.canary_id, error = %err, "canary pause failed");
}
}
}
#[derive(Clone)]
pub struct NackRateMetricSource {
pool: PgPool,
}
impl NackRateMetricSource {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
}
#[async_trait::async_trait]
impl CanaryMetricSource for NackRateMetricSource {
async fn read(
&self,
canary: &authz_entity_pb::PolicyCanary,
_window_secs: i64,
) -> canarylib::CanarySignal {
let kind = authz_entity_pb::CanaryScopeKind::try_from(canary.scope_kind)
.unwrap_or(authz_entity_pb::CanaryScopeKind::Unspecified);
let values: Vec<String> =
serde_json::from_str(canary.scope_values.trim()).unwrap_or_default();
let scope = CanaryScope::from_row(kind, &values);
self.nack_fraction_in_scope(&scope)
.await
.unwrap_or(canarylib::CanarySignal {
value: 0.0,
samples: 0,
})
}
}
impl NackRateMetricSource {
async fn nack_fraction_in_scope(
&self,
scope: &CanaryScope,
) -> Result<canarylib::CanarySignal, Status> {
let m = native_model(
"udb.core.control.entity.v1.ControlPlaneNodeState",
&["node_id", "nack_error_detail"],
);
let rel = m.relation.clone();
let rows = sqlx::query(&format!(
"SELECT {node_id} AS node_id, \
(COALESCE({nack}, '') <> '') AS nacked \
FROM {rel}",
node_id = m.q("node_id"),
nack = m.q("nack_error_detail"),
rel = rel,
))
.fetch_all(&self.pool)
.await
.map_err(|err| Status::internal(format!("read node-state ledger failed: {err}")))?;
let mut samples = 0i64;
let mut nacked = 0i64;
for row in &rows {
let node_id: String = row.try_get("node_id").unwrap_or_default();
if matches!(scope, CanaryScope::Tenants(_)) {
continue;
}
if !scope.node_in_scope(&node_id) {
continue;
}
samples += 1;
if row.try_get::<bool, _>("nacked").unwrap_or(false) {
nacked += 1;
}
}
let value = if samples > 0 {
nacked as f64 / samples as f64
} else {
0.0
};
Ok(canarylib::CanarySignal { value, samples })
}
}
fn canary_from_row(row: &sqlx::postgres::PgRow) -> authz_entity_pb::PolicyCanary {
let started_at_unix: i64 = row.try_get("started_at_unix").unwrap_or(0);
authz_entity_pb::PolicyCanary {
canary_id: row.try_get("canary_id").unwrap_or_default(),
policy_set_id: row.try_get("policy_set_id").unwrap_or_default(),
policy_version_id: row.try_get("policy_version_id").unwrap_or_default(),
scope_kind: canary_scope_kind_from_db(
&row.try_get::<String, _>("scope_kind").unwrap_or_default(),
),
scope_values: row
.try_get("scope_values")
.unwrap_or_else(|_| "[]".to_string()),
state: canary_state_from_db(&row.try_get::<String, _>("state").unwrap_or_default()),
started_at: timestamp_from_unix(started_at_unix.max(0) as u64),
success_window_secs: row.try_get("success_window_secs").unwrap_or(0),
metric_threshold: row.try_get("metric_threshold").unwrap_or(0.0),
created_by: row.try_get("created_by").unwrap_or_default(),
tenant_id: row.try_get("tenant_id").unwrap_or_default(),
project_id: row.try_get("project_id").unwrap_or_default(),
min_samples: row.try_get("min_samples").unwrap_or(1),
rollback_version_id: row.try_get("rollback_version_id").unwrap_or_default(),
outcome_reason: row.try_get("outcome_reason").unwrap_or_default(),
revision: row.try_get("revision").unwrap_or(1),
}
}
fn canary_state_to_db(state: authz_entity_pb::CanaryState) -> &'static str {
use authz_entity_pb::CanaryState as S;
match state {
S::Unspecified => "CANARY_STATE_UNSPECIFIED",
S::Active => "CANARY_STATE_ACTIVE",
S::Promoted => "CANARY_STATE_PROMOTED",
S::RolledBack => "CANARY_STATE_ROLLED_BACK",
S::Paused => "CANARY_STATE_PAUSED",
}
}
fn canary_state_from_db(value: &str) -> i32 {
use authz_entity_pb::CanaryState as S;
let v = match value {
"CANARY_STATE_ACTIVE" | "ACTIVE" => S::Active,
"CANARY_STATE_PROMOTED" | "PROMOTED" => S::Promoted,
"CANARY_STATE_ROLLED_BACK" | "ROLLED_BACK" => S::RolledBack,
"CANARY_STATE_PAUSED" | "PAUSED" => S::Paused,
_ => S::Unspecified,
};
v as i32
}
fn canary_scope_kind_to_db(kind: authz_entity_pb::CanaryScopeKind) -> &'static str {
use authz_entity_pb::CanaryScopeKind as K;
match kind {
K::Unspecified => "CANARY_SCOPE_KIND_UNSPECIFIED",
K::Node => "CANARY_SCOPE_KIND_NODE",
K::Tenant => "CANARY_SCOPE_KIND_TENANT",
K::Percent => "CANARY_SCOPE_KIND_PERCENT",
}
}
fn canary_scope_kind_from_db(value: &str) -> i32 {
use authz_entity_pb::CanaryScopeKind as K;
let v = match value {
"CANARY_SCOPE_KIND_NODE" | "NODE" => K::Node,
"CANARY_SCOPE_KIND_TENANT" | "TENANT" => K::Tenant,
"CANARY_SCOPE_KIND_PERCENT" | "PERCENT" => K::Percent,
_ => K::Unspecified,
};
v as i32
}