use std::future::Future;
use std::pin::Pin;
use std::time::Duration;
use async_trait::async_trait;
use tokio_postgres::error::SqlState;
use tokio_postgres::{Client, Config, Row, Transaction};
use super::envelope::{DeploymentKek, EnvelopeError, SealedSecret};
use super::{
ENVELOPE_CAPABILITIES, KekRef, SecretDescriptor, SecretError, SecretMaterial, SecretResolver,
SecretStore,
};
use crate::desired_state::secrets::{LifecycleTransition, SecretLifecycle, SecretOwner, SecretRef};
use crate::desired_state::{SecretId, Uuid7Generator};
const BACKEND: &str = "encrypted-postgres";
const SCHEMA_DDL: &str = include_str!("../../../sql/secret_store_v1.sql");
const MATERIAL_COLUMNS: &str = "tenant_id, project_id, lifecycle, scheme, kek_reference, wrapped_dek, dek_nonce, \
ciphertext, nonce";
#[derive(Debug, Clone)]
pub struct SecretStoreSettings {
pub schema: Option<String>,
pub create_table: bool,
pub connect_timeout: Duration,
pub operation_timeout: Duration,
}
impl Default for SecretStoreSettings {
fn default() -> Self {
Self {
schema: None,
create_table: true,
connect_timeout: Duration::from_secs(10),
operation_timeout: Duration::from_secs(30),
}
}
}
impl SecretStoreSettings {
pub fn from_config(
secret_store: &crate::config::SecretStore,
control_plane: &crate::config::ControlPlane,
) -> Self {
Self {
schema: secret_store
.schema
.as_deref()
.map(str::trim)
.filter(|schema| !schema.is_empty())
.map(str::to_owned),
create_table: secret_store.create_table,
connect_timeout: Duration::from_millis(control_plane.connect_timeout_ms),
operation_timeout: Duration::from_millis(control_plane.operation_timeout_ms),
}
}
}
pub struct PostgresSecrets {
config: Config,
settings: SecretStoreSettings,
search_path: Option<String>,
kek: DeploymentKek,
ids: Uuid7Generator,
client: tokio::sync::Mutex<Option<Client>>,
}
impl std::fmt::Debug for PostgresSecrets {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PostgresSecrets")
.field("schema", &self.search_path)
.field("kek", self.kek.reference())
.finish_non_exhaustive()
}
}
impl PostgresSecrets {
pub async fn connect(
dsn: &str,
settings: SecretStoreSettings,
kek: DeploymentKek,
) -> Result<Self, SecretError> {
let mut config: Config = dsn.parse().map_err(|error| {
denied(format!("the secret-store DSN could not be parsed: {error}"))
})?;
config.connect_timeout(settings.connect_timeout);
config.application_name(crate::telemetry::SERVICE_NAME);
let search_path = settings
.schema
.as_deref()
.map(|schema| {
crate::usage::validate_table_name(schema).map_err(denied)?;
if schema.contains('.') {
return Err(denied(format!(
"`{schema}` is not a single unqualified schema name"
)));
}
Ok(schema.to_owned())
})
.transpose()?;
let store = Self {
config,
settings,
search_path,
kek,
ids: Uuid7Generator::new(),
client: tokio::sync::Mutex::new(None),
};
let client = tokio::time::timeout(store.settings.connect_timeout, store.connect_client())
.await
.map_err(|_| unavailable_message("connection timed out"))?
.map_err(|error| {
boot_failure(
"connect to the secret store",
&error,
Writability::Unknown,
|| {
"Check the role, password, and database named by the `dsn_env` connection \
string under `[secret_store]`."
.to_owned()
},
)
})?;
let writability = Writability::of(&client).await;
writability.report();
store.prepare_schema(&client, writability).await?;
*store.client.lock().await = Some(client);
Ok(store)
}
pub fn kek_reference(&self) -> &KekRef {
self.kek.reference()
}
async fn prepare_schema(
&self,
client: &Client,
writability: Writability,
) -> Result<(), SecretError> {
if self.settings.create_table {
client.batch_execute(SCHEMA_DDL).await.map_err(|error| {
boot_failure("apply secret-store schema", &error, writability, || {
"Grant the connecting role `CREATE` on the schema and ownership of \
`axond_secret`, or apply `ops/postgres/secret_store_v1.sql` yourself and \
set `create_table = false` under `[secret_store]`."
.to_owned()
})
})?;
}
client
.query_one("SELECT count(*) FROM axond_secret WHERE false", &[])
.await
.map_err(|error| {
boot_failure(
"read the secret store's `axond_secret` table",
&error,
writability,
|| {
"Apply `ops/postgres/secret_store_v1.sql`, or set `create_table = true` \
under `[secret_store]` to let boot apply it."
.to_owned()
},
)
})?;
Ok(())
}
async fn connect_client(&self) -> Result<Client, tokio_postgres::Error> {
let (client, connection) = self.config.connect(crate::usage::tls_connector()).await?;
tokio::spawn(async move {
if let Err(error) = connection.await {
tracing::warn!(%error, "postgres secret-store connection closed");
}
});
if let Some(schema) = &self.search_path {
client
.batch_execute(&format!("SET search_path TO {schema}"))
.await?;
}
Ok(client)
}
async fn run<T>(
&self,
operation: impl for<'a> FnOnce(
&'a mut Client,
) -> Pin<
Box<dyn Future<Output = Result<T, SecretError>> + Send + 'a>,
>,
) -> Result<T, SecretError> {
let mut guard = self.client.lock().await;
if guard.as_ref().is_none_or(Client::is_closed) {
*guard = Some(
self.connect_client()
.await
.map_err(|error| unavailable("reconnect", &error))?,
);
}
let result = tokio::time::timeout(
self.settings.operation_timeout,
operation(guard.as_mut().expect("connected")),
)
.await
.map_err(|_| unavailable_message("operation timed out"))
.and_then(|result| result);
if matches!(result, Err(SecretError::Unavailable { .. })) {
*guard = None;
}
result
}
async fn insert(
client: &Transaction<'_>,
owner: SecretOwner,
reference: SecretRef,
sealed: &SealedSecret,
) -> Result<u64, SecretError> {
client
.execute(
"INSERT INTO axond_secret (secret_id, version, tenant_id, project_id, lifecycle, \
scheme, kek_reference, wrapped_dek, dek_nonce, ciphertext, nonce) \
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) \
ON CONFLICT (secret_id, version) DO NOTHING",
&[
&reference.secret.to_string(),
&version_of(reference),
&owner.tenant.to_string(),
&owner.project.map(|project| project.to_string()),
&SecretLifecycle::Staged.as_str(),
&sealed.scheme,
&sealed.kek.0,
&sealed.wrapped_dek,
&sealed.dek_nonce,
&sealed.ciphertext,
&sealed.nonce,
],
)
.await
.map_err(|error| unavailable("insert secret version", &error))
}
fn descriptor_of(
row: &Row,
owner: SecretOwner,
reference: &SecretRef,
) -> Result<SecretDescriptor, SecretError> {
let tenant: String = row.get("tenant_id");
let project: Option<String> = row.get("project_id");
if tenant != owner.tenant.to_string()
|| project != owner.project.map(|project| project.to_string())
{
return Err(SecretError::Ownership {
reference: *reference,
owner,
});
}
let stored: String = row.get("lifecycle");
let lifecycle = SecretLifecycle::parse(&stored).ok_or_else(|| {
SecretError::Invalid(format!(
"secret {reference} is stored in state `{stored}`, which this build does not read"
))
})?;
Ok(SecretDescriptor {
reference: *reference,
owner,
lifecycle,
})
}
async fn locked_descriptor(
transaction: &Transaction<'_>,
owner: SecretOwner,
reference: &SecretRef,
) -> Result<SecretDescriptor, SecretError> {
let row = transaction
.query_opt(
"SELECT tenant_id, project_id, lifecycle FROM axond_secret \
WHERE secret_id = $1 AND version = $2 FOR UPDATE",
&[&reference.secret.to_string(), &version_of(*reference)],
)
.await
.map_err(|error| unavailable("read secret version", &error))?
.ok_or(SecretError::NotFound(*reference))?;
Self::descriptor_of(&row, owner, reference)
}
}
#[async_trait]
impl SecretResolver for PostgresSecrets {
fn name(&self) -> &'static str {
BACKEND
}
fn capabilities(&self) -> crate::backends::Capabilities {
ENVELOPE_CAPABILITIES
}
async fn resolve(
&self,
owner: SecretOwner,
reference: &SecretRef,
) -> Result<SecretMaterial, SecretError> {
let reference = *reference;
let sealed = self
.run(|client| {
Box::pin(async move {
let row = client
.query_opt(
&format!(
"SELECT {MATERIAL_COLUMNS} FROM axond_secret \
WHERE secret_id = $1 AND version = $2"
),
&[&reference.secret.to_string(), &version_of(reference)],
)
.await
.map_err(|error| unavailable("read secret material", &error))?
.ok_or(SecretError::NotFound(reference))?;
let descriptor = Self::descriptor_of(&row, owner, &reference)?;
if !descriptor.permits_resolution() {
return Err(SecretError::Lifecycle {
reference,
state: descriptor.lifecycle,
});
}
sealed_of(&row, &reference)
})
})
.await?;
self.kek
.open(owner, &reference, &sealed)
.map_err(|error| unwrap_error(error, &reference, sealed.kek))
}
async fn exists(&self, owner: SecretOwner, reference: &SecretRef) -> Result<bool, SecretError> {
match self.describe(owner, reference).await {
Ok(descriptor) => Ok(descriptor.lifecycle.permits_resolution()),
Err(SecretError::NotFound(_) | SecretError::Ownership { .. }) => Ok(false),
Err(error) => Err(error),
}
}
}
#[async_trait]
impl SecretStore for PostgresSecrets {
async fn stage(
&self,
owner: SecretOwner,
material: SecretMaterial,
) -> Result<SecretDescriptor, SecretError> {
if material.is_empty() {
return Err(SecretError::Invalid("material is empty".to_owned()));
}
let reference = SecretRef::first(SecretId::new(self.ids.next()));
let sealed = self
.kek
.seal(owner, &reference, &material)
.map_err(|error| seal_error(error, &reference))?;
self.run(|client| {
Box::pin(async move {
let transaction = client
.transaction()
.await
.map_err(|error| unavailable("begin stage", &error))?;
let inserted = Self::insert(&transaction, owner, reference, &sealed).await?;
if inserted != 1 {
return Err(SecretError::Invalid(format!("{reference} already exists")));
}
transaction
.commit()
.await
.map_err(|error| unavailable("commit stage", &error))?;
Ok(SecretDescriptor {
reference,
owner,
lifecycle: SecretLifecycle::Staged,
})
})
})
.await
}
async fn rotate(
&self,
owner: SecretOwner,
reference: &SecretRef,
material: SecretMaterial,
) -> Result<SecretDescriptor, SecretError> {
if material.is_empty() {
return Err(SecretError::Invalid("material is empty".to_owned()));
}
let reference = *reference;
let rotated = reference.rotated();
let sealed = self
.kek
.seal(owner, &rotated, &material)
.map_err(|error| seal_error(error, &rotated))?;
self.run(|client| {
Box::pin(async move {
let transaction = client
.transaction()
.await
.map_err(|error| unavailable("begin rotation", &error))?;
let current = Self::locked_descriptor(&transaction, owner, &reference).await?;
if current.lifecycle.is_terminal() {
return Err(SecretError::Lifecycle {
reference,
state: current.lifecycle,
});
}
let inserted = Self::insert(&transaction, owner, rotated, &sealed).await?;
if inserted != 1 {
return Err(SecretError::Invalid(format!("{rotated} already exists")));
}
transaction
.commit()
.await
.map_err(|error| unavailable("commit rotation", &error))?;
Ok(SecretDescriptor {
reference: rotated,
owner,
lifecycle: SecretLifecycle::Staged,
})
})
})
.await
}
async fn transition(
&self,
owner: SecretOwner,
reference: &SecretRef,
next: SecretLifecycle,
) -> Result<LifecycleTransition, SecretError> {
let reference = *reference;
self.run(|client| {
Box::pin(async move {
let transaction = client
.transaction()
.await
.map_err(|error| unavailable("begin transition", &error))?;
let current = Self::locked_descriptor(&transaction, owner, &reference).await?;
let transition = current
.lifecycle
.transition_to(next)
.map_err(|source| SecretError::Transition { reference, source })?;
if let LifecycleTransition::Moved { to, .. } = transition {
let statement = if to == SecretLifecycle::Tombstoned {
"UPDATE axond_secret SET lifecycle = $3, updated_at = now(), \
destroyed_at = now(), wrapped_dek = NULL, dek_nonce = NULL, \
ciphertext = NULL, nonce = NULL \
WHERE secret_id = $1 AND version = $2"
} else {
"UPDATE axond_secret SET lifecycle = $3, updated_at = now() \
WHERE secret_id = $1 AND version = $2"
};
transaction
.execute(
statement,
&[
&reference.secret.to_string(),
&version_of(reference),
&to.as_str(),
],
)
.await
.map_err(|error| unavailable("record transition", &error))?;
}
transaction
.commit()
.await
.map_err(|error| unavailable("commit transition", &error))?;
Ok(transition)
})
})
.await
}
async fn describe(
&self,
owner: SecretOwner,
reference: &SecretRef,
) -> Result<SecretDescriptor, SecretError> {
let reference = *reference;
self.run(|client| {
Box::pin(async move {
let row = client
.query_opt(
"SELECT tenant_id, project_id, lifecycle FROM axond_secret \
WHERE secret_id = $1 AND version = $2",
&[&reference.secret.to_string(), &version_of(reference)],
)
.await
.map_err(|error| unavailable("describe secret version", &error))?
.ok_or(SecretError::NotFound(reference))?;
Self::descriptor_of(&row, owner, &reference)
})
})
.await
}
}
fn version_of(reference: SecretRef) -> i64 {
i64::try_from(reference.version.get()).unwrap_or(i64::MAX)
}
fn sealed_of(row: &Row, reference: &SecretRef) -> Result<SealedSecret, SecretError> {
let wrapped_dek: Option<Vec<u8>> = row.get("wrapped_dek");
let dek_nonce: Option<Vec<u8>> = row.get("dek_nonce");
let ciphertext: Option<Vec<u8>> = row.get("ciphertext");
let nonce: Option<Vec<u8>> = row.get("nonce");
let kek = KekRef(row.get::<_, String>("kek_reference"));
match (wrapped_dek, dek_nonce, ciphertext, nonce) {
(Some(wrapped_dek), Some(dek_nonce), Some(ciphertext), Some(nonce)) => Ok(SealedSecret {
scheme: row.get("scheme"),
kek,
wrapped_dek,
dek_nonce,
ciphertext,
nonce,
}),
_ => Err(SecretError::Unwrap {
reference: *reference,
kek,
}),
}
}
fn unwrap_error(error: EnvelopeError, reference: &SecretRef, kek: KekRef) -> SecretError {
match error {
EnvelopeError::UnknownScheme { found } => SecretError::Invalid(format!(
"secret {reference} is sealed with scheme `{found}`, which this build does not read"
)),
EnvelopeError::Random | EnvelopeError::Unopenable | EnvelopeError::Malformed { .. } => {
SecretError::Unwrap {
reference: *reference,
kek,
}
}
}
}
fn seal_error(error: EnvelopeError, reference: &SecretRef) -> SecretError {
SecretError::Denied {
backend: BACKEND,
message: format!("secret {reference} could not be sealed: {error}"),
}
}
fn denied(message: impl Into<String>) -> SecretError {
SecretError::Denied {
backend: BACKEND,
message: message.into(),
}
}
fn unavailable_message(message: impl Into<String>) -> SecretError {
SecretError::Unavailable {
backend: BACKEND,
message: message.into(),
}
}
fn boot_failure(
operation: &str,
error: &tokio_postgres::Error,
writability: Writability,
remedy: impl FnOnce() -> String,
) -> SecretError {
match error.code().filter(|code| operator_must_act(code)) {
Some(code) => denied(format!(
"could not {operation} ({}: {error}). {}",
code.code(),
remedy()
)),
None => match writability.diagnosis(error.code()) {
Some(diagnosis) => {
unavailable_message(format!("{operation} failed: {error}. {diagnosis}"))
}
None => unavailable(operation, error),
},
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Writability {
Writable,
Standby,
ReadOnly,
Unknown,
}
impl Writability {
async fn of(client: &Client) -> Self {
let Ok(row) = client
.query_one(
"SELECT pg_is_in_recovery(), current_setting('transaction_read_only') = 'on'",
&[],
)
.await
else {
return Self::Unknown;
};
match (row.get::<_, bool>(0), row.get::<_, bool>(1)) {
(true, _) => Self::Standby,
(false, true) => Self::ReadOnly,
(false, false) => Self::Writable,
}
}
fn report(self) {
match self {
Self::Standby => tracing::warn!(
"the secret store's `dsn_env` endpoint is in recovery: writes are refused while it \
is a standby. Repoint it at the primary if this is not a failover in progress."
),
Self::ReadOnly => tracing::warn!(
"the secret store's `dsn_env` endpoint refuses writes although it is not in \
recovery: check `default_transaction_read_only` on the role or the database, and \
the pooler's routing."
),
Self::Writable | Self::Unknown => {}
}
}
fn diagnosis(self, code: Option<&SqlState>) -> Option<&'static str> {
if code != Some(&SqlState::READ_ONLY_SQL_TRANSACTION) {
return None;
}
match self {
Self::Standby => Some(
"The endpoint is in recovery (`pg_is_in_recovery()`). This retries in case a \
failover is in progress; if it is not, the `dsn_env` under `[secret_store]` names \
a standby and has to be repointed at the primary.",
),
Self::ReadOnly => Some(
"The endpoint is not in recovery but still refuses writes: check \
`default_transaction_read_only` on the connecting role and the database, and the \
pooler's routing.",
),
Self::Writable | Self::Unknown => None,
}
}
}
fn operator_must_act(code: &SqlState) -> bool {
if matches!(
*code,
SqlState::DUPLICATE_TABLE | SqlState::DUPLICATE_OBJECT
) {
return false;
}
matches!(code.code().get(..2), Some("42" | "3F" | "28" | "3D"))
}
fn unavailable(operation: &str, error: &tokio_postgres::Error) -> SecretError {
SecretError::Unavailable {
backend: BACKEND,
message: format!("{operation} failed: {error}"),
}
}
#[cfg(test)]
mod tests {
use base64::Engine as _;
use base64::engine::general_purpose::STANDARD;
use super::*;
use crate::backends::{BackendFailure, Capability, FailureCategory};
use crate::desired_state::fixtures::{project_id, tenant_id};
use crate::desired_state::secrets::SecretVersion;
use crate::test_services::postgres_dsn;
const PLAINTEXT: &str = "sk-live-do-not-log";
#[test]
fn only_the_sqlstates_an_operator_can_clear_refuse_a_boot() {
for permanent in [
SqlState::INSUFFICIENT_PRIVILEGE,
SqlState::UNDEFINED_TABLE,
SqlState::INVALID_SCHEMA_NAME,
SqlState::INVALID_PASSWORD,
SqlState::INVALID_CATALOG_NAME,
] {
assert!(
operator_must_act(&permanent),
"{} needs an operator",
permanent.code()
);
}
for transient in [
SqlState::CANNOT_CONNECT_NOW,
SqlState::TOO_MANY_CONNECTIONS,
SqlState::ADMIN_SHUTDOWN,
SqlState::T_R_DEADLOCK_DETECTED,
SqlState::LOCK_NOT_AVAILABLE,
SqlState::UNIQUE_VIOLATION,
SqlState::DUPLICATE_TABLE,
SqlState::DUPLICATE_OBJECT,
SqlState::CONNECTION_FAILURE,
SqlState::READ_ONLY_SQL_TRANSACTION,
] {
assert!(
!operator_must_act(&transient),
"{} is worth retrying",
transient.code()
);
}
}
#[test]
fn a_read_only_endpoint_is_diagnosed_rather_than_refused() {
let read_only = Some(&SqlState::READ_ONLY_SQL_TRANSACTION);
let standby = Writability::Standby
.diagnosis(read_only)
.expect("a standby names itself");
assert!(standby.contains("pg_is_in_recovery"), "{standby}");
assert!(standby.contains("dsn_env"), "{standby}");
let misconfigured = Writability::ReadOnly
.diagnosis(read_only)
.expect("a read-only session names its setting");
assert!(
misconfigured.contains("default_transaction_read_only"),
"{misconfigured}"
);
assert_eq!(Writability::Writable.diagnosis(read_only), None);
assert_eq!(Writability::Unknown.diagnosis(read_only), None);
assert_eq!(
Writability::Standby.diagnosis(Some(&SqlState::TOO_MANY_CONNECTIONS)),
None
);
}
fn with_role(dsn: &str, role: &str, password: &str) -> Option<String> {
let (scheme, rest) = dsn.split_once("://")?;
if !matches!(scheme, "postgres" | "postgresql") {
return None;
}
let endpoint = match rest.split_once('/') {
Some((authority, _)) => authority
.rsplit_once('@')
.map_or(rest, |(_, host)| &rest[authority.len() - host.len()..]),
None => rest.rsplit_once('@').map_or(rest, |(_, host)| host),
};
Some(format!("{scheme}://{role}:{password}@{endpoint}"))
}
#[test]
fn a_test_dsn_is_rewritten_onto_another_role_whatever_credentials_it_carried() {
for dsn in [
"postgres://postgres:secret@127.0.0.1:5432/axond",
"postgres://127.0.0.1:5432/axond",
"postgresql://someone@127.0.0.1:5432/axond?sslmode=disable",
] {
let rewritten = with_role(dsn, "reader", "pw").expect("a URL-shaped DSN");
assert!(
rewritten.contains("://reader:pw@127.0.0.1:5432/axond"),
"{rewritten}"
);
}
assert_eq!(
with_role("host=127.0.0.1 user=postgres", "reader", "pw"),
None
);
}
#[tokio::test]
async fn a_read_only_boot_is_an_outage_carrying_its_own_diagnosis() {
let Some(dsn) = postgres_dsn() else {
return;
};
let (client, connection) = tokio_postgres::connect(&dsn, crate::usage::tls_connector())
.await
.expect("connect to the test database");
tokio::spawn(async move {
let _ = connection.await;
});
let role = format!(
"axond_ro_{}",
Uuid7Generator::new().next().to_string().replace('-', "")
);
client
.batch_execute(&format!(
"CREATE ROLE {role} LOGIN PASSWORD 'axond'; ALTER ROLE {role} SET \
default_transaction_read_only = on; GRANT CREATE, USAGE ON SCHEMA public TO \
{role}"
))
.await
.expect("a role that may create but cannot write");
let outcome = match with_role(&dsn, &role, "axond") {
Some(read_only_dsn) => Some((
schema_ddl_sqlstate(&read_only_dsn).await,
PostgresSecrets::connect(&read_only_dsn, SecretStoreSettings::default(), kek(29))
.await
.err(),
)),
None => None,
};
client
.batch_execute(&format!(
"REVOKE ALL ON SCHEMA public FROM {role}; DROP ROLE IF EXISTS {role}"
))
.await
.expect("drop the test role");
let Some((sqlstate, outcome)) = outcome else {
return;
};
assert_eq!(
sqlstate.as_ref(),
Some(&SqlState::READ_ONLY_SQL_TRANSACTION),
"this fixture is only a demotion window if the endpoint refuses the DDL for being \
read-only; a different SQLSTATE means the role or the pooler, not read-only-ness, \
is what fails here"
);
let error = outcome.expect("a read-only endpoint cannot be prepared");
assert_eq!(
error.category(),
FailureCategory::Unavailable,
"a demotion window must be retried, not refused: {error}"
);
let message = error.to_string();
assert!(
message.contains("default_transaction_read_only"),
"the outage carries the preflight's answer for the endpoint that is read-only \
without being in recovery: {message}"
);
}
async fn schema_ddl_sqlstate(dsn: &str) -> Option<SqlState> {
let (client, connection) = tokio_postgres::connect(dsn, crate::usage::tls_connector())
.await
.ok()?;
tokio::spawn(async move {
let _ = connection.await;
});
client
.batch_execute(SCHEMA_DDL)
.await
.err()?
.code()
.cloned()
}
fn kek(seed: u8) -> DeploymentKek {
DeploymentKek::parse(
KekRef("AXOND_TEST_KEK".to_owned()),
&STANDARD.encode([seed; 32]),
)
.expect("a 32-byte key")
}
fn owner() -> SecretOwner {
SecretOwner::tenant(tenant_id(1))
}
async fn store(seed: u8) -> Option<(PostgresSecrets, String)> {
let dsn = postgres_dsn()?;
let schema = format!(
"axond_secret_test_{}",
Uuid7Generator::new().next().to_string().replace('-', "")
);
let (client, connection) = tokio_postgres::connect(&dsn, crate::usage::tls_connector())
.await
.expect("connect to the test database");
tokio::spawn(async move {
let _ = connection.await;
});
client
.batch_execute(&format!("CREATE SCHEMA {schema}"))
.await
.expect("create the test schema");
let store = PostgresSecrets::connect(
&dsn,
SecretStoreSettings {
schema: Some(schema.clone()),
..SecretStoreSettings::default()
},
kek(seed),
)
.await
.expect("the store applies its own schema");
Some((store, schema))
}
async fn drop_schema(schema: &str) {
let Some(dsn) = postgres_dsn() else {
return;
};
let (client, connection) = tokio_postgres::connect(&dsn, crate::usage::tls_connector())
.await
.expect("connect to the test database");
tokio::spawn(async move {
let _ = connection.await;
});
client
.batch_execute(&format!("DROP SCHEMA IF EXISTS {schema} CASCADE"))
.await
.expect("drop the test schema");
}
#[test]
fn the_store_declares_envelope_encryption_and_never_the_request_path() {
let responsibility =
crate::backends::responsibility("SecretStore").expect("a declared responsibility");
assert!(!responsibility.path.on_request_path());
assert!(ENVELOPE_CAPABILITIES.has(Capability::EnvelopeEncryption));
}
#[tokio::test]
async fn an_unreachable_store_is_unavailable() {
let error = PostgresSecrets::connect(
"postgres://axond@127.0.0.1:1/axond?connect_timeout=1",
SecretStoreSettings {
connect_timeout: Duration::from_millis(200),
..SecretStoreSettings::default()
},
kek(1),
)
.await
.expect_err("nothing listens there");
assert_eq!(error.category(), FailureCategory::Unavailable);
assert!(!error.to_string().contains("sk-"), "{error}");
}
#[tokio::test]
async fn a_boot_against_an_absent_database_is_refused_not_an_outage() {
let Some(dsn) = postgres_dsn() else {
return;
};
let config: Config = dsn.parse().expect("the test DSN parses");
let tokio_postgres::config::Host::Tcp(host) = &config.get_hosts()[0] else {
panic!("the test DSN names a TCP host");
};
let absent = format!(
"host={host} port={} user={} password={} dbname=axond_absent_database",
config.get_ports()[0],
config.get_user().unwrap_or("postgres"),
String::from_utf8_lossy(config.get_password().unwrap_or_default()),
);
let error = PostgresSecrets::connect(&absent, SecretStoreSettings::default(), kek(1))
.await
.expect_err("there is no such database");
assert_eq!(error.category(), FailureCategory::Denied);
let message = error.to_string();
assert!(message.contains("dsn_env"), "{message}");
}
#[tokio::test]
async fn a_malformed_dsn_is_refused_without_being_echoed() {
let error = PostgresSecrets::connect("not a dsn", SecretStoreSettings::default(), kek(1))
.await
.expect_err("a DSN that does not parse");
assert_eq!(error.category(), FailureCategory::Denied);
assert!(!error.to_string().contains("not a dsn"), "{error}");
}
#[tokio::test]
async fn material_round_trips_and_the_database_holds_only_ciphertext() {
let Some((store, schema)) = store(11).await else {
return;
};
let staged = store
.stage(owner(), SecretMaterial::new(PLAINTEXT.to_owned()))
.await
.expect("staging material");
assert_eq!(staged.lifecycle, SecretLifecycle::Staged);
assert_eq!(staged.reference.version, SecretVersion::FIRST);
assert_eq!(
store
.resolve(owner(), &staged.reference)
.await
.expect("staged material resolves")
.expose(),
PLAINTEXT
);
assert!(store.exists(owner(), &staged.reference).await.unwrap());
let row = store
.run(|client| {
Box::pin(async move {
client
.query_one(
"SELECT ciphertext, scheme, kek_reference FROM axond_secret \
WHERE secret_id = $1",
&[&staged.reference.secret.to_string()],
)
.await
.map_err(|error| unavailable("read the row back", &error))
})
})
.await
.expect("one row");
let ciphertext: Vec<u8> = row.get("ciphertext");
assert!(!ciphertext.windows(6).any(|window| window == b"sk-liv"));
assert_eq!(
row.get::<_, String>("scheme"),
super::super::envelope::SCHEME
);
assert_eq!(row.get::<_, String>("kek_reference"), "AXOND_TEST_KEK");
drop_schema(&schema).await;
}
#[tokio::test]
async fn the_lifecycle_is_what_the_domain_defines() {
let Some((store, schema)) = store(12).await else {
return;
};
let first = store
.stage(owner(), SecretMaterial::new(PLAINTEXT.to_owned()))
.await
.expect("staging")
.reference;
store
.transition(owner(), &first, SecretLifecycle::Active)
.await
.expect("staged material can be put in service");
let second = store
.rotate(owner(), &first, SecretMaterial::new("sk-live-2".to_owned()))
.await
.expect("rotating");
assert_eq!(second.reference, first.rotated());
assert_eq!(second.lifecycle, SecretLifecycle::Staged);
assert_eq!(
store.resolve(owner(), &first).await.unwrap().expose(),
PLAINTEXT
);
assert_eq!(
store
.resolve(owner(), &second.reference)
.await
.unwrap()
.expose(),
"sk-live-2"
);
assert!(matches!(
store
.rotate(owner(), &first, SecretMaterial::new("sk-live-3".to_owned()))
.await,
Err(SecretError::Invalid(_))
));
assert_eq!(
store
.resolve(owner(), &second.reference)
.await
.unwrap()
.expose(),
"sk-live-2"
);
store
.transition(owner(), &first, SecretLifecycle::Disabled)
.await
.expect("disabling");
assert!(matches!(
store.resolve(owner(), &first).await,
Err(SecretError::Lifecycle {
state: SecretLifecycle::Disabled,
..
})
));
assert!(!store.exists(owner(), &first).await.unwrap());
assert_eq!(
store
.transition(owner(), &first, SecretLifecycle::Disabled)
.await
.expect("a retry is not a conflict"),
LifecycleTransition::Unchanged(SecretLifecycle::Disabled)
);
store
.transition(owner(), &first, SecretLifecycle::Active)
.await
.expect("a disabled version rolls back into service");
assert_eq!(
store.resolve(owner(), &first).await.unwrap().expose(),
PLAINTEXT
);
store
.transition(owner(), &first, SecretLifecycle::Revoked)
.await
.expect("revoking");
assert!(matches!(
store
.transition(owner(), &first, SecretLifecycle::Active)
.await,
Err(SecretError::Transition { .. })
));
assert!(matches!(
store
.rotate(owner(), &first, SecretMaterial::new("x".to_owned()))
.await,
Err(SecretError::Invalid(_) | SecretError::Lifecycle { .. })
));
store
.transition(owner(), &first, SecretLifecycle::Tombstoned)
.await
.expect("tombstoning");
assert_eq!(
store.describe(owner(), &first).await.unwrap().lifecycle,
SecretLifecycle::Tombstoned
);
let held = store
.run(|client| {
Box::pin(async move {
client
.query_one(
"SELECT ciphertext IS NOT NULL AS held, destroyed_at IS NOT NULL AS \
destroyed FROM axond_secret WHERE secret_id = $1 AND version = $2",
&[&first.secret.to_string(), &version_of(first)],
)
.await
.map_err(|error| unavailable("read the tombstoned row", &error))
})
})
.await
.expect("the row survives its material");
assert!(!held.get::<_, bool>("held"), "tombstoning destroys bytes");
assert!(held.get::<_, bool>("destroyed"));
assert!(matches!(
store.resolve(owner(), &first).await,
Err(SecretError::Lifecycle {
state: SecretLifecycle::Tombstoned,
..
})
));
assert!(matches!(
store
.rotate(owner(), &first, SecretMaterial::new("sk-live-4".to_owned()))
.await,
Err(SecretError::Lifecycle { .. })
));
drop_schema(&schema).await;
}
#[tokio::test]
async fn material_is_isolated_by_owner() {
let Some((store, schema)) = store(13).await else {
return;
};
let mine = store
.stage(owner(), SecretMaterial::new(PLAINTEXT.to_owned()))
.await
.expect("staging")
.reference;
for theirs in [
SecretOwner::tenant(tenant_id(9)),
SecretOwner::project(tenant_id(1), project_id(2)),
] {
assert!(matches!(
store.resolve(theirs, &mine).await,
Err(SecretError::Ownership { .. })
));
assert!(!store.exists(theirs, &mine).await.unwrap());
assert!(matches!(
store.describe(theirs, &mine).await,
Err(SecretError::Ownership { .. })
));
assert!(matches!(
store
.rotate(theirs, &mine, SecretMaterial::new("sk-live-x".to_owned()))
.await,
Err(SecretError::Ownership { .. })
));
assert!(matches!(
store
.transition(theirs, &mine, SecretLifecycle::Revoked)
.await,
Err(SecretError::Ownership { .. })
));
assert_eq!(
store.describe(theirs, &mine).await.unwrap_err().category(),
FailureCategory::NotFound
);
}
let absent = SecretRef::first(SecretId::new(Uuid7Generator::new().next()));
assert!(matches!(
store.resolve(owner(), &absent).await,
Err(SecretError::NotFound(_))
));
assert!(!store.exists(owner(), &absent).await.unwrap());
assert!(store.resolve(owner(), &mine).await.is_ok());
drop_schema(&schema).await;
}
#[tokio::test]
async fn a_rotated_kek_cannot_unwrap_existing_material() {
let Some((store, schema)) = store(14).await else {
return;
};
let staged = store
.stage(owner(), SecretMaterial::new(PLAINTEXT.to_owned()))
.await
.expect("staging")
.reference;
let dsn = postgres_dsn().expect("a configured store");
let rotated = PostgresSecrets::connect(
&dsn,
SecretStoreSettings {
schema: Some(schema.clone()),
create_table: false,
..SecretStoreSettings::default()
},
kek(15),
)
.await
.expect("a store with a different key still connects");
let error = rotated
.resolve(owner(), &staged)
.await
.expect_err("material does not open under another key");
assert!(matches!(error, SecretError::Unwrap { .. }));
assert_eq!(error.category(), FailureCategory::Corrupt);
assert!(!error.to_string().contains("sk-"), "{error}");
drop_schema(&schema).await;
}
#[tokio::test]
async fn empty_material_is_refused() {
let Some((store, schema)) = store(16).await else {
return;
};
assert!(matches!(
store
.stage(owner(), SecretMaterial::new(String::new()))
.await,
Err(SecretError::Invalid(_))
));
let staged = store
.stage(owner(), SecretMaterial::new(PLAINTEXT.to_owned()))
.await
.expect("staging")
.reference;
assert!(matches!(
store
.rotate(owner(), &staged, SecretMaterial::new(String::new()))
.await,
Err(SecretError::Invalid(_))
));
drop_schema(&schema).await;
}
#[tokio::test]
async fn a_schema_the_table_cannot_be_created_in_is_refused_not_an_outage() {
let Some(dsn) = postgres_dsn() else {
return;
};
let error = PostgresSecrets::connect(
&dsn,
SecretStoreSettings {
schema: Some(format!(
"axond_absent_{}",
Uuid7Generator::new().next().to_string().replace('-', "")
)),
create_table: true,
..SecretStoreSettings::default()
},
kek(19),
)
.await
.expect_err("the table cannot be created there");
assert_eq!(error.category(), FailureCategory::Denied);
let message = error.to_string();
assert!(message.contains("create_table"), "{message}");
}
#[tokio::test]
async fn rotating_from_a_revoked_version_mints_a_successor_and_leaves_it_revoked() {
let Some((store, schema)) = store(23).await else {
return;
};
let staged = store
.stage(owner(), SecretMaterial::new(PLAINTEXT.to_owned()))
.await
.expect("a staged version");
store
.transition(owner(), &staged.reference, SecretLifecycle::Revoked)
.await
.expect("a withdrawal");
let rotated = store
.rotate(
owner(),
&staged.reference,
SecretMaterial::new("sk-replacement".to_owned()),
)
.await
.expect("a replacement for withdrawn material");
assert_eq!(rotated.reference, staged.reference.rotated());
assert_eq!(rotated.lifecycle, SecretLifecycle::Staged);
assert_eq!(
store
.describe(owner(), &staged.reference)
.await
.expect("the base version")
.lifecycle,
SecretLifecycle::Revoked,
"the withdrawn version does not come back"
);
assert!(matches!(
store.resolve(owner(), &staged.reference).await,
Err(SecretError::Lifecycle { .. })
));
drop_schema(&schema).await;
}
#[tokio::test]
async fn a_missing_table_is_refused_rather_than_created() {
let Some(dsn) = postgres_dsn() else {
return;
};
let schema = format!(
"axond_secret_test_{}",
Uuid7Generator::new().next().to_string().replace('-', "")
);
let (client, connection) = tokio_postgres::connect(&dsn, crate::usage::tls_connector())
.await
.expect("connect");
tokio::spawn(async move {
let _ = connection.await;
});
client
.batch_execute(&format!("CREATE SCHEMA {schema}"))
.await
.expect("create the test schema");
let error = PostgresSecrets::connect(
&dsn,
SecretStoreSettings {
schema: Some(schema.clone()),
create_table: false,
..SecretStoreSettings::default()
},
kek(17),
)
.await
.expect_err("an empty schema with no permission to create the table");
assert_eq!(error.category(), FailureCategory::Denied);
assert!(error.to_string().contains("secret_store_v1.sql"), "{error}");
drop_schema(&schema).await;
}
}