use std::sync::Arc;
use std::time::Duration;
use tracing::{error, info, warn};
use uuid::Uuid;
use acme_proxy_core::audit::Actor;
use acme_proxy_core::audit::AuditEvent;
use acme_proxy_core::audit::AuditRecord;
use acme_proxy_core::audit::ClientContext;
use acme_proxy_core::audit::RequestContext;
use acme_proxy_core::error::Problem;
use acme_proxy_jobs::auditor::Auditor;
use acme_proxy_jobs::jobs::JobHandler;
use acme_proxy_jobs::jobs::JobOutcome;
use acme_proxy_jobs::jobs::JobQueue;
use acme_proxy_jobs::jobs::JobSpec;
use acme_proxy_jobs::notify::CertificateRevokedData;
use acme_proxy_jobs::notify::NotifyDispatcher;
use acme_proxy_jobs::notify::NotifyEvent;
use acme_proxy_signer::RevocationRoute;
use acme_proxy_signer::SignerBackend;
use acme_proxy_signer::SignerError;
use acme_proxy_store::account::Account;
use acme_proxy_store::db::Database;
use acme_proxy_store::job::Job;
use acme_proxy_store::order::Order;
pub enum Revoker<'a> {
Backend(&'a dyn SignerBackend),
Ledger {
issuer: &'a str,
jobs: &'a JobQueue,
},
Queued { jobs: &'a JobQueue, wait: Duration },
}
impl<'a> Revoker<'a> {
#[must_use]
pub fn for_route(route: &'a RevocationRoute, jobs: &'a JobQueue, wait: Duration) -> Self {
match route {
RevocationRoute::Ledger { issuer } => Revoker::Ledger { issuer, jobs },
RevocationRoute::Delegated => Revoker::Queued { jobs, wait },
}
}
}
pub enum Client<'a> {
Request(&'a RequestContext),
Resolved(ClientContext),
}
impl Client<'_> {
async fn resolve(self, audit: &Auditor) -> ClientContext {
match self {
Client::Request(request) => audit.client(request).await,
Client::Resolved(client) => client,
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum RevokeError {
#[error("{}", .0.detail())]
Refused(Problem),
#[error("no such order")]
NotFound,
#[error("order has no issued certificate")]
NotIssued,
#[error("certificate already revoked")]
AlreadyRevoked,
#[error("unsupported revocation reason code {0}")]
BadReason(u32),
#[error("signer error: {}", signer_detail(.0))]
Signer(SignerError),
#[error("database error: {0}")]
Database(sqlx::Error),
#[error("the revocation is queued as job {job} and has not completed yet")]
Pending { job: Uuid },
#[error("the revocation failed (job {job}): {reason}")]
Abandoned { job: Uuid, reason: String },
#[error("internal error: {0}")]
Internal(String),
}
pub fn signer_detail(error: &SignerError) -> String {
match error {
SignerError::Internal(detail) => detail.clone(),
SignerError::BadCsr => "unexpected badCsr from revoke".to_string(),
}
}
impl From<sqlx::Error> for RevokeError {
fn from(error: sqlx::Error) -> Self {
Self::Database(error)
}
}
impl From<RevokeError> for Problem {
fn from(error: RevokeError) -> Self {
match error {
RevokeError::Refused(problem) => problem,
RevokeError::NotFound | RevokeError::NotIssued => {
Problem::malformed("Unknown certificate")
}
RevokeError::AlreadyRevoked => Problem::already_revoked("Certificate already revoked"),
RevokeError::BadReason(reason) => Problem::bad_revocation_reason(format!(
"Unsupported revocation reason code {reason}"
)),
RevokeError::Pending { .. } => {
Problem::service_unavailable("The revocation is queued; retry to confirm it")
}
RevokeError::Signer(_)
| RevokeError::Database(_)
| RevokeError::Internal(_)
| RevokeError::Abandoned { .. } => Problem::server_internal("Revocation failed"),
}
}
}
pub struct Revocations<'a> {
pub database: &'a Arc<Database>,
pub audit: &'a Auditor,
pub notify: Option<&'a NotifyDispatcher>,
pub revoker: Revoker<'a>,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Refusals {
Audited,
Silent,
}
impl Revocations<'_> {
pub async fn revoke_certificate(
&self,
profile: &str,
cert_der: &[u8],
reason: Option<u32>,
pubkey: &[u8],
account: Option<Account>,
request: &RequestContext,
) -> Result<Order, RevokeError> {
let database = self.database;
let (serial_hex, _) = acme_proxy_core::cert::cert_serial_and_spki(cert_der).map_err(|error| {
warn!(event = "certificate_revoke_parse_failed", outcome = "failure", error = %error);
RevokeError::Refused(Problem::malformed("certificate is unparsable"))
})?;
let order = Order::find_by_cert_serial(profile, &serial_hex, database)
.await
.map_err(|error| {
error!(event = "certificate_revoke_lookup_failed", outcome = "failure", cert_serial = %serial_hex, error = %error);
RevokeError::Refused(Problem::server_internal("Certificate lookup failed"))
})?
.filter(|order| {
order
.certificate
.as_deref()
.and_then(|chain| acme_proxy_core::cert::leaf_der_from_chain(chain).ok())
.is_some_and(|leaf| leaf == cert_der)
})
.ok_or(());
let actor = match &account {
Some(cached) => Actor::acme(cached.id.to_string()),
None => Actor::acme_certificate_key(),
};
let client = Client::Request(request).resolve(self.audit).await;
let revoke_failed = |reason: &'static str, detail: &str| {
refusal(
profile,
actor.clone(),
&serial_hex,
client.clone(),
reason,
detail,
)
};
let order = match order {
Ok(order) => order,
Err(()) => {
warn!(event = "certificate_revoke_unknown_certificate", outcome = "failure", cert_serial = %serial_hex);
self.audit
.record(revoke_failed(
"malformed",
"no certificate issued here matches",
))
.await;
return Err(RevokeError::Refused(Problem::malformed(
"Unknown certificate",
)));
}
};
let cert_key_matches = order.cert_pubkey.as_deref() == Some(pubkey);
let authorized = if cert_key_matches {
true
} else {
match account {
Some(cached) => cached.id == order.account_id,
None => Account::find_by_pubkey(profile, pubkey, database)
.await
.map_err(|error| {
error!(event = "certificate_revoke_account_lookup_failed", outcome = "failure", error = %error);
RevokeError::Refused(Problem::server_internal("Account lookup failed"))
})?
.is_some_and(|found| found.id == order.account_id),
}
};
if !authorized {
warn!(event = "certificate_revoke_unauthorized", outcome = "failure", order_id = %order.id, cert_serial = %serial_hex);
self.audit
.record(
revoke_failed(
"unauthorized",
"signed by neither the order's account nor the certificate's own key",
)
.with_order(order.id, order.account_id, &order.identifiers),
)
.await;
return Err(RevokeError::Refused(Problem::unauthorized(
"Neither the order's account nor the certificate's own key signed this request",
)));
}
self.revoke(
order,
cert_der,
serial_hex,
reason,
actor,
client,
Refusals::Audited,
)
.await
}
pub async fn revoke_order(
&self,
id: &str,
reason: Option<u32>,
actor: Actor,
client: ClientContext,
) -> Result<Order, RevokeError> {
let Some(order) = Order::find_by_id(id, self.database).await? else {
return Err(RevokeError::NotFound);
};
let Some(chain) = order.certificate.as_deref() else {
return Err(RevokeError::NotIssued);
};
let cert_der = acme_proxy_core::cert::leaf_der_from_chain(chain).map_err(|error| {
RevokeError::Internal(format!("stored certificate chain is unparsable: {error}"))
})?;
let serial = match order.cert_serial.clone() {
Some(serial) => serial,
None => acme_proxy_core::cert::cert_serial_and_spki(&cert_der)
.map(|(serial, _)| serial)
.map_err(|error| {
RevokeError::Internal(format!("stored certificate is unparsable: {error}"))
})?,
};
self.revoke(
order,
&cert_der,
serial,
reason,
actor,
client,
Refusals::Silent,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn revoke(
&self,
mut order: Order,
cert_der: &[u8],
serial_hex: String,
reason: Option<u32>,
actor: Actor,
client: ClientContext,
refusals: Refusals,
) -> Result<Order, RevokeError> {
let refused = |order: &Order, reason: &'static str, detail: &str| {
refusal(
&order.profile,
actor.clone(),
&serial_hex,
client.clone(),
reason,
detail,
)
.with_order(order.id, order.account_id, &order.identifiers)
};
let audited = refusals == Refusals::Audited;
if order.revoked_at.is_some() {
if audited {
warn!(event = "certificate_revoke_already_revoked", outcome = "failure", order_id = %order.id, cert_serial = %serial_hex);
self.audit
.record(refused(&order, "alreadyRevoked", "already revoked"))
.await;
}
return Err(RevokeError::AlreadyRevoked);
}
if let Some(code) = reason
&& !acme_proxy_core::cert::is_valid_revocation_reason(code)
{
if audited {
warn!(
event = "certificate_revoke_bad_reason",
outcome = "failure",
reason = code
);
self.audit
.record(refused(
&order,
"badRevocationReason",
&format!("reason code {code}"),
))
.await;
}
return Err(RevokeError::BadReason(code));
}
match self.revoker {
Revoker::Queued { jobs, wait } => {
return self
.revoke_through_the_queue(&order, reason, &actor, &client, jobs, wait)
.await;
}
Revoker::Backend(signer) => {
if let Err(error) = signer.revoke(cert_der, reason).await {
error!(event = "certificate_revoke_signer_failed", outcome = "failure", order_id = %order.id, cert_serial = %serial_hex, error = %error);
self.audit
.record(refused(&order, "serverInternal", &error.to_string()))
.await;
return Err(RevokeError::Signer(error));
}
match order.revoke(reason.map(i64::from), self.database).await {
Ok(false) => {
if audited {
self.audit
.record(refused(&order, "alreadyRevoked", "already revoked"))
.await;
}
return Err(RevokeError::AlreadyRevoked);
}
Ok(true) => {}
Err(error) => {
error!(event = "certificate_revoke_persist_failed", outcome = "failure", order_id = %order.id, cert_serial = %serial_hex, error = %error);
self.audit
.record(refused(&order, "serverInternal", &error.to_string()))
.await;
return Err(RevokeError::Database(error));
}
}
}
Revoker::Ledger { issuer, jobs } => {
if let Err(error) = self
.record_in_ledger(&mut order, issuer, cert_der, &serial_hex, reason)
.await
{
self.audit
.record(refused(&order, "serverInternal", &error.to_string()))
.await;
return Err(error);
}
jobs.enqueue_or_log(acme_proxy_signer::local_ca::sweep::regenerate_spec(issuer))
.await;
}
}
info!(event = "certificate_revoked", outcome = "success", order_id = %order.id, cert_serial = %serial_hex);
let revoked = AuditRecord::new(AuditEvent::CertificateRevoked, &order.profile, actor)
.with_order(order.id, order.account_id, &order.identifiers)
.with_serial(&serial_hex)
.with_client(client.clone());
self.audit
.record(match reason {
Some(code) => revoked.with_reason(code.to_string()),
None => revoked,
})
.await;
if let Some(dispatcher) = self.notify {
dispatcher
.dispatch(NotifyEvent::CertificateRevoked(CertificateRevokedData {
profile: order.profile.clone(),
order_id: order.id.to_string(),
account_id: order.account_id.to_string(),
cert_serial: serial_hex,
reason,
client_ip: client.ip,
}))
.await;
}
Ok(order)
}
}
impl Revocations<'_> {
async fn revoke_through_the_queue(
&self,
order: &Order,
reason: Option<u32>,
actor: &Actor,
client: &ClientContext,
jobs: &JobQueue,
wait: Duration,
) -> Result<Order, RevokeError> {
let id = order.id.to_string();
jobs.enqueue(signer_revoke_spec(&id, reason, actor, client))
.await
.inspect_err(|error| {
error!(event = "certificate_revoke_queue_failed", outcome = "failure", order_id = %id, error = %error);
})?;
let job = acme_proxy_store::job::Job::find_latest_by_dedup(
SIGNER_REVOKE_KIND,
&id,
self.database,
)
.await?
.ok_or_else(|| {
RevokeError::Internal(format!("the revocation of order {id} was not queued"))
})?;
info!(event = "certificate_revoke_queued", outcome = "progress", order_id = %id, job_id = %job.id);
match await_job(self.database, job.id, wait).await? {
JobSettled::Done => Order::find_by_id(&id, self.database)
.await?
.ok_or(RevokeError::NotFound),
JobSettled::Failed(reason) => Err(RevokeError::Abandoned {
job: job.id,
reason,
}),
JobSettled::Cancelled => Err(RevokeError::Abandoned {
job: job.id,
reason: "cancelled by an operator".to_string(),
}),
JobSettled::Pending => Err(RevokeError::Pending { job: job.id }),
}
}
async fn record_in_ledger(
&self,
order: &mut Order,
issuer: &str,
cert_der: &[u8],
serial_hex: &str,
reason: Option<u32>,
) -> Result<(), RevokeError> {
let revoked_at = acme_proxy_store::nonce::now_secs();
let row = acme_proxy_store::revocation::Revocation {
issuer: issuer.to_string(),
serial: serial_hex.to_string(),
revoked_at,
reason,
not_after: acme_proxy_core::cert::cert_validity(cert_der)
.ok()
.map(|(_, not_after)| not_after),
};
let written = async {
let mut tx = self.database.write_transaction().await?;
if acme_proxy_store::crl::StoredCrl::find(issuer, tx.conn())
.await?
.is_none()
{
return Ok(Recorded::NoCrlYet);
}
if !Order::set_revoked(order.id, reason.map(i64::from), revoked_at, tx.conn()).await? {
return Ok(Recorded::Already);
}
row.insert_if_absent(tx.conn()).await?;
tx.commit().await?;
Ok::<Recorded, sqlx::Error>(Recorded::Yes)
}
.await;
let recorded = written.map_err(|error| {
error!(event = "certificate_revoke_persist_failed", outcome = "failure", order_id = %order.id, cert_serial = %serial_hex, error = %error);
RevokeError::Database(error)
})?;
if recorded == Recorded::Already {
warn!(event = "certificate_revoke_already_revoked", outcome = "failure", order_id = %order.id, cert_serial = %serial_hex);
return Err(RevokeError::AlreadyRevoked);
}
if recorded == Recorded::NoCrlYet {
warn!(event = "certificate_revoke_ca_uninitialized", outcome = "failure", order_id = %order.id, issuer = %issuer);
return Err(RevokeError::Internal(format!(
"the CA {issuer} has no stored CRL yet — start `acme-proxy serve` with this \
configuration once, so it can import its revocation ledger, then retry"
)));
}
order.revoked_at = Some(revoked_at);
order.revocation_reason = reason.map(i64::from);
Ok(())
}
}
#[derive(Debug, PartialEq, Eq)]
enum Recorded {
Yes,
Already,
NoCrlYet,
}
fn refusal(
profile: &str,
actor: Actor,
serial: &str,
client: ClientContext,
reason: &'static str,
detail: &str,
) -> AuditRecord {
AuditRecord::new(AuditEvent::CertificateRevokeFailed, profile, actor)
.with_serial(serial)
.with_client(client)
.with_reason(reason)
.with_detail(detail)
}
pub const SIGNER_REVOKE_KIND: &str = "signer_revoke";
#[must_use]
pub fn signer_revoke_spec(
order_id: &str,
reason: Option<u32>,
actor: &Actor,
client: &ClientContext,
) -> JobSpec {
JobSpec::now(SIGNER_REVOKE_KIND, order_id).with_payload(serde_json::json!({
"order_id": order_id,
"reason": reason,
"actor_kind": actor.kind.as_str(),
"actor_id": actor.id,
"client": client.to_json(),
}))
}
#[must_use]
pub fn request_wait(request_timeout_ms: u64) -> Duration {
Duration::from_millis(request_timeout_ms).saturating_sub(Duration::from_secs(1))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum JobSettled {
Done,
Failed(String),
Cancelled,
Pending,
}
const AWAIT_PACE: Duration = Duration::from_millis(100);
pub async fn await_job(
database: &Database,
id: Uuid,
wait: Duration,
) -> Result<JobSettled, sqlx::Error> {
let deadline = tokio::time::Instant::now() + wait;
loop {
let Some(job) = acme_proxy_store::job::Job::find_by_id(id, database).await? else {
return Ok(JobSettled::Failed(format!("job {id} disappeared")));
};
match job.status.as_str() {
"done" => return Ok(JobSettled::Done),
"failed" => {
return Ok(JobSettled::Failed(
job.last_error
.unwrap_or_else(|| "no reason recorded".to_string()),
));
}
"cancelled" => return Ok(JobSettled::Cancelled),
_ if tokio::time::Instant::now() >= deadline => return Ok(JobSettled::Pending),
_ => tokio::time::sleep(AWAIT_PACE).await,
}
}
}
pub struct SignerRevokeJob {
database: Arc<Database>,
audit: Arc<Auditor>,
signers: Vec<(String, Arc<dyn SignerBackend>)>,
notifiers: acme_proxy_jobs::notify::Notifiers,
}
impl SignerRevokeJob {
#[must_use]
pub fn new(
database: Arc<Database>,
audit: Arc<Auditor>,
signers: Vec<(String, Arc<dyn SignerBackend>)>,
notifiers: acme_proxy_jobs::notify::Notifiers,
) -> Self {
Self {
database,
audit,
signers,
notifiers,
}
}
}
fn payload_actor(payload: &serde_json::Value) -> Actor {
let id = payload["actor_id"].as_str().map(str::to_string);
let kind = match payload["actor_kind"].as_str() {
Some("admin") => acme_proxy_core::audit::ActorKind::Admin,
Some("acme") => acme_proxy_core::audit::ActorKind::Acme,
Some("system") => acme_proxy_core::audit::ActorKind::System,
_ => acme_proxy_core::audit::ActorKind::Cli,
};
Actor { kind, id }
}
#[async_trait::async_trait]
impl JobHandler for SignerRevokeJob {
fn kind(&self) -> &'static str {
SIGNER_REVOKE_KIND
}
async fn run(&self, job: &Job) -> JobOutcome {
let payload = &job.payload;
let Some(order_id) = payload["order_id"].as_str() else {
return JobOutcome::Failed("the payload names no order".to_string());
};
let reason = payload["reason"]
.as_u64()
.and_then(|code| u32::try_from(code).ok());
let order = match Order::find_by_id(order_id, &self.database).await {
Ok(Some(order)) => order,
Ok(None) => return JobOutcome::Failed("the order no longer exists".to_string()),
Err(error) => return JobOutcome::Retry(format!("reading the order failed: {error}")),
};
let Some(signer) = super::mounted(&self.signers, &order.profile) else {
return JobOutcome::Retry(format!(
"profile `{}` is not mounted by this process",
order.profile
));
};
let dispatcher = self.notifiers.get(&order.profile);
let revocations = Revocations {
database: &self.database,
audit: &self.audit,
notify: dispatcher.as_deref(),
revoker: Revoker::Backend(signer.as_ref()),
};
match revocations
.revoke_order(
order_id,
reason,
payload_actor(payload),
ClientContext::from_json(&payload["client"]),
)
.await
{
Ok(_) | Err(RevokeError::AlreadyRevoked) => JobOutcome::Done,
Err(error @ (RevokeError::Signer(_) | RevokeError::Database(_))) => {
JobOutcome::Retry(error.to_string())
}
Err(error) => JobOutcome::Failed(error.to_string()),
}
}
async fn abandon(&self, job: &Job, reason: &str) {
error!(
event = "certificate_revoke_abandoned",
outcome = "failure",
order_id = %job.dedup_key,
reason = %reason,
"a queued revocation was given up; the certificate is still trusted"
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use acme_proxy_core::identifier::Identifier;
use acme_proxy_jobs::jobs::JobHandler;
use acme_proxy_signer::IssueOutcome;
use acme_proxy_signer::RequestedValidity;
struct Failing;
#[async_trait::async_trait]
impl SignerBackend for Failing {
async fn issue(
&self,
_order_id: &str,
_csr_der: &[u8],
_identifiers: &[Identifier],
_validity: RequestedValidity,
) -> Result<IssueOutcome, SignerError> {
Err(SignerError::Internal("not here".into()))
}
async fn revoke(&self, _cert_der: &[u8], _reason: Option<u32>) -> Result<(), SignerError> {
Err(SignerError::Internal("upstream unreachable".into()))
}
}
fn handler(database: &Arc<Database>) -> SignerRevokeJob {
SignerRevokeJob::new(
database.clone(),
Arc::new(Auditor::offline(database.clone())),
vec![("default".to_string(), Arc::new(Failing))],
std::collections::HashMap::new().into(),
)
}
async fn queued(database: &Arc<Database>, order_id: &str) -> Job {
let queue = acme_proxy_jobs::testutil::idle_job_queue(database.clone());
queue
.enqueue(signer_revoke_spec(
order_id,
None,
&Actor::cli(),
&ClientContext::default(),
))
.await
.unwrap();
Job::find_live(SIGNER_REVOKE_KIND, order_id, database)
.await
.unwrap()
.unwrap()
}
#[tokio::test]
async fn a_failing_backend_is_retried() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = acme_proxy_store::testutil::account_id(&database).await;
let order = acme_proxy_store::testutil::issued_order(
&database,
"default",
account,
&["example.com"],
30,
)
.await;
let job = queued(&database, &order.id.to_string()).await;
assert!(matches!(
handler(&database).run(&job).await,
JobOutcome::Retry(_)
));
let stored = Order::find_by_id(&order.id.to_string(), &database)
.await
.unwrap()
.unwrap();
assert!(stored.revoked_at.is_none());
}
struct Succeeding;
#[async_trait::async_trait]
impl SignerBackend for Succeeding {
async fn issue(
&self,
_order_id: &str,
_csr_der: &[u8],
_identifiers: &[Identifier],
_validity: RequestedValidity,
) -> Result<IssueOutcome, SignerError> {
Err(SignerError::Internal("not here".into()))
}
async fn revoke(&self, _cert_der: &[u8], _reason: Option<u32>) -> Result<(), SignerError> {
Ok(())
}
}
fn worker(
database: &Arc<Database>,
backend: Arc<dyn SignerBackend>,
max_attempts: u32,
) -> (JobQueue, tokio::sync::watch::Sender<bool>) {
let config = acme_proxy_core::config::JobsConfig {
poll_interval_ms: 10,
max_attempts,
retry_base_seconds: 0,
retry_max_seconds: 0,
..acme_proxy_core::config::JobsConfig::default()
};
let queue = JobQueue::new(database.clone(), &config);
let mut registry = acme_proxy_jobs::jobs::JobRegistry::new();
registry
.register(Arc::new(SignerRevokeJob::new(
database.clone(),
Arc::new(Auditor::offline(database.clone())),
vec![("default".to_string(), backend)],
std::collections::HashMap::new().into(),
)))
.unwrap();
let (shutdown, receiver) = tokio::sync::watch::channel(false);
acme_proxy_jobs::jobs::spawn_runner(queue.clone(), Arc::new(registry), &config, receiver);
(queue, shutdown)
}
async fn issued(database: &Arc<Database>) -> Order {
let account = acme_proxy_store::testutil::account_id(database).await;
acme_proxy_store::testutil::issued_order(database, "default", account, &["example.com"], 30)
.await
}
fn operator_client() -> ClientContext {
ClientContext {
ip: Some("198.51.100.4".to_string()),
..ClientContext::default()
}
}
#[tokio::test]
async fn a_queued_revocation_answers_once_the_worker_has_revoked() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let order = issued(&database).await;
let (jobs, _worker) = worker(&database, Arc::new(Succeeding), 5);
let audit = Auditor::offline(database.clone());
let revocations = Revocations {
database: &database,
audit: &audit,
notify: None,
revoker: Revoker::Queued {
jobs: &jobs,
wait: Duration::from_secs(10),
},
};
let revoked = revocations
.revoke_order(
&order.id.to_string(),
Some(1),
Actor::admin("root"),
operator_client(),
)
.await
.unwrap();
assert!(revoked.revoked_at.is_some());
assert_eq!(revoked.revocation_reason, Some(1));
let query = acme_proxy_store::audit::AuditQuery {
limit: 50,
..acme_proxy_store::audit::AuditQuery::default()
};
let (rows, _) = acme_proxy_store::audit::AuditEntry::search(&query, &database)
.await
.unwrap();
assert_eq!(rows.len(), 1, "{rows:?}");
assert_eq!(rows[0].event, "certificate_revoked");
assert_eq!(rows[0].actor_kind, "admin");
assert_eq!(rows[0].client_ip.as_deref(), Some("198.51.100.4"));
}
#[tokio::test]
async fn concurrent_revocations_of_one_certificate_record_one() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let order = issued(&database).await;
let audit = Auditor::offline(database.clone());
let backend = Succeeding;
let revocations = Revocations {
database: &database,
audit: &audit,
notify: None,
revoker: Revoker::Backend(&backend),
};
let order_id = order.id.to_string();
let revoke = |reason: u32| {
revocations.revoke_order(
&order_id,
Some(reason),
Actor::admin("root"),
operator_client(),
)
};
let (first, second) = tokio::join!(revoke(1), revoke(4));
let outcomes = [first, second];
assert_eq!(
outcomes.iter().filter(|outcome| outcome.is_ok()).count(),
1,
"exactly one revocation records the withdrawal: {outcomes:?}"
);
assert!(
outcomes
.iter()
.any(|outcome| matches!(outcome, Err(RevokeError::AlreadyRevoked)))
);
let query = acme_proxy_store::audit::AuditQuery {
limit: 50,
..acme_proxy_store::audit::AuditQuery::default()
};
let (rows, _) = acme_proxy_store::audit::AuditEntry::search(&query, &database)
.await
.unwrap();
assert_eq!(
rows.iter()
.filter(|row| row.event == "certificate_revoked")
.count(),
1,
"{rows:?}"
);
let winner = outcomes
.iter()
.find_map(|outcome| outcome.as_ref().ok())
.unwrap();
let stored = Order::find_by_id(&order.id.to_string(), &database)
.await
.unwrap()
.unwrap();
assert_eq!(stored.revocation_reason, winner.revocation_reason);
assert_eq!(stored.revoked_at, winner.revoked_at);
}
#[tokio::test]
async fn an_unanswered_revocation_is_pending_and_asking_again_follows_it() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let order = issued(&database).await;
let jobs = acme_proxy_jobs::testutil::idle_job_queue(database.clone());
let audit = Auditor::offline(database.clone());
let revocations = Revocations {
database: &database,
audit: &audit,
notify: None,
revoker: Revoker::Queued {
jobs: &jobs,
wait: Duration::ZERO,
},
};
let mut seen = Vec::new();
for _ in 0..2 {
match revocations
.revoke_order(
&order.id.to_string(),
None,
Actor::cli(),
ClientContext::default(),
)
.await
{
Err(RevokeError::Pending { job }) => seen.push(job),
other => panic!("expected a pending revocation, got {other:?}"),
}
}
assert_eq!(seen[0], seen[1], "the second ask follows the first job");
assert_eq!(
Job::count_live(SIGNER_REVOKE_KIND, &database)
.await
.unwrap(),
1
);
let problem = Problem::from(RevokeError::Pending { job: seen[0] }).to_value();
assert_eq!(problem["status"], 503);
assert_eq!(problem["type"], "urn:ietf:params:acme:error:serverInternal");
}
#[tokio::test]
async fn a_revocation_the_worker_gave_up_on_is_abandoned() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let order = issued(&database).await;
let (jobs, _worker) = worker(&database, Arc::new(Failing), 1);
let audit = Auditor::offline(database.clone());
let revocations = Revocations {
database: &database,
audit: &audit,
notify: None,
revoker: Revoker::Queued {
jobs: &jobs,
wait: Duration::from_secs(10),
},
};
let error = revocations
.revoke_order(
&order.id.to_string(),
None,
Actor::cli(),
ClientContext::default(),
)
.await
.unwrap_err();
let RevokeError::Abandoned { reason, .. } = &error else {
panic!("expected an abandoned revocation, got {error:?}")
};
assert!(reason.contains("upstream unreachable"), "{reason}");
assert_eq!(
Problem::from(error).to_value()["detail"],
"Revocation failed"
);
let stored = Order::find_by_id(&order.id.to_string(), &database)
.await
.unwrap()
.unwrap();
assert!(stored.revoked_at.is_none());
}
#[tokio::test]
async fn a_cancelled_or_vanished_job_settles_the_wait() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let order = issued(&database).await;
let job = queued(&database, &order.id.to_string()).await;
Job::cancel_row(
job.id,
acme_proxy_store::status::JobStatus::Ready,
&database,
)
.await
.unwrap()
.unwrap();
assert_eq!(
await_job(&database, job.id, Duration::from_secs(10))
.await
.unwrap(),
JobSettled::Cancelled
);
let JobSettled::Failed(reason) = await_job(
&database,
acme_proxy_store::id::mint(),
Duration::from_secs(10),
)
.await
.unwrap() else {
panic!("a job that is not there has failed")
};
assert!(reason.contains("disappeared"), "{reason}");
}
#[tokio::test]
async fn the_route_picks_the_revoker_and_the_wait_fits_the_deadline() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let jobs = acme_proxy_jobs::testutil::idle_job_queue(database);
let ledger = RevocationRoute::Ledger {
issuer: "ab".to_string(),
};
assert!(matches!(
Revoker::for_route(&ledger, &jobs, Duration::ZERO),
Revoker::Ledger { issuer: "ab", .. }
));
assert!(matches!(
Revoker::for_route(&RevocationRoute::Delegated, &jobs, Duration::from_secs(3)),
Revoker::Queued { wait, .. } if wait == Duration::from_secs(3)
));
assert_eq!(request_wait(60_000), Duration::from_secs(59));
assert_eq!(request_wait(500), Duration::ZERO);
}
#[tokio::test]
async fn a_vanished_order_fails_for_good() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let job = queued(&database, &acme_proxy_store::id::mint().to_string()).await;
assert!(matches!(
handler(&database).run(&job).await,
JobOutcome::Failed(_)
));
}
}