use std::collections::{BTreeMap, BTreeSet};
use std::io::Write;
use std::time::Duration;
use flate2::Compression;
use flate2::write::GzEncoder;
use k8s_openapi::ByteString;
use k8s_openapi::api::core::v1::ConfigMap;
use k8s_openapi::apimachinery::pkg::apis::meta::v1::OwnerReference;
use kube::api::{Api, DeleteParams, ListParams, Patch, PatchParams, PostParams};
use kube::core::labels::{Expression, Selector};
use kube::{Client, Resource, ResourceExt};
use sha2::{Digest, Sha256};
use tracing::info;
use crate::crd::{
ChangeSummary, CrdReconciliationMode, LABEL_CANDIDATE, LABEL_DATABASE_IDENTITY, LABEL_PLAN,
LABEL_POLICY, PlanOrigin, PlanPhase, PlanReference, PolicyCondition, PolicyPlanRef,
PostgresPolicy, PostgresPolicyCandidate, PostgresPolicyPlan, PostgresPolicyPlanSpec,
PostgresPolicyPlanStatus, SqlCompression, SqlRef, is_retention_exempt,
};
use crate::k8s_names::{LabelValue, truncate_name_prefix};
use crate::reconciler::ReconcileError;
use pgroles_core::approval::{APPROVAL_EFFECT_ENCODING_V3, TargetIdentity};
#[derive(Debug, Clone)]
pub enum PlanCreationResult {
Created(String),
Deduplicated(String),
DeduplicatedFailed(String),
}
impl PlanCreationResult {
pub fn plan_name(&self) -> &str {
match self {
PlanCreationResult::Created(name)
| PlanCreationResult::Deduplicated(name)
| PlanCreationResult::DeduplicatedFailed(name) => name,
}
}
pub fn is_created(&self) -> bool {
matches!(self, PlanCreationResult::Created(_))
}
pub fn is_failed_backoff(&self) -> bool {
matches!(self, PlanCreationResult::DeduplicatedFailed(_))
}
}
const MAX_INLINE_SQL_BYTES: usize = 16 * 1024;
const SQL_CONFIGMAP_GZIP_KEY: &str = "plan.sql.gz";
const MAX_CONFIGMAP_SQL_BYTES: usize = 900 * 1024;
const ORPHAN_GRACE_SECS: i64 = 60;
const CLEANUP_TIMEOUT_SECS: u64 = 5;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PlanRetention {
pub applied: usize,
pub applied_ceiling: usize,
pub applied_min_age_secs: i64,
pub decided: usize,
pub superseded: usize,
}
impl Default for PlanRetention {
fn default() -> Self {
Self {
applied: 25,
applied_ceiling: 200,
applied_min_age_secs: 30 * 24 * 60 * 60,
decided: 10,
superseded: 3,
}
}
}
impl PlanRetention {
pub const ENV_APPLIED: &'static str = "PLAN_RETENTION_APPLIED";
pub const ENV_APPLIED_CEILING: &'static str = "PLAN_RETENTION_APPLIED_CEILING";
pub const ENV_APPLIED_MIN_AGE: &'static str = "PLAN_RETENTION_APPLIED_MIN_AGE";
pub const ENV_DECIDED: &'static str = "PLAN_RETENTION_DECIDED";
pub const ENV_SUPERSEDED: &'static str = "PLAN_RETENTION_SUPERSEDED";
pub fn from_env() -> Result<Self, PlanRetentionConfigError> {
let mut resolved = std::collections::BTreeMap::new();
for variable in [
Self::ENV_APPLIED,
Self::ENV_APPLIED_CEILING,
Self::ENV_APPLIED_MIN_AGE,
Self::ENV_DECIDED,
Self::ENV_SUPERSEDED,
] {
if let Some(value) = env_read(variable, std::env::var(variable))? {
resolved.insert(variable, value);
}
}
Self::from_lookup(|variable| resolved.get(variable).cloned())
}
fn from_lookup(
lookup: impl Fn(&str) -> Option<String>,
) -> Result<Self, PlanRetentionConfigError> {
let defaults = Self::default();
let retention = Self {
applied: count_bound(&lookup, Self::ENV_APPLIED, defaults.applied)?,
applied_ceiling: count_bound(
&lookup,
Self::ENV_APPLIED_CEILING,
defaults.applied_ceiling,
)?,
applied_min_age_secs: age_bound(
&lookup,
Self::ENV_APPLIED_MIN_AGE,
defaults.applied_min_age_secs,
)?,
decided: count_bound(&lookup, Self::ENV_DECIDED, defaults.decided)?,
superseded: count_bound(&lookup, Self::ENV_SUPERSEDED, defaults.superseded)?,
};
if retention.applied_ceiling < retention.applied {
return Err(PlanRetentionConfigError::CeilingBelowCount {
ceiling: retention.applied_ceiling,
count: retention.applied,
});
}
Ok(retention)
}
}
fn env_read(
variable: &'static str,
read: Result<String, std::env::VarError>,
) -> Result<Option<String>, PlanRetentionConfigError> {
match read {
Ok(value) => Ok(Some(value)),
Err(std::env::VarError::NotPresent) => Ok(None),
Err(std::env::VarError::NotUnicode(_)) => {
Err(PlanRetentionConfigError::NotUnicode { variable })
}
}
}
fn count_bound(
lookup: impl Fn(&str) -> Option<String>,
variable: &'static str,
default: usize,
) -> Result<usize, PlanRetentionConfigError> {
match lookup(variable) {
None => Ok(default),
Some(value) => value
.trim()
.parse()
.map_err(|_| PlanRetentionConfigError::InvalidCount { variable, value }),
}
}
fn age_bound(
lookup: impl Fn(&str) -> Option<String>,
variable: &'static str,
default: i64,
) -> Result<i64, PlanRetentionConfigError> {
let Some(value) = lookup(variable) else {
return Ok(default);
};
let invalid = |detail: String| PlanRetentionConfigError::InvalidDuration {
variable,
value: value.clone(),
detail,
};
let duration = crate::ephemeral::parse_duration(&value).map_err(|error| {
invalid(match error {
crate::ephemeral::EphemeralError::Invalid(detail) => detail,
other => other.to_string(),
})
})?;
i64::try_from(duration.as_secs()).map_err(|_| invalid("duration is too large".to_string()))
}
#[derive(Debug, thiserror::Error)]
pub enum PlanRetentionConfigError {
#[error("{variable} must be a non-negative integer, got {value:?}")]
InvalidCount {
variable: &'static str,
value: String,
},
#[error("{variable} is not a valid duration ({detail}), got {value:?}")]
InvalidDuration {
variable: &'static str,
value: String,
detail: String,
},
#[error(
"PLAN_RETENTION_APPLIED_CEILING ({ceiling}) is below PLAN_RETENTION_APPLIED ({count}); \
a ceiling under the count bound can never take effect"
)]
CeilingBelowCount { ceiling: usize, count: usize },
#[error("{variable} is set to a value that is not valid Unicode")]
NotUnicode { variable: &'static str },
}
const FAILED_PLAN_DEDUP_WINDOW_SECS: i64 = 120;
#[derive(Debug, Clone, PartialEq, Eq)]
enum PlanSqlArtifact {
Inline(String),
CompressedConfigMap {
configmap_name: String,
key: String,
compressed_sql: Vec<u8>,
},
TruncatedInline(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PreparedPlanSql {
artifact: PlanSqlArtifact,
redacted_sql_hash: String,
original_bytes: usize,
stored_bytes: usize,
}
impl PreparedPlanSql {
fn sql_ref(&self) -> Option<SqlRef> {
match &self.artifact {
PlanSqlArtifact::CompressedConfigMap {
configmap_name,
key,
..
} => Some(SqlRef {
name: configmap_name.clone(),
key: key.clone(),
compression: Some(SqlCompression::Gzip),
}),
PlanSqlArtifact::Inline(_) | PlanSqlArtifact::TruncatedInline(_) => None,
}
}
fn sql_inline(&self) -> Option<String> {
match &self.artifact {
PlanSqlArtifact::Inline(sql) | PlanSqlArtifact::TruncatedInline(sql) => {
Some(sql.clone())
}
PlanSqlArtifact::CompressedConfigMap { .. } => None,
}
}
fn is_truncated(&self) -> bool {
matches!(self.artifact, PlanSqlArtifact::TruncatedInline(_))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SupersedeCause {
EffectsChanged,
EffectsCleared,
ReplacedByNewerPlan,
TargetChanged(pgroles_core::approval::TargetIdentityReason),
PolicyStoppedPlanning,
SupersededByPromotion,
BaseContentChanged,
}
impl SupersedeCause {
pub fn reason(self) -> &'static str {
match self {
SupersedeCause::TargetChanged(reason) => reason.as_str(),
SupersedeCause::SupersededByPromotion => {
crate::crd::candidate_reason::SUPERSEDED_BY_PROMOTION
}
SupersedeCause::BaseContentChanged => "SupersededByBaseChange",
_ => "Superseded",
}
}
pub fn message(self) -> &'static str {
match self {
SupersedeCause::EffectsChanged => {
"the policy's effects changed since this plan was computed, so it no longer \
describes what would happen"
}
SupersedeCause::EffectsCleared => {
"the changes this plan described are no longer pending, so there is nothing left \
to execute"
}
SupersedeCause::ReplacedByNewerPlan => {
"a newer plan holds the current effects and replaces this one"
}
SupersedeCause::TargetChanged(reason) => reason.message(),
SupersedeCause::PolicyStoppedPlanning => {
"the policy no longer references this plan, so it will never be executed"
}
SupersedeCause::SupersededByPromotion => {
"another candidate's content was promoted and executed, so this plan describes a \
change against a base that no longer exists"
}
SupersedeCause::BaseContentChanged => {
"the policy's content changed after this candidate plan was computed against it; \
a fresh plan pinned to the current base replaces it and needs its own review"
}
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PlanApprovalState {
Pending,
Approved,
Rejected,
}
pub const AUTO_APPROVAL_ACTOR: &str = "system:pgroles-operator(auto-approval)";
pub fn check_plan_approval(plan: &PostgresPolicyPlan) -> PlanApprovalState {
let Some(status) = plan.status.as_ref() else {
return PlanApprovalState::Pending;
};
let decided = |condition_type: &str| {
status
.conditions
.iter()
.any(|c| c.condition_type == condition_type && c.status == "True")
};
if decided("Denied") {
return PlanApprovalState::Rejected;
}
if decided("Approved") {
return PlanApprovalState::Approved;
}
PlanApprovalState::Pending
}
#[derive(Debug, Clone, Copy)]
pub struct CandidatePlanBinding<'a> {
pub candidate: &'a PostgresPolicyCandidate,
pub content_digest: &'a str,
pub content_digest_encoding: &'a str,
pub base_content_digest: &'a str,
}
pub(crate) fn candidate_plan_origin(
binding: CandidatePlanBinding<'_>,
policy: &PostgresPolicy,
) -> PlanOrigin {
PlanOrigin {
kind: PostgresPolicyCandidate::kind(&()).to_string(),
name: binding.candidate.name_any(),
uid: binding.candidate.metadata.uid.clone().unwrap_or_default(),
content_digest: Some(binding.content_digest.to_string()),
content_digest_encoding: Some(binding.content_digest_encoding.to_string()),
policy_uid: policy.metadata.uid.clone(),
base_content_digest: Some(binding.base_content_digest.to_string()),
}
}
#[derive(Debug, Clone, Copy)]
enum PlanOwner<'a> {
Policy(&'a PostgresPolicy),
Candidate(&'a PostgresPolicyCandidate),
}
impl PlanOwner<'_> {
fn uid(&self) -> Option<&str> {
match self {
PlanOwner::Policy(policy) => policy.metadata.uid.as_deref(),
PlanOwner::Candidate(candidate) => candidate.metadata.uid.as_deref(),
}
}
fn name(&self) -> String {
match self {
PlanOwner::Policy(policy) => policy.name_any(),
PlanOwner::Candidate(candidate) => candidate.name_any(),
}
}
fn generation(&self) -> i64 {
match self {
PlanOwner::Policy(policy) => policy.metadata.generation.unwrap_or(0),
PlanOwner::Candidate(candidate) => candidate.metadata.generation.unwrap_or(0),
}
}
fn owner_reference(&self) -> OwnerReference {
match self {
PlanOwner::Policy(policy) => build_owner_reference(policy),
PlanOwner::Candidate(candidate) => OwnerReference {
api_version: PostgresPolicyCandidate::api_version(&()).to_string(),
kind: PostgresPolicyCandidate::kind(&()).to_string(),
name: candidate.name_any(),
uid: candidate.metadata.uid.clone().unwrap_or_default(),
controller: Some(true),
block_owner_deletion: Some(true),
},
}
}
fn owns<K: Resource>(&self, resource: &K) -> bool {
let Some(uid) = self.uid() else {
return false;
};
is_owned_by_uid(resource, uid)
}
}
fn plan_base_pin(plan: &PostgresPolicyPlan) -> Option<&str> {
plan.spec
.origin
.as_ref()
.and_then(|origin| origin.base_content_digest.as_deref())
}
#[allow(clippy::too_many_arguments)]
pub async fn create_or_update_plan(
client: &Client,
policy: &PostgresPolicy,
changes: &[pgroles_core::diff::Change],
sql_context: &pgroles_core::sql::SqlContext,
inspect_config: &pgroles_inspect::InspectConfig,
reconciliation_mode: CrdReconciliationMode,
database_identity: &str,
target_identity: &TargetIdentity,
change_summary: &ChangeSummary,
password_source_versions: &BTreeMap<String, String>,
plan_retention: PlanRetention,
candidate: Option<CandidatePlanBinding<'_>>,
) -> Result<PlanCreationResult, ReconcileError> {
let namespace = policy.namespace().ok_or(ReconcileError::NoNamespace)?;
let policy_name = policy.name_any();
let generation = policy.metadata.generation.unwrap_or(0);
let owner = match candidate {
Some(binding) => PlanOwner::Candidate(binding.candidate),
None => PlanOwner::Policy(policy),
};
let owner_name = owner.name();
let expected_base: Option<&str> = candidate.and_then(|binding| {
(binding.base_content_digest != binding.content_digest)
.then_some(binding.base_content_digest)
});
let full_sql = render_full_sql(changes, sql_context);
let change_digest = compute_change_digest(
changes,
reconciliation_mode,
database_identity,
target_identity,
password_source_versions,
&inspect_config.managed_roles,
&inspect_config.managed_schemas,
)?;
let sql_hash = compute_sql_hash(&full_sql);
let sql_statement_count = full_sql.lines().filter(|l| !l.trim().is_empty()).count() as i64;
let redacted_sql = render_redacted_sql(changes, sql_context);
if candidate.is_none() {
cleanup_old_plans_best_effort(client, policy, plan_retention).await;
}
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let selector = match candidate {
Some(binding) => candidate_selector(&binding.candidate.name_any()),
None => policy_selector(&policy_name),
};
let existing_plans: Vec<PostgresPolicyPlan> = plans_api
.list(&ListParams::default().labels_from(&selector))
.await?
.into_iter()
.filter(|plan| owner.owns(plan))
.collect();
for plan in &existing_plans {
if let Some(ref status) = plan.status
&& status.phase == PlanPhase::Pending
&& plan_matches_digest(status, &change_digest)
&& expected_base.is_none_or(|base| plan_base_pin(plan) == Some(base))
{
let plan_name = plan.name_any();
info!(
plan = %plan_name,
policy = %policy_name,
"existing pending plan has identical change digest, skipping creation"
);
supersede_stale_plans(
&plans_api,
&existing_plans,
&policy_name,
&plan_name,
&change_digest,
expected_base,
)
.await?;
return Ok(PlanCreationResult::Deduplicated(plan_name));
}
}
let now_ts = now_epoch_secs();
for plan in &existing_plans {
if let Some(ref status) = plan.status
&& status.phase == PlanPhase::Failed
&& plan_matches_digest(status, &change_digest)
{
let failed_ts = status
.failed_at
.as_deref()
.and_then(parse_rfc3339_epoch_secs)
.unwrap_or(0);
if failed_ts > 0 && now_ts - failed_ts < FAILED_PLAN_DEDUP_WINDOW_SECS {
let plan_name = plan.name_any();
info!(
plan = %plan_name,
policy = %policy_name,
age_secs = now_ts - failed_ts,
"recently-failed plan has identical change digest, skipping creation"
);
return Ok(PlanCreationResult::DeduplicatedFailed(plan_name));
}
}
}
let plan_name = generate_plan_name(&owner_name, &sql_hash);
let prepared_sql = prepare_plan_sql(&plan_name, &redacted_sql)?;
let sql_configmap_name = create_plan_sql_configmap(
client,
owner,
&namespace,
&policy_name,
database_identity,
&prepared_sql,
)
.await?;
let owner_ref = owner.owner_reference();
let plan = PostgresPolicyPlan::new(
&plan_name,
PostgresPolicyPlanSpec {
policy_ref: PolicyPlanRef {
name: policy_name.clone(),
},
policy_generation: generation,
reconciliation_mode,
owned_roles: inspect_config.managed_roles.clone(),
owned_schemas: inspect_config.managed_schemas.clone(),
managed_database_identity: database_identity.to_string(),
origin: candidate.map(|binding| candidate_plan_origin(binding, policy)),
scope: None,
},
);
let mut plan = plan;
plan.metadata.namespace = Some(namespace.clone());
plan.metadata.owner_references = Some(vec![owner_ref.clone()]);
let mut plan_labels = BTreeMap::from([
(LABEL_POLICY.to_string(), sanitize_label_value(&policy_name)),
(
LABEL_DATABASE_IDENTITY.to_string(),
sanitize_label_value(database_identity),
),
]);
if let Some(binding) = candidate {
plan_labels.insert(
LABEL_CANDIDATE.to_string(),
sanitize_label_value(&binding.candidate.name_any()),
);
}
plan.metadata.labels = Some(plan_labels);
let sql_preview = redacted_sql.lines().take(5).collect::<Vec<_>>().join("\n");
let summary_text = format!(
"{}R {}G {}D {}DP {}M",
change_summary.roles_created + change_summary.roles_altered,
change_summary.grants_added,
change_summary.default_privileges_set,
change_summary.roles_dropped,
change_summary.members_added,
);
plan.metadata.annotations = Some(BTreeMap::from([
("pgroles.io/sql-preview".to_string(), sql_preview),
("pgroles.io/summary".to_string(), summary_text),
(
"pgroles.io/sql-hash".to_string(),
sql_hash[..12].to_string(),
),
(
"pgroles.io/redacted-sql-hash".to_string(),
prepared_sql.redacted_sql_hash[..12].to_string(),
),
(
"pgroles.io/sql-original-bytes".to_string(),
prepared_sql.original_bytes.to_string(),
),
(
"pgroles.io/sql-stored-bytes".to_string(),
prepared_sql.stored_bytes.to_string(),
),
]));
let (created_plan, created_new_plan) =
match plans_api.create(&PostParams::default(), &plan).await {
Ok(plan) => (plan, true),
Err(kube::Error::Api(api_err)) if api_err.code == 409 => {
let existing = plans_api.get(&plan_name).await?;
if !owner.owns(&existing) {
rollback_plan_sql_configmap(client, &namespace, sql_configmap_name.as_ref())
.await;
return Err(ReconcileError::PlanSqlStorage(format!(
"plan {plan_name} already exists and is owned by another object"
)));
}
if !should_patch_existing_plan_status(&existing) {
return Ok(PlanCreationResult::Deduplicated(existing.name_any()));
}
(existing, false)
}
Err(err) => {
rollback_plan_sql_configmap(client, &namespace, sql_configmap_name.as_ref()).await;
return Err(err.into());
}
};
let plan_name = created_plan.name_any();
let computed_message = if prepared_sql.is_truncated() {
format!(
"Plan computed with {} change(s); SQL preview truncated because compressed SQL exceeded Kubernetes ConfigMap limits",
change_summary.total
)
} else {
format!("Plan computed with {} change(s)", change_summary.total)
};
let plan_status = PostgresPolicyPlanStatus {
phase: PlanPhase::Pending,
decided_by: None,
conditions: vec![
PolicyCondition {
condition_type: "Computed".to_string(),
status: "True".to_string(),
reason: Some("PlanComputed".to_string()),
message: Some(computed_message),
last_transition_time: Some(crate::crd::now_rfc3339()),
},
PolicyCondition {
condition_type: "Approved".to_string(),
status: "False".to_string(),
reason: Some("PendingApproval".to_string()),
message: Some("Plan awaiting approval".to_string()),
last_transition_time: Some(crate::crd::now_rfc3339()),
},
],
change_summary: Some(change_summary.clone()),
sql_ref: prepared_sql.sql_ref(),
sql_inline: prepared_sql.sql_inline(),
sql_truncated: prepared_sql.is_truncated(),
computed_at: Some(crate::crd::now_rfc3339()),
applied_at: None,
last_error: None,
sql_hash: Some(sql_hash),
change_digest: Some(change_digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
target_physical_identity: target_identity.physical.clone(),
target_logical_fingerprint: target_identity.logical.clone(),
physical_identity_available: Some(target_identity.has_physical()),
revalidated_generation: Some(owner.generation()),
revalidated_at: Some(crate::crd::now_rfc3339()),
applying_since: None,
failed_at: None,
sql_statements: Some(sql_statement_count),
redacted_sql_hash: Some(prepared_sql.redacted_sql_hash.clone()),
sql_original_bytes: Some(prepared_sql.original_bytes as i64),
sql_stored_bytes: Some(prepared_sql.stored_bytes as i64),
};
let status_patch = serde_json::json!({ "status": plan_status });
if let Err(err) = plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&status_patch),
)
.await
{
if created_new_plan {
delete_plan_best_effort(&plans_api, &plan_name).await;
}
rollback_plan_sql_configmap(client, &namespace, sql_configmap_name.as_ref()).await;
return Err(err.into());
}
supersede_stale_plans(
&plans_api,
&existing_plans,
&owner_name,
&plan_name,
&change_digest,
expected_base,
)
.await?;
info!(
plan = %plan_name,
policy = %policy_name,
changes = change_summary.total,
"created new plan"
);
Ok(PlanCreationResult::Created(plan_name))
}
pub(crate) fn supersedes_after_create(status: &PostgresPolicyPlanStatus, new_digest: &str) -> bool {
match status.phase {
PlanPhase::Pending => true,
PlanPhase::Approved => !plan_matches_digest(status, new_digest),
_ => false,
}
}
async fn supersede_stale_plans(
plans_api: &Api<PostgresPolicyPlan>,
existing_plans: &[PostgresPolicyPlan],
policy_name: &str,
new_plan_name: &str,
new_digest: &str,
expected_base: Option<&str>,
) -> Result<(), ReconcileError> {
for plan in existing_plans {
let Some(ref status) = plan.status else {
continue;
};
let old_plan_name = plan.name_any();
let base_stale = expected_base.is_some_and(|base| {
matches!(status.phase, PlanPhase::Pending | PlanPhase::Approved)
&& plan_base_pin(plan) != Some(base)
});
if old_plan_name == new_plan_name
|| (!supersedes_after_create(status, new_digest) && !base_stale)
{
continue;
}
let cause = if base_stale {
SupersedeCause::BaseContentChanged
} else {
SupersedeCause::ReplacedByNewerPlan
};
info!(
plan = %old_plan_name,
policy = %policy_name,
phase = ?status.phase,
"marking existing plan as Superseded"
);
let patch = serde_json::json!({ "status": superseded_status(status, cause) });
plans_api
.patch_status(
&old_plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
}
Ok(())
}
fn retry_deferred_until_window_expires(plan: &PostgresPolicyPlan, now_ts: i64) -> bool {
let Some(status) = plan.status.as_ref() else {
return false;
};
let failed_ts = status
.failed_at
.as_deref()
.and_then(parse_rfc3339_epoch_secs)
.unwrap_or(0);
if !(failed_ts > 0 && failed_ts <= now_ts && now_ts - failed_ts < FAILED_PLAN_DEDUP_WINDOW_SECS)
{
return false;
}
let applied_ts = status
.applied_at
.as_deref()
.and_then(parse_rfc3339_epoch_secs)
.unwrap_or(0);
applied_ts <= failed_ts
}
pub async fn recorded_plan_failure(client: &Client, namespace: &str, plan_name: &str) -> String {
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), namespace);
plans_api
.get(plan_name)
.await
.ok()
.and_then(|plan| plan.status.and_then(|status| status.last_error))
.unwrap_or_else(|| "no error recorded".to_string())
}
pub async fn execute_plan(
client: &Client,
plan: &PostgresPolicyPlan,
pool: &sqlx::PgPool,
sql_context: &pgroles_core::sql::SqlContext,
changes: &[pgroles_core::diff::Change],
) -> Result<(), ReconcileError> {
let namespace = plan.namespace().ok_or(ReconcileError::NoNamespace)?;
let plan_name = plan.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
if retry_deferred_until_window_expires(plan, now_epoch_secs()) {
let recorded = plan
.status
.as_ref()
.and_then(|status| status.last_error.clone())
.unwrap_or_else(|| "no error recorded".to_string());
info!(
plan = %plan_name,
"plan failed recently, deferring retry to the policy interval"
);
return Err(ReconcileError::PlanRetryDeferred(plan_name, recorded));
}
update_plan_phase(&plans_api, &plan_name, PlanPhase::Applying).await?;
let result = execute_changes_in_transaction(pool, changes, sql_context).await;
match result {
Ok(statements_executed) => {
let mut applied_status = plan.status.clone().unwrap_or_default();
applied_status.phase = PlanPhase::Applied;
applied_status.applied_at = Some(crate::crd::now_rfc3339());
applied_status.last_error = None;
set_plan_condition(
&mut applied_status.conditions,
"Approved",
"True",
"Approved",
"Plan approved and executed",
);
let patch = serde_json::json!({ "status": applied_status });
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
info!(
plan = %plan_name,
statements = statements_executed,
"plan executed successfully"
);
Ok(())
}
Err(err) => {
let error_message = err.to_string();
let mut failed_status = plan.status.clone().unwrap_or_default();
failed_status.phase = PlanPhase::Failed;
failed_status.last_error = Some(error_message);
failed_status.failed_at = Some(crate::crd::now_rfc3339());
let patch = serde_json::json!({ "status": failed_status });
if let Err(status_err) = plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await
{
tracing::warn!(
plan = %plan_name,
%status_err,
"failed to update plan status to Failed"
);
}
Err(err)
}
}
}
fn prepare_plan_sql(
plan_name: &str,
redacted_sql: &str,
) -> Result<PreparedPlanSql, ReconcileError> {
let original_bytes = redacted_sql.len();
let redacted_sql_hash = compute_sql_hash(redacted_sql);
if original_bytes <= MAX_INLINE_SQL_BYTES {
return Ok(PreparedPlanSql {
artifact: PlanSqlArtifact::Inline(redacted_sql.to_string()),
redacted_sql_hash,
original_bytes,
stored_bytes: original_bytes,
});
}
let compressed_sql = gzip_bytes(redacted_sql.as_bytes())?;
if compressed_sql.len() <= MAX_CONFIGMAP_SQL_BYTES {
let stored_bytes = compressed_sql.len();
return Ok(PreparedPlanSql {
artifact: PlanSqlArtifact::CompressedConfigMap {
configmap_name: format!("{plan_name}-sql"),
key: SQL_CONFIGMAP_GZIP_KEY.to_string(),
compressed_sql,
},
redacted_sql_hash,
original_bytes,
stored_bytes,
});
}
let truncated = truncate_utf8(
redacted_sql,
MAX_INLINE_SQL_BYTES,
"\n-- truncated: compressed SQL preview exceeded Kubernetes ConfigMap limits --",
);
let stored_bytes = truncated.len();
Ok(PreparedPlanSql {
artifact: PlanSqlArtifact::TruncatedInline(truncated),
redacted_sql_hash,
original_bytes,
stored_bytes,
})
}
fn gzip_bytes(bytes: &[u8]) -> Result<Vec<u8>, ReconcileError> {
let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
encoder
.write_all(bytes)
.map_err(|err| ReconcileError::PlanSqlStorage(err.to_string()))?;
encoder
.finish()
.map_err(|err| ReconcileError::PlanSqlStorage(err.to_string()))
}
fn truncate_utf8(text: &str, max_bytes: usize, marker: &str) -> String {
if text.len() <= max_bytes {
return text.to_string();
}
let target_len = max_bytes.saturating_sub(marker.len());
let mut end = target_len.min(text.len());
while end > 0 && !text.is_char_boundary(end) {
end -= 1;
}
let mut truncated = text[..end].to_string();
truncated.push_str(marker);
truncated
}
struct PlanSqlConfigMap {
name: String,
created: bool,
}
async fn create_plan_sql_configmap(
client: &Client,
owner: PlanOwner<'_>,
namespace: &str,
policy_name: &str,
database_identity: &str,
prepared_sql: &PreparedPlanSql,
) -> Result<Option<PlanSqlConfigMap>, ReconcileError> {
let PlanSqlArtifact::CompressedConfigMap {
configmap_name,
key: _,
compressed_sql: _,
} = &prepared_sql.artifact
else {
return Ok(None);
};
let configmap = build_plan_sql_configmap_object(
owner,
namespace,
policy_name,
database_identity,
prepared_sql,
)?;
let configmaps_api: Api<ConfigMap> = Api::namespaced(client.clone(), namespace);
match configmaps_api
.create(&PostParams::default(), &configmap)
.await
{
Ok(_) => Ok(Some(PlanSqlConfigMap {
name: configmap_name.clone(),
created: true,
})),
Err(kube::Error::Api(api_err)) if api_err.code == 409 => {
let existing = configmaps_api.get(configmap_name).await?;
if is_owned_by_another(&existing, owner) {
return Err(ReconcileError::PlanSqlStorage(format!(
"plan SQL ConfigMap {configmap_name} is owned by another object"
)));
}
validate_existing_sql_configmap(&existing, prepared_sql)?;
Ok(Some(PlanSqlConfigMap {
name: configmap_name.clone(),
created: false,
}))
}
Err(err) => Err(err.into()),
}
}
fn build_plan_sql_configmap_object(
owner: PlanOwner<'_>,
namespace: &str,
policy_name: &str,
database_identity: &str,
prepared_sql: &PreparedPlanSql,
) -> Result<ConfigMap, ReconcileError> {
let PlanSqlArtifact::CompressedConfigMap {
configmap_name,
key,
compressed_sql,
} = &prepared_sql.artifact
else {
return Err(ReconcileError::PlanSqlStorage(
"cannot build ConfigMap for inline plan SQL".to_string(),
));
};
Ok(ConfigMap {
metadata: k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta {
name: Some(configmap_name.clone()),
namespace: Some(namespace.to_string()),
owner_references: Some(vec![owner.owner_reference()]),
labels: Some(BTreeMap::from([
(LABEL_POLICY.to_string(), sanitize_label_value(policy_name)),
(
LABEL_DATABASE_IDENTITY.to_string(),
sanitize_label_value(database_identity),
),
(
LABEL_PLAN.to_string(),
plan_label_value(configmap_plan_name(configmap_name)),
),
])),
annotations: Some(BTreeMap::from([
("pgroles.io/sql-compression".to_string(), "gzip".to_string()),
(
"pgroles.io/redacted-sql-hash".to_string(),
prepared_sql.redacted_sql_hash.clone(),
),
(
"pgroles.io/sql-original-bytes".to_string(),
prepared_sql.original_bytes.to_string(),
),
(
"pgroles.io/sql-stored-bytes".to_string(),
prepared_sql.stored_bytes.to_string(),
),
])),
..Default::default()
},
binary_data: Some(BTreeMap::from([(
key.clone(),
ByteString(compressed_sql.clone()),
)])),
..Default::default()
})
}
fn configmap_plan_name(configmap_name: &str) -> &str {
configmap_name
.strip_suffix("-sql")
.unwrap_or(configmap_name)
}
fn plan_label_value(plan_name: &str) -> String {
compute_sql_hash(plan_name)[..32].to_string()
}
fn validate_existing_sql_configmap(
configmap: &ConfigMap,
prepared_sql: &PreparedPlanSql,
) -> Result<(), ReconcileError> {
let Some(annotations) = configmap.metadata.annotations.as_ref() else {
return Err(ReconcileError::PlanSqlStorage(format!(
"existing ConfigMap {} is missing SQL storage annotations",
configmap.name_any()
)));
};
let hash_matches = annotations
.get("pgroles.io/redacted-sql-hash")
.map(|hash| hash == &prepared_sql.redacted_sql_hash)
.unwrap_or(false);
if hash_matches {
Ok(())
} else {
Err(ReconcileError::PlanSqlStorage(format!(
"existing ConfigMap {} does not match computed SQL preview hash",
configmap.name_any()
)))
}
}
pub(crate) async fn execute_changes_in_transaction(
pool: &sqlx::PgPool,
changes: &[pgroles_core::diff::Change],
sql_context: &pgroles_core::sql::SqlContext,
) -> Result<usize, ReconcileError> {
let mut transaction = pool.begin().await?;
let mut statements_executed = 0usize;
for change in changes {
let is_sensitive = matches!(change, pgroles_core::diff::Change::SetPassword { .. });
for sql in pgroles_core::sql::render_statements_with_context(change, sql_context) {
if is_sensitive {
tracing::debug!("executing: ALTER ROLE ... PASSWORD [REDACTED]");
} else {
tracing::debug!(%sql, "executing");
}
sqlx::query(&sql).execute(transaction.as_mut()).await?;
statements_executed += 1;
}
}
transaction.commit().await?;
Ok(statements_executed)
}
pub async fn cleanup_old_plans_best_effort(
client: &Client,
policy: &PostgresPolicy,
retention: PlanRetention,
) {
match tokio::time::timeout(
Duration::from_secs(CLEANUP_TIMEOUT_SECS),
cleanup_old_plans(client, policy, retention),
)
.await
{
Ok(Ok(())) => {}
Ok(Err(err)) => tracing::warn!(%err, "failed to clean up old plans"),
Err(_) => tracing::warn!(
timeout_secs = CLEANUP_TIMEOUT_SECS,
"timed out cleaning up old plans"
),
}
}
pub async fn cleanup_old_plans(
client: &Client,
policy: &PostgresPolicy,
retention: PlanRetention,
) -> Result<(), ReconcileError> {
let namespace = policy.namespace().ok_or(ReconcileError::NoNamespace)?;
let policy_name = policy.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let selector = policy_selector(&policy_name);
let existing_plans: Vec<PostgresPolicyPlan> = plans_api
.list(&ListParams::default().labels_from(&selector))
.await?
.into_iter()
.filter(|plan| is_owned_by_policy(plan, policy))
.collect();
let now_ts = now_epoch_secs();
for plan in existing_plans
.iter()
.filter(|plan| is_stale_statusless_plan(plan, now_ts))
{
let plan_name = plan.name_any();
info!(
plan = %plan_name,
policy = %policy_name,
"cleaning up stale status-less plan"
);
if let Err(err) = plans_api.delete(&plan_name, &DeleteParams::default()).await {
tracing::warn!(
plan = %plan_name,
%err,
"failed to delete stale status-less plan during cleanup"
);
}
}
for plan in plans_to_evict(&existing_plans, retention, now_ts) {
let plan_name = plan.name_any();
info!(
plan = %plan_name,
policy = %policy_name,
phase = ?plan.status.as_ref().map(|status| &status.phase),
"cleaning up old plan"
);
if let Err(err) = plans_api.delete(&plan_name, &DeleteParams::default()).await {
tracing::warn!(
plan = %plan_name,
%err,
"failed to delete old plan during cleanup"
);
}
}
cleanup_orphan_sql_configmaps(
client,
&namespace,
policy,
&policy_name,
&existing_plans,
now_ts,
)
.await?;
Ok(())
}
async fn cleanup_orphan_sql_configmaps(
client: &Client,
namespace: &str,
policy: &PostgresPolicy,
policy_name: &str,
existing_plans: &[PostgresPolicyPlan],
now_ts: i64,
) -> Result<(), ReconcileError> {
let configmaps_api: Api<ConfigMap> = Api::namespaced(client.clone(), namespace);
let selector = policy_selector(policy_name);
let configmaps: Vec<ConfigMap> = configmaps_api
.list(&ListParams::default().labels_from(&selector))
.await?
.into_iter()
.filter(|configmap| is_owned_by_policy(configmap, policy))
.collect();
let known_plan_labels: BTreeSet<String> = existing_plans
.iter()
.map(|plan| plan_label_value(&plan.name_any()))
.collect();
let known_plan_names: BTreeSet<String> =
existing_plans.iter().map(ResourceExt::name_any).collect();
for configmap in configmaps {
if !is_orphan_sql_configmap(&configmap, &known_plan_names, &known_plan_labels, now_ts) {
continue;
}
let configmap_name = configmap.name_any();
info!(
configmap = %configmap_name,
policy = %policy_name,
"cleaning up orphan plan SQL ConfigMap"
);
if let Err(err) = configmaps_api
.delete(&configmap_name, &DeleteParams::default())
.await
{
tracing::warn!(
configmap = %configmap_name,
%err,
"failed to delete orphan plan SQL ConfigMap during cleanup"
);
}
}
Ok(())
}
fn is_orphan_sql_configmap(
configmap: &ConfigMap,
known_plan_names: &BTreeSet<String>,
known_plan_labels: &BTreeSet<String>,
now_ts: i64,
) -> bool {
let Some(labels) = configmap.metadata.labels.as_ref() else {
return false;
};
if !labels.contains_key(LABEL_POLICY) || !is_stale_object(configmap, now_ts) {
return false;
}
if known_plan_names.contains(configmap_plan_name(&configmap.name_any())) {
return false;
}
labels
.get(LABEL_PLAN)
.map(|plan_label| !known_plan_labels.contains(plan_label))
.unwrap_or(true)
}
fn should_patch_existing_plan_status(plan: &PostgresPolicyPlan) -> bool {
plan.status
.as_ref()
.map(|status| status.phase == PlanPhase::Pending)
.unwrap_or(true)
}
fn plans_to_evict(
plans: &[PostgresPolicyPlan],
retention: PlanRetention,
now_ts: i64,
) -> Vec<&PostgresPolicyPlan> {
let mut applied = Vec::new();
let mut decided = Vec::new();
let mut superseded = Vec::new();
for plan in plans.iter().filter(|plan| !is_retention_exempt(*plan)) {
let Some(phase) = plan.status.as_ref().map(|status| &status.phase) else {
continue;
};
match phase {
PlanPhase::Applied => applied.push(plan),
PlanPhase::Failed | PlanPhase::Rejected => decided.push(plan),
PlanPhase::Superseded => superseded.push(plan),
PlanPhase::Pending | PlanPhase::Approved | PlanPhase::Applying => {}
}
}
let mut evict = oldest_beyond(&mut decided, retention.decided);
evict.extend(oldest_beyond(&mut superseded, retention.superseded));
evict.extend(applied_plans_to_evict(applied, retention, now_ts));
evict
}
pub(crate) fn applied_plans_to_evict(
mut applied: Vec<&PostgresPolicyPlan>,
retention: PlanRetention,
now_ts: i64,
) -> Vec<&PostgresPolicyPlan> {
applied.sort_by_key(|plan| applied_epoch_secs(plan));
let mut evict = Vec::new();
let mut remaining = applied.len();
for plan in applied {
if remaining <= retention.applied {
break;
}
let inside_floor = applied_age_secs(plan, now_ts) < retention.applied_min_age_secs;
if inside_floor && remaining <= retention.applied_ceiling {
continue;
}
evict.push(plan);
remaining -= 1;
}
evict
}
fn oldest_beyond<'a>(
plans: &mut Vec<&'a PostgresPolicyPlan>,
keep: usize,
) -> Vec<&'a PostgresPolicyPlan> {
if plans.len() <= keep {
return Vec::new();
}
sort_oldest_first(plans);
plans[..plans.len() - keep].to_vec()
}
fn sort_oldest_first(plans: &mut [&PostgresPolicyPlan]) {
plans.sort_by(|a, b| {
a.metadata
.creation_timestamp
.as_ref()
.cmp(&b.metadata.creation_timestamp.as_ref())
});
}
fn applied_epoch_secs(plan: &PostgresPolicyPlan) -> Option<i64> {
plan.status
.as_ref()
.and_then(|status| status.applied_at.as_deref())
.and_then(parse_rfc3339_epoch_secs)
.or_else(|| {
plan.metadata
.creation_timestamp
.as_ref()
.map(|timestamp| timestamp.0.as_second())
})
}
fn applied_age_secs(plan: &PostgresPolicyPlan, now_ts: i64) -> i64 {
applied_epoch_secs(plan)
.map(|applied_ts| now_ts.saturating_sub(applied_ts))
.unwrap_or(0)
}
fn is_stale_statusless_plan(plan: &PostgresPolicyPlan, now_ts: i64) -> bool {
plan.status.is_none() && is_stale_object(plan, now_ts)
}
fn is_stale_object<K>(resource: &K, now_ts: i64) -> bool
where
K: Resource,
{
resource
.meta()
.creation_timestamp
.as_ref()
.map(|timestamp| now_ts.saturating_sub(timestamp.0.as_second()) > ORPHAN_GRACE_SECS)
.unwrap_or(false)
}
async fn delete_plan_best_effort(plans_api: &Api<PostgresPolicyPlan>, plan_name: &str) {
if let Err(err) = plans_api.delete(plan_name, &DeleteParams::default()).await {
tracing::warn!(
plan = %plan_name,
%err,
"failed to roll back plan after status update failure"
);
}
}
async fn rollback_plan_sql_configmap(
client: &Client,
namespace: &str,
configmap: Option<&PlanSqlConfigMap>,
) {
if let Some(configmap) = configmap
&& configmap.created
{
delete_configmap_best_effort(client, namespace, &configmap.name).await;
}
}
async fn delete_configmap_best_effort(client: &Client, namespace: &str, configmap_name: &str) {
let configmaps_api: Api<ConfigMap> = Api::namespaced(client.clone(), namespace);
if let Err(err) = configmaps_api
.delete(configmap_name, &DeleteParams::default())
.await
{
tracing::warn!(
configmap = %configmap_name,
%err,
"failed to roll back plan SQL ConfigMap"
);
}
}
pub(crate) fn render_full_sql(
changes: &[pgroles_core::diff::Change],
sql_context: &pgroles_core::sql::SqlContext,
) -> String {
changes
.iter()
.flat_map(|change| pgroles_core::sql::render_statements_with_context(change, sql_context))
.collect::<Vec<_>>()
.join("\n")
}
fn render_redacted_sql(
changes: &[pgroles_core::diff::Change],
sql_context: &pgroles_core::sql::SqlContext,
) -> String {
changes
.iter()
.flat_map(|change| {
if let pgroles_core::diff::Change::SetPassword { name, .. } = change {
vec![format!(
"ALTER ROLE {} PASSWORD '[REDACTED]';",
pgroles_core::sql::quote_ident(name)
)]
} else {
pgroles_core::sql::render_statements_with_context(change, sql_context)
}
})
.collect::<Vec<_>>()
.join("\n")
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn compute_change_digest(
changes: &[pgroles_core::diff::Change],
reconciliation_mode: CrdReconciliationMode,
database_identity: &str,
target_identity: &TargetIdentity,
password_source_versions: &BTreeMap<String, String>,
owned_roles: &[String],
owned_schemas: &[String],
) -> Result<String, ReconcileError> {
Ok(pgroles_core::approval::compute_change_digest(
changes,
&pgroles_core::approval::EffectDigestInputs {
reconciliation_mode: reconciliation_mode.into(),
target: database_identity,
target_identity,
password_source_versions,
owned_roles,
owned_schemas,
},
)?)
}
pub(crate) fn plan_target_identity(
status: &crate::crd::PostgresPolicyPlanStatus,
) -> TargetIdentity {
TargetIdentity {
physical: status.target_physical_identity.clone(),
logical: status.target_logical_fingerprint.clone(),
}
}
pub(crate) fn plan_matches_digest(
status: &crate::crd::PostgresPolicyPlanStatus,
change_digest: &str,
) -> bool {
status.change_digest_encoding.as_deref() == Some(APPROVAL_EFFECT_ENCODING_V3)
&& status.change_digest.as_deref() == Some(change_digest)
}
pub(crate) fn compute_sql_hash(sql: &str) -> String {
use std::fmt::Write as _;
let mut hasher = Sha256::new();
hasher.update(sql.as_bytes());
let digest = hasher.finalize();
let mut hex = String::with_capacity(digest.len() * 2);
for byte in digest {
write!(&mut hex, "{byte:02x}").expect("writing to a string should succeed");
}
hex
}
fn generate_plan_name(policy_name: &str, sql_hash: &str) -> String {
let timestamp = format_timestamp_compact();
let suffix = &sql_hash[..12.min(sql_hash.len())];
let max_name_len = crate::k8s_names::MAX_RESOURCE_NAME_LENGTH - 4; let max_prefix_len = max_name_len - "-plan-".len() - timestamp.len() - "-".len() - suffix.len();
let prefix = truncate_name_prefix(policy_name, max_prefix_len);
format!("{prefix}-plan-{timestamp}-{suffix}")
}
fn format_timestamp_compact() -> String {
use std::time::SystemTime;
let now = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.unwrap_or_default();
let secs = now.as_secs();
let (year, month, day) = crate::crd::days_to_date(secs / 86400);
let remaining = secs % 86400;
let hours = remaining / 3600;
let minutes = (remaining % 3600) / 60;
let seconds = remaining % 60;
format!("{year:04}{month:02}{day:02}-{hours:02}{minutes:02}{seconds:02}")
}
fn sanitize_label_value(value: &str) -> String {
LabelValue::sanitize(value).into_string()
}
fn policy_selector(policy_name: &str) -> Selector {
Expression::Equal(LABEL_POLICY.to_string(), sanitize_label_value(policy_name)).into()
}
fn candidate_selector(candidate_name: &str) -> Selector {
Expression::Equal(
LABEL_CANDIDATE.to_string(),
sanitize_label_value(candidate_name),
)
.into()
}
fn is_owned_by_policy<K: Resource>(resource: &K, policy: &PostgresPolicy) -> bool {
let Some(policy_uid) = policy.metadata.uid.as_deref() else {
return false;
};
is_owned_by_uid(resource, policy_uid)
}
pub(crate) fn is_owned_by_uid<K: Resource>(resource: &K, uid: &str) -> bool {
resource
.meta()
.owner_references
.as_deref()
.unwrap_or_default()
.iter()
.any(|owner| owner.uid == uid && owner.controller.unwrap_or(false))
}
fn is_owned_by_another<K: Resource>(resource: &K, owner_object: PlanOwner<'_>) -> bool {
let policy_uid = owner_object.uid();
resource
.meta()
.owner_references
.as_deref()
.unwrap_or_default()
.iter()
.any(|owner| owner.controller.unwrap_or(false) && Some(owner.uid.as_str()) != policy_uid)
}
fn now_epoch_secs() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs() as i64
}
fn parse_rfc3339_epoch_secs(rfc3339: &str) -> Option<i64> {
rfc3339
.parse::<jiff::Timestamp>()
.ok()
.map(|t| t.as_second())
}
fn build_owner_reference(policy: &PostgresPolicy) -> OwnerReference {
OwnerReference {
api_version: PostgresPolicy::api_version(&()).to_string(),
kind: PostgresPolicy::kind(&()).to_string(),
name: policy.name_any(),
uid: policy.metadata.uid.clone().unwrap_or_default(),
controller: Some(true),
block_owner_deletion: Some(true),
}
}
async fn update_plan_phase(
plans_api: &Api<PostgresPolicyPlan>,
plan_name: &str,
phase: PlanPhase,
) -> Result<(), ReconcileError> {
let mut patch_value = serde_json::json!({ "status": { "phase": phase } });
if phase == PlanPhase::Applying {
patch_value["status"]["applying_since"] = serde_json::json!(crate::crd::now_rfc3339());
}
plans_api
.patch_status(
plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch_value),
)
.await?;
Ok(())
}
fn set_plan_condition(
conditions: &mut Vec<PolicyCondition>,
condition_type: &str,
status: &str,
reason: &str,
message: &str,
) {
let transition_time = if let Some(existing) = conditions
.iter()
.find(|c| c.condition_type == condition_type)
{
if existing.status == status {
existing.last_transition_time.clone()
} else {
Some(crate::crd::now_rfc3339())
}
} else {
Some(crate::crd::now_rfc3339())
};
let condition = PolicyCondition {
condition_type: condition_type.to_string(),
status: status.to_string(),
reason: Some(reason.to_string()),
message: Some(message.to_string()),
last_transition_time: transition_time,
};
conditions.retain(|c| c.condition_type != condition_type);
conditions.push(condition);
}
pub async fn update_policy_plan_ref(
client: &Client,
policy: &PostgresPolicy,
plan_name: &str,
) -> Result<(), ReconcileError> {
let namespace = policy.namespace().ok_or(ReconcileError::NoNamespace)?;
let policy_api: Api<PostgresPolicy> = Api::namespaced(client.clone(), &namespace);
let patch = serde_json::json!({
"status": {
"current_plan_ref": PlanReference {
name: plan_name.to_string(),
}
}
});
policy_api
.patch_status(
&policy.name_any(),
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
Ok(())
}
pub async fn get_current_actionable_plan(
client: &Client,
policy: &PostgresPolicy,
) -> Result<Option<PostgresPolicyPlan>, ReconcileError> {
let namespace = policy.namespace().ok_or(ReconcileError::NoNamespace)?;
let policy_name = policy.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let selector = policy_selector(&policy_name);
let existing_plans: Vec<PostgresPolicyPlan> = plans_api
.list(&ListParams::default().labels_from(&selector))
.await?
.into_iter()
.filter(|plan| is_owned_by_policy(plan, policy))
.collect();
let mut pending_plans: Vec<PostgresPolicyPlan> = existing_plans
.into_iter()
.filter(|plan| {
plan.status
.as_ref()
.map(|s| matches!(s.phase, PlanPhase::Pending | PlanPhase::Approved))
.unwrap_or(false)
})
.collect();
pending_plans.sort_by(|a, b| {
let a_time = a.metadata.creation_timestamp.as_ref();
let b_time = b.metadata.creation_timestamp.as_ref();
b_time.cmp(&a_time) });
Ok(pending_plans.into_iter().next())
}
pub async fn get_plan_by_phase(
client: &Client,
policy: &PostgresPolicy,
target_phase: PlanPhase,
) -> Result<Option<PostgresPolicyPlan>, ReconcileError> {
let namespace = policy.namespace().ok_or(ReconcileError::NoNamespace)?;
let policy_name = policy.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let selector = policy_selector(&policy_name);
let existing_plans: Vec<PostgresPolicyPlan> = plans_api
.list(&ListParams::default().labels_from(&selector))
.await?
.into_iter()
.filter(|plan| is_owned_by_policy(plan, policy))
.collect();
let mut matching_plans: Vec<PostgresPolicyPlan> = existing_plans
.into_iter()
.filter(|plan| {
plan.status
.as_ref()
.map(|s| s.phase == target_phase)
.unwrap_or(false)
})
.collect();
matching_plans.sort_by(|a, b| {
let a_time = a.metadata.creation_timestamp.as_ref();
let b_time = b.metadata.creation_timestamp.as_ref();
b_time.cmp(&a_time) });
Ok(matching_plans.into_iter().next())
}
pub async fn mark_plan_failed(
client: &Client,
plan: &PostgresPolicyPlan,
error_message: &str,
) -> Result<(), ReconcileError> {
let namespace = plan.namespace().ok_or(ReconcileError::NoNamespace)?;
let plan_name = plan.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let mut status = plan.status.clone().unwrap_or_default();
status.phase = PlanPhase::Failed;
status.last_error = Some(error_message.to_string());
status.failed_at = Some(crate::crd::now_rfc3339());
let patch = serde_json::json!({ "status": status });
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
info!(
plan = %plan_name,
"marked stuck Applying plan as Failed"
);
Ok(())
}
pub async fn mark_plan_approved(
client: &Client,
plan: &PostgresPolicyPlan,
reason: &str,
message: &str,
) -> Result<(), ReconcileError> {
let namespace = plan.namespace().ok_or(ReconcileError::NoNamespace)?;
let plan_name = plan.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let existing = plan.status.clone().unwrap_or_default();
if has_terminal_decision(&existing) {
let patch = serde_json::json!({ "status": { "phase": PlanPhase::Approved } });
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
return Ok(());
}
let mut status = existing;
status.phase = PlanPhase::Approved;
set_plan_condition(&mut status.conditions, "Approved", "True", reason, message);
if status.decided_by.is_none() {
status.decided_by = Some(crate::crd::DecisionActor {
username: AUTO_APPROVAL_ACTOR.to_string(),
uid: None,
groups: Vec::new(),
});
}
let patch = serde_json::json!({ "status": status });
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
Ok(())
}
pub async fn mark_plan_rejected(
client: &Client,
plan: &PostgresPolicyPlan,
) -> Result<(), ReconcileError> {
let namespace = plan.namespace().ok_or(ReconcileError::NoNamespace)?;
let plan_name = plan.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let patch = serde_json::json!({ "status": { "phase": PlanPhase::Rejected } });
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
Ok(())
}
pub(crate) fn has_terminal_decision(status: &PostgresPolicyPlanStatus) -> bool {
status.conditions.iter().any(|c| {
(c.condition_type == "Approved" || c.condition_type == "Denied") && c.status == "True"
})
}
pub(crate) fn superseded_status(
status: &PostgresPolicyPlanStatus,
cause: SupersedeCause,
) -> PostgresPolicyPlanStatus {
let mut next = status.clone();
next.phase = PlanPhase::Superseded;
set_plan_condition(
&mut next.conditions,
crate::crd::CONDITION_SUPERSEDED,
"True",
cause.reason(),
cause.message(),
);
if !has_terminal_decision(status) {
set_plan_condition(
&mut next.conditions,
"Approved",
"False",
cause.reason(),
cause.message(),
);
}
next
}
pub async fn mark_plan_superseded(
client: &Client,
plan: &PostgresPolicyPlan,
cause: SupersedeCause,
) -> Result<(), ReconcileError> {
let namespace = plan.namespace().ok_or(ReconcileError::NoNamespace)?;
let plan_name = plan.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let status = superseded_status(&plan.status.clone().unwrap_or_default(), cause);
let patch = serde_json::json!({ "status": status });
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
Ok(())
}
pub(crate) fn needs_revalidation_record(
status: &crate::crd::PostgresPolicyPlanStatus,
generation: Option<i64>,
) -> bool {
status.revalidated_generation != generation
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum PendingPlanDecision {
Retain,
Replace,
Clear,
}
pub(crate) fn decide_pending_plan(
plan_status: Option<&crate::crd::PostgresPolicyPlanStatus>,
fresh_digest: &str,
has_changes: bool,
) -> PendingPlanDecision {
if !has_changes {
return PendingPlanDecision::Clear;
}
if plan_status.is_some_and(|status| plan_matches_digest(status, fresh_digest)) {
PendingPlanDecision::Retain
} else {
PendingPlanDecision::Replace
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ApprovedPlanDecision {
Execute,
Replace,
Clear,
}
pub(crate) fn decide_approved_plan(
plan_status: Option<&crate::crd::PostgresPolicyPlanStatus>,
fresh_digest: &str,
has_changes: bool,
) -> ApprovedPlanDecision {
if plan_status.is_some_and(|status| plan_matches_digest(status, fresh_digest)) {
return ApprovedPlanDecision::Execute;
}
if has_changes {
ApprovedPlanDecision::Replace
} else {
ApprovedPlanDecision::Clear
}
}
pub async fn record_plan_revalidation(
client: &Client,
plan: &PostgresPolicyPlan,
generation: Option<i64>,
) -> Result<(), ReconcileError> {
let namespace = plan.namespace().ok_or(ReconcileError::NoNamespace)?;
let plan_name = plan.name_any();
let plans_api: Api<PostgresPolicyPlan> = Api::namespaced(client.clone(), &namespace);
let patch = serde_json::json!({
"status": {
"revalidatedGeneration": generation,
"revalidatedAt": crate::crd::now_rfc3339(),
}
});
plans_api
.patch_status(
&plan_name,
&PatchParams::apply("pgroles-operator"),
&Patch::Merge(&patch),
)
.await?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::crd::CrdReconciliationMode;
use crate::crd::{LocalObjectReference, PolicyContent, PostgresPolicyCandidateSpec};
use base64::Engine as _;
use flate2::read::GzDecoder;
use std::io::Read;
fn test_plan_spec() -> PostgresPolicyPlanSpec {
PostgresPolicyPlanSpec {
policy_ref: PolicyPlanRef {
name: "orders".to_string(),
},
policy_generation: 1,
reconciliation_mode: CrdReconciliationMode::Authoritative,
owned_roles: Vec::new(),
owned_schemas: Vec::new(),
managed_database_identity: "default/db/DATABASE_URL".to_string(),
origin: None,
scope: None,
}
}
fn test_candidate(name: &str, uid: &str) -> PostgresPolicyCandidate {
let mut candidate = PostgresPolicyCandidate::new(
name,
PostgresPolicyCandidateSpec {
policy_ref: LocalObjectReference {
name: "orders".to_string(),
},
replaces: None,
target: None,
content: PolicyContent::default(),
},
);
candidate.metadata.namespace = Some("default".to_string());
candidate.metadata.uid = Some(uid.to_string());
candidate
}
#[test]
fn a_candidate_plan_binds_the_candidate_the_content_and_the_policy() {
let mut policy = PostgresPolicy::new("orders", test_policy_spec());
policy.metadata.uid = Some("policy-uid".to_string());
let candidate = test_candidate("orders-change-x7k2p", "candidate-uid");
let origin = candidate_plan_origin(
CandidatePlanBinding {
candidate: &candidate,
content_digest: "sha256:abc",
content_digest_encoding: pgroles_core::candidate::CANDIDATE_CONTENT_ENCODING_V1,
base_content_digest: "sha256:base",
},
&policy,
);
assert_eq!(origin.kind, "PostgresPolicyCandidate");
assert_eq!(origin.name, "orders-change-x7k2p");
assert_eq!(origin.uid, "candidate-uid");
assert_eq!(origin.content_digest.as_deref(), Some("sha256:abc"));
assert_eq!(
origin.content_digest_encoding.as_deref(),
Some(pgroles_core::candidate::CANDIDATE_CONTENT_ENCODING_V1)
);
assert_eq!(origin.policy_uid.as_deref(), Some("policy-uid"));
}
#[test]
fn a_candidate_plan_is_owned_by_the_candidate_not_the_policy() {
let mut policy = PostgresPolicy::new("orders", test_policy_spec());
policy.metadata.uid = Some("policy-uid".to_string());
let candidate = test_candidate("orders-change-x7k2p", "candidate-uid");
let owner = PlanOwner::Candidate(&candidate).owner_reference();
assert_eq!(owner.kind, "PostgresPolicyCandidate");
assert_eq!(owner.uid, "candidate-uid");
assert_eq!(owner.controller, Some(true));
assert_eq!(owner.block_owner_deletion, Some(true));
let mut plan = PostgresPolicyPlan::new("plan", test_plan_spec());
plan.metadata.owner_references = Some(vec![owner]);
assert!(PlanOwner::Candidate(&candidate).owns(&plan));
assert!(!is_owned_by_policy(&plan, &policy));
}
#[test]
fn the_keep_label_exempts_a_terminal_object_from_retention() {
let mut plan = PostgresPolicyPlan::new("plan", test_plan_spec());
assert!(!crate::crd::is_retention_exempt(&plan));
plan.metadata.labels = Some(BTreeMap::from([(
crate::crd::LABEL_KEEP.to_string(),
"true".to_string(),
)]));
assert!(crate::crd::is_retention_exempt(&plan));
}
fn retained_plan(
name: &str,
phase: PlanPhase,
age_secs: i64,
now_ts: i64,
) -> PostgresPolicyPlan {
let mut plan = PostgresPolicyPlan::new(name, test_plan_spec());
plan.metadata.creation_timestamp =
Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
jiff::Timestamp::from_second(now_ts - age_secs).expect("epoch second in range"),
));
plan.status = Some(PostgresPolicyPlanStatus {
phase,
..Default::default()
});
plan
}
fn evicted_names(
plans: &[PostgresPolicyPlan],
retention: PlanRetention,
now_ts: i64,
) -> Vec<String> {
let mut names: Vec<String> = plans_to_evict(plans, retention, now_ts)
.into_iter()
.map(|plan| plan.name_any())
.collect();
names.sort();
names
}
#[test]
fn superseded_churn_does_not_evict_applied_plans() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let year = 400 * 24 * 60 * 60;
let mut plans = vec![retained_plan("applied-old", PlanPhase::Applied, year, now)];
for i in 0..50 {
plans.push(retained_plan(
&format!("superseded-{i:03}"),
PlanPhase::Superseded,
1000 - i,
now,
));
}
let evicted = evicted_names(&plans, retention, now);
assert!(
!evicted.contains(&"applied-old".to_string()),
"an applied plan must survive any amount of replan churn"
);
assert_eq!(
evicted.len(),
50 - retention.superseded,
"every superseded plan past the small bound is evicted"
);
}
#[test]
fn each_terminal_phase_is_bounded_separately() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let year = 400 * 24 * 60 * 60;
let mut plans = Vec::new();
for (prefix, phase, count) in [
("applied", PlanPhase::Applied, retention.applied + 4),
("failed", PlanPhase::Failed, retention.decided + 4),
(
"superseded",
PlanPhase::Superseded,
retention.superseded + 4,
),
] {
for i in 0..count {
plans.push(retained_plan(
&format!("{prefix}-{i:03}"),
phase.clone(),
year + count as i64 - i as i64,
now,
));
}
}
let evicted = evicted_names(&plans, retention, now);
for prefix in ["applied", "failed", "superseded"] {
let count = evicted.iter().filter(|n| n.starts_with(prefix)).count();
assert_eq!(count, 4, "{prefix} should lose exactly its 4 excess plans");
}
}
#[test]
fn failed_and_rejected_share_the_decided_bound() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let year = 400 * 24 * 60 * 60;
let mut plans = vec![retained_plan(
"rejected",
PlanPhase::Rejected,
year * 2,
now,
)];
for i in 0..retention.decided {
plans.push(retained_plan(
&format!("failed-{i:03}"),
PlanPhase::Failed,
year - i as i64,
now,
));
}
assert_eq!(
evicted_names(&plans, retention, now),
vec!["rejected".to_string()]
);
}
#[test]
fn the_age_floor_keeps_applied_plans_past_the_count_bound() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let recent: Vec<PostgresPolicyPlan> = (0..retention.applied + 10)
.map(|i| {
retained_plan(
&format!("applied-{i:03}"),
PlanPhase::Applied,
retention.applied_min_age_secs - 1 - i as i64,
now,
)
})
.collect();
assert!(
evicted_names(&recent, retention, now).is_empty(),
"nothing inside the floor is evicted, even past the count bound"
);
let stale: Vec<PostgresPolicyPlan> = (0..retention.applied + 10)
.map(|i| {
retained_plan(
&format!("applied-{i:03}"),
PlanPhase::Applied,
retention.applied_min_age_secs + 1 + i as i64,
now,
)
})
.collect();
assert_eq!(evicted_names(&stale, retention, now).len(), 10);
}
#[test]
fn the_ceiling_overrides_the_age_floor() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let plans: Vec<PostgresPolicyPlan> = (0..retention.applied_ceiling + 40)
.map(|i| {
retained_plan(
&format!("applied-{i:04}"),
PlanPhase::Applied,
60 + i as i64,
now,
)
})
.collect();
assert_eq!(
evicted_names(&plans, retention, now).len(),
40,
"the ceiling trims back to itself, not to the count bound"
);
}
#[test]
fn live_plans_and_kept_plans_are_never_evicted() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let year = 400 * 24 * 60 * 60;
let mut plans = Vec::new();
for (i, phase) in [PlanPhase::Pending, PlanPhase::Approved, PlanPhase::Applying]
.into_iter()
.enumerate()
{
plans.push(retained_plan(&format!("live-{i}"), phase, year * 2, now));
}
for i in 0..40 {
plans.push(retained_plan(
&format!("superseded-{i:03}"),
PlanPhase::Superseded,
year,
now,
));
}
let mut kept = retained_plan("kept", PlanPhase::Superseded, year * 3, now);
kept.metadata
.labels
.get_or_insert_with(Default::default)
.insert("pgroles.io/keep".to_string(), "true".to_string());
plans.push(kept);
let evicted = evicted_names(&plans, retention, now);
assert!(!evicted.iter().any(|name| name.starts_with("live-")));
assert!(!evicted.contains(&"kept".to_string()));
assert_eq!(evicted.len(), 40 - retention.superseded);
}
fn retention_from(vars: &[(&str, &str)]) -> Result<PlanRetention, PlanRetentionConfigError> {
PlanRetention::from_lookup(|variable| {
vars.iter()
.find(|(name, _)| *name == variable)
.map(|(_, value)| value.to_string())
})
}
#[test]
fn unset_retention_variables_keep_the_defaults() {
assert_eq!(
retention_from(&[]).expect("an empty environment is valid"),
PlanRetention::default()
);
}
#[test]
fn each_retention_variable_overrides_its_bound() {
let resolved = retention_from(&[
(PlanRetention::ENV_APPLIED, "7"),
(PlanRetention::ENV_APPLIED_CEILING, "70"),
(PlanRetention::ENV_APPLIED_MIN_AGE, "36h"),
(PlanRetention::ENV_DECIDED, "4"),
(PlanRetention::ENV_SUPERSEDED, "1"),
])
.expect("all five values are valid");
let expected = PlanRetention {
applied: 7,
applied_ceiling: 70,
applied_min_age_secs: 36 * 60 * 60,
decided: 4,
superseded: 1,
};
assert_ne!(expected, PlanRetention::default());
assert_eq!(resolved, expected);
}
#[test]
fn an_invalid_retention_count_is_rejected_naming_the_variable() {
let error = retention_from(&[(PlanRetention::ENV_DECIDED, "many")])
.expect_err("a non-numeric count must be rejected, not defaulted");
assert!(
error.to_string().contains(PlanRetention::ENV_DECIDED),
"the error must name the variable to fix: {error}"
);
}
#[test]
fn an_invalid_retention_min_age_is_rejected_naming_the_variable() {
let error = retention_from(&[(PlanRetention::ENV_APPLIED_MIN_AGE, "30d")])
.expect_err("an unsupported duration unit must be rejected, not defaulted");
assert!(
error
.to_string()
.contains(PlanRetention::ENV_APPLIED_MIN_AGE),
"the error must name the variable to fix: {error}"
);
}
#[test]
fn the_applied_floor_runs_from_when_the_plan_applied_not_when_it_was_created() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let excess = 10;
assert!(
retention.applied + excess <= retention.applied_ceiling,
"this test must not lean on the ceiling to pass"
);
let plans: Vec<PostgresPolicyPlan> = (0..retention.applied + excess)
.map(|i| {
let mut plan = retained_plan(
&format!("applied-{i:03}"),
PlanPhase::Applied,
retention.applied_min_age_secs * 3 + i as i64,
now,
);
plan.status
.as_mut()
.expect("retained_plan always sets a status")
.applied_at = Some(
jiff::Timestamp::from_second(now - 60 - i as i64)
.expect("epoch second in range")
.to_string(),
);
plan
})
.collect();
assert!(
evicted_names(&plans, retention, now).is_empty(),
"a plan applied inside the floor is kept, however old the object is"
);
}
#[test]
fn above_the_ceiling_eviction_order_follows_when_plans_applied_not_when_created() {
let now = 1_700_000_000;
let retention = PlanRetention::default();
let excess = 5;
let count = retention.applied_ceiling + excess;
let plans: Vec<PostgresPolicyPlan> = (0..count)
.map(|i| {
let mut plan = retained_plan(
&format!("applied-{i:04}"),
PlanPhase::Applied,
retention.applied_min_age_secs * 3 + (count - i) as i64,
now,
);
plan.status
.as_mut()
.expect("retained_plan always sets a status")
.applied_at = Some(
jiff::Timestamp::from_second(now - 100 - i as i64)
.expect("epoch second in range")
.to_string(),
);
plan
})
.collect();
let evicted = evicted_names(&plans, retention, now);
assert_eq!(
evicted.len(),
excess,
"the ceiling trims exactly the excess"
);
for i in 0..excess {
let earliest_executed = format!("applied-{:04}", count - 1 - i);
assert!(
evicted.contains(&earliest_executed),
"the earliest executions go first, whatever their objects' creation order"
);
}
assert!(
!evicted.contains(&"applied-0000".to_string()),
"the most recent execution must survive even as the oldest object by creation"
);
}
#[test]
fn a_non_unicode_environment_value_is_rejected_not_defaulted() {
use std::os::unix::ffi::OsStringExt;
let read = Err(std::env::VarError::NotUnicode(
std::ffi::OsString::from_vec(vec![b'2', b'5', 0xff]),
));
let error = env_read(PlanRetention::ENV_APPLIED, read)
.expect_err("a set-but-non-Unicode value must be rejected, not read as unset");
assert!(
error.to_string().contains(PlanRetention::ENV_APPLIED),
"the error must name the variable to fix: {error}"
);
assert_eq!(
env_read(
PlanRetention::ENV_APPLIED,
Err(std::env::VarError::NotPresent)
)
.expect("absent is not an error"),
None
);
assert_eq!(
env_read(PlanRetention::ENV_APPLIED, Ok("25".to_string()))
.expect("a Unicode value is not an error"),
Some("25".to_string())
);
}
#[test]
fn a_ceiling_below_the_applied_count_is_rejected() {
retention_from(&[
(PlanRetention::ENV_APPLIED, "50"),
(PlanRetention::ENV_APPLIED_CEILING, "10"),
])
.expect_err("a ceiling below the count bound must be rejected");
retention_from(&[
(PlanRetention::ENV_APPLIED, "50"),
(PlanRetention::ENV_APPLIED_CEILING, "50"),
])
.expect("a ceiling equal to the count bound is valid");
}
fn plan_failed_at(phase: PlanPhase, age_secs: i64, now_ts: i64) -> PostgresPolicyPlan {
let mut plan = PostgresPolicyPlan::new("plan", test_plan_spec());
plan.status = Some(PostgresPolicyPlanStatus {
phase,
last_error: Some("permission denied to create role".to_string()),
failed_at: Some(
jiff::Timestamp::from_second(now_ts - age_secs)
.expect("epoch second in range")
.to_string(),
),
..Default::default()
});
plan
}
#[test]
fn a_plan_that_just_failed_is_not_retried_until_the_window_expires() {
let now = 1_700_000_000;
assert!(
retry_deferred_until_window_expires(&plan_failed_at(PlanPhase::Failed, 1, now), now),
"a failure one second old must not be retried immediately"
);
assert!(
retry_deferred_until_window_expires(
&plan_failed_at(PlanPhase::Failed, FAILED_PLAN_DEDUP_WINDOW_SECS - 1, now),
now
),
"still inside the window"
);
assert!(
!retry_deferred_until_window_expires(
&plan_failed_at(PlanPhase::Failed, FAILED_PLAN_DEDUP_WINDOW_SECS, now),
now
),
"the window must expire so a fixed environment recovers"
);
}
#[test]
fn a_failure_timestamp_in_the_future_does_not_defer_the_retry() {
let now = 1_700_000_000;
for skew in [1, 60, FAILED_PLAN_DEDUP_WINDOW_SECS, 86_400] {
assert!(
!retry_deferred_until_window_expires(
&plan_failed_at(PlanPhase::Failed, -skew, now),
now
),
"a failure {skew}s in the future must not block execution"
);
}
assert!(retry_deferred_until_window_expires(
&plan_failed_at(PlanPhase::Failed, 0, now),
now
));
}
#[test]
fn the_backoff_survives_the_phase_being_rewritten() {
let now = 1_700_000_000;
for phase in [
PlanPhase::Pending,
PlanPhase::Approved,
PlanPhase::Applying,
PlanPhase::Failed,
] {
assert!(
retry_deferred_until_window_expires(&plan_failed_at(phase.clone(), 1, now), now),
"{phase:?} with a one-second-old failure must still defer"
);
}
}
#[test]
fn a_plan_that_applied_after_failing_is_not_deferred() {
let now = 1_700_000_000;
let mut recovered = plan_failed_at(PlanPhase::Applied, 30, now);
recovered.status.as_mut().unwrap().applied_at =
Some(jiff::Timestamp::from_second(now - 10).unwrap().to_string());
assert!(!retry_deferred_until_window_expires(&recovered, now));
let mut failed_again = plan_failed_at(PlanPhase::Failed, 10, now);
failed_again.status.as_mut().unwrap().applied_at =
Some(jiff::Timestamp::from_second(now - 30).unwrap().to_string());
assert!(retry_deferred_until_window_expires(&failed_again, now));
}
#[test]
fn an_apply_and_a_failure_in_the_same_second_still_defer() {
let now = 1_700_000_000;
let same_second = jiff::Timestamp::from_second(now - 5).unwrap().to_string();
let mut plan = plan_failed_at(PlanPhase::Failed, 5, now);
let status = plan.status.as_mut().unwrap();
status.applied_at = Some(same_second.clone());
assert_eq!(status.failed_at.as_deref(), Some(same_second.as_str()));
assert!(
retry_deferred_until_window_expires(&plan, now),
"an unorderable pair must resolve to the bounded side"
);
}
#[test]
fn a_plan_with_no_recorded_failure_runs() {
let now = 1_700_000_000;
let bare = PostgresPolicyPlan::new("plan", test_plan_spec());
assert!(!retry_deferred_until_window_expires(&bare, now));
let mut no_timestamp = PostgresPolicyPlan::new("plan", test_plan_spec());
no_timestamp.status = Some(PostgresPolicyPlanStatus {
phase: PlanPhase::Failed,
..Default::default()
});
assert!(!retry_deferred_until_window_expires(&no_timestamp, now));
}
fn test_policy_spec() -> crate::crd::PostgresPolicySpec {
crate::crd::PostgresPolicySpec {
connection: crate::crd::ConnectionSpec {
secret_ref: Some(crate::crd::SecretReference {
name: "db-credentials".to_string(),
}),
secret_key: Some("DATABASE_URL".to_string()),
params: None,
require_physical_identity: None,
},
interval: "5m".to_string(),
suspend: false,
mode: crate::crd::PolicyMode::Apply,
reconciliation_mode: CrdReconciliationMode::default(),
default_owner: None,
profiles: Default::default(),
schemas: Vec::new(),
roles: Vec::new(),
grants: Vec::new(),
default_privileges: Vec::new(),
memberships: Vec::new(),
retirements: Vec::new(),
approval: None,
}
}
fn policy_with_uid(name: &str, uid: &str) -> PostgresPolicy {
let mut policy = PostgresPolicy::new(name, test_policy_spec());
policy.metadata.namespace = Some("default".to_string());
policy.metadata.uid = Some(uid.to_string());
policy
}
#[test]
fn colliding_label_values_are_separated_by_owner_uid() {
let prefix = "a".repeat(63);
let first = policy_with_uid(&format!("{prefix}-one"), "uid-one");
let second = policy_with_uid(&format!("{prefix}-two"), "uid-two");
assert_eq!(
sanitize_label_value(&first.name_any()),
sanitize_label_value(&second.name_any()),
);
let mut plan = test_plan("plan-1", PlanPhase::Pending, None);
plan.metadata.owner_references = Some(vec![build_owner_reference(&first)]);
assert!(is_owned_by_policy(&plan, &first));
assert!(
!is_owned_by_policy(&plan, &second),
"a colliding policy must not claim another policy's plan"
);
}
#[test]
fn only_a_rival_controller_owner_blocks_adoption() {
let mine = policy_with_uid("orders", "uid-mine");
let theirs = policy_with_uid("orders-other", "uid-theirs");
let mut ours = test_plan("plan-1", PlanPhase::Pending, None);
ours.metadata.owner_references = Some(vec![build_owner_reference(&mine)]);
assert!(!is_owned_by_another(&ours, PlanOwner::Policy(&mine)));
let mut rival = test_plan("plan-2", PlanPhase::Pending, None);
rival.metadata.owner_references = Some(vec![build_owner_reference(&theirs)]);
assert!(is_owned_by_another(&rival, PlanOwner::Policy(&mine)));
let orphan = test_plan("plan-3", PlanPhase::Pending, None);
assert!(!is_owned_by_another(&orphan, PlanOwner::Policy(&mine)));
let mut non_controller = build_owner_reference(&theirs);
non_controller.controller = Some(false);
let mut referenced = test_plan("plan-4", PlanPhase::Pending, None);
referenced.metadata.owner_references = Some(vec![non_controller]);
assert!(!is_owned_by_another(&referenced, PlanOwner::Policy(&mine)));
}
#[test]
fn a_policy_without_a_uid_can_adopt_nothing_owned() {
let mut no_uid = policy_with_uid("orders", "uid-orders");
no_uid.metadata.uid = None;
let owner = policy_with_uid("orders", "uid-orders");
let mut claimed = test_plan("plan-1", PlanPhase::Pending, None);
claimed.metadata.owner_references = Some(vec![build_owner_reference(&owner)]);
assert!(is_owned_by_another(&claimed, PlanOwner::Policy(&no_uid)));
let orphan = test_plan("plan-2", PlanPhase::Pending, None);
assert!(!is_owned_by_another(&orphan, PlanOwner::Policy(&no_uid)));
}
#[test]
fn ownership_requires_a_controller_owner_reference() {
let policy = policy_with_uid("orders", "uid-orders");
let plan = test_plan("plan-1", PlanPhase::Pending, None);
assert!(!is_owned_by_policy(&plan, &policy));
let mut non_controller = build_owner_reference(&policy);
non_controller.controller = Some(false);
let mut plan = test_plan("plan-2", PlanPhase::Pending, None);
plan.metadata.owner_references = Some(vec![non_controller]);
assert!(!is_owned_by_policy(&plan, &policy));
}
#[test]
fn policy_without_uid_owns_nothing() {
let mut policy = PostgresPolicy::new("orders", test_policy_spec());
policy.metadata.namespace = Some("default".to_string());
policy.metadata.uid = None;
let mut plan = test_plan("plan-1", PlanPhase::Pending, None);
plan.metadata.owner_references = Some(vec![build_owner_reference(&policy)]);
assert!(!is_owned_by_policy(&plan, &policy));
}
#[test]
fn recreated_policy_does_not_inherit_previous_plans() {
let original = policy_with_uid("orders", "uid-original");
let recreated = policy_with_uid("orders", "uid-recreated");
let mut plan = test_plan("plan-1", PlanPhase::Applied, None);
plan.metadata.owner_references = Some(vec![build_owner_reference(&original)]);
assert!(is_owned_by_policy(&plan, &original));
assert!(!is_owned_by_policy(&plan, &recreated));
}
#[test]
fn configmap_ownership_uses_owner_uid() {
let prefix = "a".repeat(63);
let mine = policy_with_uid(&format!("{prefix}-one"), "uid-one");
let theirs = policy_with_uid(&format!("{prefix}-two"), "uid-two");
let configmap = ConfigMap {
metadata: k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta {
name: Some("plan-1-sql".to_string()),
owner_references: Some(vec![build_owner_reference(&theirs)]),
..Default::default()
},
..Default::default()
};
assert!(is_owned_by_policy(&configmap, &theirs));
assert!(
!is_owned_by_policy(&configmap, &mine),
"cleanup must not treat another policy's ConfigMap as its own"
);
}
fn test_plan_with_decisions(
name: &str,
phase: PlanPhase,
decisions: &[(&str, &str)],
) -> PostgresPolicyPlan {
let mut plan = test_plan(name, phase, None);
let status = plan.status.as_mut().expect("status");
status.conditions = decisions
.iter()
.map(|(condition_type, condition_status)| PolicyCondition {
condition_type: (*condition_type).to_string(),
status: (*condition_status).to_string(),
reason: Some("DecidedByReviewer".to_string()),
message: None,
last_transition_time: Some(crate::crd::now_rfc3339()),
})
.collect();
if decisions.iter().any(|(_, s)| *s == "True") {
status.decided_by = Some(crate::crd::DecisionActor {
username: "reviewer@example.com".to_string(),
uid: Some("uid-1".to_string()),
groups: vec!["platform".to_string()],
});
}
plan
}
fn test_plan(
name: &str,
phase: PlanPhase,
annotations: Option<BTreeMap<String, String>>,
) -> PostgresPolicyPlan {
let mut plan = PostgresPolicyPlan::new(
name,
PostgresPolicyPlanSpec {
policy_ref: PolicyPlanRef {
name: "test-policy".to_string(),
},
policy_generation: 1,
reconciliation_mode: CrdReconciliationMode::Authoritative,
owned_roles: vec!["role-a".to_string()],
owned_schemas: vec!["public".to_string()],
managed_database_identity: "default/db/DATABASE_URL".to_string(),
origin: None,
scope: None,
},
);
plan.metadata.namespace = Some("default".to_string());
plan.metadata.annotations = annotations;
plan.status = Some(PostgresPolicyPlanStatus {
phase,
..Default::default()
});
plan
}
#[test]
fn a_plan_with_no_recorded_decision_is_pending() {
let plan = test_plan("plan-1", PlanPhase::Pending, None);
assert_eq!(check_plan_approval(&plan), PlanApprovalState::Pending);
let mut statusless = test_plan("plan-2", PlanPhase::Pending, None);
statusless.status = None;
assert_eq!(check_plan_approval(&statusless), PlanApprovalState::Pending);
}
#[test]
fn a_decision_is_read_from_the_status_conditions() {
let approved =
test_plan_with_decisions("plan-1", PlanPhase::Pending, &[("Approved", "True")]);
assert_eq!(check_plan_approval(&approved), PlanApprovalState::Approved);
let denied = test_plan_with_decisions("plan-1", PlanPhase::Pending, &[("Denied", "True")]);
assert_eq!(check_plan_approval(&denied), PlanApprovalState::Rejected);
}
#[test]
fn a_false_decision_condition_is_not_a_decision() {
for conditions in [
&[("Approved", "False")][..],
&[("Denied", "False")][..],
&[("Approved", "False"), ("Denied", "False")][..],
] {
let plan = test_plan_with_decisions("plan-1", PlanPhase::Pending, conditions);
assert_eq!(
check_plan_approval(&plan),
PlanApprovalState::Pending,
"conditions {conditions:?} must not read as a decision"
);
}
}
#[test]
fn a_contradictory_decision_never_executes() {
let plan = test_plan_with_decisions(
"plan-1",
PlanPhase::Pending,
&[("Approved", "True"), ("Denied", "True")],
);
assert_eq!(check_plan_approval(&plan), PlanApprovalState::Rejected);
}
#[test]
fn the_retired_approval_annotation_grants_nothing() {
let annotations = BTreeMap::from([
("pgroles.io/approved".to_string(), "true".to_string()),
("pgroles.io/rejected".to_string(), "true".to_string()),
]);
let plan = test_plan("plan-1", PlanPhase::Pending, Some(annotations));
assert_eq!(check_plan_approval(&plan), PlanApprovalState::Pending);
}
#[test]
fn every_supersede_cause_names_its_own_reason() {
use pgroles_core::approval::TargetIdentityReason;
let causes = [
SupersedeCause::EffectsChanged,
SupersedeCause::EffectsCleared,
SupersedeCause::ReplacedByNewerPlan,
SupersedeCause::PolicyStoppedPlanning,
SupersedeCause::SupersededByPromotion,
SupersedeCause::BaseContentChanged,
SupersedeCause::TargetChanged(TargetIdentityReason::TargetChanged),
];
for cause in causes {
match cause {
SupersedeCause::EffectsChanged
| SupersedeCause::EffectsCleared
| SupersedeCause::ReplacedByNewerPlan
| SupersedeCause::PolicyStoppedPlanning
| SupersedeCause::SupersededByPromotion
| SupersedeCause::BaseContentChanged
| SupersedeCause::TargetChanged(_) => {}
}
}
let mut messages = std::collections::BTreeSet::new();
for cause in causes {
let message = cause.message();
assert!(!message.is_empty(), "{cause:?} has no message");
assert!(
messages.insert(message),
"{cause:?} reuses another cause's message"
);
assert!(
!message.contains("Database state changed since plan was approved"),
"{cause:?} still carries the old catch-all message"
);
}
}
#[test]
fn supersede_reasons_stay_stable_except_for_a_moved_target() {
use pgroles_core::approval::TargetIdentityReason;
for cause in [
SupersedeCause::EffectsChanged,
SupersedeCause::EffectsCleared,
SupersedeCause::ReplacedByNewerPlan,
SupersedeCause::PolicyStoppedPlanning,
] {
assert_eq!(cause.reason(), "Superseded", "{cause:?}");
}
assert_eq!(
SupersedeCause::SupersededByPromotion.reason(),
crate::crd::candidate_reason::SUPERSEDED_BY_PROMOTION
);
assert_eq!(
SupersedeCause::TargetChanged(TargetIdentityReason::TargetIdentityUnavailable).reason(),
"TargetIdentityUnavailable"
);
}
fn cel_admits(
old: &crate::crd::PostgresPolicyPlanStatus,
new: &crate::crd::PostgresPolicyPlanStatus,
) -> Result<(), &'static str> {
let true_decisions = |status: &crate::crd::PostgresPolicyPlanStatus| -> Vec<String> {
status
.conditions
.iter()
.filter(|c| {
(c.condition_type == "Approved" || c.condition_type == "Denied")
&& c.status == "True"
})
.map(|c| c.condition_type.clone())
.collect()
};
let old_decisions = true_decisions(old);
if !old_decisions.is_empty() && old_decisions != true_decisions(new) {
return Err("plan decisions are terminal");
}
if old.decided_by.is_some()
&& new.decided_by.as_ref().map(|d| &d.username)
!= old.decided_by.as_ref().map(|d| &d.username)
{
return Err("decision identity is write-once");
}
if true_decisions(new).is_empty() != new.decided_by.is_none() {
return Err(
"a terminal plan decision and decidedBy identity must be recorded together",
);
}
Ok(())
}
#[test]
fn superseding_an_approved_plan_preserves_the_decision() {
let plan = test_plan_with_decisions("plan-1", PlanPhase::Approved, &[("Approved", "True")]);
let old = plan.status.clone().expect("status");
let new = superseded_status(&old, SupersedeCause::EffectsChanged);
assert_eq!(new.phase, PlanPhase::Superseded);
let approved = new
.conditions
.iter()
.find(|c| c.condition_type == "Approved")
.expect("Approved condition preserved");
assert_eq!(approved.status, "True");
assert_eq!(
new.decided_by.as_ref().map(|d| d.username.as_str()),
Some("reviewer@example.com")
);
let superseded = new
.conditions
.iter()
.find(|c| c.condition_type == crate::crd::CONDITION_SUPERSEDED)
.expect("Superseded condition recorded");
assert_eq!(superseded.status, "True");
assert_eq!(
superseded.message.as_deref(),
Some(SupersedeCause::EffectsChanged.message())
);
assert_eq!(cel_admits(&old, &new), Ok(()));
}
#[test]
fn superseding_a_denied_plan_preserves_the_decision() {
let plan = test_plan_with_decisions("plan-1", PlanPhase::Rejected, &[("Denied", "True")]);
let old = plan.status.clone().expect("status");
let new = superseded_status(&old, SupersedeCause::SupersededByPromotion);
assert!(
new.conditions
.iter()
.any(|c| c.condition_type == "Denied" && c.status == "True")
);
assert!(
!new.conditions
.iter()
.any(|c| c.condition_type == "Approved")
);
assert_eq!(cel_admits(&old, &new), Ok(()));
}
#[test]
fn superseding_a_pending_plan_still_records_the_cause_on_approved() {
let plan = test_plan_with_decisions("plan-1", PlanPhase::Pending, &[("Approved", "False")]);
let old = plan.status.clone().expect("status");
let new = superseded_status(&old, SupersedeCause::ReplacedByNewerPlan);
let approved = new
.conditions
.iter()
.find(|c| c.condition_type == "Approved")
.expect("Approved condition");
assert_eq!(approved.status, "False");
assert_eq!(
approved.message.as_deref(),
Some(SupersedeCause::ReplacedByNewerPlan.message())
);
assert!(
new.conditions
.iter()
.any(|c| c.condition_type == crate::crd::CONDITION_SUPERSEDED)
);
assert_eq!(cel_admits(&old, &new), Ok(()));
}
#[test]
fn voiding_an_approval_by_flipping_the_condition_is_rejected() {
let plan = test_plan_with_decisions("plan-1", PlanPhase::Approved, &[("Approved", "True")]);
let old = plan.status.clone().expect("status");
let mut new = old.clone();
new.phase = PlanPhase::Superseded;
set_plan_condition(
&mut new.conditions,
"Approved",
"False",
SupersedeCause::EffectsChanged.reason(),
SupersedeCause::EffectsChanged.message(),
);
assert_eq!(cel_admits(&old, &new), Err("plan decisions are terminal"));
}
#[test]
fn a_superseded_plan_is_not_actionable_even_when_approved() {
let plan = test_plan_with_decisions("plan-1", PlanPhase::Approved, &[("Approved", "True")]);
let superseded = superseded_status(
&plan.status.clone().expect("status"),
SupersedeCause::SupersededByPromotion,
);
let mut retired = plan.clone();
retired.status = Some(superseded.clone());
assert_eq!(check_plan_approval(&retired), PlanApprovalState::Approved);
assert!(!matches!(
superseded.phase,
PlanPhase::Pending | PlanPhase::Approved
));
assert!(!supersedes_after_create(&superseded, "some-other-digest"));
}
#[test]
fn compute_sql_hash_is_deterministic() {
let sql = "CREATE ROLE test LOGIN;\nGRANT SELECT ON ALL TABLES IN SCHEMA public TO test;";
let hash1 = compute_sql_hash(sql);
let hash2 = compute_sql_hash(sql);
assert_eq!(hash1, hash2);
assert_eq!(hash1.len(), 64); }
#[test]
fn compute_sql_hash_differs_for_different_sql() {
let hash1 = compute_sql_hash("CREATE ROLE a;");
let hash2 = compute_sql_hash("CREATE ROLE b;");
assert_ne!(hash1, hash2);
}
#[test]
fn compute_sql_hash_matches_pinned_fixture() {
assert_eq!(
compute_sql_hash("CREATE ROLE app LOGIN;"),
"12a9743285d98ce73cfa9c840e943fc627d1fcbce22c5206fda1b21c84c1ac9c"
);
}
#[test]
fn generate_plan_name_has_expected_format() {
let hash = "abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789";
let name = generate_plan_name("my-policy", hash);
assert!(name.starts_with("my-policy-plan-"));
assert!(name.ends_with("-abcdef012345"));
let suffix = name.strip_prefix("my-policy-plan-").unwrap();
assert_eq!(suffix.len(), 28);
assert_eq!(&suffix[8..9], "-");
assert_eq!(&suffix[15..16], "-");
}
#[test]
fn generate_plan_name_is_idempotent_for_same_hash_in_same_second() {
let hash = "abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789";
let name1 = generate_plan_name("my-policy", hash);
let name2 = generate_plan_name("my-policy", hash);
assert_eq!(name1, name2);
}
#[test]
fn generate_plan_name_truncates_on_utf8_boundary() {
let hash = "abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789";
let name = generate_plan_name(&"é".repeat(140), hash);
assert!(name.len() <= 249);
assert!(name.ends_with("-abcdef012345"));
}
#[test]
fn plan_sql_configmap_identifiers_are_valid_at_the_name_limit() {
use crate::k8s_names::{is_valid_label_value, is_valid_resource_name};
let hash = "abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789";
let identity = format!(
"prod/params\0literal={}\0literal=appdb\05432",
"h".repeat(80)
);
for policy_name in [
"orders".to_string(),
"a.very-long-policy.name-with-dots.and-dashes.at-the-limit.xxxxx".to_string(),
format!("{}.{}", "a".repeat(214), "b".repeat(38)),
"a".repeat(253),
] {
let plan_name = generate_plan_name(&policy_name, hash);
let configmap_name = format!("{plan_name}-sql");
assert!(
is_valid_resource_name(&configmap_name),
"invalid ConfigMap name for policy {policy_name:?}: {configmap_name}"
);
assert!(
configmap_name.len() <= crate::k8s_names::MAX_RESOURCE_NAME_LENGTH,
"ConfigMap name over the limit: {} bytes",
configmap_name.len()
);
assert_eq!(configmap_plan_name(&configmap_name), plan_name);
for label in [
sanitize_label_value(&policy_name),
sanitize_label_value(&identity),
plan_label_value(&plan_name),
] {
assert!(
is_valid_label_value(&label),
"invalid label value {label:?} for policy {policy_name:?}"
);
}
}
}
#[test]
fn plan_label_value_is_stable_and_label_safe_for_long_names() {
let plan_name = "very-long-policy-name-".repeat(20);
let label = plan_label_value(&plan_name);
assert_eq!(label, plan_label_value(&plan_name));
assert_eq!(label.len(), 32);
assert!(label.chars().all(|ch| ch.is_ascii_hexdigit()));
}
#[test]
fn existing_non_pending_plan_status_is_not_repatched_on_create_conflict() {
let approved = test_plan("plan-1", PlanPhase::Approved, None);
let applying = test_plan("plan-1", PlanPhase::Applying, None);
let applied = test_plan("plan-1", PlanPhase::Applied, None);
assert!(!should_patch_existing_plan_status(&approved));
assert!(!should_patch_existing_plan_status(&applying));
assert!(!should_patch_existing_plan_status(&applied));
}
#[test]
fn existing_pending_or_statusless_plan_can_be_patched_on_create_conflict() {
let pending = test_plan("plan-1", PlanPhase::Pending, None);
let mut statusless = pending.clone();
statusless.status = None;
assert!(should_patch_existing_plan_status(&pending));
assert!(should_patch_existing_plan_status(&statusless));
}
#[test]
fn prepare_plan_sql_keeps_small_sql_inline() {
let prepared = prepare_plan_sql("plan-1", "CREATE ROLE app LOGIN;").unwrap();
assert!(matches!(prepared.artifact, PlanSqlArtifact::Inline(_)));
assert_eq!(
prepared.sql_inline(),
Some("CREATE ROLE app LOGIN;".to_string())
);
assert!(prepared.sql_ref().is_none());
assert!(!prepared.is_truncated());
}
#[test]
fn prepare_plan_sql_compresses_large_brownfield_sized_sql() {
let sql = brownfield_sized_sql();
assert!(sql.len() > 1_048_576);
let prepared = prepare_plan_sql("policy-plan-20260506-000000-abcdef012345", &sql).unwrap();
let PlanSqlArtifact::CompressedConfigMap {
key,
compressed_sql,
..
} = &prepared.artifact
else {
panic!("expected compressed ConfigMap artifact");
};
assert_eq!(key, SQL_CONFIGMAP_GZIP_KEY);
assert!(compressed_sql.len() < MAX_CONFIGMAP_SQL_BYTES);
assert_eq!(gunzip(compressed_sql), sql);
assert_eq!(
prepared.sql_ref().unwrap().compression,
Some(SqlCompression::Gzip)
);
assert_eq!(prepared.original_bytes, sql.len());
assert_eq!(prepared.stored_bytes, compressed_sql.len());
}
#[test]
fn configmap_binary_data_serializes_with_one_base64_layer() {
let sql = brownfield_sized_sql();
let prepared = prepare_plan_sql("policy-plan-20260506-000000-abcdef012345", &sql).unwrap();
let PlanSqlArtifact::CompressedConfigMap {
key,
compressed_sql,
..
} = &prepared.artifact
else {
panic!("expected compressed ConfigMap artifact");
};
let configmap = ConfigMap {
binary_data: Some(BTreeMap::from([(
key.clone(),
ByteString(compressed_sql.clone()),
)])),
..Default::default()
};
let encoded = serde_json::to_value(&configmap).unwrap()["binaryData"][key]
.as_str()
.unwrap()
.to_string();
let decoded = base64::engine::general_purpose::STANDARD
.decode(encoded)
.unwrap();
assert_eq!(decoded, *compressed_sql);
assert_eq!(gunzip(&decoded), sql);
}
#[test]
fn prepare_plan_sql_truncates_when_compressed_sql_is_still_too_large() {
let sql = deterministic_incompressible_sql(1_400_000);
let prepared = prepare_plan_sql("policy-plan-20260506-000000-abcdef012345", &sql).unwrap();
let PlanSqlArtifact::TruncatedInline(preview) = &prepared.artifact else {
panic!("expected truncated inline artifact");
};
assert!(preview.len() <= MAX_INLINE_SQL_BYTES);
assert!(preview.contains("truncated"));
assert!(prepared.sql_ref().is_none());
assert!(prepared.is_truncated());
}
#[test]
fn sanitize_label_value_replaces_slashes() {
let sanitized = sanitize_label_value("default/db-creds/DATABASE_URL");
assert!(!sanitized.contains('/'));
assert_eq!(sanitized, "default_db-creds_DATABASE_URL");
}
#[test]
fn sanitize_label_value_truncates_to_63_chars() {
let long_value = "a".repeat(100);
let sanitized = sanitize_label_value(&long_value);
assert!(sanitized.len() <= 63);
}
#[test]
fn stale_policy_sql_configmap_without_plan_label_is_orphan() {
let configmap = ConfigMap {
metadata: k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta {
labels: Some(BTreeMap::from([(
LABEL_POLICY.to_string(),
sanitize_label_value("test-policy"),
)])),
creation_timestamp: Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
jiff::Timestamp::from_second(0).unwrap(),
)),
..Default::default()
},
..Default::default()
};
assert!(is_orphan_sql_configmap(
&configmap,
&BTreeSet::new(),
&BTreeSet::new(),
ORPHAN_GRACE_SECS + 1
));
}
#[test]
fn stale_policy_sql_configmap_with_current_plan_name_is_not_orphan() {
let plan_name = "test-policy-plan-20260506-000000-abcdef012345";
let configmap = ConfigMap {
metadata: k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta {
name: Some(format!("{plan_name}-sql")),
labels: Some(BTreeMap::from([
(
LABEL_POLICY.to_string(),
sanitize_label_value("test-policy"),
),
(
LABEL_PLAN.to_string(),
sanitize_label_value("legacy-colliding-label"),
),
])),
creation_timestamp: Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
jiff::Timestamp::from_second(0).unwrap(),
)),
..Default::default()
},
..Default::default()
};
assert!(!is_orphan_sql_configmap(
&configmap,
&BTreeSet::from([plan_name.to_string()]),
&BTreeSet::new(),
ORPHAN_GRACE_SECS + 1
));
}
#[test]
fn stale_policy_sql_configmap_with_known_hash_plan_label_is_not_orphan() {
let plan_name = "test-policy-plan-20260506-000000-abcdef012345";
let plan_label = plan_label_value(plan_name);
let configmap = ConfigMap {
metadata: k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta {
name: Some("different-plan-sql".to_string()),
labels: Some(BTreeMap::from([
(
LABEL_POLICY.to_string(),
sanitize_label_value("test-policy"),
),
(LABEL_PLAN.to_string(), plan_label.clone()),
])),
creation_timestamp: Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
jiff::Timestamp::from_second(0).unwrap(),
)),
..Default::default()
},
..Default::default()
};
assert!(!is_orphan_sql_configmap(
&configmap,
&BTreeSet::new(),
&BTreeSet::from([plan_label]),
ORPHAN_GRACE_SECS + 1
));
}
#[test]
fn stale_policy_sql_configmap_with_only_legacy_colliding_label_is_orphan() {
let plan_name =
"very-long-policy-name-that-would-have-collided-plan-20260506-000000-abcdef012345";
let legacy_label = sanitize_label_value(plan_name);
let configmap = ConfigMap {
metadata: k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta {
name: Some("deleted-historical-plan-sql".to_string()),
labels: Some(BTreeMap::from([
(
LABEL_POLICY.to_string(),
sanitize_label_value("test-policy"),
),
(LABEL_PLAN.to_string(), legacy_label.clone()),
])),
creation_timestamp: Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
jiff::Timestamp::from_second(0).unwrap(),
)),
..Default::default()
},
..Default::default()
};
assert!(is_orphan_sql_configmap(
&configmap,
&BTreeSet::new(),
&BTreeSet::new(),
ORPHAN_GRACE_SECS + 1
));
}
#[test]
fn render_redacted_sql_masks_passwords() {
let changes = vec![
pgroles_core::diff::Change::CreateRole {
name: "app".to_string(),
state: pgroles_core::model::RoleState {
login: true,
..pgroles_core::model::RoleState::default()
},
},
pgroles_core::diff::Change::SetPassword {
name: "app".to_string(),
password: "super_secret".to_string(),
},
];
let ctx = pgroles_core::sql::SqlContext::default();
let redacted = render_redacted_sql(&changes, &ctx);
assert!(redacted.contains("[REDACTED]"));
assert!(!redacted.contains("super_secret"));
assert!(redacted.contains("CREATE ROLE"));
}
#[test]
fn render_redacted_sql_password_only_plan() {
let changes = vec![pgroles_core::diff::Change::SetPassword {
name: "db-user".to_string(),
password: "my_secret_pw".to_string(),
}];
let ctx = pgroles_core::sql::SqlContext::default();
let redacted = render_redacted_sql(&changes, &ctx);
assert!(redacted.contains("[REDACTED]"));
assert!(!redacted.contains("my_secret_pw"));
}
#[test]
fn render_full_sql_includes_passwords() {
let changes = vec![pgroles_core::diff::Change::SetPassword {
name: "app".to_string(),
password: "super_secret".to_string(),
}];
let ctx = pgroles_core::sql::SqlContext::default();
let full = render_full_sql(&changes, &ctx);
assert!(full.contains("super_secret") || full.contains("SCRAM-SHA-256"));
}
#[test]
fn now_epoch_secs_returns_plausible_value() {
let now = now_epoch_secs();
let y2025 = 1_735_689_600_i64;
let y2100 = 4_102_444_800_i64;
assert!(
now > y2025 && now < y2100,
"epoch secs {now} should be between 2025 and 2100"
);
}
fn brownfield_sized_sql() -> String {
let mut sql = String::new();
for schema in 0..33 {
for profile in ["reader", "writer", "owner", "cdc"] {
let role = format!("schema_{schema}_{profile}");
sql.push_str(&format!(
"CREATE ROLE \"{role}\" LOGIN;\nCOMMENT ON ROLE \"{role}\" IS 'Generated from profile {profile} for brownfield migration schema {schema} with cdc ownership directives and review metadata';\n"
));
for relkind in ["TABLES", "SEQUENCES", "FUNCTIONS"] {
sql.push_str(&format!(
"GRANT SELECT ON ALL {relkind} IN SCHEMA \"schema_{schema}\" TO \"{role}\";\n"
));
}
for owner in 0..20 {
sql.push_str(&format!(
"ALTER DEFAULT PRIVILEGES FOR ROLE \"owner_{owner}\" IN SCHEMA \"schema_{schema}\" GRANT SELECT ON TABLES TO \"{role}\";\n"
));
}
}
}
for member in 0..70 {
sql.push_str(&format!(
"GRANT \"group_{member}\" TO \"service_login_{}\";\n",
member % 20
));
}
while sql.len() <= 1_100_000 {
sql.push_str("-- brownfield migration padding for large plan regression\n");
}
sql
}
fn deterministic_incompressible_sql(target_bytes: usize) -> String {
let mut state = 0x1234_5678_u64;
let mut sql = String::with_capacity(target_bytes);
while sql.len() < target_bytes {
state ^= state << 13;
state ^= state >> 7;
state ^= state << 17;
let value = (state % 62) as u8;
let ch = match value {
0..=9 => b'0' + value,
10..=35 => b'a' + (value - 10),
_ => b'A' + (value - 36),
};
sql.push(ch as char);
if sql.len().is_multiple_of(120) {
sql.push('\n');
}
}
sql
}
fn gunzip(bytes: &[u8]) -> String {
let mut decoder = GzDecoder::new(bytes);
let mut decoded = String::new();
decoder.read_to_string(&mut decoded).unwrap();
decoded
}
fn password_versions(role: &str, version: &str) -> BTreeMap<String, String> {
BTreeMap::from([(role.to_string(), version.to_string())])
}
fn set_password(name: &str, verifier: &str) -> pgroles_core::diff::Change {
pgroles_core::diff::Change::SetPassword {
name: name.to_string(),
password: verifier.to_string(),
}
}
fn grant_change(role: &str) -> pgroles_core::diff::Change {
pgroles_core::diff::Change::Grant {
role: pgroles_core::model::Grantee::parse(role),
privileges: [pgroles_core::manifest::Privilege::Select]
.into_iter()
.collect(),
object_type: pgroles_core::manifest::ObjectType::Table,
schema: Some("inventory".to_string()),
name: Some("orders".to_string()),
}
}
fn test_target_identity() -> TargetIdentity {
TargetIdentity {
physical: Some("7412330000000000001".to_string()),
logical: Some("sha256:endpoint".to_string()),
}
}
fn digest_for(
changes: &[pgroles_core::diff::Change],
versions: &BTreeMap<String, String>,
) -> String {
compute_change_digest(
changes,
CrdReconciliationMode::default(),
"default/db-credentials:DATABASE_URL",
&test_target_identity(),
versions,
&[],
&[],
)
.expect("digest")
}
#[test]
fn password_plans_keep_one_approval_identity_across_reconciles() {
let versions = password_versions("app", "role-passwords:app:7");
let first = digest_for(
&[set_password("app", "SCRAM-SHA-256$4096:aaa$sA:vA")],
&versions,
);
let second = digest_for(
&[set_password("app", "SCRAM-SHA-256$4096:bbb$sB:vB")],
&versions,
);
assert_eq!(first, second);
let ctx = pgroles_core::sql::SqlContext::default();
assert_ne!(
compute_sql_hash(&render_full_sql(
&[set_password("app", "SCRAM-SHA-256$4096:aaa$sA:vA")],
&ctx
)),
compute_sql_hash(&render_full_sql(
&[set_password("app", "SCRAM-SHA-256$4096:bbb$sB:vB")],
&ctx
)),
);
}
#[test]
fn rotating_a_password_source_is_a_new_approval() {
let change = [set_password("app", "SCRAM-SHA-256$4096:aaa$sA:vA")];
assert_ne!(
digest_for(&change, &password_versions("app", "role-passwords:app:7")),
digest_for(&change, &password_versions("app", "role-passwords:app:8")),
);
}
#[test]
fn revalidation_is_recorded_once_per_generation() {
let confirmed_at_3 = PostgresPolicyPlanStatus {
revalidated_generation: Some(3),
..Default::default()
};
assert!(!needs_revalidation_record(&confirmed_at_3, Some(3)));
assert!(needs_revalidation_record(&confirmed_at_3, Some(4)));
assert!(needs_revalidation_record(&confirmed_at_3, Some(2)));
let legacy = PostgresPolicyPlanStatus::default();
assert!(needs_revalidation_record(&legacy, Some(1)));
assert!(!needs_revalidation_record(&legacy, None));
}
#[test]
fn a_pending_plan_is_retained_only_while_its_effects_hold() {
let versions = password_versions("app", "role-passwords:app:7");
let original = [grant_change("reporting")];
let edited = [grant_change("analytics")];
let planned = digest_for(&original, &versions);
let pending = PostgresPolicyPlanStatus {
phase: PlanPhase::Pending,
change_digest: Some(planned.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
revalidated_generation: Some(1),
..Default::default()
};
assert!(plan_matches_digest(
&pending,
&digest_for(&original, &versions)
));
assert!(needs_revalidation_record(&pending, Some(2)));
assert!(!plan_matches_digest(
&pending,
&digest_for(&edited, &versions)
));
}
#[test]
fn creating_a_replacement_retires_the_pending_plan_it_replaces() {
let versions = password_versions("app", "role-passwords:app:1");
let old_digest = digest_for(&[grant_change("app")], &versions);
let new_digest = digest_for(&[grant_change("reporting")], &versions);
let pending = PostgresPolicyPlanStatus {
phase: PlanPhase::Pending,
change_digest: Some(old_digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert!(supersedes_after_create(&pending, &new_digest));
assert!(supersedes_after_create(&pending, &old_digest));
}
#[test]
fn creating_a_replacement_voids_an_approval_that_no_longer_describes_the_effects() {
let versions = password_versions("app", "role-passwords:app:1");
let approved_digest = digest_for(&[grant_change("app")], &versions);
let fresh_digest = digest_for(&[grant_change("reporting")], &versions);
let approved = PostgresPolicyPlanStatus {
phase: PlanPhase::Approved,
change_digest: Some(approved_digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert!(supersedes_after_create(&approved, &fresh_digest));
assert!(!supersedes_after_create(&approved, &approved_digest));
}
#[test]
fn creating_a_plan_never_disturbs_a_settled_one() {
for phase in [
PlanPhase::Applied,
PlanPhase::Failed,
PlanPhase::Rejected,
PlanPhase::Superseded,
PlanPhase::Applying,
] {
let status = PostgresPolicyPlanStatus {
phase,
..Default::default()
};
assert!(!supersedes_after_create(&status, "any-digest"));
}
}
#[test]
fn a_recently_failed_identical_plan_is_reported_apart_from_a_pending_one() {
let pending = PlanCreationResult::Deduplicated("plan-a".to_string());
assert_eq!(pending.plan_name(), "plan-a");
assert!(!pending.is_created());
assert!(!pending.is_failed_backoff());
let failed = PlanCreationResult::DeduplicatedFailed("plan-b".to_string());
assert_eq!(failed.plan_name(), "plan-b");
assert!(!failed.is_created());
assert!(failed.is_failed_backoff());
let created = PlanCreationResult::Created("plan-c".to_string());
assert!(created.is_created());
assert!(!created.is_failed_backoff());
}
#[test]
fn a_candidate_plans_provenance_is_the_candidates_own_generation() {
let mut policy = PostgresPolicy::new(
"orders",
serde_json::from_value(serde_json::json!({
"connection": { "secretRef": { "name": "db" } },
}))
.expect("minimal policy spec"),
);
policy.metadata.generation = Some(5);
let mut candidate = PostgresPolicyCandidate::new(
"orders-change-x7k2p",
crate::crd::PostgresPolicyCandidateSpec {
policy_ref: crate::crd::LocalObjectReference {
name: "orders".to_string(),
},
replaces: None,
target: None,
content: Default::default(),
},
);
candidate.metadata.generation = Some(1);
assert_eq!(PlanOwner::Policy(&policy).generation(), 5);
assert_eq!(PlanOwner::Candidate(&candidate).generation(), 1);
}
#[test]
fn a_newly_created_plan_records_its_own_generation_as_confirmed() {
let fresh = PostgresPolicyPlanStatus {
revalidated_generation: Some(7),
..Default::default()
};
assert!(!needs_revalidation_record(&fresh, Some(7)));
}
#[test]
fn a_plan_matches_only_its_own_digest_under_the_current_encoding() {
let versions = password_versions("app", "role-passwords:app:7");
let digest = digest_for(
&[set_password("app", "SCRAM-SHA-256$4096:aaa$sA:vA")],
&versions,
);
let current = PostgresPolicyPlanStatus {
change_digest: Some(digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert!(plan_matches_digest(¤t, &digest));
let legacy = PostgresPolicyPlanStatus {
sql_hash: Some("deadbeef".to_string()),
..Default::default()
};
assert!(!plan_matches_digest(&legacy, &digest));
let other_encoding = PostgresPolicyPlanStatus {
change_digest: Some(digest.clone()),
change_digest_encoding: Some("pgroles.io/approval-effect/v0".to_string()),
..Default::default()
};
assert!(!plan_matches_digest(&other_encoding, &digest));
let different_effects = PostgresPolicyPlanStatus {
change_digest: Some("sha256:0000".to_string()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert!(!plan_matches_digest(&different_effects, &digest));
let previous_encoding = PostgresPolicyPlanStatus {
change_digest: Some(digest.clone()),
change_digest_encoding: Some(
pgroles_core::approval::APPROVAL_EFFECT_ENCODING_V2.to_string(),
),
..Default::default()
};
assert!(!plan_matches_digest(&previous_encoding, &digest));
}
#[test]
fn a_pending_plan_whose_effects_vanish_is_cleared_not_replaced() {
let versions = password_versions("app", "role-passwords:app:7");
let planned = digest_for(&[grant_change("reporting")], &versions);
let pending = PostgresPolicyPlanStatus {
phase: PlanPhase::Pending,
change_digest: Some(planned),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
let empty_digest = digest_for(&[], &versions);
assert_eq!(
decide_pending_plan(Some(&pending), &empty_digest, false),
PendingPlanDecision::Clear,
);
let empty_plan = PostgresPolicyPlanStatus {
phase: PlanPhase::Pending,
change_digest: Some(empty_digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert_eq!(
decide_pending_plan(Some(&empty_plan), &empty_digest, false),
PendingPlanDecision::Clear,
);
}
#[test]
fn a_pending_plan_is_replaced_only_when_real_effects_moved() {
let versions = password_versions("app", "role-passwords:app:7");
let original = [grant_change("reporting")];
let edited = [grant_change("analytics")];
let pending = PostgresPolicyPlanStatus {
phase: PlanPhase::Pending,
change_digest: Some(digest_for(&original, &versions)),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert_eq!(
decide_pending_plan(Some(&pending), &digest_for(&original, &versions), true),
PendingPlanDecision::Retain,
);
assert_eq!(
decide_pending_plan(Some(&pending), &digest_for(&edited, &versions), true),
PendingPlanDecision::Replace,
);
assert_eq!(
decide_pending_plan(None, &digest_for(&original, &versions), true),
PendingPlanDecision::Replace,
);
assert_eq!(
decide_pending_plan(
Some(&PostgresPolicyPlanStatus::default()),
&digest_for(&original, &versions),
true
),
PendingPlanDecision::Replace,
);
}
#[test]
fn an_approval_executes_only_the_effects_it_was_given_for() {
let versions = password_versions("app", "role-passwords:app:7");
let approved_changes = [grant_change("reporting")];
let approved_digest = digest_for(&approved_changes, &versions);
let approved = PostgresPolicyPlanStatus {
phase: PlanPhase::Approved,
change_digest: Some(approved_digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert_eq!(
decide_approved_plan(Some(&approved), &approved_digest, true),
ApprovedPlanDecision::Execute,
);
assert_eq!(
decide_approved_plan(
Some(&approved),
&digest_for(&[grant_change("analytics")], &versions),
true
),
ApprovedPlanDecision::Replace,
);
for unprovable in [
None,
Some(&PostgresPolicyPlanStatus::default()),
Some(&PostgresPolicyPlanStatus {
change_digest: Some(approved_digest.clone()),
change_digest_encoding: Some("pgroles.io/approval-effect/v0".to_string()),
..Default::default()
}),
] {
assert_ne!(
decide_approved_plan(unprovable, &approved_digest, true),
ApprovedPlanDecision::Execute,
);
}
}
#[test]
fn an_approval_does_not_carry_over_to_a_moved_target() {
let versions = BTreeMap::new();
let changes = [grant_change("reporting")];
let approved = PostgresPolicyPlanStatus {
change_digest: Some(digest_for(&changes, &versions)),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
target_physical_identity: test_target_identity().physical,
target_logical_fingerprint: test_target_identity().logical,
physical_identity_available: Some(true),
..Default::default()
};
let moved = TargetIdentity {
logical: Some("sha256:other-endpoint".to_string()),
..test_target_identity()
};
let fresh_digest = compute_change_digest(
&changes,
CrdReconciliationMode::default(),
"default/db-credentials:DATABASE_URL",
&moved,
&versions,
&[],
&[],
)
.expect("digest");
assert_eq!(
decide_approved_plan(Some(&approved), &fresh_digest, true),
ApprovedPlanDecision::Replace,
);
assert_eq!(
pgroles_core::approval::evaluate_target_identity(
&plan_target_identity(&approved),
&moved,
false,
),
pgroles_core::approval::TargetIdentityVerdict::Superseded(
pgroles_core::approval::TargetIdentityReason::TargetChanged
),
);
}
#[test]
fn a_plan_without_recorded_identities_matches_nothing_observed() {
let legacy = PostgresPolicyPlanStatus::default();
assert_eq!(plan_target_identity(&legacy), TargetIdentity::default());
assert_ne!(
pgroles_core::approval::evaluate_target_identity(
&plan_target_identity(&legacy),
&test_target_identity(),
false,
),
pgroles_core::approval::TargetIdentityVerdict::Proceed,
);
}
#[test]
fn an_approved_plan_whose_effects_vanish_is_cleared_not_replaced() {
let versions = password_versions("app", "role-passwords:app:7");
let approved = PostgresPolicyPlanStatus {
phase: PlanPhase::Approved,
change_digest: Some(digest_for(&[grant_change("reporting")], &versions)),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
let empty_digest = digest_for(&[], &versions);
assert_eq!(
decide_approved_plan(Some(&approved), &empty_digest, false),
ApprovedPlanDecision::Clear,
);
let approved_empty = PostgresPolicyPlanStatus {
phase: PlanPhase::Approved,
change_digest: Some(empty_digest.clone()),
change_digest_encoding: Some(APPROVAL_EFFECT_ENCODING_V3.to_string()),
..Default::default()
};
assert_eq!(
decide_approved_plan(Some(&approved_empty), &empty_digest, false),
ApprovedPlanDecision::Execute,
);
}
#[test]
fn setting_a_condition_leaves_exactly_one_of_that_type() {
let mut conditions = vec![
PolicyCondition {
condition_type: "Computed".to_string(),
status: "True".to_string(),
reason: None,
message: None,
last_transition_time: None,
},
PolicyCondition {
condition_type: "Approved".to_string(),
status: "False".to_string(),
reason: Some("PendingApproval".to_string()),
message: None,
last_transition_time: None,
},
PolicyCondition {
condition_type: "Approved".to_string(),
status: "True".to_string(),
reason: Some("ApprovedByReviewer".to_string()),
message: None,
last_transition_time: None,
},
];
set_plan_condition(&mut conditions, "Approved", "True", "Reason", "Message");
let approved: Vec<_> = conditions
.iter()
.filter(|c| c.condition_type == "Approved")
.collect();
assert_eq!(approved.len(), 1, "duplicate Approved conditions survived");
assert_eq!(approved[0].status, "True");
assert!(conditions.iter().any(|c| c.condition_type == "Computed"));
}
#[test]
fn the_digest_carries_no_password_material() {
let verifier = "SCRAM-SHA-256$4096:c2FsdA==$c3RvcmVk:c2VydmVy";
let changes = [set_password("app", verifier)];
let versions = password_versions("app", "role-passwords:app:7");
let bytes = pgroles_core::approval::canonical_change_set_bytes(
&changes,
&pgroles_core::approval::EffectDigestInputs {
reconciliation_mode: CrdReconciliationMode::default().into(),
target: "default/db-credentials:DATABASE_URL",
target_identity: &test_target_identity(),
password_source_versions: &versions,
owned_roles: &[],
owned_schemas: &[],
},
)
.expect("canonical bytes");
let encoded = String::from_utf8(bytes).expect("canonical bytes are UTF-8");
assert!(
!encoded.contains(verifier),
"verifier reached the digest input: {encoded}"
);
assert!(
!encoded.contains("c3RvcmVk"),
"stored key reached the digest input: {encoded}"
);
assert!(encoded.contains("role-passwords:app:7"));
assert!(digest_for(&changes, &versions).starts_with("sha256:"));
}
}