use kube::{CustomResource, KubeSchema};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet};
use pgroles_core::bounds::*;
use pgroles_core::manifest::{
DefaultPrivilege, Ensure, Grant, Membership, ObjectType, Privilege, RoleRetirement,
SchemaBinding,
};
pub const VALID_SSL_MODES: &[&str] = &[
"disable",
"allow",
"prefer",
"require",
"verify-ca",
"verify-full",
];
#[derive(CustomResource, KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[kube(
group = "pgroles.io",
version = "v1alpha1",
kind = "PostgresPolicy",
namespaced,
status = "PostgresPolicyStatus",
shortname = "pgr",
category = "pgroles",
printcolumn = r#"{"name":"Ready","type":"string","jsonPath":".status.conditions[?(@.type==\"Ready\")].status"}"#,
printcolumn = r#"{"name":"Mode","type":"string","jsonPath":".spec.mode"}"#,
printcolumn = r#"{"name":"Recon","type":"string","jsonPath":".spec.reconciliation_mode","priority":1}"#,
printcolumn = r#"{"name":"Drift","type":"string","jsonPath":".status.conditions[?(@.type==\"Drifted\")].status"}"#,
printcolumn = r#"{"name":"Changes","type":"integer","jsonPath":".status.change_summary.total"}"#,
printcolumn = r#"{"name":"Last Reconcile","type":"date","jsonPath":".status.last_successful_reconcile_time"}"#,
printcolumn = r#"{"name":"Age","type":"date","jsonPath":".metadata.creationTimestamp"}"#
)]
pub struct PostgresPolicySpec {
pub connection: ConnectionSpec,
#[serde(default = "default_interval")]
pub interval: String,
#[serde(default)]
pub suspend: bool,
#[serde(default)]
pub mode: PolicyMode,
#[serde(default)]
pub reconciliation_mode: CrdReconciliationMode,
#[serde(default)]
#[schemars(length(min = 1, max = MAX_IDENTIFIER))]
pub default_owner: Option<String>,
#[serde(default)]
#[schemars(extend("maxProperties" = MAX_PROFILES))]
pub profiles: std::collections::HashMap<String, ProfileSpec>,
#[serde(default)]
#[schemars(length(max = MAX_SCHEMAS))]
#[x_kube(merge_strategy = ListMerge::Map(vec!["name".into()]))]
pub schemas: Vec<SchemaBinding>,
#[serde(default)]
#[schemars(length(max = MAX_ROLES))]
#[x_kube(merge_strategy = ListMerge::Map(vec!["name".into()]))]
pub roles: Vec<RoleSpec>,
#[serde(default)]
#[schemars(length(max = MAX_GRANTS))]
pub grants: Vec<Grant>,
#[serde(default)]
#[schemars(length(max = MAX_DEFAULT_PRIVILEGES))]
pub default_privileges: Vec<DefaultPrivilege>,
#[serde(default)]
#[schemars(length(max = MAX_MEMBERSHIPS))]
pub memberships: Vec<Membership>,
#[serde(default)]
#[schemars(length(max = MAX_RETIREMENTS))]
#[x_kube(merge_strategy = ListMerge::Map(vec!["role".into()]))]
pub retirements: Vec<RoleRetirement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub approval: Option<ApprovalMode>,
}
fn default_interval() -> String {
"5m".to_string()
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "lowercase")]
pub enum PolicyMode {
#[default]
Apply,
Observe,
Plan,
}
impl PolicyMode {
pub fn never_executes(self) -> bool {
matches!(self, PolicyMode::Observe | PolicyMode::Plan)
}
pub fn is_deprecated_spelling(self) -> bool {
matches!(self, PolicyMode::Plan)
}
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum CrdReconciliationMode {
#[default]
Authoritative,
Additive,
Adopt,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
pub enum ApprovalMode {
#[serde(rename = "manual")]
Manual,
#[serde(rename = "auto")]
Auto,
}
impl PostgresPolicySpec {
pub fn effective_approval(&self) -> ApprovalMode {
match &self.approval {
Some(mode) => mode.clone(),
None => match self.mode {
PolicyMode::Apply => ApprovalMode::Auto,
PolicyMode::Observe | PolicyMode::Plan => ApprovalMode::Manual,
},
}
}
}
pub const REQUESTED_RECONCILE_ANNOTATION: &str = "reconcile.pgroles.io/requestedAt";
pub const LABEL_POLICY: &str = "pgroles.io/policy";
pub const LABEL_DATABASE_IDENTITY: &str = "pgroles.io/database-identity";
pub const LABEL_PLAN: &str = "pgroles.io/plan";
pub const LABEL_CANDIDATE: &str = "pgroles.io/candidate";
pub const LABEL_KEEP: &str = "pgroles.io/keep";
pub fn is_retention_exempt<K: kube::Resource>(resource: &K) -> bool {
resource
.meta()
.labels
.as_ref()
.and_then(|labels| labels.get(LABEL_KEEP))
.is_some_and(|value| value == "true")
}
pub const LABEL_ACCESS_POLICY_UID: &str = "pgroles.io/access-policy-uid";
pub const LABEL_TARGET_POLICY_UID: &str = "pgroles.io/target-policy-uid";
impl From<CrdReconciliationMode> for pgroles_core::diff::ReconciliationMode {
fn from(crd: CrdReconciliationMode) -> Self {
match crd {
CrdReconciliationMode::Authoritative => {
pgroles_core::diff::ReconciliationMode::Authoritative
}
CrdReconciliationMode::Additive => pgroles_core::diff::ReconciliationMode::Additive,
CrdReconciliationMode::Adopt => pgroles_core::diff::ReconciliationMode::Adopt,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct ConnectionSpec {
#[serde(default)]
pub secret_ref: Option<SecretReference>,
#[serde(default)]
pub secret_key: Option<String>,
#[serde(default)]
pub params: Option<ConnectionParams>,
#[serde(default)]
pub require_physical_identity: Option<bool>,
}
impl ConnectionSpec {
pub fn requires_physical_identity(&self) -> bool {
self.require_physical_identity.unwrap_or(false)
}
pub fn effective_secret_key(&self) -> &str {
self.secret_key.as_deref().unwrap_or("DATABASE_URL")
}
pub fn collect_secret_names(&self, names: &mut BTreeSet<String>) {
if let Some(ref secret_ref) = self.secret_ref {
names.insert(secret_ref.name.clone());
}
if let Some(ref params) = self.params {
for sel in [
¶ms.host_secret,
¶ms.port_secret,
¶ms.dbname_secret,
¶ms.username_secret,
¶ms.password_secret,
¶ms.ssl_mode_secret,
]
.into_iter()
.flatten()
{
names.insert(sel.name.clone());
}
}
}
pub fn identity_key(&self) -> String {
if let Some(ref secret_ref) = self.secret_ref {
format!("{}/{}", secret_ref.name, self.effective_secret_key())
} else if let Some(ref params) = self.params {
let port_part = params
.port
.as_ref()
.map(|p| format!("literal={p}"))
.or_else(|| {
params
.port_secret
.as_ref()
.map(|s| format!("secret={}\0{}", s.name, s.key))
})
.unwrap_or_else(|| "5432".to_string());
format!(
"params\0{}\0{}\0{}",
field_identity_repr(¶ms.host, ¶ms.host_secret),
field_identity_repr(¶ms.dbname, ¶ms.dbname_secret),
port_part,
)
} else {
"invalid-connection".to_string()
}
}
pub fn cache_key(&self, namespace: &str) -> String {
if let Some(ref params) = self.params {
let user_part = field_identity_repr(¶ms.username, ¶ms.username_secret);
let pass_part = field_identity_repr(¶ms.password, ¶ms.password_secret);
let auth_part = params
.auth
.as_ref()
.map(ConnectionAuth::cache_key)
.unwrap_or_default();
let ssl_part = params
.ssl_mode
.as_ref()
.map(|v| format!("literal={v}"))
.or_else(|| {
params
.ssl_mode_secret
.as_ref()
.map(|s| format!("secret={}\0{}", s.name, s.key))
})
.unwrap_or_default();
let role_part = params.set_role.as_deref().unwrap_or("");
format!(
"{namespace}/{}\0user={user_part}\0pass={pass_part}\0auth={auth_part}\0ssl={ssl_part}\0role={role_part}",
self.identity_key()
)
} else {
format!("{namespace}/{}", self.identity_key())
}
}
}
fn field_identity_repr(literal: &Option<String>, secret: &Option<SecretKeySelector>) -> String {
if let Some(value) = literal {
format!("literal={value}")
} else if let Some(sel) = secret {
format!("secret={}\0{}", sel.name, sel.key)
} else {
String::new()
}
}
pub const DEFAULT_GCP_CLOUD_SQL_LOGIN_SCOPE: &str =
"https://www.googleapis.com/auth/sqlservice.login";
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct ConnectionParams {
#[serde(default)]
pub host: Option<String>,
#[serde(default)]
pub host_secret: Option<SecretKeySelector>,
#[serde(default)]
pub port: Option<u16>,
#[serde(default)]
pub port_secret: Option<SecretKeySelector>,
#[serde(default)]
pub dbname: Option<String>,
#[serde(default)]
pub dbname_secret: Option<SecretKeySelector>,
#[serde(default)]
pub username: Option<String>,
#[serde(default)]
pub username_secret: Option<SecretKeySelector>,
#[serde(default)]
pub password: Option<String>,
#[serde(default)]
pub password_secret: Option<SecretKeySelector>,
#[serde(default)]
pub auth: Option<ConnectionAuth>,
#[serde(default)]
pub ssl_mode: Option<String>,
#[serde(default)]
pub ssl_mode_secret: Option<SecretKeySelector>,
#[serde(default)]
#[schemars(regex(pattern = SET_ROLE_PATTERN))]
pub set_role: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(tag = "type")]
pub enum ConnectionAuth {
#[serde(rename = "gcp_workload_identity", rename_all = "camelCase")]
GcpWorkloadIdentity {
#[serde(default)]
impersonate_service_account: Option<String>,
#[serde(default)]
scope: Option<String>,
},
}
impl ConnectionAuth {
pub fn gcp_scope(&self) -> &str {
match self {
Self::GcpWorkloadIdentity { scope, .. } => scope
.as_deref()
.unwrap_or(DEFAULT_GCP_CLOUD_SQL_LOGIN_SCOPE),
}
}
pub fn gcp_impersonate_service_account(&self) -> Option<&str> {
match self {
Self::GcpWorkloadIdentity {
impersonate_service_account,
..
} => impersonate_service_account.as_deref(),
}
}
fn cache_key(&self) -> String {
match self {
Self::GcpWorkloadIdentity {
impersonate_service_account,
scope,
} => format!(
"gcp_workload_identity\0impersonate={}\0scope={}",
impersonate_service_account.as_deref().unwrap_or_default(),
scope
.as_deref()
.unwrap_or(DEFAULT_GCP_CLOUD_SQL_LOGIN_SCOPE)
),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct SecretKeySelector {
#[schemars(length(min = 1, max = MAX_K8S_NAME))]
pub name: String,
#[schemars(length(min = 1, max = MAX_SECRET_KEY))]
pub key: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct SecretReference {
#[schemars(length(min = 1, max = MAX_K8S_NAME))]
pub name: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct ProfileSpec {
#[serde(default)]
pub login: Option<bool>,
#[serde(default)]
pub inherit: Option<bool>,
#[serde(default)]
#[schemars(length(max = MAX_PROFILE_GRANTS))]
pub grants: Vec<ProfileGrantSpec>,
#[serde(default)]
#[schemars(length(max = MAX_PROFILE_DEFAULT_PRIVILEGES))]
pub default_privileges: Vec<DefaultPrivilegeGrantSpec>,
#[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
#[schemars(extend("maxProperties" = MAX_CONFIG_ENTRIES))]
pub config: std::collections::BTreeMap<String, pgroles_core::manifest::ConfigValue>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct ProfileGrantSpec {
#[schemars(length(min = 1, max = MAX_PRIVILEGES))]
pub privileges: Vec<Privilege>,
#[serde(alias = "on")]
pub object: ProfileObjectTargetSpec,
#[serde(default, skip_serializing_if = "Ensure::is_present")]
pub ensure: Ensure,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct ProfileObjectTargetSpec {
#[serde(rename = "type")]
pub object_type: ObjectType,
#[serde(default)]
#[schemars(length(min = 1, max = MAX_OBJECT_NAME))]
pub name: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct DefaultPrivilegeGrantSpec {
#[serde(default)]
#[schemars(length(min = 1, max = MAX_IDENTIFIER))]
pub role: Option<String>,
#[schemars(length(min = 1, max = MAX_PRIVILEGES))]
pub privileges: Vec<Privilege>,
pub on_type: ObjectType,
#[serde(default, skip_serializing_if = "Ensure::is_present")]
pub ensure: Ensure,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct RoleSpec {
#[schemars(length(min = 1, max = MAX_IDENTIFIER))]
pub name: String,
#[serde(default)]
pub external: bool,
#[serde(default)]
pub login: Option<bool>,
#[serde(default)]
pub superuser: Option<bool>,
#[serde(default)]
pub createdb: Option<bool>,
#[serde(default)]
pub createrole: Option<bool>,
#[serde(default)]
pub inherit: Option<bool>,
#[serde(default)]
pub replication: Option<bool>,
#[serde(default)]
pub bypassrls: Option<bool>,
#[serde(default)]
pub connection_limit: Option<i32>,
#[serde(default)]
#[schemars(length(max = MAX_OBJECT_NAME))]
pub comment: Option<String>,
#[serde(default)]
pub password: Option<PasswordSpec>,
#[serde(default)]
#[schemars(length(max = MAX_TIMESTAMP))]
pub password_valid_until: Option<String>,
#[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
#[schemars(extend("maxProperties" = MAX_CONFIG_ENTRIES))]
pub config: std::collections::BTreeMap<String, pgroles_core::manifest::ConfigValue>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct PasswordSpec {
#[serde(default)]
pub secret_ref: Option<SecretReference>,
#[serde(default)]
#[schemars(length(min = 1, max = MAX_SECRET_KEY))]
pub secret_key: Option<String>,
#[serde(default)]
pub generate: Option<GeneratePasswordSpec>,
}
impl PasswordSpec {
pub fn is_secret_ref(&self) -> bool {
self.secret_ref.is_some()
}
pub fn is_generate(&self) -> bool {
self.generate.is_some()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct GeneratePasswordSpec {
#[serde(default)]
pub length: Option<u32>,
#[serde(default)]
#[schemars(length(min = 1, max = MAX_K8S_NAME))]
pub secret_name: Option<String>,
#[serde(default)]
#[schemars(length(min = 1, max = MAX_SECRET_KEY))]
pub secret_key: Option<String>,
}
#[derive(Debug, Clone, thiserror::Error)]
pub enum PasswordValidationError {
#[error("role \"{role}\" has a password but login is not enabled")]
PasswordWithoutLogin { role: String },
#[error("role \"{role}\" password must set exactly one of secretRef or generate")]
InvalidPasswordMode { role: String },
#[error("role \"{role}\" password.generate.length must be between {min} and {max}")]
InvalidGeneratedLength { role: String, min: u32, max: u32 },
#[error(
"role \"{role}\" password.generate.secretName \"{name}\" is not a valid Kubernetes Secret name"
)]
InvalidGeneratedSecretName { role: String, name: String },
#[error("role \"{role}\" password {field} \"{key}\" is not a valid Kubernetes Secret data key")]
InvalidSecretKey {
role: String,
field: &'static str,
key: String,
},
#[error(
"role \"{role}\" password.generate.secretKey \"{key}\" is reserved for the SCRAM verifier"
)]
ReservedGeneratedSecretKey { role: String, key: String },
}
#[derive(Debug, Clone, thiserror::Error)]
pub enum ConnectionValidationError {
#[error("connection: exactly one of secretRef or params must be set, but both were provided")]
BothModesSet,
#[error("connection: exactly one of secretRef or params must be set, but neither was provided")]
NeitherModeSet,
#[error("connection.params.{field}: secret {detail}")]
EmptySecretKeyRef { field: String, detail: String },
#[error(
"connection.params.sslMode: \"{value}\" is not valid (expected one of: disable, allow, prefer, require, verify-ca, verify-full)"
)]
InvalidSslMode { value: String },
#[error("connection.params.{field}: literal value must not be empty or whitespace-only")]
EmptyLiteral { field: String },
#[error("connection.params: exactly one of {field} or {field}Secret must be set")]
NeitherFieldSet { field: String },
#[error(
"connection.params: only one of {field} or {field}Secret may be set, but both were provided"
)]
BothFieldsSet { field: String },
#[error("connection.params.auth: {field} must not be empty or whitespace-only")]
EmptyAuthField { field: String },
#[error("connection.params: password/passwordSecret are mutually exclusive with auth")]
AuthWithPassword,
#[error(
"connection.params.setRole: \"{value}\" is not a valid PostgreSQL role identifier (must match {pattern})",
pattern = SET_ROLE_PATTERN,
)]
InvalidRoleName { value: String },
}
pub(crate) const SET_ROLE_PATTERN: &str = "^[A-Za-z_][A-Za-z0-9_$-]*$";
pub(crate) fn is_valid_set_role_identifier(s: &str) -> bool {
let mut bytes = s.bytes();
match bytes.next() {
Some(b) if b.is_ascii_alphabetic() || b == b'_' => {}
_ => return false,
}
bytes.all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'$' || b == b'-')
}
fn is_valid_secret_name(name: &str) -> bool {
crate::k8s_names::is_valid_resource_name(name)
}
fn is_valid_secret_key(key: &str) -> bool {
!key.is_empty()
&& key
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.'))
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
pub struct PostgresPolicyStatus {
#[serde(default)]
pub conditions: Vec<PolicyCondition>,
#[serde(default)]
pub observed_generation: Option<i64>,
#[serde(default)]
pub last_attempted_generation: Option<i64>,
#[serde(default)]
pub last_successful_reconcile_time: Option<String>,
#[serde(default, rename = "lastHandledReconcileAt")]
pub last_handled_reconcile_at: Option<String>,
#[serde(default)]
pub change_summary: Option<ChangeSummary>,
#[serde(default)]
pub last_reconcile_mode: Option<PolicyMode>,
#[serde(default)]
pub managed_database_identity: Option<String>,
#[serde(default)]
pub owned_roles: Vec<String>,
#[serde(default)]
pub owned_schemas: Vec<String>,
#[serde(default)]
pub last_error: Option<String>,
#[serde(default)]
pub applied_password_source_versions: BTreeMap<String, String>,
#[serde(default)]
pub transient_failure_count: i32,
#[serde(default)]
pub current_plan_ref: Option<PlanReference>,
#[serde(default)]
pub content_digest: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct PolicyCondition {
#[serde(rename = "type")]
pub condition_type: String,
pub status: String,
#[serde(default)]
pub reason: Option<String>,
#[serde(default)]
pub message: Option<String>,
#[serde(default)]
pub last_transition_time: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct PlanReference {
pub name: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
#[serde(default)]
pub struct ChangeSummary {
#[serde(default)]
pub roles_created: i32,
#[serde(default)]
pub roles_altered: i32,
#[serde(default)]
pub schemas_created: i32,
#[serde(default)]
pub schema_owners_altered: i32,
#[serde(default)]
pub roles_dropped: i32,
#[serde(default)]
pub sessions_terminated: i32,
#[serde(default)]
pub grants_added: i32,
#[serde(default)]
pub grants_revoked: i32,
#[serde(default)]
pub default_privileges_set: i32,
#[serde(default)]
pub default_privileges_revoked: i32,
#[serde(default)]
pub members_added: i32,
#[serde(default)]
pub members_removed: i32,
#[serde(default)]
pub passwords_set: i32,
#[serde(default)]
pub total: i32,
}
#[derive(CustomResource, KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[kube(
group = "pgroles.io",
version = "v1alpha1",
kind = "PostgresPolicyPlan",
namespaced,
status = "PostgresPolicyPlanStatus",
shortname = "pgplan",
category = "pgroles",
printcolumn = r#"{"name":"Policy","type":"string","jsonPath":".spec.policyRef.name"}"#,
printcolumn = r#"{"name":"Mode","type":"string","jsonPath":".spec.reconciliationMode"}"#,
printcolumn = r#"{"name":"Approved","type":"string","jsonPath":".status.conditions[?(@.type==\"Approved\")].status"}"#,
printcolumn = r#"{"name":"Changes","type":"integer","jsonPath":".status.changeSummary.total"}"#,
printcolumn = r#"{"name":"SQL Stmts","type":"integer","jsonPath":".status.sqlStatements","priority":1}"#,
printcolumn = r#"{"name":"Phase","type":"string","jsonPath":".status.phase"}"#,
printcolumn = r#"{"name":"SQL","type":"string","jsonPath":".status.sqlRef.name","priority":1}"#,
printcolumn = r#"{"name":"Digest","type":"string","jsonPath":".status.changeDigest","priority":1}"#,
printcolumn = r#"{"name":"Hash","type":"string","jsonPath":".status.sqlHash","priority":1}"#,
printcolumn = r#"{"name":"Age","type":"date","jsonPath":".metadata.creationTimestamp"}"#
)]
#[x_kube(validation = Rule::new(
"!has(oldSelf.origin) || (has(self.origin) && self.origin == oldSelf.origin)"
)
.message("plan origin is immutable once set"))]
#[serde(rename_all = "camelCase")]
pub struct PostgresPolicyPlanSpec {
pub policy_ref: PolicyPlanRef,
pub policy_generation: i64,
pub reconciliation_mode: CrdReconciliationMode,
#[serde(default)]
pub owned_roles: Vec<String>,
#[serde(default)]
pub owned_schemas: Vec<String>,
pub managed_database_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub origin: Option<PlanOrigin>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub scope: Option<PlanScope>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct PlanOrigin {
#[schemars(length(max = 63))]
pub kind: String,
#[schemars(length(max = 253))]
pub name: String,
#[schemars(length(max = 63))]
pub uid: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 128))]
pub content_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub content_digest_encoding: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 63))]
pub policy_uid: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 128))]
pub base_content_digest: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct PlanScope {
pub kind: String,
pub operation: ScopedPlanOperation,
pub bundle_hash: String,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
pub enum ScopedPlanOperation {
Activate,
Revoke,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct PolicyPlanRef {
pub name: String,
}
#[derive(KubeSchema, Debug, Clone, Default, Serialize, Deserialize)]
#[x_kube(
validation = Rule::new(
"!(self.conditions.exists(c, c.type == 'Approved' && c.status == 'True') && self.conditions.exists(c, c.type == 'Denied' && c.status == 'True'))"
).message("Approved=True and Denied=True are mutually exclusive"),
// Compares the *decision types* that are true, not whole condition
// objects. The operator rewrites this array on every status write, so
// comparing objects would reject a later write that only refreshed a
// timestamp or message — while the decision itself was unchanged. What
// must not change is which decisions are true.
validation = Rule::new(
"oldSelf.conditions.filter(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True').map(c, c.type) == self.conditions.filter(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True').map(c, c.type) || oldSelf.conditions.filter(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True').size() == 0"
).message("plan decisions are terminal"),
validation = Rule::new(
"!has(oldSelf.decidedBy) || (has(self.decidedBy) && self.decidedBy == oldSelf.decidedBy)"
).message("decision identity is write-once"),
validation = Rule::new(
"self.conditions.exists(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True') == has(self.decidedBy)"
).message("a terminal plan decision and decidedBy identity must be recorded together"),
// The approval-binding fields are write-once. Execution gates on the
// recorded digest matching freshly recomputed effects, and the identity
// downgrade check reads these values — so anyone holding plans/status
// patch (every plan approver does) who could rewrite them could point an
// existing approval at different effects or a different server. The
// operator writes each of these exactly once, when it first materialises
// the plan's status; a legitimate rewrite repeats the same value, which
// equality admits.
validation = Rule::new(
"!has(oldSelf.changeDigest) || (has(self.changeDigest) && self.changeDigest == oldSelf.changeDigest)"
).message("changeDigest is write-once"),
validation = Rule::new(
"!has(oldSelf.changeDigestEncoding) || (has(self.changeDigestEncoding) && self.changeDigestEncoding == oldSelf.changeDigestEncoding)"
).message("changeDigestEncoding is write-once"),
validation = Rule::new(
"!has(oldSelf.targetPhysicalIdentity) || (has(self.targetPhysicalIdentity) && self.targetPhysicalIdentity == oldSelf.targetPhysicalIdentity)"
).message("targetPhysicalIdentity is write-once"),
validation = Rule::new(
"!has(oldSelf.targetLogicalFingerprint) || (has(self.targetLogicalFingerprint) && self.targetLogicalFingerprint == oldSelf.targetLogicalFingerprint)"
).message("targetLogicalFingerprint is write-once"),
validation = Rule::new(
"!has(oldSelf.physicalIdentityAvailable) || (has(self.physicalIdentityAvailable) && self.physicalIdentityAvailable == oldSelf.physicalIdentityAvailable)"
).message("physicalIdentityAvailable is write-once")
)]
#[serde(rename_all = "camelCase")]
pub struct PostgresPolicyPlanStatus {
#[serde(default)]
pub phase: PlanPhase,
#[serde(default)]
#[schemars(length(max = 16))]
pub conditions: Vec<PolicyCondition>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub decided_by: Option<DecisionActor>,
#[serde(default)]
pub change_summary: Option<ChangeSummary>,
#[serde(default)]
pub sql_ref: Option<SqlRef>,
#[serde(default)]
pub sql_inline: Option<String>,
#[serde(default)]
pub sql_truncated: bool,
#[serde(default)]
pub computed_at: Option<String>,
#[serde(default)]
pub applied_at: Option<String>,
#[serde(default)]
pub last_error: Option<String>,
#[serde(default)]
pub sql_hash: Option<String>,
#[serde(default)]
pub change_digest: Option<String>,
#[serde(default)]
pub change_digest_encoding: Option<String>,
#[serde(default)]
pub target_physical_identity: Option<String>,
#[serde(default)]
pub target_logical_fingerprint: Option<String>,
#[serde(default)]
pub physical_identity_available: Option<bool>,
#[serde(default)]
pub revalidated_generation: Option<i64>,
#[serde(default)]
pub revalidated_at: Option<String>,
#[serde(default)]
pub applying_since: Option<String>,
#[serde(default)]
pub failed_at: Option<String>,
#[serde(default)]
pub sql_statements: Option<i64>,
#[serde(default)]
pub redacted_sql_hash: Option<String>,
#[serde(default)]
pub sql_original_bytes: Option<i64>,
#[serde(default)]
pub sql_stored_bytes: Option<i64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct SqlRef {
pub name: String,
pub key: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub compression: Option<SqlCompression>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum SqlCompression {
Gzip,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
pub enum PlanPhase {
#[default]
Pending,
Approved,
Applying,
Applied,
Failed,
Superseded,
Rejected,
}
impl std::fmt::Display for PlanPhase {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
PlanPhase::Pending => write!(f, "Pending"),
PlanPhase::Approved => write!(f, "Approved"),
PlanPhase::Applying => write!(f, "Applying"),
PlanPhase::Applied => write!(f, "Applied"),
PlanPhase::Failed => write!(f, "Failed"),
PlanPhase::Superseded => write!(f, "Superseded"),
PlanPhase::Rejected => write!(f, "Rejected"),
}
}
}
pub const EPHEMERAL_BUNDLE_ENCODING_V1: &str = "pgroles.io/ephemeral-membership-bundle-v1";
pub const EPHEMERAL_MEMBERSHIP_SEMANTICS_V1: &str =
"postgres-membership-v1-admin-false-set-server-default";
#[derive(CustomResource, KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[kube(
group = "pgroles.io",
version = "v1alpha1",
kind = "EphemeralAccessPolicy",
namespaced,
status = "EphemeralAccessPolicyStatus",
shortname = "pgeap",
category = "pgroles",
printcolumn = r#"{"name":"Target","type":"string","jsonPath":".spec.postgresPolicyRef.name"}"#,
printcolumn = r#"{"name":"Accepted","type":"string","jsonPath":".status.conditions[?(@.type==\"Accepted\")].status"}"#,
printcolumn = r#"{"name":"Suspended","type":"boolean","jsonPath":".spec.suspend"}"#,
printcolumn = r#"{"name":"Age","type":"date","jsonPath":".metadata.creationTimestamp"}"#
)]
#[serde(rename_all = "camelCase")]
pub struct EphemeralAccessPolicySpec {
pub postgres_policy_ref: LocalObjectReference,
#[schemars(length(min = 1, max = 32))]
pub memberships: Vec<EphemeralMembership>,
#[schemars(length(max = 64), regex(pattern = r"^([0-9]+[smh])+$"))]
pub maximum_duration: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64), regex(pattern = r"^([0-9]+[smh])+$"))]
pub default_duration: Option<String>,
#[serde(rename = "pendingRequestTTL", default = "default_pending_request_ttl")]
#[schemars(length(max = 64), regex(pattern = r"^([0-9]+[smh])+$"))]
pub pending_request_ttl: String,
pub justification: EphemeralJustificationPolicy,
pub approval: EphemeralApprovalPolicy,
#[serde(default)]
pub suspend: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 128))]
pub display_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 2048))]
pub description: Option<String>,
}
fn default_pending_request_ttl() -> String {
"15m".to_string()
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct LocalObjectReference {
#[schemars(length(min = 1, max = 253))]
pub name: String,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct EphemeralMembership {
#[schemars(length(min = 1, max = 63))]
pub role: String,
pub inherit: bool,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize)]
pub struct EphemeralJustificationPolicy {
pub required: bool,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize)]
pub struct EphemeralApprovalPolicy {
pub mode: EphemeralApprovalMode,
}
#[derive(JsonSchema, Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
pub enum EphemeralApprovalMode {
Automatic,
Required,
}
#[derive(KubeSchema, Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct EphemeralAccessPolicyStatus {
#[serde(default)]
pub observed_generation: Option<i64>,
#[serde(default)]
#[schemars(length(max = 16))]
pub conditions: Vec<EphemeralAccessCondition>,
#[serde(default)]
#[schemars(length(max = 32), inner(length(min = 1, max = 63)))]
pub resolved_roles: Vec<String>,
}
#[derive(CustomResource, KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[kube(
group = "pgroles.io",
version = "v1alpha1",
kind = "EphemeralAccessRequest",
namespaced,
status = "EphemeralAccessRequestStatus",
shortname = "pgear",
category = "pgroles",
printcolumn = r#"{"name":"Access Policy","type":"string","jsonPath":".spec.accessPolicyRef.name"}"#,
printcolumn = r#"{"name":"Subject","type":"string","jsonPath":".spec.subject.role"}"#,
printcolumn = r#"{"name":"Phase","type":"string","jsonPath":".status.phase"}"#,
printcolumn = r#"{"name":"Expires","type":"date","jsonPath":".status.expiresAt"}"#,
printcolumn = r#"{"name":"Age","type":"date","jsonPath":".metadata.creationTimestamp"}"#
)]
#[x_kube(
validation = Rule::new("self == oldSelf").message("request spec is immutable")
)]
#[serde(rename_all = "camelCase")]
pub struct EphemeralAccessRequestSpec {
pub access_policy_ref: LocalObjectReference,
pub subject: EphemeralAccessSubject,
pub requested_by: DecisionActor,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64), regex(pattern = r"^([0-9]+[smh])+$"))]
pub requested_duration: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 2048))]
pub justification: Option<String>,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize)]
pub struct EphemeralAccessSubject {
#[schemars(length(min = 1, max = 63))]
pub role: String,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct DecisionActor {
#[schemars(length(min = 1, max = 512))]
pub username: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 128))]
pub uid: Option<String>,
#[serde(default)]
#[schemars(length(max = 64), inner(length(min = 1, max = 256)))]
pub groups: Vec<String>,
}
#[derive(KubeSchema, Debug, Clone, Default, Serialize, Deserialize)]
#[x_kube(
validation = Rule::new(
"!has(oldSelf.resolvedAccess) || (has(self.resolvedAccess) && self.resolvedAccess == oldSelf.resolvedAccess)"
).message("resolvedAccess is write-once"),
validation = Rule::new(
"!(self.conditions.exists(c, c.type == 'Approved' && c.status == 'True') && self.conditions.exists(c, c.type == 'Denied' && c.status == 'True'))"
).message("Approved=True and Denied=True are mutually exclusive"),
validation = Rule::new(
"oldSelf.conditions.filter(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True').size() == 0 || self.conditions.filter(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True') == oldSelf.conditions.filter(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True')"
).message("approval decisions are terminal"),
validation = Rule::new(
"!has(oldSelf.decidedBy) || (has(self.decidedBy) && self.decidedBy == oldSelf.decidedBy)"
).message("decision identity is write-once"),
validation = Rule::new(
"self.conditions.exists(c, (c.type == 'Approved' || c.type == 'Denied') && c.status == 'True') == has(self.decidedBy)"
).message("a terminal approval decision and decidedBy identity must be recorded together"),
validation = Rule::new(
"self.conditions.all(c, c.type in ['Approved', 'Denied', 'Resolved', 'Ready', 'Applied'])"
).message("request conditions must use a declared lifecycle or decision type")
)]
#[serde(rename_all = "camelCase")]
pub struct EphemeralAccessRequestStatus {
#[serde(default)]
pub phase: EphemeralAccessRequestPhase,
#[serde(default)]
#[schemars(length(max = 8))]
pub conditions: Vec<EphemeralAccessCondition>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub resolved_access: Option<ResolvedEphemeralAccess>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub decided_by: Option<DecisionActor>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub approval_expires_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub activated_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub expires_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub ended_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 4096))]
pub last_error: Option<String>,
#[serde(default)]
#[schemars(length(max = 32))]
pub retained_memberships: Vec<ResolvedEphemeralMembership>,
}
#[derive(JsonSchema, Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
pub enum EphemeralAccessRequestPhase {
#[default]
Pending,
PendingApproval,
Applying,
Active,
Revoking,
Ended,
Revoked,
Cancelled,
Denied,
ApprovalExpired,
Failed,
}
impl std::fmt::Display for EphemeralAccessRequestPhase {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{self:?}")
}
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct EphemeralAccessCondition {
#[serde(rename = "type")]
#[schemars(length(min = 1, max = 32))]
pub condition_type: String,
#[schemars(length(min = 1, max = 16))]
pub status: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 128))]
pub reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 2048))]
pub message: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub last_transition_time: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 71))]
pub bundle_hash: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 64))]
pub granted_duration: Option<String>,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct ResolvedEphemeralAccess {
#[schemars(length(max = 128))]
pub access_policy_uid: String,
pub access_policy_generation: i64,
#[schemars(length(max = 128))]
pub target_policy_uid: String,
pub target_policy_generation: i64,
#[schemars(length(max = 71))]
pub target_database_fingerprint: String,
#[schemars(length(max = 64))]
pub granted_duration: String,
#[schemars(length(max = 128))]
pub bundle_encoding: String,
#[schemars(length(max = 71))]
pub bundle_hash: String,
#[schemars(length(max = 32))]
pub memberships: Vec<ResolvedEphemeralMembership>,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)]
#[serde(rename_all = "camelCase")]
pub struct ResolvedEphemeralMembership {
#[schemars(length(min = 1, max = 63))]
pub role: String,
#[schemars(length(min = 1, max = 63))]
pub member: String,
pub inherit: bool,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct CanonicalEphemeralBundle<'a> {
bundle_encoding: &'a str,
membership_semantics: &'a str,
target_database_fingerprint: &'a str,
memberships: &'a [ResolvedEphemeralMembership],
}
impl ResolvedEphemeralAccess {
pub fn canonical_bundle_bytes(&self) -> Vec<u8> {
let mut memberships = self.memberships.clone();
memberships.sort();
serde_json::to_vec(&CanonicalEphemeralBundle {
bundle_encoding: EPHEMERAL_BUNDLE_ENCODING_V1,
membership_semantics: EPHEMERAL_MEMBERSHIP_SEMANTICS_V1,
target_database_fingerprint: &self.target_database_fingerprint,
memberships: &memberships,
})
.expect("canonical ephemeral bundle is serializable")
}
pub fn compute_bundle_hash(&self) -> String {
use sha2::{Digest, Sha256};
use std::fmt::Write;
let digest = Sha256::digest(self.canonical_bundle_bytes());
let mut hash = String::with_capacity(7 + digest.len() * 2);
hash.push_str("sha256:");
for byte in digest {
write!(&mut hash, "{byte:02x}").expect("writing to a String cannot fail");
}
hash
}
pub fn has_valid_bundle_hash(&self) -> bool {
self.bundle_encoding == EPHEMERAL_BUNDLE_ENCODING_V1
&& self.bundle_hash == self.compute_bundle_hash()
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub struct DatabaseIdentity(String);
impl DatabaseIdentity {
pub fn from_connection(namespace: &str, connection: &ConnectionSpec) -> Self {
Self(format!("{namespace}/{}", connection.identity_key()))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct OwnershipClaims {
pub roles: BTreeSet<String>,
pub schemas: BTreeSet<String>,
pub public_databases: BTreeSet<String>,
}
impl OwnershipClaims {
pub fn overlaps(&self, other: &Self) -> bool {
!self.roles.is_disjoint(&other.roles)
|| !self.schemas.is_disjoint(&other.schemas)
|| !self.public_databases.is_disjoint(&other.public_databases)
}
pub fn overlap_summary(&self, other: &Self) -> String {
let overlapping_roles: Vec<_> = self.roles.intersection(&other.roles).cloned().collect();
let overlapping_schemas: Vec<_> =
self.schemas.intersection(&other.schemas).cloned().collect();
let overlapping_databases: Vec<_> = self
.public_databases
.intersection(&other.public_databases)
.cloned()
.collect();
let mut parts = Vec::new();
if !overlapping_roles.is_empty() {
parts.push(format!("roles: {}", overlapping_roles.join(", ")));
}
if !overlapping_schemas.is_empty() {
parts.push(format!("schemas: {}", overlapping_schemas.join(", ")));
}
if !overlapping_databases.is_empty() {
parts.push(format!(
"PUBLIC on databases: {}",
overlapping_databases.join(", ")
));
}
parts.join("; ")
}
}
impl PostgresPolicySpec {
pub fn validate_password_specs(
&self,
policy_name: &str,
) -> Result<(), PasswordValidationError> {
for role in &self.roles {
let Some(password) = &role.password else {
continue;
};
if role.login != Some(true) {
return Err(PasswordValidationError::PasswordWithoutLogin {
role: role.name.clone(),
});
}
match (&password.secret_ref, &password.generate) {
(Some(_), None) => {
let secret_key = password.secret_key.as_deref().unwrap_or(&role.name);
if !is_valid_secret_key(secret_key) {
return Err(PasswordValidationError::InvalidSecretKey {
role: role.name.clone(),
field: "secretKey",
key: secret_key.to_string(),
});
}
}
(None, Some(generate)) => {
if let Some(length) = generate.length
&& !(crate::password::MIN_PASSWORD_LENGTH
..=crate::password::MAX_PASSWORD_LENGTH)
.contains(&length)
{
return Err(PasswordValidationError::InvalidGeneratedLength {
role: role.name.clone(),
min: crate::password::MIN_PASSWORD_LENGTH,
max: crate::password::MAX_PASSWORD_LENGTH,
});
}
let secret_name =
crate::password::generated_secret_name(policy_name, &role.name, generate);
if !is_valid_secret_name(&secret_name) {
return Err(PasswordValidationError::InvalidGeneratedSecretName {
role: role.name.clone(),
name: secret_name,
});
}
let secret_key = crate::password::generated_secret_key(generate);
if !is_valid_secret_key(&secret_key) {
return Err(PasswordValidationError::InvalidSecretKey {
role: role.name.clone(),
field: "generate.secretKey",
key: secret_key,
});
}
if secret_key == crate::password::GENERATED_VERIFIER_KEY {
return Err(PasswordValidationError::ReservedGeneratedSecretKey {
role: role.name.clone(),
key: secret_key,
});
}
}
_ => {
return Err(PasswordValidationError::InvalidPasswordMode {
role: role.name.clone(),
});
}
}
}
Ok(())
}
pub fn validate_connection_spec(&self) -> Result<(), ConnectionValidationError> {
let conn = &self.connection;
match (&conn.secret_ref, &conn.params) {
(Some(_), None) => {
Ok(())
}
(None, Some(params)) => {
fn validate_required_field(
field: &str,
literal: &Option<String>,
secret: &Option<SecretKeySelector>,
) -> Result<(), ConnectionValidationError> {
match (literal, secret) {
(Some(_), Some(_)) => {
return Err(ConnectionValidationError::BothFieldsSet {
field: field.to_string(),
});
}
(None, None) => {
return Err(ConnectionValidationError::NeitherFieldSet {
field: field.to_string(),
});
}
(Some(s), None) => {
if s.trim().is_empty() {
return Err(ConnectionValidationError::EmptyLiteral {
field: field.to_string(),
});
}
}
(None, Some(sel)) => {
validate_secret_selector(field, sel)?;
}
}
Ok(())
}
fn validate_optional_field(
field: &str,
literal: &Option<impl AsRef<str>>,
secret: &Option<SecretKeySelector>,
) -> Result<(), ConnectionValidationError> {
let has_literal = literal.is_some();
if has_literal && secret.is_some() {
return Err(ConnectionValidationError::BothFieldsSet {
field: field.to_string(),
});
}
if let Some(s) = literal
&& s.as_ref().trim().is_empty()
{
return Err(ConnectionValidationError::EmptyLiteral {
field: field.to_string(),
});
}
if let Some(sel) = secret {
validate_secret_selector(field, sel)?;
}
Ok(())
}
fn validate_secret_selector(
field: &str,
sel: &SecretKeySelector,
) -> Result<(), ConnectionValidationError> {
if sel.name.trim().is_empty() {
return Err(ConnectionValidationError::EmptySecretKeyRef {
field: field.to_string(),
detail: "name must not be empty".to_string(),
});
}
if sel.key.trim().is_empty() {
return Err(ConnectionValidationError::EmptySecretKeyRef {
field: field.to_string(),
detail: "key must not be empty".to_string(),
});
}
Ok(())
}
validate_required_field("host", ¶ms.host, ¶ms.host_secret)?;
validate_required_field("dbname", ¶ms.dbname, ¶ms.dbname_secret)?;
validate_required_field("username", ¶ms.username, ¶ms.username_secret)?;
if let Some(auth) = ¶ms.auth {
if params.password.is_some() || params.password_secret.is_some() {
return Err(ConnectionValidationError::AuthWithPassword);
}
match auth {
ConnectionAuth::GcpWorkloadIdentity {
impersonate_service_account,
scope,
} => {
if let Some(value) = impersonate_service_account
&& value.trim().is_empty()
{
return Err(ConnectionValidationError::EmptyAuthField {
field: "impersonateServiceAccount".to_string(),
});
}
if let Some(value) = scope
&& value.trim().is_empty()
{
return Err(ConnectionValidationError::EmptyAuthField {
field: "scope".to_string(),
});
}
}
}
} else {
validate_required_field("password", ¶ms.password, ¶ms.password_secret)?;
}
let port_str = params.port.map(|p| p.to_string());
validate_optional_field("port", &port_str, ¶ms.port_secret)?;
validate_optional_field("sslMode", ¶ms.ssl_mode, ¶ms.ssl_mode_secret)?;
if let Some(value) = ¶ms.ssl_mode
&& !VALID_SSL_MODES.contains(&value.as_str())
{
return Err(ConnectionValidationError::InvalidSslMode {
value: value.clone(),
});
}
if let Some(value) = ¶ms.set_role {
if value.trim().is_empty() {
return Err(ConnectionValidationError::EmptyLiteral {
field: "setRole".to_string(),
});
}
if !is_valid_set_role_identifier(value) {
return Err(ConnectionValidationError::InvalidRoleName {
value: value.clone(),
});
}
}
Ok(())
}
(Some(_), Some(_)) => Err(ConnectionValidationError::BothModesSet),
(None, None) => Err(ConnectionValidationError::NeitherModeSet),
}
}
pub fn referenced_secret_names(&self, policy_name: &str) -> BTreeSet<String> {
let mut names = BTreeSet::new();
self.connection.collect_secret_names(&mut names);
for role in &self.roles {
if let Some(pw) = &role.password {
if let Some(secret_ref) = &pw.secret_ref {
names.insert(secret_ref.name.clone());
}
if let Some(gen_spec) = &pw.generate {
let secret_name =
crate::password::generated_secret_name(policy_name, &role.name, gen_spec);
names.insert(secret_name);
}
}
}
names
}
}
#[allow(clippy::too_many_arguments)]
fn build_policy_manifest<'a>(
default_owner: Option<&str>,
profiles: impl Iterator<Item = (&'a String, &'a ProfileSpec)>,
schemas: &[SchemaBinding],
roles: &[RoleSpec],
grants: &[Grant],
default_privileges: &[DefaultPrivilege],
memberships: &[Membership],
retirements: &[RoleRetirement],
) -> pgroles_core::manifest::PolicyManifest {
use pgroles_core::manifest::{
DefaultPrivilegeGrant, MemberSpec, PolicyManifest, Profile, ProfileGrant,
ProfileObjectTarget, RoleDefinition,
};
let profiles = profiles
.map(|(name, spec)| {
let profile = Profile {
login: spec.login,
inherit: spec.inherit,
grants: spec
.grants
.iter()
.map(|g| ProfileGrant {
privileges: g.privileges.clone(),
object: ProfileObjectTarget {
object_type: g.object.object_type,
name: g.object.name.clone(),
},
ensure: g.ensure,
})
.collect(),
default_privileges: spec
.default_privileges
.iter()
.map(|dp| DefaultPrivilegeGrant {
role: dp.role.clone(),
privileges: dp.privileges.clone(),
on_type: dp.on_type,
ensure: dp.ensure,
})
.collect(),
config: spec.config.clone(),
};
(name.clone(), profile)
})
.collect();
let roles = roles
.iter()
.map(|r| RoleDefinition {
name: r.name.clone(),
external: r.external,
login: r.login,
superuser: r.superuser,
createdb: r.createdb,
createrole: r.createrole,
inherit: r.inherit,
replication: r.replication,
bypassrls: r.bypassrls,
connection_limit: r.connection_limit,
comment: r.comment.clone(),
password: None, password_valid_until: r.password_valid_until.clone(),
config: r.config.clone(),
})
.collect();
let memberships = memberships
.iter()
.map(|m| pgroles_core::manifest::Membership {
role: m.role.clone(),
members: m
.members
.iter()
.map(|ms| MemberSpec {
name: ms.name.clone(),
inherit: ms.inherit,
admin: ms.admin,
})
.collect(),
})
.collect();
PolicyManifest {
default_owner: default_owner.map(str::to_string),
auth_providers: Vec::new(),
profiles,
schemas: schemas.to_vec(),
roles,
grants: grants.to_vec(),
default_privileges: default_privileges.to_vec(),
memberships,
retirements: retirements.to_vec(),
}
}
impl PolicyContent {
pub fn to_policy_manifest(&self) -> pgroles_core::manifest::PolicyManifest {
build_policy_manifest(
self.default_owner.as_deref(),
self.profiles.iter(),
&self.schemas,
&self.roles,
&self.grants,
&self.default_privileges,
&self.memberships,
&self.retirements,
)
}
pub fn content_digest(&self) -> String {
pgroles_core::candidate::compute_content_digest(self)
}
}
impl PostgresPolicySpec {
pub fn policy_content(&self) -> PolicyContent {
PolicyContent {
reconciliation_mode: self.reconciliation_mode,
default_owner: self.default_owner.clone(),
profiles: self
.profiles
.iter()
.map(|(name, spec)| (name.clone(), spec.clone()))
.collect(),
schemas: self.schemas.clone(),
roles: self.roles.clone(),
grants: self.grants.clone(),
default_privileges: self.default_privileges.clone(),
memberships: self.memberships.clone(),
retirements: self.retirements.clone(),
}
}
pub fn content_digest(&self) -> String {
self.policy_content().content_digest()
}
pub fn to_policy_manifest(&self) -> pgroles_core::manifest::PolicyManifest {
build_policy_manifest(
self.default_owner.as_deref(),
self.profiles.iter(),
&self.schemas,
&self.roles,
&self.grants,
&self.default_privileges,
&self.memberships,
&self.retirements,
)
}
pub fn ownership_claims(
&self,
) -> Result<OwnershipClaims, pgroles_core::manifest::ManifestError> {
let manifest = self.to_policy_manifest();
let expanded = pgroles_core::manifest::expand_manifest(&manifest)?;
let mut roles: BTreeSet<String> = expanded.roles.into_iter().map(|r| r.name).collect();
let mut schemas: BTreeSet<String> = self.schemas.iter().map(|s| s.name.clone()).collect();
let public_databases: BTreeSet<String> = manifest
.grants
.iter()
.filter(|g| g.object.object_type == ObjectType::Database && g.role == "PUBLIC")
.filter_map(|g| g.object.name.clone())
.collect();
roles.extend(manifest.retirements.into_iter().map(|r| r.role));
roles.extend(
manifest
.grants
.iter()
.map(|g| g.role.clone())
.filter(|role| role != "PUBLIC"),
);
roles.extend(
manifest
.default_privileges
.iter()
.flat_map(|dp| dp.grant.iter().filter_map(|grant| grant.role.clone()))
.filter(|role| role != "PUBLIC"),
);
roles.extend(manifest.memberships.iter().map(|m| m.role.clone()));
roles.extend(
manifest
.memberships
.iter()
.flat_map(|m| m.members.iter().map(|member| member.name.clone())),
);
schemas.extend(
manifest
.grants
.iter()
.filter_map(|g| match g.object.object_type {
ObjectType::Database => None,
ObjectType::Schema => g.object.name.clone(),
_ => g.object.schema.clone(),
}),
);
for dp in &manifest.default_privileges {
match dp.resolved_scope() {
Ok(scope) => match scope.schema() {
Some(schema) => {
schemas.insert(schema.to_string());
}
None => {
if let Some(owner) = dp.owner.as_ref().or(manifest.default_owner.as_ref()) {
roles.insert(owner.clone());
}
}
},
Err(_) => continue,
}
}
Ok(OwnershipClaims {
roles,
schemas,
public_databases,
})
}
}
impl PostgresPolicyStatus {
pub fn set_condition(&mut self, new: PolicyCondition) {
if let Some(existing) = self
.conditions
.iter()
.find(|c| c.condition_type == new.condition_type)
&& existing.status == new.status
{
let mut updated = new;
updated.last_transition_time = existing.last_transition_time.clone();
self.conditions
.retain(|c| c.condition_type != updated.condition_type);
self.conditions.push(updated);
return;
}
self.conditions
.retain(|c| c.condition_type != new.condition_type);
self.conditions.push(new);
}
}
pub fn now_rfc3339() -> String {
use std::time::SystemTime;
let now = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.unwrap_or_default();
let secs = now.as_secs();
let days = secs / 86400;
let remaining = secs % 86400;
let hours = remaining / 3600;
let minutes = (remaining % 3600) / 60;
let seconds = remaining % 60;
let (year, month, day) = days_to_date(days);
format!("{year:04}-{month:02}-{day:02}T{hours:02}:{minutes:02}:{seconds:02}Z")
}
pub fn days_to_date(days_since_epoch: u64) -> (u64, u64, u64) {
let z = days_since_epoch + 719468;
let era = z / 146097;
let doe = z - era * 146097;
let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146096) / 365;
let y = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let d = doy - (153 * mp + 2) / 5 + 1;
let m = if mp < 10 { mp + 3 } else { mp - 9 };
let y = if m <= 2 { y + 1 } else { y };
(y, m, d)
}
pub fn ready_condition(status: bool, reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: "Ready".to_string(),
status: if status { "True" } else { "False" }.to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn reconciling_condition(message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: "Reconciling".to_string(),
status: "True".to_string(),
reason: Some("Reconciling".to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn degraded_condition(reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: "Degraded".to_string(),
status: "True".to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn paused_condition(message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: "Paused".to_string(),
status: "True".to_string(),
reason: Some("Suspended".to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub const CONDITION_TARGET_IDENTITY_BLOCKED: &str = "TargetIdentityBlocked";
pub fn target_identity_blocked_condition(reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: CONDITION_TARGET_IDENTITY_BLOCKED.to_string(),
status: "True".to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn conflict_condition(reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: "Conflict".to_string(),
status: "True".to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub const CONDITION_APPROVAL_UNSET: &str = "ApprovalUnset";
pub const CONDITION_MODE_VALUE_DEPRECATED: &str = "ModeValueDeprecated";
pub const CONDITION_ABSENCE_ASSERTIONS_IGNORED: &str = "AbsenceAssertionsIgnored";
pub const CONDITION_APPROVAL_IGNORED: &str = "ApprovalIgnored";
pub fn absence_assertions_ignored_condition() -> PolicyCondition {
PolicyCondition {
condition_type: CONDITION_ABSENCE_ASSERTIONS_IGNORED.to_string(),
status: "True".to_string(),
reason: Some("AdditiveModeNeverRevokes".to_string()),
message: Some(
"spec.reconciliation_mode is `additive`, so every `ensure: absent` assertion is \
ignored. Use `adopt` or `authoritative` to enforce absence."
.to_string(),
),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn approval_ignored_condition(plan_name: &str) -> PolicyCondition {
PolicyCondition {
condition_type: CONDITION_APPROVAL_IGNORED.to_string(),
status: "True".to_string(),
reason: Some("ObserveModeNeverExecutes".to_string()),
message: Some(format!(
"Plan {plan_name} is approved, but spec.mode is `observe`, so it will never execute and \
no SQL will run. For a reviewed apply use `mode: apply` with `approval: manual`."
)),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn approval_unset_condition(inferred: ApprovalMode) -> PolicyCondition {
let inferred = match inferred {
ApprovalMode::Auto => "auto",
ApprovalMode::Manual => "manual",
};
PolicyCondition {
condition_type: CONDITION_APPROVAL_UNSET.to_string(),
status: "True".to_string(),
reason: Some("InferredFromMode".to_string()),
message: Some(format!(
"spec.approval is not set and is currently inferred as {inferred} from spec.mode. \
This inference is deprecated and will become an error in a future release: set \
`approval: {inferred}` explicitly to keep the current behaviour."
)),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn mode_value_deprecated_condition() -> PolicyCondition {
PolicyCondition {
condition_type: CONDITION_MODE_VALUE_DEPRECATED.to_string(),
status: "True".to_string(),
reason: Some("PlanSpelledObserve".to_string()),
message: Some(
"spec.mode is `plan`, the deprecated spelling of `observe`. Behaviour is identical: \
plans are computed and published, nothing executes. Change the manifest to \
`mode: observe` — a future release removes the `plan` value."
.to_string(),
),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn drifted_condition(status: bool, reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: "Drifted".to_string(),
status: if status { "True" } else { "False" }.to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
pub struct PolicyContent {
#[serde(default)]
pub reconciliation_mode: CrdReconciliationMode,
#[serde(default)]
#[schemars(length(min = 1, max = MAX_IDENTIFIER))]
pub default_owner: Option<String>,
#[serde(default)]
#[schemars(extend("maxProperties" = MAX_PROFILES))]
pub profiles: std::collections::BTreeMap<String, ProfileSpec>,
#[serde(default)]
#[schemars(length(max = MAX_SCHEMAS))]
pub schemas: Vec<SchemaBinding>,
#[serde(default)]
#[schemars(length(max = MAX_ROLES))]
pub roles: Vec<RoleSpec>,
#[serde(default)]
#[schemars(length(max = MAX_GRANTS))]
pub grants: Vec<Grant>,
#[serde(default)]
#[schemars(length(max = MAX_DEFAULT_PRIVILEGES))]
pub default_privileges: Vec<DefaultPrivilege>,
#[serde(default)]
#[schemars(length(max = MAX_MEMBERSHIPS))]
pub memberships: Vec<Membership>,
#[serde(default)]
#[schemars(length(max = MAX_RETIREMENTS))]
pub retirements: Vec<RoleRetirement>,
}
#[derive(CustomResource, KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[kube(
group = "pgroles.io",
version = "v1alpha1",
kind = "PostgresPolicyCandidate",
namespaced,
status = "PostgresPolicyCandidateStatus",
shortname = "pgcand",
category = "pgroles",
printcolumn = r#"{"name":"Policy","type":"string","jsonPath":".spec.policyRef.name"}"#,
printcolumn = r#"{"name":"Phase","type":"string","jsonPath":".status.phase"}"#,
printcolumn = r#"{"name":"Plan","type":"string","jsonPath":".status.planRef.name"}"#,
printcolumn = r#"{"name":"Digest","type":"string","jsonPath":".status.contentDigest","priority":1}"#,
printcolumn = r#"{"name":"Age","type":"date","jsonPath":".metadata.creationTimestamp"}"#
)]
#[x_kube(validation = Rule::new("self == oldSelf").message("candidate spec is immutable"))]
#[serde(rename_all = "camelCase")]
pub struct PostgresPolicyCandidateSpec {
pub policy_ref: LocalObjectReference,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(min = 1, max = MAX_K8S_NAME))]
pub replaces: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub target: Option<CandidateTarget>,
pub content: PolicyContent,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct CandidateTarget {
pub connection_ref: CandidateConnectionRef,
}
#[derive(KubeSchema, Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct CandidateConnectionRef {
#[schemars(length(min = 1, max = MAX_K8S_NAME))]
pub secret_name: String,
#[schemars(length(min = 1, max = MAX_SECRET_KEY))]
pub key: String,
}
#[derive(KubeSchema, Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct PostgresPolicyCandidateStatus {
#[serde(default)]
pub phase: CandidatePhase,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[schemars(length(max = 128))]
pub content_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub plan_ref: Option<PlanReference>,
#[serde(default)]
#[schemars(length(max = 16))]
pub conditions: Vec<PolicyCondition>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub observed_generation: Option<i64>,
}
#[derive(JsonSchema, Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
pub enum CandidatePhase {
#[default]
Pending,
Planned,
Promoted,
Superseded,
Stale,
}
impl CandidatePhase {
pub fn is_terminal(self) -> bool {
matches!(self, CandidatePhase::Promoted | CandidatePhase::Superseded)
}
}
pub mod candidate_reason {
pub const PLANNED: &str = "Planned";
pub const NO_EFFECTS: &str = "NoEffects";
pub const BLOCKED_BY_ACTIVE_POLICY: &str = "BlockedByActivePolicy";
pub const OVERLAY_OVERLAP: &str = "OverlayOverlap";
pub const PLANNING_FAILED: &str = "PlanningFailed";
pub const OVER_BUDGET: &str = "CandidateBudgetExceeded";
pub const EXPIRED: &str = "Expired";
pub const REPLACED: &str = "Replaced";
pub const EFFECTS_CHANGED: &str = "EffectsChanged";
pub const PLAN_DENIED: &str = "PlanDenied";
pub const PROMOTED: &str = "Promoted";
pub const PROMOTED_WITHOUT_APPROVAL: &str = "PromotedWithoutApproval";
pub const PROMOTION_DIGEST_MISMATCH: &str = "PromotionDigestMismatch";
pub const PROMOTION_BASE_CHANGED: &str = "PromotionBaseChanged";
pub const PROMOTION_NOT_EXECUTED: &str = "PromotionNotExecuted";
pub const SUPERSEDED_BY_PROMOTION: &str = "SupersededByPromotion";
}
pub const CONDITION_SUPERSEDED: &str = "Superseded";
pub const CONDITION_PROMOTED: &str = "Promoted";
pub fn promoted_condition(promoted: bool, reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: CONDITION_PROMOTED.to_string(),
status: if promoted { "True" } else { "False" }.to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
pub fn set_condition_in(conditions: &mut Vec<PolicyCondition>, new: PolicyCondition) {
if let Some(existing) = conditions
.iter()
.find(|c| c.condition_type == new.condition_type)
&& existing.status == new.status
{
let mut updated = new;
updated.last_transition_time = existing.last_transition_time.clone();
conditions.retain(|c| c.condition_type != updated.condition_type);
conditions.push(updated);
return;
}
conditions.retain(|c| c.condition_type != new.condition_type);
conditions.push(new);
}
pub fn superseded_condition(reason: &str, message: &str) -> PolicyCondition {
PolicyCondition {
condition_type: CONDITION_SUPERSEDED.to_string(),
status: "True".to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: Some(now_rfc3339()),
}
}
impl std::fmt::Display for CandidatePhase {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let name = match self {
CandidatePhase::Pending => "Pending",
CandidatePhase::Planned => "Planned",
CandidatePhase::Promoted => "Promoted",
CandidatePhase::Superseded => "Superseded",
CandidatePhase::Stale => "Stale",
};
f.write_str(name)
}
}
pub fn postgres_policy_candidate_crd()
-> k8s_openapi::apiextensions_apiserver::pkg::apis::apiextensions::v1::CustomResourceDefinition {
use kube::CustomResourceExt;
let mut crd = PostgresPolicyCandidate::crd();
for version in &mut crd.spec.versions {
let Some(content) = version
.schema
.as_mut()
.and_then(|s| s.open_api_v3_schema.as_mut())
.and_then(|s| s.properties.as_mut())
.and_then(|p| p.get_mut("spec"))
.and_then(|s| s.properties.as_mut())
.and_then(|p| p.get_mut("content"))
else {
continue;
};
strip_schema_defaults(content);
}
crd
}
fn strip_schema_defaults(
schema: &mut k8s_openapi::apiextensions_apiserver::pkg::apis::apiextensions::v1::JSONSchemaProps,
) {
use k8s_openapi::apiextensions_apiserver::pkg::apis::apiextensions::v1::{
JSONSchemaPropsOrArray, JSONSchemaPropsOrBool,
};
schema.default = None;
if let Some(properties) = schema.properties.as_mut() {
for child in properties.values_mut() {
strip_schema_defaults(child);
}
}
if let Some(JSONSchemaPropsOrBool::Schema(child)) = schema.additional_properties.as_mut() {
strip_schema_defaults(child);
}
match schema.items.as_mut() {
Some(JSONSchemaPropsOrArray::Schema(child)) => strip_schema_defaults(child),
Some(JSONSchemaPropsOrArray::Schemas(children)) => {
children.iter_mut().for_each(strip_schema_defaults)
}
None => {}
}
for branch in [
schema.all_of.as_mut(),
schema.any_of.as_mut(),
schema.one_of.as_mut(),
]
.into_iter()
.flatten()
{
branch.iter_mut().for_each(strip_schema_defaults);
}
if let Some(child) = schema.not.as_mut() {
strip_schema_defaults(child);
}
}
#[cfg(test)]
mod tests {
use super::*;
use kube::CustomResourceExt;
#[test]
fn named_spec_arrays_declare_list_map_keys() {
let crd = PostgresPolicy::crd();
let schema = crd.spec.versions[0]
.schema
.as_ref()
.and_then(|s| s.open_api_v3_schema.as_ref())
.expect("CRD should carry an OpenAPI schema");
let spec_props = schema
.properties
.as_ref()
.and_then(|p| p.get("spec"))
.and_then(|s| s.properties.as_ref())
.expect("spec should have properties");
for (field, key) in [
("schemas", "name"),
("roles", "name"),
("retirements", "role"),
] {
let prop = spec_props
.get(field)
.unwrap_or_else(|| panic!("spec.{field} should exist"));
assert_eq!(
prop.x_kubernetes_list_type.as_deref(),
Some("map"),
"spec.{field} should be a map-list"
);
assert_eq!(
prop.x_kubernetes_list_map_keys.as_deref(),
Some([key.to_string()].as_slice()),
"spec.{field} should be keyed by {key}"
);
}
for field in ["memberships", "grants", "default_privileges"] {
let prop = spec_props
.get(field)
.unwrap_or_else(|| panic!("spec.{field} should exist"));
assert!(
prop.x_kubernetes_list_type.is_none(),
"spec.{field} should stay a plain array until its merge semantics are designed"
);
}
}
#[test]
fn crd_generates_valid_schema() {
let crd = PostgresPolicy::crd();
let yaml = serde_yaml::to_string(&crd).expect("CRD should serialize to YAML");
assert!(yaml.contains("pgroles.io"), "group should be pgroles.io");
assert!(yaml.contains("v1alpha1"), "version should be v1alpha1");
assert!(
yaml.contains("PostgresPolicy"),
"kind should be PostgresPolicy"
);
assert!(
yaml.contains("\"mode\"") || yaml.contains(" mode:"),
"schema should declare spec.mode"
);
assert!(
yaml.contains("\"object\"") || yaml.contains(" object:"),
"schema should declare grant object targets using object"
);
}
#[test]
fn spec_to_policy_manifest_roundtrip() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-secret".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: Some("app_owner".to_string()),
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "analytics".to_string(),
external: true,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: Some("test role".to_string()),
password: None,
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![RoleRetirement {
role: "legacy-app".to_string(),
reassign_owned_to: Some("app_owner".to_string()),
drop_owned: true,
terminate_sessions: true,
}],
approval: None,
};
let manifest = spec.to_policy_manifest();
assert_eq!(manifest.default_owner, Some("app_owner".to_string()));
assert_eq!(manifest.roles.len(), 1);
assert_eq!(manifest.roles[0].name, "analytics");
assert!(manifest.roles[0].external);
assert_eq!(manifest.roles[0].login, Some(true));
assert_eq!(manifest.roles[0].comment, Some("test role".to_string()));
assert_eq!(manifest.retirements.len(), 1);
assert_eq!(manifest.retirements[0].role, "legacy-app");
assert_eq!(
manifest.retirements[0].reassign_owned_to.as_deref(),
Some("app_owner")
);
assert!(manifest.retirements[0].drop_owned);
assert!(manifest.retirements[0].terminate_sessions);
}
#[test]
fn spec_to_policy_manifest_preserves_profile_inherit() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-secret".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::from([(
"editor".to_string(),
ProfileSpec {
login: Some(false),
inherit: Some(false),
grants: vec![],
default_privileges: vec![],
config: Default::default(),
},
)]),
schemas: vec![],
roles: vec![],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
let manifest = spec.to_policy_manifest();
assert_eq!(manifest.profiles["editor"].login, Some(false));
assert_eq!(manifest.profiles["editor"].inherit, Some(false));
}
#[test]
fn status_set_condition_replaces_existing() {
let mut status = PostgresPolicyStatus::default();
status.set_condition(ready_condition(false, "Pending", "Initial"));
assert_eq!(status.conditions.len(), 1);
assert_eq!(status.conditions[0].status, "False");
status.set_condition(ready_condition(true, "Reconciled", "All good"));
assert_eq!(status.conditions.len(), 1);
assert_eq!(status.conditions[0].status, "True");
assert_eq!(status.conditions[0].reason.as_deref(), Some("Reconciled"));
}
#[test]
fn status_set_condition_adds_new_type() {
let mut status = PostgresPolicyStatus::default();
status.set_condition(ready_condition(true, "OK", "ready"));
status.set_condition(degraded_condition("Error", "something broke"));
assert_eq!(status.conditions.len(), 2);
}
#[test]
fn paused_condition_has_expected_shape() {
let paused = paused_condition("paused by spec");
assert_eq!(paused.condition_type, "Paused");
assert_eq!(paused.status, "True");
assert_eq!(paused.reason.as_deref(), Some("Suspended"));
}
#[test]
fn ownership_claims_include_expanded_roles_and_schemas() {
let mut profiles = std::collections::HashMap::new();
profiles.insert(
"editor".to_string(),
ProfileSpec {
login: Some(false),
inherit: None,
grants: vec![],
default_privileges: vec![],
config: Default::default(),
},
);
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-secret".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles,
schemas: vec![SchemaBinding {
name: "inventory".to_string(),
profiles: vec!["editor".to_string()],
role_pattern: "{schema}-{profile}".to_string(),
owner: None,
}],
roles: vec![RoleSpec {
name: "app-service".to_string(),
external: false,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: None,
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![RoleRetirement {
role: "legacy-app".to_string(),
reassign_owned_to: None,
drop_owned: false,
terminate_sessions: false,
}],
approval: None,
};
let claims = spec.ownership_claims().unwrap();
assert!(claims.roles.contains("inventory-editor"));
assert!(claims.roles.contains("app-service"));
assert!(claims.roles.contains("legacy-app"));
assert!(claims.schemas.contains("inventory"));
}
#[test]
fn a_global_default_privilege_claims_the_implicit_default_owner() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-secret".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: Some("app_owner".to_string()),
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![],
grants: vec![],
default_privileges: vec![DefaultPrivilege {
owner: None,
schema: None,
scope: Some(pgroles_core::manifest::DefaultPrivilegeScopeSpec {
scope_type: pgroles_core::manifest::DefaultPrivilegeScopeType::Global,
schema: None,
}),
grant: vec![pgroles_core::manifest::DefaultPrivilegeGrant {
role: Some("reader".to_string()),
privileges: vec![pgroles_core::manifest::Privilege::Select],
on_type: ObjectType::Table,
ensure: pgroles_core::manifest::Ensure::Present,
}],
}],
memberships: vec![],
retirements: vec![],
approval: None,
};
let claims = spec.ownership_claims().unwrap();
assert!(
claims.roles.contains("app_owner"),
"expected the resolved default owner to be claimed, got {:?}",
claims.roles
);
}
#[test]
fn ownership_overlap_summary_reports_roles_and_schemas() {
let mut left = OwnershipClaims::default();
left.roles.insert("analytics".to_string());
left.schemas.insert("reporting".to_string());
let mut right = OwnershipClaims::default();
right.roles.insert("analytics".to_string());
right.schemas.insert("reporting".to_string());
right.schemas.insert("other".to_string());
assert!(left.overlaps(&right));
let summary = left.overlap_summary(&right);
assert!(summary.contains("roles: analytics"));
assert!(summary.contains("schemas: reporting"));
}
#[test]
fn database_identity_uses_namespace_and_identity_key() {
let conn = ConnectionSpec {
secret_ref: Some(SecretReference {
name: "db-creds".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
};
let identity = DatabaseIdentity::from_connection("prod", &conn);
assert_eq!(identity.as_str(), "prod/db-creds/DATABASE_URL");
}
#[test]
fn identity_key_same_database_different_users_are_equal() {
let user_a = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("my-host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("mydb".into()),
dbname_secret: None,
username: Some("alice".into()),
username_secret: None,
password: Some("pass-a".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
};
let user_b = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("my-host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("mydb".into()),
dbname_secret: None,
username: Some("bob".into()),
username_secret: None,
password: Some("pass-b".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
};
assert_eq!(
user_a.identity_key(),
user_b.identity_key(),
"same database with different users should have the same identity key"
);
assert_ne!(
user_a.cache_key("default"),
user_b.cache_key("default"),
"different credentials should produce different cache keys"
);
}
#[test]
fn cache_key_no_collision_between_literal_and_secret_username() {
let literal_conn = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("my-host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("mydb".into()),
dbname_secret: None,
username: Some("secret=creds\0password".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
};
let secret_conn = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("my-host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("mydb".into()),
dbname_secret: None,
username: None,
username_secret: Some(SecretKeySelector {
name: "creds".into(),
key: "password".into(),
}),
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
};
assert_ne!(
literal_conn.cache_key("default"),
secret_conn.cache_key("default"),
"literal and secret ref should produce different cache keys"
);
}
#[test]
fn cache_key_includes_ssl_mode() {
let conn_no_ssl = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
};
let conn_with_ssl = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: Some("require".into()),
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
};
assert_ne!(
conn_no_ssl.cache_key("ns"),
conn_with_ssl.cache_key("ns"),
"cache key should differ when sslMode is present"
);
}
#[test]
fn validate_connection_rejects_empty_literal_host() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("mydb".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
let err = spec.validate_connection_spec().unwrap_err();
assert!(
matches!(err, ConnectionValidationError::EmptyLiteral { ref field } if field == "host"),
"expected EmptyLiteral for host, got: {err}"
);
}
#[test]
fn validate_connection_rejects_whitespace_literal_dbname() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some(" ".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
let err = spec.validate_connection_spec().unwrap_err();
assert!(
matches!(err, ConnectionValidationError::EmptyLiteral { ref field } if field == "dbname"),
"expected EmptyLiteral for dbname, got: {err}"
);
}
fn spec_with_connection(connection: ConnectionSpec) -> PostgresPolicySpec {
PostgresPolicySpec {
connection,
interval: "5m".into(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: Default::default(),
schemas: vec![],
roles: vec![],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
}
}
fn url_mode_connection() -> ConnectionSpec {
ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-creds".into(),
}),
secret_key: Some("DATABASE_URL".into()),
params: None,
require_physical_identity: None,
}
}
fn params_mode_connection() -> ConnectionSpec {
ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("my-postgres".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("mydb".into()),
dbname_secret: None,
username: None,
username_secret: Some(SecretKeySelector {
name: "pg-creds".into(),
key: "username".into(),
}),
password: None,
password_secret: Some(SecretKeySelector {
name: "pg-creds".into(),
key: "password".into(),
}),
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
}
}
#[test]
fn validate_connection_accepts_url_mode() {
let spec = spec_with_connection(url_mode_connection());
assert!(spec.validate_connection_spec().is_ok());
}
#[test]
fn validate_connection_accepts_params_mode() {
let spec = spec_with_connection(params_mode_connection());
assert!(spec.validate_connection_spec().is_ok());
}
#[test]
fn validate_connection_rejects_both_modes_set() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-creds".into(),
}),
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(matches!(
spec.validate_connection_spec(),
Err(ConnectionValidationError::BothModesSet)
));
}
#[test]
fn validate_connection_rejects_neither_mode_set() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: None,
require_physical_identity: None,
});
assert!(spec.validate_connection_spec().is_err());
}
#[test]
fn validate_connection_rejects_invalid_ssl_mode() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: Some("invalid-mode".into()),
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(spec.validate_connection_spec().is_err());
}
fn params_with_set_role(set_role: Option<String>) -> ConnectionParams {
ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role,
}
}
#[test]
fn validate_connection_accepts_valid_set_role() {
for role in [
"cloudsqlsuperuser",
"_underscore_start",
"role-with-dash",
"role_with$dollar",
"Mixed_Case_Role",
"r2d2",
] {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(params_with_set_role(Some(role.into()))),
require_physical_identity: None,
});
assert!(
spec.validate_connection_spec().is_ok(),
"expected {role} to be accepted"
);
}
}
#[test]
fn validate_connection_rejects_invalid_set_role() {
for role in [
"1leading_digit",
"has space",
"has\"quote",
"has;semicolon",
"has'singlequote",
"ünicode",
] {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(params_with_set_role(Some(role.into()))),
require_physical_identity: None,
});
let err = spec
.validate_connection_spec()
.expect_err(&format!("expected {role} to be rejected"));
assert!(
matches!(err, ConnectionValidationError::InvalidRoleName { ref value } if value == role),
"unexpected error for {role}: {err:?}",
);
}
}
#[test]
fn validate_connection_rejects_empty_set_role() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(params_with_set_role(Some(" ".into()))),
require_physical_identity: None,
});
assert!(matches!(
spec.validate_connection_spec(),
Err(ConnectionValidationError::EmptyLiteral { ref field }) if field == "setRole"
));
}
#[test]
fn cache_key_includes_set_role() {
let conn_no_role = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(params_with_set_role(None)),
require_physical_identity: None,
};
let conn_with_role = ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(params_with_set_role(Some("cloudsqlsuperuser".into()))),
require_physical_identity: None,
};
assert_ne!(
conn_no_role.cache_key("ns"),
conn_with_role.cache_key("ns"),
"cache key should differ when setRole is present"
);
}
#[test]
fn validate_connection_accepts_gcp_workload_identity_without_password() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("10.0.0.5".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("discovery".into()),
dbname_secret: None,
username: Some("pgroles-operator@my-project.iam".into()),
username_secret: None,
password: None,
password_secret: None,
auth: Some(ConnectionAuth::GcpWorkloadIdentity {
impersonate_service_account: None,
scope: None,
}),
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(spec.validate_connection_spec().is_ok());
assert!(spec.referenced_secret_names("policy").is_empty());
}
#[test]
fn validate_connection_rejects_gcp_workload_identity_with_password() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("10.0.0.5".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("discovery".into()),
dbname_secret: None,
username: Some("pgroles-operator@my-project.iam".into()),
username_secret: None,
password: Some("static-password".into()),
password_secret: None,
auth: Some(ConnectionAuth::GcpWorkloadIdentity {
impersonate_service_account: None,
scope: None,
}),
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(matches!(
spec.validate_connection_spec(),
Err(ConnectionValidationError::AuthWithPassword)
));
}
#[test]
fn validate_connection_rejects_empty_gcp_auth_fields() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("10.0.0.5".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("discovery".into()),
dbname_secret: None,
username: Some("pgroles-operator@my-project.iam".into()),
username_secret: None,
password: None,
password_secret: None,
auth: Some(ConnectionAuth::GcpWorkloadIdentity {
impersonate_service_account: Some(" ".into()),
scope: None,
}),
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(matches!(
spec.validate_connection_spec(),
Err(ConnectionValidationError::EmptyAuthField { ref field })
if field == "impersonateServiceAccount"
));
}
#[test]
fn validate_connection_accepts_valid_ssl_modes() {
for mode in &[
"disable",
"allow",
"prefer",
"require",
"verify-ca",
"verify-full",
] {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: Some((*mode).into()),
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(
spec.validate_connection_spec().is_ok(),
"sslMode '{mode}' should be accepted"
);
}
}
#[test]
fn validate_connection_rejects_empty_secret_name() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: None,
username_secret: Some(SecretKeySelector {
name: "".into(),
key: "username".into(),
}),
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(spec.validate_connection_spec().is_err());
}
#[test]
fn validate_connection_rejects_both_literal_and_secret_for_same_field() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: Some("host".into()),
host_secret: Some(SecretKeySelector {
name: "s".into(),
key: "k".into(),
}),
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(matches!(
spec.validate_connection_spec(),
Err(ConnectionValidationError::BothFieldsSet { ref field }) if field == "host"
));
}
#[test]
fn validate_connection_rejects_neither_literal_nor_secret_for_required_field() {
let spec = spec_with_connection(ConnectionSpec {
secret_ref: None,
secret_key: None,
params: Some(ConnectionParams {
host: None,
host_secret: None,
port: None,
port_secret: None,
dbname: Some("db".into()),
dbname_secret: None,
username: Some("user".into()),
username_secret: None,
password: Some("pass".into()),
password_secret: None,
auth: None,
ssl_mode: None,
ssl_mode_secret: None,
set_role: None,
}),
require_physical_identity: None,
});
assert!(matches!(
spec.validate_connection_spec(),
Err(ConnectionValidationError::NeitherFieldSet { ref field }) if field == "host"
));
}
#[test]
fn connection_spec_backward_compat_url_mode() {
let yaml = r#"
secretRef:
name: pg-creds
secretKey: DATABASE_URL
"#;
let conn: ConnectionSpec = serde_yaml::from_str(yaml).unwrap();
assert!(conn.secret_ref.is_some());
assert_eq!(conn.effective_secret_key(), "DATABASE_URL");
assert!(conn.params.is_none());
}
#[test]
fn connection_spec_backward_compat_default_secret_key() {
let yaml = r#"
secretRef:
name: pg-creds
"#;
let conn: ConnectionSpec = serde_yaml::from_str(yaml).unwrap();
assert_eq!(conn.effective_secret_key(), "DATABASE_URL");
}
#[test]
fn connection_spec_params_mode_deserializes_keycloak_style() {
let yaml = r#"
params:
host: my-postgres
port: 5432
dbname: mydb
usernameSecret:
name: creds
key: username
passwordSecret:
name: creds
key: password
sslMode: require
"#;
let conn: ConnectionSpec = serde_yaml::from_str(yaml).unwrap();
assert!(conn.secret_ref.is_none());
let params = conn.params.unwrap();
assert_eq!(params.host.as_deref(), Some("my-postgres"));
assert_eq!(params.port, Some(5432));
assert!(params.username_secret.is_some());
assert_eq!(params.username_secret.as_ref().unwrap().name, "creds");
assert_eq!(params.ssl_mode.as_deref(), Some("require"));
}
#[test]
fn connection_spec_params_mode_deserializes_gcp_workload_identity_auth() {
let yaml = r#"
params:
host: 10.0.0.5
port: 5432
dbname: discovery
username: pgroles-operator@my-project.iam
auth:
type: gcp_workload_identity
impersonateServiceAccount: target@other-project.iam.gserviceaccount.com
scope: https://example.com/custom-scope
"#;
let conn: ConnectionSpec = serde_yaml::from_str(yaml).unwrap();
let params = conn.params.as_ref().unwrap();
let auth = params.auth.as_ref().expect("auth should deserialize");
assert!(params.password.is_none());
assert_eq!(auth.gcp_scope(), "https://example.com/custom-scope");
assert_eq!(
auth.gcp_impersonate_service_account(),
Some("target@other-project.iam.gserviceaccount.com")
);
assert!(conn.cache_key("prod").contains("gcp_workload_identity"));
let spec = spec_with_connection(conn);
assert!(spec.validate_connection_spec().is_ok());
}
#[test]
fn connection_spec_params_mode_all_secrets() {
let yaml = r#"
params:
hostSecret:
name: cluster-app
key: host
portSecret:
name: cluster-app
key: port
dbnameSecret:
name: cluster-app
key: dbname
usernameSecret:
name: cluster-app
key: user
passwordSecret:
name: cluster-app
key: password
"#;
let conn: ConnectionSpec = serde_yaml::from_str(yaml).unwrap();
let params = conn.params.unwrap();
assert!(params.host.is_none());
assert!(params.host_secret.is_some());
assert_eq!(params.host_secret.as_ref().unwrap().name, "cluster-app");
assert!(params.port.is_none());
assert!(params.port_secret.is_some());
}
#[test]
fn referenced_secret_names_includes_params_secrets() {
let spec = spec_with_connection(params_mode_connection());
let names = spec.referenced_secret_names("test-policy");
assert!(
names.contains("pg-creds"),
"should include the credential secret from params"
);
}
#[test]
fn referenced_secret_names_deduplicates_across_modes() {
let mut spec = spec_with_connection(params_mode_connection());
spec.roles = vec![RoleSpec {
name: "app".into(),
external: false,
login: Some(true),
password: Some(PasswordSpec {
secret_ref: Some(SecretReference {
name: "pg-creds".into(),
}),
secret_key: Some("app-password".into()),
generate: None,
}),
password_valid_until: None,
config: Default::default(),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
}];
let names = spec.referenced_secret_names("test-policy");
assert_eq!(
names.iter().filter(|n| *n == "pg-creds").count(),
1,
"BTreeSet should deduplicate"
);
}
#[test]
fn connection_params_port_defaults_to_none() {
let yaml = r#"
params:
host: my-host
dbname: mydb
username: user
password: pass
"#;
let conn: ConnectionSpec = serde_yaml::from_str(yaml).unwrap();
let params = conn.params.unwrap();
assert!(
params.port.is_none(),
"port should default to None (resolved as 5432 at runtime)"
);
assert!(
params.port_secret.is_none(),
"portSecret should also default to None"
);
}
#[test]
fn now_rfc3339_produces_valid_format() {
let ts = now_rfc3339();
assert!(ts.len() == 20, "expected 20 chars, got {}: {ts}", ts.len());
assert!(ts.ends_with('Z'), "should end with Z: {ts}");
assert_eq!(&ts[4..5], "-", "should have dash at pos 4: {ts}");
assert_eq!(&ts[10..11], "T", "should have T at pos 10: {ts}");
}
#[test]
fn ready_condition_true_has_expected_shape() {
let cond = ready_condition(true, "Reconciled", "All changes applied");
assert_eq!(cond.condition_type, "Ready");
assert_eq!(cond.status, "True");
assert_eq!(cond.reason.as_deref(), Some("Reconciled"));
assert_eq!(cond.message.as_deref(), Some("All changes applied"));
assert!(cond.last_transition_time.is_some());
}
#[test]
fn ready_condition_false_has_expected_shape() {
let cond = ready_condition(false, "InvalidSpec", "bad manifest");
assert_eq!(cond.condition_type, "Ready");
assert_eq!(cond.status, "False");
assert_eq!(cond.reason.as_deref(), Some("InvalidSpec"));
assert_eq!(cond.message.as_deref(), Some("bad manifest"));
}
#[test]
fn degraded_condition_has_expected_shape() {
let cond = degraded_condition("InvalidSpec", "expansion failed");
assert_eq!(cond.condition_type, "Degraded");
assert_eq!(cond.status, "True");
assert_eq!(cond.reason.as_deref(), Some("InvalidSpec"));
assert_eq!(cond.message.as_deref(), Some("expansion failed"));
assert!(cond.last_transition_time.is_some());
}
#[test]
fn reconciling_condition_has_expected_shape() {
let cond = reconciling_condition("Reconciliation in progress");
assert_eq!(cond.condition_type, "Reconciling");
assert_eq!(cond.status, "True");
assert_eq!(cond.reason.as_deref(), Some("Reconciling"));
assert_eq!(cond.message.as_deref(), Some("Reconciliation in progress"));
assert!(cond.last_transition_time.is_some());
}
#[test]
fn conflict_condition_has_expected_shape() {
let cond = conflict_condition("ConflictingPolicy", "overlaps with ns/other");
assert_eq!(cond.condition_type, "Conflict");
assert_eq!(cond.status, "True");
assert_eq!(cond.reason.as_deref(), Some("ConflictingPolicy"));
assert_eq!(cond.message.as_deref(), Some("overlaps with ns/other"));
assert!(cond.last_transition_time.is_some());
}
#[test]
fn approval_unset_condition_names_the_inferred_mode() {
for (mode, expected) in [
(ApprovalMode::Auto, "auto"),
(ApprovalMode::Manual, "manual"),
] {
let cond = approval_unset_condition(mode);
assert_eq!(cond.condition_type, CONDITION_APPROVAL_UNSET);
assert_eq!(cond.status, "True");
assert_eq!(cond.reason.as_deref(), Some("InferredFromMode"));
assert!(cond.last_transition_time.is_some());
let message = cond.message.expect("condition should carry a message");
assert!(
message.contains(&format!("inferred as {expected}")),
"message should name the inferred mode, got: {message}"
);
assert!(
message.contains(&format!("`approval: {expected}`")),
"message should show the fix to apply, got: {message}"
);
assert!(
message.contains("deprecated"),
"message should say the inference is deprecated, got: {message}"
);
}
}
#[test]
fn ownership_claims_no_overlap() {
let mut left = OwnershipClaims::default();
left.roles.insert("analytics".to_string());
left.schemas.insert("reporting".to_string());
let mut right = OwnershipClaims::default();
right.roles.insert("billing".to_string());
right.schemas.insert("payments".to_string());
assert!(!left.overlaps(&right));
let summary = left.overlap_summary(&right);
assert!(summary.is_empty());
}
fn public_connect_spec(
database: &str,
ensure: pgroles_core::manifest::Ensure,
) -> PostgresPolicySpec {
PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-secret".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![],
grants: vec![Grant {
role: "PUBLIC".to_string(),
privileges: vec![pgroles_core::manifest::Privilege::Connect],
object: pgroles_core::manifest::ObjectTarget {
object_type: ObjectType::Database,
schema: None,
name: Some(database.to_string()),
},
ensure,
}],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
}
}
#[test]
fn contradictory_public_database_rules_are_detected_as_conflicting() {
use pgroles_core::manifest::Ensure;
let present = public_connect_spec("mydb", Ensure::Present)
.ownership_claims()
.unwrap();
let absent = public_connect_spec("mydb", Ensure::Absent)
.ownership_claims()
.unwrap();
assert!(present.roles.is_disjoint(&absent.roles));
assert!(present.schemas.is_disjoint(&absent.schemas));
assert!(
present.overlaps(&absent),
"two policies writing PUBLIC on the same database must conflict"
);
assert!(present.overlap_summary(&absent).contains("mydb"));
}
#[test]
fn public_rules_on_different_databases_do_not_conflict() {
use pgroles_core::manifest::Ensure;
let left = public_connect_spec("orders", Ensure::Present)
.ownership_claims()
.unwrap();
let right = public_connect_spec("billing", Ensure::Present)
.ownership_claims()
.unwrap();
assert!(!left.overlaps(&right));
}
#[test]
fn ownership_claims_partial_role_overlap() {
let mut left = OwnershipClaims::default();
left.roles.insert("analytics".to_string());
left.roles.insert("reporting-viewer".to_string());
let mut right = OwnershipClaims::default();
right.roles.insert("analytics".to_string());
right.roles.insert("other-role".to_string());
assert!(left.overlaps(&right));
let summary = left.overlap_summary(&right);
assert!(summary.contains("roles: analytics"));
assert!(!summary.contains("schemas"));
}
#[test]
fn ownership_claims_empty_is_disjoint() {
let left = OwnershipClaims::default();
let right = OwnershipClaims::default();
assert!(!left.overlaps(&right));
}
#[test]
fn database_identity_equality() {
let conn_a = ConnectionSpec {
secret_ref: Some(SecretReference {
name: "db-creds".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
};
let a = DatabaseIdentity::from_connection("prod", &conn_a);
let b = DatabaseIdentity::from_connection("prod", &conn_a);
let c = DatabaseIdentity::from_connection("staging", &conn_a);
assert_eq!(a, b);
assert_ne!(a, c);
}
#[test]
fn database_identity_different_key() {
let conn_a = ConnectionSpec {
secret_ref: Some(SecretReference {
name: "db-creds".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
};
let conn_b = ConnectionSpec {
secret_ref: Some(SecretReference {
name: "db-creds".to_string(),
}),
secret_key: Some("CUSTOM_URL".to_string()),
params: None,
require_physical_identity: None,
};
let a = DatabaseIdentity::from_connection("prod", &conn_a);
let b = DatabaseIdentity::from_connection("prod", &conn_b);
assert_ne!(a, b);
}
#[test]
fn status_default_has_empty_conditions() {
let status = PostgresPolicyStatus::default();
assert!(status.conditions.is_empty());
assert!(status.observed_generation.is_none());
assert!(status.last_attempted_generation.is_none());
assert!(status.last_successful_reconcile_time.is_none());
assert!(status.change_summary.is_none());
assert!(status.managed_database_identity.is_none());
assert!(status.owned_roles.is_empty());
assert!(status.owned_schemas.is_empty());
assert!(status.last_error.is_none());
assert!(status.applied_password_source_versions.is_empty());
}
#[test]
fn status_degraded_workflow_sets_ready_false_and_degraded_true() {
let mut status = PostgresPolicyStatus::default();
status.set_condition(ready_condition(false, "InvalidSpec", "bad manifest"));
status.set_condition(degraded_condition("InvalidSpec", "bad manifest"));
status
.conditions
.retain(|c| c.condition_type != "Reconciling" && c.condition_type != "Paused");
status.change_summary = None;
status.last_error = Some("bad manifest".to_string());
let ready = status
.conditions
.iter()
.find(|c| c.condition_type == "Ready")
.expect("should have Ready condition");
assert_eq!(ready.status, "False");
assert_eq!(ready.reason.as_deref(), Some("InvalidSpec"));
let degraded = status
.conditions
.iter()
.find(|c| c.condition_type == "Degraded")
.expect("should have Degraded condition");
assert_eq!(degraded.status, "True");
assert_eq!(degraded.reason.as_deref(), Some("InvalidSpec"));
assert_eq!(status.last_error.as_deref(), Some("bad manifest"));
}
#[test]
fn status_conflict_workflow() {
let mut status = PostgresPolicyStatus::default();
let msg = "policy ownership overlaps with staging/other on database target prod/db/URL";
status.set_condition(ready_condition(false, "ConflictingPolicy", msg));
status.set_condition(conflict_condition("ConflictingPolicy", msg));
status.set_condition(degraded_condition("ConflictingPolicy", msg));
status
.conditions
.retain(|c| c.condition_type != "Reconciling");
status.last_error = Some(msg.to_string());
let conflict = status
.conditions
.iter()
.find(|c| c.condition_type == "Conflict")
.expect("should have Conflict condition");
assert_eq!(conflict.status, "True");
assert_eq!(conflict.reason.as_deref(), Some("ConflictingPolicy"));
let ready = status
.conditions
.iter()
.find(|c| c.condition_type == "Ready")
.expect("should have Ready condition");
assert_eq!(ready.status, "False");
let degraded = status
.conditions
.iter()
.find(|c| c.condition_type == "Degraded")
.expect("should have Degraded condition");
assert_eq!(degraded.status, "True");
}
#[test]
fn status_successful_reconcile_records_generation_and_time() {
let mut status = PostgresPolicyStatus::default();
let generation = Some(3_i64);
let summary = ChangeSummary {
roles_created: 2,
total: 2,
..Default::default()
};
status.set_condition(ready_condition(true, "Reconciled", "All changes applied"));
status.conditions.retain(|c| {
c.condition_type != "Reconciling"
&& c.condition_type != "Degraded"
&& c.condition_type != "Conflict"
&& c.condition_type != "Paused"
});
status.observed_generation = generation;
status.last_attempted_generation = generation;
status.last_successful_reconcile_time = Some(now_rfc3339());
status.change_summary = Some(summary);
status.last_error = None;
let ready = status
.conditions
.iter()
.find(|c| c.condition_type == "Ready")
.expect("should have Ready condition");
assert_eq!(ready.status, "True");
assert_eq!(ready.reason.as_deref(), Some("Reconciled"));
assert_eq!(status.observed_generation, Some(3));
assert_eq!(status.last_attempted_generation, Some(3));
assert!(status.last_successful_reconcile_time.is_some());
let summary = status.change_summary.as_ref().unwrap();
assert_eq!(summary.roles_created, 2);
assert_eq!(summary.total, 2);
assert!(status.last_error.is_none());
assert!(
status
.conditions
.iter()
.all(|c| c.condition_type != "Degraded"
&& c.condition_type != "Conflict"
&& c.condition_type != "Paused"
&& c.condition_type != "Reconciling")
);
}
#[test]
fn status_suspended_workflow() {
let mut status = PostgresPolicyStatus::default();
let generation = Some(2_i64);
status.set_condition(paused_condition("Reconciliation suspended by spec"));
status.set_condition(ready_condition(
false,
"Suspended",
"Reconciliation suspended by spec",
));
status
.conditions
.retain(|c| c.condition_type != "Reconciling");
status.last_attempted_generation = generation;
status.last_error = None;
let paused = status
.conditions
.iter()
.find(|c| c.condition_type == "Paused")
.expect("should have Paused condition");
assert_eq!(paused.status, "True");
let ready = status
.conditions
.iter()
.find(|c| c.condition_type == "Ready")
.expect("should have Ready condition");
assert_eq!(ready.status, "False");
assert_eq!(ready.reason.as_deref(), Some("Suspended"));
assert!(
!status
.conditions
.iter()
.any(|c| c.condition_type == "Reconciling")
);
}
#[test]
fn status_transitions_from_degraded_to_ready() {
let mut status = PostgresPolicyStatus::default();
status.set_condition(ready_condition(false, "InvalidSpec", "error"));
status.set_condition(degraded_condition("InvalidSpec", "error"));
status.last_error = Some("error".to_string());
assert_eq!(status.conditions.len(), 2);
status.set_condition(ready_condition(true, "Reconciled", "All changes applied"));
status.conditions.retain(|c| {
c.condition_type != "Reconciling"
&& c.condition_type != "Degraded"
&& c.condition_type != "Conflict"
&& c.condition_type != "Paused"
});
status.last_error = None;
let ready = status
.conditions
.iter()
.find(|c| c.condition_type == "Ready")
.expect("should have Ready condition");
assert_eq!(ready.status, "True");
assert!(
!status
.conditions
.iter()
.any(|c| c.condition_type == "Degraded")
);
assert_eq!(status.conditions.len(), 1);
assert!(status.last_error.is_none());
}
#[test]
fn change_summary_default_is_all_zero() {
let summary = ChangeSummary::default();
assert_eq!(summary.roles_created, 0);
assert_eq!(summary.roles_altered, 0);
assert_eq!(summary.roles_dropped, 0);
assert_eq!(summary.sessions_terminated, 0);
assert_eq!(summary.grants_added, 0);
assert_eq!(summary.grants_revoked, 0);
assert_eq!(summary.default_privileges_set, 0);
assert_eq!(summary.default_privileges_revoked, 0);
assert_eq!(summary.members_added, 0);
assert_eq!(summary.members_removed, 0);
assert_eq!(summary.total, 0);
}
#[test]
fn status_serializes_to_json() {
let mut status = PostgresPolicyStatus::default();
status.set_condition(ready_condition(true, "Reconciled", "done"));
status.observed_generation = Some(5);
status.managed_database_identity = Some("ns/secret/key".to_string());
status.owned_roles = vec!["role-a".to_string(), "role-b".to_string()];
status.owned_schemas = vec!["public".to_string()];
status.change_summary = Some(ChangeSummary {
roles_created: 1,
total: 1,
..Default::default()
});
let json = serde_json::to_string(&status).expect("should serialize");
assert!(json.contains("\"Reconciled\""));
assert!(json.contains("\"observed_generation\":5"));
assert!(json.contains("\"role-a\""));
assert!(json.contains("\"ns/secret/key\""));
}
#[test]
fn crd_spec_deserializes_from_yaml() {
let yaml = r#"
connection:
secretRef:
name: pg-credentials
interval: "10m"
default_owner: app_owner
profiles:
editor:
grants:
- privileges: [USAGE]
object: { type: schema }
- privileges: [SELECT, INSERT, UPDATE, DELETE]
object: { type: table, name: "*" }
default_privileges:
- privileges: [SELECT, INSERT, UPDATE, DELETE]
on_type: table
schemas:
- name: inventory
profiles: [editor]
roles:
- name: analytics
login: true
grants:
- role: analytics
privileges: [CONNECT]
object: { type: database, name: mydb }
memberships:
- role: inventory-editor
members:
- name: analytics
retirements:
- role: legacy-app
reassign_owned_to: app_owner
drop_owned: true
terminate_sessions: true
"#;
let spec: PostgresPolicySpec = serde_yaml::from_str(yaml).expect("should deserialize");
assert_eq!(spec.interval, "10m");
assert_eq!(spec.default_owner, Some("app_owner".to_string()));
assert_eq!(spec.profiles.len(), 1);
assert!(spec.profiles.contains_key("editor"));
assert_eq!(spec.schemas.len(), 1);
assert_eq!(spec.roles.len(), 1);
assert_eq!(spec.grants.len(), 1);
assert_eq!(spec.memberships.len(), 1);
assert_eq!(spec.retirements.len(), 1);
assert_eq!(spec.retirements[0].role, "legacy-app");
assert!(spec.retirements[0].terminate_sessions);
}
#[test]
fn referenced_secret_names_includes_connection_secret() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
let names = spec.referenced_secret_names("test-policy");
assert!(names.contains("pg-conn"));
assert_eq!(names.len(), 1);
}
#[test]
fn referenced_secret_names_includes_password_secrets() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![
RoleSpec {
name: "role-a".to_string(),
external: false,
login: Some(true),
password: Some(PasswordSpec {
secret_ref: Some(SecretReference {
name: "role-passwords".to_string(),
}),
secret_key: Some("role-a".to_string()),
generate: None,
}),
password_valid_until: None,
config: Default::default(),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
},
RoleSpec {
name: "role-b".to_string(),
external: false,
login: Some(true),
password: Some(PasswordSpec {
secret_ref: Some(SecretReference {
name: "other-secret".to_string(),
}),
secret_key: None,
generate: None,
}),
password_valid_until: None,
config: Default::default(),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
},
RoleSpec {
name: "role-c".to_string(),
external: false,
login: None,
password: None,
password_valid_until: None,
config: Default::default(),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
},
],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
let names = spec.referenced_secret_names("test-policy");
assert!(
names.contains("pg-conn"),
"should include connection secret"
);
assert!(
names.contains("role-passwords"),
"should include role-a password secret"
);
assert!(
names.contains("other-secret"),
"should include role-b password secret"
);
assert_eq!(names.len(), 3);
}
#[test]
fn validate_password_specs_rejects_password_without_login() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: Some(false),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: Some(SecretReference {
name: "role-passwords".to_string(),
}),
secret_key: None,
generate: None,
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::PasswordWithoutLogin { ref role }) if role == "app-user"
));
}
#[test]
fn validate_password_specs_rejects_password_with_login_omitted() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: None, superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: Some(SecretReference {
name: "role-passwords".to_string(),
}),
secret_key: None,
generate: None,
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::PasswordWithoutLogin { ref role }) if role == "app-user"
));
}
#[test]
fn validate_password_specs_rejects_invalid_password_mode() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: Some(SecretReference {
name: "role-passwords".to_string(),
}),
secret_key: None,
generate: Some(GeneratePasswordSpec {
length: Some(32),
secret_name: None,
secret_key: None,
}),
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::InvalidPasswordMode { ref role }) if role == "app-user"
));
}
#[test]
fn validate_password_specs_rejects_invalid_generated_length() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: None,
secret_key: None,
generate: Some(GeneratePasswordSpec {
length: Some(8),
secret_name: None,
secret_key: None,
}),
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::InvalidGeneratedLength { ref role, .. }) if role == "app-user"
));
}
#[test]
fn validate_password_specs_rejects_invalid_generated_secret_key() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: None,
secret_key: None,
generate: Some(GeneratePasswordSpec {
length: Some(32),
secret_name: None,
secret_key: Some("bad/key".to_string()),
}),
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::InvalidSecretKey { ref role, field, .. })
if role == "app-user" && field == "generate.secretKey"
));
}
#[test]
fn validate_password_specs_rejects_invalid_generated_secret_name() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: None,
secret_key: None,
generate: Some(GeneratePasswordSpec {
length: Some(32),
secret_name: Some("Bad_Name".to_string()),
secret_key: None,
}),
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::InvalidGeneratedSecretName { ref role, .. }) if role == "app-user"
));
}
#[test]
fn validate_password_specs_rejects_reserved_generated_secret_key() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "pg-conn".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: std::collections::HashMap::new(),
schemas: vec![],
roles: vec![RoleSpec {
name: "app-user".to_string(),
external: false,
login: Some(true),
superuser: None,
createdb: None,
createrole: None,
inherit: None,
replication: None,
bypassrls: None,
connection_limit: None,
comment: None,
password: Some(PasswordSpec {
secret_ref: None,
secret_key: None,
generate: Some(GeneratePasswordSpec {
length: Some(32),
secret_name: None,
secret_key: Some("verifier".to_string()),
}),
}),
password_valid_until: None,
config: Default::default(),
}],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert!(matches!(
spec.validate_password_specs("test-policy"),
Err(PasswordValidationError::ReservedGeneratedSecretKey { ref role, ref key })
if role == "app-user" && key == "verifier"
));
}
#[test]
fn plan_crd_generates_valid_schema() {
let crd = PostgresPolicyPlan::crd();
let yaml = serde_yaml::to_string(&crd).expect("CRD should serialize to YAML");
assert!(yaml.contains("pgroles.io"), "group should be pgroles.io");
assert!(yaml.contains("v1alpha1"), "version should be v1alpha1");
assert!(
yaml.contains("PostgresPolicyPlan"),
"kind should be PostgresPolicyPlan"
);
assert!(yaml.contains("pgplan"), "should have shortname pgplan");
}
#[test]
fn plan_phase_display() {
assert_eq!(PlanPhase::Pending.to_string(), "Pending");
assert_eq!(PlanPhase::Approved.to_string(), "Approved");
assert_eq!(PlanPhase::Applying.to_string(), "Applying");
assert_eq!(PlanPhase::Applied.to_string(), "Applied");
assert_eq!(PlanPhase::Failed.to_string(), "Failed");
assert_eq!(PlanPhase::Superseded.to_string(), "Superseded");
}
#[test]
fn plan_phase_default_is_pending() {
assert_eq!(PlanPhase::default(), PlanPhase::Pending);
}
#[test]
fn effective_approval_infers_from_mode() {
let base = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "test".into(),
}),
secret_key: Some("DATABASE_URL".into()),
params: None,
require_physical_identity: None,
},
interval: "5m".into(),
suspend: false,
mode: PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::Authoritative,
default_owner: None,
profiles: Default::default(),
schemas: vec![],
roles: vec![],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: None,
};
assert_eq!(base.effective_approval(), ApprovalMode::Auto);
let plan = PostgresPolicySpec {
mode: PolicyMode::Observe,
..base.clone()
};
assert_eq!(plan.effective_approval(), ApprovalMode::Manual);
let explicit = PostgresPolicySpec {
approval: Some(ApprovalMode::Manual),
..base.clone()
};
assert_eq!(explicit.effective_approval(), ApprovalMode::Manual);
}
#[test]
fn approval_mode_serde_roundtrip() {
let manual: ApprovalMode = serde_json::from_str("\"manual\"").unwrap();
assert_eq!(manual, ApprovalMode::Manual);
let auto: ApprovalMode = serde_json::from_str("\"auto\"").unwrap();
assert_eq!(auto, ApprovalMode::Auto);
let manual_json = serde_json::to_value(&ApprovalMode::Manual).unwrap();
assert_eq!(manual_json, serde_json::Value::String("manual".to_string()));
let auto_json = serde_json::to_value(&ApprovalMode::Auto).unwrap();
assert_eq!(auto_json, serde_json::Value::String("auto".to_string()));
}
#[test]
fn the_legacy_plan_mode_value_is_accepted_and_behaves_as_observe() {
let legacy: PolicyMode = serde_json::from_str("\"plan\"").unwrap();
assert!(legacy.never_executes());
assert!(legacy.is_deprecated_spelling());
assert!(PolicyMode::Observe.never_executes());
assert!(!PolicyMode::Observe.is_deprecated_spelling());
assert!(!PolicyMode::Apply.never_executes());
assert_eq!(
serde_json::to_value(PolicyMode::Plan).unwrap(),
serde_json::Value::String("plan".to_string())
);
let schema = serde_json::to_value(schemars::schema_for!(PolicyMode)).unwrap();
let rendered = schema.to_string();
assert!(rendered.contains("observe"));
assert!(
rendered.contains("\"plan\""),
"the schema must keep accepting the deprecated value during the \
deprecation window: {rendered}"
);
}
#[test]
fn plan_status_default_is_empty() {
let status = PostgresPolicyPlanStatus::default();
assert_eq!(status.phase, PlanPhase::Pending);
assert!(status.conditions.is_empty());
assert!(status.change_summary.is_none());
assert!(status.sql_ref.is_none());
assert!(status.sql_inline.is_none());
assert!(status.computed_at.is_none());
assert!(status.applied_at.is_none());
assert!(status.last_error.is_none());
}
#[test]
fn spec_without_approval_field_deserializes_as_none() {
let json = serde_json::json!({
"connection": {
"secretRef": { "name": "pg-secret" },
"secretKey": "DATABASE_URL"
},
"interval": "5m",
"suspend": false,
"mode": "apply",
"reconciliation_mode": "authoritative"
});
let spec: PostgresPolicySpec =
serde_json::from_value(json).expect("should deserialize without approval field");
assert!(
spec.approval.is_none(),
"approval should be None when omitted"
);
assert_eq!(
spec.effective_approval(),
ApprovalMode::Auto,
"effective_approval should infer Auto from apply mode"
);
}
#[test]
fn status_without_current_plan_ref_deserializes_as_none() {
let json = serde_json::json!({
"conditions": [],
"owned_roles": [],
"owned_schemas": []
});
let status: PostgresPolicyStatus =
serde_json::from_value(json).expect("should deserialize without current_plan_ref");
assert!(
status.current_plan_ref.is_none(),
"current_plan_ref should be None when omitted"
);
}
#[test]
fn effective_approval_explicit_auto_overrides_plan_mode() {
let spec = PostgresPolicySpec {
connection: ConnectionSpec {
secret_ref: Some(SecretReference {
name: "test".into(),
}),
secret_key: Some("DATABASE_URL".into()),
params: None,
require_physical_identity: None,
},
interval: "5m".into(),
suspend: false,
mode: PolicyMode::Observe,
reconciliation_mode: CrdReconciliationMode::Authoritative,
default_owner: None,
profiles: Default::default(),
schemas: vec![],
roles: vec![],
grants: vec![],
default_privileges: vec![],
memberships: vec![],
retirements: vec![],
approval: Some(ApprovalMode::Auto),
};
assert_eq!(
spec.effective_approval(),
ApprovalMode::Auto,
"explicit Auto should override Observe mode's default of Manual"
);
}
#[test]
fn plan_phase_rejected_display() {
assert_eq!(PlanPhase::Rejected.to_string(), "Rejected");
}
#[test]
fn plan_phase_all_variants_display() {
let variants = [
PlanPhase::Pending,
PlanPhase::Approved,
PlanPhase::Applying,
PlanPhase::Applied,
PlanPhase::Failed,
PlanPhase::Superseded,
PlanPhase::Rejected,
];
for variant in &variants {
let display = variant.to_string();
assert!(
!display.is_empty(),
"PlanPhase::{variant:?} should have non-empty Display output"
);
}
}
#[test]
fn plan_status_defaults() {
let status = PostgresPolicyPlanStatus::default();
assert_eq!(status.phase, PlanPhase::Pending);
assert!(status.conditions.is_empty());
assert!(status.sql_ref.is_none());
assert!(status.sql_hash.is_none());
assert!(status.sql_inline.is_none());
assert!(!status.sql_truncated);
assert!(status.redacted_sql_hash.is_none());
assert!(status.sql_original_bytes.is_none());
assert!(status.sql_stored_bytes.is_none());
assert!(status.change_summary.is_none());
assert!(status.computed_at.is_none());
assert!(status.applied_at.is_none());
assert!(status.last_error.is_none());
}
#[test]
fn sql_ref_missing_compression_deserializes_as_uncompressed_legacy_shape() {
let json = serde_json::json!({
"name": "legacy-plan-sql",
"key": "plan.sql"
});
let sql_ref: SqlRef = serde_json::from_value(json).expect("legacy SqlRef should decode");
assert_eq!(sql_ref.name, "legacy-plan-sql");
assert_eq!(sql_ref.key, "plan.sql");
assert_eq!(sql_ref.compression, None);
}
#[test]
fn plan_spec_camel_case_serialization() {
let spec = PostgresPolicyPlanSpec {
policy_ref: PolicyPlanRef {
name: "my-policy".into(),
},
policy_generation: 3,
reconciliation_mode: CrdReconciliationMode::Authoritative,
owned_roles: vec!["role-a".into()],
owned_schemas: vec!["public".into()],
managed_database_identity: "ns/secret/key".into(),
origin: None,
scope: None,
};
let json = serde_json::to_value(&spec).expect("should serialize to JSON");
let obj = json.as_object().expect("should be a JSON object");
assert!(
obj.contains_key("policyRef"),
"should use camelCase: policyRef"
);
assert!(
obj.contains_key("policyGeneration"),
"should use camelCase: policyGeneration"
);
assert!(
obj.contains_key("reconciliationMode"),
"should use camelCase: reconciliationMode"
);
assert!(
obj.contains_key("ownedRoles"),
"should use camelCase: ownedRoles"
);
assert!(
obj.contains_key("ownedSchemas"),
"should use camelCase: ownedSchemas"
);
assert!(
obj.contains_key("managedDatabaseIdentity"),
"should use camelCase: managedDatabaseIdentity"
);
}
fn resolved_bundle(memberships: Vec<ResolvedEphemeralMembership>) -> ResolvedEphemeralAccess {
ResolvedEphemeralAccess {
access_policy_uid: "access-uid".into(),
access_policy_generation: 1,
target_policy_uid: "target-uid".into(),
target_policy_generation: 2,
target_database_fingerprint: "sha256:database".into(),
granted_duration: "1800s".into(),
bundle_encoding: EPHEMERAL_BUNDLE_ENCODING_V1.into(),
bundle_hash: String::new(),
memberships,
}
}
#[test]
fn ephemeral_bundle_hash_is_order_independent_and_uid_independent() {
let first = ResolvedEphemeralMembership {
role: "editor".into(),
member: "alice@example.com".into(),
inherit: false,
};
let second = ResolvedEphemeralMembership {
role: "auditor".into(),
member: "alice@example.com".into(),
inherit: true,
};
let original = resolved_bundle(vec![first.clone(), second.clone()]);
let mut reordered = resolved_bundle(vec![second, first]);
reordered.access_policy_uid = "replacement-uid".into();
reordered.target_policy_generation = 99;
assert_eq!(
original.compute_bundle_hash(),
reordered.compute_bundle_hash()
);
}
#[test]
fn ephemeral_bundle_hash_covers_membership_options() {
let original = resolved_bundle(vec![ResolvedEphemeralMembership {
role: "editor".into(),
member: "alice@example.com".into(),
inherit: false,
}]);
let changed = resolved_bundle(vec![ResolvedEphemeralMembership {
role: "editor".into(),
member: "alice@example.com".into(),
inherit: true,
}]);
assert_ne!(
original.compute_bundle_hash(),
changed.compute_bundle_hash()
);
}
#[test]
fn ephemeral_bundle_hash_covers_database_target() {
let original = resolved_bundle(Vec::new());
let mut retargeted = original.clone();
retargeted.target_database_fingerprint = "sha256:other-database".into();
assert_ne!(
original.compute_bundle_hash(),
retargeted.compute_bundle_hash()
);
}
#[test]
fn ephemeral_request_crd_contains_immutability_rules() {
let json = serde_json::to_string(&EphemeralAccessRequest::crd())
.expect("request CRD should serialize");
assert!(json.contains("request spec is immutable"));
assert!(json.contains("resolvedAccess is write-once"));
assert!(json.contains("approval decisions are terminal"));
assert!(json.contains("decision identity is write-once"));
assert!(json.contains(
"a terminal approval decision and decidedBy identity must be recorded together"
));
assert!(json.contains(r#""requestedBy""#));
assert!(json.contains(r#""decidedBy""#));
assert!(json.contains(r#""maxItems":8"#));
assert!(json.contains(r#""pattern":"^([0-9]+[smh])+$""#));
}
fn assert_bounded_strings_and_collections(schema: &serde_json::Value, path: &str) {
match schema {
serde_json::Value::Object(object) => {
match object.get("type").and_then(serde_json::Value::as_str) {
Some("string") if !object.contains_key("enum") => {
assert!(
object.contains_key("maxLength"),
"unbounded string schema at {path}"
);
}
Some("array") => {
assert!(
object.contains_key("maxItems"),
"unbounded collection schema at {path}"
);
}
_ => {}
}
for (name, child) in object {
assert_bounded_strings_and_collections(child, &format!("{path}.{name}"));
}
}
serde_json::Value::Array(items) => {
for (index, child) in items.iter().enumerate() {
assert_bounded_strings_and_collections(child, &format!("{path}[{index}]"));
}
}
_ => {}
}
}
#[test]
fn ephemeral_crds_bound_every_string_and_collection() {
for (kind, crd) in [
("EphemeralAccessPolicy", EphemeralAccessPolicy::crd()),
("EphemeralAccessRequest", EphemeralAccessRequest::crd()),
] {
let value = serde_json::to_value(crd).expect("CRD should serialize");
let schema = &value["spec"]["versions"][0]["schema"]["openAPIV3Schema"];
assert_bounded_strings_and_collections(schema, kind);
}
}
#[test]
fn candidate_spec_bounds_every_string_collection_and_map() {
let value =
serde_json::to_value(postgres_policy_candidate_crd()).expect("CRD should serialize");
let spec = &value["spec"]["versions"][0]["schema"]["openAPIV3Schema"]["properties"]["spec"];
assert_bounded_strings_and_collections(spec, "PostgresPolicyCandidate.spec");
assert_bounded_maps(spec, "PostgresPolicyCandidate.spec");
}
fn assert_bounded_maps(schema: &serde_json::Value, path: &str) {
if let serde_json::Value::Object(object) = schema {
if object.get("type").and_then(serde_json::Value::as_str) == Some("object")
&& object.contains_key("additionalProperties")
{
assert!(
object.contains_key("maxProperties"),
"unbounded map schema at {path}"
);
}
for (name, child) in object {
assert_bounded_maps(child, &format!("{path}.{name}"));
}
}
}
#[test]
fn candidate_content_emits_no_openapi_defaults() {
let value =
serde_json::to_value(postgres_policy_candidate_crd()).expect("CRD should serialize");
let content = &value["spec"]["versions"][0]["schema"]["openAPIV3Schema"]["properties"]["spec"]
["properties"]["content"];
assert!(content.is_object(), "spec.content should be in the schema");
fn find_default(schema: &serde_json::Value, path: &str, found: &mut Vec<String>) {
match schema {
serde_json::Value::Object(object) => {
if object.contains_key("default") {
found.push(path.to_string());
}
for (name, child) in object {
find_default(child, &format!("{path}.{name}"), found);
}
}
serde_json::Value::Array(items) => {
for (index, child) in items.iter().enumerate() {
find_default(child, &format!("{path}[{index}]"), found);
}
}
_ => {}
}
}
let mut found = Vec::new();
find_default(content, "spec.content", &mut found);
assert!(
found.is_empty(),
"spec.content must emit no OpenAPI defaults, found: {found:?}"
);
}
#[test]
fn candidate_content_keeps_serde_defaults() {
let content: PolicyContent =
serde_json::from_str("{}").expect("empty content should deserialize");
assert_eq!(
content.reconciliation_mode,
CrdReconciliationMode::default()
);
assert!(content.roles.is_empty());
assert!(content.grants.is_empty());
assert!(content.profiles.is_empty());
}
#[test]
fn promoting_candidate_content_yields_the_candidates_digest() {
let content: PolicyContent = serde_json::from_value(serde_json::json!({
"reconciliation_mode": "additive",
"default_owner": "app_owner",
"profiles": { "reader": { "grants": [] } },
"schemas": [{ "name": "app", "profiles": ["reader"] }],
"roles": [{ "name": "reporting-reader", "login": true }],
"grants": [{
"role": "reporting-reader",
"privileges": ["CONNECT"],
"object": { "type": "database", "name": "orders" }
}],
"memberships": [{ "role": "reporting-reader", "members": [{ "name": "app_owner" }] }],
}))
.expect("content fixture");
let mut spec_json = serde_json::to_value(&content).expect("content serializes");
let object = spec_json.as_object_mut().expect("content is an object");
object.insert(
"connection".to_string(),
serde_json::json!({ "secretRef": { "name": "db" } }),
);
object.insert("interval".to_string(), serde_json::json!("30s"));
object.insert("mode".to_string(), serde_json::json!("apply"));
object.insert("approval".to_string(), serde_json::json!("manual"));
object.insert("suspend".to_string(), serde_json::json!(false));
let spec: PostgresPolicySpec =
serde_json::from_value(spec_json).expect("promoted policy spec");
assert_eq!(spec.content_digest(), content.content_digest());
assert_eq!(
spec.content_digest(),
pgroles_core::candidate::compute_content_digest(&content),
);
let mut other = spec.clone();
other.interval = "1h".to_string();
other.suspend = true;
assert_eq!(other.content_digest(), spec.content_digest());
let mut edited = spec.clone();
edited.roles[0].login = Some(false);
assert_ne!(edited.content_digest(), spec.content_digest());
}
#[test]
fn candidate_content_matches_policy_content() {
let policy = serde_json::to_value(PostgresPolicy::crd()).expect("CRD should serialize");
let policy_spec = &policy["spec"]["versions"][0]["schema"]["openAPIV3Schema"]["properties"]
["spec"]["properties"];
let candidate =
serde_json::to_value(postgres_policy_candidate_crd()).expect("CRD should serialize");
let content = &candidate["spec"]["versions"][0]["schema"]["openAPIV3Schema"]["properties"]
["spec"]["properties"]["content"]["properties"];
let execution = ["connection", "interval", "mode", "suspend", "approval"];
let policy_content: Vec<&String> = policy_spec
.as_object()
.expect("policy spec has properties")
.keys()
.filter(|k| !execution.contains(&k.as_str()))
.collect();
let mut candidate_content: Vec<&String> = content
.as_object()
.expect("candidate content has properties")
.keys()
.collect();
candidate_content.sort();
let mut policy_content = policy_content;
policy_content.sort();
assert_eq!(policy_content, candidate_content);
for field in &candidate_content {
let policy_field = &policy_spec[field.as_str()];
let candidate_field = &content[field.as_str()];
assert_eq!(
policy_field["type"], candidate_field["type"],
"spec.{field} and spec.content.{field} must have the same type"
);
assert_eq!(
policy_field["items"]["properties"]
.as_object()
.map(|o| o.keys().collect::<Vec<_>>()),
candidate_field["items"]["properties"]
.as_object()
.map(|o| o.keys().collect::<Vec<_>>()),
"spec.{field} and spec.content.{field} must have the same item fields"
);
}
}
#[test]
fn candidate_crd_exposes_operational_columns() {
let crd = postgres_policy_candidate_crd();
let columns: Vec<&str> = crd.spec.versions[0]
.additional_printer_columns
.as_deref()
.unwrap_or_default()
.iter()
.map(|c| c.name.as_str())
.collect();
assert_eq!(columns, ["Policy", "Phase", "Plan", "Digest", "Age"]);
assert_eq!(
crd.spec.names.short_names.as_deref(),
Some(["pgcand".to_string()].as_slice())
);
assert_eq!(
crd.spec.names.categories.as_deref(),
Some(["pgroles".to_string()].as_slice())
);
assert!(
crd.spec.versions[0]
.subresources
.as_ref()
.and_then(|s| s.status.as_ref())
.is_some()
);
}
#[test]
fn ephemeral_policy_crd_exposes_operational_columns() {
let json = serde_json::to_string(&EphemeralAccessPolicy::crd())
.expect("policy CRD should serialize");
assert!(json.contains("postgresPolicyRef"));
assert!(json.contains("maximumDuration"));
assert!(json.contains("Accepted"));
assert!(json.contains(r#""minItems":1"#));
assert!(json.contains(r#""pattern":"^([0-9]+[smh])+$""#));
}
}