use std::net::IpAddr;
use std::sync::Arc;
use base64::prelude::*;
use tracing::{error, info, warn};
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::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::JobSpec;
use acme_proxy_signer::IssueOutcome;
use acme_proxy_signer::RequestedValidity;
use acme_proxy_signer::SignerBackend;
use acme_proxy_signer::SignerError;
use acme_proxy_signer::issuance::IssuanceError;
use acme_proxy_signer::issuance::announce_issuance;
use acme_proxy_signer::issuance::record_issuance;
use acme_proxy_signer::issuance::record_issue_failure;
use acme_proxy_store::account::Account;
use acme_proxy_store::authz::Authorization;
use acme_proxy_store::db::Database;
use acme_proxy_store::job::Job;
use acme_proxy_store::order::Order;
use acme_proxy_store::status::AuthzStatus;
use acme_proxy_store::status::OrderStatus;
pub const SIGNER_ISSUE_KIND: &str = "signer_issue";
#[must_use]
pub fn signer_issue_spec(
order: &Order,
csr_der: &[u8],
client: &ClientContext,
client_ip: Option<IpAddr>,
) -> JobSpec {
JobSpec::now(SIGNER_ISSUE_KIND, order.id.to_string())
.with_payload(serde_json::json!({
"order_id": order.id.to_string(),
"profile": order.profile,
"csr": BASE64_URL_SAFE_NO_PAD.encode(csr_der),
"client": client.to_json(),
"client_ip": client_ip.map(|ip| acme_proxy_core::client::canonical(ip).to_string()),
}))
.with_deadline(Some(order.expires))
}
pub struct SignerIssueJob {
database: Arc<Database>,
audit: Arc<Auditor>,
signers: Vec<(String, Arc<dyn SignerBackend>)>,
notifiers: acme_proxy_jobs::notify::Notifiers,
}
impl SignerIssueJob {
#[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 signer(&self, profile: &str) -> Option<&Arc<dyn SignerBackend>> {
super::mounted(&self.signers, profile)
}
}
async fn authority_withdrawn(
order: &Order,
database: &Database,
) -> Result<Option<&'static str>, sqlx::Error> {
let account =
Account::find_by_id(&order.profile, &order.account_id.to_string(), database).await?;
if account.is_none_or(|account| account.is_deactivated()) {
return Ok(Some(
"the account was deactivated before the certificate was issued",
));
}
let authzs = Authorization::find_by_order(order.id, database).await?;
if authzs.len() != order.identifiers.len()
|| authzs
.iter()
.any(|authz| authz.status != AuthzStatus::Valid)
{
return Ok(Some(
"an authorization was deactivated before the certificate was issued",
));
}
Ok(None)
}
fn payload_client(job: &Job) -> ClientContext {
ClientContext::from_json(&job.payload["client"])
}
#[async_trait::async_trait]
impl JobHandler for SignerIssueJob {
fn kind(&self) -> &'static str {
SIGNER_ISSUE_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 Some(csr_der) = payload["csr"]
.as_str()
.and_then(|csr| BASE64_URL_SAFE_NO_PAD.decode(csr).ok())
else {
return JobOutcome::Failed("the payload carries no readable CSR".to_string());
};
let mut 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}")),
};
if order.status != OrderStatus::Processing {
return JobOutcome::Done;
}
match authority_withdrawn(&order, &self.database).await {
Ok(None) => {}
Ok(Some(detail)) => {
warn!(event = "order_finalize_authority_withdrawn", outcome = "failure", order_id = %order_id, detail = %detail);
let problem = Problem::unauthorized(detail);
if let Err(error) = order.mark_invalid(problem.to_value(), &self.database).await {
error!(event = "order_mark_invalid_failed", outcome = "failure", order_id = %order_id, error = %error);
return JobOutcome::Retry(format!("recording the refusal failed: {error}"));
}
self.audit
.record(
AuditRecord::new(
AuditEvent::CertificateIssueFailed,
&order.profile,
Actor::acme(order.account_id),
)
.with_order(order.id, order.account_id, &order.identifiers)
.with_client(payload_client(job))
.with_reason("unauthorized")
.with_detail(detail),
)
.await;
return JobOutcome::Done;
}
Err(error) => {
return JobOutcome::Retry(format!("reading the order's authority failed: {error}"));
}
}
let Some(signer) = self.signer(&order.profile) else {
return JobOutcome::Retry(format!(
"profile `{}` is not mounted by this process",
order.profile
));
};
let client = payload_client(job);
let validity = RequestedValidity {
not_before: order.not_before,
not_after: order.not_after,
};
let issued = signer
.issue(order_id, &csr_der, &order.identifiers, validity)
.await;
match issued {
Ok(IssueOutcome::Issued(chain)) => {
match record_issuance(&mut order, chain, &self.database).await {
Ok(serial) => {
info!(event = "order_finalized", outcome = "success", order_id = %order_id, cert_serial = %serial);
let dispatcher = self.notifiers.get(&order.profile);
announce_issuance(
&order,
&serial,
job.created_at,
Actor::acme(order.account_id.to_string()),
client,
payload["client_ip"].as_str().map(str::to_string),
&self.audit,
dispatcher.as_deref(),
)
.await;
JobOutcome::Done
}
Err(IssuanceError::Chain(error)) => {
error!(event = "order_finalize_chain_unparsable", outcome = "failure", order_id = %order_id, error = %error);
JobOutcome::Failed(format!("the issued chain is unparsable: {error}"))
}
Err(IssuanceError::Leaf(error)) => {
error!(event = "order_finalize_leaf_unparsable", outcome = "failure", order_id = %order_id, error = %error);
JobOutcome::Failed(format!("the issued certificate is unparsable: {error}"))
}
Err(IssuanceError::Persist(error)) => {
error!(
event = "order_finalize_persistence_failed",
outcome = "failure",
order_id = %order_id,
error = %error
);
JobOutcome::Retry(format!("recording the certificate failed: {error}"))
}
}
}
Ok(IssueOutcome::Processing) => {
if let Err(error) = acme_proxy_store::upstream_order::UpstreamOrder::set_client(
order_id,
&client,
&self.database,
)
.await
{
warn!(
event = "upstream_order_client_context_failed",
outcome = "failure",
order_id = %order_id,
error = %error
);
}
info!(event = "order_finalize_delegated", outcome = "success", order_id = %order_id);
JobOutcome::Done
}
Err(SignerError::BadCsr) => {
warn!(event = "order_finalize_bad_csr", outcome = "failure", order_id = %order_id);
let problem = Problem::bad_csr("CSR invalid or does not match order");
if let Err(error) = order.mark_invalid(problem.to_value(), &self.database).await {
error!(event = "order_mark_invalid_failed", outcome = "failure", order_id = %order_id, error = %error);
return JobOutcome::Retry(format!("recording the refusal failed: {error}"));
}
self.audit
.record(
AuditRecord::new(
AuditEvent::CertificateIssueFailed,
&order.profile,
Actor::acme(order.account_id),
)
.with_order(order.id, order.account_id, &order.identifiers)
.with_client(client)
.with_reason("badCSR")
.with_detail("the signer backend rejected the CSR"),
)
.await;
JobOutcome::Done
}
Err(SignerError::Internal(detail)) => {
error!(
event = "order_finalize_issuance_failed",
outcome = "failure",
order_id = %order_id,
detail = %detail
);
JobOutcome::Retry(detail)
}
}
}
async fn abandon(&self, job: &Job, reason: &str) {
let Some(order_id) = job.payload["order_id"].as_str() else {
return;
};
warn!(
event = "order_finalize_abandoned",
outcome = "failure",
order_id = %order_id,
attempts = job.attempts,
reason = %reason,
"a queued issuance was given up; its order is marked invalid"
);
let mut order = match Order::find_by_id(order_id, &self.database).await {
Ok(Some(order)) if order.status == OrderStatus::Processing => order,
Ok(_) => return,
Err(error) => {
error!(event = "order_mark_invalid_failed", outcome = "failure", order_id = %order_id, error = %error);
return;
}
};
let account = order.account_id;
if let Err(error) = record_issue_failure(
&mut order,
&Problem::server_internal("Certificate issuance failed"),
reason,
Actor::acme(account),
payload_client(job),
&self.audit,
&self.database,
)
.await
{
error!(event = "order_mark_invalid_failed", outcome = "failure", order_id = %order_id, error = %error);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::acme::order::tests::{account, ready_order};
use acme_proxy_core::identifier::Identifier;
enum Answer {
BadCsr,
Internal,
Chain(&'static str),
Deferred,
}
struct Scripted(Answer);
#[async_trait::async_trait]
impl SignerBackend for Scripted {
async fn issue(
&self,
_order_id: &str,
_csr_der: &[u8],
_identifiers: &[Identifier],
_validity: RequestedValidity,
) -> Result<IssueOutcome, SignerError> {
match &self.0 {
Answer::BadCsr => Err(SignerError::BadCsr),
Answer::Internal => Err(SignerError::Internal("the token is gone".into())),
Answer::Chain(chain) => Ok(IssueOutcome::Issued((*chain).to_string())),
Answer::Deferred => Ok(IssueOutcome::Processing),
}
}
async fn revoke(&self, _cert_der: &[u8], _reason: Option<u32>) -> Result<(), SignerError> {
Ok(())
}
}
fn client() -> ClientContext {
ClientContext {
ip: Some("203.0.113.9".to_string()),
ptr: Some("client.example.net".to_string()),
user_agent: Some("certbot/9".to_string()),
request_id: Some("req-1".to_string()),
}
}
async fn claimed(database: &Arc<Database>) -> (Order, Job) {
let account = account(database).await;
let (mut order, csr) = ready_order(database, &account).await;
assert!(order.claim_for_finalize(database).await.unwrap());
let spec = signer_issue_spec(
&order,
&BASE64_URL_SAFE_NO_PAD.decode(csr).unwrap(),
&client(),
Some("203.0.113.9".parse().unwrap()),
);
let job = Job {
kind: SIGNER_ISSUE_KIND.to_string(),
dedup_key: spec.key.clone(),
payload: spec.payload.clone(),
..acme_proxy_store::testutil::job_fixture()
};
(order, job)
}
fn handler(database: &Arc<Database>, signer: Arc<dyn SignerBackend>) -> SignerIssueJob {
let (_tx, notifiers) = acme_proxy_jobs::notify::notifiers_channel(
acme_proxy_jobs::notify::DispatcherMap::new(),
);
SignerIssueJob::new(
database.clone(),
Arc::new(Auditor::offline(database.clone())),
vec![("default".to_string(), signer)],
notifiers,
)
}
async fn reload(database: &Database, order: &Order) -> Order {
Order::find_by_id(&order.id.to_string(), database)
.await
.unwrap()
.unwrap()
}
async fn audit_rows(database: &Database) -> Vec<acme_proxy_store::audit::AuditEntry> {
let query = acme_proxy_store::audit::AuditQuery {
limit: 50,
..acme_proxy_store::audit::AuditQuery::default()
};
acme_proxy_store::audit::AuditEntry::search(&query, database)
.await
.unwrap()
.0
}
#[tokio::test]
async fn a_signed_certificate_settles_the_order_and_names_the_client() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
let ca = acme_proxy_signer::local_ca::LocalCa::generate_in_memory(
"ecdsa-p256",
90,
database.clone(),
)
.unwrap();
assert!(matches!(
handler(&database, Arc::new(ca)).run(&job).await,
JobOutcome::Done
));
let stored = reload(&database, &order).await;
assert_eq!(stored.status, OrderStatus::Valid);
assert!(stored.certificate.is_some());
let rows = audit_rows(&database).await;
assert_eq!(rows.len(), 1, "{rows:?}");
assert_eq!(rows[0].event, "certificate_issued");
assert_eq!(rows[0].cert_serial, stored.cert_serial);
assert_eq!(rows[0].client_ip.as_deref(), Some("203.0.113.9"));
assert_eq!(rows[0].client_ptr.as_deref(), Some("client.example.net"));
assert_eq!(rows[0].request_id.as_deref(), Some("req-1"));
}
#[tokio::test]
async fn a_csr_the_backend_rejects_invalidates_the_order_once() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
assert!(matches!(
handler(&database, Arc::new(Scripted(Answer::BadCsr)))
.run(&job)
.await,
JobOutcome::Done
));
let stored = reload(&database, &order).await;
assert_eq!(stored.status, OrderStatus::Invalid);
let error = stored.error.unwrap();
assert_eq!(error["type"], "urn:ietf:params:acme:error:badCSR");
assert_eq!(error["detail"], "CSR invalid or does not match order");
let rows = audit_rows(&database).await;
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].event, "certificate_issue_failed");
assert_eq!(rows[0].reason.as_deref(), Some("badCSR"));
assert_eq!(rows[0].client_ip.as_deref(), Some("203.0.113.9"));
}
#[tokio::test]
async fn a_backend_failure_retries_and_is_abandoned_once() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
let handler = handler(&database, Arc::new(Scripted(Answer::Internal)));
let JobOutcome::Retry(reason) = handler.run(&job).await else {
panic!("an internal failure is retried")
};
assert_eq!(reason, "the token is gone");
assert_eq!(
reload(&database, &order).await.status,
OrderStatus::Processing
);
assert!(audit_rows(&database).await.is_empty());
handler.abandon(&job, &reason).await;
let stored = reload(&database, &order).await;
assert_eq!(stored.status, OrderStatus::Invalid);
assert_eq!(
stored.error.unwrap()["detail"],
"Certificate issuance failed"
);
let rows = audit_rows(&database).await;
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].reason.as_deref(), Some("serverInternal"));
assert_eq!(rows[0].detail.as_deref(), Some("the token is gone"));
handler.abandon(&job, &reason).await;
assert_eq!(audit_rows(&database).await.len(), 1);
}
#[tokio::test]
async fn an_unreadable_chain_fails_the_row() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
let outcome = handler(&database, Arc::new(Scripted(Answer::Chain("not a chain"))))
.run(&job)
.await;
let JobOutcome::Failed(reason) = outcome else {
panic!("an unreadable chain is permanent")
};
assert!(reason.contains("unparsable"), "{reason}");
let stored = reload(&database, &order).await;
assert!(stored.certificate.is_none());
}
#[tokio::test]
async fn a_deferred_issuance_leaves_the_order_to_its_backend() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
assert!(matches!(
handler(&database, Arc::new(Scripted(Answer::Deferred)))
.run(&job)
.await,
JobOutcome::Done
));
assert_eq!(
reload(&database, &order).await.status,
OrderStatus::Processing
);
assert!(audit_rows(&database).await.is_empty());
}
#[tokio::test]
async fn an_order_no_longer_processing_is_not_issued_again() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (mut order, job) = claimed(&database).await;
let settled = serde_json::json!({"type": "urn:ietf:params:acme:error:serverInternal"});
assert!(
order
.mark_invalid(settled.clone(), &database)
.await
.unwrap()
);
assert!(matches!(
handler(&database, Arc::new(Scripted(Answer::Internal)))
.run(&job)
.await,
JobOutcome::Done
));
let stored = reload(&database, &order).await;
assert_eq!(stored.status, OrderStatus::Invalid);
assert_eq!(stored.error, Some(settled));
}
#[tokio::test]
async fn a_deactivated_account_is_not_issued_for() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
Account::find_by_id("default", &order.account_id.to_string(), &database)
.await
.unwrap()
.unwrap()
.deactivate(&database)
.await
.unwrap();
assert!(matches!(
handler(&database, Arc::new(Scripted(Answer::Internal)))
.run(&job)
.await,
JobOutcome::Done
));
let stored = reload(&database, &order).await;
assert_eq!(stored.status, OrderStatus::Invalid);
assert_eq!(
stored.error.unwrap()["type"],
"urn:ietf:params:acme:error:unauthorized"
);
let rows = audit_rows(&database).await;
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].reason.as_deref(), Some("unauthorized"));
}
#[tokio::test]
async fn an_authorization_deactivated_after_the_claim_is_not_issued_for() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (order, job) = claimed(&database).await;
let authz = &Authorization::find_by_order(order.id, &database)
.await
.unwrap()[0];
assert!(
Authorization::set_deactivated(authz.id, database.raw_pool())
.await
.unwrap()
);
assert!(matches!(
handler(&database, Arc::new(Scripted(Answer::Internal)))
.run(&job)
.await,
JobOutcome::Done
));
assert_eq!(reload(&database, &order).await.status, OrderStatus::Invalid);
}
#[tokio::test]
async fn an_unmounted_profile_retries_and_a_broken_row_fails() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let (_, job) = claimed(&database).await;
let (_tx, notifiers) = acme_proxy_jobs::notify::notifiers_channel(
acme_proxy_jobs::notify::DispatcherMap::new(),
);
let elsewhere = SignerIssueJob::new(
database.clone(),
Arc::new(Auditor::offline(database.clone())),
Vec::new(),
notifiers,
);
assert!(matches!(elsewhere.run(&job).await, JobOutcome::Retry(_)));
for payload in [
serde_json::json!({}),
serde_json::json!({ "order_id": job.payload["order_id"], "csr": "!!" }),
serde_json::json!({ "order_id": uuid::Uuid::nil().to_string(), "csr": "MAA" }),
] {
let broken = Job {
payload,
..job.clone()
};
assert!(matches!(
elsewhere.run(&broken).await,
JobOutcome::Failed(_)
));
}
elsewhere
.abandon(
&Job {
payload: serde_json::json!({}),
..job
},
"gone",
)
.await;
}
}