use std::sync::Arc;
use std::time::Duration;
use axum::http::StatusCode;
use base64::Engine as _;
use base64::engine::general_purpose::STANDARD;
use tower::util::ServiceExt as _;
use super::harness::{
FakeProvider, PROVIDER_MATERIAL, ROTATED_MATERIAL, Replica, chat_request, material, owner,
state_pinning, sweep,
};
use crate::backends::secrets::envelope::DeploymentKek;
use crate::backends::secrets::postgres::{PostgresSecrets, SecretStoreSettings};
use crate::backends::secrets::{KekRef, SecretStore as _};
use crate::convergence::Outcome;
use crate::desired_state::{
ResourceVersionNumber, SecretLifecycle, SecretOwner, SecretRef, Uuid7Generator, fixtures,
};
use crate::routes::router;
async fn store() -> Option<(Arc<PostgresSecrets>, String)> {
let dsn = crate::test_services::postgres_dsn()?;
let schema = format!(
"axond_secret_runtime_{}",
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 kek = DeploymentKek::parse(
KekRef("AXOND_TEST_KEK".to_owned()),
&STANDARD.encode([7u8; 32]),
)
.expect("a 32-byte key");
let store = PostgresSecrets::connect(
&dsn,
SecretStoreSettings {
schema: Some(schema.clone()),
connect_timeout: Duration::from_secs(5),
..SecretStoreSettings::default()
},
kek,
)
.await
.expect("the store applies its own schema");
Some((Arc::new(store), schema))
}
async fn drop_schema(schema: &str) {
let Some(dsn) = crate::test_services::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");
}
async fn dump(schema: &str) -> String {
let dsn = crate::test_services::postgres_dsn().expect("a configured DSN");
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 rows = client
.query(
&format!(r#"SELECT t::text FROM "{schema}"."axond_secret" t"#),
&[],
)
.await
.expect("the store's rows as text");
assert!(
!rows.is_empty(),
"the store holds no rows, so the sweep would be vacuous"
);
rows.into_iter()
.map(|row| row.get::<_, Option<String>>(0).unwrap_or_default())
.collect::<Vec<_>>()
.join("\n")
}
async fn in_service(secrets: &PostgresSecrets, plaintext: &str) -> SecretRef {
let staged = secrets
.stage(owner(), material(plaintext))
.await
.expect("the store accepts material");
activate(secrets, staged.reference).await
}
async fn rotated_into_service(
secrets: &PostgresSecrets,
reference: SecretRef,
plaintext: &str,
) -> SecretRef {
let staged = secrets
.rotate(owner(), &reference, material(plaintext))
.await
.expect("the store accepts the next version");
assert_eq!(
staged.reference,
reference.rotated(),
"a rotation mints the next version of the same secret"
);
activate(secrets, staged.reference).await
}
async fn activate(secrets: &PostgresSecrets, reference: SecretRef) -> SecretRef {
secrets
.transition(owner(), &reference, SecretLifecycle::Active)
.await
.expect("staged material may be activated");
reference
}
#[tokio::test]
async fn the_credential_lifecycle_rotates_and_rolls_back_without_a_restart() {
let Some((secrets, schema)) = store().await else {
return;
};
let provider = FakeProvider::serving().await;
let replica = Replica::backed_by(&provider, Arc::clone(&secrets), Vec::new());
let sweep = sweep();
let first = in_service(&secrets, PROVIDER_MATERIAL).await;
replica
.publish("first", state_pinning(first, ResourceVersionNumber::FIRST))
.await;
let outcome = replica.converge().await;
assert!(matches!(outcome, Outcome::Published { .. }), "{outcome:?}");
let response = router(replica.state.clone())
.oneshot(chat_request())
.await
.expect("a response");
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
provider.presented().last().expect("a served request"),
&format!("Bearer {PROVIDER_MATERIAL}"),
"the shipped projection authenticated the call with the staged material"
);
let rotated = rotated_into_service(&secrets, first, ROTATED_MATERIAL).await;
replica
.publish(
"rotation",
state_pinning(rotated, ResourceVersionNumber::FIRST.next()),
)
.await;
let outcome = replica.converge().await;
assert!(matches!(outcome, Outcome::Published { .. }), "{outcome:?}");
let response = router(replica.state.clone())
.oneshot(chat_request())
.await
.expect("a response");
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
provider.presented().last().expect("a served request"),
&format!("Bearer {ROTATED_MATERIAL}"),
);
replica
.publish(
"rollback",
state_pinning(first, ResourceVersionNumber::FIRST.next().next()),
)
.await;
let outcome = replica.converge().await;
assert!(matches!(outcome, Outcome::Published { .. }), "{outcome:?}");
let response = router(replica.state.clone())
.oneshot(chat_request())
.await
.expect("a response");
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
provider.presented().last().expect("a served request"),
&format!("Bearer {PROVIDER_MATERIAL}"),
);
assert_eq!(replica.compiler.resolutions(), 3);
sweep.assert_absent("the secret store's rows", &dump(&schema).await);
drop_schema(&schema).await;
}
#[tokio::test]
async fn a_revoked_version_refuses_the_candidate_and_leaves_the_snapshot_serving() {
let Some((secrets, schema)) = store().await else {
return;
};
let provider = FakeProvider::serving().await;
let replica = Replica::backed_by(&provider, Arc::clone(&secrets), Vec::new());
let first = in_service(&secrets, PROVIDER_MATERIAL).await;
replica
.publish("first", state_pinning(first, ResourceVersionNumber::FIRST))
.await;
replica.converge().await;
assert_eq!(replica.generation(), 1);
let rotated = rotated_into_service(&secrets, first, ROTATED_MATERIAL).await;
secrets
.transition(owner(), &rotated, SecretLifecycle::Revoked)
.await
.expect("active material may be revoked");
replica
.publish(
"rotation",
state_pinning(rotated, ResourceVersionNumber::FIRST.next()),
)
.await;
let outcome = replica.converge().await;
assert!(
matches!(outcome, Outcome::Rejected { reason, .. } if reason == "secret"),
"{outcome:?}"
);
let report = replica.reconciler.report();
let rejection = report.last_rejection.as_ref().expect("a recorded refusal");
let sweep = sweep();
sweep.assert_absent("a refusal's detail", &rejection.detail);
assert!(
rejection.detail.contains(&rotated.to_string()),
"the refusal names the version an operator has to act on: {}",
rejection.detail
);
assert_eq!(replica.generation(), 1);
let response = router(replica.state.clone())
.oneshot(chat_request())
.await
.expect("a response");
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
provider.presented().last().expect("a served request"),
&format!("Bearer {PROVIDER_MATERIAL}"),
);
drop_schema(&schema).await;
}
#[tokio::test]
async fn another_owners_material_never_reaches_a_pool() {
let Some((secrets, schema)) = store().await else {
return;
};
let provider = FakeProvider::serving().await;
let replica = Replica::backed_by(&provider, Arc::clone(&secrets), Vec::new());
let foreign = SecretOwner::tenant(fixtures::tenant_id(9));
let staged = secrets
.stage(foreign, material(PROVIDER_MATERIAL))
.await
.expect("the store accepts material");
secrets
.transition(foreign, &staged.reference, SecretLifecycle::Active)
.await
.expect("staged material may be activated");
replica
.publish(
"cross-owner",
state_pinning(staged.reference, ResourceVersionNumber::FIRST),
)
.await;
let outcome = replica.converge().await;
assert!(
matches!(outcome, Outcome::Rejected { reason, .. } if reason == "secret"),
"{outcome:?}"
);
assert_eq!(replica.generation(), 0);
sweep().assert_absent(
"a cross-owner refusal",
&format!("{:?}", replica.reconciler.report()),
);
drop_schema(&schema).await;
}