use std::collections::BTreeMap;
use std::sync::Arc;
use std::sync::Mutex;
use std::time::{Duration, Instant};
use arc_swap::ArcSwap;
use sha2::{Digest, Sha256};
use sqlx::{PgPool, Row};
use tonic::{Request, Response, Status};
use uuid::Uuid;
use crate::proto::udb::core::authz::entity::v1 as authz_entity_pb;
use crate::proto::udb::core::authz::services::v1 as authz_pb;
use authz_pb::authz_service_server::AuthzService;
use super::mappings::{
authz_principal_to_runtime, decision_to_pb, effect_to_entity, entity_effect_to_runtime,
page_response, policy_from_pg_row, policy_to_rule_pb, resource_to_runtime, scopes_to_db,
timestamp_from_unix,
};
use super::now_unix;
use crate::ir::{
ComparisonOp, ConflictStrategy, LogicalAssignment, LogicalDelete, LogicalFilter, LogicalRecord,
LogicalUpdate, LogicalValue,
};
use crate::runtime::DataBrokerRuntime;
use crate::runtime::authz::Effect;
use crate::runtime::authz::{
AuthzPolicy, AuthzQuery, AuthzSnapshot, Decision, PolicyEngine, Principal, RelationshipTuple,
ResourceRef, RoleBinding,
};
use crate::runtime::channels::{ChannelManager, ChannelPermit, OperationChannel};
use crate::runtime::native_catalog::{NativeModel, native_model};
use super::events::{self, AuthEvent, AuthEventSink, topics};
mod audit;
mod tuples;
mod governance;
mod governance_activate;
mod governance_drafts;
mod governance_logic;
mod governance_sim;
mod governance_store;
fn policy_rule_model() -> NativeModel {
native_model(
"udb.core.authz.entity.v1.PolicyRule",
&[
"policy_id",
"subject",
"domain",
"object",
"action",
"effect",
"condition",
"description",
"is_active",
"created_by",
"deleted_at",
"tenant_id",
"deleted_by",
"project_id",
"resource_type",
"attributes_json",
],
)
}
fn policy_tuple_model() -> NativeModel {
native_model(
"udb.core.authz.entity.v1.PolicyTuple",
&[
"tuple_kind",
"subject",
"domain",
"object",
"action",
"effect",
"condition",
"tenant_id",
"project_id",
],
)
}
fn role_model() -> NativeModel {
native_model(
"udb.core.authz.entity.v1.Role",
&[
"role_id",
"name",
"description",
"is_system",
"is_active",
"created_by",
"deleted_at",
"tenant_id",
"deleted_by",
"role_code",
"domain",
"project_id",
"scope_type",
"access_surface",
"metadata_json",
],
)
}
fn user_role_model() -> NativeModel {
native_model(
"udb.core.authz.entity.v1.UserRole",
&[
"user_role_id",
"user_id",
"role_id",
"domain",
"assigned_by",
"expires_at",
"created_by",
"tenant_id",
],
)
}
fn access_decision_audit_model() -> NativeModel {
native_model(
"udb.core.authz.entity.v1.AccessDecisionAudit",
&[
"decision_audit_id",
"user_id",
"domain",
"object",
"action",
"effect",
"decision_source",
"matched_rule",
"reason",
"ip_address",
"correlation_id",
"decided_at",
"tenant_id",
"decision_id",
"policy_version",
"relationship_version",
"purpose",
"scopes",
"matched_policy_ids",
"project_id",
"actor_kind",
"resource_type",
"trace_id",
"span_id",
"user_agent_hash",
"decision_input",
],
)
}
#[derive(Default, Clone)]
pub(super) struct AuditContext {
pub source_ip: String,
pub correlation_id: String,
pub trace_id: String,
pub span_id: String,
pub user_agent: String,
pub purpose: String,
pub decision_input: String,
}
impl AuditContext {
pub(super) fn from_attributes(
attributes: &std::collections::BTreeMap<String, String>,
purpose: &str,
) -> Self {
let get = |k: &str| attributes.get(k).cloned().unwrap_or_default();
let source_ip = {
let raw = get("ip_address");
let raw = if raw.is_empty() {
get("source_ip")
} else {
raw
};
events::mask_source_ip(&raw)
};
let user_agent = events::hash_user_agent(&get("user_agent"));
let correlation_id = {
let c = get("correlation_id");
if c.is_empty() {
get("x-correlation-id")
} else {
c
}
};
let mut input = serde_json::Map::new();
for (k, v) in attributes {
if matches!(
k.as_str(),
"ip_address" | "source_ip" | "user_agent" | "correlation_id" | "x-correlation-id"
) {
continue;
}
input.insert(k.clone(), serde_json::Value::String(v.clone()));
}
let mut input_value = serde_json::Value::Object(input);
events::redact_auth_payload(&mut input_value);
Self {
source_ip,
correlation_id,
trace_id: get("trace_id"),
span_id: get("span_id"),
user_agent,
purpose: purpose.to_string(),
decision_input: input_value.to_string(),
}
}
}
fn role_select_projection(model: &NativeModel) -> String {
[
model.text("role_id"),
model.select("name"),
model.text_or_empty("description"),
model.select("is_system"),
model.select("is_active"),
model.text_or_empty("created_by"),
model.text_or_empty_as("tenant_id", "tenant"),
model.text_or_empty_as("project_id", "project"),
model.text_or_empty("role_code"),
model.text_or_empty("domain"),
model.select("scope_type"),
model.text_or_empty("access_surface"),
model.json_text_as("metadata_json", "metadata_json"),
model.text_or_empty("deleted_by"),
]
.join(", ")
}
fn user_role_select_projection(model: &NativeModel) -> String {
[
model.text("user_role_id"),
model.text("user_id"),
model.text("role_id"),
model.text_or_empty("domain"),
model.text_or_empty("assigned_by"),
model.timestamp_unix_as("expires_at", "expires_at_unix"),
model.text_or_empty_as("tenant_id", "tenant"),
model.text_or_empty("created_by"),
]
.join(", ")
}
fn policy_rule_select_projection(model: &NativeModel) -> String {
[
model.text("policy_id"),
model.text_or_empty("subject"),
model.text_or_empty("domain"),
model.text_or_empty("object"),
model.text_or_empty("action"),
model.select("effect"),
model.text_or_empty("condition"),
model.text_or_empty("description"),
model.select("is_active"),
model.text_or_empty("created_by"),
model.text_or_empty("tenant_id"),
model.text_or_empty("deleted_by"),
model.text_or_empty("project_id"),
model.text_or_empty("resource_type"),
model.json_text_as("attributes_json", "attributes_json"),
]
.join(", ")
}
fn stable_audit_user_uuid(principal: &Principal) -> Uuid {
let subject = [
principal.subject.as_str(),
principal.user_id.as_str(),
principal.principal_id.as_str(),
principal.service_identity.as_str(),
]
.into_iter()
.find(|value| !value.trim().is_empty())
.unwrap_or("anonymous");
stable_uuid_from_subject(subject)
}
pub(super) fn stable_uuid_from_subject(subject: &str) -> Uuid {
if let Ok(uuid) = Uuid::parse_str(subject) {
return uuid;
}
use sha2::{Digest, Sha256};
let digest = Sha256::digest(subject.as_bytes());
let mut bytes = [0u8; 16];
bytes.copy_from_slice(&digest[..16]);
Uuid::from_bytes(bytes)
}
pub(super) fn parse_uuid_field(field_name: &str, value: &str) -> Result<Uuid, Status> {
Uuid::parse_str(value)
.map_err(|_| Status::invalid_argument(format!("{field_name} must be a UUID")))
}
pub(super) fn timestamp_unix_field(
field_name: &str,
value: Option<prost_types::Timestamp>,
) -> Result<Option<i64>, Status> {
let Some(value) = value else {
return Ok(None);
};
if value.seconds <= 0 {
return Err(Status::invalid_argument(format!(
"{field_name} must be a positive unix timestamp"
)));
}
Ok(Some(value.seconds))
}
fn tenant_from_domain(tenant_id: &str, domain: &str) -> String {
if !tenant_id.trim().is_empty() {
tenant_id.to_string()
} else if let Some((prefix, suffix)) = domain.split_once(':') {
if matches!(prefix, "tenant" | "project" | "resource") && !suffix.trim().is_empty() {
suffix.to_string()
} else {
domain.to_string()
}
} else {
domain.to_string()
}
}
fn tuple_condition_expired(condition: &str, now: u64) -> bool {
if condition.trim().is_empty() {
return false;
}
let Ok(value) = serde_json::from_str::<serde_json::Value>(condition) else {
return false;
};
let Some(expires_at) = value
.get("expires_at_unix")
.and_then(serde_json::Value::as_i64)
else {
return false;
};
expires_at > 0 && (expires_at as u64) <= now
}
fn enrich_resource(resource: &mut ResourceRef) {
if !resource.table.trim().is_empty() {
return;
}
let key = if !resource.message_type.trim().is_empty() {
resource.message_type.clone()
} else {
resource.resource_name.clone()
};
if key.trim().is_empty() {
return;
}
if let Some((schema, table)) = crate::runtime::native_catalog::native_relation(&key) {
resource.schema = schema;
resource.table = table;
if resource.backend.trim().is_empty() {
resource.backend = "postgres".to_string();
}
}
}
fn role_scope_type_to_db(scope_type: i32) -> &'static str {
match authz_entity_pb::RoleScopeType::try_from(scope_type).unwrap_or_default() {
authz_entity_pb::RoleScopeType::Global => "GLOBAL",
authz_entity_pb::RoleScopeType::Tenant => "TENANT",
authz_entity_pb::RoleScopeType::Project => "PROJECT",
authz_entity_pb::RoleScopeType::Resource => "RESOURCE",
authz_entity_pb::RoleScopeType::External => "EXTERNAL",
authz_entity_pb::RoleScopeType::Unspecified => "UNSPECIFIED",
}
}
fn role_scope_type_from_db(value: &str) -> i32 {
match value {
"GLOBAL" | "ROLE_SCOPE_TYPE_GLOBAL" => authz_entity_pb::RoleScopeType::Global as i32,
"TENANT" | "ROLE_SCOPE_TYPE_TENANT" => authz_entity_pb::RoleScopeType::Tenant as i32,
"PROJECT" | "ROLE_SCOPE_TYPE_PROJECT" => authz_entity_pb::RoleScopeType::Project as i32,
"RESOURCE" | "ROLE_SCOPE_TYPE_RESOURCE" => authz_entity_pb::RoleScopeType::Resource as i32,
"EXTERNAL" | "ROLE_SCOPE_TYPE_EXTERNAL" => authz_entity_pb::RoleScopeType::External as i32,
_ => authz_entity_pb::RoleScopeType::Unspecified as i32,
}
}
fn effect_to_db(effect: Effect) -> &'static str {
match effect {
Effect::Allow => "ALLOW",
Effect::Deny => "DENY",
}
}
pub(super) fn effect_from_db(value: &str) -> i32 {
match value {
"ALLOW" | "allow" | "POLICY_EFFECT_ALLOW" => authz_entity_pb::PolicyEffect::Allow as i32,
"DENY" | "deny" | "POLICY_EFFECT_DENY" => authz_entity_pb::PolicyEffect::Deny as i32,
_ => authz_entity_pb::PolicyEffect::Unspecified as i32,
}
}
fn role_from_row(row: &sqlx::postgres::PgRow) -> Result<authz_entity_pb::Role, Status> {
let map = |e: sqlx::Error| Status::internal(format!("decode role failed: {e}"));
Ok(authz_entity_pb::Role {
role_id: row.try_get("role_id").map_err(map)?,
name: row.try_get("name").map_err(map)?,
description: row.try_get("description").map_err(map)?,
is_system: row.try_get("is_system").map_err(map)?,
is_active: row.try_get("is_active").map_err(map)?,
created_by: row.try_get("created_by").map_err(map)?,
created_at: None,
updated_at: None,
deleted_at: None,
tenant_id: row.try_get("tenant").map_err(map)?,
deleted_by: row.try_get("deleted_by").map_err(map)?,
role_code: row.try_get("role_code").map_err(map)?,
domain: row.try_get("domain").map_err(map)?,
project_id: row.try_get("project").map_err(map)?,
scope_type: role_scope_type_from_db(&row.try_get::<String, _>("scope_type").map_err(map)?),
access_surface: row.try_get("access_surface").map_err(map)?,
metadata_json: row.try_get("metadata_json").map_err(map)?,
})
}
fn user_role_from_row(row: &sqlx::postgres::PgRow) -> Result<authz_entity_pb::UserRole, Status> {
let map = |e: sqlx::Error| Status::internal(format!("decode user role failed: {e}"));
Ok(authz_entity_pb::UserRole {
user_role_id: row.try_get("user_role_id").map_err(map)?,
user_id: row.try_get("user_id").map_err(map)?,
role_id: row.try_get("role_id").map_err(map)?,
domain: row.try_get("domain").map_err(map)?,
assigned_by: row.try_get("assigned_by").map_err(map)?,
assigned_at: None,
expires_at: timestamp_from_unix(
row.try_get::<i64, _>("expires_at_unix")
.map_err(map)?
.max(0) as u64,
),
created_at: None,
updated_at: None,
created_by: row.try_get("created_by").map_err(map)?,
tenant_id: row.try_get("tenant").map_err(map)?,
})
}
fn policy_rule_from_row(
row: &sqlx::postgres::PgRow,
) -> Result<authz_entity_pb::PolicyRule, Status> {
let map = |e: sqlx::Error| Status::internal(format!("decode policy rule failed: {e}"));
Ok(authz_entity_pb::PolicyRule {
policy_id: row.try_get("policy_id").map_err(map)?,
subject: row.try_get("subject").map_err(map)?,
domain: row.try_get("domain").map_err(map)?,
object: row.try_get("object").map_err(map)?,
action: row.try_get("action").map_err(map)?,
effect: effect_from_db(&row.try_get::<String, _>("effect").map_err(map)?),
condition: row.try_get("condition").map_err(map)?,
description: row.try_get("description").map_err(map)?,
is_active: row.try_get("is_active").map_err(map)?,
created_by: row.try_get("created_by").map_err(map)?,
created_at: None,
updated_at: None,
deleted_at: None,
tenant_id: row.try_get("tenant_id").map_err(map)?,
deleted_by: row.try_get("deleted_by").map_err(map)?,
project_id: row.try_get("project_id").map_err(map)?,
resource_type: row.try_get("resource_type").map_err(map)?,
attributes_json: row.try_get("attributes_json").map_err(map)?,
})
}
#[derive(Clone)]
pub struct AuthzServiceImpl {
snapshot: Arc<ArcSwap<AuthzSnapshot>>,
snapshot_loaded_at: Arc<Mutex<Option<Instant>>>,
snapshot_ttl: Duration,
pg_pool: Option<PgPool>,
event_sink: Arc<dyn AuthEventSink>,
metrics: Arc<dyn crate::metrics::MetricsRecorder>,
channels: Option<ChannelManager>,
runtime: Option<Arc<DataBrokerRuntime>>,
}
impl AuthzServiceImpl {
#[cfg(test)]
pub fn new(snapshot: AuthzSnapshot) -> Self {
Self {
snapshot: Arc::new(ArcSwap::from_pointee(snapshot)),
snapshot_loaded_at: Arc::new(Mutex::new(Some(Instant::now()))),
snapshot_ttl: authz_snapshot_ttl(),
pg_pool: None,
event_sink: events::noop_sink(),
metrics: Arc::new(crate::metrics::NoopMetrics),
channels: None,
runtime: None,
}
}
pub fn shared(snapshot: Arc<ArcSwap<AuthzSnapshot>>) -> Self {
Self {
snapshot,
snapshot_loaded_at: Arc::new(Mutex::new(None)),
snapshot_ttl: authz_snapshot_ttl(),
pg_pool: None,
event_sink: events::noop_sink(),
metrics: Arc::new(crate::metrics::NoopMetrics),
channels: None,
runtime: None,
}
}
pub(crate) fn with_channels(mut self, channels: Option<ChannelManager>) -> Self {
self.channels = channels;
self
}
pub(crate) fn with_runtime(mut self, runtime: Option<Arc<DataBrokerRuntime>>) -> Self {
self.runtime = runtime;
self
}
async fn admit(&self, tenant: &str) -> Result<Option<ChannelPermit>, Status> {
let Some(channels) = self.channels.as_ref() else {
return Ok(None);
};
let op = OperationChannel::Admin;
let tenant_hash = {
use sha2::{Digest, Sha256};
let digest = Sha256::digest(tenant.as_bytes());
digest[..8]
.iter()
.map(|b| format!("{b:02x}"))
.collect::<String>()
};
match channels
.acquire_fair_with_backpressure(op, Some(tenant), None, None, None, op.default_cost())
.await
{
Ok(permit) => {
self.metrics.record_fair_admission(
"default",
&tenant_hash,
"authz",
"default",
op.as_str(),
"accepted",
);
self.metrics.add_fair_cost(
"default",
&tenant_hash,
"authz",
"default",
op.as_str(),
f64::from(op.default_cost()),
);
Ok(Some(permit))
}
Err(err) => {
self.metrics.inc_channel_rejected(op.as_str());
self.metrics.record_fair_admission(
"default",
&tenant_hash,
"authz",
"default",
op.as_str(),
"rejected",
);
Err(err)
}
}
}
pub fn with_postgres(mut self, pool: Option<PgPool>) -> Self {
let has_pool = pool.is_some();
self.pg_pool = pool;
if has_pool {
self.invalidate_snapshot_cache();
}
self
}
pub(super) fn invalidate_snapshot_cache(&self) {
if let Ok(mut guard) = self.snapshot_loaded_at.lock() {
*guard = None;
}
}
pub(crate) fn snapshot_ttl(&self) -> Duration {
self.snapshot_ttl
}
pub(crate) async fn warm_shared_snapshot(&self) {
self.invalidate_snapshot_cache();
if let Err(err) = self.current_snapshot().await {
tracing::warn!(
error = %err,
"authz snapshot warm reload failed; retaining last good snapshot"
);
}
}
pub(crate) fn with_event_sink(mut self, sink: Arc<dyn AuthEventSink>) -> Self {
self.event_sink = sink;
self
}
pub(crate) fn with_metrics(
mut self,
metrics: Arc<dyn crate::metrics::MetricsRecorder>,
) -> Self {
self.metrics = metrics;
self
}
pub(super) async fn emit_event(&self, event: AuthEvent) {
let topic = event.topic;
if let Err(err) = self.event_sink.emit(event).await {
tracing::warn!(topic, error = %err, "failed to publish authz event");
}
}
pub(super) async fn decide_with_snapshot(
&self,
snapshot: &AuthzSnapshot,
principal: &Principal,
resource: &ResourceRef,
action: &str,
purpose: &str,
attributes: &BTreeMap<String, String>,
) -> Decision {
PolicyEngine::decide(
snapshot,
&AuthzQuery {
principal,
resource,
action,
purpose,
attributes,
},
)
.await
}
pub(super) fn policies_model(&self) -> NativeModel {
policy_rule_model()
}
pub(super) fn relationship_tuples_model(&self) -> NativeModel {
policy_tuple_model()
}
pub(super) fn roles_model(&self) -> NativeModel {
role_model()
}
pub(super) fn user_roles_model(&self) -> NativeModel {
user_role_model()
}
pub(super) fn audits_model(&self) -> NativeModel {
access_decision_audit_model()
}
pub(super) async fn role_code_for(
&self,
role_id: Uuid,
domain: &str,
) -> Result<Option<String>, Status> {
let pool = self.require_pool()?;
let m = self.roles_model();
let rel = m.relation.clone();
let row = sqlx::query(&format!(
"SELECT COALESCE(NULLIF({role_code}, ''), {name}) AS role FROM {rel} \
WHERE {role_id} = $1::UUID AND {deleted_at} IS NULL AND ($2 = '' OR {domain_col} = $2) LIMIT 1",
role_code = m.q("role_code"),
name = m.q("name"),
role_id = m.q("role_id"),
deleted_at = m.q("deleted_at"),
domain_col = m.q("domain"),
))
.bind(role_id)
.bind(domain)
.fetch_optional(pool)
.await
.map_err(|err| Status::internal(format!("resolve role code failed: {err}")))?;
Ok(row.and_then(|r| r.try_get::<String, _>("role").ok()))
}
pub(super) fn require_pool(&self) -> Result<&PgPool, Status> {
self.pg_pool.as_ref().ok_or_else(|| {
Status::failed_precondition(
"this operation requires a Postgres-backed auth store (no PG pool configured)",
)
})
}
pub(super) async fn write_decision_audit(
&self,
principal: &Principal,
resource: &ResourceRef,
action: &str,
decision: &Decision,
ctx: &AuditContext,
) {
if !decision.allowed {
self.metrics.record_authz_deny();
}
if decision.allowed && !decision.audit_required {
return;
}
let Some(runtime) = &self.runtime else {
tracing::warn!(
"skipping access-decision audit write: authz runtime handle is not wired"
);
return;
};
let user_uuid = stable_audit_user_uuid(principal);
let effect = if decision.allowed { "ALLOW" } else { "DENY" };
let source = if decision.matched_policy_ids.is_empty() {
"NO_MATCH"
} else if decision.allowed && decision.via_role {
"ROLE_POLICY"
} else {
"DIRECT_POLICY"
};
let domain = if principal.project_id.trim().is_empty() {
principal.tenant_id.clone()
} else {
principal.project_id.clone()
};
let actor_kind = if !principal.user_id.trim().is_empty() {
"user"
} else if !principal.service_identity.trim().is_empty() {
"service"
} else {
"external"
};
let scopes = decision.required_scopes.join(",");
let decision_input = serde_json::from_str(&ctx.decision_input)
.unwrap_or_else(|_| serde_json::Value::String(ctx.decision_input.clone()));
let mut record = LogicalRecord::new();
record.insert(
"decision_audit_id".to_string(),
LogicalValue::String(Uuid::new_v4().to_string()),
);
record.insert(
"user_id".to_string(),
LogicalValue::String(user_uuid.to_string()),
);
record.insert("domain".to_string(), LogicalValue::String(domain));
record.insert(
"object".to_string(),
LogicalValue::String(resource.resource_name.clone()),
);
record.insert(
"action".to_string(),
LogicalValue::String(action.to_string()),
);
record.insert(
"effect".to_string(),
LogicalValue::String(effect.to_string()),
);
record.insert(
"decision_source".to_string(),
LogicalValue::String(source.to_string()),
);
record.insert(
"matched_rule".to_string(),
LogicalValue::String(
decision
.matched_policy_ids
.first()
.cloned()
.unwrap_or_default(),
),
);
record.insert(
"reason".to_string(),
LogicalValue::String(decision.deny_reason.clone()),
);
record.insert(
"ip_address".to_string(),
LogicalValue::String(ctx.source_ip.clone()),
);
record.insert(
"correlation_id".to_string(),
LogicalValue::String(ctx.correlation_id.clone()),
);
record.insert(
"tenant_id".to_string(),
LogicalValue::String(principal.tenant_id.clone()),
);
record.insert(
"decision_id".to_string(),
LogicalValue::String(decision.decision_id.clone()),
);
record.insert(
"policy_version".to_string(),
LogicalValue::String(decision.policy_version.clone()),
);
record.insert(
"relationship_version".to_string(),
LogicalValue::String(decision.relationship_version.clone()),
);
record.insert(
"purpose".to_string(),
LogicalValue::String(ctx.purpose.clone()),
);
record.insert("scopes".to_string(), LogicalValue::String(scopes));
record.insert(
"matched_policy_ids".to_string(),
LogicalValue::Array(
decision
.matched_policy_ids
.iter()
.cloned()
.map(LogicalValue::String)
.collect(),
),
);
record.insert(
"project_id".to_string(),
LogicalValue::String(principal.project_id.clone()),
);
record.insert(
"actor_kind".to_string(),
LogicalValue::String(actor_kind.to_string()),
);
record.insert(
"resource_type".to_string(),
LogicalValue::String(resource.resource_type.clone()),
);
record.insert(
"trace_id".to_string(),
LogicalValue::String(ctx.trace_id.clone()),
);
record.insert(
"span_id".to_string(),
LogicalValue::String(ctx.span_id.clone()),
);
record.insert(
"user_agent_hash".to_string(),
LogicalValue::String(ctx.user_agent.clone()),
);
record.insert(
"decision_input".to_string(),
LogicalValue::Json(decision_input),
);
let audit_context = crate::RequestContext {
tenant_id: principal.tenant_id.clone(),
project_id: principal.project_id.clone(),
..crate::RequestContext::default()
};
let result = runtime
.native_entity_write_for_service(
"authz",
&audit_context,
"udb.core.authz.entity.v1.AccessDecisionAudit",
record,
ConflictStrategy::Error,
)
.await;
if let Err(err) = result {
tracing::warn!(error = %err, "failed to write access-decision audit");
}
if decision.allowed && decision.audit_required {
let actor = if principal.subject.trim().is_empty() {
principal.principal_id.clone()
} else {
principal.subject.clone()
};
let correlation = if ctx.correlation_id.trim().is_empty() {
format!("audit_allow:{}:{}", actor, resource.resource_name)
} else {
ctx.correlation_id.clone()
};
self.emit_event(
AuthEvent::new(
topics::ACCESS_AUDIT_REQUIRED_ALLOW,
actor.clone(),
principal.tenant_id.clone(),
serde_json::json!({
"subject": actor.clone(),
"resource": resource.resource_name.clone(),
"action": action,
"purpose": ctx.purpose.clone(),
"matched_policy_ids": decision.matched_policy_ids.clone(),
}),
)
.with_correlation(correlation)
.with_compliance(events::ComplianceEnvelope {
actor: actor.clone(),
actor_project: principal.project_id.clone(),
target_resource: resource.resource_name.clone(),
operation: action.to_string(),
outcome: "allow".to_string(),
reason_code: "audit_required_allow".to_string(),
decision_id: decision.decision_id.clone(),
policy_version: decision.policy_version.clone(),
relationship_version: decision.relationship_version.clone(),
trace_id: ctx.trace_id.clone(),
span_id: ctx.span_id.clone(),
..events::ComplianceEnvelope::default()
}),
)
.await;
}
}
pub(super) fn require_snapshot_fallback(&self) -> Result<(), Status> {
Err(Status::failed_precondition(
"native authz requires a Postgres-backed auth store",
))
}
async fn authz_revision_fingerprint(&self) -> Result<String, Status> {
let Some(pool) = &self.pg_pool else {
return Ok(String::new());
};
let m = self.authz_revisions_model();
let rows = sqlx::query(&format!(
"SELECT DISTINCT ON ({tenant_id}, {project_id}) \
{tenant_id}::TEXT AS tenant_id, COALESCE({project_id}, '') AS project_id, \
{policy_revision} AS policy_revision, \
{relationship_revision} AS relationship_revision, \
COALESCE({content_hash}, '') AS content_hash \
FROM {rel} \
ORDER BY {tenant_id}, {project_id}, {policy_revision} DESC, \
{relationship_revision} DESC, {changed_at} DESC",
rel = m.relation.clone(),
tenant_id = m.q("tenant_id"),
project_id = m.q("project_id"),
policy_revision = m.q("policy_revision"),
relationship_revision = m.q("relationship_revision"),
content_hash = m.q("content_hash"),
changed_at = m.q("changed_at"),
))
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("read authz revision fence failed: {err}")))?;
let mut hasher = Sha256::new();
for row in rows {
let tenant_id: String = row
.try_get("tenant_id")
.map_err(|err| Status::internal(format!("decode authz fence failed: {err}")))?;
let project_id: String = row
.try_get("project_id")
.map_err(|err| Status::internal(format!("decode authz fence failed: {err}")))?;
let policy_revision: i64 = row
.try_get("policy_revision")
.map_err(|err| Status::internal(format!("decode authz fence failed: {err}")))?;
let relationship_revision: i64 = row
.try_get("relationship_revision")
.map_err(|err| Status::internal(format!("decode authz fence failed: {err}")))?;
let content_hash: String = row
.try_get("content_hash")
.map_err(|err| Status::internal(format!("decode authz fence failed: {err}")))?;
hasher.update(tenant_id.as_bytes());
hasher.update([0]);
hasher.update(project_id.as_bytes());
hasher.update([0]);
hasher.update(policy_revision.to_be_bytes());
hasher.update(relationship_revision.to_be_bytes());
hasher.update(content_hash.as_bytes());
hasher.update([0xff]);
}
Ok(format!("{:x}", hasher.finalize()))
}
async fn load_snapshot_from_postgres(&self) -> Result<Option<AuthzSnapshot>, Status> {
let Some(pool) = &self.pg_pool else {
return Ok(None);
};
let policy = self.policies_model();
let binding = self.user_roles_model();
let role = self.roles_model();
let tuple = self.relationship_tuples_model();
let policy_rows = sqlx::query(&format!(
"SELECT {policy_id_text} AS id, COALESCE(NULLIF({attributes_json}->>'priority', '')::INT, 0) AS priority, {is_active} AS enabled, {effect}, {tenant_id} AS tenant, COALESCE({project_id}, '') AS project, {subject}, \
COALESCE({attributes_json}->>'role', '') AS role, {action}, {object_col} AS resource, COALESCE({attributes_json}->>'purpose', '') AS purpose, \
COALESCE({attributes_json}->>'relationship', '') AS relationship, {attributes_json} AS conditions, COALESCE({attributes_json}->>'required_scopes', '') AS required_scopes \
FROM {policy_rel} \
WHERE {deleted_at} IS NULL AND {is_active} = TRUE \
ORDER BY priority DESC, {policy_id} ASC",
policy_rel = policy.relation.clone(),
policy_id_text = format!("{}::TEXT", policy.q("policy_id")),
policy_id = policy.q("policy_id"),
attributes_json = policy.q("attributes_json"),
is_active = policy.q("is_active"),
effect = policy.q("effect"),
tenant_id = policy.q("tenant_id"),
project_id = policy.q("project_id"),
subject = policy.q("subject"),
action = policy.q("action"),
object_col = policy.q("object"),
deleted_at = policy.q("deleted_at"),
))
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("load authz policies failed: {err}")))?;
let mut policies = Vec::with_capacity(policy_rows.len());
for row in &policy_rows {
policies.push(
policy_from_pg_row(row).map_err(|err| {
Status::internal(format!("decode authz policy failed: {err}"))
})?,
);
}
let binding_rows = sqlx::query(&format!(
"SELECT ur.{user_id}::TEXT AS subject, COALESCE(NULLIF(r.{role_code}, ''), NULLIF(r.{name}, ''), ur.{role_id}::TEXT) AS role, ur.{tenant_id} AS tenant, COALESCE(r.{project_id}, '') AS project \
FROM {binding_rel} ur \
LEFT JOIN {role_rel} r ON r.{role_role_id} = ur.{role_id} \
WHERE (ur.{expires_at} IS NULL OR ur.{expires_at} > NOW()) \
AND r.{deleted_at} IS NULL \
AND (r.{is_active} IS NULL OR r.{is_active} = TRUE)",
binding_rel = binding.relation.clone(),
role_rel = role.relation.clone(),
user_id = binding.q("user_id"),
role_id = binding.q("role_id"),
tenant_id = binding.q("tenant_id"),
expires_at = binding.q("expires_at"),
role_role_id = role.q("role_id"),
role_code = role.q("role_code"),
name = role.q("name"),
project_id = role.q("project_id"),
deleted_at = role.q("deleted_at"),
is_active = role.q("is_active"),
))
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("load role bindings failed: {err}")))?;
let mut role_bindings = Vec::with_capacity(binding_rows.len());
for row in binding_rows {
role_bindings.push(RoleBinding {
subject: row.try_get("subject").map_err(|err| {
Status::internal(format!("decode role binding failed: {err}"))
})?,
role: row.try_get("role").map_err(|err| {
Status::internal(format!("decode role binding failed: {err}"))
})?,
tenant: row.try_get("tenant").map_err(|err| {
Status::internal(format!("decode role binding failed: {err}"))
})?,
project: row.try_get("project").map_err(|err| {
Status::internal(format!("decode role binding failed: {err}"))
})?,
});
}
let grouping_rows = sqlx::query(&format!(
"SELECT {subject}, {action} AS role, {tenant_id} AS tenant, COALESCE({project_id}, '') AS project, {condition} \
FROM {tuple_rel} \
WHERE {tuple_kind} = 'grouping'",
tuple_rel = tuple.relation.clone(),
subject = tuple.q("subject"),
action = tuple.q("action"),
tenant_id = tuple.q("tenant_id"),
project_id = tuple.q("project_id"),
condition = tuple.q("condition"),
tuple_kind = tuple.q("tuple_kind"),
))
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("load grouping tuples failed: {err}")))?;
let now = now_unix();
for row in grouping_rows {
let condition: String = row
.try_get("condition")
.map_err(|err| Status::internal(format!("decode grouping tuple failed: {err}")))?;
if tuple_condition_expired(&condition, now) {
continue;
}
role_bindings.push(RoleBinding {
subject: row.try_get("subject").map_err(|err| {
Status::internal(format!("decode grouping tuple failed: {err}"))
})?,
role: row.try_get("role").map_err(|err| {
Status::internal(format!("decode grouping tuple failed: {err}"))
})?,
tenant: row.try_get("tenant").map_err(|err| {
Status::internal(format!("decode grouping tuple failed: {err}"))
})?,
project: row.try_get("project").map_err(|err| {
Status::internal(format!("decode grouping tuple failed: {err}"))
})?,
});
}
let tuple_rows = sqlx::query(&format!(
"SELECT {subject}, {action} AS relation, {object_col}, {tenant_id} AS tenant, COALESCE({project_id}, '') AS project, {condition} FROM {tuple_rel} \
WHERE {tuple_kind} = 'relationship'",
tuple_rel = tuple.relation.clone(),
subject = tuple.q("subject"),
action = tuple.q("action"),
object_col = tuple.q("object"),
tenant_id = tuple.q("tenant_id"),
project_id = tuple.q("project_id"),
condition = tuple.q("condition"),
tuple_kind = tuple.q("tuple_kind"),
))
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("load relationship tuples failed: {err}")))?;
let mut tuples = Vec::with_capacity(tuple_rows.len());
for row in tuple_rows {
let condition: String = row.try_get("condition").map_err(|err| {
Status::internal(format!("decode relationship tuple failed: {err}"))
})?;
if tuple_condition_expired(&condition, now) {
continue;
}
tuples.push(RelationshipTuple {
subject: row.try_get("subject").map_err(|err| {
Status::internal(format!("decode relationship tuple failed: {err}"))
})?,
relation: row.try_get("relation").map_err(|err| {
Status::internal(format!("decode relationship tuple failed: {err}"))
})?,
object: row.try_get("object").map_err(|err| {
Status::internal(format!("decode relationship tuple failed: {err}"))
})?,
tenant: row.try_get("tenant").map_err(|err| {
Status::internal(format!("decode relationship tuple failed: {err}"))
})?,
project: row.try_get("project").map_err(|err| {
Status::internal(format!("decode relationship tuple failed: {err}"))
})?,
});
}
Ok(Some(AuthzSnapshot {
version: policy_content_version(&policies),
relationship_version: tuple_content_version(&tuples),
policies,
role_bindings,
tuples,
default_allow: false,
}))
}
pub(super) async fn current_snapshot(&self) -> Result<Arc<AuthzSnapshot>, Status> {
let cached_is_fresh = self
.snapshot_loaded_at
.lock()
.map(|guard| {
guard
.map(|t| t.elapsed() < self.snapshot_ttl)
.unwrap_or(false)
})
.unwrap_or(false);
if cached_is_fresh {
return Ok(self.snapshot.load_full());
}
let reload_started = Instant::now();
let revision_before = self.authz_revision_fingerprint().await?;
let loaded = self.load_snapshot_from_postgres().await?;
let revision_after = self.authz_revision_fingerprint().await?;
self.metrics
.observe_policy_reload_seconds(reload_started.elapsed().as_secs_f64());
if revision_before != revision_after {
return Err(Status::aborted(
"authz revision changed while loading snapshot; retry snapshot load",
));
}
if let Some(snapshot) = loaded {
self.snapshot.store(Arc::new(snapshot));
if let Ok(mut guard) = self.snapshot_loaded_at.lock() {
*guard = Some(Instant::now());
}
return Ok(self.snapshot.load_full());
}
self.require_snapshot_fallback()?;
if let Ok(mut guard) = self.snapshot_loaded_at.lock() {
*guard = Some(Instant::now());
}
Ok(self.snapshot.load_full())
}
}
fn authz_snapshot_ttl() -> Duration {
let secs = std::env::var("UDB_AUTHZ_SNAPSHOT_TTL_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(5)
.max(1);
Duration::from_secs(secs)
}
#[tonic::async_trait]
impl AuthzService for AuthzServiceImpl {
async fn authorize(
&self,
request: Request<authz_pb::AuthzRequest>,
) -> Result<Response<authz_pb::AuthzResponse>, Status> {
let started = Instant::now();
let req = request.into_inner();
let requested_scopes = req
.principal
.as_ref()
.map(|p| p.scopes.clone())
.unwrap_or_default();
let mut principal = req
.principal
.as_ref()
.map(authz_principal_to_runtime)
.unwrap_or_default();
if principal.tenant_id.trim().is_empty() {
principal.tenant_id = if req.tenant_id.trim().is_empty() {
req.domain.clone()
} else {
req.tenant_id.clone()
};
}
if principal.project_id.trim().is_empty() {
principal.project_id = req.project_id.clone();
}
if principal.subject.trim().is_empty() {
principal.subject = if !principal.user_id.trim().is_empty() {
principal.user_id.clone()
} else {
principal.principal_id.clone()
};
}
if principal.principal_id.trim().is_empty() {
principal.principal_id = principal.subject.clone();
}
if crate::runtime::service::method_security::claim_context_present() {
let ctx = crate::runtime::service::method_security::current_claim_context();
crate::runtime::service::method_security::enforce_body_tenant_matches_claim(
&ctx,
&req.tenant_id,
&req.project_id,
)?;
if !ctx.is_cross_tenant_admin() {
let base = ctx.to_principal();
principal.subject = base.subject;
principal.tenant_id = base.tenant_id;
principal.project_id = base.project_id;
principal.scopes = base.scopes;
if principal.principal_id.trim().is_empty() {
principal.principal_id = principal.subject.clone();
}
}
}
let mut resource = req
.resource
.as_ref()
.map(resource_to_runtime)
.unwrap_or_default();
enrich_resource(&mut resource);
let mut attributes: BTreeMap<String, String> = req.attributes.into_iter().collect();
if let Some(ctx) = req.context {
attributes.extend(ctx.attributes.into_iter());
}
let _admit = self.admit(&principal.tenant_id).await?;
let snap = self.current_snapshot().await?;
let decision = self
.decide_with_snapshot(
&snap,
&principal,
&resource,
&req.action,
&req.purpose,
&attributes,
)
.await;
tracing::debug!(
decision_id = %decision.decision_id,
allowed = decision.allowed,
action = %req.action,
latency_us = started.elapsed().as_micros() as u64,
"authz decision",
);
if !requested_scopes.is_empty() {
attributes.insert("requested_scopes".to_string(), requested_scopes.join(","));
}
let audit_ctx = AuditContext::from_attributes(&attributes, &req.purpose);
self.write_decision_audit(&principal, &resource, &req.action, &decision, &audit_ctx)
.await;
if !decision.allowed {
let subject = if principal.user_id.trim().is_empty() {
principal.subject.clone()
} else {
principal.user_id.clone()
};
self.emit_event(
AuthEvent::new(
topics::ACCESS_DENIED,
subject.clone(),
principal.tenant_id.clone(),
serde_json::json!({
"user_id": principal.user_id.clone(),
"subject": subject.clone(),
"tenant_id": principal.tenant_id.clone(),
"resource": resource.resource_name.clone(),
"action": req.action.clone(),
"deny_reason": decision.deny_reason.clone(),
"decision_id": decision.decision_id.clone(),
"requested_scopes": requested_scopes.clone(),
"effective_scopes": principal.scopes.clone(),
}),
)
.with_correlation(if decision.decision_id.trim().is_empty() {
format!("deny:{}:{}", subject, resource.resource_name)
} else {
decision.decision_id.clone()
})
.with_compliance(events::ComplianceEnvelope {
actor: subject.clone(),
actor_project: principal.project_id.clone(),
target_resource: resource.resource_name.clone(),
operation: req.action.clone(),
outcome: "deny".to_string(),
reason_code: if decision.deny_reason.trim().is_empty() {
"access_denied".to_string()
} else {
decision.deny_reason.clone()
},
decision_id: decision.decision_id.clone(),
policy_version: decision.policy_version.clone(),
relationship_version: decision.relationship_version.clone(),
source_ip: audit_ctx.source_ip.clone(),
trace_id: audit_ctx.trace_id.clone(),
span_id: audit_ctx.span_id.clone(),
..events::ComplianceEnvelope::default()
}),
)
.await;
}
Ok(Response::new(authz_pb::AuthzResponse {
decision: Some(decision_to_pb(&decision)),
}))
}
async fn put_role_binding(
&self,
request: Request<authz_pb::PutRoleBindingRequest>,
) -> Result<Response<authz_pb::AuthMutationResponse>, Status> {
self.put_role_binding_impl(request).await
}
async fn put_relationship(
&self,
request: Request<authz_pb::PutRelationshipRequest>,
) -> Result<Response<authz_pb::AuthMutationResponse>, Status> {
self.put_relationship_impl(request).await
}
async fn put_authz_policy(
&self,
request: Request<authz_pb::PutAuthzPolicyRequest>,
) -> Result<Response<authz_pb::AuthMutationResponse>, Status> {
if governance::governed_mode_enabled() {
return Err(Status::failed_precondition(
"governed mode: direct PutAuthzPolicy is disabled; create a policy draft and activate it (or use break-glass governance)",
));
}
let p = request
.into_inner()
.policy
.ok_or_else(|| Status::invalid_argument("policy is required"))?;
if p.id.trim().is_empty() {
return Err(Status::invalid_argument("policy id is required"));
}
let effect = if p.effect.eq_ignore_ascii_case("deny") {
Effect::Deny
} else if p.effect.eq_ignore_ascii_case("allow") {
Effect::Allow
} else {
return Err(Status::invalid_argument(format!(
"policy effect must be 'allow' or 'deny', got '{}'",
p.effect
)));
};
let policy = AuthzPolicy {
id: p.id,
priority: p.priority,
enabled: p.enabled,
effect,
tenant: p.tenant,
project: p.project,
subject: p.subject,
role: p.role,
action: p.action,
resource: p.resource,
purpose: p.purpose,
relationship: p.relationship,
conditions: p.conditions.into_iter().collect(),
required_scopes: p.required_scopes,
};
if self.pg_pool.is_some() {
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition(
"native authz requires runtime-backed policy persistence",
)
})?;
let policy_id = parse_uuid_field("policy.id", &policy.id)?;
let mut attributes = serde_json::Map::new();
for (key, value) in &policy.conditions {
attributes.insert(key.clone(), serde_json::Value::String(value.clone()));
}
attributes.insert(
"priority".to_string(),
serde_json::Value::String(policy.priority.to_string()),
);
attributes.insert(
"role".to_string(),
serde_json::Value::String(policy.role.clone()),
);
attributes.insert(
"purpose".to_string(),
serde_json::Value::String(policy.purpose.clone()),
);
attributes.insert(
"relationship".to_string(),
serde_json::Value::String(policy.relationship.clone()),
);
attributes.insert(
"required_scopes".to_string(),
serde_json::Value::String(scopes_to_db(&policy.required_scopes)),
);
let mut record = LogicalRecord::new();
record.insert(
"policy_id".to_string(),
LogicalValue::String(policy_id.to_string()),
);
record.insert(
"subject".to_string(),
LogicalValue::String(policy.subject.clone()),
);
record.insert(
"domain".to_string(),
LogicalValue::String(policy.tenant.clone()),
);
record.insert(
"object".to_string(),
LogicalValue::String(policy.resource.clone()),
);
record.insert(
"action".to_string(),
LogicalValue::String(policy.action.clone()),
);
record.insert(
"effect".to_string(),
LogicalValue::String(effect_to_db(policy.effect).to_string()),
);
record.insert("condition".to_string(), LogicalValue::String(String::new()));
record.insert(
"description".to_string(),
LogicalValue::String(String::new()),
);
record.insert("is_active".to_string(), LogicalValue::Bool(policy.enabled));
record.insert(
"tenant_id".to_string(),
LogicalValue::String(policy.tenant.clone()),
);
record.insert(
"project_id".to_string(),
LogicalValue::String(policy.project.clone()),
);
record.insert(
"attributes_json".to_string(),
LogicalValue::Json(serde_json::Value::Object(attributes)),
);
let context = crate::RequestContext {
tenant_id: policy.tenant.clone(),
project_id: policy.project.clone(),
..crate::RequestContext::default()
};
runtime
.native_entity_write_for_service(
"authz",
&context,
"udb.core.authz.entity.v1.PolicyRule",
record,
ConflictStrategy::update(vec![
"subject".to_string(),
"domain".to_string(),
"object".to_string(),
"action".to_string(),
"effect".to_string(),
"is_active".to_string(),
"tenant_id".to_string(),
"project_id".to_string(),
"attributes_json".to_string(),
]),
)
.await
.map_err(|err| Status::internal(format!("store authz policy failed: {err}")))?;
let _ = self
.bump_authz_revision(
&policy.tenant,
&policy.project,
authz_entity_pb::AuthzChangeType::Policy,
"policy-put",
"",
)
.await;
} else {
self.require_snapshot_fallback()?;
}
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::AuthMutationResponse {
ok: true,
message: "authz policy stored".to_string(),
}))
}
async fn lint_authz_policies(
&self,
_request: Request<authz_pb::LintAuthzPoliciesRequest>,
) -> Result<Response<authz_pb::LintAuthzPoliciesResponse>, Status> {
let snap = self.current_snapshot().await?;
let findings = PolicyEngine::lint(snap.as_ref())
.await
.into_iter()
.map(|f| format!("[{}] {}: {}", f.severity, f.category, f.message))
.collect();
Ok(Response::new(authz_pb::LintAuthzPoliciesResponse {
findings,
}))
}
async fn check_access(
&self,
request: Request<authz_pb::CheckAccessRequest>,
) -> Result<Response<authz_pb::CheckAccessResponse>, Status> {
let req = request.into_inner();
if req.user_id.trim().is_empty() {
return Err(Status::invalid_argument("user_id is required"));
}
if req.object.trim().is_empty() {
return Err(Status::invalid_argument("object is required"));
}
if req.action.trim().is_empty() {
return Err(Status::invalid_argument("action is required"));
}
let mut principal = req
.principal
.as_ref()
.map(authz_principal_to_runtime)
.unwrap_or_default();
if principal.principal_id.trim().is_empty() {
principal.principal_id = req.user_id.clone();
}
if principal.subject.trim().is_empty() {
principal.subject = req.user_id.clone();
}
if principal.user_id.trim().is_empty() {
principal.user_id = req.user_id.clone();
}
if principal.tenant_id.trim().is_empty() {
principal.tenant_id = if req.tenant_id.trim().is_empty() {
req.domain.clone()
} else {
req.tenant_id.clone()
};
}
if principal.project_id.trim().is_empty() {
principal.project_id = req.project_id.clone();
}
if crate::runtime::service::method_security::claim_context_present() {
let ctx = crate::runtime::service::method_security::current_claim_context();
crate::runtime::service::method_security::enforce_body_tenant_matches_claim(
&ctx,
&req.tenant_id,
&req.project_id,
)?;
if !ctx.is_cross_tenant_admin() {
let base = ctx.to_principal();
principal.subject = base.subject;
principal.user_id = base.user_id;
principal.tenant_id = base.tenant_id;
principal.project_id = base.project_id;
principal.scopes = base.scopes;
if principal.principal_id.trim().is_empty() {
principal.principal_id = principal.subject.clone();
}
}
}
let mut resource = req
.resource
.as_ref()
.map(resource_to_runtime)
.unwrap_or_default();
if resource.resource_name.trim().is_empty() {
resource.resource_name = req.object.clone();
}
if resource.message_type.trim().is_empty() {
resource.message_type = req.object.clone();
}
enrich_resource(&mut resource);
let mut attributes: BTreeMap<String, String> = req.attributes.into_iter().collect();
if let Some(ctx) = req.context {
attributes.extend(ctx.attributes.into_iter());
}
let _admit = self.admit(&principal.tenant_id).await?;
let snap = self.current_snapshot().await?;
let decision = self
.decide_with_snapshot(
&snap,
&principal,
&resource,
&req.action,
&req.purpose,
&attributes,
)
.await;
let audit_ctx = AuditContext::from_attributes(&attributes, &req.purpose);
self.write_decision_audit(&principal, &resource, &req.action, &decision, &audit_ctx)
.await;
Ok(Response::new(authz_pb::CheckAccessResponse {
allowed: decision.allowed,
effect: effect_to_entity(decision.effect),
matched_rule: decision
.matched_policy_ids
.first()
.cloned()
.unwrap_or_default(),
reason: decision.deny_reason.clone(),
decision: Some(decision_to_pb(&decision)),
}))
}
async fn create_role(
&self,
request: Request<authz_pb::CreateRoleRequest>,
) -> Result<Response<authz_pb::CreateRoleResponse>, Status> {
governance::guard_governed_role_mutation("CreateRole")?;
let req = request.into_inner();
if req.name.trim().is_empty() {
return Err(Status::invalid_argument("name is required"));
}
let created_by = if crate::runtime::service::method_security::claim_context_present() {
let ctx = crate::runtime::service::method_security::current_claim_context();
let claim_id = stable_uuid_from_subject(&ctx.subject);
if req.created_by.trim().is_empty() {
claim_id
} else {
let supplied = parse_uuid_field("created_by", &req.created_by)?;
if supplied != claim_id && !ctx.is_cross_tenant_admin() {
return Err(Status::permission_denied(
"created_by must match the authenticated caller",
));
}
supplied
}
} else {
if req.created_by.trim().is_empty() {
return Err(Status::invalid_argument("created_by is required"));
}
parse_uuid_field("created_by", &req.created_by)?
};
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition("native authz requires runtime-backed role persistence")
})?;
let role_id = Uuid::new_v4().to_string();
let tenant_id = tenant_from_domain(&req.tenant_id, &req.domain);
if tenant_id.trim().is_empty() {
return Err(Status::invalid_argument("tenant_id or domain is required"));
}
let metadata_json =
serde_json::to_string(&req.metadata).unwrap_or_else(|_| "{}".to_string());
let mut record = LogicalRecord::new();
record.insert("role_id".to_string(), LogicalValue::String(role_id.clone()));
record.insert("name".to_string(), LogicalValue::String(req.name.clone()));
record.insert(
"description".to_string(),
LogicalValue::String(req.description.clone()),
);
record.insert("is_system".to_string(), LogicalValue::Bool(false));
record.insert("is_active".to_string(), LogicalValue::Bool(true));
record.insert(
"created_by".to_string(),
LogicalValue::String(created_by.to_string()),
);
record.insert(
"tenant_id".to_string(),
LogicalValue::String(tenant_id.clone()),
);
record.insert(
"project_id".to_string(),
LogicalValue::String(req.project_id.clone()),
);
record.insert(
"role_code".to_string(),
LogicalValue::String(req.role_code.clone()),
);
record.insert(
"domain".to_string(),
LogicalValue::String(req.domain.clone()),
);
record.insert(
"scope_type".to_string(),
LogicalValue::String(role_scope_type_to_db(req.scope_type).to_string()),
);
record.insert(
"access_surface".to_string(),
LogicalValue::String(req.access_surface.clone()),
);
record.insert(
"metadata_json".to_string(),
LogicalValue::Json(
serde_json::from_str(&metadata_json)
.unwrap_or_else(|_| serde_json::Value::Object(Default::default())),
),
);
let context = crate::RequestContext {
tenant_id: tenant_id.clone(),
project_id: req.project_id.clone(),
..crate::RequestContext::default()
};
runtime
.native_entity_write_for_service(
"authz",
&context,
"udb.core.authz.entity.v1.Role",
record,
ConflictStrategy::Error,
)
.await
.map_err(|err| {
crate::runtime::executor_utils::prefix_status("create role failed", err)
})?;
self.emit_event(
AuthEvent::new(
topics::ROLE_CREATED,
role_id.clone(),
tenant_id.clone(),
serde_json::json!({
"role_id": role_id.clone(),
"role_code": req.role_code.clone(),
"tenant_id": tenant_id.clone(),
"project_id": req.project_id.clone(),
"created_by": req.created_by.clone(),
}),
)
.with_correlation(format!("role_create:{role_id}"))
.with_compliance(events::ComplianceEnvelope {
actor: if req.created_by.trim().is_empty() {
created_by.to_string()
} else {
req.created_by.clone()
},
actor_project: req.project_id.clone(),
target_resource: format!("role:{role_id}"),
operation: "role_create".to_string(),
outcome: "success".to_string(),
reason_code: "role_created".to_string(),
..events::ComplianceEnvelope::default()
}),
)
.await;
let _ = self
.bump_authz_revision(
&tenant_id,
&req.project_id,
authz_entity_pb::AuthzChangeType::Role,
"role-created",
&req.created_by,
)
.await;
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::CreateRoleResponse {
role: Some(authz_entity_pb::Role {
role_id,
name: req.name,
description: req.description,
is_system: false,
is_active: true,
created_by: req.created_by,
created_at: None,
updated_at: None,
deleted_at: None,
tenant_id,
deleted_by: String::new(),
role_code: req.role_code,
domain: req.domain,
project_id: req.project_id,
scope_type: req.scope_type,
access_surface: req.access_surface,
metadata_json,
}),
}))
}
async fn assign_role(
&self,
request: Request<authz_pb::AssignRoleRequest>,
) -> Result<Response<authz_pb::AssignRoleResponse>, Status> {
governance::guard_governed_role_mutation("AssignRole")?;
let req = request.into_inner();
use authz_entity_pb::PrincipalKind;
let principal_kind = PrincipalKind::try_from(req.principal_kind).unwrap_or_default();
let principal_ref = if !req.principal_id.trim().is_empty() {
req.principal_id.clone()
} else {
req.user_id.clone()
};
if principal_ref.trim().is_empty() || req.role_id.trim().is_empty() {
return Err(Status::invalid_argument(
"user_id (or principal_id) and role_id are required",
));
}
if matches!(principal_kind, PrincipalKind::Group) && req.principal_id.trim().is_empty() {
return Err(Status::invalid_argument(
"group role bindings require an explicit principal_id (IdP/SCIM group mapping)",
));
}
let assigned_by_uuid = if crate::runtime::service::method_security::claim_context_present()
{
let ctx = crate::runtime::service::method_security::current_claim_context();
let claim_id = stable_uuid_from_subject(&ctx.subject);
if req.assigned_by.trim().is_empty() {
claim_id
} else {
let supplied = parse_uuid_field("assigned_by", &req.assigned_by)?;
if supplied != claim_id && !ctx.is_cross_tenant_admin() {
return Err(Status::permission_denied(
"assigned_by must match the authenticated caller",
));
}
supplied
}
} else {
if req.assigned_by.trim().is_empty() {
return Err(Status::invalid_argument("assigned_by is required"));
}
parse_uuid_field("assigned_by", &req.assigned_by)?
};
let user_id = match principal_kind {
PrincipalKind::User | PrincipalKind::Unspecified => {
Uuid::parse_str(principal_ref.trim())
.unwrap_or_else(|_| stable_uuid_from_subject(principal_ref.trim()))
}
_ => stable_uuid_from_subject(principal_ref.trim()),
};
let role_id = parse_uuid_field("role_id", &req.role_id)?;
let assigned_by = assigned_by_uuid;
let expires_at_unix = timestamp_unix_field("expires_at", req.expires_at.clone())?;
let tenant_id = tenant_from_domain(&req.tenant_id, &req.domain);
if tenant_id.trim().is_empty() {
return Err(Status::invalid_argument("tenant_id or domain is required"));
}
let is_literal_principal = matches!(
principal_kind,
PrincipalKind::ServiceAccount
| PrincipalKind::Workload
| PrincipalKind::Group
| PrincipalKind::ExternalSubject
) || (matches!(
principal_kind,
PrincipalKind::Unspecified | PrincipalKind::User
) && Uuid::parse_str(principal_ref.trim()).is_err());
if is_literal_principal {
let role_code = self
.role_code_for(role_id, &req.domain)
.await?
.unwrap_or_else(|| req.role_id.clone());
let condition = serde_json::json!({
"source": "assign_role",
"principal_kind": req.principal_kind,
"expires_at_unix": expires_at_unix.unwrap_or(0),
})
.to_string();
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition(
"native authz requires runtime-backed tuple persistence",
)
})?;
let mut record = LogicalRecord::new();
record.insert(
"tuple_kind".to_string(),
LogicalValue::String("grouping".to_string()),
);
record.insert(
"subject".to_string(),
LogicalValue::String(principal_ref.trim().to_string()),
);
record.insert(
"domain".to_string(),
LogicalValue::String(tenant_id.clone()),
);
record.insert("object".to_string(), LogicalValue::String(String::new()));
record.insert(
"action".to_string(),
LogicalValue::String(role_code.clone()),
);
record.insert("effect".to_string(), LogicalValue::String(String::new()));
record.insert(
"condition".to_string(),
LogicalValue::String(condition.clone()),
);
record.insert(
"tenant_id".to_string(),
LogicalValue::String(tenant_id.clone()),
);
record.insert(
"project_id".to_string(),
LogicalValue::String(req.project_id.clone()),
);
let context = crate::RequestContext {
tenant_id: tenant_id.clone(),
project_id: req.project_id.clone(),
..crate::RequestContext::default()
};
runtime
.native_entity_write_for_service(
"authz",
&context,
"udb.core.authz.entity.v1.PolicyTuple",
record,
ConflictStrategy::update_on(
vec!["condition".to_string(), "tenant_id".to_string()],
vec![
"tuple_kind".to_string(),
"subject".to_string(),
"domain".to_string(),
"object".to_string(),
"action".to_string(),
"effect".to_string(),
],
),
)
.await
.map_err(|err| {
Status::internal(format!("assign role (principal) failed: {err}"))
})?;
self.emit_event(
AuthEvent::new(
topics::ROLE_ASSIGNED,
principal_ref.clone(),
tenant_id.clone(),
serde_json::json!({
"principal_id": principal_ref.clone(),
"principal_kind": req.principal_kind,
"role_id": req.role_id.clone(),
"role_code": role_code.clone(),
"tenant_id": tenant_id.clone(),
"domain": req.domain.clone(),
"assigned_by": req.assigned_by.clone(),
}),
)
.with_correlation(format!("role_assign:{principal_ref}:{}", req.role_id))
.with_compliance(events::ComplianceEnvelope {
actor: if req.assigned_by.trim().is_empty() {
principal_ref.clone()
} else {
req.assigned_by.clone()
},
actor_project: req.project_id.clone(),
target_resource: principal_ref.clone(),
operation: "role_assign".to_string(),
outcome: "success".to_string(),
reason_code: "role_assigned".to_string(),
..events::ComplianceEnvelope::default()
}),
)
.await;
let _ = self
.bump_authz_revision(
&tenant_id,
&req.project_id,
authz_entity_pb::AuthzChangeType::RoleAssignment,
"role-assignment",
&req.assigned_by,
)
.await;
self.invalidate_snapshot_cache();
return Ok(Response::new(authz_pb::AssignRoleResponse {
user_role: Some(authz_entity_pb::UserRole {
user_role_id: stable_uuid_from_subject(&format!(
"{principal_ref}:{}:{}",
req.role_id, req.domain
))
.to_string(),
user_id: principal_ref,
role_id: req.role_id,
domain: req.domain,
assigned_by: req.assigned_by.clone(),
assigned_at: None,
expires_at: req.expires_at,
created_at: None,
updated_at: None,
created_by: req.assigned_by,
tenant_id,
}),
}));
}
let new_user_role_id = Uuid::new_v4().to_string();
let expires_value = match expires_at_unix {
Some(seconds) if seconds > 0 => chrono::DateTime::from_timestamp(seconds, 0)
.map(LogicalValue::Timestamp)
.unwrap_or(LogicalValue::Null),
_ => LogicalValue::Null,
};
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition(
"native authz requires runtime-backed user-role persistence",
)
})?;
let mut record = LogicalRecord::new();
record.insert(
"user_role_id".to_string(),
LogicalValue::String(new_user_role_id.clone()),
);
record.insert(
"user_id".to_string(),
LogicalValue::String(user_id.to_string()),
);
record.insert(
"role_id".to_string(),
LogicalValue::String(role_id.to_string()),
);
record.insert(
"domain".to_string(),
LogicalValue::String(req.domain.clone()),
);
record.insert(
"assigned_by".to_string(),
LogicalValue::String(assigned_by.to_string()),
);
record.insert(
"tenant_id".to_string(),
LogicalValue::String(tenant_id.clone()),
);
record.insert(
"created_by".to_string(),
LogicalValue::String(req.assigned_by.clone()),
);
record.insert("expires_at".to_string(), expires_value);
let context = crate::RequestContext {
tenant_id: tenant_id.clone(),
project_id: req.project_id.clone(),
..crate::RequestContext::default()
};
let returned = runtime
.native_entity_write_for_service_returning(
"authz",
&context,
"udb.core.authz.entity.v1.UserRole",
record,
ConflictStrategy::update_on(
vec![
"assigned_by".to_string(),
"tenant_id".to_string(),
"created_by".to_string(),
"expires_at".to_string(),
],
vec![
"user_id".to_string(),
"role_id".to_string(),
"domain".to_string(),
],
),
vec!["user_role_id".to_string()],
)
.await
.map_err(|err| Status::internal(format!("assign role failed: {err}")))?;
let user_role_id = returned
.first()
.and_then(|r| r.get("user_role_id"))
.and_then(|v| v.as_str())
.map(|s| s.to_string())
.unwrap_or(new_user_role_id);
self.emit_event(
AuthEvent::new(
topics::ROLE_ASSIGNED,
req.user_id.clone(),
tenant_id.clone(),
serde_json::json!({
"user_role_id": user_role_id.clone(),
"user_id": req.user_id.clone(),
"role_id": req.role_id.clone(),
"tenant_id": tenant_id.clone(),
"domain": req.domain.clone(),
"assigned_by": req.assigned_by.clone(),
}),
)
.with_correlation(format!("role_assign:{}:{}", req.user_id, req.role_id))
.with_compliance(events::ComplianceEnvelope {
actor: if req.assigned_by.trim().is_empty() {
req.user_id.clone()
} else {
req.assigned_by.clone()
},
actor_project: req.project_id.clone(),
target_resource: req.user_id.clone(),
operation: "role_assign".to_string(),
outcome: "success".to_string(),
reason_code: "role_assigned".to_string(),
..events::ComplianceEnvelope::default()
}),
)
.await;
let _ = self
.bump_authz_revision(
&tenant_id,
&req.project_id,
authz_entity_pb::AuthzChangeType::RoleAssignment,
"role-assignment",
&req.assigned_by,
)
.await;
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::AssignRoleResponse {
user_role: Some(authz_entity_pb::UserRole {
user_role_id,
user_id: req.user_id,
role_id: req.role_id,
domain: req.domain,
assigned_by: req.assigned_by.clone(),
assigned_at: None,
expires_at: req.expires_at,
created_at: None,
updated_at: None,
created_by: req.assigned_by,
tenant_id,
}),
}))
}
async fn create_policy_rule(
&self,
request: Request<authz_pb::CreatePolicyRuleRequest>,
) -> Result<Response<authz_pb::CreatePolicyRuleResponse>, Status> {
if governance::governed_mode_enabled() {
return Err(Status::failed_precondition(
"governed mode: direct CreatePolicyRule is disabled; create a policy draft and activate it (or use break-glass governance)",
));
}
let req = request.into_inner();
if req.subject.trim().is_empty() {
return Err(Status::invalid_argument("subject is required"));
}
if req.domain.trim().is_empty() {
return Err(Status::invalid_argument("domain is required"));
}
if req.object.trim().is_empty() {
return Err(Status::invalid_argument("object is required"));
}
if req.action.trim().is_empty() {
return Err(Status::invalid_argument("action is required"));
}
let created_by = if crate::runtime::service::method_security::claim_context_present() {
let ctx = crate::runtime::service::method_security::current_claim_context();
let claim_id = stable_uuid_from_subject(&ctx.subject);
if req.created_by.trim().is_empty() {
claim_id
} else {
let supplied = parse_uuid_field("created_by", &req.created_by)?;
if supplied != claim_id && !ctx.is_cross_tenant_admin() {
return Err(Status::permission_denied(
"created_by must match the authenticated caller",
));
}
supplied
}
} else {
if req.created_by.trim().is_empty() {
return Err(Status::invalid_argument("created_by is required"));
}
parse_uuid_field("created_by", &req.created_by)?
};
let effect = entity_effect_to_runtime(req.effect)?;
let policy = AuthzPolicy {
id: Uuid::new_v4().to_string(),
priority: 0,
enabled: true,
effect,
tenant: if req.tenant_id.trim().is_empty() {
req.domain.clone()
} else {
req.tenant_id.clone()
},
project: req.project_id.clone(),
subject: req.subject.clone(),
role: String::new(),
action: req.action.clone(),
resource: req.object.clone(),
purpose: String::new(),
relationship: String::new(),
conditions: req.attributes.into_iter().collect(),
required_scopes: Vec::new(),
};
let policy_rule = policy_to_rule_pb(&policy);
if self.pg_pool.is_some() {
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition(
"native authz requires runtime-backed policy persistence",
)
})?;
let attributes = serde_json::to_value(&policy.conditions).map_err(|err| {
Status::internal(format!("encode policy conditions failed: {err}"))
})?;
let mut record = LogicalRecord::new();
record.insert(
"policy_id".to_string(),
LogicalValue::String(policy.id.clone()),
);
record.insert(
"subject".to_string(),
LogicalValue::String(policy.subject.clone()),
);
record.insert(
"domain".to_string(),
LogicalValue::String(req.domain.clone()),
);
record.insert(
"object".to_string(),
LogicalValue::String(policy.resource.clone()),
);
record.insert(
"action".to_string(),
LogicalValue::String(policy.action.clone()),
);
record.insert(
"effect".to_string(),
LogicalValue::String(effect_to_db(policy.effect).to_string()),
);
record.insert(
"condition".to_string(),
LogicalValue::String(req.condition.clone()),
);
record.insert(
"description".to_string(),
LogicalValue::String(req.description.clone()),
);
record.insert("is_active".to_string(), LogicalValue::Bool(true));
record.insert(
"created_by".to_string(),
LogicalValue::String(created_by.to_string()),
);
record.insert(
"tenant_id".to_string(),
LogicalValue::String(policy.tenant.clone()),
);
record.insert(
"project_id".to_string(),
LogicalValue::String(policy.project.clone()),
);
record.insert(
"resource_type".to_string(),
LogicalValue::String(req.resource_type.clone()),
);
record.insert(
"attributes_json".to_string(),
LogicalValue::Json(attributes),
);
let context = crate::RequestContext {
tenant_id: policy.tenant.clone(),
project_id: String::new(),
..crate::RequestContext::default()
};
let returned = runtime
.native_entity_write_for_service_returning(
"authz",
&context,
"udb.core.authz.entity.v1.PolicyRule",
record,
ConflictStrategy::Error,
vec!["policy_id".to_string()],
)
.await
.map_err(|err| Status::internal(format!("create policy rule failed: {err}")))?;
let created_policy_id = returned
.first()
.and_then(|r| r.get("policy_id"))
.and_then(|v| v.as_str())
.ok_or_else(|| {
Status::internal("create policy rule returned no persisted id".to_string())
})?
.to_string();
if created_policy_id != policy.id {
return Err(Status::internal(
"create policy rule returned mismatched policy_id",
));
}
let _ = self
.bump_authz_revision(
&policy.tenant,
&policy.project,
authz_entity_pb::AuthzChangeType::Policy,
"policy-created",
&req.created_by,
)
.await;
} else {
self.require_snapshot_fallback()?;
}
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::CreatePolicyRuleResponse {
policy: Some(authz_entity_pb::PolicyRule {
domain: req.domain,
description: req.description,
created_by: req.created_by,
resource_type: req.resource_type,
condition: req.condition,
..policy_rule
}),
}))
}
async fn list_user_permissions(
&self,
request: Request<authz_pb::ListUserPermissionsRequest>,
) -> Result<Response<authz_pb::ListUserPermissionsResponse>, Status> {
let req = request.into_inner();
if req.user_id.trim().is_empty() {
return Err(Status::invalid_argument("user_id is required"));
}
let snap = self.current_snapshot().await?;
let mut roles = Vec::new();
for binding in &snap.role_bindings {
if binding.subject == req.user_id
&& (req.domain.trim().is_empty()
|| binding.tenant == req.domain
|| binding.project == req.domain)
&& !roles.contains(&binding.role)
{
roles.push(binding.role.clone());
}
}
let mut permissions = Vec::new();
for policy in &snap.policies {
if !policy.enabled || policy.effect != Effect::Allow {
continue;
}
let domain_matches = req.domain.trim().is_empty()
|| policy.tenant == req.domain
|| policy.project == req.domain;
let subject_matches =
policy.subject.is_empty() || policy.subject == "*" || policy.subject == req.user_id;
let role_matches =
!policy.role.trim().is_empty() && roles.iter().any(|r| r == &policy.role);
if domain_matches && (subject_matches || role_matches) {
permissions.push(authz_pb::EffectivePermission {
object: policy.resource.clone(),
action: policy.action.clone(),
via_role: if role_matches {
policy.role.clone()
} else {
String::new()
},
resource_type: String::new(),
domain: if policy.tenant.trim().is_empty() {
policy.project.clone()
} else {
policy.tenant.clone()
},
});
}
}
Ok(Response::new(authz_pb::ListUserPermissionsResponse {
permissions,
}))
}
async fn list_access_decision_audits(
&self,
request: Request<authz_pb::ListAccessDecisionAuditsRequest>,
) -> Result<Response<authz_pb::ListAccessDecisionAuditsResponse>, Status> {
self.list_access_decision_audits_impl(request).await
}
async fn revoke_role(
&self,
request: Request<authz_pb::RevokeRoleRequest>,
) -> Result<Response<authz_pb::RevokeRoleResponse>, Status> {
governance::guard_governed_role_mutation("RevokeRole")?;
let req = request.into_inner();
if req.user_role_id.trim().is_empty() {
return Err(Status::invalid_argument("user_role_id is required"));
}
let user_role_id = parse_uuid_field("user_role_id", &req.user_role_id)?;
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition(
"native authz requires runtime-backed user-role persistence",
)
})?;
let op = LogicalDelete {
message_type: "udb.core.authz.entity.v1.UserRole".to_string(),
filter: LogicalFilter::Comparison {
field: "user_role_id".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String(user_role_id.to_string()),
},
return_fields: vec!["tenant_id".to_string()],
};
let context = crate::RequestContext::default();
let deleted_rows = runtime
.native_entity_delete_rows_for_service("authz", &context, op)
.await
.map_err(|err| Status::internal(format!("revoke role failed: {err}")))?;
let revoked = !deleted_rows.is_empty();
if let Some(row) = deleted_rows.first() {
let tenant: String = row
.get("tenant_id")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string();
let project = String::new();
self.emit_event(
AuthEvent::new(
topics::ROLE_REVOKED,
req.user_id.clone(),
tenant.clone(),
serde_json::json!({
"user_role_id": req.user_role_id.clone(),
"user_id": req.user_id.clone(),
"reason": req.reason.clone(),
"revoked_by": req.revoked_by.clone(),
}),
)
.with_correlation(format!("role_revoke:{}", req.user_role_id))
.with_compliance(events::ComplianceEnvelope {
actor: if req.revoked_by.trim().is_empty() {
req.user_id.clone()
} else {
req.revoked_by.clone()
},
target_resource: req.user_id.clone(),
operation: "role_revoke".to_string(),
outcome: "success".to_string(),
reason_code: if req.reason.trim().is_empty() {
"role_revoked".to_string()
} else {
req.reason.clone()
},
..events::ComplianceEnvelope::default()
}),
)
.await;
if !tenant.trim().is_empty() {
let _ = self
.bump_authz_revision(
&tenant,
&project,
authz_entity_pb::AuthzChangeType::RoleAssignment,
"role-revoked",
&req.revoked_by,
)
.await;
}
}
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::RevokeRoleResponse { revoked }))
}
async fn list_user_roles(
&self,
request: Request<authz_pb::ListUserRolesRequest>,
) -> Result<Response<authz_pb::ListUserRolesResponse>, Status> {
let req = request.into_inner();
if req.user_id.trim().is_empty() {
return Err(Status::invalid_argument("user_id is required"));
}
let user_id = parse_uuid_field("user_id", &req.user_id)?;
let pool = self.require_pool()?;
let user_role_model = self.user_roles_model();
let rel = user_role_model.relation.clone();
let projection = user_role_select_projection(&user_role_model);
let rows = sqlx::query(&format!(
"SELECT {projection} \
FROM {rel} \
WHERE {user_id} = $1::UUID \
AND ($2 = '' OR {domain_col} = $2) \
AND (NOT $3 OR {expires_at} IS NULL OR {expires_at} > NOW())",
user_id = user_role_model.q("user_id"),
domain_col = user_role_model.q("domain"),
expires_at = user_role_model.q("expires_at"),
))
.bind(user_id)
.bind(&req.domain)
.bind(req.active_only)
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("list user roles failed: {err}")))?;
let mut user_roles = Vec::with_capacity(rows.len());
for row in &rows {
user_roles.push(user_role_from_row(row)?);
}
Ok(Response::new(authz_pb::ListUserRolesResponse {
user_roles,
}))
}
async fn get_role(
&self,
request: Request<authz_pb::GetRoleRequest>,
) -> Result<Response<authz_pb::GetRoleResponse>, Status> {
let req = request.into_inner();
if req.role_id.trim().is_empty() && req.role_code.trim().is_empty() {
return Err(Status::invalid_argument("role_id or role_code is required"));
}
let role_id_filter = if req.role_id.trim().is_empty() {
None
} else {
Some(parse_uuid_field("role_id", &req.role_id)?)
};
let pool = self.require_pool()?;
let role_model = self.roles_model();
let rel = role_model.relation.clone();
let projection = role_select_projection(&role_model);
let row = sqlx::query(&format!(
"SELECT {projection} \
FROM {rel} \
WHERE {deleted_at} IS NULL \
AND (($1::UUID IS NOT NULL AND {role_id} = $1) OR ($1::UUID IS NULL AND {role_code} = $2)) \
AND ($3 = '' OR {domain_col} = $3) \
LIMIT 1",
role_id = role_model.q("role_id"),
role_code = role_model.q("role_code"),
domain_col = role_model.q("domain"),
deleted_at = role_model.q("deleted_at"),
))
.bind(role_id_filter)
.bind(&req.role_code)
.bind(&req.domain)
.fetch_optional(pool)
.await
.map_err(|err| Status::internal(format!("get role failed: {err}")))?;
match row {
Some(row) => Ok(Response::new(authz_pb::GetRoleResponse {
role: Some(role_from_row(&row)?),
})),
None => Err(Status::not_found("role not found")),
}
}
async fn list_roles(
&self,
request: Request<authz_pb::ListRolesRequest>,
) -> Result<Response<authz_pb::ListRolesResponse>, Status> {
let req = request.into_inner();
let pool = self.require_pool()?;
let role_model = self.roles_model();
let rel = role_model.relation.clone();
let projection = role_select_projection(&role_model);
let rows = sqlx::query(&format!(
"SELECT {projection} \
FROM {rel} \
WHERE {deleted_at} IS NULL \
AND ($1 = '' OR {domain_col} = $1) \
AND (NOT $2 OR {is_active} = TRUE) \
ORDER BY {name} ASC",
domain_col = role_model.q("domain"),
is_active = role_model.q("is_active"),
name = role_model.q("name"),
deleted_at = role_model.q("deleted_at"),
))
.bind(&req.domain)
.bind(req.active_only)
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("list roles failed: {err}")))?;
let mut all = Vec::with_capacity(rows.len());
for row in &rows {
all.push(role_from_row(row)?);
}
let page = req.page.as_ref();
let page_number = page.map(|p| p.page).filter(|p| *p > 0).unwrap_or(1) as usize;
let page_size = page
.map(|p| p.page_size)
.filter(|s| *s > 0)
.unwrap_or(all.len().max(1) as i32) as usize;
let start = page_number.saturating_sub(1).saturating_mul(page_size);
let total = all.len();
let roles = all.into_iter().skip(start).take(page_size).collect();
Ok(Response::new(authz_pb::ListRolesResponse {
page: Some(page_response(total, page)),
roles,
}))
}
async fn batch_check_permissions(
&self,
request: Request<authz_pb::BatchCheckPermissionsRequest>,
) -> Result<Response<authz_pb::BatchCheckPermissionsResponse>, Status> {
let req = request.into_inner();
if req.user_id.trim().is_empty() {
return Err(Status::invalid_argument("user_id is required"));
}
let attributes: BTreeMap<String, String> = req
.context
.map(|ctx| ctx.attributes.into_iter().collect())
.unwrap_or_default();
let principal = Principal {
principal_id: req.user_id.clone(),
subject: req.user_id.clone(),
user_id: req.user_id.clone(),
tenant_id: req.domain.clone(),
..Default::default()
};
let _admit = self.admit(&principal.tenant_id).await?;
let snap = self.current_snapshot().await?;
let audit_ctx = AuditContext::from_attributes(&attributes, "");
let mut results = std::collections::HashMap::new();
for check in req.checks {
let mut resource = ResourceRef {
resource_name: check.object.clone(),
message_type: check.object.clone(),
..Default::default()
};
enrich_resource(&mut resource);
let decision = self
.decide_with_snapshot(&snap, &principal, &resource, &check.action, "", &attributes)
.await;
self.write_decision_audit(&principal, &resource, &check.action, &decision, &audit_ctx)
.await;
results.insert(
format!("{}:{}", check.object, check.action),
decision.allowed,
);
}
Ok(Response::new(authz_pb::BatchCheckPermissionsResponse {
results,
}))
}
async fn update_role(
&self,
request: Request<authz_pb::UpdateRoleRequest>,
) -> Result<Response<authz_pb::UpdateRoleResponse>, Status> {
governance::guard_governed_role_mutation("UpdateRole")?;
let req = request.into_inner();
if req.role_id.trim().is_empty() {
return Err(Status::invalid_argument("role_id is required"));
}
let role_id = parse_uuid_field("role_id", &req.role_id)?;
if req.updated_by.trim().is_empty() {
return Err(Status::invalid_argument("updated_by is required"));
}
let pool = self.require_pool()?;
let role_model = self.roles_model();
let rel = role_model.relation.clone();
let projection = role_select_projection(&role_model);
let row = sqlx::query(&format!(
"UPDATE {rel} SET \
{name} = COALESCE(NULLIF($2, ''), {name}), \
{description} = COALESCE(NULLIF($3, ''), {description}), \
{is_active} = COALESCE($4, {is_active}) \
WHERE {role_id} = $1::UUID AND {deleted_at} IS NULL \
RETURNING {projection}",
name = role_model.q("name"),
description = role_model.q("description"),
is_active = role_model.q("is_active"),
role_id = role_model.q("role_id"),
deleted_at = role_model.q("deleted_at"),
))
.bind(role_id)
.bind(&req.name)
.bind(&req.description)
.bind(req.is_active)
.fetch_optional(pool)
.await
.map_err(|err| {
crate::runtime::executor_utils::sqlx_error_to_status("update role failed", &err)
})?;
let Some(row) = row else {
return Err(Status::not_found("role not found"));
};
let role = role_from_row(&row)?;
self.emit_event(
AuthEvent::new(
topics::ROLE_UPDATED,
role.role_id.clone(),
role.tenant_id.clone(),
serde_json::json!({
"role_id": role.role_id.clone(),
"role_code": role.role_code.clone(),
"tenant_id": role.tenant_id.clone(),
"updated_by": req.updated_by.clone(),
}),
)
.with_correlation(format!("role_update:{}", role.role_id))
.with_compliance(events::ComplianceEnvelope {
actor: if req.updated_by.trim().is_empty() {
role.role_id.clone()
} else {
req.updated_by.clone()
},
actor_project: role.project_id.clone(),
target_resource: format!("role:{}", role.role_id),
operation: "role_update".to_string(),
outcome: "success".to_string(),
reason_code: "role_updated".to_string(),
..events::ComplianceEnvelope::default()
}),
)
.await;
let _ = self
.bump_authz_revision(
&role.tenant_id,
&role.project_id,
authz_entity_pb::AuthzChangeType::Role,
"role-updated",
&req.updated_by,
)
.await;
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::UpdateRoleResponse {
role: Some(role),
}))
}
async fn delete_role(
&self,
request: Request<authz_pb::DeleteRoleRequest>,
) -> Result<Response<authz_pb::DeleteRoleResponse>, Status> {
governance::guard_governed_role_mutation("DeleteRole")?;
let req = request.into_inner();
if req.role_id.trim().is_empty() {
return Err(Status::invalid_argument("role_id is required"));
}
if req.deleted_by.trim().is_empty() {
return Err(Status::invalid_argument("deleted_by is required"));
}
let role_id = parse_uuid_field("role_id", &req.role_id)?;
let deleted_by = parse_uuid_field("deleted_by", &req.deleted_by)?;
let pool = self.require_pool()?;
let role_model = self.roles_model();
let role_rel = role_model.relation.clone();
let scope_row = sqlx::query(&format!(
"SELECT {tenant_id}::TEXT AS tenant, COALESCE({project_id}, '') AS project FROM {role_rel} WHERE {role_id} = $1::UUID",
tenant_id = role_model.q("tenant_id"),
project_id = role_model.q("project_id"),
role_id = role_model.q("role_id"),
))
.bind(role_id)
.fetch_optional(pool)
.await
.map_err(|err| Status::internal(format!("read role scope failed: {err}")))?;
let (role_tenant, role_project) = scope_row
.map(|r| {
(
r.try_get::<String, _>("tenant").unwrap_or_default(),
r.try_get::<String, _>("project").unwrap_or_default(),
)
})
.unwrap_or_default();
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition("native authz requires runtime-backed role persistence")
})?;
let mut assignments = std::collections::BTreeMap::new();
assignments.insert("deleted_at".to_string(), LogicalAssignment::ServerNow);
assignments.insert(
"deleted_by".to_string(),
LogicalAssignment::Set {
value: LogicalValue::String(deleted_by.to_string()),
},
);
assignments.insert(
"is_active".to_string(),
LogicalAssignment::Set {
value: LogicalValue::Bool(false),
},
);
let (affected, _) = runtime
.native_entity_update_for_service(
"authz",
&crate::RequestContext::default(),
LogicalUpdate {
message_type: "udb.core.authz.entity.v1.Role".to_string(),
filter: LogicalFilter::And(vec![
LogicalFilter::Comparison {
field: "role_id".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String(role_id.to_string()),
},
LogicalFilter::IsNull("deleted_at".to_string()),
]),
assignments,
return_fields: Vec::new(),
require_affected: false,
},
)
.await
.map_err(|err| Status::internal(format!("delete role failed: {err}")))?;
if affected > 0 {
runtime
.native_entity_delete_for_service(
"authz",
&crate::RequestContext::default(),
LogicalDelete {
message_type: "udb.core.authz.entity.v1.UserRole".to_string(),
filter: LogicalFilter::Comparison {
field: "role_id".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String(role_id.to_string()),
},
return_fields: Vec::new(),
},
)
.await
.map_err(|err| {
Status::internal(format!("delete role assignments failed: {err}"))
})?;
}
if affected > 0 && !role_tenant.trim().is_empty() {
let _ = self
.bump_authz_revision(
&role_tenant,
&role_project,
authz_entity_pb::AuthzChangeType::Role,
"role-deleted",
&req.deleted_by,
)
.await;
}
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::DeleteRoleResponse {
deleted: affected > 0,
}))
}
async fn get_policy_rule(
&self,
request: Request<authz_pb::GetPolicyRuleRequest>,
) -> Result<Response<authz_pb::GetPolicyRuleResponse>, Status> {
let req = request.into_inner();
if req.policy_id.trim().is_empty() {
return Err(Status::invalid_argument("policy_id is required"));
}
if let Some(pool) = &self.pg_pool {
let policy_id = parse_uuid_field("policy_id", &req.policy_id)?;
let policy_model = self.policies_model();
let rel = policy_model.relation.clone();
let projection = policy_rule_select_projection(&policy_model);
let row = sqlx::query(&format!(
"SELECT {projection} \
FROM {rel} \
WHERE {policy_id} = $1::UUID AND {deleted_at} IS NULL \
LIMIT 1",
policy_id = policy_model.q("policy_id"),
deleted_at = policy_model.q("deleted_at"),
))
.bind(policy_id)
.fetch_optional(pool)
.await
.map_err(|err| Status::internal(format!("get policy rule failed: {err}")))?;
return match row {
Some(row) => Ok(Response::new(authz_pb::GetPolicyRuleResponse {
policy: Some(policy_rule_from_row(&row)?),
})),
None => Err(Status::not_found("policy rule not found")),
};
}
let snap = self.current_snapshot().await?;
let policy = snap
.policies
.iter()
.find(|p| p.id == req.policy_id)
.map(policy_to_rule_pb);
Ok(Response::new(authz_pb::GetPolicyRuleResponse { policy }))
}
async fn list_policy_rules(
&self,
request: Request<authz_pb::ListPolicyRulesRequest>,
) -> Result<Response<authz_pb::ListPolicyRulesResponse>, Status> {
let req = request.into_inner();
if let Some(pool) = &self.pg_pool {
let policy_model = self.policies_model();
let rel = policy_model.relation.clone();
let projection = policy_rule_select_projection(&policy_model);
let rows = sqlx::query(&format!(
"SELECT {projection} \
FROM {rel} \
WHERE {deleted_at} IS NULL \
AND ($1 = '' OR {domain_col} = $1 OR {tenant_id} = $1 OR {project_id} = $1) \
AND ($2 = '' OR {subject} = $2) \
AND ($3 = '' OR {object_col} = $3) \
AND (NOT $4 OR {is_active} = TRUE) \
ORDER BY {is_active} DESC, {policy_id} ASC",
deleted_at = policy_model.q("deleted_at"),
domain_col = policy_model.q("domain"),
tenant_id = policy_model.q("tenant_id"),
project_id = policy_model.q("project_id"),
subject = policy_model.q("subject"),
object_col = policy_model.q("object"),
is_active = policy_model.q("is_active"),
policy_id = policy_model.q("policy_id"),
))
.bind(&req.domain)
.bind(&req.subject)
.bind(&req.object)
.bind(req.active_only)
.fetch_all(pool)
.await
.map_err(|err| Status::internal(format!("list policy rules failed: {err}")))?;
let all = rows
.iter()
.map(policy_rule_from_row)
.collect::<Result<Vec<_>, _>>()?;
let page = req.page.as_ref();
let page_number = page.map(|p| p.page).filter(|p| *p > 0).unwrap_or(1) as usize;
let page_size = page
.map(|p| p.page_size)
.filter(|s| *s > 0)
.unwrap_or(all.len().max(1) as i32) as usize;
let start = page_number.saturating_sub(1).saturating_mul(page_size);
let policies = all.into_iter().skip(start).take(page_size).collect();
return Ok(Response::new(authz_pb::ListPolicyRulesResponse {
page: Some(page_response(rows.len(), page)),
policies,
}));
}
let snap = self.current_snapshot().await?;
let policies: Vec<_> = snap
.policies
.iter()
.filter(|p| !req.active_only || p.enabled)
.filter(|p| {
req.domain.trim().is_empty() || p.tenant == req.domain || p.project == req.domain
})
.filter(|p| req.subject.trim().is_empty() || p.subject == req.subject)
.filter(|p| req.object.trim().is_empty() || p.resource == req.object)
.map(policy_to_rule_pb)
.collect();
Ok(Response::new(authz_pb::ListPolicyRulesResponse {
page: Some(page_response(policies.len(), req.page.as_ref())),
policies,
}))
}
async fn delete_policy_rule(
&self,
request: Request<authz_pb::DeletePolicyRuleRequest>,
) -> Result<Response<authz_pb::DeletePolicyRuleResponse>, Status> {
let req = request.into_inner();
if req.policy_id.trim().is_empty() {
return Err(Status::invalid_argument("policy_id is required"));
}
let mut deleted = false;
if self.pg_pool.is_some() {
let runtime = self.runtime.as_ref().ok_or_else(|| {
Status::failed_precondition(
"native authz requires runtime-backed policy persistence",
)
})?;
let policy_id = parse_uuid_field("policy_id", &req.policy_id)?;
let deleted_by = if req.deleted_by.trim().is_empty() {
None
} else {
Some(parse_uuid_field("deleted_by", &req.deleted_by)?)
};
let mut assignments = std::collections::BTreeMap::new();
assignments.insert("deleted_at".to_string(), LogicalAssignment::ServerNow);
assignments.insert(
"deleted_by".to_string(),
LogicalAssignment::Set {
value: match &deleted_by {
Some(u) => LogicalValue::String(u.to_string()),
None => LogicalValue::Null,
},
},
);
assignments.insert(
"is_active".to_string(),
LogicalAssignment::Set {
value: LogicalValue::Bool(false),
},
);
let op = LogicalUpdate {
message_type: "udb.core.authz.entity.v1.PolicyRule".to_string(),
filter: LogicalFilter::And(vec![
LogicalFilter::Comparison {
field: "policy_id".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String(policy_id.to_string()),
},
LogicalFilter::IsNull("deleted_at".to_string()),
]),
assignments,
return_fields: Vec::new(),
require_affected: false,
};
let context = crate::RequestContext::default();
let (affected, _) = runtime
.native_entity_update_for_service("authz", &context, op)
.await
.map_err(|err| Status::internal(format!("delete policy rule failed: {err}")))?;
deleted = affected > 0;
} else {
self.require_snapshot_fallback()?;
}
self.invalidate_snapshot_cache();
Ok(Response::new(authz_pb::DeletePolicyRuleResponse {
deleted,
}))
}
async fn get_native_access(
&self,
request: Request<authz_pb::NativeAccessRequest>,
) -> Result<Response<authz_pb::NativeAccessResponse>, Status> {
use crate::runtime::authz::native_access::NativeAccessConfig;
let req = request.into_inner();
let requested_scopes = req.requested_scopes.clone();
let mut principal = req
.principal
.as_ref()
.map(authz_principal_to_runtime)
.unwrap_or_default();
if principal.tenant_id.trim().is_empty() {
principal.tenant_id = req.tenant_id.clone();
}
if principal.project_id.trim().is_empty() {
principal.project_id = req.project_id.clone();
}
if crate::runtime::service::method_security::claim_context_present() {
let ctx = crate::runtime::service::method_security::current_claim_context();
crate::runtime::service::method_security::enforce_body_tenant_matches_claim(
&ctx,
&req.tenant_id,
&req.project_id,
)?;
if !ctx.is_cross_tenant_admin() {
let base = ctx.to_principal();
principal.subject = base.subject;
principal.user_id = base.user_id;
principal.tenant_id = base.tenant_id;
principal.project_id = base.project_id;
principal.scopes = base.scopes;
}
if !req.requested_scopes.is_empty() {
principal.scopes.retain(|scope| {
req.requested_scopes
.iter()
.any(|requested| requested == scope)
});
}
} else if principal.scopes.is_empty() {
principal.scopes = req.requested_scopes.clone();
}
if principal.subject.trim().is_empty() {
principal.subject = if !principal.user_id.trim().is_empty() {
principal.user_id.clone()
} else {
principal.principal_id.clone()
};
}
if principal.principal_id.trim().is_empty() {
principal.principal_id = principal.subject.clone();
}
if principal.tenant_id.trim().is_empty() {
return Err(Status::invalid_argument("tenant_id is required"));
}
let mut resource = req
.resource
.as_ref()
.map(resource_to_runtime)
.unwrap_or_default();
if !req.backend.trim().is_empty() && resource.backend.trim().is_empty() {
resource.backend = req.backend.clone();
}
enrich_resource(&mut resource);
let mut attributes: BTreeMap<String, String> = req.attributes.into_iter().collect();
if let Some(ctx) = req.context {
attributes.extend(ctx.attributes.into_iter());
}
let snap = self.current_snapshot().await?;
let decision = self
.decide_with_snapshot(
&snap,
&principal,
&resource,
&req.action,
&req.purpose,
&attributes,
)
.await;
let audit_ctx = AuditContext::from_attributes(&attributes, &req.purpose);
self.write_decision_audit(&principal, &resource, &req.action, &decision, &audit_ctx)
.await;
let grant = NativeAccessConfig::from_env()
.mint(
&principal,
&resource,
&req.action,
&req.purpose,
&decision,
now_unix() as i64,
)
.map(|g| authz_pb::NativeAccessGrant {
dsn: g.dsn,
role: g.role,
backend: g.backend,
database: g.database,
schema: g.schema,
session_variables: g.session_variables.into_iter().collect(),
expires_at_unix: g.expires_at_unix,
ttl_seconds: g.ttl_seconds,
});
let actor = if principal.subject.trim().is_empty() {
principal.principal_id.clone()
} else {
principal.subject.clone()
};
let correlation = if audit_ctx.correlation_id.trim().is_empty() {
format!("native_access:{}:{}", actor, resource.resource_name)
} else {
audit_ctx.correlation_id.clone()
};
let issued = decision.allowed && grant.is_some();
let (grant_topic, grant_op, grant_outcome, grant_reason) = if issued {
(
topics::NATIVE_ACCESS_GRANT_ISSUED,
"native_access_grant",
"allow",
"grant_issued".to_string(),
)
} else {
(
topics::NATIVE_ACCESS_GRANT_DENIED,
"native_access_grant",
"deny",
if decision.deny_reason.trim().is_empty() {
"no_grant_minted".to_string()
} else {
decision.deny_reason.clone()
},
)
};
self.emit_event(
AuthEvent::new(
grant_topic,
actor.clone(),
principal.tenant_id.clone(),
serde_json::json!({
"subject": actor.clone(),
"resource": resource.resource_name.clone(),
"backend": resource.backend.clone(),
"action": req.action.clone(),
"purpose": req.purpose.clone(),
"allowed": decision.allowed,
"granted": issued,
"requested_scopes": requested_scopes.clone(),
"effective_scopes": principal.scopes.clone(),
}),
)
.with_correlation(correlation)
.with_compliance(events::ComplianceEnvelope {
actor: actor.clone(),
actor_project: principal.project_id.clone(),
target_resource: resource.resource_name.clone(),
operation: grant_op.to_string(),
outcome: grant_outcome.to_string(),
reason_code: grant_reason,
decision_id: decision.decision_id.clone(),
policy_version: decision.policy_version.clone(),
relationship_version: decision.relationship_version.clone(),
..events::ComplianceEnvelope::default()
}),
)
.await;
Ok(Response::new(authz_pb::NativeAccessResponse {
decision: Some(decision_to_pb(&decision)),
grant,
}))
}
async fn get_policy_bundle(
&self,
request: Request<authz_pb::PolicyBundleRequest>,
) -> Result<Response<authz_pb::PolicyBundleResponse>, Status> {
use crate::runtime::authz::bundle::PolicyBundleConfig;
let req = request.into_inner();
if req.tenant_id.trim().is_empty() {
return Err(Status::invalid_argument(
"tenant_id is required for a policy bundle",
));
}
let cfg = PolicyBundleConfig::from_env();
if !cfg.enabled() {
return Err(Status::failed_precondition(
"policy bundle signing is not configured; set UDB_POLICY_BUNDLE_SECRET \
(or UDB_SESSION_HASH_SECRET)",
));
}
let snap = self.current_snapshot().await?;
let tenant = if req.tenant_id.trim().is_empty() {
req.domain.clone()
} else {
req.tenant_id.clone()
};
let now = now_unix() as i64;
let signed = PolicyEngine::bundle(snap.as_ref(), &cfg, &tenant, &req.project_id, now)
.await
.ok_or_else(|| Status::internal("failed to sign policy bundle"))?;
let bundle_actor = format!("tenant:{tenant}");
self.emit_event(
AuthEvent::new(
topics::POLICY_BUNDLE_ISSUED,
format!("{tenant}/{}", req.project_id),
tenant.clone(),
serde_json::json!({
"tenant_id": tenant.clone(),
"project_id": req.project_id.clone(),
"key_id": signed.key_id.clone(),
"policy_version": signed.policy_version.clone(),
"relationship_version": signed.relationship_version.clone(),
}),
)
.with_correlation(format!("policy_bundle:{tenant}/{}", req.project_id))
.with_compliance(events::ComplianceEnvelope {
actor: bundle_actor,
target_resource: format!("policy_bundle:{tenant}/{}", req.project_id),
operation: "policy_bundle_issue".to_string(),
outcome: "success".to_string(),
reason_code: "bundle_signed".to_string(),
policy_version: signed.policy_version.clone(),
relationship_version: signed.relationship_version.clone(),
..events::ComplianceEnvelope::default()
}),
)
.await;
Ok(Response::new(authz_pb::PolicyBundleResponse {
bundle: Some(authz_pb::SignedPolicyBundle {
bundle: signed.bundle,
signature: signed.signature,
key_id: signed.key_id,
algorithm: signed.algorithm,
policy_version: signed.policy_version,
relationship_version: signed.relationship_version,
issued_at_unix: signed.issued_at_unix,
expires_at_unix: signed.expires_at_unix,
ttl_seconds: signed.ttl_seconds,
}),
}))
}
async fn create_policy_draft(
&self,
request: Request<authz_pb::CreatePolicyDraftRequest>,
) -> Result<Response<authz_pb::PolicyDraftResponse>, Status> {
self.create_policy_draft_impl(request).await
}
async fn update_policy_draft(
&self,
request: Request<authz_pb::UpdatePolicyDraftRequest>,
) -> Result<Response<authz_pb::PolicyDraftResponse>, Status> {
self.update_policy_draft_impl(request).await
}
async fn diff_policy_draft(
&self,
request: Request<authz_pb::DiffPolicyDraftRequest>,
) -> Result<Response<authz_pb::DiffPolicyDraftResponse>, Status> {
self.diff_policy_draft_impl(request).await
}
async fn submit_policy_draft(
&self,
request: Request<authz_pb::SubmitPolicyDraftRequest>,
) -> Result<Response<authz_pb::PolicyDraftResponse>, Status> {
self.submit_policy_draft_impl(request).await
}
async fn approve_policy_draft(
&self,
request: Request<authz_pb::ApprovePolicyDraftRequest>,
) -> Result<Response<authz_pb::PolicyApprovalResponse>, Status> {
self.approve_policy_draft_impl(request).await
}
async fn reject_policy_draft(
&self,
request: Request<authz_pb::RejectPolicyDraftRequest>,
) -> Result<Response<authz_pb::PolicyApprovalResponse>, Status> {
self.reject_policy_draft_impl(request).await
}
async fn activate_policy_version(
&self,
request: Request<authz_pb::ActivatePolicyVersionRequest>,
) -> Result<Response<authz_pb::ActivationResponse>, Status> {
self.activate_policy_version_impl(request).await
}
async fn rollback_policy_version(
&self,
request: Request<authz_pb::RollbackPolicyVersionRequest>,
) -> Result<Response<authz_pb::ActivationResponse>, Status> {
self.rollback_policy_version_impl(request).await
}
async fn activate_canary(
&self,
request: Request<authz_pb::ActivateCanaryRequest>,
) -> Result<Response<authz_pb::CanaryResponse>, Status> {
self.activate_canary_impl(request).await
}
async fn promote_canary(
&self,
request: Request<authz_pb::PromoteCanaryRequest>,
) -> Result<Response<authz_pb::CanaryResponse>, Status> {
self.promote_canary_impl(request).await
}
async fn get_canary_status(
&self,
request: Request<authz_pb::GetCanaryStatusRequest>,
) -> Result<Response<authz_pb::GetCanaryStatusResponse>, Status> {
self.get_canary_status_impl(request).await
}
async fn list_policy_versions(
&self,
request: Request<authz_pb::ListPolicyVersionsRequest>,
) -> Result<Response<authz_pb::ListPolicyVersionsResponse>, Status> {
self.list_policy_versions_impl(request).await
}
async fn simulate_policy(
&self,
request: Request<authz_pb::SimulatePolicyRequest>,
) -> Result<Response<authz_pb::SimulatePolicyResponse>, Status> {
self.simulate_policy_impl(request).await
}
async fn explain_policy(
&self,
request: Request<authz_pb::ExplainPolicyRequest>,
) -> Result<Response<authz_pb::ExplainPolicyResponse>, Status> {
self.explain_policy_impl(request).await
}
async fn get_authz_revision(
&self,
request: Request<authz_pb::GetAuthzRevisionRequest>,
) -> Result<Response<authz_pb::GetAuthzRevisionResponse>, Status> {
self.get_authz_revision_impl(request).await
}
async fn invalidate_policy_bundles(
&self,
request: Request<authz_pb::InvalidatePolicyBundlesRequest>,
) -> Result<Response<authz_pb::InvalidatePolicyBundlesResponse>, Status> {
self.invalidate_policy_bundles_impl(request).await
}
async fn seed_builtin_roles(
&self,
request: Request<authz_pb::SeedBuiltinRolesRequest>,
) -> Result<Response<authz_pb::SeedBuiltinRolesResponse>, Status> {
self.seed_builtin_roles_impl(request).await
}
async fn migrate_legacy_policies(
&self,
request: Request<authz_pb::MigrateLegacyPoliciesRequest>,
) -> Result<Response<authz_pb::MigrateLegacyPoliciesResponse>, Status> {
self.migrate_legacy_policies_impl(request).await
}
}
fn policy_content_version(policies: &[AuthzPolicy]) -> String {
use std::hash::{Hash, Hasher};
let mut entries: Vec<String> = policies
.iter()
.map(|p| {
let mut scopes = p.required_scopes.clone();
scopes.sort();
format!(
"{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{:?}\u{1f}{}\u{1f}{:?}",
p.id, p.priority, p.enabled, p.effect.as_str(), p.tenant, p.project, p.subject,
p.role, p.action, p.resource, p.purpose, p.conditions, p.relationship, scopes,
)
})
.collect();
entries.sort();
let mut hasher = std::collections::hash_map::DefaultHasher::new();
entries.len().hash(&mut hasher);
for entry in &entries {
entry.hash(&mut hasher);
}
format!("pg-{}-{:016x}", policies.len(), hasher.finish())
}
fn tuple_content_version(tuples: &[RelationshipTuple]) -> String {
use std::hash::{Hash, Hasher};
let mut entries: Vec<String> = tuples
.iter()
.map(|t| {
format!(
"{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}",
t.subject, t.relation, t.object, t.tenant, t.project,
)
})
.collect();
entries.sort();
let mut hasher = std::collections::hash_map::DefaultHasher::new();
entries.len().hash(&mut hasher);
for entry in &entries {
entry.hash(&mut hasher);
}
format!("pg-{}-{:016x}", tuples.len(), hasher.finish())
}
#[cfg(test)]
mod version_tests {
use super::*;
fn policy(id: &str, effect: Effect) -> AuthzPolicy {
AuthzPolicy {
id: id.to_string(),
priority: 0,
enabled: true,
effect,
tenant: "acme".to_string(),
project: String::new(),
subject: "u1".to_string(),
role: String::new(),
action: "read".to_string(),
resource: "doc".to_string(),
purpose: String::new(),
relationship: String::new(),
conditions: Default::default(),
required_scopes: Vec::new(),
}
}
#[test]
fn version_changes_on_in_place_edit_same_count() {
let before = vec![policy("p1", Effect::Allow)];
let after = vec![policy("p1", Effect::Deny)];
assert_ne!(
policy_content_version(&before),
policy_content_version(&after),
"an in-place effect change must bump the version"
);
}
#[test]
fn version_is_order_independent_and_stable() {
let a = vec![policy("p1", Effect::Allow), policy("p2", Effect::Deny)];
let b = vec![policy("p2", Effect::Deny), policy("p1", Effect::Allow)];
assert_eq!(policy_content_version(&a), policy_content_version(&b));
assert_eq!(policy_content_version(&a), policy_content_version(&a));
}
}
#[cfg(test)]
mod fair_admission_tests {
use super::*;
#[tokio::test]
async fn admit_acquires_permit_when_channels_wired() {
let bare = AuthzServiceImpl::new(AuthzSnapshot::default());
let none = bare
.admit("tenant-a")
.await
.expect("bare admit must never reject");
assert!(
none.is_none(),
"no channels wired ⇒ admit must return None (no-op, callers still admitted)"
);
let svc = AuthzServiceImpl::new(AuthzSnapshot::default())
.with_channels(Some(ChannelManager::from_env()));
let permit = svc
.admit("tenant-a")
.await
.expect("Admin channel has capacity ⇒ permit must be granted");
assert!(
permit.is_some(),
"channels wired + capacity available ⇒ admit must hold a ChannelPermit"
);
}
}